test_wp10_tenant_http_postgres.py 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278
  1. from __future__ import annotations
  2. import os
  3. import uuid
  4. from concurrent.futures import ThreadPoolExecutor
  5. import pytest
  6. from sqlalchemy import create_engine, text
  7. from sqlalchemy.engine import make_url
  8. from sqlalchemy.exc import SQLAlchemyError
  9. from app.core.system.tenant_control_repository import SqlAlchemyTenantControlRepository
  10. from app.core.system.tenant_control_service import (
  11. TenantControlError,
  12. TenantControlService,
  13. )
  14. pytestmark = pytest.mark.integration
  15. def test_wp10_real_flask_claim_route_uses_dedicated_control_login_and_commits(monkeypatch):
  16. database_url = os.getenv("TEST_DATABASE_URL")
  17. if not database_url:
  18. pytest.skip("TEST_DATABASE_URL is required")
  19. tenant_id = f"wp10-http-{uuid.uuid4().hex[:14]}"
  20. concurrent_tenant_id = f"wp10-http-concurrent-{uuid.uuid4().hex[:10]}"
  21. tenant_ids = (tenant_id, concurrent_tenant_id)
  22. principal_id = str(uuid.uuid4())
  23. concurrent_principal_id = str(uuid.uuid4())
  24. control_name = f"wp10_control_{uuid.uuid4().hex[:12]}"
  25. control_password = f"wp10_{uuid.uuid4().hex}"
  26. approver_name = f"wp10_approver_{uuid.uuid4().hex[:12]}"
  27. approver_password = f"wp10_{uuid.uuid4().hex}"
  28. admin = create_engine(database_url, pool_pre_ping=True)
  29. try:
  30. with admin.begin() as connection:
  31. quoted_control = connection.dialect.identifier_preparer.quote(control_name)
  32. connection.execute(text(
  33. f"CREATE ROLE {quoted_control} LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION PASSWORD :password"
  34. ), {"password": control_password})
  35. connection.execute(text(f"GRANT dataops_tenant_control TO {quoted_control}"))
  36. quoted_approver = connection.dialect.identifier_preparer.quote(approver_name)
  37. connection.execute(text(f"CREATE ROLE {quoted_approver} LOGIN NOSUPERUSER NOCREATEDB NOCREATEROLE NOREPLICATION PASSWORD :password"), {"password": approver_password})
  38. connection.execute(text(f"GRANT dataops_tenant_approver TO {quoted_approver}"))
  39. control_url = make_url(database_url).set(username=control_name, password=control_password).render_as_string(hide_password=False)
  40. approver_url = make_url(database_url).set(username=approver_name, password=approver_password).render_as_string(hide_password=False)
  41. monkeypatch.setenv("TENANT_CONTROL_DATABASE_URL", control_url)
  42. monkeypatch.setenv("TENANT_APPROVER_DATABASE_URL", approver_url)
  43. monkeypatch.setenv("TENANT_APPROVER_ISSUER_ID", str(uuid.uuid4()))
  44. monkeypatch.setenv("PRIVATE_SINGLE_TENANT_ID", tenant_id)
  45. monkeypatch.setenv("PRIVATE_SINGLE_TENANT_HOST", "private.example")
  46. monkeypatch.setenv("TRUSTED_TENANT_ROUTE_HOST", "private.example")
  47. monkeypatch.setattr(
  48. "app.core.system.auth.load_identity_from_token",
  49. lambda token, secret: {"id": principal_id, "roles": [token.removeprefix("wp10-")]} if token in {"wp10-editor", "wp10-admin"} else None,
  50. )
  51. from app import create_app
  52. app = create_app()
  53. app.config.update(TESTING=True)
  54. client = app.test_client()
  55. headers = {"Authorization": "Bearer wp10-editor", "Host": "private.example"}
  56. # Diagnostic RED boundary: this is the same dedicated control login
  57. # used by Flask, so a failure below exposes PostgreSQL's class/code.
  58. with pytest.raises(TenantControlError, match="tenant_control_provision_denied"):
  59. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private(
  60. principal_id="not-a-uuid", host="ignored.example", request={"idempotency_key": "diagnostic-invalid"}
  61. )
  62. provision_key = f"provision-{tenant_id}"
  63. provision = client.post("/api/system/tenant/provisions", json={"idempotency_key": provision_key}, headers={"Authorization": "Bearer wp10-admin", "Host": "private.example"})
  64. assert provision.status_code == 201
  65. assert provision.headers["Cache-Control"] == "no-store"
  66. assert provision.get_json()["data"]["tenant_id"] == tenant_id
  67. # A fresh Flask application process must observe the durable exact
  68. # replay, not create a second membership or trust the HTTP Host.
  69. replay_app = create_app()
  70. replay_app.config.update(TESTING=True)
  71. replay = replay_app.test_client().post(
  72. "/api/system/tenant/provisions",
  73. json={"idempotency_key": provision_key},
  74. headers={"Authorization": "Bearer wp10-admin", "Host": "attacker.example"},
  75. )
  76. assert replay.status_code == 201
  77. assert replay.get_json()["data"]["replayed"] is True
  78. # The same durable key cannot be rebound to a different trusted
  79. # principal, host, or configured private tenant.
  80. with pytest.raises(TenantControlError, match="tenant_control_provision_denied"):
  81. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private(
  82. principal_id=str(uuid.uuid4()), host="ignored.example", request={"idempotency_key": provision_key}
  83. )
  84. monkeypatch.setenv("PRIVATE_SINGLE_TENANT_HOST", "other.example")
  85. with pytest.raises(TenantControlError, match="tenant_control_provision_denied"):
  86. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private(
  87. principal_id=principal_id, host="ignored.example", request={"idempotency_key": provision_key}
  88. )
  89. monkeypatch.setenv("PRIVATE_SINGLE_TENANT_HOST", "private.example")
  90. monkeypatch.setenv("PRIVATE_SINGLE_TENANT_ID", f"other-{tenant_id}")
  91. with pytest.raises(TenantControlError, match="tenant_control_provision_denied"):
  92. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).provision_private(
  93. principal_id=principal_id, host="ignored.example", request={"idempotency_key": provision_key}
  94. )
  95. monkeypatch.setenv("PRIVATE_SINGLE_TENANT_ID", tenant_id)
  96. # Two independent control connections race the same first provision.
  97. # The durable idempotency record/advisory lock leaves one creation and
  98. # one exact replay; it must never add a second membership.
  99. concurrent_key = f"concurrent-{uuid.uuid4().hex}"
  100. def concurrent_provision():
  101. return SqlAlchemyTenantControlRepository(control_url).provision(
  102. tenant_id=concurrent_tenant_id,
  103. principal_id=concurrent_principal_id,
  104. host="private.example",
  105. idempotency_key=concurrent_key,
  106. )
  107. with ThreadPoolExecutor(max_workers=2) as pool:
  108. concurrent_results = list(pool.map(lambda _: concurrent_provision(), range(2)))
  109. assert {item["tenant_id"] for item in concurrent_results} == {concurrent_tenant_id}
  110. assert sum(bool(item["replayed"]) for item in concurrent_results) == 1
  111. with admin.connect() as check:
  112. assert check.execute(
  113. text("SELECT count(*) FROM public.tenant_control_memberships WHERE tenant_id=:tenant_id"),
  114. {"tenant_id": concurrent_tenant_id},
  115. ).scalar_one() == 1
  116. rejected = client.post("/api/system/tenant/quota-claims", json={"tenant_id": "other", "quota_name": "records", "amount": "1", "idempotency_key": "http-1"}, headers=headers)
  117. assert rejected.status_code == 400
  118. created = client.post("/api/system/tenant/quota-claims", json={"quota_name": "records", "amount": "1", "idempotency_key": "http-1"}, headers=headers)
  119. assert created.status_code == 201
  120. assert created.headers["Cache-Control"] == "no-store"
  121. assert created.get_json()["data"]["tenant_id"] == tenant_id
  122. propagation = created.get_json()["data"]["propagation"]
  123. assert propagation["task"]["tenant_id"] == tenant_id
  124. assert propagation["event"]["event_type"] == "tenant.quota.claim"
  125. assert all(value.startswith(f"{tenant_id}/") for value in propagation["namespaces"].values())
  126. with admin.connect() as check:
  127. 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()
  128. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  129. 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())
  130. assert count == 1
  131. monkeypatch.setenv("TENANT_CONTROL_WORKER_ID", "tenant-worker-1")
  132. with create_engine(control_url).connect() as legacy_gateway, pytest.raises(SQLAlchemyError):
  133. 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())})
  134. worker = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=outbox_uid)
  135. assert worker["tenant_id"] == tenant_id
  136. assert worker["propagation"]["event"]["tenant_id"] == tenant_id
  137. assert all(value.startswith(f"{tenant_id}/") for value in worker["propagation"]["namespaces"].values())
  138. settled = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(
  139. outbox_uid=outbox_uid, lease_fence=worker["lease_fence"], payload_digest=worker["payload_digest"], outcome="complete"
  140. )
  141. assert settled["state"] == "completed"
  142. assert TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(
  143. outbox_uid=outbox_uid, lease_fence=worker["lease_fence"], payload_digest=worker["payload_digest"], outcome="complete"
  144. )["replayed"] is True
  145. with admin.connect() as check:
  146. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  147. 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"
  148. retry_claim = client.post("/api/system/tenant/quota-claims", json={"quota_name": "background_tasks", "amount": "1", "idempotency_key": "http-2"}, headers=headers)
  149. assert retry_claim.status_code == 201
  150. with admin.connect() as check:
  151. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  152. 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())
  153. retry_fences = []
  154. for attempt in range(3):
  155. retry_worker = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=retry_outbox_uid)
  156. retry_fences.append(retry_worker["lease_fence"])
  157. state = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(
  158. outbox_uid=retry_outbox_uid,
  159. lease_fence=retry_worker["lease_fence"],
  160. payload_digest=retry_worker["payload_digest"],
  161. outcome="fail",
  162. )["state"]
  163. assert state == ("failed" if attempt == 2 else "queued")
  164. assert retry_fences == sorted(retry_fences)
  165. with admin.connect() as check:
  166. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  167. 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"
  168. with pytest.raises(TenantControlError, match="tenant_worker_scope_denied"):
  169. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=retry_outbox_uid)
  170. # A DB-clock-expired lease cannot complete; a fresh repository/engine
  171. # acts as a worker restart and can reclaim only with a higher fence.
  172. expiry_claim = client.post("/api/system/tenant/quota-claims", json={"quota_name": "events", "amount": "1", "idempotency_key": "http-expiry"}, headers=headers)
  173. assert expiry_claim.status_code == 201
  174. with admin.connect() as check:
  175. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  176. 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())
  177. expired_worker = SqlAlchemyTenantControlRepository(control_url).claim_outbox(outbox_uid=expiry_outbox_uid, worker_id="tenant-worker-1", lease_seconds="1")
  178. with admin.begin() as wait_for_db_clock:
  179. wait_for_db_clock.execute(text("SELECT pg_sleep(1.1)"))
  180. with pytest.raises(TenantControlError, match="tenant_worker_settle_denied"):
  181. 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")
  182. restarted_worker = TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=expiry_outbox_uid)
  183. assert restarted_worker["lease_fence"] > expired_worker["lease_fence"]
  184. with pytest.raises(TenantControlError, match="tenant_worker_settle_denied"):
  185. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).settle_outbox(outbox_uid=expiry_outbox_uid, lease_fence=restarted_worker["lease_fence"], payload_digest="0" * 64, outcome="complete")
  186. 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"
  187. worker_race_claim = client.post("/api/system/tenant/quota-claims", json={"quota_name": "models", "amount": "1", "idempotency_key": "http-worker-race"}, headers=headers)
  188. assert worker_race_claim.status_code == 201
  189. with admin.connect() as check:
  190. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  191. 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())
  192. def race_worker_claim():
  193. try:
  194. return "claimed", SqlAlchemyTenantControlRepository(control_url).claim_outbox(
  195. outbox_uid=worker_race_outbox_uid, worker_id="tenant-worker-race"
  196. )
  197. except TenantControlError:
  198. return "denied", None
  199. with ThreadPoolExecutor(max_workers=2) as pool:
  200. worker_race_results = list(pool.map(lambda _: race_worker_claim(), range(2)))
  201. assert [outcome for outcome, _ in worker_race_results].count("claimed") == 1
  202. claimed_worker = next(value for outcome, value in worker_race_results if outcome == "claimed")
  203. monkeypatch.setenv("TENANT_CONTROL_WORKER_ID", "tenant-worker-race")
  204. 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"
  205. monkeypatch.setenv("TENANT_CONTROL_WORKER_ID", "tenant-worker-1")
  206. admin_headers = {"Authorization": "Bearer wp10-admin", "Host": "private.example"}
  207. frozen = client.post("/api/system/tenant/lifecycle/freeze", json={"expected_fence": "0", "idempotency_key": "freeze-1"}, headers=admin_headers)
  208. assert frozen.status_code == 200
  209. assert frozen.headers["Cache-Control"] == "no-store"
  210. with admin.connect() as check:
  211. check.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": tenant_id})
  212. 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())
  213. with pytest.raises(TenantControlError, match="tenant_worker_scope_denied"):
  214. TenantControlService(SqlAlchemyTenantControlRepository(control_url)).claim_outbox(outbox_uid=frozen_outbox_uid)
  215. status = client.get("/api/system/tenant/status", headers=admin_headers)
  216. assert status.status_code == 200
  217. assert status.get_json()["data"]["state"] == "frozen"
  218. assert client.get("/api/system/tenant/audit", headers=admin_headers).status_code == 200
  219. assert client.get("/api/system/tenant/manifests", headers=admin_headers).get_json()["data"]["manifests"] == []
  220. def approval(operation, fence, key, backup=None, retention=None):
  221. body = {"operation": operation, "expected_fence": str(fence), "idempotency_key": key}
  222. if backup is not None:
  223. body["backup_digest"] = backup
  224. if retention is not None:
  225. body["retention_seconds"] = retention
  226. response = client.post("/api/system/tenant/approvals", json=body, headers=admin_headers)
  227. assert response.status_code == 201
  228. assert response.headers["Cache-Control"] == "no-store"
  229. return response.get_json()["data"]["approval_ref"]
  230. recover_approval = approval("begin_recovery", 1, "recover-1")
  231. 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
  232. 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
  233. candidate = {"expected_fence": "3", "idempotency_key": "candidate-1", "approval_ref": "approval-3", "backup_digest": "a" * 64, "retention_seconds": "60"}
  234. assert client.post("/api/system/tenant/lifecycle/deletion-candidate", json=candidate, headers=admin_headers).status_code == 200
  235. assert client.post("/api/system/tenant/lifecycle/rollback", json={"expected_fence": "4", "idempotency_key": "rollback-1"}, headers=admin_headers).status_code == 200
  236. recover_second_approval = approval("begin_recovery", 5, "recover-2")
  237. 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
  238. 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
  239. candidate["expected_fence"] = "7"
  240. candidate["idempotency_key"] = "candidate-2"
  241. candidate["approval_ref"] = "approval-6"
  242. assert client.post("/api/system/tenant/lifecycle/deletion-candidate", json=candidate, headers=admin_headers).status_code == 200
  243. delete_approval = approval("delete", 8, "delete-1", "b" * 64, "120")
  244. deleted = {"expected_fence": "8", "idempotency_key": "delete-1", "approval_ref": delete_approval, "backup_digest": "b" * 64, "retention_seconds": "120"}
  245. assert client.post("/api/system/tenant/lifecycle/delete", json=deleted, headers=admin_headers).status_code == 200
  246. assert client.post("/api/system/tenant/lifecycle/rollback", json={"expected_fence": "9", "idempotency_key": "rollback-after-delete"}, headers=admin_headers).status_code == 400
  247. finally:
  248. with admin.begin() as cleanup:
  249. for scoped_tenant_id in tenant_ids:
  250. cleanup.execute(text("SELECT set_config('dataops.tenant_id', :tenant_id, true)"), {"tenant_id": scoped_tenant_id})
  251. cleanup.execute(text("DELETE FROM public.tenant_control_outbox_scope_lookup WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  252. cleanup.execute(text("DELETE FROM public.tenant_control_outbox WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  253. cleanup.execute(text("DELETE FROM public.tenant_control_claims WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  254. cleanup.execute(text("DELETE FROM public.tenant_control_memberships WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  255. cleanup.execute(text("DELETE FROM public.tenant_provision_operations WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  256. cleanup.execute(text("DELETE FROM public.tenant_quota_reservations WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  257. cleanup.execute(text("DELETE FROM public.tenant_quotas WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  258. cleanup.execute(text("DELETE FROM public.tenant_lifecycle_events WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  259. cleanup.execute(text("DELETE FROM public.tenant_approval_facts WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  260. cleanup.execute(text("DELETE FROM public.tenant_audit_events WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  261. cleanup.execute(text("DELETE FROM public.tenants WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id})
  262. 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"):
  263. for scoped_tenant_id in tenant_ids:
  264. assert cleanup.execute(text(f"SELECT count(*) FROM public.{table} WHERE tenant_id=:tenant_id"), {"tenant_id": scoped_tenant_id}).scalar_one() == 0
  265. cleanup.execute(text(f"DROP ROLE IF EXISTS {cleanup.dialect.identifier_preparer.quote(control_name)}"))
  266. cleanup.execute(text(f"DROP ROLE IF EXISTS {cleanup.dialect.identifier_preparer.quote(approver_name)}"))
  267. admin.dispose()