"""Durable SQLite task and outbound-event queue for an offline edge gateway.""" from __future__ import annotations import hashlib import json import re import sqlite3 import uuid from collections.abc import Callable, Iterator, Mapping from contextlib import contextmanager from dataclasses import dataclass from datetime import UTC, datetime, timedelta from pathlib import Path from types import MappingProxyType from .contracts import ( EdgeContractError, EdgeEventContract, EdgeTaskContract, SignedTaskEnvelope, canonical_sha256, canonical_timestamp, strict_json_bytes, ) from .policy import EdgeEgressPolicy, EdgePolicyError class EdgeQueueError(RuntimeError): """Base error for durable edge queue operations.""" class EdgeQueueConflictError(EdgeQueueError): """Raised when an existing stable ID is replayed with different content.""" class EdgeQueueLeaseError(EdgeQueueError): """Raised when a mutation is attempted without the current lease token.""" class EdgeQueueSchemaError(EdgeQueueError): """Raised when an edge queue database schema cannot be trusted.""" SCHEMA_VERSION = 2 _V1_TASK_TABLE_SQL = """CREATE TABLE edge_tasks ( task_id TEXT PRIMARY KEY, digest TEXT NOT NULL, task_json TEXT NOT NULL, status TEXT NOT NULL CHECK ( status IN ('pending', 'leased', 'completed', 'failed', 'cancelled') ), attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0), available_at TEXT NOT NULL, deadline_at TEXT NOT NULL, deadline_epoch_us INTEGER NOT NULL, lease_owner TEXT, lease_token TEXT, lease_expires_at TEXT, error_code TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL )""" _TASK_TABLE_SQL = _V1_TASK_TABLE_SQL.rstrip()[:-1].rstrip() + ",\n" + ( " authority_key_id TEXT,\n" " authority_digest TEXT,\n" " authority_json TEXT,\n" " remote_lease_token_digest TEXT,\n" " remote_lease_ciphertext TEXT,\n" " remote_lease_expires_at TEXT,\n" " remote_attempt INTEGER CHECK (remote_attempt IS NULL OR remote_attempt >= 1)\n" ")" ) _TASK_INDEX_SQL = """CREATE INDEX edge_tasks_claim_idx ON edge_tasks( status, deadline_epoch_us, available_at, lease_expires_at, created_at )""" _V1_EVENT_TABLE_SQL = """CREATE TABLE edge_outbound_events ( event_id TEXT PRIMARY KEY, digest TEXT NOT NULL, event_json TEXT NOT NULL, status TEXT NOT NULL CHECK ( status IN ('pending', 'sending', 'acknowledged', 'failed') ), attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0), available_at TEXT NOT NULL, lease_owner TEXT, lease_token TEXT, lease_expires_at TEXT, error_code TEXT, acknowledgement_json TEXT, ack_digest TEXT, lease_token_digest TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL )""" _EVENT_TABLE_SQL = _V1_EVENT_TABLE_SQL.rstrip()[:-1].rstrip() + ",\n" + ( " task_digest TEXT,\n" " remote_lease_token_digest TEXT,\n" " control_lease_token_digest TEXT\n" ")" ) _EVENT_INDEX_SQL = """CREATE INDEX edge_events_claim_idx ON edge_outbound_events( status, available_at, lease_expires_at, created_at )""" _ARTIFACT_TABLE_SQL = """CREATE TABLE edge_local_artifacts ( artifact_digest TEXT PRIMARY KEY, artifact_ref TEXT, artifact_ref_hash TEXT NOT NULL, classification TEXT NOT NULL CHECK ( classification IN ('raw', 'recent_detail', 'restricted', 'desensitized_metadata', 'statistics', 'lineage', 'evidence') ), retention_until TEXT NOT NULL, task_id TEXT NOT NULL, event_id TEXT NOT NULL, cleanup_status TEXT NOT NULL CHECK ( cleanup_status IN ('pending', 'deleting', 'deleted', 'failed') ), attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0), available_at TEXT NOT NULL, lease_owner TEXT, lease_token TEXT, lease_expires_at TEXT, error_code TEXT, deletion_receipt_json TEXT, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, FOREIGN KEY(task_id) REFERENCES edge_tasks(task_id), FOREIGN KEY(event_id) REFERENCES edge_outbound_events(event_id) )""" _ARTIFACT_INDEX_SQL = """CREATE INDEX edge_artifacts_task_idx ON edge_local_artifacts( cleanup_status, retention_until, available_at, lease_expires_at, created_at )""" _RELEASE_TABLE_SQL = """CREATE TABLE edge_release_state ( release_id TEXT PRIMARY KEY, manifest_digest TEXT NOT NULL, version TEXT NOT NULL, rollback_version TEXT NOT NULL, status TEXT NOT NULL CHECK ( status IN ('offered', 'accepted', 'candidate', 'installed', 'failed', 'rolled_back') ), previous_version TEXT, current_version TEXT, updated_at TEXT NOT NULL )""" _RELEASE_INDEX_SQL = """CREATE INDEX edge_release_status_idx ON edge_release_state( status, updated_at, release_id )""" _OUTCOME_TABLE_SQL = """CREATE TABLE edge_task_outcomes ( task_id TEXT PRIMARY KEY, digest TEXT NOT NULL, outcome_json TEXT NOT NULL, status TEXT NOT NULL CHECK ( status IN ('pending', 'sending', 'acknowledged', 'failed') ), attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0), available_at TEXT NOT NULL, lease_owner TEXT, lease_token TEXT, lease_expires_at TEXT, error_code TEXT, acknowledgement_json TEXT, ack_digest TEXT, remote_lease_token_digest TEXT NOT NULL, created_at TEXT NOT NULL, updated_at TEXT NOT NULL, FOREIGN KEY(task_id) REFERENCES edge_tasks(task_id) )""" _OUTCOME_INDEX_SQL = """CREATE INDEX edge_outcomes_claim_idx ON edge_task_outcomes( status, available_at, lease_expires_at, created_at )""" _EXPECTED_TABLE_INFO = { "edge_tasks": ( ("task_id", "TEXT", 0, None, 1), ("digest", "TEXT", 1, None, 0), ("task_json", "TEXT", 1, None, 0), ("status", "TEXT", 1, None, 0), ("attempt_count", "INTEGER", 1, "0", 0), ("available_at", "TEXT", 1, None, 0), ("deadline_at", "TEXT", 1, None, 0), ("deadline_epoch_us", "INTEGER", 1, None, 0), ("lease_owner", "TEXT", 0, None, 0), ("lease_token", "TEXT", 0, None, 0), ("lease_expires_at", "TEXT", 0, None, 0), ("error_code", "TEXT", 0, None, 0), ("created_at", "TEXT", 1, None, 0), ("updated_at", "TEXT", 1, None, 0), ("authority_key_id", "TEXT", 0, None, 0), ("authority_digest", "TEXT", 0, None, 0), ("authority_json", "TEXT", 0, None, 0), ("remote_lease_token_digest", "TEXT", 0, None, 0), ("remote_lease_ciphertext", "TEXT", 0, None, 0), ("remote_lease_expires_at", "TEXT", 0, None, 0), ("remote_attempt", "INTEGER", 0, None, 0), ), "edge_outbound_events": ( ("event_id", "TEXT", 0, None, 1), ("digest", "TEXT", 1, None, 0), ("event_json", "TEXT", 1, None, 0), ("status", "TEXT", 1, None, 0), ("attempt_count", "INTEGER", 1, "0", 0), ("available_at", "TEXT", 1, None, 0), ("lease_owner", "TEXT", 0, None, 0), ("lease_token", "TEXT", 0, None, 0), ("lease_expires_at", "TEXT", 0, None, 0), ("error_code", "TEXT", 0, None, 0), ("acknowledgement_json", "TEXT", 0, None, 0), ("ack_digest", "TEXT", 0, None, 0), ("lease_token_digest", "TEXT", 0, None, 0), ("created_at", "TEXT", 1, None, 0), ("updated_at", "TEXT", 1, None, 0), ("task_digest", "TEXT", 0, None, 0), ("remote_lease_token_digest", "TEXT", 0, None, 0), ("control_lease_token_digest", "TEXT", 0, None, 0), ), "edge_local_artifacts": ( ("artifact_digest", "TEXT", 0, None, 1), ("artifact_ref", "TEXT", 0, None, 0), ("artifact_ref_hash", "TEXT", 1, None, 0), ("classification", "TEXT", 1, None, 0), ("retention_until", "TEXT", 1, None, 0), ("task_id", "TEXT", 1, None, 0), ("event_id", "TEXT", 1, None, 0), ("cleanup_status", "TEXT", 1, None, 0), ("attempt_count", "INTEGER", 1, "0", 0), ("available_at", "TEXT", 1, None, 0), ("lease_owner", "TEXT", 0, None, 0), ("lease_token", "TEXT", 0, None, 0), ("lease_expires_at", "TEXT", 0, None, 0), ("error_code", "TEXT", 0, None, 0), ("deletion_receipt_json", "TEXT", 0, None, 0), ("created_at", "TEXT", 1, None, 0), ("updated_at", "TEXT", 1, None, 0), ), "edge_release_state": ( ("release_id", "TEXT", 0, None, 1), ("manifest_digest", "TEXT", 1, None, 0), ("version", "TEXT", 1, None, 0), ("rollback_version", "TEXT", 1, None, 0), ("status", "TEXT", 1, None, 0), ("previous_version", "TEXT", 0, None, 0), ("current_version", "TEXT", 0, None, 0), ("updated_at", "TEXT", 1, None, 0), ), "edge_task_outcomes": ( ("task_id", "TEXT", 0, None, 1), ("digest", "TEXT", 1, None, 0), ("outcome_json", "TEXT", 1, None, 0), ("status", "TEXT", 1, None, 0), ("attempt_count", "INTEGER", 1, "0", 0), ("available_at", "TEXT", 1, None, 0), ("lease_owner", "TEXT", 0, None, 0), ("lease_token", "TEXT", 0, None, 0), ("lease_expires_at", "TEXT", 0, None, 0), ("error_code", "TEXT", 0, None, 0), ("acknowledgement_json", "TEXT", 0, None, 0), ("ack_digest", "TEXT", 0, None, 0), ("remote_lease_token_digest", "TEXT", 1, None, 0), ("created_at", "TEXT", 1, None, 0), ("updated_at", "TEXT", 1, None, 0), ), } _EXPECTED_INDEXES = { "edge_tasks": { "edge_tasks_claim_idx": ( 0, "c", 0, ("status", "deadline_epoch_us", "available_at", "lease_expires_at", "created_at"), ), }, "edge_outbound_events": { "edge_events_claim_idx": ( 0, "c", 0, ("status", "available_at", "lease_expires_at", "created_at"), ), }, "edge_local_artifacts": { "edge_artifacts_task_idx": ( 0, "c", 0, ("cleanup_status", "retention_until", "available_at", "lease_expires_at", "created_at"), ), }, "edge_release_state": { "edge_release_status_idx": (0, "c", 0, ("status", "updated_at", "release_id")), }, "edge_task_outcomes": { "edge_outcomes_claim_idx": ( 0, "c", 0, ("status", "available_at", "lease_expires_at", "created_at") ), }, } _EXPECTED_TABLE_SQL = { "edge_tasks": _TASK_TABLE_SQL, "edge_outbound_events": _EVENT_TABLE_SQL, "edge_local_artifacts": _ARTIFACT_TABLE_SQL, "edge_release_state": _RELEASE_TABLE_SQL, "edge_task_outcomes": _OUTCOME_TABLE_SQL, } _V1_EXPECTED_TABLE_INFO = { table: tuple(row for row in info if row[0] not in { "authority_key_id", "authority_digest", "authority_json", "remote_lease_token_digest", "remote_lease_ciphertext", "remote_lease_expires_at", "remote_attempt", "task_digest", "control_lease_token_digest", }) for table, info in _EXPECTED_TABLE_INFO.items() if table not in {"edge_local_artifacts", "edge_release_state", "edge_task_outcomes"} } _V1_EXPECTED_TABLE_SQL = { "edge_tasks": _V1_TASK_TABLE_SQL, "edge_outbound_events": _V1_EVENT_TABLE_SQL, } @dataclass(frozen=True, slots=True) class QueuedTask: task_id: str digest: str status: str task: Mapping[str, object] attempt_count: int available_at: str lease_owner: str | None lease_token: str | None lease_expires_at: str | None error_code: str | None authority_key_id: str | None authority_digest: str | None remote_lease_token_digest: str | None remote_lease_expires_at: str | None remote_attempt: int | None created_at: str updated_at: str @dataclass(frozen=True, slots=True) class QueuedEvent: event_id: str digest: str status: str event: Mapping[str, object] attempt_count: int available_at: str lease_owner: str | None lease_token: str | None lease_expires_at: str | None error_code: str | None acknowledgement: Mapping[str, object] | None task_digest: str | None remote_lease_token_digest: str | None control_lease_token_digest: str | None created_at: str updated_at: str @dataclass(frozen=True, slots=True) class LocalArtifact: artifact_digest: str artifact_ref: str | None artifact_ref_hash: str classification: str retention_until: str task_id: str event_id: str cleanup_status: str attempt_count: int available_at: str lease_owner: str | None lease_token: str | None lease_expires_at: str | None error_code: str | None deletion_receipt: Mapping[str, object] | None created_at: str updated_at: str @dataclass(frozen=True, slots=True) class ReleaseState: release_id: str manifest_digest: str version: str rollback_version: str status: str previous_version: str | None current_version: str | None updated_at: str @dataclass(frozen=True, slots=True) class QueuedOutcome: task_id: str digest: str status: str outcome: Mapping[str, object] attempt_count: int available_at: str lease_owner: str | None lease_token: str | None lease_expires_at: str | None error_code: str | None acknowledgement: Mapping[str, object] | None remote_lease_token_digest: str created_at: str updated_at: str _ERROR_CODE = re.compile(r"^[a-z][a-z0-9_.-]{0,127}$") _SENSITIVE_ERROR_PARTS = ("password", "passwd", "secret", "token", "credential") def _utc_now(clock: Callable[[], object]) -> datetime: value = clock() if isinstance(value, bool): raise ValueError("queue clock returned an invalid value") if isinstance(value, (int, float)): value = datetime.fromtimestamp(value, tz=UTC) if not isinstance(value, datetime): raise ValueError("queue clock must return a datetime or timestamp") if value.tzinfo is None or value.utcoffset() is None: raise ValueError("queue clock must return a timezone-aware value") return value.astimezone(UTC) def _timestamp(value: datetime) -> str: utc = value.astimezone(UTC) canonical = utc.strftime("%Y-%m-%dT%H:%M:%S") if utc.microsecond: canonical += f".{utc.microsecond:06d}".rstrip("0") return f"{canonical}Z" def _epoch_microseconds(value: datetime) -> int: epoch = datetime(1970, 1, 1, tzinfo=UTC) delta = value.astimezone(UTC) - epoch return ( (delta.days * 86_400 + delta.seconds) * 1_000_000 + delta.microseconds ) def _normalize_schema_sql(value: str | None) -> str: if not isinstance(value, str): return "" normalized = re.sub(r"\s+", " ", value.strip()).lower() return re.sub(r"\s*([(),])\s*", r"\1", normalized) def _freeze_mapping(value: Mapping[str, object]) -> Mapping[str, object]: def freeze(item: object) -> object: if isinstance(item, dict): return MappingProxyType({key: freeze(child) for key, child in item.items()}) if isinstance(item, list): return tuple(freeze(child) for child in item) return item return freeze(dict(value)) # type: ignore[return-value] class SqliteEdgeQueue: """An on-disk, WAL-backed queue with fenced leases and exact replay.""" def __init__( self, db_path: str | Path, *, clock: Callable[[], object] | None = None, lease_seconds: int = 60, max_attempts: int = 5, retry_base_seconds: int = 2, retry_max_seconds: int = 300, busy_timeout_ms: int = 5_000, egress_policy: EdgeEgressPolicy | None = None, random_source: Callable[[], float] | None = None, remote_lease_encoder: Callable[[str], str] | None = None, remote_lease_decoder: Callable[[str], str] | None = None, ) -> None: if not isinstance(db_path, (str, Path)) or not str(db_path).strip(): raise ValueError("an explicit SQLite database path is required") path_text = str(db_path) if path_text == ":memory:" or "mode=memory" in path_text: raise ValueError("an in-memory database is not allowed") for label, value in { "lease_seconds": lease_seconds, "max_attempts": max_attempts, "retry_base_seconds": retry_base_seconds, "retry_max_seconds": retry_max_seconds, "busy_timeout_ms": busy_timeout_ms, }.items(): if isinstance(value, bool) or not isinstance(value, int) or value < 1: raise ValueError(f"{label} must be a positive integer") if retry_base_seconds > retry_max_seconds: raise ValueError("retry_base_seconds cannot exceed retry_max_seconds") self.db_path = str(Path(db_path).expanduser().resolve()) self._clock = clock or (lambda: datetime.now(UTC)) self.lease_seconds = lease_seconds self.max_attempts = max_attempts self.retry_base_seconds = retry_base_seconds self.retry_max_seconds = retry_max_seconds self.busy_timeout_ms = busy_timeout_ms self.egress_policy = egress_policy or EdgeEgressPolicy( allowed_control_hosts=set() ) # The injectable source makes retry timing testable; 0.5 is neutral jitter. self._random = random_source or (lambda: 0.5) if (remote_lease_encoder is None) != (remote_lease_decoder is None): raise ValueError("remote lease encoder and decoder must be paired") self._remote_lease_encoder = remote_lease_encoder self._remote_lease_decoder = remote_lease_decoder Path(self.db_path).parent.mkdir(parents=True, exist_ok=True) self._initialize() def _connect(self) -> sqlite3.Connection: connection = sqlite3.connect( self.db_path, timeout=self.busy_timeout_ms / 1_000, isolation_level=None, ) connection.row_factory = sqlite3.Row connection.execute(f"PRAGMA busy_timeout = {self.busy_timeout_ms}") connection.execute("PRAGMA foreign_keys = ON") return connection def _initialize(self) -> None: connection = self._connect() try: version = connection.execute("PRAGMA user_version").fetchone()[0] edge_tables = { row[0] for row in connection.execute( """ SELECT name FROM sqlite_master WHERE type = 'table' AND name LIKE 'edge_%' """ ).fetchall() } if version == 0 and edge_tables: raise EdgeQueueSchemaError("unversioned edge queue schema is unsafe") if version not in {0, 1, SCHEMA_VERSION}: raise EdgeQueueSchemaError("edge queue schema version is unsupported") connection.execute("PRAGMA journal_mode = WAL") connection.execute("PRAGMA synchronous = FULL") connection.execute("BEGIN IMMEDIATE") if version == 0: for statement in ( _TASK_TABLE_SQL, _TASK_INDEX_SQL, _EVENT_TABLE_SQL, _EVENT_INDEX_SQL, _ARTIFACT_TABLE_SQL, _ARTIFACT_INDEX_SQL, _RELEASE_TABLE_SQL, _RELEASE_INDEX_SQL, _OUTCOME_TABLE_SQL, _OUTCOME_INDEX_SQL, ): connection.execute(statement) connection.execute(f"PRAGMA user_version = {SCHEMA_VERSION}") elif version == 1: self._validate_schema( connection, expected_info=_V1_EXPECTED_TABLE_INFO, expected_sql=_V1_EXPECTED_TABLE_SQL, expected_indexes={ key: _EXPECTED_INDEXES[key] for key in ("edge_tasks", "edge_outbound_events") }, ) active_unsigned = connection.execute( """SELECT EXISTS( SELECT 1 FROM edge_tasks WHERE status IN ('pending','leased') )""" ).fetchone()[0] unacknowledged_events = connection.execute( """SELECT EXISTS( SELECT 1 FROM edge_outbound_events WHERE status <> 'acknowledged' )""" ).fetchone()[0] if active_unsigned or unacknowledged_events: raise EdgeQueueSchemaError( "v1 edge queue must drain active tasks and unacknowledged events" ) for statement in ( "ALTER TABLE edge_tasks ADD COLUMN authority_key_id TEXT", "ALTER TABLE edge_tasks ADD COLUMN authority_digest TEXT", "ALTER TABLE edge_tasks ADD COLUMN authority_json TEXT", "ALTER TABLE edge_tasks ADD COLUMN remote_lease_token_digest TEXT", "ALTER TABLE edge_tasks ADD COLUMN remote_lease_ciphertext TEXT", "ALTER TABLE edge_tasks ADD COLUMN remote_lease_expires_at TEXT", "ALTER TABLE edge_tasks ADD COLUMN remote_attempt INTEGER CHECK (remote_attempt IS NULL OR remote_attempt >= 1)", "ALTER TABLE edge_outbound_events ADD COLUMN task_digest TEXT", "ALTER TABLE edge_outbound_events ADD COLUMN remote_lease_token_digest TEXT", "ALTER TABLE edge_outbound_events ADD COLUMN control_lease_token_digest TEXT", _ARTIFACT_TABLE_SQL, _ARTIFACT_INDEX_SQL, _RELEASE_TABLE_SQL, _RELEASE_INDEX_SQL, _OUTCOME_TABLE_SQL, _OUTCOME_INDEX_SQL, ): connection.execute(statement) connection.execute(f"PRAGMA user_version = {SCHEMA_VERSION}") self._validate_schema(connection) if connection.in_transaction: connection.execute("COMMIT") except Exception: if connection.in_transaction: connection.execute("ROLLBACK") raise finally: connection.close() @staticmethod def _validate_schema( connection: sqlite3.Connection, *, expected_info: Mapping[str, tuple] = _EXPECTED_TABLE_INFO, expected_sql: Mapping[str, str] = _EXPECTED_TABLE_SQL, expected_indexes: Mapping[str, Mapping[str, tuple]] = _EXPECTED_INDEXES, ) -> None: primary_keys = { "edge_tasks": "task_id", "edge_outbound_events": "event_id", "edge_local_artifacts": "artifact_digest", "edge_release_state": "release_id", "edge_task_outcomes": "task_id", } for table, table_expected_info in expected_info.items(): rows = connection.execute(f"PRAGMA table_info({table})").fetchall() actual_info = tuple( (row[1], row[2].upper(), row[3], row[4], row[5]) for row in rows ) if actual_info != table_expected_info: raise EdgeQueueSchemaError("edge queue schema columns are invalid") table_sql = connection.execute( "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = ?", (table,), ).fetchone() if table_sql is None or _normalize_schema_sql( table_sql[0] ) != _normalize_schema_sql(expected_sql[table]): raise EdgeQueueSchemaError( "edge queue schema constraints are invalid" ) index_rows = connection.execute(f"PRAGMA index_list({table})").fetchall() custom_rows = {row[1]: row for row in index_rows if row[3] == "c"} table_expected_indexes = expected_indexes[table] if set(custom_rows) != set(table_expected_indexes): raise EdgeQueueSchemaError("edge queue schema indexes are invalid") for name, (unique, origin, partial, columns) in table_expected_indexes.items(): row = custom_rows[name] actual_columns = tuple( item[2] for item in connection.execute( f"PRAGMA index_info({name})" ).fetchall() ) if (row[2], row[3], row[4], actual_columns) != ( unique, origin, partial, columns, ): raise EdgeQueueSchemaError( "edge queue schema index definition is invalid" ) primary_indexes = [row for row in index_rows if row[3] == "pk"] if len(primary_indexes) != 1: raise EdgeQueueSchemaError("edge queue schema primary key is invalid") primary = primary_indexes[0] primary_columns = tuple( item[2] for item in connection.execute( f"PRAGMA index_info({primary[1]})" ).fetchall() ) if (primary[2], primary[4], primary_columns) != ( 1, 0, (primary_keys[table],), ): raise EdgeQueueSchemaError("edge queue schema primary key is invalid") @contextmanager def _transaction(self) -> Iterator[sqlite3.Connection]: connection = self._connect() try: connection.execute("BEGIN IMMEDIATE") yield connection connection.execute("COMMIT") except Exception: if connection.in_transaction: connection.execute("ROLLBACK") raise finally: connection.close() def enqueue(self, task: EdgeTaskContract | Mapping[str, object]) -> QueuedTask: contract = EdgeTaskContract.from_mapping( task.to_mapping() if isinstance(task, EdgeTaskContract) else task ) return self._enqueue_contract(contract, authority=None) def enqueue_signed(self, envelope: SignedTaskEnvelope) -> QueuedTask: """Atomically persist a previously verified authority envelope and task.""" if not isinstance(envelope, SignedTaskEnvelope) or not envelope.is_verified: raise EdgeContractError("a verified SignedTaskEnvelope is required") # Reconstruct the task so object.__setattr__ cannot bypass validation. contract = EdgeTaskContract.from_mapping(envelope.task.to_mapping()) if envelope.contract_digest != contract.digest: raise EdgeContractError("signed authority no longer binds the task") return self._enqueue_contract(contract, authority=envelope) def accept_signed_task( self, envelope: SignedTaskEnvelope, *, remote_lease_token: str, remote_lease_expires_at: str, remote_attempt: int, ) -> QueuedTask: """Atomically accept authority and an encrypted recoverable control lease.""" if self._remote_lease_encoder is None: raise EdgeQueueLeaseError("secure remote lease codec is required") if not isinstance(envelope, SignedTaskEnvelope) or not envelope.is_verified: raise EdgeContractError("a verified SignedTaskEnvelope is required") return self._enqueue_contract( EdgeTaskContract.from_mapping(envelope.task.to_mapping()), authority=envelope, remote_binding=self._normalize_remote_binding( remote_lease_token, remote_lease_expires_at, remote_attempt ), ) def _enqueue_contract( self, contract: EdgeTaskContract, *, authority: SignedTaskEnvelope | None, remote_binding: tuple[str, str | None, str, int] | None = None, ) -> QueuedTask: task_mapping = contract.to_mapping() encoded = strict_json_bytes(task_mapping).decode("utf-8") digest = canonical_sha256(task_mapping) if remote_binding is not None and remote_binding[3] != contract.attempt: raise EdgeQueueLeaseError( "remote lease attempt does not bind the task contract" ) with self._transaction() as connection: existing = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (contract.task_id,) ).fetchone() if existing is not None: if existing["digest"] != digest: raise EdgeQueueConflictError( "task ID already exists with a different digest" ) if authority is not None and existing["authority_digest"] != authority.digest: raise EdgeQueueConflictError( "task ID already exists with a different authority digest" ) if remote_binding is not None: now = _utc_now(self._clock) old_expiry_text = existing["remote_lease_expires_at"] old_expiry = ( datetime.fromisoformat(old_expiry_text.replace("Z", "+00:00")) if old_expiry_text is not None else None ) if ( existing["remote_lease_token_digest"] != remote_binding[0] and old_expiry is not None and old_expiry > now ): raise EdgeQueueConflictError("active remote lease binding conflicts") new_expiry = datetime.fromisoformat(remote_binding[2].replace("Z", "+00:00")) if new_expiry <= now: raise EdgeQueueLeaseError("remote lease has expired") if ( existing["status"] in {"failed", "cancelled"} and not ( existing["status"] == "failed" and existing["error_code"] == "execution_lease_expired" ) ): raise EdgeQueueConflictError( "terminal task cannot be rebound to a remote lease" ) changed_remote_lease = ( existing["remote_lease_token_digest"] != remote_binding[0] ) if changed_remote_lease: event_rows = connection.execute( """SELECT * FROM edge_outbound_events WHERE status IN ('pending','sending') AND json_extract(event_json,'$.task_id')=?""", (contract.task_id,), ).fetchall() for event_row in event_rows: event_mapping = json.loads(event_row["event_json"]) if ( event_row["task_digest"] != digest or event_row["remote_lease_token_digest"] != existing["remote_lease_token_digest"] or event_mapping.get("task_id") != contract.task_id or event_mapping.get("gateway_id") != contract.gateway_id or canonical_sha256(event_mapping) != event_row["digest"] ): raise EdgeQueueConflictError( "pending event does not bind the task lease" ) connection.execute( """UPDATE edge_outbound_events SET remote_lease_token_digest=?,updated_at=? WHERE status IN ('pending','sending') AND json_extract(event_json,'$.task_id')=?""", ( remote_binding[0], _timestamp(now), contract.task_id, ), ) recoverable_failure = ( existing["status"] == "failed" and existing["error_code"] == "execution_lease_expired" ) connection.execute( """UPDATE edge_tasks SET remote_lease_token_digest=?, remote_lease_ciphertext=?,remote_lease_expires_at=?,remote_attempt=?, status=CASE WHEN ? THEN 'pending' ELSE status END, available_at=CASE WHEN ? THEN ? ELSE available_at END, error_code=CASE WHEN ? THEN NULL ELSE error_code END, updated_at=? WHERE task_id=? AND digest=? AND authority_digest=?""", ( remote_binding[0], remote_binding[1], remote_binding[2], remote_binding[3], recoverable_failure, recoverable_failure, _timestamp(now), recoverable_failure, _timestamp(now), contract.task_id, digest, authority.digest if authority is not None else None, ), ) existing = connection.execute( "SELECT * FROM edge_tasks WHERE task_id=?", (contract.task_id,) ).fetchone() return self._task_record(existing) now = _utc_now(self._clock) deadline = datetime.fromisoformat( contract.deadline_at.replace("Z", "+00:00") ).astimezone(UTC) if deadline <= now: raise EdgeContractError("task deadline has expired") if remote_binding is not None: remote_expiry = datetime.fromisoformat(remote_binding[2].replace("Z", "+00:00")) if remote_expiry <= now or remote_expiry > deadline: raise EdgeQueueLeaseError("remote lease validity window is invalid") now_text = _timestamp(now) connection.execute( """ INSERT INTO edge_tasks ( task_id, digest, task_json, status, attempt_count, available_at, deadline_at, deadline_epoch_us, authority_key_id, authority_digest, authority_json, remote_lease_token_digest, remote_lease_ciphertext, remote_lease_expires_at, remote_attempt, created_at, updated_at ) VALUES (?, ?, ?, 'pending', 0, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( contract.task_id, digest, encoded, now_text, contract.deadline_at, _epoch_microseconds(deadline), authority.authority_key_id if authority is not None else None, authority.digest if authority is not None else None, strict_json_bytes(authority.to_mapping()).decode("utf-8") if authority is not None else None, remote_binding[0] if remote_binding is not None else None, remote_binding[1] if remote_binding is not None else None, remote_binding[2] if remote_binding is not None else None, remote_binding[3] if remote_binding is not None else None, now_text, now_text, ), ) row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (contract.task_id,) ).fetchone() return self._task_record(row) def claim(self, lease_owner: str) -> QueuedTask | None: owner = self._lease_owner(lease_owner) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) now_epoch_us = _epoch_microseconds(now) lease_token = str(uuid.uuid4()) expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds)) connection.execute( """ UPDATE edge_tasks SET status = 'failed', lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, error_code = 'task_deadline_expired', updated_at = ? WHERE status IN ('pending', 'leased') AND deadline_epoch_us <= ? """, (now_text, now_epoch_us), ) connection.execute( """ UPDATE edge_tasks SET status = 'failed', lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, error_code = 'lease_attempts_exhausted', updated_at = ? WHERE status = 'leased' AND julianday(lease_expires_at) <= julianday(?) AND attempt_count >= ? """, (now_text, now_text, self.max_attempts), ) row = connection.execute( """ SELECT * FROM edge_tasks WHERE deadline_epoch_us > ? AND attempt_count < ? AND ( (status = 'pending' AND julianday(available_at) <= julianday(?)) OR (status = 'leased' AND julianday(lease_expires_at) <= julianday(?)) ) ORDER BY created_at, task_id LIMIT 1 """, (now_epoch_us, self.max_attempts, now_text, now_text), ).fetchone() if row is None: return None updated = connection.execute( """ UPDATE edge_tasks SET status = 'leased', attempt_count = attempt_count + 1, lease_owner = ?, lease_token = ?, lease_expires_at = ?, error_code = NULL, updated_at = ? WHERE task_id = ? AND deadline_epoch_us > ? AND attempt_count = ? AND ( (status = 'pending' AND julianday(available_at) <= julianday(?)) OR (status = 'leased' AND julianday(lease_expires_at) <= julianday(?)) ) """, ( owner, lease_token, expires_at, now_text, row["task_id"], now_epoch_us, row["attempt_count"], now_text, now_text, ), ) if updated.rowcount != 1: return None claimed = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (row["task_id"],) ).fetchone() return self._task_record(claimed) def bind_remote_lease( self, task_id: str, *, lease_token: str, remote_lease_token: str, remote_lease_expires_at: str, remote_attempt: int, ) -> QueuedTask: """Bind control-plane authority to the signed contract attempt.""" digest, ciphertext, expiry_text, remote_attempt = self._normalize_remote_binding( remote_lease_token, remote_lease_expires_at, remote_attempt ) expiry = datetime.fromisoformat(expiry_text.replace("Z", "+00:00")) with self._transaction() as connection: row = self._leased_row(connection, "edge_tasks", "task_id", task_id, lease_token) now = _utc_now(self._clock) if expiry <= now: raise EdgeQueueLeaseError("remote lease has expired") task_mapping = json.loads(row["task_json"]) if remote_attempt != task_mapping["attempt"]: raise EdgeQueueLeaseError("remote lease attempt does not bind task contract") existing = row["remote_lease_token_digest"] if existing is not None and ( existing != digest or row["remote_lease_expires_at"] != expiry_text or row["remote_attempt"] != remote_attempt ): raise EdgeQueueConflictError("remote lease binding conflicts") connection.execute( """ UPDATE edge_tasks SET remote_lease_token_digest=?, remote_lease_ciphertext=?, remote_lease_expires_at=?, remote_attempt=?, updated_at=? WHERE task_id=? AND status='leased' AND lease_token=? """, (digest, ciphertext, expiry_text, remote_attempt, _timestamp(now), task_id, lease_token), ) updated = connection.execute( "SELECT * FROM edge_tasks WHERE task_id=?", (task_id,) ).fetchone() return self._task_record(updated) def recover_remote_lease(self, task_id: str) -> str: """Recover a control lease through the injected secure decoder without exposing it in records.""" if self._remote_lease_decoder is None: raise EdgeQueueLeaseError("secure remote lease codec is required") connection = self._connect() try: row = connection.execute( "SELECT remote_lease_ciphertext,remote_lease_token_digest FROM edge_tasks WHERE task_id=?", (task_id,), ).fetchone() finally: connection.close() if row is None or row["remote_lease_ciphertext"] is None: raise EdgeQueueLeaseError("recoverable remote lease does not exist") try: token = self._remote_lease_decoder(row["remote_lease_ciphertext"]) except Exception as exc: raise EdgeQueueLeaseError("remote lease recovery failed") from exc if not isinstance(token, str) or self._opaque_token_digest(token) != row["remote_lease_token_digest"]: raise EdgeQueueLeaseError("recovered remote lease failed integrity validation") return token def complete_with_event( self, task_id: str, *, lease_token: str, remote_lease_token: str, event: EdgeEventContract | Mapping[str, object], artifacts: list[Mapping[str, object]] | tuple[Mapping[str, object], ...] = (), ) -> tuple[QueuedTask, QueuedEvent]: """Atomically commit a safe event, local artifact metadata, and task terminal state.""" contract = EdgeEventContract.from_mapping( event.to_mapping() if isinstance(event, EdgeEventContract) else event ) approved = self.egress_policy.approve_event(contract.to_mapping()) encoded_event = strict_json_bytes(approved).decode("utf-8") normalized_artifacts = [self._validate_local_artifact(item) for item in artifacts] if not isinstance(remote_lease_token, str) or not remote_lease_token: raise EdgeQueueLeaseError("remote lease token is invalid") remote_digest = self._opaque_token_digest(remote_lease_token) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) row = self._leased_row(connection, "edge_tasks", "task_id", task_id, lease_token) if row["deadline_epoch_us"] <= _epoch_microseconds(now): raise EdgeQueueLeaseError("task deadline has expired") if row["remote_lease_token_digest"] != remote_digest: raise EdgeQueueLeaseError("current remote lease token is required") remote_expiry = row["remote_lease_expires_at"] if remote_expiry is None or datetime.fromisoformat( remote_expiry.replace("Z", "+00:00") ) <= now: raise EdgeQueueLeaseError("remote lease has expired") task_mapping = json.loads(row["task_json"]) if row["remote_attempt"] != task_mapping["attempt"]: raise EdgeQueueLeaseError("remote lease attempt is stale") for field in ("task_id", "gateway_id", "environment", "network_zone", "purpose", "policy_digest"): if contract.to_mapping()[field] != task_mapping[field]: raise EdgeQueueConflictError(f"event {field} does not bind the task") if contract.attempt != task_mapping["attempt"]: raise EdgeQueueConflictError("event attempt does not bind the task") existing = connection.execute( "SELECT digest FROM edge_outbound_events WHERE event_id=?", (contract.event_id,) ).fetchone() if existing is not None: raise EdgeQueueConflictError("event ID already exists before task completion") connection.execute( """ INSERT INTO edge_outbound_events ( event_id,digest,event_json,status,attempt_count,available_at, task_digest,remote_lease_token_digest,created_at,updated_at ) VALUES (?,?,?,'pending',0,?,?,?,?,?) """, ( contract.event_id, contract.digest, encoded_event, now_text, row["digest"], remote_digest, now_text, now_text, ), ) for artifact in normalized_artifacts: connection.execute( """ INSERT INTO edge_local_artifacts ( artifact_digest,artifact_ref,artifact_ref_hash,classification, retention_until,task_id,event_id,cleanup_status,attempt_count, available_at,created_at,updated_at ) VALUES (?,?,?,?,?,?,?,'pending',0,?,?,?) """, ( artifact["artifact_digest"], artifact["artifact_ref"], artifact["artifact_ref_hash"], artifact["classification"], artifact["retention_until"], task_id, contract.event_id, now_text, now_text, now_text, ), ) updated = connection.execute( """ UPDATE edge_tasks SET status='completed', lease_owner=NULL, lease_token=NULL, lease_expires_at=NULL, error_code=NULL, updated_at=? WHERE task_id=? AND status='leased' AND lease_token=? """, (now_text, task_id, lease_token), ) if updated.rowcount != 1: raise EdgeQueueLeaseError("task completion lost its lease") task_row = connection.execute("SELECT * FROM edge_tasks WHERE task_id=?", (task_id,)).fetchone() event_row = connection.execute("SELECT * FROM edge_outbound_events WHERE event_id=?", (contract.event_id,)).fetchone() return self._task_record(task_row), self._event_record(event_row) def list_local_artifacts(self, *, task_id: str) -> tuple[LocalArtifact, ...]: connection = self._connect() try: rows = connection.execute( "SELECT * FROM edge_local_artifacts WHERE task_id=? ORDER BY created_at,artifact_digest", (task_id,), ).fetchall() finally: connection.close() return tuple(self._artifact_record(row) for row in rows) def claim_artifact_cleanup(self, lease_owner: str) -> LocalArtifact | None: owner = self._lease_owner(lease_owner) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds)) connection.execute( """UPDATE edge_local_artifacts SET cleanup_status='failed', lease_owner=NULL,lease_token=NULL,lease_expires_at=NULL, error_code='lease_attempts_exhausted',updated_at=? WHERE cleanup_status='deleting' AND julianday(lease_expires_at)<=julianday(?) AND attempt_count>=?""", (now_text, now_text, self.max_attempts), ) row = connection.execute( """SELECT * FROM edge_local_artifacts WHERE attempt_count LocalArtifact: required = {"artifact_digest", "artifact_ref_hash", "deleted_at", "status"} with self._transaction() as connection: row = self._leased_row( connection, "edge_local_artifacts", "artifact_digest", artifact_digest, lease_token, status="deleting", status_column="cleanup_status", ) if ( not isinstance(receipt, Mapping) or set(receipt) != required or receipt.get("artifact_digest") != artifact_digest or receipt.get("artifact_ref_hash") != row["artifact_ref_hash"] or receipt.get("status") != "deleted" ): raise EdgePolicyError("artifact deletion receipt is invalid") try: deleted_at = canonical_timestamp(receipt["deleted_at"], "deleted_at") except EdgeContractError as exc: raise EdgePolicyError("artifact deletion receipt is invalid") from exc normalized = { "artifact_digest": artifact_digest, "artifact_ref_hash": row["artifact_ref_hash"], "deleted_at": deleted_at, "status": "deleted", } now_text = _timestamp(_utc_now(self._clock)) connection.execute( """UPDATE edge_local_artifacts SET artifact_ref=NULL, cleanup_status='deleted',deletion_receipt_json=?,lease_owner=NULL, lease_token=NULL,lease_expires_at=NULL,error_code=NULL,updated_at=? WHERE artifact_digest=? AND cleanup_status='deleting' AND lease_token=?""", ( strict_json_bytes(normalized).decode("utf-8"), now_text, artifact_digest, lease_token, ), ) updated = connection.execute( "SELECT * FROM edge_local_artifacts WHERE artifact_digest=?", (artifact_digest,), ).fetchone() return self._artifact_record(updated) def fail_artifact_cleanup( self, artifact_digest: str, *, lease_token: str, error_code: str ) -> LocalArtifact: safe_error = self._error_code(error_code) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) row = self._leased_row( connection, "edge_local_artifacts", "artifact_digest", artifact_digest, lease_token, status="deleting", status_column="cleanup_status", ) terminal = row["attempt_count"] >= self.max_attempts status = "failed" if terminal else "pending" available_at = now_text if terminal else _timestamp( now + timedelta(seconds=self._retry_delay(row["attempt_count"])) ) connection.execute( """UPDATE edge_local_artifacts SET cleanup_status=?,available_at=?, lease_owner=NULL,lease_token=NULL,lease_expires_at=NULL,error_code=?, updated_at=? WHERE artifact_digest=? AND cleanup_status='deleting' AND lease_token=?""", ( status, available_at, safe_error, now_text, artifact_digest, lease_token, ), ) updated = connection.execute( "SELECT * FROM edge_local_artifacts WHERE artifact_digest=?", (artifact_digest,), ).fetchone() return self._artifact_record(updated) def put_release_state( self, *, release_id: str, manifest_digest: str, version: str, rollback_version: str, ) -> ReleaseState: release_id = self._bounded_identifier(release_id, "release_id") manifest_digest = self._sha256(manifest_digest, "manifest_digest") version = self._bounded_identifier(version, "version") rollback_version = self._bounded_identifier(rollback_version, "rollback_version") with self._transaction() as connection: existing = connection.execute( "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,) ).fetchone() if existing is not None: if ( existing["manifest_digest"], existing["version"], existing["rollback_version"] ) != (manifest_digest, version, rollback_version): raise EdgeQueueConflictError("release ID has different immutable content") return ReleaseState(**dict(existing)) now_text = _timestamp(_utc_now(self._clock)) connection.execute( """INSERT INTO edge_release_state (release_id,manifest_digest,version,rollback_version,status,updated_at) VALUES (?,?,?,?,'offered',?)""", (release_id, manifest_digest, version, rollback_version, now_text), ) row = connection.execute( "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,) ).fetchone() return ReleaseState(**dict(row)) def compare_and_set_release_state( self, release_id: str, *, expected_status: str, target_status: str, previous_version: str | None = None, current_version: str | None = None, ) -> ReleaseState: release_id = self._bounded_identifier(release_id, "release_id") allowed = {"offered", "accepted", "candidate", "installed", "failed", "rolled_back"} if expected_status not in allowed or target_status not in allowed: raise EdgeQueueConflictError("release status is invalid") transitions = { "offered": {"accepted", "failed"}, "accepted": {"candidate", "failed"}, "candidate": {"installed", "failed", "rolled_back"}, "installed": {"rolled_back", "failed"}, "failed": {"rolled_back"}, "rolled_back": set(), } if target_status not in transitions[expected_status]: raise EdgeQueueConflictError("release state transition is not approved") if previous_version is not None: previous_version = self._bounded_identifier(previous_version, "previous_version") if current_version is not None: current_version = self._bounded_identifier(current_version, "current_version") with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) updated = connection.execute( """UPDATE edge_release_state SET status=?,previous_version=COALESCE(?,previous_version), current_version=COALESCE(?,current_version),updated_at=? WHERE release_id=? AND status=?""", ( target_status, previous_version, current_version, now_text, release_id, expected_status, ), ) if updated.rowcount != 1: row = connection.execute( "SELECT status FROM edge_release_state WHERE release_id=?", (release_id,) ).fetchone() if row is None: raise EdgeQueueError("release does not exist") raise EdgeQueueConflictError("release state compare-and-set lost") row = connection.execute( "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,) ).fetchone() return ReleaseState(**dict(row)) def get_release_state(self, release_id: str) -> ReleaseState | None: release_id = self._bounded_identifier(release_id, "release_id") connection = self._connect() try: row = connection.execute( "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,) ).fetchone() finally: connection.close() return ReleaseState(**dict(row)) if row is not None else None def list_release_states(self) -> tuple[ReleaseState, ...]: connection = self._connect() try: rows = connection.execute( "SELECT * FROM edge_release_state ORDER BY updated_at,release_id" ).fetchall() finally: connection.close() return tuple(ReleaseState(**dict(row)) for row in rows) def complete(self, task_id: str, *, lease_token: str) -> QueuedTask: return self._finish_task( task_id, lease_token=lease_token, target_status="completed", error_code=None, ) def fail( self, task_id: str, *, lease_token: str, error_code: str, retryable: bool = True, ) -> QueuedTask: safe_error = self._error_code(error_code) if not isinstance(retryable, bool): raise ValueError("retryable must be a boolean") with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) row = self._leased_row( connection, "edge_tasks", "task_id", task_id, lease_token ) terminal = not retryable or row["attempt_count"] >= self.max_attempts status = "failed" if terminal else "pending" available_at = now_text if not terminal: available_at = _timestamp( now + timedelta(seconds=self._retry_delay(row["attempt_count"])) ) connection.execute( """ UPDATE edge_tasks SET status = ?, available_at = ?, lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, error_code = ?, updated_at = ? WHERE task_id = ? AND status = 'leased' AND lease_token = ? """, (status, available_at, safe_error, now_text, task_id, lease_token), ) updated = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,) ).fetchone() return self._task_record(updated) def _finish_task( self, task_id: str, *, lease_token: str, target_status: str, error_code: str | None, ) -> QueuedTask: with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) self._leased_row( connection, "edge_tasks", "task_id", task_id, lease_token ) connection.execute( """ UPDATE edge_tasks SET status = ?, lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, error_code = ?, updated_at = ? WHERE task_id = ? AND status = 'leased' AND lease_token = ? """, (target_status, error_code, now_text, task_id, lease_token), ) row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,) ).fetchone() return self._task_record(row) def cancel(self, task_id: str) -> QueuedTask: with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,) ).fetchone() if row is None: raise EdgeQueueError("task does not exist") if row["status"] in {"pending", "leased"}: connection.execute( """ UPDATE edge_tasks SET status = 'cancelled', lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, updated_at = ? WHERE task_id = ? AND status IN ('pending', 'leased') """, (now_text, task_id), ) row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,) ).fetchone() return self._task_record(row) def cancel_with_outcome(self, task_id: str) -> QueuedOutcome: """Atomically fence local work and persist the cancelled control outcome.""" task_id = self._bounded_identifier(task_id, "task_id") with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,) ).fetchone() if row is None: raise EdgeQueueError("task does not exist") if row["status"] in {"pending", "leased"}: connection.execute( """UPDATE edge_tasks SET status='cancelled',lease_owner=NULL, lease_token=NULL,lease_expires_at=NULL,updated_at=? WHERE task_id=? AND status IN ('pending','leased')""", (now_text, task_id), ) row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id=?", (task_id,) ).fetchone() if row["status"] != "cancelled": raise EdgeQueueConflictError("only a cancelled task can emit cancellation") remote_digest = row["remote_lease_token_digest"] if not isinstance(remote_digest, str): raise EdgeQueueLeaseError("cancel outcome requires a remote lease binding") outcome = { "task_id": task_id, "outcome": "cancelled", "safe_summary": {"task_status": "cancelled"}, } digest = canonical_sha256(outcome) encoded = strict_json_bytes(outcome).decode("utf-8") existing = connection.execute( "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,) ).fetchone() if existing is not None: if ( existing["digest"] != digest or existing["remote_lease_token_digest"] != remote_digest ): raise EdgeQueueConflictError("cancel outcome replay conflicts") return self._outcome_record(existing) connection.execute( """INSERT INTO edge_task_outcomes( task_id,digest,outcome_json,status,attempt_count,available_at, remote_lease_token_digest,created_at,updated_at ) VALUES (?,?,?,'pending',0,?,?,?,?)""", (task_id, digest, encoded, now_text, remote_digest, now_text, now_text), ) inserted = connection.execute( "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,) ).fetchone() return self._outcome_record(inserted) def claim_outcome(self, lease_owner: str) -> QueuedOutcome | None: owner = self._lease_owner(lease_owner) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds)) connection.execute( """UPDATE edge_task_outcomes SET status='failed',lease_owner=NULL, lease_token=NULL,lease_expires_at=NULL,error_code='lease_attempts_exhausted', updated_at=? WHERE status='sending' AND julianday(lease_expires_at)<=julianday(?) AND attempt_count>=?""", (now_text, now_text, self.max_attempts), ) row = connection.execute( """SELECT * FROM edge_task_outcomes WHERE attempt_count QueuedOutcome: required = {"task_id", "status", "replayed"} if ( not isinstance(acknowledgement, Mapping) or set(acknowledgement) != required or acknowledgement.get("task_id") != task_id or acknowledgement.get("status") != "cancelled" or not isinstance(acknowledgement.get("replayed"), bool) ): raise EdgePolicyError("cancel outcome acknowledgement is invalid") encoded = strict_json_bytes(acknowledgement).decode("utf-8") ack_digest = canonical_sha256(acknowledgement) with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) self._leased_row( connection, "edge_task_outcomes", "task_id", task_id, lease_token, status="sending", ) connection.execute( """UPDATE edge_task_outcomes SET status='acknowledged', acknowledgement_json=?,ack_digest=?,lease_owner=NULL,lease_token=NULL, lease_expires_at=NULL,updated_at=? WHERE task_id=? AND status='sending' AND lease_token=?""", (encoded, ack_digest, now_text, task_id, lease_token), ) row = connection.execute( "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,) ).fetchone() return self._outcome_record(row) def fail_outcome( self, task_id: str, *, lease_token: str, error_code: str ) -> QueuedOutcome: safe_error = self._error_code(error_code) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) row = self._leased_row( connection, "edge_task_outcomes", "task_id", task_id, lease_token, status="sending", ) terminal = row["attempt_count"] >= self.max_attempts status = "failed" if terminal else "pending" available_at = now_text if terminal else _timestamp( now + timedelta(seconds=self._retry_delay(row["attempt_count"])) ) connection.execute( """UPDATE edge_task_outcomes SET status=?,available_at=?,lease_owner=NULL, lease_token=NULL,lease_expires_at=NULL,error_code=?,updated_at=? WHERE task_id=? AND status='sending' AND lease_token=?""", ( status, available_at, safe_error, now_text, task_id, lease_token, ), ) updated = connection.execute( "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,) ).fetchone() return self._outcome_record(updated) def get_outcome(self, task_id: str) -> QueuedOutcome | None: connection = self._connect() try: row = connection.execute( "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,) ).fetchone() finally: connection.close() return self._outcome_record(row) if row is not None else None def get_task(self, task_id: str) -> QueuedTask | None: connection = self._connect() try: row = connection.execute( "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,) ).fetchone() finally: connection.close() return self._task_record(row) if row is not None else None def persist_event( self, event: EdgeEventContract | Mapping[str, object] ) -> QueuedEvent: raw = event.to_mapping() if isinstance(event, EdgeEventContract) else dict(event) try: contract = EdgeEventContract.from_mapping(raw) except EdgeContractError as exc: raw_id = raw.get("event_id") if isinstance(raw_id, str) and raw_id: with self._transaction() as connection: existing = connection.execute( "SELECT 1 FROM edge_outbound_events WHERE event_id = ?", (raw_id,), ).fetchone() if existing is not None: raise EdgeQueueConflictError( "event ID already exists with a different digest" ) from exc raise normalized = contract.to_mapping() digest = contract.digest with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) existing = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (contract.event_id,), ).fetchone() if existing is not None: if existing["digest"] != digest: raise EdgeQueueConflictError( "event ID already exists with a different digest" ) return self._event_record(existing) approved = self.egress_policy.approve_event(normalized) encoded = strict_json_bytes(approved).decode("utf-8") connection.execute( """ INSERT INTO edge_outbound_events ( event_id, digest, event_json, status, attempt_count, available_at, created_at, updated_at ) VALUES (?, ?, ?, 'pending', 0, ?, ?, ?) """, (contract.event_id, digest, encoded, now_text, now_text, now_text), ) row = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (contract.event_id,), ).fetchone() return self._event_record(row) def claim_event(self, lease_owner: str) -> QueuedEvent | None: owner = self._lease_owner(lease_owner) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) lease_token = str(uuid.uuid4()) expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds)) connection.execute( """ UPDATE edge_outbound_events SET status = 'failed', lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, error_code = 'lease_attempts_exhausted', updated_at = ? WHERE status = 'sending' AND julianday(lease_expires_at) <= julianday(?) AND attempt_count >= ? """, (now_text, now_text, self.max_attempts), ) row = connection.execute( """ SELECT * FROM edge_outbound_events WHERE attempt_count < ? AND ( (status = 'pending' AND julianday(available_at) <= julianday(?)) OR (status = 'sending' AND julianday(lease_expires_at) <= julianday(?)) ) ORDER BY created_at, event_id LIMIT 1 """, (self.max_attempts, now_text, now_text), ).fetchone() if row is None: return None connection.execute( """ UPDATE edge_outbound_events SET status = 'sending', attempt_count = attempt_count + 1, lease_owner = ?, lease_token = ?, lease_expires_at = ?, error_code = NULL, updated_at = ? WHERE event_id = ? AND attempt_count = ? AND ( (status = 'pending' AND julianday(available_at) <= julianday(?)) OR (status = 'sending' AND julianday(lease_expires_at) <= julianday(?)) ) """, ( owner, lease_token, expires_at, now_text, row["event_id"], row["attempt_count"], now_text, now_text, ), ) claimed = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (row["event_id"],), ).fetchone() return self._event_record(claimed) def acknowledge_event( self, event_id: str, *, lease_token: str, acknowledgement: Mapping[str, object], ) -> QueuedEvent: token_digest = canonical_sha256(lease_token) with self._transaction() as connection: now_text = _timestamp(_utc_now(self._clock)) existing = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,) ).fetchone() if existing is None: raise EdgeQueueError("event does not exist") normalized_ack = self._validate_acknowledgement( acknowledgement, event_id=event_id, event_digest=existing["digest"], remote_lease_token_digest=existing["remote_lease_token_digest"], ) encoded = strict_json_bytes(normalized_ack) if len(encoded) > 4_096: raise EdgePolicyError("acknowledgement exceeds the byte limit") ack_digest = canonical_sha256(normalized_ack) if existing["status"] == "acknowledged": if existing["lease_token_digest"] != token_digest: raise EdgeQueueLeaseError( "acknowledgement replay requires the original lease token" ) if existing["ack_digest"] != ack_digest: raise EdgeQueueConflictError( "acknowledgement replay has a different digest" ) return self._event_record(existing) self._leased_row( connection, "edge_outbound_events", "event_id", event_id, lease_token, status="sending", ) connection.execute( """ UPDATE edge_outbound_events SET status = 'acknowledged', acknowledgement_json = ?, ack_digest = ?, lease_token_digest = ?, control_lease_token_digest = ?, lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, updated_at = ? WHERE event_id = ? AND status = 'sending' AND lease_token = ? """, ( encoded.decode("utf-8"), ack_digest, token_digest, token_digest, now_text, event_id, lease_token, ), ) row = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,) ).fetchone() return self._event_record(row) def fail_event( self, event_id: str, *, lease_token: str, error_code: str ) -> QueuedEvent: safe_error = self._error_code(error_code) with self._transaction() as connection: now = _utc_now(self._clock) now_text = _timestamp(now) row = self._leased_row( connection, "edge_outbound_events", "event_id", event_id, lease_token, status="sending", ) terminal = row["attempt_count"] >= self.max_attempts status = "failed" if terminal else "pending" available_at = now_text if not terminal: available_at = _timestamp( now + timedelta(seconds=self._retry_delay(row["attempt_count"])) ) connection.execute( """ UPDATE edge_outbound_events SET status = ?, available_at = ?, lease_owner = NULL, lease_token = NULL, lease_expires_at = NULL, error_code = ?, updated_at = ? WHERE event_id = ? AND status = 'sending' AND lease_token = ? """, (status, available_at, safe_error, now_text, event_id, lease_token), ) updated = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,) ).fetchone() return self._event_record(updated) def get_event(self, event_id: str) -> QueuedEvent | None: connection = self._connect() try: row = connection.execute( "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,) ).fetchone() finally: connection.close() return self._event_record(row) if row is not None else None @staticmethod def _validate_local_artifact(value: Mapping[str, object]) -> dict[str, str]: required = { "artifact_digest", "artifact_ref", "artifact_ref_hash", "classification", "retention_until", } if not isinstance(value, Mapping) or set(value) != required: raise EdgePolicyError("local artifact requires exact metadata fields") digest = value["artifact_digest"] reference = value["artifact_ref"] reference_hash = value["artifact_ref_hash"] classification = value["classification"] if not isinstance(digest, str) or not re.fullmatch(r"[0-9a-f]{64}", digest): raise EdgePolicyError("local artifact digest is invalid") if ( not isinstance(reference, str) or not Path(reference).is_absolute() or len(reference.encode("utf-8")) > 2_048 or not isinstance(reference_hash, str) or not re.fullmatch(r"[0-9a-f]{64}", reference_hash) or reference_hash != canonical_sha256(reference) ): raise EdgePolicyError("local artifact reference is invalid") if classification not in { "raw", "recent_detail", "restricted", "desensitized_metadata", "statistics", "lineage", "evidence", }: raise EdgePolicyError("local artifact classification is invalid") try: retention = canonical_timestamp(value["retention_until"], "retention_until") except EdgeContractError as exc: raise EdgePolicyError("local artifact retention is invalid") from exc return { "artifact_digest": digest, "artifact_ref": reference, # type: ignore[dict-item] "artifact_ref_hash": reference_hash, # type: ignore[dict-item] "classification": classification, # type: ignore[dict-item] "retention_until": retention, } def _normalize_remote_binding( self, token: object, expires_at: object, attempt: object ) -> tuple[str, str | None, str, int]: if ( not isinstance(token, str) or not token.strip() or len(token.encode("utf-8")) > 1024 ): raise EdgeQueueLeaseError("remote lease token is invalid") if isinstance(attempt, bool) or not isinstance(attempt, int) or attempt < 1: raise EdgeQueueLeaseError("remote lease attempt is invalid") try: expiry = canonical_timestamp(expires_at, "remote_lease_expires_at") except EdgeContractError as exc: raise EdgeQueueLeaseError("remote lease expiry is invalid") from exc ciphertext: str | None = None if self._remote_lease_encoder is not None: try: ciphertext = self._remote_lease_encoder(token) except Exception as exc: raise EdgeQueueLeaseError("remote lease protection failed") from exc if ( not isinstance(ciphertext, str) or not ciphertext or ciphertext == token or len(ciphertext.encode("utf-8")) > 4096 ): raise EdgeQueueLeaseError("remote lease protection produced invalid ciphertext") return self._opaque_token_digest(token), ciphertext, expiry, attempt @staticmethod def _opaque_token_digest(token: str) -> str: return hashlib.sha256(token.encode("utf-8")).hexdigest() def _retry_delay(self, attempt_count: int) -> int: base = min( self.retry_max_seconds, self.retry_base_seconds * (2 ** max(0, attempt_count - 1)), ) sample = self._random() if isinstance(sample, bool) or not isinstance(sample, (int, float)) or not 0 <= sample <= 1: raise EdgeQueueError("retry random source returned an invalid value") return min( self.retry_max_seconds, max(1, round(base * (0.5 + float(sample)))), ) @staticmethod def _bounded_identifier(value: object, label: str) -> str: if ( not isinstance(value, str) or not value or value.strip() != value or len(value.encode("utf-8")) > 255 or not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._:-]{0,254}", value) ): raise EdgeQueueError(f"{label} is invalid") return value @staticmethod def _sha256(value: object, label: str) -> str: if not isinstance(value, str) or not re.fullmatch(r"[0-9a-f]{64}", value): raise EdgeQueueError(f"{label} is invalid") return value @staticmethod def _lease_owner(value: object) -> str: if ( not isinstance(value, str) or not value or len(value.encode("utf-8")) > 255 or value.strip() != value ): raise ValueError("lease_owner is invalid") return value @staticmethod def _error_code(value: object) -> str: if not isinstance(value, str) or not _ERROR_CODE.fullmatch(value): raise EdgePolicyError("error_code must be a bounded code") if any(part in value for part in _SENSITIVE_ERROR_PARTS): raise EdgePolicyError("error_code contains sensitive content") return value def _leased_row( self, connection: sqlite3.Connection, table: str, id_column: str, identity: str, lease_token: str, *, status: str = "leased", status_column: str = "status", ) -> sqlite3.Row: row = connection.execute( f"SELECT * FROM {table} WHERE {id_column} = ?", # noqa: S608 (identity,), ).fetchone() if row is None: raise EdgeQueueError("queue item does not exist") if row[status_column] != status or row["lease_token"] != lease_token: raise EdgeQueueLeaseError("current lease token is required") lease_expires_at = row["lease_expires_at"] if lease_expires_at is None: raise EdgeQueueLeaseError("lease has expired") parsed_expiry = datetime.fromisoformat(lease_expires_at.replace("Z", "+00:00")) if parsed_expiry <= _utc_now(self._clock): raise EdgeQueueLeaseError("lease has expired") return row @staticmethod def _validate_acknowledgement( value: Mapping[str, object], *, event_id: str, event_digest: str, remote_lease_token_digest: str | None, ) -> dict[str, str]: if remote_lease_token_digest is not None: required = {"event_id", "event_digest", "remote_lease_digest", "received_at", "status"} if not isinstance(value, Mapping) or set(value) != required: raise EdgePolicyError("bound acknowledgement requires exact approved fields") if value["event_id"] != event_id or value["event_digest"] != event_digest: raise EdgePolicyError("acknowledgement does not bind the event") if value["remote_lease_digest"] != remote_lease_token_digest: raise EdgePolicyError("acknowledgement remote lease is invalid") if value["status"] != "accepted": raise EdgePolicyError("acknowledgement status is invalid") try: received_at = canonical_timestamp(value["received_at"], "received_at") except EdgeContractError as exc: raise EdgePolicyError("acknowledgement received_at is invalid") from exc return { "event_id": event_id, "event_digest": event_digest, "remote_lease_digest": remote_lease_token_digest, "received_at": received_at, "status": "accepted", } if not isinstance(value, Mapping) or set(value) != { "message_id", "received_at", "status", }: raise EdgePolicyError("acknowledgement requires exact approved fields") message_id = value["message_id"] status = value["status"] if ( not isinstance(message_id, str) or not message_id or message_id.strip() != message_id or len(message_id.encode("utf-8")) > 255 ): raise EdgePolicyError("acknowledgement message_id is invalid") if status != "accepted": raise EdgePolicyError("acknowledgement status is invalid") try: received_at = canonical_timestamp(value["received_at"], "received_at") except EdgeContractError as exc: raise EdgePolicyError("acknowledgement received_at is invalid") from exc return { "message_id": message_id, "received_at": received_at, "status": status, } @staticmethod def _task_record(row: sqlite3.Row) -> QueuedTask: task = json.loads(row["task_json"]) return QueuedTask( task_id=row["task_id"], digest=row["digest"], status=row["status"], task=_freeze_mapping(task), attempt_count=row["attempt_count"], available_at=row["available_at"], lease_owner=row["lease_owner"], lease_token=row["lease_token"], lease_expires_at=row["lease_expires_at"], error_code=row["error_code"], authority_key_id=row["authority_key_id"], authority_digest=row["authority_digest"], remote_lease_token_digest=row["remote_lease_token_digest"], remote_lease_expires_at=row["remote_lease_expires_at"], remote_attempt=row["remote_attempt"], created_at=row["created_at"], updated_at=row["updated_at"], ) @staticmethod def _event_record(row: sqlite3.Row) -> QueuedEvent: event = json.loads(row["event_json"]) acknowledgement = ( json.loads(row["acknowledgement_json"]) if row["acknowledgement_json"] is not None else None ) return QueuedEvent( event_id=row["event_id"], digest=row["digest"], status=row["status"], event=_freeze_mapping(event), attempt_count=row["attempt_count"], available_at=row["available_at"], lease_owner=row["lease_owner"], lease_token=row["lease_token"], lease_expires_at=row["lease_expires_at"], error_code=row["error_code"], acknowledgement=( _freeze_mapping(acknowledgement) if acknowledgement is not None else None ), task_digest=row["task_digest"], remote_lease_token_digest=row["remote_lease_token_digest"], control_lease_token_digest=row["control_lease_token_digest"], created_at=row["created_at"], updated_at=row["updated_at"], ) @staticmethod def _artifact_record(row: sqlite3.Row) -> LocalArtifact: receipt = ( json.loads(row["deletion_receipt_json"]) if row["deletion_receipt_json"] is not None else None ) return LocalArtifact( artifact_digest=row["artifact_digest"], artifact_ref=row["artifact_ref"], artifact_ref_hash=row["artifact_ref_hash"], classification=row["classification"], retention_until=row["retention_until"], task_id=row["task_id"], event_id=row["event_id"], cleanup_status=row["cleanup_status"], attempt_count=row["attempt_count"], available_at=row["available_at"], lease_owner=row["lease_owner"], lease_token=row["lease_token"], lease_expires_at=row["lease_expires_at"], error_code=row["error_code"], deletion_receipt=_freeze_mapping(receipt) if receipt is not None else None, created_at=row["created_at"], updated_at=row["updated_at"], ) @staticmethod def _outcome_record(row: sqlite3.Row) -> QueuedOutcome: outcome = json.loads(row["outcome_json"]) acknowledgement = ( json.loads(row["acknowledgement_json"]) if row["acknowledgement_json"] is not None else None ) return QueuedOutcome( task_id=row["task_id"], digest=row["digest"], status=row["status"], outcome=_freeze_mapping(outcome), attempt_count=row["attempt_count"], available_at=row["available_at"], lease_owner=row["lease_owner"], lease_token=row["lease_token"], lease_expires_at=row["lease_expires_at"], error_code=row["error_code"], acknowledgement=( _freeze_mapping(acknowledgement) if acknowledgement is not None else None ), remote_lease_token_digest=row["remote_lease_token_digest"], created_at=row["created_at"], updated_at=row["updated_at"], )