from __future__ import annotations import hashlib import os import subprocess import uuid from concurrent.futures import ThreadPoolExecutor from datetime import UTC, datetime, timedelta from pathlib import Path from threading import Barrier, Event from types import SimpleNamespace from urllib.parse import quote 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 sqlalchemy import create_engine, inspect, text from sqlalchemy.engine import make_url from sqlalchemy.exc import DBAPIError from sqlalchemy.orm import Session from app.config.config import validate_production_database_identity from app.core.edge_gateway.contracts import ( EdgeEventContract, SignedTaskEnvelope, canonical_json_bytes, canonical_sha256, stable_event_id, ) from app.core.edge_gateway.repository import EdgeGatewayRepository from app.core.edge_gateway.service import ( EdgeGatewayAuthenticationError, EdgeGatewayConfigurationError, EdgeGatewayConflictError, EdgeGatewayError, EdgeGatewayService, EdgeGatewayValidationError, ) from app.edge_gateway.agent import EdgeAgent, EdgeRunnerAdapter, SignedReleaseManager from app.edge_gateway.bootstrap import EdgeBootstrapConfig ROOT = Path(__file__).resolve().parents[2] DATABASE_URL = os.environ.get("TEST_DATABASE_URL") pytestmark = pytest.mark.integration _EVIDENCE_TRIGGERS = { "edge_gateway_audits": ( "trg_edge_audit_append_only", "trg_edge_audit_no_truncate", ), "edge_gateway_failure_windows": ( "trg_edge_failure_window_protected", "trg_edge_failure_window_no_truncate", ), } _SIGNING_PRIVATE_KEY = Ed25519PrivateKey.generate() _SIGNING_KEY_ID = "edge-control-test-key-1" def _service(repository): return EdgeGatewayService( repository, signing_private_key=_SIGNING_PRIVATE_KEY, signing_key_id=_SIGNING_KEY_ID, ) def _uid() -> str: return str(uuid.uuid4()) def _alembic(command: str, target: str, database_url: str = None) -> None: assert DATABASE_URL result = subprocess.run( [str(ROOT / ".venv/bin/alembic"), "-c", str(ROOT / "alembic.ini"), command, target], cwd=ROOT, env={ **os.environ, "MIGRATION_DATABASE_URL": database_url or DATABASE_URL, "SQLALCHEMY_DATABASE_URI": "", "DATABASE_URL": "", }, check=False, capture_output=True, text=True, ) if result.returncode: raise AssertionError(result.stderr) def _clear_edge_rows_for_test(engine) -> None: tables = set(inspect(engine).get_table_names(schema="public")) targets = [] for table in ( "edge_gateway_failure_windows", "edge_gateway_audits", "edge_gateway_events", "edge_gateway_releases", "edge_gateway_tasks", "edge_gateway_rotation_requests", "edge_gateway_credentials", "edge_gateways", "edge_gateway_enrollments", ): if table in tables: targets.append(f"public.{table}") if not targets: return with engine.begin() as connection: existing_triggers = set(connection.execute(text(""" SELECT tgname FROM pg_trigger WHERE tgrelid IN ( 'public.edge_gateway_audits'::regclass, to_regclass('public.edge_gateway_failure_windows') ) AND NOT tgisinternal """)).scalars()) if "edge_gateway_audits" in tables else set() disabled = [] for table, triggers in _EVIDENCE_TRIGGERS.items(): if table not in tables: continue for trigger in triggers: if trigger in existing_triggers: connection.execute(text(f"ALTER TABLE public.{table} DISABLE TRIGGER {trigger}")) disabled.append((table, trigger)) connection.execute(text(f"TRUNCATE TABLE {','.join(targets)} CASCADE")) for table, trigger in disabled: connection.execute(text(f"ALTER TABLE public.{table} ENABLE TRIGGER {trigger}")) @pytest.fixture(scope="module") def engine(wp06_postgres_identities): if not DATABASE_URL: pytest.skip("TEST_DATABASE_URL is required") bootstrap = create_engine(DATABASE_URL, pool_pre_ping=True) _clear_edge_rows_for_test(bootstrap) bootstrap.dispose() _alembic("downgrade", "20260802_476") # WP06's trusted-delivery migrations explicitly reject a superuser or # CREATEROLE migrator. Keep the older WP04 exercise at the same current # head by using the isolated no-superuser migrator identity. _alembic("upgrade", "head", wp06_postgres_identities.migration_url) value = create_engine(DATABASE_URL, pool_pre_ping=True) yield value value.dispose() @pytest.fixture(autouse=True) def clear_edge_rows(engine): _clear_edge_rows_for_test(engine) yield _clear_edge_rows_for_test(engine) @pytest.fixture def service(engine): session = Session(engine) yield _service(EdgeGatewayRepository(session)) session.close() def _enrollment(service, **overrides): values = { "gateway_name": "plant-a-edge", "environment": "staging", "network_zone": "manufacturing-zone-a", "policy_digest": "a" * 64, "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": ["proxy.example.test"], "expected_certificate_sha256": "b" * 64, "ttl_seconds": 600, "actor_uid": _uid(), } values.update(overrides) return service.create_enrollment(**values) def _register(service, enrollment, certificate_sha256="b" * 64): return service.register( enrollment_token=enrollment["enrollment_token"], certificate_sha256=certificate_sha256, gateway_id=enrollment["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", policy_digest="a" * 64, allowed_control_hosts=["control.example.test"], allowed_proxy_hosts=["proxy.example.test"], version="3.0.0", ) def _event(task, gateway, classification="statistics", payload=None): identity = { "task_id": task["task_id"], "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": classification, "contract_version": 1, "occurred_at": datetime.now(UTC).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "evt-key-1", "policy_digest": "a" * 64, "payload": payload or {"asset_count": 12}, } return {"event_id": stable_event_id(identity), **identity} def test_migration_is_479_head_and_schema_is_constrained(engine): with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479" functions = connection.execute(text(""" SELECT pg_get_function_identity_arguments(p.oid),r.rolname,p.prosecdef,p.proconfig FROM pg_proc p JOIN pg_namespace n ON n.oid=p.pronamespace JOIN pg_roles r ON r.oid=p.proowner WHERE n.nspname='public' AND p.proname='record_edge_gateway_failure' """)).all() constraints = connection.execute(text(""" SELECT conname FROM pg_constraint WHERE conrelid IN ('public.edge_gateways'::regclass, 'public.edge_gateway_credentials'::regclass) """)).scalars().all() assert "ck_edge_gateway_status" in constraints assert "ck_edge_credential_status" in constraints assert "trg_edge_gateway_binding" in constraints assert "trg_edge_credential_binding" in constraints assert functions == [( "p_audit_uid uuid, p_fingerprint character, p_gateway uuid, p_event_type character varying, p_reason_code character varying", "dataops_edge_evidence_owner", True, ["search_path=pg_catalog"], )] enrollment_columns = {column["name"] for column in inspect(engine).get_columns("edge_gateway_enrollments", schema="public")} assert "expected_certificate_sha256" in enrollment_columns task_columns = { column["name"] for column in inspect(engine).get_columns("edge_gateway_tasks", schema="public") } release_columns = { column["name"] for column in inspect(engine).get_columns("edge_gateway_releases", schema="public") } assert "authority" in task_columns assert { "artifact_name", "deadline_at", "signature_algorithm", "key_id", "manifest_digest", "signature", "signed_manifest", } <= release_columns with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(""" INSERT INTO edge_gateway_enrollments (uid,gateway_id,gateway_name,environment,network_zone,policy_digest, allowed_control_hosts,allowed_proxy_hosts,token_hash,status,expires_at,created_by) VALUES (CAST(:uid AS uuid),CAST(:gateway AS uuid),'invalid-host','staging','zone-a', :policy,ARRAY['control.example.test','CONTROL.example.test'],ARRAY[]::text[], :token,'pending',CURRENT_TIMESTAMP+INTERVAL '10 minutes',CAST(:actor AS uuid)) """), {"uid": _uid(), "gateway": _uid(), "actor": _uid(), "policy": "1" * 64, "token": "2" * 64}) def test_migration_controlled_477_downgrade_and_476_upgrade(engine): try: _alembic("downgrade", "20260802_476") with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260802_476" assert not inspect(connection).has_table("edge_gateways", schema="public") finally: _alembic("upgrade", "head") with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479" def test_479_downgrade_and_upgrade_are_reversible_without_signed_rows(engine): try: _alembic("downgrade", "20260809_478") with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_478" assert "authority" not in { column["name"] for column in inspect(connection).get_columns("edge_gateway_tasks", schema="public") } finally: _alembic("upgrade", "head") with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479" def test_478_downgrade_never_restores_unsafe_function(engine): try: _alembic("downgrade", "20260809_477") with engine.connect() as connection: signatures = connection.execute(text(""" SELECT pg_get_function_identity_arguments(p.oid) FROM pg_proc p JOIN pg_namespace n ON n.oid=p.pronamespace WHERE n.nspname='public' AND p.proname='record_edge_gateway_failure' ORDER BY 1 """)).scalars().all() assert signatures == [ "p_audit_uid uuid, p_fingerprint character, p_gateway uuid, p_event_type character varying, p_reason_code character varying" ] finally: _alembic("upgrade", "head") def test_migration_rejects_downgrade_with_task2_data(service, engine): _enrollment(service) with pytest.raises(subprocess.CalledProcessError): _alembic("downgrade", "20260802_476") with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479" _clear_edge_rows_for_test(engine) with engine.begin() as connection: connection.execute(text(""" SELECT public.record_edge_gateway_audit( CAST(:uid AS uuid),NULL,'enrollment_created',NULL,FALSE, '{"reason_code":"controlled_downgrade_probe"}'::jsonb ) """), {"uid": _uid()}) with pytest.raises(subprocess.CalledProcessError): _alembic("downgrade", "20260802_476") with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479" def test_enrollment_registration_hash_only_replay_and_failed_binding(service, engine): enrollment = _enrollment(service) assert enrollment["returned_once"] is True token = enrollment["enrollment_token"] with engine.connect() as connection: row = connection.execute(text("SELECT token_hash, consumed_at FROM edge_gateway_enrollments")).mappings().one() assert row["token_hash"] == hashlib.sha256(token.encode()).hexdigest() assert token not in str(row) with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(""" UPDATE edge_gateway_enrollments SET expected_certificate_sha256=:certificate WHERE gateway_id=CAST(:gateway AS uuid) """), {"certificate": "c" * 64, "gateway": enrollment["gateway_id"]}) with pytest.raises(EdgeGatewayAuthenticationError): _register(service, enrollment, certificate_sha256="c" * 64) with engine.connect() as connection: assert connection.execute(text("SELECT consumed_at FROM edge_gateway_enrollments")).scalar_one() is None failures = connection.execute(text("SELECT safe_detail FROM edge_gateway_audits WHERE success=false")).scalars().all() assert failures and all(token not in str(detail) for detail in failures) with pytest.raises(EdgeGatewayAuthenticationError): service.register( enrollment_token=token, certificate_sha256="b" * 64, gateway_id=enrollment["gateway_id"], environment="production", network_zone="manufacturing-zone-a", policy_digest="a" * 64, allowed_control_hosts=["control.example.test"], allowed_proxy_hosts=["proxy.example.test"], version="3.0.0", ) with engine.connect() as connection: assert connection.execute(text("SELECT consumed_at FROM edge_gateway_enrollments")).scalar_one() is None gateway = _register(service, enrollment) assert gateway["credential_returned_once"] is True with engine.connect() as connection: credential = connection.execute(text("SELECT credential_hash, certificate_sha256, generation FROM edge_gateway_credentials")).mappings().one() assert gateway["credential"] not in str(credential) assert credential["certificate_sha256"] == "b" * 64 assert credential["generation"] == 1 with pytest.raises(EdgeGatewayAuthenticationError): _register(service, enrollment) def test_registration_failures_are_windowed_rate_limited_and_do_not_consume(service, engine): enrollment = _enrollment(service) for _ in range(9): with pytest.raises(EdgeGatewayAuthenticationError): _register(service, enrollment, certificate_sha256="c" * 64) with pytest.raises(EdgeGatewayError, match="rate limit") as limited: _register(service, enrollment, certificate_sha256="c" * 64) assert limited.value.status_code == 429 with engine.connect() as connection: failure = connection.execute(text(""" SELECT occurrence_count FROM edge_gateway_failure_windows WHERE event_type='untrusted_request_rejected' """)).scalar_one() audit = connection.execute(text(""" SELECT count(*),string_agg(safe_detail::text,'') FROM edge_gateway_audits WHERE event_type='untrusted_request_rejected' """)).one() consumed_at = connection.execute(text(""" SELECT consumed_at FROM edge_gateway_enrollments WHERE gateway_id=CAST(:gateway_id AS uuid) """), {"gateway_id": enrollment["gateway_id"]}).scalar_one() assert failure == 10 assert audit[0] == 1 assert enrollment["enrollment_token"] not in (audit[1] or "") assert consumed_at is None assert _register(service, enrollment)["generation"] == 1 def test_mixed_unauthenticated_failures_share_one_bounded_unknown_bucket(service, engine): limited = None for attempt in range(10): gateway_id = _uid() try: if attempt % 2: service.register( enrollment_token=f"dope_invalid_{attempt}", certificate_sha256="b" * 64, gateway_id=gateway_id, environment="staging", network_zone=f"random-zone-{attempt}", policy_digest="a" * 64, allowed_control_hosts=["control.example.test"], allowed_proxy_hosts=[], version="3.0.0", ) else: service.authenticate( credential=f"dopg_invalid_{attempt}", certificate_sha256="b" * 64, gateway_id=gateway_id, environment="staging", network_zone=f"random-zone-{attempt}", generation=1, ) except EdgeGatewayError as exc: limited = exc assert limited is not None and limited.status_code == 429 with engine.connect() as connection: windows = connection.execute(text(""" SELECT failure_fingerprint,gateway_id,event_type,occurrence_count FROM edge_gateway_failure_windows """)).mappings().all() audits = connection.execute(text(""" SELECT gateway_id,event_type,safe_detail FROM edge_gateway_audits WHERE success=false """)).mappings().all() assert len(windows) == 1 and windows[0]["occurrence_count"] == 10 assert windows[0]["gateway_id"] is None assert windows[0]["event_type"] == "untrusted_request_rejected" assert len(audits) == 1 and audits[0]["gateway_id"] is None assert audits[0]["safe_detail"] == {"reason_code": "untrusted_failure"} def test_database_rejects_pending_enrollment_and_status_only_gateway_activation(service, engine): enrollment = _enrollment(service) values = { "gateway_id": enrollment["gateway_id"], "policy": "a" * 64, } with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(""" INSERT INTO edge_gateways (uid,name,environment,network_zone,status,current_generation,policy_digest, allowed_control_hosts,allowed_proxy_hosts,runtime_version) VALUES (CAST(:gateway_id AS uuid),'plant-a-edge','staging','manufacturing-zone-a', 'active',1,:policy,ARRAY['control.example.test'], ARRAY['proxy.example.test'],'3.0.0') """), values) with engine.begin() as connection: connection.execute(text(""" INSERT INTO edge_gateways (uid,name,environment,network_zone,status,current_generation,policy_digest, allowed_control_hosts,allowed_proxy_hosts,runtime_version) VALUES (CAST(:gateway_id AS uuid),'plant-a-edge','staging','manufacturing-zone-a', 'pending',1,:policy,ARRAY['control.example.test'], ARRAY['proxy.example.test'],'3.0.0') """), values) with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(""" UPDATE edge_gateways SET status='active' WHERE uid=CAST(:gateway_id AS uuid) """), values) def test_authentication_rotation_revocation_heartbeat_and_active_binding(service, engine): gateway = _register(service, _enrollment(service)) identity = service.authenticate( credential=gateway["credential"], certificate_sha256="b" * 64, gateway_id=gateway["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", generation=1, ) assert identity["gateway_id"] == gateway["gateway_id"] for statement, parameters in ( ("UPDATE edge_gateway_credentials SET certificate_sha256=:value WHERE gateway_id=CAST(:gateway AS uuid)", {"value": "c" * 64}), ("UPDATE edge_gateway_credentials SET credential_hash=:value WHERE gateway_id=CAST(:gateway AS uuid)", {"value": "d" * 64}), ("UPDATE edge_gateways SET name='attacker' WHERE uid=CAST(:gateway AS uuid)", {}), ("UPDATE edge_gateways SET environment='production' WHERE uid=CAST(:gateway AS uuid)", {}), ("UPDATE edge_gateways SET network_zone='attacker-zone' WHERE uid=CAST(:gateway AS uuid)", {}), ("UPDATE edge_gateways SET policy_digest=:value WHERE uid=CAST(:gateway AS uuid)", {"value": "c" * 64}), ("UPDATE edge_gateways SET allowed_control_hosts=ARRAY['attacker.test'] WHERE uid=CAST(:gateway AS uuid)", {}), ("UPDATE edge_gateways SET allowed_proxy_hosts=ARRAY['attacker.test'] WHERE uid=CAST(:gateway AS uuid)", {}), ("""UPDATE edge_gateway_credentials SET generation=2,certificate_sha256=:certificate, credential_hash=:credential WHERE gateway_id=CAST(:gateway AS uuid); UPDATE edge_gateways SET current_generation=2 WHERE uid=CAST(:gateway AS uuid)""", {"certificate": "c" * 64, "credential": "d" * 64}), ): with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(statement), {"gateway": gateway["gateway_id"], **parameters}) with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(""" UPDATE edge_gateway_credentials SET status='rotated',revoked_at=CURRENT_TIMESTAMP WHERE gateway_id=CAST(:gateway_id AS uuid) AND generation=1 """), {"gateway_id": gateway["gateway_id"]}) with pytest.raises(EdgeGatewayAuthenticationError): service.authenticate( credential=gateway["credential"], certificate_sha256="c" * 64, gateway_id=gateway["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", generation=1, ) with pytest.raises(EdgeGatewayValidationError): service.heartbeat( credential=gateway["credential"], certificate_sha256="b" * 64, gateway_id=gateway["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", generation=1, version="3.0.0", safe_summary={"password": "must-not-persist"}, ) with engine.connect() as connection: assert connection.execute(text("SELECT last_heartbeat_at FROM edge_gateways")).scalar_one() is None rotation_actor = _uid() rotated = service.rotate( gateway["gateway_id"], certificate_sha256="c" * 64, request_id="rotation-request-1", actor_uid=rotation_actor, ) rotation_replay = service.rotate( gateway["gateway_id"], certificate_sha256="c" * 64, request_id="rotation-request-1", actor_uid=rotation_actor, ) assert rotation_replay == { "gateway_id": gateway["gateway_id"], "generation": 2, "credential": None, "credential_returned_once": False, "replayed": True, } with pytest.raises(EdgeGatewayConflictError): service.rotate( gateway["gateway_id"], certificate_sha256="d" * 64, request_id="rotation-request-1", actor_uid=rotation_actor, ) with pytest.raises(EdgeGatewayAuthenticationError): service.authenticate( credential=gateway["credential"], certificate_sha256="b" * 64, gateway_id=gateway["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", generation=1, ) assert service.heartbeat( credential=rotated["credential"], certificate_sha256="c" * 64, gateway_id=gateway["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", generation=2, version="3.0.1", safe_summary={"queue_depth": 1}, )["status"] == "online" assert service.revoke(gateway["gateway_id"], actor_uid=_uid()) is True with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text("UPDATE edge_gateways SET status='active' WHERE uid=CAST(:gateway AS uuid)"), {"gateway": gateway["gateway_id"]}) with pytest.raises(EdgeGatewayAuthenticationError): service.authenticate( credential=rotated["credential"], certificate_sha256="c" * 64, gateway_id=gateway["gateway_id"], environment="staging", network_zone="manufacturing-zone-a", generation=2, ) with engine.connect() as connection: assert "untrusted_request_rejected" in set(connection.execute(text( "SELECT event_type FROM edge_gateway_audits WHERE success=false" )).scalars()) def test_database_rejects_second_active_credential_combination_injection(service, engine): gateway = _register(service, _enrollment(service)) with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(""" INSERT INTO edge_gateway_credentials (uid,gateway_id,generation,credential_hash,certificate_sha256,status,expires_at) VALUES (CAST(:uid AS uuid),CAST(:gateway AS uuid),2,:credential,:certificate, 'active',CURRENT_TIMESTAMP+INTERVAL '365 days'); UPDATE edge_gateways SET current_generation=2 WHERE uid=CAST(:gateway AS uuid) """), { "uid": _uid(), "gateway": gateway["gateway_id"], "credential": "9" * 64, "certificate": "8" * 64, }) with engine.connect() as connection: row = connection.execute(text(""" SELECT g.current_generation,count(*) FILTER (WHERE c.status='active') AS active_count FROM edge_gateways g JOIN edge_gateway_credentials c ON c.gateway_id=g.uid GROUP BY g.current_generation """)).mappings().one() assert row == {"current_generation": 1, "active_count": 1} def test_failure_evidence_tables_reject_direct_mutation_and_truncation(service, engine): with pytest.raises(EdgeGatewayAuthenticationError): service.authenticate( credential="dopg_invalid", certificate_sha256="b" * 64, gateway_id=_uid(), environment="staging", network_zone="zone-a", generation=1, ) with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text("TRUNCATE TABLE edge_gateway_audits")) for statement in ( "UPDATE edge_gateway_failure_windows SET occurrence_count=occurrence_count+1,last_seen_at=CURRENT_TIMESTAMP", "UPDATE edge_gateway_failure_windows SET occurrence_count=occurrence_count+5", "UPDATE edge_gateway_failure_windows SET last_seen_at=last_seen_at-INTERVAL '1 second'", "UPDATE edge_gateway_failure_windows SET first_seen_at=first_seen_at+INTERVAL '1 second'", "UPDATE edge_gateway_failure_windows SET failure_fingerprint=:fingerprint", "DELETE FROM edge_gateway_failure_windows", "TRUNCATE TABLE edge_gateway_failure_windows", ): with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text(statement), {"fingerprint": "f" * 64}) with engine.connect() as connection: evidence = connection.execute(text(""" SELECT occurrence_count,first_seen_at<=last_seen_at AS monotonic FROM edge_gateway_failure_windows """)).one() assert evidence == (1, True) assert connection.execute(text("SELECT count(*) FROM edge_gateway_audits")).scalar_one() == 1 def test_task_pull_cancel_event_idempotency_policy_and_release(service, engine): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } task = service.issue_task({ "task_id": "task-1", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-key-1", "policy_digest": "a" * 64, }, actor_uid=_uid()) conflicting_task = { "task_id": "task-conflicting", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "different-purpose", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-key-1", "policy_digest": "a" * 64, } with pytest.raises(EdgeGatewayConflictError): service.issue_task(conflicting_task, actor_uid=_uid()) with pytest.raises(EdgeGatewayValidationError): service.issue_task({ "task_id": "task-wrong-binding", "gateway_id": gateway["gateway_id"], "environment": "production", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-wrong-binding-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) pulled = service.pull_task(**auth) assert pulled["task"]["task_id"] == task["task_id"] approved = _event(task, gateway) with pytest.raises(EdgeGatewayAuthenticationError): service.accept_event(event=approved, lease_token="dopl_wrong", **auth) ack1 = service.accept_event(event=approved, lease_token=pulled["lease_token"], **auth) ack2 = service.accept_event(event=approved, lease_token=pulled["lease_token"], **auth) assert ack1 == ack2 assert set(ack1) == { "event_id", "event_digest", "remote_lease_digest", "received_at", "status", } assert ack1["event_digest"] == EdgeEventContract.from_mapping(approved).digest assert ack1["remote_lease_digest"] == hashlib.sha256( pulled["lease_token"].encode("utf-8") ).hexdigest() assert pulled["lease_token"] not in str(ack1) changed = dict(approved) changed["payload"] = {"asset_count": 13} with pytest.raises((EdgeGatewayConflictError, EdgeGatewayValidationError)): service.accept_event(event=changed, **auth) raw = dict(approved) raw.update(classification="raw", payload={"password": "super-secret-value"}, event_id="evt_" + "0" * 64) with pytest.raises(EdgeGatewayValidationError): service.accept_event(event=raw, **auth) wrong_purpose = dict(approved) wrong_purpose["purpose"] = "different-purpose" identity = {key: value for key, value in wrong_purpose.items() if key != "event_id"} wrong_purpose["event_id"] = stable_event_id(identity) with pytest.raises((EdgeGatewayValidationError, EdgeGatewayAuthenticationError)): service.accept_event(event=wrong_purpose, lease_token=pulled["lease_token"], **auth) assert service.cancel_task(task["task_id"], actor_uid=_uid()) is False assert service.reconcile(**auth)["cancelled_task_ids"] == [] assert service.accept_event( event=approved, lease_token=pulled["lease_token"], **auth, ) == ack1 with engine.connect() as connection: assert connection.execute(text("SELECT count(*) FROM edge_gateway_events")).scalar_one() == 1 release = service.offer_release( gateway["gateway_id"], version="3.1.0", artifact_digest="d" * 64, rollback_version="3.0.0", request_id="release-request-1", actor_uid=_uid(), ) assert service.offer_release( gateway["gateway_id"], version="3.1.0", artifact_digest="d" * 64, rollback_version="3.0.0", request_id="release-request-1", actor_uid=_uid(), ) == release with pytest.raises(EdgeGatewayConflictError): service.offer_release( gateway["gateway_id"], version="3.1.1", artifact_digest="d" * 64, rollback_version="3.0.0", request_id="release-request-1", actor_uid=_uid(), ) with pytest.raises(EdgeGatewayConflictError): service.acknowledge_release( release_id=release["release_id"], outcome="rollback", safe_summary={"reason_code": "health_check_failed"}, **auth, ) accepted = service.acknowledge_release( release_id=release["release_id"], outcome="accepted", safe_summary={"reason_code": "accepted"}, **auth, ) pending = service.reconcile(**auth)["release_offers"] assert [(item["release_id"], item["status"]) for item in pending] == [(release["release_id"], "accepted")] assert service.acknowledge_release( release_id=release["release_id"], outcome="accepted", safe_summary={"reason_code": "accepted"}, **auth, ) == accepted service.acknowledge_release( release_id=release["release_id"], outcome="failed", safe_summary={"reason_code": "health_check_failed"}, **auth, ) acked = service.acknowledge_release( release_id=release["release_id"], outcome="rollback", safe_summary={"reason_code": "health_check_failed"}, **auth, ) assert acked["status"] == "rolled_back" installed_release = service.offer_release( gateway["gateway_id"], version="3.2.0", artifact_digest="e" * 64, rollback_version="3.1.0", request_id="release-request-2", actor_uid=_uid(), ) service.acknowledge_release( release_id=installed_release["release_id"], outcome="accepted", safe_summary={"reason_code": "accepted"}, **auth, ) service.acknowledge_release( release_id=installed_release["release_id"], outcome="installed", safe_summary={"reason_code": "installed"}, **auth, ) installed_rollback = service.acknowledge_release( release_id=installed_release["release_id"], outcome="rollback", safe_summary={"reason_code": "post_install_health_failed"}, **auth, ) assert installed_rollback["status"] == "rolled_back" assert service.acknowledge_release( release_id=installed_release["release_id"], outcome="rollback", safe_summary={"reason_code": "post_install_health_failed"}, **auth, ) == installed_rollback with engine.connect() as connection: audit_rows = connection.execute(text("SELECT event_type,safe_detail FROM edge_gateway_audits")).mappings().all() details = [row["safe_detail"] for row in audit_rows] assert all("credential" not in str(detail).lower() and "token" not in str(detail).lower() for detail in details) assert all("super-secret-value" not in str(detail) for detail in details) assert {"idempotency_conflict", "binding_rejected", "policy_rejected", "release_rejected"} <= { row["event_type"] for row in audit_rows } with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text("UPDATE edge_gateway_audits SET success=true")) with pytest.raises(DBAPIError), engine.begin() as connection: connection.execute(text("DELETE FROM edge_gateway_audits")) task_failed = service.issue_task({ "task_id": "task-failed", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "quality", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-failed-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) failed_lease = service.pull_task(**auth) failed = service.task_outcome( task_failed["task_id"], outcome="failed", lease_token=failed_lease["lease_token"], safe_summary={"reason_code": "controlled_failure"}, **auth, ) assert failed["status"] == "failed" assert service.task_outcome( task_failed["task_id"], outcome="failed", lease_token=failed_lease["lease_token"], safe_summary={"reason_code": "controlled_failure"}, **auth, )["replayed"] is True task_cancel = service.issue_task({ "task_id": "task-cancelled", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-cancelled-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) cancel_lease = service.pull_task(**auth) assert service.cancel_task(task_cancel["task_id"], actor_uid=_uid()) is True assert service.task_outcome( task_cancel["task_id"], outcome="cancelled", lease_token=cancel_lease["lease_token"], safe_summary={"reason_code": "cancel_acknowledged"}, **auth, )["status"] == "cancelled" def test_control_plane_signs_task_authority_and_complete_release_state(service, engine): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } deadline = (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z") task = service.issue_task({ "task_id": "task-signed-authority", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": deadline, "attempt": 1, "idempotency_key": "task-signed-authority-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) pulled = service.pull_task(**auth) envelope = pulled["signed_task_envelope"] assert pulled["task"] == envelope["task"] assert envelope["task"]["task_id"] == task["task_id"] assert set(envelope) == { "task", "authority_key_id", "signature_algorithm", "contract_digest", "gateway_id", "environment", "network_zone", "policy_digest", "purpose", "issued_at", "expires_at", "signature", } assert envelope["signature_algorithm"] == "Ed25519" assert envelope["authority_key_id"] == _SIGNING_KEY_ID assert envelope["contract_digest"] == canonical_sha256(envelope["task"]) unsigned_authority = {key: value for key, value in envelope.items() if key != "signature"} _SIGNING_PRIVATE_KEY.public_key().verify( bytes.fromhex(envelope["signature"]), SignedTaskEnvelope.canonical_unsigned_bytes(unsigned_authority), ) release_deadline = (datetime.now(UTC) + timedelta(hours=1)).isoformat().replace("+00:00", "Z") release = service.offer_release( gateway["gateway_id"], version="6.0.0", artifact_digest="6" * 64, artifact_name="edge-agent-6.0.0.bin", deadline_at=release_deadline, rollback_version="3.0.0", request_id="signed-release-1", actor_uid=_uid(), ) manifest = release["manifest"] with pytest.raises(EdgeGatewayValidationError): service.offer_release( gateway["gateway_id"], version="release-six", artifact_digest="6" * 64, rollback_version="3.0.0", request_id="invalid-release-version", actor_uid=_uid(), ) with pytest.raises(EdgeGatewayValidationError): service.offer_release( gateway["gateway_id"], version="2.9.0", artifact_digest="6" * 64, rollback_version="3.0.0", request_id="release-downgrade", actor_uid=_uid(), ) assert set(manifest) == { "release_id", "version", "rollback_version", "artifact_digest", "artifact_name", "deadline_at", "signature_algorithm", "key_id", "manifest_digest", "signature", "status", } assert manifest["status"] == "offered" unsigned_manifest = { key: value for key, value in manifest.items() if key not in {"manifest_digest", "signature"} } assert manifest["manifest_digest"] == canonical_sha256(unsigned_manifest) _SIGNING_PRIVATE_KEY.public_key().verify( bytes.fromhex(manifest["signature"]), canonical_json_bytes(unsigned_manifest), ) for outcome, expected in (("accepted", "accepted"), ("failed", "failed"), ("rollback", "rolled_back")): service.acknowledge_release( release_id=release["release_id"], outcome=outcome, safe_summary={"release_status": expected, "version": "3.0.0"}, **auth, ) response = service.reconcile(**auth) candidates = list(response["release_offers"]) if response["release_baseline"] is not None: candidates.append(response["release_baseline"]) reconciled = { item["release_id"]: item for item in candidates }[release["release_id"]] assert reconciled["status"] == expected unsigned = { key: value for key, value in reconciled.items() if key not in {"manifest_digest", "signature"} } assert reconciled["manifest_digest"] == canonical_sha256(unsigned) _SIGNING_PRIVATE_KEY.public_key().verify( bytes.fromhex(reconciled["signature"]), canonical_json_bytes(unsigned), ) with engine.connect() as connection: task_row = connection.execute(text(""" SELECT authority,contract FROM edge_gateway_tasks WHERE task_id='task-signed-authority' """)).mappings().one() release_row = connection.execute(text(""" SELECT signed_manifest FROM edge_gateway_releases WHERE uid=CAST(:uid AS uuid) """), {"uid": release["release_id"]}).scalar_one() assert dict(task_row["authority"]) == envelope assert dict(release_row)["status"] == "rolled_back" persisted = f"{task_row['authority']}{release_row}" assert "PRIVATE KEY" not in persisted and "signing_private_key" not in persisted @pytest.mark.parametrize("seconds", [30, 119, 120]) def test_pull_lease_never_extends_beyond_task_deadline(service, seconds): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } deadline = (datetime.now(UTC) + timedelta(seconds=seconds)).isoformat().replace( "+00:00", "Z" ) service.issue_task( { "task_id": f"deadline-boundary-{seconds}", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": deadline, "attempt": 1, "idempotency_key": f"deadline-boundary-{seconds}-key", "policy_digest": "a" * 64, }, actor_uid=_uid(), ) pulled = service.pull_task(**auth) assert datetime.fromisoformat(pulled["lease_expires_at"]) <= datetime.fromisoformat( deadline.replace("Z", "+00:00") ) def test_cancel_reconcile_uses_stable_cursor_so_more_than_fifty_never_starve( service, ): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } expected = [] deadline = (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace( "+00:00", "Z" ) for index in range(52): task_id = f"cancel-page-{index:02d}" service.issue_task( { "task_id": task_id, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": deadline, "attempt": 1, "idempotency_key": f"cancel-page-key-{index:02d}", "policy_digest": "a" * 64, }, actor_uid=_uid(), ) assert service.pull_task(**auth)["task"]["task_id"] == task_id assert service.cancel_task(task_id, actor_uid=_uid()) is True expected.append(task_id) first = service.reconcile(limit=50, cancel_cursor=None, release_cursor=None, **auth) assert first["cancelled_task_ids"] == expected[:50] assert first["cancel_next_cursor"] == expected[49] second = service.reconcile( limit=50, cancel_cursor=first["cancel_next_cursor"], release_cursor=None, **auth, ) assert second["cancelled_task_ids"] == expected[50:] assert second["cancel_next_cursor"] is None def test_pull_terminalizes_expired_task_without_issuing_lease(service, engine): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } task_id = "already-expired-task" service.issue_task( { "task_id": task_id, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(milliseconds=200)) .isoformat() .replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "already-expired-task-key", "policy_digest": "a" * 64, }, actor_uid=_uid(), ) with engine.connect() as connection: connection.execute(text("SELECT pg_sleep(0.3)")) assert service.pull_task(**auth) == {"task": None} with engine.connect() as connection: row = connection.execute( text( "SELECT status,lease_token_hash,lease_expires_at,result_summary " "FROM edge_gateway_tasks WHERE task_id=:task_id" ), {"task_id": task_id}, ).mappings().one() assert row["status"] == "failed" assert row["lease_token_hash"] is None and row["lease_expires_at"] is None assert dict(row["result_summary"]) == {"reason_code": "task_deadline_expired"} def test_reconcile_prioritizes_new_actionable_release_over_terminal_history(service): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } for index in range(51): release = service.offer_release( gateway["gateway_id"], version=f"4.{index}.0", rollback_version="3.0.0", artifact_digest=f"{index:064x}", request_id=f"terminal-release-{index}", actor_uid=_uid(), ) service.acknowledge_release( release_id=release["release_id"], outcome="accepted", safe_summary={"release_status": "accepted", "version": f"4.{index}.0"}, **auth, ) service.acknowledge_release( release_id=release["release_id"], outcome="installed", safe_summary={"release_status": "installed", "version": f"4.{index}.0"}, **auth, ) actionable = service.offer_release( gateway["gateway_id"], version="9.0.0", rollback_version="3.0.0", artifact_digest="f" * 64, request_id="new-actionable-release", actor_uid=_uid(), ) reconciled = service.reconcile(limit=50, **auth) releases = reconciled["release_offers"] assert len(releases) == 1 assert releases[0]["release_id"] == actionable["release_id"] assert releases[0]["status"] == "offered" assert reconciled["release_next_cursor"] is None assert reconciled["release_baseline"]["status"] == "installed" assert reconciled["release_baseline"]["version"] == "4.50.0" def test_signed_release_offer_installs_through_real_postgres_control_plane( service, engine, tmp_path, ): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } artifact = b"postgres-backed-signed-edge-release" offered = service.offer_release( gateway["gateway_id"], version="3.1.0", rollback_version="3.0.0", artifact_digest=hashlib.sha256(artifact).hexdigest(), artifact_name="edge-agent-3.1.0.bin", deadline_at=(datetime.now(UTC) + timedelta(hours=1)) .isoformat() .replace("+00:00", "Z"), request_id="postgres-agent-release-chain", actor_uid=_uid(), ) class PostgresControlTransport: def reconcile(self): return service.reconcile(**auth) def acknowledge_release(self, release_id, outcome, summary): return service.acknowledge_release( release_id=release_id, outcome=outcome, safe_summary=summary, **auth, ) def pull_task(self): return None public_key = _SIGNING_PRIVATE_KEY.public_key().public_bytes( encoding=serialization.Encoding.Raw, format=serialization.PublicFormat.Raw, ).hex() activated = [] agent = EdgeAgent( EdgeBootstrapConfig( queue_path=str(tmp_path / "postgres-release-agent.sqlite3"), artifact_root=str(tmp_path), gateway_id=gateway["gateway_id"], credential=gateway["credential"], certificate_sha256="b" * 64, client_certificate_path=str(tmp_path / "client.pem"), client_private_key_path=str(tmp_path / "client.key"), ca_bundle_path=str(tmp_path / "ca.pem"), generation=1, environment="staging", network_zone="manufacturing-zone-a", policy_digest="a" * 64, control_url="https://control.example.test", proxy_url=None, allowed_control_hosts=frozenset({"control.example.test"}), allowed_proxy_hosts=frozenset(), trusted_release_keys={_SIGNING_KEY_ID: public_key}, trusted_task_keys={_SIGNING_KEY_ID: public_key}, version="3.0.0", ), PostgresControlTransport(), EdgeRunnerAdapter( lambda _node, _request, _cancel: { "classification": "statistics", "payload": {"metric_count": 1}, } ), release_manager=SignedReleaseManager( current_version="3.0.0", trusted_keys={_SIGNING_KEY_ID: public_key}, artifact_loader=lambda _manifest: artifact, activator=lambda version, body: activated.append((version, body)) or True, rollback=lambda _version: False, ), ) agent.run_once() with engine.connect() as connection: persisted = connection.execute( text("SELECT status,signed_manifest FROM edge_gateway_releases WHERE uid=CAST(:uid AS uuid)"), {"uid": offered["release_id"]}, ).mappings().one() assert persisted["status"] == "installed" assert dict(persisted["signed_manifest"])["status"] == "installed" assert activated == [("3.1.0", artifact)] assert agent.queue.get_release_state(offered["release_id"]).status == "installed" def test_signing_operations_fail_closed_without_configured_private_key(service, engine): gateway = _register(service, _enrollment(service)) session = Session(engine) unsigned_service = EdgeGatewayService(EdgeGatewayRepository(session)) try: with pytest.raises(EdgeGatewayConfigurationError): unsigned_service.issue_task({ "task_id": "task-no-signer", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-no-signer-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) finally: session.close() with pytest.raises(EdgeGatewayValidationError): service.issue_task({ "task_id": "task-expired-authority", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) - timedelta(seconds=1)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-expired-authority-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) def test_task_outcome_audit_is_atomic_bounded_and_replay_idempotent(service, engine): gateway = _register(service, _enrollment(service)) auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} task = service.issue_task({ "task_id": "task-outcome-audit", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "quality", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-outcome-audit-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) lease = service.pull_task(**auth) with pytest.raises(EdgeGatewayValidationError): service.task_outcome( task["task_id"], outcome="failed", lease_token=lease["lease_token"], safe_summary={"reason_code": "x" * 40_000}, **auth, ) with engine.begin() as connection: connection.execute(text(""" CREATE OR REPLACE FUNCTION test_reject_task_outcome_audit() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN IF NEW.event_type='task_outcome_recorded' THEN RAISE EXCEPTION 'test audit failure'; END IF; RETURN NEW; END $$; CREATE TRIGGER test_reject_task_outcome_audit BEFORE INSERT ON edge_gateway_audits FOR EACH ROW EXECUTE FUNCTION test_reject_task_outcome_audit(); """)) session_usable = False try: with pytest.raises(DBAPIError): service.task_outcome( task["task_id"], outcome="failed", lease_token=lease["lease_token"], safe_summary={"reason_code": "controlled_failure"}, **auth, ) service.repository.session.execute(text("SELECT 1")).scalar_one() session_usable = True finally: service.repository.rollback() with engine.begin() as connection: connection.execute(text("DROP TRIGGER test_reject_task_outcome_audit ON edge_gateway_audits")) connection.execute(text("DROP FUNCTION test_reject_task_outcome_audit()")) assert session_usable is True with engine.connect() as connection: assert connection.execute(text(""" SELECT status FROM edge_gateway_tasks WHERE task_id=:task_id """), {"task_id": task["task_id"]}).scalar_one() == "leased" assert connection.execute(text(""" SELECT count(*) FROM edge_gateway_audits WHERE event_type='task_outcome_recorded' """)).scalar_one() == 0 result = service.task_outcome( task["task_id"], outcome="failed", lease_token=lease["lease_token"], safe_summary={"reason_code": "controlled_failure"}, **auth, ) assert result["status"] == "failed" and result["replayed"] is False replay = service.task_outcome( task["task_id"], outcome="failed", lease_token=lease["lease_token"], safe_summary={"reason_code": "controlled_failure"}, **auth, ) assert replay["status"] == "failed" and replay["replayed"] is True with engine.connect() as connection: audits = connection.execute(text(""" SELECT safe_detail FROM edge_gateway_audits WHERE event_type='task_outcome_recorded' """)).scalars().all() assert audits == [{"task_id": task["task_id"], "outcome": "failed"}] def test_cross_gateway_event_replay_never_discloses_foreign_ack(service): gateway_a = _register(service, _enrollment(service, gateway_name="gateway-a")) gateway_b = _register(service, _enrollment(service, gateway_name="gateway-b"), certificate_sha256="b" * 64) auth_a = {"credential": gateway_a["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway_a["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} auth_b = {"credential": gateway_b["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway_b["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} task = service.issue_task({ "task_id": "task-owner-a", "gateway_id": gateway_a["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-owner-a-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) pulled = service.pull_task(**auth_a) event = _event(task, gateway_a) ack = service.accept_event(event=event, lease_token=pulled["lease_token"], **auth_a) with pytest.raises(EdgeGatewayAuthenticationError): service.accept_event(event=event, lease_token=pulled["lease_token"], **auth_b) assert ack["event_id"] == event["event_id"] def test_release_offer_is_concurrently_idempotent_and_conflicts_are_controlled(service, engine): gateway = _register(service, _enrollment(service)) barrier = Barrier(2) def offer_once(): session = Session(engine) try: barrier.wait(timeout=5) return _service(EdgeGatewayRepository(session)).offer_release( gateway["gateway_id"], version="5.0.0", artifact_digest="5" * 64, rollback_version="4.9.0", request_id="concurrent-release", actor_uid=_uid(), ) finally: session.close() with ThreadPoolExecutor(max_workers=2) as pool: results = [future.result() for future in (pool.submit(offer_once), pool.submit(offer_once))] assert results[0] == results[1] with pytest.raises(EdgeGatewayConflictError): service.offer_release( gateway["gateway_id"], version="5.0.1", artifact_digest="6" * 64, rollback_version="4.9.0", request_id="concurrent-release", actor_uid=_uid(), ) with pytest.raises(EdgeGatewayConflictError): service.offer_release( gateway["gateway_id"], version="5.0.0", artifact_digest="5" * 64, rollback_version="4.9.0", request_id="different-request", actor_uid=_uid(), ) service.offer_release( gateway["gateway_id"], version="5.1.0", artifact_digest="7" * 64, rollback_version="5.0.0", request_id="release-b", actor_uid=_uid(), ) with pytest.raises(EdgeGatewayConflictError): service.offer_release( gateway["gateway_id"], version="5.1.0", artifact_digest="7" * 64, rollback_version="5.0.0", request_id="concurrent-release", actor_uid=_uid(), ) assert service.list_gateways()[0]["gateway_id"] == gateway["gateway_id"] def test_evidence_owner_and_runtime_role_are_really_isolated(engine): runtime_login = f"edge_test_runtime_{uuid.uuid4().hex[:12]}" runtime_password = f"runtime-{uuid.uuid4().hex}" runtime_engine = None runtime_session = None with engine.begin() as connection: connection.execute(text(f'CREATE ROLE "{runtime_login}" LOGIN PASSWORD :password'), { "password": runtime_password, }) connection.execute(text(f'GRANT dataops_app_runtime TO "{runtime_login}"')) try: runtime_url = make_url(DATABASE_URL).set( username=runtime_login, password=runtime_password, ) runtime_engine = create_engine(runtime_url, pool_pre_ping=True) validate_production_database_identity( SimpleNamespace(config={"FLASK_ENV": "production"}), runtime_engine, ) runtime_session = Session(runtime_engine) runtime_service = _service(EdgeGatewayRepository(runtime_session)) enrollment = _enrollment(runtime_service) gateway = _register(runtime_service, enrollment) assert gateway["generation"] == 1 auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } assert runtime_service.heartbeat( version="3.0.1", safe_summary={"queue_depth": 0}, **auth, )["status"] == "online" task = runtime_service.issue_task({ "task_id": "runtime-role-task", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "runtime-role-task-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) assert runtime_service.pull_task(**auth)["task"]["task_id"] == task["task_id"] with runtime_engine.connect() as connection: identity = connection.execute(text("SELECT current_user,session_user")).one() assert identity == (runtime_login, runtime_login) owners = dict(connection.execute(text(""" SELECT c.relname,r.rolname FROM pg_class c JOIN pg_roles r ON r.oid=c.relowner WHERE c.relname IN ('edge_gateway_audits','edge_gateway_failure_windows') """)).all()) assert owners == { "edge_gateway_audits": "dataops_edge_evidence_owner", "edge_gateway_failure_windows": "dataops_edge_evidence_owner", } assert connection.execute(text(""" SELECT count(*)=0 FROM information_schema.routine_privileges WHERE routine_schema='public' AND routine_name IN ('record_edge_gateway_failure','record_edge_gateway_audit') AND grantee='PUBLIC' AND privilege_type='EXECUTE' """)).scalar_one() is True attacks = ( "INSERT INTO edge_gateway_audits(uid,event_type,success,safe_detail) VALUES (gen_random_uuid(),'authentication_rejected',false,'{}')", "UPDATE edge_gateway_audits SET success=true", "DELETE FROM edge_gateway_failure_windows", "TRUNCATE TABLE edge_gateway_audits,edge_gateway_failure_windows", "ALTER TABLE edge_gateway_audits DISABLE TRIGGER ALL", "SET ROLE dataops_edge_evidence_owner", ) for statement in attacks: with pytest.raises(DBAPIError), runtime_engine.begin() as connection: connection.execute(text(statement)) with runtime_engine.begin() as connection: connection.execute(text("SELECT set_config('dataops.edge_failure_writer','477-controlled',false)")) with pytest.raises(DBAPIError): connection.execute(text(""" INSERT INTO edge_gateway_failure_windows (failure_fingerprint,window_started_at,event_type,occurrence_count,first_seen_at,last_seen_at) VALUES (:fingerprint,CURRENT_TIMESTAMP,'authentication_rejected',1,CURRENT_TIMESTAMP,CURRENT_TIMESTAMP) """), {"fingerprint": "f" * 64}) finally: if runtime_session is not None: runtime_session.close() if runtime_engine is not None: runtime_engine.dispose() with engine.begin() as connection: connection.execute(text(f'DROP ROLE IF EXISTS "{runtime_login}"')) def test_production_owner_database_identity_and_direct_entry_fail_closed(engine): with pytest.raises(RuntimeError, match="least-privilege runtime identity"): validate_production_database_identity( SimpleNamespace(config={"FLASK_ENV": "production"}), engine, ) result = subprocess.run( [str(ROOT / ".venv/bin/python"), "-c", "from app import create_app; create_app()"], cwd=ROOT, env={ **os.environ, "PYTHONPATH": str(ROOT), "FLASK_ENV": "production", "APP_ENV_FILE": "/nonexistent/dataops-test.env", "DATABASE_URL": DATABASE_URL, "NEO4J_URI": "bolt://neo4j.example.test:7687", "NEO4J_HTTP_URI": "http://neo4j.example.test:7474", "NEO4J_USER": "dataops_graph", "NEO4J_PASSWORD": "graph-test-secret", "MINIO_HOST": "minio.example.test:9000", "MINIO_USER": "dataops_objects", "MINIO_PASSWORD": "object-test-secret", "MINIO_BUCKET": "dataops-test", }, capture_output=True, text=True, ) assert result.returncode != 0 assert "least-privilege runtime identity" in result.stderr assert "dataops-test-password" not in result.stderr def test_preprovisioned_no_createrole_migrator_can_apply_477_and_478(engine): migrator = f"edge_test_migrator_{uuid.uuid4().hex[:12]}" password = f"migration-{uuid.uuid4().hex}" _clear_edge_rows_for_test(engine) _alembic("downgrade", "20260802_476") try: with engine.begin() as connection: connection.execute(text(f'CREATE ROLE "{migrator}" LOGIN NOSUPERUSER NOCREATEROLE NOCREATEDB PASSWORD :password'), {"password": password}) connection.execute(text(f'GRANT dataops_edge_evidence_owner TO "{migrator}"')) connection.execute(text('GRANT USAGE,CREATE ON SCHEMA public TO dataops_edge_evidence_owner')) connection.execute(text(f'GRANT USAGE,CREATE ON SCHEMA public TO "{migrator}"')) connection.execute(text(f'GRANT SELECT,INSERT,UPDATE,DELETE ON alembic_version TO "{migrator}"')) migrator_url = make_url(DATABASE_URL).set( username=migrator, password=password, ).render_as_string(hide_password=False) try: _alembic("upgrade", "head", migrator_url) except subprocess.CalledProcessError as exc: pytest.fail(exc.stderr) with engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479" assert connection.execute(text("SELECT NOT rolcreaterole AND NOT rolsuper FROM pg_roles WHERE rolname=:role"), {"role": migrator}).scalar_one() is True finally: _clear_edge_rows_for_test(engine) _alembic("downgrade", "20260802_476") with engine.begin() as connection: connection.execute(text(f'REVOKE dataops_edge_evidence_owner FROM "{migrator}"')) connection.execute(text(f'DROP OWNED BY "{migrator}"')) connection.execute(text(f'DROP ROLE "{migrator}"')) _alembic("upgrade", "head") def test_backend_runtime_configuration_never_contains_migrator_credentials(): compose = (ROOT / "deploy/docker/docker-compose.yml").read_text() backend_dockerfile = (ROOT / "deploy/docker/backend.Dockerfile").read_text() runner_dockerfile = (ROOT / "deploy/docker/runner.Dockerfile").read_text() run_script = (ROOT / "scripts/run_dataops.sh").read_text() assert "db-role-init:" in compose and "db-migrate:" in compose assert "DATAOPS_RUNTIME_PASSWORD" in compose assert "postgresql://dataops_app:" in compose backend = compose.split(" backend:", 1)[1].split(" frontend:", 1)[0] assert "postgresql://dataops:dataops-test-password" not in backend assert "MIGRATION_DATABASE_URL" not in backend assert "db-migrate:\n condition: service_completed_successfully" in backend backend_runtime = backend_dockerfile.split("FROM dependencies AS runtime", 1)[1] assert "alembic.ini" not in backend_runtime and "COPY migrations/" not in backend_runtime assert "pip uninstall -y alembic" in backend_runtime assert "alembic.ini" not in runner_dockerfile and "COPY migrations/" not in runner_dockerfile assert "pip uninstall -y alembic" in runner_dockerfile assert 'role_init_url="${DB_ROLE_INIT_DATABASE_URL:' in run_script assert 'migration_url="${MIGRATION_DATABASE_URL:' in run_script assert "unset DB_ROLE_INIT_DATABASE_URL MIGRATION_DATABASE_URL" in run_script assert "MIGRATION_DATABASE_URL=" in (ROOT / "deployment/dataops.env").read_text() def test_task_pull_and_exact_event_replay_are_concurrency_safe(service, engine): gateway = _register(service, _enrollment(service)) auth = { "credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1, } task = service.issue_task({ "task_id": "task-concurrent", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-concurrent-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) pull_barrier = Barrier(2) def pull_once(): session = Session(engine) try: worker = _service(EdgeGatewayRepository(session)) pull_barrier.wait() return worker.pull_task(**auth) finally: session.close() with ThreadPoolExecutor(max_workers=2) as pool: pulls = [future.result() for future in (pool.submit(pull_once), pool.submit(pull_once))] assert sorted(result["task"] is not None for result in pulls) == [False, True] lease_token = next(result["lease_token"] for result in pulls if result["task"] is not None) event = _event(task, gateway) event_barrier = Barrier(2) def accept_once(): session = Session(engine) try: worker = _service(EdgeGatewayRepository(session)) event_barrier.wait() return worker.accept_event(event=event, lease_token=lease_token, **auth) finally: session.close() with ThreadPoolExecutor(max_workers=2) as pool: acks = [future.result() for future in (pool.submit(accept_once), pool.submit(accept_once))] assert acks[0] == acks[1] with engine.connect() as connection: assert connection.execute(text("SELECT count(*) FROM edge_gateway_events WHERE event_id=:event_id"), {"event_id": event["event_id"]}).scalar_one() == 1 assert service.accept_event(event=event, lease_token=lease_token, **auth) == acks[0] def test_cancel_and_event_barrier_has_single_linearized_terminal_state(service, engine): gateway = _register(service, _enrollment(service)) auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} task = service.issue_task({ "task_id": "task-race", "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": "task-race-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) lease = service.pull_task(**auth) event = _event(task, gateway) barrier = Barrier(2) def accept(): session = Session(engine) try: barrier.wait() return _service(EdgeGatewayRepository(session)).accept_event( event=event, lease_token=lease["lease_token"], **auth, ) except EdgeGatewayError: return None finally: session.close() def cancel(): session = Session(engine) try: barrier.wait() return _service(EdgeGatewayRepository(session)).cancel_task(task["task_id"], actor_uid=_uid()) finally: session.close() with ThreadPoolExecutor(max_workers=2) as pool: accepted, cancelled = pool.submit(accept), pool.submit(cancel) accepted, cancelled = accepted.result(), cancelled.result() with engine.connect() as connection: row = connection.execute(text(""" SELECT status,(SELECT count(*) FROM edge_gateway_events WHERE task_id=:task_id) AS events FROM edge_gateway_tasks WHERE task_id=:task_id """), {"task_id": task["task_id"]}).mappings().one() assert (row["status"], row["events"]) in {("completed", 1), ("cancel_requested", 0)} assert bool(accepted) != bool(cancelled) def test_cancel_event_barriers_prove_both_lock_orders(service, engine): gateway = _register(service, _enrollment(service)) auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} def issue_and_pull(task_id): task = service.issue_task({ "task_id": task_id, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": f"{task_id}-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) return task, service.pull_task(**auth) cancel_task, cancel_lease = issue_and_pull("task-cancel-wins") cancel_has_gateway = Event() cancel_barrier = Barrier(2) class CancelFirstRepository(EdgeGatewayRepository): def cancel_task(self, task_id, actor_uid): self.session.execute(text(""" SELECT g.uid FROM edge_gateway_tasks t JOIN edge_gateways g ON g.uid=t.gateway_id WHERE t.task_id=:task_id FOR UPDATE OF g """), {"task_id": task_id}).scalar_one() cancel_has_gateway.set() cancel_barrier.wait(timeout=5) return super().cancel_task(task_id, actor_uid) def cancel_first(): session = Session(engine) try: return _service(CancelFirstRepository(session)).cancel_task( cancel_task["task_id"], actor_uid=_uid(), ) finally: session.close() def event_after_cancel_lock(): assert cancel_has_gateway.wait(timeout=5) cancel_barrier.wait(timeout=5) session = Session(engine) try: return _service(EdgeGatewayRepository(session)).accept_event( event=_event(cancel_task, gateway), lease_token=cancel_lease["lease_token"], **auth, ) except EdgeGatewayError: return None finally: session.close() with ThreadPoolExecutor(max_workers=2) as pool: cancel_result = pool.submit(cancel_first) event_result = pool.submit(event_after_cancel_lock) assert cancel_result.result() is True assert event_result.result() is None event_task, event_lease = issue_and_pull("task-event-wins") event_has_gateway = Event() event_barrier = Barrier(2) class EventFirstRepository(EdgeGatewayRepository): def authenticate(self, credential_hash, *, lock=False): identity = super().authenticate(credential_hash, lock=lock) if lock: event_has_gateway.set() event_barrier.wait(timeout=5) return identity def event_first(): session = Session(engine) try: return _service(EventFirstRepository(session)).accept_event( event=_event(event_task, gateway), lease_token=event_lease["lease_token"], **auth, ) finally: session.close() def cancel_after_event_lock(): assert event_has_gateway.wait(timeout=5) event_barrier.wait(timeout=5) session = Session(engine) try: return _service(EdgeGatewayRepository(session)).cancel_task( event_task["task_id"], actor_uid=_uid(), ) finally: session.close() with ThreadPoolExecutor(max_workers=2) as pool: event_result = pool.submit(event_first) cancel_result = pool.submit(cancel_after_event_lock) assert event_result.result()["status"] == "accepted" assert cancel_result.result() is False with engine.connect() as connection: rows = dict(connection.execute(text(""" SELECT task_id,status FROM edge_gateway_tasks WHERE task_id IN ('task-cancel-wins','task-event-wins') """)).all()) assert rows == {"task-cancel-wins": "cancel_requested", "task-event-wins": "completed"} def test_machine_operations_and_rotate_are_gateway_fenced_with_barrier(service, engine): gateway = _register(service, _enrollment(service)) auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} release = service.offer_release( gateway["gateway_id"], version="4.0.0", artifact_digest="7" * 64, rollback_version="3.0.0", request_id="barrier-release", actor_uid=_uid(), ) barrier = Barrier(3) def run_machine(kind): session = Session(engine) worker = _service(EdgeGatewayRepository(session)) try: barrier.wait() if kind == "reconcile": return worker.reconcile(**auth) return worker.acknowledge_release( release_id=release["release_id"], outcome="accepted", safe_summary={"reason_code": "accepted"}, **auth, ) except EdgeGatewayError: return None finally: session.close() def run_rotate(): session = Session(engine) try: barrier.wait() return _service(EdgeGatewayRepository(session)).rotate( gateway["gateway_id"], certificate_sha256="c" * 64, request_id="barrier-rotate", actor_uid=_uid(), ) finally: session.close() with ThreadPoolExecutor(max_workers=3) as pool: results = [future.result() for future in ( pool.submit(run_machine, "reconcile"), pool.submit(run_machine, "ack"), pool.submit(run_rotate), )] assert results[2]["generation"] == 2 with pytest.raises(EdgeGatewayAuthenticationError): service.reconcile(**auth) def test_pull_event_and_revoke_are_gateway_fenced_with_barrier(service, engine): gateway = _register(service, _enrollment(service)) auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "generation": 1} def issue(task_id): return service.issue_task({ "task_id": task_id, "gateway_id": gateway["gateway_id"], "environment": "staging", "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics", "task_type": "profile", "contract_version": 1, "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"), "attempt": 1, "idempotency_key": f"{task_id}-key", "policy_digest": "a" * 64, }, actor_uid=_uid()) event_task = issue("task-revoke-event") event_lease = service.pull_task(**auth) issue("task-revoke-pull") event = _event(event_task, gateway) barrier = Barrier(3) def machine(kind): session = Session(engine) worker = _service(EdgeGatewayRepository(session)) try: barrier.wait() if kind == "event": return worker.accept_event(event=event, lease_token=event_lease["lease_token"], **auth) return worker.pull_task(**auth) except EdgeGatewayError: return None finally: session.close() def revoke(): session = Session(engine) try: barrier.wait() return _service(EdgeGatewayRepository(session)).revoke(gateway["gateway_id"], actor_uid=_uid()) finally: session.close() with ThreadPoolExecutor(max_workers=3) as pool: results = [future.result() for future in ( pool.submit(machine, "event"), pool.submit(machine, "pull"), pool.submit(revoke), )] assert results[2] is True with pytest.raises(EdgeGatewayAuthenticationError): service.pull_task(**auth) def test_failure_window_uses_lock_time_not_stale_transaction_time(engine): fingerprint = hashlib.sha256(f"stale-clock:{_uid()}".encode()).hexdigest() statement = text(""" SELECT occurrence_count FROM public.record_edge_gateway_failure( CAST(:audit_uid AS uuid), :fingerprint, NULL, 'untrusted_request_rejected', 'untrusted_failure' ) """) older = engine.connect() transaction = older.begin() try: older.execute(text("SELECT CURRENT_TIMESTAMP,pg_sleep(0.02)")) with engine.begin() as newer: first = newer.execute(statement, { "audit_uid": _uid(), "fingerprint": fingerprint, }).scalar_one() second = older.execute(statement, { "audit_uid": _uid(), "fingerprint": fingerprint, }).scalar_one() transaction.commit() finally: if transaction.is_active: transaction.rollback() older.close() assert (first, second) == (1, 2) def test_edge_flask_postgres_machine_binding_and_human_permissions(monkeypatch, engine): monkeypatch.setenv("DATABASE_URL", DATABASE_URL) actor = _uid() monkeypatch.setattr( "app.core.system.permissions.authenticate_request", lambda: {"id": actor, "sub": actor, "roles": ["admin"]}, ) from app import create_app app = create_app() app.config.update( TESTING=True, EDGE_GATEWAY_SIGNING_PRIVATE_KEY=_SIGNING_PRIVATE_KEY, EDGE_GATEWAY_SIGNING_KEY_ID=_SIGNING_KEY_ID, EDGE_MTLS_TRUSTED_PROXY_IPS=("127.0.0.1",), ) certificate_key = rsa.generate_private_key(public_exponent=65537, key_size=2048) certificate_name = x509.Name( [x509.NameAttribute(NameOID.COMMON_NAME, "plant-route-edge")] ) now = datetime.now(UTC) certificate = ( x509.CertificateBuilder() .subject_name(certificate_name) .issuer_name(certificate_name) .public_key(certificate_key.public_key()) .serial_number(x509.random_serial_number()) .not_valid_before(now - timedelta(minutes=1)) .not_valid_after(now + timedelta(minutes=10)) .sign(certificate_key, hashes.SHA256()) ) certificate_pem = certificate.public_bytes(serialization.Encoding.PEM).decode() certificate_sha256 = certificate.fingerprint(hashes.SHA256()).hex() def mtls_headers(**values): return { "X-DataOps-Edge-Client-Cert": quote(certificate_pem, safe=""), "X-DataOps-Edge-Client-Verify": "SUCCESS", "X-Edge-Certificate-SHA256": certificate_sha256, **values, } client = app.test_client() oversized = client.post("/api/datasource/edge/enrollments", json={"padding": "x" * 270_000}) assert oversized.status_code == 413 enrollment_response = client.post("/api/datasource/edge/enrollments", json={ "gateway_name": "plant-route", "environment": "staging", "network_zone": "zone-route", "policy_digest": "e" * 64, "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [], "expected_certificate_sha256": certificate_sha256, "ttl_seconds": 600, }) assert enrollment_response.status_code == 201 assert enrollment_response.headers["Cache-Control"] == "no-store" enrollment = enrollment_response.get_json()["data"] missing_tls_proof = client.post("/api/datasource/edge/register", headers={ "X-Edge-Enrollment": enrollment["enrollment_token"], "X-Edge-Certificate-SHA256": certificate_sha256, }, json={ "gateway_id": enrollment["gateway_id"], "environment": "staging", "network_zone": "zone-route", "policy_digest": "e" * 64, "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [], "version": "3.0.0", }) assert missing_tls_proof.status_code == 401 spoofed_tls_proof = client.post( "/api/datasource/edge/register", headers=mtls_headers(**{ "X-Edge-Enrollment": enrollment["enrollment_token"], }), environ_overrides={"REMOTE_ADDR": "203.0.113.7"}, json={ "gateway_id": enrollment["gateway_id"], "environment": "staging", "network_zone": "zone-route", "policy_digest": "e" * 64, "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [], "version": "3.0.0", }, ) assert spoofed_tls_proof.status_code == 401 register = client.post("/api/datasource/edge/register", headers=mtls_headers(**{ "X-Edge-Enrollment": enrollment["enrollment_token"], }), environ_overrides={"REMOTE_ADDR": "127.0.0.1"}, json={ "gateway_id": enrollment["gateway_id"], "environment": "staging", "network_zone": "zone-route", "policy_digest": "e" * 64, "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [], "version": "3.0.0", }) assert register.status_code == 201 assert register.headers["Cache-Control"] == "no-store" registered = register.get_json()["data"] bad = client.post( f"/api/datasource/edge/gateways/{registered['gateway_id']}/heartbeat", headers=mtls_headers(**{"X-Edge-Credential": registered["credential"] + "wrong"}), environ_overrides={"REMOTE_ADDR": "127.0.0.1"}, json={"environment": "staging", "network_zone": "zone-route", "generation": 1, "version": "3.0.0", "safe_summary": {}}, ) assert bad.status_code == 401 assert bad.headers["Cache-Control"] == "no-store" assert registered["credential"] not in bad.get_data(as_text=True) limited = None for _ in range(9): limited = client.post( f"/api/datasource/edge/gateways/{registered['gateway_id']}/heartbeat", headers=mtls_headers(**{"X-Edge-Credential": registered["credential"] + "wrong"}), environ_overrides={"REMOTE_ADDR": "127.0.0.1"}, json={"environment": "staging", "network_zone": "zone-route", "generation": 1, "version": "3.0.0", "safe_summary": {}}, ) assert limited is not None and limited.status_code == 429 assert limited.headers["Cache-Control"] == "no-store" with engine.connect() as connection: window = connection.execute(text("SELECT occurrence_count FROM edge_gateway_failure_windows")).scalar_one() failure_rows = connection.execute(text("SELECT count(*) FROM edge_gateway_audits WHERE success=false")).scalar_one() assert window == 10 assert failure_rows == 1 listed = client.get("/api/datasource/edge/gateways") assert listed.status_code == 200 assert listed.get_json()["data"]["gateways"][0]["gateway_id"] == registered["gateway_id"] assert client.get("/api/datasource/edge/gateways?limit=101").status_code == 400 oversized_reconcile = client.post( f"/api/datasource/edge/gateways/{registered['gateway_id']}/reconcile", headers=mtls_headers(**{"X-Edge-Credential": registered["credential"]}), environ_overrides={"REMOTE_ADDR": "127.0.0.1"}, json={"environment": "staging", "network_zone": "zone-route", "generation": 1, "limit": 101}, ) assert oversized_reconcile.status_code == 400 app.config["EDGE_GATEWAY_SIGNING_PRIVATE_KEY"] = "" unavailable_signer = client.post( f"/api/datasource/edge/gateways/{registered['gateway_id']}/releases", json={ "version": "7.0.0", "artifact_digest": "7" * 64, "rollback_version": "6.0.0", "request_id": "no-production-signer", }, ) assert unavailable_signer.status_code == 503 assert unavailable_signer.get_json()["error"]["code"] == "EDGE_GATEWAY_SIGNING_UNAVAILABLE"