20260817_514_tenant_outbox_scope_lookup.py 5.5 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980
  1. """Resolve a durable worker scope before entering FORCE RLS outbox rows."""
  2. from alembic import op
  3. revision = "20260817_514"
  4. down_revision = "20260817_513"
  5. branch_labels = None
  6. depends_on = None
  7. def upgrade() -> None:
  8. op.execute(r"""
  9. CREATE TABLE public.tenant_control_outbox_scope_lookup (
  10. outbox_uid uuid PRIMARY KEY,
  11. tenant_id text NOT NULL REFERENCES public.tenants(tenant_id) ON DELETE RESTRICT,
  12. created_at timestamptz NOT NULL DEFAULT clock_timestamp()
  13. );
  14. ALTER TABLE public.tenant_control_outbox_scope_lookup OWNER TO dataops_tenant_foundation_owner;
  15. REVOKE ALL ON public.tenant_control_outbox_scope_lookup FROM PUBLIC,dataops_app_runtime,dataops_tenant_control;
  16. ALTER TABLE public.tenant_control_outbox NO FORCE ROW LEVEL SECURITY;
  17. INSERT INTO public.tenant_control_outbox_scope_lookup(outbox_uid,tenant_id)
  18. SELECT outbox_uid,tenant_id FROM public.tenant_control_outbox
  19. ON CONFLICT(outbox_uid) DO NOTHING;
  20. ALTER TABLE public.tenant_control_outbox FORCE ROW LEVEL SECURITY;
  21. CREATE OR REPLACE FUNCTION public.tenant_control_claim_outbox_enqueue()
  22. RETURNS trigger LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $tenant_outbox_enqueue$
  23. DECLARE v_outbox uuid;
  24. BEGIN
  25. v_outbox:=gen_random_uuid();
  26. INSERT INTO public.tenant_control_outbox_scope_lookup(outbox_uid,tenant_id) VALUES(v_outbox,NEW.tenant_id);
  27. PERFORM set_config('dataops.tenant_id',NEW.tenant_id,true);
  28. INSERT INTO public.tenant_control_outbox(outbox_uid,tenant_id,claim_uid,event_type,payload_digest)
  29. VALUES(v_outbox,NEW.tenant_id,NEW.claim_uid,
  30. CASE WHEN NEW.action='lifecycle_transition' THEN 'tenant.lifecycle.claim' ELSE 'tenant.quota.claim' END,
  31. NEW.request_digest);
  32. INSERT INTO public.tenant_audit_events(tenant_id,event_type,payload_digest,lease_fence)
  33. VALUES(NEW.tenant_id,'tenant.outbox.queued',NEW.request_digest,0);
  34. RETURN NEW;
  35. END; $tenant_outbox_enqueue$;
  36. ALTER FUNCTION public.tenant_control_claim_outbox_enqueue() OWNER TO dataops_tenant_foundation_owner;
  37. CREATE OR REPLACE FUNCTION public.tenant_control_worker_claim(p_payload jsonb)
  38. RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $tenant_worker_claim$
  39. DECLARE v_outbox uuid; v_worker text; v_seconds integer; v_row public.tenant_control_outbox%ROWTYPE; v_tenant record; v_scope text;
  40. BEGIN
  41. IF NOT pg_has_role(session_user,'dataops_tenant_control','MEMBER') THEN RAISE EXCEPTION 'tenant_control_identity_required'; END IF;
  42. IF jsonb_typeof(p_payload)<>'object' OR NOT (p_payload ?& ARRAY['outbox_uid','worker_id','lease_seconds']) OR p_payload-ARRAY['outbox_uid','worker_id','lease_seconds']<>'{}'::jsonb THEN RAISE EXCEPTION 'tenant_payload_closed'; END IF;
  43. IF coalesce(p_payload->>'outbox_uid','') !~* '^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$'
  44. OR coalesce(p_payload->>'worker_id','') !~ '^[A-Za-z0-9._:-]{1,120}$'
  45. OR coalesce(p_payload->>'lease_seconds','') !~ '^[1-9][0-9]{0,3}$' THEN RAISE EXCEPTION 'tenant_payload_invalid'; END IF;
  46. v_outbox:=(p_payload->>'outbox_uid')::uuid; v_worker:=p_payload->>'worker_id'; v_seconds:=(p_payload->>'lease_seconds')::integer;
  47. SELECT tenant_id INTO v_scope FROM public.tenant_control_outbox_scope_lookup WHERE outbox_uid=v_outbox;
  48. IF NOT FOUND THEN RAISE EXCEPTION 'tenant_outbox_missing'; END IF;
  49. PERFORM set_config('dataops.tenant_id',v_scope,true);
  50. SELECT * INTO v_row FROM public.tenant_control_outbox WHERE outbox_uid=v_outbox FOR UPDATE;
  51. IF NOT FOUND OR v_row.tenant_id<>v_scope THEN RAISE EXCEPTION 'tenant_outbox_missing'; END IF;
  52. SELECT state,lease_fence INTO v_tenant FROM public.tenants WHERE tenant_id=v_row.tenant_id FOR UPDATE;
  53. IF NOT FOUND OR v_tenant.state<>'active' THEN RAISE EXCEPTION 'tenant_worker_scope_denied'; END IF;
  54. IF v_row.state='completed' THEN RAISE EXCEPTION 'tenant_outbox_replayed'; END IF;
  55. IF v_row.state='leased' AND v_row.lease_expires_at>clock_timestamp() THEN RAISE EXCEPTION 'tenant_outbox_leased'; END IF;
  56. UPDATE public.tenant_control_outbox SET state='leased',worker_id=v_worker,leased_at=clock_timestamp(),lease_expires_at=clock_timestamp()+make_interval(secs=>v_seconds),lease_fence=lease_fence+1
  57. WHERE outbox_uid=v_outbox RETURNING * INTO v_row;
  58. INSERT INTO public.tenant_audit_events(tenant_id,event_type,payload_digest,lease_fence)
  59. VALUES(v_row.tenant_id,'tenant.outbox.leased',v_row.payload_digest,v_row.lease_fence);
  60. RETURN jsonb_build_object('outbox_uid',v_row.outbox_uid,'tenant_id',v_row.tenant_id,'event_type',v_row.event_type,'payload_digest',v_row.payload_digest,'lease_fence',v_row.lease_fence,'lease_expires_at',v_row.lease_expires_at);
  61. END; $tenant_worker_claim$;
  62. ALTER FUNCTION public.tenant_control_worker_claim(jsonb) OWNER TO dataops_tenant_foundation_owner;
  63. REVOKE ALL ON FUNCTION public.tenant_control_worker_claim(jsonb) FROM PUBLIC,dataops_app_runtime;
  64. GRANT EXECUTE ON FUNCTION public.tenant_control_worker_claim(jsonb) TO dataops_tenant_control;
  65. """)
  66. def downgrade() -> None:
  67. op.execute("""
  68. DO $$ BEGIN
  69. IF EXISTS (SELECT 1 FROM public.tenant_control_outbox_scope_lookup LIMIT 1) THEN RAISE EXCEPTION 'tenant outbox scope downgrade refused while work exists'; END IF;
  70. END $$;
  71. DROP TABLE public.tenant_control_outbox_scope_lookup;
  72. """)