diff --git "a/gdw_workspace.py" "b/gdw_workspace.py" --- "a/gdw_workspace.py" +++ "b/gdw_workspace.py" @@ -1,61 +1,372 @@ -"""Concurrency-safe SQLite state, idempotency, and receipt storage for GDW.""" +"""Tenant-safe SQLite state, idempotency, quota, and outbox storage for GDW.""" import contextlib import hashlib import json import os -import secrets +import re import sqlite3 import threading +import uuid +from dataclasses import dataclass from datetime import datetime, timedelta, timezone from pathlib import Path from typing import Any, Dict, Iterator, Optional, Tuple +SCHEMA_VERSION = 3 _SCHEMA_LOCK = threading.RLock() _PROCESS_WRITE_LOCK = threading.RLock() -_INITIALISED_PATHS = set() -_OBJECT_TYPES = {"request", "session"} -_EFFECT_KINDS = {"receipt_projection", "proof_export"} +_IDENTIFIER = re.compile(r"^[A-Za-z0-9._:-]{1,128}$") +_TRUTHY = {"1", "true", "yes", "on"} +_JOURNAL_MODES = {"DELETE", "WAL"} +_SYNCHRONOUS_MODES = {"FULL", "NORMAL"} +_V1_TABLES = ( + "session_state", + "requests", + "receipts", + "proof_outbox", + "effect_outbox", +) -def _canonical_json(value: Any) -> str: - return json.dumps( - value, sort_keys=True, separators=(",", ":"), ensure_ascii=True - ) +class GDWWorkspaceError(RuntimeError): + """Base class for workspace failures that callers may map to safe responses.""" + + +class GDWConfigurationError(GDWWorkspaceError): + """The persistent workspace was not configured safely.""" + + +class GDWSchemaError(GDWWorkspaceError): + """The database schema cannot be used without an explicit operator action.""" + + +class GDWLegacyMigrationRequired(GDWSchemaError): + """Legacy rows cannot be assigned to an owner implicitly.""" + + +class GDWQuotaExceeded(GDWWorkspaceError): + """A transactional owner or global resource ceiling was reached.""" + + def __init__(self, code: str): + super().__init__(code) + self.code = code + + +class GDWLifecycleError(GDWWorkspaceError): + """An expired or tombstoned object cannot be mutated or replayed.""" + + +def _env_int(name: str, default: int, minimum: int = 1) -> int: + raw = os.environ.get(name) + if raw is None: + return default + try: + value = int(raw) + except ValueError as exc: + raise GDWConfigurationError(f"{name} must be an integer") from exc + if value < minimum: + raise GDWConfigurationError(f"{name} must be at least {minimum}") + return value + + +def _utc_now() -> datetime: + return datetime.now(timezone.utc) + + +def _normalise_time(value: Optional[Any] = None) -> datetime: + if value is None: + return _utc_now() + if isinstance(value, datetime): + parsed = value + else: + parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00")) + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc) -def _sha256_json(value: Any) -> str: - return hashlib.sha256(_canonical_json(value).encode("utf-8")).hexdigest() +def _text_time(value: Optional[Any] = None) -> str: + return _normalise_time(value).isoformat() + + +def _json_text(value: Any) -> str: + return json.dumps(value, sort_keys=True, separators=(",", ":")) + + +def _byte_len(value: Optional[str]) -> int: + return len(value.encode("utf-8")) if value is not None else 0 + + +def _checked_identity(value: Optional[str], field: str) -> str: + candidate = (value or "").strip() + if not _IDENTIFIER.fullmatch(candidate): + raise GDWConfigurationError( + f"{field} must contain 1-128 canonical identifier characters" + ) + return candidate + + +@dataclass(frozen=True) +class GDWQuotaPolicy: + """Hard limits enforced while the same SQLite write transaction is held.""" + + owner_active_sessions: int = 1_000 + owner_active_requests: int = 100_000 + owner_pending_effects: int = 10_000 + owner_stored_bytes: int = 256 * 1024 * 1024 + global_active_sessions: int = 10_000 + global_active_requests: int = 1_000_000 + global_pending_effects: int = 100_000 + global_stored_bytes: int = 2 * 1024 * 1024 * 1024 + + @classmethod + def from_environment(cls) -> "GDWQuotaPolicy": + return cls( + owner_active_sessions=_env_int( + "GDW_OWNER_MAX_ACTIVE_SESSIONS", cls.owner_active_sessions + ), + owner_active_requests=_env_int( + "GDW_OWNER_MAX_ACTIVE_REQUESTS", cls.owner_active_requests + ), + owner_pending_effects=_env_int( + "GDW_OWNER_MAX_PENDING_EFFECTS", cls.owner_pending_effects + ), + owner_stored_bytes=_env_int( + "GDW_OWNER_MAX_STORED_BYTES", cls.owner_stored_bytes + ), + global_active_sessions=_env_int( + "GDW_GLOBAL_MAX_ACTIVE_SESSIONS", cls.global_active_sessions + ), + global_active_requests=_env_int( + "GDW_GLOBAL_MAX_ACTIVE_REQUESTS", cls.global_active_requests + ), + global_pending_effects=_env_int( + "GDW_GLOBAL_MAX_PENDING_EFFECTS", cls.global_pending_effects + ), + global_stored_bytes=_env_int( + "GDW_GLOBAL_MAX_STORED_BYTES", cls.global_stored_bytes + ), + ) -def _effect_identity_key( - *, - generation_id: str, - owner_id: str, - request_id: str, - request_digest: str, - kind: str, - canonical_identity: str, - payload_sha256: str, -) -> str: - return hashlib.sha256( - ( - f"{generation_id}:{owner_id}:{request_id}:{request_digest}:" - f"{kind}:{canonical_identity}:{payload_sha256}" - ).encode("utf-8") - ).hexdigest() +_SCHEMA_STATEMENTS = ( + """ + CREATE TABLE schema_meta ( + schema_name TEXT PRIMARY KEY, + schema_version INTEGER NOT NULL, + database_generation_id TEXT NOT NULL, + created_at TEXT NOT NULL, + upgraded_at TEXT NOT NULL + ) + """, + """ + CREATE TABLE usage ( + namespace TEXT NOT NULL, + owner_id TEXT NOT NULL, + active_sessions INTEGER NOT NULL DEFAULT 0 CHECK(active_sessions >= 0), + active_requests INTEGER NOT NULL DEFAULT 0 CHECK(active_requests >= 0), + pending_effects INTEGER NOT NULL DEFAULT 0 CHECK(pending_effects >= 0), + stored_bytes INTEGER NOT NULL DEFAULT 0 CHECK(stored_bytes >= 0), + updated_at TEXT NOT NULL, + PRIMARY KEY(namespace, owner_id) + ) + """, + """ + CREATE TABLE session_state ( + namespace TEXT NOT NULL, + owner_id TEXT NOT NULL, + session_id TEXT NOT NULL, + step INTEGER NOT NULL CHECK(step >= 0), + state_json TEXT, + state_hash TEXT NOT NULL, + lifecycle TEXT NOT NULL + CHECK(lifecycle IN ('ACTIVE', 'EXPIRING', 'TOMBSTONED')), + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL, + last_accessed_at TEXT NOT NULL, + expires_at TEXT, + tombstoned_at TEXT, + PRIMARY KEY(namespace, owner_id, session_id) + ) + """, + """ + CREATE TABLE requests ( + namespace TEXT NOT NULL, + owner_id TEXT NOT NULL, + request_id TEXT NOT NULL, + request_digest TEXT NOT NULL, + session_id TEXT NOT NULL, + response_json TEXT, + response_hash TEXT NOT NULL, + lifecycle TEXT NOT NULL CHECK(lifecycle IN ('ACTIVE', 'TOMBSTONED')), + created_at TEXT NOT NULL, + expires_at TEXT, + tombstoned_at TEXT, + PRIMARY KEY(namespace, owner_id, request_id) + ) + """, + """ + CREATE TABLE receipts ( + namespace TEXT NOT NULL, + owner_id TEXT NOT NULL, + receipt_hash TEXT NOT NULL, + request_id TEXT NOT NULL, + session_id TEXT NOT NULL, + step INTEGER NOT NULL, + receipt_json TEXT, + created_at TEXT NOT NULL, + tombstoned_at TEXT, + PRIMARY KEY(namespace, owner_id, receipt_hash), + UNIQUE(namespace, owner_id, request_id), + FOREIGN KEY(namespace, owner_id, request_id) + REFERENCES requests(namespace, owner_id, request_id) + DEFERRABLE INITIALLY DEFERRED + ) + """, + """ + CREATE TABLE proof_outbox ( + namespace TEXT NOT NULL, + owner_id TEXT NOT NULL, + proposal_id TEXT NOT NULL, + payload_json TEXT, + payload_sha256 TEXT NOT NULL, + status TEXT NOT NULL + CHECK(status IN ('PENDING', 'EXPORTED', 'DEAD_LETTER')), + artifact_json TEXT, + created_at TEXT NOT NULL, + exported_at TEXT, + tombstoned_at TEXT, + PRIMARY KEY(namespace, owner_id, proposal_id) + ) + """, + """ + CREATE TABLE effect_outbox ( + namespace TEXT NOT NULL, + owner_id TEXT NOT NULL, + idempotency_key TEXT NOT NULL, + database_generation_id TEXT NOT NULL, + request_id TEXT NOT NULL, + kind TEXT NOT NULL CHECK(kind IN ('receipt_projection', 'proof_export')), + receipt_hash TEXT, + payload_json TEXT, + payload_sha256 TEXT NOT NULL, + intent_sha256 TEXT NOT NULL, + status TEXT NOT NULL + CHECK(status IN ('PENDING', 'CLAIMED', 'EXPORTED', 'DEAD_LETTER')), + attempts INTEGER NOT NULL DEFAULT 0 CHECK(attempts >= 0), + max_attempts INTEGER NOT NULL CHECK(max_attempts > 0), + next_attempt_at TEXT NOT NULL, + lease_owner TEXT, + lease_until TEXT, + claim_generation INTEGER NOT NULL DEFAULT 0 + CHECK(claim_generation >= 0), + last_error TEXT, + artifact_json TEXT, + created_at TEXT NOT NULL, + exported_at TEXT, + tombstoned_at TEXT, + PRIMARY KEY(namespace, owner_id, idempotency_key), + UNIQUE(namespace, owner_id, request_id, kind), + FOREIGN KEY(namespace, owner_id, request_id) + REFERENCES requests(namespace, owner_id, request_id) + DEFERRABLE INITIALLY DEFERRED, + FOREIGN KEY(namespace, owner_id, receipt_hash) + REFERENCES receipts(namespace, owner_id, receipt_hash) + DEFERRABLE INITIALLY DEFERRED + ) + """, + """ + CREATE INDEX idx_sessions_lifecycle + ON session_state(namespace, owner_id, lifecycle, expires_at) + """, + """ + CREATE INDEX idx_requests_session + ON requests(namespace, owner_id, session_id, lifecycle) + """, + """ + CREATE INDEX idx_receipts_session + ON receipts(namespace, owner_id, session_id, step) + """, + """ + CREATE INDEX idx_proof_outbox_status + ON proof_outbox(namespace, owner_id, status, created_at) + """, + """ + CREATE INDEX idx_effect_outbox_status + ON effect_outbox( + namespace, owner_id, status, next_attempt_at, lease_until, created_at + ) + """, +) class GDWWorkspace: - def __init__(self, path: Optional[str] = None): + """SQLite workspace bound to one stable owner and namespace by default.""" + + def __init__( + self, + path: Optional[str] = None, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + production: Optional[bool] = None, + quota_policy: Optional[GDWQuotaPolicy] = None, + ): + self.production = ( + production + if production is not None + else os.environ.get("GDW_PRODUCTION_MODE", "").lower() in _TRUTHY + ) + explicit_path = path is not None or bool(os.environ.get("GDW_DB_PATH")) configured = path or os.environ.get("GDW_DB_PATH", "output/gdw/gdw.sqlite3") self.path = Path(configured).resolve() - self.journal_mode = os.environ.get("GDW_SQLITE_JOURNAL", "WAL").strip().upper() - if self.journal_mode not in {"DELETE", "WAL"}: - raise RuntimeError("GDW_SQLITE_JOURNAL must be DELETE or WAL") + if self.production and not explicit_path: + raise GDWConfigurationError( + "GDW_DB_PATH must be explicit in production mode" + ) + if self.production and not self.path.is_file(): + raise GDWConfigurationError( + "production GDW database must already exist and be provisioned" + ) + + configured_namespace = namespace or os.environ.get("GDW_NAMESPACE") + configured_owner = ( + owner_id + or os.environ.get("GDW_OWNER_ID") + or os.environ.get("GDW_SERVICE_OWNER_ID") + ) + if self.production and (not configured_namespace or not configured_owner): + raise GDWConfigurationError( + "stable GDW_NAMESPACE and GDW_OWNER_ID are required in production" + ) + self.namespace = _checked_identity(configured_namespace or "a11oy", "namespace") + self.owner_id = _checked_identity(configured_owner or "local-owner", "owner_id") + self.journal_mode = os.environ.get("GDW_SQLITE_JOURNAL", "WAL").upper() + if self.journal_mode not in _JOURNAL_MODES: + raise GDWConfigurationError( + "GDW_SQLITE_JOURNAL must be DELETE or WAL" + ) + self.synchronous_mode = os.environ.get( + "GDW_SQLITE_SYNCHRONOUS", "NORMAL" + ).upper() + if self.synchronous_mode not in _SYNCHRONOUS_MODES: + raise GDWConfigurationError( + "GDW_SQLITE_SYNCHRONOUS must be FULL or NORMAL" + ) + self.quota_policy = quota_policy or GDWQuotaPolicy.from_environment() + self.retention_seconds = _env_int("GDW_RETENTION_SECONDS", 30 * 24 * 60 * 60) + self.tombstone_seconds = _env_int("GDW_TOMBSTONE_SECONDS", 90 * 24 * 60 * 60) + self.effect_max_attempts = _env_int("GDW_EFFECT_MAX_ATTEMPTS", 8) + self.effect_backoff_seconds = _env_int("GDW_EFFECT_BACKOFF_SECONDS", 5) + self.database_generation_id = "" self._initialise() + @staticmethod + def schema_version() -> int: + return SCHEMA_VERSION + def _connect(self) -> sqlite3.Connection: connection = sqlite3.connect( str(self.path), @@ -65,161 +376,535 @@ class GDWWorkspace: ) connection.row_factory = sqlite3.Row connection.execute("PRAGMA busy_timeout=30000") - connection.execute("PRAGMA synchronous=NORMAL") + observed_journal = str( + connection.execute( + f"PRAGMA journal_mode={self.journal_mode}" + ).fetchone()[0] + ).upper() + if observed_journal != self.journal_mode: + connection.close() + raise GDWConfigurationError( + "SQLite journal mode does not match GDW_SQLITE_JOURNAL" + ) + connection.execute(f"PRAGMA synchronous={self.synchronous_mode}") connection.execute("PRAGMA foreign_keys=ON") return connection + @staticmethod + def _table_names(connection: sqlite3.Connection) -> set: + return { + row[0] + for row in connection.execute( + "SELECT name FROM sqlite_master " + "WHERE type = 'table' AND name NOT LIKE 'sqlite_%'" + ).fetchall() + } + + @staticmethod + def _create_schema( + connection: sqlite3.Connection, + timestamp: str, + generation_id: Optional[str] = None, + ) -> str: + generation = generation_id or uuid.uuid4().hex + for statement in _SCHEMA_STATEMENTS: + connection.execute(statement) + connection.execute( + """ + INSERT INTO schema_meta( + schema_name, schema_version, database_generation_id, + created_at, upgraded_at + ) VALUES ('gdw', ?, ?, ?, ?) + """, + (SCHEMA_VERSION, generation, timestamp, timestamp), + ) + return generation + + @staticmethod + def _legacy_row_count(connection: sqlite3.Connection, tables: set) -> int: + return sum( + int(connection.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]) + for table in _V1_TABLES + if table in tables + ) + + def _migrate_v1(self, connection: sqlite3.Connection, tables: set) -> None: + del tables + timestamp = _text_time() + connection.execute("PRAGMA foreign_keys=OFF") + connection.execute("BEGIN IMMEDIATE") + try: + locked_tables = self._table_names(connection) + nonempty = self._legacy_row_count( + connection, locked_tables + ) > 0 + if nonempty: + raise GDWLegacyMigrationRequired( + "nonempty v1 workspace has no durable owner/provenance " + "binding; an offline evidence-preserving migration is required" + ) + for index in ( + "idx_requests_session", + "idx_receipts_session", + "idx_proof_outbox_status", + "idx_effect_outbox_status", + ): + connection.execute(f"DROP INDEX IF EXISTS {index}") + for table in _V1_TABLES: + if table in locked_tables: + connection.execute(f"ALTER TABLE {table} RENAME TO {table}_v1") + self._create_schema(connection, timestamp) + + for table in reversed(_V1_TABLES): + if table in locked_tables: + connection.execute(f"DROP TABLE {table}_v1") + self._reconcile_usage(connection) + connection.execute("COMMIT") + except Exception: + connection.execute("ROLLBACK") + raise + finally: + connection.execute("PRAGMA foreign_keys=ON") + + def _migrate_v2(self, connection: sqlite3.Connection) -> None: + """Transactionally bind valid v2 effects to a new database generation.""" + + from gdw_proofs import build_proof_payload + + timestamp = _text_time() + generation = uuid.uuid4().hex + connection.execute("PRAGMA foreign_keys=OFF") + connection.execute("BEGIN IMMEDIATE") + try: + version = connection.execute( + "SELECT schema_version FROM schema_meta WHERE schema_name = 'gdw'" + ).fetchone() + if version is None or int(version[0]) != 2: + raise GDWSchemaError("GDW v2 migration source changed under lock") + required = { + "schema_meta", + "usage", + "session_state", + "requests", + "receipts", + "proof_outbox", + "effect_outbox", + } + missing = required - self._table_names(connection) + if missing: + raise GDWSchemaError( + "GDW v2 schema is incomplete: " + ", ".join(sorted(missing)) + ) + connection.execute("DROP INDEX IF EXISTS idx_effect_outbox_status") + connection.execute( + "ALTER TABLE effect_outbox RENAME TO effect_outbox_v2" + ) + connection.execute(_SCHEMA_STATEMENTS[6]) + connection.execute( + "ALTER TABLE schema_meta " + "ADD COLUMN database_generation_id TEXT" + ) + connection.execute( + """ + UPDATE schema_meta + SET schema_version = ?, database_generation_id = ?, + upgraded_at = ? + WHERE schema_name = 'gdw' + """, + (SCHEMA_VERSION, generation, timestamp), + ) + self.database_generation_id = generation + online_migration_blockers = { + "session_state": connection.execute( + "SELECT COUNT(*) FROM session_state " + "WHERE state_json IS NOT NULL" + ).fetchone()[0], + "receipts": connection.execute( + "SELECT COUNT(*) FROM receipts " + "WHERE receipt_json IS NOT NULL" + ).fetchone()[0], + "proof_outbox": connection.execute( + "SELECT COUNT(*) FROM proof_outbox " + "WHERE payload_json IS NOT NULL" + ).fetchone()[0], + "exported_effects": connection.execute( + "SELECT COUNT(*) FROM effect_outbox_v2 " + "WHERE status = 'EXPORTED' OR artifact_json IS NOT NULL" + ).fetchone()[0], + "receipt_effects": connection.execute( + "SELECT COUNT(*) FROM effect_outbox_v2 " + "WHERE kind = 'receipt_projection'" + ).fetchone()[0], + } + blocked = sorted( + name + for name, count in online_migration_blockers.items() + if int(count) + ) + if blocked: + raise GDWLegacyMigrationRequired( + "stateful v2 records require an offline " + "evidence-preserving migration: " + ",".join(blocked) + ) + rows = connection.execute( + "SELECT * FROM effect_outbox_v2 ORDER BY namespace, owner_id, " + "created_at, idempotency_key" + ).fetchall() + for source in rows: + if source["payload_json"] is None: + raise GDWLegacyMigrationRequired( + "compacted v2 effects require an offline " + "evidence-preserving migration" + ) + try: + payload = json.loads(source["payload_json"]) + except (TypeError, json.JSONDecodeError) as exc: + raise GDWLegacyMigrationRequired( + "v2 effect payload is not canonical JSON" + ) from exc + namespace = str(source["namespace"]) + owner_id = str(source["owner_id"]) + request_id = str(source["request_id"]) + kind = str(source["kind"]) + source_payload_sha256 = self._effect_payload_digest(kind, payload) + if source["payload_sha256"] != source_payload_sha256: + raise GDWLegacyMigrationRequired( + "v2 effect payload digest is invalid" + ) + source_key = self.scoped_effect_key( + namespace, + owner_id, + request_id, + kind, + source_payload_sha256, + ) + if source["idempotency_key"] != source_key: + raise GDWLegacyMigrationRequired( + "v2 effect idempotency binding is invalid" + ) + if kind != "proof_export": + raise GDWLegacyMigrationRequired( + "v2 receipt effects require an offline " + "evidence-preserving migration" + ) + request_row = connection.execute( + """ + SELECT request_digest, session_id, response_json, + response_hash + FROM requests + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (namespace, owner_id, request_id), + ).fetchone() + if request_row is None or request_row["response_json"] is None: + raise GDWLegacyMigrationRequired( + "v2 effect request anchor is unavailable" + ) + try: + response = json.loads(request_row["response_json"]) + governance = response["audit"]["governance"] + except (KeyError, TypeError, json.JSONDecodeError) as exc: + raise GDWLegacyMigrationRequired( + "v2 effect request anchor is invalid" + ) from exc + if ( + hashlib.sha256( + _json_text(response).encode("utf-8") + ).hexdigest() + != request_row["response_hash"] + ): + raise GDWLegacyMigrationRequired( + "v2 effect request response hash is invalid" + ) + principal = response.get("principal") + if not isinstance(principal, dict): + principal = { + "namespace": namespace, + "owner_id": owner_id, + "key_id": "UNAVAILABLE_V2", + } + if ( + principal.get("namespace") != namespace + or principal.get("owner_id") != owner_id + ): + raise GDWLegacyMigrationRequired( + "v2 effect principal binding is invalid" + ) + legacy_proof = build_proof_payload( + proposal_id=response["proposal_id"], + request_id=response["request_id"], + request_digest=str(request_row["request_digest"]), + namespace=namespace, + owner_id=owner_id, + database_generation_id=str( + response.get("database_generation_id") or "" + ), + step=int(response["step"]), + before_hash=response["state_before_hash"], + after_hash=response["state_hash"], + decision=response["decision"], + scheduler_mode=response["scheduler_mode"], + receipt_hash=response.get("receipt_hash") or "", + dry_run=bool(response["dry_run"]), + governance=governance, + ) + if "database_generation_id" not in payload: + for field in ( + "request_digest", + "namespace", + "owner_id", + "database_generation_id", + ): + legacy_proof.pop(field, None) + legacy_proof.pop("payload_sha256", None) + legacy_proof["payload_sha256"] = hashlib.sha256( + _json_text(legacy_proof).encode("utf-8") + ).hexdigest() + if payload != legacy_proof: + raise GDWLegacyMigrationRequired( + "v2 proof effect differs from persisted request" + ) + rebound_response = dict(response) + rebound_response["request_digest"] = str( + request_row["request_digest"] + ) + rebound_response["database_generation_id"] = generation + rebound_response["principal"] = principal + rebound_response["proposal_id"] = hashlib.sha256( + _json_text( + { + "schema": "szl.gdw.proposal-identity/v1", + "database_generation_id": generation, + "namespace": namespace, + "owner_id": owner_id, + "request_id": request_id, + "request_digest": str( + request_row["request_digest"] + ), + "state_before_hash": response[ + "state_before_hash" + ], + "governance_evidence_sha256": hashlib.sha256( + _json_text(governance).encode("utf-8") + ).hexdigest(), + } + ).encode("utf-8") + ).hexdigest() + rebound_response_text = _json_text(rebound_response) + connection.execute( + """ + UPDATE requests + SET response_json = ?, response_hash = ? + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + ( + rebound_response_text, + hashlib.sha256( + rebound_response_text.encode("utf-8") + ).hexdigest(), + namespace, + owner_id, + request_id, + ), + ) + request_anchor = self._request_anchor( + connection, namespace, owner_id, request_id + ) + payload = self._expected_proof_payload(request_anchor) + payload_sha256 = self._effect_payload_digest(kind, payload) + canonical_key = self.scoped_effect_key( + namespace, + owner_id, + request_id, + kind, + payload_sha256, + ) + receipt_hash = None + intent_sha256 = self._canonical_effect_intent( + request_anchor, + namespace=namespace, + owner_id=owner_id, + request_id=request_id, + kind=kind, + payload_sha256=payload_sha256, + receipt_hash=receipt_hash, + ) + status = source["status"] + if status == "CLAIMED": + status = ( + "DEAD_LETTER" + if int(source["attempts"]) >= int(source["max_attempts"]) + else "PENDING" + ) + connection.execute( + """ + INSERT INTO effect_outbox( + namespace, owner_id, idempotency_key, + database_generation_id, request_id, kind, receipt_hash, + payload_json, payload_sha256, intent_sha256, status, + attempts, max_attempts, next_attempt_at, lease_owner, + lease_until, claim_generation, last_error, artifact_json, + created_at, exported_at, tombstoned_at + ) VALUES ( + ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, + 0, ?, ?, ?, ?, ? + ) + """, + ( + namespace, + owner_id, + canonical_key, + generation, + request_id, + kind, + receipt_hash, + _json_text(payload), + payload_sha256, + intent_sha256, + status, + int(source["attempts"]), + int(source["max_attempts"]), + source["next_attempt_at"], + ( + source["lease_owner"] + if status == "CLAIMED" + else None + ), + ( + source["lease_until"] + if status == "CLAIMED" + else None + ), + source["last_error"], + source["artifact_json"], + source["created_at"], + source["exported_at"], + source["tombstoned_at"], + ), + ) + connection.execute("DROP TABLE effect_outbox_v2") + connection.execute(_SCHEMA_STATEMENTS[-1]) + self._reconcile_usage(connection) + connection.execute("COMMIT") + except Exception: + connection.execute("ROLLBACK") + raise + finally: + connection.execute("PRAGMA foreign_keys=ON") + + @staticmethod + def _validate_schema(connection: sqlite3.Connection) -> str: + row = connection.execute( + """ + SELECT schema_version, database_generation_id + FROM schema_meta WHERE schema_name = 'gdw' + """ + ).fetchone() + if row is None or int(row[0]) != SCHEMA_VERSION: + observed = None if row is None else int(row[0]) + raise GDWSchemaError( + f"GDW schema version {observed!r} is unsupported; " + f"expected {SCHEMA_VERSION}" + ) + required = { + "schema_meta", + "usage", + "session_state", + "requests", + "receipts", + "proof_outbox", + "effect_outbox", + } + missing = required - GDWWorkspace._table_names(connection) + if missing: + raise GDWSchemaError( + "GDW schema is incomplete: " + ", ".join(sorted(missing)) + ) + expected_columns = { + "session_state": {"namespace", "owner_id", "lifecycle", "expires_at"}, + "requests": {"namespace", "owner_id", "lifecycle", "expires_at"}, + "effect_outbox": { + "namespace", + "owner_id", + "max_attempts", + "next_attempt_at", + "database_generation_id", + "intent_sha256", + "claim_generation", + }, + } + for table, expected in expected_columns.items(): + columns = { + row["name"] for row in connection.execute(f"PRAGMA table_info({table})") + } + if not expected.issubset(columns): + raise GDWSchemaError( + f"GDW schema table {table} is missing required columns" + ) + generation_id = str(row["database_generation_id"] or "") + if not re.fullmatch(r"[0-9a-f]{32}", generation_id): + raise GDWSchemaError("GDW database generation id is invalid") + return generation_id + def _initialise(self) -> None: - key = str(self.path) - if key in _INITIALISED_PATHS: - return with _SCHEMA_LOCK: - if key in _INITIALISED_PATHS: - return - self.path.parent.mkdir(parents=True, exist_ok=True) - connection = sqlite3.connect(str(self.path), timeout=30.0) + if not self.path.exists(): + if self.production: + raise GDWConfigurationError( + "production GDW database must already exist" + ) + self.path.parent.mkdir(parents=True, exist_ok=True) + connection = sqlite3.connect( + str(self.path), timeout=30.0, isolation_level=None + ) + connection.row_factory = sqlite3.Row try: - observed_journal = str( + connection.execute("PRAGMA busy_timeout=30000") + selected_journal = str( connection.execute( f"PRAGMA journal_mode={self.journal_mode}" ).fetchone()[0] ).upper() - if observed_journal != self.journal_mode: - raise RuntimeError( - "configured SQLite journal mode did not converge" + if selected_journal != self.journal_mode: + raise GDWConfigurationError( + "SQLite journal mode does not match GDW_SQLITE_JOURNAL" ) - connection.execute("PRAGMA synchronous=NORMAL") + connection.execute(f"PRAGMA synchronous={self.synchronous_mode}") connection.execute("PRAGMA foreign_keys=ON") - self._ensure_schema(connection) - connection.commit() - _INITIALISED_PATHS.add(key) + tables = self._table_names(connection) + if not tables: + if self.production: + raise GDWSchemaError( + "production GDW database has no provisioned schema" + ) + connection.execute("BEGIN IMMEDIATE") + try: + self._create_schema(connection, _text_time()) + connection.execute("COMMIT") + except Exception: + connection.execute("ROLLBACK") + raise + elif "schema_meta" not in tables: + self._migrate_v1(connection, tables) + else: + version = connection.execute( + "SELECT schema_version FROM schema_meta " + "WHERE schema_name = 'gdw'" + ).fetchone() + if version is not None and int(version[0]) == 2: + self._migrate_v2(connection) + self.database_generation_id = self._validate_schema(connection) finally: connection.close() - @staticmethod - def _ensure_schema(connection: sqlite3.Connection) -> None: - connection.executescript( - """ - CREATE TABLE IF NOT EXISTS session_state ( - session_id TEXT PRIMARY KEY, - step INTEGER NOT NULL CHECK(step >= 0), - state_json TEXT NOT NULL, - state_hash TEXT NOT NULL, - updated_at TEXT NOT NULL - ); - CREATE TABLE IF NOT EXISTS requests ( - request_id TEXT PRIMARY KEY, - request_digest TEXT NOT NULL, - session_id TEXT NOT NULL, - response_json TEXT NOT NULL, - response_hash TEXT NOT NULL, - created_at TEXT NOT NULL - ); - CREATE TABLE IF NOT EXISTS receipts ( - receipt_hash TEXT PRIMARY KEY, - request_id TEXT NOT NULL UNIQUE, - session_id TEXT NOT NULL, - step INTEGER NOT NULL, - receipt_json TEXT NOT NULL, - created_at TEXT NOT NULL, - FOREIGN KEY(request_id) REFERENCES requests(request_id) - DEFERRABLE INITIALLY DEFERRED - ); - CREATE TABLE IF NOT EXISTS proof_outbox ( - proposal_id TEXT PRIMARY KEY, - payload_json TEXT NOT NULL, - payload_sha256 TEXT NOT NULL, - status TEXT NOT NULL CHECK(status IN ('PENDING', 'EXPORTED')), - artifact_json TEXT, - created_at TEXT NOT NULL, - exported_at TEXT - ); - CREATE TABLE IF NOT EXISTS effect_outbox ( - idempotency_key TEXT PRIMARY KEY, - request_id TEXT NOT NULL, - kind TEXT NOT NULL - CHECK(kind IN ('receipt_projection', 'proof_export')), - generation_id TEXT NOT NULL, - owner_id TEXT NOT NULL, - canonical_identity TEXT NOT NULL, - payload_json TEXT NOT NULL, - payload_sha256 TEXT NOT NULL, - status TEXT NOT NULL - CHECK(status IN ('PENDING', 'CLAIMED', 'EXPORTED')), - attempts INTEGER NOT NULL DEFAULT 0 CHECK(attempts >= 0), - lease_owner TEXT, - lease_until TEXT, - claim_token TEXT, - last_error TEXT, - artifact_json TEXT, - created_at TEXT NOT NULL, - exported_at TEXT, - UNIQUE(request_id, kind), - FOREIGN KEY(request_id) REFERENCES requests(request_id) - DEFERRABLE INITIALLY DEFERRED - ); - CREATE TABLE IF NOT EXISTS workspace_meta ( - key TEXT PRIMARY KEY, - value TEXT NOT NULL - ); - CREATE TABLE IF NOT EXISTS object_owners ( - object_type TEXT NOT NULL - CHECK(object_type IN ('request', 'session')), - object_id TEXT NOT NULL, - owner_id TEXT NOT NULL, - created_at TEXT NOT NULL, - expires_at TEXT NOT NULL, - PRIMARY KEY(object_type, object_id) - ); - CREATE TABLE IF NOT EXISTS evidence_intents ( - canonical_identity TEXT PRIMARY KEY, - request_id TEXT NOT NULL, - kind TEXT NOT NULL - CHECK(kind IN ('receipt_projection', 'proof_export')), - generation_id TEXT NOT NULL, - owner_id TEXT NOT NULL, - payload_json TEXT NOT NULL, - payload_sha256 TEXT NOT NULL, - created_at TEXT NOT NULL, - UNIQUE(request_id, kind), - FOREIGN KEY(request_id) REFERENCES requests(request_id) - DEFERRABLE INITIALLY DEFERRED - ); - CREATE INDEX IF NOT EXISTS idx_requests_session ON requests(session_id); - CREATE INDEX IF NOT EXISTS idx_receipts_session ON receipts(session_id, step); - CREATE INDEX IF NOT EXISTS idx_proof_outbox_status - ON proof_outbox(status, created_at); - CREATE INDEX IF NOT EXISTS idx_effect_outbox_status - ON effect_outbox(status, lease_until, created_at); - CREATE INDEX IF NOT EXISTS idx_object_owners_owner - ON object_owners(owner_id, object_type, expires_at); - """ - ) - # Forward-only compatibility for a local database first opened by the - # unmerged predecessor. Legacy rows remain visibly unbound and make - # integrity fail closed; only the table shape is migrated. - columns = { - row[1] - for row in connection.execute("PRAGMA table_info(effect_outbox)") - } - migrations = { - "generation_id": "TEXT", - "owner_id": "TEXT", - "canonical_identity": "TEXT", - "claim_token": "TEXT", - } - for name, declaration in migrations.items(): - if name not in columns: - connection.execute( - f"ALTER TABLE effect_outbox ADD COLUMN {name} {declaration}" - ) - connection.execute( - "CREATE INDEX IF NOT EXISTS idx_effect_outbox_owner " - "ON effect_outbox(owner_id, status, created_at)" - ) - connection.execute( - "INSERT OR IGNORE INTO workspace_meta(key, value) VALUES (?, ?)", - ("generation_id", secrets.token_hex(16)), + def _identity( + self, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + ) -> Tuple[str, str]: + return ( + _checked_identity(namespace or self.namespace, "namespace"), + _checked_identity(owner_id or self.owner_id, "owner_id"), ) @contextlib.contextmanager @@ -236,315 +921,236 @@ class GDWWorkspace: finally: connection.close() - def generation_id(self) -> str: - connection = self._connect() - try: - row = connection.execute( - "SELECT value FROM workspace_meta WHERE key = 'generation_id'" - ).fetchone() - if row is None or len(row["value"]) != 32: - raise RuntimeError("workspace generation identity is unavailable") - return str(row["value"]) - finally: - connection.close() - - @staticmethod - def object_owner( - connection: sqlite3.Connection, - object_type: str, - object_id: str, - ) -> Optional[str]: - if object_type not in _OBJECT_TYPES: - raise ValueError("unsupported object owner type") - row = connection.execute( - "SELECT owner_id FROM object_owners " - "WHERE object_type = ? AND object_id = ?", - (object_type, object_id), - ).fetchone() - return None if row is None else str(row["owner_id"]) - - @classmethod - def require_object_owner( - cls, - connection: sqlite3.Connection, - object_type: str, - object_id: str, - owner_id: str, - ) -> None: - persisted = cls.object_owner(connection, object_type, object_id) - if persisted is None or persisted != owner_id: - raise PermissionError(f"{object_type} is not owned by this principal") - - @classmethod - def claim_object_owner( - cls, - connection: sqlite3.Connection, - object_type: str, - object_id: str, - owner_id: str, - created_at: str, - expires_at: str, - ) -> None: - persisted = cls.object_owner(connection, object_type, object_id) - if persisted is not None: - if persisted != owner_id: - raise PermissionError( - f"{object_type} is not owned by this principal" - ) - connection.execute( - "UPDATE object_owners SET expires_at = ? " - "WHERE object_type = ? AND object_id = ?", - (expires_at, object_type, object_id), - ) - return + def _usage_row( + self, connection: sqlite3.Connection, namespace: str, owner_id: str + ) -> sqlite3.Row: connection.execute( - "INSERT INTO object_owners(" - "object_type, object_id, owner_id, created_at, expires_at" - ") VALUES (?, ?, ?, ?, ?)", - (object_type, object_id, owner_id, created_at, expires_at), + """ + INSERT OR IGNORE INTO usage( + namespace, owner_id, active_sessions, active_requests, + pending_effects, stored_bytes, updated_at + ) VALUES (?, ?, 0, 0, 0, 0, ?) + """, + (namespace, owner_id, _text_time()), ) + return connection.execute( + "SELECT * FROM usage WHERE namespace = ? AND owner_id = ?", + (namespace, owner_id), + ).fetchone() - @staticmethod - def reclaim_expired( - connection: sqlite3.Connection, - now_text: str, - ) -> Dict[str, int]: - reclaimed = {"requests": 0, "sessions": 0} - requests = connection.execute( - "SELECT object_id FROM object_owners " - "WHERE object_type = 'request' AND expires_at <= ?", - (now_text,), - ).fetchall() - for row in requests: - request_id = str(row["object_id"]) - pending = int( - connection.execute( - "SELECT COUNT(*) FROM effect_outbox " - "WHERE request_id = ? AND status != 'EXPORTED'", - (request_id,), - ).fetchone()[0] - ) - if pending: - continue - connection.execute( - "DELETE FROM effect_outbox WHERE request_id = ?", (request_id,) - ) - connection.execute( - "DELETE FROM evidence_intents WHERE request_id = ?", (request_id,) - ) - connection.execute( - "DELETE FROM receipts WHERE request_id = ?", (request_id,) - ) - connection.execute( - "DELETE FROM requests WHERE request_id = ?", (request_id,) - ) - connection.execute( - "DELETE FROM object_owners " - "WHERE object_type = 'request' AND object_id = ?", - (request_id,), - ) - reclaimed["requests"] += 1 - - sessions = connection.execute( - "SELECT object_id FROM object_owners " - "WHERE object_type = 'session' AND expires_at <= ?", - (now_text,), - ).fetchall() - for row in sessions: - session_id = str(row["object_id"]) - retained_requests = int( - connection.execute( - """ - SELECT COUNT(*) FROM requests r - JOIN object_owners o - ON o.object_type = 'request' - AND o.object_id = r.request_id - WHERE r.session_id = ? - AND ( - o.expires_at > ? - OR EXISTS ( - SELECT 1 FROM effect_outbox e - WHERE e.request_id = r.request_id - AND e.status != 'EXPORTED' - ) - ) - """, - (session_id, now_text), - ).fetchone()[0] - ) - if retained_requests: - continue - connection.execute( - "DELETE FROM session_state WHERE session_id = ?", (session_id,) - ) - connection.execute( - "DELETE FROM object_owners " - "WHERE object_type = 'session' AND object_id = ?", - (session_id,), - ) - reclaimed["sessions"] += 1 - return reclaimed - - @classmethod - def admit_request( - cls, + def _reserve_usage( + self, connection: sqlite3.Connection, - *, + namespace: str, owner_id: str, - request_id: str, - session_id: str, - mutates: bool, - created_at: str, - expires_at: str, - limits: Dict[str, int], + *, + sessions: int = 0, + requests: int = 0, + pending_effects: int = 0, + stored_bytes: int = 0, ) -> None: - cls.reclaim_expired(connection, created_at) - if connection.execute( - "SELECT 1 FROM requests WHERE request_id = ?", (request_id,) - ).fetchone(): - cls.require_object_owner( - connection, "request", request_id, owner_id - ) - return - - owner_requests = int( - connection.execute( - "SELECT COUNT(*) FROM object_owners " - "WHERE object_type = 'request' AND owner_id = ?", - (owner_id,), - ).fetchone()[0] - ) - global_requests = int( - connection.execute( - "SELECT COUNT(*) FROM object_owners " - "WHERE object_type = 'request'" - ).fetchone()[0] - ) - if owner_requests >= limits["owner_requests"]: - raise OverflowError("per-owner request quota exceeded") - if global_requests >= limits["global_requests"]: - raise OverflowError("global request quota exceeded") - - existing_session = connection.execute( - "SELECT 1 FROM session_state WHERE session_id = ?", (session_id,) - ).fetchone() - if existing_session is not None: - cls.claim_object_owner( - connection, - "session", - session_id, - owner_id, - created_at, - expires_at, - ) - elif mutates: - owner_sessions = int( - connection.execute( - "SELECT COUNT(*) FROM object_owners " - "WHERE object_type = 'session' AND owner_id = ?", - (owner_id,), - ).fetchone()[0] - ) - global_sessions = int( - connection.execute( - "SELECT COUNT(*) FROM object_owners " - "WHERE object_type = 'session'" - ).fetchone()[0] - ) - if owner_sessions >= limits["owner_sessions"]: - raise OverflowError("per-owner session quota exceeded") - if global_sessions >= limits["global_sessions"]: - raise OverflowError("global session quota exceeded") - cls.claim_object_owner( - connection, - "session", - session_id, + owner = self._usage_row(connection, namespace, owner_id) + current = { + "sessions": int(owner["active_sessions"]), + "requests": int(owner["active_requests"]), + "pending_effects": int(owner["pending_effects"]), + "stored_bytes": int(owner["stored_bytes"]), + } + deltas = { + "sessions": sessions, + "requests": requests, + "pending_effects": pending_effects, + "stored_bytes": stored_bytes, + } + owner_limits = { + "sessions": self.quota_policy.owner_active_sessions, + "requests": self.quota_policy.owner_active_requests, + "pending_effects": self.quota_policy.owner_pending_effects, + "stored_bytes": self.quota_policy.owner_stored_bytes, + } + global_limits = { + "sessions": self.quota_policy.global_active_sessions, + "requests": self.quota_policy.global_active_requests, + "pending_effects": self.quota_policy.global_pending_effects, + "stored_bytes": self.quota_policy.global_stored_bytes, + } + totals_row = connection.execute( + """ + SELECT COALESCE(SUM(active_sessions), 0) AS sessions, + COALESCE(SUM(active_requests), 0) AS requests, + COALESCE(SUM(pending_effects), 0) AS pending_effects, + COALESCE(SUM(stored_bytes), 0) AS stored_bytes + FROM usage + """ + ).fetchone() + for name, delta in deltas.items(): + proposed_owner = current[name] + delta + if proposed_owner < 0: + raise GDWSchemaError(f"usage counter underflow: {name}") + if delta > 0 and proposed_owner > owner_limits[name]: + raise GDWQuotaExceeded(f"OWNER_{name.upper()}_QUOTA") + proposed_global = int(totals_row[name]) + delta + if proposed_global < 0: + raise GDWSchemaError(f"global usage counter underflow: {name}") + if delta > 0 and proposed_global > global_limits[name]: + raise GDWQuotaExceeded(f"GLOBAL_{name.upper()}_QUOTA") + connection.execute( + """ + UPDATE usage + SET active_sessions = active_sessions + ?, + active_requests = active_requests + ?, + pending_effects = pending_effects + ?, + stored_bytes = stored_bytes + ?, + updated_at = ? + WHERE namespace = ? AND owner_id = ? + """, + ( + sessions, + requests, + pending_effects, + stored_bytes, + _text_time(), + namespace, owner_id, - created_at, - expires_at, - ) - - cls.claim_object_owner( - connection, - "request", - request_id, - owner_id, - created_at, - expires_at, + ), ) - @staticmethod def cached_request( + self, connection: sqlite3.Connection, request_id: str, - owner_id: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, ) -> Optional[Tuple[str, Dict[str, Any]]]: + ns, owner = self._identity(namespace, owner_id) row = connection.execute( - "SELECT request_digest, response_json FROM requests WHERE request_id = ?", - (request_id,), + """ + SELECT request_digest, response_json, lifecycle + FROM requests + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (ns, owner, request_id), ).fetchone() if row is None: return None - GDWWorkspace.require_object_owner( - connection, "request", request_id, owner_id - ) + if row["lifecycle"] != "ACTIVE" or row["response_json"] is None: + raise GDWLifecycleError("idempotency record is outside its replay window") return row["request_digest"], json.loads(row["response_json"]) - @staticmethod def session_state( + self, connection: sqlite3.Connection, session_id: str, + *, + namespace: Optional[str] = None, owner_id: Optional[str] = None, ) -> Optional[Dict[str, Any]]: - if owner_id is not None: - GDWWorkspace.require_object_owner( - connection, "session", session_id, owner_id - ) + ns, owner = self._identity(namespace, owner_id) row = connection.execute( - "SELECT step, state_json, state_hash, updated_at " - "FROM session_state WHERE session_id = ?", - (session_id,), + """ + SELECT step, state_json, state_hash, updated_at, lifecycle, + expires_at + FROM session_state + WHERE namespace = ? AND owner_id = ? AND session_id = ? + """, + (ns, owner, session_id), ).fetchone() - if row is None: + if row is None or row["lifecycle"] != "ACTIVE" or row["state_json"] is None: return None + connection.execute( + """ + UPDATE session_state SET last_accessed_at = ? + WHERE namespace = ? AND owner_id = ? AND session_id = ? + """, + (_text_time(), ns, owner, session_id), + ) return { + "namespace": ns, + "owner_id": owner, "session_id": session_id, + "database_generation_id": self.database_generation_id, "step": int(row["step"]), "state": json.loads(row["state_json"]), "state_hash": row["state_hash"], "updated_at": row["updated_at"], + "expires_at": row["expires_at"], + "lifecycle": row["lifecycle"], } - @staticmethod def save_state( + self, connection: sqlite3.Connection, session_id: str, step: int, state: Dict[str, Any], state_hash: str, updated_at: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + expires_at: Optional[str] = None, ) -> None: + ns, owner = self._identity(namespace, owner_id) + if state.get("database_generation_id") != self.database_generation_id: + raise GDWConfigurationError( + "state database generation does not match workspace" + ) + state_text = _json_text(state) + timestamp = _text_time(updated_at) + expiry = ( + expires_at + or ( + _normalise_time(timestamp) + timedelta(seconds=self.retention_seconds) + ).isoformat() + ) + existing = connection.execute( + """ + SELECT state_json, lifecycle FROM session_state + WHERE namespace = ? AND owner_id = ? AND session_id = ? + """, + (ns, owner, session_id), + ).fetchone() + if existing is not None and existing["lifecycle"] != "ACTIVE": + raise GDWLifecycleError("session is expired or tombstoned") + self._reserve_usage( + connection, + ns, + owner, + sessions=1 if existing is None else 0, + stored_bytes=_byte_len(state_text) + - _byte_len(existing["state_json"] if existing else None), + ) connection.execute( """ - INSERT INTO session_state(session_id, step, state_json, state_hash, updated_at) - VALUES (?, ?, ?, ?, ?) - ON CONFLICT(session_id) DO UPDATE SET + INSERT INTO session_state( + namespace, owner_id, session_id, step, state_json, state_hash, + lifecycle, created_at, updated_at, last_accessed_at, expires_at + ) VALUES (?, ?, ?, ?, ?, ?, 'ACTIVE', ?, ?, ?, ?) + ON CONFLICT(namespace, owner_id, session_id) DO UPDATE SET step = excluded.step, state_json = excluded.state_json, state_hash = excluded.state_hash, - updated_at = excluded.updated_at + updated_at = excluded.updated_at, + last_accessed_at = excluded.last_accessed_at, + expires_at = excluded.expires_at """, ( + ns, + owner, session_id, step, - json.dumps(state, sort_keys=True, separators=(",", ":")), + state_text, state_hash, - updated_at, + timestamp, + timestamp, + timestamp, + expiry, ), ) - @staticmethod def save_request( + self, connection: sqlite3.Connection, request_id: str, request_digest: str, @@ -552,26 +1158,56 @@ class GDWWorkspace: response: Dict[str, Any], response_hash: str, created_at: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + expires_at: Optional[str] = None, ) -> None: + ns, owner = self._identity(namespace, owner_id) + response_text = _json_text(response) + canonical_response_hash = hashlib.sha256( + response_text.encode("utf-8") + ).hexdigest() + if response_hash != canonical_response_hash: + raise GDWConfigurationError( + "request response hash does not match canonical response" + ) + timestamp = _text_time(created_at) + expiry = ( + expires_at + or ( + _normalise_time(timestamp) + timedelta(seconds=self.retention_seconds) + ).isoformat() + ) + self._reserve_usage( + connection, + ns, + owner, + requests=1, + stored_bytes=_byte_len(response_text), + ) connection.execute( """ INSERT INTO requests( - request_id, request_digest, session_id, response_json, - response_hash, created_at - ) VALUES (?, ?, ?, ?, ?, ?) + namespace, owner_id, request_id, request_digest, session_id, + response_json, response_hash, lifecycle, created_at, expires_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, 'ACTIVE', ?, ?) """, ( + ns, + owner, request_id, request_digest, session_id, - json.dumps(response, sort_keys=True, separators=(",", ":")), - response_hash, - created_at, + response_text, + canonical_response_hash, + timestamp, + expiry, ), ) - @staticmethod def save_receipt( + self, connection: sqlite3.Connection, receipt_hash: str, request_id: str, @@ -579,205 +1215,702 @@ class GDWWorkspace: step: int, receipt: Dict[str, Any], created_at: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, ) -> None: + ns, owner = self._identity(namespace, owner_id) + receipt_text = _json_text(receipt) + claimed_receipt_hash = str(receipt.get("receipt_hash") or "") + unsigned_receipt = dict(receipt) + unsigned_receipt.pop("receipt_hash", None) + canonical_receipt_hash = hashlib.sha256( + _json_text(unsigned_receipt).encode("utf-8") + ).hexdigest() + if ( + receipt_hash != claimed_receipt_hash + or receipt_hash != canonical_receipt_hash + ): + raise GDWConfigurationError( + "receipt hash does not match canonical receipt" + ) + request_anchor = self._request_anchor( + connection, ns, owner, request_id + ) + expected_identity = { + "namespace": ns, + "owner_id": owner, + "request_id": request_id, + "session_id": session_id, + "step": int(step), + "database_generation_id": self.database_generation_id, + } + for field, expected in expected_identity.items(): + if receipt.get(field) != expected: + raise GDWConfigurationError( + f"receipt {field} does not match persisted identity" + ) + if request_anchor["session_id"] != session_id: + raise GDWConfigurationError( + "receipt session does not match persisted request" + ) + if request_anchor["response"].get("receipt_hash") != receipt_hash: + raise GDWConfigurationError( + "receipt hash does not match persisted response" + ) + self._reserve_usage(connection, ns, owner, stored_bytes=_byte_len(receipt_text)) connection.execute( """ INSERT INTO receipts( - receipt_hash, request_id, session_id, step, receipt_json, created_at - ) VALUES (?, ?, ?, ?, ?, ?) + namespace, owner_id, receipt_hash, request_id, session_id, + step, receipt_json, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) """, ( - receipt_hash, + ns, + owner, + canonical_receipt_hash, request_id, session_id, step, - json.dumps(receipt, sort_keys=True, separators=(",", ":")), - created_at, + receipt_text, + _text_time(created_at), ), ) - @staticmethod def save_proof_outbox( + self, connection: sqlite3.Connection, proposal_id: str, payload: Dict[str, Any], payload_sha256: str, created_at: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, ) -> None: + ns, owner = self._identity(namespace, owner_id) + payload_text = _json_text(payload) + self._reserve_usage( + connection, + ns, + owner, + pending_effects=1, + stored_bytes=_byte_len(payload_text), + ) connection.execute( """ INSERT INTO proof_outbox( - proposal_id, payload_json, payload_sha256, status, created_at - ) VALUES (?, ?, ?, 'PENDING', ?) + namespace, owner_id, proposal_id, payload_json, + payload_sha256, status, created_at + ) VALUES (?, ?, ?, ?, ?, 'PENDING', ?) """, ( + ns, + owner, proposal_id, - json.dumps(payload, sort_keys=True, separators=(",", ":")), + payload_text, payload_sha256, - created_at, + _text_time(created_at), ), ) @staticmethod + def scoped_effect_key( + namespace: str, + owner_id: str, + request_id: str, + kind: str, + payload_sha256: str, + ) -> str: + if ( + len(payload_sha256) != 64 + or any(ch not in "0123456789abcdef" for ch in payload_sha256) + ): + raise GDWConfigurationError( + "effect payload digest must be a lowercase SHA-256 value" + ) + material = ( + f"{namespace}\x1f{owner_id}\x1f{request_id}\x1f" + f"{kind}\x1f{payload_sha256}" + ) + return hashlib.sha256(material.encode("utf-8")).hexdigest() + + @staticmethod + def _effect_payload_digest(kind: str, payload: Dict[str, Any]) -> str: + if kind == "proof_export": + claimed_digest = str(payload.get("payload_sha256") or "") + unsigned = dict(payload) + unsigned.pop("payload_sha256", None) + observed_digest = hashlib.sha256( + _json_text(unsigned).encode("utf-8") + ).hexdigest() + if claimed_digest != observed_digest: + raise GDWConfigurationError( + "proof payload digest does not match canonical payload" + ) + return claimed_digest + if kind == "receipt_projection": + return hashlib.sha256( + _json_text(payload).encode("utf-8") + ).hexdigest() + raise GDWConfigurationError("unsupported effect kind") + + @staticmethod + def artifact_binding_errors( + row: Dict[str, Any], + artifact: Any, + ) -> list[str]: + errors = [] + if not isinstance(artifact, dict): + return ["artifact_not_object"] + identity = str(row.get("intent_sha256") or "") + owner_id = str(row.get("owner_id") or "") + kind = str(row.get("kind") or "") + if artifact.get("artifact_identity") != identity: + errors.append("artifact_identity_mismatch") + if artifact.get("immutable") is not True: + errors.append("artifact_not_immutable") + owner_scope = hashlib.sha256(owner_id.encode("utf-8")).hexdigest()[:32] + if artifact.get("owner_scope") != owner_scope: + errors.append("artifact_owner_scope_mismatch") + configured_root = ( + os.environ.get("GDW_PROOF_DIR", "output/proofs") + if kind == "proof_export" + else os.environ.get( + "GDW_RECEIPT_PROJECTION_DIR", + "output/gdw/receipts", + ) + ) + if kind not in {"proof_export", "receipt_projection"}: + errors.append("artifact_kind_invalid") + return errors + expected_path = ( + Path(configured_root).resolve() + / owner_scope + / f"{identity}.json" + ) + try: + observed_path = Path(str(artifact.get("path") or "")).resolve() + except (OSError, RuntimeError, ValueError): + errors.append("artifact_path_invalid") + return errors + if observed_path != expected_path: + errors.append("artifact_path_mismatch") + if not observed_path.is_file(): + errors.append("artifact_missing") + return errors + try: + observed_sha256 = hashlib.sha256( + observed_path.read_bytes() + ).hexdigest() + except OSError: + errors.append("artifact_unreadable") + return errors + if artifact.get("sha256") != observed_sha256: + errors.append("artifact_digest_mismatch") + return errors + + def _request_anchor( + self, + connection: sqlite3.Connection, + namespace: str, + owner_id: str, + request_id: str, + ) -> Dict[str, Any]: + row = connection.execute( + """ + SELECT request_digest, session_id, response_json, response_hash, + lifecycle + FROM requests + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (namespace, owner_id, request_id), + ).fetchone() + if row is None or row["response_json"] is None: + raise GDWConfigurationError("effect request anchor is unavailable") + try: + response = json.loads(row["response_json"]) + except (TypeError, json.JSONDecodeError) as exc: + raise GDWConfigurationError( + "effect request anchor is not canonical JSON" + ) from exc + observed_hash = hashlib.sha256( + _json_text(response).encode("utf-8") + ).hexdigest() + if observed_hash != row["response_hash"]: + raise GDWConfigurationError("effect request response hash is invalid") + if response.get("request_digest") != row["request_digest"]: + raise GDWConfigurationError( + "effect request digest identity is invalid" + ) + if response.get("database_generation_id") != self.database_generation_id: + raise GDWConfigurationError( + "effect request database generation is invalid" + ) + try: + proposal_material = { + "schema": "szl.gdw.proposal-identity/v1", + "database_generation_id": self.database_generation_id, + "namespace": response["principal"]["namespace"], + "owner_id": response["principal"]["owner_id"], + "request_id": response["request_id"], + "request_digest": row["request_digest"], + "state_before_hash": response["state_before_hash"], + "governance_evidence_sha256": hashlib.sha256( + _json_text(response["audit"]["governance"]).encode("utf-8") + ).hexdigest(), + } + except (KeyError, TypeError) as exc: + raise GDWConfigurationError( + "effect request proposal identity is unavailable" + ) from exc + expected_proposal_id = hashlib.sha256( + _json_text(proposal_material).encode("utf-8") + ).hexdigest() + if response.get("proposal_id") != expected_proposal_id: + raise GDWConfigurationError( + "effect request proposal identity is invalid" + ) + return { + "request_digest": row["request_digest"], + "session_id": row["session_id"], + "response_hash": row["response_hash"], + "response": response, + "lifecycle": row["lifecycle"], + } + + def _receipt_anchor( + self, + connection: sqlite3.Connection, + namespace: str, + owner_id: str, + request_id: str, + ) -> Dict[str, Any]: + row = connection.execute( + """ + SELECT receipt_hash, request_id, session_id, step, receipt_json + FROM receipts + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (namespace, owner_id, request_id), + ).fetchone() + if row is None or row["receipt_json"] is None: + raise GDWConfigurationError("effect receipt anchor is unavailable") + try: + receipt = json.loads(row["receipt_json"]) + except (TypeError, json.JSONDecodeError) as exc: + raise GDWConfigurationError( + "effect receipt anchor is not canonical JSON" + ) from exc + claimed = str(receipt.get("receipt_hash") or "") + unsigned = dict(receipt) + unsigned.pop("receipt_hash", None) + observed = hashlib.sha256( + _json_text(unsigned).encode("utf-8") + ).hexdigest() + if claimed != row["receipt_hash"] or observed != row["receipt_hash"]: + raise GDWConfigurationError("effect receipt anchor hash is invalid") + expected_identity = { + "namespace": namespace, + "owner_id": owner_id, + "request_id": row["request_id"], + "session_id": row["session_id"], + "step": int(row["step"]), + "database_generation_id": self.database_generation_id, + } + for field, expected in expected_identity.items(): + if receipt.get(field) != expected: + raise GDWConfigurationError( + f"effect receipt {field} identity is invalid" + ) + return {"receipt_hash": row["receipt_hash"], "receipt": receipt} + + @staticmethod + def _expected_proof_payload( + request_anchor: Dict[str, Any], + ) -> Dict[str, Any]: + from gdw_proofs import build_proof_payload + + response = request_anchor["response"] + try: + return build_proof_payload( + proposal_id=response["proposal_id"], + request_id=response["request_id"], + request_digest=request_anchor["request_digest"], + namespace=response["principal"]["namespace"], + owner_id=response["principal"]["owner_id"], + database_generation_id=response["database_generation_id"], + step=int(response["step"]), + before_hash=response["state_before_hash"], + after_hash=response["state_hash"], + decision=response["decision"], + scheduler_mode=response["scheduler_mode"], + receipt_hash=response.get("receipt_hash") or "", + dry_run=bool(response["dry_run"]), + governance=response["audit"]["governance"], + ) + except (KeyError, TypeError, ValueError) as exc: + raise GDWConfigurationError( + "persisted request cannot reconstruct proof intent" + ) from exc + + def _canonical_effect_intent( + self, + request_anchor: Dict[str, Any], + *, + namespace: str, + owner_id: str, + request_id: str, + kind: str, + payload_sha256: str, + receipt_hash: Optional[str], + ) -> str: + material = { + "schema": "szl.gdw.effect-intent/v1", + "database_generation_id": self.database_generation_id, + "namespace": namespace, + "owner_id": owner_id, + "request_id": request_id, + "request_digest": request_anchor["request_digest"], + "session_id": request_anchor["session_id"], + "response_hash": request_anchor["response_hash"], + "kind": kind, + "payload_sha256": payload_sha256, + "receipt_hash": receipt_hash, + } + return hashlib.sha256(_json_text(material).encode("utf-8")).hexdigest() + + @classmethod + def effect_binding_errors(cls, row: Dict[str, Any]) -> list[str]: + """Validate row-local bindings; persisted anchors are checked separately.""" + + errors = [] + payload = row.get("payload") + if not isinstance(payload, dict): + return ["payload_not_object"] + kind = str(row.get("kind") or "") + request_id = str(row.get("request_id") or "") + namespace = str(row.get("namespace") or "") + owner_id = str(row.get("owner_id") or "") + if kind == "proof_export": + claimed_digest = str(payload.get("payload_sha256") or "") + unsigned = dict(payload) + unsigned.pop("payload_sha256", None) + observed_digest = hashlib.sha256( + _json_text(unsigned).encode("utf-8") + ).hexdigest() + if claimed_digest != observed_digest: + errors.append("proof_payload_digest_invalid") + payload_digest = claimed_digest + principal = payload.get("governance", {}).get("principal", {}) + elif kind == "receipt_projection": + payload_digest = hashlib.sha256( + _json_text(payload).encode("utf-8") + ).hexdigest() + principal = payload + else: + return ["unsupported_effect_kind"] + if row.get("payload_sha256") != payload_digest: + errors.append("row_payload_digest_mismatch") + if payload.get("request_id") != request_id: + errors.append("request_identity_mismatch") + if payload.get("database_generation_id") != row.get( + "database_generation_id" + ): + errors.append("database_generation_payload_mismatch") + if not isinstance(principal, dict) or ( + principal.get("namespace") != namespace + or principal.get("owner_id") != owner_id + ): + errors.append("principal_identity_mismatch") + try: + expected_key = cls.scoped_effect_key( + namespace, + owner_id, + request_id, + kind, + payload_digest, + ) + except GDWConfigurationError: + errors.append("payload_digest_invalid") + else: + if row.get("idempotency_key") != expected_key: + errors.append("idempotency_key_mismatch") + return errors + + def _connection_effect_binding_errors( + self, + connection: sqlite3.Connection, + row: Dict[str, Any], + ) -> list[str]: + errors = self.effect_binding_errors(row) + try: + request_anchor = self._request_anchor( + connection, + str(row.get("namespace") or ""), + str(row.get("owner_id") or ""), + str(row.get("request_id") or ""), + ) + except GDWConfigurationError as exc: + return sorted(set(errors + [str(exc)])) + payload = row.get("payload") + kind = str(row.get("kind") or "") + receipt_hash = None + try: + if kind == "receipt_projection": + receipt_anchor = self._receipt_anchor( + connection, + row["namespace"], + row["owner_id"], + row["request_id"], + ) + receipt_hash = receipt_anchor["receipt_hash"] + if payload != receipt_anchor["receipt"]: + errors.append("receipt_payload_anchor_mismatch") + elif kind == "proof_export": + expected_proof = self._expected_proof_payload(request_anchor) + if payload != expected_proof: + errors.append("proof_payload_anchor_mismatch") + claimed_receipt_hash = str( + expected_proof.get("delta_update_receipt_hash") or "" + ) + if claimed_receipt_hash: + receipt_anchor = self._receipt_anchor( + connection, + row["namespace"], + row["owner_id"], + row["request_id"], + ) + receipt_hash = receipt_anchor["receipt_hash"] + if receipt_hash != claimed_receipt_hash: + errors.append("proof_receipt_anchor_mismatch") + else: + errors.append("unsupported_effect_kind") + except GDWConfigurationError as exc: + errors.append(str(exc)) + if row.get("receipt_hash") != receipt_hash: + errors.append("receipt_hash_anchor_mismatch") + if row.get("database_generation_id") != self.database_generation_id: + errors.append("database_generation_mismatch") + expected_intent = self._canonical_effect_intent( + request_anchor, + namespace=row["namespace"], + owner_id=row["owner_id"], + request_id=row["request_id"], + kind=kind, + payload_sha256=str(row.get("payload_sha256") or ""), + receipt_hash=receipt_hash, + ) + if row.get("intent_sha256") != expected_intent: + errors.append("effect_intent_mismatch") + return sorted(set(errors)) + + def effect_binding_errors_for_row(self, row: Dict[str, Any]) -> list[str]: + connection = self._connect() + try: + return self._connection_effect_binding_errors(connection, row) + finally: + connection.close() + def save_effect_outbox( + self, connection: sqlite3.Connection, request_id: str, kind: str, - generation_id: str, - owner_id: str, - canonical_identity: str, payload: Dict[str, Any], payload_sha256: str, - idempotency_key: str, + idempotency_key: Optional[str], created_at: str, - ) -> None: - if kind not in _EFFECT_KINDS: - raise ValueError("unsupported effect kind") - canonical_payload = _canonical_json(payload) - actual_payload_sha256 = hashlib.sha256( - canonical_payload.encode("utf-8") - ).hexdigest() - if payload_sha256 != actual_payload_sha256: - raise ValueError("effect payload_sha256 does not match payload") - expected_identity = ( - payload.get("receipt_hash") - if kind == "receipt_projection" - else payload.get("payload_sha256") + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + max_attempts: Optional[int] = None, + ) -> str: + ns, owner = self._identity(namespace, owner_id) + payload_text = _json_text(payload) + canonical_payload_sha256 = self._effect_payload_digest(kind, payload) + if payload_sha256 != canonical_payload_sha256: + raise GDWConfigurationError( + "effect payload digest does not match canonical payload" + ) + request_anchor = self._request_anchor( + connection, ns, owner, request_id ) - if canonical_identity != expected_identity: - raise ValueError("effect canonical identity does not match payload") - request = connection.execute( - "SELECT request_digest FROM requests WHERE request_id = ?", - (request_id,), - ).fetchone() - if request is None: - raise ValueError("effect request is missing") - if ( - payload.get("request_id") != request_id - or payload.get("request_digest") != request["request_digest"] - or payload.get("owner_id") != owner_id - or payload.get("generation_id") != generation_id - ): - raise ValueError("effect payload identity is not request-bound") - if GDWWorkspace.object_owner( - connection, "request", request_id - ) != owner_id: - raise ValueError("effect owner is not request-bound") - expected_key = _effect_identity_key( - generation_id=generation_id, - owner_id=owner_id, + receipt_hash = None + if kind == "receipt_projection": + receipt_anchor = self._receipt_anchor( + connection, ns, owner, request_id + ) + receipt_hash = receipt_anchor["receipt_hash"] + if payload != receipt_anchor["receipt"]: + raise GDWConfigurationError( + "receipt projection differs from persisted receipt" + ) + elif kind == "proof_export": + expected_proof = self._expected_proof_payload(request_anchor) + if payload != expected_proof: + raise GDWConfigurationError( + "proof export differs from persisted request intent" + ) + claimed_receipt_hash = str( + expected_proof.get("delta_update_receipt_hash") or "" + ) + if claimed_receipt_hash: + receipt_anchor = self._receipt_anchor( + connection, ns, owner, request_id + ) + receipt_hash = receipt_anchor["receipt_hash"] + if receipt_hash != claimed_receipt_hash: + raise GDWConfigurationError( + "proof receipt differs from persisted receipt" + ) + else: + raise GDWConfigurationError("unsupported effect kind") + intent_sha256 = self._canonical_effect_intent( + request_anchor, + namespace=ns, + owner_id=owner, request_id=request_id, - request_digest=request["request_digest"], kind=kind, - canonical_identity=canonical_identity, - payload_sha256=payload_sha256, + payload_sha256=canonical_payload_sha256, + receipt_hash=receipt_hash, ) - if idempotency_key != expected_key: - raise ValueError("effect idempotency identity mismatch") - connection.execute( - """ - INSERT INTO evidence_intents( - canonical_identity, request_id, kind, generation_id, owner_id, - payload_json, payload_sha256, created_at - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?) - """, - ( - canonical_identity, - request_id, - kind, - generation_id, - owner_id, - canonical_payload, - payload_sha256, - created_at, - ), + canonical_key = self.scoped_effect_key( + ns, + owner, + request_id, + kind, + canonical_payload_sha256, + ) + if idempotency_key and self.production and idempotency_key != canonical_key: + raise GDWConfigurationError( + "effect idempotency key is not bound to identity and payload" + ) + key = canonical_key + bounded_attempts = max_attempts or self.effect_max_attempts + if bounded_attempts < 1: + raise GDWConfigurationError("max_attempts must be positive") + timestamp = _text_time(created_at) + self._reserve_usage( + connection, + ns, + owner, + pending_effects=1, + stored_bytes=_byte_len(payload_text), ) connection.execute( """ INSERT INTO effect_outbox( - idempotency_key, request_id, kind, generation_id, owner_id, - canonical_identity, payload_json, payload_sha256, status, - created_at - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'PENDING', ?) + namespace, owner_id, idempotency_key, database_generation_id, + request_id, kind, receipt_hash, payload_json, payload_sha256, + intent_sha256, status, attempts, max_attempts, + next_attempt_at, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'PENDING', 0, ?, ?, ?) """, ( - idempotency_key, + ns, + owner, + key, + self.database_generation_id, request_id, kind, - generation_id, - owner_id, - canonical_identity, - canonical_payload, - payload_sha256, - created_at, + receipt_hash, + payload_text, + canonical_payload_sha256, + intent_sha256, + bounded_attempts, + timestamp, + timestamp, ), ) + return key def claim_effects( self, worker_id: str, limit: int = 100, lease_seconds: int = 300, - max_attempts: Optional[int] = None, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + now: Optional[Any] = None, ) -> list: if not worker_id or len(worker_id) > 128: raise ValueError("worker_id must contain 1-128 characters") - bounded = max(1, min(int(limit), 10000)) - lease = max(1, min(int(lease_seconds), 3600)) - attempt_ceiling = max( - 1, - min( - int( - max_attempts - if max_attempts is not None - else os.environ.get("GDW_MAX_EFFECT_ATTEMPTS", "20") - ), - 100, - ), - ) - now = datetime.now(timezone.utc) - now_text = now.isoformat() - lease_until = (now + timedelta(seconds=lease)).isoformat() + ns, owner = self._identity(namespace, owner_id) + bounded = max(1, min(int(limit), 10_000)) + lease = max(1, min(int(lease_seconds), 3_600)) + current = _normalise_time(now) + now_text = current.isoformat() + lease_until = (current + timedelta(seconds=lease)).isoformat() claimed = [] with self.transaction() as connection: + expired_terminal = int( + connection.execute( + """ + SELECT COUNT(*) FROM effect_outbox + WHERE namespace = ? AND owner_id = ? AND status = 'CLAIMED' + AND lease_until <= ? AND attempts >= max_attempts + """, + (ns, owner, now_text), + ).fetchone()[0] + ) + if expired_terminal: + connection.execute( + """ + UPDATE effect_outbox + SET status = 'DEAD_LETTER', lease_owner = NULL, + lease_until = NULL, last_error = 'LEASE_EXPIRED_FINAL_ATTEMPT' + WHERE namespace = ? AND owner_id = ? AND status = 'CLAIMED' + AND lease_until <= ? AND attempts >= max_attempts + """, + (ns, owner, now_text), + ) + self._reserve_usage( + connection, + ns, + owner, + pending_effects=-expired_terminal, + ) connection.execute( """ UPDATE effect_outbox - SET status = 'PENDING', lease_owner = NULL, lease_until = NULL, - claim_token = NULL - WHERE status = 'CLAIMED' AND lease_until <= ? + SET status = 'PENDING', lease_owner = NULL, lease_until = NULL + WHERE namespace = ? AND owner_id = ? AND status = 'CLAIMED' + AND lease_until <= ? AND attempts < max_attempts """, - (now_text,), + (ns, owner, now_text), ) rows = connection.execute( """ - SELECT idempotency_key, request_id, kind, generation_id, - owner_id, canonical_identity, payload_json, - payload_sha256, attempts + SELECT idempotency_key, database_generation_id, request_id, + kind, receipt_hash, payload_json, payload_sha256, + intent_sha256, attempts, max_attempts, claim_generation FROM effect_outbox - WHERE status = 'PENDING' AND attempts < ? - AND generation_id IS NOT NULL - AND owner_id IS NOT NULL - AND canonical_identity IS NOT NULL - ORDER BY created_at, idempotency_key + WHERE namespace = ? AND owner_id = ? AND status = 'PENDING' + AND next_attempt_at <= ? AND attempts < max_attempts + ORDER BY next_attempt_at, created_at, idempotency_key LIMIT ? """, - (attempt_ceiling, bounded), + (ns, owner, now_text, bounded), ).fetchall() for row in rows: - claim_token = secrets.token_hex(16) updated = connection.execute( """ UPDATE effect_outbox SET status = 'CLAIMED', lease_owner = ?, lease_until = ?, - claim_token = ?, attempts = attempts + 1, + attempts = attempts + 1, + claim_generation = claim_generation + 1, last_error = NULL - WHERE idempotency_key = ? AND status = 'PENDING' + WHERE namespace = ? AND owner_id = ? + AND idempotency_key = ? AND status = 'PENDING' """, ( worker_id, lease_until, - claim_token, + ns, + owner, row["idempotency_key"], ), ) @@ -785,148 +1918,76 @@ class GDWWorkspace: continue claimed.append( { + "namespace": ns, + "owner_id": owner, "idempotency_key": row["idempotency_key"], + "database_generation_id": row["database_generation_id"], "request_id": row["request_id"], - "kind": row["kind"], - "generation_id": row["generation_id"], - "owner_id": row["owner_id"], - "canonical_identity": row["canonical_identity"], - "payload": json.loads(row["payload_json"]), - "payload_sha256": row["payload_sha256"], - "attempt": int(row["attempts"]) + 1, - "lease_owner": worker_id, - "lease_until": lease_until, - "claim_token": claim_token, - } - ) - return claimed - - @staticmethod - def _validate_effect_binding( - connection: sqlite3.Connection, - row: sqlite3.Row, - ) -> None: - payload = json.loads(row["payload_json"]) - payload_json = _canonical_json(payload) - payload_sha256 = hashlib.sha256( - payload_json.encode("utf-8") - ).hexdigest() - if payload_sha256 != row["payload_sha256"]: - raise ValueError("effect outbox payload digest mismatch") - expected_identity = ( - payload.get("receipt_hash") - if row["kind"] == "receipt_projection" - else payload.get("payload_sha256") - ) - if expected_identity != row["canonical_identity"]: - raise ValueError("effect outbox canonical identity mismatch") - - intent = connection.execute( - """ - SELECT request_id, kind, generation_id, owner_id, payload_json, - payload_sha256 - FROM evidence_intents WHERE canonical_identity = ? - """, - (row["canonical_identity"],), - ).fetchone() - if intent is None: - raise ValueError("canonical evidence intent is missing") - expected = ( - row["request_id"], - row["kind"], - row["generation_id"], - row["owner_id"], - payload_json, - row["payload_sha256"], - ) - actual = ( - intent["request_id"], - intent["kind"], - intent["generation_id"], - intent["owner_id"], - intent["payload_json"], - intent["payload_sha256"], - ) - if actual != expected: - raise ValueError("effect outbox diverges from canonical intent") - - owner = GDWWorkspace.object_owner( - connection, "request", row["request_id"] - ) - if owner != row["owner_id"]: - raise ValueError("effect owner diverges from request owner") - request = connection.execute( - "SELECT request_digest, response_json FROM requests " - "WHERE request_id = ?", - (row["request_id"],), - ).fetchone() - if request is None: - raise ValueError("effect request is missing") - response = json.loads(request["response_json"]) - if ( - payload.get("request_id") != row["request_id"] - or payload.get("request_digest") != request["request_digest"] - or payload.get("owner_id") != row["owner_id"] - or payload.get("generation_id") != row["generation_id"] - ): - raise ValueError("effect payload identity diverges from request") - expected_key = _effect_identity_key( - generation_id=row["generation_id"], - owner_id=row["owner_id"], - request_id=row["request_id"], - request_digest=request["request_digest"], - kind=row["kind"], - canonical_identity=row["canonical_identity"], - payload_sha256=row["payload_sha256"], - ) - if row["idempotency_key"] != expected_key: - raise ValueError("effect idempotency identity mismatch") - generation = connection.execute( - "SELECT value FROM workspace_meta WHERE key = 'generation_id'" - ).fetchone() - if generation is None or generation["value"] != row["generation_id"]: - raise ValueError("effect generation identity mismatch") - - if row["kind"] == "receipt_projection": - receipt = connection.execute( - "SELECT receipt_json FROM receipts " - "WHERE receipt_hash = ? AND request_id = ?", - (row["canonical_identity"], row["request_id"]), - ).fetchone() - if receipt is None or receipt["receipt_json"] != payload_json: - raise ValueError("receipt projection diverges from receipt ledger") - if response.get("receipt_hash") != row["canonical_identity"]: - raise ValueError("receipt projection diverges from response") - else: - unsigned_payload = dict(payload) - embedded_digest = unsigned_payload.pop("payload_sha256", None) - if ( - embedded_digest != row["canonical_identity"] - or _sha256_json(unsigned_payload) != embedded_digest - or response.get("proof", {}).get("canonical_identity") - != row["canonical_identity"] - or response.get("proof", {}).get("idempotency_key") - != row["idempotency_key"] - ): - raise ValueError("proof export diverges from canonical response") + "kind": row["kind"], + "receipt_hash": row["receipt_hash"], + "payload": json.loads(row["payload_json"]), + "payload_sha256": row["payload_sha256"], + "intent_sha256": row["intent_sha256"], + "attempt": int(row["attempts"]) + 1, + "max_attempts": int(row["max_attempts"]), + "lease_owner": worker_id, + "lease_until": lease_until, + "claim_generation": int(row["claim_generation"]) + 1, + } + ) + return claimed - def validate_claimed_effect(self, claim: Dict[str, Any]) -> None: + def assert_effect_claim( + self, + idempotency_key: str, + worker_id: str, + claim_generation: int, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + now: Optional[Any] = None, + ) -> Dict[str, Any]: + ns, owner = self._identity(namespace, owner_id) + current = _text_time(now) connection = self._connect() try: row = connection.execute( - "SELECT * FROM effect_outbox WHERE idempotency_key = ?", - (claim["idempotency_key"],), + """ + SELECT namespace, owner_id, idempotency_key, + database_generation_id, request_id, kind, receipt_hash, + payload_json, payload_sha256, intent_sha256, + claim_generation, lease_owner, lease_until + FROM effect_outbox + WHERE namespace = ? AND owner_id = ? AND idempotency_key = ? + AND status = 'CLAIMED' AND lease_owner = ? + AND claim_generation = ? AND lease_until > ? + """, + ( + ns, + owner, + idempotency_key, + worker_id, + int(claim_generation), + current, + ), ).fetchone() if row is None: - raise ValueError("effect outbox row is missing") - if ( - row["status"] != "CLAIMED" - or row["lease_owner"] != claim["lease_owner"] - or row["claim_token"] != claim["claim_token"] - or row["lease_until"] <= datetime.now(timezone.utc).isoformat() - ): - raise RuntimeError("effect claim is stale or fenced") - self._validate_effect_binding(connection, row) + raise RuntimeError( + "effect claim is absent, expired, or owned elsewhere" + ) + candidate = dict(row) + try: + candidate["payload"] = json.loads(row["payload_json"]) + except (TypeError, json.JSONDecodeError) as exc: + raise RuntimeError("effect claim payload is invalid") from exc + errors = self._connection_effect_binding_errors( + connection, candidate + ) + if errors: + raise RuntimeError( + "effect claim binding is invalid: " + ",".join(errors) + ) + return candidate finally: connection.close() @@ -934,109 +1995,185 @@ class GDWWorkspace: self, idempotency_key: str, worker_id: str, - claim_token: str, + claim_generation: int, artifact: Dict[str, Any], exported_at: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, ) -> None: - if artifact.get("artifact_identity") != idempotency_key: - raise ValueError("artifact identity does not match effect key") - if artifact.get("immutable") is not True: - raise ValueError("artifact is not marked immutable") - path = Path(str(artifact.get("path") or "")).resolve() - if not path.is_file(): - raise ValueError("exported artifact is missing") - artifact_sha256 = hashlib.sha256(path.read_bytes()).hexdigest() - if artifact_sha256 != artifact.get("sha256"): - raise ValueError("exported artifact digest mismatch") - now_text = datetime.now(timezone.utc).isoformat() + ns, owner = self._identity(namespace, owner_id) + completed_at = _text_time(exported_at) + observed_now = _text_time() with self.transaction() as connection: row = connection.execute( - "SELECT * FROM effect_outbox WHERE idempotency_key = ?", - (idempotency_key,), + """ + SELECT namespace, owner_id, idempotency_key, + database_generation_id, request_id, kind, receipt_hash, + payload_json, payload_sha256, intent_sha256, + status, artifact_json + FROM effect_outbox + WHERE namespace = ? AND owner_id = ? AND idempotency_key = ? + AND status = 'CLAIMED' AND lease_owner = ? + AND claim_generation = ? AND lease_until > ? + """, + ( + ns, + owner, + idempotency_key, + worker_id, + int(claim_generation), + observed_now, + ), ).fetchone() if row is None: - raise RuntimeError("effect claim is absent") - self._validate_effect_binding(connection, row) - owner_scope = hashlib.sha256( - row["owner_id"].encode("utf-8") - ).hexdigest()[:32] - configured_root = ( - os.environ.get("GDW_PROOF_DIR", "output/proofs") - if row["kind"] == "proof_export" - else os.environ.get( - "GDW_RECEIPT_PROJECTION_DIR", "output/gdw/receipts" + raise RuntimeError( + "effect claim is absent, expired, or owned elsewhere" ) + candidate = dict(row) + try: + candidate["payload"] = json.loads(row["payload_json"]) + except (TypeError, json.JSONDecodeError) as exc: + raise GDWConfigurationError( + "effect claim payload is invalid" + ) from exc + binding_errors = self._connection_effect_binding_errors( + connection, + candidate, ) - expected_parent = (Path(configured_root).resolve() / owner_scope) - if ( - artifact.get("owner_scope") != owner_scope - or path.parent != expected_parent - or path.name != f"{idempotency_key}.json" - ): - raise ValueError("artifact path is not owner- and effect-bound") + artifact_errors = self.artifact_binding_errors( + candidate, + artifact, + ) + errors = sorted(set(binding_errors + artifact_errors)) + if errors: + raise GDWConfigurationError( + "effect completion is invalid: " + ",".join(errors) + ) + artifact_text = _json_text(artifact) updated = connection.execute( """ UPDATE effect_outbox SET status = 'EXPORTED', artifact_json = ?, exported_at = ?, - lease_owner = NULL, lease_until = NULL, claim_token = NULL, - last_error = NULL - WHERE idempotency_key = ? AND status = 'CLAIMED' - AND lease_owner = ? AND claim_token = ? - AND lease_until > ? + lease_owner = NULL, lease_until = NULL, last_error = NULL + WHERE namespace = ? AND owner_id = ? AND idempotency_key = ? + AND status = 'CLAIMED' AND lease_owner = ? + AND claim_generation = ? AND lease_until > ? """, ( - json.dumps(artifact, sort_keys=True, separators=(",", ":")), - exported_at, + artifact_text, + completed_at, + ns, + owner, idempotency_key, worker_id, - claim_token, - now_text, + int(claim_generation), + observed_now, ), ) if updated.rowcount != 1: - raise RuntimeError("effect claim is absent, expired, or owned elsewhere") + raise RuntimeError( + "effect claim is absent, expired, or owned elsewhere" + ) + self._reserve_usage( + connection, + ns, + owner, + pending_effects=-1, + stored_bytes=_byte_len(artifact_text), + ) def release_effect( self, idempotency_key: str, worker_id: str, - claim_token: str, + claim_generation: int, error: str, - ) -> bool: + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + now: Optional[Any] = None, + ) -> str: + ns, owner = self._identity(namespace, owner_id) + current = _normalise_time(now) with self.transaction() as connection: - updated = connection.execute( + row = connection.execute( """ - UPDATE effect_outbox SET status = 'PENDING', - lease_owner = NULL, lease_until = NULL, claim_token = NULL, - last_error = ? - WHERE idempotency_key = ? AND status = 'CLAIMED' - AND lease_owner = ? AND claim_token = ? + SELECT attempts, max_attempts + FROM effect_outbox + WHERE namespace = ? AND owner_id = ? AND idempotency_key = ? + AND status = 'CLAIMED' AND lease_owner = ? + AND claim_generation = ? AND lease_until > ? + """, + ( + ns, + owner, + idempotency_key, + worker_id, + int(claim_generation), + current.isoformat(), + ), + ).fetchone() + if row is None: + raise RuntimeError( + "effect claim is absent, expired, or owned elsewhere" + ) + attempts = int(row["attempts"]) + terminal = attempts >= int(row["max_attempts"]) + status = "DEAD_LETTER" if terminal else "PENDING" + delay = self.effect_backoff_seconds * (2 ** max(0, attempts - 1)) + delay = min(delay, 24 * 60 * 60) + next_attempt = (current + timedelta(seconds=delay)).isoformat() + connection.execute( + """ + UPDATE effect_outbox + SET status = ?, lease_owner = NULL, lease_until = NULL, + last_error = ?, next_attempt_at = ? + WHERE namespace = ? AND owner_id = ? AND idempotency_key = ? + AND status = 'CLAIMED' AND lease_owner = ? + AND claim_generation = ? AND lease_until > ? """, ( + status, str(error)[:1024], + next_attempt, + ns, + owner, idempotency_key, worker_id, - claim_token, + int(claim_generation), + current.isoformat(), ), ) - return updated.rowcount == 1 + if terminal: + self._reserve_usage(connection, ns, owner, pending_effects=-1) + return status - def pending_proofs(self, limit: int = 100) -> list: - bounded = max(1, min(int(limit), 10000)) + def pending_proofs( + self, + limit: int = 100, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + ) -> list: + ns, owner = self._identity(namespace, owner_id) + bounded = max(1, min(int(limit), 10_000)) connection = self._connect() try: rows = connection.execute( """ SELECT proposal_id, payload_json, payload_sha256 FROM proof_outbox - WHERE status = 'PENDING' + WHERE namespace = ? AND owner_id = ? AND status = 'PENDING' ORDER BY created_at, proposal_id LIMIT ? """, - (bounded,), + (ns, owner, bounded), ).fetchall() return [ { + "namespace": ns, + "owner_id": owner, "proposal_id": row["proposal_id"], "payload": json.loads(row["payload_json"]), "payload_sha256": row["payload_sha256"], @@ -1051,241 +2188,513 @@ class GDWWorkspace: proposal_id: str, artifact: Dict[str, Any], exported_at: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, ) -> None: + ns, owner = self._identity(namespace, owner_id) + artifact_text = _json_text(artifact) with self.transaction() as connection: - connection.execute( + updated = connection.execute( """ UPDATE proof_outbox SET status = 'EXPORTED', artifact_json = ?, exported_at = ? - WHERE proposal_id = ? AND status = 'PENDING' + WHERE namespace = ? AND owner_id = ? AND proposal_id = ? + AND status = 'PENDING' """, ( - json.dumps(artifact, sort_keys=True, separators=(",", ":")), - exported_at, + artifact_text, + _text_time(exported_at), + ns, + owner, proposal_id, ), ) + if updated.rowcount != 1: + raise RuntimeError("proof is absent, exported, or owned elsewhere") + self._reserve_usage( + connection, + ns, + owner, + pending_effects=-1, + stored_bytes=_byte_len(artifact_text), + ) def read_session( self, session_id: str, - owner_id: str, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, ) -> Optional[Dict[str, Any]]: connection = self._connect() try: - row = connection.execute( - "SELECT 1 FROM session_state WHERE session_id = ?", - (session_id,), - ).fetchone() - if row is None: - return None - self.require_object_owner( - connection, "session", session_id, owner_id + return self.session_state( + connection, + session_id, + namespace=namespace, + owner_id=owner_id, ) - return self.session_state(connection, session_id) finally: connection.close() - def integrity(self) -> Dict[str, Any]: + def pending_effect_identities(self) -> list[Tuple[str, str]]: + """Return principals with pending work for the trusted drain supervisor.""" + connection = self._connect() try: - check = connection.execute("PRAGMA integrity_check").fetchone()[0] - counts = {} - for table in ( - "session_state", - "requests", - "receipts", - "proof_outbox", - "effect_outbox", - "object_owners", - "evidence_intents", - ): - counts[table] = int( - connection.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0] - ) - pending_proofs = int( + rows = connection.execute( + """ + SELECT namespace, owner_id FROM effect_outbox + WHERE status IN ('PENDING', 'CLAIMED') + UNION + SELECT namespace, owner_id FROM proof_outbox + WHERE status = 'PENDING' + ORDER BY namespace, owner_id + """ + ).fetchall() + return [(row["namespace"], row["owner_id"]) for row in rows] + finally: + connection.close() + + @staticmethod + def _reconcile_usage(connection: sqlite3.Connection) -> None: + timestamp = _text_time() + identities = connection.execute( + """ + SELECT namespace, owner_id FROM session_state + UNION SELECT namespace, owner_id FROM requests + UNION SELECT namespace, owner_id FROM receipts + UNION SELECT namespace, owner_id FROM proof_outbox + UNION SELECT namespace, owner_id FROM effect_outbox + """ + ).fetchall() + connection.execute("DELETE FROM usage") + for identity in identities: + ns, owner = identity["namespace"], identity["owner_id"] + sessions = int( connection.execute( - "SELECT COUNT(*) FROM proof_outbox WHERE status = 'PENDING'" + """ + SELECT COUNT(*) FROM session_state + WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE' + """, + (ns, owner), ).fetchone()[0] ) - pending_effects = int( + requests = int( connection.execute( - "SELECT COUNT(*) FROM effect_outbox " - "WHERE status IN ('PENDING', 'CLAIMED')" + """ + SELECT COUNT(*) FROM requests + WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE' + """, + (ns, owner), ).fetchone()[0] ) - claimed_effects = int( + pending = int( connection.execute( - "SELECT COUNT(*) FROM effect_outbox WHERE status = 'CLAIMED'" + """ + SELECT + (SELECT COUNT(*) FROM effect_outbox + WHERE namespace = ? AND owner_id = ? + AND status IN ('PENDING', 'CLAIMED')) + + (SELECT COUNT(*) FROM proof_outbox + WHERE namespace = ? AND owner_id = ? AND status = 'PENDING') + """, + (ns, owner, ns, owner), ).fetchone()[0] ) - orphan_receipts = int( + stored = int( connection.execute( """ - SELECT COUNT(*) FROM receipts r - LEFT JOIN requests q ON q.request_id = r.request_id - WHERE q.request_id IS NULL - """ + SELECT + COALESCE((SELECT SUM(LENGTH(CAST(state_json AS BLOB))) + FROM session_state + WHERE namespace = ? AND owner_id = ?), 0) + + COALESCE((SELECT SUM(LENGTH(CAST(response_json AS BLOB))) + FROM requests + WHERE namespace = ? AND owner_id = ?), 0) + + COALESCE((SELECT SUM(LENGTH(CAST(receipt_json AS BLOB))) + FROM receipts + WHERE namespace = ? AND owner_id = ?), 0) + + COALESCE((SELECT SUM( + LENGTH(CAST(payload_json AS BLOB)) + + COALESCE(LENGTH(CAST(artifact_json AS BLOB)), 0)) + FROM proof_outbox + WHERE namespace = ? AND owner_id = ?), 0) + + COALESCE((SELECT SUM( + LENGTH(CAST(payload_json AS BLOB)) + + COALESCE(LENGTH(CAST(artifact_json AS BLOB)), 0)) + FROM effect_outbox + WHERE namespace = ? AND owner_id = ?), 0) + """, + (ns, owner, ns, owner, ns, owner, ns, owner, ns, owner), ).fetchone()[0] ) - violations = { - "orphan_receipts": orphan_receipts, - "unowned_sessions": int( - connection.execute( - """ - SELECT COUNT(*) FROM session_state s - LEFT JOIN object_owners o - ON o.object_type = 'session' - AND o.object_id = s.session_id - WHERE o.object_id IS NULL - """ - ).fetchone()[0] - ), - "unowned_requests": int( - connection.execute( - """ - SELECT COUNT(*) FROM requests r - LEFT JOIN object_owners o - ON o.object_type = 'request' - AND o.object_id = r.request_id - WHERE o.object_id IS NULL - """ - ).fetchone()[0] - ), - "legacy_pending_proofs": pending_proofs, - "state_digest_mismatches": 0, - "response_digest_mismatches": 0, - "receipt_digest_mismatches": 0, - "intent_digest_mismatches": 0, - "effect_binding_mismatches": 0, - "artifact_mismatches": 0, - "dead_effects": 0, + connection.execute( + """ + INSERT INTO usage( + namespace, owner_id, active_sessions, active_requests, + pending_effects, stored_bytes, updated_at + ) VALUES (?, ?, ?, ?, ?, ?, ?) + """, + (ns, owner, sessions, requests, pending, stored, timestamp), + ) + + def reconcile_usage(self) -> Dict[str, int]: + with self.transaction() as connection: + self._reconcile_usage(connection) + row = self._usage_row(connection, self.namespace, self.owner_id) + return { + "active_sessions": int(row["active_sessions"]), + "active_requests": int(row["active_requests"]), + "pending_effects": int(row["pending_effects"]), + "stored_bytes": int(row["stored_bytes"]), } - for row in connection.execute( - "SELECT state_json, state_hash FROM session_state" - ): - try: - if _sha256_json(json.loads(row["state_json"])) != row["state_hash"]: - violations["state_digest_mismatches"] += 1 - except Exception: - violations["state_digest_mismatches"] += 1 + def collect_garbage( + self, + *, + now: Optional[Any] = None, + limit: int = 100, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + ) -> Dict[str, int]: + """Compact expired objects; unexported effects are never compacted or deleted.""" - for row in connection.execute( - "SELECT response_json, response_hash FROM requests" - ): - try: - if ( - _sha256_json(json.loads(row["response_json"])) - != row["response_hash"] - ): - violations["response_digest_mismatches"] += 1 - except Exception: - violations["response_digest_mismatches"] += 1 + ns, owner = self._identity(namespace, owner_id) + current = _normalise_time(now) + now_text = current.isoformat() + purge_before = (current - timedelta(seconds=self.tombstone_seconds)).isoformat() + exported_before = ( + current - timedelta(seconds=self.retention_seconds) + ).isoformat() + bounded = max(1, min(int(limit), 10_000)) + result = { + "sessions_tombstoned": 0, + "requests_tombstoned": 0, + "effects_compacted": 0, + "proofs_compacted": 0, + "tombstones_purged": 0, + } + with self.transaction() as connection: + session_rows = connection.execute( + """ + SELECT session_id FROM session_state + WHERE namespace = ? AND owner_id = ? AND lifecycle = 'ACTIVE' + AND expires_at IS NOT NULL AND expires_at <= ? + ORDER BY expires_at, session_id LIMIT ? + """, + (ns, owner, now_text, bounded), + ).fetchall() + for row in session_rows: + connection.execute( + """ + UPDATE session_state + SET lifecycle = 'TOMBSTONED', state_json = NULL, + tombstoned_at = ? + WHERE namespace = ? AND owner_id = ? AND session_id = ? + """, + (now_text, ns, owner, row["session_id"]), + ) + result["sessions_tombstoned"] = len(session_rows) - for row in connection.execute( - "SELECT receipt_hash, receipt_json FROM receipts" - ): - try: - payload = json.loads(row["receipt_json"]) - embedded = payload.pop("receipt_hash") - if ( - embedded != row["receipt_hash"] - or _sha256_json(payload) != row["receipt_hash"] - ): - violations["receipt_digest_mismatches"] += 1 - except Exception: - violations["receipt_digest_mismatches"] += 1 + request_rows = connection.execute( + """ + SELECT r.request_id + FROM requests r + WHERE r.namespace = ? AND r.owner_id = ? + AND r.lifecycle = 'ACTIVE' AND r.expires_at IS NOT NULL + AND r.expires_at <= ? + AND NOT EXISTS ( + SELECT 1 FROM effect_outbox e + WHERE e.namespace = r.namespace + AND e.owner_id = r.owner_id + AND e.request_id = r.request_id + AND e.status != 'EXPORTED' + ) + ORDER BY r.expires_at, r.request_id LIMIT ? + """, + (ns, owner, now_text, bounded), + ).fetchall() + for row in request_rows: + connection.execute( + """ + UPDATE requests + SET lifecycle = 'TOMBSTONED', response_json = NULL, + tombstoned_at = ? + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (now_text, ns, owner, row["request_id"]), + ) + connection.execute( + """ + UPDATE receipts SET receipt_json = NULL, tombstoned_at = ? + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (now_text, ns, owner, row["request_id"]), + ) + result["requests_tombstoned"] = len(request_rows) - for row in connection.execute("SELECT * FROM evidence_intents"): - try: - payload = json.loads(row["payload_json"]) - full_digest = _sha256_json(payload) - identity = ( - payload.get("receipt_hash") - if row["kind"] == "receipt_projection" - else payload.get("payload_sha256") - ) - if ( - full_digest != row["payload_sha256"] - or identity != row["canonical_identity"] - ): - violations["intent_digest_mismatches"] += 1 - except Exception: - violations["intent_digest_mismatches"] += 1 + effects = connection.execute( + """ + SELECT idempotency_key FROM effect_outbox + WHERE namespace = ? AND owner_id = ? AND status = 'EXPORTED' + AND tombstoned_at IS NULL AND exported_at <= ? + ORDER BY exported_at, idempotency_key LIMIT ? + """, + (ns, owner, exported_before, bounded), + ).fetchall() + for row in effects: + connection.execute( + """ + UPDATE effect_outbox + SET payload_json = NULL, artifact_json = NULL, + tombstoned_at = ? + WHERE namespace = ? AND owner_id = ? AND idempotency_key = ? + AND status = 'EXPORTED' + """, + (now_text, ns, owner, row["idempotency_key"]), + ) + result["effects_compacted"] = len(effects) + + proofs = connection.execute( + """ + SELECT proposal_id FROM proof_outbox + WHERE namespace = ? AND owner_id = ? AND status = 'EXPORTED' + AND tombstoned_at IS NULL AND exported_at <= ? + ORDER BY exported_at, proposal_id LIMIT ? + """, + (ns, owner, exported_before, bounded), + ).fetchall() + for row in proofs: + connection.execute( + """ + UPDATE proof_outbox + SET payload_json = NULL, artifact_json = NULL, + tombstoned_at = ? + WHERE namespace = ? AND owner_id = ? AND proposal_id = ? + AND status = 'EXPORTED' + """, + (now_text, ns, owner, row["proposal_id"]), + ) + result["proofs_compacted"] = len(proofs) + + purged = 0 + purged += connection.execute( + """ + DELETE FROM proof_outbox + WHERE namespace = ? AND owner_id = ? AND status = 'EXPORTED' + AND tombstoned_at IS NOT NULL AND tombstoned_at <= ? + """, + (ns, owner, purge_before), + ).rowcount + purged += connection.execute( + """ + DELETE FROM effect_outbox + WHERE namespace = ? AND owner_id = ? AND status = 'EXPORTED' + AND tombstoned_at IS NOT NULL AND tombstoned_at <= ? + """, + (ns, owner, purge_before), + ).rowcount + request_ids = [ + row["request_id"] + for row in connection.execute( + """ + SELECT request_id FROM requests + WHERE namespace = ? AND owner_id = ? + AND lifecycle = 'TOMBSTONED' + AND tombstoned_at <= ? + AND NOT EXISTS ( + SELECT 1 FROM effect_outbox e + WHERE e.namespace = requests.namespace + AND e.owner_id = requests.owner_id + AND e.request_id = requests.request_id + ) + LIMIT ? + """, + (ns, owner, purge_before, bounded), + ).fetchall() + ] + for request_id in request_ids: + connection.execute( + """ + DELETE FROM receipts + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (ns, owner, request_id), + ) + purged += connection.execute( + """ + DELETE FROM requests + WHERE namespace = ? AND owner_id = ? AND request_id = ? + """, + (ns, owner, request_id), + ).rowcount + purged += connection.execute( + """ + DELETE FROM session_state + WHERE namespace = ? AND owner_id = ? + AND lifecycle = 'TOMBSTONED' AND tombstoned_at <= ? + """, + (ns, owner, purge_before), + ).rowcount + result["tombstones_purged"] = purged + self._reconcile_usage(connection) + return result - attempt_ceiling = max( - 1, - min(int(os.environ.get("GDW_MAX_EFFECT_ATTEMPTS", "20")), 100), + def integrity( + self, + *, + namespace: Optional[str] = None, + owner_id: Optional[str] = None, + global_scope: bool = False, + ) -> Dict[str, Any]: + ns, owner = self._identity(namespace, owner_id) + connection = self._connect() + try: + check = connection.execute("PRAGMA integrity_check").fetchone()[0] + predicate = "" if global_scope else " WHERE namespace = ? AND owner_id = ?" + params = () if global_scope else (ns, owner) + counts = { + table: int( + connection.execute( + f"SELECT COUNT(*) FROM {table}{predicate}", params + ).fetchone()[0] + ) + for table in _V1_TABLES + } + effect_predicate = ( + "status IN ('PENDING', 'CLAIMED')" + if global_scope + else "namespace = ? AND owner_id = ? " + "AND status IN ('PENDING', 'CLAIMED')" + ) + effect_params = () if global_scope else (ns, owner) + pending_effects = int( + connection.execute( + f"SELECT COUNT(*) FROM effect_outbox WHERE {effect_predicate}", + effect_params, + ).fetchone()[0] + ) + claimed_predicate = ( + "status = 'CLAIMED'" + if global_scope + else "namespace = ? AND owner_id = ? AND status = 'CLAIMED'" + ) + claimed_effects = int( + connection.execute( + f"SELECT COUNT(*) FROM effect_outbox WHERE {claimed_predicate}", + effect_params, + ).fetchone()[0] + ) + dead_letter_predicate = ( + "status = 'DEAD_LETTER'" + if global_scope + else "namespace = ? AND owner_id = ? AND status = 'DEAD_LETTER'" + ) + dead_letter_effects = int( + connection.execute( + "SELECT COUNT(*) FROM effect_outbox " + f"WHERE {dead_letter_predicate}", + effect_params, + ).fetchone()[0] + ) + proof_predicate = ( + "status = 'PENDING'" + if global_scope + else "namespace = ? AND owner_id = ? AND status = 'PENDING'" + ) + pending_proofs = int( + connection.execute( + f"SELECT COUNT(*) FROM proof_outbox WHERE {proof_predicate}", + effect_params, + ).fetchone()[0] + ) + orphan_predicate = ( + "" if global_scope else "AND r.namespace = ? AND r.owner_id = ?" + ) + orphan_receipts = int( + connection.execute( + f""" + SELECT COUNT(*) FROM receipts r + LEFT JOIN requests q + ON q.namespace = r.namespace + AND q.owner_id = r.owner_id + AND q.request_id = r.request_id + WHERE q.request_id IS NULL {orphan_predicate} + """, + params, + ).fetchone()[0] ) - for row in connection.execute("SELECT * FROM effect_outbox"): + effect_rows = connection.execute( + """ + SELECT namespace, owner_id, idempotency_key, + database_generation_id, request_id, kind, receipt_hash, + payload_json, payload_sha256, intent_sha256, + status, artifact_json + FROM effect_outbox + WHERE payload_json IS NOT NULL + """ + + ( + "" + if global_scope + else " AND namespace = ? AND owner_id = ?" + ), + params, + ).fetchall() + invalid_effect_bindings = 0 + invalid_exported_artifacts = 0 + for row in effect_rows: try: - self._validate_effect_binding(connection, row) - except Exception: - violations["effect_binding_mismatches"] += 1 - if ( - row["status"] != "EXPORTED" - and int(row["attempts"]) >= attempt_ceiling + payload = json.loads(row["payload_json"]) + except (TypeError, json.JSONDecodeError): + invalid_effect_bindings += 1 + continue + candidate = dict(row) + candidate["payload"] = payload + if self._connection_effect_binding_errors( + connection, candidate ): - violations["dead_effects"] += 1 + invalid_effect_bindings += 1 if row["status"] == "EXPORTED": try: artifact = json.loads(row["artifact_json"]) - path = Path(str(artifact["path"])).resolve() - owner_scope = hashlib.sha256( - row["owner_id"].encode("utf-8") - ).hexdigest()[:32] - configured_root = ( - os.environ.get("GDW_PROOF_DIR", "output/proofs") - if row["kind"] == "proof_export" - else os.environ.get( - "GDW_RECEIPT_PROJECTION_DIR", - "output/gdw/receipts", - ) - ) - if ( - artifact.get("immutable") is not True - or artifact.get("owner_scope") != owner_scope - or artifact.get("artifact_identity") - != row["idempotency_key"] - or path.parent - != Path(configured_root).resolve() / owner_scope - or path.name - != f"{row['idempotency_key']}.json" - or not path.is_file() - or hashlib.sha256(path.read_bytes()).hexdigest() - != artifact.get("sha256") - ): - violations["artifact_mismatches"] += 1 - except Exception: - violations["artifact_mismatches"] += 1 - - generation = connection.execute( - "SELECT value FROM workspace_meta WHERE key = 'generation_id'" - ).fetchone() - generation_id = str(generation["value"]) if generation else "" - if len(generation_id) != 32: - violations["effect_binding_mismatches"] += 1 - observed_journal = str( - connection.execute("PRAGMA journal_mode").fetchone()[0] - ).upper() - if observed_journal != self.journal_mode: - violations["effect_binding_mismatches"] += 1 - violation_count = sum(int(value) for value in violations.values()) - return { - "ok": check == "ok" and violation_count == 0, + except (TypeError, json.JSONDecodeError): + invalid_exported_artifacts += 1 + else: + if self.artifact_binding_errors(candidate, artifact): + invalid_exported_artifacts += 1 + elif row["artifact_json"] is not None: + invalid_exported_artifacts += 1 + result = { + "ok": ( + check == "ok" + and orphan_receipts == 0 + and invalid_effect_bindings == 0 + and invalid_exported_artifacts == 0 + ), + "schema_version": SCHEMA_VERSION, + "database_generation_id": self.database_generation_id, "sqlite_integrity": check, "orphan_receipts": orphan_receipts, "pending_proofs": pending_proofs, "pending_effects": pending_effects, "claimed_effects": claimed_effects, + "dead_letter_effects": dead_letter_effects, + "invalid_effect_bindings": invalid_effect_bindings, + "invalid_exported_artifacts": invalid_exported_artifacts, "counts": counts, - "violations": violations, - "generation_id": generation_id, - "path": str(self.path), - "journal_mode": observed_journal, - "wal": observed_journal == "WAL", - "synchronous": "NORMAL", + "scope": "global" if global_scope else "owner", + "journal_mode": str( + connection.execute("PRAGMA journal_mode").fetchone()[0] + ).upper(), + "synchronous": self.synchronous_mode, } + if global_scope: + result["path"] = str(self.path) + else: + result["namespace"] = ns + result["owner_id"] = owner + return result finally: connection.close()