| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189 |
- """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<?
- AND julianday(retention_until)<=julianday(?) AND (
- (cleanup_status='pending' AND julianday(available_at)<=julianday(?))
- OR (cleanup_status='deleting'
- AND julianday(lease_expires_at)<=julianday(?)))
- ORDER BY retention_until,created_at,artifact_digest LIMIT 1""",
- (self.max_attempts, now_text, now_text, now_text),
- ).fetchone()
- if row is None:
- return None
- lease_token = str(uuid.uuid4())
- connection.execute(
- """UPDATE edge_local_artifacts SET cleanup_status='deleting',
- attempt_count=attempt_count+1,lease_owner=?,lease_token=?,
- lease_expires_at=?,error_code=NULL,updated_at=?
- WHERE artifact_digest=? AND attempt_count=? AND (
- (cleanup_status='pending' AND julianday(available_at)<=julianday(?))
- OR (cleanup_status='deleting'
- AND julianday(lease_expires_at)<=julianday(?)))""",
- (
- owner, lease_token, expires_at, now_text, row["artifact_digest"],
- row["attempt_count"], now_text, now_text,
- ),
- )
- claimed = connection.execute(
- "SELECT * FROM edge_local_artifacts WHERE artifact_digest=?",
- (row["artifact_digest"],),
- ).fetchone()
- return self._artifact_record(claimed)
- def acknowledge_artifact_cleanup(
- self,
- artifact_digest: str,
- *,
- lease_token: str,
- receipt: Mapping[str, object],
- ) -> 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<? AND (
- (status='pending' AND julianday(available_at)<=julianday(?)) OR
- (status='sending' AND julianday(lease_expires_at)<=julianday(?)))
- ORDER BY created_at,task_id LIMIT 1""",
- (self.max_attempts, now_text, now_text),
- ).fetchone()
- if row is None:
- return None
- lease_token = str(uuid.uuid4())
- connection.execute(
- """UPDATE edge_task_outcomes SET status='sending',
- attempt_count=attempt_count+1,lease_owner=?,lease_token=?,
- lease_expires_at=?,error_code=NULL,updated_at=?
- WHERE task_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["task_id"],
- row["attempt_count"], now_text, now_text,
- ),
- )
- claimed = connection.execute(
- "SELECT * FROM edge_task_outcomes WHERE task_id=?", (row["task_id"],)
- ).fetchone()
- return self._outcome_record(claimed)
- def acknowledge_outcome(
- self,
- task_id: str,
- *,
- lease_token: str,
- acknowledgement: Mapping[str, object],
- ) -> 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"],
- )
|