| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273 |
- """Durable pull-only edge agent, closed runner adapter, and release verifier."""
- from __future__ import annotations
- import base64
- import hashlib
- import os
- import re
- import sqlite3
- import stat
- from collections.abc import Callable, Mapping
- from datetime import UTC, datetime, timedelta
- from pathlib import Path, PurePosixPath
- from types import MappingProxyType
- from cryptography.exceptions import InvalidSignature
- from cryptography.hazmat.primitives import hashes
- from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
- from cryptography.hazmat.primitives.ciphers.aead import AESGCM
- from cryptography.hazmat.primitives.kdf.hkdf import HKDF
- from app.core.edge_gateway.contracts import (
- EdgeContractError,
- EdgeEventContract,
- EdgeTaskContract,
- SignedTaskEnvelope,
- canonical_json_bytes,
- canonical_sha256,
- canonical_timestamp,
- stable_event_id,
- )
- from app.core.edge_gateway.policy import EdgeEgressPolicy, EdgePolicyError
- from app.core.edge_gateway.queue import (
- EdgeQueueConflictError,
- EdgeQueueLeaseError,
- SqliteEdgeQueue,
- )
- from app.edge_gateway.bootstrap import EdgeBootstrapConfig
- from app.edge_gateway.transport import EdgeAuthenticationStopped, EdgeTransportError
- _SHA256 = re.compile(r"^[0-9a-f]{64}$")
- _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})$")
- _RELEASE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._:-]{0,254}$")
- EDGE_NODE_BY_OPERATION = MappingProxyType(
- {
- "collect": "edge.collect",
- "profile": "edge.profile",
- "quality": "quality.check",
- "lineage": "edge.lineage",
- "controlled_query": "edge.controlled_query",
- }
- )
- EVENT_CLASSIFICATION_BY_OPERATION = MappingProxyType(
- {
- "collect": "desensitized_metadata",
- "profile": "statistics",
- "quality": "evidence",
- "lineage": "lineage",
- "controlled_query": "statistics",
- }
- )
- DEFAULT_GOVERNED_PURPOSES = frozenset(
- {
- "inventory",
- "governed-inventory",
- "metadata-inventory",
- "profile-statistics",
- "quality-evaluation",
- "lineage-capture",
- "controlled-query",
- }
- )
- _RUNNER_REQUEST_FIELDS = frozenset(
- {
- "task_id",
- "operation",
- "purpose",
- "classification",
- "environment",
- "network_zone",
- "idempotency_key",
- "deadline_at",
- }
- )
- _RUNNER_RESULT_FIELDS = frozenset(
- {"classification", "payload", "local_artifact_digest", "local_artifact_ref"}
- )
- class EdgeAgentError(RuntimeError):
- """Bounded edge runtime failure."""
- class EdgeTaskCancelled(EdgeAgentError):
- """Cooperative cancellation was observed before a terminal event."""
- class EdgeReleaseError(EdgeAgentError):
- """A candidate release failed a closed verification or activation gate."""
- class LocalArtifactStore:
- """Validate and delete local artifacts without exposing their paths outbound."""
- def __init__(self, root: str | Path, *, clock: Callable[[], object] | None = None):
- candidate = Path(root)
- if not candidate.is_absolute() or not candidate.is_dir() or candidate.is_symlink():
- raise ValueError("artifact root must be an existing absolute directory")
- self.root = candidate.resolve()
- self.clock = clock or (lambda: datetime.now(UTC))
- def inspect(self, reference: object, expected_digest: object) -> dict[str, str]:
- if (
- not isinstance(reference, str)
- or not reference
- or "\x00" in reference
- or len(reference.encode("utf-8")) > 2_048
- or not isinstance(expected_digest, str)
- or not _SHA256.fullmatch(expected_digest)
- ):
- raise EdgePolicyError("local artifact reference is invalid")
- raw = Path(reference)
- path = raw if raw.is_absolute() else self.root / raw
- try:
- resolved = path.resolve(strict=True)
- resolved.relative_to(self.root)
- current = self.root
- for part in resolved.relative_to(self.root).parts:
- current = current / part
- if current.is_symlink():
- raise EdgePolicyError("local artifact symlink is forbidden")
- metadata = resolved.lstat()
- except (OSError, ValueError) as exc:
- raise EdgePolicyError("local artifact escapes its configured root") from exc
- if not stat.S_ISREG(metadata.st_mode):
- raise EdgePolicyError("local artifact must be a regular file")
- digest = hashlib.sha256()
- try:
- with resolved.open("rb") as stream:
- for chunk in iter(lambda: stream.read(1_048_576), b""):
- digest.update(chunk)
- except OSError as exc:
- raise EdgePolicyError("local artifact cannot be read") from exc
- if digest.hexdigest() != expected_digest:
- raise EdgePolicyError("local artifact digest does not match content")
- canonical_ref = str(resolved)
- return {
- "artifact_digest": expected_digest,
- "artifact_ref": canonical_ref,
- "artifact_ref_hash": canonical_sha256(canonical_ref),
- }
- def delete(self, reference: object, expected_digest: object) -> dict[str, str]:
- inspected = self.inspect(reference, expected_digest)
- try:
- Path(inspected["artifact_ref"]).unlink()
- except OSError as exc:
- raise EdgePolicyError("local artifact deletion failed") from exc
- return {
- "artifact_digest": inspected["artifact_digest"],
- "artifact_ref_hash": inspected["artifact_ref_hash"],
- "deleted_at": canonical_timestamp(
- _now(self.clock).isoformat().replace("+00:00", "Z")
- ),
- "status": "deleted",
- }
- def receipt_for_missing(
- self, reference: object, expected_digest: object, expected_ref_hash: object
- ) -> dict[str, str]:
- if (
- not isinstance(reference, str)
- or not isinstance(expected_digest, str)
- or not _SHA256.fullmatch(expected_digest)
- or not isinstance(expected_ref_hash, str)
- or not _SHA256.fullmatch(expected_ref_hash)
- ):
- raise EdgePolicyError("missing artifact recovery metadata is invalid")
- path = Path(reference)
- try:
- path.parent.resolve(strict=True).relative_to(self.root)
- canonical_ref = str(path.resolve(strict=False))
- except (OSError, ValueError) as exc:
- raise EdgePolicyError("missing artifact escapes its configured root") from exc
- if path.exists() or path.is_symlink() or canonical_sha256(canonical_ref) != expected_ref_hash:
- raise EdgePolicyError("missing artifact recovery proof is invalid")
- return {
- "artifact_digest": expected_digest,
- "artifact_ref_hash": expected_ref_hash,
- "deleted_at": canonical_timestamp(
- _now(self.clock).isoformat().replace("+00:00", "Z")
- ),
- "status": "deleted",
- }
- def _now(clock: Callable[[], object]) -> datetime:
- value = clock()
- if not isinstance(value, datetime) or value.tzinfo is None or value.utcoffset() is None:
- raise ValueError("edge clock must return a timezone-aware datetime")
- return value.astimezone(UTC)
- def _timestamp(clock: Callable[[], object]) -> str:
- return canonical_timestamp(_now(clock).isoformat().replace("+00:00", "Z"))
- def _remote_lease_codec(config: EdgeBootstrapConfig):
- """Protect recoverable control leases without persisting the credential."""
- key = HKDF(
- algorithm=hashes.SHA256(),
- length=32,
- salt=bytes.fromhex(config.certificate_sha256),
- info=(
- f"dataops-edge-remote-lease-v1:{config.gateway_id}:"
- f"{config.environment}:{config.network_zone}"
- ).encode(),
- ).derive(config.credential.encode())
- cipher = AESGCM(key)
- aad = f"{config.gateway_id}:{config.policy_digest}".encode()
- def encode(token: str) -> str:
- nonce = os.urandom(12)
- protected = cipher.encrypt(nonce, token.encode("utf-8"), aad)
- return base64.urlsafe_b64encode(nonce + protected).decode("ascii")
- def decode(value: str) -> str:
- raw = base64.b64decode(value.encode("ascii"), altchars=b"-_", validate=True)
- if len(raw) < 29:
- raise ValueError("protected remote lease is invalid")
- return cipher.decrypt(raw[:12], raw[12:], aad).decode("utf-8")
- return encode, decode
- class EdgeRunnerAdapter:
- """Map governed operations to exact local runner handlers.
- No control-plane supplied parameters, module names, commands, SQL, paths, or
- network destinations are accepted by this adapter. Local handlers resolve
- their governed configuration using the immutable task identity.
- """
- def __init__(
- self,
- executor: Callable[[str, Mapping[str, object], Callable[[], bool]], object],
- *,
- policy: EdgeEgressPolicy | None = None,
- allowed_purposes: set[str] | frozenset[str] = DEFAULT_GOVERNED_PURPOSES,
- ) -> None:
- if not callable(executor):
- raise ValueError("an explicit edge runner executor is required")
- self._executor = executor
- self.policy = policy or EdgeEgressPolicy(allowed_control_hosts=set())
- if not isinstance(allowed_purposes, (set, frozenset)) or not allowed_purposes:
- raise ValueError("explicit governed purposes are required")
- if any(
- not isinstance(item, str)
- or not item
- or len(item.encode("utf-8")) > 255
- for item in allowed_purposes
- ):
- raise ValueError("governed purpose allowlist is invalid")
- self.allowed_purposes = frozenset(allowed_purposes)
- def execute(
- self, task: EdgeTaskContract, cancel_requested: Callable[[], bool]
- ) -> dict[str, object]:
- task = EdgeTaskContract.from_mapping(task.to_mapping())
- node = EDGE_NODE_BY_OPERATION.get(task.operation)
- if node is None:
- raise EdgePolicyError("edge task operation is not runner-approved")
- if task.purpose not in self.allowed_purposes:
- raise EdgePolicyError("edge task purpose is not runner-approved")
- if cancel_requested():
- raise EdgeTaskCancelled("edge task was cancelled")
- request = {
- "task_id": task.task_id,
- "operation": task.operation,
- "purpose": task.purpose,
- "classification": task.classification,
- "environment": task.environment,
- "network_zone": task.network_zone,
- "idempotency_key": task.idempotency_key,
- "deadline_at": task.deadline_at,
- }
- if set(request) != _RUNNER_REQUEST_FIELDS: # pragma: no cover - invariant
- raise EdgePolicyError("edge runner request schema is invalid")
- try:
- raw = self._executor(node, MappingProxyType(request), cancel_requested)
- except EdgeTaskCancelled:
- raise
- except (MemoryError, RecursionError) as exc:
- raise EdgePolicyError("edge runner result exceeds safe limits") from exc
- if cancel_requested():
- raise EdgeTaskCancelled("edge task was cancelled")
- if not isinstance(raw, Mapping):
- raise EdgePolicyError("edge runner result must be a mapping")
- unknown = set(raw) - _RUNNER_RESULT_FIELDS
- missing = {"classification", "payload"} - set(raw)
- if unknown or missing:
- raise EdgePolicyError("edge runner result schema is invalid")
- classification = raw["classification"]
- expected = EVENT_CLASSIFICATION_BY_OPERATION[task.operation]
- if classification != expected:
- raise EdgePolicyError("edge runner result classification is invalid")
- payload = raw["payload"]
- if not isinstance(payload, Mapping):
- raise EdgePolicyError("edge runner result payload must be a mapping")
- try:
- self.policy.validate_approved_payload(expected, payload)
- except (EdgePolicyError, EdgeContractError, MemoryError, RecursionError) as exc:
- raise EdgePolicyError("edge runner result is not approved for egress") from exc
- artifact_digest = raw.get("local_artifact_digest")
- artifact_ref = raw.get("local_artifact_ref")
- if (artifact_digest is None) != (artifact_ref is None):
- raise EdgePolicyError("local artifact digest and reference must be paired")
- if artifact_digest is not None:
- if not isinstance(artifact_digest, str) or not _SHA256.fullmatch(artifact_digest):
- raise EdgePolicyError("local artifact digest is invalid")
- if (
- not isinstance(artifact_ref, str)
- or not artifact_ref
- or "\x00" in artifact_ref
- or len(artifact_ref.encode("utf-8")) > 2_048
- ):
- raise EdgePolicyError("local artifact reference is invalid")
- return {
- "classification": expected,
- "payload": dict(payload),
- "local_artifact_digest": artifact_digest,
- "local_artifact_ref": artifact_ref,
- }
- class SignedReleaseManager:
- """Verify Ed25519 manifests and activate candidates through injected primitives."""
- _MANIFEST_FIELDS = frozenset(
- {
- "release_id",
- "version",
- "rollback_version",
- "artifact_digest",
- "artifact_name",
- "deadline_at",
- "signature_algorithm",
- "key_id",
- "manifest_digest",
- "signature",
- "status",
- }
- )
- def __init__(
- self,
- *,
- current_version: str,
- trusted_keys: Mapping[str, str],
- artifact_loader: Callable[[Mapping[str, object]], bytes],
- activator: Callable[[str, bytes], bool],
- rollback: Callable[[str], bool],
- clock: Callable[[], object] | None = None,
- ) -> None:
- if not _VERSION.fullmatch(current_version):
- raise ValueError("current release version is invalid")
- if any(int(part) > 2_147_483_647 for part in current_version.split(".")):
- raise ValueError("current release version is invalid")
- if not trusted_keys:
- raise ValueError("trusted release keys are required")
- self.current_version = current_version
- self.trusted_keys = dict(trusted_keys)
- self.artifact_loader = artifact_loader
- self.activator = activator
- self.rollback = rollback
- self.clock = clock or (lambda: datetime.now(UTC))
- self._seen: dict[str, str] = {}
- @staticmethod
- def _version_tuple(value: str) -> tuple[int, ...]:
- return tuple(int(item) for item in value.split("."))
- def verify(
- self,
- value: Mapping[str, object],
- *,
- allow_recovery: bool = False,
- ) -> dict[str, object]:
- if not isinstance(value, Mapping) or set(value) != self._MANIFEST_FIELDS:
- raise EdgeReleaseError("release manifest schema is invalid")
- manifest = dict(value)
- release_id = manifest["release_id"]
- version = manifest["version"]
- rollback_version = manifest["rollback_version"]
- if not isinstance(release_id, str) or not _RELEASE_ID.fullmatch(release_id):
- raise EdgeReleaseError("release identifier is invalid")
- status = manifest.get("status")
- if status not in {"offered", "accepted", "failed", "installed", "rolled_back"}:
- raise EdgeReleaseError("release state is invalid")
- if (
- not isinstance(version, str)
- or not isinstance(rollback_version, str)
- or not _VERSION.fullmatch(version)
- or not _VERSION.fullmatch(rollback_version)
- or any(int(part) > 2_147_483_647 for part in version.split("."))
- or any(int(part) > 2_147_483_647 for part in rollback_version.split("."))
- ):
- raise EdgeReleaseError("release downgrade or rollback binding was rejected")
- if not allow_recovery and status in {"offered", "accepted"} and (
- rollback_version != self.current_version
- or self._version_tuple(version) <= self._version_tuple(self.current_version)
- ):
- raise EdgeReleaseError("release downgrade or rollback binding was rejected")
- artifact_name = manifest["artifact_name"]
- if (
- not isinstance(artifact_name, str)
- or not artifact_name
- or len(artifact_name.encode("utf-8")) > 255
- or PurePosixPath(artifact_name).name != artifact_name
- or artifact_name in {".", ".."}
- or any(character in artifact_name for character in ("/", "\\", "\x00"))
- ):
- raise EdgeReleaseError("release artifact path was rejected")
- if manifest["signature_algorithm"] != "Ed25519":
- raise EdgeReleaseError("release signature algorithm is not approved")
- key_id = manifest["key_id"]
- key_hex = self.trusted_keys.get(key_id) if isinstance(key_id, str) else None
- if key_hex is None:
- raise EdgeReleaseError("release signing key is not trusted")
- for field in ("artifact_digest", "manifest_digest"):
- if not isinstance(manifest[field], str) or not _SHA256.fullmatch(manifest[field]):
- raise EdgeReleaseError("release digest is invalid")
- try:
- deadline = datetime.fromisoformat(str(manifest["deadline_at"]).replace("Z", "+00:00"))
- except ValueError as exc:
- raise EdgeReleaseError("release deadline is invalid") from exc
- if deadline.tzinfo is None or deadline.astimezone(UTC) <= _now(self.clock):
- raise EdgeReleaseError("release deadline has expired")
- try:
- canonical_deadline = canonical_timestamp(manifest["deadline_at"], "deadline_at")
- except EdgeContractError as exc:
- raise EdgeReleaseError("release deadline is invalid") from exc
- if canonical_deadline != manifest["deadline_at"]:
- raise EdgeReleaseError("release deadline must use canonical UTC form")
- unsigned = {
- key: manifest[key]
- for key in sorted(self._MANIFEST_FIELDS - {"manifest_digest", "signature"})
- }
- expected_digest = canonical_sha256(unsigned)
- if expected_digest != manifest["manifest_digest"]:
- raise EdgeReleaseError("release manifest digest does not match")
- content_digest = canonical_sha256(
- {
- key: manifest[key]
- for key in (
- "release_id",
- "version",
- "rollback_version",
- "artifact_digest",
- "artifact_name",
- "deadline_at",
- )
- }
- )
- prior = self._seen.get(release_id)
- if prior is not None and prior != content_digest:
- raise EdgeReleaseError("release identifier replay changed its digest")
- try:
- signature = bytes.fromhex(str(manifest["signature"]))
- Ed25519PublicKey.from_public_bytes(bytes.fromhex(key_hex)).verify(
- signature, canonical_json_bytes(unsigned)
- )
- except (ValueError, InvalidSignature) as exc:
- raise EdgeReleaseError("release signature verification failed") from exc
- self._seen[release_id] = content_digest
- return manifest
- def install(self, value: Mapping[str, object]) -> dict[str, object]:
- manifest = self.verify(value)
- release_id = str(manifest["release_id"])
- previous = self.current_version
- try:
- artifact = self.artifact_loader(MappingProxyType(manifest))
- if not isinstance(artifact, bytes) or hashlib.sha256(artifact).hexdigest() != manifest["artifact_digest"]:
- raise EdgeReleaseError("release artifact digest does not match")
- if not self.activator(str(manifest["version"]), artifact):
- raise EdgeReleaseError("release candidate health confirmation failed")
- except Exception:
- rolled_back = bool(self.rollback(previous))
- return {
- "release_id": release_id,
- "outcome": "rollback" if rolled_back else "failed",
- "summary": {
- "release_status": "rolled_back" if rolled_back else "failed",
- "version": previous,
- },
- }
- self.current_version = str(manifest["version"])
- return {
- "release_id": release_id,
- "outcome": "installed",
- "summary": {"release_status": "installed", "version": self.current_version},
- }
- class EdgeAgent:
- """One durable worker for pull, execute, persist, and acknowledge cycles."""
- def __init__(
- self,
- config: EdgeBootstrapConfig,
- transport,
- runner: EdgeRunnerAdapter,
- *,
- worker_id: str = "edge-worker-1",
- clock: Callable[[], object] | None = None,
- random: Callable[[], float] | None = None,
- release_manager: SignedReleaseManager | None = None,
- artifact_store: LocalArtifactStore | None = None,
- ) -> None:
- if not isinstance(config, EdgeBootstrapConfig):
- raise ValueError("validated edge bootstrap configuration is required")
- if not isinstance(runner, EdgeRunnerAdapter):
- raise ValueError("validated edge runner adapter is required")
- if not isinstance(worker_id, str) or not worker_id or len(worker_id.encode()) > 255:
- raise ValueError("worker_id is invalid")
- self.config = config
- self.transport = transport
- self.runner = runner
- self.worker_id = worker_id
- self.clock = clock or (lambda: datetime.now(UTC))
- self.random = random or __import__("random").random
- self.policy = EdgeEgressPolicy(
- allowed_control_hosts=set(config.allowed_control_hosts),
- allowed_proxy_hosts=set(config.allowed_proxy_hosts),
- )
- lease_encoder, lease_decoder = _remote_lease_codec(config)
- self.queue = SqliteEdgeQueue(
- config.queue_path,
- clock=self.clock,
- egress_policy=self.policy,
- random_source=self.random,
- remote_lease_encoder=lease_encoder,
- remote_lease_decoder=lease_decoder,
- )
- self.artifact_store = artifact_store or LocalArtifactStore(
- config.artifact_root, clock=self.clock
- )
- self.release_manager = release_manager
- self._cancelled: set[str] = set()
- self._stopped = False
- self._needs_remote_lease_task: str | None = None
- self._trusted_task_keys = {
- key_id: bytes.fromhex(public_key)
- for key_id, public_key in config.trusted_task_keys.items()
- }
- self.gateway_id_digest = hashlib.sha256(config.gateway_id.encode()).hexdigest()
- def retry_delay(self, attempt: int) -> float:
- if isinstance(attempt, bool) or not isinstance(attempt, int) or attempt < 1:
- raise ValueError("retry attempt is invalid")
- base = min(300.0, float(2 ** min(attempt - 1, 8)))
- jitter = self.random()
- if not isinstance(jitter, (int, float)) or not 0 <= jitter <= 1:
- raise ValueError("random source returned an invalid value")
- return min(300.0, base * (0.75 + 0.5 * float(jitter)))
- def _validate_binding(self, contract: EdgeTaskContract) -> None:
- expected = (
- self.config.gateway_id,
- self.config.environment,
- self.config.network_zone,
- self.config.policy_digest,
- )
- actual = (
- contract.gateway_id,
- contract.environment,
- contract.network_zone,
- contract.policy_digest,
- )
- if actual != expected:
- raise EdgePolicyError("pulled task binding does not match this edge gateway")
- def accept_pulled_task(
- self,
- envelope_value: Mapping[str, object],
- lease_token: str,
- lease_expires_at: str,
- ):
- if not isinstance(lease_token, str) or not lease_token.startswith("dopl_") or len(lease_token) > 128:
- raise EdgePolicyError("control task lease is invalid")
- try:
- envelope = SignedTaskEnvelope.verify_mapping(
- envelope_value,
- authority_keys=self._trusted_task_keys,
- now=_now(self.clock),
- allowed_future_skew_seconds=(
- self.config.task_authority_clock_skew_seconds
- ),
- )
- except EdgeContractError as exc:
- raise EdgePolicyError("signed task authority was rejected") from exc
- self._validate_binding(envelope.task)
- return self.queue.accept_signed_task(
- envelope,
- remote_lease_token=lease_token,
- remote_lease_expires_at=lease_expires_at,
- remote_attempt=envelope.task.attempt,
- )
- def _cancel_requested(self, task_id: str) -> bool:
- if task_id in self._cancelled:
- return True
- probe = getattr(self.transport, "cancel_requested", None)
- if callable(probe):
- try:
- if bool(probe(task_id)):
- self._cancelled.add(task_id)
- return True
- except EdgeAuthenticationStopped:
- self._stopped = True
- raise
- except EdgeTransportError:
- return task_id in self._cancelled
- return False
- def _event(self, task: EdgeTaskContract, result: Mapping[str, object]) -> EdgeEventContract:
- identity = {
- "task_id": task.task_id,
- "gateway_id": task.gateway_id,
- "environment": task.environment,
- "network_zone": task.network_zone,
- "purpose": task.purpose,
- "classification": result["classification"],
- "contract_version": task.contract_version,
- "occurred_at": _timestamp(self.clock),
- "attempt": task.attempt,
- "idempotency_key": task.idempotency_key,
- "policy_digest": task.policy_digest,
- "payload": result["payload"],
- }
- return EdgeEventContract(event_id=stable_event_id(identity), **identity) # type: ignore[arg-type]
- def _artifact_metadata(
- self,
- task: EdgeTaskContract,
- result: Mapping[str, object],
- ) -> tuple[dict[str, object], ...]:
- digest = result.get("local_artifact_digest")
- reference = result.get("local_artifact_ref")
- if digest is None or reference is None:
- return ()
- inspected = self.artifact_store.inspect(reference, digest)
- retention_class = (
- task.classification
- if task.classification in {"raw", "recent_detail", "evidence"}
- else "raw"
- if task.classification == "restricted"
- else "metadata"
- )
- retention_until = _now(self.clock) + timedelta(
- days=self.policy.retention_days(retention_class)
- )
- return (
- {
- "artifact_digest": inspected["artifact_digest"],
- "artifact_ref": inspected["artifact_ref"],
- "artifact_ref_hash": inspected["artifact_ref_hash"],
- "classification": task.classification,
- "retention_until": canonical_timestamp(
- retention_until.isoformat().replace("+00:00", "Z")
- ),
- },
- )
- def _drain_artifact_cleanup(self) -> str | None:
- artifact = self.queue.claim_artifact_cleanup(self.worker_id)
- if artifact is None:
- return None
- if artifact.artifact_ref is None:
- self.queue.fail_artifact_cleanup(
- artifact.artifact_digest,
- lease_token=str(artifact.lease_token),
- error_code="artifact_reference_missing",
- )
- return None
- try:
- receipt = self.artifact_store.delete(
- artifact.artifact_ref, artifact.artifact_digest
- )
- self.queue.acknowledge_artifact_cleanup(
- artifact.artifact_digest,
- lease_token=str(artifact.lease_token),
- receipt=receipt,
- )
- return artifact.artifact_digest
- except EdgePolicyError:
- try:
- receipt = self.artifact_store.receipt_for_missing(
- artifact.artifact_ref,
- artifact.artifact_digest,
- artifact.artifact_ref_hash,
- )
- self.queue.acknowledge_artifact_cleanup(
- artifact.artifact_digest,
- lease_token=str(artifact.lease_token),
- receipt=receipt,
- )
- return artifact.artifact_digest
- except EdgePolicyError:
- pass
- self.queue.fail_artifact_cleanup(
- artifact.artifact_digest,
- lease_token=str(artifact.lease_token),
- error_code="artifact_cleanup_failed",
- )
- return None
- def _drain_event(self) -> str | None:
- event = self.queue.claim_event(self.worker_id)
- if event is None:
- return None
- task_id = str(event.event["task_id"])
- task = self.queue.get_task(task_id)
- if task is None:
- self.queue.fail_event(
- event.event_id,
- lease_token=str(event.lease_token),
- error_code="control_lease_unavailable",
- )
- return None
- expiry = task.remote_lease_expires_at
- expired = (
- expiry is None
- or datetime.fromisoformat(expiry.replace("Z", "+00:00")) <= _now(self.clock)
- )
- if expired:
- self._needs_remote_lease_task = task_id
- self.queue.fail_event(
- event.event_id,
- lease_token=str(event.lease_token),
- error_code="remote_lease_expired",
- )
- return None
- try:
- control_lease = self.queue.recover_remote_lease(task_id)
- except EdgeQueueLeaseError:
- self._needs_remote_lease_task = task_id
- self.queue.fail_event(
- event.event_id,
- lease_token=str(event.lease_token),
- error_code="control_lease_unavailable",
- )
- return None
- try:
- ack = self.transport.send_event(event.event, control_lease)
- self.queue.acknowledge_event(
- event.event_id,
- lease_token=str(event.lease_token),
- acknowledgement=ack,
- )
- return event.event_id
- except EdgeAuthenticationStopped:
- if not expired:
- self._stopped = True
- else:
- self._needs_remote_lease_task = task_id
- self.queue.fail_event(
- event.event_id,
- lease_token=str(event.lease_token),
- error_code=(
- "credential_stopped" if not expired else "remote_lease_expired"
- ),
- )
- if not expired:
- raise
- return None
- except EdgeTransportError:
- self.queue.fail_event(
- event.event_id,
- lease_token=str(event.lease_token),
- error_code="control_unavailable",
- )
- return None
- def _has_unresolved_events(self) -> bool:
- connection = sqlite3.connect(f"file:{self.queue.db_path}?mode=ro", uri=True)
- try:
- count = connection.execute(
- "SELECT COUNT(*) FROM edge_outbound_events WHERE status IN ('pending','sending')"
- ).fetchone()[0]
- finally:
- connection.close()
- return bool(count)
- def _expired_pending_remote_task(self) -> str | None:
- connection = sqlite3.connect(f"file:{self.queue.db_path}?mode=ro", uri=True)
- try:
- row = connection.execute(
- """SELECT task_id FROM edge_tasks
- WHERE status='pending' AND (
- remote_lease_expires_at IS NULL
- OR julianday(remote_lease_expires_at) <= julianday(?)
- ) ORDER BY created_at,task_id LIMIT 1""",
- (_timestamp(self.clock),),
- ).fetchone()
- finally:
- connection.close()
- return str(row[0]) if row is not None else None
- def _apply_reconcile(self, value: Mapping[str, object]) -> tuple[list[str], list[dict]]:
- legacy = {"cancelled_task_ids", "release_offers"}
- paged = {
- "cancelled_task_ids", "cancel_next_cursor", "release_offers",
- "release_next_cursor", "release_baseline",
- }
- if (
- not isinstance(value, Mapping)
- or (set(value) != legacy and set(value) != paged)
- ):
- raise EdgeTransportError("reconcile response schema is invalid")
- cancelled = value["cancelled_task_ids"]
- releases = value["release_offers"]
- if not isinstance(cancelled, list) or not isinstance(releases, list) or len(cancelled) > 100 or len(releases) > 100:
- raise EdgeTransportError("reconcile response exceeds safe limits")
- applied: list[str] = []
- for task_id in cancelled:
- if not isinstance(task_id, str) or len(task_id.encode()) > 255:
- raise EdgeTransportError("cancel reconciliation is invalid")
- self._cancelled.add(task_id)
- queued = self.queue.get_task(task_id)
- if queued is not None and queued.status in {"pending", "leased"}:
- self.queue.cancel_with_outcome(task_id)
- applied.append(task_id)
- release_values: list[dict] = []
- for release in releases:
- if not isinstance(release, Mapping):
- raise EdgeTransportError("release reconciliation is invalid")
- release_values.append(dict(release))
- if set(value) == paged and value["release_baseline"] is not None:
- baseline = value["release_baseline"]
- if not isinstance(baseline, Mapping):
- raise EdgeTransportError("release baseline is invalid")
- if all(
- item.get("release_id") != baseline.get("release_id")
- for item in release_values
- ):
- release_values.append(dict(baseline))
- return applied, release_values
- def _drain_outcome(self) -> str | None:
- outcome = self.queue.claim_outcome(self.worker_id)
- if outcome is None:
- return None
- try:
- control_lease = self.queue.recover_remote_lease(outcome.task_id)
- sender = getattr(self.transport, "task_outcome", None)
- if not callable(sender):
- raise EdgeTransportError("task outcome transport is unavailable")
- acknowledgement = sender(
- outcome.task_id,
- outcome.outcome["outcome"],
- control_lease,
- outcome.outcome["safe_summary"],
- )
- self.queue.acknowledge_outcome(
- outcome.task_id,
- lease_token=str(outcome.lease_token),
- acknowledgement=acknowledgement,
- )
- return outcome.task_id
- except EdgeAuthenticationStopped:
- self._stopped = True
- self.queue.fail_outcome(
- outcome.task_id,
- lease_token=str(outcome.lease_token),
- error_code="credential_stopped",
- )
- raise
- except (EdgeTransportError, EdgeQueueLeaseError):
- self.queue.fail_outcome(
- outcome.task_id,
- lease_token=str(outcome.lease_token),
- error_code="control_unavailable",
- )
- return None
- def _process_releases(self, releases: list[dict]) -> None:
- if not releases:
- return
- if self.release_manager is None:
- raise EdgeReleaseError("signed release manager is not configured")
- for release in releases:
- verified = self.release_manager.verify(release, allow_recovery=True)
- release_id = str(verified["release_id"])
- version = str(verified["version"])
- rollback_version = str(verified["rollback_version"])
- content_digest = canonical_sha256(
- {
- key: verified[key]
- for key in (
- "release_id",
- "version",
- "rollback_version",
- "artifact_digest",
- "artifact_name",
- "deadline_at",
- )
- }
- )
- existing = self.queue.get_release_state(release_id)
- state = self.queue.put_release_state(
- release_id=release_id,
- manifest_digest=content_digest,
- version=version,
- rollback_version=rollback_version,
- )
- server_status = str(verified["status"])
- if (
- existing is None
- and server_status == "offered"
- and (
- rollback_version != self.release_manager.current_version
- or self.release_manager._version_tuple(version)
- <= self.release_manager._version_tuple(
- self.release_manager.current_version
- )
- )
- ):
- raise EdgeReleaseError(
- "release downgrade or rollback binding was rejected"
- )
- if server_status == "rolled_back":
- while state.status != "rolled_back":
- target = (
- "accepted"
- if state.status == "offered"
- else "candidate"
- if state.status == "accepted"
- else "failed"
- if state.status in {"candidate", "installed"}
- else "rolled_back"
- )
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status=state.status,
- target_status=target,
- previous_version=rollback_version,
- current_version=rollback_version,
- )
- self.release_manager.current_version = rollback_version
- continue
- if server_status == "installed":
- if state.status in {"offered", "accepted"}:
- if state.status == "offered":
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="offered",
- target_status="accepted",
- previous_version=rollback_version,
- )
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="accepted",
- target_status="candidate",
- previous_version=rollback_version,
- )
- if state.status == "candidate":
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="candidate",
- target_status="installed",
- previous_version=rollback_version,
- current_version=version,
- )
- self.release_manager.current_version = version
- continue
- if server_status == "offered":
- if state.status == "offered":
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="offered",
- target_status="accepted",
- previous_version=rollback_version,
- )
- self.transport.acknowledge_release(
- release_id,
- "accepted",
- {"release_status": "accepted", "version": version},
- )
- if state.status == "installed":
- self.transport.acknowledge_release(
- release_id,
- "installed",
- {"release_status": "installed", "version": version},
- )
- self.release_manager.current_version = version
- continue
- if state.status == "rolled_back":
- if server_status == "accepted":
- self.transport.acknowledge_release(
- release_id,
- "failed",
- {"release_status": "failed", "version": version},
- )
- self.transport.acknowledge_release(
- release_id,
- "rollback",
- {"release_status": "rolled_back", "version": rollback_version},
- )
- continue
- if state.status == "failed" or server_status == "failed":
- if state.status != "failed":
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status=state.status,
- target_status="failed",
- previous_version=rollback_version,
- )
- if not self.release_manager.rollback(rollback_version):
- raise EdgeReleaseError("release rollback health confirmation failed")
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="failed",
- target_status="rolled_back",
- previous_version=rollback_version,
- current_version=rollback_version,
- )
- self.release_manager.current_version = rollback_version
- if server_status != "failed":
- self.transport.acknowledge_release(
- release_id,
- "failed",
- {"release_status": "failed", "version": version},
- )
- self.transport.acknowledge_release(
- release_id,
- "rollback",
- {"release_status": "rolled_back", "version": rollback_version},
- )
- continue
- if state.status == "candidate":
- if not self.release_manager.rollback(rollback_version):
- raise EdgeReleaseError("uncertain candidate rollback failed")
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="candidate",
- target_status="failed",
- previous_version=rollback_version,
- )
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="failed",
- target_status="rolled_back",
- previous_version=rollback_version,
- current_version=rollback_version,
- )
- self.transport.acknowledge_release(
- release_id,
- "failed",
- {"release_status": "failed", "version": version},
- )
- self.transport.acknowledge_release(
- release_id,
- "rollback",
- {"release_status": "rolled_back", "version": rollback_version},
- )
- continue
- state = self.queue.compare_and_set_release_state(
- release_id,
- expected_status="accepted",
- target_status="candidate",
- previous_version=rollback_version,
- )
- result = self.release_manager.install(release)
- if result["outcome"] == "installed":
- self.queue.compare_and_set_release_state(
- release_id,
- expected_status="candidate",
- target_status="installed",
- previous_version=rollback_version,
- current_version=version,
- )
- self.transport.acknowledge_release(
- release_id, "installed", result["summary"]
- )
- else:
- self.queue.compare_and_set_release_state(
- release_id,
- expected_status="candidate",
- target_status="failed",
- previous_version=rollback_version,
- )
- self.queue.compare_and_set_release_state(
- release_id,
- expected_status="failed",
- target_status="rolled_back",
- previous_version=rollback_version,
- current_version=rollback_version,
- )
- self.transport.acknowledge_release(
- release_id,
- "failed",
- {"release_status": "failed", "version": version},
- )
- self.transport.acknowledge_release(
- release_id, "rollback", result["summary"]
- )
- def heartbeat_once(self) -> Mapping[str, object]:
- """Send only the policy-approved bounded health summary."""
- summary = self.safe_diagnostic()
- self.policy.validate_approved_payload("health_summary", summary)
- return self.transport.heartbeat(summary)
- def run_once(self) -> dict[str, object]:
- if self._stopped:
- raise EdgeAuthenticationStopped("edge agent is stopped")
- self._drain_artifact_cleanup()
- try:
- cancelled, releases = self._apply_reconcile(self.transport.reconcile())
- self._drain_outcome()
- self._process_releases(releases)
- sent = self._drain_event()
- if sent is not None:
- return {"cancelled": cancelled, "sent": sent, "executed": None}
- if self._has_unresolved_events() and self._needs_remote_lease_task is None:
- return {"cancelled": cancelled, "sent": None, "executed": None}
- pulled = self.transport.pull_task()
- if pulled is not None:
- if not isinstance(pulled, tuple) or len(pulled) != 3:
- raise EdgeTransportError("pulled task response is invalid")
- accepted = self.accept_pulled_task(pulled[0], pulled[1], pulled[2])
- if self._needs_remote_lease_task is not None:
- if accepted.task_id == self._needs_remote_lease_task:
- self._needs_remote_lease_task = None
- return {
- "cancelled": cancelled,
- "sent": None,
- "executed": None,
- }
- elif self._needs_remote_lease_task is not None:
- return {"cancelled": cancelled, "sent": None, "executed": None}
- expired_pending = self._expired_pending_remote_task()
- if expired_pending is not None:
- self._needs_remote_lease_task = expired_pending
- return {"cancelled": cancelled, "sent": None, "executed": None}
- except EdgeAuthenticationStopped:
- self._stopped = True
- raise
- task_row = self.queue.claim(self.worker_id)
- if task_row is None:
- return {"cancelled": cancelled, "sent": None, "executed": None}
- contract = EdgeTaskContract.from_mapping(task_row.task)
- now = _now(self.clock)
- if (
- datetime.fromisoformat(contract.deadline_at.replace("Z", "+00:00")) <= now
- or task_row.lease_expires_at is None
- or datetime.fromisoformat(
- task_row.lease_expires_at.replace("Z", "+00:00")
- ) <= now
- or task_row.remote_lease_expires_at is None
- or datetime.fromisoformat(
- task_row.remote_lease_expires_at.replace("Z", "+00:00")
- ) <= now
- ):
- self.queue.fail(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- error_code="execution_lease_expired",
- retryable=False,
- )
- self._needs_remote_lease_task = contract.task_id
- return {"cancelled": cancelled, "sent": None, "executed": None}
- try:
- control_lease = self.queue.recover_remote_lease(contract.task_id)
- except EdgeQueueLeaseError:
- self.queue.fail(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- error_code="control_lease_unavailable",
- )
- return {"cancelled": cancelled, "sent": None, "executed": None}
- try:
- result = self.runner.execute(
- contract, lambda: self._cancel_requested(contract.task_id)
- )
- if self._cancel_requested(contract.task_id):
- raise EdgeTaskCancelled("edge task was cancelled")
- event = self._event(contract, result)
- self.queue.complete_with_event(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- remote_lease_token=control_lease,
- event=event,
- artifacts=self._artifact_metadata(contract, result),
- )
- self._drain_artifact_cleanup()
- except EdgeTaskCancelled:
- self.queue.cancel_with_outcome(contract.task_id)
- self._drain_outcome()
- return {"cancelled": sorted(set(cancelled) | {contract.task_id}), "sent": None, "executed": None}
- except (EdgePolicyError, EdgeContractError, EdgeQueueConflictError):
- self.queue.fail(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- error_code="execution_boundary_rejected",
- retryable=False,
- )
- raise
- except EdgeQueueLeaseError:
- current = self.queue.get_task(contract.task_id)
- if current is not None and current.status == "cancelled":
- return {
- "cancelled": sorted(set(cancelled) | {contract.task_id}),
- "sent": None,
- "executed": None,
- }
- if current is not None and current.status == "leased":
- self.queue.fail(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- error_code="execution_lease_expired",
- retryable=False,
- )
- self._needs_remote_lease_task = contract.task_id
- return {"cancelled": cancelled, "sent": None, "executed": None}
- except EdgeAuthenticationStopped:
- self._stopped = True
- self.queue.fail(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- error_code="credential_stopped",
- )
- raise
- except Exception as exc:
- self.queue.fail(
- contract.task_id,
- lease_token=str(task_row.lease_token),
- error_code="runner_execution_failed",
- retryable=False,
- )
- raise EdgeAgentError("edge runner execution failed") from exc
- sent = self._drain_event()
- return {"cancelled": cancelled, "sent": sent, "executed": contract.task_id}
- def safe_diagnostic(self) -> dict[str, object]:
- connection = sqlite3.connect(f"file:{self.queue.db_path}?mode=ro", uri=True)
- try:
- pending = connection.execute(
- "SELECT COUNT(*) FROM edge_tasks WHERE status IN ('pending','leased')"
- ).fetchone()[0]
- events = connection.execute(
- "SELECT COUNT(*) FROM edge_outbound_events WHERE status != 'acknowledged'"
- ).fetchone()[0]
- finally:
- connection.close()
- return {
- "agent_status": "stopped" if self._stopped else "ready",
- "gateway_id_digest": self.gateway_id_digest,
- "queue_pending_count": int(pending),
- "queue_unacknowledged_count": int(events),
- "version": self.config.version,
- }
|