"""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, }