| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186 |
- """Targeted durable edge-agent tests without enterprise data access."""
- from __future__ import annotations
- from datetime import UTC, datetime, timedelta
- import pytest
- from cryptography import x509
- from cryptography.hazmat.primitives import hashes, serialization
- from cryptography.hazmat.primitives.asymmetric import rsa
- from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
- from cryptography.x509.oid import NameOID
- from app.core.edge_gateway.contracts import (
- SignedTaskEnvelope,
- canonical_json_bytes,
- canonical_sha256,
- canonical_timestamp,
- )
- from app.core.edge_gateway.policy import EdgePolicyError
- from app.core.edge_gateway.queue import EdgeQueueConflictError
- from app.edge_gateway.agent import (
- EdgeAgent,
- EdgeReleaseError,
- EdgeRunnerAdapter,
- LocalArtifactStore,
- SignedReleaseManager,
- )
- from app.edge_gateway.bootstrap import EdgeBootstrapConfig
- from app.edge_gateway.transport import EdgeTransportError
- TASK_SIGNING_PRIVATE_KEY = Ed25519PrivateKey.generate()
- TASK_SIGNING_PUBLIC_KEY = TASK_SIGNING_PRIVATE_KEY.public_key().public_bytes(
- encoding=serialization.Encoding.Raw,
- format=serialization.PublicFormat.Raw,
- ).hex()
- def config(tmp_path, **changes):
- cert_path = tmp_path / "edge-cert.pem"
- key_path = tmp_path / "edge-key.pem"
- ca_path = tmp_path / "ca.pem"
- artifact_root = tmp_path / "artifacts"
- artifact_root.mkdir(exist_ok=True)
- if not cert_path.exists():
- key = rsa.generate_private_key(public_exponent=65537, key_size=2048)
- name = x509.Name([x509.NameAttribute(NameOID.COMMON_NAME, "gateway-1")])
- now = datetime.now(UTC)
- cert = (
- x509.CertificateBuilder()
- .subject_name(name)
- .issuer_name(name)
- .public_key(key.public_key())
- .serial_number(x509.random_serial_number())
- .not_valid_before(now - timedelta(minutes=1))
- .not_valid_after(now + timedelta(days=1))
- .sign(key, hashes.SHA256())
- )
- cert_pem = cert.public_bytes(serialization.Encoding.PEM)
- cert_path.write_bytes(cert_pem)
- ca_path.write_bytes(cert_pem)
- key_path.write_bytes(
- key.private_bytes(
- serialization.Encoding.PEM,
- serialization.PrivateFormat.PKCS8,
- serialization.NoEncryption(),
- )
- )
- key_path.chmod(0o600)
- fingerprint = x509.load_pem_x509_certificate(cert_path.read_bytes()).fingerprint(
- hashes.SHA256()
- ).hex()
- value = {
- "queue_path": str(tmp_path / "edge.sqlite3"),
- "artifact_root": str(artifact_root),
- "gateway_id": "gateway-1",
- "credential": "dopg_local-test-credential",
- "certificate_sha256": fingerprint,
- "client_certificate_path": str(cert_path),
- "client_private_key_path": str(key_path),
- "ca_bundle_path": str(ca_path),
- "generation": 1,
- "environment": "production",
- "network_zone": "zone-a",
- "policy_digest": "a" * 64,
- "control_url": "https://control.enterprise.test",
- "proxy_url": None,
- "allowed_control_hosts": ["control.enterprise.test"],
- "allowed_proxy_hosts": [],
- "trusted_release_keys": {"release-key-1": "01" * 32},
- "trusted_task_keys": {"task-key-1": TASK_SIGNING_PUBLIC_KEY},
- "task_authority_clock_skew_seconds": 60,
- "version": "3.0.0",
- }
- value.update(changes)
- return EdgeBootstrapConfig.from_mapping(value)
- @pytest.mark.parametrize(
- "changes",
- [
- {"queue_path": ":memory:"},
- {"credential": "changeme"},
- {"trusted_release_keys": {}},
- {"control_url": "http://control.enterprise.test"},
- {"gateway_id": "TBD"},
- ],
- )
- def test_bootstrap_fails_closed_for_missing_placeholder_or_memory_configuration(tmp_path, changes):
- with pytest.raises(ValueError):
- config(tmp_path, **changes)
- def _task():
- return {
- "task_id": "task-1",
- "gateway_id": "gateway-1",
- "environment": "production",
- "network_zone": "zone-a",
- "purpose": "governed-inventory",
- "classification": "statistics",
- "task_type": "profile",
- "contract_version": 1,
- "deadline_at": canonical_timestamp(
- (datetime.now(UTC) + timedelta(minutes=5)).isoformat().replace("+00:00", "Z")
- ),
- "attempt": 1,
- "idempotency_key": "idem-1",
- "policy_digest": "a" * 64,
- }
- def _signed_task_envelope(
- task=None,
- *,
- key_id="task-key-1",
- private_key=TASK_SIGNING_PRIVATE_KEY,
- issued_at=None,
- ):
- task = task or _task()
- issued_at = issued_at or canonical_timestamp(
- (datetime.now(UTC) - timedelta(seconds=1)).isoformat().replace("+00:00", "Z")
- )
- unsigned = {
- "task": task,
- "authority_key_id": key_id,
- "signature_algorithm": "Ed25519",
- "contract_digest": canonical_sha256(task),
- "gateway_id": task["gateway_id"],
- "environment": task["environment"],
- "network_zone": task["network_zone"],
- "policy_digest": task["policy_digest"],
- "purpose": task["purpose"],
- "issued_at": issued_at,
- "expires_at": task["deadline_at"],
- }
- return {
- **unsigned,
- "signature": private_key.sign(
- SignedTaskEnvelope.canonical_unsigned_bytes(unsigned)
- ).hex(),
- }
- def test_agent_applies_configured_task_authority_future_clock_skew(tmp_path):
- now = datetime(2026, 8, 9, 8, 0, tzinfo=UTC)
- envelope = _signed_task_envelope(
- task={
- **_task(),
- "deadline_at": "2026-08-09T08:05:00Z",
- },
- issued_at="2026-08-09T08:01:00Z",
- )
- agent = EdgeAgent(
- config(tmp_path),
- SignedTransport(),
- EdgeRunnerAdapter(
- lambda _node, _request, _cancel: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- clock=lambda: now,
- )
- accepted = agent.accept_pulled_task(
- envelope,
- "dopl_clock-skew-lease",
- "2026-08-09T08:02:00Z",
- )
- assert accepted.task_id == "task-1"
- class SignedTransport:
- def __init__(self, *, remote_expiry=None):
- self.events = []
- self.cancelled = []
- self.heartbeats = []
- self.tasks = [(
- _signed_task_envelope(),
- "dopl_remote-lease-1",
- remote_expiry or canonical_timestamp(
- (datetime.now(UTC) + timedelta(minutes=2)).isoformat().replace("+00:00", "Z")
- ),
- )]
- def reconcile(self):
- return {"cancelled_task_ids": list(self.cancelled), "release_offers": []}
- def pull_task(self):
- return self.tasks.pop(0) if self.tasks else None
- def heartbeat(self, summary):
- self.heartbeats.append(dict(summary))
- return {"status": "online"}
- def send_event(self, event, lease_token):
- self.events.append((dict(event), lease_token))
- return {
- "event_id": event["event_id"],
- "event_digest": canonical_sha256(event),
- "remote_lease_digest": __import__("hashlib").sha256(
- lease_token.encode("utf-8")
- ).hexdigest(),
- "received_at": canonical_timestamp(datetime.now(UTC).isoformat().replace("+00:00", "Z")),
- "status": "accepted",
- }
- def test_agent_requires_valid_signed_task_authority_before_atomic_accept(tmp_path):
- transport = SignedTransport()
- envelope, token, expiry = transport.tasks[0]
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _n, _r, _c: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- )
- accepted = agent.accept_pulled_task(envelope, token, expiry)
- assert accepted.authority_key_id == "task-key-1"
- tampered = {**envelope, "task": {**envelope["task"], "purpose": "controlled-query"}}
- with pytest.raises(EdgePolicyError):
- agent.accept_pulled_task(tampered, token, expiry)
- unknown = _signed_task_envelope(key_id="retired-task-key")
- with pytest.raises(EdgePolicyError):
- agent.accept_pulled_task(unknown, token, expiry)
- def test_second_agent_recovers_encrypted_remote_lease_after_restart(tmp_path):
- transport = SignedTransport()
- first = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(lambda _n, _r, _c: {"classification": "statistics", "payload": {"metric_count": 1}}),
- worker_id="first",
- )
- first.accept_pulled_task(*transport.tasks.pop(0))
- executions = []
- restarted = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _n, request, _c: executions.append(request["task_id"])
- or {"classification": "statistics", "payload": {"metric_count": 1}}
- ),
- worker_id="restarted",
- )
- result = restarted.run_once()
- assert result["executed"] == "task-1"
- assert executions == ["task-1"]
- assert transport.events[0][1] == "dopl_remote-lease-1"
- assert restarted.queue.get_event(transport.events[0][0]["event_id"]).status == "acknowledged"
- def test_second_agent_reclaims_crashed_local_attempt_with_new_remote_lease(tmp_path):
- clock = MutableClock()
- first_expiry = canonical_timestamp(
- (clock.value + timedelta(seconds=20)).isoformat().replace("+00:00", "Z")
- )
- transport = SignedTransport(remote_expiry=first_expiry)
- envelope = transport.tasks[0][0]
- first = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _node, _request, _cancel: (_ for _ in ()).throw(
- AssertionError("crashed agent must not execute")
- )
- ),
- worker_id="crashed-agent",
- clock=clock,
- )
- first.accept_pulled_task(*transport.tasks.pop(0))
- crashed_claim = first.queue.claim("crashed-agent")
- assert crashed_claim.attempt_count == 1
- clock.advance(61)
- transport.tasks.append(
- (
- envelope,
- "dopl_remote-lease-2",
- canonical_timestamp(
- (clock.value + timedelta(minutes=2))
- .isoformat()
- .replace("+00:00", "Z")
- ),
- )
- )
- executions = []
- replacement = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _node, request, _cancel: executions.append(request["task_id"])
- or {"classification": "statistics", "payload": {"metric_count": 1}}
- ),
- worker_id="replacement-agent",
- clock=clock,
- )
- result = replacement.run_once()
- assert result["executed"] == "task-1"
- assert executions == ["task-1"]
- assert transport.events == [(transport.events[0][0], "dopl_remote-lease-2")]
- assert replacement.queue.get_task("task-1").attempt_count == 2
- with __import__("sqlite3").connect(replacement.queue.db_path) as connection:
- assert connection.execute(
- "SELECT count(*) FROM edge_outbound_events"
- ).fetchone()[0] == 1
- class LostAckThenRenewedLeaseTransport(SignedTransport):
- def __init__(self, *, remote_expiry, clock):
- super().__init__(remote_expiry=remote_expiry)
- self.clock = clock
- self.send_tokens = []
- self.lost = False
- def send_event(self, event, lease_token):
- self.send_tokens.append(lease_token)
- if not self.lost:
- self.lost = True
- self.events.append((dict(event), lease_token))
- raise EdgeTransportError("event acknowledgement was lost")
- if lease_token != "dopl_remote-lease-2":
- raise AssertionError("expired remote lease must not be retried")
- return super().send_event(event, lease_token)
- def test_ack_loss_rebinds_completed_event_after_remote_lease_expiry(tmp_path):
- clock = MutableClock()
- expiry = canonical_timestamp(
- (clock.value + timedelta(seconds=1)).isoformat().replace("+00:00", "Z")
- )
- transport = LostAckThenRenewedLeaseTransport(
- remote_expiry=expiry,
- clock=clock,
- )
- envelope = transport.tasks[0][0]
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _node, _request, _cancel: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- clock=clock,
- random=lambda: 0.0,
- )
- first = agent.run_once()
- assert first["executed"] == "task-1"
- event_id = transport.events[0][0]["event_id"]
- clock.advance(2)
- transport.tasks.append(
- (
- envelope,
- "dopl_remote-lease-2",
- canonical_timestamp(
- (clock.value + timedelta(minutes=2))
- .isoformat()
- .replace("+00:00", "Z")
- ),
- )
- )
- agent.run_once()
- clock.advance(5)
- agent.run_once()
- assert transport.send_tokens == ["dopl_remote-lease-1", "dopl_remote-lease-2"]
- assert agent.queue.get_event(event_id).status == "acknowledged"
- def test_expired_remote_lease_is_reacquired_before_first_execution(tmp_path):
- clock = MutableClock()
- expiry = canonical_timestamp(
- (clock.value + timedelta(seconds=1)).isoformat().replace("+00:00", "Z")
- )
- transport = SignedTransport(remote_expiry=expiry)
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _n, _r, _c: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- clock=clock,
- )
- envelope, token, remote_expiry = transport.tasks.pop(0)
- agent.accept_pulled_task(envelope, token, remote_expiry)
- clock.advance(2)
- transport.tasks.append((
- envelope,
- "dopl_remote-lease-2",
- canonical_timestamp(
- (clock.value + timedelta(minutes=2)).isoformat().replace("+00:00", "Z")
- ),
- ))
- result = agent.run_once()
- assert result["executed"] == "task-1"
- assert transport.events[0][1] == "dopl_remote-lease-2"
- class WrongAckTransport(SignedTransport):
- def send_event(self, event, lease_token):
- self.events.append((dict(event), lease_token))
- return {
- "event_id": "evt_wrong",
- "event_digest": "f" * 64,
- "remote_lease_digest": __import__("hashlib").sha256(
- lease_token.encode("utf-8")
- ).hexdigest(),
- "received_at": canonical_timestamp(
- datetime.now(UTC).isoformat().replace("+00:00", "Z")
- ),
- "status": "accepted",
- }
- def test_wrong_event_or_digest_ack_never_marks_event_acknowledged(tmp_path):
- transport = WrongAckTransport()
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _n, _r, _c: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- )
- with pytest.raises(EdgePolicyError):
- agent.run_once()
- event_id = transport.events[0][0]["event_id"]
- assert agent.queue.get_event(event_id).status != "acknowledged"
- class FakeTransport:
- def __init__(self):
- self.tasks = [(
- _signed_task_envelope(),
- "dopl_lease-1",
- canonical_timestamp(
- (datetime.now(UTC) + timedelta(minutes=2)).isoformat().replace("+00:00", "Z")
- ),
- )]
- self.events = []
- self.outcomes = []
- self.cancelled = []
- self.heartbeats = []
- def reconcile(self):
- return {"cancelled_task_ids": list(self.cancelled), "release_offers": []}
- def pull_task(self):
- return self.tasks.pop(0) if self.tasks else None
- def send_event(self, event, lease_token):
- self.events.append((dict(event), lease_token))
- return {
- "event_id": event["event_id"],
- "event_digest": canonical_sha256(event),
- "remote_lease_digest": __import__("hashlib").sha256(
- lease_token.encode("utf-8")
- ).hexdigest(),
- "received_at": canonical_timestamp(
- datetime.now(UTC).isoformat().replace("+00:00", "Z")
- ),
- "status": "accepted",
- }
- def task_outcome(self, task_id, outcome, lease_token, summary):
- self.outcomes.append((task_id, outcome, lease_token, dict(summary)))
- return {"task_id": task_id, "status": outcome, "replayed": False}
- def heartbeat(self, summary):
- self.heartbeats.append(dict(summary))
- return {"status": "online"}
- def test_agent_persists_before_execute_and_send_and_reconnects_exactly_once(tmp_path):
- transport = FakeTransport()
- observations = []
- def execute(_node, request, _cancel):
- observations.append((request["task_id"], agent.queue.get_task(request["task_id"]).status))
- return {"classification": "statistics", "payload": {"metric_count": 3}}
- agent = EdgeAgent(config(tmp_path), transport, EdgeRunnerAdapter(execute))
- first = agent.run_once()
- assert first["executed"] == "task-1"
- assert observations == [("task-1", "leased")]
- assert transport.events and agent.queue.get_event(transport.events[0][0]["event_id"]).status == "acknowledged"
- second = agent.run_once()
- assert second["executed"] is None
- assert len(transport.events) == 1
- assert agent.safe_diagnostic() == {
- "agent_status": "ready",
- "gateway_id_digest": agent.gateway_id_digest,
- "queue_pending_count": 0,
- "queue_unacknowledged_count": 0,
- "version": "3.0.0",
- }
- def test_two_agents_do_not_execute_same_sqlite_task(tmp_path):
- transport = FakeTransport()
- count = 0
- def execute(_node, _request, _cancel):
- nonlocal count
- count += 1
- return {"classification": "statistics", "payload": {"metric_count": 1}}
- first = EdgeAgent(config(tmp_path), transport, EdgeRunnerAdapter(execute), worker_id="one")
- second = EdgeAgent(config(tmp_path), transport, EdgeRunnerAdapter(execute), worker_id="two")
- first.run_once()
- second.run_once()
- assert count == 1
- def test_offline_cancel_converges_before_execution(tmp_path):
- transport = FakeTransport()
- called = False
- def execute(_node, _request, _cancel):
- nonlocal called
- called = True
- return {"classification": "statistics", "payload": {"metric_count": 1}}
- agent = EdgeAgent(config(tmp_path), transport, EdgeRunnerAdapter(execute))
- agent.accept_pulled_task(*transport.pull_task())
- transport.cancelled = ["task-1"]
- result = agent.run_once()
- assert result["cancelled"] == ["task-1"]
- assert not called
- assert agent.queue.get_task("task-1").status == "cancelled"
- assert transport.outcomes == [
- ("task-1", "cancelled", "dopl_lease-1", {"task_status": "cancelled"})
- ]
- assert agent.queue.get_outcome("task-1").status == "acknowledged"
- class CancelOutcomeAckLossTransport(FakeTransport):
- def __init__(self):
- super().__init__()
- self.outcome_attempts = 0
- def task_outcome(self, task_id, outcome, lease_token, summary):
- self.outcome_attempts += 1
- self.outcomes.append((task_id, outcome, lease_token, dict(summary)))
- if self.outcome_attempts == 1:
- raise EdgeTransportError("cancel acknowledgement lost")
- return {"task_id": task_id, "status": outcome, "replayed": True}
- def test_cancel_outcome_survives_ack_loss_disconnect_restart_and_replays_exactly(tmp_path):
- clock = MutableClock()
- transport = CancelOutcomeAckLossTransport()
- first = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _node, _request, _cancel: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- worker_id="first",
- clock=clock,
- )
- first.accept_pulled_task(*transport.pull_task())
- transport.cancelled = ["task-1"]
- assert first.run_once()["cancelled"] == ["task-1"]
- pending = first.queue.get_outcome("task-1")
- assert pending.status == "pending"
- stable_digest = pending.digest
- clock.advance(3)
- restarted = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda *_args: (_ for _ in ()).throw(
- AssertionError("cancelled task must never execute")
- )
- ),
- worker_id="restarted",
- clock=clock,
- )
- restarted.run_once()
- acknowledged = restarted.queue.get_outcome("task-1")
- assert acknowledged.status == "acknowledged"
- assert acknowledged.digest == stable_digest
- assert [attempt[:2] for attempt in transport.outcomes] == [
- ("task-1", "cancelled"),
- ("task-1", "cancelled"),
- ]
- def test_cancel_outcome_sqlite_claim_is_single_sender_across_agents(tmp_path):
- transport = FakeTransport()
- first = EdgeAgent(
- config(tmp_path), transport, EdgeRunnerAdapter(lambda *_args: {}), worker_id="one"
- )
- first.accept_pulled_task(*transport.pull_task())
- first.queue.cancel_with_outcome("task-1")
- second = EdgeAgent(
- config(tmp_path), transport, EdgeRunnerAdapter(lambda *_args: {}), worker_id="two"
- )
- claimed = first.queue.claim_outcome("one")
- assert claimed is not None
- assert second.queue.claim_outcome("two") is None
- class MutableClock:
- def __init__(self):
- self.value = datetime.now(UTC)
- def __call__(self):
- return self.value
- def advance(self, seconds):
- self.value += timedelta(seconds=seconds)
- class AckLossTransport(FakeTransport):
- def __init__(self):
- super().__init__()
- self.attempts = 0
- def send_event(self, event, lease_token):
- self.attempts += 1
- self.events.append((dict(event), lease_token))
- if self.attempts == 1:
- raise EdgeTransportError("ack lost")
- return {
- "event_id": event["event_id"],
- "event_digest": canonical_sha256(event),
- "remote_lease_digest": __import__("hashlib").sha256(
- lease_token.encode("utf-8")
- ).hexdigest(),
- "received_at": canonical_timestamp(
- datetime.now(UTC).isoformat().replace("+00:00", "Z")
- ),
- "status": "accepted",
- }
- def test_ack_loss_keeps_oldest_event_durable_and_replays_same_identity(tmp_path):
- clock = MutableClock()
- transport = AckLossTransport()
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _node, _request, _cancel: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- clock=clock,
- )
- first = agent.run_once()
- assert first["executed"] == "task-1" and first["sent"] is None
- event_id = transport.events[0][0]["event_id"]
- assert agent.queue.get_event(event_id).status == "pending"
- # Backoff has not elapsed: unresolved oldest event blocks all new work.
- task_two = {**_task(), "task_id": "task-2", "idempotency_key": "idem-2"}
- transport.tasks.append((
- _signed_task_envelope(task_two),
- "dopl_lease-2",
- canonical_timestamp(
- (clock.value + timedelta(minutes=2)).isoformat().replace("+00:00", "Z")
- ),
- ))
- assert agent.run_once()["executed"] is None
- assert len(transport.tasks) == 1
- clock.advance(3)
- assert agent.run_once()["sent"] == event_id
- assert [event[0]["event_id"] for event in transport.events] == [event_id, event_id]
- assert agent.queue.get_event(event_id).status == "acknowledged"
- def test_changed_event_digest_is_rejected_without_replacing_durable_event(tmp_path):
- transport = FakeTransport()
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _node, _request, _cancel: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- )
- agent.run_once()
- original = transport.events[0][0]
- changed = {**original, "payload": {"metric_count": 2}}
- with pytest.raises(EdgeQueueConflictError):
- agent.queue.persist_event(changed)
- assert agent.queue.get_event(original["event_id"]).event["payload"]["metric_count"] == 1
- class CancelRaceTransport(FakeTransport):
- def __init__(self):
- super().__init__()
- self.probes = 0
- def cancel_requested(self, _task_id):
- self.probes += 1
- return self.probes >= 2
- def test_cooperative_cancel_race_never_emits_success_event(tmp_path):
- transport = CancelRaceTransport()
- called = 0
- def execute(_node, _request, cancel):
- nonlocal called
- called += 1
- assert cancel()
- return {"classification": "statistics", "payload": {"metric_count": 1}}
- agent = EdgeAgent(config(tmp_path), transport, EdgeRunnerAdapter(execute))
- result = agent.run_once()
- assert result["executed"] is None
- assert result["cancelled"] == ["task-1"]
- assert called == 1 and transport.events == []
- assert agent.queue.get_task("task-1").status == "cancelled"
- def _signed_manifest(private_key, **changes):
- unsigned = {
- "release_id": "release-1",
- "version": "3.1.0",
- "rollback_version": "3.0.0",
- "artifact_digest": "0" * 64,
- "artifact_name": "edge-agent-3.1.0.bin",
- "deadline_at": canonical_timestamp(
- (datetime.now(UTC) + timedelta(minutes=5))
- .isoformat()
- .replace("+00:00", "Z")
- ),
- "signature_algorithm": "Ed25519",
- "key_id": "release-key-1",
- "status": "offered",
- }
- unsigned.update(changes)
- return {
- **unsigned,
- "manifest_digest": canonical_sha256(unsigned),
- "signature": private_key.sign(canonical_json_bytes(unsigned)).hex(),
- }
- def _release_manager(private_key, artifact, activator, rollback):
- public = private_key.public_key().public_bytes(
- encoding=serialization.Encoding.Raw,
- format=serialization.PublicFormat.Raw,
- ).hex()
- return SignedReleaseManager(
- current_version="3.0.0",
- trusted_keys={"release-key-1": public},
- artifact_loader=lambda _manifest: artifact,
- activator=activator,
- rollback=rollback,
- )
- def test_signed_candidate_activates_only_after_digest_and_health_confirmation():
- private = Ed25519PrivateKey.generate()
- artifact = b"verified-edge-agent"
- activated = []
- manager = _release_manager(
- private,
- artifact,
- lambda version, body: activated.append((version, body)) or True,
- lambda _version: False,
- )
- manifest = _signed_manifest(private, artifact_digest=__import__("hashlib").sha256(artifact).hexdigest())
- result = manager.install(manifest)
- assert result["outcome"] == "installed"
- assert activated == [("3.1.0", artifact)]
- assert manager.current_version == "3.1.0"
- def test_candidate_health_failure_rolls_back_previous_verified_version():
- private = Ed25519PrivateKey.generate()
- artifact = b"verified-edge-agent"
- rolled_back = []
- manager = _release_manager(
- private,
- artifact,
- lambda _version, _body: False,
- lambda version: rolled_back.append(version) or True,
- )
- result = manager.install(
- _signed_manifest(private, artifact_digest=__import__("hashlib").sha256(artifact).hexdigest())
- )
- assert result["outcome"] == "rollback"
- assert rolled_back == ["3.0.0"]
- assert manager.current_version == "3.0.0"
- @pytest.mark.parametrize(
- "change",
- [
- {"artifact_name": "../agent.bin"},
- {"artifact_name": "/tmp/agent.bin"},
- {"version": "2.9.0"},
- {"rollback_version": "2.8.0"},
- {"signature_algorithm": "RSA"},
- ],
- )
- def test_release_manifest_rejects_path_traversal_downgrade_and_unapproved_algorithm(change):
- private = Ed25519PrivateKey.generate()
- manager = _release_manager(private, b"body", lambda _v, _b: True, lambda _v: True)
- with pytest.raises(EdgeReleaseError):
- manager.verify(_signed_manifest(private, **change))
- def test_release_id_changed_digest_replay_fails_closed():
- private = Ed25519PrivateKey.generate()
- manager = _release_manager(private, b"body", lambda _v, _b: True, lambda _v: True)
- first = _signed_manifest(private)
- manager.verify(first)
- with pytest.raises(EdgeReleaseError, match="replay"):
- manager.verify(_signed_manifest(private, artifact_digest="1" * 64))
- def test_safe_diagnostic_is_approved_and_contains_no_runtime_secrets(tmp_path):
- agent = EdgeAgent(
- config(tmp_path),
- FakeTransport(),
- EdgeRunnerAdapter(lambda _n, _r, _c: {"classification": "statistics", "payload": {"metric_count": 1}}),
- )
- diagnostic = agent.safe_diagnostic()
- agent.policy.validate_approved_payload("diagnostic_summary", diagnostic)
- encoded = str(diagnostic)
- assert agent.config.credential not in encoded
- assert agent.config.certificate_sha256 not in encoded
- assert agent.heartbeat_once() == {"status": "online"}
- assert agent.transport.heartbeats == [diagnostic]
- def test_local_raw_artifact_reference_never_enters_outbound_event(tmp_path):
- transport = FakeTransport()
- envelope, lease, expiry = transport.tasks[0]
- raw_task = {**envelope["task"], "classification": "raw"}
- transport.tasks[0] = (_signed_task_envelope(raw_task), lease, expiry)
- body = b"local raw rows never leave the edge"
- digest = __import__("hashlib").sha256(body).hexdigest()
- artifact_root = tmp_path / "artifacts"
- artifact_root.mkdir()
- artifact_path = artifact_root / digest
- artifact_path.write_bytes(body)
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _n, _r, _c: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- "local_artifact_digest": digest,
- "local_artifact_ref": str(artifact_path),
- }
- ),
- )
- agent.run_once()
- event = transport.events[0][0]
- assert event["payload"] == {"metric_count": 1}
- assert "artifact" not in str(event).lower()
- assert not artifact_path.exists()
- artifact = agent.queue.list_local_artifacts(task_id="task-1")[0]
- assert artifact.artifact_ref is None
- assert artifact.artifact_ref_hash == canonical_sha256(str(artifact_path.resolve()))
- assert artifact.cleanup_status == "deleted"
- assert artifact.deletion_receipt["status"] == "deleted"
- @pytest.mark.parametrize("case", ["escape", "symlink", "digest_mismatch"])
- def test_local_artifact_store_rejects_escape_symlink_and_digest_mismatch(
- tmp_path, case
- ):
- root = tmp_path / "artifacts"
- root.mkdir()
- outside = tmp_path / "outside.bin"
- outside.write_bytes(b"enterprise raw data")
- expected = __import__("hashlib").sha256(outside.read_bytes()).hexdigest()
- reference = outside
- digest = expected
- if case == "symlink":
- reference = root / "linked.bin"
- reference.symlink_to(outside)
- elif case == "digest_mismatch":
- reference = root / "artifact.bin"
- reference.write_bytes(outside.read_bytes())
- digest = "0" * 64
- store = LocalArtifactStore(root)
- with pytest.raises(EdgePolicyError):
- store.inspect(str(reference), digest)
- assert outside.exists()
- class SimulatedProcessCrash(BaseException):
- pass
- class CrashAfterArtifactDeleteStore(LocalArtifactStore):
- def delete(self, reference, expected_digest):
- super().delete(reference, expected_digest)
- raise SimulatedProcessCrash("simulated crash after unlink before sqlite receipt")
- def test_artifact_cleanup_recovers_when_crash_happens_after_unlink_before_receipt(
- tmp_path,
- ):
- clock = MutableClock()
- transport = FakeTransport()
- envelope, lease, expiry = transport.tasks[0]
- raw_task = {**envelope["task"], "classification": "raw"}
- transport.tasks[0] = (_signed_task_envelope(raw_task), lease, expiry)
- root = tmp_path / "artifacts"
- root.mkdir()
- body = b"raw artifact crash recovery"
- digest = __import__("hashlib").sha256(body).hexdigest()
- artifact_path = root / digest
- artifact_path.write_bytes(body)
- first = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda *_args: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- "local_artifact_digest": digest,
- "local_artifact_ref": str(artifact_path),
- }
- ),
- clock=clock,
- artifact_store=CrashAfterArtifactDeleteStore(root, clock=clock),
- )
- with pytest.raises(SimulatedProcessCrash):
- first.run_once()
- assert not artifact_path.exists()
- assert first.queue.list_local_artifacts(task_id="task-1")[0].cleanup_status == "deleting"
- clock.advance(61)
- restarted = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(lambda *_args: {}),
- clock=clock,
- )
- restarted.run_once()
- recovered = restarted.queue.list_local_artifacts(task_id="task-1")[0]
- assert recovered.cleanup_status == "deleted"
- assert recovered.artifact_ref is None
- def test_artifact_cleanup_sqlite_claim_is_single_deleter_across_agents(tmp_path):
- clock = MutableClock()
- transport = FakeTransport()
- root = tmp_path / "artifacts"
- root.mkdir()
- body = b"retained metadata artifact"
- digest = __import__("hashlib").sha256(body).hexdigest()
- path = root / digest
- path.write_bytes(body)
- first = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda *_args: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- "local_artifact_digest": digest,
- "local_artifact_ref": str(path),
- }
- ),
- worker_id="one",
- clock=clock,
- )
- first.run_once()
- assert path.exists()
- clock.advance(1_096 * 86_400)
- second = EdgeAgent(
- config(tmp_path), transport, EdgeRunnerAdapter(lambda *_args: {}),
- worker_id="two", clock=clock,
- )
- claimed = first.queue.claim_artifact_cleanup("one")
- assert claimed is not None
- assert second.queue.claim_artifact_cleanup("two") is None
- class ReleaseTransport(FakeTransport):
- def __init__(self, manifest):
- super().__init__()
- self.tasks = []
- self.manifest = manifest
- self.acks = []
- def reconcile(self):
- manifest, self.manifest = self.manifest, None
- return {
- "cancelled_task_ids": [],
- "release_offers": [manifest] if manifest else [],
- }
- def acknowledge_release(self, release_id, outcome, summary):
- self.acks.append((release_id, outcome, dict(summary)))
- return {"release_id": release_id, "status": outcome}
- def test_agent_release_state_sequence_matches_control_plane_transitions(tmp_path):
- private = Ed25519PrivateKey.generate()
- artifact = b"verified-edge-agent"
- manifest = _signed_manifest(
- private, artifact_digest=__import__("hashlib").sha256(artifact).hexdigest()
- )
- manager = _release_manager(
- private,
- artifact,
- lambda _version, _body: False,
- lambda _version: True,
- )
- transport = ReleaseTransport(manifest)
- agent = EdgeAgent(
- config(tmp_path),
- transport,
- EdgeRunnerAdapter(
- lambda _n, _r, _c: {
- "classification": "statistics",
- "payload": {"metric_count": 1},
- }
- ),
- release_manager=manager,
- )
- agent.run_once()
- assert [ack[1] for ack in transport.acks] == ["accepted", "failed", "rollback"]
- def _manifest_state(private_key, manifest, status):
- unsigned = {
- key: value
- for key, value in manifest.items()
- if key not in {"manifest_digest", "signature", "signature_algorithm", "key_id"}
- }
- unsigned["status"] = status
- unsigned["signature_algorithm"] = "Ed25519"
- unsigned["key_id"] = "release-key-1"
- return {
- **unsigned,
- "manifest_digest": canonical_sha256(unsigned),
- "signature": private_key.sign(canonical_json_bytes(unsigned)).hex(),
- }
- class StatefulReleaseTransport(FakeTransport):
- def __init__(self, private_key, manifest, *, fail_once=None):
- super().__init__()
- self.tasks = []
- self.private_key = private_key
- self.manifest = manifest
- self.fail_once = fail_once
- self.acks = []
- def reconcile(self):
- return {"cancelled_task_ids": [], "release_offers": [self.manifest]}
- def acknowledge_release(self, release_id, outcome, summary):
- self.acks.append((release_id, outcome, dict(summary)))
- if self.fail_once == outcome:
- self.fail_once = None
- raise EdgeTransportError("release acknowledgement lost")
- status = "rolled_back" if outcome == "rollback" else outcome
- self.manifest = _manifest_state(self.private_key, self.manifest, status)
- return {"release_id": release_id, "status": status}
- @pytest.mark.parametrize("lost_ack", ["accepted", "installed"])
- def test_release_ack_loss_and_agent_restart_converge_without_reactivation(tmp_path, lost_ack):
- private = Ed25519PrivateKey.generate()
- artifact = b"verified-edge-agent"
- manifest = _signed_manifest(
- private, artifact_digest=__import__("hashlib").sha256(artifact).hexdigest()
- )
- activations = []
- transport = StatefulReleaseTransport(private, manifest, fail_once=lost_ack)
- manager = _release_manager(
- private,
- artifact,
- lambda version, _body: activations.append(version) or True,
- lambda _version: True,
- )
- first = EdgeAgent(
- config(tmp_path), transport,
- EdgeRunnerAdapter(lambda _n, _r, _c: {"classification": "statistics", "payload": {"metric_count": 1}}),
- release_manager=manager,
- )
- with pytest.raises(EdgeTransportError):
- first.run_once()
- restarted_activations = []
- restarted_manager = _release_manager(
- private,
- artifact,
- lambda version, _body: restarted_activations.append(version) or True,
- lambda _version: True,
- )
- restarted = EdgeAgent(
- config(tmp_path), transport,
- EdgeRunnerAdapter(lambda _n, _r, _c: {"classification": "statistics", "payload": {"metric_count": 1}}),
- release_manager=restarted_manager,
- )
- restarted.run_once()
- assert transport.manifest["status"] == "installed"
- assert activations + restarted_activations == ["3.1.0"]
- assert restarted_activations == (["3.1.0"] if lost_ack == "accepted" else [])
- assert restarted.queue.get_release_state("release-1").status == "installed"
- def test_failed_ack_loss_restart_resends_state_without_duplicate_rollback(tmp_path):
- private = Ed25519PrivateKey.generate()
- artifact = b"verified-edge-agent"
- manifest = _signed_manifest(
- private, artifact_digest=__import__("hashlib").sha256(artifact).hexdigest()
- )
- rollbacks = []
- transport = StatefulReleaseTransport(private, manifest, fail_once="failed")
- manager = _release_manager(
- private,
- artifact,
- lambda _version, _body: False,
- lambda version: rollbacks.append(version) or True,
- )
- first = EdgeAgent(
- config(tmp_path), transport,
- EdgeRunnerAdapter(lambda _n, _r, _c: {"classification": "statistics", "payload": {"metric_count": 1}}),
- release_manager=manager,
- )
- with pytest.raises(EdgeTransportError):
- first.run_once()
- assert first.queue.get_release_state("release-1").status == "rolled_back"
- restarted = EdgeAgent(
- config(tmp_path), transport,
- EdgeRunnerAdapter(lambda _n, _r, _c: {"classification": "statistics", "payload": {"metric_count": 1}}),
- release_manager=_release_manager(
- private, artifact,
- lambda _version, _body: False,
- lambda _version: (_ for _ in ()).throw(AssertionError("must not rollback twice")),
- ),
- )
- restarted.run_once()
- assert transport.manifest["status"] == "rolled_back"
- assert rollbacks == ["3.0.0"]
|