from __future__ import annotations import os import subprocess import uuid from concurrent.futures import ThreadPoolExecutor from datetime import UTC, datetime from pathlib import Path import pytest from sqlalchemy import create_engine, text from sqlalchemy.exc import DBAPIError pytestmark = pytest.mark.integration ROOT = Path(__file__).resolve().parents[2] def _url(): value = os.environ.get("TEST_DATABASE_URL") if not value: pytest.skip("TEST_DATABASE_URL is required") return value def _alembic(command: str, revision: str): environment = dict(os.environ) environment["MIGRATION_DATABASE_URL"] = os.environ.get("TEST_MIGRATION_DATABASE_URL", _url()) result = subprocess.run( [str(ROOT / ".venv/bin/alembic"), "-c", str(ROOT / "alembic.ini"), command, revision], cwd=ROOT, env=environment, capture_output=True, text=True, ) if result.returncode: raise AssertionError(result.stderr) def _cleanup(engine): with engine.begin() as connection: # The ledger is append-only to normal application identities. The test # migration operator temporarily bypasses only this trigger to remove # its own `wp06-test-*` fixtures; no non-namespaced evidence is touched. connection.execute(text("ALTER TABLE public.trusted_delivery_receipts DISABLE TRIGGER USER")) 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-%')")) connection.execute(text("ALTER TABLE public.trusted_delivery_receipts ENABLE TRIGGER USER")) connection.execute(text("""DO $$ BEGIN IF to_regclass('public.trusted_delivery_subscription_deliveries') IS NOT NULL THEN 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-%')); DELETE FROM public.trusted_delivery_subscriptions WHERE grant_uid IN (SELECT uid FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%'); END IF; END $$;""")) 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-%')")) connection.execute(text("DELETE FROM public.trusted_delivery_hold_releases WHERE idempotency_key LIKE 'wp06-test-%'")) connection.execute(text("DELETE FROM public.trusted_delivery_policy_transitions WHERE idempotency_key LIKE 'wp06-test-%'")) 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_%')")) connection.execute(text("DELETE FROM public.trusted_delivery_grants WHERE idempotency_key LIKE 'wp06-test-%'")) 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-%'")) connection.execute(text("DELETE FROM public.trusted_delivery_policy_versions WHERE code LIKE 'WP06_TEST_%'")) connection.execute(text("ALTER TABLE public.trusted_delivery_evidence DISABLE TRIGGER USER")) 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-%'")) connection.execute(text("ALTER TABLE public.trusted_delivery_evidence ENABLE TRIGGER USER")) if connection.execute(text("SELECT to_regclass('public.trusted_delivery_approval_facts')")).scalar_one() is not None: connection.execute(text("DELETE FROM public.trusted_delivery_approval_facts WHERE approval_ref LIKE 'wp06-test-%'")) connection.execute(text("DELETE FROM public.users WHERE username LIKE 'wp06-test-%'")) def _approval_fact(connection, *, ref, digest, actor, operation, receipt): connection.execute(text("""INSERT INTO public.trusted_delivery_approval_facts (approval_ref,approval_digest,actor_uid,operation,status,expires_at,signed_receipt_digest) VALUES(:ref,:digest,CAST(:actor AS uuid),:operation,'approved',clock_timestamp()+interval '1 day',:receipt)"""), { "ref": ref, "digest": digest, "actor": actor, "operation": operation, "receipt": receipt, }) @pytest.fixture(scope="module") def pg_engine(wp06_head): engine = create_engine(_url(), pool_pre_ping=True) _cleanup(engine) yield engine _cleanup(engine) engine.dispose() _alembic("upgrade", "head") def test_trusted_delivery_481_482_round_trip_is_empty_namespace_safe(pg_engine): """482 may only downgrade after its lifecycle records have been archived.""" _alembic("downgrade", "20260811_481") with pg_engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_481" assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_policy_transitions')")).scalar_one() is None assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_hold_releases')")).scalar_one() is None _alembic("upgrade", "20260811_482") with pg_engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_482" assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_policy_transitions')")).scalar_one() == "trusted_delivery_policy_transitions" assert connection.execute(text("SELECT to_regclass('public.trusted_delivery_hold_releases')")).scalar_one() == "trusted_delivery_hold_releases" reclaimed_constraint = connection.execute(text("SELECT 1 FROM pg_constraint WHERE conname='trusted_delivery_grants_status_check' AND pg_get_constraintdef(oid) LIKE '%reclaimed%'")) assert reclaimed_constraint.scalar_one() == 1 def test_trusted_delivery_480_to_head_round_trip_uses_the_restricted_migrator(pg_engine): """The fresh NO-SUPERUSER migrator must also restore the full WP06 head.""" _alembic("downgrade", "20260811_480") with pg_engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_480" _alembic("upgrade", "head") with pg_engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_495" def test_trusted_delivery_494_495_adjacent_round_trip_uses_the_restricted_migrator(pg_engine): """The final evidence-only hardening migration is independently reversible.""" _alembic("downgrade", "20260811_494") with pg_engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_494" _alembic("upgrade", "20260811_495") with pg_engine.connect() as connection: assert connection.execute(text("SELECT version_num FROM alembic_version")).scalar_one() == "20260811_495" def test_trusted_delivery_482_concurrency_restart_hold_release_and_reclaim(pg_engine): from app.core.system.trusted_delivery import ( ClosedProviderRegistry, TrustedDeliveryService, ) from app.core.system.trusted_delivery_repository import ( SqlAlchemyTrustedDeliveryRepository, ) actor = str(uuid.uuid4()) approver = str(uuid.uuid4()) hold_approver_one = str(uuid.uuid4()) hold_approver_two = str(uuid.uuid4()) policy_one = str(uuid.uuid4()) policy_two = str(uuid.uuid4()) asset = str(uuid.uuid4()) with pg_engine.begin() as connection: connection.execute(text("""INSERT INTO public.users (id,username,display_name,password_hash,status) VALUES (CAST(:id AS uuid),:name,:name,'test','active'), (CAST(:approver AS uuid),:approver_name,:approver_name,'test','active'), (CAST(:hold_one AS uuid),:hold_one_name,:hold_one_name,'test','active'), (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]}"}) _approval_fact(connection, ref="wp06-test-concurrent-hold-one", digest="f" * 64, actor=hold_approver_one, operation="hold", receipt="f" * 64) _approval_fact(connection, ref="wp06-test-concurrent-hold-two", digest="f" * 64, actor=hold_approver_two, operation="hold", receipt="f" * 64) _approval_fact(connection, ref="wp06-test-concurrent-release", digest="2" * 64, actor=approver, operation="release_hold", receipt="2" * 64) repository = SqlAlchemyTrustedDeliveryRepository(connection) for policy_uid, version, status in ((policy_one, "1.0.0", "active"), (policy_two, "2.0.0", "draft")): repository.create_policy_version({"uid": policy_uid, "code": "WP06_TEST_CONCURRENT", "version": version, "status": status, "selector": {}, "resource_rules": {}, "created_by": actor, "current_version": 1}) def activate(policy_uid, version, key): with pg_engine.begin() as connection: 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) with ThreadPoolExecutor(max_workers=2) as pool: 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")))) assert {result["status"] for result in results} == {"active"} with pg_engine.begin() as connection: repository = SqlAlchemyTrustedDeliveryRepository(connection) assert connection.execute(text("SELECT count(*) FROM public.trusted_delivery_policy_versions WHERE code='WP06_TEST_CONCURRENT' AND status='active'")).scalar_one() == 1 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) assert replay["uid"] == policy_two with pytest.raises(RuntimeError, match="idempotency conflict"): repository.transition_policy("WP06_TEST_CONCURRENT", "1.0.0", "activate", "wp06-test-transition-two", "d" * 64, approver, "approval:changed", "c" * 64) grant_uid = str(uuid.uuid4()) 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}) hold_uid = str(uuid.uuid4()) 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}) class Provider: calls = [] def revoke(self, envelope): self.calls.append(envelope) return {"status": "revoked", "receipt_code": "wp06-revoked", "response_digest": "1" * 64} provider = Provider() 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)) with pytest.raises(PermissionError, match="active legal hold"): service.execute_reclaim({"grant_uid": grant_uid, "idempotency_key": "wp06-test-reclaim"}, actor_uid=actor) with pytest.raises(PermissionError, match="independent"): 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) 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) assert released["status"] == "released" reclaimed = service.execute_reclaim({"grant_uid": grant_uid, "idempotency_key": "wp06-test-reclaim"}, actor_uid=actor) assert reclaimed["status"] == "reclaimed" assert len(provider.calls) == 1 restart_grant = str(uuid.uuid4()) 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}) delivery = repository.enqueue_delivery(restart_grant, "wp06-test-restart-delivery", "4" * 64) assert repository.claim_delivery(delivery["uid"], "worker-before-restart")["lease_fence"] == 1 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"]}) with pg_engine.begin() as visible, pg_engine.begin() as restarted: visible_status = visible.execute(text("SELECT status FROM public.trusted_delivery_grants WHERE idempotency_key='wp06-test-reclaim-grant'")) assert visible_status.scalar_one() == "reclaimed" reclaimed_delivery = SqlAlchemyTrustedDeliveryRepository(restarted).claim_delivery(delivery["uid"], "worker-after-restart") assert reclaimed_delivery["lease_fence"] == 2 def test_trusted_delivery_postgres_replay_fencing_expiry_hold_and_append_only_receipt(pg_engine): from app.core.system.trusted_delivery_repository import ( SqlAlchemyTrustedDeliveryRepository, ) actor = str(uuid.uuid4()) hold_approver_one = str(uuid.uuid4()) hold_approver_two = str(uuid.uuid4()) policy_uid = str(uuid.uuid4()) grant_uid = str(uuid.uuid4()) hold_uid = str(uuid.uuid4()) with pg_engine.begin() as connection: connection.execute(text("""INSERT INTO public.users (id,username,display_name,password_hash,status) VALUES (CAST(:id AS uuid),:name,:name,'test','active'), (CAST(:hold_one AS uuid),:hold_one_name,:hold_one_name,'test','active'), (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]}"}) _approval_fact(connection, ref="wp06-test-replay-hold-one", digest="f" * 64, actor=hold_approver_one, operation="hold", receipt="f" * 64) _approval_fact(connection, ref="wp06-test-replay-hold-two", digest="f" * 64, actor=hold_approver_two, operation="hold", receipt="f" * 64) repository = SqlAlchemyTrustedDeliveryRepository(connection) policy = repository.create_policy_version({ "uid": policy_uid, "code": "WP06_TEST_POLICY", "version": "1.0.0", "status": "active", "selector": {}, "resource_rules": {}, "created_by": actor, "current_version": 1, }) assert policy["code"] == "WP06_TEST_POLICY" grant = { "uid": grant_uid, "policy_uid": policy_uid, "subject_uid": actor, "asset_uid": str(uuid.uuid4()), "provider": "database", "idempotency_key": "wp06-test-grant-1", "request_digest": "a" * 64, "expires_at": "2026-08-11T10:00:00+00:00", "status": "active", "reason_code": "approval_bound", "current_version": 1, "created_by": actor, } assert repository.enqueue_grant(grant)["uid"] == grant_uid assert repository.enqueue_grant(grant)["uid"] == grant_uid conflict = {**grant, "request_digest": "b" * 64} with pytest.raises(RuntimeError, match="idempotency conflict"): repository.enqueue_grant(conflict) delivery = repository.enqueue_delivery(grant_uid, "wp06-test-delivery-1", "c" * 64) claimed = repository.claim_delivery(delivery["uid"], "worker-a") assert claimed["lease_fence"] == 1 with pytest.raises(RuntimeError, match="lease conflict"): repository.complete_delivery(delivery["uid"], "worker-b", 1, "applied", "receipt-1", "d" * 64, "e" * 64) completed = repository.complete_delivery(delivery["uid"], "worker-a", 1, "applied", "receipt-1", "d" * 64, "e" * 64) assert completed["status"] == "applied" savepoint = connection.begin_nested() with pytest.raises(DBAPIError): connection.execute(text("UPDATE public.trusted_delivery_receipts SET receipt_code='mutated' WHERE delivery_uid=CAST(:uid AS uuid)"), {"uid": delivery["uid"]}) savepoint.rollback() 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}) assert hold["status"] == "active" assert repository.expire_grants("2026-08-12T00:00:00+00:00", actor) == []