-- Copyright (c) Microsoft Corporation. -- Licensed under the PostgreSQL License. /* */ /* This file is auto generated by pgrx. The ordering of items is not stable, it is driven by a dependency graph. */ /* */ /* */ -- src/lib.rs:75 CREATE SCHEMA IF NOT EXISTS df; /* pg_durable::df */ /* */ /* */ -- src/lib.rs:82 -- requires: -- df -- Table to store function nodes (SQL steps, THEN chains, etc.) CREATE TABLE IF NOT EXISTS df.nodes ( id VARCHAR(8) PRIMARY KEY, instance_id VARCHAR(8), node_type TEXT NOT NULL, query TEXT, result_name TEXT, left_node VARCHAR(8), right_node VARCHAR(8), status TEXT DEFAULT 'pending', result JSONB, error TEXT, submitted_by REGROLE, login_role REGROLE, database TEXT, created_at TIMESTAMPTZ DEFAULT now(), updated_at TIMESTAMPTZ DEFAULT now() ); COMMENT ON COLUMN df.nodes.submitted_by IS 'Effective role (outer user) for privilege isolation. Set by df.start() when node is linked to an instance.'; COMMENT ON COLUMN df.nodes.login_role IS 'Authenticated role (session user) for connection authentication. Set by df.start() when node is linked to an instance.'; -- Table to store function instances CREATE TABLE IF NOT EXISTS df.instances ( id VARCHAR(8) PRIMARY KEY, label TEXT, root_node VARCHAR(8) NOT NULL, status TEXT DEFAULT 'pending', submitted_by REGROLE NOT NULL, login_role REGROLE NOT NULL, database TEXT, created_at TIMESTAMPTZ DEFAULT now(), updated_at TIMESTAMPTZ DEFAULT now(), completed_at TIMESTAMPTZ ); COMMENT ON COLUMN df.instances.submitted_by IS 'Effective role (outer user) when df.start() was called - used for SET ROLE during execution'; COMMENT ON COLUMN df.instances.login_role IS 'Authenticated role (session user) when df.start() was called - used for connection authentication'; -- Index for finding pending instances CREATE INDEX IF NOT EXISTS idx_instances_status ON df.instances(status); -- Index for finding nodes by instance CREATE INDEX IF NOT EXISTS idx_nodes_instance ON df.nodes(instance_id); -- Table to store workflow variables (captured at df.start()) CREATE TABLE IF NOT EXISTS df.vars ( name TEXT PRIMARY KEY, value TEXT ); -- Sentinel table: the background worker writes its epoch_id here after -- initialising. If the extension is DROP-ed and re-CREATEd between -- two poll ticks the epoch row disappears, so the worker detects the -- recreation even though the extension is always "present" in pg_extension. CREATE TABLE IF NOT EXISTS df._worker_epoch ( epoch_id UUID PRIMARY KEY, started_at TIMESTAMPTZ DEFAULT now(), last_seen_at TIMESTAMPTZ DEFAULT now() ); /* */ /* */ -- src/lib.rs:157 -- requires: -- create_tables -- Enable RLS on df.instances (no FORCE — superuser/table-owner bypasses RLS) ALTER TABLE df.instances ENABLE ROW LEVEL SECURITY; CREATE POLICY instances_user_isolation ON df.instances FOR ALL USING (submitted_by = current_user::regrole) WITH CHECK (submitted_by = current_user::regrole); -- Enable RLS on df.nodes ALTER TABLE df.nodes ENABLE ROW LEVEL SECURITY; CREATE POLICY nodes_user_isolation ON df.nodes FOR ALL USING (submitted_by = current_user::regrole) WITH CHECK (submitted_by = current_user::regrole); -- Auto-grant permissions to PUBLIC (safe with RLS enabled) GRANT USAGE ON SCHEMA df TO PUBLIC; GRANT EXECUTE ON ALL FUNCTIONS IN SCHEMA df TO PUBLIC; -- Users need INSERT for df.start(), SELECT for df.status()/result() -- Column-level UPDATE on instances: only status + updated_at (for df.cancel()) -- No UPDATE on identity columns (submitted_by, login_role) or structural columns (root_node) -- No DELETE — instance/node deletion should happen via admin API or TTL GRANT SELECT, INSERT ON df.instances TO PUBLIC; GRANT UPDATE (status, updated_at) ON df.instances TO PUBLIC; GRANT SELECT, INSERT ON df.nodes TO PUBLIC; GRANT SELECT, INSERT, UPDATE, DELETE ON df.vars TO PUBLIC; -- Validate that the worker role is a superuser. -- The background worker must bypass RLS to manage all users' instances/nodes. -- If the worker role is not a superuser, workflows will silently fail because -- RLS will filter out rows the worker needs to read/update. DO $$ DECLARE wrole TEXT; is_super BOOLEAN; BEGIN wrole := current_setting('pg_durable.worker_role', true); IF wrole IS NULL OR wrole = '' THEN wrole := 'azuresu'; END IF; SELECT rolsuper INTO is_super FROM pg_roles WHERE rolname = wrole; IF is_super IS NULL THEN RAISE WARNING 'pg_durable: worker role "%" does not exist. The background worker will not be able to process workflows. Create the role as a superuser before using pg_durable.', wrole; ELSIF NOT is_super THEN RAISE WARNING 'pg_durable: worker role "%" is not a superuser. The background worker must be a superuser to bypass RLS and manage all users'' instances. Grant superuser or BYPASSRLS to this role.', wrole; END IF; END $$; /* */ /* */ -- src/dsl.rs:258 -- pg_durable::dsl::break CREATE FUNCTION df."break"( "value" TEXT DEFAULT NULL /* core::option::Option<&str> */ ) RETURNS TEXT /* alloc::string::String */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'break_fn_wrapper'; /* */ /* */ -- src/dsl.rs:312 -- pg_durable::dsl::join3 CREATE FUNCTION df."join3"( "a" TEXT, /* &str */ "b" TEXT, /* &str */ "c" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'join3_wrapper'; /* */ /* */ -- src/explain.rs:33 -- pg_durable::explain::explain CREATE FUNCTION df."explain"( "input" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_wrapper'; /* */ /* */ -- src/dsl.rs:28 -- pg_durable::dsl::version CREATE FUNCTION df."version"() RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'version_wrapper'; /* */ /* */ -- src/dsl.rs:334 -- pg_durable::dsl::race CREATE FUNCTION df."race"( "a" TEXT, /* &str */ "b" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'race_wrapper'; /* */ /* */ -- src/dsl.rs:731 -- pg_durable::dsl::result CREATE FUNCTION df."result"( "instance_id" TEXT /* &str */ ) RETURNS TEXT /* core::option::Option */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'result_wrapper'; /* */ /* */ -- src/dsl.rs:102 -- pg_durable::dsl::clearvars CREATE FUNCTION df."clearvars"() RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'clearvars_wrapper'; /* */ /* */ -- src/dsl.rs:84 -- pg_durable::dsl::unsetvar CREATE FUNCTION df."unsetvar"( "name" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'unsetvar_wrapper'; /* */ /* */ -- src/dsl.rs:721 -- pg_durable::dsl::run CREATE FUNCTION df."run"( "instance_id" TEXT DEFAULT NULL /* core::option::Option<&str> */ ) RETURNS TEXT /* alloc::string::String */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'run_wrapper'; /* */ /* */ -- src/dsl.rs:133 -- pg_durable::dsl::seq CREATE FUNCTION df."seq"( "a" TEXT, /* &str */ "b" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'then_fn_wrapper'; /* */ /* */ -- src/dsl.rs:274 -- pg_durable::dsl::if CREATE FUNCTION df."if"( "condition" TEXT, /* &str */ "then_branch" TEXT, /* &str */ "else_branch" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'if_fn_wrapper'; /* */ /* */ -- src/monitoring.rs:265 -- pg_durable::monitoring::metrics CREATE FUNCTION df."metrics"() RETURNS TABLE ( "total_instances" bigint, /* i64 */ "running_instances" bigint, /* i64 */ "completed_instances" bigint, /* i64 */ "failed_instances" bigint, /* i64 */ "total_executions" bigint, /* i64 */ "total_events" bigint /* i64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'metrics_wrapper'; /* */ /* */ -- src/dsl.rs:492 -- pg_durable::dsl::start CREATE FUNCTION df."start"( "fut" TEXT, /* &str */ "label" TEXT DEFAULT NULL, /* core::option::Option<&str> */ "database" TEXT DEFAULT NULL /* core::option::Option<&str> */ ) RETURNS TEXT /* alloc::string::String */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'start_wrapper'; /* */ /* */ -- src/monitoring.rs:314 -- pg_durable::monitoring::instance_nodes CREATE FUNCTION df."instance_nodes"( "instance_id_param" TEXT, /* &str */ "last_n_executions" INT DEFAULT 5 /* i32 */ ) RETURNS TABLE ( "execution_id" bigint, /* i64 */ "node_id" TEXT, /* alloc::string::String */ "node_type" TEXT, /* alloc::string::String */ "query" TEXT, /* core::option::Option */ "result_name" TEXT, /* core::option::Option */ "left_node" TEXT, /* core::option::Option */ "right_node" TEXT, /* core::option::Option */ "status" TEXT, /* core::option::Option */ "result" TEXT, /* core::option::Option */ "updated_at" timestamp with time zone /* core::option::Option */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'instance_nodes_wrapper'; /* */ /* */ -- src/dsl.rs:120 -- pg_durable::dsl::sql CREATE FUNCTION df."sql"( "query" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'sql_wrapper'; /* */ /* */ -- src/dsl.rs:681 -- pg_durable::dsl::cancel CREATE FUNCTION df."cancel"( "instance_id" TEXT, /* &str */ "reason" TEXT DEFAULT 'Cancelled by user' /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'cancel_wrapper'; /* */ /* */ -- src/dsl.rs:360 -- pg_durable::dsl::http CREATE FUNCTION df."http"( "url" TEXT, /* &str */ "method" TEXT DEFAULT 'POST', /* &str */ "body" TEXT DEFAULT NULL, /* core::option::Option<&str> */ "headers" jsonb DEFAULT NULL, /* core::option::Option */ "timeout_seconds" INT DEFAULT 30 /* i32 */ ) RETURNS TEXT /* alloc::string::String */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'http_wrapper'; /* */ /* */ -- src/dsl.rs:54 -- pg_durable::dsl::setvar CREATE FUNCTION df."setvar"( "name" TEXT, /* &str */ "value" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'setvar_wrapper'; /* */ /* */ -- src/dsl.rs:414 -- pg_durable::dsl::wait_for_signal CREATE FUNCTION df."wait_for_signal"( "name" TEXT, /* &str */ "timeout_seconds" INT DEFAULT NULL /* core::option::Option */ ) RETURNS TEXT /* alloc::string::String */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'wait_for_signal_wrapper'; /* */ /* */ -- src/dsl.rs:176 -- pg_durable::dsl::wait_for_schedule CREATE FUNCTION df."wait_for_schedule"( "cron_expr" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'wait_for_schedule_wrapper'; /* */ /* */ -- src/dsl.rs:160 -- pg_durable::dsl::sleep CREATE FUNCTION df."sleep"( "seconds" bigint /* i64 */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'sleep_wrapper'; /* */ /* */ -- src/dsl.rs:223 -- pg_durable::dsl::loop CREATE FUNCTION df."loop"( "body" TEXT, /* &str */ "condition" TEXT DEFAULT NULL /* core::option::Option<&str> */ ) RETURNS TEXT /* alloc::string::String */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'loop_fn_wrapper'; /* */ /* */ -- src/types.rs:70 -- pg_durable::types::target_database CREATE FUNCTION df."target_database"() RETURNS TEXT /* alloc::string::String */ IMMUTABLE STRICT PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'target_database_wrapper'; /* */ /* */ -- src/lib.rs:223 -- requires: -- df -- target_database -- Validate that CREATE EXTENSION is run in the correct database -- The background worker connects to one specific database (determined by -- the pg_durable.database GUC, defaults to "postgres"). -- The extension must be created in that database for workflows to execute. DO $$ DECLARE current_db TEXT; target_db TEXT; BEGIN -- Get the current database SELECT current_database() INTO current_db; -- Get the target database that the background worker will connect to SELECT df.target_database() INTO target_db; IF current_db != target_db THEN RAISE EXCEPTION 'pg_durable extension must be created in database "%" (currently in "%"). The background worker only processes functions in the database specified by the pg_durable.database GUC (defaults to "postgres").', target_db, current_db USING HINT = 'Connect to the correct database and run: CREATE EXTENSION pg_durable;'; END IF; END $$; /* */ /* */ -- src/lib.rs:270 -- requires: -- validate_database -- BEGIN duroxide-pg-opt migrations (checked-in copy) CREATE SCHEMA IF NOT EXISTS duroxide; SET LOCAL search_path TO duroxide; -- Migration 0001: 0001_initial_schema.sql -- Migration: 0001_initial_schema.sql -- Description: Complete Duroxide PostgreSQL provider schema -- All tables, indexes, and stored procedures in a single migration. -- All timestamps are NOT NULL without defaults - provider must supply values. -- ============================================================================ -- Tables -- ============================================================================ -- Instance metadata CREATE TABLE IF NOT EXISTS instances ( instance_id TEXT PRIMARY KEY, orchestration_name TEXT NOT NULL, orchestration_version TEXT, -- NULLable, set by runtime via ack_orchestration_item metadata current_execution_id BIGINT NOT NULL DEFAULT 1, created_at TIMESTAMPTZ NOT NULL, updated_at TIMESTAMPTZ NOT NULL ); -- Multi-execution support CREATE TABLE IF NOT EXISTS executions ( instance_id TEXT NOT NULL, execution_id BIGINT NOT NULL, status TEXT NOT NULL DEFAULT 'Running', output TEXT, started_at TIMESTAMPTZ NOT NULL, completed_at TIMESTAMPTZ, -- NULL until completed/failed PRIMARY KEY (instance_id, execution_id) ); -- Event history (append-only) CREATE TABLE IF NOT EXISTS history ( instance_id TEXT NOT NULL, execution_id BIGINT NOT NULL, event_id BIGINT NOT NULL, event_type TEXT NOT NULL, event_data TEXT NOT NULL, -- JSON serialized Event created_at TIMESTAMPTZ NOT NULL, PRIMARY KEY (instance_id, execution_id, event_id) ); -- Orchestrator queue CREATE TABLE IF NOT EXISTS orchestrator_queue ( id BIGSERIAL PRIMARY KEY, instance_id TEXT NOT NULL, work_item TEXT NOT NULL, -- JSON serialized WorkItem visible_at TIMESTAMPTZ NOT NULL, lock_token TEXT, locked_until BIGINT, -- Unix timestamp in milliseconds created_at TIMESTAMPTZ NOT NULL, attempt_count INTEGER NOT NULL DEFAULT 0 ); -- Worker queue CREATE TABLE IF NOT EXISTS worker_queue ( id BIGSERIAL PRIMARY KEY, work_item TEXT NOT NULL, -- JSON serialized WorkItem visible_at TIMESTAMPTZ NOT NULL, -- When the item becomes available for processing lock_token TEXT, locked_until BIGINT, -- Unix timestamp in milliseconds created_at TIMESTAMPTZ NOT NULL, attempt_count INTEGER NOT NULL DEFAULT 0 ); -- Instance-level locks for concurrent dispatcher coordination CREATE TABLE IF NOT EXISTS instance_locks ( instance_id TEXT PRIMARY KEY, lock_token TEXT NOT NULL, locked_until BIGINT NOT NULL, -- Unix timestamp in milliseconds locked_at BIGINT NOT NULL -- Unix timestamp in milliseconds ); -- ============================================================================ -- Indexes -- ============================================================================ CREATE INDEX IF NOT EXISTS idx_orch_visible ON orchestrator_queue(visible_at, lock_token); CREATE INDEX IF NOT EXISTS idx_orch_instance ON orchestrator_queue(instance_id); CREATE INDEX IF NOT EXISTS idx_orch_lock ON orchestrator_queue(lock_token); CREATE INDEX IF NOT EXISTS idx_worker_visible ON worker_queue(visible_at, lock_token); CREATE INDEX IF NOT EXISTS idx_worker_available ON worker_queue(lock_token, id); CREATE INDEX IF NOT EXISTS idx_instance_locks_locked_until ON instance_locks(locked_until); CREATE INDEX IF NOT EXISTS idx_history_lookup ON history(instance_id, execution_id, event_id); -- Migration tracking table (create in each schema) CREATE TABLE IF NOT EXISTS _duroxide_migrations ( version BIGINT PRIMARY KEY, name TEXT NOT NULL, applied_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); -- ============================================================================ -- NOTIFY Triggers for Long-Polling -- Triggers fire on INSERT to notify waiting dispatchers of new work. -- Payload contains visible_at as epoch milliseconds for timer scheduling. -- NOTE: We use TG_TABLE_SCHEMA (the schema of the table being modified) -- instead of current_schema() because current_schema() returns the first -- schema in the session's search_path, which may not be our schema. -- ============================================================================ -- Trigger function for orchestrator queue DROP FUNCTION IF EXISTS notify_orch_work() CASCADE; CREATE OR REPLACE FUNCTION notify_orch_work() RETURNS TRIGGER AS $$ BEGIN PERFORM pg_notify( TG_TABLE_SCHEMA || '_orch_work', (EXTRACT(EPOCH FROM NEW.visible_at) * 1000)::BIGINT::TEXT ); RETURN NEW; END; $$ LANGUAGE plpgsql; -- Trigger function for worker queue -- Worker queue items now have visible_at for delayed visibility -- Send visible_at timestamp for timer scheduling DROP FUNCTION IF EXISTS notify_worker_work() CASCADE; CREATE OR REPLACE FUNCTION notify_worker_work() RETURNS TRIGGER AS $$ BEGIN PERFORM pg_notify( TG_TABLE_SCHEMA || '_worker_work', (EXTRACT(EPOCH FROM NEW.visible_at) * 1000)::BIGINT::TEXT ); RETURN NEW; END; $$ LANGUAGE plpgsql; -- Attach triggers to queues DROP TRIGGER IF EXISTS trg_notify_orch_work ON orchestrator_queue; CREATE TRIGGER trg_notify_orch_work AFTER INSERT ON orchestrator_queue FOR EACH ROW EXECUTE FUNCTION notify_orch_work(); DROP TRIGGER IF EXISTS trg_notify_worker_work ON worker_queue; CREATE TRIGGER trg_notify_worker_work AFTER INSERT ON worker_queue FOR EACH ROW EXECUTE FUNCTION notify_worker_work(); -- ============================================================================ -- Stored Procedures -- ============================================================================ DO $$ DECLARE v_schema_name TEXT := current_schema(); BEGIN -- ============================================================================ -- Schema Management Procedures -- ============================================================================ -- Procedure: cleanup_schema -- Drops all tables AND functions in the schema (for testing only) -- Functions must be dropped because PostgreSQL cannot change return types with CREATE OR REPLACE EXECUTE format('DROP FUNCTION IF EXISTS %I.cleanup_schema()', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.cleanup_schema() RETURNS VOID AS $cleanup$ BEGIN -- Drop tables first DROP TABLE IF EXISTS %I.instances CASCADE; DROP TABLE IF EXISTS %I.executions CASCADE; DROP TABLE IF EXISTS %I.history CASCADE; DROP TABLE IF EXISTS %I.orchestrator_queue CASCADE; DROP TABLE IF EXISTS %I.worker_queue CASCADE; DROP TABLE IF EXISTS %I.instance_locks CASCADE; DROP TABLE IF EXISTS %I._duroxide_migrations CASCADE; -- Drop all stored procedures (required because return type changes cannot use CREATE OR REPLACE) DROP FUNCTION IF EXISTS %I.cleanup_schema(); DROP FUNCTION IF EXISTS %I.list_instances(); DROP FUNCTION IF EXISTS %I.list_executions(TEXT); DROP FUNCTION IF EXISTS %I.latest_execution_id(TEXT); DROP FUNCTION IF EXISTS %I.list_instances_by_status(TEXT); DROP FUNCTION IF EXISTS %I.get_instance_info(TEXT); DROP FUNCTION IF EXISTS %I.get_execution_info(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.get_system_metrics(); DROP FUNCTION IF EXISTS %I.get_queue_depths(BIGINT); DROP FUNCTION IF EXISTS %I.enqueue_worker_work(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.renew_work_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_work_item(BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.abandon_work_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.enqueue_orchestrator_work(TEXT, TEXT, TIMESTAMPTZ, BIGINT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_orchestration_item(BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB); DROP FUNCTION IF EXISTS %I.abandon_orchestration_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.renew_orchestration_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_history(TEXT); DROP FUNCTION IF EXISTS %I.fetch_history_with_execution(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.append_history(TEXT, BIGINT, JSONB, BIGINT); -- Drop trigger functions (not schema-qualified, they use search_path) -- CASCADE is required because triggers depend on these functions DROP FUNCTION IF EXISTS notify_orch_work() CASCADE; DROP FUNCTION IF EXISTS notify_worker_work() CASCADE; END; $cleanup$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Simple Query Procedures -- ============================================================================ -- Procedure: list_instances EXECUTE format('DROP FUNCTION IF EXISTS %I.list_instances()', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.list_instances() RETURNS TABLE(instance_id TEXT) AS $list_inst$ BEGIN RETURN QUERY SELECT i.instance_id FROM %I.instances i ORDER BY i.created_at DESC; END; $list_inst$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: list_executions EXECUTE format('DROP FUNCTION IF EXISTS %I.list_executions(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.list_executions(p_instance_id TEXT) RETURNS TABLE(execution_id BIGINT) AS $list_exec$ BEGIN RETURN QUERY SELECT e.execution_id FROM %I.executions e WHERE e.instance_id = p_instance_id ORDER BY e.execution_id; END; $list_exec$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: latest_execution_id EXECUTE format('DROP FUNCTION IF EXISTS %I.latest_execution_id(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.latest_execution_id(p_instance_id TEXT) RETURNS BIGINT AS $latest_exec$ DECLARE v_execution_id BIGINT; BEGIN SELECT i.current_execution_id INTO v_execution_id FROM %I.instances i WHERE i.instance_id = p_instance_id; RETURN v_execution_id; END; $latest_exec$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: list_instances_by_status EXECUTE format('DROP FUNCTION IF EXISTS %I.list_instances_by_status(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.list_instances_by_status(p_status TEXT) RETURNS TABLE(instance_id TEXT) AS $list_by_status$ BEGIN RETURN QUERY SELECT i.instance_id FROM %I.instances i JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE e.status = p_status ORDER BY i.created_at DESC; END; $list_by_status$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Join and Aggregate Query Procedures -- ============================================================================ -- Procedure: get_instance_info EXECUTE format('DROP FUNCTION IF EXISTS %I.get_instance_info(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_instance_info(p_instance_id TEXT) RETURNS TABLE( instance_id TEXT, orchestration_name TEXT, orchestration_version TEXT, current_execution_id BIGINT, created_at TIMESTAMPTZ, updated_at TIMESTAMPTZ, status TEXT, output TEXT ) AS $get_inst_info$ BEGIN RETURN QUERY SELECT i.instance_id, i.orchestration_name, COALESCE(i.orchestration_version, ''unknown'') as orchestration_version, i.current_execution_id, i.created_at, i.updated_at, e.status, e.output FROM %I.instances i LEFT JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE i.instance_id = p_instance_id; END; $get_inst_info$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- Procedure: get_execution_info EXECUTE format('DROP FUNCTION IF EXISTS %I.get_execution_info(TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_execution_info( p_instance_id TEXT, p_execution_id BIGINT ) RETURNS TABLE( execution_id BIGINT, status TEXT, output TEXT, started_at TIMESTAMPTZ, completed_at TIMESTAMPTZ, event_count BIGINT ) AS $get_exec_info$ BEGIN RETURN QUERY SELECT e.execution_id, e.status, e.output, e.started_at, e.completed_at, COALESCE(COUNT(h.event_id), 0)::BIGINT as event_count FROM %I.executions e LEFT JOIN %I.history h ON e.instance_id = h.instance_id AND e.execution_id = h.execution_id WHERE e.instance_id = p_instance_id AND e.execution_id = p_execution_id GROUP BY e.execution_id, e.status, e.output, e.started_at, e.completed_at; END; $get_exec_info$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- Procedure: get_system_metrics EXECUTE format('DROP FUNCTION IF EXISTS %I.get_system_metrics()', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_system_metrics() RETURNS TABLE( total_instances BIGINT, total_executions BIGINT, running_instances BIGINT, completed_instances BIGINT, failed_instances BIGINT, total_events BIGINT ) AS $get_metrics$ BEGIN RETURN QUERY SELECT (SELECT COUNT(*)::BIGINT FROM %I.instances) as total_instances, (SELECT COUNT(*)::BIGINT FROM %I.executions) as total_executions, (SELECT COUNT(DISTINCT i.instance_id)::BIGINT FROM %I.instances i JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE e.status = ''Running'') as running_instances, (SELECT COUNT(DISTINCT i.instance_id)::BIGINT FROM %I.instances i JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE e.status = ''Completed'') as completed_instances, (SELECT COUNT(DISTINCT i.instance_id)::BIGINT FROM %I.instances i JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE e.status = ''Failed'') as failed_instances, (SELECT COUNT(*)::BIGINT FROM %I.history) as total_events; END; $get_metrics$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- Procedure: get_queue_depths -- Returns count of items available for processing (visible and unlocked/lock expired) EXECUTE format('DROP FUNCTION IF EXISTS %I.get_queue_depths(BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_queue_depths(p_now_ms BIGINT) RETURNS TABLE( orchestrator_queue BIGINT, worker_queue BIGINT ) AS $get_queue_depths$ BEGIN RETURN QUERY SELECT (SELECT COUNT(*)::BIGINT FROM %I.orchestrator_queue WHERE visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (lock_token IS NULL OR locked_until <= p_now_ms)) as orchestrator_queue, (SELECT COUNT(*)::BIGINT FROM %I.worker_queue WHERE visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (lock_token IS NULL OR locked_until <= p_now_ms)) as worker_queue; END; $get_queue_depths$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Queue Operation Procedures -- All procedures accept p_now_ms for timestamp generation (Rust clock only) -- ============================================================================ -- Procedure: enqueue_worker_work EXECUTE format('DROP FUNCTION IF EXISTS %I.enqueue_worker_work(TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.enqueue_worker_work( p_work_item TEXT, p_now_ms BIGINT ) RETURNS VOID AS $enq_worker$ DECLARE v_now_ts TIMESTAMPTZ; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); INSERT INTO %I.worker_queue (work_item, visible_at, created_at) VALUES (p_work_item, v_now_ts, v_now_ts); END; $enq_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: ack_worker -- When p_completion_json is NULL, only delete from worker_queue (no enqueue) -- This is used when the orchestration is terminal or missing EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_worker( p_lock_token TEXT, p_instance_id TEXT, p_completion_json TEXT, p_now_ms BIGINT ) RETURNS VOID AS $ack_worker$ DECLARE v_rows_affected INTEGER; v_now_ts TIMESTAMPTZ; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); DELETE FROM %I.worker_queue WHERE lock_token = p_lock_token; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Worker queue item not found or already processed''; END IF; -- Validate: if completion provided, instance_id must also be provided IF p_completion_json IS NOT NULL AND p_instance_id IS NULL THEN RAISE EXCEPTION ''instance_id required when completion_json is provided''; END IF; -- Only enqueue completion if provided (not NULL) IF p_completion_json IS NOT NULL THEN INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (p_instance_id, p_completion_json, v_now_ts, v_now_ts); END IF; END; $ack_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- Procedure: renew_work_item_lock -- Returns execution_status for cancellation support -- Note: DROP first because return type changed from VOID to TEXT EXECUTE format('DROP FUNCTION IF EXISTS %I.renew_work_item_lock(TEXT, BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.renew_work_item_lock( p_lock_token TEXT, p_now_ms BIGINT, p_extend_ms BIGINT ) RETURNS TEXT AS $renew_lock$ DECLARE v_rows_affected INTEGER; v_work_item_json JSONB; v_instance_id TEXT; v_execution_id BIGINT; v_execution_status TEXT; BEGIN -- Get the work item before updating SELECT work_item::JSONB INTO v_work_item_json FROM %I.worker_queue WHERE lock_token = p_lock_token AND locked_until > p_now_ms; IF NOT FOUND THEN RAISE EXCEPTION ''Lock token invalid, expired, or already acked''; END IF; -- Check execution status BEFORE extending the lock -- Per provider contract: lock can only be renewed if execution is Running IF v_work_item_json ? ''ActivityExecute'' THEN v_instance_id := v_work_item_json->''ActivityExecute''->>''instance''; v_execution_id := (v_work_item_json->''ActivityExecute''->>''execution_id'')::BIGINT; -- Get execution status directly from executions table -- Note: We check executions table, not instances table, because: -- 1. Instance record may not exist if ack_orchestration_item was called with NULL version -- 2. Execution record is the authoritative source for execution state SELECT e.status INTO v_execution_status FROM %I.executions e WHERE e.instance_id = v_instance_id AND e.execution_id = v_execution_id; IF v_execution_status IS NULL THEN -- Execution record missing - return NULL to signal Missing state -- Do NOT extend lock per contract RETURN NULL; END IF; IF v_execution_status <> ''Running'' THEN -- Execution is terminal - return status but do NOT extend lock per contract RETURN v_execution_status; END IF; END IF; -- Only extend lock if execution is Running (or non-ActivityExecute item) UPDATE %I.worker_queue SET locked_until = GREATEST(locked_until, p_now_ms) + p_extend_ms WHERE lock_token = p_lock_token AND locked_until > p_now_ms; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Lock token invalid, expired, or already acked''; END IF; RETURN v_execution_status; END; $renew_lock$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- Procedure: fetch_work_item -- Item is available if: -- 1. visible_at <= now (not delayed) -- 2. AND (lock_token IS NULL OR locked_until <= now) (not locked or lock expired) -- Returns execution_status from the execution table for cancellation support -- Note: DROP first because return type changed to include out_execution_status EXECUTE format('DROP FUNCTION IF EXISTS %I.fetch_work_item(BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.fetch_work_item( p_now_ms BIGINT, p_lock_timeout_ms BIGINT ) RETURNS TABLE( out_work_item TEXT, out_lock_token TEXT, out_attempt_count INTEGER, out_execution_status TEXT ) AS $fetch_worker$ DECLARE v_id BIGINT; v_work_item_json JSONB; v_instance_id TEXT; v_execution_id BIGINT; BEGIN SELECT q.id INTO v_id FROM %I.worker_queue q WHERE q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (q.lock_token IS NULL OR q.locked_until <= p_now_ms) ORDER BY q.id LIMIT 1 FOR UPDATE OF q SKIP LOCKED; IF NOT FOUND THEN RETURN; END IF; out_lock_token := ''lock_'' || gen_random_uuid()::TEXT; UPDATE %I.worker_queue SET lock_token = out_lock_token, locked_until = p_now_ms + p_lock_timeout_ms, attempt_count = attempt_count + 1 WHERE id = v_id; SELECT work_item, attempt_count INTO out_work_item, out_attempt_count FROM %I.worker_queue WHERE id = v_id; -- Parse work item to get instance and execution_id for status lookup v_work_item_json := out_work_item::JSONB; IF v_work_item_json ? ''ActivityExecute'' THEN v_instance_id := v_work_item_json->''ActivityExecute''->>''instance''; v_execution_id := (v_work_item_json->''ActivityExecute''->>''execution_id'')::BIGINT; SELECT e.status INTO out_execution_status FROM %I.executions e WHERE e.instance_id = v_instance_id AND e.execution_id = v_execution_id; ELSE out_execution_status := NULL; END IF; RETURN NEXT; END; $fetch_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- Procedure: abandon_work_item -- Always clear lock_token and locked_until when abandoning. -- Use visible_at to control when item becomes available again. EXECUTE format('DROP FUNCTION IF EXISTS %I.abandon_work_item(TEXT, BIGINT, BIGINT, BOOLEAN)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.abandon_work_item( p_lock_token TEXT, p_now_ms BIGINT, p_delay_ms BIGINT DEFAULT NULL, p_ignore_attempt BOOLEAN DEFAULT FALSE ) RETURNS VOID AS $abandon_worker$ DECLARE v_rows_affected INTEGER; v_visible_at TIMESTAMPTZ; BEGIN -- Calculate visible_at based on delay using Rust-provided time IF p_delay_ms IS NOT NULL AND p_delay_ms > 0 THEN v_visible_at := TO_TIMESTAMP((p_now_ms + p_delay_ms) / 1000.0); ELSE v_visible_at := TO_TIMESTAMP(p_now_ms / 1000.0); END IF; -- Always clear lock_token and locked_until when abandoning -- Use visible_at to control when item becomes available again IF p_ignore_attempt THEN UPDATE %I.worker_queue SET lock_token = NULL, locked_until = NULL, visible_at = v_visible_at, attempt_count = GREATEST(0, attempt_count - 1) WHERE lock_token = p_lock_token; ELSE UPDATE %I.worker_queue SET lock_token = NULL, locked_until = NULL, visible_at = v_visible_at WHERE lock_token = p_lock_token; END IF; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Invalid lock token or already acked''; END IF; END; $abandon_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- Procedure: enqueue_orchestrator_work EXECUTE format('DROP FUNCTION IF EXISTS %I.enqueue_orchestrator_work(TEXT, TEXT, TIMESTAMPTZ, BIGINT, TEXT, TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.enqueue_orchestrator_work( p_instance_id TEXT, p_work_item TEXT, p_visible_at TIMESTAMPTZ, p_now_ms BIGINT, p_orchestration_name TEXT DEFAULT NULL, p_orchestration_version TEXT DEFAULT NULL, p_execution_id BIGINT DEFAULT NULL ) RETURNS VOID AS $enq_orch$ BEGIN -- Parameters p_orchestration_name, p_orchestration_version, p_execution_id are ignored -- Instance creation happens ONLY via ack_orchestration_item metadata INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (p_instance_id, p_work_item, p_visible_at, TO_TIMESTAMP(p_now_ms / 1000.0)); END; $enq_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: fetch_orchestration_item EXECUTE format('DROP FUNCTION IF EXISTS %I.fetch_orchestration_item(BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.fetch_orchestration_item( p_now_ms BIGINT, p_lock_timeout_ms BIGINT ) RETURNS TABLE( out_instance_id TEXT, out_orchestration_name TEXT, out_orchestration_version TEXT, out_execution_id BIGINT, out_history JSONB, out_messages JSONB, out_lock_token TEXT, out_attempt_count INTEGER ) AS $fetch_orch$ DECLARE v_instance_id TEXT; v_lock_token TEXT; v_locked_until BIGINT; v_orchestration_name TEXT; v_orchestration_version TEXT; v_current_execution_id BIGINT; v_history JSONB; v_messages JSONB; v_lock_acquired INTEGER; v_max_attempt_count INTEGER; BEGIN -- Phase 1: Find a candidate instance (no FOR UPDATE yet) SELECT q.instance_id INTO v_instance_id FROM %I.orchestrator_queue q WHERE q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND NOT EXISTS ( SELECT 1 FROM %I.instance_locks il WHERE il.instance_id = q.instance_id AND il.locked_until > p_now_ms ) ORDER BY q.visible_at, q.id LIMIT 1; IF NOT FOUND THEN RETURN; END IF; -- Phase 2: Acquire instance-level advisory lock PERFORM pg_advisory_xact_lock(hashtext(v_instance_id)); -- Phase 3: Re-verify with FOR UPDATE SELECT q.instance_id INTO v_instance_id FROM %I.orchestrator_queue q WHERE q.instance_id = v_instance_id AND q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND NOT EXISTS ( SELECT 1 FROM %I.instance_locks il WHERE il.instance_id = q.instance_id AND il.locked_until > p_now_ms ) FOR UPDATE OF q SKIP LOCKED; IF NOT FOUND THEN RETURN; END IF; v_lock_token := ''lock_'' || gen_random_uuid()::TEXT; v_locked_until := p_now_ms + p_lock_timeout_ms; INSERT INTO %I.instance_locks (instance_id, lock_token, locked_until, locked_at) VALUES (v_instance_id, v_lock_token, v_locked_until, p_now_ms) ON CONFLICT(instance_id) DO UPDATE SET lock_token = EXCLUDED.lock_token, locked_until = EXCLUDED.locked_until, locked_at = EXCLUDED.locked_at WHERE %I.instance_locks.locked_until <= p_now_ms; GET DIAGNOSTICS v_lock_acquired = ROW_COUNT; IF v_lock_acquired = 0 THEN RETURN; END IF; UPDATE %I.orchestrator_queue q SET lock_token = v_lock_token, locked_until = v_locked_until, attempt_count = q.attempt_count + 1 WHERE q.instance_id = v_instance_id AND q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (q.lock_token IS NULL OR q.locked_until <= p_now_ms); SELECT COALESCE(JSONB_AGG(q.work_item::JSONB ORDER BY q.id), ''[]''::JSONB), COALESCE(MAX(q.attempt_count), 1) INTO v_messages, v_max_attempt_count FROM %I.orchestrator_queue q WHERE q.lock_token = v_lock_token; SELECT i.orchestration_name, i.orchestration_version, i.current_execution_id INTO v_orchestration_name, v_orchestration_version, v_current_execution_id FROM %I.instances i WHERE i.instance_id = v_instance_id; IF FOUND THEN SELECT COALESCE(JSONB_AGG(h.event_data::JSONB ORDER BY h.event_id), ''[]''::JSONB) INTO v_history FROM %I.history h WHERE h.instance_id = v_instance_id AND h.execution_id = v_current_execution_id; v_orchestration_version := COALESCE(v_orchestration_version, ''unknown''); ELSE SELECT COALESCE(JSONB_AGG(h.event_data::JSONB ORDER BY h.execution_id, h.event_id), ''[]''::JSONB) INTO v_history FROM %I.history h WHERE h.instance_id = v_instance_id; IF JSONB_ARRAY_LENGTH(v_history) > 0 AND v_history->0 ? ''OrchestrationStarted'' THEN v_orchestration_name := v_history->0->''OrchestrationStarted''->>''name''; v_orchestration_version := v_history->0->''OrchestrationStarted''->>''version''; v_current_execution_id := 1; ELSIF JSONB_ARRAY_LENGTH(v_messages) > 0 AND v_messages->0 ? ''StartOrchestration'' THEN v_orchestration_name := v_messages->0->''StartOrchestration''->>''orchestration''; v_orchestration_version := COALESCE(v_messages->0->''StartOrchestration''->>''version'', ''unknown''); v_current_execution_id := COALESCE((v_messages->0->''StartOrchestration''->>''execution_id'')::BIGINT, 1); ELSIF JSONB_ARRAY_LENGTH(v_messages) > 0 AND v_messages->0 ? ''ContinueAsNew'' THEN v_orchestration_name := v_messages->0->''ContinueAsNew''->>''orchestration''; v_orchestration_version := COALESCE(v_messages->0->''ContinueAsNew''->>''version'', ''unknown''); v_current_execution_id := 1; ELSE v_orchestration_name := ''Unknown''; v_orchestration_version := ''unknown''; v_current_execution_id := 1; END IF; END IF; RETURN QUERY SELECT v_instance_id, v_orchestration_name, v_orchestration_version, v_current_execution_id, v_history, v_messages, v_lock_token, v_max_attempt_count; END; $fetch_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- Procedure: ack_orchestration_item EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_orchestration_item( p_lock_token TEXT, p_now_ms BIGINT, p_execution_id BIGINT, p_history_delta JSONB, p_worker_items JSONB, p_orchestrator_items JSONB, p_metadata JSONB, p_cancelled_activities JSONB DEFAULT ''[]''::JSONB ) RETURNS VOID AS $ack_orch$ DECLARE v_instance_id TEXT; v_now_ts TIMESTAMPTZ; v_orchestration_name TEXT; v_orchestration_version TEXT; v_status TEXT; v_output TEXT; v_completed_at TIMESTAMPTZ; v_elem JSONB; v_visible_at TIMESTAMPTZ; v_fire_at_ms BIGINT; v_item_instance_id TEXT; v_cancelled JSONB; v_cancelled_execution_id BIGINT; v_cancelled_activity_id BIGINT; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); SELECT il.instance_id INTO v_instance_id FROM %I.instance_locks il WHERE il.lock_token = p_lock_token AND il.locked_until > p_now_ms; IF NOT FOUND THEN RAISE EXCEPTION ''Invalid lock token''; END IF; v_orchestration_name := p_metadata->>''orchestration_name''; v_orchestration_version := p_metadata->>''orchestration_version''; v_status := p_metadata->>''status''; v_output := p_metadata->>''output''; IF v_orchestration_name IS NOT NULL AND v_orchestration_version IS NOT NULL THEN INSERT INTO %I.instances (instance_id, orchestration_name, orchestration_version, current_execution_id, created_at, updated_at) VALUES (v_instance_id, v_orchestration_name, v_orchestration_version, p_execution_id, v_now_ts, v_now_ts) ON CONFLICT (instance_id) DO NOTHING; UPDATE %I.instances i SET orchestration_name = v_orchestration_name, orchestration_version = v_orchestration_version, updated_at = v_now_ts WHERE i.instance_id = v_instance_id; END IF; INSERT INTO %I.executions (instance_id, execution_id, status, started_at) VALUES (v_instance_id, p_execution_id, ''Running'', v_now_ts) ON CONFLICT (instance_id, execution_id) DO NOTHING; UPDATE %I.instances i SET current_execution_id = GREATEST(i.current_execution_id, p_execution_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; IF p_history_delta IS NOT NULL AND JSONB_ARRAY_LENGTH(p_history_delta) > 0 THEN INSERT INTO %I.history (instance_id, execution_id, event_id, event_type, event_data, created_at) SELECT v_instance_id, p_execution_id, (elem->>''event_id'')::BIGINT, elem->>''event_type'', elem->>''event_data'', v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_history_delta) AS elem; END IF; IF v_status IS NOT NULL THEN v_completed_at := CASE WHEN v_status IN (''Completed'', ''Failed'') THEN v_now_ts ELSE NULL END; UPDATE %I.executions e SET status = v_status, output = v_output, completed_at = v_completed_at WHERE e.instance_id = v_instance_id AND e.execution_id = p_execution_id; END IF; IF p_worker_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_worker_items) > 0 THEN INSERT INTO %I.worker_queue (work_item, visible_at, created_at) SELECT elem::TEXT, v_now_ts, v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_worker_items) AS elem; END IF; IF p_orchestrator_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_orchestrator_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_orchestrator_items) LOOP IF v_elem ? ''StartOrchestration'' THEN v_item_instance_id := v_elem->''StartOrchestration''->>''instance''; ELSIF v_elem ? ''ContinueAsNew'' THEN v_item_instance_id := v_elem->''ContinueAsNew''->>''instance''; ELSIF v_elem ? ''TimerFired'' THEN v_item_instance_id := v_elem->''TimerFired''->>''instance''; v_fire_at_ms := (v_elem->''TimerFired''->>''fire_at_ms'')::BIGINT; ELSIF v_elem ? ''ActivityCompleted'' THEN v_item_instance_id := v_elem->''ActivityCompleted''->>''instance''; ELSIF v_elem ? ''ActivityFailed'' THEN v_item_instance_id := v_elem->''ActivityFailed''->>''instance''; ELSIF v_elem ? ''ExternalRaised'' THEN v_item_instance_id := v_elem->''ExternalRaised''->>''instance''; ELSIF v_elem ? ''CancelInstance'' THEN v_item_instance_id := v_elem->''CancelInstance''->>''instance''; ELSIF v_elem ? ''SubOrchCompleted'' THEN v_item_instance_id := v_elem->''SubOrchCompleted''->>''parent_instance''; ELSIF v_elem ? ''SubOrchFailed'' THEN v_item_instance_id := v_elem->''SubOrchFailed''->>''parent_instance''; ELSE v_item_instance_id := v_instance_id; END IF; IF v_elem ? ''TimerFired'' AND v_fire_at_ms IS NOT NULL AND v_fire_at_ms > 0 THEN v_visible_at := TO_TIMESTAMP(v_fire_at_ms / 1000.0); ELSE v_visible_at := v_now_ts; END IF; INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (v_item_instance_id, v_elem::TEXT, v_visible_at, v_now_ts); v_fire_at_ms := NULL; END LOOP; END IF; -- ================================================================ -- Lock-Stealing: Delete worker queue entries for cancelled activities -- Uses v_instance_id (from lock_token) for instance constraint -- Uses execution_id and activity_id from JSON to identify activities -- ================================================================ IF p_cancelled_activities IS NOT NULL AND JSONB_ARRAY_LENGTH(p_cancelled_activities) > 0 THEN FOR v_cancelled IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_cancelled_activities) LOOP v_cancelled_execution_id := (v_cancelled->>''execution_id'')::BIGINT; v_cancelled_activity_id := (v_cancelled->>''activity_id'')::BIGINT; -- Delete matching ActivityExecute items from worker_queue -- The work_item JSON contains ActivityExecute with instance, execution_id, and id (activity_id) DELETE FROM %I.worker_queue wq WHERE wq.work_item::JSONB ? ''ActivityExecute'' AND (wq.work_item::JSONB->''ActivityExecute''->>''instance'') = v_instance_id AND (wq.work_item::JSONB->''ActivityExecute''->>''execution_id'')::BIGINT = v_cancelled_execution_id AND (wq.work_item::JSONB->''ActivityExecute''->>''id'')::BIGINT = v_cancelled_activity_id; END LOOP; END IF; DELETE FROM %I.orchestrator_queue q WHERE q.lock_token = p_lock_token; DELETE FROM %I.instance_locks il WHERE il.instance_id = v_instance_id AND il.lock_token = p_lock_token; END; $ack_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- Procedure: abandon_orchestration_item EXECUTE format('DROP FUNCTION IF EXISTS %I.abandon_orchestration_item(TEXT, BIGINT, BIGINT, BOOLEAN)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.abandon_orchestration_item( p_lock_token TEXT, p_now_ms BIGINT, p_delay_ms BIGINT DEFAULT NULL, p_ignore_attempt BOOLEAN DEFAULT FALSE ) RETURNS TEXT AS $abandon_orch$ DECLARE v_instance_id TEXT; v_visible_at TIMESTAMPTZ; BEGIN SELECT il.instance_id INTO v_instance_id FROM %I.instance_locks il WHERE il.lock_token = p_lock_token; IF NOT FOUND THEN RAISE EXCEPTION ''Invalid lock token''; END IF; IF p_delay_ms IS NOT NULL AND p_delay_ms > 0 THEN v_visible_at := TO_TIMESTAMP((p_now_ms + p_delay_ms) / 1000.0); IF p_ignore_attempt THEN UPDATE %I.orchestrator_queue SET lock_token = NULL, locked_until = NULL, visible_at = v_visible_at, attempt_count = GREATEST(0, attempt_count - 1) WHERE lock_token = p_lock_token; ELSE UPDATE %I.orchestrator_queue SET lock_token = NULL, locked_until = NULL, visible_at = v_visible_at WHERE lock_token = p_lock_token; END IF; ELSE IF p_ignore_attempt THEN UPDATE %I.orchestrator_queue SET lock_token = NULL, locked_until = NULL, attempt_count = GREATEST(0, attempt_count - 1) WHERE lock_token = p_lock_token; ELSE UPDATE %I.orchestrator_queue SET lock_token = NULL, locked_until = NULL WHERE lock_token = p_lock_token; END IF; END IF; DELETE FROM %I.instance_locks WHERE lock_token = p_lock_token; RETURN v_instance_id; END; $abandon_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- Procedure: renew_orchestration_item_lock EXECUTE format('DROP FUNCTION IF EXISTS %I.renew_orchestration_item_lock(TEXT, BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.renew_orchestration_item_lock( p_lock_token TEXT, p_now_ms BIGINT, p_extend_ms BIGINT ) RETURNS VOID AS $renew_orch_lock$ DECLARE v_rows_affected INTEGER; BEGIN UPDATE %I.instance_locks SET locked_until = GREATEST(locked_until, p_now_ms) + p_extend_ms WHERE lock_token = p_lock_token AND locked_until > p_now_ms; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Lock token invalid, expired, or already released''; END IF; UPDATE %I.orchestrator_queue SET locked_until = GREATEST(locked_until, p_now_ms) + p_extend_ms WHERE lock_token = p_lock_token; END; $renew_orch_lock$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- History Procedures -- ============================================================================ -- Procedure: fetch_history EXECUTE format('DROP FUNCTION IF EXISTS %I.fetch_history(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.fetch_history( p_instance_id TEXT ) RETURNS TABLE(out_event_data TEXT) AS $fetch_history$ DECLARE v_execution_id BIGINT; BEGIN SELECT COALESCE(MAX(execution_id), 1) INTO v_execution_id FROM %I.executions WHERE instance_id = p_instance_id; RETURN QUERY SELECT h.event_data FROM %I.history h WHERE h.instance_id = p_instance_id AND h.execution_id = v_execution_id ORDER BY h.event_id; END; $fetch_history$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- Procedure: fetch_history_with_execution EXECUTE format('DROP FUNCTION IF EXISTS %I.fetch_history_with_execution(TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.fetch_history_with_execution( p_instance_id TEXT, p_execution_id BIGINT ) RETURNS TABLE(out_event_data TEXT) AS $fetch_history_exec$ BEGIN RETURN QUERY SELECT h.event_data FROM %I.history h WHERE h.instance_id = p_instance_id AND h.execution_id = p_execution_id ORDER BY h.event_id; END; $fetch_history_exec$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: append_history EXECUTE format('DROP FUNCTION IF EXISTS %I.append_history(TEXT, BIGINT, JSONB, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.append_history( p_instance_id TEXT, p_execution_id BIGINT, p_events JSONB, p_now_ms BIGINT ) RETURNS VOID AS $append_hist$ DECLARE v_now_ts TIMESTAMPTZ; BEGIN IF p_events IS NULL OR JSONB_ARRAY_LENGTH(p_events) = 0 THEN RETURN; END IF; v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); IF EXISTS ( SELECT 1 FROM JSONB_ARRAY_ELEMENTS(p_events) elem WHERE COALESCE((elem->>''event_id'')::BIGINT, 0) <= 0 ) THEN RAISE EXCEPTION ''Invalid event_id in append_history''; END IF; INSERT INTO %I.history (instance_id, execution_id, event_id, event_type, event_data, created_at) SELECT p_instance_id, p_execution_id, (elem->>''event_id'')::BIGINT, elem->>''event_type'', elem->>''event_data'', v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_events) AS elem ON CONFLICT (instance_id, execution_id, event_id) DO NOTHING; END; $append_hist$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); END $$; INSERT INTO _duroxide_migrations(version, name) VALUES (1, '0001_initial_schema.sql') ON CONFLICT (version) DO NOTHING; -- Migration 0002: 0002_add_deletion_and_pruning_support.sql -- Migration 0002: Add deletion and pruning support -- This migration adds: -- 1. parent_instance_id column to instances table (for cascade deletion) -- 2. Updates get_instance_info to include parent_instance_id -- 3. Updates ack_orchestration_item to store parent_instance_id from metadata -- 4. New stored procedures for ProviderAdmin methods: -- - list_children -- - get_parent_id -- - delete_instances_atomic -- - prune_executions -- Add parent_instance_id column to instances table ALTER TABLE instances ADD COLUMN IF NOT EXISTS parent_instance_id TEXT; -- Add index for efficient child lookups CREATE INDEX IF NOT EXISTS idx_instances_parent ON instances(parent_instance_id); -- Get the current schema name (set by migration runner) DO $$ DECLARE v_schema_name TEXT := current_schema(); BEGIN -- ============================================================================ -- Update get_instance_info to include parent_instance_id -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.get_instance_info(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_instance_info(p_instance_id TEXT) RETURNS TABLE( instance_id TEXT, orchestration_name TEXT, orchestration_version TEXT, current_execution_id BIGINT, created_at TIMESTAMPTZ, updated_at TIMESTAMPTZ, status TEXT, output TEXT, parent_instance_id TEXT ) AS $get_inst_info$ BEGIN RETURN QUERY SELECT i.instance_id, i.orchestration_name, COALESCE(i.orchestration_version, ''unknown'') as orchestration_version, i.current_execution_id, i.created_at, i.updated_at, e.status, e.output, i.parent_instance_id FROM %I.instances i LEFT JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE i.instance_id = p_instance_id; END; $get_inst_info$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Update ack_orchestration_item to store parent_instance_id from metadata -- This is critical for hierarchy tracking (cascade deletion, list_children, etc) -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_orchestration_item( p_lock_token TEXT, p_now_ms BIGINT, p_execution_id BIGINT, p_history_delta JSONB, p_worker_items JSONB, p_orchestrator_items JSONB, p_metadata JSONB, p_cancelled_activities JSONB DEFAULT ''[]''::JSONB ) RETURNS VOID AS $ack_orch$ DECLARE v_instance_id TEXT; v_now_ts TIMESTAMPTZ; v_orchestration_name TEXT; v_orchestration_version TEXT; v_parent_instance_id TEXT; v_status TEXT; v_output TEXT; v_completed_at TIMESTAMPTZ; v_elem JSONB; v_visible_at TIMESTAMPTZ; v_fire_at_ms BIGINT; v_item_instance_id TEXT; v_cancelled JSONB; v_cancelled_execution_id BIGINT; v_cancelled_activity_id BIGINT; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); SELECT il.instance_id INTO v_instance_id FROM %I.instance_locks il WHERE il.lock_token = p_lock_token AND il.locked_until > p_now_ms; IF NOT FOUND THEN RAISE EXCEPTION ''Invalid lock token''; END IF; v_orchestration_name := p_metadata->>''orchestration_name''; v_orchestration_version := p_metadata->>''orchestration_version''; v_parent_instance_id := p_metadata->>''parent_instance_id''; v_status := p_metadata->>''status''; v_output := p_metadata->>''output''; IF v_orchestration_name IS NOT NULL AND v_orchestration_version IS NOT NULL THEN INSERT INTO %I.instances (instance_id, orchestration_name, orchestration_version, current_execution_id, parent_instance_id, created_at, updated_at) VALUES (v_instance_id, v_orchestration_name, v_orchestration_version, p_execution_id, v_parent_instance_id, v_now_ts, v_now_ts) ON CONFLICT (instance_id) DO NOTHING; UPDATE %I.instances i SET orchestration_name = v_orchestration_name, orchestration_version = v_orchestration_version, parent_instance_id = COALESCE(i.parent_instance_id, v_parent_instance_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; END IF; INSERT INTO %I.executions (instance_id, execution_id, status, started_at) VALUES (v_instance_id, p_execution_id, ''Running'', v_now_ts) ON CONFLICT (instance_id, execution_id) DO NOTHING; UPDATE %I.instances i SET current_execution_id = GREATEST(i.current_execution_id, p_execution_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; IF p_history_delta IS NOT NULL AND JSONB_ARRAY_LENGTH(p_history_delta) > 0 THEN INSERT INTO %I.history (instance_id, execution_id, event_id, event_type, event_data, created_at) SELECT v_instance_id, p_execution_id, (elem->>''event_id'')::BIGINT, elem->>''event_type'', elem->>''event_data'', v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_history_delta) AS elem; END IF; IF v_status IS NOT NULL THEN v_completed_at := CASE WHEN v_status IN (''Completed'', ''Failed'') THEN v_now_ts ELSE NULL END; UPDATE %I.executions e SET status = v_status, output = v_output, completed_at = v_completed_at WHERE e.instance_id = v_instance_id AND e.execution_id = p_execution_id; END IF; IF p_worker_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_worker_items) > 0 THEN INSERT INTO %I.worker_queue (work_item, visible_at, created_at) SELECT elem::TEXT, v_now_ts, v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_worker_items) AS elem; END IF; IF p_orchestrator_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_orchestrator_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_orchestrator_items) LOOP IF v_elem ? ''StartOrchestration'' THEN v_item_instance_id := v_elem->''StartOrchestration''->>''instance''; ELSIF v_elem ? ''ContinueAsNew'' THEN v_item_instance_id := v_elem->''ContinueAsNew''->>''instance''; ELSIF v_elem ? ''TimerFired'' THEN v_item_instance_id := v_elem->''TimerFired''->>''instance''; v_fire_at_ms := (v_elem->''TimerFired''->>''fire_at_ms'')::BIGINT; ELSIF v_elem ? ''ActivityCompleted'' THEN v_item_instance_id := v_elem->''ActivityCompleted''->>''instance''; ELSIF v_elem ? ''ActivityFailed'' THEN v_item_instance_id := v_elem->''ActivityFailed''->>''instance''; ELSIF v_elem ? ''ExternalRaised'' THEN v_item_instance_id := v_elem->''ExternalRaised''->>''instance''; ELSIF v_elem ? ''CancelInstance'' THEN v_item_instance_id := v_elem->''CancelInstance''->>''instance''; ELSIF v_elem ? ''SubOrchCompleted'' THEN v_item_instance_id := v_elem->''SubOrchCompleted''->>''parent_instance''; ELSIF v_elem ? ''SubOrchFailed'' THEN v_item_instance_id := v_elem->''SubOrchFailed''->>''parent_instance''; ELSE v_item_instance_id := v_instance_id; END IF; IF v_elem ? ''TimerFired'' AND v_fire_at_ms IS NOT NULL AND v_fire_at_ms > 0 THEN v_visible_at := TO_TIMESTAMP(v_fire_at_ms / 1000.0); ELSE v_visible_at := v_now_ts; END IF; INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (v_item_instance_id, v_elem::TEXT, v_visible_at, v_now_ts); v_fire_at_ms := NULL; END LOOP; END IF; -- Lock-Stealing: Delete worker queue entries for cancelled activities IF p_cancelled_activities IS NOT NULL AND JSONB_ARRAY_LENGTH(p_cancelled_activities) > 0 THEN FOR v_cancelled IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_cancelled_activities) LOOP v_cancelled_execution_id := (v_cancelled->>''execution_id'')::BIGINT; v_cancelled_activity_id := (v_cancelled->>''activity_id'')::BIGINT; DELETE FROM %I.worker_queue wq WHERE wq.work_item::JSONB ? ''ActivityExecute'' AND (wq.work_item::JSONB->''ActivityExecute''->>''instance'') = v_instance_id AND (wq.work_item::JSONB->''ActivityExecute''->>''execution_id'')::BIGINT = v_cancelled_execution_id AND (wq.work_item::JSONB->''ActivityExecute''->>''id'')::BIGINT = v_cancelled_activity_id; END LOOP; END IF; DELETE FROM %I.orchestrator_queue q WHERE q.lock_token = p_lock_token; DELETE FROM %I.instance_locks il WHERE il.instance_id = v_instance_id AND il.lock_token = p_lock_token; END; $ack_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Hierarchy Primitive Procedures -- ============================================================================ -- Procedure: list_children EXECUTE format(' CREATE OR REPLACE FUNCTION %I.list_children(p_instance_id TEXT) RETURNS TABLE(child_instance_id TEXT) AS $list_children$ BEGIN RETURN QUERY SELECT i.instance_id FROM %I.instances i WHERE i.parent_instance_id = p_instance_id ORDER BY i.created_at; END; $list_children$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- Procedure: get_parent_id EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_parent_id(p_instance_id TEXT) RETURNS TEXT AS $get_parent$ DECLARE v_parent_id TEXT; BEGIN SELECT i.parent_instance_id INTO v_parent_id FROM %I.instances i WHERE i.instance_id = p_instance_id; IF NOT FOUND THEN RAISE EXCEPTION ''Instance not found: %%'', p_instance_id; END IF; RETURN v_parent_id; END; $get_parent$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- ============================================================================ -- Deletion Procedures -- ============================================================================ -- Procedure: delete_instances_atomic EXECUTE format(' CREATE OR REPLACE FUNCTION %I.delete_instances_atomic( p_instance_ids TEXT[], p_force BOOLEAN ) RETURNS TABLE( instances_deleted BIGINT, executions_deleted BIGINT, events_deleted BIGINT, queue_messages_deleted BIGINT ) AS $delete_atomic$ DECLARE v_instance_id TEXT; v_orphan_id TEXT; v_instances_deleted BIGINT := 0; v_executions_deleted BIGINT := 0; v_events_deleted BIGINT := 0; v_queue_deleted BIGINT := 0; v_count BIGINT; BEGIN IF p_instance_ids IS NULL OR array_length(p_instance_ids, 1) IS NULL THEN instances_deleted := 0; executions_deleted := 0; events_deleted := 0; queue_messages_deleted := 0; RETURN NEXT; RETURN; END IF; IF NOT p_force THEN SELECT i.instance_id INTO v_instance_id FROM %I.instances i JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE i.instance_id = ANY(p_instance_ids) AND e.status = ''Running'' LIMIT 1; IF v_instance_id IS NOT NULL THEN RAISE EXCEPTION ''Instance %% is Running. Use force=true to delete.'', v_instance_id; END IF; END IF; PERFORM 1 FROM %I.instances WHERE instance_id = ANY(p_instance_ids) FOR UPDATE; SELECT i.instance_id INTO v_orphan_id FROM %I.instances i WHERE i.parent_instance_id = ANY(p_instance_ids) AND NOT (i.instance_id = ANY(p_instance_ids)) LIMIT 1; IF v_orphan_id IS NOT NULL THEN RAISE EXCEPTION ''Orphan detected: instance %% has parent in delete list but is not included'', v_orphan_id; END IF; DELETE FROM %I.history WHERE instance_id = ANY(p_instance_ids); GET DIAGNOSTICS v_count = ROW_COUNT; v_events_deleted := v_count; DELETE FROM %I.executions WHERE instance_id = ANY(p_instance_ids); GET DIAGNOSTICS v_count = ROW_COUNT; v_executions_deleted := v_count; DELETE FROM %I.orchestrator_queue WHERE instance_id = ANY(p_instance_ids); GET DIAGNOSTICS v_count = ROW_COUNT; v_queue_deleted := v_count; DELETE FROM %I.worker_queue WHERE work_item::JSONB ? ''ActivityExecute'' AND (work_item::JSONB->''ActivityExecute''->>''instance'') = ANY(p_instance_ids); GET DIAGNOSTICS v_count = ROW_COUNT; v_queue_deleted := v_queue_deleted + v_count; DELETE FROM %I.instance_locks WHERE instance_id = ANY(p_instance_ids); DELETE FROM %I.instances WHERE instance_id = ANY(p_instance_ids); GET DIAGNOSTICS v_count = ROW_COUNT; v_instances_deleted := v_count; instances_deleted := v_instances_deleted; executions_deleted := v_executions_deleted; events_deleted := v_events_deleted; queue_messages_deleted := v_queue_deleted; RETURN NEXT; END; $delete_atomic$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Pruning Procedures -- ============================================================================ -- Procedure: prune_executions EXECUTE format(' CREATE OR REPLACE FUNCTION %I.prune_executions( p_instance_id TEXT, p_keep_last INTEGER DEFAULT NULL, p_completed_before_ms BIGINT DEFAULT NULL ) RETURNS TABLE( instances_processed BIGINT, executions_deleted BIGINT, events_deleted BIGINT ) AS $prune_exec$ DECLARE v_current_execution_id BIGINT; v_executions_deleted BIGINT := 0; v_events_deleted BIGINT := 0; v_count BIGINT; v_exec_ids_to_delete BIGINT[]; BEGIN SELECT i.current_execution_id INTO v_current_execution_id FROM %I.instances i WHERE i.instance_id = p_instance_id; IF NOT FOUND THEN RAISE EXCEPTION ''Instance %% not found'', p_instance_id; END IF; SELECT array_agg(e.execution_id) INTO v_exec_ids_to_delete FROM %I.executions e WHERE e.instance_id = p_instance_id AND e.execution_id != v_current_execution_id AND e.status != ''Running'' AND (p_completed_before_ms IS NULL OR e.completed_at < TO_TIMESTAMP(p_completed_before_ms / 1000.0)) AND (p_keep_last IS NULL OR e.execution_id NOT IN ( SELECT e2.execution_id FROM %I.executions e2 WHERE e2.instance_id = p_instance_id ORDER BY e2.execution_id DESC LIMIT p_keep_last )); IF v_exec_ids_to_delete IS NULL OR array_length(v_exec_ids_to_delete, 1) IS NULL THEN instances_processed := 1; executions_deleted := 0; events_deleted := 0; RETURN NEXT; RETURN; END IF; DELETE FROM %I.history h WHERE h.instance_id = p_instance_id AND h.execution_id = ANY(v_exec_ids_to_delete); GET DIAGNOSTICS v_count = ROW_COUNT; v_events_deleted := v_count; DELETE FROM %I.executions e WHERE e.instance_id = p_instance_id AND e.execution_id = ANY(v_exec_ids_to_delete); GET DIAGNOSTICS v_count = ROW_COUNT; v_executions_deleted := v_count; instances_processed := 1; executions_deleted := v_executions_deleted; events_deleted := v_events_deleted; RETURN NEXT; END; $prune_exec$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); END $$; INSERT INTO _duroxide_migrations(version, name) VALUES (2, '0002_add_deletion_and_pruning_support.sql') ON CONFLICT (version) DO NOTHING; -- Migration 0003: 0003_add_capability_filtering.sql -- Migration 0003: Add capability filtering support -- This migration adds: -- 1. duroxide_version_major/minor/patch columns to executions table -- 2. Updates fetch_orchestration_item to accept version filter parameters -- 3. Updates ack_orchestration_item to store pinned_duroxide_version from metadata -- 4. Updates cleanup_schema to drop new function signatures -- Add version columns to executions table ALTER TABLE executions ADD COLUMN IF NOT EXISTS duroxide_version_major INTEGER; ALTER TABLE executions ADD COLUMN IF NOT EXISTS duroxide_version_minor INTEGER; ALTER TABLE executions ADD COLUMN IF NOT EXISTS duroxide_version_patch INTEGER; -- Get the current schema name (set by migration runner) DO $$ DECLARE v_schema_name TEXT := current_schema(); BEGIN -- ============================================================================ -- Update fetch_orchestration_item to accept version filter parameters -- Drop old 2-param signature, create new 4-param signature -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.fetch_orchestration_item(BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.fetch_orchestration_item( p_now_ms BIGINT, p_lock_timeout_ms BIGINT, p_min_version_packed BIGINT DEFAULT NULL, p_max_version_packed BIGINT DEFAULT NULL ) RETURNS TABLE( out_instance_id TEXT, out_orchestration_name TEXT, out_orchestration_version TEXT, out_execution_id BIGINT, out_history JSONB, out_messages JSONB, out_lock_token TEXT, out_attempt_count INTEGER ) AS $fetch_orch$ DECLARE v_instance_id TEXT; v_lock_token TEXT; v_locked_until BIGINT; v_orchestration_name TEXT; v_orchestration_version TEXT; v_current_execution_id BIGINT; v_history JSONB; v_messages JSONB; v_lock_acquired INTEGER; v_max_attempt_count INTEGER; BEGIN -- Phase 1: Find a candidate instance (no FOR UPDATE yet) IF p_min_version_packed IS NOT NULL THEN -- Version-filtered path: join to instances and executions to check pinned version SELECT q.instance_id INTO v_instance_id FROM %I.orchestrator_queue q LEFT JOIN %I.instances i ON q.instance_id = i.instance_id LEFT JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND NOT EXISTS ( SELECT 1 FROM %I.instance_locks il WHERE il.instance_id = q.instance_id AND il.locked_until > p_now_ms ) AND ( e.duroxide_version_major IS NULL OR (e.duroxide_version_major * 1000000 + e.duroxide_version_minor * 1000 + e.duroxide_version_patch) BETWEEN p_min_version_packed AND p_max_version_packed ) ORDER BY q.visible_at, q.id LIMIT 1; ELSE -- No filter: original behavior SELECT q.instance_id INTO v_instance_id FROM %I.orchestrator_queue q WHERE q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND NOT EXISTS ( SELECT 1 FROM %I.instance_locks il WHERE il.instance_id = q.instance_id AND il.locked_until > p_now_ms ) ORDER BY q.visible_at, q.id LIMIT 1; END IF; IF NOT FOUND THEN RETURN; END IF; -- Phase 2: Acquire instance-level advisory lock PERFORM pg_advisory_xact_lock(hashtext(v_instance_id)); -- Phase 3: Re-verify with FOR UPDATE (include version filter if applicable) IF p_min_version_packed IS NOT NULL THEN SELECT q.instance_id INTO v_instance_id FROM %I.orchestrator_queue q LEFT JOIN %I.instances i ON q.instance_id = i.instance_id LEFT JOIN %I.executions e ON i.instance_id = e.instance_id AND i.current_execution_id = e.execution_id WHERE q.instance_id = v_instance_id AND q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND NOT EXISTS ( SELECT 1 FROM %I.instance_locks il WHERE il.instance_id = q.instance_id AND il.locked_until > p_now_ms ) AND ( e.duroxide_version_major IS NULL OR (e.duroxide_version_major * 1000000 + e.duroxide_version_minor * 1000 + e.duroxide_version_patch) BETWEEN p_min_version_packed AND p_max_version_packed ) FOR UPDATE OF q SKIP LOCKED; ELSE SELECT q.instance_id INTO v_instance_id FROM %I.orchestrator_queue q WHERE q.instance_id = v_instance_id AND q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND NOT EXISTS ( SELECT 1 FROM %I.instance_locks il WHERE il.instance_id = q.instance_id AND il.locked_until > p_now_ms ) FOR UPDATE OF q SKIP LOCKED; END IF; IF NOT FOUND THEN RETURN; END IF; v_lock_token := ''lock_'' || gen_random_uuid()::TEXT; v_locked_until := p_now_ms + p_lock_timeout_ms; INSERT INTO %I.instance_locks (instance_id, lock_token, locked_until, locked_at) VALUES (v_instance_id, v_lock_token, v_locked_until, p_now_ms) ON CONFLICT(instance_id) DO UPDATE SET lock_token = EXCLUDED.lock_token, locked_until = EXCLUDED.locked_until, locked_at = EXCLUDED.locked_at WHERE %I.instance_locks.locked_until <= p_now_ms; GET DIAGNOSTICS v_lock_acquired = ROW_COUNT; IF v_lock_acquired = 0 THEN RETURN; END IF; UPDATE %I.orchestrator_queue q SET lock_token = v_lock_token, locked_until = v_locked_until, attempt_count = q.attempt_count + 1 WHERE q.instance_id = v_instance_id AND q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (q.lock_token IS NULL OR q.locked_until <= p_now_ms); SELECT COALESCE(JSONB_AGG(q.work_item::JSONB ORDER BY q.id), ''[]''::JSONB), COALESCE(MAX(q.attempt_count), 1) INTO v_messages, v_max_attempt_count FROM %I.orchestrator_queue q WHERE q.lock_token = v_lock_token; SELECT i.orchestration_name, i.orchestration_version, i.current_execution_id INTO v_orchestration_name, v_orchestration_version, v_current_execution_id FROM %I.instances i WHERE i.instance_id = v_instance_id; IF FOUND THEN SELECT COALESCE(JSONB_AGG(h.event_data::JSONB ORDER BY h.event_id), ''[]''::JSONB) INTO v_history FROM %I.history h WHERE h.instance_id = v_instance_id AND h.execution_id = v_current_execution_id; v_orchestration_version := COALESCE(v_orchestration_version, ''unknown''); ELSE SELECT COALESCE(JSONB_AGG(h.event_data::JSONB ORDER BY h.execution_id, h.event_id), ''[]''::JSONB) INTO v_history FROM %I.history h WHERE h.instance_id = v_instance_id; IF JSONB_ARRAY_LENGTH(v_history) > 0 AND v_history->0 ? ''OrchestrationStarted'' THEN v_orchestration_name := v_history->0->''OrchestrationStarted''->>''name''; v_orchestration_version := v_history->0->''OrchestrationStarted''->>''version''; v_current_execution_id := 1; ELSIF JSONB_ARRAY_LENGTH(v_messages) > 0 AND v_messages->0 ? ''StartOrchestration'' THEN v_orchestration_name := v_messages->0->''StartOrchestration''->>''orchestration''; v_orchestration_version := COALESCE(v_messages->0->''StartOrchestration''->>''version'', ''unknown''); v_current_execution_id := COALESCE((v_messages->0->''StartOrchestration''->>''execution_id'')::BIGINT, 1); ELSIF JSONB_ARRAY_LENGTH(v_messages) > 0 AND v_messages->0 ? ''ContinueAsNew'' THEN v_orchestration_name := v_messages->0->''ContinueAsNew''->>''orchestration''; v_orchestration_version := COALESCE(v_messages->0->''ContinueAsNew''->>''version'', ''unknown''); v_current_execution_id := 1; ELSE v_orchestration_name := ''Unknown''; v_orchestration_version := ''unknown''; v_current_execution_id := 1; END IF; END IF; RETURN QUERY SELECT v_instance_id, v_orchestration_name, v_orchestration_version, v_current_execution_id, v_history, v_messages, v_lock_token, v_max_attempt_count; END; $fetch_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Update ack_orchestration_item to handle pinned_duroxide_version -- Drop old signature, create new one with version handling -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_orchestration_item( p_lock_token TEXT, p_now_ms BIGINT, p_execution_id BIGINT, p_history_delta JSONB, p_worker_items JSONB, p_orchestrator_items JSONB, p_metadata JSONB, p_cancelled_activities JSONB DEFAULT ''[]''::JSONB ) RETURNS VOID AS $ack_orch$ DECLARE v_instance_id TEXT; v_now_ts TIMESTAMPTZ; v_orchestration_name TEXT; v_orchestration_version TEXT; v_parent_instance_id TEXT; v_status TEXT; v_output TEXT; v_completed_at TIMESTAMPTZ; v_elem JSONB; v_visible_at TIMESTAMPTZ; v_fire_at_ms BIGINT; v_item_instance_id TEXT; v_cancelled JSONB; v_cancelled_execution_id BIGINT; v_cancelled_activity_id BIGINT; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); SELECT il.instance_id INTO v_instance_id FROM %I.instance_locks il WHERE il.lock_token = p_lock_token AND il.locked_until > p_now_ms; IF NOT FOUND THEN RAISE EXCEPTION ''Invalid lock token''; END IF; v_orchestration_name := p_metadata->>''orchestration_name''; v_orchestration_version := p_metadata->>''orchestration_version''; v_parent_instance_id := p_metadata->>''parent_instance_id''; v_status := p_metadata->>''status''; v_output := p_metadata->>''output''; IF v_orchestration_name IS NOT NULL AND v_orchestration_version IS NOT NULL THEN INSERT INTO %I.instances (instance_id, orchestration_name, orchestration_version, current_execution_id, parent_instance_id, created_at, updated_at) VALUES (v_instance_id, v_orchestration_name, v_orchestration_version, p_execution_id, v_parent_instance_id, v_now_ts, v_now_ts) ON CONFLICT (instance_id) DO NOTHING; UPDATE %I.instances i SET orchestration_name = v_orchestration_name, orchestration_version = v_orchestration_version, parent_instance_id = COALESCE(i.parent_instance_id, v_parent_instance_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; END IF; INSERT INTO %I.executions (instance_id, execution_id, status, started_at) VALUES (v_instance_id, p_execution_id, ''Running'', v_now_ts) ON CONFLICT (instance_id, execution_id) DO NOTHING; UPDATE %I.instances i SET current_execution_id = GREATEST(i.current_execution_id, p_execution_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; IF p_history_delta IS NOT NULL AND JSONB_ARRAY_LENGTH(p_history_delta) > 0 THEN INSERT INTO %I.history (instance_id, execution_id, event_id, event_type, event_data, created_at) SELECT v_instance_id, p_execution_id, (elem->>''event_id'')::BIGINT, elem->>''event_type'', elem->>''event_data'', v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_history_delta) AS elem; END IF; IF v_status IS NOT NULL THEN v_completed_at := CASE WHEN v_status IN (''Completed'', ''Failed'') THEN v_now_ts ELSE NULL END; UPDATE %I.executions e SET status = v_status, output = v_output, completed_at = v_completed_at WHERE e.instance_id = v_instance_id AND e.execution_id = p_execution_id; END IF; -- Store pinned duroxide version on execution if provided IF p_metadata ? ''pinned_duroxide_version'' AND p_metadata->''pinned_duroxide_version'' IS NOT NULL AND p_metadata->>''pinned_duroxide_version'' != ''null'' THEN UPDATE %I.executions SET duroxide_version_major = (p_metadata->''pinned_duroxide_version''->>''major'')::INTEGER, duroxide_version_minor = (p_metadata->''pinned_duroxide_version''->>''minor'')::INTEGER, duroxide_version_patch = (p_metadata->''pinned_duroxide_version''->>''patch'')::INTEGER WHERE instance_id = v_instance_id AND execution_id = p_execution_id; END IF; IF p_worker_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_worker_items) > 0 THEN INSERT INTO %I.worker_queue (work_item, visible_at, created_at) SELECT elem::TEXT, v_now_ts, v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_worker_items) AS elem; END IF; IF p_orchestrator_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_orchestrator_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_orchestrator_items) LOOP IF v_elem ? ''StartOrchestration'' THEN v_item_instance_id := v_elem->''StartOrchestration''->>''instance''; ELSIF v_elem ? ''ContinueAsNew'' THEN v_item_instance_id := v_elem->''ContinueAsNew''->>''instance''; ELSIF v_elem ? ''TimerFired'' THEN v_item_instance_id := v_elem->''TimerFired''->>''instance''; v_fire_at_ms := (v_elem->''TimerFired''->>''fire_at_ms'')::BIGINT; ELSIF v_elem ? ''ActivityCompleted'' THEN v_item_instance_id := v_elem->''ActivityCompleted''->>''instance''; ELSIF v_elem ? ''ActivityFailed'' THEN v_item_instance_id := v_elem->''ActivityFailed''->>''instance''; ELSIF v_elem ? ''ExternalRaised'' THEN v_item_instance_id := v_elem->''ExternalRaised''->>''instance''; ELSIF v_elem ? ''CancelInstance'' THEN v_item_instance_id := v_elem->''CancelInstance''->>''instance''; ELSIF v_elem ? ''SubOrchCompleted'' THEN v_item_instance_id := v_elem->''SubOrchCompleted''->>''parent_instance''; ELSIF v_elem ? ''SubOrchFailed'' THEN v_item_instance_id := v_elem->''SubOrchFailed''->>''parent_instance''; ELSE v_item_instance_id := v_instance_id; END IF; IF v_elem ? ''TimerFired'' AND v_fire_at_ms IS NOT NULL AND v_fire_at_ms > 0 THEN v_visible_at := TO_TIMESTAMP(v_fire_at_ms / 1000.0); ELSE v_visible_at := v_now_ts; END IF; INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (v_item_instance_id, v_elem::TEXT, v_visible_at, v_now_ts); v_fire_at_ms := NULL; END LOOP; END IF; -- Lock-Stealing: Delete worker queue entries for cancelled activities IF p_cancelled_activities IS NOT NULL AND JSONB_ARRAY_LENGTH(p_cancelled_activities) > 0 THEN FOR v_cancelled IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_cancelled_activities) LOOP v_cancelled_execution_id := (v_cancelled->>''execution_id'')::BIGINT; v_cancelled_activity_id := (v_cancelled->>''activity_id'')::BIGINT; DELETE FROM %I.worker_queue wq WHERE wq.work_item::JSONB ? ''ActivityExecute'' AND (wq.work_item::JSONB->''ActivityExecute''->>''instance'') = v_instance_id AND (wq.work_item::JSONB->''ActivityExecute''->>''execution_id'')::BIGINT = v_cancelled_execution_id AND (wq.work_item::JSONB->''ActivityExecute''->>''id'')::BIGINT = v_cancelled_activity_id; END LOOP; END IF; DELETE FROM %I.orchestrator_queue q WHERE q.lock_token = p_lock_token; DELETE FROM %I.instance_locks il WHERE il.instance_id = v_instance_id AND il.lock_token = p_lock_token; END; $ack_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Update cleanup_schema to drop new function signatures -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.cleanup_schema()', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.cleanup_schema() RETURNS VOID AS $cleanup$ BEGIN -- Drop tables first DROP TABLE IF EXISTS %I.instances CASCADE; DROP TABLE IF EXISTS %I.executions CASCADE; DROP TABLE IF EXISTS %I.history CASCADE; DROP TABLE IF EXISTS %I.orchestrator_queue CASCADE; DROP TABLE IF EXISTS %I.worker_queue CASCADE; DROP TABLE IF EXISTS %I.instance_locks CASCADE; DROP TABLE IF EXISTS %I._duroxide_migrations CASCADE; -- Drop all stored procedures (required because return type changes cannot use CREATE OR REPLACE) DROP FUNCTION IF EXISTS %I.cleanup_schema(); DROP FUNCTION IF EXISTS %I.list_instances(); DROP FUNCTION IF EXISTS %I.list_executions(TEXT); DROP FUNCTION IF EXISTS %I.latest_execution_id(TEXT); DROP FUNCTION IF EXISTS %I.list_instances_by_status(TEXT); DROP FUNCTION IF EXISTS %I.get_instance_info(TEXT); DROP FUNCTION IF EXISTS %I.get_execution_info(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.get_system_metrics(); DROP FUNCTION IF EXISTS %I.get_queue_depths(BIGINT); DROP FUNCTION IF EXISTS %I.enqueue_worker_work(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.renew_work_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_work_item(BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.abandon_work_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.enqueue_orchestrator_work(TEXT, TEXT, TIMESTAMPTZ, BIGINT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_orchestration_item(BIGINT, BIGINT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB); DROP FUNCTION IF EXISTS %I.abandon_orchestration_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.renew_orchestration_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_history(TEXT); DROP FUNCTION IF EXISTS %I.fetch_history_with_execution(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.append_history(TEXT, BIGINT, JSONB, BIGINT); DROP FUNCTION IF EXISTS %I.list_children(TEXT); DROP FUNCTION IF EXISTS %I.get_parent_id(TEXT); DROP FUNCTION IF EXISTS %I.delete_instances_atomic(TEXT[], BOOLEAN); DROP FUNCTION IF EXISTS %I.prune_executions(TEXT, INTEGER, BIGINT); -- Drop trigger functions (not schema-qualified, they use search_path) -- CASCADE is required because triggers depend on these functions DROP FUNCTION IF EXISTS notify_orch_work() CASCADE; DROP FUNCTION IF EXISTS notify_worker_work() CASCADE; END; $cleanup$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); END $$; INSERT INTO _duroxide_migrations(version, name) VALUES (3, '0003_add_capability_filtering.sql') ON CONFLICT (version) DO NOTHING; -- Migration 0004: 0004_add_session_support.sql -- Migration 0004: Add session affinity support -- This migration adds: -- 1. session_id column to worker_queue -- 2. sessions table for tracking session ownership -- 3. Session-aware fetch_work_item with routing logic -- 4. Session piggybacking in ack_worker and renew_work_item_lock -- 5. renew_session_lock and cleanup_orphaned_sessions stored procedures -- 6. Updated enqueue_worker_work and ack_orchestration_item for session_id -- Schema changes using unqualified names (search_path set by migration runner) ALTER TABLE worker_queue ADD COLUMN IF NOT EXISTS session_id TEXT; CREATE INDEX IF NOT EXISTS idx_worker_queue_session_id ON worker_queue(session_id); CREATE TABLE IF NOT EXISTS sessions ( session_id TEXT PRIMARY KEY, worker_id TEXT NOT NULL, locked_until BIGINT NOT NULL, last_activity_at BIGINT NOT NULL ); -- Stored procedure changes using schema-qualified names DO $$ DECLARE v_schema_name TEXT := current_schema(); BEGIN -- ============================================================================ -- Part 1: Update enqueue_worker_work to accept session_id -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.enqueue_worker_work(TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.enqueue_worker_work( p_work_item TEXT, p_now_ms BIGINT, p_session_id TEXT DEFAULT NULL ) RETURNS VOID AS $enq_worker$ DECLARE v_now_ts TIMESTAMPTZ; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); INSERT INTO %I.worker_queue (work_item, visible_at, created_at, session_id) VALUES (p_work_item, v_now_ts, v_now_ts, p_session_id); END; $enq_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- ============================================================================ -- Part 2: Update fetch_work_item with session routing -- Drop old 2-param version, create new 4-param version returning 3 columns -- (out_execution_status removed - cancellation now via lock-stealing) -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.fetch_work_item(BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.fetch_work_item( p_now_ms BIGINT, p_lock_timeout_ms BIGINT, p_owner_id TEXT DEFAULT NULL, p_session_lock_timeout_ms BIGINT DEFAULT NULL ) RETURNS TABLE( out_work_item TEXT, out_lock_token TEXT, out_attempt_count INTEGER ) AS $fetch_worker$ DECLARE v_id BIGINT; v_session_id TEXT; v_session_locked_until BIGINT; BEGIN IF p_owner_id IS NOT NULL THEN -- Session-aware fetch: find eligible items considering session routing -- Eligible items are: -- 1. Non-session items (q.session_id IS NULL) -- 2. Items for sessions owned by this worker (s.worker_id = p_owner_id AND s.locked_until > p_now_ms) -- 3. Items for claimable sessions (no active session row, or expired lock) SELECT q.id, q.session_id INTO v_id, v_session_id FROM %I.worker_queue q LEFT JOIN %I.sessions s ON s.session_id = q.session_id AND s.locked_until > p_now_ms WHERE q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (q.lock_token IS NULL OR q.locked_until <= p_now_ms) AND ( q.session_id IS NULL OR s.worker_id = p_owner_id OR s.session_id IS NULL ) ORDER BY q.id LIMIT 1 FOR UPDATE OF q SKIP LOCKED; ELSE -- Non-session fetch: only non-session items SELECT q.id, q.session_id INTO v_id, v_session_id FROM %I.worker_queue q WHERE q.visible_at <= TO_TIMESTAMP(p_now_ms / 1000.0) AND (q.lock_token IS NULL OR q.locked_until <= p_now_ms) AND q.session_id IS NULL ORDER BY q.id LIMIT 1 FOR UPDATE OF q SKIP LOCKED; END IF; IF NOT FOUND THEN RETURN; END IF; out_lock_token := ''lock_'' || gen_random_uuid()::TEXT; -- Increment attempt_count and lock the item UPDATE %I.worker_queue SET lock_token = out_lock_token, locked_until = p_now_ms + p_lock_timeout_ms, attempt_count = attempt_count + 1 WHERE id = v_id; SELECT work_item, attempt_count INTO out_work_item, out_attempt_count FROM %I.worker_queue WHERE id = v_id; -- If session-bound, upsert the sessions row IF v_session_id IS NOT NULL AND p_owner_id IS NOT NULL THEN v_session_locked_until := p_now_ms + COALESCE(p_session_lock_timeout_ms, p_lock_timeout_ms); INSERT INTO %I.sessions (session_id, worker_id, locked_until, last_activity_at) VALUES (v_session_id, p_owner_id, v_session_locked_until, p_now_ms) ON CONFLICT (session_id) DO UPDATE SET worker_id = p_owner_id, locked_until = v_session_locked_until, last_activity_at = p_now_ms WHERE %I.sessions.locked_until <= p_now_ms OR %I.sessions.worker_id = p_owner_id; -- If upsert affected 0 rows, another worker owns this session. -- Roll back: clear lock so item can be retried. IF NOT FOUND THEN UPDATE %I.worker_queue SET lock_token = NULL, locked_until = NULL, attempt_count = attempt_count - 1 WHERE id = v_id; RETURN; END IF; END IF; RETURN NEXT; END; $fetch_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 3: Update ack_worker to piggyback session last_activity_at -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_worker( p_lock_token TEXT, p_instance_id TEXT, p_completion_json TEXT, p_now_ms BIGINT ) RETURNS VOID AS $ack_worker$ DECLARE v_rows_affected INTEGER; v_now_ts TIMESTAMPTZ; v_session_id TEXT; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); -- Capture session_id before deleting SELECT session_id INTO v_session_id FROM %I.worker_queue WHERE lock_token = p_lock_token; -- Delete the worker queue item DELETE FROM %I.worker_queue WHERE lock_token = p_lock_token; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Worker queue item not found or already processed''; END IF; -- Validate: if completion provided, instance_id must also be provided IF p_completion_json IS NOT NULL AND p_instance_id IS NULL THEN RAISE EXCEPTION ''instance_id required when completion_json is provided''; END IF; -- Only enqueue completion if provided (not NULL) IF p_completion_json IS NOT NULL THEN INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (p_instance_id, p_completion_json, v_now_ts, v_now_ts); END IF; -- Piggyback: update session last_activity_at IF v_session_id IS NOT NULL AND p_now_ms IS NOT NULL THEN UPDATE %I.sessions SET last_activity_at = p_now_ms WHERE session_id = v_session_id AND locked_until > p_now_ms; END IF; END; $ack_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 4: Update renew_work_item_lock to piggyback session last_activity_at -- Also change return type from TEXT to VOID (execution_status removed in 0.1.8) -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.renew_work_item_lock(TEXT, BIGINT, BIGINT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.renew_work_item_lock( p_lock_token TEXT, p_now_ms BIGINT, p_extend_ms BIGINT ) RETURNS VOID AS $renew_lock$ DECLARE v_rows_affected INTEGER; v_session_id TEXT; BEGIN -- Read session_id before updating SELECT session_id INTO v_session_id FROM %I.worker_queue WHERE lock_token = p_lock_token; -- Update lock timeout only if lock is still valid UPDATE %I.worker_queue SET locked_until = GREATEST(locked_until, p_now_ms) + p_extend_ms WHERE lock_token = p_lock_token AND locked_until > p_now_ms; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Lock token invalid, expired, or already acked''; END IF; -- Piggyback: update session last_activity_at IF v_session_id IS NOT NULL THEN UPDATE %I.sessions SET last_activity_at = p_now_ms WHERE session_id = v_session_id AND locked_until > p_now_ms; END IF; END; $renew_lock$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 5: Create renew_session_lock stored procedure -- ============================================================================ EXECUTE format(' CREATE OR REPLACE FUNCTION %I.renew_session_lock( p_owner_ids TEXT[], p_now_ms BIGINT, p_extend_ms BIGINT, p_idle_timeout_ms BIGINT ) RETURNS BIGINT AS $renew_session$ DECLARE v_count BIGINT; BEGIN UPDATE %I.sessions SET locked_until = p_now_ms + p_extend_ms WHERE worker_id = ANY(p_owner_ids) AND locked_until > p_now_ms AND last_activity_at > (p_now_ms - p_idle_timeout_ms); GET DIAGNOSTICS v_count = ROW_COUNT; RETURN v_count; END; $renew_session$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- ============================================================================ -- Part 6: Create cleanup_orphaned_sessions stored procedure -- ============================================================================ EXECUTE format(' CREATE OR REPLACE FUNCTION %I.cleanup_orphaned_sessions( p_now_ms BIGINT ) RETURNS BIGINT AS $cleanup_sessions$ DECLARE v_count BIGINT; BEGIN DELETE FROM %I.sessions WHERE locked_until < p_now_ms AND NOT EXISTS (SELECT 1 FROM %I.worker_queue WHERE worker_queue.session_id = sessions.session_id); GET DIAGNOSTICS v_count = ROW_COUNT; RETURN v_count; END; $cleanup_sessions$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 7: Update ack_orchestration_item to extract session_id from worker items -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_orchestration_item( p_lock_token TEXT, p_now_ms BIGINT, p_execution_id BIGINT, p_history_delta JSONB, p_worker_items JSONB, p_orchestrator_items JSONB, p_metadata JSONB, p_cancelled_activities JSONB DEFAULT ''[]''::JSONB ) RETURNS VOID AS $ack_orch$ DECLARE v_instance_id TEXT; v_now_ts TIMESTAMPTZ; v_orchestration_name TEXT; v_orchestration_version TEXT; v_parent_instance_id TEXT; v_status TEXT; v_output TEXT; v_completed_at TIMESTAMPTZ; v_elem JSONB; v_visible_at TIMESTAMPTZ; v_fire_at_ms BIGINT; v_item_instance_id TEXT; v_item_session_id TEXT; v_cancelled JSONB; v_cancelled_execution_id BIGINT; v_cancelled_activity_id BIGINT; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); SELECT il.instance_id INTO v_instance_id FROM %I.instance_locks il WHERE il.lock_token = p_lock_token AND il.locked_until > p_now_ms; IF NOT FOUND THEN RAISE EXCEPTION ''Invalid lock token''; END IF; v_orchestration_name := p_metadata->>''orchestration_name''; v_orchestration_version := p_metadata->>''orchestration_version''; v_parent_instance_id := p_metadata->>''parent_instance_id''; v_status := p_metadata->>''status''; v_output := p_metadata->>''output''; IF v_orchestration_name IS NOT NULL AND v_orchestration_version IS NOT NULL THEN INSERT INTO %I.instances (instance_id, orchestration_name, orchestration_version, current_execution_id, parent_instance_id, created_at, updated_at) VALUES (v_instance_id, v_orchestration_name, v_orchestration_version, p_execution_id, v_parent_instance_id, v_now_ts, v_now_ts) ON CONFLICT (instance_id) DO NOTHING; UPDATE %I.instances i SET orchestration_name = v_orchestration_name, orchestration_version = v_orchestration_version, parent_instance_id = COALESCE(i.parent_instance_id, v_parent_instance_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; END IF; INSERT INTO %I.executions (instance_id, execution_id, status, started_at) VALUES (v_instance_id, p_execution_id, ''Running'', v_now_ts) ON CONFLICT (instance_id, execution_id) DO NOTHING; UPDATE %I.instances i SET current_execution_id = GREATEST(i.current_execution_id, p_execution_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; IF p_history_delta IS NOT NULL AND JSONB_ARRAY_LENGTH(p_history_delta) > 0 THEN INSERT INTO %I.history (instance_id, execution_id, event_id, event_type, event_data, created_at) SELECT v_instance_id, p_execution_id, (elem->>''event_id'')::BIGINT, elem->>''event_type'', elem->>''event_data'', v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_history_delta) AS elem; END IF; IF v_status IS NOT NULL THEN v_completed_at := CASE WHEN v_status IN (''Completed'', ''Failed'') THEN v_now_ts ELSE NULL END; UPDATE %I.executions e SET status = v_status, output = v_output, completed_at = v_completed_at WHERE e.instance_id = v_instance_id AND e.execution_id = p_execution_id; END IF; -- Store pinned duroxide version on execution if provided IF p_metadata ? ''pinned_duroxide_version'' AND p_metadata->''pinned_duroxide_version'' IS NOT NULL AND p_metadata->>''pinned_duroxide_version'' != ''null'' THEN UPDATE %I.executions SET duroxide_version_major = (p_metadata->''pinned_duroxide_version''->>''major'')::INTEGER, duroxide_version_minor = (p_metadata->''pinned_duroxide_version''->>''minor'')::INTEGER, duroxide_version_patch = (p_metadata->''pinned_duroxide_version''->>''patch'')::INTEGER WHERE instance_id = v_instance_id AND execution_id = p_execution_id; END IF; -- Enqueue worker items with session_id support IF p_worker_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_worker_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_worker_items) LOOP IF v_elem ? ''ActivityExecute'' THEN v_item_session_id := v_elem->''ActivityExecute''->>''session_id''; ELSE v_item_session_id := NULL; END IF; INSERT INTO %I.worker_queue (work_item, visible_at, created_at, session_id) VALUES (v_elem::TEXT, v_now_ts, v_now_ts, v_item_session_id); END LOOP; END IF; IF p_orchestrator_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_orchestrator_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_orchestrator_items) LOOP IF v_elem ? ''StartOrchestration'' THEN v_item_instance_id := v_elem->''StartOrchestration''->>''instance''; ELSIF v_elem ? ''ContinueAsNew'' THEN v_item_instance_id := v_elem->''ContinueAsNew''->>''instance''; ELSIF v_elem ? ''TimerFired'' THEN v_item_instance_id := v_elem->''TimerFired''->>''instance''; v_fire_at_ms := (v_elem->''TimerFired''->>''fire_at_ms'')::BIGINT; ELSIF v_elem ? ''ActivityCompleted'' THEN v_item_instance_id := v_elem->''ActivityCompleted''->>''instance''; ELSIF v_elem ? ''ActivityFailed'' THEN v_item_instance_id := v_elem->''ActivityFailed''->>''instance''; ELSIF v_elem ? ''ExternalRaised'' THEN v_item_instance_id := v_elem->''ExternalRaised''->>''instance''; ELSIF v_elem ? ''CancelInstance'' THEN v_item_instance_id := v_elem->''CancelInstance''->>''instance''; ELSIF v_elem ? ''SubOrchCompleted'' THEN v_item_instance_id := v_elem->''SubOrchCompleted''->>''parent_instance''; ELSIF v_elem ? ''SubOrchFailed'' THEN v_item_instance_id := v_elem->''SubOrchFailed''->>''parent_instance''; ELSE v_item_instance_id := v_instance_id; END IF; IF v_elem ? ''TimerFired'' AND v_fire_at_ms IS NOT NULL AND v_fire_at_ms > 0 THEN v_visible_at := TO_TIMESTAMP(v_fire_at_ms / 1000.0); ELSE v_visible_at := v_now_ts; END IF; INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (v_item_instance_id, v_elem::TEXT, v_visible_at, v_now_ts); v_fire_at_ms := NULL; END LOOP; END IF; -- Lock-Stealing: Delete worker queue entries for cancelled activities IF p_cancelled_activities IS NOT NULL AND JSONB_ARRAY_LENGTH(p_cancelled_activities) > 0 THEN FOR v_cancelled IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_cancelled_activities) LOOP v_cancelled_execution_id := (v_cancelled->>''execution_id'')::BIGINT; v_cancelled_activity_id := (v_cancelled->>''activity_id'')::BIGINT; DELETE FROM %I.worker_queue wq WHERE wq.work_item::JSONB ? ''ActivityExecute'' AND (wq.work_item::JSONB->''ActivityExecute''->>''instance'') = v_instance_id AND (wq.work_item::JSONB->''ActivityExecute''->>''execution_id'')::BIGINT = v_cancelled_execution_id AND (wq.work_item::JSONB->''ActivityExecute''->>''id'')::BIGINT = v_cancelled_activity_id; END LOOP; END IF; DELETE FROM %I.orchestrator_queue q WHERE q.lock_token = p_lock_token; DELETE FROM %I.instance_locks il WHERE il.instance_id = v_instance_id AND il.lock_token = p_lock_token; END; $ack_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 8: Update cleanup_schema to drop new functions and sessions table -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.cleanup_schema()', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.cleanup_schema() RETURNS VOID AS $cleanup$ BEGIN -- Drop tables first DROP TABLE IF EXISTS %I.sessions CASCADE; DROP TABLE IF EXISTS %I.instances CASCADE; DROP TABLE IF EXISTS %I.executions CASCADE; DROP TABLE IF EXISTS %I.history CASCADE; DROP TABLE IF EXISTS %I.orchestrator_queue CASCADE; DROP TABLE IF EXISTS %I.worker_queue CASCADE; DROP TABLE IF EXISTS %I.instance_locks CASCADE; DROP TABLE IF EXISTS %I._duroxide_migrations CASCADE; -- Drop all stored procedures DROP FUNCTION IF EXISTS %I.cleanup_schema(); DROP FUNCTION IF EXISTS %I.list_instances(); DROP FUNCTION IF EXISTS %I.list_executions(TEXT); DROP FUNCTION IF EXISTS %I.latest_execution_id(TEXT); DROP FUNCTION IF EXISTS %I.list_instances_by_status(TEXT); DROP FUNCTION IF EXISTS %I.get_instance_info(TEXT); DROP FUNCTION IF EXISTS %I.get_execution_info(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.get_system_metrics(); DROP FUNCTION IF EXISTS %I.get_queue_depths(BIGINT); DROP FUNCTION IF EXISTS %I.enqueue_worker_work(TEXT, BIGINT, TEXT); DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.renew_work_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_work_item(BIGINT, BIGINT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.abandon_work_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.enqueue_orchestrator_work(TEXT, TEXT, TIMESTAMPTZ, BIGINT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_orchestration_item(BIGINT, BIGINT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB); DROP FUNCTION IF EXISTS %I.abandon_orchestration_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.renew_orchestration_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_history(TEXT); DROP FUNCTION IF EXISTS %I.fetch_history_with_execution(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.append_history(TEXT, BIGINT, JSONB, BIGINT); DROP FUNCTION IF EXISTS %I.list_children(TEXT); DROP FUNCTION IF EXISTS %I.get_parent_id(TEXT); DROP FUNCTION IF EXISTS %I.delete_instances_atomic(TEXT[], BOOLEAN); DROP FUNCTION IF EXISTS %I.prune_executions(TEXT, INTEGER, BIGINT); DROP FUNCTION IF EXISTS %I.renew_session_lock(TEXT[], BIGINT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.cleanup_orphaned_sessions(BIGINT); -- Drop trigger functions (not schema-qualified, they use search_path) DROP FUNCTION IF EXISTS notify_orch_work() CASCADE; DROP FUNCTION IF EXISTS notify_worker_work() CASCADE; END; $cleanup$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); RAISE NOTICE 'Migration 0004: Added session support (sessions table, session routing in fetch, session piggybacking in ack/renew)'; END $$; INSERT INTO _duroxide_migrations(version, name) VALUES (4, '0004_add_session_support.sql') ON CONFLICT (version) DO NOTHING; -- Migration 0005: 0005_add_custom_status.sql -- Migration 0005: Add custom status support -- Description: Adds custom status support for orchestration instances. -- Adds custom_status and custom_status_version columns to instances table, -- updates ack_orchestration_item to handle custom status updates, -- and adds get_custom_status stored procedure for polling. -- Schema changes using unqualified names (search_path set by migration runner) ALTER TABLE instances ADD COLUMN IF NOT EXISTS custom_status TEXT; ALTER TABLE instances ADD COLUMN IF NOT EXISTS custom_status_version INTEGER NOT NULL DEFAULT 0; -- Stored procedure changes using schema-qualified names DO $$ DECLARE v_schema_name TEXT := current_schema(); BEGIN -- ============================================================================ -- Part 1: Update ack_orchestration_item to handle custom_status -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_orchestration_item( p_lock_token TEXT, p_now_ms BIGINT, p_execution_id BIGINT, p_history_delta JSONB, p_worker_items JSONB, p_orchestrator_items JSONB, p_metadata JSONB, p_cancelled_activities JSONB DEFAULT ''[]''::JSONB ) RETURNS VOID AS $ack_orch$ DECLARE v_instance_id TEXT; v_now_ts TIMESTAMPTZ; v_orchestration_name TEXT; v_orchestration_version TEXT; v_parent_instance_id TEXT; v_status TEXT; v_output TEXT; v_completed_at TIMESTAMPTZ; v_elem JSONB; v_visible_at TIMESTAMPTZ; v_fire_at_ms BIGINT; v_item_instance_id TEXT; v_item_session_id TEXT; v_cancelled JSONB; v_cancelled_execution_id BIGINT; v_cancelled_activity_id BIGINT; v_custom_status_action TEXT; v_custom_status_value TEXT; BEGIN v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); SELECT il.instance_id INTO v_instance_id FROM %I.instance_locks il WHERE il.lock_token = p_lock_token AND il.locked_until > p_now_ms; IF NOT FOUND THEN RAISE EXCEPTION ''Invalid lock token''; END IF; v_orchestration_name := p_metadata->>''orchestration_name''; v_orchestration_version := p_metadata->>''orchestration_version''; v_parent_instance_id := p_metadata->>''parent_instance_id''; v_status := p_metadata->>''status''; v_output := p_metadata->>''output''; IF v_orchestration_name IS NOT NULL AND v_orchestration_version IS NOT NULL THEN INSERT INTO %I.instances (instance_id, orchestration_name, orchestration_version, current_execution_id, parent_instance_id, created_at, updated_at) VALUES (v_instance_id, v_orchestration_name, v_orchestration_version, p_execution_id, v_parent_instance_id, v_now_ts, v_now_ts) ON CONFLICT (instance_id) DO NOTHING; UPDATE %I.instances i SET orchestration_name = v_orchestration_name, orchestration_version = v_orchestration_version, parent_instance_id = COALESCE(i.parent_instance_id, v_parent_instance_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; END IF; INSERT INTO %I.executions (instance_id, execution_id, status, started_at) VALUES (v_instance_id, p_execution_id, ''Running'', v_now_ts) ON CONFLICT (instance_id, execution_id) DO NOTHING; UPDATE %I.instances i SET current_execution_id = GREATEST(i.current_execution_id, p_execution_id), updated_at = v_now_ts WHERE i.instance_id = v_instance_id; IF p_history_delta IS NOT NULL AND JSONB_ARRAY_LENGTH(p_history_delta) > 0 THEN INSERT INTO %I.history (instance_id, execution_id, event_id, event_type, event_data, created_at) SELECT v_instance_id, p_execution_id, (elem->>''event_id'')::BIGINT, elem->>''event_type'', elem->>''event_data'', v_now_ts FROM JSONB_ARRAY_ELEMENTS(p_history_delta) AS elem; END IF; IF v_status IS NOT NULL THEN v_completed_at := CASE WHEN v_status IN (''Completed'', ''Failed'') THEN v_now_ts ELSE NULL END; UPDATE %I.executions e SET status = v_status, output = v_output, completed_at = v_completed_at WHERE e.instance_id = v_instance_id AND e.execution_id = p_execution_id; END IF; -- Store pinned duroxide version on execution if provided IF p_metadata ? ''pinned_duroxide_version'' AND p_metadata->''pinned_duroxide_version'' IS NOT NULL AND p_metadata->>''pinned_duroxide_version'' != ''null'' THEN UPDATE %I.executions SET duroxide_version_major = (p_metadata->''pinned_duroxide_version''->>''major'')::INTEGER, duroxide_version_minor = (p_metadata->''pinned_duroxide_version''->>''minor'')::INTEGER, duroxide_version_patch = (p_metadata->''pinned_duroxide_version''->>''patch'')::INTEGER WHERE instance_id = v_instance_id AND execution_id = p_execution_id; END IF; -- Handle custom_status update on instances table v_custom_status_action := p_metadata->>''custom_status_action''; IF v_custom_status_action = ''set'' THEN v_custom_status_value := p_metadata->>''custom_status_value''; UPDATE %I.instances SET custom_status = v_custom_status_value, custom_status_version = custom_status_version + 1 WHERE instance_id = v_instance_id; ELSIF v_custom_status_action = ''clear'' THEN UPDATE %I.instances SET custom_status = NULL, custom_status_version = custom_status_version + 1 WHERE instance_id = v_instance_id; END IF; -- Enqueue worker items with session_id support IF p_worker_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_worker_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_worker_items) LOOP IF v_elem ? ''ActivityExecute'' THEN v_item_session_id := v_elem->''ActivityExecute''->>''session_id''; ELSE v_item_session_id := NULL; END IF; INSERT INTO %I.worker_queue (work_item, visible_at, created_at, session_id) VALUES (v_elem::TEXT, v_now_ts, v_now_ts, v_item_session_id); END LOOP; END IF; IF p_orchestrator_items IS NOT NULL AND JSONB_ARRAY_LENGTH(p_orchestrator_items) > 0 THEN FOR v_elem IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_orchestrator_items) LOOP IF v_elem ? ''StartOrchestration'' THEN v_item_instance_id := v_elem->''StartOrchestration''->>''instance''; ELSIF v_elem ? ''ContinueAsNew'' THEN v_item_instance_id := v_elem->''ContinueAsNew''->>''instance''; ELSIF v_elem ? ''TimerFired'' THEN v_item_instance_id := v_elem->''TimerFired''->>''instance''; v_fire_at_ms := (v_elem->''TimerFired''->>''fire_at_ms'')::BIGINT; ELSIF v_elem ? ''ActivityCompleted'' THEN v_item_instance_id := v_elem->''ActivityCompleted''->>''instance''; ELSIF v_elem ? ''ActivityFailed'' THEN v_item_instance_id := v_elem->''ActivityFailed''->>''instance''; ELSIF v_elem ? ''ExternalRaised'' THEN v_item_instance_id := v_elem->''ExternalRaised''->>''instance''; ELSIF v_elem ? ''CancelInstance'' THEN v_item_instance_id := v_elem->''CancelInstance''->>''instance''; ELSIF v_elem ? ''SubOrchCompleted'' THEN v_item_instance_id := v_elem->''SubOrchCompleted''->>''parent_instance''; ELSIF v_elem ? ''SubOrchFailed'' THEN v_item_instance_id := v_elem->''SubOrchFailed''->>''parent_instance''; ELSIF v_elem ? ''QueueMessage'' THEN v_item_instance_id := v_elem->''QueueMessage''->>''instance''; ELSE v_item_instance_id := v_instance_id; END IF; IF v_elem ? ''TimerFired'' AND v_fire_at_ms IS NOT NULL AND v_fire_at_ms > 0 THEN v_visible_at := TO_TIMESTAMP(v_fire_at_ms / 1000.0); ELSE v_visible_at := v_now_ts; END IF; INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (v_item_instance_id, v_elem::TEXT, v_visible_at, v_now_ts); v_fire_at_ms := NULL; END LOOP; END IF; -- Lock-Stealing: Delete worker queue entries for cancelled activities IF p_cancelled_activities IS NOT NULL AND JSONB_ARRAY_LENGTH(p_cancelled_activities) > 0 THEN FOR v_cancelled IN SELECT value FROM JSONB_ARRAY_ELEMENTS(p_cancelled_activities) LOOP v_cancelled_execution_id := (v_cancelled->>''execution_id'')::BIGINT; v_cancelled_activity_id := (v_cancelled->>''activity_id'')::BIGINT; DELETE FROM %I.worker_queue wq WHERE wq.work_item::JSONB ? ''ActivityExecute'' AND (wq.work_item::JSONB->''ActivityExecute''->>''instance'') = v_instance_id AND (wq.work_item::JSONB->''ActivityExecute''->>''execution_id'')::BIGINT = v_cancelled_execution_id AND (wq.work_item::JSONB->''ActivityExecute''->>''id'')::BIGINT = v_cancelled_activity_id; END LOOP; END IF; DELETE FROM %I.orchestrator_queue q WHERE q.lock_token = p_lock_token; DELETE FROM %I.instance_locks il WHERE il.instance_id = v_instance_id AND il.lock_token = p_lock_token; END; $ack_orch$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 2: Update ack_worker to validate lock expiry -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT)', v_schema_name); EXECUTE format('DROP FUNCTION IF EXISTS %I.ack_worker(TEXT)', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.ack_worker( p_lock_token TEXT, p_instance_id TEXT DEFAULT NULL, p_completion_json TEXT DEFAULT NULL, p_now_ms BIGINT DEFAULT NULL ) RETURNS VOID AS $ack_worker$ DECLARE v_rows_affected INTEGER; v_now_ts TIMESTAMPTZ; v_session_id TEXT; BEGIN -- Capture session_id before deleting SELECT session_id INTO v_session_id FROM %I.worker_queue WHERE lock_token = p_lock_token; -- Delete the worker queue item, only if lock is still valid IF p_now_ms IS NOT NULL THEN DELETE FROM %I.worker_queue WHERE lock_token = p_lock_token AND locked_until > p_now_ms; ELSE DELETE FROM %I.worker_queue WHERE lock_token = p_lock_token; END IF; GET DIAGNOSTICS v_rows_affected = ROW_COUNT; IF v_rows_affected = 0 THEN RAISE EXCEPTION ''Worker queue item not found or already processed''; END IF; -- Only enqueue completion if provided (NULL means cancelled activity) IF p_completion_json IS NOT NULL THEN -- Validate required parameters for completion IF p_instance_id IS NULL THEN RAISE EXCEPTION ''p_instance_id is required when p_completion_json is provided''; END IF; IF p_now_ms IS NULL THEN RAISE EXCEPTION ''p_now_ms is required when p_completion_json is provided''; END IF; v_now_ts := TO_TIMESTAMP(p_now_ms / 1000.0); INSERT INTO %I.orchestrator_queue (instance_id, work_item, visible_at, created_at) VALUES (p_instance_id, p_completion_json, v_now_ts, v_now_ts); END IF; -- Piggyback: update session last_activity_at IF v_session_id IS NOT NULL AND p_now_ms IS NOT NULL THEN UPDATE %I.sessions SET last_activity_at = p_now_ms WHERE session_id = v_session_id AND locked_until > p_now_ms; END IF; END; $ack_worker$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); -- ============================================================================ -- Part 3: Create get_custom_status stored procedure -- ============================================================================ EXECUTE format(' CREATE OR REPLACE FUNCTION %I.get_custom_status( p_instance_id TEXT, p_last_seen_version BIGINT ) RETURNS TABLE( out_custom_status TEXT, out_custom_status_version BIGINT ) AS $get_cs$ BEGIN RETURN QUERY SELECT i.custom_status, i.custom_status_version::BIGINT FROM %I.instances i WHERE i.instance_id = p_instance_id AND i.custom_status_version > p_last_seen_version; END; $get_cs$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name); -- ============================================================================ -- Part 4: Update cleanup_schema to drop new function -- ============================================================================ EXECUTE format('DROP FUNCTION IF EXISTS %I.cleanup_schema()', v_schema_name); EXECUTE format(' CREATE OR REPLACE FUNCTION %I.cleanup_schema() RETURNS VOID AS $cleanup$ BEGIN -- Drop tables first DROP TABLE IF EXISTS %I.sessions CASCADE; DROP TABLE IF EXISTS %I.instances CASCADE; DROP TABLE IF EXISTS %I.executions CASCADE; DROP TABLE IF EXISTS %I.history CASCADE; DROP TABLE IF EXISTS %I.orchestrator_queue CASCADE; DROP TABLE IF EXISTS %I.worker_queue CASCADE; DROP TABLE IF EXISTS %I.instance_locks CASCADE; DROP TABLE IF EXISTS %I._duroxide_migrations CASCADE; -- Drop all stored procedures DROP FUNCTION IF EXISTS %I.cleanup_schema(); DROP FUNCTION IF EXISTS %I.list_instances(); DROP FUNCTION IF EXISTS %I.list_executions(TEXT); DROP FUNCTION IF EXISTS %I.latest_execution_id(TEXT); DROP FUNCTION IF EXISTS %I.list_instances_by_status(TEXT); DROP FUNCTION IF EXISTS %I.get_instance_info(TEXT); DROP FUNCTION IF EXISTS %I.get_execution_info(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.get_system_metrics(); DROP FUNCTION IF EXISTS %I.get_queue_depths(BIGINT); DROP FUNCTION IF EXISTS %I.enqueue_worker_work(TEXT, BIGINT, TEXT); DROP FUNCTION IF EXISTS %I.ack_worker(TEXT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.renew_work_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_work_item(BIGINT, BIGINT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.abandon_work_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.enqueue_orchestrator_work(TEXT, TEXT, TIMESTAMPTZ, BIGINT, TEXT, TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_orchestration_item(BIGINT, BIGINT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.ack_orchestration_item(TEXT, BIGINT, BIGINT, JSONB, JSONB, JSONB, JSONB, JSONB); DROP FUNCTION IF EXISTS %I.abandon_orchestration_item(TEXT, BIGINT, BIGINT, BOOLEAN); DROP FUNCTION IF EXISTS %I.renew_orchestration_item_lock(TEXT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.fetch_history(TEXT); DROP FUNCTION IF EXISTS %I.fetch_history_with_execution(TEXT, BIGINT); DROP FUNCTION IF EXISTS %I.append_history(TEXT, BIGINT, JSONB, BIGINT); DROP FUNCTION IF EXISTS %I.list_children(TEXT); DROP FUNCTION IF EXISTS %I.get_parent_id(TEXT); DROP FUNCTION IF EXISTS %I.delete_instances_atomic(TEXT[], BOOLEAN); DROP FUNCTION IF EXISTS %I.prune_executions(TEXT, INTEGER, BIGINT); DROP FUNCTION IF EXISTS %I.renew_session_lock(TEXT[], BIGINT, BIGINT, BIGINT); DROP FUNCTION IF EXISTS %I.cleanup_orphaned_sessions(BIGINT); DROP FUNCTION IF EXISTS %I.get_custom_status(TEXT, BIGINT); -- Drop trigger functions (not schema-qualified, they use search_path) DROP FUNCTION IF EXISTS notify_orch_work() CASCADE; DROP FUNCTION IF EXISTS notify_worker_work() CASCADE; END; $cleanup$ LANGUAGE plpgsql; ', v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name, v_schema_name); RAISE NOTICE 'Migration 0005: Added custom status support (custom_status/custom_status_version columns, get_custom_status function, updated ack_orchestration_item)'; END $$; INSERT INTO _duroxide_migrations(version, name) VALUES (5, '0005_add_custom_status.sql') ON CONFLICT (version) DO NOTHING; SET LOCAL search_path TO @extschema@; -- END duroxide-pg-opt migrations (checked-in copy) /* */ /* */ -- src/monitoring.rs:17 -- pg_durable::monitoring::list_instances CREATE FUNCTION df."list_instances"( "status_filter" TEXT DEFAULT NULL, /* core::option::Option<&str> */ "limit_count" INT DEFAULT 100 /* i32 */ ) RETURNS TABLE ( "instance_id" TEXT, /* alloc::string::String */ "label" TEXT, /* core::option::Option */ "function_name" TEXT, /* alloc::string::String */ "status" TEXT, /* alloc::string::String */ "execution_count" bigint, /* i64 */ "output" TEXT /* core::option::Option */ ) LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'list_instances_wrapper'; /* */ /* */ -- src/dsl.rs:757 -- pg_durable::dsl::wait_for_completion CREATE FUNCTION df."wait_for_completion"( "instance_id" TEXT, /* &str */ "timeout_seconds" INT DEFAULT 30 /* i32 */ ) RETURNS TEXT /* core::result::Result> */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'wait_for_completion_wrapper'; /* */ /* */ -- src/dsl.rs:74 -- pg_durable::dsl::getvar CREATE FUNCTION df."getvar"( "name" TEXT /* &str */ ) RETURNS TEXT /* core::option::Option */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'getvar_wrapper'; /* */ /* */ -- src/dsl.rs:151 -- pg_durable::dsl::as CREATE FUNCTION df."as"( "fut" TEXT, /* &str */ "name" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'as_named_wrapper'; /* */ /* */ -- src/dsl.rs:296 -- pg_durable::dsl::join CREATE FUNCTION df."join"( "a" TEXT, /* &str */ "b" TEXT /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'join_wrapper'; /* */ /* */ -- src/lib.rs:280 -- requires: -- dsl::then_fn -- dsl::as_named -- dsl::join -- dsl::race -- dsl::if_fn -- dsl::loop_fn -- Operator ~> for sequencing: a ~> b means "run a, then run b" CREATE OPERATOR ~> ( FUNCTION = df.seq, LEFTARG = text, RIGHTARG = text ); -- Operator |=> for naming: fut |=> 'name' means "name this result as $name" CREATE OR REPLACE FUNCTION df.as_op(fut text, name text) RETURNS text AS $$ SELECT df.as(fut, name); $$ LANGUAGE SQL IMMUTABLE; CREATE OPERATOR |=> ( FUNCTION = df.as_op, LEFTARG = text, RIGHTARG = text ); -- Operator & for parallel join: a & b means "run a and b in parallel, wait for both" CREATE OPERATOR & ( FUNCTION = df.join, LEFTARG = text, RIGHTARG = text ); -- Operator | for race: a | b means "run a and b in parallel, first wins" CREATE OPERATOR | ( FUNCTION = df.race, LEFTARG = text, RIGHTARG = text ); -- Operators ?> and !> for if-then-else: cond ?> then_branch !> else_branch -- We need helper functions to build the if node incrementally -- Helper: cond ?> then creates a partial if (stores condition and then branch) CREATE OR REPLACE FUNCTION df.if_then_op(condition text, then_branch text) RETURNS text AS $$ DECLARE cond_fut jsonb; then_fut jsonb; result_obj jsonb; BEGIN -- Ensure both are durofuts cond_fut := df.ensure_durofut(condition)::jsonb; then_fut := df.ensure_durofut(then_branch)::jsonb; -- Return a special marker object for the partial if result_obj := jsonb_build_object( '_partial_if', true, 'condition', cond_fut, 'then_branch', then_fut ); RETURN result_obj::text; END; $$ LANGUAGE plpgsql IMMUTABLE; -- Helper: partial_if !> else completes the if node CREATE OR REPLACE FUNCTION df.if_else_op(partial_if text, else_branch text) RETURNS text AS $$ DECLARE partial jsonb; else_fut text; cond_text text; then_text text; BEGIN partial := partial_if::jsonb; -- Check if it's a partial if IF partial->>'_partial_if' IS NULL THEN RAISE EXCEPTION 'Invalid if-then-else: left side of !> must be a ?> expression'; END IF; cond_text := partial->'condition'::text; then_text := partial->'then_branch'::text; else_fut := df.ensure_durofut(else_branch); -- Now call the real df.if function RETURN df.if(cond_text, then_text, else_fut); END; $$ LANGUAGE plpgsql IMMUTABLE; -- Helper to ensure a value is a durofut (returns JSON string) -- Rejects JSON with unknown node_type values. -- NOTE: The valid node type list here must be kept in sync with -- VALID_NODE_TYPES in src/types.rs (the Rust constant is the canonical source). CREATE OR REPLACE FUNCTION df.ensure_durofut(val text) RETURNS text AS $$ DECLARE node_type_val text; BEGIN -- Try to parse as JSON to check if it's already a durofut BEGIN node_type_val := (val::jsonb)->>'node_type'; IF node_type_val IS NOT NULL THEN -- Has a node_type - validate it IF node_type_val NOT IN ('SQL', 'THEN', 'IF', 'JOIN', 'LOOP', 'BREAK', 'RACE', 'SLEEP', 'WAIT_SCHEDULE', 'HTTP', 'SIGNAL') THEN RAISE EXCEPTION 'Unknown node_type ''%''. Valid types: SQL, THEN, IF, JOIN, LOOP, BREAK, RACE, SLEEP, WAIT_SCHEDULE, HTTP, SIGNAL', node_type_val; END IF; RETURN val; END IF; EXCEPTION WHEN invalid_text_representation THEN -- Not valid JSON, treat as SQL NULL; WHEN raise_exception THEN -- Re-raise our validation error RAISE; WHEN OTHERS THEN -- Not valid JSON, treat as SQL NULL; END; -- It's plain SQL, wrap it RETURN df.sql(val); END; $$ LANGUAGE plpgsql IMMUTABLE; CREATE OPERATOR ?> ( FUNCTION = df.if_then_op, LEFTARG = text, RIGHTARG = text ); CREATE OPERATOR !> ( FUNCTION = df.if_else_op, LEFTARG = text, RIGHTARG = text ); -- Operator @> for loop: @> body means "repeat body forever" -- This is a PREFIX operator with lowest precedence CREATE OR REPLACE FUNCTION df.loop_prefix_op(body text) RETURNS text AS $$ SELECT df.loop(body); $$ LANGUAGE SQL IMMUTABLE; CREATE OPERATOR @> ( FUNCTION = df.loop_prefix_op, RIGHTARG = text ); /* */ /* */ -- src/monitoring.rs:184 -- pg_durable::monitoring::instance_executions CREATE FUNCTION df."instance_executions"( "instance_id" TEXT, /* &str */ "limit_count" INT DEFAULT 5 /* i32 */ ) RETURNS TABLE ( "execution_id" bigint, /* i64 */ "status" TEXT, /* alloc::string::String */ "event_count" bigint, /* i64 */ "duration_ms" bigint, /* i64 */ "output" TEXT /* core::option::Option */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'instance_executions_wrapper'; /* */ /* */ -- src/dsl.rs:448 -- pg_durable::dsl::signal CREATE FUNCTION df."signal"( "instance_id" TEXT, /* &str */ "signal_name" TEXT, /* &str */ "signal_data" TEXT DEFAULT '{}' /* &str */ ) RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'signal_wrapper'; /* */ /* */ -- src/dsl.rs:38 -- pg_durable::dsl::debug_connection CREATE FUNCTION df."debug_connection"() RETURNS TEXT /* alloc::string::String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'debug_connection_wrapper'; /* */ /* */ -- src/monitoring.rs:109 -- pg_durable::monitoring::instance_info CREATE FUNCTION df."instance_info"( "instance_id" TEXT /* &str */ ) RETURNS TABLE ( "instance_id" TEXT, /* alloc::string::String */ "label" TEXT, /* core::option::Option */ "function_name" TEXT, /* alloc::string::String */ "function_version" TEXT, /* alloc::string::String */ "current_execution_id" bigint, /* i64 */ "status" TEXT, /* alloc::string::String */ "output" TEXT /* core::option::Option */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'instance_info_wrapper'; /* */ /* */ -- src/dsl.rs:714 -- pg_durable::dsl::status CREATE FUNCTION df."status"( "instance_id" TEXT /* &str */ ) RETURNS TEXT /* core::option::Option */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'status_wrapper'; /* */