"""Make WP13 lifecycle, claims and fixed fixture execution database-enforced.""" from alembic import op revision = "20260818_554" down_revision = "20260818_553" 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 AND NOT rolsuper AND NOT rolcreaterole) 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$; ALTER TABLE public.plugin_registry_versions ADD COLUMN lifecycle_fence bigint NOT NULL DEFAULT 0 CHECK(lifecycle_fence>=0), ADD COLUMN failure_count integer NOT NULL DEFAULT 0 CHECK(failure_count>=0), ADD COLUMN incident_uid uuid, ADD COLUMN updated_at timestamptz NOT NULL DEFAULT clock_timestamp(), ADD COLUMN capabilities jsonb NOT NULL DEFAULT '[]'::jsonb, ADD COLUMN resource jsonb NOT NULL DEFAULT '{"timeout_ms":5000,"max_output_bytes":65536,"max_concurrency":1,"max_retries":3}'::jsonb; ALTER TABLE public.plugin_runs DROP CONSTRAINT plugin_runs_state_check; ALTER TABLE public.plugin_runs ADD CONSTRAINT plugin_runs_state_check CHECK(state IN ('pending','running','succeeded','failed','dead_letter')); ALTER TABLE public.plugin_runs ADD COLUMN attempt_count integer NOT NULL DEFAULT 0 CHECK(attempt_count BETWEEN 0 AND 3), ADD COLUMN next_attempt_at timestamptz NOT NULL DEFAULT clock_timestamp(), ADD COLUMN lease_owner text, ADD COLUMN lease_expires_at timestamptz; CREATE TABLE public.plugin_request_claims ( claim_uid uuid PRIMARY KEY DEFAULT gen_random_uuid(), plugin_uid text NOT NULL, version text NOT NULL, action_name text NOT NULL CHECK(action_name='invoke'), tenant_ref text NOT NULL CHECK(tenant_ref ~ '^[a-z][a-z0-9-]{0,62}$'), domain_ref text NOT NULL CHECK(domain_ref ~ '^[a-z][a-z0-9-]{0,62}$'), principal_ref text NOT NULL CHECK(principal_ref ~ '^[A-Za-z0-9_.:-]{1,120}$'), request_digest char(64) NOT NULL CHECK(request_digest ~ '^[0-9a-f]{64}$'), expires_at timestamptz NOT NULL, consumed_at timestamptz, created_at timestamptz NOT NULL DEFAULT clock_timestamp() ); CREATE TABLE public.plugin_runtime_breakers ( plugin_uid text NOT NULL, version text NOT NULL, tenant_ref text NOT NULL, failure_count integer NOT NULL DEFAULT 0 CHECK(failure_count>=0), opened_until timestamptz, updated_at timestamptz NOT NULL DEFAULT clock_timestamp(), PRIMARY KEY(plugin_uid,version,tenant_ref), FOREIGN KEY(plugin_uid,version) REFERENCES public.plugin_registry_versions(plugin_uid,version) ON DELETE RESTRICT ); ALTER TABLE public.plugin_request_claims OWNER TO dataops_plugin_platform_owner; ALTER TABLE public.plugin_runtime_breakers OWNER TO dataops_plugin_platform_owner; REVOKE ALL ON TABLE public.plugin_request_claims,public.plugin_runtime_breakers FROM PUBLIC,dataops_app,dataops_app_runtime; CREATE FUNCTION public.plugin_platform_control_v2(p jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $control$ DECLARE action text; r record; approval uuid; target text; incident uuid; fence bigint; incident_exists boolean; BEGIN IF NOT pg_has_role(session_user,'dataops_plugin_platform_control','MEMBER') OR jsonb_typeof(p)<>'object' OR p ?& ARRAY['action','plugin_uid','version','actor_ref'] IS FALSE OR p->>'plugin_uid' !~ '^[a-z][a-z0-9-]{2,62}$' OR p->>'actor_ref' !~ '^[A-Za-z0-9_.:-]{1,120}$' THEN RAISE EXCEPTION 'plugin_control_denied'; END IF; action:=p->>'action'; IF action='register' THEN IF p-ARRAY['action','plugin_uid','version','actor_ref','plugin_type','manifest_digest','artifact_digest','signature_digest','sbom_digest','license_digest','vulnerability_digest','provenance_digest','capabilities','resource']<>'{}'::jsonb OR p->>'version' !~ '^[0-9]+[.][0-9]+[.][0-9]+([-.][A-Za-z0-9.]+)?$' OR p->>'plugin_type' NOT IN ('connector','parser','quality','notification','approval','agent_mcp') OR p->>'manifest_digest' !~ '^[0-9a-f]{64}$' OR p->>'artifact_digest' !~ '^[0-9a-f]{64}$' OR p->>'signature_digest' <> encode(sha256(convert_to('local-fixture-key-v1|' || (p->>'artifact_digest'),'utf8')),'hex') OR jsonb_typeof(p->'capabilities')<>'array' OR jsonb_array_length(p->'capabilities') NOT BETWEEN 1 AND 8 OR p->'resource' <> '{"timeout_ms":100,"max_output_bytes":4096,"max_concurrency":1,"max_retries":1}'::jsonb THEN RAISE EXCEPTION 'plugin_register_closed'; END IF; INSERT INTO public.plugin_registry_versions(plugin_uid,version,plugin_type,manifest_digest,artifact_digest,trust_store_key_id,signature_digest,sbom_digest,license_digest,vulnerability_digest,provenance_digest,fixture_id,permissions,capabilities,resource,state,review_actor) VALUES(p->>'plugin_uid',p->>'version',p->>'plugin_type',p->>'manifest_digest',p->>'artifact_digest','local-fixture-key-v1',p->>'signature_digest',p->>'sbom_digest',p->>'license_digest',p->>'vulnerability_digest',p->>'provenance_digest','ENGINEERING_EVIDENCE_ONLY','{"file":false,"network":false,"secret":false,"child_process":false}'::jsonb,p->'capabilities',p->'resource','draft',p->>'actor_ref'); INSERT INTO public.plugin_audit_outbox(plugin_uid,version,tenant_ref,event_type,actor_ref,payload_digest) VALUES(p->>'plugin_uid',p->>'version','local-engineering','registered',p->>'actor_ref',p->>'manifest_digest'); RETURN jsonb_build_object('state','draft','plugin_uid',p->>'plugin_uid','version',p->>'version'); ELSIF action='review' THEN UPDATE public.plugin_registry_versions SET state='reviewed',review_actor=p->>'actor_ref',lifecycle_fence=lifecycle_fence+1,updated_at=clock_timestamp() WHERE plugin_uid=p->>'plugin_uid' AND version=p->>'version' AND state='draft' RETURNING lifecycle_fence INTO fence; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_review_transition_denied'; END IF; RETURN jsonb_build_object('state','reviewed','fence',fence); ELSIF action='describe' THEN IF p-ARRAY['action','plugin_uid','version','actor_ref']<>'{}'::jsonb THEN RAISE EXCEPTION 'plugin_describe_closed'; END IF; SELECT * INTO r FROM public.plugin_registry_versions WHERE plugin_uid=p->>'plugin_uid' AND version=p->>'version'; IF NOT FOUND OR r.state NOT IN ('canary','active') THEN RAISE EXCEPTION 'plugin_not_active'; END IF; RETURN jsonb_build_object('plugin_type',r.plugin_type,'capabilities',r.capabilities,'resource',r.resource,'manifest_digest',r.manifest_digest); ELSIF action='issue_approval' THEN IF p-ARRAY['action','plugin_uid','version','actor_ref','reviewer_ref','approval_action','tenant_ref','domain_ref','expires_in_seconds']<>'{}'::jsonb OR p->>'reviewer_ref'=p->>'actor_ref' OR p->>'approval_action' NOT IN ('approve','canary','activate','pause','rollback','revoke','recover') OR p->>'tenant_ref' !~ '^[a-z][a-z0-9-]{0,62}$' OR p->>'domain_ref' !~ '^[a-z][a-z0-9-]{0,62}$' OR (p->>'expires_in_seconds')::int NOT BETWEEN 1 AND 600 THEN RAISE EXCEPTION 'plugin_approval_closed'; END IF; SELECT * INTO r FROM public.plugin_registry_versions WHERE plugin_uid=p->>'plugin_uid' AND version=p->>'version'; IF NOT FOUND OR r.review_actor=p->>'actor_ref' THEN RAISE EXCEPTION 'plugin_approval_review_denied'; END IF; INSERT INTO public.plugin_approvals(plugin_uid,version,action_name,actor_ref,reviewer_ref,tenant_ref,domain_ref,manifest_digest,expires_at) VALUES(r.plugin_uid,r.version,p->>'approval_action',p->>'actor_ref',p->>'reviewer_ref',p->>'tenant_ref',p->>'domain_ref',r.manifest_digest,clock_timestamp()+make_interval(secs=>(p->>'expires_in_seconds')::int)) RETURNING approval_uid INTO approval; RETURN jsonb_build_object('approval_uid',approval::text,'manifest_digest',r.manifest_digest); ELSIF action='transition' THEN IF p-ARRAY['action','plugin_uid','version','actor_ref','tenant_ref','domain_ref','approval_uid','target_state','expected_fence','incident_uid']<>'{}'::jsonb OR p->>'target_state' NOT IN ('approved','canary','active','paused','rolled_back','revoked','recovery') THEN RAISE EXCEPTION 'plugin_transition_closed'; END IF; target:=p->>'target_state'; incident:=NULLIF(p->>'incident_uid','')::uuid; SELECT * INTO r FROM public.plugin_registry_versions WHERE plugin_uid=p->>'plugin_uid' AND version=p->>'version' FOR UPDATE; IF NOT FOUND OR r.lifecycle_fence<>(p->>'expected_fence')::bigint THEN RAISE EXCEPTION 'plugin_lifecycle_fence_conflict'; END IF; IF (target='approved' AND r.state='reviewed') OR (target='canary' AND r.state='approved') OR (target='active' AND r.state IN ('approved','canary','recovery')) OR (target IN ('paused','rolled_back','revoked') AND r.state IN ('active','canary')) OR (target='recovery' AND r.state IN ('paused','revoked','rolled_back')) THEN NULL; ELSE RAISE EXCEPTION 'plugin_transition_denied'; END IF; IF target IN ('paused','revoked','rolled_back') AND incident IS NULL THEN RAISE EXCEPTION 'plugin_incident_required'; END IF; IF target IN ('paused','revoked','rolled_back') AND 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; UPDATE public.plugin_approvals SET consumed_at=clock_timestamp() WHERE approval_uid=(p->>'approval_uid')::uuid AND consumed_at IS NULL AND expires_at>clock_timestamp() AND plugin_uid=r.plugin_uid AND version=r.version AND actor_ref=p->>'actor_ref' AND tenant_ref=p->>'tenant_ref' AND domain_ref=p->>'domain_ref' AND manifest_digest=r.manifest_digest AND action_name=(CASE target WHEN 'approved' THEN 'approve' WHEN 'active' THEN 'activate' WHEN 'rolled_back' THEN 'rollback' ELSE target END) AND reviewer_ref<>actor_ref; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_approval_binding_denied'; END IF; UPDATE public.plugin_registry_versions SET state=target,incident_uid=COALESCE(incident,incident_uid),lifecycle_fence=lifecycle_fence+1,updated_at=clock_timestamp(),activated_at=CASE WHEN target='active' THEN clock_timestamp() ELSE activated_at END WHERE plugin_uid=r.plugin_uid AND version=r.version AND lifecycle_fence=r.lifecycle_fence RETURNING lifecycle_fence INTO fence; INSERT INTO public.plugin_audit_outbox(plugin_uid,version,tenant_ref,event_type,actor_ref,payload_digest) VALUES(r.plugin_uid,r.version,p->>'tenant_ref',target,p->>'actor_ref',r.manifest_digest); RETURN jsonb_build_object('state',target,'fence',fence); END IF; RAISE EXCEPTION 'plugin_control_action_denied'; END $control$; CREATE FUNCTION public.plugin_platform_issue_claim(p jsonb) RETURNS jsonb LANGUAGE plpgsql SECURITY DEFINER SET search_path=pg_catalog,public AS $claim$ DECLARE claim uuid; d text; BEGIN IF NOT pg_has_role(session_user,'dataops_plugin_platform_control','MEMBER') OR jsonb_typeof(p)<>'object' OR p-ARRAY['plugin_uid','version','tenant_ref','domain_ref','principal_ref','request_digest']<>'{}'::jsonb OR p->>'request_digest' !~ '^[0-9a-f]{64}$' THEN RAISE EXCEPTION 'plugin_claim_denied'; END IF; IF NOT EXISTS(SELECT 1 FROM public.plugin_registry_versions WHERE plugin_uid=p->>'plugin_uid' AND version=p->>'version' AND state IN ('canary','active')) THEN RAISE EXCEPTION 'plugin_not_active'; END IF; INSERT INTO public.plugin_request_claims(plugin_uid,version,action_name,tenant_ref,domain_ref,principal_ref,request_digest,expires_at) VALUES(p->>'plugin_uid',p->>'version','invoke',p->>'tenant_ref',p->>'domain_ref',p->>'principal_ref',p->>'request_digest',clock_timestamp()+interval '30 seconds') RETURNING claim_uid INTO claim; RETURN jsonb_build_object('claim_uid',claim::text); END $claim$; CREATE 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; outd text; 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; 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); 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; RETURN jsonb_build_object('state',run.state,'attempt_count',run.attempt_count,'fence',run.lease_fence); ELSIF a='recover' THEN UPDATE public.plugin_runs SET state='pending',attempt_count=0,next_attempt_at=clock_timestamp(),lease_owner=NULL,lease_expires_at=NULL WHERE run_uid=(p->>'run_uid')::uuid AND state='dead_letter' RETURNING * INTO run; IF NOT FOUND THEN RAISE EXCEPTION 'plugin_recovery_denied'; END IF; DELETE FROM public.plugin_dead_letters WHERE plugin_dead_letters.run_uid=run.run_uid; 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; RETURN jsonb_build_object('state','pending'); END IF; RAISE EXCEPTION 'plugin_runtime_action_denied'; END $runtime$; ALTER FUNCTION public.plugin_platform_control_v2(jsonb) OWNER TO dataops_plugin_platform_owner; ALTER FUNCTION public.plugin_platform_issue_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_control_v2(jsonb),public.plugin_platform_issue_claim(jsonb),public.plugin_platform_runtime_v2(jsonb) FROM PUBLIC,dataops_app,dataops_app_runtime; GRANT EXECUTE ON FUNCTION public.plugin_platform_control_v2(jsonb),public.plugin_platform_issue_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_request_claims LIMIT 1) OR EXISTS(SELECT 1 FROM public.plugin_runtime_breakers LIMIT 1) OR EXISTS(SELECT 1 FROM public.plugin_runs LIMIT 1) OR EXISTS(SELECT 1 FROM public.plugin_approvals LIMIT 1) OR EXISTS(SELECT 1 FROM public.plugin_registry_versions LIMIT 1)").scalar(): raise RuntimeError("downgrade refused: WP13 lifecycle facts are nonempty") bind.exec_driver_sql(""" REVOKE ALL ON FUNCTION public.plugin_platform_control_v2(jsonb),public.plugin_platform_issue_claim(jsonb),public.plugin_platform_runtime_v2(jsonb) FROM PUBLIC,dataops_app,dataops_app_runtime,dataops_plugin_platform_control; DROP FUNCTION public.plugin_platform_runtime_v2(jsonb); DROP FUNCTION public.plugin_platform_issue_claim(jsonb); DROP FUNCTION public.plugin_platform_control_v2(jsonb); DROP TABLE public.plugin_runtime_breakers,public.plugin_request_claims; ALTER TABLE public.plugin_runs DROP COLUMN lease_expires_at,DROP COLUMN lease_owner,DROP COLUMN next_attempt_at,DROP COLUMN attempt_count; ALTER TABLE public.plugin_runs DROP CONSTRAINT plugin_runs_state_check; ALTER TABLE public.plugin_runs ADD CONSTRAINT plugin_runs_state_check CHECK(state IN ('running','succeeded','failed','dead_letter')); ALTER TABLE public.plugin_registry_versions DROP COLUMN resource,DROP COLUMN capabilities,DROP COLUMN updated_at,DROP COLUMN incident_uid,DROP COLUMN failure_count,DROP COLUMN lifecycle_fence; """)