test_phase3_wp06_trusted_delivery_postgres.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249
  1. from __future__ import annotations
  2. import os
  3. import subprocess
  4. import uuid
  5. from concurrent.futures import ThreadPoolExecutor
  6. from datetime import UTC, datetime
  7. from pathlib import Path
  8. import pytest
  9. from sqlalchemy import create_engine, text
  10. from sqlalchemy.exc import DBAPIError
  11. pytestmark = pytest.mark.integration
  12. ROOT = Path(__file__).resolve().parents[2]
  13. def _url():
  14. value = os.environ.get("TEST_DATABASE_URL")
  15. if not value:
  16. pytest.skip("TEST_DATABASE_URL is required")
  17. return value
  18. def _alembic(command: str, revision: str):
  19. environment = dict(os.environ)
  20. environment["MIGRATION_DATABASE_URL"] = os.environ.get("TEST_MIGRATION_DATABASE_URL", _url())
  21. result = subprocess.run(
  22. [str(ROOT / ".venv/bin/alembic"), "-c", str(ROOT / "alembic.ini"), command, revision],
  23. cwd=ROOT, env=environment, capture_output=True, text=True,
  24. )
  25. if result.returncode:
  26. raise AssertionError(result.stderr)
  27. def _cleanup(engine):
  28. with engine.begin() as connection:
  29. # The ledger is append-only to normal application identities. The test
  30. # migration operator temporarily bypasses only this trigger to remove
  31. # its own `wp06-test-*` fixtures; no non-namespaced evidence is touched.
  32. connection.execute(text("ALTER TABLE public.trusted_delivery_receipts DISABLE TRIGGER USER"))
  33. connection.execute(text("DELETE FROM public.trusted_delivery_receipts WHERE grant_uid IN (SELECT uid FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%')"))
  34. connection.execute(text("ALTER TABLE public.trusted_delivery_receipts ENABLE TRIGGER USER"))
  35. connection.execute(text("""DO $$ BEGIN
  36. IF to_regclass('public.trusted_delivery_subscription_deliveries') IS NOT NULL THEN
  37. DELETE FROM public.trusted_delivery_subscription_deliveries WHERE subscription_uid IN (SELECT uid FROM public.trusted_delivery_subscriptions WHERE grant_uid IN (SELECT uid FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%'));
  38. DELETE FROM public.trusted_delivery_subscriptions WHERE grant_uid IN (SELECT uid FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%');
  39. END IF;
  40. END $$;"""))
  41. connection.execute(text("DELETE FROM public.trusted_delivery_deliveries WHERE grant_uid IN (SELECT uid FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%')"))
  42. connection.execute(text("DELETE FROM public.trusted_delivery_hold_releases WHERE idempotency_key LIKE 'wp06-test-%'"))
  43. connection.execute(text("DELETE FROM public.trusted_delivery_policy_transitions WHERE idempotency_key LIKE 'wp06-test-%'"))
  44. connection.execute(text("DELETE FROM public.trusted_delivery_policy_transitions WHERE idempotency_key LIKE 'bootstrap:%' AND policy_uid IN (SELECT uid FROM public.trusted_delivery_policy_versions WHERE code LIKE 'WP06_TEST_%')"))
  45. connection.execute(text("DELETE FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%'"))
  46. connection.execute(text("DELETE FROM public.trusted_delivery_legal_holds h USING public.users u WHERE h.created_by=u.id AND u.username LIKE 'wp06-test-%'"))
  47. connection.execute(text("DELETE FROM public.trusted_delivery_policy_versions WHERE code LIKE 'WP06_TEST_%'"))
  48. connection.execute(text("ALTER TABLE public.trusted_delivery_evidence DISABLE TRIGGER USER"))
  49. connection.execute(text("DELETE FROM public.trusted_delivery_evidence e USING public.users u WHERE e.actor_uid=u.id AND u.username LIKE 'wp06-test-%'"))
  50. connection.execute(text("ALTER TABLE public.trusted_delivery_evidence ENABLE TRIGGER USER"))
  51. if connection.execute(text("SELECT to_regclass('public.trusted_delivery_approval_facts')")).scalar_one() is not None:
  52. connection.execute(text("DELETE FROM public.trusted_delivery_approval_facts WHERE approval_ref LIKE 'wp06-test-%'"))
  53. connection.execute(text("DELETE FROM public.users WHERE username LIKE 'wp06-test-%'"))
  54. def _approval_fact(connection, *, ref, digest, actor, operation, receipt):
  55. connection.execute(text("""INSERT INTO public.trusted_delivery_approval_facts
  56. (approval_ref,approval_digest,actor_uid,operation,status,expires_at,signed_receipt_digest)
  57. VALUES(:ref,:digest,CAST(:actor AS uuid),:operation,'approved',clock_timestamp()+interval '1 day',:receipt)"""), {
  58. "ref": ref, "digest": digest, "actor": actor, "operation": operation,
  59. "receipt": receipt,
  60. })
  61. @pytest.fixture(scope="module")
  62. def pg_engine(wp06_head):
  63. engine = create_engine(_url(), pool_pre_ping=True)
  64. _cleanup(engine)
  65. yield engine
  66. _cleanup(engine)
  67. engine.dispose()
  68. _alembic("upgrade", "head")
  69. def test_trusted_delivery_481_482_round_trip_is_empty_namespace_safe(pg_engine):
  70. """482 may only downgrade after its lifecycle records have been archived."""
  71. _alembic("downgrade", "20260811_481")
  72. with pg_engine.connect() as connection:
  73. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_481"
  74. assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_policy_transitions')")).scalar_one() is None
  75. assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_hold_releases')")).scalar_one() is None
  76. _alembic("upgrade", "20260811_482")
  77. with pg_engine.connect() as connection:
  78. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_482"
  79. assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_policy_transitions')")).scalar_one() == "trusted_delivery_policy_transitions"
  80. assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_hold_releases')")).scalar_one() == "trusted_delivery_hold_releases"
  81. reclaimed_constraint = connection.execute(text("SELECT 1 FROM pg_constraint WHERE conname='trusted_delivery_grants_status_check' AND pg_get_constraintdef(oid) LIKE '%reclaimed%'"))
  82. assert reclaimed_constraint.scalar_one() == 1
  83. def test_trusted_delivery_480_to_head_round_trip_uses_the_restricted_migrator(pg_engine):
  84. """The fresh NO-SUPERUSER migrator must also restore the full WP06 head."""
  85. _alembic("downgrade", "20260811_480")
  86. with pg_engine.connect() as connection:
  87. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_480"
  88. _alembic("upgrade", "head")
  89. with pg_engine.connect() as connection:
  90. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_495"
  91. def test_trusted_delivery_494_495_adjacent_round_trip_uses_the_restricted_migrator(pg_engine):
  92. """The final evidence-only hardening migration is independently reversible."""
  93. _alembic("downgrade", "20260811_494")
  94. with pg_engine.connect() as connection:
  95. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_494"
  96. _alembic("upgrade", "20260811_495")
  97. with pg_engine.connect() as connection:
  98. assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_495"
  99. def test_trusted_delivery_482_concurrency_restart_hold_release_and_reclaim(pg_engine):
  100. from app.core.system.trusted_delivery import (
  101. ClosedProviderRegistry,
  102. TrustedDeliveryService,
  103. )
  104. from app.core.system.trusted_delivery_repository import (
  105. SqlAlchemyTrustedDeliveryRepository,
  106. )
  107. actor = str(uuid.uuid4())
  108. approver = str(uuid.uuid4())
  109. hold_approver_one = str(uuid.uuid4())
  110. hold_approver_two = str(uuid.uuid4())
  111. policy_one = str(uuid.uuid4())
  112. policy_two = str(uuid.uuid4())
  113. asset = str(uuid.uuid4())
  114. with pg_engine.begin() as connection:
  115. connection.execute(text("""INSERT INTO public.users (id,username,display_name,password_hash,status) VALUES
  116. (CAST(:id AS uuid),:name,:name,'test','active'),
  117. (CAST(:approver AS uuid),:approver_name,:approver_name,'test','active'),
  118. (CAST(:hold_one AS uuid),:hold_one_name,:hold_one_name,'test','active'),
  119. (CAST(:hold_two AS uuid),:hold_two_name,:hold_two_name,'test','active')"""), {"id": actor, "name": f"wp06-test-{actor[:8]}", "approver": approver, "approver_name": f"wp06-test-{approver[:8]}", "hold_one": hold_approver_one, "hold_one_name": f"wp06-test-{hold_approver_one[:8]}", "hold_two": hold_approver_two, "hold_two_name": f"wp06-test-{hold_approver_two[:8]}"})
  120. _approval_fact(connection, ref="wp06-test-concurrent-hold-one", digest="f" * 64, actor=hold_approver_one, operation="hold", receipt="f" * 64)
  121. _approval_fact(connection, ref="wp06-test-concurrent-hold-two", digest="f" * 64, actor=hold_approver_two, operation="hold", receipt="f" * 64)
  122. _approval_fact(connection, ref="wp06-test-concurrent-release", digest="2" * 64, actor=approver, operation="release_hold", receipt="2" * 64)
  123. repository = SqlAlchemyTrustedDeliveryRepository(connection)
  124. for policy_uid, version, status in ((policy_one, "1.0.0", "active"), (policy_two, "2.0.0", "draft")):
  125. repository.create_policy_version({"uid": policy_uid, "code": "WP06_TEST_CONCURRENT", "version": version, "status": status, "selector": {}, "resource_rules": {}, "created_by": actor, "current_version": 1})
  126. def activate(policy_uid, version, key):
  127. with pg_engine.begin() as connection:
  128. return SqlAlchemyTrustedDeliveryRepository(connection).transition_policy("WP06_TEST_CONCURRENT", version, "activate", key, "a" * 64 if version == "1.0.0" else "b" * 64, approver, f"approval:{key}", "c" * 64)
  129. with ThreadPoolExecutor(max_workers=2) as pool:
  130. results = list(pool.map(lambda values: activate(*values), ((policy_one, "1.0.0", "wp06-test-transition-one"), (policy_two, "2.0.0", "wp06-test-transition-two"))))
  131. assert {result["status"] for result in results} == {"active"}
  132. with pg_engine.begin() as connection:
  133. repository = SqlAlchemyTrustedDeliveryRepository(connection)
  134. assert connection.execute(text("SELECT count(*) FROM public.trusted_delivery_policy_versions WHERE code='WP06_TEST_CONCURRENT' AND status='active'")).scalar_one() == 1
  135. replay = repository.transition_policy("WP06_TEST_CONCURRENT", "2.0.0", "activate", "wp06-test-transition-two", "b" * 64, approver, "approval:wp06-test-transition-two", "c" * 64)
  136. assert replay["uid"] == policy_two
  137. with pytest.raises(RuntimeError, match="idempotency conflict"):
  138. repository.transition_policy("WP06_TEST_CONCURRENT", "1.0.0", "activate", "wp06-test-transition-two", "d" * 64, approver, "approval:changed", "c" * 64)
  139. grant_uid = str(uuid.uuid4())
  140. repository.enqueue_grant({"uid": grant_uid, "policy_uid": policy_two, "subject_uid": actor, "asset_uid": asset, "provider": "database", "idempotency_key": "wp06-test-reclaim-grant", "request_digest": "e" * 64, "expires_at": "2026-08-11T10:00:00+00:00", "status": "active", "reason_code": "approval_bound", "current_version": 1, "created_by": actor})
  141. hold_uid = str(uuid.uuid4())
  142. repository.create_legal_hold({"uid": hold_uid, "asset_uid": asset, "operation": "freeze", "status": "active", "approver_ref_one": "wp06-test-concurrent-hold-one", "approver_ref_two": "wp06-test-concurrent-hold-two", "evidence_digest": "f" * 64, "created_by": actor, "current_version": 1})
  143. class Provider:
  144. calls = []
  145. def revoke(self, envelope):
  146. self.calls.append(envelope)
  147. return {"status": "revoked", "receipt_code": "wp06-revoked", "response_digest": "1" * 64}
  148. provider = Provider()
  149. service = TrustedDeliveryService(repository, provider_registry=ClosedProviderRegistry.for_tests({"database": provider}), uid_factory=lambda: str(uuid.uuid4()), now_factory=lambda: datetime(2026, 8, 12, tzinfo=UTC))
  150. with pytest.raises(PermissionError, match="active legal hold"):
  151. service.execute_reclaim({"grant_uid": grant_uid, "idempotency_key": "wp06-test-reclaim"}, actor_uid=actor)
  152. with pytest.raises(PermissionError, match="independent"):
  153. service.release_legal_hold({"hold_uid": hold_uid, "approval_ref": "wp06-test-concurrent-release", "approval_digest": "2" * 64, "idempotency_key": "wp06-test-release"}, actor_uid=actor)
  154. released = service.release_legal_hold({"hold_uid": hold_uid, "approval_ref": "wp06-test-concurrent-release", "approval_digest": "2" * 64, "idempotency_key": "wp06-test-release"}, actor_uid=approver)
  155. assert released["status"] == "released"
  156. reclaimed = service.execute_reclaim({"grant_uid": grant_uid, "idempotency_key": "wp06-test-reclaim"}, actor_uid=actor)
  157. assert reclaimed["status"] == "reclaimed"
  158. assert len(provider.calls) == 1
  159. restart_grant = str(uuid.uuid4())
  160. repository.enqueue_grant({"uid": restart_grant, "policy_uid": policy_two, "subject_uid": actor, "asset_uid": str(uuid.uuid4()), "provider": "database", "idempotency_key": "wp06-test-restart-grant", "request_digest": "3" * 64, "expires_at": "2026-08-12T10:00:00+00:00", "status": "active", "reason_code": "approval_bound", "current_version": 1, "created_by": actor})
  161. delivery = repository.enqueue_delivery(restart_grant, "wp06-test-restart-delivery", "4" * 64)
  162. assert repository.claim_delivery(delivery["uid"], "worker-before-restart")["lease_fence"] == 1
  163. connection.execute(text("UPDATE public.trusted_delivery_deliveries SET lease_expires_at=clock_timestamp()-interval '1 second' WHERE uid=CAST(:uid AS uuid)"), {"uid": delivery["uid"]})
  164. with pg_engine.begin() as visible, pg_engine.begin() as restarted:
  165. visible_status = visible.execute(text("SELECT status FROM public.trusted_delivery_grants WHERE idempotency_key='wp06-test-reclaim-grant'"))
  166. assert visible_status.scalar_one() == "reclaimed"
  167. reclaimed_delivery = SqlAlchemyTrustedDeliveryRepository(restarted).claim_delivery(delivery["uid"], "worker-after-restart")
  168. assert reclaimed_delivery["lease_fence"] == 2
  169. def test_trusted_delivery_postgres_replay_fencing_expiry_hold_and_append_only_receipt(pg_engine):
  170. from app.core.system.trusted_delivery_repository import (
  171. SqlAlchemyTrustedDeliveryRepository,
  172. )
  173. actor = str(uuid.uuid4())
  174. hold_approver_one = str(uuid.uuid4())
  175. hold_approver_two = str(uuid.uuid4())
  176. policy_uid = str(uuid.uuid4())
  177. grant_uid = str(uuid.uuid4())
  178. hold_uid = str(uuid.uuid4())
  179. with pg_engine.begin() as connection:
  180. connection.execute(text("""INSERT INTO public.users (id,username,display_name,password_hash,status) VALUES
  181. (CAST(:id AS uuid),:name,:name,'test','active'),
  182. (CAST(:hold_one AS uuid),:hold_one_name,:hold_one_name,'test','active'),
  183. (CAST(:hold_two AS uuid),:hold_two_name,:hold_two_name,'test','active')"""), {"id": actor, "name": f"wp06-test-{actor[:8]}", "hold_one": hold_approver_one, "hold_one_name": f"wp06-test-{hold_approver_one[:8]}", "hold_two": hold_approver_two, "hold_two_name": f"wp06-test-{hold_approver_two[:8]}"})
  184. _approval_fact(connection, ref="wp06-test-replay-hold-one", digest="f" * 64, actor=hold_approver_one, operation="hold", receipt="f" * 64)
  185. _approval_fact(connection, ref="wp06-test-replay-hold-two", digest="f" * 64, actor=hold_approver_two, operation="hold", receipt="f" * 64)
  186. repository = SqlAlchemyTrustedDeliveryRepository(connection)
  187. policy = repository.create_policy_version({
  188. "uid": policy_uid, "code": "WP06_TEST_POLICY", "version": "1.0.0", "status": "active",
  189. "selector": {}, "resource_rules": {}, "created_by": actor, "current_version": 1,
  190. })
  191. assert policy["code"] == "WP06_TEST_POLICY"
  192. grant = {
  193. "uid": grant_uid, "policy_uid": policy_uid, "subject_uid": actor, "asset_uid": str(uuid.uuid4()),
  194. "provider": "database", "idempotency_key": "wp06-test-grant-1", "request_digest": "a" * 64,
  195. "expires_at": "2026-08-11T10:00:00+00:00", "status": "active", "reason_code": "approval_bound",
  196. "current_version": 1, "created_by": actor,
  197. }
  198. assert repository.enqueue_grant(grant)["uid"] == grant_uid
  199. assert repository.enqueue_grant(grant)["uid"] == grant_uid
  200. conflict = {**grant, "request_digest": "b" * 64}
  201. with pytest.raises(RuntimeError, match="idempotency conflict"):
  202. repository.enqueue_grant(conflict)
  203. delivery = repository.enqueue_delivery(grant_uid, "wp06-test-delivery-1", "c" * 64)
  204. claimed = repository.claim_delivery(delivery["uid"], "worker-a")
  205. assert claimed["lease_fence"] == 1
  206. with pytest.raises(RuntimeError, match="lease conflict"):
  207. repository.complete_delivery(delivery["uid"], "worker-b", 1, "applied", "receipt-1", "d" * 64, "e" * 64)
  208. completed = repository.complete_delivery(delivery["uid"], "worker-a", 1, "applied", "receipt-1", "d" * 64, "e" * 64)
  209. assert completed["status"] == "applied"
  210. savepoint = connection.begin_nested()
  211. with pytest.raises(DBAPIError):
  212. connection.execute(text("UPDATE public.trusted_delivery_receipts SET receipt_code='mutated' WHERE delivery_uid=CAST(:uid AS uuid)"), {"uid": delivery["uid"]})
  213. savepoint.rollback()
  214. hold = repository.create_legal_hold({"uid": hold_uid, "asset_uid": grant["asset_uid"], "operation": "freeze", "status": "active", "approver_ref_one": "wp06-test-replay-hold-one", "approver_ref_two": "wp06-test-replay-hold-two", "evidence_digest": "f" * 64, "created_by": actor, "current_version": 1})
  215. assert hold["status"] == "active"
  216. assert repository.expire_grants("2026-08-12T00:00:00+00:00", actor) == []