from __future__ import annotations import os import uuid from concurrent.futures import ThreadPoolExecutor import pytest from sqlalchemy import create_engine, text from sqlalchemy.engine import make_url from sqlalchemy.exc import SQLAlchemyError from app.core.system.tenant_control_repository import SqlAlchemyTenantControlRepository from app.core.system.tenant_control_service import ( TenantControlError, TenantControlService, ) pytestmark = pytest.mark.integration def test_wp10_real_flask_claim_route_uses_dedicated_control_login_and_commits(monkeypatch): database_url = os.getenv("TEST_DATABASE_URL") if not database_url: pytest.skip("TEST_DATABASE_URL is required") tenant_id = f"wp10-http-{uuid.uuid4().hex[:14]}" concurrent_tenant_id = f"wp10-http-concurrent-{uuid.uuid4().hex[:10]}" tenant_ids = (tenant_id, concurrent_tenant_id) principal_id = str(uuid.uuid4()) concurrent_principal_id = str(uuid.uuid4()) control_name = f"wp10_control_{uuid.uuid4().hex[:12]}" control_password = f"wp10_{uuid.uuid4().hex}" approver_name = f"wp10_approver_{uuid.uuid4().hex[:12]}" approver_password = f"wp10_{uuid.uuid4().hex}" admin = create_engine(database_url, pool_pre_ping=True) try: with admin.begin() as connection: quoted_control = connection.dialect.identifier_preparer.quote(control_name) connection.execute(text( f"CREATE ROLE {quoted_control} LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION PASSWORD :password" ), {"password": control_password}) connection.execute(text(f"GRANT dataops_tenant_control TO {quoted_control}")) quoted_approver = connection.dialect.identifier_preparer.quote(approver_name) connection.execute(text(f"CREATE ROLE {quoted_approver} LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION PASSWORD :password"), {"password": approver_password}) connection.execute(text(f"GRANT dataops_tenant_approver TO {quoted_approver}")) control_url = make_url(database_url).set(username=control_name, password=control_password).render_as_string(hide_password=False) approver_url = make_url(database_url).set(username=approver_name, password=approver_password).render_as_string(hide_password=False) monkeypatch.setenv("TENANT_CONTROL_DATABASE_URL", control_url) monkeypatch.setenv("TENANT_APPROVER_DATABASE_URL", approver_url) monkeypatch.setenv("TENANT_APPROVER_ISSUER_ID", str(uuid.uuid4())) monkeypatch.setenv("PRIVATE_SINGLE_TENANT_ID", tenant_id) monkeypatch.setenv("PRIVATE_SINGLE_TENANT_HOST", "private.example") monkeypatch.setenv("TRUSTED_TENANT_ROUTE_HOST", "private.example") monkeypatch.setattr( "app.core.system.auth.load_identity_from_token", lambda token, secret: {"id": principal_id, "roles": [token.removeprefix("wp10-")]} if token in {"wp10-editor", "wp10-admin"} else None, ) from app import create_app app = create_app() app.config.update(TESTING=True) client = app.test_client() headers = {"Authorization": "Bearer wp10-editor", "Host": "private.example"} # Diagnostic RED boundary: this is the same dedicated control login # used by Flask, so a failure below exposes PostgreSQL's class/code. with pytest.raises(TenantControlError, match="tenant_control_provision_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private( principal_id="not-a-uuid", host="ignored.example", request={"idempotency_key": "diagnostic-invalid"} ) provision_key = f"provision-{tenant_id}" provision = client.post("/api/system/tenant/provisions", json={"idempotency_key": provision_key}, headers={"Authorization": "Bearer wp10-admin", "Host": "private.example"}) assert provision.status_code == 201 assert provision.headers["Cache-Control"] == "no-store" assert provision.get_json()["data"]["tenant_id"] == tenant_id # A fresh Flask application process must observe the durable exact # replay, not create a second membership or trust the HTTP Host. replay_app = create_app() replay_app.config.update(TESTING=True) replay = replay_app.test_client().post( "/api/system/tenant/provisions", json={"idempotency_key": provision_key}, headers={"Authorization": "Bearer wp10-admin", "Host": "attacker.example"}, ) assert replay.status_code == 201 assert replay.get_json()["data"]["replayed"] is True # The same durable key cannot be rebound to a different trusted # principal, host, or configured private tenant. with pytest.raises(TenantControlError, match="tenant_control_provision_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private( principal_id=str(uuid.uuid4()), host="ignored.example", request={"idempotency_key": provision_key} ) monkeypatch.setenv("PRIVATE_SINGLE_TENANT_HOST", "other.example") with pytest.raises(TenantControlError, match="tenant_control_provision_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private( principal_id=principal_id, host="ignored.example", request={"idempotency_key": provision_key} ) monkeypatch.setenv("PRIVATE_SINGLE_TENANT_HOST", "private.example") monkeypatch.setenv("PRIVATE_SINGLE_TENANT_ID", f"other-{tenant_id}") with pytest.raises(TenantControlError, match="tenant_control_provision_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private( principal_id=principal_id, host="ignored.example", request={"idempotency_key": provision_key} ) monkeypatch.setenv("PRIVATE_SINGLE_TENANT_ID", tenant_id) # Two independent control connections race the same first provision. # The durable idempotency record/advisory lock leaves one creation and # one exact replay; it must never add a second membership. concurrent_key = f"concurrent-{uuid.uuid4().hex}" def concurrent_provision(): return SqlAlchemyTenantControlRepository(control_url).provision( tenant_id=concurrent_tenant_id, principal_id=concurrent_principal_id, host="private.example", idempotency_key=concurrent_key, ) with ThreadPoolExecutor(max_workers=2) as pool: concurrent_results = list(pool.map(lambda _: concurrent_provision(), range(2))) assert {item["tenant_id"] for item in concurrent_results} == {concurrent_tenant_id} assert sum(bool(item["replayed"]) for item in concurrent_results) == 1 with admin.connect() as check: assert check.execute( text("SELECT count(*) FROM public.tenant_control_memberships WHERE tenant_id=:tenant_id"), {"tenant_id": concurrent_tenant_id}, ).scalar_one() == 1 rejected = client.post("/api/system/tenant/quota-claims", json={"tenant_id": "other", "quota_name": "records", "amount": "1", "idempotency_key": "http-1"}, headers=headers) assert rejected.status_code == 400 created = client.post("/api/system/tenant/quota-claims", json={"quota_name": "records", "amount": "1", "idempotency_key": "http-1"}, headers=headers) assert created.status_code == 201 assert created.headers["Cache-Control"] == "no-store" assert created.get_json()["data"]["tenant_id"] == tenant_id propagation = created.get_json()["data"]["propagation"] assert propagation["task"]["tenant_id"] == tenant_id assert propagation["event"]["event_type"] == "tenant.quota.claim" assert all(value.startswith(f"{tenant_id}/") for value in propagation["namespaces"].values()) with admin.connect() as check: count = check.execute(text("SELECT count(*) FROM public.tenant_control_claims WHERE tenant_id=:tenant_id AND status='issued'"), {"tenant_id": tenant_id}).scalar_one() check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) outbox_uid = str(check.execute(text("SELECT outbox_uid FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id AND event_type='tenant.quota.claim'"), {"tenant_id": tenant_id}).scalar_one()) assert count == 1 monkeypatch.setenv("TENANT_CONTROL_WORKER_ID", "tenant-worker-1") with create_engine(control_url).connect() as legacy_gateway, pytest.raises(SQLAlchemyError): legacy_gateway.execute(text("SELECT public.tenant_control_lifecycle_mutate(CAST(:claim_uid AS uuid),CAST(:nonce AS uuid),'{}'::jsonb)"), {"claim_uid": str(uuid.uuid4()), "nonce": str(uuid.uuid4())}) worker = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=outbox_uid) assert worker["tenant_id"] == tenant_id assert worker["propagation"]["event"]["tenant_id"] == tenant_id assert all(value.startswith(f"{tenant_id}/") for value in worker["propagation"]["namespaces"].values()) settled = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox( outbox_uid=outbox_uid, lease_fence=worker["lease_fence"], payload_digest=worker["payload_digest"], outcome="complete" ) assert settled["state"] == "completed" assert TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox( outbox_uid=outbox_uid, lease_fence=worker["lease_fence"], payload_digest=worker["payload_digest"], outcome="complete" )["replayed"] is True with admin.connect() as check: check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) assert check.execute(text("SELECT status FROM public.tenant_quota_reservations WHERE tenant_id=:tenant_id AND quota_name='records' AND idempotency_key='http-1'"), {"tenant_id": tenant_id}).scalar_one() == "settled" retry_claim = client.post("/api/system/tenant/quota-claims", json={"quota_name": "background_tasks", "amount": "1", "idempotency_key": "http-2"}, headers=headers) assert retry_claim.status_code == 201 with admin.connect() as check: check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) retry_outbox_uid = str(check.execute(text("SELECT outbox_uid FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id AND state='queued' ORDER BY created_at DESC LIMIT 1"), {"tenant_id": tenant_id}).scalar_one()) retry_fences = [] for attempt in range(3): retry_worker = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=retry_outbox_uid) retry_fences.append(retry_worker["lease_fence"]) state = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox( outbox_uid=retry_outbox_uid, lease_fence=retry_worker["lease_fence"], payload_digest=retry_worker["payload_digest"], outcome="fail", )["state"] assert state == ("failed" if attempt == 2 else "queued") assert retry_fences == sorted(retry_fences) with admin.connect() as check: check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) assert check.execute(text("SELECT status FROM public.tenant_quota_reservations WHERE tenant_id=:tenant_id AND quota_name='background_tasks' AND idempotency_key='http-2'"), {"tenant_id": tenant_id}).scalar_one() == "released" with pytest.raises(TenantControlError, match="tenant_worker_scope_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=retry_outbox_uid) # A DB-clock-expired lease cannot complete; a fresh repository/engine # acts as a worker restart and can reclaim only with a higher fence. expiry_claim = client.post("/api/system/tenant/quota-claims", json={"quota_name": "events", "amount": "1", "idempotency_key": "http-expiry"}, headers=headers) assert expiry_claim.status_code == 201 with admin.connect() as check: check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) expiry_outbox_uid = str(check.execute(text("SELECT outbox_uid FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id AND state='queued' ORDER BY created_at DESC LIMIT 1"), {"tenant_id": tenant_id}).scalar_one()) expired_worker = SqlAlchemyTenantControlRepository(control_url).claim_outbox(outbox_uid=expiry_outbox_uid, worker_id="tenant-worker-1", lease_seconds="1") with admin.begin() as wait_for_db_clock: wait_for_db_clock.execute(text("SELECT pg_sleep(1.1)")) with pytest.raises(TenantControlError, match="tenant_worker_settle_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(outbox_uid=expiry_outbox_uid, lease_fence=expired_worker["lease_fence"], payload_digest=expired_worker["payload_digest"], outcome="complete") restarted_worker = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=expiry_outbox_uid) assert restarted_worker["lease_fence"] > expired_worker["lease_fence"] with pytest.raises(TenantControlError, match="tenant_worker_settle_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(outbox_uid=expiry_outbox_uid, lease_fence=restarted_worker["lease_fence"], payload_digest="0" * 64, outcome="complete") assert TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(outbox_uid=expiry_outbox_uid, lease_fence=restarted_worker["lease_fence"], payload_digest=restarted_worker["payload_digest"], outcome="complete")["state"] == "completed" worker_race_claim = client.post("/api/system/tenant/quota-claims", json={"quota_name": "models", "amount": "1", "idempotency_key": "http-worker-race"}, headers=headers) assert worker_race_claim.status_code == 201 with admin.connect() as check: check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) worker_race_outbox_uid = str(check.execute(text("SELECT outbox_uid FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id AND state='queued' ORDER BY created_at DESC LIMIT 1"), {"tenant_id": tenant_id}).scalar_one()) def race_worker_claim(): try: return "claimed", SqlAlchemyTenantControlRepository(control_url).claim_outbox( outbox_uid=worker_race_outbox_uid, worker_id="tenant-worker-race" ) except TenantControlError: return "denied", None with ThreadPoolExecutor(max_workers=2) as pool: worker_race_results = list(pool.map(lambda _: race_worker_claim(), range(2))) assert [outcome for outcome, _ in worker_race_results].count("claimed") == 1 claimed_worker = next(value for outcome, value in worker_race_results if outcome == "claimed") monkeypatch.setenv("TENANT_CONTROL_WORKER_ID", "tenant-worker-race") assert TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(outbox_uid=worker_race_outbox_uid, lease_fence=claimed_worker["lease_fence"], payload_digest=claimed_worker["payload_digest"], outcome="complete")["state"] == "completed" monkeypatch.setenv("TENANT_CONTROL_WORKER_ID", "tenant-worker-1") admin_headers = {"Authorization": "Bearer wp10-admin", "Host": "private.example"} frozen = client.post("/api/system/tenant/lifecycle/freeze", json={"expected_fence": "0", "idempotency_key": "freeze-1"}, headers=admin_headers) assert frozen.status_code == 200 assert frozen.headers["Cache-Control"] == "no-store" with admin.connect() as check: check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id}) frozen_outbox_uid = str(check.execute(text("SELECT outbox_uid FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id AND event_type='tenant.lifecycle.claim' ORDER BY created_at DESC LIMIT 1"), {"tenant_id": tenant_id}).scalar_one()) with pytest.raises(TenantControlError, match="tenant_worker_scope_denied"): TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=frozen_outbox_uid) status = client.get("/api/system/tenant/status", headers=admin_headers) assert status.status_code == 200 assert status.get_json()["data"]["state"] == "frozen" assert client.get("/api/system/tenant/audit", headers=admin_headers).status_code == 200 assert client.get("/api/system/tenant/manifests", headers=admin_headers).get_json()["data"]["manifests"] == [] def approval(operation, fence, key, backup=None, retention=None): body = {"operation": operation, "expected_fence": str(fence), "idempotency_key": key} if backup is not None: body["backup_digest"] = backup if retention is not None: body["retention_seconds"] = retention response = client.post("/api/system/tenant/approvals", json=body, headers=admin_headers) assert response.status_code == 201 assert response.headers["Cache-Control"] == "no-store" return response.get_json()["data"]["approval_ref"] recover_approval = approval("begin_recovery", 1, "recover-1") assert client.post("/api/system/tenant/lifecycle/recover", json={"expected_fence": "1", "idempotency_key": "recover-1", "approval_ref": recover_approval}, headers=admin_headers).status_code == 200 assert client.post("/api/system/tenant/lifecycle/activate", json={"expected_fence": "2", "idempotency_key": "activate-1", "approval_ref": "approval-2"}, headers=admin_headers).status_code == 200 candidate = {"expected_fence": "3", "idempotency_key": "candidate-1", "approval_ref": "approval-3", "backup_digest": "a" * 64, "retention_seconds": "60"} assert client.post("/api/system/tenant/lifecycle/deletion-candidate", json=candidate, headers=admin_headers).status_code == 200 assert client.post("/api/system/tenant/lifecycle/rollback", json={"expected_fence": "4", "idempotency_key": "rollback-1"}, headers=admin_headers).status_code == 200 recover_second_approval = approval("begin_recovery", 5, "recover-2") assert client.post("/api/system/tenant/lifecycle/recover", json={"expected_fence": "5", "idempotency_key": "recover-2", "approval_ref": recover_second_approval}, headers=admin_headers).status_code == 200 assert client.post("/api/system/tenant/lifecycle/activate", json={"expected_fence": "6", "idempotency_key": "activate-2", "approval_ref": "approval-5"}, headers=admin_headers).status_code == 200 candidate["expected_fence"] = "7" candidate["idempotency_key"] = "candidate-2" candidate["approval_ref"] = "approval-6" assert client.post("/api/system/tenant/lifecycle/deletion-candidate", json=candidate, headers=admin_headers).status_code == 200 delete_approval = approval("delete", 8, "delete-1", "b" * 64, "120") deleted = {"expected_fence": "8", "idempotency_key": "delete-1", "approval_ref": delete_approval, "backup_digest": "b" * 64, "retention_seconds": "120"} assert client.post("/api/system/tenant/lifecycle/delete", json=deleted, headers=admin_headers).status_code == 200 assert client.post("/api/system/tenant/lifecycle/rollback", json={"expected_fence": "9", "idempotency_key": "rollback-after-delete"}, headers=admin_headers).status_code == 400 finally: with admin.begin() as cleanup: for scoped_tenant_id in tenant_ids: cleanup.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_control_outbox_scope_lookup WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_control_claims WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_control_memberships WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_provision_operations WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_quota_reservations WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_quotas WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_lifecycle_events WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_approval_facts WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenant_audit_events WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) cleanup.execute(text("DELETE FROM public.tenants WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}) for table in ("tenant_control_outbox_scope_lookup", "tenant_control_outbox", "tenant_control_claims", "tenant_control_memberships", "tenant_provision_operations", "tenant_quota_reservations", "tenant_quotas", "tenant_lifecycle_events", "tenant_approval_facts", "tenant_audit_events", "tenants"): for scoped_tenant_id in tenant_ids: assert cleanup.execute(text(f"SELECT count(*) FROM public.{table} WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}).scalar_one() == 0 cleanup.execute(text(f"DROP ROLE IF EXISTS {cleanup.dialect.identifier_preparer.quote(control_name)}")) cleanup.execute(text(f"DROP ROLE IF EXISTS {cleanup.dialect.identifier_preparer.quote(approver_name)}")) admin.dispose()