| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249 |
- 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) == []
|