queue.py 92 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189
  1. """Durable SQLite task and outbound-event queue for an offline edge gateway."""
  2. from __future__ import annotations
  3. import hashlib
  4. import json
  5. import re
  6. import sqlite3
  7. import uuid
  8. from collections.abc import Callable, Iterator, Mapping
  9. from contextlib import contextmanager
  10. from dataclasses import dataclass
  11. from datetime import UTC, datetime, timedelta
  12. from pathlib import Path
  13. from types import MappingProxyType
  14. from .contracts import (
  15. EdgeContractError,
  16. EdgeEventContract,
  17. EdgeTaskContract,
  18. SignedTaskEnvelope,
  19. canonical_sha256,
  20. canonical_timestamp,
  21. strict_json_bytes,
  22. )
  23. from .policy import EdgeEgressPolicy, EdgePolicyError
  24. class EdgeQueueError(RuntimeError):
  25. """Base error for durable edge queue operations."""
  26. class EdgeQueueConflictError(EdgeQueueError):
  27. """Raised when an existing stable ID is replayed with different content."""
  28. class EdgeQueueLeaseError(EdgeQueueError):
  29. """Raised when a mutation is attempted without the current lease token."""
  30. class EdgeQueueSchemaError(EdgeQueueError):
  31. """Raised when an edge queue database schema cannot be trusted."""
  32. SCHEMA_VERSION = 2
  33. _V1_TASK_TABLE_SQL = """CREATE TABLE edge_tasks (
  34. task_id TEXT PRIMARY KEY,
  35. digest TEXT NOT NULL,
  36. task_json TEXT NOT NULL,
  37. status TEXT NOT NULL CHECK (
  38. status IN ('pending', 'leased', 'completed', 'failed', 'cancelled')
  39. ),
  40. attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0),
  41. available_at TEXT NOT NULL,
  42. deadline_at TEXT NOT NULL,
  43. deadline_epoch_us INTEGER NOT NULL,
  44. lease_owner TEXT,
  45. lease_token TEXT,
  46. lease_expires_at TEXT,
  47. error_code TEXT,
  48. created_at TEXT NOT NULL,
  49. updated_at TEXT NOT NULL
  50. )"""
  51. _TASK_TABLE_SQL = _V1_TASK_TABLE_SQL.rstrip()[:-1].rstrip() + ",\n" + (
  52. " authority_key_id TEXT,\n"
  53. " authority_digest TEXT,\n"
  54. " authority_json TEXT,\n"
  55. " remote_lease_token_digest TEXT,\n"
  56. " remote_lease_ciphertext TEXT,\n"
  57. " remote_lease_expires_at TEXT,\n"
  58. " remote_attempt INTEGER CHECK (remote_attempt IS NULL OR remote_attempt >= 1)\n"
  59. ")"
  60. )
  61. _TASK_INDEX_SQL = """CREATE INDEX edge_tasks_claim_idx ON edge_tasks(
  62. status, deadline_epoch_us, available_at, lease_expires_at, created_at
  63. )"""
  64. _V1_EVENT_TABLE_SQL = """CREATE TABLE edge_outbound_events (
  65. event_id TEXT PRIMARY KEY,
  66. digest TEXT NOT NULL,
  67. event_json TEXT NOT NULL,
  68. status TEXT NOT NULL CHECK (
  69. status IN ('pending', 'sending', 'acknowledged', 'failed')
  70. ),
  71. attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0),
  72. available_at TEXT NOT NULL,
  73. lease_owner TEXT,
  74. lease_token TEXT,
  75. lease_expires_at TEXT,
  76. error_code TEXT,
  77. acknowledgement_json TEXT,
  78. ack_digest TEXT,
  79. lease_token_digest TEXT,
  80. created_at TEXT NOT NULL,
  81. updated_at TEXT NOT NULL
  82. )"""
  83. _EVENT_TABLE_SQL = _V1_EVENT_TABLE_SQL.rstrip()[:-1].rstrip() + ",\n" + (
  84. " task_digest TEXT,\n"
  85. " remote_lease_token_digest TEXT,\n"
  86. " control_lease_token_digest TEXT\n"
  87. ")"
  88. )
  89. _EVENT_INDEX_SQL = """CREATE INDEX edge_events_claim_idx ON edge_outbound_events(
  90. status, available_at, lease_expires_at, created_at
  91. )"""
  92. _ARTIFACT_TABLE_SQL = """CREATE TABLE edge_local_artifacts (
  93. artifact_digest TEXT PRIMARY KEY,
  94. artifact_ref TEXT,
  95. artifact_ref_hash TEXT NOT NULL,
  96. classification TEXT NOT NULL CHECK (
  97. classification IN ('raw', 'recent_detail', 'restricted',
  98. 'desensitized_metadata', 'statistics', 'lineage', 'evidence')
  99. ),
  100. retention_until TEXT NOT NULL,
  101. task_id TEXT NOT NULL,
  102. event_id TEXT NOT NULL,
  103. cleanup_status TEXT NOT NULL CHECK (
  104. cleanup_status IN ('pending', 'deleting', 'deleted', 'failed')
  105. ),
  106. attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0),
  107. available_at TEXT NOT NULL,
  108. lease_owner TEXT,
  109. lease_token TEXT,
  110. lease_expires_at TEXT,
  111. error_code TEXT,
  112. deletion_receipt_json TEXT,
  113. created_at TEXT NOT NULL,
  114. updated_at TEXT NOT NULL,
  115. FOREIGN KEY(task_id) REFERENCES edge_tasks(task_id),
  116. FOREIGN KEY(event_id) REFERENCES edge_outbound_events(event_id)
  117. )"""
  118. _ARTIFACT_INDEX_SQL = """CREATE INDEX edge_artifacts_task_idx ON edge_local_artifacts(
  119. cleanup_status, retention_until, available_at, lease_expires_at, created_at
  120. )"""
  121. _RELEASE_TABLE_SQL = """CREATE TABLE edge_release_state (
  122. release_id TEXT PRIMARY KEY,
  123. manifest_digest TEXT NOT NULL,
  124. version TEXT NOT NULL,
  125. rollback_version TEXT NOT NULL,
  126. status TEXT NOT NULL CHECK (
  127. status IN ('offered', 'accepted', 'candidate', 'installed', 'failed', 'rolled_back')
  128. ),
  129. previous_version TEXT,
  130. current_version TEXT,
  131. updated_at TEXT NOT NULL
  132. )"""
  133. _RELEASE_INDEX_SQL = """CREATE INDEX edge_release_status_idx ON edge_release_state(
  134. status, updated_at, release_id
  135. )"""
  136. _OUTCOME_TABLE_SQL = """CREATE TABLE edge_task_outcomes (
  137. task_id TEXT PRIMARY KEY,
  138. digest TEXT NOT NULL,
  139. outcome_json TEXT NOT NULL,
  140. status TEXT NOT NULL CHECK (
  141. status IN ('pending', 'sending', 'acknowledged', 'failed')
  142. ),
  143. attempt_count INTEGER NOT NULL DEFAULT 0 CHECK (attempt_count >= 0),
  144. available_at TEXT NOT NULL,
  145. lease_owner TEXT,
  146. lease_token TEXT,
  147. lease_expires_at TEXT,
  148. error_code TEXT,
  149. acknowledgement_json TEXT,
  150. ack_digest TEXT,
  151. remote_lease_token_digest TEXT NOT NULL,
  152. created_at TEXT NOT NULL,
  153. updated_at TEXT NOT NULL,
  154. FOREIGN KEY(task_id) REFERENCES edge_tasks(task_id)
  155. )"""
  156. _OUTCOME_INDEX_SQL = """CREATE INDEX edge_outcomes_claim_idx ON edge_task_outcomes(
  157. status, available_at, lease_expires_at, created_at
  158. )"""
  159. _EXPECTED_TABLE_INFO = {
  160. "edge_tasks": (
  161. ("task_id", "TEXT", 0, None, 1),
  162. ("digest", "TEXT", 1, None, 0),
  163. ("task_json", "TEXT", 1, None, 0),
  164. ("status", "TEXT", 1, None, 0),
  165. ("attempt_count", "INTEGER", 1, "0", 0),
  166. ("available_at", "TEXT", 1, None, 0),
  167. ("deadline_at", "TEXT", 1, None, 0),
  168. ("deadline_epoch_us", "INTEGER", 1, None, 0),
  169. ("lease_owner", "TEXT", 0, None, 0),
  170. ("lease_token", "TEXT", 0, None, 0),
  171. ("lease_expires_at", "TEXT", 0, None, 0),
  172. ("error_code", "TEXT", 0, None, 0),
  173. ("created_at", "TEXT", 1, None, 0),
  174. ("updated_at", "TEXT", 1, None, 0),
  175. ("authority_key_id", "TEXT", 0, None, 0),
  176. ("authority_digest", "TEXT", 0, None, 0),
  177. ("authority_json", "TEXT", 0, None, 0),
  178. ("remote_lease_token_digest", "TEXT", 0, None, 0),
  179. ("remote_lease_ciphertext", "TEXT", 0, None, 0),
  180. ("remote_lease_expires_at", "TEXT", 0, None, 0),
  181. ("remote_attempt", "INTEGER", 0, None, 0),
  182. ),
  183. "edge_outbound_events": (
  184. ("event_id", "TEXT", 0, None, 1),
  185. ("digest", "TEXT", 1, None, 0),
  186. ("event_json", "TEXT", 1, None, 0),
  187. ("status", "TEXT", 1, None, 0),
  188. ("attempt_count", "INTEGER", 1, "0", 0),
  189. ("available_at", "TEXT", 1, None, 0),
  190. ("lease_owner", "TEXT", 0, None, 0),
  191. ("lease_token", "TEXT", 0, None, 0),
  192. ("lease_expires_at", "TEXT", 0, None, 0),
  193. ("error_code", "TEXT", 0, None, 0),
  194. ("acknowledgement_json", "TEXT", 0, None, 0),
  195. ("ack_digest", "TEXT", 0, None, 0),
  196. ("lease_token_digest", "TEXT", 0, None, 0),
  197. ("created_at", "TEXT", 1, None, 0),
  198. ("updated_at", "TEXT", 1, None, 0),
  199. ("task_digest", "TEXT", 0, None, 0),
  200. ("remote_lease_token_digest", "TEXT", 0, None, 0),
  201. ("control_lease_token_digest", "TEXT", 0, None, 0),
  202. ),
  203. "edge_local_artifacts": (
  204. ("artifact_digest", "TEXT", 0, None, 1),
  205. ("artifact_ref", "TEXT", 0, None, 0),
  206. ("artifact_ref_hash", "TEXT", 1, None, 0),
  207. ("classification", "TEXT", 1, None, 0),
  208. ("retention_until", "TEXT", 1, None, 0),
  209. ("task_id", "TEXT", 1, None, 0),
  210. ("event_id", "TEXT", 1, None, 0),
  211. ("cleanup_status", "TEXT", 1, None, 0),
  212. ("attempt_count", "INTEGER", 1, "0", 0),
  213. ("available_at", "TEXT", 1, None, 0),
  214. ("lease_owner", "TEXT", 0, None, 0),
  215. ("lease_token", "TEXT", 0, None, 0),
  216. ("lease_expires_at", "TEXT", 0, None, 0),
  217. ("error_code", "TEXT", 0, None, 0),
  218. ("deletion_receipt_json", "TEXT", 0, None, 0),
  219. ("created_at", "TEXT", 1, None, 0),
  220. ("updated_at", "TEXT", 1, None, 0),
  221. ),
  222. "edge_release_state": (
  223. ("release_id", "TEXT", 0, None, 1),
  224. ("manifest_digest", "TEXT", 1, None, 0),
  225. ("version", "TEXT", 1, None, 0),
  226. ("rollback_version", "TEXT", 1, None, 0),
  227. ("status", "TEXT", 1, None, 0),
  228. ("previous_version", "TEXT", 0, None, 0),
  229. ("current_version", "TEXT", 0, None, 0),
  230. ("updated_at", "TEXT", 1, None, 0),
  231. ),
  232. "edge_task_outcomes": (
  233. ("task_id", "TEXT", 0, None, 1),
  234. ("digest", "TEXT", 1, None, 0),
  235. ("outcome_json", "TEXT", 1, None, 0),
  236. ("status", "TEXT", 1, None, 0),
  237. ("attempt_count", "INTEGER", 1, "0", 0),
  238. ("available_at", "TEXT", 1, None, 0),
  239. ("lease_owner", "TEXT", 0, None, 0),
  240. ("lease_token", "TEXT", 0, None, 0),
  241. ("lease_expires_at", "TEXT", 0, None, 0),
  242. ("error_code", "TEXT", 0, None, 0),
  243. ("acknowledgement_json", "TEXT", 0, None, 0),
  244. ("ack_digest", "TEXT", 0, None, 0),
  245. ("remote_lease_token_digest", "TEXT", 1, None, 0),
  246. ("created_at", "TEXT", 1, None, 0),
  247. ("updated_at", "TEXT", 1, None, 0),
  248. ),
  249. }
  250. _EXPECTED_INDEXES = {
  251. "edge_tasks": {
  252. "edge_tasks_claim_idx": (
  253. 0,
  254. "c",
  255. 0,
  256. ("status", "deadline_epoch_us", "available_at", "lease_expires_at", "created_at"),
  257. ),
  258. },
  259. "edge_outbound_events": {
  260. "edge_events_claim_idx": (
  261. 0,
  262. "c",
  263. 0,
  264. ("status", "available_at", "lease_expires_at", "created_at"),
  265. ),
  266. },
  267. "edge_local_artifacts": {
  268. "edge_artifacts_task_idx": (
  269. 0, "c", 0,
  270. ("cleanup_status", "retention_until", "available_at", "lease_expires_at", "created_at"),
  271. ),
  272. },
  273. "edge_release_state": {
  274. "edge_release_status_idx": (0, "c", 0, ("status", "updated_at", "release_id")),
  275. },
  276. "edge_task_outcomes": {
  277. "edge_outcomes_claim_idx": (
  278. 0, "c", 0, ("status", "available_at", "lease_expires_at", "created_at")
  279. ),
  280. },
  281. }
  282. _EXPECTED_TABLE_SQL = {
  283. "edge_tasks": _TASK_TABLE_SQL,
  284. "edge_outbound_events": _EVENT_TABLE_SQL,
  285. "edge_local_artifacts": _ARTIFACT_TABLE_SQL,
  286. "edge_release_state": _RELEASE_TABLE_SQL,
  287. "edge_task_outcomes": _OUTCOME_TABLE_SQL,
  288. }
  289. _V1_EXPECTED_TABLE_INFO = {
  290. table: tuple(row for row in info if row[0] not in {
  291. "authority_key_id", "authority_digest", "authority_json",
  292. "remote_lease_token_digest", "remote_lease_ciphertext", "remote_lease_expires_at", "remote_attempt",
  293. "task_digest", "control_lease_token_digest",
  294. })
  295. for table, info in _EXPECTED_TABLE_INFO.items()
  296. if table not in {"edge_local_artifacts", "edge_release_state", "edge_task_outcomes"}
  297. }
  298. _V1_EXPECTED_TABLE_SQL = {
  299. "edge_tasks": _V1_TASK_TABLE_SQL,
  300. "edge_outbound_events": _V1_EVENT_TABLE_SQL,
  301. }
  302. @dataclass(frozen=True, slots=True)
  303. class QueuedTask:
  304. task_id: str
  305. digest: str
  306. status: str
  307. task: Mapping[str, object]
  308. attempt_count: int
  309. available_at: str
  310. lease_owner: str | None
  311. lease_token: str | None
  312. lease_expires_at: str | None
  313. error_code: str | None
  314. authority_key_id: str | None
  315. authority_digest: str | None
  316. remote_lease_token_digest: str | None
  317. remote_lease_expires_at: str | None
  318. remote_attempt: int | None
  319. created_at: str
  320. updated_at: str
  321. @dataclass(frozen=True, slots=True)
  322. class QueuedEvent:
  323. event_id: str
  324. digest: str
  325. status: str
  326. event: Mapping[str, object]
  327. attempt_count: int
  328. available_at: str
  329. lease_owner: str | None
  330. lease_token: str | None
  331. lease_expires_at: str | None
  332. error_code: str | None
  333. acknowledgement: Mapping[str, object] | None
  334. task_digest: str | None
  335. remote_lease_token_digest: str | None
  336. control_lease_token_digest: str | None
  337. created_at: str
  338. updated_at: str
  339. @dataclass(frozen=True, slots=True)
  340. class LocalArtifact:
  341. artifact_digest: str
  342. artifact_ref: str | None
  343. artifact_ref_hash: str
  344. classification: str
  345. retention_until: str
  346. task_id: str
  347. event_id: str
  348. cleanup_status: str
  349. attempt_count: int
  350. available_at: str
  351. lease_owner: str | None
  352. lease_token: str | None
  353. lease_expires_at: str | None
  354. error_code: str | None
  355. deletion_receipt: Mapping[str, object] | None
  356. created_at: str
  357. updated_at: str
  358. @dataclass(frozen=True, slots=True)
  359. class ReleaseState:
  360. release_id: str
  361. manifest_digest: str
  362. version: str
  363. rollback_version: str
  364. status: str
  365. previous_version: str | None
  366. current_version: str | None
  367. updated_at: str
  368. @dataclass(frozen=True, slots=True)
  369. class QueuedOutcome:
  370. task_id: str
  371. digest: str
  372. status: str
  373. outcome: Mapping[str, object]
  374. attempt_count: int
  375. available_at: str
  376. lease_owner: str | None
  377. lease_token: str | None
  378. lease_expires_at: str | None
  379. error_code: str | None
  380. acknowledgement: Mapping[str, object] | None
  381. remote_lease_token_digest: str
  382. created_at: str
  383. updated_at: str
  384. _ERROR_CODE = re.compile(r"^[a-z][a-z0-9_.-]{0,127}$")
  385. _SENSITIVE_ERROR_PARTS = ("password", "passwd", "secret", "token", "credential")
  386. def _utc_now(clock: Callable[[], object]) -> datetime:
  387. value = clock()
  388. if isinstance(value, bool):
  389. raise ValueError("queue clock returned an invalid value")
  390. if isinstance(value, (int, float)):
  391. value = datetime.fromtimestamp(value, tz=UTC)
  392. if not isinstance(value, datetime):
  393. raise ValueError("queue clock must return a datetime or timestamp")
  394. if value.tzinfo is None or value.utcoffset() is None:
  395. raise ValueError("queue clock must return a timezone-aware value")
  396. return value.astimezone(UTC)
  397. def _timestamp(value: datetime) -> str:
  398. utc = value.astimezone(UTC)
  399. canonical = utc.strftime("%Y-%m-%dT%H:%M:%S")
  400. if utc.microsecond:
  401. canonical += f".{utc.microsecond:06d}".rstrip("0")
  402. return f"{canonical}Z"
  403. def _epoch_microseconds(value: datetime) -> int:
  404. epoch = datetime(1970, 1, 1, tzinfo=UTC)
  405. delta = value.astimezone(UTC) - epoch
  406. return (
  407. (delta.days * 86_400 + delta.seconds) * 1_000_000
  408. + delta.microseconds
  409. )
  410. def _normalize_schema_sql(value: str | None) -> str:
  411. if not isinstance(value, str):
  412. return ""
  413. normalized = re.sub(r"\s+", " ", value.strip()).lower()
  414. return re.sub(r"\s*([(),])\s*", r"\1", normalized)
  415. def _freeze_mapping(value: Mapping[str, object]) -> Mapping[str, object]:
  416. def freeze(item: object) -> object:
  417. if isinstance(item, dict):
  418. return MappingProxyType({key: freeze(child) for key, child in item.items()})
  419. if isinstance(item, list):
  420. return tuple(freeze(child) for child in item)
  421. return item
  422. return freeze(dict(value)) # type: ignore[return-value]
  423. class SqliteEdgeQueue:
  424. """An on-disk, WAL-backed queue with fenced leases and exact replay."""
  425. def __init__(
  426. self,
  427. db_path: str | Path,
  428. *,
  429. clock: Callable[[], object] | None = None,
  430. lease_seconds: int = 60,
  431. max_attempts: int = 5,
  432. retry_base_seconds: int = 2,
  433. retry_max_seconds: int = 300,
  434. busy_timeout_ms: int = 5_000,
  435. egress_policy: EdgeEgressPolicy | None = None,
  436. random_source: Callable[[], float] | None = None,
  437. remote_lease_encoder: Callable[[str], str] | None = None,
  438. remote_lease_decoder: Callable[[str], str] | None = None,
  439. ) -> None:
  440. if not isinstance(db_path, (str, Path)) or not str(db_path).strip():
  441. raise ValueError("an explicit SQLite database path is required")
  442. path_text = str(db_path)
  443. if path_text == ":memory:" or "mode=memory" in path_text:
  444. raise ValueError("an in-memory database is not allowed")
  445. for label, value in {
  446. "lease_seconds": lease_seconds,
  447. "max_attempts": max_attempts,
  448. "retry_base_seconds": retry_base_seconds,
  449. "retry_max_seconds": retry_max_seconds,
  450. "busy_timeout_ms": busy_timeout_ms,
  451. }.items():
  452. if isinstance(value, bool) or not isinstance(value, int) or value < 1:
  453. raise ValueError(f"{label} must be a positive integer")
  454. if retry_base_seconds > retry_max_seconds:
  455. raise ValueError("retry_base_seconds cannot exceed retry_max_seconds")
  456. self.db_path = str(Path(db_path).expanduser().resolve())
  457. self._clock = clock or (lambda: datetime.now(UTC))
  458. self.lease_seconds = lease_seconds
  459. self.max_attempts = max_attempts
  460. self.retry_base_seconds = retry_base_seconds
  461. self.retry_max_seconds = retry_max_seconds
  462. self.busy_timeout_ms = busy_timeout_ms
  463. self.egress_policy = egress_policy or EdgeEgressPolicy(
  464. allowed_control_hosts=set()
  465. )
  466. # The injectable source makes retry timing testable; 0.5 is neutral jitter.
  467. self._random = random_source or (lambda: 0.5)
  468. if (remote_lease_encoder is None) != (remote_lease_decoder is None):
  469. raise ValueError("remote lease encoder and decoder must be paired")
  470. self._remote_lease_encoder = remote_lease_encoder
  471. self._remote_lease_decoder = remote_lease_decoder
  472. Path(self.db_path).parent.mkdir(parents=True, exist_ok=True)
  473. self._initialize()
  474. def _connect(self) -> sqlite3.Connection:
  475. connection = sqlite3.connect(
  476. self.db_path,
  477. timeout=self.busy_timeout_ms / 1_000,
  478. isolation_level=None,
  479. )
  480. connection.row_factory = sqlite3.Row
  481. connection.execute(f"PRAGMA busy_timeout = {self.busy_timeout_ms}")
  482. connection.execute("PRAGMA foreign_keys = ON")
  483. return connection
  484. def _initialize(self) -> None:
  485. connection = self._connect()
  486. try:
  487. version = connection.execute("PRAGMA user_version").fetchone()[0]
  488. edge_tables = {
  489. row[0]
  490. for row in connection.execute(
  491. """
  492. SELECT name FROM sqlite_master
  493. WHERE type = 'table' AND name LIKE 'edge_%'
  494. """
  495. ).fetchall()
  496. }
  497. if version == 0 and edge_tables:
  498. raise EdgeQueueSchemaError("unversioned edge queue schema is unsafe")
  499. if version not in {0, 1, SCHEMA_VERSION}:
  500. raise EdgeQueueSchemaError("edge queue schema version is unsupported")
  501. connection.execute("PRAGMA journal_mode = WAL")
  502. connection.execute("PRAGMA synchronous = FULL")
  503. connection.execute("BEGIN IMMEDIATE")
  504. if version == 0:
  505. for statement in (
  506. _TASK_TABLE_SQL,
  507. _TASK_INDEX_SQL,
  508. _EVENT_TABLE_SQL,
  509. _EVENT_INDEX_SQL,
  510. _ARTIFACT_TABLE_SQL,
  511. _ARTIFACT_INDEX_SQL,
  512. _RELEASE_TABLE_SQL,
  513. _RELEASE_INDEX_SQL,
  514. _OUTCOME_TABLE_SQL,
  515. _OUTCOME_INDEX_SQL,
  516. ):
  517. connection.execute(statement)
  518. connection.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
  519. elif version == 1:
  520. self._validate_schema(
  521. connection,
  522. expected_info=_V1_EXPECTED_TABLE_INFO,
  523. expected_sql=_V1_EXPECTED_TABLE_SQL,
  524. expected_indexes={
  525. key: _EXPECTED_INDEXES[key]
  526. for key in ("edge_tasks", "edge_outbound_events")
  527. },
  528. )
  529. active_unsigned = connection.execute(
  530. """SELECT EXISTS(
  531. SELECT 1 FROM edge_tasks WHERE status IN ('pending','leased')
  532. )"""
  533. ).fetchone()[0]
  534. unacknowledged_events = connection.execute(
  535. """SELECT EXISTS(
  536. SELECT 1 FROM edge_outbound_events WHERE status <> 'acknowledged'
  537. )"""
  538. ).fetchone()[0]
  539. if active_unsigned or unacknowledged_events:
  540. raise EdgeQueueSchemaError(
  541. "v1 edge queue must drain active tasks and unacknowledged events"
  542. )
  543. for statement in (
  544. "ALTER TABLE edge_tasks ADD COLUMN authority_key_id TEXT",
  545. "ALTER TABLE edge_tasks ADD COLUMN authority_digest TEXT",
  546. "ALTER TABLE edge_tasks ADD COLUMN authority_json TEXT",
  547. "ALTER TABLE edge_tasks ADD COLUMN remote_lease_token_digest TEXT",
  548. "ALTER TABLE edge_tasks ADD COLUMN remote_lease_ciphertext TEXT",
  549. "ALTER TABLE edge_tasks ADD COLUMN remote_lease_expires_at TEXT",
  550. "ALTER TABLE edge_tasks ADD COLUMN remote_attempt INTEGER CHECK (remote_attempt IS NULL OR remote_attempt >= 1)",
  551. "ALTER TABLE edge_outbound_events ADD COLUMN task_digest TEXT",
  552. "ALTER TABLE edge_outbound_events ADD COLUMN remote_lease_token_digest TEXT",
  553. "ALTER TABLE edge_outbound_events ADD COLUMN control_lease_token_digest TEXT",
  554. _ARTIFACT_TABLE_SQL,
  555. _ARTIFACT_INDEX_SQL,
  556. _RELEASE_TABLE_SQL,
  557. _RELEASE_INDEX_SQL,
  558. _OUTCOME_TABLE_SQL,
  559. _OUTCOME_INDEX_SQL,
  560. ):
  561. connection.execute(statement)
  562. connection.execute(f"PRAGMA user_version = {SCHEMA_VERSION}")
  563. self._validate_schema(connection)
  564. if connection.in_transaction:
  565. connection.execute("COMMIT")
  566. except Exception:
  567. if connection.in_transaction:
  568. connection.execute("ROLLBACK")
  569. raise
  570. finally:
  571. connection.close()
  572. @staticmethod
  573. def _validate_schema(
  574. connection: sqlite3.Connection,
  575. *,
  576. expected_info: Mapping[str, tuple] = _EXPECTED_TABLE_INFO,
  577. expected_sql: Mapping[str, str] = _EXPECTED_TABLE_SQL,
  578. expected_indexes: Mapping[str, Mapping[str, tuple]] = _EXPECTED_INDEXES,
  579. ) -> None:
  580. primary_keys = {
  581. "edge_tasks": "task_id",
  582. "edge_outbound_events": "event_id",
  583. "edge_local_artifacts": "artifact_digest",
  584. "edge_release_state": "release_id",
  585. "edge_task_outcomes": "task_id",
  586. }
  587. for table, table_expected_info in expected_info.items():
  588. rows = connection.execute(f"PRAGMA table_info({table})").fetchall()
  589. actual_info = tuple(
  590. (row[1], row[2].upper(), row[3], row[4], row[5]) for row in rows
  591. )
  592. if actual_info != table_expected_info:
  593. raise EdgeQueueSchemaError("edge queue schema columns are invalid")
  594. table_sql = connection.execute(
  595. "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = ?",
  596. (table,),
  597. ).fetchone()
  598. if table_sql is None or _normalize_schema_sql(
  599. table_sql[0]
  600. ) != _normalize_schema_sql(expected_sql[table]):
  601. raise EdgeQueueSchemaError(
  602. "edge queue schema constraints are invalid"
  603. )
  604. index_rows = connection.execute(f"PRAGMA index_list({table})").fetchall()
  605. custom_rows = {row[1]: row for row in index_rows if row[3] == "c"}
  606. table_expected_indexes = expected_indexes[table]
  607. if set(custom_rows) != set(table_expected_indexes):
  608. raise EdgeQueueSchemaError("edge queue schema indexes are invalid")
  609. for name, (unique, origin, partial, columns) in table_expected_indexes.items():
  610. row = custom_rows[name]
  611. actual_columns = tuple(
  612. item[2]
  613. for item in connection.execute(
  614. f"PRAGMA index_info({name})"
  615. ).fetchall()
  616. )
  617. if (row[2], row[3], row[4], actual_columns) != (
  618. unique,
  619. origin,
  620. partial,
  621. columns,
  622. ):
  623. raise EdgeQueueSchemaError(
  624. "edge queue schema index definition is invalid"
  625. )
  626. primary_indexes = [row for row in index_rows if row[3] == "pk"]
  627. if len(primary_indexes) != 1:
  628. raise EdgeQueueSchemaError("edge queue schema primary key is invalid")
  629. primary = primary_indexes[0]
  630. primary_columns = tuple(
  631. item[2]
  632. for item in connection.execute(
  633. f"PRAGMA index_info({primary[1]})"
  634. ).fetchall()
  635. )
  636. if (primary[2], primary[4], primary_columns) != (
  637. 1,
  638. 0,
  639. (primary_keys[table],),
  640. ):
  641. raise EdgeQueueSchemaError("edge queue schema primary key is invalid")
  642. @contextmanager
  643. def _transaction(self) -> Iterator[sqlite3.Connection]:
  644. connection = self._connect()
  645. try:
  646. connection.execute("BEGIN IMMEDIATE")
  647. yield connection
  648. connection.execute("COMMIT")
  649. except Exception:
  650. if connection.in_transaction:
  651. connection.execute("ROLLBACK")
  652. raise
  653. finally:
  654. connection.close()
  655. def enqueue(self, task: EdgeTaskContract | Mapping[str, object]) -> QueuedTask:
  656. contract = EdgeTaskContract.from_mapping(
  657. task.to_mapping() if isinstance(task, EdgeTaskContract) else task
  658. )
  659. return self._enqueue_contract(contract, authority=None)
  660. def enqueue_signed(self, envelope: SignedTaskEnvelope) -> QueuedTask:
  661. """Atomically persist a previously verified authority envelope and task."""
  662. if not isinstance(envelope, SignedTaskEnvelope) or not envelope.is_verified:
  663. raise EdgeContractError("a verified SignedTaskEnvelope is required")
  664. # Reconstruct the task so object.__setattr__ cannot bypass validation.
  665. contract = EdgeTaskContract.from_mapping(envelope.task.to_mapping())
  666. if envelope.contract_digest != contract.digest:
  667. raise EdgeContractError("signed authority no longer binds the task")
  668. return self._enqueue_contract(contract, authority=envelope)
  669. def accept_signed_task(
  670. self,
  671. envelope: SignedTaskEnvelope,
  672. *,
  673. remote_lease_token: str,
  674. remote_lease_expires_at: str,
  675. remote_attempt: int,
  676. ) -> QueuedTask:
  677. """Atomically accept authority and an encrypted recoverable control lease."""
  678. if self._remote_lease_encoder is None:
  679. raise EdgeQueueLeaseError("secure remote lease codec is required")
  680. if not isinstance(envelope, SignedTaskEnvelope) or not envelope.is_verified:
  681. raise EdgeContractError("a verified SignedTaskEnvelope is required")
  682. return self._enqueue_contract(
  683. EdgeTaskContract.from_mapping(envelope.task.to_mapping()),
  684. authority=envelope,
  685. remote_binding=self._normalize_remote_binding(
  686. remote_lease_token, remote_lease_expires_at, remote_attempt
  687. ),
  688. )
  689. def _enqueue_contract(
  690. self,
  691. contract: EdgeTaskContract,
  692. *,
  693. authority: SignedTaskEnvelope | None,
  694. remote_binding: tuple[str, str | None, str, int] | None = None,
  695. ) -> QueuedTask:
  696. task_mapping = contract.to_mapping()
  697. encoded = strict_json_bytes(task_mapping).decode("utf-8")
  698. digest = canonical_sha256(task_mapping)
  699. if remote_binding is not None and remote_binding[3] != contract.attempt:
  700. raise EdgeQueueLeaseError(
  701. "remote lease attempt does not bind the task contract"
  702. )
  703. with self._transaction() as connection:
  704. existing = connection.execute(
  705. "SELECT * FROM edge_tasks WHERE task_id = ?", (contract.task_id,)
  706. ).fetchone()
  707. if existing is not None:
  708. if existing["digest"] != digest:
  709. raise EdgeQueueConflictError(
  710. "task ID already exists with a different digest"
  711. )
  712. if authority is not None and existing["authority_digest"] != authority.digest:
  713. raise EdgeQueueConflictError(
  714. "task ID already exists with a different authority digest"
  715. )
  716. if remote_binding is not None:
  717. now = _utc_now(self._clock)
  718. old_expiry_text = existing["remote_lease_expires_at"]
  719. old_expiry = (
  720. datetime.fromisoformat(old_expiry_text.replace("Z", "+00:00"))
  721. if old_expiry_text is not None else None
  722. )
  723. if (
  724. existing["remote_lease_token_digest"] != remote_binding[0]
  725. and old_expiry is not None and old_expiry > now
  726. ):
  727. raise EdgeQueueConflictError("active remote lease binding conflicts")
  728. new_expiry = datetime.fromisoformat(remote_binding[2].replace("Z", "+00:00"))
  729. if new_expiry <= now:
  730. raise EdgeQueueLeaseError("remote lease has expired")
  731. if (
  732. existing["status"] in {"failed", "cancelled"}
  733. and not (
  734. existing["status"] == "failed"
  735. and existing["error_code"] == "execution_lease_expired"
  736. )
  737. ):
  738. raise EdgeQueueConflictError(
  739. "terminal task cannot be rebound to a remote lease"
  740. )
  741. changed_remote_lease = (
  742. existing["remote_lease_token_digest"] != remote_binding[0]
  743. )
  744. if changed_remote_lease:
  745. event_rows = connection.execute(
  746. """SELECT * FROM edge_outbound_events
  747. WHERE status IN ('pending','sending')
  748. AND json_extract(event_json,'$.task_id')=?""",
  749. (contract.task_id,),
  750. ).fetchall()
  751. for event_row in event_rows:
  752. event_mapping = json.loads(event_row["event_json"])
  753. if (
  754. event_row["task_digest"] != digest
  755. or event_row["remote_lease_token_digest"]
  756. != existing["remote_lease_token_digest"]
  757. or event_mapping.get("task_id") != contract.task_id
  758. or event_mapping.get("gateway_id") != contract.gateway_id
  759. or canonical_sha256(event_mapping) != event_row["digest"]
  760. ):
  761. raise EdgeQueueConflictError(
  762. "pending event does not bind the task lease"
  763. )
  764. connection.execute(
  765. """UPDATE edge_outbound_events
  766. SET remote_lease_token_digest=?,updated_at=?
  767. WHERE status IN ('pending','sending')
  768. AND json_extract(event_json,'$.task_id')=?""",
  769. (
  770. remote_binding[0],
  771. _timestamp(now),
  772. contract.task_id,
  773. ),
  774. )
  775. recoverable_failure = (
  776. existing["status"] == "failed"
  777. and existing["error_code"] == "execution_lease_expired"
  778. )
  779. connection.execute(
  780. """UPDATE edge_tasks SET remote_lease_token_digest=?,
  781. remote_lease_ciphertext=?,remote_lease_expires_at=?,remote_attempt=?,
  782. status=CASE WHEN ? THEN 'pending' ELSE status END,
  783. available_at=CASE WHEN ? THEN ? ELSE available_at END,
  784. error_code=CASE WHEN ? THEN NULL ELSE error_code END,
  785. updated_at=?
  786. WHERE task_id=? AND digest=? AND authority_digest=?""",
  787. (
  788. remote_binding[0], remote_binding[1], remote_binding[2],
  789. remote_binding[3], recoverable_failure,
  790. recoverable_failure, _timestamp(now), recoverable_failure,
  791. _timestamp(now), contract.task_id,
  792. digest, authority.digest if authority is not None else None,
  793. ),
  794. )
  795. existing = connection.execute(
  796. "SELECT * FROM edge_tasks WHERE task_id=?", (contract.task_id,)
  797. ).fetchone()
  798. return self._task_record(existing)
  799. now = _utc_now(self._clock)
  800. deadline = datetime.fromisoformat(
  801. contract.deadline_at.replace("Z", "+00:00")
  802. ).astimezone(UTC)
  803. if deadline <= now:
  804. raise EdgeContractError("task deadline has expired")
  805. if remote_binding is not None:
  806. remote_expiry = datetime.fromisoformat(remote_binding[2].replace("Z", "+00:00"))
  807. if remote_expiry <= now or remote_expiry > deadline:
  808. raise EdgeQueueLeaseError("remote lease validity window is invalid")
  809. now_text = _timestamp(now)
  810. connection.execute(
  811. """
  812. INSERT INTO edge_tasks (
  813. task_id, digest, task_json, status, attempt_count,
  814. available_at, deadline_at, deadline_epoch_us,
  815. authority_key_id, authority_digest, authority_json,
  816. remote_lease_token_digest, remote_lease_ciphertext,
  817. remote_lease_expires_at, remote_attempt,
  818. created_at, updated_at
  819. ) VALUES (?, ?, ?, 'pending', 0, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
  820. """,
  821. (
  822. contract.task_id,
  823. digest,
  824. encoded,
  825. now_text,
  826. contract.deadline_at,
  827. _epoch_microseconds(deadline),
  828. authority.authority_key_id if authority is not None else None,
  829. authority.digest if authority is not None else None,
  830. strict_json_bytes(authority.to_mapping()).decode("utf-8")
  831. if authority is not None else None,
  832. remote_binding[0] if remote_binding is not None else None,
  833. remote_binding[1] if remote_binding is not None else None,
  834. remote_binding[2] if remote_binding is not None else None,
  835. remote_binding[3] if remote_binding is not None else None,
  836. now_text,
  837. now_text,
  838. ),
  839. )
  840. row = connection.execute(
  841. "SELECT * FROM edge_tasks WHERE task_id = ?", (contract.task_id,)
  842. ).fetchone()
  843. return self._task_record(row)
  844. def claim(self, lease_owner: str) -> QueuedTask | None:
  845. owner = self._lease_owner(lease_owner)
  846. with self._transaction() as connection:
  847. now = _utc_now(self._clock)
  848. now_text = _timestamp(now)
  849. now_epoch_us = _epoch_microseconds(now)
  850. lease_token = str(uuid.uuid4())
  851. expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds))
  852. connection.execute(
  853. """
  854. UPDATE edge_tasks
  855. SET status = 'failed', lease_owner = NULL, lease_token = NULL,
  856. lease_expires_at = NULL,
  857. error_code = 'task_deadline_expired', updated_at = ?
  858. WHERE status IN ('pending', 'leased') AND deadline_epoch_us <= ?
  859. """,
  860. (now_text, now_epoch_us),
  861. )
  862. connection.execute(
  863. """
  864. UPDATE edge_tasks
  865. SET status = 'failed', lease_owner = NULL, lease_token = NULL,
  866. lease_expires_at = NULL,
  867. error_code = 'lease_attempts_exhausted', updated_at = ?
  868. WHERE status = 'leased'
  869. AND julianday(lease_expires_at) <= julianday(?)
  870. AND attempt_count >= ?
  871. """,
  872. (now_text, now_text, self.max_attempts),
  873. )
  874. row = connection.execute(
  875. """
  876. SELECT * FROM edge_tasks
  877. WHERE deadline_epoch_us > ? AND attempt_count < ? AND (
  878. (status = 'pending' AND julianday(available_at) <= julianday(?))
  879. OR (status = 'leased'
  880. AND julianday(lease_expires_at) <= julianday(?))
  881. )
  882. ORDER BY created_at, task_id
  883. LIMIT 1
  884. """,
  885. (now_epoch_us, self.max_attempts, now_text, now_text),
  886. ).fetchone()
  887. if row is None:
  888. return None
  889. updated = connection.execute(
  890. """
  891. UPDATE edge_tasks
  892. SET status = 'leased', attempt_count = attempt_count + 1,
  893. lease_owner = ?, lease_token = ?, lease_expires_at = ?,
  894. error_code = NULL, updated_at = ?
  895. WHERE task_id = ? AND deadline_epoch_us > ?
  896. AND attempt_count = ? AND (
  897. (status = 'pending' AND julianday(available_at) <= julianday(?))
  898. OR (status = 'leased'
  899. AND julianday(lease_expires_at) <= julianday(?))
  900. )
  901. """,
  902. (
  903. owner,
  904. lease_token,
  905. expires_at,
  906. now_text,
  907. row["task_id"],
  908. now_epoch_us,
  909. row["attempt_count"],
  910. now_text,
  911. now_text,
  912. ),
  913. )
  914. if updated.rowcount != 1:
  915. return None
  916. claimed = connection.execute(
  917. "SELECT * FROM edge_tasks WHERE task_id = ?", (row["task_id"],)
  918. ).fetchone()
  919. return self._task_record(claimed)
  920. def bind_remote_lease(
  921. self,
  922. task_id: str,
  923. *,
  924. lease_token: str,
  925. remote_lease_token: str,
  926. remote_lease_expires_at: str,
  927. remote_attempt: int,
  928. ) -> QueuedTask:
  929. """Bind control-plane authority to the signed contract attempt."""
  930. digest, ciphertext, expiry_text, remote_attempt = self._normalize_remote_binding(
  931. remote_lease_token, remote_lease_expires_at, remote_attempt
  932. )
  933. expiry = datetime.fromisoformat(expiry_text.replace("Z", "+00:00"))
  934. with self._transaction() as connection:
  935. row = self._leased_row(connection, "edge_tasks", "task_id", task_id, lease_token)
  936. now = _utc_now(self._clock)
  937. if expiry <= now:
  938. raise EdgeQueueLeaseError("remote lease has expired")
  939. task_mapping = json.loads(row["task_json"])
  940. if remote_attempt != task_mapping["attempt"]:
  941. raise EdgeQueueLeaseError("remote lease attempt does not bind task contract")
  942. existing = row["remote_lease_token_digest"]
  943. if existing is not None and (
  944. existing != digest
  945. or row["remote_lease_expires_at"] != expiry_text
  946. or row["remote_attempt"] != remote_attempt
  947. ):
  948. raise EdgeQueueConflictError("remote lease binding conflicts")
  949. connection.execute(
  950. """
  951. UPDATE edge_tasks
  952. SET remote_lease_token_digest=?, remote_lease_ciphertext=?, remote_lease_expires_at=?,
  953. remote_attempt=?, updated_at=?
  954. WHERE task_id=? AND status='leased' AND lease_token=?
  955. """,
  956. (digest, ciphertext, expiry_text, remote_attempt, _timestamp(now), task_id, lease_token),
  957. )
  958. updated = connection.execute(
  959. "SELECT * FROM edge_tasks WHERE task_id=?", (task_id,)
  960. ).fetchone()
  961. return self._task_record(updated)
  962. def recover_remote_lease(self, task_id: str) -> str:
  963. """Recover a control lease through the injected secure decoder without exposing it in records."""
  964. if self._remote_lease_decoder is None:
  965. raise EdgeQueueLeaseError("secure remote lease codec is required")
  966. connection = self._connect()
  967. try:
  968. row = connection.execute(
  969. "SELECT remote_lease_ciphertext,remote_lease_token_digest FROM edge_tasks WHERE task_id=?",
  970. (task_id,),
  971. ).fetchone()
  972. finally:
  973. connection.close()
  974. if row is None or row["remote_lease_ciphertext"] is None:
  975. raise EdgeQueueLeaseError("recoverable remote lease does not exist")
  976. try:
  977. token = self._remote_lease_decoder(row["remote_lease_ciphertext"])
  978. except Exception as exc:
  979. raise EdgeQueueLeaseError("remote lease recovery failed") from exc
  980. if not isinstance(token, str) or self._opaque_token_digest(token) != row["remote_lease_token_digest"]:
  981. raise EdgeQueueLeaseError("recovered remote lease failed integrity validation")
  982. return token
  983. def complete_with_event(
  984. self,
  985. task_id: str,
  986. *,
  987. lease_token: str,
  988. remote_lease_token: str,
  989. event: EdgeEventContract | Mapping[str, object],
  990. artifacts: list[Mapping[str, object]] | tuple[Mapping[str, object], ...] = (),
  991. ) -> tuple[QueuedTask, QueuedEvent]:
  992. """Atomically commit a safe event, local artifact metadata, and task terminal state."""
  993. contract = EdgeEventContract.from_mapping(
  994. event.to_mapping() if isinstance(event, EdgeEventContract) else event
  995. )
  996. approved = self.egress_policy.approve_event(contract.to_mapping())
  997. encoded_event = strict_json_bytes(approved).decode("utf-8")
  998. normalized_artifacts = [self._validate_local_artifact(item) for item in artifacts]
  999. if not isinstance(remote_lease_token, str) or not remote_lease_token:
  1000. raise EdgeQueueLeaseError("remote lease token is invalid")
  1001. remote_digest = self._opaque_token_digest(remote_lease_token)
  1002. with self._transaction() as connection:
  1003. now = _utc_now(self._clock)
  1004. now_text = _timestamp(now)
  1005. row = self._leased_row(connection, "edge_tasks", "task_id", task_id, lease_token)
  1006. if row["deadline_epoch_us"] <= _epoch_microseconds(now):
  1007. raise EdgeQueueLeaseError("task deadline has expired")
  1008. if row["remote_lease_token_digest"] != remote_digest:
  1009. raise EdgeQueueLeaseError("current remote lease token is required")
  1010. remote_expiry = row["remote_lease_expires_at"]
  1011. if remote_expiry is None or datetime.fromisoformat(
  1012. remote_expiry.replace("Z", "+00:00")
  1013. ) <= now:
  1014. raise EdgeQueueLeaseError("remote lease has expired")
  1015. task_mapping = json.loads(row["task_json"])
  1016. if row["remote_attempt"] != task_mapping["attempt"]:
  1017. raise EdgeQueueLeaseError("remote lease attempt is stale")
  1018. for field in ("task_id", "gateway_id", "environment", "network_zone", "purpose", "policy_digest"):
  1019. if contract.to_mapping()[field] != task_mapping[field]:
  1020. raise EdgeQueueConflictError(f"event {field} does not bind the task")
  1021. if contract.attempt != task_mapping["attempt"]:
  1022. raise EdgeQueueConflictError("event attempt does not bind the task")
  1023. existing = connection.execute(
  1024. "SELECT digest FROM edge_outbound_events WHERE event_id=?", (contract.event_id,)
  1025. ).fetchone()
  1026. if existing is not None:
  1027. raise EdgeQueueConflictError("event ID already exists before task completion")
  1028. connection.execute(
  1029. """
  1030. INSERT INTO edge_outbound_events (
  1031. event_id,digest,event_json,status,attempt_count,available_at,
  1032. task_digest,remote_lease_token_digest,created_at,updated_at
  1033. ) VALUES (?,?,?,'pending',0,?,?,?,?,?)
  1034. """,
  1035. (
  1036. contract.event_id, contract.digest, encoded_event, now_text,
  1037. row["digest"], remote_digest, now_text, now_text,
  1038. ),
  1039. )
  1040. for artifact in normalized_artifacts:
  1041. connection.execute(
  1042. """
  1043. INSERT INTO edge_local_artifacts (
  1044. artifact_digest,artifact_ref,artifact_ref_hash,classification,
  1045. retention_until,task_id,event_id,cleanup_status,attempt_count,
  1046. available_at,created_at,updated_at
  1047. ) VALUES (?,?,?,?,?,?,?,'pending',0,?,?,?)
  1048. """,
  1049. (
  1050. artifact["artifact_digest"], artifact["artifact_ref"],
  1051. artifact["artifact_ref_hash"], artifact["classification"],
  1052. artifact["retention_until"], task_id, contract.event_id,
  1053. now_text, now_text, now_text,
  1054. ),
  1055. )
  1056. updated = connection.execute(
  1057. """
  1058. UPDATE edge_tasks SET status='completed', lease_owner=NULL,
  1059. lease_token=NULL, lease_expires_at=NULL, error_code=NULL, updated_at=?
  1060. WHERE task_id=? AND status='leased' AND lease_token=?
  1061. """,
  1062. (now_text, task_id, lease_token),
  1063. )
  1064. if updated.rowcount != 1:
  1065. raise EdgeQueueLeaseError("task completion lost its lease")
  1066. task_row = connection.execute("SELECT * FROM edge_tasks WHERE task_id=?", (task_id,)).fetchone()
  1067. event_row = connection.execute("SELECT * FROM edge_outbound_events WHERE event_id=?", (contract.event_id,)).fetchone()
  1068. return self._task_record(task_row), self._event_record(event_row)
  1069. def list_local_artifacts(self, *, task_id: str) -> tuple[LocalArtifact, ...]:
  1070. connection = self._connect()
  1071. try:
  1072. rows = connection.execute(
  1073. "SELECT * FROM edge_local_artifacts WHERE task_id=? ORDER BY created_at,artifact_digest",
  1074. (task_id,),
  1075. ).fetchall()
  1076. finally:
  1077. connection.close()
  1078. return tuple(self._artifact_record(row) for row in rows)
  1079. def claim_artifact_cleanup(self, lease_owner: str) -> LocalArtifact | None:
  1080. owner = self._lease_owner(lease_owner)
  1081. with self._transaction() as connection:
  1082. now = _utc_now(self._clock)
  1083. now_text = _timestamp(now)
  1084. expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds))
  1085. connection.execute(
  1086. """UPDATE edge_local_artifacts SET cleanup_status='failed',
  1087. lease_owner=NULL,lease_token=NULL,lease_expires_at=NULL,
  1088. error_code='lease_attempts_exhausted',updated_at=?
  1089. WHERE cleanup_status='deleting'
  1090. AND julianday(lease_expires_at)<=julianday(?) AND attempt_count>=?""",
  1091. (now_text, now_text, self.max_attempts),
  1092. )
  1093. row = connection.execute(
  1094. """SELECT * FROM edge_local_artifacts WHERE attempt_count<?
  1095. AND julianday(retention_until)<=julianday(?) AND (
  1096. (cleanup_status='pending' AND julianday(available_at)<=julianday(?))
  1097. OR (cleanup_status='deleting'
  1098. AND julianday(lease_expires_at)<=julianday(?)))
  1099. ORDER BY retention_until,created_at,artifact_digest LIMIT 1""",
  1100. (self.max_attempts, now_text, now_text, now_text),
  1101. ).fetchone()
  1102. if row is None:
  1103. return None
  1104. lease_token = str(uuid.uuid4())
  1105. connection.execute(
  1106. """UPDATE edge_local_artifacts SET cleanup_status='deleting',
  1107. attempt_count=attempt_count+1,lease_owner=?,lease_token=?,
  1108. lease_expires_at=?,error_code=NULL,updated_at=?
  1109. WHERE artifact_digest=? AND attempt_count=? AND (
  1110. (cleanup_status='pending' AND julianday(available_at)<=julianday(?))
  1111. OR (cleanup_status='deleting'
  1112. AND julianday(lease_expires_at)<=julianday(?)))""",
  1113. (
  1114. owner, lease_token, expires_at, now_text, row["artifact_digest"],
  1115. row["attempt_count"], now_text, now_text,
  1116. ),
  1117. )
  1118. claimed = connection.execute(
  1119. "SELECT * FROM edge_local_artifacts WHERE artifact_digest=?",
  1120. (row["artifact_digest"],),
  1121. ).fetchone()
  1122. return self._artifact_record(claimed)
  1123. def acknowledge_artifact_cleanup(
  1124. self,
  1125. artifact_digest: str,
  1126. *,
  1127. lease_token: str,
  1128. receipt: Mapping[str, object],
  1129. ) -> LocalArtifact:
  1130. required = {"artifact_digest", "artifact_ref_hash", "deleted_at", "status"}
  1131. with self._transaction() as connection:
  1132. row = self._leased_row(
  1133. connection, "edge_local_artifacts", "artifact_digest",
  1134. artifact_digest, lease_token, status="deleting",
  1135. status_column="cleanup_status",
  1136. )
  1137. if (
  1138. not isinstance(receipt, Mapping)
  1139. or set(receipt) != required
  1140. or receipt.get("artifact_digest") != artifact_digest
  1141. or receipt.get("artifact_ref_hash") != row["artifact_ref_hash"]
  1142. or receipt.get("status") != "deleted"
  1143. ):
  1144. raise EdgePolicyError("artifact deletion receipt is invalid")
  1145. try:
  1146. deleted_at = canonical_timestamp(receipt["deleted_at"], "deleted_at")
  1147. except EdgeContractError as exc:
  1148. raise EdgePolicyError("artifact deletion receipt is invalid") from exc
  1149. normalized = {
  1150. "artifact_digest": artifact_digest,
  1151. "artifact_ref_hash": row["artifact_ref_hash"],
  1152. "deleted_at": deleted_at,
  1153. "status": "deleted",
  1154. }
  1155. now_text = _timestamp(_utc_now(self._clock))
  1156. connection.execute(
  1157. """UPDATE edge_local_artifacts SET artifact_ref=NULL,
  1158. cleanup_status='deleted',deletion_receipt_json=?,lease_owner=NULL,
  1159. lease_token=NULL,lease_expires_at=NULL,error_code=NULL,updated_at=?
  1160. WHERE artifact_digest=? AND cleanup_status='deleting' AND lease_token=?""",
  1161. (
  1162. strict_json_bytes(normalized).decode("utf-8"), now_text,
  1163. artifact_digest, lease_token,
  1164. ),
  1165. )
  1166. updated = connection.execute(
  1167. "SELECT * FROM edge_local_artifacts WHERE artifact_digest=?",
  1168. (artifact_digest,),
  1169. ).fetchone()
  1170. return self._artifact_record(updated)
  1171. def fail_artifact_cleanup(
  1172. self, artifact_digest: str, *, lease_token: str, error_code: str
  1173. ) -> LocalArtifact:
  1174. safe_error = self._error_code(error_code)
  1175. with self._transaction() as connection:
  1176. now = _utc_now(self._clock)
  1177. now_text = _timestamp(now)
  1178. row = self._leased_row(
  1179. connection, "edge_local_artifacts", "artifact_digest",
  1180. artifact_digest, lease_token, status="deleting",
  1181. status_column="cleanup_status",
  1182. )
  1183. terminal = row["attempt_count"] >= self.max_attempts
  1184. status = "failed" if terminal else "pending"
  1185. available_at = now_text if terminal else _timestamp(
  1186. now + timedelta(seconds=self._retry_delay(row["attempt_count"]))
  1187. )
  1188. connection.execute(
  1189. """UPDATE edge_local_artifacts SET cleanup_status=?,available_at=?,
  1190. lease_owner=NULL,lease_token=NULL,lease_expires_at=NULL,error_code=?,
  1191. updated_at=? WHERE artifact_digest=? AND cleanup_status='deleting'
  1192. AND lease_token=?""",
  1193. (
  1194. status, available_at, safe_error, now_text,
  1195. artifact_digest, lease_token,
  1196. ),
  1197. )
  1198. updated = connection.execute(
  1199. "SELECT * FROM edge_local_artifacts WHERE artifact_digest=?",
  1200. (artifact_digest,),
  1201. ).fetchone()
  1202. return self._artifact_record(updated)
  1203. def put_release_state(
  1204. self,
  1205. *,
  1206. release_id: str,
  1207. manifest_digest: str,
  1208. version: str,
  1209. rollback_version: str,
  1210. ) -> ReleaseState:
  1211. release_id = self._bounded_identifier(release_id, "release_id")
  1212. manifest_digest = self._sha256(manifest_digest, "manifest_digest")
  1213. version = self._bounded_identifier(version, "version")
  1214. rollback_version = self._bounded_identifier(rollback_version, "rollback_version")
  1215. with self._transaction() as connection:
  1216. existing = connection.execute(
  1217. "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,)
  1218. ).fetchone()
  1219. if existing is not None:
  1220. if (
  1221. existing["manifest_digest"], existing["version"], existing["rollback_version"]
  1222. ) != (manifest_digest, version, rollback_version):
  1223. raise EdgeQueueConflictError("release ID has different immutable content")
  1224. return ReleaseState(**dict(existing))
  1225. now_text = _timestamp(_utc_now(self._clock))
  1226. connection.execute(
  1227. """INSERT INTO edge_release_state
  1228. (release_id,manifest_digest,version,rollback_version,status,updated_at)
  1229. VALUES (?,?,?,?,'offered',?)""",
  1230. (release_id, manifest_digest, version, rollback_version, now_text),
  1231. )
  1232. row = connection.execute(
  1233. "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,)
  1234. ).fetchone()
  1235. return ReleaseState(**dict(row))
  1236. def compare_and_set_release_state(
  1237. self,
  1238. release_id: str,
  1239. *,
  1240. expected_status: str,
  1241. target_status: str,
  1242. previous_version: str | None = None,
  1243. current_version: str | None = None,
  1244. ) -> ReleaseState:
  1245. release_id = self._bounded_identifier(release_id, "release_id")
  1246. allowed = {"offered", "accepted", "candidate", "installed", "failed", "rolled_back"}
  1247. if expected_status not in allowed or target_status not in allowed:
  1248. raise EdgeQueueConflictError("release status is invalid")
  1249. transitions = {
  1250. "offered": {"accepted", "failed"},
  1251. "accepted": {"candidate", "failed"},
  1252. "candidate": {"installed", "failed", "rolled_back"},
  1253. "installed": {"rolled_back", "failed"},
  1254. "failed": {"rolled_back"},
  1255. "rolled_back": set(),
  1256. }
  1257. if target_status not in transitions[expected_status]:
  1258. raise EdgeQueueConflictError("release state transition is not approved")
  1259. if previous_version is not None:
  1260. previous_version = self._bounded_identifier(previous_version, "previous_version")
  1261. if current_version is not None:
  1262. current_version = self._bounded_identifier(current_version, "current_version")
  1263. with self._transaction() as connection:
  1264. now_text = _timestamp(_utc_now(self._clock))
  1265. updated = connection.execute(
  1266. """UPDATE edge_release_state
  1267. SET status=?,previous_version=COALESCE(?,previous_version),
  1268. current_version=COALESCE(?,current_version),updated_at=?
  1269. WHERE release_id=? AND status=?""",
  1270. (
  1271. target_status, previous_version, current_version, now_text,
  1272. release_id, expected_status,
  1273. ),
  1274. )
  1275. if updated.rowcount != 1:
  1276. row = connection.execute(
  1277. "SELECT status FROM edge_release_state WHERE release_id=?", (release_id,)
  1278. ).fetchone()
  1279. if row is None:
  1280. raise EdgeQueueError("release does not exist")
  1281. raise EdgeQueueConflictError("release state compare-and-set lost")
  1282. row = connection.execute(
  1283. "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,)
  1284. ).fetchone()
  1285. return ReleaseState(**dict(row))
  1286. def get_release_state(self, release_id: str) -> ReleaseState | None:
  1287. release_id = self._bounded_identifier(release_id, "release_id")
  1288. connection = self._connect()
  1289. try:
  1290. row = connection.execute(
  1291. "SELECT * FROM edge_release_state WHERE release_id=?", (release_id,)
  1292. ).fetchone()
  1293. finally:
  1294. connection.close()
  1295. return ReleaseState(**dict(row)) if row is not None else None
  1296. def list_release_states(self) -> tuple[ReleaseState, ...]:
  1297. connection = self._connect()
  1298. try:
  1299. rows = connection.execute(
  1300. "SELECT * FROM edge_release_state ORDER BY updated_at,release_id"
  1301. ).fetchall()
  1302. finally:
  1303. connection.close()
  1304. return tuple(ReleaseState(**dict(row)) for row in rows)
  1305. def complete(self, task_id: str, *, lease_token: str) -> QueuedTask:
  1306. return self._finish_task(
  1307. task_id,
  1308. lease_token=lease_token,
  1309. target_status="completed",
  1310. error_code=None,
  1311. )
  1312. def fail(
  1313. self,
  1314. task_id: str,
  1315. *,
  1316. lease_token: str,
  1317. error_code: str,
  1318. retryable: bool = True,
  1319. ) -> QueuedTask:
  1320. safe_error = self._error_code(error_code)
  1321. if not isinstance(retryable, bool):
  1322. raise ValueError("retryable must be a boolean")
  1323. with self._transaction() as connection:
  1324. now = _utc_now(self._clock)
  1325. now_text = _timestamp(now)
  1326. row = self._leased_row(
  1327. connection, "edge_tasks", "task_id", task_id, lease_token
  1328. )
  1329. terminal = not retryable or row["attempt_count"] >= self.max_attempts
  1330. status = "failed" if terminal else "pending"
  1331. available_at = now_text
  1332. if not terminal:
  1333. available_at = _timestamp(
  1334. now + timedelta(seconds=self._retry_delay(row["attempt_count"]))
  1335. )
  1336. connection.execute(
  1337. """
  1338. UPDATE edge_tasks
  1339. SET status = ?, available_at = ?, lease_owner = NULL,
  1340. lease_token = NULL, lease_expires_at = NULL,
  1341. error_code = ?, updated_at = ?
  1342. WHERE task_id = ? AND status = 'leased' AND lease_token = ?
  1343. """,
  1344. (status, available_at, safe_error, now_text, task_id, lease_token),
  1345. )
  1346. updated = connection.execute(
  1347. "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,)
  1348. ).fetchone()
  1349. return self._task_record(updated)
  1350. def _finish_task(
  1351. self,
  1352. task_id: str,
  1353. *,
  1354. lease_token: str,
  1355. target_status: str,
  1356. error_code: str | None,
  1357. ) -> QueuedTask:
  1358. with self._transaction() as connection:
  1359. now_text = _timestamp(_utc_now(self._clock))
  1360. self._leased_row(
  1361. connection, "edge_tasks", "task_id", task_id, lease_token
  1362. )
  1363. connection.execute(
  1364. """
  1365. UPDATE edge_tasks
  1366. SET status = ?, lease_owner = NULL, lease_token = NULL,
  1367. lease_expires_at = NULL, error_code = ?, updated_at = ?
  1368. WHERE task_id = ? AND status = 'leased' AND lease_token = ?
  1369. """,
  1370. (target_status, error_code, now_text, task_id, lease_token),
  1371. )
  1372. row = connection.execute(
  1373. "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,)
  1374. ).fetchone()
  1375. return self._task_record(row)
  1376. def cancel(self, task_id: str) -> QueuedTask:
  1377. with self._transaction() as connection:
  1378. now_text = _timestamp(_utc_now(self._clock))
  1379. row = connection.execute(
  1380. "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,)
  1381. ).fetchone()
  1382. if row is None:
  1383. raise EdgeQueueError("task does not exist")
  1384. if row["status"] in {"pending", "leased"}:
  1385. connection.execute(
  1386. """
  1387. UPDATE edge_tasks
  1388. SET status = 'cancelled', lease_owner = NULL,
  1389. lease_token = NULL, lease_expires_at = NULL,
  1390. updated_at = ?
  1391. WHERE task_id = ? AND status IN ('pending', 'leased')
  1392. """,
  1393. (now_text, task_id),
  1394. )
  1395. row = connection.execute(
  1396. "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,)
  1397. ).fetchone()
  1398. return self._task_record(row)
  1399. def cancel_with_outcome(self, task_id: str) -> QueuedOutcome:
  1400. """Atomically fence local work and persist the cancelled control outcome."""
  1401. task_id = self._bounded_identifier(task_id, "task_id")
  1402. with self._transaction() as connection:
  1403. now_text = _timestamp(_utc_now(self._clock))
  1404. row = connection.execute(
  1405. "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,)
  1406. ).fetchone()
  1407. if row is None:
  1408. raise EdgeQueueError("task does not exist")
  1409. if row["status"] in {"pending", "leased"}:
  1410. connection.execute(
  1411. """UPDATE edge_tasks SET status='cancelled',lease_owner=NULL,
  1412. lease_token=NULL,lease_expires_at=NULL,updated_at=?
  1413. WHERE task_id=? AND status IN ('pending','leased')""",
  1414. (now_text, task_id),
  1415. )
  1416. row = connection.execute(
  1417. "SELECT * FROM edge_tasks WHERE task_id=?", (task_id,)
  1418. ).fetchone()
  1419. if row["status"] != "cancelled":
  1420. raise EdgeQueueConflictError("only a cancelled task can emit cancellation")
  1421. remote_digest = row["remote_lease_token_digest"]
  1422. if not isinstance(remote_digest, str):
  1423. raise EdgeQueueLeaseError("cancel outcome requires a remote lease binding")
  1424. outcome = {
  1425. "task_id": task_id,
  1426. "outcome": "cancelled",
  1427. "safe_summary": {"task_status": "cancelled"},
  1428. }
  1429. digest = canonical_sha256(outcome)
  1430. encoded = strict_json_bytes(outcome).decode("utf-8")
  1431. existing = connection.execute(
  1432. "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,)
  1433. ).fetchone()
  1434. if existing is not None:
  1435. if (
  1436. existing["digest"] != digest
  1437. or existing["remote_lease_token_digest"] != remote_digest
  1438. ):
  1439. raise EdgeQueueConflictError("cancel outcome replay conflicts")
  1440. return self._outcome_record(existing)
  1441. connection.execute(
  1442. """INSERT INTO edge_task_outcomes(
  1443. task_id,digest,outcome_json,status,attempt_count,available_at,
  1444. remote_lease_token_digest,created_at,updated_at
  1445. ) VALUES (?,?,?,'pending',0,?,?,?,?)""",
  1446. (task_id, digest, encoded, now_text, remote_digest, now_text, now_text),
  1447. )
  1448. inserted = connection.execute(
  1449. "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,)
  1450. ).fetchone()
  1451. return self._outcome_record(inserted)
  1452. def claim_outcome(self, lease_owner: str) -> QueuedOutcome | None:
  1453. owner = self._lease_owner(lease_owner)
  1454. with self._transaction() as connection:
  1455. now = _utc_now(self._clock)
  1456. now_text = _timestamp(now)
  1457. expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds))
  1458. connection.execute(
  1459. """UPDATE edge_task_outcomes SET status='failed',lease_owner=NULL,
  1460. lease_token=NULL,lease_expires_at=NULL,error_code='lease_attempts_exhausted',
  1461. updated_at=? WHERE status='sending'
  1462. AND julianday(lease_expires_at)<=julianday(?) AND attempt_count>=?""",
  1463. (now_text, now_text, self.max_attempts),
  1464. )
  1465. row = connection.execute(
  1466. """SELECT * FROM edge_task_outcomes WHERE attempt_count<? AND (
  1467. (status='pending' AND julianday(available_at)<=julianday(?)) OR
  1468. (status='sending' AND julianday(lease_expires_at)<=julianday(?)))
  1469. ORDER BY created_at,task_id LIMIT 1""",
  1470. (self.max_attempts, now_text, now_text),
  1471. ).fetchone()
  1472. if row is None:
  1473. return None
  1474. lease_token = str(uuid.uuid4())
  1475. connection.execute(
  1476. """UPDATE edge_task_outcomes SET status='sending',
  1477. attempt_count=attempt_count+1,lease_owner=?,lease_token=?,
  1478. lease_expires_at=?,error_code=NULL,updated_at=?
  1479. WHERE task_id=? AND attempt_count=? AND (
  1480. (status='pending' AND julianday(available_at)<=julianday(?)) OR
  1481. (status='sending' AND julianday(lease_expires_at)<=julianday(?)))""",
  1482. (
  1483. owner, lease_token, expires_at, now_text, row["task_id"],
  1484. row["attempt_count"], now_text, now_text,
  1485. ),
  1486. )
  1487. claimed = connection.execute(
  1488. "SELECT * FROM edge_task_outcomes WHERE task_id=?", (row["task_id"],)
  1489. ).fetchone()
  1490. return self._outcome_record(claimed)
  1491. def acknowledge_outcome(
  1492. self,
  1493. task_id: str,
  1494. *,
  1495. lease_token: str,
  1496. acknowledgement: Mapping[str, object],
  1497. ) -> QueuedOutcome:
  1498. required = {"task_id", "status", "replayed"}
  1499. if (
  1500. not isinstance(acknowledgement, Mapping)
  1501. or set(acknowledgement) != required
  1502. or acknowledgement.get("task_id") != task_id
  1503. or acknowledgement.get("status") != "cancelled"
  1504. or not isinstance(acknowledgement.get("replayed"), bool)
  1505. ):
  1506. raise EdgePolicyError("cancel outcome acknowledgement is invalid")
  1507. encoded = strict_json_bytes(acknowledgement).decode("utf-8")
  1508. ack_digest = canonical_sha256(acknowledgement)
  1509. with self._transaction() as connection:
  1510. now_text = _timestamp(_utc_now(self._clock))
  1511. self._leased_row(
  1512. connection, "edge_task_outcomes", "task_id", task_id, lease_token,
  1513. status="sending",
  1514. )
  1515. connection.execute(
  1516. """UPDATE edge_task_outcomes SET status='acknowledged',
  1517. acknowledgement_json=?,ack_digest=?,lease_owner=NULL,lease_token=NULL,
  1518. lease_expires_at=NULL,updated_at=?
  1519. WHERE task_id=? AND status='sending' AND lease_token=?""",
  1520. (encoded, ack_digest, now_text, task_id, lease_token),
  1521. )
  1522. row = connection.execute(
  1523. "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,)
  1524. ).fetchone()
  1525. return self._outcome_record(row)
  1526. def fail_outcome(
  1527. self, task_id: str, *, lease_token: str, error_code: str
  1528. ) -> QueuedOutcome:
  1529. safe_error = self._error_code(error_code)
  1530. with self._transaction() as connection:
  1531. now = _utc_now(self._clock)
  1532. now_text = _timestamp(now)
  1533. row = self._leased_row(
  1534. connection, "edge_task_outcomes", "task_id", task_id, lease_token,
  1535. status="sending",
  1536. )
  1537. terminal = row["attempt_count"] >= self.max_attempts
  1538. status = "failed" if terminal else "pending"
  1539. available_at = now_text if terminal else _timestamp(
  1540. now + timedelta(seconds=self._retry_delay(row["attempt_count"]))
  1541. )
  1542. connection.execute(
  1543. """UPDATE edge_task_outcomes SET status=?,available_at=?,lease_owner=NULL,
  1544. lease_token=NULL,lease_expires_at=NULL,error_code=?,updated_at=?
  1545. WHERE task_id=? AND status='sending' AND lease_token=?""",
  1546. (
  1547. status, available_at, safe_error, now_text, task_id, lease_token,
  1548. ),
  1549. )
  1550. updated = connection.execute(
  1551. "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,)
  1552. ).fetchone()
  1553. return self._outcome_record(updated)
  1554. def get_outcome(self, task_id: str) -> QueuedOutcome | None:
  1555. connection = self._connect()
  1556. try:
  1557. row = connection.execute(
  1558. "SELECT * FROM edge_task_outcomes WHERE task_id=?", (task_id,)
  1559. ).fetchone()
  1560. finally:
  1561. connection.close()
  1562. return self._outcome_record(row) if row is not None else None
  1563. def get_task(self, task_id: str) -> QueuedTask | None:
  1564. connection = self._connect()
  1565. try:
  1566. row = connection.execute(
  1567. "SELECT * FROM edge_tasks WHERE task_id = ?", (task_id,)
  1568. ).fetchone()
  1569. finally:
  1570. connection.close()
  1571. return self._task_record(row) if row is not None else None
  1572. def persist_event(
  1573. self, event: EdgeEventContract | Mapping[str, object]
  1574. ) -> QueuedEvent:
  1575. raw = event.to_mapping() if isinstance(event, EdgeEventContract) else dict(event)
  1576. try:
  1577. contract = EdgeEventContract.from_mapping(raw)
  1578. except EdgeContractError as exc:
  1579. raw_id = raw.get("event_id")
  1580. if isinstance(raw_id, str) and raw_id:
  1581. with self._transaction() as connection:
  1582. existing = connection.execute(
  1583. "SELECT 1 FROM edge_outbound_events WHERE event_id = ?",
  1584. (raw_id,),
  1585. ).fetchone()
  1586. if existing is not None:
  1587. raise EdgeQueueConflictError(
  1588. "event ID already exists with a different digest"
  1589. ) from exc
  1590. raise
  1591. normalized = contract.to_mapping()
  1592. digest = contract.digest
  1593. with self._transaction() as connection:
  1594. now_text = _timestamp(_utc_now(self._clock))
  1595. existing = connection.execute(
  1596. "SELECT * FROM edge_outbound_events WHERE event_id = ?",
  1597. (contract.event_id,),
  1598. ).fetchone()
  1599. if existing is not None:
  1600. if existing["digest"] != digest:
  1601. raise EdgeQueueConflictError(
  1602. "event ID already exists with a different digest"
  1603. )
  1604. return self._event_record(existing)
  1605. approved = self.egress_policy.approve_event(normalized)
  1606. encoded = strict_json_bytes(approved).decode("utf-8")
  1607. connection.execute(
  1608. """
  1609. INSERT INTO edge_outbound_events (
  1610. event_id, digest, event_json, status, attempt_count,
  1611. available_at, created_at, updated_at
  1612. ) VALUES (?, ?, ?, 'pending', 0, ?, ?, ?)
  1613. """,
  1614. (contract.event_id, digest, encoded, now_text, now_text, now_text),
  1615. )
  1616. row = connection.execute(
  1617. "SELECT * FROM edge_outbound_events WHERE event_id = ?",
  1618. (contract.event_id,),
  1619. ).fetchone()
  1620. return self._event_record(row)
  1621. def claim_event(self, lease_owner: str) -> QueuedEvent | None:
  1622. owner = self._lease_owner(lease_owner)
  1623. with self._transaction() as connection:
  1624. now = _utc_now(self._clock)
  1625. now_text = _timestamp(now)
  1626. lease_token = str(uuid.uuid4())
  1627. expires_at = _timestamp(now + timedelta(seconds=self.lease_seconds))
  1628. connection.execute(
  1629. """
  1630. UPDATE edge_outbound_events
  1631. SET status = 'failed', lease_owner = NULL, lease_token = NULL,
  1632. lease_expires_at = NULL,
  1633. error_code = 'lease_attempts_exhausted', updated_at = ?
  1634. WHERE status = 'sending'
  1635. AND julianday(lease_expires_at) <= julianday(?)
  1636. AND attempt_count >= ?
  1637. """,
  1638. (now_text, now_text, self.max_attempts),
  1639. )
  1640. row = connection.execute(
  1641. """
  1642. SELECT * FROM edge_outbound_events
  1643. WHERE attempt_count < ? AND (
  1644. (status = 'pending' AND julianday(available_at) <= julianday(?))
  1645. OR (status = 'sending'
  1646. AND julianday(lease_expires_at) <= julianday(?))
  1647. )
  1648. ORDER BY created_at, event_id
  1649. LIMIT 1
  1650. """,
  1651. (self.max_attempts, now_text, now_text),
  1652. ).fetchone()
  1653. if row is None:
  1654. return None
  1655. connection.execute(
  1656. """
  1657. UPDATE edge_outbound_events
  1658. SET status = 'sending', attempt_count = attempt_count + 1,
  1659. lease_owner = ?, lease_token = ?, lease_expires_at = ?,
  1660. error_code = NULL, updated_at = ?
  1661. WHERE event_id = ? AND attempt_count = ? AND (
  1662. (status = 'pending' AND julianday(available_at) <= julianday(?))
  1663. OR (status = 'sending'
  1664. AND julianday(lease_expires_at) <= julianday(?))
  1665. )
  1666. """,
  1667. (
  1668. owner,
  1669. lease_token,
  1670. expires_at,
  1671. now_text,
  1672. row["event_id"],
  1673. row["attempt_count"],
  1674. now_text,
  1675. now_text,
  1676. ),
  1677. )
  1678. claimed = connection.execute(
  1679. "SELECT * FROM edge_outbound_events WHERE event_id = ?",
  1680. (row["event_id"],),
  1681. ).fetchone()
  1682. return self._event_record(claimed)
  1683. def acknowledge_event(
  1684. self,
  1685. event_id: str,
  1686. *,
  1687. lease_token: str,
  1688. acknowledgement: Mapping[str, object],
  1689. ) -> QueuedEvent:
  1690. token_digest = canonical_sha256(lease_token)
  1691. with self._transaction() as connection:
  1692. now_text = _timestamp(_utc_now(self._clock))
  1693. existing = connection.execute(
  1694. "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,)
  1695. ).fetchone()
  1696. if existing is None:
  1697. raise EdgeQueueError("event does not exist")
  1698. normalized_ack = self._validate_acknowledgement(
  1699. acknowledgement,
  1700. event_id=event_id,
  1701. event_digest=existing["digest"],
  1702. remote_lease_token_digest=existing["remote_lease_token_digest"],
  1703. )
  1704. encoded = strict_json_bytes(normalized_ack)
  1705. if len(encoded) > 4_096:
  1706. raise EdgePolicyError("acknowledgement exceeds the byte limit")
  1707. ack_digest = canonical_sha256(normalized_ack)
  1708. if existing["status"] == "acknowledged":
  1709. if existing["lease_token_digest"] != token_digest:
  1710. raise EdgeQueueLeaseError(
  1711. "acknowledgement replay requires the original lease token"
  1712. )
  1713. if existing["ack_digest"] != ack_digest:
  1714. raise EdgeQueueConflictError(
  1715. "acknowledgement replay has a different digest"
  1716. )
  1717. return self._event_record(existing)
  1718. self._leased_row(
  1719. connection,
  1720. "edge_outbound_events",
  1721. "event_id",
  1722. event_id,
  1723. lease_token,
  1724. status="sending",
  1725. )
  1726. connection.execute(
  1727. """
  1728. UPDATE edge_outbound_events
  1729. SET status = 'acknowledged', acknowledgement_json = ?,
  1730. ack_digest = ?, lease_token_digest = ?,
  1731. control_lease_token_digest = ?,
  1732. lease_owner = NULL, lease_token = NULL,
  1733. lease_expires_at = NULL, updated_at = ?
  1734. WHERE event_id = ? AND status = 'sending' AND lease_token = ?
  1735. """,
  1736. (
  1737. encoded.decode("utf-8"),
  1738. ack_digest,
  1739. token_digest,
  1740. token_digest,
  1741. now_text,
  1742. event_id,
  1743. lease_token,
  1744. ),
  1745. )
  1746. row = connection.execute(
  1747. "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,)
  1748. ).fetchone()
  1749. return self._event_record(row)
  1750. def fail_event(
  1751. self, event_id: str, *, lease_token: str, error_code: str
  1752. ) -> QueuedEvent:
  1753. safe_error = self._error_code(error_code)
  1754. with self._transaction() as connection:
  1755. now = _utc_now(self._clock)
  1756. now_text = _timestamp(now)
  1757. row = self._leased_row(
  1758. connection,
  1759. "edge_outbound_events",
  1760. "event_id",
  1761. event_id,
  1762. lease_token,
  1763. status="sending",
  1764. )
  1765. terminal = row["attempt_count"] >= self.max_attempts
  1766. status = "failed" if terminal else "pending"
  1767. available_at = now_text
  1768. if not terminal:
  1769. available_at = _timestamp(
  1770. now + timedelta(seconds=self._retry_delay(row["attempt_count"]))
  1771. )
  1772. connection.execute(
  1773. """
  1774. UPDATE edge_outbound_events
  1775. SET status = ?, available_at = ?, lease_owner = NULL,
  1776. lease_token = NULL, lease_expires_at = NULL,
  1777. error_code = ?, updated_at = ?
  1778. WHERE event_id = ? AND status = 'sending' AND lease_token = ?
  1779. """,
  1780. (status, available_at, safe_error, now_text, event_id, lease_token),
  1781. )
  1782. updated = connection.execute(
  1783. "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,)
  1784. ).fetchone()
  1785. return self._event_record(updated)
  1786. def get_event(self, event_id: str) -> QueuedEvent | None:
  1787. connection = self._connect()
  1788. try:
  1789. row = connection.execute(
  1790. "SELECT * FROM edge_outbound_events WHERE event_id = ?", (event_id,)
  1791. ).fetchone()
  1792. finally:
  1793. connection.close()
  1794. return self._event_record(row) if row is not None else None
  1795. @staticmethod
  1796. def _validate_local_artifact(value: Mapping[str, object]) -> dict[str, str]:
  1797. required = {
  1798. "artifact_digest", "artifact_ref", "artifact_ref_hash",
  1799. "classification", "retention_until",
  1800. }
  1801. if not isinstance(value, Mapping) or set(value) != required:
  1802. raise EdgePolicyError("local artifact requires exact metadata fields")
  1803. digest = value["artifact_digest"]
  1804. reference = value["artifact_ref"]
  1805. reference_hash = value["artifact_ref_hash"]
  1806. classification = value["classification"]
  1807. if not isinstance(digest, str) or not re.fullmatch(r"[0-9a-f]{64}", digest):
  1808. raise EdgePolicyError("local artifact digest is invalid")
  1809. if (
  1810. not isinstance(reference, str)
  1811. or not Path(reference).is_absolute()
  1812. or len(reference.encode("utf-8")) > 2_048
  1813. or not isinstance(reference_hash, str)
  1814. or not re.fullmatch(r"[0-9a-f]{64}", reference_hash)
  1815. or reference_hash != canonical_sha256(reference)
  1816. ):
  1817. raise EdgePolicyError("local artifact reference is invalid")
  1818. if classification not in {
  1819. "raw", "recent_detail", "restricted", "desensitized_metadata",
  1820. "statistics", "lineage", "evidence",
  1821. }:
  1822. raise EdgePolicyError("local artifact classification is invalid")
  1823. try:
  1824. retention = canonical_timestamp(value["retention_until"], "retention_until")
  1825. except EdgeContractError as exc:
  1826. raise EdgePolicyError("local artifact retention is invalid") from exc
  1827. return {
  1828. "artifact_digest": digest,
  1829. "artifact_ref": reference, # type: ignore[dict-item]
  1830. "artifact_ref_hash": reference_hash, # type: ignore[dict-item]
  1831. "classification": classification, # type: ignore[dict-item]
  1832. "retention_until": retention,
  1833. }
  1834. def _normalize_remote_binding(
  1835. self, token: object, expires_at: object, attempt: object
  1836. ) -> tuple[str, str | None, str, int]:
  1837. if (
  1838. not isinstance(token, str) or not token.strip()
  1839. or len(token.encode("utf-8")) > 1024
  1840. ):
  1841. raise EdgeQueueLeaseError("remote lease token is invalid")
  1842. if isinstance(attempt, bool) or not isinstance(attempt, int) or attempt < 1:
  1843. raise EdgeQueueLeaseError("remote lease attempt is invalid")
  1844. try:
  1845. expiry = canonical_timestamp(expires_at, "remote_lease_expires_at")
  1846. except EdgeContractError as exc:
  1847. raise EdgeQueueLeaseError("remote lease expiry is invalid") from exc
  1848. ciphertext: str | None = None
  1849. if self._remote_lease_encoder is not None:
  1850. try:
  1851. ciphertext = self._remote_lease_encoder(token)
  1852. except Exception as exc:
  1853. raise EdgeQueueLeaseError("remote lease protection failed") from exc
  1854. if (
  1855. not isinstance(ciphertext, str) or not ciphertext
  1856. or ciphertext == token or len(ciphertext.encode("utf-8")) > 4096
  1857. ):
  1858. raise EdgeQueueLeaseError("remote lease protection produced invalid ciphertext")
  1859. return self._opaque_token_digest(token), ciphertext, expiry, attempt
  1860. @staticmethod
  1861. def _opaque_token_digest(token: str) -> str:
  1862. return hashlib.sha256(token.encode("utf-8")).hexdigest()
  1863. def _retry_delay(self, attempt_count: int) -> int:
  1864. base = min(
  1865. self.retry_max_seconds,
  1866. self.retry_base_seconds * (2 ** max(0, attempt_count - 1)),
  1867. )
  1868. sample = self._random()
  1869. if isinstance(sample, bool) or not isinstance(sample, (int, float)) or not 0 <= sample <= 1:
  1870. raise EdgeQueueError("retry random source returned an invalid value")
  1871. return min(
  1872. self.retry_max_seconds,
  1873. max(1, round(base * (0.5 + float(sample)))),
  1874. )
  1875. @staticmethod
  1876. def _bounded_identifier(value: object, label: str) -> str:
  1877. if (
  1878. not isinstance(value, str) or not value or value.strip() != value
  1879. or len(value.encode("utf-8")) > 255
  1880. or not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._:-]{0,254}", value)
  1881. ):
  1882. raise EdgeQueueError(f"{label} is invalid")
  1883. return value
  1884. @staticmethod
  1885. def _sha256(value: object, label: str) -> str:
  1886. if not isinstance(value, str) or not re.fullmatch(r"[0-9a-f]{64}", value):
  1887. raise EdgeQueueError(f"{label} is invalid")
  1888. return value
  1889. @staticmethod
  1890. def _lease_owner(value: object) -> str:
  1891. if (
  1892. not isinstance(value, str)
  1893. or not value
  1894. or len(value.encode("utf-8")) > 255
  1895. or value.strip() != value
  1896. ):
  1897. raise ValueError("lease_owner is invalid")
  1898. return value
  1899. @staticmethod
  1900. def _error_code(value: object) -> str:
  1901. if not isinstance(value, str) or not _ERROR_CODE.fullmatch(value):
  1902. raise EdgePolicyError("error_code must be a bounded code")
  1903. if any(part in value for part in _SENSITIVE_ERROR_PARTS):
  1904. raise EdgePolicyError("error_code contains sensitive content")
  1905. return value
  1906. def _leased_row(
  1907. self,
  1908. connection: sqlite3.Connection,
  1909. table: str,
  1910. id_column: str,
  1911. identity: str,
  1912. lease_token: str,
  1913. *,
  1914. status: str = "leased",
  1915. status_column: str = "status",
  1916. ) -> sqlite3.Row:
  1917. row = connection.execute(
  1918. f"SELECT * FROM {table} WHERE {id_column} = ?", # noqa: S608
  1919. (identity,),
  1920. ).fetchone()
  1921. if row is None:
  1922. raise EdgeQueueError("queue item does not exist")
  1923. if row[status_column] != status or row["lease_token"] != lease_token:
  1924. raise EdgeQueueLeaseError("current lease token is required")
  1925. lease_expires_at = row["lease_expires_at"]
  1926. if lease_expires_at is None:
  1927. raise EdgeQueueLeaseError("lease has expired")
  1928. parsed_expiry = datetime.fromisoformat(lease_expires_at.replace("Z", "+00:00"))
  1929. if parsed_expiry <= _utc_now(self._clock):
  1930. raise EdgeQueueLeaseError("lease has expired")
  1931. return row
  1932. @staticmethod
  1933. def _validate_acknowledgement(
  1934. value: Mapping[str, object],
  1935. *,
  1936. event_id: str,
  1937. event_digest: str,
  1938. remote_lease_token_digest: str | None,
  1939. ) -> dict[str, str]:
  1940. if remote_lease_token_digest is not None:
  1941. required = {"event_id", "event_digest", "remote_lease_digest", "received_at", "status"}
  1942. if not isinstance(value, Mapping) or set(value) != required:
  1943. raise EdgePolicyError("bound acknowledgement requires exact approved fields")
  1944. if value["event_id"] != event_id or value["event_digest"] != event_digest:
  1945. raise EdgePolicyError("acknowledgement does not bind the event")
  1946. if value["remote_lease_digest"] != remote_lease_token_digest:
  1947. raise EdgePolicyError("acknowledgement remote lease is invalid")
  1948. if value["status"] != "accepted":
  1949. raise EdgePolicyError("acknowledgement status is invalid")
  1950. try:
  1951. received_at = canonical_timestamp(value["received_at"], "received_at")
  1952. except EdgeContractError as exc:
  1953. raise EdgePolicyError("acknowledgement received_at is invalid") from exc
  1954. return {
  1955. "event_id": event_id,
  1956. "event_digest": event_digest,
  1957. "remote_lease_digest": remote_lease_token_digest,
  1958. "received_at": received_at,
  1959. "status": "accepted",
  1960. }
  1961. if not isinstance(value, Mapping) or set(value) != {
  1962. "message_id",
  1963. "received_at",
  1964. "status",
  1965. }:
  1966. raise EdgePolicyError("acknowledgement requires exact approved fields")
  1967. message_id = value["message_id"]
  1968. status = value["status"]
  1969. if (
  1970. not isinstance(message_id, str)
  1971. or not message_id
  1972. or message_id.strip() != message_id
  1973. or len(message_id.encode("utf-8")) > 255
  1974. ):
  1975. raise EdgePolicyError("acknowledgement message_id is invalid")
  1976. if status != "accepted":
  1977. raise EdgePolicyError("acknowledgement status is invalid")
  1978. try:
  1979. received_at = canonical_timestamp(value["received_at"], "received_at")
  1980. except EdgeContractError as exc:
  1981. raise EdgePolicyError("acknowledgement received_at is invalid") from exc
  1982. return {
  1983. "message_id": message_id,
  1984. "received_at": received_at,
  1985. "status": status,
  1986. }
  1987. @staticmethod
  1988. def _task_record(row: sqlite3.Row) -> QueuedTask:
  1989. task = json.loads(row["task_json"])
  1990. return QueuedTask(
  1991. task_id=row["task_id"],
  1992. digest=row["digest"],
  1993. status=row["status"],
  1994. task=_freeze_mapping(task),
  1995. attempt_count=row["attempt_count"],
  1996. available_at=row["available_at"],
  1997. lease_owner=row["lease_owner"],
  1998. lease_token=row["lease_token"],
  1999. lease_expires_at=row["lease_expires_at"],
  2000. error_code=row["error_code"],
  2001. authority_key_id=row["authority_key_id"],
  2002. authority_digest=row["authority_digest"],
  2003. remote_lease_token_digest=row["remote_lease_token_digest"],
  2004. remote_lease_expires_at=row["remote_lease_expires_at"],
  2005. remote_attempt=row["remote_attempt"],
  2006. created_at=row["created_at"],
  2007. updated_at=row["updated_at"],
  2008. )
  2009. @staticmethod
  2010. def _event_record(row: sqlite3.Row) -> QueuedEvent:
  2011. event = json.loads(row["event_json"])
  2012. acknowledgement = (
  2013. json.loads(row["acknowledgement_json"])
  2014. if row["acknowledgement_json"] is not None
  2015. else None
  2016. )
  2017. return QueuedEvent(
  2018. event_id=row["event_id"],
  2019. digest=row["digest"],
  2020. status=row["status"],
  2021. event=_freeze_mapping(event),
  2022. attempt_count=row["attempt_count"],
  2023. available_at=row["available_at"],
  2024. lease_owner=row["lease_owner"],
  2025. lease_token=row["lease_token"],
  2026. lease_expires_at=row["lease_expires_at"],
  2027. error_code=row["error_code"],
  2028. acknowledgement=(
  2029. _freeze_mapping(acknowledgement)
  2030. if acknowledgement is not None
  2031. else None
  2032. ),
  2033. task_digest=row["task_digest"],
  2034. remote_lease_token_digest=row["remote_lease_token_digest"],
  2035. control_lease_token_digest=row["control_lease_token_digest"],
  2036. created_at=row["created_at"],
  2037. updated_at=row["updated_at"],
  2038. )
  2039. @staticmethod
  2040. def _artifact_record(row: sqlite3.Row) -> LocalArtifact:
  2041. receipt = (
  2042. json.loads(row["deletion_receipt_json"])
  2043. if row["deletion_receipt_json"] is not None
  2044. else None
  2045. )
  2046. return LocalArtifact(
  2047. artifact_digest=row["artifact_digest"],
  2048. artifact_ref=row["artifact_ref"],
  2049. artifact_ref_hash=row["artifact_ref_hash"],
  2050. classification=row["classification"],
  2051. retention_until=row["retention_until"],
  2052. task_id=row["task_id"],
  2053. event_id=row["event_id"],
  2054. cleanup_status=row["cleanup_status"],
  2055. attempt_count=row["attempt_count"],
  2056. available_at=row["available_at"],
  2057. lease_owner=row["lease_owner"],
  2058. lease_token=row["lease_token"],
  2059. lease_expires_at=row["lease_expires_at"],
  2060. error_code=row["error_code"],
  2061. deletion_receipt=_freeze_mapping(receipt) if receipt is not None else None,
  2062. created_at=row["created_at"],
  2063. updated_at=row["updated_at"],
  2064. )
  2065. @staticmethod
  2066. def _outcome_record(row: sqlite3.Row) -> QueuedOutcome:
  2067. outcome = json.loads(row["outcome_json"])
  2068. acknowledgement = (
  2069. json.loads(row["acknowledgement_json"])
  2070. if row["acknowledgement_json"] is not None
  2071. else None
  2072. )
  2073. return QueuedOutcome(
  2074. task_id=row["task_id"],
  2075. digest=row["digest"],
  2076. status=row["status"],
  2077. outcome=_freeze_mapping(outcome),
  2078. attempt_count=row["attempt_count"],
  2079. available_at=row["available_at"],
  2080. lease_owner=row["lease_owner"],
  2081. lease_token=row["lease_token"],
  2082. lease_expires_at=row["lease_expires_at"],
  2083. error_code=row["error_code"],
  2084. acknowledgement=(
  2085. _freeze_mapping(acknowledgement)
  2086. if acknowledgement is not None
  2087. else None
  2088. ),
  2089. remote_lease_token_digest=row["remote_lease_token_digest"],
  2090. created_at=row["created_at"],
  2091. updated_at=row["updated_at"],
  2092. )