"""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"]