"""Seal WP13 legacy runtime execution and bind dead-letter recovery.""" from alembic import op revision = "20260818_555" down_revision = "20260818_554" branch_labels = None depends_on = None def upgrade() -> None: op.get_bind().exec_driver_sql(r''' DO $roles$ BEGIN IF NOT EXISTS(SELECT 1 FROM pg_roles WHERE rolname='dataops_plugin_platform_owner' AND NOT rolcanlogin) OR NOT EXISTS(SELECT 1 FROM pg_roles WHERE rolname='dataops_plugin_platform_control' AND NOT rolcanlogin) OR NOT EXISTS(SELECT 1 FROM pg_roles WHERE rolname='dataops_app_runtime') THEN RAISE EXCEPTION 'WP13 role-init prerequisite missing'; END IF; END $roles$; CREATE TABLE public.plugin_recovery_claims ( recovery_claim_uid uuid PRIMARY KEY DEFAULT gen_random_uuid(), run_uid uuid NOT NULL REFERENCES public.plugin_runs(run_uid) ON DELETE RESTRICT, plugin_uid text NOT NULL, version text NOT NULL, tenant_ref text NOT NULL, domain_ref text NOT NULL, actor_ref text NOT NULL, incident_uid uuid NOT NULL, approval_uid uuid NOT NULL UNIQUE REFERENCES public.plugin_approvals(approval_uid) ON DELETE RESTRICT, expected_fence bigint NOT NULL CHECK(expected_fence>=0), expires_at timestamptz NOT NULL, consumed_at timestamptz, created_at timestamptz NOT NULL DEFAULT clock_timestamp(), FOREIGN KEY(plugin_uid,version) REFERENCES public.plugin_registry_versions(plugin_uid,version) ON DELETE RESTRICT ); CREATE UNIQUE INDEX plugin_recovery_one_live_claim ON public.plugin_recovery_claims(run_uid) WHERE consumed_at IS NULL; ALTER TABLE public.plugin_recovery_claims OWNER TO dataops_plugin_platform_owner; REVOKE ALL ON TABLE public.plugin_recovery_claims FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_plugin_platform_control; -- 553's direct writer predates claims, leases, capability checking and audit. -- Remove it rather than leaving an ACL-dependent bypass surface. DO $legacy$ BEGIN IF to_regprocedure('public.plugin_platform_runtime_execute(jsonb)') IS NOT NULL THEN REVOKE ALL ON FUNCTION public.plugin_platform_runtime_execute(jsonb) FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_plugin_platform_control; DROP FUNCTION public.plugin_platform_runtime_execute(jsonb); END IF; END $legacy$; CREATE FUNCTION public.plugin_platform_issue_recovery_claim(p jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $control$ DECLARE r record; approval record; claim uuid; incident_exists boolean; incident uuid; BEGIN IF NOT pg_has_role(session_user,'dataops_plugin_platform_control','MEMBER') OR jsonb_typeof(p)<>'object' OR p-ARRAY['plugin_uid','version','actor_ref','tenant_ref','domain_ref','run_uid','approval_uid','incident_uid','expected_fence']<>'{}'::jsonb OR p->>'plugin_uid' !~ '^[a-z][a-z0-9-]{2,62}$' OR p->>'actor_ref' !~ '^[A-Za-z0-9_.:-]{1,120}$' OR p->>'tenant_ref' !~ '^[a-z][a-z0-9-]{0,62}$' OR p->>'domain_ref' !~ '^[a-z][a-z0-9-]{0,62}$' THEN RAISE EXCEPTION 'plugin_recovery_claim_denied'; END IF; incident:=NULLIF(p->>'incident_uid','')::uuid; IF incident IS NULL THEN RAISE EXCEPTION 'plugin_recovery_incident_required'; END IF; IF to_regclass('public.data_incidents') IS NOT NULL THEN EXECUTE 'SELECT EXISTS(SELECT 1 FROM public.data_incidents WHERE uid=$1)' INTO incident_exists USING incident; IF NOT incident_exists THEN RAISE EXCEPTION 'plugin_incident_reference_denied'; END IF; END IF; SELECT q.*,v.manifest_digest INTO r FROM public.plugin_runs q JOIN public.plugin_registry_versions v ON (v.plugin_uid=q.plugin_uid AND v.version=q.version) WHERE q.run_uid=(p->>'run_uid')::uuid FOR UPDATE; IF NOT FOUND OR r.plugin_uid<>p->>'plugin_uid' OR r.version<>p->>'version' OR r.tenant_ref<>p->>'tenant_ref' OR r.domain_ref<>p->>'domain_ref' OR r.state<>'dead_letter' OR r.lease_fence<>(p->>'expected_fence')::bigint THEN RAISE EXCEPTION 'plugin_recovery_scope_or_fence_denied'; END IF; SELECT * INTO approval FROM public.plugin_approvals WHERE approval_uid=(p->>'approval_uid')::uuid FOR UPDATE; IF NOT FOUND OR approval.consumed_at IS NOT NULL OR approval.expires_at<=clock_timestamp() OR approval.plugin_uid<>r.plugin_uid OR approval.version<>r.version OR approval.action_name<>'recover' OR approval.actor_ref<>p->>'actor_ref' OR approval.tenant_ref<>p->>'tenant_ref' OR approval.domain_ref<>p->>'domain_ref' OR approval.manifest_digest<>r.manifest_digest OR approval.reviewer_ref=approval.actor_ref THEN RAISE EXCEPTION 'plugin_recovery_approval_denied'; END IF; UPDATE public.plugin_recovery_claims SET consumed_at=clock_timestamp() WHERE run_uid=r.run_uid AND consumed_at IS NULL AND expires_at<=clock_timestamp(); IF EXISTS(SELECT 1 FROM public.plugin_recovery_claims WHERE run_uid=r.run_uid AND consumed_at IS NULL) THEN RAISE EXCEPTION 'plugin_recovery_claim_live'; END IF; UPDATE public.plugin_approvals SET consumed_at=clock_timestamp() WHERE approval_uid=approval.approval_uid; INSERT INTO public.plugin_recovery_claims(run_uid,plugin_uid,version,tenant_ref,domain_ref,actor_ref,incident_uid,approval_uid,expected_fence,expires_at) VALUES(r.run_uid,r.plugin_uid,r.version,r.tenant_ref,r.domain_ref,p->>'actor_ref',incident,approval.approval_uid,r.lease_fence,clock_timestamp()+interval '30 seconds') RETURNING recovery_claim_uid INTO claim; RETURN jsonb_build_object('claim_uid',claim::text,'run_uid',r.run_uid::text,'fence',r.lease_fence); END $control$; CREATE OR REPLACE FUNCTION public.plugin_platform_runtime_v2(p jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $runtime$ DECLARE a text; c record; run record; recovery record; BEGIN IF NOT pg_has_role(session_user,'dataops_app_runtime','MEMBER') OR jsonb_typeof(p)<>'object' OR NOT(p ? 'action') THEN RAISE EXCEPTION 'plugin_runtime_denied'; END IF; a:=p->>'action'; IF a='enqueue' THEN IF p-ARRAY['action','claim_uid','idempotency_key','input_digest','operation_name']<>'{}'::jsonb OR p->>'input_digest' !~ '^[0-9a-f]{64}$' THEN RAISE EXCEPTION 'plugin_enqueue_closed'; END IF; UPDATE public.plugin_request_claims SET consumed_at=clock_timestamp() WHERE claim_uid=(p->>'claim_uid')::uuid AND consumed_at IS NULL AND expires_at>clock_timestamp() RETURNING * INTO c; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_claim_consumed'; END IF; INSERT INTO public.plugin_runs(plugin_uid,version,tenant_ref,domain_ref,principal_ref,operation_name,idempotency_key,input_digest,request_digest,lease_fence,state) VALUES(c.plugin_uid,c.version,c.tenant_ref,c.domain_ref,c.principal_ref,p->>'operation_name',p->>'idempotency_key',p->>'input_digest',c.request_digest,0,'pending') ON CONFLICT(plugin_uid,version,tenant_ref,idempotency_key) DO NOTHING; SELECT * INTO run FROM public.plugin_runs WHERE plugin_uid=c.plugin_uid AND version=c.version AND tenant_ref=c.tenant_ref AND idempotency_key=p->>'idempotency_key'; IF run.request_digest<>c.request_digest OR run.input_digest<>p->>'input_digest' THEN RAISE EXCEPTION 'plugin_replay_conflict'; END IF; RETURN jsonb_build_object('run_uid',run.run_uid::text,'state',run.state,'replay',run.created_at'{}'::jsonb THEN RAISE EXCEPTION 'plugin_claim_closed'; END IF; IF EXISTS(SELECT 1 FROM public.plugin_runtime_breakers b JOIN public.plugin_runs q ON q.plugin_uid=b.plugin_uid AND q.version=b.version AND q.tenant_ref=b.tenant_ref WHERE q.run_uid=(p->>'run_uid')::uuid AND b.opened_until>clock_timestamp()) THEN RAISE EXCEPTION 'plugin_breaker_open'; END IF; UPDATE public.plugin_runs SET state='running',lease_owner=p->>'worker',lease_fence=lease_fence+1,lease_expires_at=clock_timestamp()+interval '10 seconds' WHERE run_uid=(p->>'run_uid')::uuid AND (state='pending' OR (state='running' AND lease_expires_at<=clock_timestamp())) AND next_attempt_at<=clock_timestamp() RETURNING * INTO run; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_lease_unavailable'; END IF; RETURN jsonb_build_object('run_uid',run.run_uid::text,'fence',run.lease_fence,'lease_owner',run.lease_owner); ELSIF a='settle' THEN IF p-ARRAY['action','run_uid','worker','fence','success','output_digest','failure_digest']<>'{}'::jsonb OR p->>'output_digest' !~ '^[0-9a-f]{64}$' OR p->>'failure_digest' !~ '^[0-9a-f]{64}$' THEN RAISE EXCEPTION 'plugin_settle_closed'; END IF; UPDATE public.plugin_runs SET output_digest=CASE WHEN (p->>'success')::boolean THEN p->>'output_digest' ELSE output_digest END,attempt_count=CASE WHEN (p->>'success')::boolean THEN attempt_count ELSE attempt_count+1 END,state=CASE WHEN (p->>'success')::boolean THEN 'succeeded' WHEN attempt_count+1>=3 THEN 'dead_letter' ELSE 'pending' END,lease_owner=NULL,lease_expires_at=NULL,next_attempt_at=CASE WHEN (p->>'success')::boolean THEN next_attempt_at ELSE clock_timestamp()+(attempt_count+1)*interval '1 second' END,completed_at=CASE WHEN (p->>'success')::boolean THEN clock_timestamp() ELSE completed_at END WHERE run_uid=(p->>'run_uid')::uuid AND state='running' AND lease_owner=p->>'worker' AND lease_fence=(p->>'fence')::bigint AND lease_expires_at>clock_timestamp() RETURNING * INTO run; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_late_worker_fenced'; END IF; IF run.state='dead_letter' THEN INSERT INTO public.plugin_dead_letters(run_uid,failure_digest) VALUES(run.run_uid,p->>'failure_digest') ON CONFLICT(run_uid) DO NOTHING; END IF; IF (p->>'success')::boolean THEN INSERT INTO public.plugin_runtime_breakers(plugin_uid,version,tenant_ref,failure_count,opened_until) VALUES(run.plugin_uid,run.version,run.tenant_ref,0,NULL) ON CONFLICT(plugin_uid,version,tenant_ref) DO UPDATE SET failure_count=0,opened_until=NULL,updated_at=clock_timestamp(); ELSE INSERT INTO public.plugin_runtime_breakers(plugin_uid,version,tenant_ref,failure_count,opened_until) VALUES(run.plugin_uid,run.version,run.tenant_ref,1,NULL) ON CONFLICT(plugin_uid,version,tenant_ref) DO UPDATE SET failure_count=plugin_runtime_breakers.failure_count+1,opened_until=CASE WHEN plugin_runtime_breakers.failure_count+1>=3 THEN clock_timestamp()+interval '60 seconds' ELSE NULL END,updated_at=clock_timestamp(); END IF; INSERT INTO public.plugin_audit_outbox(plugin_uid,version,tenant_ref,event_type,actor_ref,payload_digest) VALUES(run.plugin_uid,run.version,run.tenant_ref,CASE WHEN run.state='dead_letter' THEN 'dead_letter' ELSE 'invoked' END,p->>'worker',CASE WHEN (p->>'success')::boolean THEN p->>'output_digest' ELSE p->>'failure_digest' END); RETURN jsonb_build_object('state',run.state,'attempt_count',run.attempt_count,'fence',run.lease_fence); ELSIF a='recover' THEN IF p-ARRAY['action','claim_uid','worker']<>'{}'::jsonb THEN RAISE EXCEPTION 'plugin_recovery_closed'; END IF; SELECT * INTO recovery FROM public.plugin_recovery_claims WHERE recovery_claim_uid=(p->>'claim_uid')::uuid FOR UPDATE; IF NOT FOUND OR recovery.consumed_at IS NOT NULL OR recovery.expires_at<=clock_timestamp() THEN RAISE EXCEPTION 'plugin_recovery_claim_denied'; END IF; UPDATE public.plugin_runs SET state='pending',attempt_count=0,next_attempt_at=clock_timestamp(),lease_owner=NULL,lease_expires_at=NULL,lease_fence=lease_fence+1 WHERE run_uid=recovery.run_uid AND plugin_uid=recovery.plugin_uid AND version=recovery.version AND tenant_ref=recovery.tenant_ref AND domain_ref=recovery.domain_ref AND state='dead_letter' AND lease_fence=recovery.expected_fence RETURNING * INTO run; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_recovery_fenced'; END IF; UPDATE public.plugin_recovery_claims SET consumed_at=clock_timestamp() WHERE recovery_claim_uid=recovery.recovery_claim_uid AND consumed_at IS NULL; UPDATE public.plugin_runtime_breakers SET failure_count=0,opened_until=NULL,updated_at=clock_timestamp() WHERE plugin_uid=run.plugin_uid AND version=run.version AND tenant_ref=run.tenant_ref; DELETE FROM public.plugin_dead_letters WHERE plugin_dead_letters.run_uid=run.run_uid; INSERT INTO public.plugin_audit_outbox(plugin_uid,version,tenant_ref,event_type,actor_ref,payload_digest) VALUES(run.plugin_uid,run.version,run.tenant_ref,'recovery',recovery.actor_ref,encode(sha256(convert_to(recovery.incident_uid::text,'utf8')),'hex')); RETURN jsonb_build_object('state','pending','fence',run.lease_fence); END IF; RAISE EXCEPTION 'plugin_runtime_action_denied'; END $runtime$; ALTER FUNCTION public.plugin_platform_issue_recovery_claim(jsonb) OWNER TO dataops_plugin_platform_owner; ALTER FUNCTION public.plugin_platform_runtime_v2(jsonb) OWNER TO dataops_plugin_platform_owner; REVOKE ALL ON FUNCTION public.plugin_platform_issue_recovery_claim(jsonb),public.plugin_platform_runtime_v2(jsonb) FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_plugin_platform_control; GRANT EXECUTE ON FUNCTION public.plugin_platform_issue_recovery_claim(jsonb) TO dataops_plugin_platform_control; GRANT EXECUTE ON FUNCTION public.plugin_platform_runtime_v2(jsonb) TO dataops_app_runtime; ''') def downgrade() -> None: bind = op.get_bind() if bind.exec_driver_sql("SELECT EXISTS(SELECT 1 FROM public.plugin_recovery_claims LIMIT 1)").scalar(): raise RuntimeError("downgrade refused: WP13 recovery claims are nonempty") bind.exec_driver_sql(""" REVOKE ALL ON FUNCTION public.plugin_platform_issue_recovery_claim(jsonb) FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_plugin_platform_control; DROP FUNCTION public.plugin_platform_issue_recovery_claim(jsonb); DROP TABLE public.plugin_recovery_claims; """)