agent.py 52 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273
  1. """Durable pull-only edge agent, closed runner adapter, and release verifier."""
  2. from __future__ import annotations
  3. import base64
  4. import hashlib
  5. import os
  6. import re
  7. import sqlite3
  8. import stat
  9. from collections.abc import Callable, Mapping
  10. from datetime import UTC, datetime, timedelta
  11. from pathlib import Path, PurePosixPath
  12. from types import MappingProxyType
  13. from cryptography.exceptions import InvalidSignature
  14. from cryptography.hazmat.primitives import hashes
  15. from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
  16. from cryptography.hazmat.primitives.ciphers.aead import AESGCM
  17. from cryptography.hazmat.primitives.kdf.hkdf import HKDF
  18. from app.core.edge_gateway.contracts import (
  19. EdgeContractError,
  20. EdgeEventContract,
  21. EdgeTaskContract,
  22. SignedTaskEnvelope,
  23. canonical_json_bytes,
  24. canonical_sha256,
  25. canonical_timestamp,
  26. stable_event_id,
  27. )
  28. from app.core.edge_gateway.policy import EdgeEgressPolicy, EdgePolicyError
  29. from app.core.edge_gateway.queue import (
  30. EdgeQueueConflictError,
  31. EdgeQueueLeaseError,
  32. SqliteEdgeQueue,
  33. )
  34. from app.edge_gateway.bootstrap import EdgeBootstrapConfig
  35. from app.edge_gateway.transport import EdgeAuthenticationStopped, EdgeTransportError
  36. _SHA256 = re.compile(r"^[0-9a-f]{64}$")
  37. _VERSION = re.compile(r"^(0|[1-9][0-9]{0,9})\.(0|[1-9][0-9]{0,9})\.(0|[1-9][0-9]{0,9})$")
  38. _RELEASE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,254}$")
  39. EDGE_NODE_BY_OPERATION = MappingProxyType(
  40. {
  41. "collect": "edge.collect",
  42. "profile": "edge.profile",
  43. "quality": "quality.check",
  44. "lineage": "edge.lineage",
  45. "controlled_query": "edge.controlled_query",
  46. }
  47. )
  48. EVENT_CLASSIFICATION_BY_OPERATION = MappingProxyType(
  49. {
  50. "collect": "desensitized_metadata",
  51. "profile": "statistics",
  52. "quality": "evidence",
  53. "lineage": "lineage",
  54. "controlled_query": "statistics",
  55. }
  56. )
  57. DEFAULT_GOVERNED_PURPOSES = frozenset(
  58. {
  59. "inventory",
  60. "governed-inventory",
  61. "metadata-inventory",
  62. "profile-statistics",
  63. "quality-evaluation",
  64. "lineage-capture",
  65. "controlled-query",
  66. }
  67. )
  68. _RUNNER_REQUEST_FIELDS = frozenset(
  69. {
  70. "task_id",
  71. "operation",
  72. "purpose",
  73. "classification",
  74. "environment",
  75. "network_zone",
  76. "idempotency_key",
  77. "deadline_at",
  78. }
  79. )
  80. _RUNNER_RESULT_FIELDS = frozenset(
  81. {"classification", "payload", "local_artifact_digest", "local_artifact_ref"}
  82. )
  83. class EdgeAgentError(RuntimeError):
  84. """Bounded edge runtime failure."""
  85. class EdgeTaskCancelled(EdgeAgentError):
  86. """Cooperative cancellation was observed before a terminal event."""
  87. class EdgeReleaseError(EdgeAgentError):
  88. """A candidate release failed a closed verification or activation gate."""
  89. class LocalArtifactStore:
  90. """Validate and delete local artifacts without exposing their paths outbound."""
  91. def __init__(self, root: str | Path, *, clock: Callable[[], object] | None = None):
  92. candidate = Path(root)
  93. if not candidate.is_absolute() or not candidate.is_dir() or candidate.is_symlink():
  94. raise ValueError("artifact root must be an existing absolute directory")
  95. self.root = candidate.resolve()
  96. self.clock = clock or (lambda: datetime.now(UTC))
  97. def inspect(self, reference: object, expected_digest: object) -> dict[str, str]:
  98. if (
  99. not isinstance(reference, str)
  100. or not reference
  101. or "\x00" in reference
  102. or len(reference.encode("utf-8")) > 2_048
  103. or not isinstance(expected_digest, str)
  104. or not _SHA256.fullmatch(expected_digest)
  105. ):
  106. raise EdgePolicyError("local artifact reference is invalid")
  107. raw = Path(reference)
  108. path = raw if raw.is_absolute() else self.root / raw
  109. try:
  110. resolved = path.resolve(strict=True)
  111. resolved.relative_to(self.root)
  112. current = self.root
  113. for part in resolved.relative_to(self.root).parts:
  114. current = current / part
  115. if current.is_symlink():
  116. raise EdgePolicyError("local artifact symlink is forbidden")
  117. metadata = resolved.lstat()
  118. except (OSError, ValueError) as exc:
  119. raise EdgePolicyError("local artifact escapes its configured root") from exc
  120. if not stat.S_ISREG(metadata.st_mode):
  121. raise EdgePolicyError("local artifact must be a regular file")
  122. digest = hashlib.sha256()
  123. try:
  124. with resolved.open("rb") as stream:
  125. for chunk in iter(lambda: stream.read(1_048_576), b""):
  126. digest.update(chunk)
  127. except OSError as exc:
  128. raise EdgePolicyError("local artifact cannot be read") from exc
  129. if digest.hexdigest() != expected_digest:
  130. raise EdgePolicyError("local artifact digest does not match content")
  131. canonical_ref = str(resolved)
  132. return {
  133. "artifact_digest": expected_digest,
  134. "artifact_ref": canonical_ref,
  135. "artifact_ref_hash": canonical_sha256(canonical_ref),
  136. }
  137. def delete(self, reference: object, expected_digest: object) -> dict[str, str]:
  138. inspected = self.inspect(reference, expected_digest)
  139. try:
  140. Path(inspected["artifact_ref"]).unlink()
  141. except OSError as exc:
  142. raise EdgePolicyError("local artifact deletion failed") from exc
  143. return {
  144. "artifact_digest": inspected["artifact_digest"],
  145. "artifact_ref_hash": inspected["artifact_ref_hash"],
  146. "deleted_at": canonical_timestamp(
  147. _now(self.clock).isoformat().replace("+00:00", "Z")
  148. ),
  149. "status": "deleted",
  150. }
  151. def receipt_for_missing(
  152. self, reference: object, expected_digest: object, expected_ref_hash: object
  153. ) -> dict[str, str]:
  154. if (
  155. not isinstance(reference, str)
  156. or not isinstance(expected_digest, str)
  157. or not _SHA256.fullmatch(expected_digest)
  158. or not isinstance(expected_ref_hash, str)
  159. or not _SHA256.fullmatch(expected_ref_hash)
  160. ):
  161. raise EdgePolicyError("missing artifact recovery metadata is invalid")
  162. path = Path(reference)
  163. try:
  164. path.parent.resolve(strict=True).relative_to(self.root)
  165. canonical_ref = str(path.resolve(strict=False))
  166. except (OSError, ValueError) as exc:
  167. raise EdgePolicyError("missing artifact escapes its configured root") from exc
  168. if path.exists() or path.is_symlink() or canonical_sha256(canonical_ref) != expected_ref_hash:
  169. raise EdgePolicyError("missing artifact recovery proof is invalid")
  170. return {
  171. "artifact_digest": expected_digest,
  172. "artifact_ref_hash": expected_ref_hash,
  173. "deleted_at": canonical_timestamp(
  174. _now(self.clock).isoformat().replace("+00:00", "Z")
  175. ),
  176. "status": "deleted",
  177. }
  178. def _now(clock: Callable[[], object]) -> datetime:
  179. value = clock()
  180. if not isinstance(value, datetime) or value.tzinfo is None or value.utcoffset() is None:
  181. raise ValueError("edge clock must return a timezone-aware datetime")
  182. return value.astimezone(UTC)
  183. def _timestamp(clock: Callable[[], object]) -> str:
  184. return canonical_timestamp(_now(clock).isoformat().replace("+00:00", "Z"))
  185. def _remote_lease_codec(config: EdgeBootstrapConfig):
  186. """Protect recoverable control leases without persisting the credential."""
  187. key = HKDF(
  188. algorithm=hashes.SHA256(),
  189. length=32,
  190. salt=bytes.fromhex(config.certificate_sha256),
  191. info=(
  192. f"dataops-edge-remote-lease-v1:{config.gateway_id}:"
  193. f"{config.environment}:{config.network_zone}"
  194. ).encode(),
  195. ).derive(config.credential.encode())
  196. cipher = AESGCM(key)
  197. aad = f"{config.gateway_id}:{config.policy_digest}".encode()
  198. def encode(token: str) -> str:
  199. nonce = os.urandom(12)
  200. protected = cipher.encrypt(nonce, token.encode("utf-8"), aad)
  201. return base64.urlsafe_b64encode(nonce + protected).decode("ascii")
  202. def decode(value: str) -> str:
  203. raw = base64.b64decode(value.encode("ascii"), altchars=b"-_", validate=True)
  204. if len(raw) < 29:
  205. raise ValueError("protected remote lease is invalid")
  206. return cipher.decrypt(raw[:12], raw[12:], aad).decode("utf-8")
  207. return encode, decode
  208. class EdgeRunnerAdapter:
  209. """Map governed operations to exact local runner handlers.
  210. No control-plane supplied parameters, module names, commands, SQL, paths, or
  211. network destinations are accepted by this adapter. Local handlers resolve
  212. their governed configuration using the immutable task identity.
  213. """
  214. def __init__(
  215. self,
  216. executor: Callable[[str, Mapping[str, object], Callable[[], bool]], object],
  217. *,
  218. policy: EdgeEgressPolicy | None = None,
  219. allowed_purposes: set[str] | frozenset[str] = DEFAULT_GOVERNED_PURPOSES,
  220. ) -> None:
  221. if not callable(executor):
  222. raise ValueError("an explicit edge runner executor is required")
  223. self._executor = executor
  224. self.policy = policy or EdgeEgressPolicy(allowed_control_hosts=set())
  225. if not isinstance(allowed_purposes, (set, frozenset)) or not allowed_purposes:
  226. raise ValueError("explicit governed purposes are required")
  227. if any(
  228. not isinstance(item, str)
  229. or not item
  230. or len(item.encode("utf-8")) > 255
  231. for item in allowed_purposes
  232. ):
  233. raise ValueError("governed purpose allowlist is invalid")
  234. self.allowed_purposes = frozenset(allowed_purposes)
  235. def execute(
  236. self, task: EdgeTaskContract, cancel_requested: Callable[[], bool]
  237. ) -> dict[str, object]:
  238. task = EdgeTaskContract.from_mapping(task.to_mapping())
  239. node = EDGE_NODE_BY_OPERATION.get(task.operation)
  240. if node is None:
  241. raise EdgePolicyError("edge task operation is not runner-approved")
  242. if task.purpose not in self.allowed_purposes:
  243. raise EdgePolicyError("edge task purpose is not runner-approved")
  244. if cancel_requested():
  245. raise EdgeTaskCancelled("edge task was cancelled")
  246. request = {
  247. "task_id": task.task_id,
  248. "operation": task.operation,
  249. "purpose": task.purpose,
  250. "classification": task.classification,
  251. "environment": task.environment,
  252. "network_zone": task.network_zone,
  253. "idempotency_key": task.idempotency_key,
  254. "deadline_at": task.deadline_at,
  255. }
  256. if set(request) != _RUNNER_REQUEST_FIELDS: # pragma: no cover - invariant
  257. raise EdgePolicyError("edge runner request schema is invalid")
  258. try:
  259. raw = self._executor(node, MappingProxyType(request), cancel_requested)
  260. except EdgeTaskCancelled:
  261. raise
  262. except (MemoryError, RecursionError) as exc:
  263. raise EdgePolicyError("edge runner result exceeds safe limits") from exc
  264. if cancel_requested():
  265. raise EdgeTaskCancelled("edge task was cancelled")
  266. if not isinstance(raw, Mapping):
  267. raise EdgePolicyError("edge runner result must be a mapping")
  268. unknown = set(raw) - _RUNNER_RESULT_FIELDS
  269. missing = {"classification", "payload"} - set(raw)
  270. if unknown or missing:
  271. raise EdgePolicyError("edge runner result schema is invalid")
  272. classification = raw["classification"]
  273. expected = EVENT_CLASSIFICATION_BY_OPERATION[task.operation]
  274. if classification != expected:
  275. raise EdgePolicyError("edge runner result classification is invalid")
  276. payload = raw["payload"]
  277. if not isinstance(payload, Mapping):
  278. raise EdgePolicyError("edge runner result payload must be a mapping")
  279. try:
  280. self.policy.validate_approved_payload(expected, payload)
  281. except (EdgePolicyError, EdgeContractError, MemoryError, RecursionError) as exc:
  282. raise EdgePolicyError("edge runner result is not approved for egress") from exc
  283. artifact_digest = raw.get("local_artifact_digest")
  284. artifact_ref = raw.get("local_artifact_ref")
  285. if (artifact_digest is None) != (artifact_ref is None):
  286. raise EdgePolicyError("local artifact digest and reference must be paired")
  287. if artifact_digest is not None:
  288. if not isinstance(artifact_digest, str) or not _SHA256.fullmatch(artifact_digest):
  289. raise EdgePolicyError("local artifact digest is invalid")
  290. if (
  291. not isinstance(artifact_ref, str)
  292. or not artifact_ref
  293. or "\x00" in artifact_ref
  294. or len(artifact_ref.encode("utf-8")) > 2_048
  295. ):
  296. raise EdgePolicyError("local artifact reference is invalid")
  297. return {
  298. "classification": expected,
  299. "payload": dict(payload),
  300. "local_artifact_digest": artifact_digest,
  301. "local_artifact_ref": artifact_ref,
  302. }
  303. class SignedReleaseManager:
  304. """Verify Ed25519 manifests and activate candidates through injected primitives."""
  305. _MANIFEST_FIELDS = frozenset(
  306. {
  307. "release_id",
  308. "version",
  309. "rollback_version",
  310. "artifact_digest",
  311. "artifact_name",
  312. "deadline_at",
  313. "signature_algorithm",
  314. "key_id",
  315. "manifest_digest",
  316. "signature",
  317. "status",
  318. }
  319. )
  320. def __init__(
  321. self,
  322. *,
  323. current_version: str,
  324. trusted_keys: Mapping[str, str],
  325. artifact_loader: Callable[[Mapping[str, object]], bytes],
  326. activator: Callable[[str, bytes], bool],
  327. rollback: Callable[[str], bool],
  328. clock: Callable[[], object] | None = None,
  329. ) -> None:
  330. if not _VERSION.fullmatch(current_version):
  331. raise ValueError("current release version is invalid")
  332. if any(int(part) > 2_147_483_647 for part in current_version.split(".")):
  333. raise ValueError("current release version is invalid")
  334. if not trusted_keys:
  335. raise ValueError("trusted release keys are required")
  336. self.current_version = current_version
  337. self.trusted_keys = dict(trusted_keys)
  338. self.artifact_loader = artifact_loader
  339. self.activator = activator
  340. self.rollback = rollback
  341. self.clock = clock or (lambda: datetime.now(UTC))
  342. self._seen: dict[str, str] = {}
  343. @staticmethod
  344. def _version_tuple(value: str) -> tuple[int, ...]:
  345. return tuple(int(item) for item in value.split("."))
  346. def verify(
  347. self,
  348. value: Mapping[str, object],
  349. *,
  350. allow_recovery: bool = False,
  351. ) -> dict[str, object]:
  352. if not isinstance(value, Mapping) or set(value) != self._MANIFEST_FIELDS:
  353. raise EdgeReleaseError("release manifest schema is invalid")
  354. manifest = dict(value)
  355. release_id = manifest["release_id"]
  356. version = manifest["version"]
  357. rollback_version = manifest["rollback_version"]
  358. if not isinstance(release_id, str) or not _RELEASE_ID.fullmatch(release_id):
  359. raise EdgeReleaseError("release identifier is invalid")
  360. status = manifest.get("status")
  361. if status not in {"offered", "accepted", "failed", "installed", "rolled_back"}:
  362. raise EdgeReleaseError("release state is invalid")
  363. if (
  364. not isinstance(version, str)
  365. or not isinstance(rollback_version, str)
  366. or not _VERSION.fullmatch(version)
  367. or not _VERSION.fullmatch(rollback_version)
  368. or any(int(part) > 2_147_483_647 for part in version.split("."))
  369. or any(int(part) > 2_147_483_647 for part in rollback_version.split("."))
  370. ):
  371. raise EdgeReleaseError("release downgrade or rollback binding was rejected")
  372. if not allow_recovery and status in {"offered", "accepted"} and (
  373. rollback_version != self.current_version
  374. or self._version_tuple(version) <= self._version_tuple(self.current_version)
  375. ):
  376. raise EdgeReleaseError("release downgrade or rollback binding was rejected")
  377. artifact_name = manifest["artifact_name"]
  378. if (
  379. not isinstance(artifact_name, str)
  380. or not artifact_name
  381. or len(artifact_name.encode("utf-8")) > 255
  382. or PurePosixPath(artifact_name).name != artifact_name
  383. or artifact_name in {".", ".."}
  384. or any(character in artifact_name for character in ("/", "\\", "\x00"))
  385. ):
  386. raise EdgeReleaseError("release artifact path was rejected")
  387. if manifest["signature_algorithm"] != "Ed25519":
  388. raise EdgeReleaseError("release signature algorithm is not approved")
  389. key_id = manifest["key_id"]
  390. key_hex = self.trusted_keys.get(key_id) if isinstance(key_id, str) else None
  391. if key_hex is None:
  392. raise EdgeReleaseError("release signing key is not trusted")
  393. for field in ("artifact_digest", "manifest_digest"):
  394. if not isinstance(manifest[field], str) or not _SHA256.fullmatch(manifest[field]):
  395. raise EdgeReleaseError("release digest is invalid")
  396. try:
  397. deadline = datetime.fromisoformat(str(manifest["deadline_at"]).replace("Z", "+00:00"))
  398. except ValueError as exc:
  399. raise EdgeReleaseError("release deadline is invalid") from exc
  400. if deadline.tzinfo is None or deadline.astimezone(UTC) <= _now(self.clock):
  401. raise EdgeReleaseError("release deadline has expired")
  402. try:
  403. canonical_deadline = canonical_timestamp(manifest["deadline_at"], "deadline_at")
  404. except EdgeContractError as exc:
  405. raise EdgeReleaseError("release deadline is invalid") from exc
  406. if canonical_deadline != manifest["deadline_at"]:
  407. raise EdgeReleaseError("release deadline must use canonical UTC form")
  408. unsigned = {
  409. key: manifest[key]
  410. for key in sorted(self._MANIFEST_FIELDS - {"manifest_digest", "signature"})
  411. }
  412. expected_digest = canonical_sha256(unsigned)
  413. if expected_digest != manifest["manifest_digest"]:
  414. raise EdgeReleaseError("release manifest digest does not match")
  415. content_digest = canonical_sha256(
  416. {
  417. key: manifest[key]
  418. for key in (
  419. "release_id",
  420. "version",
  421. "rollback_version",
  422. "artifact_digest",
  423. "artifact_name",
  424. "deadline_at",
  425. )
  426. }
  427. )
  428. prior = self._seen.get(release_id)
  429. if prior is not None and prior != content_digest:
  430. raise EdgeReleaseError("release identifier replay changed its digest")
  431. try:
  432. signature = bytes.fromhex(str(manifest["signature"]))
  433. Ed25519PublicKey.from_public_bytes(bytes.fromhex(key_hex)).verify(
  434. signature, canonical_json_bytes(unsigned)
  435. )
  436. except (ValueError, InvalidSignature) as exc:
  437. raise EdgeReleaseError("release signature verification failed") from exc
  438. self._seen[release_id] = content_digest
  439. return manifest
  440. def install(self, value: Mapping[str, object]) -> dict[str, object]:
  441. manifest = self.verify(value)
  442. release_id = str(manifest["release_id"])
  443. previous = self.current_version
  444. try:
  445. artifact = self.artifact_loader(MappingProxyType(manifest))
  446. if not isinstance(artifact, bytes) or hashlib.sha256(artifact).hexdigest() != manifest["artifact_digest"]:
  447. raise EdgeReleaseError("release artifact digest does not match")
  448. if not self.activator(str(manifest["version"]), artifact):
  449. raise EdgeReleaseError("release candidate health confirmation failed")
  450. except Exception:
  451. rolled_back = bool(self.rollback(previous))
  452. return {
  453. "release_id": release_id,
  454. "outcome": "rollback" if rolled_back else "failed",
  455. "summary": {
  456. "release_status": "rolled_back" if rolled_back else "failed",
  457. "version": previous,
  458. },
  459. }
  460. self.current_version = str(manifest["version"])
  461. return {
  462. "release_id": release_id,
  463. "outcome": "installed",
  464. "summary": {"release_status": "installed", "version": self.current_version},
  465. }
  466. class EdgeAgent:
  467. """One durable worker for pull, execute, persist, and acknowledge cycles."""
  468. def __init__(
  469. self,
  470. config: EdgeBootstrapConfig,
  471. transport,
  472. runner: EdgeRunnerAdapter,
  473. *,
  474. worker_id: str = "edge-worker-1",
  475. clock: Callable[[], object] | None = None,
  476. random: Callable[[], float] | None = None,
  477. release_manager: SignedReleaseManager | None = None,
  478. artifact_store: LocalArtifactStore | None = None,
  479. ) -> None:
  480. if not isinstance(config, EdgeBootstrapConfig):
  481. raise ValueError("validated edge bootstrap configuration is required")
  482. if not isinstance(runner, EdgeRunnerAdapter):
  483. raise ValueError("validated edge runner adapter is required")
  484. if not isinstance(worker_id, str) or not worker_id or len(worker_id.encode()) > 255:
  485. raise ValueError("worker_id is invalid")
  486. self.config = config
  487. self.transport = transport
  488. self.runner = runner
  489. self.worker_id = worker_id
  490. self.clock = clock or (lambda: datetime.now(UTC))
  491. self.random = random or __import__("random").random
  492. self.policy = EdgeEgressPolicy(
  493. allowed_control_hosts=set(config.allowed_control_hosts),
  494. allowed_proxy_hosts=set(config.allowed_proxy_hosts),
  495. )
  496. lease_encoder, lease_decoder = _remote_lease_codec(config)
  497. self.queue = SqliteEdgeQueue(
  498. config.queue_path,
  499. clock=self.clock,
  500. egress_policy=self.policy,
  501. random_source=self.random,
  502. remote_lease_encoder=lease_encoder,
  503. remote_lease_decoder=lease_decoder,
  504. )
  505. self.artifact_store = artifact_store or LocalArtifactStore(
  506. config.artifact_root, clock=self.clock
  507. )
  508. self.release_manager = release_manager
  509. self._cancelled: set[str] = set()
  510. self._stopped = False
  511. self._needs_remote_lease_task: str | None = None
  512. self._trusted_task_keys = {
  513. key_id: bytes.fromhex(public_key)
  514. for key_id, public_key in config.trusted_task_keys.items()
  515. }
  516. self.gateway_id_digest = hashlib.sha256(config.gateway_id.encode()).hexdigest()
  517. def retry_delay(self, attempt: int) -> float:
  518. if isinstance(attempt, bool) or not isinstance(attempt, int) or attempt < 1:
  519. raise ValueError("retry attempt is invalid")
  520. base = min(300.0, float(2 ** min(attempt - 1, 8)))
  521. jitter = self.random()
  522. if not isinstance(jitter, (int, float)) or not 0 <= jitter <= 1:
  523. raise ValueError("random source returned an invalid value")
  524. return min(300.0, base * (0.75 + 0.5 * float(jitter)))
  525. def _validate_binding(self, contract: EdgeTaskContract) -> None:
  526. expected = (
  527. self.config.gateway_id,
  528. self.config.environment,
  529. self.config.network_zone,
  530. self.config.policy_digest,
  531. )
  532. actual = (
  533. contract.gateway_id,
  534. contract.environment,
  535. contract.network_zone,
  536. contract.policy_digest,
  537. )
  538. if actual != expected:
  539. raise EdgePolicyError("pulled task binding does not match this edge gateway")
  540. def accept_pulled_task(
  541. self,
  542. envelope_value: Mapping[str, object],
  543. lease_token: str,
  544. lease_expires_at: str,
  545. ):
  546. if not isinstance(lease_token, str) or not lease_token.startswith("dopl_") or len(lease_token) > 128:
  547. raise EdgePolicyError("control task lease is invalid")
  548. try:
  549. envelope = SignedTaskEnvelope.verify_mapping(
  550. envelope_value,
  551. authority_keys=self._trusted_task_keys,
  552. now=_now(self.clock),
  553. allowed_future_skew_seconds=(
  554. self.config.task_authority_clock_skew_seconds
  555. ),
  556. )
  557. except EdgeContractError as exc:
  558. raise EdgePolicyError("signed task authority was rejected") from exc
  559. self._validate_binding(envelope.task)
  560. return self.queue.accept_signed_task(
  561. envelope,
  562. remote_lease_token=lease_token,
  563. remote_lease_expires_at=lease_expires_at,
  564. remote_attempt=envelope.task.attempt,
  565. )
  566. def _cancel_requested(self, task_id: str) -> bool:
  567. if task_id in self._cancelled:
  568. return True
  569. probe = getattr(self.transport, "cancel_requested", None)
  570. if callable(probe):
  571. try:
  572. if bool(probe(task_id)):
  573. self._cancelled.add(task_id)
  574. return True
  575. except EdgeAuthenticationStopped:
  576. self._stopped = True
  577. raise
  578. except EdgeTransportError:
  579. return task_id in self._cancelled
  580. return False
  581. def _event(self, task: EdgeTaskContract, result: Mapping[str, object]) -> EdgeEventContract:
  582. identity = {
  583. "task_id": task.task_id,
  584. "gateway_id": task.gateway_id,
  585. "environment": task.environment,
  586. "network_zone": task.network_zone,
  587. "purpose": task.purpose,
  588. "classification": result["classification"],
  589. "contract_version": task.contract_version,
  590. "occurred_at": _timestamp(self.clock),
  591. "attempt": task.attempt,
  592. "idempotency_key": task.idempotency_key,
  593. "policy_digest": task.policy_digest,
  594. "payload": result["payload"],
  595. }
  596. return EdgeEventContract(event_id=stable_event_id(identity), **identity) # type: ignore[arg-type]
  597. def _artifact_metadata(
  598. self,
  599. task: EdgeTaskContract,
  600. result: Mapping[str, object],
  601. ) -> tuple[dict[str, object], ...]:
  602. digest = result.get("local_artifact_digest")
  603. reference = result.get("local_artifact_ref")
  604. if digest is None or reference is None:
  605. return ()
  606. inspected = self.artifact_store.inspect(reference, digest)
  607. retention_class = (
  608. task.classification
  609. if task.classification in {"raw", "recent_detail", "evidence"}
  610. else "raw"
  611. if task.classification == "restricted"
  612. else "metadata"
  613. )
  614. retention_until = _now(self.clock) + timedelta(
  615. days=self.policy.retention_days(retention_class)
  616. )
  617. return (
  618. {
  619. "artifact_digest": inspected["artifact_digest"],
  620. "artifact_ref": inspected["artifact_ref"],
  621. "artifact_ref_hash": inspected["artifact_ref_hash"],
  622. "classification": task.classification,
  623. "retention_until": canonical_timestamp(
  624. retention_until.isoformat().replace("+00:00", "Z")
  625. ),
  626. },
  627. )
  628. def _drain_artifact_cleanup(self) -> str | None:
  629. artifact = self.queue.claim_artifact_cleanup(self.worker_id)
  630. if artifact is None:
  631. return None
  632. if artifact.artifact_ref is None:
  633. self.queue.fail_artifact_cleanup(
  634. artifact.artifact_digest,
  635. lease_token=str(artifact.lease_token),
  636. error_code="artifact_reference_missing",
  637. )
  638. return None
  639. try:
  640. receipt = self.artifact_store.delete(
  641. artifact.artifact_ref, artifact.artifact_digest
  642. )
  643. self.queue.acknowledge_artifact_cleanup(
  644. artifact.artifact_digest,
  645. lease_token=str(artifact.lease_token),
  646. receipt=receipt,
  647. )
  648. return artifact.artifact_digest
  649. except EdgePolicyError:
  650. try:
  651. receipt = self.artifact_store.receipt_for_missing(
  652. artifact.artifact_ref,
  653. artifact.artifact_digest,
  654. artifact.artifact_ref_hash,
  655. )
  656. self.queue.acknowledge_artifact_cleanup(
  657. artifact.artifact_digest,
  658. lease_token=str(artifact.lease_token),
  659. receipt=receipt,
  660. )
  661. return artifact.artifact_digest
  662. except EdgePolicyError:
  663. pass
  664. self.queue.fail_artifact_cleanup(
  665. artifact.artifact_digest,
  666. lease_token=str(artifact.lease_token),
  667. error_code="artifact_cleanup_failed",
  668. )
  669. return None
  670. def _drain_event(self) -> str | None:
  671. event = self.queue.claim_event(self.worker_id)
  672. if event is None:
  673. return None
  674. task_id = str(event.event["task_id"])
  675. task = self.queue.get_task(task_id)
  676. if task is None:
  677. self.queue.fail_event(
  678. event.event_id,
  679. lease_token=str(event.lease_token),
  680. error_code="control_lease_unavailable",
  681. )
  682. return None
  683. expiry = task.remote_lease_expires_at
  684. expired = (
  685. expiry is None
  686. or datetime.fromisoformat(expiry.replace("Z", "+00:00")) <= _now(self.clock)
  687. )
  688. if expired:
  689. self._needs_remote_lease_task = task_id
  690. self.queue.fail_event(
  691. event.event_id,
  692. lease_token=str(event.lease_token),
  693. error_code="remote_lease_expired",
  694. )
  695. return None
  696. try:
  697. control_lease = self.queue.recover_remote_lease(task_id)
  698. except EdgeQueueLeaseError:
  699. self._needs_remote_lease_task = task_id
  700. self.queue.fail_event(
  701. event.event_id,
  702. lease_token=str(event.lease_token),
  703. error_code="control_lease_unavailable",
  704. )
  705. return None
  706. try:
  707. ack = self.transport.send_event(event.event, control_lease)
  708. self.queue.acknowledge_event(
  709. event.event_id,
  710. lease_token=str(event.lease_token),
  711. acknowledgement=ack,
  712. )
  713. return event.event_id
  714. except EdgeAuthenticationStopped:
  715. if not expired:
  716. self._stopped = True
  717. else:
  718. self._needs_remote_lease_task = task_id
  719. self.queue.fail_event(
  720. event.event_id,
  721. lease_token=str(event.lease_token),
  722. error_code=(
  723. "credential_stopped" if not expired else "remote_lease_expired"
  724. ),
  725. )
  726. if not expired:
  727. raise
  728. return None
  729. except EdgeTransportError:
  730. self.queue.fail_event(
  731. event.event_id,
  732. lease_token=str(event.lease_token),
  733. error_code="control_unavailable",
  734. )
  735. return None
  736. def _has_unresolved_events(self) -> bool:
  737. connection = sqlite3.connect(f"file:{self.queue.db_path}?mode=ro", uri=True)
  738. try:
  739. count = connection.execute(
  740. "SELECT COUNT(*) FROM edge_outbound_events WHERE status IN ('pending','sending')"
  741. ).fetchone()[0]
  742. finally:
  743. connection.close()
  744. return bool(count)
  745. def _expired_pending_remote_task(self) -> str | None:
  746. connection = sqlite3.connect(f"file:{self.queue.db_path}?mode=ro", uri=True)
  747. try:
  748. row = connection.execute(
  749. """SELECT task_id FROM edge_tasks
  750. WHERE status='pending' AND (
  751. remote_lease_expires_at IS NULL
  752. OR julianday(remote_lease_expires_at) <= julianday(?)
  753. ) ORDER BY created_at,task_id LIMIT 1""",
  754. (_timestamp(self.clock),),
  755. ).fetchone()
  756. finally:
  757. connection.close()
  758. return str(row[0]) if row is not None else None
  759. def _apply_reconcile(self, value: Mapping[str, object]) -> tuple[list[str], list[dict]]:
  760. legacy = {"cancelled_task_ids", "release_offers"}
  761. paged = {
  762. "cancelled_task_ids", "cancel_next_cursor", "release_offers",
  763. "release_next_cursor", "release_baseline",
  764. }
  765. if (
  766. not isinstance(value, Mapping)
  767. or (set(value) != legacy and set(value) != paged)
  768. ):
  769. raise EdgeTransportError("reconcile response schema is invalid")
  770. cancelled = value["cancelled_task_ids"]
  771. releases = value["release_offers"]
  772. if not isinstance(cancelled, list) or not isinstance(releases, list) or len(cancelled) > 100 or len(releases) > 100:
  773. raise EdgeTransportError("reconcile response exceeds safe limits")
  774. applied: list[str] = []
  775. for task_id in cancelled:
  776. if not isinstance(task_id, str) or len(task_id.encode()) > 255:
  777. raise EdgeTransportError("cancel reconciliation is invalid")
  778. self._cancelled.add(task_id)
  779. queued = self.queue.get_task(task_id)
  780. if queued is not None and queued.status in {"pending", "leased"}:
  781. self.queue.cancel_with_outcome(task_id)
  782. applied.append(task_id)
  783. release_values: list[dict] = []
  784. for release in releases:
  785. if not isinstance(release, Mapping):
  786. raise EdgeTransportError("release reconciliation is invalid")
  787. release_values.append(dict(release))
  788. if set(value) == paged and value["release_baseline"] is not None:
  789. baseline = value["release_baseline"]
  790. if not isinstance(baseline, Mapping):
  791. raise EdgeTransportError("release baseline is invalid")
  792. if all(
  793. item.get("release_id") != baseline.get("release_id")
  794. for item in release_values
  795. ):
  796. release_values.append(dict(baseline))
  797. return applied, release_values
  798. def _drain_outcome(self) -> str | None:
  799. outcome = self.queue.claim_outcome(self.worker_id)
  800. if outcome is None:
  801. return None
  802. try:
  803. control_lease = self.queue.recover_remote_lease(outcome.task_id)
  804. sender = getattr(self.transport, "task_outcome", None)
  805. if not callable(sender):
  806. raise EdgeTransportError("task outcome transport is unavailable")
  807. acknowledgement = sender(
  808. outcome.task_id,
  809. outcome.outcome["outcome"],
  810. control_lease,
  811. outcome.outcome["safe_summary"],
  812. )
  813. self.queue.acknowledge_outcome(
  814. outcome.task_id,
  815. lease_token=str(outcome.lease_token),
  816. acknowledgement=acknowledgement,
  817. )
  818. return outcome.task_id
  819. except EdgeAuthenticationStopped:
  820. self._stopped = True
  821. self.queue.fail_outcome(
  822. outcome.task_id,
  823. lease_token=str(outcome.lease_token),
  824. error_code="credential_stopped",
  825. )
  826. raise
  827. except (EdgeTransportError, EdgeQueueLeaseError):
  828. self.queue.fail_outcome(
  829. outcome.task_id,
  830. lease_token=str(outcome.lease_token),
  831. error_code="control_unavailable",
  832. )
  833. return None
  834. def _process_releases(self, releases: list[dict]) -> None:
  835. if not releases:
  836. return
  837. if self.release_manager is None:
  838. raise EdgeReleaseError("signed release manager is not configured")
  839. for release in releases:
  840. verified = self.release_manager.verify(release, allow_recovery=True)
  841. release_id = str(verified["release_id"])
  842. version = str(verified["version"])
  843. rollback_version = str(verified["rollback_version"])
  844. content_digest = canonical_sha256(
  845. {
  846. key: verified[key]
  847. for key in (
  848. "release_id",
  849. "version",
  850. "rollback_version",
  851. "artifact_digest",
  852. "artifact_name",
  853. "deadline_at",
  854. )
  855. }
  856. )
  857. existing = self.queue.get_release_state(release_id)
  858. state = self.queue.put_release_state(
  859. release_id=release_id,
  860. manifest_digest=content_digest,
  861. version=version,
  862. rollback_version=rollback_version,
  863. )
  864. server_status = str(verified["status"])
  865. if (
  866. existing is None
  867. and server_status == "offered"
  868. and (
  869. rollback_version != self.release_manager.current_version
  870. or self.release_manager._version_tuple(version)
  871. <= self.release_manager._version_tuple(
  872. self.release_manager.current_version
  873. )
  874. )
  875. ):
  876. raise EdgeReleaseError(
  877. "release downgrade or rollback binding was rejected"
  878. )
  879. if server_status == "rolled_back":
  880. while state.status != "rolled_back":
  881. target = (
  882. "accepted"
  883. if state.status == "offered"
  884. else "candidate"
  885. if state.status == "accepted"
  886. else "failed"
  887. if state.status in {"candidate", "installed"}
  888. else "rolled_back"
  889. )
  890. state = self.queue.compare_and_set_release_state(
  891. release_id,
  892. expected_status=state.status,
  893. target_status=target,
  894. previous_version=rollback_version,
  895. current_version=rollback_version,
  896. )
  897. self.release_manager.current_version = rollback_version
  898. continue
  899. if server_status == "installed":
  900. if state.status in {"offered", "accepted"}:
  901. if state.status == "offered":
  902. state = self.queue.compare_and_set_release_state(
  903. release_id,
  904. expected_status="offered",
  905. target_status="accepted",
  906. previous_version=rollback_version,
  907. )
  908. state = self.queue.compare_and_set_release_state(
  909. release_id,
  910. expected_status="accepted",
  911. target_status="candidate",
  912. previous_version=rollback_version,
  913. )
  914. if state.status == "candidate":
  915. state = self.queue.compare_and_set_release_state(
  916. release_id,
  917. expected_status="candidate",
  918. target_status="installed",
  919. previous_version=rollback_version,
  920. current_version=version,
  921. )
  922. self.release_manager.current_version = version
  923. continue
  924. if server_status == "offered":
  925. if state.status == "offered":
  926. state = self.queue.compare_and_set_release_state(
  927. release_id,
  928. expected_status="offered",
  929. target_status="accepted",
  930. previous_version=rollback_version,
  931. )
  932. self.transport.acknowledge_release(
  933. release_id,
  934. "accepted",
  935. {"release_status": "accepted", "version": version},
  936. )
  937. if state.status == "installed":
  938. self.transport.acknowledge_release(
  939. release_id,
  940. "installed",
  941. {"release_status": "installed", "version": version},
  942. )
  943. self.release_manager.current_version = version
  944. continue
  945. if state.status == "rolled_back":
  946. if server_status == "accepted":
  947. self.transport.acknowledge_release(
  948. release_id,
  949. "failed",
  950. {"release_status": "failed", "version": version},
  951. )
  952. self.transport.acknowledge_release(
  953. release_id,
  954. "rollback",
  955. {"release_status": "rolled_back", "version": rollback_version},
  956. )
  957. continue
  958. if state.status == "failed" or server_status == "failed":
  959. if state.status != "failed":
  960. state = self.queue.compare_and_set_release_state(
  961. release_id,
  962. expected_status=state.status,
  963. target_status="failed",
  964. previous_version=rollback_version,
  965. )
  966. if not self.release_manager.rollback(rollback_version):
  967. raise EdgeReleaseError("release rollback health confirmation failed")
  968. state = self.queue.compare_and_set_release_state(
  969. release_id,
  970. expected_status="failed",
  971. target_status="rolled_back",
  972. previous_version=rollback_version,
  973. current_version=rollback_version,
  974. )
  975. self.release_manager.current_version = rollback_version
  976. if server_status != "failed":
  977. self.transport.acknowledge_release(
  978. release_id,
  979. "failed",
  980. {"release_status": "failed", "version": version},
  981. )
  982. self.transport.acknowledge_release(
  983. release_id,
  984. "rollback",
  985. {"release_status": "rolled_back", "version": rollback_version},
  986. )
  987. continue
  988. if state.status == "candidate":
  989. if not self.release_manager.rollback(rollback_version):
  990. raise EdgeReleaseError("uncertain candidate rollback failed")
  991. state = self.queue.compare_and_set_release_state(
  992. release_id,
  993. expected_status="candidate",
  994. target_status="failed",
  995. previous_version=rollback_version,
  996. )
  997. state = self.queue.compare_and_set_release_state(
  998. release_id,
  999. expected_status="failed",
  1000. target_status="rolled_back",
  1001. previous_version=rollback_version,
  1002. current_version=rollback_version,
  1003. )
  1004. self.transport.acknowledge_release(
  1005. release_id,
  1006. "failed",
  1007. {"release_status": "failed", "version": version},
  1008. )
  1009. self.transport.acknowledge_release(
  1010. release_id,
  1011. "rollback",
  1012. {"release_status": "rolled_back", "version": rollback_version},
  1013. )
  1014. continue
  1015. state = self.queue.compare_and_set_release_state(
  1016. release_id,
  1017. expected_status="accepted",
  1018. target_status="candidate",
  1019. previous_version=rollback_version,
  1020. )
  1021. result = self.release_manager.install(release)
  1022. if result["outcome"] == "installed":
  1023. self.queue.compare_and_set_release_state(
  1024. release_id,
  1025. expected_status="candidate",
  1026. target_status="installed",
  1027. previous_version=rollback_version,
  1028. current_version=version,
  1029. )
  1030. self.transport.acknowledge_release(
  1031. release_id, "installed", result["summary"]
  1032. )
  1033. else:
  1034. self.queue.compare_and_set_release_state(
  1035. release_id,
  1036. expected_status="candidate",
  1037. target_status="failed",
  1038. previous_version=rollback_version,
  1039. )
  1040. self.queue.compare_and_set_release_state(
  1041. release_id,
  1042. expected_status="failed",
  1043. target_status="rolled_back",
  1044. previous_version=rollback_version,
  1045. current_version=rollback_version,
  1046. )
  1047. self.transport.acknowledge_release(
  1048. release_id,
  1049. "failed",
  1050. {"release_status": "failed", "version": version},
  1051. )
  1052. self.transport.acknowledge_release(
  1053. release_id, "rollback", result["summary"]
  1054. )
  1055. def heartbeat_once(self) -> Mapping[str, object]:
  1056. """Send only the policy-approved bounded health summary."""
  1057. summary = self.safe_diagnostic()
  1058. self.policy.validate_approved_payload("health_summary", summary)
  1059. return self.transport.heartbeat(summary)
  1060. def run_once(self) -> dict[str, object]:
  1061. if self._stopped:
  1062. raise EdgeAuthenticationStopped("edge agent is stopped")
  1063. self._drain_artifact_cleanup()
  1064. try:
  1065. cancelled, releases = self._apply_reconcile(self.transport.reconcile())
  1066. self._drain_outcome()
  1067. self._process_releases(releases)
  1068. sent = self._drain_event()
  1069. if sent is not None:
  1070. return {"cancelled": cancelled, "sent": sent, "executed": None}
  1071. if self._has_unresolved_events() and self._needs_remote_lease_task is None:
  1072. return {"cancelled": cancelled, "sent": None, "executed": None}
  1073. pulled = self.transport.pull_task()
  1074. if pulled is not None:
  1075. if not isinstance(pulled, tuple) or len(pulled) != 3:
  1076. raise EdgeTransportError("pulled task response is invalid")
  1077. accepted = self.accept_pulled_task(pulled[0], pulled[1], pulled[2])
  1078. if self._needs_remote_lease_task is not None:
  1079. if accepted.task_id == self._needs_remote_lease_task:
  1080. self._needs_remote_lease_task = None
  1081. return {
  1082. "cancelled": cancelled,
  1083. "sent": None,
  1084. "executed": None,
  1085. }
  1086. elif self._needs_remote_lease_task is not None:
  1087. return {"cancelled": cancelled, "sent": None, "executed": None}
  1088. expired_pending = self._expired_pending_remote_task()
  1089. if expired_pending is not None:
  1090. self._needs_remote_lease_task = expired_pending
  1091. return {"cancelled": cancelled, "sent": None, "executed": None}
  1092. except EdgeAuthenticationStopped:
  1093. self._stopped = True
  1094. raise
  1095. task_row = self.queue.claim(self.worker_id)
  1096. if task_row is None:
  1097. return {"cancelled": cancelled, "sent": None, "executed": None}
  1098. contract = EdgeTaskContract.from_mapping(task_row.task)
  1099. now = _now(self.clock)
  1100. if (
  1101. datetime.fromisoformat(contract.deadline_at.replace("Z", "+00:00")) <= now
  1102. or task_row.lease_expires_at is None
  1103. or datetime.fromisoformat(
  1104. task_row.lease_expires_at.replace("Z", "+00:00")
  1105. ) <= now
  1106. or task_row.remote_lease_expires_at is None
  1107. or datetime.fromisoformat(
  1108. task_row.remote_lease_expires_at.replace("Z", "+00:00")
  1109. ) <= now
  1110. ):
  1111. self.queue.fail(
  1112. contract.task_id,
  1113. lease_token=str(task_row.lease_token),
  1114. error_code="execution_lease_expired",
  1115. retryable=False,
  1116. )
  1117. self._needs_remote_lease_task = contract.task_id
  1118. return {"cancelled": cancelled, "sent": None, "executed": None}
  1119. try:
  1120. control_lease = self.queue.recover_remote_lease(contract.task_id)
  1121. except EdgeQueueLeaseError:
  1122. self.queue.fail(
  1123. contract.task_id,
  1124. lease_token=str(task_row.lease_token),
  1125. error_code="control_lease_unavailable",
  1126. )
  1127. return {"cancelled": cancelled, "sent": None, "executed": None}
  1128. try:
  1129. result = self.runner.execute(
  1130. contract, lambda: self._cancel_requested(contract.task_id)
  1131. )
  1132. if self._cancel_requested(contract.task_id):
  1133. raise EdgeTaskCancelled("edge task was cancelled")
  1134. event = self._event(contract, result)
  1135. self.queue.complete_with_event(
  1136. contract.task_id,
  1137. lease_token=str(task_row.lease_token),
  1138. remote_lease_token=control_lease,
  1139. event=event,
  1140. artifacts=self._artifact_metadata(contract, result),
  1141. )
  1142. self._drain_artifact_cleanup()
  1143. except EdgeTaskCancelled:
  1144. self.queue.cancel_with_outcome(contract.task_id)
  1145. self._drain_outcome()
  1146. return {"cancelled": sorted(set(cancelled) | {contract.task_id}), "sent": None, "executed": None}
  1147. except (EdgePolicyError, EdgeContractError, EdgeQueueConflictError):
  1148. self.queue.fail(
  1149. contract.task_id,
  1150. lease_token=str(task_row.lease_token),
  1151. error_code="execution_boundary_rejected",
  1152. retryable=False,
  1153. )
  1154. raise
  1155. except EdgeQueueLeaseError:
  1156. current = self.queue.get_task(contract.task_id)
  1157. if current is not None and current.status == "cancelled":
  1158. return {
  1159. "cancelled": sorted(set(cancelled) | {contract.task_id}),
  1160. "sent": None,
  1161. "executed": None,
  1162. }
  1163. if current is not None and current.status == "leased":
  1164. self.queue.fail(
  1165. contract.task_id,
  1166. lease_token=str(task_row.lease_token),
  1167. error_code="execution_lease_expired",
  1168. retryable=False,
  1169. )
  1170. self._needs_remote_lease_task = contract.task_id
  1171. return {"cancelled": cancelled, "sent": None, "executed": None}
  1172. except EdgeAuthenticationStopped:
  1173. self._stopped = True
  1174. self.queue.fail(
  1175. contract.task_id,
  1176. lease_token=str(task_row.lease_token),
  1177. error_code="credential_stopped",
  1178. )
  1179. raise
  1180. except Exception as exc:
  1181. self.queue.fail(
  1182. contract.task_id,
  1183. lease_token=str(task_row.lease_token),
  1184. error_code="runner_execution_failed",
  1185. retryable=False,
  1186. )
  1187. raise EdgeAgentError("edge runner execution failed") from exc
  1188. sent = self._drain_event()
  1189. return {"cancelled": cancelled, "sent": sent, "executed": contract.task_id}
  1190. def safe_diagnostic(self) -> dict[str, object]:
  1191. connection = sqlite3.connect(f"file:{self.queue.db_path}?mode=ro", uri=True)
  1192. try:
  1193. pending = connection.execute(
  1194. "SELECT COUNT(*) FROM edge_tasks WHERE status IN ('pending','leased')"
  1195. ).fetchone()[0]
  1196. events = connection.execute(
  1197. "SELECT COUNT(*) FROM edge_outbound_events WHERE status != 'acknowledged'"
  1198. ).fetchone()[0]
  1199. finally:
  1200. connection.close()
  1201. return {
  1202. "agent_status": "stopped" if self._stopped else "ready",
  1203. "gateway_id_digest": self.gateway_id_digest,
  1204. "queue_pending_count": int(pending),
  1205. "queue_unacknowledged_count": int(events),
  1206. "version": self.config.version,
  1207. }