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