Files
m4 8376f56ab4 feat: native sandbox execution, dynamic review middleware, and workspace files
Adds native sandbox execution runtime, dynamic review middleware, and
workspace file handling, with supporting stream events, prompt, and scope
registry changes plus architecture docs.
2026-08-19 20:00:09 +08:00

1306 lines
49 KiB
Python

"""Persistent ownership registry for conversation workspaces.
The LangGraph checkpoint store is not an authority for filesystem ownership:
threads, runs and cron records can be created independently and the filesystem
cannot participate in their transactions. This module keeps the small,
deployment-local registry that binds all of them to one conversation scope.
The v1 implementation intentionally uses SQLite. It is safe for the supported
single-host deployment topology and keeps the registry outside every agent
workspace. Callers must use the public methods below rather than addressing the
database directly so a future PostgreSQL adapter has one replacement point.
"""
from __future__ import annotations
import os
import secrets
import sqlite3
import threading
import uuid
from collections.abc import Iterator
from contextlib import contextmanager
from dataclasses import dataclass
from datetime import UTC, datetime, timedelta
from pathlib import Path
from typing import Literal
from . import paths
ScopeState = Literal["provisioning", "draft", "active", "deleting", "deleted"]
OwnerState = Literal[
"reserved", "active", "draining", "terminal", "failed", "quarantined"
]
class ScopeRegistryError(RuntimeError):
"""Base class for registry failures."""
class ScopeNotFoundError(ScopeRegistryError):
"""Raised when no matching deployment scope exists."""
class ScopeConflictError(ScopeRegistryError):
"""Raised for stale revisions or incompatible ownership mappings."""
class ScopeIdempotencyConflictError(ScopeConflictError):
"""Raised when one run request id is reused for another payload."""
class ScopeInterruptResolvedError(ScopeConflictError):
"""Raised when an interrupt already has a persisted resolution."""
class ScopeAccessError(ScopeRegistryError):
"""Raised when a runtime does not own the requested scope."""
@dataclass(frozen=True, slots=True)
class ScopeRecord:
deployment_id: str
scope_id: str
primary_thread_id: str
state: ScopeState
revision: int
primary_owner_id: str
created_at: str
updated_at: str
@dataclass(frozen=True, slots=True)
class OwnerRecord:
deployment_id: str
owner_id: str
scope_id: str
owner_type: str
resource_id: str | None
parent_owner_id: str | None
state: OwnerState
created_at: str
updated_at: str
@dataclass(frozen=True, slots=True)
class DeploymentLock:
deployment_id: str
lock_name: str
operation_id: str
expires_at: str
@dataclass(frozen=True, slots=True)
class ScopeOperation:
deployment_id: str
operation_id: str
scope_id: str | None
kind: str
state: str
result_sha256: str | None
last_error_code: str | None
@dataclass(frozen=True, slots=True)
class RunReservation:
deployment_id: str
scope_id: str
run_request_id: str
turn_id: str
interrupt_key: str | None
request_hash: str
run_owner_id: str
run_id: str | None
state: str
# Kept as a source-compatible type name for callers that have not yet switched
# to the run-request terminology. A turn is a logical conversation unit; a
# reservation belongs to one concrete run request.
TurnReservation = RunReservation
_TERMINAL_SCOPE_STATES = frozenset({"deleted"})
_TERMINAL_OWNER_STATES = frozenset({"terminal", "quarantined"})
_SCOPE_TRANSITIONS: dict[str, frozenset[str]] = {
"provisioning": frozenset({"draft", "deleting", "deleted"}),
"draft": frozenset({"active", "deleting", "deleted"}),
"active": frozenset({"deleting"}),
"deleting": frozenset({"deleted"}),
"deleted": frozenset(),
}
def _utc_now() -> str:
return datetime.now(UTC).isoformat()
def _ensure_uuid(value: str, field: str) -> str:
try:
return str(uuid.UUID(value))
except (TypeError, ValueError) as exc:
raise ScopeRegistryError(f"{field} must be a UUID") from exc
def deployment_id_for_workspace(workspace_root: Path | str | None = None) -> str:
"""Return the stable deployment identifier for a workspace root."""
configured = os.getenv("EVOSCIENTIST_DEPLOYMENT_ID", "").strip()
if configured:
return configured
# Deploy resolves the workspace before starting LangGraph. Avoid resolve()
# here because this function also runs on the agent's async execution path.
root = Path(workspace_root or paths.WORKSPACE_ROOT).expanduser()
return str(uuid.uuid5(uuid.NAMESPACE_URL, f"evoscientist:{root}"))
def default_registry_path(workspace_root: Path | str | None = None) -> Path:
root = Path(workspace_root or paths.WORKSPACE_ROOT).expanduser()
return root / ".evoscientist" / "control" / "scope-registry.sqlite3"
def scope_service_token_path(workspace_root: Path | str | None = None) -> Path:
"""Location of the local backend-to-BFF registry credential."""
configured = os.getenv("EVOSCIENTIST_CONTROL_DIR", "").strip()
control_dir = (
Path(configured).expanduser()
if configured
else Path.home() / ".evoscientist" / "control"
)
return control_dir / "scope-service-token"
def get_scope_service_token(workspace_root: Path | str | None = None) -> str:
"""Return the durable host-local token shared by deploy and the WebUI.
It lives beside the control-plane database, never in a conversation scope
or browser-delivered configuration. Creation is atomic so a concurrent
launcher cannot replace a running backend's credential.
"""
path = scope_service_token_path(workspace_root)
path.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
try:
path.parent.chmod(0o700)
except OSError:
pass
try:
token = path.read_text(encoding="utf-8").strip()
except FileNotFoundError:
token = ""
if token:
return token
generated = secrets.token_urlsafe(32)
try:
fd = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_EXCL, 0o600)
except FileExistsError:
token = path.read_text(encoding="utf-8").strip()
if token:
return token
raise ScopeRegistryError("workspace scope service token is empty") from None
with os.fdopen(fd, "w", encoding="utf-8") as token_file:
token_file.write(generated)
token_file.write("\n")
try:
path.chmod(0o600)
except OSError:
pass
return generated
class ScopeRegistry:
"""SQLite-backed registry with revision-checked state transitions."""
def __init__(self, database_path: Path | str):
self.path = Path(database_path).expanduser()
self._init_lock = threading.Lock()
self._initialized = False
def initialize(self) -> None:
with self._init_lock:
if self._initialized:
return
self.path.parent.mkdir(mode=0o700, parents=True, exist_ok=True)
try:
self.path.parent.chmod(0o700)
except OSError:
pass
with self._connect() as conn:
self._migrate(conn)
try:
self.path.chmod(0o600)
except OSError:
pass
self._initialized = True
def _connect(self) -> sqlite3.Connection:
conn = sqlite3.connect(self.path, timeout=30, isolation_level=None)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA foreign_keys = ON")
conn.execute("PRAGMA journal_mode = WAL")
conn.execute("PRAGMA busy_timeout = 30000")
return conn
@staticmethod
def _migrate(conn: sqlite3.Connection) -> None:
version = int(conn.execute("PRAGMA user_version").fetchone()[0])
if version > 2:
raise ScopeRegistryError("scope registry schema is newer than this binary")
if version == 0:
conn.executescript(
"""
CREATE TABLE scopes (
deployment_id TEXT NOT NULL,
scope_id TEXT NOT NULL,
primary_thread_id TEXT NOT NULL,
state TEXT NOT NULL,
revision INTEGER NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
deleted_at TEXT,
PRIMARY KEY (deployment_id, scope_id),
UNIQUE (deployment_id, primary_thread_id)
);
CREATE TABLE scope_owners (
deployment_id TEXT NOT NULL,
owner_id TEXT NOT NULL,
scope_id TEXT NOT NULL,
owner_type TEXT NOT NULL,
resource_id TEXT,
parent_owner_id TEXT,
state TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
terminal_at TEXT,
PRIMARY KEY (deployment_id, owner_id),
FOREIGN KEY (deployment_id, scope_id)
REFERENCES scopes(deployment_id, scope_id)
);
CREATE UNIQUE INDEX scope_owner_resource_unique
ON scope_owners(deployment_id, owner_type, resource_id)
WHERE resource_id IS NOT NULL;
CREATE INDEX scope_owners_scope_state
ON scope_owners(deployment_id, scope_id, state);
CREATE TABLE scope_run_requests (
deployment_id TEXT NOT NULL,
scope_id TEXT NOT NULL,
run_request_id TEXT NOT NULL,
turn_id TEXT NOT NULL,
interrupt_key TEXT,
request_hash TEXT NOT NULL,
run_owner_id TEXT NOT NULL,
run_id TEXT,
state TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (deployment_id, scope_id, run_request_id),
FOREIGN KEY (deployment_id, scope_id)
REFERENCES scopes(deployment_id, scope_id)
);
CREATE INDEX scope_run_requests_turn
ON scope_run_requests(deployment_id, scope_id, turn_id);
CREATE UNIQUE INDEX scope_run_requests_interrupt_unique
ON scope_run_requests(deployment_id, scope_id, interrupt_key)
WHERE interrupt_key IS NOT NULL;
CREATE TABLE scope_operations (
deployment_id TEXT NOT NULL,
operation_id TEXT NOT NULL,
scope_id TEXT,
kind TEXT NOT NULL,
expected_revision INTEGER,
state TEXT NOT NULL,
external_resource_id TEXT,
result_sha256 TEXT,
last_error_code TEXT,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (deployment_id, operation_id)
);
CREATE INDEX scope_operations_scope_state
ON scope_operations(deployment_id, scope_id, state);
CREATE TABLE deployment_locks (
deployment_id TEXT NOT NULL,
lock_name TEXT NOT NULL,
operation_id TEXT NOT NULL,
expires_at TEXT NOT NULL,
created_at TEXT NOT NULL,
PRIMARY KEY (deployment_id, lock_name)
);
"""
)
conn.execute("PRAGMA user_version = 2")
return
if version == 1:
# SQLite cannot change a composite primary key in place. Historical
# reservations used turn_id as the idempotency key, so copy each row
# with run_request_id = turn_id. New resume runs may then retain the
# logical turn while receiving distinct request ids.
conn.executescript(
"""
CREATE TABLE scope_run_requests (
deployment_id TEXT NOT NULL,
scope_id TEXT NOT NULL,
run_request_id TEXT NOT NULL,
turn_id TEXT NOT NULL,
interrupt_key TEXT,
request_hash TEXT NOT NULL,
run_owner_id TEXT NOT NULL,
run_id TEXT,
state TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
PRIMARY KEY (deployment_id, scope_id, run_request_id),
FOREIGN KEY (deployment_id, scope_id)
REFERENCES scopes(deployment_id, scope_id)
);
INSERT INTO scope_run_requests(
deployment_id, scope_id, run_request_id, turn_id,
interrupt_key, request_hash, run_owner_id, run_id, state,
created_at, updated_at
)
SELECT deployment_id, scope_id, turn_id, turn_id,
NULL, request_hash, run_owner_id, run_id, state,
created_at, updated_at
FROM scope_turns;
CREATE INDEX scope_run_requests_turn
ON scope_run_requests(deployment_id, scope_id, turn_id);
CREATE UNIQUE INDEX scope_run_requests_interrupt_unique
ON scope_run_requests(deployment_id, scope_id, interrupt_key)
WHERE interrupt_key IS NOT NULL;
DROP TABLE scope_turns;
"""
)
conn.execute("PRAGMA user_version = 2")
@contextmanager
def _transaction(self) -> Iterator[sqlite3.Connection]:
self.initialize()
with self._connect() as conn:
conn.execute("BEGIN IMMEDIATE")
try:
yield conn
except Exception:
conn.rollback()
raise
else:
conn.commit()
@staticmethod
def _scope_from_row(row: sqlite3.Row, owner_id: str) -> ScopeRecord:
return ScopeRecord(
deployment_id=str(row["deployment_id"]),
scope_id=str(row["scope_id"]),
primary_thread_id=str(row["primary_thread_id"]),
state=str(row["state"]), # type: ignore[arg-type]
revision=int(row["revision"]),
primary_owner_id=owner_id,
created_at=str(row["created_at"]),
updated_at=str(row["updated_at"]),
)
@staticmethod
def _owner_from_row(row: sqlite3.Row) -> OwnerRecord:
return OwnerRecord(
deployment_id=str(row["deployment_id"]),
owner_id=str(row["owner_id"]),
scope_id=str(row["scope_id"]),
owner_type=str(row["owner_type"]),
resource_id=(str(row["resource_id"]) if row["resource_id"] else None),
parent_owner_id=(
str(row["parent_owner_id"]) if row["parent_owner_id"] else None
),
state=str(row["state"]), # type: ignore[arg-type]
created_at=str(row["created_at"]),
updated_at=str(row["updated_at"]),
)
@staticmethod
def _primary_owner(
conn: sqlite3.Connection, deployment_id: str, scope_id: str
) -> str:
row = conn.execute(
"""
SELECT owner_id FROM scope_owners
WHERE deployment_id = ? AND scope_id = ? AND owner_type = 'primary_thread'
LIMIT 1
""",
(deployment_id, scope_id),
).fetchone()
if row is None:
raise ScopeRegistryError("scope has no primary owner")
return str(row["owner_id"])
@staticmethod
def _workspace_mutation_lock_active(
conn: sqlite3.Connection,
deployment_id: str,
operation_id: str | None = None,
) -> bool:
row = conn.execute(
"""
SELECT 1 FROM deployment_locks
WHERE deployment_id = ?
AND lock_name IN ('workspace-cutover', 'workspace-lifecycle')
AND expires_at > ?
AND (? IS NULL OR operation_id != ?)
""",
(deployment_id, _utc_now(), operation_id, operation_id),
).fetchone()
return row is not None
def provision(
self,
deployment_id: str,
primary_thread_id: str,
*,
scope_id: str | None = None,
state: ScopeState = "draft",
operation_id: str | None = None,
lock_operation_id: str | None = None,
) -> ScopeRecord:
"""Reserve one immutable scope for a primary thread.
Retrying the same primary thread returns its existing mapping. Passing a
different explicit scope for an existing thread is a conflict rather than
an opportunity to silently remap its files.
"""
if state not in {"provisioning", "draft"}:
raise ScopeRegistryError("new scopes must start provisioning or draft")
requested_scope = (
_ensure_uuid(scope_id, "scope_id") if scope_id else str(uuid.uuid4())
)
now = _utc_now()
operation_id = (
_ensure_uuid(operation_id, "operation_id")
if operation_id
else str(uuid.uuid4())
)
if lock_operation_id is not None:
lock_operation_id = _ensure_uuid(lock_operation_id, "lock_operation_id")
with self._transaction() as conn:
existing = conn.execute(
"""
SELECT * FROM scopes WHERE deployment_id = ? AND primary_thread_id = ?
""",
(deployment_id, primary_thread_id),
).fetchone()
if existing is not None:
if scope_id and str(existing["scope_id"]) != requested_scope:
raise ScopeConflictError("thread already belongs to another scope")
if str(existing["state"]) == "deleted":
conn.execute(
"""
UPDATE scopes
SET state = ?, revision = revision + 1, updated_at = ?, deleted_at = NULL
WHERE deployment_id = ? AND scope_id = ?
""",
(state, now, deployment_id, str(existing["scope_id"])),
)
conn.execute(
"""
UPDATE scope_owners
SET state = 'active', updated_at = ?, terminal_at = NULL
WHERE deployment_id = ? AND scope_id = ?
AND owner_type = 'primary_thread'
""",
(now, deployment_id, str(existing["scope_id"])),
)
existing = conn.execute(
"""
SELECT * FROM scopes
WHERE deployment_id = ? AND primary_thread_id = ?
""",
(deployment_id, primary_thread_id),
).fetchone()
assert existing is not None
primary_owner = self._primary_owner(
conn, deployment_id, str(existing["scope_id"])
)
return self._scope_from_row(existing, primary_owner)
if self._workspace_mutation_lock_active(
conn, deployment_id, lock_operation_id
):
raise ScopeAccessError(
"workspace cutover or maintenance is in progress"
)
conn.execute(
"""
INSERT INTO scopes(
deployment_id, scope_id, primary_thread_id, state, revision,
created_at, updated_at
) VALUES (?, ?, ?, ?, 1, ?, ?)
""",
(deployment_id, requested_scope, primary_thread_id, state, now, now),
)
primary_owner = str(uuid.uuid4())
conn.execute(
"""
INSERT INTO scope_owners(
deployment_id, owner_id, scope_id, owner_type, resource_id,
parent_owner_id, state, created_at, updated_at
) VALUES (?, ?, ?, 'primary_thread', ?, NULL, 'active', ?, ?)
""",
(
deployment_id,
primary_owner,
requested_scope,
primary_thread_id,
now,
now,
),
)
conn.execute(
"""
INSERT OR REPLACE INTO scope_operations(
deployment_id, operation_id, scope_id, kind, state, created_at, updated_at
) VALUES (?, ?, ?, 'provision', 'completed', ?, ?)
""",
(deployment_id, operation_id, requested_scope, now, now),
)
return ScopeRecord(
deployment_id=deployment_id,
scope_id=requested_scope,
primary_thread_id=primary_thread_id,
state=state,
revision=1,
primary_owner_id=primary_owner,
created_at=now,
updated_at=now,
)
def get_by_thread(self, deployment_id: str, thread_id: str) -> ScopeRecord:
self.initialize()
with self._connect() as conn:
row = conn.execute(
"SELECT * FROM scopes WHERE deployment_id = ? AND primary_thread_id = ?",
(deployment_id, thread_id),
).fetchone()
if row is None:
raise ScopeNotFoundError("no workspace scope for thread")
return self._scope_from_row(
row, self._primary_owner(conn, deployment_id, str(row["scope_id"]))
)
def get(self, deployment_id: str, scope_id: str) -> ScopeRecord:
scope_id = _ensure_uuid(scope_id, "scope_id")
self.initialize()
with self._connect() as conn:
row = conn.execute(
"SELECT * FROM scopes WHERE deployment_id = ? AND scope_id = ?",
(deployment_id, scope_id),
).fetchone()
if row is None:
raise ScopeNotFoundError("workspace scope not found")
return self._scope_from_row(
row, self._primary_owner(conn, deployment_id, scope_id)
)
def list_scopes(
self, deployment_id: str, *, state: ScopeState | None = None
) -> list[ScopeRecord]:
"""List deployment scopes for administrative maintenance only."""
if state is not None and state not in _SCOPE_TRANSITIONS:
raise ScopeRegistryError("invalid workspace scope state")
self.initialize()
with self._connect() as conn:
if state is None:
rows = conn.execute(
"""
SELECT * FROM scopes WHERE deployment_id = ? ORDER BY created_at ASC
""",
(deployment_id,),
).fetchall()
else:
rows = conn.execute(
"""
SELECT * FROM scopes WHERE deployment_id = ? AND state = ?
ORDER BY created_at ASC
""",
(deployment_id, state),
).fetchall()
return [
self._scope_from_row(
row,
self._primary_owner(conn, deployment_id, str(row["scope_id"])),
)
for row in rows
]
def transition_scope(
self,
deployment_id: str,
scope_id: str,
*,
expected_revision: int,
state: ScopeState,
) -> ScopeRecord:
scope_id = _ensure_uuid(scope_id, "scope_id")
now = _utc_now()
with self._transaction() as conn:
current = conn.execute(
"SELECT * FROM scopes WHERE deployment_id = ? AND scope_id = ?",
(deployment_id, scope_id),
).fetchone()
if current is None:
raise ScopeNotFoundError("workspace scope not found")
current_state = str(current["state"])
if state not in _SCOPE_TRANSITIONS.get(current_state, frozenset()):
raise ScopeConflictError(
f"cannot transition scope from {current_state} to {state}"
)
deleted_at = now if state == "deleted" else None
updated = conn.execute(
"""
UPDATE scopes
SET state = ?, revision = revision + 1, updated_at = ?,
deleted_at = COALESCE(?, deleted_at)
WHERE deployment_id = ? AND scope_id = ? AND revision = ?
""",
(state, now, deleted_at, deployment_id, scope_id, expected_revision),
)
if updated.rowcount != 1:
raise ScopeConflictError("scope revision changed")
row = conn.execute(
"SELECT * FROM scopes WHERE deployment_id = ? AND scope_id = ?",
(deployment_id, scope_id),
).fetchone()
assert row is not None
return self._scope_from_row(
row, self._primary_owner(conn, deployment_id, scope_id)
)
def register_owner(
self,
deployment_id: str,
scope_id: str,
*,
owner_type: str,
resource_id: str | None = None,
parent_owner_id: str | None = None,
owner_id: str | None = None,
state: OwnerState = "reserved",
) -> OwnerRecord:
scope_id = _ensure_uuid(scope_id, "scope_id")
owner_id = _ensure_uuid(owner_id, "owner_id") if owner_id else str(uuid.uuid4())
now = _utc_now()
with self._transaction() as conn:
scope = conn.execute(
"SELECT state FROM scopes WHERE deployment_id = ? AND scope_id = ?",
(deployment_id, scope_id),
).fetchone()
if scope is None:
raise ScopeNotFoundError("workspace scope not found")
if self._workspace_mutation_lock_active(conn, deployment_id):
raise ScopeAccessError(
"workspace cutover or maintenance is in progress"
)
if (
str(scope["state"]) in _TERMINAL_SCOPE_STATES
or str(scope["state"]) == "deleting"
):
raise ScopeAccessError("scope does not accept new owners")
if parent_owner_id:
parent = conn.execute(
"""
SELECT state FROM scope_owners
WHERE deployment_id = ? AND owner_id = ? AND scope_id = ?
""",
(deployment_id, parent_owner_id, scope_id),
).fetchone()
if parent is None or str(parent["state"]) != "active":
raise ScopeAccessError("parent owner is not active")
try:
conn.execute(
"""
INSERT INTO scope_owners(
deployment_id, owner_id, scope_id, owner_type, resource_id,
parent_owner_id, state, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
deployment_id,
owner_id,
scope_id,
owner_type,
resource_id,
parent_owner_id,
state,
now,
now,
),
)
except sqlite3.IntegrityError as exc:
raise ScopeConflictError(
"owner or external resource already exists"
) from exc
return OwnerRecord(
deployment_id=deployment_id,
owner_id=owner_id,
scope_id=scope_id,
owner_type=owner_type,
resource_id=resource_id,
parent_owner_id=parent_owner_id,
state=state,
created_at=now,
updated_at=now,
)
def bind_owner(
self,
deployment_id: str,
scope_id: str,
owner_id: str,
resource_id: str,
*,
state: OwnerState = "active",
) -> OwnerRecord:
scope_id = _ensure_uuid(scope_id, "scope_id")
owner_id = _ensure_uuid(owner_id, "owner_id")
now = _utc_now()
with self._transaction() as conn:
try:
result = conn.execute(
"""
UPDATE scope_owners
SET resource_id = ?, state = ?, updated_at = ?
WHERE deployment_id = ? AND scope_id = ? AND owner_id = ?
AND state NOT IN ('terminal', 'quarantined')
""",
(resource_id, state, now, deployment_id, scope_id, owner_id),
)
except sqlite3.IntegrityError as exc:
raise ScopeConflictError(
"external resource belongs to another scope"
) from exc
if result.rowcount != 1:
raise ScopeAccessError("owner cannot be bound")
row = conn.execute(
"""
SELECT * FROM scope_owners
WHERE deployment_id = ? AND scope_id = ? AND owner_id = ?
""",
(deployment_id, scope_id, owner_id),
).fetchone()
assert row is not None
return self._owner_from_row(row)
def assert_runtime(
self,
deployment_id: str,
scope_id: str,
thread_id: str,
owner_id: str,
) -> ScopeRecord:
"""Verify a graph/tool runtime is an active owner of this scope."""
scope_id = _ensure_uuid(scope_id, "workspace_scope_id")
owner_id = _ensure_uuid(owner_id, "workspace_scope_owner_id")
self.initialize()
with self._connect() as conn:
scope = conn.execute(
"SELECT * FROM scopes WHERE deployment_id = ? AND scope_id = ?",
(deployment_id, scope_id),
).fetchone()
if scope is None:
raise ScopeNotFoundError("workspace scope not found")
if str(scope["state"]) not in {"draft", "active"}:
raise ScopeAccessError("workspace scope is not active")
if self._workspace_mutation_lock_active(conn, deployment_id):
raise ScopeAccessError(
"workspace cutover or maintenance is in progress"
)
owner = conn.execute(
"""
SELECT * FROM scope_owners
WHERE deployment_id = ? AND scope_id = ? AND owner_id = ?
""",
(deployment_id, scope_id, owner_id),
).fetchone()
if owner is None or (
str(owner["state"]) != "active"
and not (
str(owner["owner_type"]) == "primary_run"
and str(owner["state"]) == "reserved"
)
):
raise ScopeAccessError("workspace owner is not active")
if str(owner["owner_type"]) == "primary_thread":
if (
str(scope["primary_thread_id"]) != thread_id
or str(owner["resource_id"]) != thread_id
):
raise ScopeAccessError("primary thread does not own this scope")
elif str(owner["owner_type"]) == "primary_run":
if str(scope["primary_thread_id"]) != thread_id or str(
owner["parent_owner_id"]
) != self._primary_owner(conn, deployment_id, scope_id):
raise ScopeAccessError("primary run does not own this scope")
elif str(owner["owner_type"]) != "schedule" and str(
owner["resource_id"]
) not in {None, thread_id}:
raise ScopeAccessError("derived thread does not own this scope")
return self._scope_from_row(
scope, self._primary_owner(conn, deployment_id, scope_id)
)
def owners(self, deployment_id: str, scope_id: str) -> list[OwnerRecord]:
scope_id = _ensure_uuid(scope_id, "scope_id")
self.initialize()
with self._connect() as conn:
rows = conn.execute(
"""
SELECT * FROM scope_owners
WHERE deployment_id = ? AND scope_id = ? ORDER BY created_at ASC
""",
(deployment_id, scope_id),
).fetchall()
return [self._owner_from_row(row) for row in rows]
def get_owner_by_resource(
self, deployment_id: str, resource_id: str
) -> OwnerRecord:
"""Return the durable owner for one external child resource."""
self.initialize()
with self._connect() as conn:
row = conn.execute(
"""
SELECT * FROM scope_owners
WHERE deployment_id = ? AND resource_id = ?
""",
(deployment_id, resource_id),
).fetchone()
if row is None:
raise ScopeNotFoundError("workspace owner not found")
return self._owner_from_row(row)
def reserve_run(
self,
deployment_id: str,
scope_id: str,
run_request_id: str,
turn_id: str,
request_hash: str,
*,
interrupt_key: str | None = None,
) -> RunReservation:
"""Reserve one idempotent run request before creating it remotely.
``turn_id`` groups a user message and all of its approval resumes.
``run_request_id`` identifies one external ``runs.create`` call, so a
lost response can be retried without turning a valid resume into a
conflict with the initial run.
"""
scope_id = _ensure_uuid(scope_id, "scope_id")
run_request_id = _ensure_uuid(run_request_id, "run_request_id")
turn_id = _ensure_uuid(turn_id, "turn_id")
if interrupt_key is not None:
if (
not isinstance(interrupt_key, str)
or not interrupt_key
or len(interrupt_key) > 256
):
raise ScopeRegistryError(
"interrupt_key must be a non-empty string up to 256 characters"
)
now = _utc_now()
with self._transaction() as conn:
existing = conn.execute(
"""
SELECT * FROM scope_run_requests
WHERE deployment_id = ? AND scope_id = ? AND run_request_id = ?
""",
(deployment_id, scope_id, run_request_id),
).fetchone()
if existing is not None:
existing_interrupt_key = (
str(existing["interrupt_key"])
if existing["interrupt_key"]
else None
)
if (
str(existing["request_hash"]) != request_hash
or str(existing["turn_id"]) != turn_id
or existing_interrupt_key != interrupt_key
):
raise ScopeIdempotencyConflictError(
"run_request_id was reused with another request"
)
return RunReservation(
deployment_id=deployment_id,
scope_id=scope_id,
run_request_id=str(existing["run_request_id"]),
turn_id=turn_id,
interrupt_key=existing_interrupt_key,
request_hash=request_hash,
run_owner_id=str(existing["run_owner_id"]),
run_id=str(existing["run_id"]) if existing["run_id"] else None,
state=str(existing["state"]),
)
if interrupt_key is not None:
resolved_interrupt = conn.execute(
"""
SELECT run_request_id FROM scope_run_requests
WHERE deployment_id = ? AND scope_id = ? AND interrupt_key = ?
""",
(deployment_id, scope_id, interrupt_key),
).fetchone()
if resolved_interrupt is not None:
raise ScopeInterruptResolvedError(
"interrupt already has a resume request"
)
scope = conn.execute(
"SELECT state FROM scopes WHERE deployment_id = ? AND scope_id = ?",
(deployment_id, scope_id),
).fetchone()
if scope is None or str(scope["state"]) not in {"draft", "active"}:
raise ScopeAccessError("workspace scope does not accept a run")
if self._workspace_mutation_lock_active(conn, deployment_id):
raise ScopeAccessError(
"workspace cutover or maintenance is in progress"
)
parent_owner = self._primary_owner(conn, deployment_id, scope_id)
run_owner_id = str(uuid.uuid4())
conn.execute(
"""
INSERT INTO scope_owners(
deployment_id, owner_id, scope_id, owner_type, resource_id,
parent_owner_id, state, created_at, updated_at
) VALUES (?, ?, ?, 'primary_run', NULL, ?, 'reserved', ?, ?)
""",
(deployment_id, run_owner_id, scope_id, parent_owner, now, now),
)
conn.execute(
"""
INSERT INTO scope_run_requests(
deployment_id, scope_id, run_request_id, turn_id, interrupt_key,
request_hash, run_owner_id, run_id, state, created_at, updated_at
) VALUES (?, ?, ?, ?, ?, ?, ?, NULL, 'reserved', ?, ?)
""",
(
deployment_id,
scope_id,
run_request_id,
turn_id,
interrupt_key,
request_hash,
run_owner_id,
now,
now,
),
)
return RunReservation(
deployment_id=deployment_id,
scope_id=scope_id,
run_request_id=run_request_id,
turn_id=turn_id,
interrupt_key=interrupt_key,
request_hash=request_hash,
run_owner_id=run_owner_id,
run_id=None,
state="reserved",
)
def bind_run(
self,
deployment_id: str,
scope_id: str,
run_request_id: str,
run_id: str,
) -> RunReservation:
"""Attach the remote run id to a prior reservation exactly once."""
scope_id = _ensure_uuid(scope_id, "scope_id")
run_request_id = _ensure_uuid(run_request_id, "run_request_id")
now = _utc_now()
with self._transaction() as conn:
row = conn.execute(
"""
SELECT * FROM scope_run_requests
WHERE deployment_id = ? AND scope_id = ? AND run_request_id = ?
""",
(deployment_id, scope_id, run_request_id),
).fetchone()
if row is None:
raise ScopeNotFoundError("workspace run reservation not found")
if row["run_id"] and str(row["run_id"]) != run_id:
raise ScopeConflictError("run request is already bound to another run")
owner_id = str(row["run_owner_id"])
conn.execute(
"""
UPDATE scope_owners
SET resource_id = ?, state = 'active', updated_at = ?
WHERE deployment_id = ? AND scope_id = ? AND owner_id = ?
AND state IN ('reserved', 'active')
""",
(run_id, now, deployment_id, scope_id, owner_id),
)
conn.execute(
"""
UPDATE scope_run_requests SET run_id = ?, state = 'active', updated_at = ?
WHERE deployment_id = ? AND scope_id = ? AND run_request_id = ?
""",
(run_id, now, deployment_id, scope_id, run_request_id),
)
return RunReservation(
deployment_id=deployment_id,
scope_id=scope_id,
run_request_id=run_request_id,
turn_id=str(row["turn_id"]),
interrupt_key=(
str(row["interrupt_key"]) if row["interrupt_key"] else None
),
request_hash=str(row["request_hash"]),
run_owner_id=owner_id,
run_id=run_id,
state="active",
)
def reserve_turn(
self,
deployment_id: str,
scope_id: str,
turn_id: str,
request_hash: str,
) -> RunReservation:
"""Backward-compatible reservation for pre-run-request callers."""
return self.reserve_run(
deployment_id,
scope_id,
turn_id,
turn_id,
request_hash,
)
def bind_turn(
self,
deployment_id: str,
scope_id: str,
turn_id: str,
run_id: str,
) -> RunReservation:
"""Backward-compatible binding for pre-run-request callers."""
return self.bind_run(deployment_id, scope_id, turn_id, run_id)
def acquire_lock(
self,
deployment_id: str,
lock_name: str,
operation_id: str,
*,
lease_seconds: int = 60,
) -> DeploymentLock:
operation_id = _ensure_uuid(operation_id, "operation_id")
now = datetime.now(UTC)
expires = now + timedelta(seconds=max(1, lease_seconds))
now_text, expires_text = now.isoformat(), expires.isoformat()
with self._transaction() as conn:
row = conn.execute(
"""
SELECT operation_id, expires_at FROM deployment_locks
WHERE deployment_id = ? AND lock_name = ?
""",
(deployment_id, lock_name),
).fetchone()
if (
row is not None
and str(row["operation_id"]) != operation_id
and str(row["expires_at"]) > now_text
):
raise ScopeConflictError("deployment lock is held")
conn.execute(
"""
INSERT INTO deployment_locks(
deployment_id, lock_name, operation_id, expires_at, created_at
) VALUES (?, ?, ?, ?, ?)
ON CONFLICT(deployment_id, lock_name) DO UPDATE SET
operation_id = excluded.operation_id,
expires_at = excluded.expires_at,
created_at = excluded.created_at
""",
(deployment_id, lock_name, operation_id, expires_text, now_text),
)
return DeploymentLock(deployment_id, lock_name, operation_id, expires_text)
def release_lock(
self, deployment_id: str, lock_name: str, operation_id: str
) -> None:
with self._transaction() as conn:
result = conn.execute(
"""
DELETE FROM deployment_locks
WHERE deployment_id = ? AND lock_name = ? AND operation_id = ?
""",
(deployment_id, lock_name, operation_id),
)
if result.rowcount != 1:
raise ScopeAccessError("deployment lock is not held by this operation")
def renew_lock(
self,
deployment_id: str,
lock_name: str,
operation_id: str,
*,
lease_seconds: int = 60,
) -> DeploymentLock:
"""Renew a lease only while this operation still owns an active lock."""
operation_id = _ensure_uuid(operation_id, "operation_id")
now = datetime.now(UTC)
now_text = now.isoformat()
expires_text = (now + timedelta(seconds=max(1, lease_seconds))).isoformat()
with self._transaction() as conn:
result = conn.execute(
"""
UPDATE deployment_locks SET expires_at = ?
WHERE deployment_id = ? AND lock_name = ? AND operation_id = ?
AND expires_at > ?
""",
(expires_text, deployment_id, lock_name, operation_id, now_text),
)
if result.rowcount != 1:
raise ScopeAccessError(
"deployment lock expired or belongs to another operation"
)
return DeploymentLock(deployment_id, lock_name, operation_id, expires_text)
def active_lock(self, deployment_id: str, lock_name: str) -> DeploymentLock | None:
"""Return an unexpired deployment lock without mutating its lease."""
self.initialize()
now = _utc_now()
with self._connect() as conn:
row = conn.execute(
"""
SELECT operation_id, expires_at FROM deployment_locks
WHERE deployment_id = ? AND lock_name = ? AND expires_at > ?
""",
(deployment_id, lock_name, now),
).fetchone()
if row is None:
return None
return DeploymentLock(
deployment_id=deployment_id,
lock_name=lock_name,
operation_id=str(row["operation_id"]),
expires_at=str(row["expires_at"]),
)
def begin_operation(
self,
deployment_id: str,
operation_id: str,
*,
kind: str,
scope_id: str | None = None,
) -> None:
operation_id = _ensure_uuid(operation_id, "operation_id")
if scope_id is not None:
scope_id = _ensure_uuid(scope_id, "scope_id")
now = _utc_now()
with self._transaction() as conn:
existing = conn.execute(
"""
SELECT kind, state FROM scope_operations
WHERE deployment_id = ? AND operation_id = ?
""",
(deployment_id, operation_id),
).fetchone()
if existing is not None:
if str(existing["kind"]) != kind:
raise ScopeConflictError("operation id belongs to another kind")
return
conn.execute(
"""
INSERT INTO scope_operations(
deployment_id, operation_id, scope_id, kind, state, created_at, updated_at
) VALUES (?, ?, ?, ?, 'running', ?, ?)
""",
(deployment_id, operation_id, scope_id, kind, now, now),
)
def finish_operation(
self,
deployment_id: str,
operation_id: str,
*,
state: str,
result_sha256: str | None = None,
last_error_code: str | None = None,
) -> None:
if state not in {"completed", "failed"}:
raise ScopeRegistryError(
"operation terminal state must be completed or failed"
)
operation_id = _ensure_uuid(operation_id, "operation_id")
with self._transaction() as conn:
result = conn.execute(
"""
UPDATE scope_operations
SET state = ?, result_sha256 = ?, last_error_code = ?, updated_at = ?
WHERE deployment_id = ? AND operation_id = ? AND state = 'running'
""",
(
state,
result_sha256,
last_error_code,
_utc_now(),
deployment_id,
operation_id,
),
)
if result.rowcount != 1:
raise ScopeAccessError("operation is not running")
def get_operation(self, deployment_id: str, operation_id: str) -> ScopeOperation:
operation_id = _ensure_uuid(operation_id, "operation_id")
self.initialize()
with self._connect() as conn:
row = conn.execute(
"""
SELECT scope_id, kind, state, result_sha256, last_error_code
FROM scope_operations WHERE deployment_id = ? AND operation_id = ?
""",
(deployment_id, operation_id),
).fetchone()
if row is None:
raise ScopeNotFoundError("scope operation not found")
return ScopeOperation(
deployment_id=deployment_id,
operation_id=operation_id,
scope_id=str(row["scope_id"]) if row["scope_id"] else None,
kind=str(row["kind"]),
state=str(row["state"]),
result_sha256=(
str(row["result_sha256"]) if row["result_sha256"] else None
),
last_error_code=(
str(row["last_error_code"]) if row["last_error_code"] else None
),
)
_registry_cache: dict[Path, ScopeRegistry] = {}
_registry_cache_lock = threading.Lock()
def get_scope_registry(workspace_root: Path | str | None = None) -> ScopeRegistry:
path = default_registry_path(workspace_root)
with _registry_cache_lock:
registry = _registry_cache.get(path)
if registry is None:
registry = ScopeRegistry(path)
_registry_cache[path] = registry
registry.initialize()
return registry