Wl6adams's picture
Organize private Nexum release into Lite, Universal, and Expanded profiles
9a70a84
Raw
History Blame Contribute Delete
5.46 kB
"""Durable execution intents and exactly-once replay protection."""
from __future__ import annotations
import fcntl
import hashlib
import hmac
import json
import os
import secrets
import time
from dataclasses import asdict, dataclass
from pathlib import Path
from typing import Any
from .contracts import ToolCall, ToolExecutionResult
from .sandbox import control_root
@dataclass(frozen=True)
class ExecutionIntent:
key: str
action_sha256: str
session_sha256: str
tool_call_id: str
tool: str
status: str
created_unix: int
completed_unix: int = 0
result: dict[str, Any] | None = None
class IdempotencyStore:
def __init__(self, workspace: str | Path) -> None:
root = control_root(workspace) / "executions"
root.mkdir(parents=True, exist_ok=True, mode=0o700)
os.chmod(root, 0o700)
self.root = root
@staticmethod
def action_sha256(session_id: str, call: ToolCall) -> str:
canonical = json.dumps(
{
"session_id": session_id,
"call_id": call.call_id,
"name": call.name,
"args": call.args,
},
sort_keys=True,
separators=(",", ":"),
).encode("utf-8")
return hashlib.sha256(canonical).hexdigest()
@classmethod
def key(cls, session_id: str, call: ToolCall) -> str:
supplied = call.idempotency_key.strip()
if supplied:
digest = hashlib.sha256(supplied.encode("utf-8")).hexdigest()
else:
digest = cls.action_sha256(session_id, call)
return digest
def _path(self, key: str) -> Path:
if len(key) != 64 or any(char not in "0123456789abcdef" for char in key):
raise ValueError("execution key is invalid")
return self.root / f"{key}.json"
@staticmethod
def _load(path: Path) -> ExecutionIntent:
payload = json.loads(path.read_text(encoding="utf-8"))
if not isinstance(payload, dict):
raise RuntimeError("execution intent is invalid")
return ExecutionIntent(**payload)
@staticmethod
def _write(path: Path, intent: ExecutionIntent) -> None:
temporary = path.with_name(f".{path.name}.{secrets.token_hex(8)}.tmp")
temporary.write_text(
json.dumps(asdict(intent), sort_keys=True, separators=(",", ":")),
encoding="utf-8",
)
temporary.chmod(0o600)
with temporary.open("rb") as handle:
os.fsync(handle.fileno())
os.replace(temporary, path)
directory_fd = os.open(path.parent, os.O_RDONLY)
try:
os.fsync(directory_fd)
finally:
os.close(directory_fd)
def begin(self, session_id: str, call: ToolCall) -> tuple[ExecutionIntent, bool]:
key = self.key(session_id, call)
path = self._path(key)
lock_path = path.with_suffix(".lock")
with lock_path.open("a+b") as lock_handle:
lock_path.chmod(0o600)
fcntl.flock(lock_handle.fileno(), fcntl.LOCK_EX)
try:
action_sha256 = self.action_sha256(session_id, call)
if path.is_file():
existing = self._load(path)
if not hmac.compare_digest(existing.action_sha256, action_sha256):
raise RuntimeError(
"idempotency key is already bound to another action"
)
return existing, False
intent = ExecutionIntent(
key=key,
action_sha256=action_sha256,
session_sha256=hashlib.sha256(
session_id.encode("utf-8")
).hexdigest(),
tool_call_id=call.call_id,
tool=call.name,
status="started",
created_unix=int(time.time()),
)
self._write(path, intent)
return intent, True
finally:
fcntl.flock(lock_handle.fileno(), fcntl.LOCK_UN)
def complete(
self,
session_id: str,
call: ToolCall,
result: ToolExecutionResult,
) -> ExecutionIntent:
key = self.key(session_id, call)
path = self._path(key)
lock_path = path.with_suffix(".lock")
with lock_path.open("a+b") as lock_handle:
lock_path.chmod(0o600)
fcntl.flock(lock_handle.fileno(), fcntl.LOCK_EX)
try:
existing = self._load(path)
expected = self.action_sha256(session_id, call)
if not hmac.compare_digest(existing.action_sha256, expected):
raise RuntimeError("execution intent does not match result")
if existing.status == "completed":
return existing
completed = ExecutionIntent(
**{
**asdict(existing),
"status": "completed",
"completed_unix": int(time.time()),
"result": result.to_dict(),
}
)
self._write(path, completed)
return completed
finally:
fcntl.flock(lock_handle.fileno(), fcntl.LOCK_UN)
__all__ = ["ExecutionIntent", "IdempotencyStore"]