test_phase3_wp04_edge_gateway_postgres.py 90 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950
  1. from __future__ import annotations
  2. import hashlib
  3. import os
  4. import subprocess
  5. import uuid
  6. from concurrent.futures import ThreadPoolExecutor
  7. from datetime import UTC, datetime, timedelta
  8. from pathlib import Path
  9. from threading import Barrier, Event
  10. from types import SimpleNamespace
  11. from urllib.parse import quote
  12. import pytest
  13. from cryptography import x509
  14. from cryptography.hazmat.primitives import hashes, serialization
  15. from cryptography.hazmat.primitives.asymmetric import rsa
  16. from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey
  17. from cryptography.x509.oid import NameOID
  18. from sqlalchemy import create_engine, inspect, text
  19. from sqlalchemy.engine import make_url
  20. from sqlalchemy.exc import DBAPIError
  21. from sqlalchemy.orm import Session
  22. from app.config.config import validate_production_database_identity
  23. from app.core.edge_gateway.contracts import (
  24. EdgeEventContract,
  25. SignedTaskEnvelope,
  26. canonical_json_bytes,
  27. canonical_sha256,
  28. stable_event_id,
  29. )
  30. from app.core.edge_gateway.repository import EdgeGatewayRepository
  31. from app.core.edge_gateway.service import (
  32. EdgeGatewayAuthenticationError,
  33. EdgeGatewayConfigurationError,
  34. EdgeGatewayConflictError,
  35. EdgeGatewayError,
  36. EdgeGatewayService,
  37. EdgeGatewayValidationError,
  38. )
  39. from app.edge_gateway.agent import EdgeAgent, EdgeRunnerAdapter, SignedReleaseManager
  40. from app.edge_gateway.bootstrap import EdgeBootstrapConfig
  41. ROOT = Path(__file__).resolve().parents[2]
  42. DATABASE_URL = os.environ.get("TEST_DATABASE_URL")
  43. pytestmark = pytest.mark.integration
  44. _EVIDENCE_TRIGGERS = {
  45. "edge_gateway_audits": (
  46. "trg_edge_audit_append_only",
  47. "trg_edge_audit_no_truncate",
  48. ),
  49. "edge_gateway_failure_windows": (
  50. "trg_edge_failure_window_protected",
  51. "trg_edge_failure_window_no_truncate",
  52. ),
  53. }
  54. _SIGNING_PRIVATE_KEY = Ed25519PrivateKey.generate()
  55. _SIGNING_KEY_ID = "edge-control-test-key-1"
  56. def _service(repository):
  57. return EdgeGatewayService(
  58. repository,
  59. signing_private_key=_SIGNING_PRIVATE_KEY,
  60. signing_key_id=_SIGNING_KEY_ID,
  61. )
  62. def _uid() -> str:
  63. return str(uuid.uuid4())
  64. def _alembic(command: str, target: str, database_url: str = None) -> None:
  65. assert DATABASE_URL
  66. result = subprocess.run(
  67. [str(ROOT / ".venv/bin/alembic"), "-c", str(ROOT / "alembic.ini"), command, target],
  68. cwd=ROOT,
  69. env={
  70. **os.environ,
  71. "MIGRATION_DATABASE_URL": database_url or DATABASE_URL,
  72. "SQLALCHEMY_DATABASE_URI": "",
  73. "DATABASE_URL": "",
  74. },
  75. check=False,
  76. capture_output=True,
  77. text=True,
  78. )
  79. if result.returncode:
  80. raise AssertionError(result.stderr)
  81. def _clear_edge_rows_for_test(engine) -> None:
  82. tables = set(inspect(engine).get_table_names(schema="public"))
  83. targets = []
  84. for table in (
  85. "edge_gateway_failure_windows", "edge_gateway_audits", "edge_gateway_events",
  86. "edge_gateway_releases", "edge_gateway_tasks", "edge_gateway_rotation_requests",
  87. "edge_gateway_credentials", "edge_gateways", "edge_gateway_enrollments",
  88. ):
  89. if table in tables:
  90. targets.append(f"public.{table}")
  91. if not targets:
  92. return
  93. with engine.begin() as connection:
  94. existing_triggers = set(connection.execute(text("""
  95. SELECT tgname FROM pg_trigger
  96. WHERE tgrelid IN (
  97. 'public.edge_gateway_audits'::regclass,
  98. to_regclass('public.edge_gateway_failure_windows')
  99. ) AND NOT tgisinternal
  100. """)).scalars()) if "edge_gateway_audits" in tables else set()
  101. disabled = []
  102. for table, triggers in _EVIDENCE_TRIGGERS.items():
  103. if table not in tables:
  104. continue
  105. for trigger in triggers:
  106. if trigger in existing_triggers:
  107. connection.execute(text(f"ALTER TABLE public.{table} DISABLE TRIGGER {trigger}"))
  108. disabled.append((table, trigger))
  109. connection.execute(text(f"TRUNCATE TABLE {','.join(targets)} CASCADE"))
  110. for table, trigger in disabled:
  111. connection.execute(text(f"ALTER TABLE public.{table} ENABLE TRIGGER {trigger}"))
  112. @pytest.fixture(scope="module")
  113. def engine(wp06_postgres_identities):
  114. if not DATABASE_URL:
  115. pytest.skip("TEST_DATABASE_URL is required")
  116. bootstrap = create_engine(DATABASE_URL, pool_pre_ping=True)
  117. _clear_edge_rows_for_test(bootstrap)
  118. bootstrap.dispose()
  119. _alembic("downgrade", "20260802_476")
  120. # WP06's trusted-delivery migrations explicitly reject a superuser or
  121. # CREATEROLE migrator. Keep the older WP04 exercise at the same current
  122. # head by using the isolated no-superuser migrator identity.
  123. _alembic("upgrade", "head", wp06_postgres_identities.migration_url)
  124. value = create_engine(DATABASE_URL, pool_pre_ping=True)
  125. yield value
  126. value.dispose()
  127. @pytest.fixture(autouse=True)
  128. def clear_edge_rows(engine):
  129. _clear_edge_rows_for_test(engine)
  130. yield
  131. _clear_edge_rows_for_test(engine)
  132. @pytest.fixture
  133. def service(engine):
  134. session = Session(engine)
  135. yield _service(EdgeGatewayRepository(session))
  136. session.close()
  137. def _enrollment(service, **overrides):
  138. values = {
  139. "gateway_name": "plant-a-edge",
  140. "environment": "staging",
  141. "network_zone": "manufacturing-zone-a",
  142. "policy_digest": "a" * 64,
  143. "allowed_control_hosts": ["control.example.test"],
  144. "allowed_proxy_hosts": ["proxy.example.test"],
  145. "expected_certificate_sha256": "b" * 64,
  146. "ttl_seconds": 600,
  147. "actor_uid": _uid(),
  148. }
  149. values.update(overrides)
  150. return service.create_enrollment(**values)
  151. def _register(service, enrollment, certificate_sha256="b" * 64):
  152. return service.register(
  153. enrollment_token=enrollment["enrollment_token"],
  154. certificate_sha256=certificate_sha256,
  155. gateway_id=enrollment["gateway_id"],
  156. environment="staging",
  157. network_zone="manufacturing-zone-a",
  158. policy_digest="a" * 64,
  159. allowed_control_hosts=["control.example.test"],
  160. allowed_proxy_hosts=["proxy.example.test"],
  161. version="3.0.0",
  162. )
  163. def _event(task, gateway, classification="statistics", payload=None):
  164. identity = {
  165. "task_id": task["task_id"], "gateway_id": gateway["gateway_id"],
  166. "environment": "staging", "network_zone": "manufacturing-zone-a",
  167. "purpose": "inventory", "classification": classification,
  168. "contract_version": 1,
  169. "occurred_at": datetime.now(UTC).isoformat().replace("+00:00", "Z"),
  170. "attempt": 1, "idempotency_key": "evt-key-1", "policy_digest": "a" * 64,
  171. "payload": payload or {"asset_count": 12},
  172. }
  173. return {"event_id": stable_event_id(identity), **identity}
  174. def test_migration_is_479_head_and_schema_is_constrained(engine):
  175. with engine.connect() as connection:
  176. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479"
  177. functions = connection.execute(text("""
  178. SELECT pg_get_function_identity_arguments(p.oid),r.rolname,p.prosecdef,p.proconfig
  179. FROM pg_proc p JOIN pg_namespace n ON n.oid=p.pronamespace
  180. JOIN pg_roles r ON r.oid=p.proowner
  181. WHERE n.nspname='public' AND p.proname='record_edge_gateway_failure'
  182. """)).all()
  183. constraints = connection.execute(text("""
  184. SELECT conname FROM pg_constraint
  185. WHERE conrelid IN ('public.edge_gateways'::regclass,
  186. 'public.edge_gateway_credentials'::regclass)
  187. """)).scalars().all()
  188. assert "ck_edge_gateway_status" in constraints
  189. assert "ck_edge_credential_status" in constraints
  190. assert "trg_edge_gateway_binding" in constraints
  191. assert "trg_edge_credential_binding" in constraints
  192. assert functions == [(
  193. "p_audit_uid uuid, p_fingerprint character, p_gateway uuid, p_event_type character varying, p_reason_code character varying",
  194. "dataops_edge_evidence_owner", True, ["search_path=pg_catalog"],
  195. )]
  196. enrollment_columns = {column["name"] for column in inspect(engine).get_columns("edge_gateway_enrollments", schema="public")}
  197. assert "expected_certificate_sha256" in enrollment_columns
  198. task_columns = {
  199. column["name"]
  200. for column in inspect(engine).get_columns("edge_gateway_tasks", schema="public")
  201. }
  202. release_columns = {
  203. column["name"]
  204. for column in inspect(engine).get_columns("edge_gateway_releases", schema="public")
  205. }
  206. assert "authority" in task_columns
  207. assert {
  208. "artifact_name", "deadline_at", "signature_algorithm", "key_id",
  209. "manifest_digest", "signature", "signed_manifest",
  210. } <= release_columns
  211. with pytest.raises(DBAPIError), engine.begin() as connection:
  212. connection.execute(text("""
  213. INSERT INTO edge_gateway_enrollments
  214. (uid,gateway_id,gateway_name,environment,network_zone,policy_digest,
  215. allowed_control_hosts,allowed_proxy_hosts,token_hash,status,expires_at,created_by)
  216. VALUES (CAST(:uid AS uuid),CAST(:gateway AS uuid),'invalid-host','staging','zone-a',
  217. :policy,ARRAY['control.example.test','CONTROL.example.test'],ARRAY[]::text[],
  218. :token,'pending',CURRENT_TIMESTAMP+INTERVAL '10 minutes',CAST(:actor AS uuid))
  219. """), {"uid": _uid(), "gateway": _uid(), "actor": _uid(), "policy": "1" * 64, "token": "2" * 64})
  220. def test_migration_controlled_477_downgrade_and_476_upgrade(engine):
  221. try:
  222. _alembic("downgrade", "20260802_476")
  223. with engine.connect() as connection:
  224. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260802_476"
  225. assert not inspect(connection).has_table("edge_gateways", schema="public")
  226. finally:
  227. _alembic("upgrade", "head")
  228. with engine.connect() as connection:
  229. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479"
  230. def test_479_downgrade_and_upgrade_are_reversible_without_signed_rows(engine):
  231. try:
  232. _alembic("downgrade", "20260809_478")
  233. with engine.connect() as connection:
  234. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_478"
  235. assert "authority" not in {
  236. column["name"]
  237. for column in inspect(connection).get_columns("edge_gateway_tasks", schema="public")
  238. }
  239. finally:
  240. _alembic("upgrade", "head")
  241. with engine.connect() as connection:
  242. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479"
  243. def test_478_downgrade_never_restores_unsafe_function(engine):
  244. try:
  245. _alembic("downgrade", "20260809_477")
  246. with engine.connect() as connection:
  247. signatures = connection.execute(text("""
  248. SELECT pg_get_function_identity_arguments(p.oid)
  249. FROM pg_proc p JOIN pg_namespace n ON n.oid=p.pronamespace
  250. WHERE n.nspname='public' AND p.proname='record_edge_gateway_failure'
  251. ORDER BY 1
  252. """)).scalars().all()
  253. assert signatures == [
  254. "p_audit_uid uuid, p_fingerprint character, p_gateway uuid, p_event_type character varying, p_reason_code character varying"
  255. ]
  256. finally:
  257. _alembic("upgrade", "head")
  258. def test_migration_rejects_downgrade_with_task2_data(service, engine):
  259. _enrollment(service)
  260. with pytest.raises(subprocess.CalledProcessError):
  261. _alembic("downgrade", "20260802_476")
  262. with engine.connect() as connection:
  263. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479"
  264. _clear_edge_rows_for_test(engine)
  265. with engine.begin() as connection:
  266. connection.execute(text("""
  267. SELECT public.record_edge_gateway_audit(
  268. CAST(:uid AS uuid),NULL,'enrollment_created',NULL,FALSE,
  269. '{"reason_code":"controlled_downgrade_probe"}'::jsonb
  270. )
  271. """), {"uid": _uid()})
  272. with pytest.raises(subprocess.CalledProcessError):
  273. _alembic("downgrade", "20260802_476")
  274. with engine.connect() as connection:
  275. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479"
  276. def test_enrollment_registration_hash_only_replay_and_failed_binding(service, engine):
  277. enrollment = _enrollment(service)
  278. assert enrollment["returned_once"] is True
  279. token = enrollment["enrollment_token"]
  280. with engine.connect() as connection:
  281. row = connection.execute(text("SELECT token_hash, consumed_at FROM edge_gateway_enrollments")).mappings().one()
  282. assert row["token_hash"] == hashlib.sha256(token.encode()).hexdigest()
  283. assert token not in str(row)
  284. with pytest.raises(DBAPIError), engine.begin() as connection:
  285. connection.execute(text("""
  286. UPDATE edge_gateway_enrollments SET expected_certificate_sha256=:certificate
  287. WHERE gateway_id=CAST(:gateway AS uuid)
  288. """), {"certificate": "c" * 64, "gateway": enrollment["gateway_id"]})
  289. with pytest.raises(EdgeGatewayAuthenticationError):
  290. _register(service, enrollment, certificate_sha256="c" * 64)
  291. with engine.connect() as connection:
  292. assert connection.execute(text("SELECT consumed_at FROM edge_gateway_enrollments")).scalar_one() is None
  293. failures = connection.execute(text("SELECT safe_detail FROM edge_gateway_audits WHERE success=false")).scalars().all()
  294. assert failures and all(token not in str(detail) for detail in failures)
  295. with pytest.raises(EdgeGatewayAuthenticationError):
  296. service.register(
  297. enrollment_token=token, certificate_sha256="b" * 64,
  298. gateway_id=enrollment["gateway_id"], environment="production",
  299. network_zone="manufacturing-zone-a", policy_digest="a" * 64,
  300. allowed_control_hosts=["control.example.test"],
  301. allowed_proxy_hosts=["proxy.example.test"], version="3.0.0",
  302. )
  303. with engine.connect() as connection:
  304. assert connection.execute(text("SELECT consumed_at FROM edge_gateway_enrollments")).scalar_one() is None
  305. gateway = _register(service, enrollment)
  306. assert gateway["credential_returned_once"] is True
  307. with engine.connect() as connection:
  308. credential = connection.execute(text("SELECT credential_hash, certificate_sha256, generation FROM edge_gateway_credentials")).mappings().one()
  309. assert gateway["credential"] not in str(credential)
  310. assert credential["certificate_sha256"] == "b" * 64
  311. assert credential["generation"] == 1
  312. with pytest.raises(EdgeGatewayAuthenticationError):
  313. _register(service, enrollment)
  314. def test_registration_failures_are_windowed_rate_limited_and_do_not_consume(service, engine):
  315. enrollment = _enrollment(service)
  316. for _ in range(9):
  317. with pytest.raises(EdgeGatewayAuthenticationError):
  318. _register(service, enrollment, certificate_sha256="c" * 64)
  319. with pytest.raises(EdgeGatewayError, match="rate limit") as limited:
  320. _register(service, enrollment, certificate_sha256="c" * 64)
  321. assert limited.value.status_code == 429
  322. with engine.connect() as connection:
  323. failure = connection.execute(text("""
  324. SELECT occurrence_count FROM edge_gateway_failure_windows
  325. WHERE event_type='untrusted_request_rejected'
  326. """)).scalar_one()
  327. audit = connection.execute(text("""
  328. SELECT count(*),string_agg(safe_detail::text,'')
  329. FROM edge_gateway_audits WHERE event_type='untrusted_request_rejected'
  330. """)).one()
  331. consumed_at = connection.execute(text("""
  332. SELECT consumed_at FROM edge_gateway_enrollments
  333. WHERE gateway_id=CAST(:gateway_id AS uuid)
  334. """), {"gateway_id": enrollment["gateway_id"]}).scalar_one()
  335. assert failure == 10
  336. assert audit[0] == 1
  337. assert enrollment["enrollment_token"] not in (audit[1] or "")
  338. assert consumed_at is None
  339. assert _register(service, enrollment)["generation"] == 1
  340. def test_mixed_unauthenticated_failures_share_one_bounded_unknown_bucket(service, engine):
  341. limited = None
  342. for attempt in range(10):
  343. gateway_id = _uid()
  344. try:
  345. if attempt % 2:
  346. service.register(
  347. enrollment_token=f"dope_invalid_{attempt}",
  348. certificate_sha256="b" * 64, gateway_id=gateway_id,
  349. environment="staging", network_zone=f"random-zone-{attempt}",
  350. policy_digest="a" * 64,
  351. allowed_control_hosts=["control.example.test"],
  352. allowed_proxy_hosts=[], version="3.0.0",
  353. )
  354. else:
  355. service.authenticate(
  356. credential=f"dopg_invalid_{attempt}", certificate_sha256="b" * 64,
  357. gateway_id=gateway_id, environment="staging",
  358. network_zone=f"random-zone-{attempt}", generation=1,
  359. )
  360. except EdgeGatewayError as exc:
  361. limited = exc
  362. assert limited is not None and limited.status_code == 429
  363. with engine.connect() as connection:
  364. windows = connection.execute(text("""
  365. SELECT failure_fingerprint,gateway_id,event_type,occurrence_count
  366. FROM edge_gateway_failure_windows
  367. """)).mappings().all()
  368. audits = connection.execute(text("""
  369. SELECT gateway_id,event_type,safe_detail FROM edge_gateway_audits WHERE success=false
  370. """)).mappings().all()
  371. assert len(windows) == 1 and windows[0]["occurrence_count"] == 10
  372. assert windows[0]["gateway_id"] is None
  373. assert windows[0]["event_type"] == "untrusted_request_rejected"
  374. assert len(audits) == 1 and audits[0]["gateway_id"] is None
  375. assert audits[0]["safe_detail"] == {"reason_code": "untrusted_failure"}
  376. def test_database_rejects_pending_enrollment_and_status_only_gateway_activation(service, engine):
  377. enrollment = _enrollment(service)
  378. values = {
  379. "gateway_id": enrollment["gateway_id"], "policy": "a" * 64,
  380. }
  381. with pytest.raises(DBAPIError), engine.begin() as connection:
  382. connection.execute(text("""
  383. INSERT INTO edge_gateways
  384. (uid,name,environment,network_zone,status,current_generation,policy_digest,
  385. allowed_control_hosts,allowed_proxy_hosts,runtime_version)
  386. VALUES (CAST(:gateway_id AS uuid),'plant-a-edge','staging','manufacturing-zone-a',
  387. 'active',1,:policy,ARRAY['control.example.test'],
  388. ARRAY['proxy.example.test'],'3.0.0')
  389. """), values)
  390. with engine.begin() as connection:
  391. connection.execute(text("""
  392. INSERT INTO edge_gateways
  393. (uid,name,environment,network_zone,status,current_generation,policy_digest,
  394. allowed_control_hosts,allowed_proxy_hosts,runtime_version)
  395. VALUES (CAST(:gateway_id AS uuid),'plant-a-edge','staging','manufacturing-zone-a',
  396. 'pending',1,:policy,ARRAY['control.example.test'],
  397. ARRAY['proxy.example.test'],'3.0.0')
  398. """), values)
  399. with pytest.raises(DBAPIError), engine.begin() as connection:
  400. connection.execute(text("""
  401. UPDATE edge_gateways SET status='active'
  402. WHERE uid=CAST(:gateway_id AS uuid)
  403. """), values)
  404. def test_authentication_rotation_revocation_heartbeat_and_active_binding(service, engine):
  405. gateway = _register(service, _enrollment(service))
  406. identity = service.authenticate(
  407. credential=gateway["credential"], certificate_sha256="b" * 64,
  408. gateway_id=gateway["gateway_id"], environment="staging",
  409. network_zone="manufacturing-zone-a", generation=1,
  410. )
  411. assert identity["gateway_id"] == gateway["gateway_id"]
  412. for statement, parameters in (
  413. ("UPDATE edge_gateway_credentials SET certificate_sha256=:value WHERE gateway_id=CAST(:gateway AS uuid)", {"value": "c" * 64}),
  414. ("UPDATE edge_gateway_credentials SET credential_hash=:value WHERE gateway_id=CAST(:gateway AS uuid)", {"value": "d" * 64}),
  415. ("UPDATE edge_gateways SET name='attacker' WHERE uid=CAST(:gateway AS uuid)", {}),
  416. ("UPDATE edge_gateways SET environment='production' WHERE uid=CAST(:gateway AS uuid)", {}),
  417. ("UPDATE edge_gateways SET network_zone='attacker-zone' WHERE uid=CAST(:gateway AS uuid)", {}),
  418. ("UPDATE edge_gateways SET policy_digest=:value WHERE uid=CAST(:gateway AS uuid)", {"value": "c" * 64}),
  419. ("UPDATE edge_gateways SET allowed_control_hosts=ARRAY['attacker.test'] WHERE uid=CAST(:gateway AS uuid)", {}),
  420. ("UPDATE edge_gateways SET allowed_proxy_hosts=ARRAY['attacker.test'] WHERE uid=CAST(:gateway AS uuid)", {}),
  421. ("""UPDATE edge_gateway_credentials SET generation=2,certificate_sha256=:certificate,
  422. credential_hash=:credential WHERE gateway_id=CAST(:gateway AS uuid);
  423. UPDATE edge_gateways SET current_generation=2 WHERE uid=CAST(:gateway AS uuid)""",
  424. {"certificate": "c" * 64, "credential": "d" * 64}),
  425. ):
  426. with pytest.raises(DBAPIError), engine.begin() as connection:
  427. connection.execute(text(statement), {"gateway": gateway["gateway_id"], **parameters})
  428. with pytest.raises(DBAPIError), engine.begin() as connection:
  429. connection.execute(text("""
  430. UPDATE edge_gateway_credentials SET status='rotated',revoked_at=CURRENT_TIMESTAMP
  431. WHERE gateway_id=CAST(:gateway_id AS uuid) AND generation=1
  432. """), {"gateway_id": gateway["gateway_id"]})
  433. with pytest.raises(EdgeGatewayAuthenticationError):
  434. service.authenticate(
  435. credential=gateway["credential"], certificate_sha256="c" * 64,
  436. gateway_id=gateway["gateway_id"], environment="staging",
  437. network_zone="manufacturing-zone-a", generation=1,
  438. )
  439. with pytest.raises(EdgeGatewayValidationError):
  440. service.heartbeat(
  441. credential=gateway["credential"], certificate_sha256="b" * 64,
  442. gateway_id=gateway["gateway_id"], environment="staging",
  443. network_zone="manufacturing-zone-a", generation=1, version="3.0.0",
  444. safe_summary={"password": "must-not-persist"},
  445. )
  446. with engine.connect() as connection:
  447. assert connection.execute(text("SELECT last_heartbeat_at FROM edge_gateways")).scalar_one() is None
  448. rotation_actor = _uid()
  449. rotated = service.rotate(
  450. gateway["gateway_id"], certificate_sha256="c" * 64,
  451. request_id="rotation-request-1", actor_uid=rotation_actor,
  452. )
  453. rotation_replay = service.rotate(
  454. gateway["gateway_id"], certificate_sha256="c" * 64,
  455. request_id="rotation-request-1", actor_uid=rotation_actor,
  456. )
  457. assert rotation_replay == {
  458. "gateway_id": gateway["gateway_id"], "generation": 2,
  459. "credential": None, "credential_returned_once": False, "replayed": True,
  460. }
  461. with pytest.raises(EdgeGatewayConflictError):
  462. service.rotate(
  463. gateway["gateway_id"], certificate_sha256="d" * 64,
  464. request_id="rotation-request-1", actor_uid=rotation_actor,
  465. )
  466. with pytest.raises(EdgeGatewayAuthenticationError):
  467. service.authenticate(
  468. credential=gateway["credential"], certificate_sha256="b" * 64,
  469. gateway_id=gateway["gateway_id"], environment="staging",
  470. network_zone="manufacturing-zone-a", generation=1,
  471. )
  472. assert service.heartbeat(
  473. credential=rotated["credential"], certificate_sha256="c" * 64,
  474. gateway_id=gateway["gateway_id"], environment="staging",
  475. network_zone="manufacturing-zone-a", generation=2, version="3.0.1",
  476. safe_summary={"queue_depth": 1},
  477. )["status"] == "online"
  478. assert service.revoke(gateway["gateway_id"], actor_uid=_uid()) is True
  479. with pytest.raises(DBAPIError), engine.begin() as connection:
  480. connection.execute(text("UPDATE edge_gateways SET status='active' WHERE uid=CAST(:gateway AS uuid)"), {"gateway": gateway["gateway_id"]})
  481. with pytest.raises(EdgeGatewayAuthenticationError):
  482. service.authenticate(
  483. credential=rotated["credential"], certificate_sha256="c" * 64,
  484. gateway_id=gateway["gateway_id"], environment="staging",
  485. network_zone="manufacturing-zone-a", generation=2,
  486. )
  487. with engine.connect() as connection:
  488. assert "untrusted_request_rejected" in set(connection.execute(text(
  489. "SELECT event_type FROM edge_gateway_audits WHERE success=false"
  490. )).scalars())
  491. def test_database_rejects_second_active_credential_combination_injection(service, engine):
  492. gateway = _register(service, _enrollment(service))
  493. with pytest.raises(DBAPIError), engine.begin() as connection:
  494. connection.execute(text("""
  495. INSERT INTO edge_gateway_credentials
  496. (uid,gateway_id,generation,credential_hash,certificate_sha256,status,expires_at)
  497. VALUES (CAST(:uid AS uuid),CAST(:gateway AS uuid),2,:credential,:certificate,
  498. 'active',CURRENT_TIMESTAMP+INTERVAL '365 days');
  499. UPDATE edge_gateways SET current_generation=2 WHERE uid=CAST(:gateway AS uuid)
  500. """), {
  501. "uid": _uid(), "gateway": gateway["gateway_id"],
  502. "credential": "9" * 64, "certificate": "8" * 64,
  503. })
  504. with engine.connect() as connection:
  505. row = connection.execute(text("""
  506. SELECT g.current_generation,count(*) FILTER (WHERE c.status='active') AS active_count
  507. FROM edge_gateways g JOIN edge_gateway_credentials c ON c.gateway_id=g.uid
  508. GROUP BY g.current_generation
  509. """)).mappings().one()
  510. assert row == {"current_generation": 1, "active_count": 1}
  511. def test_failure_evidence_tables_reject_direct_mutation_and_truncation(service, engine):
  512. with pytest.raises(EdgeGatewayAuthenticationError):
  513. service.authenticate(
  514. credential="dopg_invalid", certificate_sha256="b" * 64,
  515. gateway_id=_uid(), environment="staging",
  516. network_zone="zone-a", generation=1,
  517. )
  518. with pytest.raises(DBAPIError), engine.begin() as connection:
  519. connection.execute(text("TRUNCATE TABLE edge_gateway_audits"))
  520. for statement in (
  521. "UPDATE edge_gateway_failure_windows SET occurrence_count=occurrence_count+1,last_seen_at=CURRENT_TIMESTAMP",
  522. "UPDATE edge_gateway_failure_windows SET occurrence_count=occurrence_count+5",
  523. "UPDATE edge_gateway_failure_windows SET last_seen_at=last_seen_at-INTERVAL '1 second'",
  524. "UPDATE edge_gateway_failure_windows SET first_seen_at=first_seen_at+INTERVAL '1 second'",
  525. "UPDATE edge_gateway_failure_windows SET failure_fingerprint=:fingerprint",
  526. "DELETE FROM edge_gateway_failure_windows",
  527. "TRUNCATE TABLE edge_gateway_failure_windows",
  528. ):
  529. with pytest.raises(DBAPIError), engine.begin() as connection:
  530. connection.execute(text(statement), {"fingerprint": "f" * 64})
  531. with engine.connect() as connection:
  532. evidence = connection.execute(text("""
  533. SELECT occurrence_count,first_seen_at<=last_seen_at AS monotonic
  534. FROM edge_gateway_failure_windows
  535. """)).one()
  536. assert evidence == (1, True)
  537. assert connection.execute(text("SELECT count(*) FROM edge_gateway_audits")).scalar_one() == 1
  538. def test_task_pull_cancel_event_idempotency_policy_and_release(service, engine):
  539. gateway = _register(service, _enrollment(service))
  540. auth = {
  541. "credential": gateway["credential"], "certificate_sha256": "b" * 64,
  542. "gateway_id": gateway["gateway_id"], "environment": "staging",
  543. "network_zone": "manufacturing-zone-a", "generation": 1,
  544. }
  545. task = service.issue_task({
  546. "task_id": "task-1", "gateway_id": gateway["gateway_id"],
  547. "environment": "staging", "network_zone": "manufacturing-zone-a",
  548. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  549. "contract_version": 1,
  550. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  551. "attempt": 1, "idempotency_key": "task-key-1", "policy_digest": "a" * 64,
  552. }, actor_uid=_uid())
  553. conflicting_task = {
  554. "task_id": "task-conflicting", "gateway_id": gateway["gateway_id"],
  555. "environment": "staging", "network_zone": "manufacturing-zone-a",
  556. "purpose": "different-purpose", "classification": "statistics", "task_type": "profile",
  557. "contract_version": 1,
  558. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  559. "attempt": 1, "idempotency_key": "task-key-1", "policy_digest": "a" * 64,
  560. }
  561. with pytest.raises(EdgeGatewayConflictError):
  562. service.issue_task(conflicting_task, actor_uid=_uid())
  563. with pytest.raises(EdgeGatewayValidationError):
  564. service.issue_task({
  565. "task_id": "task-wrong-binding", "gateway_id": gateway["gateway_id"],
  566. "environment": "production", "network_zone": "manufacturing-zone-a",
  567. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  568. "contract_version": 1,
  569. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  570. "attempt": 1, "idempotency_key": "task-wrong-binding-key", "policy_digest": "a" * 64,
  571. }, actor_uid=_uid())
  572. pulled = service.pull_task(**auth)
  573. assert pulled["task"]["task_id"] == task["task_id"]
  574. approved = _event(task, gateway)
  575. with pytest.raises(EdgeGatewayAuthenticationError):
  576. service.accept_event(event=approved, lease_token="dopl_wrong", **auth)
  577. ack1 = service.accept_event(event=approved, lease_token=pulled["lease_token"], **auth)
  578. ack2 = service.accept_event(event=approved, lease_token=pulled["lease_token"], **auth)
  579. assert ack1 == ack2
  580. assert set(ack1) == {
  581. "event_id", "event_digest", "remote_lease_digest", "received_at", "status",
  582. }
  583. assert ack1["event_digest"] == EdgeEventContract.from_mapping(approved).digest
  584. assert ack1["remote_lease_digest"] == hashlib.sha256(
  585. pulled["lease_token"].encode("utf-8")
  586. ).hexdigest()
  587. assert pulled["lease_token"] not in str(ack1)
  588. changed = dict(approved)
  589. changed["payload"] = {"asset_count": 13}
  590. with pytest.raises((EdgeGatewayConflictError, EdgeGatewayValidationError)):
  591. service.accept_event(event=changed, **auth)
  592. raw = dict(approved)
  593. raw.update(classification="raw", payload={"password": "super-secret-value"}, event_id="evt_" + "0" * 64)
  594. with pytest.raises(EdgeGatewayValidationError):
  595. service.accept_event(event=raw, **auth)
  596. wrong_purpose = dict(approved)
  597. wrong_purpose["purpose"] = "different-purpose"
  598. identity = {key: value for key, value in wrong_purpose.items() if key != "event_id"}
  599. wrong_purpose["event_id"] = stable_event_id(identity)
  600. with pytest.raises((EdgeGatewayValidationError, EdgeGatewayAuthenticationError)):
  601. service.accept_event(event=wrong_purpose, lease_token=pulled["lease_token"], **auth)
  602. assert service.cancel_task(task["task_id"], actor_uid=_uid()) is False
  603. assert service.reconcile(**auth)["cancelled_task_ids"] == []
  604. assert service.accept_event(
  605. event=approved, lease_token=pulled["lease_token"], **auth,
  606. ) == ack1
  607. with engine.connect() as connection:
  608. assert connection.execute(text("SELECT count(*) FROM edge_gateway_events")).scalar_one() == 1
  609. release = service.offer_release(
  610. gateway["gateway_id"], version="3.1.0", artifact_digest="d" * 64,
  611. rollback_version="3.0.0", request_id="release-request-1", actor_uid=_uid(),
  612. )
  613. assert service.offer_release(
  614. gateway["gateway_id"], version="3.1.0", artifact_digest="d" * 64,
  615. rollback_version="3.0.0", request_id="release-request-1", actor_uid=_uid(),
  616. ) == release
  617. with pytest.raises(EdgeGatewayConflictError):
  618. service.offer_release(
  619. gateway["gateway_id"], version="3.1.1", artifact_digest="d" * 64,
  620. rollback_version="3.0.0", request_id="release-request-1", actor_uid=_uid(),
  621. )
  622. with pytest.raises(EdgeGatewayConflictError):
  623. service.acknowledge_release(
  624. release_id=release["release_id"], outcome="rollback",
  625. safe_summary={"reason_code": "health_check_failed"}, **auth,
  626. )
  627. accepted = service.acknowledge_release(
  628. release_id=release["release_id"], outcome="accepted",
  629. safe_summary={"reason_code": "accepted"}, **auth,
  630. )
  631. pending = service.reconcile(**auth)["release_offers"]
  632. assert [(item["release_id"], item["status"]) for item in pending] == [(release["release_id"], "accepted")]
  633. assert service.acknowledge_release(
  634. release_id=release["release_id"], outcome="accepted",
  635. safe_summary={"reason_code": "accepted"}, **auth,
  636. ) == accepted
  637. service.acknowledge_release(
  638. release_id=release["release_id"], outcome="failed",
  639. safe_summary={"reason_code": "health_check_failed"}, **auth,
  640. )
  641. acked = service.acknowledge_release(
  642. release_id=release["release_id"], outcome="rollback",
  643. safe_summary={"reason_code": "health_check_failed"}, **auth,
  644. )
  645. assert acked["status"] == "rolled_back"
  646. installed_release = service.offer_release(
  647. gateway["gateway_id"], version="3.2.0", artifact_digest="e" * 64,
  648. rollback_version="3.1.0", request_id="release-request-2", actor_uid=_uid(),
  649. )
  650. service.acknowledge_release(
  651. release_id=installed_release["release_id"], outcome="accepted",
  652. safe_summary={"reason_code": "accepted"}, **auth,
  653. )
  654. service.acknowledge_release(
  655. release_id=installed_release["release_id"], outcome="installed",
  656. safe_summary={"reason_code": "installed"}, **auth,
  657. )
  658. installed_rollback = service.acknowledge_release(
  659. release_id=installed_release["release_id"], outcome="rollback",
  660. safe_summary={"reason_code": "post_install_health_failed"}, **auth,
  661. )
  662. assert installed_rollback["status"] == "rolled_back"
  663. assert service.acknowledge_release(
  664. release_id=installed_release["release_id"], outcome="rollback",
  665. safe_summary={"reason_code": "post_install_health_failed"}, **auth,
  666. ) == installed_rollback
  667. with engine.connect() as connection:
  668. audit_rows = connection.execute(text("SELECT event_type,safe_detail FROM edge_gateway_audits")).mappings().all()
  669. details = [row["safe_detail"] for row in audit_rows]
  670. assert all("credential" not in str(detail).lower() and "token" not in str(detail).lower() for detail in details)
  671. assert all("super-secret-value" not in str(detail) for detail in details)
  672. assert {"idempotency_conflict", "binding_rejected", "policy_rejected", "release_rejected"} <= {
  673. row["event_type"] for row in audit_rows
  674. }
  675. with pytest.raises(DBAPIError), engine.begin() as connection:
  676. connection.execute(text("UPDATE edge_gateway_audits SET success=true"))
  677. with pytest.raises(DBAPIError), engine.begin() as connection:
  678. connection.execute(text("DELETE FROM edge_gateway_audits"))
  679. task_failed = service.issue_task({
  680. "task_id": "task-failed", "gateway_id": gateway["gateway_id"],
  681. "environment": "staging", "network_zone": "manufacturing-zone-a",
  682. "purpose": "inventory", "classification": "statistics", "task_type": "quality",
  683. "contract_version": 1,
  684. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  685. "attempt": 1, "idempotency_key": "task-failed-key", "policy_digest": "a" * 64,
  686. }, actor_uid=_uid())
  687. failed_lease = service.pull_task(**auth)
  688. failed = service.task_outcome(
  689. task_failed["task_id"], outcome="failed", lease_token=failed_lease["lease_token"],
  690. safe_summary={"reason_code": "controlled_failure"}, **auth,
  691. )
  692. assert failed["status"] == "failed"
  693. assert service.task_outcome(
  694. task_failed["task_id"], outcome="failed", lease_token=failed_lease["lease_token"],
  695. safe_summary={"reason_code": "controlled_failure"}, **auth,
  696. )["replayed"] is True
  697. task_cancel = service.issue_task({
  698. "task_id": "task-cancelled", "gateway_id": gateway["gateway_id"],
  699. "environment": "staging", "network_zone": "manufacturing-zone-a",
  700. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  701. "contract_version": 1,
  702. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  703. "attempt": 1, "idempotency_key": "task-cancelled-key", "policy_digest": "a" * 64,
  704. }, actor_uid=_uid())
  705. cancel_lease = service.pull_task(**auth)
  706. assert service.cancel_task(task_cancel["task_id"], actor_uid=_uid()) is True
  707. assert service.task_outcome(
  708. task_cancel["task_id"], outcome="cancelled", lease_token=cancel_lease["lease_token"],
  709. safe_summary={"reason_code": "cancel_acknowledged"}, **auth,
  710. )["status"] == "cancelled"
  711. def test_control_plane_signs_task_authority_and_complete_release_state(service, engine):
  712. gateway = _register(service, _enrollment(service))
  713. auth = {
  714. "credential": gateway["credential"],
  715. "certificate_sha256": "b" * 64,
  716. "gateway_id": gateway["gateway_id"],
  717. "environment": "staging",
  718. "network_zone": "manufacturing-zone-a",
  719. "generation": 1,
  720. }
  721. deadline = (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z")
  722. task = service.issue_task({
  723. "task_id": "task-signed-authority", "gateway_id": gateway["gateway_id"],
  724. "environment": "staging", "network_zone": "manufacturing-zone-a",
  725. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  726. "contract_version": 1, "deadline_at": deadline, "attempt": 1,
  727. "idempotency_key": "task-signed-authority-key", "policy_digest": "a" * 64,
  728. }, actor_uid=_uid())
  729. pulled = service.pull_task(**auth)
  730. envelope = pulled["signed_task_envelope"]
  731. assert pulled["task"] == envelope["task"]
  732. assert envelope["task"]["task_id"] == task["task_id"]
  733. assert set(envelope) == {
  734. "task", "authority_key_id", "signature_algorithm", "contract_digest",
  735. "gateway_id", "environment",
  736. "network_zone", "policy_digest", "purpose", "issued_at", "expires_at",
  737. "signature",
  738. }
  739. assert envelope["signature_algorithm"] == "Ed25519"
  740. assert envelope["authority_key_id"] == _SIGNING_KEY_ID
  741. assert envelope["contract_digest"] == canonical_sha256(envelope["task"])
  742. unsigned_authority = {key: value for key, value in envelope.items() if key != "signature"}
  743. _SIGNING_PRIVATE_KEY.public_key().verify(
  744. bytes.fromhex(envelope["signature"]),
  745. SignedTaskEnvelope.canonical_unsigned_bytes(unsigned_authority),
  746. )
  747. release_deadline = (datetime.now(UTC) + timedelta(hours=1)).isoformat().replace("+00:00", "Z")
  748. release = service.offer_release(
  749. gateway["gateway_id"], version="6.0.0", artifact_digest="6" * 64,
  750. artifact_name="edge-agent-6.0.0.bin", deadline_at=release_deadline,
  751. rollback_version="3.0.0", request_id="signed-release-1", actor_uid=_uid(),
  752. )
  753. manifest = release["manifest"]
  754. with pytest.raises(EdgeGatewayValidationError):
  755. service.offer_release(
  756. gateway["gateway_id"], version="release-six", artifact_digest="6" * 64,
  757. rollback_version="3.0.0", request_id="invalid-release-version",
  758. actor_uid=_uid(),
  759. )
  760. with pytest.raises(EdgeGatewayValidationError):
  761. service.offer_release(
  762. gateway["gateway_id"], version="2.9.0", artifact_digest="6" * 64,
  763. rollback_version="3.0.0", request_id="release-downgrade",
  764. actor_uid=_uid(),
  765. )
  766. assert set(manifest) == {
  767. "release_id", "version", "rollback_version", "artifact_digest", "artifact_name",
  768. "deadline_at", "signature_algorithm", "key_id", "manifest_digest", "signature",
  769. "status",
  770. }
  771. assert manifest["status"] == "offered"
  772. unsigned_manifest = {
  773. key: value for key, value in manifest.items()
  774. if key not in {"manifest_digest", "signature"}
  775. }
  776. assert manifest["manifest_digest"] == canonical_sha256(unsigned_manifest)
  777. _SIGNING_PRIVATE_KEY.public_key().verify(
  778. bytes.fromhex(manifest["signature"]), canonical_json_bytes(unsigned_manifest),
  779. )
  780. for outcome, expected in (("accepted", "accepted"), ("failed", "failed"), ("rollback", "rolled_back")):
  781. service.acknowledge_release(
  782. release_id=release["release_id"], outcome=outcome,
  783. safe_summary={"release_status": expected, "version": "3.0.0"}, **auth,
  784. )
  785. response = service.reconcile(**auth)
  786. candidates = list(response["release_offers"])
  787. if response["release_baseline"] is not None:
  788. candidates.append(response["release_baseline"])
  789. reconciled = {
  790. item["release_id"]: item for item in candidates
  791. }[release["release_id"]]
  792. assert reconciled["status"] == expected
  793. unsigned = {
  794. key: value for key, value in reconciled.items()
  795. if key not in {"manifest_digest", "signature"}
  796. }
  797. assert reconciled["manifest_digest"] == canonical_sha256(unsigned)
  798. _SIGNING_PRIVATE_KEY.public_key().verify(
  799. bytes.fromhex(reconciled["signature"]), canonical_json_bytes(unsigned),
  800. )
  801. with engine.connect() as connection:
  802. task_row = connection.execute(text("""
  803. SELECT authority,contract FROM edge_gateway_tasks WHERE task_id='task-signed-authority'
  804. """)).mappings().one()
  805. release_row = connection.execute(text("""
  806. SELECT signed_manifest FROM edge_gateway_releases WHERE uid=CAST(:uid AS uuid)
  807. """), {"uid": release["release_id"]}).scalar_one()
  808. assert dict(task_row["authority"]) == envelope
  809. assert dict(release_row)["status"] == "rolled_back"
  810. persisted = f"{task_row['authority']}{release_row}"
  811. assert "PRIVATE KEY" not in persisted and "signing_private_key" not in persisted
  812. @pytest.mark.parametrize("seconds", [30, 119, 120])
  813. def test_pull_lease_never_extends_beyond_task_deadline(service, seconds):
  814. gateway = _register(service, _enrollment(service))
  815. auth = {
  816. "credential": gateway["credential"],
  817. "certificate_sha256": "b" * 64,
  818. "gateway_id": gateway["gateway_id"],
  819. "environment": "staging",
  820. "network_zone": "manufacturing-zone-a",
  821. "generation": 1,
  822. }
  823. deadline = (datetime.now(UTC) + timedelta(seconds=seconds)).isoformat().replace(
  824. "+00:00", "Z"
  825. )
  826. service.issue_task(
  827. {
  828. "task_id": f"deadline-boundary-{seconds}",
  829. "gateway_id": gateway["gateway_id"],
  830. "environment": "staging",
  831. "network_zone": "manufacturing-zone-a",
  832. "purpose": "inventory",
  833. "classification": "statistics",
  834. "task_type": "profile",
  835. "contract_version": 1,
  836. "deadline_at": deadline,
  837. "attempt": 1,
  838. "idempotency_key": f"deadline-boundary-{seconds}-key",
  839. "policy_digest": "a" * 64,
  840. },
  841. actor_uid=_uid(),
  842. )
  843. pulled = service.pull_task(**auth)
  844. assert datetime.fromisoformat(pulled["lease_expires_at"]) <= datetime.fromisoformat(
  845. deadline.replace("Z", "+00:00")
  846. )
  847. def test_cancel_reconcile_uses_stable_cursor_so_more_than_fifty_never_starve(
  848. service,
  849. ):
  850. gateway = _register(service, _enrollment(service))
  851. auth = {
  852. "credential": gateway["credential"],
  853. "certificate_sha256": "b" * 64,
  854. "gateway_id": gateway["gateway_id"],
  855. "environment": "staging",
  856. "network_zone": "manufacturing-zone-a",
  857. "generation": 1,
  858. }
  859. expected = []
  860. deadline = (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace(
  861. "+00:00", "Z"
  862. )
  863. for index in range(52):
  864. task_id = f"cancel-page-{index:02d}"
  865. service.issue_task(
  866. {
  867. "task_id": task_id,
  868. "gateway_id": gateway["gateway_id"],
  869. "environment": "staging",
  870. "network_zone": "manufacturing-zone-a",
  871. "purpose": "inventory",
  872. "classification": "statistics",
  873. "task_type": "profile",
  874. "contract_version": 1,
  875. "deadline_at": deadline,
  876. "attempt": 1,
  877. "idempotency_key": f"cancel-page-key-{index:02d}",
  878. "policy_digest": "a" * 64,
  879. },
  880. actor_uid=_uid(),
  881. )
  882. assert service.pull_task(**auth)["task"]["task_id"] == task_id
  883. assert service.cancel_task(task_id, actor_uid=_uid()) is True
  884. expected.append(task_id)
  885. first = service.reconcile(limit=50, cancel_cursor=None, release_cursor=None, **auth)
  886. assert first["cancelled_task_ids"] == expected[:50]
  887. assert first["cancel_next_cursor"] == expected[49]
  888. second = service.reconcile(
  889. limit=50,
  890. cancel_cursor=first["cancel_next_cursor"],
  891. release_cursor=None,
  892. **auth,
  893. )
  894. assert second["cancelled_task_ids"] == expected[50:]
  895. assert second["cancel_next_cursor"] is None
  896. def test_pull_terminalizes_expired_task_without_issuing_lease(service, engine):
  897. gateway = _register(service, _enrollment(service))
  898. auth = {
  899. "credential": gateway["credential"],
  900. "certificate_sha256": "b" * 64,
  901. "gateway_id": gateway["gateway_id"],
  902. "environment": "staging",
  903. "network_zone": "manufacturing-zone-a",
  904. "generation": 1,
  905. }
  906. task_id = "already-expired-task"
  907. service.issue_task(
  908. {
  909. "task_id": task_id,
  910. "gateway_id": gateway["gateway_id"],
  911. "environment": "staging",
  912. "network_zone": "manufacturing-zone-a",
  913. "purpose": "inventory",
  914. "classification": "statistics",
  915. "task_type": "profile",
  916. "contract_version": 1,
  917. "deadline_at": (datetime.now(UTC) + timedelta(milliseconds=200))
  918. .isoformat()
  919. .replace("+00:00", "Z"),
  920. "attempt": 1,
  921. "idempotency_key": "already-expired-task-key",
  922. "policy_digest": "a" * 64,
  923. },
  924. actor_uid=_uid(),
  925. )
  926. with engine.connect() as connection:
  927. connection.execute(text("SELECT pg_sleep(0.3)"))
  928. assert service.pull_task(**auth) == {"task": None}
  929. with engine.connect() as connection:
  930. row = connection.execute(
  931. text(
  932. "SELECT status,lease_token_hash,lease_expires_at,result_summary "
  933. "FROM edge_gateway_tasks WHERE task_id=:task_id"
  934. ),
  935. {"task_id": task_id},
  936. ).mappings().one()
  937. assert row["status"] == "failed"
  938. assert row["lease_token_hash"] is None and row["lease_expires_at"] is None
  939. assert dict(row["result_summary"]) == {"reason_code": "task_deadline_expired"}
  940. def test_reconcile_prioritizes_new_actionable_release_over_terminal_history(service):
  941. gateway = _register(service, _enrollment(service))
  942. auth = {
  943. "credential": gateway["credential"],
  944. "certificate_sha256": "b" * 64,
  945. "gateway_id": gateway["gateway_id"],
  946. "environment": "staging",
  947. "network_zone": "manufacturing-zone-a",
  948. "generation": 1,
  949. }
  950. for index in range(51):
  951. release = service.offer_release(
  952. gateway["gateway_id"],
  953. version=f"4.{index}.0",
  954. rollback_version="3.0.0",
  955. artifact_digest=f"{index:064x}",
  956. request_id=f"terminal-release-{index}",
  957. actor_uid=_uid(),
  958. )
  959. service.acknowledge_release(
  960. release_id=release["release_id"],
  961. outcome="accepted",
  962. safe_summary={"release_status": "accepted", "version": f"4.{index}.0"},
  963. **auth,
  964. )
  965. service.acknowledge_release(
  966. release_id=release["release_id"],
  967. outcome="installed",
  968. safe_summary={"release_status": "installed", "version": f"4.{index}.0"},
  969. **auth,
  970. )
  971. actionable = service.offer_release(
  972. gateway["gateway_id"],
  973. version="9.0.0",
  974. rollback_version="3.0.0",
  975. artifact_digest="f" * 64,
  976. request_id="new-actionable-release",
  977. actor_uid=_uid(),
  978. )
  979. reconciled = service.reconcile(limit=50, **auth)
  980. releases = reconciled["release_offers"]
  981. assert len(releases) == 1
  982. assert releases[0]["release_id"] == actionable["release_id"]
  983. assert releases[0]["status"] == "offered"
  984. assert reconciled["release_next_cursor"] is None
  985. assert reconciled["release_baseline"]["status"] == "installed"
  986. assert reconciled["release_baseline"]["version"] == "4.50.0"
  987. def test_signed_release_offer_installs_through_real_postgres_control_plane(
  988. service, engine, tmp_path,
  989. ):
  990. gateway = _register(service, _enrollment(service))
  991. auth = {
  992. "credential": gateway["credential"],
  993. "certificate_sha256": "b" * 64,
  994. "gateway_id": gateway["gateway_id"],
  995. "environment": "staging",
  996. "network_zone": "manufacturing-zone-a",
  997. "generation": 1,
  998. }
  999. artifact = b"postgres-backed-signed-edge-release"
  1000. offered = service.offer_release(
  1001. gateway["gateway_id"],
  1002. version="3.1.0",
  1003. rollback_version="3.0.0",
  1004. artifact_digest=hashlib.sha256(artifact).hexdigest(),
  1005. artifact_name="edge-agent-3.1.0.bin",
  1006. deadline_at=(datetime.now(UTC) + timedelta(hours=1))
  1007. .isoformat()
  1008. .replace("+00:00", "Z"),
  1009. request_id="postgres-agent-release-chain",
  1010. actor_uid=_uid(),
  1011. )
  1012. class PostgresControlTransport:
  1013. def reconcile(self):
  1014. return service.reconcile(**auth)
  1015. def acknowledge_release(self, release_id, outcome, summary):
  1016. return service.acknowledge_release(
  1017. release_id=release_id,
  1018. outcome=outcome,
  1019. safe_summary=summary,
  1020. **auth,
  1021. )
  1022. def pull_task(self):
  1023. return None
  1024. public_key = _SIGNING_PRIVATE_KEY.public_key().public_bytes(
  1025. encoding=serialization.Encoding.Raw,
  1026. format=serialization.PublicFormat.Raw,
  1027. ).hex()
  1028. activated = []
  1029. agent = EdgeAgent(
  1030. EdgeBootstrapConfig(
  1031. queue_path=str(tmp_path / "postgres-release-agent.sqlite3"),
  1032. artifact_root=str(tmp_path),
  1033. gateway_id=gateway["gateway_id"],
  1034. credential=gateway["credential"],
  1035. certificate_sha256="b" * 64,
  1036. client_certificate_path=str(tmp_path / "client.pem"),
  1037. client_private_key_path=str(tmp_path / "client.key"),
  1038. ca_bundle_path=str(tmp_path / "ca.pem"),
  1039. generation=1,
  1040. environment="staging",
  1041. network_zone="manufacturing-zone-a",
  1042. policy_digest="a" * 64,
  1043. control_url="https://control.example.test",
  1044. proxy_url=None,
  1045. allowed_control_hosts=frozenset({"control.example.test"}),
  1046. allowed_proxy_hosts=frozenset(),
  1047. trusted_release_keys={_SIGNING_KEY_ID: public_key},
  1048. trusted_task_keys={_SIGNING_KEY_ID: public_key},
  1049. version="3.0.0",
  1050. ),
  1051. PostgresControlTransport(),
  1052. EdgeRunnerAdapter(
  1053. lambda _node, _request, _cancel: {
  1054. "classification": "statistics",
  1055. "payload": {"metric_count": 1},
  1056. }
  1057. ),
  1058. release_manager=SignedReleaseManager(
  1059. current_version="3.0.0",
  1060. trusted_keys={_SIGNING_KEY_ID: public_key},
  1061. artifact_loader=lambda _manifest: artifact,
  1062. activator=lambda version, body: activated.append((version, body)) or True,
  1063. rollback=lambda _version: False,
  1064. ),
  1065. )
  1066. agent.run_once()
  1067. with engine.connect() as connection:
  1068. persisted = connection.execute(
  1069. text("SELECT status,signed_manifest FROM edge_gateway_releases WHERE uid=CAST(:uid AS uuid)"),
  1070. {"uid": offered["release_id"]},
  1071. ).mappings().one()
  1072. assert persisted["status"] == "installed"
  1073. assert dict(persisted["signed_manifest"])["status"] == "installed"
  1074. assert activated == [("3.1.0", artifact)]
  1075. assert agent.queue.get_release_state(offered["release_id"]).status == "installed"
  1076. def test_signing_operations_fail_closed_without_configured_private_key(service, engine):
  1077. gateway = _register(service, _enrollment(service))
  1078. session = Session(engine)
  1079. unsigned_service = EdgeGatewayService(EdgeGatewayRepository(session))
  1080. try:
  1081. with pytest.raises(EdgeGatewayConfigurationError):
  1082. unsigned_service.issue_task({
  1083. "task_id": "task-no-signer", "gateway_id": gateway["gateway_id"],
  1084. "environment": "staging", "network_zone": "manufacturing-zone-a",
  1085. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  1086. "contract_version": 1,
  1087. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1088. "attempt": 1, "idempotency_key": "task-no-signer-key", "policy_digest": "a" * 64,
  1089. }, actor_uid=_uid())
  1090. finally:
  1091. session.close()
  1092. with pytest.raises(EdgeGatewayValidationError):
  1093. service.issue_task({
  1094. "task_id": "task-expired-authority", "gateway_id": gateway["gateway_id"],
  1095. "environment": "staging", "network_zone": "manufacturing-zone-a",
  1096. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  1097. "contract_version": 1,
  1098. "deadline_at": (datetime.now(UTC) - timedelta(seconds=1)).isoformat().replace("+00:00", "Z"),
  1099. "attempt": 1, "idempotency_key": "task-expired-authority-key", "policy_digest": "a" * 64,
  1100. }, actor_uid=_uid())
  1101. def test_task_outcome_audit_is_atomic_bounded_and_replay_idempotent(service, engine):
  1102. gateway = _register(service, _enrollment(service))
  1103. auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1104. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1105. "network_zone": "manufacturing-zone-a", "generation": 1}
  1106. task = service.issue_task({
  1107. "task_id": "task-outcome-audit", "gateway_id": gateway["gateway_id"],
  1108. "environment": "staging", "network_zone": "manufacturing-zone-a",
  1109. "purpose": "inventory", "classification": "statistics", "task_type": "quality",
  1110. "contract_version": 1,
  1111. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1112. "attempt": 1, "idempotency_key": "task-outcome-audit-key", "policy_digest": "a" * 64,
  1113. }, actor_uid=_uid())
  1114. lease = service.pull_task(**auth)
  1115. with pytest.raises(EdgeGatewayValidationError):
  1116. service.task_outcome(
  1117. task["task_id"], outcome="failed", lease_token=lease["lease_token"],
  1118. safe_summary={"reason_code": "x" * 40_000}, **auth,
  1119. )
  1120. with engine.begin() as connection:
  1121. connection.execute(text("""
  1122. CREATE OR REPLACE FUNCTION test_reject_task_outcome_audit()
  1123. RETURNS trigger LANGUAGE plpgsql AS $$
  1124. BEGIN
  1125. IF NEW.event_type='task_outcome_recorded' THEN
  1126. RAISE EXCEPTION 'test audit failure';
  1127. END IF;
  1128. RETURN NEW;
  1129. END $$;
  1130. CREATE TRIGGER test_reject_task_outcome_audit
  1131. BEFORE INSERT ON edge_gateway_audits FOR EACH ROW
  1132. EXECUTE FUNCTION test_reject_task_outcome_audit();
  1133. """))
  1134. session_usable = False
  1135. try:
  1136. with pytest.raises(DBAPIError):
  1137. service.task_outcome(
  1138. task["task_id"], outcome="failed", lease_token=lease["lease_token"],
  1139. safe_summary={"reason_code": "controlled_failure"}, **auth,
  1140. )
  1141. service.repository.session.execute(text("SELECT 1")).scalar_one()
  1142. session_usable = True
  1143. finally:
  1144. service.repository.rollback()
  1145. with engine.begin() as connection:
  1146. connection.execute(text("DROP TRIGGER test_reject_task_outcome_audit ON edge_gateway_audits"))
  1147. connection.execute(text("DROP FUNCTION test_reject_task_outcome_audit()"))
  1148. assert session_usable is True
  1149. with engine.connect() as connection:
  1150. assert connection.execute(text("""
  1151. SELECT status FROM edge_gateway_tasks WHERE task_id=:task_id
  1152. """), {"task_id": task["task_id"]}).scalar_one() == "leased"
  1153. assert connection.execute(text("""
  1154. SELECT count(*) FROM edge_gateway_audits WHERE event_type='task_outcome_recorded'
  1155. """)).scalar_one() == 0
  1156. result = service.task_outcome(
  1157. task["task_id"], outcome="failed", lease_token=lease["lease_token"],
  1158. safe_summary={"reason_code": "controlled_failure"}, **auth,
  1159. )
  1160. assert result["status"] == "failed" and result["replayed"] is False
  1161. replay = service.task_outcome(
  1162. task["task_id"], outcome="failed", lease_token=lease["lease_token"],
  1163. safe_summary={"reason_code": "controlled_failure"}, **auth,
  1164. )
  1165. assert replay["status"] == "failed" and replay["replayed"] is True
  1166. with engine.connect() as connection:
  1167. audits = connection.execute(text("""
  1168. SELECT safe_detail FROM edge_gateway_audits WHERE event_type='task_outcome_recorded'
  1169. """)).scalars().all()
  1170. assert audits == [{"task_id": task["task_id"], "outcome": "failed"}]
  1171. def test_cross_gateway_event_replay_never_discloses_foreign_ack(service):
  1172. gateway_a = _register(service, _enrollment(service, gateway_name="gateway-a"))
  1173. gateway_b = _register(service, _enrollment(service, gateway_name="gateway-b"), certificate_sha256="b" * 64)
  1174. 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}
  1175. 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}
  1176. task = service.issue_task({
  1177. "task_id": "task-owner-a", "gateway_id": gateway_a["gateway_id"], "environment": "staging",
  1178. "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics",
  1179. "task_type": "profile", "contract_version": 1,
  1180. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1181. "attempt": 1, "idempotency_key": "task-owner-a-key", "policy_digest": "a" * 64,
  1182. }, actor_uid=_uid())
  1183. pulled = service.pull_task(**auth_a)
  1184. event = _event(task, gateway_a)
  1185. ack = service.accept_event(event=event, lease_token=pulled["lease_token"], **auth_a)
  1186. with pytest.raises(EdgeGatewayAuthenticationError):
  1187. service.accept_event(event=event, lease_token=pulled["lease_token"], **auth_b)
  1188. assert ack["event_id"] == event["event_id"]
  1189. def test_release_offer_is_concurrently_idempotent_and_conflicts_are_controlled(service, engine):
  1190. gateway = _register(service, _enrollment(service))
  1191. barrier = Barrier(2)
  1192. def offer_once():
  1193. session = Session(engine)
  1194. try:
  1195. barrier.wait(timeout=5)
  1196. return _service(EdgeGatewayRepository(session)).offer_release(
  1197. gateway["gateway_id"], version="5.0.0", artifact_digest="5" * 64,
  1198. rollback_version="4.9.0", request_id="concurrent-release", actor_uid=_uid(),
  1199. )
  1200. finally:
  1201. session.close()
  1202. with ThreadPoolExecutor(max_workers=2) as pool:
  1203. results = [future.result() for future in (pool.submit(offer_once), pool.submit(offer_once))]
  1204. assert results[0] == results[1]
  1205. with pytest.raises(EdgeGatewayConflictError):
  1206. service.offer_release(
  1207. gateway["gateway_id"], version="5.0.1", artifact_digest="6" * 64,
  1208. rollback_version="4.9.0", request_id="concurrent-release", actor_uid=_uid(),
  1209. )
  1210. with pytest.raises(EdgeGatewayConflictError):
  1211. service.offer_release(
  1212. gateway["gateway_id"], version="5.0.0", artifact_digest="5" * 64,
  1213. rollback_version="4.9.0", request_id="different-request", actor_uid=_uid(),
  1214. )
  1215. service.offer_release(
  1216. gateway["gateway_id"], version="5.1.0", artifact_digest="7" * 64,
  1217. rollback_version="5.0.0", request_id="release-b", actor_uid=_uid(),
  1218. )
  1219. with pytest.raises(EdgeGatewayConflictError):
  1220. service.offer_release(
  1221. gateway["gateway_id"], version="5.1.0", artifact_digest="7" * 64,
  1222. rollback_version="5.0.0", request_id="concurrent-release", actor_uid=_uid(),
  1223. )
  1224. assert service.list_gateways()[0]["gateway_id"] == gateway["gateway_id"]
  1225. def test_evidence_owner_and_runtime_role_are_really_isolated(engine):
  1226. runtime_login = f"edge_test_runtime_{uuid.uuid4().hex[:12]}"
  1227. runtime_password = f"runtime-{uuid.uuid4().hex}"
  1228. runtime_engine = None
  1229. runtime_session = None
  1230. with engine.begin() as connection:
  1231. connection.execute(text(f'CREATE ROLE "{runtime_login}" LOGIN PASSWORD :password'), {
  1232. "password": runtime_password,
  1233. })
  1234. connection.execute(text(f'GRANT dataops_app_runtime TO "{runtime_login}"'))
  1235. try:
  1236. runtime_url = make_url(DATABASE_URL).set(
  1237. username=runtime_login, password=runtime_password,
  1238. )
  1239. runtime_engine = create_engine(runtime_url, pool_pre_ping=True)
  1240. validate_production_database_identity(
  1241. SimpleNamespace(config={"FLASK_ENV": "production"}), runtime_engine,
  1242. )
  1243. runtime_session = Session(runtime_engine)
  1244. runtime_service = _service(EdgeGatewayRepository(runtime_session))
  1245. enrollment = _enrollment(runtime_service)
  1246. gateway = _register(runtime_service, enrollment)
  1247. assert gateway["generation"] == 1
  1248. auth = {
  1249. "credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1250. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1251. "network_zone": "manufacturing-zone-a", "generation": 1,
  1252. }
  1253. assert runtime_service.heartbeat(
  1254. version="3.0.1", safe_summary={"queue_depth": 0}, **auth,
  1255. )["status"] == "online"
  1256. task = runtime_service.issue_task({
  1257. "task_id": "runtime-role-task", "gateway_id": gateway["gateway_id"],
  1258. "environment": "staging", "network_zone": "manufacturing-zone-a",
  1259. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  1260. "contract_version": 1,
  1261. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1262. "attempt": 1, "idempotency_key": "runtime-role-task-key", "policy_digest": "a" * 64,
  1263. }, actor_uid=_uid())
  1264. assert runtime_service.pull_task(**auth)["task"]["task_id"] == task["task_id"]
  1265. with runtime_engine.connect() as connection:
  1266. identity = connection.execute(text("SELECT current_user,session_user")).one()
  1267. assert identity == (runtime_login, runtime_login)
  1268. owners = dict(connection.execute(text("""
  1269. SELECT c.relname,r.rolname FROM pg_class c
  1270. JOIN pg_roles r ON r.oid=c.relowner
  1271. WHERE c.relname IN ('edge_gateway_audits','edge_gateway_failure_windows')
  1272. """)).all())
  1273. assert owners == {
  1274. "edge_gateway_audits": "dataops_edge_evidence_owner",
  1275. "edge_gateway_failure_windows": "dataops_edge_evidence_owner",
  1276. }
  1277. assert connection.execute(text("""
  1278. SELECT count(*)=0 FROM information_schema.routine_privileges
  1279. WHERE routine_schema='public'
  1280. AND routine_name IN ('record_edge_gateway_failure','record_edge_gateway_audit')
  1281. AND grantee='PUBLIC' AND privilege_type='EXECUTE'
  1282. """)).scalar_one() is True
  1283. attacks = (
  1284. "INSERT INTO edge_gateway_audits(uid,event_type,success,safe_detail) VALUES (gen_random_uuid(),'authentication_rejected',false,'{}')",
  1285. "UPDATE edge_gateway_audits SET success=true",
  1286. "DELETE FROM edge_gateway_failure_windows",
  1287. "TRUNCATE TABLE edge_gateway_audits,edge_gateway_failure_windows",
  1288. "ALTER TABLE edge_gateway_audits DISABLE TRIGGER ALL",
  1289. "SET ROLE dataops_edge_evidence_owner",
  1290. )
  1291. for statement in attacks:
  1292. with pytest.raises(DBAPIError), runtime_engine.begin() as connection:
  1293. connection.execute(text(statement))
  1294. with runtime_engine.begin() as connection:
  1295. connection.execute(text("SELECT set_config('dataops.edge_failure_writer','477-controlled',false)"))
  1296. with pytest.raises(DBAPIError):
  1297. connection.execute(text("""
  1298. INSERT INTO edge_gateway_failure_windows
  1299. (failure_fingerprint,window_started_at,event_type,occurrence_count,first_seen_at,last_seen_at)
  1300. VALUES (:fingerprint,CURRENT_TIMESTAMP,'authentication_rejected',1,CURRENT_TIMESTAMP,CURRENT_TIMESTAMP)
  1301. """), {"fingerprint": "f" * 64})
  1302. finally:
  1303. if runtime_session is not None:
  1304. runtime_session.close()
  1305. if runtime_engine is not None:
  1306. runtime_engine.dispose()
  1307. with engine.begin() as connection:
  1308. connection.execute(text(f'DROP ROLE IF EXISTS "{runtime_login}"'))
  1309. def test_production_owner_database_identity_and_direct_entry_fail_closed(engine):
  1310. with pytest.raises(RuntimeError, match="least-privilege runtime identity"):
  1311. validate_production_database_identity(
  1312. SimpleNamespace(config={"FLASK_ENV": "production"}), engine,
  1313. )
  1314. result = subprocess.run(
  1315. [str(ROOT / ".venv/bin/python"), "-c", "from app import create_app; create_app()"],
  1316. cwd=ROOT,
  1317. env={
  1318. **os.environ,
  1319. "PYTHONPATH": str(ROOT),
  1320. "FLASK_ENV": "production",
  1321. "APP_ENV_FILE": "/nonexistent/dataops-test.env",
  1322. "DATABASE_URL": DATABASE_URL,
  1323. "NEO4J_URI": "bolt://neo4j.example.test:7687",
  1324. "NEO4J_HTTP_URI": "http://neo4j.example.test:7474",
  1325. "NEO4J_USER": "dataops_graph",
  1326. "NEO4J_PASSWORD": "graph-test-secret",
  1327. "MINIO_HOST": "minio.example.test:9000",
  1328. "MINIO_USER": "dataops_objects",
  1329. "MINIO_PASSWORD": "object-test-secret",
  1330. "MINIO_BUCKET": "dataops-test",
  1331. },
  1332. capture_output=True,
  1333. text=True,
  1334. )
  1335. assert result.returncode != 0
  1336. assert "least-privilege runtime identity" in result.stderr
  1337. assert "dataops-test-password" not in result.stderr
  1338. def test_preprovisioned_no_createrole_migrator_can_apply_477_and_478(engine):
  1339. migrator = f"edge_test_migrator_{uuid.uuid4().hex[:12]}"
  1340. password = f"migration-{uuid.uuid4().hex}"
  1341. _clear_edge_rows_for_test(engine)
  1342. _alembic("downgrade", "20260802_476")
  1343. try:
  1344. with engine.begin() as connection:
  1345. connection.execute(text(f'CREATE ROLE "{migrator}" LOGIN NOSUPERUSER NOCREATEROLE NOCREATEDB PASSWORD :password'), {"password": password})
  1346. connection.execute(text(f'GRANT dataops_edge_evidence_owner TO "{migrator}"'))
  1347. connection.execute(text('GRANT USAGE,CREATE ON SCHEMA public TO dataops_edge_evidence_owner'))
  1348. connection.execute(text(f'GRANT USAGE,CREATE ON SCHEMA public TO "{migrator}"'))
  1349. connection.execute(text(f'GRANT SELECT,INSERT,UPDATE,DELETE ON alembic_version TO "{migrator}"'))
  1350. migrator_url = make_url(DATABASE_URL).set(
  1351. username=migrator, password=password,
  1352. ).render_as_string(hide_password=False)
  1353. try:
  1354. _alembic("upgrade", "head", migrator_url)
  1355. except subprocess.CalledProcessError as exc:
  1356. pytest.fail(exc.stderr)
  1357. with engine.connect() as connection:
  1358. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260809_479"
  1359. assert connection.execute(text("SELECT NOT rolcreaterole AND NOT rolsuper FROM pg_roles WHERE rolname=:role"), {"role": migrator}).scalar_one() is True
  1360. finally:
  1361. _clear_edge_rows_for_test(engine)
  1362. _alembic("downgrade", "20260802_476")
  1363. with engine.begin() as connection:
  1364. connection.execute(text(f'REVOKE dataops_edge_evidence_owner FROM "{migrator}"'))
  1365. connection.execute(text(f'DROP OWNED BY "{migrator}"'))
  1366. connection.execute(text(f'DROP ROLE "{migrator}"'))
  1367. _alembic("upgrade", "head")
  1368. def test_backend_runtime_configuration_never_contains_migrator_credentials():
  1369. compose = (ROOT / "deploy/docker/docker-compose.yml").read_text()
  1370. backend_dockerfile = (ROOT / "deploy/docker/backend.Dockerfile").read_text()
  1371. runner_dockerfile = (ROOT / "deploy/docker/runner.Dockerfile").read_text()
  1372. run_script = (ROOT / "scripts/run_dataops.sh").read_text()
  1373. assert "db-role-init:" in compose and "db-migrate:" in compose
  1374. assert "DATAOPS_RUNTIME_PASSWORD" in compose
  1375. assert "postgresql://dataops_app:" in compose
  1376. backend = compose.split(" backend:", 1)[1].split(" frontend:", 1)[0]
  1377. assert "postgresql://dataops:dataops-test-password" not in backend
  1378. assert "MIGRATION_DATABASE_URL" not in backend
  1379. assert "db-migrate:\n condition: service_completed_successfully" in backend
  1380. backend_runtime = backend_dockerfile.split("FROM dependencies AS runtime", 1)[1]
  1381. assert "alembic.ini" not in backend_runtime and "COPY migrations/" not in backend_runtime
  1382. assert "pip uninstall -y alembic" in backend_runtime
  1383. assert "alembic.ini" not in runner_dockerfile and "COPY migrations/" not in runner_dockerfile
  1384. assert "pip uninstall -y alembic" in runner_dockerfile
  1385. assert 'role_init_url="${DB_ROLE_INIT_DATABASE_URL:' in run_script
  1386. assert 'migration_url="${MIGRATION_DATABASE_URL:' in run_script
  1387. assert "unset DB_ROLE_INIT_DATABASE_URL MIGRATION_DATABASE_URL" in run_script
  1388. assert "MIGRATION_DATABASE_URL=" in (ROOT / "deployment/dataops.env").read_text()
  1389. def test_task_pull_and_exact_event_replay_are_concurrency_safe(service, engine):
  1390. gateway = _register(service, _enrollment(service))
  1391. auth = {
  1392. "credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1393. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1394. "network_zone": "manufacturing-zone-a", "generation": 1,
  1395. }
  1396. task = service.issue_task({
  1397. "task_id": "task-concurrent", "gateway_id": gateway["gateway_id"],
  1398. "environment": "staging", "network_zone": "manufacturing-zone-a",
  1399. "purpose": "inventory", "classification": "statistics", "task_type": "profile",
  1400. "contract_version": 1,
  1401. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1402. "attempt": 1, "idempotency_key": "task-concurrent-key", "policy_digest": "a" * 64,
  1403. }, actor_uid=_uid())
  1404. pull_barrier = Barrier(2)
  1405. def pull_once():
  1406. session = Session(engine)
  1407. try:
  1408. worker = _service(EdgeGatewayRepository(session))
  1409. pull_barrier.wait()
  1410. return worker.pull_task(**auth)
  1411. finally:
  1412. session.close()
  1413. with ThreadPoolExecutor(max_workers=2) as pool:
  1414. pulls = [future.result() for future in (pool.submit(pull_once), pool.submit(pull_once))]
  1415. assert sorted(result["task"] is not None for result in pulls) == [False, True]
  1416. lease_token = next(result["lease_token"] for result in pulls if result["task"] is not None)
  1417. event = _event(task, gateway)
  1418. event_barrier = Barrier(2)
  1419. def accept_once():
  1420. session = Session(engine)
  1421. try:
  1422. worker = _service(EdgeGatewayRepository(session))
  1423. event_barrier.wait()
  1424. return worker.accept_event(event=event, lease_token=lease_token, **auth)
  1425. finally:
  1426. session.close()
  1427. with ThreadPoolExecutor(max_workers=2) as pool:
  1428. acks = [future.result() for future in (pool.submit(accept_once), pool.submit(accept_once))]
  1429. assert acks[0] == acks[1]
  1430. with engine.connect() as connection:
  1431. assert connection.execute(text("SELECT count(*) FROM edge_gateway_events WHERE event_id=:event_id"), {"event_id": event["event_id"]}).scalar_one() == 1
  1432. assert service.accept_event(event=event, lease_token=lease_token, **auth) == acks[0]
  1433. def test_cancel_and_event_barrier_has_single_linearized_terminal_state(service, engine):
  1434. gateway = _register(service, _enrollment(service))
  1435. auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1436. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1437. "network_zone": "manufacturing-zone-a", "generation": 1}
  1438. task = service.issue_task({
  1439. "task_id": "task-race", "gateway_id": gateway["gateway_id"], "environment": "staging",
  1440. "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics",
  1441. "task_type": "profile", "contract_version": 1,
  1442. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1443. "attempt": 1, "idempotency_key": "task-race-key", "policy_digest": "a" * 64,
  1444. }, actor_uid=_uid())
  1445. lease = service.pull_task(**auth)
  1446. event = _event(task, gateway)
  1447. barrier = Barrier(2)
  1448. def accept():
  1449. session = Session(engine)
  1450. try:
  1451. barrier.wait()
  1452. return _service(EdgeGatewayRepository(session)).accept_event(
  1453. event=event, lease_token=lease["lease_token"], **auth,
  1454. )
  1455. except EdgeGatewayError:
  1456. return None
  1457. finally:
  1458. session.close()
  1459. def cancel():
  1460. session = Session(engine)
  1461. try:
  1462. barrier.wait()
  1463. return _service(EdgeGatewayRepository(session)).cancel_task(task["task_id"], actor_uid=_uid())
  1464. finally:
  1465. session.close()
  1466. with ThreadPoolExecutor(max_workers=2) as pool:
  1467. accepted, cancelled = pool.submit(accept), pool.submit(cancel)
  1468. accepted, cancelled = accepted.result(), cancelled.result()
  1469. with engine.connect() as connection:
  1470. row = connection.execute(text("""
  1471. SELECT status,(SELECT count(*) FROM edge_gateway_events WHERE task_id=:task_id) AS events
  1472. FROM edge_gateway_tasks WHERE task_id=:task_id
  1473. """), {"task_id": task["task_id"]}).mappings().one()
  1474. assert (row["status"], row["events"]) in {("completed", 1), ("cancel_requested", 0)}
  1475. assert bool(accepted) != bool(cancelled)
  1476. def test_cancel_event_barriers_prove_both_lock_orders(service, engine):
  1477. gateway = _register(service, _enrollment(service))
  1478. auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1479. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1480. "network_zone": "manufacturing-zone-a", "generation": 1}
  1481. def issue_and_pull(task_id):
  1482. task = service.issue_task({
  1483. "task_id": task_id, "gateway_id": gateway["gateway_id"], "environment": "staging",
  1484. "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics",
  1485. "task_type": "profile", "contract_version": 1,
  1486. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1487. "attempt": 1, "idempotency_key": f"{task_id}-key", "policy_digest": "a" * 64,
  1488. }, actor_uid=_uid())
  1489. return task, service.pull_task(**auth)
  1490. cancel_task, cancel_lease = issue_and_pull("task-cancel-wins")
  1491. cancel_has_gateway = Event()
  1492. cancel_barrier = Barrier(2)
  1493. class CancelFirstRepository(EdgeGatewayRepository):
  1494. def cancel_task(self, task_id, actor_uid):
  1495. self.session.execute(text("""
  1496. SELECT g.uid FROM edge_gateway_tasks t JOIN edge_gateways g ON g.uid=t.gateway_id
  1497. WHERE t.task_id=:task_id FOR UPDATE OF g
  1498. """), {"task_id": task_id}).scalar_one()
  1499. cancel_has_gateway.set()
  1500. cancel_barrier.wait(timeout=5)
  1501. return super().cancel_task(task_id, actor_uid)
  1502. def cancel_first():
  1503. session = Session(engine)
  1504. try:
  1505. return _service(CancelFirstRepository(session)).cancel_task(
  1506. cancel_task["task_id"], actor_uid=_uid(),
  1507. )
  1508. finally:
  1509. session.close()
  1510. def event_after_cancel_lock():
  1511. assert cancel_has_gateway.wait(timeout=5)
  1512. cancel_barrier.wait(timeout=5)
  1513. session = Session(engine)
  1514. try:
  1515. return _service(EdgeGatewayRepository(session)).accept_event(
  1516. event=_event(cancel_task, gateway), lease_token=cancel_lease["lease_token"], **auth,
  1517. )
  1518. except EdgeGatewayError:
  1519. return None
  1520. finally:
  1521. session.close()
  1522. with ThreadPoolExecutor(max_workers=2) as pool:
  1523. cancel_result = pool.submit(cancel_first)
  1524. event_result = pool.submit(event_after_cancel_lock)
  1525. assert cancel_result.result() is True
  1526. assert event_result.result() is None
  1527. event_task, event_lease = issue_and_pull("task-event-wins")
  1528. event_has_gateway = Event()
  1529. event_barrier = Barrier(2)
  1530. class EventFirstRepository(EdgeGatewayRepository):
  1531. def authenticate(self, credential_hash, *, lock=False):
  1532. identity = super().authenticate(credential_hash, lock=lock)
  1533. if lock:
  1534. event_has_gateway.set()
  1535. event_barrier.wait(timeout=5)
  1536. return identity
  1537. def event_first():
  1538. session = Session(engine)
  1539. try:
  1540. return _service(EventFirstRepository(session)).accept_event(
  1541. event=_event(event_task, gateway), lease_token=event_lease["lease_token"], **auth,
  1542. )
  1543. finally:
  1544. session.close()
  1545. def cancel_after_event_lock():
  1546. assert event_has_gateway.wait(timeout=5)
  1547. event_barrier.wait(timeout=5)
  1548. session = Session(engine)
  1549. try:
  1550. return _service(EdgeGatewayRepository(session)).cancel_task(
  1551. event_task["task_id"], actor_uid=_uid(),
  1552. )
  1553. finally:
  1554. session.close()
  1555. with ThreadPoolExecutor(max_workers=2) as pool:
  1556. event_result = pool.submit(event_first)
  1557. cancel_result = pool.submit(cancel_after_event_lock)
  1558. assert event_result.result()["status"] == "accepted"
  1559. assert cancel_result.result() is False
  1560. with engine.connect() as connection:
  1561. rows = dict(connection.execute(text("""
  1562. SELECT task_id,status FROM edge_gateway_tasks
  1563. WHERE task_id IN ('task-cancel-wins','task-event-wins')
  1564. """)).all())
  1565. assert rows == {"task-cancel-wins": "cancel_requested", "task-event-wins": "completed"}
  1566. def test_machine_operations_and_rotate_are_gateway_fenced_with_barrier(service, engine):
  1567. gateway = _register(service, _enrollment(service))
  1568. auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1569. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1570. "network_zone": "manufacturing-zone-a", "generation": 1}
  1571. release = service.offer_release(
  1572. gateway["gateway_id"], version="4.0.0", artifact_digest="7" * 64,
  1573. rollback_version="3.0.0", request_id="barrier-release", actor_uid=_uid(),
  1574. )
  1575. barrier = Barrier(3)
  1576. def run_machine(kind):
  1577. session = Session(engine)
  1578. worker = _service(EdgeGatewayRepository(session))
  1579. try:
  1580. barrier.wait()
  1581. if kind == "reconcile":
  1582. return worker.reconcile(**auth)
  1583. return worker.acknowledge_release(
  1584. release_id=release["release_id"], outcome="accepted",
  1585. safe_summary={"reason_code": "accepted"}, **auth,
  1586. )
  1587. except EdgeGatewayError:
  1588. return None
  1589. finally:
  1590. session.close()
  1591. def run_rotate():
  1592. session = Session(engine)
  1593. try:
  1594. barrier.wait()
  1595. return _service(EdgeGatewayRepository(session)).rotate(
  1596. gateway["gateway_id"], certificate_sha256="c" * 64,
  1597. request_id="barrier-rotate", actor_uid=_uid(),
  1598. )
  1599. finally:
  1600. session.close()
  1601. with ThreadPoolExecutor(max_workers=3) as pool:
  1602. results = [future.result() for future in (
  1603. pool.submit(run_machine, "reconcile"), pool.submit(run_machine, "ack"), pool.submit(run_rotate),
  1604. )]
  1605. assert results[2]["generation"] == 2
  1606. with pytest.raises(EdgeGatewayAuthenticationError):
  1607. service.reconcile(**auth)
  1608. def test_pull_event_and_revoke_are_gateway_fenced_with_barrier(service, engine):
  1609. gateway = _register(service, _enrollment(service))
  1610. auth = {"credential": gateway["credential"], "certificate_sha256": "b" * 64,
  1611. "gateway_id": gateway["gateway_id"], "environment": "staging",
  1612. "network_zone": "manufacturing-zone-a", "generation": 1}
  1613. def issue(task_id):
  1614. return service.issue_task({
  1615. "task_id": task_id, "gateway_id": gateway["gateway_id"], "environment": "staging",
  1616. "network_zone": "manufacturing-zone-a", "purpose": "inventory", "classification": "statistics",
  1617. "task_type": "profile", "contract_version": 1,
  1618. "deadline_at": (datetime.now(UTC) + timedelta(minutes=10)).isoformat().replace("+00:00", "Z"),
  1619. "attempt": 1, "idempotency_key": f"{task_id}-key", "policy_digest": "a" * 64,
  1620. }, actor_uid=_uid())
  1621. event_task = issue("task-revoke-event")
  1622. event_lease = service.pull_task(**auth)
  1623. issue("task-revoke-pull")
  1624. event = _event(event_task, gateway)
  1625. barrier = Barrier(3)
  1626. def machine(kind):
  1627. session = Session(engine)
  1628. worker = _service(EdgeGatewayRepository(session))
  1629. try:
  1630. barrier.wait()
  1631. if kind == "event":
  1632. return worker.accept_event(event=event, lease_token=event_lease["lease_token"], **auth)
  1633. return worker.pull_task(**auth)
  1634. except EdgeGatewayError:
  1635. return None
  1636. finally:
  1637. session.close()
  1638. def revoke():
  1639. session = Session(engine)
  1640. try:
  1641. barrier.wait()
  1642. return _service(EdgeGatewayRepository(session)).revoke(gateway["gateway_id"], actor_uid=_uid())
  1643. finally:
  1644. session.close()
  1645. with ThreadPoolExecutor(max_workers=3) as pool:
  1646. results = [future.result() for future in (
  1647. pool.submit(machine, "event"), pool.submit(machine, "pull"), pool.submit(revoke),
  1648. )]
  1649. assert results[2] is True
  1650. with pytest.raises(EdgeGatewayAuthenticationError):
  1651. service.pull_task(**auth)
  1652. def test_failure_window_uses_lock_time_not_stale_transaction_time(engine):
  1653. fingerprint = hashlib.sha256(f"stale-clock:{_uid()}".encode()).hexdigest()
  1654. statement = text("""
  1655. SELECT occurrence_count
  1656. FROM public.record_edge_gateway_failure(
  1657. CAST(:audit_uid AS uuid), :fingerprint, NULL,
  1658. 'untrusted_request_rejected', 'untrusted_failure'
  1659. )
  1660. """)
  1661. older = engine.connect()
  1662. transaction = older.begin()
  1663. try:
  1664. older.execute(text("SELECT CURRENT_TIMESTAMP,pg_sleep(0.02)"))
  1665. with engine.begin() as newer:
  1666. first = newer.execute(statement, {
  1667. "audit_uid": _uid(), "fingerprint": fingerprint,
  1668. }).scalar_one()
  1669. second = older.execute(statement, {
  1670. "audit_uid": _uid(), "fingerprint": fingerprint,
  1671. }).scalar_one()
  1672. transaction.commit()
  1673. finally:
  1674. if transaction.is_active:
  1675. transaction.rollback()
  1676. older.close()
  1677. assert (first, second) == (1, 2)
  1678. def test_edge_flask_postgres_machine_binding_and_human_permissions(monkeypatch, engine):
  1679. monkeypatch.setenv("DATABASE_URL", DATABASE_URL)
  1680. actor = _uid()
  1681. monkeypatch.setattr(
  1682. "app.core.system.permissions.authenticate_request",
  1683. lambda: {"id": actor, "sub": actor, "roles": ["admin"]},
  1684. )
  1685. from app import create_app
  1686. app = create_app()
  1687. app.config.update(
  1688. TESTING=True,
  1689. EDGE_GATEWAY_SIGNING_PRIVATE_KEY=_SIGNING_PRIVATE_KEY,
  1690. EDGE_GATEWAY_SIGNING_KEY_ID=_SIGNING_KEY_ID,
  1691. EDGE_MTLS_TRUSTED_PROXY_IPS=("127.0.0.1",),
  1692. )
  1693. certificate_key = rsa.generate_private_key(public_exponent=65537, key_size=2048)
  1694. certificate_name = x509.Name(
  1695. [x509.NameAttribute(NameOID.COMMON_NAME, "plant-route-edge")]
  1696. )
  1697. now = datetime.now(UTC)
  1698. certificate = (
  1699. x509.CertificateBuilder()
  1700. .subject_name(certificate_name)
  1701. .issuer_name(certificate_name)
  1702. .public_key(certificate_key.public_key())
  1703. .serial_number(x509.random_serial_number())
  1704. .not_valid_before(now - timedelta(minutes=1))
  1705. .not_valid_after(now + timedelta(minutes=10))
  1706. .sign(certificate_key, hashes.SHA256())
  1707. )
  1708. certificate_pem = certificate.public_bytes(serialization.Encoding.PEM).decode()
  1709. certificate_sha256 = certificate.fingerprint(hashes.SHA256()).hex()
  1710. def mtls_headers(**values):
  1711. return {
  1712. "X-DataOps-Edge-Client-Cert": quote(certificate_pem, safe=""),
  1713. "X-DataOps-Edge-Client-Verify": "SUCCESS",
  1714. "X-Edge-Certificate-SHA256": certificate_sha256,
  1715. **values,
  1716. }
  1717. client = app.test_client()
  1718. oversized = client.post("/api/datasource/edge/enrollments", json={"padding": "x" * 270_000})
  1719. assert oversized.status_code == 413
  1720. enrollment_response = client.post("/api/datasource/edge/enrollments", json={
  1721. "gateway_name": "plant-route", "environment": "staging",
  1722. "network_zone": "zone-route", "policy_digest": "e" * 64,
  1723. "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [],
  1724. "expected_certificate_sha256": certificate_sha256,
  1725. "ttl_seconds": 600,
  1726. })
  1727. assert enrollment_response.status_code == 201
  1728. assert enrollment_response.headers["Cache-Control"] == "no-store"
  1729. enrollment = enrollment_response.get_json()["data"]
  1730. missing_tls_proof = client.post("/api/datasource/edge/register", headers={
  1731. "X-Edge-Enrollment": enrollment["enrollment_token"],
  1732. "X-Edge-Certificate-SHA256": certificate_sha256,
  1733. }, json={
  1734. "gateway_id": enrollment["gateway_id"], "environment": "staging",
  1735. "network_zone": "zone-route", "policy_digest": "e" * 64,
  1736. "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [],
  1737. "version": "3.0.0",
  1738. })
  1739. assert missing_tls_proof.status_code == 401
  1740. spoofed_tls_proof = client.post(
  1741. "/api/datasource/edge/register",
  1742. headers=mtls_headers(**{
  1743. "X-Edge-Enrollment": enrollment["enrollment_token"],
  1744. }),
  1745. environ_overrides={"REMOTE_ADDR": "203.0.113.7"},
  1746. json={
  1747. "gateway_id": enrollment["gateway_id"], "environment": "staging",
  1748. "network_zone": "zone-route", "policy_digest": "e" * 64,
  1749. "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [],
  1750. "version": "3.0.0",
  1751. },
  1752. )
  1753. assert spoofed_tls_proof.status_code == 401
  1754. register = client.post("/api/datasource/edge/register", headers=mtls_headers(**{
  1755. "X-Edge-Enrollment": enrollment["enrollment_token"],
  1756. }), environ_overrides={"REMOTE_ADDR": "127.0.0.1"}, json={
  1757. "gateway_id": enrollment["gateway_id"], "environment": "staging",
  1758. "network_zone": "zone-route", "policy_digest": "e" * 64,
  1759. "allowed_control_hosts": ["control.example.test"], "allowed_proxy_hosts": [],
  1760. "version": "3.0.0",
  1761. })
  1762. assert register.status_code == 201
  1763. assert register.headers["Cache-Control"] == "no-store"
  1764. registered = register.get_json()["data"]
  1765. bad = client.post(
  1766. f"/api/datasource/edge/gateways/{registered['gateway_id']}/heartbeat",
  1767. headers=mtls_headers(**{"X-Edge-Credential": registered["credential"] + "wrong"}),
  1768. environ_overrides={"REMOTE_ADDR": "127.0.0.1"},
  1769. json={"environment": "staging", "network_zone": "zone-route", "generation": 1, "version": "3.0.0", "safe_summary": {}},
  1770. )
  1771. assert bad.status_code == 401
  1772. assert bad.headers["Cache-Control"] == "no-store"
  1773. assert registered["credential"] not in bad.get_data(as_text=True)
  1774. limited = None
  1775. for _ in range(9):
  1776. limited = client.post(
  1777. f"/api/datasource/edge/gateways/{registered['gateway_id']}/heartbeat",
  1778. headers=mtls_headers(**{"X-Edge-Credential": registered["credential"] + "wrong"}),
  1779. environ_overrides={"REMOTE_ADDR": "127.0.0.1"},
  1780. json={"environment": "staging", "network_zone": "zone-route", "generation": 1, "version": "3.0.0", "safe_summary": {}},
  1781. )
  1782. assert limited is not None and limited.status_code == 429
  1783. assert limited.headers["Cache-Control"] == "no-store"
  1784. with engine.connect() as connection:
  1785. window = connection.execute(text("SELECT occurrence_count FROM edge_gateway_failure_windows")).scalar_one()
  1786. failure_rows = connection.execute(text("SELECT count(*) FROM edge_gateway_audits WHERE success=false")).scalar_one()
  1787. assert window == 10
  1788. assert failure_rows == 1
  1789. listed = client.get("/api/datasource/edge/gateways")
  1790. assert listed.status_code == 200
  1791. assert listed.get_json()["data"]["gateways"][0]["gateway_id"] == registered["gateway_id"]
  1792. assert client.get("/api/datasource/edge/gateways?limit=101").status_code == 400
  1793. oversized_reconcile = client.post(
  1794. f"/api/datasource/edge/gateways/{registered['gateway_id']}/reconcile",
  1795. headers=mtls_headers(**{"X-Edge-Credential": registered["credential"]}),
  1796. environ_overrides={"REMOTE_ADDR": "127.0.0.1"},
  1797. json={"environment": "staging", "network_zone": "zone-route", "generation": 1, "limit": 101},
  1798. )
  1799. assert oversized_reconcile.status_code == 400
  1800. app.config["EDGE_GATEWAY_SIGNING_PRIVATE_KEY"] = ""
  1801. unavailable_signer = client.post(
  1802. f"/api/datasource/edge/gateways/{registered['gateway_id']}/releases",
  1803. json={
  1804. "version": "7.0.0", "artifact_digest": "7" * 64,
  1805. "rollback_version": "6.0.0", "request_id": "no-production-signer",
  1806. },
  1807. )
  1808. assert unavailable_signer.status_code == 503
  1809. assert unavailable_signer.get_json()["error"]["code"] == "EDGE_GATEWAY_SIGNING_UNAVAILABLE"