test_phase3_wp04_edge_gateway.py 62 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675
  1. from __future__ import annotations
  2. import hashlib
  3. import json
  4. import sqlite3
  5. from concurrent.futures import ThreadPoolExecutor
  6. from dataclasses import FrozenInstanceError
  7. from datetime import UTC, datetime
  8. from threading import Barrier
  9. import pytest
  10. from cryptography.hazmat.primitives import serialization
  11. from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
  12. from app.core.edge_gateway import (
  13. CONTROL_PLANE_ALLOWED,
  14. EDGE_ONLY,
  15. RETENTION_DAYS,
  16. SCHEMA_VERSION,
  17. EdgeContractError,
  18. EdgeEgressPolicy,
  19. EdgeEventContract,
  20. EdgePolicyError,
  21. EdgeQueueConflictError,
  22. EdgeQueueLeaseError,
  23. EdgeQueueSchemaError,
  24. EdgeTaskContract,
  25. SignedTaskEnvelope,
  26. SqliteEdgeQueue,
  27. canonical_sha256,
  28. stable_event_id,
  29. )
  30. from app.core.edge_gateway import queue as edge_queue_module
  31. POLICY_DIGEST = "a" * 64
  32. APPROVED_TASK = {
  33. "task_id": "task-20260809-001",
  34. "gateway_id": "gateway-enterprise-001",
  35. "environment": "production",
  36. "network_zone": "enterprise-zone-a",
  37. "purpose": "governed-data-quality",
  38. "classification": "raw",
  39. "task_type": "profile",
  40. "contract_version": 1,
  41. "deadline_at": "2099-08-09T12:00:00Z",
  42. "attempt": 1,
  43. "idempotency_key": "profile:source-1:20260809",
  44. "policy_digest": POLICY_DIGEST,
  45. }
  46. def _event(
  47. *,
  48. classification: str = "statistics",
  49. payload: dict[str, object] | None = None,
  50. occurred_at: str = "2026-08-09T08:00:00Z",
  51. ) -> dict[str, object]:
  52. content = payload or {"metric_count": 3, "null_ratio": 0.125}
  53. event = {
  54. "task_id": APPROVED_TASK["task_id"],
  55. "gateway_id": APPROVED_TASK["gateway_id"],
  56. "environment": APPROVED_TASK["environment"],
  57. "network_zone": APPROVED_TASK["network_zone"],
  58. "purpose": APPROVED_TASK["purpose"],
  59. "classification": classification,
  60. "contract_version": 1,
  61. "occurred_at": occurred_at,
  62. "attempt": 1,
  63. "idempotency_key": f"event:{APPROVED_TASK['task_id']}:{classification}",
  64. "policy_digest": POLICY_DIGEST,
  65. "payload": content,
  66. }
  67. return {"event_id": stable_event_id(event), **event}
  68. def _ack(*, received_at: str = "2026-08-09T09:00:00+01:00") -> dict[str, str]:
  69. return {
  70. "message_id": "control-message-1",
  71. "received_at": received_at,
  72. "status": "accepted",
  73. }
  74. def _signed_task(
  75. private_key: Ed25519PrivateKey,
  76. *,
  77. task: dict[str, object] | None = None,
  78. issued_at: str = "2026-08-09T07:59:00Z",
  79. expires_at: str = "2099-08-09T12:00:00Z",
  80. ) -> dict[str, object]:
  81. payload = {
  82. "task": task or APPROVED_TASK,
  83. "authority_key_id": "control-authority-2026-01",
  84. "signature_algorithm": "Ed25519",
  85. "contract_digest": canonical_sha256(task or APPROVED_TASK),
  86. "gateway_id": (task or APPROVED_TASK)["gateway_id"],
  87. "environment": (task or APPROVED_TASK)["environment"],
  88. "network_zone": (task or APPROVED_TASK)["network_zone"],
  89. "policy_digest": (task or APPROVED_TASK)["policy_digest"],
  90. "purpose": (task or APPROVED_TASK)["purpose"],
  91. "issued_at": issued_at,
  92. "expires_at": expires_at,
  93. }
  94. unsigned = SignedTaskEnvelope.canonical_unsigned_bytes(payload)
  95. return {**payload, "signature": private_key.sign(unsigned).hex()}
  96. def test_task_contract_is_immutable_strict_and_allows_only_supported_operations():
  97. task = EdgeTaskContract.from_mapping(APPROVED_TASK)
  98. assert task.task_type == "profile"
  99. assert task.operation == "profile"
  100. assert task.to_mapping() == APPROVED_TASK
  101. with pytest.raises(FrozenInstanceError):
  102. task.task_id = "changed" # type: ignore[misc]
  103. with pytest.raises(EdgeContractError, match="unknown properties"):
  104. EdgeTaskContract.from_mapping({**APPROVED_TASK, "parameters": {}})
  105. for operation in ("collect", "profile", "quality", "lineage", "controlled_query"):
  106. assert EdgeTaskContract.from_mapping(
  107. {**APPROVED_TASK, "task_type": operation}
  108. ).task_type == operation
  109. with pytest.raises(EdgeContractError, match="operation"):
  110. EdgeTaskContract.from_mapping({**APPROVED_TASK, "task_type": "shell"})
  111. def test_direct_task_construction_cannot_bypass_shared_validation():
  112. fields = {
  113. **APPROVED_TASK,
  114. "operation": APPROVED_TASK["task_type"],
  115. }
  116. fields.pop("task_type")
  117. valid = EdgeTaskContract(**fields)
  118. assert valid.deadline_at == APPROVED_TASK["deadline_at"]
  119. for changed in (
  120. {"operation": "shell"},
  121. {"classification": "unknown"},
  122. {"contract_version": 2},
  123. {"attempt": 6},
  124. {"policy_digest": "bad"},
  125. {"deadline_at": "not-a-time"},
  126. ):
  127. with pytest.raises(EdgeContractError):
  128. EdgeTaskContract(**{**fields, **changed})
  129. def test_direct_event_construction_and_queue_instance_input_are_revalidated(tmp_path):
  130. mapping = _event()
  131. fields = dict(mapping)
  132. valid = EdgeEventContract(**fields)
  133. assert valid.payload == mapping["payload"]
  134. with pytest.raises(EdgeContractError):
  135. EdgeEventContract(**{**fields, "event_id": "evt_invalid"})
  136. with pytest.raises(EdgeContractError):
  137. EdgeEventContract(**{**fields, "attempt": 6})
  138. task = EdgeTaskContract.from_mapping(APPROVED_TASK)
  139. object.__setattr__(task, "operation", "shell")
  140. with pytest.raises(EdgeContractError):
  141. SqliteEdgeQueue(tmp_path / "edge.db").enqueue(task)
  142. @pytest.mark.parametrize(
  143. ("field", "value"),
  144. [
  145. ("classification", "unknown"),
  146. ("contract_version", 0),
  147. ("attempt", True),
  148. ("attempt", 0),
  149. ("deadline_at", "2099-08-09 12:00:00"),
  150. ("deadline_at", "2099-08-09T12:00:00"),
  151. ("policy_digest", "not-a-sha256"),
  152. ],
  153. )
  154. def test_task_contract_fails_closed_for_malformed_enums_numbers_and_timestamps(
  155. field, value
  156. ):
  157. with pytest.raises(EdgeContractError):
  158. EdgeTaskContract.from_mapping({**APPROVED_TASK, field: value})
  159. def test_equivalent_rfc3339_deadlines_canonicalize_to_one_utc_digest():
  160. offset = EdgeTaskContract.from_mapping(
  161. {**APPROVED_TASK, "deadline_at": "2026-08-09T09:00:05+01:00"}
  162. )
  163. utc = EdgeTaskContract.from_mapping(
  164. {**APPROVED_TASK, "deadline_at": "2026-08-09T08:00:05Z"}
  165. )
  166. fractional = EdgeTaskContract.from_mapping(
  167. {**APPROVED_TASK, "deadline_at": "2026-08-09T08:00:05.120000+00:00"}
  168. )
  169. assert offset.deadline_at == utc.deadline_at == "2026-08-09T08:00:05Z"
  170. assert offset.digest == utc.digest
  171. assert fractional.deadline_at == "2026-08-09T08:00:05.12Z"
  172. def test_equivalent_event_timestamps_canonicalize_before_id_and_digest_validation():
  173. utc_mapping = _event(occurred_at="2026-08-09T08:00:05Z")
  174. offset_mapping = _event(occurred_at="2026-08-09T09:00:05+01:00")
  175. offset_mapping["event_id"] = utc_mapping["event_id"]
  176. utc = EdgeEventContract.from_mapping(utc_mapping)
  177. offset = EdgeEventContract.from_mapping(offset_mapping)
  178. assert offset.occurred_at == utc.occurred_at == "2026-08-09T08:00:05Z"
  179. assert offset.event_id == utc.event_id
  180. assert offset.digest == utc.digest
  181. def test_canonical_sha256_and_event_identifier_are_stable_and_order_independent():
  182. left = {"task_id": "task-1", "payload": {"b": 2, "a": [1, "中"]}}
  183. right = {"payload": {"a": [1, "中"], "b": 2}, "task_id": "task-1"}
  184. assert canonical_sha256(left) == canonical_sha256(right)
  185. assert len(canonical_sha256(left)) == 64
  186. assert stable_event_id(left) == stable_event_id(right)
  187. assert stable_event_id(left).startswith("evt_")
  188. def test_canonical_protocol_has_explicit_numeric_semantics_and_type_boundaries():
  189. assert canonical_sha256(1) == canonical_sha256(1.0)
  190. assert canonical_sha256(0) == canonical_sha256(-0.0)
  191. assert canonical_sha256("1") != canonical_sha256(1)
  192. assert canonical_sha256(True) != canonical_sha256(1)
  193. assert canonical_sha256(None) != canonical_sha256("null")
  194. assert canonical_sha256({"é": 1, "a": 2}) == canonical_sha256(
  195. {"a": 2.0, "é": 1.0}
  196. )
  197. def test_stable_event_identity_binds_every_safe_envelope_dimension():
  198. event = _event()
  199. identity = {key: value for key, value in event.items() if key != "event_id"}
  200. baseline = stable_event_id(identity)
  201. for field, value in {
  202. "task_id": "task-other",
  203. "gateway_id": "gateway-other",
  204. "environment": "staging",
  205. "network_zone": "zone-other",
  206. "purpose": "other-purpose",
  207. "classification": "evidence",
  208. "contract_version": 2,
  209. "attempt": 2,
  210. "idempotency_key": "other-key",
  211. "policy_digest": "b" * 64,
  212. "occurred_at": "2026-08-09T08:00:01Z",
  213. "payload": {"metric_count": 4},
  214. }.items():
  215. assert stable_event_id({**identity, field: value}) != baseline
  216. def test_stable_event_id_normalizes_equivalent_rfc3339_identity_timestamps():
  217. base = {
  218. "task_id": "task-1",
  219. "gateway_id": "gateway-1",
  220. "classification": "statistics",
  221. "payload": {"count": 1},
  222. }
  223. assert stable_event_id(
  224. {**base, "occurred_at": "2026-08-09T09:00:05+01:00"}
  225. ) == stable_event_id({**base, "occurred_at": "2026-08-09T08:00:05Z"})
  226. def test_event_contract_requires_stable_identifier_and_approved_classification():
  227. event = EdgeEventContract.from_mapping(_event())
  228. assert event.event_id == _event()["event_id"]
  229. assert event.classification == "statistics"
  230. with pytest.raises(EdgeContractError, match="stable event identifier"):
  231. EdgeEventContract.from_mapping({**_event(), "event_id": "evt_changed"})
  232. with pytest.raises(EdgeContractError, match="classification"):
  233. EdgeEventContract.from_mapping(_event(classification="raw"))
  234. with pytest.raises(EdgeContractError, match="timestamp"):
  235. EdgeEventContract.from_mapping(_event(occurred_at="yesterday"))
  236. with pytest.raises(EdgeContractError, match="unknown properties"):
  237. EdgeEventContract.from_mapping({**_event(), "raw_result": []})
  238. def test_contract_and_policy_preflight_reject_cycles_and_extreme_depth():
  239. cyclic: dict[str, object] = {}
  240. cyclic["self"] = cyclic
  241. deep: dict[str, object] = {"leaf": 1}
  242. for _ in range(1_100):
  243. deep = {"next": deep}
  244. with pytest.raises(EdgeContractError, match="cycle"):
  245. canonical_sha256(cyclic)
  246. with pytest.raises(EdgeContractError, match="depth"):
  247. canonical_sha256(deep)
  248. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  249. with pytest.raises(EdgePolicyError):
  250. policy.approve_event({"classification": "statistics", "payload": cyclic})
  251. with pytest.raises(EdgePolicyError):
  252. policy.approve_event({"classification": "statistics", "payload": deep})
  253. def test_preflight_rejects_oversized_strings_and_containers_before_encoding():
  254. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  255. with pytest.raises(EdgePolicyError):
  256. policy.approve_event(
  257. {"classification": "statistics", "payload": {"value": "x" * 200_000}}
  258. )
  259. with pytest.raises(EdgePolicyError):
  260. policy.approve_event(
  261. {
  262. "classification": "statistics",
  263. "payload": {"values": list(range(20_000))},
  264. }
  265. )
  266. def test_raw_and_recent_detail_never_cross_the_control_plane_boundary():
  267. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  268. assert {"raw", "recent_detail", "restricted"} == EDGE_ONLY
  269. assert {
  270. "desensitized_metadata",
  271. "statistics",
  272. "lineage",
  273. "evidence",
  274. "health_summary",
  275. "diagnostic_summary",
  276. } == CONTROL_PLANE_ALLOWED
  277. for classification in EDGE_ONLY:
  278. with pytest.raises(EdgePolicyError, match="edge-only"):
  279. policy.approve_event(
  280. {"classification": classification, "payload": {"count": 1}}
  281. )
  282. @pytest.mark.parametrize(
  283. "payload",
  284. [
  285. {"api_key": "abc"},
  286. {"access_token": "abc"},
  287. {"nested": {"password": "abc"}},
  288. {"rows": [{"id": 1}]},
  289. {"customer_rows": [{"id": 1}]},
  290. {"raw_records": [{"id": 1}]},
  291. {"sql": "SELECT * FROM customer"},
  292. {"query_text": "SELECT * FROM customer"},
  293. {"diagnostic": "DELETE FROM customer"},
  294. {"key": "-----BEGIN PRIVATE KEY-----\nabc"},
  295. {"authorization": "Bearer abcdefghijklmnopqrstuvwxyz"},
  296. ],
  297. )
  298. def test_policy_recursively_rejects_secret_raw_detail_sql_and_private_keys(payload):
  299. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  300. with pytest.raises(EdgePolicyError, match="sensitive"):
  301. policy.approve_event({"classification": "statistics", "payload": payload})
  302. @pytest.mark.parametrize(
  303. "payload",
  304. [
  305. {"diagnostic": "connection failed: client_secret = hunter2; retry later"},
  306. {
  307. "diagnostic": (
  308. "planner output follows: SELECT customer_id FROM customers "
  309. "WHERE active = 1"
  310. )
  311. },
  312. {"result": [{"customer_id": "customer-1", "email": "raw@example.com"}]},
  313. ],
  314. )
  315. def test_policy_rejects_reviewed_sensitive_values_and_raw_row_containers(payload):
  316. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  317. with pytest.raises(EdgePolicyError, match="sensitive"):
  318. policy.approve_event({"classification": "diagnostic_summary", "payload": payload})
  319. @pytest.mark.parametrize(
  320. "payload",
  321. [
  322. {"accessToken": "secret-value"},
  323. {"pass-word": "secret-value"},
  324. {"rawRows": [{"id": 1}]},
  325. {"row": {"id": 1}},
  326. {"item": {"id": 1}},
  327. {"diagnostic": "prefix SEL/**/ECT id FR/**/OM customer suffix"},
  328. ],
  329. )
  330. def test_policy_rejects_camelcase_punctuation_single_rows_and_obfuscated_sql(payload):
  331. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  332. with pytest.raises(EdgePolicyError, match="sensitive"):
  333. policy.approve_event({"classification": "diagnostic_summary", "payload": payload})
  334. @pytest.mark.parametrize(
  335. "payload",
  336. [
  337. {"password": "secret-value"},
  338. {"sqlText": "harmless-looking-value"},
  339. {"diagnostic": "prefix EX/**/EC dbo.rotate_secret suffix"},
  340. {"diagnostic": "prefix EXECUTE\n dbo.rotate_secret suffix"},
  341. {"diagnostic": "prefix CALL customer_refresh() suffix"},
  342. ],
  343. )
  344. def test_policy_nfkc_normalizes_keys_and_rejects_stored_procedure_sql(payload):
  345. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  346. with pytest.raises(EdgePolicyError, match="sensitive"):
  347. policy.approve_event({"classification": "diagnostic_summary", "payload": payload})
  348. def test_policy_does_not_misclassify_ordinary_execution_text():
  349. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  350. policy.approve_event(
  351. {
  352. "classification": "diagnostic_summary",
  353. "payload": {
  354. "text": "Execution completed and callback scheduling succeeded."
  355. },
  356. }
  357. )
  358. def test_policy_allows_lineage_graphs_and_aggregate_scalar_containers():
  359. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  360. policy.approve_event(
  361. {
  362. "classification": "lineage",
  363. "payload": {
  364. "nodes": [{"id": "source-a"}, {"id": "product-b"}],
  365. "edges": [{"source": "source-a", "target": "product-b"}],
  366. },
  367. }
  368. )
  369. policy.approve_event(
  370. {
  371. "classification": "statistics",
  372. "payload": {
  373. "rows": 12,
  374. "records": 12,
  375. "result": 12,
  376. "data": "aggregated",
  377. "items": [1, 2],
  378. },
  379. }
  380. )
  381. def test_policy_has_independent_byte_count_depth_and_string_limits():
  382. policy = EdgeEgressPolicy(
  383. allowed_control_hosts={"control.example.com"},
  384. byte_limits={"statistics": 96, "evidence": 512},
  385. count_limits={"statistics": 2, "evidence": 8},
  386. max_depth=3,
  387. max_string_bytes=32,
  388. )
  389. policy.approve_event(
  390. {"classification": "evidence", "payload": {"items": [1, 2, 3]}}
  391. )
  392. with pytest.raises(EdgePolicyError, match="byte limit"):
  393. policy.approve_event(
  394. {"classification": "statistics", "payload": {"value": "x" * 90}}
  395. )
  396. with pytest.raises(EdgePolicyError, match="item count"):
  397. policy.approve_event(
  398. {"classification": "statistics", "payload": {"items": [1, 2, 3]}}
  399. )
  400. with pytest.raises(EdgePolicyError, match="depth"):
  401. policy.approve_event(
  402. {"classification": "evidence", "payload": {"a": {"b": {"c": 1}}}}
  403. )
  404. with pytest.raises(EdgePolicyError, match="string limit"):
  405. policy.approve_event(
  406. {"classification": "evidence", "payload": {"detail": "x" * 33}}
  407. )
  408. @pytest.mark.parametrize(
  409. "url",
  410. [
  411. "http://control.example.com/events",
  412. "https://evil.example.com/events",
  413. "https://control.example.com.evil.test/events",
  414. "https://user:pass@control.example.com/events",
  415. "https://control.example.com.:443/events",
  416. ],
  417. )
  418. def test_destination_requires_https_and_exact_approved_control_host(url):
  419. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  420. with pytest.raises(EdgePolicyError, match="destination"):
  421. policy.validate_destination(url)
  422. def test_proxy_is_optional_but_if_present_requires_exact_https_allowlist_match():
  423. policy = EdgeEgressPolicy(
  424. allowed_control_hosts={"control.example.com"},
  425. allowed_proxy_hosts={"proxy.enterprise.example"},
  426. allowed_proxy_origins={"https://proxy.enterprise.example:8443"},
  427. )
  428. policy.validate_destination(
  429. "https://control.example.com/v1/edge/events",
  430. proxy_url="https://proxy.enterprise.example:8443",
  431. )
  432. with pytest.raises(EdgePolicyError, match="proxy"):
  433. policy.validate_destination(
  434. "https://control.example.com/v1/edge/events",
  435. proxy_url="https://proxy.enterprise.example.evil.test",
  436. )
  437. def test_retention_classes_are_explicit_and_malformed_values_fail_closed():
  438. policy = EdgeEgressPolicy(allowed_control_hosts={"control.example.com"})
  439. assert RETENTION_DAYS == {
  440. "raw": 0,
  441. "recent_detail": 365,
  442. "metadata": 1095,
  443. "evidence": 2190,
  444. }
  445. assert policy.retention_days("raw") == 0
  446. assert policy.retention_days("recent_detail") == 365
  447. with pytest.raises(EdgePolicyError, match="retention"):
  448. policy.retention_days("forever")
  449. def test_queue_requires_explicit_on_disk_path_and_enables_wal(tmp_path):
  450. with pytest.raises(ValueError, match="explicit"):
  451. SqliteEdgeQueue("")
  452. with pytest.raises(ValueError, match="memory"):
  453. SqliteEdgeQueue(":memory:")
  454. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  455. with sqlite3.connect(queue.db_path) as connection:
  456. assert connection.execute("PRAGMA journal_mode").fetchone()[0].lower() == "wal"
  457. assert connection.execute("PRAGMA busy_timeout").fetchone()[0] > 0
  458. assert connection.execute("PRAGMA user_version").fetchone()[0] == SCHEMA_VERSION
  459. def test_queue_fails_closed_for_unknown_or_ambiguous_sqlite_schema(tmp_path):
  460. unknown = tmp_path / "unknown.db"
  461. with sqlite3.connect(unknown) as connection:
  462. connection.execute("PRAGMA user_version = 99")
  463. with pytest.raises(EdgeQueueSchemaError, match="version"):
  464. SqliteEdgeQueue(unknown)
  465. ambiguous = tmp_path / "ambiguous.db"
  466. with sqlite3.connect(ambiguous) as connection:
  467. connection.execute("CREATE TABLE edge_tasks (task_id TEXT PRIMARY KEY)")
  468. with pytest.raises(EdgeQueueSchemaError, match="unversioned"):
  469. SqliteEdgeQueue(ambiguous)
  470. damaged = tmp_path / "damaged.db"
  471. with sqlite3.connect(damaged) as connection:
  472. connection.execute("CREATE TABLE edge_tasks (task_id TEXT PRIMARY KEY)")
  473. connection.execute("PRAGMA user_version = 1")
  474. with pytest.raises(EdgeQueueSchemaError, match="schema"):
  475. SqliteEdgeQueue(damaged)
  476. def test_queue_rejects_versioned_schema_with_wrong_types_constraints_and_indexes(
  477. tmp_path,
  478. ):
  479. db_path = tmp_path / "forged-v1.db"
  480. task_columns = (
  481. "task_id", "digest", "task_json", "status", "attempt_count",
  482. "available_at", "deadline_at", "deadline_epoch_us", "lease_owner",
  483. "lease_token", "lease_expires_at", "error_code", "created_at", "updated_at",
  484. )
  485. event_columns = (
  486. "event_id", "digest", "event_json", "status", "attempt_count",
  487. "available_at", "lease_owner", "lease_token", "lease_expires_at",
  488. "error_code", "acknowledgement_json", "ack_digest", "lease_token_digest",
  489. "created_at", "updated_at",
  490. )
  491. with sqlite3.connect(db_path) as connection:
  492. connection.execute(
  493. "CREATE TABLE edge_tasks ("
  494. + ",".join(f"{column} TEXT" for column in task_columns)
  495. + ")"
  496. )
  497. connection.execute(
  498. "CREATE TABLE edge_outbound_events ("
  499. + ",".join(f"{column} TEXT" for column in event_columns)
  500. + ")"
  501. )
  502. connection.execute(
  503. "CREATE INDEX edge_tasks_claim_idx ON edge_tasks(task_id)"
  504. )
  505. connection.execute(
  506. "CREATE INDEX edge_events_claim_idx ON edge_outbound_events(event_id)"
  507. )
  508. connection.execute("PRAGMA user_version = 1")
  509. with pytest.raises(EdgeQueueSchemaError, match="schema"):
  510. SqliteEdgeQueue(db_path)
  511. def test_queue_exact_task_replay_returns_terminal_and_changed_digest_conflicts(tmp_path):
  512. queue = SqliteEdgeQueue(tmp_path / "edge.db", lease_seconds=30)
  513. first = queue.enqueue(APPROVED_TASK)
  514. claimed = queue.claim("worker-a")
  515. assert claimed is not None
  516. queue.complete(first.task_id, lease_token=claimed.lease_token)
  517. replay = queue.enqueue(APPROVED_TASK)
  518. assert replay.status == "completed"
  519. with pytest.raises(EdgeQueueConflictError, match="digest"):
  520. queue.enqueue({**APPROVED_TASK, "purpose": "different-purpose"})
  521. def test_exact_terminal_replay_wins_after_the_original_task_deadline(tmp_path):
  522. now = [datetime(2099, 8, 9, 11, 59, tzinfo=UTC)]
  523. queue = SqliteEdgeQueue(tmp_path / "edge.db", clock=lambda: now[0])
  524. queue.enqueue(APPROVED_TASK)
  525. claimed = queue.claim("worker")
  526. queue.complete(claimed.task_id, lease_token=claimed.lease_token)
  527. now[0] = datetime(2099, 8, 9, 12, 1, tzinfo=UTC)
  528. replay = queue.enqueue(APPROVED_TASK)
  529. assert replay.status == "completed"
  530. with pytest.raises(EdgeQueueConflictError, match="digest"):
  531. queue.enqueue({**APPROVED_TASK, "purpose": "changed-after-deadline"})
  532. def test_two_connections_atomically_claim_one_task_only(tmp_path):
  533. db_path = tmp_path / "edge.db"
  534. first = SqliteEdgeQueue(db_path)
  535. second = SqliteEdgeQueue(db_path)
  536. first.enqueue(APPROVED_TASK)
  537. claims = [first.claim("worker-a"), second.claim("worker-b")]
  538. assert sum(claim is not None for claim in claims) == 1
  539. claimed = next(claim for claim in claims if claim is not None)
  540. assert claimed.lease_owner in {"worker-a", "worker-b"}
  541. assert claimed.lease_token
  542. assert claimed.lease_expires_at is not None
  543. def test_two_threads_racing_independent_task_connections_have_one_winner(tmp_path):
  544. db_path = tmp_path / "edge.db"
  545. first = SqliteEdgeQueue(db_path)
  546. second = SqliteEdgeQueue(db_path)
  547. first.enqueue(APPROVED_TASK)
  548. barrier = Barrier(2)
  549. def claim(queue, owner):
  550. barrier.wait(timeout=5)
  551. return queue.claim(owner)
  552. with ThreadPoolExecutor(max_workers=2) as pool:
  553. results = list(
  554. pool.map(claim, (first, second), ("worker-a", "worker-b"))
  555. )
  556. assert sum(result is not None for result in results) == 1
  557. def test_claim_atomically_fails_a_task_whose_deadline_passed_while_queued(tmp_path):
  558. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  559. task = {**APPROVED_TASK, "deadline_at": "2026-08-09T08:00:05Z"}
  560. queue = SqliteEdgeQueue(tmp_path / "edge.db", clock=lambda: now[0])
  561. queue.enqueue(task)
  562. now[0] = datetime(2026, 8, 9, 8, 0, 6, tzinfo=UTC)
  563. assert queue.claim("worker") is None
  564. expired = queue.get_task(task["task_id"])
  565. assert expired.status == "failed"
  566. assert expired.error_code == "task_deadline_expired"
  567. def test_claim_fails_an_expired_offset_deadline_instead_of_leasing_it(tmp_path):
  568. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  569. task = {
  570. **APPROVED_TASK,
  571. "deadline_at": "2026-08-09T09:00:05+01:00",
  572. "task_id": "task-offset-deadline",
  573. }
  574. queue = SqliteEdgeQueue(tmp_path / "edge.db", clock=lambda: now[0])
  575. queued = queue.enqueue(task)
  576. assert queued.task["deadline_at"] == "2026-08-09T08:00:05Z"
  577. now[0] = datetime(2026, 8, 9, 8, 0, 6, tzinfo=UTC)
  578. assert queue.claim("worker") is None
  579. expired = queue.get_task(task["task_id"])
  580. assert expired.status == "failed"
  581. assert expired.error_code == "task_deadline_expired"
  582. def test_stale_task_lease_is_reclaimed_and_old_owner_is_fenced(tmp_path):
  583. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  584. queue = SqliteEdgeQueue(
  585. tmp_path / "edge.db", clock=lambda: now[0], lease_seconds=10
  586. )
  587. queue.enqueue(APPROVED_TASK)
  588. old = queue.claim("worker-old")
  589. assert old is not None
  590. now[0] = datetime(2026, 8, 9, 8, 0, 11, tzinfo=UTC)
  591. reclaimed = queue.claim("worker-new")
  592. assert reclaimed is not None
  593. assert reclaimed.lease_token != old.lease_token
  594. with pytest.raises(EdgeQueueLeaseError, match="lease"):
  595. queue.complete(old.task_id, lease_token=old.lease_token)
  596. def test_expired_lease_is_fenced_even_before_another_worker_reclaims_it(tmp_path):
  597. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  598. queue = SqliteEdgeQueue(
  599. tmp_path / "edge.db", clock=lambda: now[0], lease_seconds=10
  600. )
  601. queue.enqueue(APPROVED_TASK)
  602. claimed = queue.claim("worker-old")
  603. now[0] = datetime(2026, 8, 9, 8, 0, 11, tzinfo=UTC)
  604. with pytest.raises(EdgeQueueLeaseError, match="expired"):
  605. queue.complete(claimed.task_id, lease_token=claimed.lease_token)
  606. def test_expired_final_attempt_is_terminal_instead_of_remaining_leased(tmp_path):
  607. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  608. queue = SqliteEdgeQueue(
  609. tmp_path / "edge.db",
  610. clock=lambda: now[0],
  611. lease_seconds=10,
  612. max_attempts=1,
  613. )
  614. queue.enqueue(APPROVED_TASK)
  615. queue.claim("worker")
  616. now[0] = datetime(2026, 8, 9, 8, 0, 11, tzinfo=UTC)
  617. assert queue.claim("next-worker") is None
  618. assert queue.get_task(APPROVED_TASK["task_id"]).status == "failed"
  619. def test_fail_uses_bounded_exponential_retry_then_becomes_terminal(tmp_path):
  620. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  621. queue = SqliteEdgeQueue(
  622. tmp_path / "edge.db",
  623. clock=lambda: now[0],
  624. lease_seconds=10,
  625. max_attempts=3,
  626. retry_base_seconds=2,
  627. retry_max_seconds=3,
  628. )
  629. queue.enqueue(APPROVED_TASK)
  630. first = queue.claim("worker")
  631. retried = queue.fail(
  632. first.task_id, lease_token=first.lease_token, error_code="source_unavailable"
  633. )
  634. assert retried.status == "pending"
  635. assert retried.available_at == "2026-08-09T08:00:02Z"
  636. assert queue.claim("too-early") is None
  637. now[0] = datetime(2026, 8, 9, 8, 0, 2, tzinfo=UTC)
  638. second = queue.claim("worker")
  639. queue.fail(
  640. second.task_id, lease_token=second.lease_token, error_code="source_unavailable"
  641. )
  642. now[0] = datetime(2026, 8, 9, 8, 0, 5, tzinfo=UTC)
  643. third = queue.claim("worker")
  644. terminal = queue.fail(
  645. third.task_id, lease_token=third.lease_token, error_code="source_unavailable"
  646. )
  647. assert terminal.status == "failed"
  648. assert terminal.attempt_count == 3
  649. assert queue.enqueue(APPROVED_TASK).status == "failed"
  650. def test_cancel_is_atomic_and_prevents_lease_completion(tmp_path):
  651. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  652. queue.enqueue(APPROVED_TASK)
  653. claimed = queue.claim("worker")
  654. cancelled = queue.cancel(claimed.task_id)
  655. assert cancelled.status == "cancelled"
  656. assert queue.cancel(claimed.task_id).status == "cancelled"
  657. with pytest.raises(EdgeQueueLeaseError):
  658. queue.complete(claimed.task_id, lease_token=claimed.lease_token)
  659. def test_cancel_and_complete_race_reaches_one_consistent_terminal_state(tmp_path):
  660. db_path = tmp_path / "edge.db"
  661. first = SqliteEdgeQueue(db_path)
  662. second = SqliteEdgeQueue(db_path)
  663. first.enqueue(APPROVED_TASK)
  664. claimed = first.claim("worker")
  665. barrier = Barrier(2)
  666. def cancel():
  667. barrier.wait(timeout=5)
  668. return second.cancel(claimed.task_id).status
  669. def complete():
  670. barrier.wait(timeout=5)
  671. try:
  672. return first.complete(
  673. claimed.task_id, lease_token=claimed.lease_token
  674. ).status
  675. except EdgeQueueLeaseError:
  676. return "lease_rejected"
  677. with ThreadPoolExecutor(max_workers=2) as pool:
  678. results = [pool.submit(cancel), pool.submit(complete)]
  679. outcomes = {result.result(timeout=5) for result in results}
  680. final = first.get_task(claimed.task_id)
  681. assert final.status in {"cancelled", "completed"}
  682. assert "completed" not in outcomes or final.status == "completed"
  683. def test_transaction_exception_rolls_back_and_queue_remains_usable(tmp_path):
  684. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  685. queue.enqueue(APPROVED_TASK)
  686. with (
  687. pytest.raises(RuntimeError, match="injected"),
  688. queue._transaction() as connection,
  689. ):
  690. connection.execute(
  691. "UPDATE edge_tasks SET status = 'failed' WHERE task_id = ?",
  692. (APPROVED_TASK["task_id"],),
  693. )
  694. raise RuntimeError("injected lock-path failure")
  695. assert queue.get_task(APPROVED_TASK["task_id"]).status == "pending"
  696. assert queue.claim("worker") is not None
  697. def test_outbound_event_is_persisted_before_claim_and_exact_ack_replays(tmp_path):
  698. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  699. event = _event()
  700. persisted = queue.persist_event(event)
  701. assert persisted.status == "pending"
  702. with sqlite3.connect(queue.db_path) as connection:
  703. stored = connection.execute(
  704. "SELECT event_json FROM edge_outbound_events WHERE event_id = ?",
  705. (event["event_id"],),
  706. ).fetchone()
  707. assert json.loads(stored[0]) == event
  708. sending = queue.claim_event("sender-a")
  709. acknowledged = queue.acknowledge_event(
  710. sending.event_id,
  711. lease_token=sending.lease_token,
  712. acknowledgement=_ack(),
  713. )
  714. assert acknowledged.status == "acknowledged"
  715. replay = queue.persist_event(event)
  716. assert replay.status == "acknowledged"
  717. assert replay.acknowledgement == {
  718. "message_id": "control-message-1",
  719. "received_at": "2026-08-09T08:00:00Z",
  720. "status": "accepted",
  721. }
  722. def test_acknowledgement_exact_replay_binds_ack_and_original_lease_token(tmp_path):
  723. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  724. queue.persist_event(_event())
  725. sending = queue.claim_event("sender")
  726. first = queue.acknowledge_event(
  727. sending.event_id,
  728. lease_token=sending.lease_token,
  729. acknowledgement=_ack(),
  730. )
  731. replay = queue.acknowledge_event(
  732. sending.event_id,
  733. lease_token=sending.lease_token,
  734. acknowledgement=_ack(received_at="2026-08-09T08:00:00Z"),
  735. )
  736. assert replay.acknowledgement == first.acknowledgement
  737. with pytest.raises(EdgeQueueLeaseError):
  738. queue.acknowledge_event(
  739. sending.event_id,
  740. lease_token="different-token",
  741. acknowledgement=_ack(received_at="2026-08-09T08:00:00Z"),
  742. )
  743. with pytest.raises(EdgeQueueConflictError):
  744. queue.acknowledge_event(
  745. sending.event_id,
  746. lease_token=sending.lease_token,
  747. acknowledgement={
  748. **_ack(received_at="2026-08-09T08:00:00Z"),
  749. "message_id": "different-message",
  750. },
  751. )
  752. with sqlite3.connect(queue.db_path) as connection:
  753. row = connection.execute(
  754. """
  755. SELECT lease_token, ack_digest, lease_token_digest
  756. FROM edge_outbound_events WHERE event_id = ?
  757. """,
  758. (sending.event_id,),
  759. ).fetchone()
  760. assert row[0] is None
  761. assert row[1] == canonical_sha256(dict(first.acknowledgement))
  762. assert row[2] == canonical_sha256(sending.lease_token)
  763. def test_acknowledgement_requires_exact_fields_timestamp_and_status(tmp_path):
  764. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  765. queue.persist_event(_event())
  766. sending = queue.claim_event("sender")
  767. for acknowledgement in (
  768. {"message_id": "id", "status": "accepted"},
  769. {**_ack(), "extra": "no"},
  770. {**_ack(), "received_at": "yesterday"},
  771. {**_ack(), "status": "unknown"},
  772. ):
  773. with pytest.raises(EdgePolicyError):
  774. queue.acknowledge_event(
  775. sending.event_id,
  776. lease_token=sending.lease_token,
  777. acknowledgement=acknowledgement,
  778. )
  779. def test_event_persistence_digests_normalized_contract_before_exact_replay(tmp_path):
  780. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  781. utc_event = _event(occurred_at="2026-08-09T08:00:05Z")
  782. offset_event = _event(occurred_at="2026-08-09T09:00:05+01:00")
  783. offset_event["event_id"] = utc_event["event_id"]
  784. contract = EdgeEventContract.from_mapping(offset_event)
  785. persisted = queue.persist_event(offset_event)
  786. replay = queue.persist_event(utc_event)
  787. assert replay.event_id == persisted.event_id
  788. assert replay.digest == persisted.digest == contract.digest
  789. assert replay.event["occurred_at"] == "2026-08-09T08:00:05Z"
  790. with sqlite3.connect(queue.db_path) as connection:
  791. stored_digest = connection.execute(
  792. "SELECT digest FROM edge_outbound_events WHERE event_id = ?",
  793. (contract.event_id,),
  794. ).fetchone()[0]
  795. assert stored_digest == contract.digest
  796. def test_outbound_event_changed_digest_conflicts_and_unsafe_payload_never_persists(
  797. tmp_path,
  798. ):
  799. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  800. event = _event()
  801. queue.persist_event(event)
  802. changed = {**event, "payload": {"metric_count": 4}}
  803. with pytest.raises(EdgeQueueConflictError, match="digest"):
  804. queue.persist_event(changed)
  805. unsafe = _event(payload={"rows": [{"password": "secret"}]})
  806. with pytest.raises((EdgeContractError, EdgePolicyError)):
  807. queue.persist_event(unsafe)
  808. with sqlite3.connect(queue.db_path) as connection:
  809. serialized = " ".join(
  810. str(value)
  811. for row in connection.execute(
  812. "SELECT event_json FROM edge_outbound_events"
  813. ).fetchall()
  814. for value in row
  815. ).lower()
  816. assert "password" not in serialized
  817. assert "secret" not in serialized
  818. def test_reviewed_sensitive_payloads_are_rejected_before_event_persistence(tmp_path):
  819. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  820. payloads = [
  821. {"diagnostic": "upstream says api_key: sk-sensitive-value"},
  822. {"diagnostic": "prefix INSERT INTO customer VALUES (1) suffix"},
  823. {"data": [{"customer_id": "customer-1", "name": "Raw Name"}]},
  824. ]
  825. for payload in payloads:
  826. with pytest.raises(EdgePolicyError, match="sensitive"):
  827. queue.persist_event(
  828. _event(classification="diagnostic_summary", payload=payload)
  829. )
  830. with sqlite3.connect(queue.db_path) as connection:
  831. assert connection.execute(
  832. "SELECT COUNT(*) FROM edge_outbound_events"
  833. ).fetchone()[0] == 0
  834. def test_two_connections_atomically_claim_one_outbound_event_only(tmp_path):
  835. db_path = tmp_path / "edge.db"
  836. first = SqliteEdgeQueue(db_path)
  837. second = SqliteEdgeQueue(db_path)
  838. first.persist_event(_event())
  839. claims = [first.claim_event("sender-a"), second.claim_event("sender-b")]
  840. assert sum(claim is not None for claim in claims) == 1
  841. claimed = next(claim for claim in claims if claim is not None)
  842. assert claimed.lease_owner in {"sender-a", "sender-b"}
  843. def test_two_threads_racing_independent_event_connections_have_one_winner(tmp_path):
  844. db_path = tmp_path / "edge.db"
  845. first = SqliteEdgeQueue(db_path)
  846. second = SqliteEdgeQueue(db_path)
  847. first.persist_event(_event())
  848. barrier = Barrier(2)
  849. def claim(queue, owner):
  850. barrier.wait(timeout=5)
  851. return queue.claim_event(owner)
  852. with ThreadPoolExecutor(max_workers=2) as pool:
  853. results = list(
  854. pool.map(claim, (first, second), ("sender-a", "sender-b"))
  855. )
  856. assert sum(result is not None for result in results) == 1
  857. def test_event_retry_is_bounded_and_uses_independent_event_attempts(tmp_path):
  858. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  859. queue = SqliteEdgeQueue(
  860. tmp_path / "edge.db",
  861. clock=lambda: now[0],
  862. max_attempts=2,
  863. retry_base_seconds=1,
  864. retry_max_seconds=2,
  865. )
  866. queue.persist_event(_event())
  867. first = queue.claim_event("sender")
  868. retry = queue.fail_event(
  869. first.event_id, lease_token=first.lease_token, error_code="network_unavailable"
  870. )
  871. assert retry.status == "pending"
  872. now[0] = datetime(2026, 8, 9, 8, 0, 1, tzinfo=UTC)
  873. second = queue.claim_event("sender")
  874. failed = queue.fail_event(
  875. second.event_id,
  876. lease_token=second.lease_token,
  877. error_code="network_unavailable",
  878. )
  879. assert failed.status == "failed"
  880. assert failed.attempt_count == 2
  881. def test_signed_task_envelope_verifies_canonical_ed25519_authority_and_bindings():
  882. private = Ed25519PrivateKey.generate()
  883. public = private.public_key().public_bytes(
  884. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  885. )
  886. raw = _signed_task(private)
  887. envelope = SignedTaskEnvelope.verify_mapping(
  888. raw,
  889. authority_keys={"control-authority-2026-01": public},
  890. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  891. )
  892. assert envelope.task.to_mapping() == APPROVED_TASK
  893. assert envelope.contract_digest == canonical_sha256(APPROVED_TASK)
  894. assert envelope.to_mapping() == raw
  895. for changed in (
  896. {**raw, "purpose": "other-purpose"},
  897. {**raw, "contract_digest": "b" * 64},
  898. {**raw, "signature": "00" * 64},
  899. {**raw, "extra": "rejected"},
  900. ):
  901. with pytest.raises(EdgeContractError):
  902. SignedTaskEnvelope.verify_mapping(
  903. changed,
  904. authority_keys={"control-authority-2026-01": public},
  905. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  906. )
  907. with pytest.raises(EdgeContractError, match="expired"):
  908. SignedTaskEnvelope.verify_mapping(
  909. _signed_task(private, expires_at="2026-08-09T07:59:59Z"),
  910. authority_keys={"control-authority-2026-01": public},
  911. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  912. )
  913. def test_signed_task_authority_clock_skew_only_tolerates_bounded_future_issue_time():
  914. private = Ed25519PrivateKey.generate()
  915. public = private.public_key().public_bytes(
  916. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  917. )
  918. now = datetime(2026, 8, 9, 8, 0, tzinfo=UTC)
  919. accepted = SignedTaskEnvelope.verify_mapping(
  920. _signed_task(private, issued_at="2026-08-09T08:01:00Z"),
  921. authority_keys={"control-authority-2026-01": public},
  922. now=now,
  923. allowed_future_skew_seconds=60,
  924. )
  925. assert accepted.issued_at == "2026-08-09T08:01:00Z"
  926. with pytest.raises(EdgeContractError, match="not yet valid"):
  927. SignedTaskEnvelope.verify_mapping(
  928. _signed_task(private, issued_at="2026-08-09T08:01:01Z"),
  929. authority_keys={"control-authority-2026-01": public},
  930. now=now,
  931. allowed_future_skew_seconds=60,
  932. )
  933. with pytest.raises(EdgeContractError, match="expired"):
  934. SignedTaskEnvelope.verify_mapping(
  935. _signed_task(
  936. private,
  937. issued_at="2026-08-09T07:59:00Z",
  938. expires_at="2026-08-09T08:00:00Z",
  939. ),
  940. authority_keys={"control-authority-2026-01": public},
  941. now=now,
  942. allowed_future_skew_seconds=60,
  943. )
  944. for invalid in (-1, 301, True):
  945. with pytest.raises(EdgeContractError, match="clock skew"):
  946. SignedTaskEnvelope.verify_mapping(
  947. _signed_task(private),
  948. authority_keys={"control-authority-2026-01": public},
  949. now=now,
  950. allowed_future_skew_seconds=invalid,
  951. )
  952. @pytest.mark.parametrize("task_status", ["pending", "leased"])
  953. def test_v1_queue_refuses_active_unsigned_tasks_without_partial_schema_change(
  954. tmp_path, task_status
  955. ):
  956. db_path = tmp_path / "edge.db"
  957. with sqlite3.connect(db_path) as connection:
  958. connection.execute(edge_queue_module._V1_TASK_TABLE_SQL)
  959. connection.execute(edge_queue_module._TASK_INDEX_SQL)
  960. connection.execute(edge_queue_module._V1_EVENT_TABLE_SQL)
  961. connection.execute(edge_queue_module._EVENT_INDEX_SQL)
  962. connection.execute(
  963. """INSERT INTO edge_tasks (
  964. task_id,digest,task_json,status,attempt_count,available_at,deadline_at,
  965. deadline_epoch_us,lease_owner,lease_token,lease_expires_at,created_at,updated_at
  966. ) VALUES (?,?,?,?,0,?,?,?,?,?,?,?,?)""",
  967. (
  968. APPROVED_TASK["task_id"], canonical_sha256(APPROVED_TASK),
  969. json.dumps(APPROVED_TASK, sort_keys=True, separators=(",", ":")),
  970. task_status,
  971. "2026-08-09T08:00:00Z", APPROVED_TASK["deadline_at"],
  972. 4090060800000000,
  973. "worker" if task_status == "leased" else None,
  974. "local-lease" if task_status == "leased" else None,
  975. "2026-08-09T08:01:00Z" if task_status == "leased" else None,
  976. "2026-08-09T08:00:00Z", "2026-08-09T08:00:00Z",
  977. ),
  978. )
  979. connection.execute("PRAGMA user_version = 1")
  980. with pytest.raises(EdgeQueueSchemaError, match="drain"):
  981. SqliteEdgeQueue(db_path)
  982. with sqlite3.connect(db_path) as connection:
  983. assert connection.execute("PRAGMA user_version").fetchone()[0] == 1
  984. assert {
  985. row[1] for row in connection.execute("PRAGMA table_info(edge_tasks)")
  986. } == {
  987. row[0] for row in edge_queue_module._V1_EXPECTED_TABLE_INFO["edge_tasks"]
  988. }
  989. assert connection.execute(
  990. "SELECT COUNT(*) FROM sqlite_master WHERE name='edge_local_artifacts'"
  991. ).fetchone()[0] == 0
  992. @pytest.mark.parametrize("event_status", ["pending", "sending", "failed"])
  993. def test_v1_queue_refuses_unacknowledged_event_without_partial_migration(
  994. tmp_path, event_status
  995. ):
  996. db_path = tmp_path / "edge.db"
  997. event = _event(occurred_at="2026-08-09T08:00:00Z")
  998. with sqlite3.connect(db_path) as connection:
  999. connection.execute(edge_queue_module._V1_TASK_TABLE_SQL)
  1000. connection.execute(edge_queue_module._TASK_INDEX_SQL)
  1001. connection.execute(edge_queue_module._V1_EVENT_TABLE_SQL)
  1002. connection.execute(edge_queue_module._EVENT_INDEX_SQL)
  1003. connection.execute(
  1004. """INSERT INTO edge_outbound_events (
  1005. event_id,digest,event_json,status,attempt_count,available_at,
  1006. created_at,updated_at
  1007. ) VALUES (?,?,?,?,0,?,?,?)""",
  1008. (
  1009. event["event_id"],
  1010. canonical_sha256(event),
  1011. json.dumps(event, sort_keys=True, separators=(",", ":")),
  1012. event_status,
  1013. "2026-08-09T08:00:00Z",
  1014. "2026-08-09T08:00:00Z",
  1015. "2026-08-09T08:00:00Z",
  1016. ),
  1017. )
  1018. connection.execute("PRAGMA user_version = 1")
  1019. with pytest.raises(EdgeQueueSchemaError, match="drain"):
  1020. SqliteEdgeQueue(db_path)
  1021. with sqlite3.connect(db_path) as connection:
  1022. assert connection.execute("PRAGMA user_version").fetchone()[0] == 1
  1023. assert connection.execute(
  1024. "SELECT COUNT(*) FROM pragma_table_info('edge_outbound_events') "
  1025. "WHERE name='task_digest'"
  1026. ).fetchone()[0] == 0
  1027. def test_v1_queue_migrates_terminal_tasks_and_acknowledged_events_then_accepts_signed(
  1028. tmp_path,
  1029. ):
  1030. db_path = tmp_path / "edge.db"
  1031. event = _event(occurred_at="2026-08-09T08:00:00Z")
  1032. with sqlite3.connect(db_path) as connection:
  1033. connection.execute(edge_queue_module._V1_TASK_TABLE_SQL)
  1034. connection.execute(edge_queue_module._TASK_INDEX_SQL)
  1035. connection.execute(edge_queue_module._V1_EVENT_TABLE_SQL)
  1036. connection.execute(edge_queue_module._EVENT_INDEX_SQL)
  1037. connection.execute(
  1038. """INSERT INTO edge_tasks (
  1039. task_id,digest,task_json,status,attempt_count,available_at,deadline_at,
  1040. deadline_epoch_us,created_at,updated_at
  1041. ) VALUES (?,?,?,'completed',1,?,?,?,?,?)""",
  1042. (
  1043. APPROVED_TASK["task_id"], canonical_sha256(APPROVED_TASK),
  1044. json.dumps(APPROVED_TASK, sort_keys=True, separators=(",", ":")),
  1045. "2026-08-09T08:00:00Z", APPROVED_TASK["deadline_at"],
  1046. 4090060800000000, "2026-08-09T08:00:00Z", "2026-08-09T08:00:00Z",
  1047. ),
  1048. )
  1049. connection.execute(
  1050. """INSERT INTO edge_outbound_events (
  1051. event_id,digest,event_json,status,attempt_count,available_at,
  1052. acknowledgement_json,ack_digest,created_at,updated_at
  1053. ) VALUES (?,?,?,'acknowledged',1,?,?,?,?,?)""",
  1054. (
  1055. event["event_id"],
  1056. canonical_sha256(event),
  1057. json.dumps(event, sort_keys=True, separators=(",", ":")),
  1058. "2026-08-09T08:00:00Z",
  1059. '{"status":"accepted"}',
  1060. "b" * 64,
  1061. "2026-08-09T08:00:00Z",
  1062. "2026-08-09T08:00:00Z",
  1063. ),
  1064. )
  1065. connection.execute("PRAGMA user_version = 1")
  1066. migrated = SqliteEdgeQueue(db_path)
  1067. assert SCHEMA_VERSION == 2
  1068. assert migrated.get_task(APPROVED_TASK["task_id"]).status == "completed"
  1069. assert migrated.get_event(event["event_id"]).status == "acknowledged"
  1070. private = Ed25519PrivateKey.generate()
  1071. public = private.public_key().public_bytes(
  1072. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  1073. )
  1074. task = {**APPROVED_TASK, "task_id": "task-signed-v2"}
  1075. envelope = SignedTaskEnvelope.verify_mapping(
  1076. _signed_task(private, task=task),
  1077. authority_keys={"control-authority-2026-01": public},
  1078. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  1079. )
  1080. queued = migrated.enqueue_signed(envelope)
  1081. assert queued.authority_key_id == "control-authority-2026-01"
  1082. with sqlite3.connect(db_path) as connection:
  1083. row = connection.execute(
  1084. "SELECT authority_key_id,authority_digest,authority_json FROM edge_tasks WHERE task_id=?",
  1085. (task["task_id"],),
  1086. ).fetchone()
  1087. assert row[0] == "control-authority-2026-01"
  1088. assert row[1] == envelope.digest
  1089. assert "signature" in row[2]
  1090. def test_remote_lease_binding_and_atomic_event_artifact_completion(tmp_path):
  1091. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1092. queue = SqliteEdgeQueue(tmp_path / "edge.db", clock=lambda: now[0])
  1093. queue.enqueue(APPROVED_TASK)
  1094. claimed = queue.claim("worker")
  1095. bound = queue.bind_remote_lease(
  1096. claimed.task_id,
  1097. lease_token=claimed.lease_token,
  1098. remote_lease_token="remote-control-lease-1",
  1099. remote_lease_expires_at="2026-08-09T08:05:00Z",
  1100. remote_attempt=claimed.attempt_count,
  1101. )
  1102. assert bound.remote_lease_token_digest == hashlib.sha256(b"remote-control-lease-1").hexdigest()
  1103. assert bound.remote_lease_expires_at == "2026-08-09T08:05:00Z"
  1104. event = _event(occurred_at="2026-08-09T08:00:00Z")
  1105. artifact_ref = str((tmp_path / "artifact.bin").resolve())
  1106. task, persisted = queue.complete_with_event(
  1107. claimed.task_id,
  1108. lease_token=claimed.lease_token,
  1109. remote_lease_token="remote-control-lease-1",
  1110. event=event,
  1111. artifacts=[{
  1112. "artifact_digest": "f" * 64,
  1113. "artifact_ref": artifact_ref,
  1114. "artifact_ref_hash": canonical_sha256(artifact_ref),
  1115. "classification": "raw",
  1116. "retention_until": "2026-08-09T08:00:00Z",
  1117. }],
  1118. )
  1119. assert task.status == "completed"
  1120. assert persisted.status == "pending"
  1121. artifacts = queue.list_local_artifacts(task_id=claimed.task_id)
  1122. assert artifacts[0].artifact_ref == artifact_ref
  1123. assert artifacts[0].cleanup_status == "pending"
  1124. with sqlite3.connect(queue.db_path) as connection:
  1125. columns = {row[1] for row in connection.execute("PRAGMA table_info(edge_local_artifacts)")}
  1126. assert "body" not in columns and "content" not in columns
  1127. def test_atomic_completion_rechecks_cancel_deadline_and_remote_lease(tmp_path):
  1128. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1129. queue = SqliteEdgeQueue(tmp_path / "edge.db", clock=lambda: now[0], lease_seconds=600)
  1130. queue.enqueue({**APPROVED_TASK, "deadline_at": "2026-08-09T08:00:05Z"})
  1131. claimed = queue.claim("worker")
  1132. queue.bind_remote_lease(
  1133. claimed.task_id, lease_token=claimed.lease_token,
  1134. remote_lease_token="remote", remote_lease_expires_at="2026-08-09T08:00:04Z",
  1135. remote_attempt=1,
  1136. )
  1137. now[0] = datetime(2026, 8, 9, 8, 0, 5, tzinfo=UTC)
  1138. with pytest.raises(EdgeQueueLeaseError):
  1139. queue.complete_with_event(
  1140. claimed.task_id, lease_token=claimed.lease_token,
  1141. remote_lease_token="remote", event=_event(), artifacts=[],
  1142. )
  1143. assert queue.get_event(_event()["event_id"]) is None
  1144. second_task = {**APPROVED_TASK, "task_id": "task-cancelled"}
  1145. queue.enqueue(second_task)
  1146. second = queue.claim("worker")
  1147. queue.bind_remote_lease(
  1148. second.task_id, lease_token=second.lease_token,
  1149. remote_lease_token="remote-2", remote_lease_expires_at="2099-01-01T00:00:00Z",
  1150. remote_attempt=1,
  1151. )
  1152. queue.cancel(second.task_id)
  1153. cancelled_event = _event()
  1154. cancelled_event["task_id"] = second.task_id
  1155. cancelled_event["event_id"] = stable_event_id({k: v for k, v in cancelled_event.items() if k != "event_id"})
  1156. with pytest.raises(EdgeQueueLeaseError):
  1157. queue.complete_with_event(
  1158. second.task_id, lease_token=second.lease_token,
  1159. remote_lease_token="remote-2", event=cancelled_event, artifacts=[],
  1160. )
  1161. assert queue.get_event(cancelled_event["event_id"]) is None
  1162. def test_bound_event_ack_requires_event_digest_and_remote_control_leases(tmp_path):
  1163. queue = SqliteEdgeQueue(tmp_path / "edge.db")
  1164. queue.enqueue(APPROVED_TASK)
  1165. claimed = queue.claim("worker")
  1166. queue.bind_remote_lease(
  1167. claimed.task_id, lease_token=claimed.lease_token,
  1168. remote_lease_token="remote", remote_lease_expires_at="2099-01-01T00:00:00Z",
  1169. remote_attempt=1,
  1170. )
  1171. _, persisted = queue.complete_with_event(
  1172. claimed.task_id, lease_token=claimed.lease_token,
  1173. remote_lease_token="remote", event=_event(), artifacts=[],
  1174. )
  1175. sending = queue.claim_event("sender")
  1176. acknowledgement = {
  1177. "event_id": sending.event_id,
  1178. "event_digest": sending.digest,
  1179. "remote_lease_digest": hashlib.sha256(b"remote").hexdigest(),
  1180. "received_at": "2026-08-09T08:00:00Z",
  1181. "status": "accepted",
  1182. }
  1183. acknowledged = queue.acknowledge_event(
  1184. persisted.event_id, lease_token=sending.lease_token,
  1185. acknowledgement=acknowledgement,
  1186. )
  1187. assert acknowledged.status == "acknowledged"
  1188. with pytest.raises(EdgePolicyError):
  1189. queue.acknowledge_event(
  1190. persisted.event_id, lease_token=sending.lease_token,
  1191. acknowledgement={**acknowledgement, "remote_lease_digest": "0" * 64},
  1192. )
  1193. def test_retry_jitter_is_injected_and_used_for_task_and_event(tmp_path):
  1194. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1195. queue = SqliteEdgeQueue(
  1196. tmp_path / "edge.db", clock=lambda: now[0], retry_base_seconds=10,
  1197. retry_max_seconds=100, random_source=lambda: 0.0,
  1198. )
  1199. queue.enqueue(APPROVED_TASK)
  1200. task = queue.claim("worker")
  1201. retried = queue.fail(task.task_id, lease_token=task.lease_token, error_code="temporary")
  1202. assert retried.available_at == "2026-08-09T08:00:05Z"
  1203. queue.persist_event(_event())
  1204. event = queue.claim_event("sender")
  1205. event_retry = queue.fail_event(event.event_id, lease_token=event.lease_token, error_code="temporary")
  1206. assert event_retry.available_at == "2026-08-09T08:00:05Z"
  1207. def test_signed_acceptance_encrypts_recoverable_remote_lease_atomically(tmp_path):
  1208. private = Ed25519PrivateKey.generate()
  1209. public = private.public_key().public_bytes(
  1210. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  1211. )
  1212. envelope = SignedTaskEnvelope.verify_mapping(
  1213. _signed_task(private),
  1214. authority_keys={"control-authority-2026-01": public},
  1215. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  1216. )
  1217. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1218. queue = SqliteEdgeQueue(
  1219. tmp_path / "edge.db",
  1220. clock=lambda: now[0],
  1221. remote_lease_encoder=lambda token: "cipher:" + token[::-1],
  1222. remote_lease_decoder=lambda value: value.removeprefix("cipher:")[::-1],
  1223. )
  1224. queued = queue.accept_signed_task(
  1225. envelope,
  1226. remote_lease_token="remote-secret",
  1227. remote_lease_expires_at="2026-08-09T08:05:00Z",
  1228. remote_attempt=1,
  1229. )
  1230. assert queued.remote_lease_token_digest == hashlib.sha256(b"remote-secret").hexdigest()
  1231. assert "remote-secret" not in repr(queued)
  1232. assert queue.recover_remote_lease(queued.task_id) == "remote-secret"
  1233. with sqlite3.connect(queue.db_path) as connection:
  1234. stored = connection.execute(
  1235. "SELECT remote_lease_ciphertext FROM edge_tasks WHERE task_id=?",
  1236. (queued.task_id,),
  1237. ).fetchone()[0]
  1238. assert stored != "remote-secret"
  1239. now[0] = datetime(2026, 8, 9, 8, 6, tzinfo=UTC)
  1240. renewed = queue.accept_signed_task(
  1241. envelope,
  1242. remote_lease_token="renewed-secret",
  1243. remote_lease_expires_at="2026-08-09T08:10:00Z",
  1244. remote_attempt=1,
  1245. )
  1246. assert renewed.authority_digest == queued.authority_digest
  1247. assert queue.recover_remote_lease(queued.task_id) == "renewed-secret"
  1248. def test_remote_attempt_remains_contract_attempt_across_local_crash_recovery(tmp_path):
  1249. private = Ed25519PrivateKey.generate()
  1250. public = private.public_key().public_bytes(
  1251. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  1252. )
  1253. envelope = SignedTaskEnvelope.verify_mapping(
  1254. _signed_task(private),
  1255. authority_keys={"control-authority-2026-01": public},
  1256. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  1257. )
  1258. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1259. queue = SqliteEdgeQueue(
  1260. tmp_path / "edge.db",
  1261. clock=lambda: now[0],
  1262. lease_seconds=5,
  1263. remote_lease_encoder=lambda token: "cipher:" + token[::-1],
  1264. remote_lease_decoder=lambda value: value.removeprefix("cipher:")[::-1],
  1265. )
  1266. queue.accept_signed_task(
  1267. envelope,
  1268. remote_lease_token="remote-a",
  1269. remote_lease_expires_at="2026-08-09T08:00:04Z",
  1270. remote_attempt=1,
  1271. )
  1272. crashed_claim = queue.claim("crashed-agent")
  1273. assert crashed_claim.attempt_count == 1
  1274. now[0] = datetime(2026, 8, 9, 8, 0, 6, tzinfo=UTC)
  1275. queue.accept_signed_task(
  1276. envelope,
  1277. remote_lease_token="remote-b",
  1278. remote_lease_expires_at="2026-08-09T08:05:00Z",
  1279. remote_attempt=1,
  1280. )
  1281. recovered_claim = queue.claim("replacement-agent")
  1282. assert recovered_claim.attempt_count == 2
  1283. assert recovered_claim.remote_attempt == 1
  1284. completed, event = queue.complete_with_event(
  1285. recovered_claim.task_id,
  1286. lease_token=recovered_claim.lease_token,
  1287. remote_lease_token="remote-b",
  1288. event=_event(occurred_at="2026-08-09T08:00:06Z"),
  1289. )
  1290. assert completed.status == "completed"
  1291. assert event.attempt_count == 0
  1292. with sqlite3.connect(queue.db_path) as connection:
  1293. assert connection.execute(
  1294. "SELECT count(*) FROM edge_outbound_events"
  1295. ).fetchone()[0] == 1
  1296. def test_expired_remote_lease_rebinds_pending_event_for_exact_ack_replay(tmp_path):
  1297. private = Ed25519PrivateKey.generate()
  1298. public = private.public_key().public_bytes(
  1299. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  1300. )
  1301. envelope = SignedTaskEnvelope.verify_mapping(
  1302. _signed_task(private),
  1303. authority_keys={"control-authority-2026-01": public},
  1304. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  1305. )
  1306. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1307. queue = SqliteEdgeQueue(
  1308. tmp_path / "edge.db",
  1309. clock=lambda: now[0],
  1310. remote_lease_encoder=lambda token: "cipher:" + token[::-1],
  1311. remote_lease_decoder=lambda value: value.removeprefix("cipher:")[::-1],
  1312. )
  1313. queue.accept_signed_task(
  1314. envelope,
  1315. remote_lease_token="remote-a",
  1316. remote_lease_expires_at="2026-08-09T08:00:05Z",
  1317. remote_attempt=1,
  1318. )
  1319. task = queue.claim("worker")
  1320. _, pending = queue.complete_with_event(
  1321. task.task_id,
  1322. lease_token=task.lease_token,
  1323. remote_lease_token="remote-a",
  1324. event=_event(),
  1325. )
  1326. first_send = queue.claim_event("sender-a")
  1327. queue.fail_event(
  1328. first_send.event_id,
  1329. lease_token=first_send.lease_token,
  1330. error_code="ack_lost",
  1331. )
  1332. now[0] = datetime(2026, 8, 9, 8, 0, 6, tzinfo=UTC)
  1333. queue.accept_signed_task(
  1334. envelope,
  1335. remote_lease_token="remote-b",
  1336. remote_lease_expires_at="2026-08-09T08:05:00Z",
  1337. remote_attempt=1,
  1338. )
  1339. rebound = queue.get_event(pending.event_id)
  1340. assert rebound.remote_lease_token_digest == hashlib.sha256(b"remote-b").hexdigest()
  1341. second_send = queue.claim_event("sender-b")
  1342. acknowledged = queue.acknowledge_event(
  1343. second_send.event_id,
  1344. lease_token=second_send.lease_token,
  1345. acknowledgement={
  1346. "event_id": second_send.event_id,
  1347. "event_digest": second_send.digest,
  1348. "remote_lease_digest": hashlib.sha256(b"remote-b").hexdigest(),
  1349. "received_at": "2026-08-09T08:00:06Z",
  1350. "status": "accepted",
  1351. },
  1352. )
  1353. assert acknowledged.status == "acknowledged"
  1354. def test_remote_reacquire_revives_only_execution_lease_expiry(tmp_path):
  1355. private = Ed25519PrivateKey.generate()
  1356. public = private.public_key().public_bytes(
  1357. serialization.Encoding.Raw, serialization.PublicFormat.Raw
  1358. )
  1359. envelope = SignedTaskEnvelope.verify_mapping(
  1360. _signed_task(private),
  1361. authority_keys={"control-authority-2026-01": public},
  1362. now=datetime(2026, 8, 9, 8, 0, tzinfo=UTC),
  1363. )
  1364. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1365. def new_queue(path):
  1366. return SqliteEdgeQueue(
  1367. path,
  1368. clock=lambda: now[0],
  1369. remote_lease_encoder=lambda token: "cipher:" + token[::-1],
  1370. remote_lease_decoder=lambda value: value.removeprefix("cipher:")[::-1],
  1371. )
  1372. recoverable = new_queue(tmp_path / "recoverable.db")
  1373. recoverable.accept_signed_task(
  1374. envelope,
  1375. remote_lease_token="remote-a",
  1376. remote_lease_expires_at="2026-08-09T08:00:05Z",
  1377. remote_attempt=1,
  1378. )
  1379. claimed = recoverable.claim("worker")
  1380. recoverable.fail(
  1381. claimed.task_id,
  1382. lease_token=claimed.lease_token,
  1383. error_code="execution_lease_expired",
  1384. retryable=False,
  1385. )
  1386. now[0] = datetime(2026, 8, 9, 8, 0, 6, tzinfo=UTC)
  1387. revived = recoverable.accept_signed_task(
  1388. envelope,
  1389. remote_lease_token="remote-b",
  1390. remote_lease_expires_at="2026-08-09T08:05:00Z",
  1391. remote_attempt=1,
  1392. )
  1393. assert revived.status == "pending"
  1394. assert revived.error_code is None
  1395. business = new_queue(tmp_path / "business.db")
  1396. business.accept_signed_task(
  1397. envelope,
  1398. remote_lease_token="remote-a",
  1399. remote_lease_expires_at="2026-08-09T08:00:07Z",
  1400. remote_attempt=1,
  1401. )
  1402. business_claim = business.claim("worker")
  1403. business.fail(
  1404. business_claim.task_id,
  1405. lease_token=business_claim.lease_token,
  1406. error_code="runner_execution_failed",
  1407. retryable=False,
  1408. )
  1409. now[0] = datetime(2026, 8, 9, 8, 0, 8, tzinfo=UTC)
  1410. with pytest.raises(EdgeQueueConflictError, match="terminal task"):
  1411. business.accept_signed_task(
  1412. envelope,
  1413. remote_lease_token="remote-b",
  1414. remote_lease_expires_at="2026-08-09T08:05:00Z",
  1415. remote_attempt=1,
  1416. )
  1417. def test_retry_jitter_is_capped_after_jitter_for_task_and_event(tmp_path):
  1418. now = [datetime(2026, 8, 9, 8, 0, tzinfo=UTC)]
  1419. queue = SqliteEdgeQueue(
  1420. tmp_path / "edge.db",
  1421. clock=lambda: now[0],
  1422. retry_base_seconds=10,
  1423. retry_max_seconds=10,
  1424. random_source=lambda: 1.0,
  1425. )
  1426. queue.enqueue(APPROVED_TASK)
  1427. task = queue.claim("worker")
  1428. task_retry = queue.fail(
  1429. task.task_id, lease_token=task.lease_token, error_code="temporary"
  1430. )
  1431. assert task_retry.available_at == "2026-08-09T08:00:10Z"
  1432. queue.persist_event(_event())
  1433. event = queue.claim_event("sender")
  1434. event_retry = queue.fail_event(
  1435. event.event_id, lease_token=event.lease_token, error_code="temporary"
  1436. )
  1437. assert event_retry.available_at == "2026-08-09T08:00:10Z"
  1438. def test_release_state_is_durable_and_transitions_with_atomic_cas(tmp_path):
  1439. db_path = tmp_path / "edge.db"
  1440. first = SqliteEdgeQueue(db_path)
  1441. state = first.put_release_state(
  1442. release_id="release-1", manifest_digest="d" * 64,
  1443. version="3.1.0", rollback_version="3.0.0",
  1444. )
  1445. assert state.status == "offered"
  1446. accepted = first.compare_and_set_release_state(
  1447. "release-1", expected_status="offered", target_status="accepted",
  1448. previous_version="3.0.0",
  1449. )
  1450. assert accepted.status == "accepted"
  1451. restarted = SqliteEdgeQueue(db_path)
  1452. assert restarted.get_release_state("release-1") == accepted
  1453. with pytest.raises(EdgeQueueConflictError, match="compare-and-set"):
  1454. restarted.compare_and_set_release_state(
  1455. "release-1", expected_status="offered", target_status="accepted"
  1456. )
  1457. candidate = restarted.compare_and_set_release_state(
  1458. "release-1", expected_status="accepted", target_status="candidate"
  1459. )
  1460. installed = restarted.compare_and_set_release_state(
  1461. "release-1", expected_status="candidate", target_status="installed",
  1462. current_version="3.1.0",
  1463. )
  1464. assert candidate.status == "candidate"
  1465. assert installed.current_version == "3.1.0"
  1466. assert restarted.list_release_states() == (installed,)