/* */ /* This file is auto generated by pgrx. The ordering of items is not stable, it is driven by a dependency graph. */ /* */ /* */ -- src/lib.rs:249 -- bootstrap -- Extension schemas CREATE SCHEMA IF NOT EXISTS pgtrickle; CREATE SCHEMA IF NOT EXISTS pgtrickle_changes; -- F51: Restrict change buffer schema access to prevent unauthorized -- injection of bogus changes that would be applied on next refresh. REVOKE ALL ON SCHEMA pgtrickle_changes FROM PUBLIC; -- User-declared refresh groups for snapshot consistency CREATE TABLE IF NOT EXISTS pgtrickle.pgt_refresh_groups ( group_id SERIAL PRIMARY KEY, group_name TEXT NOT NULL UNIQUE, member_oids OID[] NOT NULL, isolation TEXT NOT NULL DEFAULT 'read_committed' CHECK (isolation IN ('read_committed', 'repeatable_read')), created_at TIMESTAMPTZ NOT NULL DEFAULT now() ); -- Core ST metadata CREATE TABLE IF NOT EXISTS pgtrickle.pgt_stream_tables ( pgt_id BIGSERIAL PRIMARY KEY, pgt_relid OID NOT NULL UNIQUE, pgt_name TEXT NOT NULL, pgt_schema TEXT NOT NULL, defining_query TEXT NOT NULL, original_query TEXT, schedule TEXT, refresh_mode TEXT NOT NULL DEFAULT 'DIFFERENTIAL' CHECK (refresh_mode IN ('FULL', 'DIFFERENTIAL', 'IMMEDIATE')), requested_refresh_mode TEXT NOT NULL DEFAULT 'DIFFERENTIAL' CHECK (requested_refresh_mode IN ('AUTO', 'FULL', 'DIFFERENTIAL', 'IMMEDIATE')), status TEXT NOT NULL DEFAULT 'INITIALIZING' CHECK (status IN ('INITIALIZING', 'ACTIVE', 'SUSPENDED', 'ERROR')), is_populated BOOLEAN NOT NULL DEFAULT FALSE, data_timestamp TIMESTAMPTZ, frontier JSONB, tentative_frontier JSONB, last_refresh_at TIMESTAMPTZ, consecutive_errors INT NOT NULL DEFAULT 0, needs_reinit BOOLEAN NOT NULL DEFAULT FALSE, auto_threshold DOUBLE PRECISION, last_full_ms DOUBLE PRECISION, functions_used TEXT[], topk_limit INT, topk_order_by TEXT, topk_offset INT, diamond_consistency TEXT NOT NULL DEFAULT 'atomic' CHECK (diamond_consistency IN ('none', 'atomic')), diamond_schedule_policy TEXT NOT NULL DEFAULT 'fastest' CHECK (diamond_schedule_policy IN ('fastest', 'slowest')), has_keyless_source BOOLEAN NOT NULL DEFAULT FALSE, function_hashes TEXT, requested_cdc_mode TEXT CHECK (requested_cdc_mode IN ('auto', 'trigger', 'wal')), is_append_only BOOLEAN NOT NULL DEFAULT FALSE, scc_id INT, last_fixpoint_iterations INT, max_differential_joins INT CHECK (max_differential_joins IS NULL OR max_differential_joins >= 0), max_delta_fraction DOUBLE PRECISION CHECK (max_delta_fraction IS NULL OR (max_delta_fraction >= 0.0 AND max_delta_fraction <= 1.0)), pooler_compatibility_mode BOOLEAN NOT NULL DEFAULT FALSE, refresh_tier TEXT NOT NULL DEFAULT 'hot' CHECK (refresh_tier IN ('hot', 'warm', 'cold', 'frozen')), effective_refresh_mode TEXT, fuse_mode TEXT NOT NULL DEFAULT 'off' CHECK (fuse_mode IN ('off', 'on', 'auto')), fuse_state TEXT NOT NULL DEFAULT 'armed' CHECK (fuse_state IN ('armed', 'blown', 'disabled')), fuse_ceiling BIGINT, fuse_sensitivity INT, blown_at TIMESTAMPTZ, blow_reason TEXT, last_error_message TEXT, last_error_at TIMESTAMPTZ, self_heal_work_mem_percent SMALLINT NOT NULL DEFAULT 100 CHECK (self_heal_work_mem_percent BETWEEN 25 AND 100), self_heal_lock_backoff_exponent SMALLINT NOT NULL DEFAULT 0 CHECK (self_heal_lock_backoff_exponent BETWEEN 0 AND 6), self_heal_success_streak SMALLINT NOT NULL DEFAULT 0 CHECK (self_heal_success_streak BETWEEN 0 AND 3), last_error_code TEXT CHECK (last_error_code IS NULL OR last_error_code IN ('LOCK_TIMEOUT', 'STATEMENT_TIMEOUT', 'DEADLOCK', 'SERIALIZATION', 'OUT_OF_MEMORY', 'CANCELLED', 'PERMANENT', 'UNKNOWN_RETRYABLE')), last_error_retryable BOOLEAN, downstream_publication_name TEXT, freshness_deadline_ms BIGINT, target_freshness_mode TEXT CHECK (target_freshness_mode IS NULL OR target_freshness_mode IN ('INTERVAL', 'ON_COMMIT', 'MANUAL')), refresh_reason TEXT, refresh_reason_detail TEXT, st_partition_key TEXT, in_shadow_build BOOLEAN NOT NULL DEFAULT FALSE, shadow_table_name TEXT, -- CITUS-3: Placement of this stream table's storage in a Citus cluster. st_placement TEXT NOT NULL DEFAULT 'local', -- v0.36.0: temporal IVM flag (CORR-1/UX-1) temporal_mode BOOLEAN NOT NULL DEFAULT FALSE, -- v0.36.0: columnar storage backend (CORR-2/UX-3) storage_backend TEXT NOT NULL DEFAULT 'heap', -- v0.47.0: post-refresh action hooks (VP-1/VP-2) post_refresh_action TEXT NOT NULL DEFAULT 'none' CHECK (post_refresh_action IN ('none', 'analyze', 'reindex', 'reindex_if_drift')), reindex_drift_threshold DOUBLE PRECISION CHECK (reindex_drift_threshold IS NULL OR (reindex_drift_threshold > 0 AND reindex_drift_threshold <= 1.0)), rows_changed_since_last_reindex BIGINT NOT NULL DEFAULT 0, last_reindex_at TIMESTAMPTZ, -- v0.36.0: column lineage metadata (F12) column_lineage JSONB, -- v0.59.0 PERF-2: hash of defining_query to skip recomputation on every refresh defining_query_hash BIGINT NOT NULL DEFAULT 0, -- v0.73.0 HOT-1: fillfactor for storage heap (NULL = PG default 100). Range 10-100. storage_fillfactor INT DEFAULT NULL CHECK (storage_fillfactor IS NULL OR (storage_fillfactor >= 10 AND storage_fillfactor <= 100)), -- v0.78.0 P-2: OpTree-derived complexity label, back-filled lazily on first refresh. query_complexity_class TEXT, -- v0.83.0: Composite row-identity encoding version. NULL is an -- unclassified pre-upgrade row; fresh objects are written explicitly. row_identity_version SMALLINT, -- v0.87.16: bounded identity probe encoding version. row_probe_version SMALLINT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), -- v0.87.7 LSEC-3: exact search_path defining_query was resolved under -- (bare $user already expanded). Set at CREATE and on any ALTER that -- changes the query; preserved by configuration-only ALTERs and -- ownership transfer. defining_search_path TEXT NOT NULL ); CREATE INDEX IF NOT EXISTS idx_pgt_status ON pgtrickle.pgt_stream_tables (status); CREATE UNIQUE INDEX IF NOT EXISTS idx_pgt_name ON pgtrickle.pgt_stream_tables (pgt_schema, pgt_name); -- PERF-4: Scheduler hot‐path lookup by relation OID. CREATE INDEX IF NOT EXISTS idx_pgt_relid ON pgtrickle.pgt_stream_tables (pgt_relid); -- v0.87.12: Immutable provenance for downstream publications. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_publication_bindings ( pgt_id BIGINT PRIMARY KEY REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, stream_relid OID NOT NULL, publication_oid OID NOT NULL UNIQUE, publication_name TEXT NOT NULL UNIQUE, publication_owner_oid OID NOT NULL, expected_relation_oids OID[] NOT NULL, CONSTRAINT pgt_publication_binding_relations_check CHECK (expected_relation_oids = ARRAY[stream_relid]) ); REVOKE ALL ON TABLE pgtrickle.pgt_publication_bindings FROM PUBLIC; -- v0.83.0: Durable private state registry for set-operation state. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_set_operation_states ( pgt_id BIGINT NOT NULL REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, node_ordinal INTEGER NOT NULL, operation TEXT NOT NULL CHECK (operation IN ('INTERSECT', 'EXCEPT')), is_all BOOLEAN NOT NULL, state_relid OID NOT NULL, schema_version SMALLINT NOT NULL, PRIMARY KEY (pgt_id, node_ordinal) ); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_set_operation_states', ''); -- Snapshot metadata catalog CREATE TABLE IF NOT EXISTS pgtrickle.pgt_snapshots ( snapshot_id BIGSERIAL PRIMARY KEY, pgt_id BIGINT NOT NULL REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, snapshot_schema TEXT NOT NULL, snapshot_table TEXT NOT NULL, snapshot_version TEXT NOT NULL, frontier JSONB, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), snapshot_relid OID, snapshot_provenance_token TEXT, created_by_role_oid OID, CONSTRAINT uq_snapshot_table UNIQUE (snapshot_schema, snapshot_table) ); CREATE INDEX IF NOT EXISTS idx_pgt_snapshots_pgt_id ON pgtrickle.pgt_snapshots (pgt_id); CREATE UNIQUE INDEX IF NOT EXISTS idx_pgt_snapshots_snapshot_relid ON pgtrickle.pgt_snapshots (snapshot_relid); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_snapshots', ''); -- Durable NOTIFY subscriptions CREATE TABLE IF NOT EXISTS pgtrickle.pgt_subscriptions ( stream_table TEXT NOT NULL, channel TEXT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (stream_table, channel) ); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_subscriptions', ''); -- DAG edges CREATE TABLE IF NOT EXISTS pgtrickle.pgt_dependencies ( pgt_id BIGINT NOT NULL REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, source_relid OID NOT NULL, source_type TEXT NOT NULL CHECK (source_type IN ('TABLE', 'STREAM_TABLE', 'VIEW', 'MATVIEW', 'FOREIGN_TABLE')), columns_used TEXT[], column_snapshot JSONB, schema_fingerprint TEXT, cdc_mode TEXT NOT NULL DEFAULT 'TRIGGER' CHECK (cdc_mode IN ('TRIGGER', 'TRANSITIONING', 'WAL')), slot_name TEXT, decoder_confirmed_lsn PG_LSN, transition_started_at TIMESTAMPTZ, cutover_target TEXT CHECK (cutover_target IN ('TRIGGER', 'WAL')), cutover_lsn PG_LSN, -- CITUS-3: Stable name for the source table (v0.32.0+). NULL = pre-upgrade row. source_stable_name TEXT, -- CITUS-3: Source placement in a Citus cluster: 'local', 'reference', 'distributed'. source_placement TEXT NOT NULL DEFAULT 'local', PRIMARY KEY (pgt_id, source_relid) ); CREATE INDEX IF NOT EXISTS idx_deps_source ON pgtrickle.pgt_dependencies (source_relid); -- PERF-4: Fast lookup by pgt_id (non‐PK prefix for multi‐column PK). CREATE INDEX IF NOT EXISTS idx_deps_pgt_id ON pgtrickle.pgt_dependencies (pgt_id); -- Refresh history / audit log CREATE TABLE IF NOT EXISTS pgtrickle.pgt_refresh_history ( refresh_id BIGSERIAL PRIMARY KEY, pgt_id BIGINT NOT NULL REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, data_timestamp TIMESTAMPTZ NOT NULL, start_time TIMESTAMPTZ NOT NULL, end_time TIMESTAMPTZ, action TEXT NOT NULL CHECK (action IN ('NO_DATA', 'FULL', 'DIFFERENTIAL', 'REINITIALIZE', 'SKIP')), rows_inserted BIGINT DEFAULT 0, rows_updated BIGINT NOT NULL DEFAULT 0, rows_deleted BIGINT DEFAULT 0, delta_row_count BIGINT DEFAULT 0, merge_strategy_used TEXT, was_full_fallback BOOLEAN NOT NULL DEFAULT FALSE, refresh_reason TEXT, refresh_reason_detail TEXT, error_message TEXT, status TEXT NOT NULL CHECK (status IN ('RUNNING', 'COMPLETED', 'FAILED', 'SKIPPED')), initiated_by TEXT CHECK (initiated_by IN ('SCHEDULER', 'MANUAL', 'INITIAL', 'SELF_MONITOR', 'SCHEDULER_FUSED')), freshness_deadline TIMESTAMPTZ, tick_watermark_lsn PG_LSN, fixpoint_iteration INT , error_code TEXT CHECK (error_code IS NULL OR error_code IN ('LOCK_TIMEOUT', 'STATEMENT_TIMEOUT', 'DEADLOCK', 'SERIALIZATION', 'OUT_OF_MEMORY', 'CANCELLED', 'PERMANENT', 'UNKNOWN_RETRYABLE')), error_sqlstate TEXT, retryable BOOLEAN ); CREATE INDEX IF NOT EXISTS idx_hist_pgt_ts ON pgtrickle.pgt_refresh_history (pgt_id, data_timestamp); -- PERF-1: Fast lookup by (pgt_id, start_time) for self-monitoring and scheduler_overhead queries. CREATE INDEX IF NOT EXISTS idx_hist_pgt_start ON pgtrickle.pgt_refresh_history (pgt_id, start_time); CREATE INDEX IF NOT EXISTS idx_hist_start_time ON pgtrickle.pgt_refresh_history (start_time, refresh_id); CREATE INDEX IF NOT EXISTS idx_hist_pgt_stats_window ON pgtrickle.pgt_refresh_history (pgt_id, start_time, status, action); -- v0.73.0 PERF-001: Incremental summary table for refresh-history metrics. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_refresh_summary ( pgt_id BIGINT PRIMARY KEY REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, total_refreshes BIGINT NOT NULL DEFAULT 0, successful_refreshes BIGINT NOT NULL DEFAULT 0, failed_refreshes BIGINT NOT NULL DEFAULT 0, total_rows_inserted BIGINT NOT NULL DEFAULT 0, total_rows_updated BIGINT NOT NULL DEFAULT 0, total_rows_deleted BIGINT NOT NULL DEFAULT 0, total_duration_ms BIGINT NOT NULL DEFAULT 0, last_refresh_action TEXT, last_refresh_status TEXT, last_refresh_at TIMESTAMPTZ, stats_reset_at TIMESTAMPTZ NOT NULL DEFAULT now(), total_full_refreshes BIGINT NOT NULL DEFAULT 0, total_diff_refreshes BIGINT NOT NULL DEFAULT 0, total_delta_rows_processed BIGINT NOT NULL DEFAULT 0, last_full_reason TEXT, last_full_reason_detail TEXT, updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); -- v0.78.0 P-3: Per-stream-table cost model summary. -- Populated by batch_update_cost_model_summary() on each scheduler tick. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_cost_model_summary ( pgt_id BIGINT NOT NULL REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, avg_full_ms DOUBLE PRECISION, avg_diff_ms DOUBLE PRECISION, sample_count INTEGER NOT NULL DEFAULT 0, p95_ms DOUBLE PRECISION, p99_ms DOUBLE PRECISION, updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), CONSTRAINT pgt_cost_model_summary_pk PRIMARY KEY (pgt_id) ); -- v0.73.0 ARCH-002 / REL-002: Persistent cleanup retry/backpressure status. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_cleanup_status ( source_relid OID PRIMARY KEY, buffer_table TEXT NOT NULL, attempt_count INT NOT NULL DEFAULT 0, blocked BOOLEAN NOT NULL DEFAULT false, last_error TEXT, last_operation TEXT, last_attempt_at TIMESTAMPTZ, next_retry_at TIMESTAMPTZ, backlog_rows BIGINT NOT NULL DEFAULT 0, updated_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX IF NOT EXISTS idx_cleanup_status_next_retry ON pgtrickle.pgt_cleanup_status (blocked, next_retry_at); -- Per-source CDC slot tracking CREATE TABLE IF NOT EXISTS pgtrickle.pgt_change_tracking ( source_relid OID PRIMARY KEY, slot_name TEXT NOT NULL, last_consumed_lsn PG_LSN, tracked_by_pgt_ids BIGINT[], -- CITUS-3: Stable hash name used for all pg_trickle-managed objects (v0.32.0+). -- NULL = pre-upgrade row using OID-based object names (STAB-1 fallback). source_stable_name TEXT, -- CITUS-3: How the source is placed in a Citus cluster. source_placement TEXT NOT NULL DEFAULT 'local', -- CITUS-7: Per-node WAL frontier for distributed sources. NULL for local sources. frontier_per_node JSONB ); -- v0.82.0: Registry for validated base-table and stream-table change buffers. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_change_buffers ( buffer_key TEXT PRIMARY KEY, source_kind TEXT NOT NULL CHECK (source_kind IN ('BASE', 'STREAM_TABLE')), source_id BIGINT NOT NULL, durability TEXT NOT NULL CHECK (durability IN ('logged', 'unlogged', 'sync')), sentinel_token BIGINT NOT NULL, -- v0.83.0: Composite row-identity encoding used by buffer writers. row_identity_version SMALLINT, -- v0.87.16: bounded identity probe encoding version. row_probe_version SMALLINT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), UNIQUE (source_kind, source_id) ); -- Scheduler job table for parallel refresh dispatch CREATE TABLE IF NOT EXISTS pgtrickle.pgt_scheduler_jobs ( job_id BIGSERIAL PRIMARY KEY, dag_version BIGINT NOT NULL, unit_key TEXT NOT NULL, unit_kind TEXT NOT NULL CHECK (unit_kind IN ('singleton', 'atomic_group', 'immediate_closure', 'cyclic_scc', 'repeatable_read_group', 'fused_chain')), member_pgt_ids BIGINT[] NOT NULL, root_pgt_id BIGINT NOT NULL, status TEXT NOT NULL DEFAULT 'QUEUED' CHECK (status IN ('QUEUED', 'RUNNING', 'SUCCEEDED', 'RETRYABLE_FAILED', 'PERMANENT_FAILED', 'CANCELLED')), scheduler_pid INT NOT NULL, worker_pid INT, attempt_no INT NOT NULL DEFAULT 1, enqueued_at TIMESTAMPTZ NOT NULL DEFAULT now(), started_at TIMESTAMPTZ, finished_at TIMESTAMPTZ, outcome_detail TEXT, retryable BOOLEAN, dispatch_tick_id BIGINT, tick_watermark_lsn PG_LSN , outcome_code TEXT CHECK (outcome_code IS NULL OR outcome_code IN ('LOCK_TIMEOUT', 'STATEMENT_TIMEOUT', 'DEADLOCK', 'SERIALIZATION', 'OUT_OF_MEMORY', 'CANCELLED', 'PERMANENT', 'UNKNOWN_RETRYABLE')), outcome_sqlstate TEXT, worker_slot_generation BIGINT ); CREATE INDEX IF NOT EXISTS idx_sched_jobs_status_enqueued ON pgtrickle.pgt_scheduler_jobs (status, enqueued_at); CREATE INDEX IF NOT EXISTS idx_sched_jobs_unit_status ON pgtrickle.pgt_scheduler_jobs (unit_key, status); CREATE INDEX IF NOT EXISTS idx_sched_jobs_finished ON pgtrickle.pgt_scheduler_jobs (finished_at) WHERE finished_at IS NOT NULL; CREATE INDEX IF NOT EXISTS idx_sched_jobs_terminal_finished ON pgtrickle.pgt_scheduler_jobs (finished_at, job_id) WHERE status IN ('SUCCEEDED', 'RETRYABLE_FAILED', 'PERMANENT_FAILED', 'CANCELLED'); -- Bootstrap source gates (v0.5.0, Phase 3) -- Records which source tables are currently "gated" (bootstrapping in progress). -- When a source is gated, all stream tables that depend on it are skipped by -- the scheduler until pgtrickle.ungate_source() is called. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_source_gates ( source_relid OID PRIMARY KEY, gated BOOLEAN NOT NULL DEFAULT true, gated_at TIMESTAMPTZ NOT NULL DEFAULT now(), ungated_at TIMESTAMPTZ, gated_by TEXT ); -- Per-source watermark state: tracks how far each external source has been loaded. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_watermarks ( source_relid OID PRIMARY KEY, watermark TIMESTAMPTZ NOT NULL, updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), advanced_by TEXT, wal_lsn_at_advance TEXT ); -- Watermark groups: declare that a set of sources must be temporally aligned. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_watermark_groups ( group_id SERIAL PRIMARY KEY, group_name TEXT UNIQUE NOT NULL, source_relids OID[] NOT NULL, tolerance_secs DOUBLE PRECISION NOT NULL DEFAULT 0, created_at TIMESTAMPTZ NOT NULL DEFAULT now() ); -- DB-3: Schema version tracking table. -- Records which schema migration versions have been applied to this database. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_schema_version ( version TEXT PRIMARY KEY, applied_at TIMESTAMPTZ NOT NULL DEFAULT now(), description TEXT ); INSERT INTO pgtrickle.pgt_schema_version (version, description) VALUES ('0.19.0', 'Initial schema version tracking') ON CONFLICT (version) DO NOTHING; INSERT INTO pgtrickle.pgt_schema_version (version, description) VALUES ( '0.84.0', 'Bootstrap catalog parity repair and manifest tooling baseline' ) ON CONFLICT (version) DO NOTHING; INSERT INTO pgtrickle.pgt_schema_version (version, description) VALUES ( '0.85.0', 'Scheduler and resource resilience gate' ) ON CONFLICT (version) DO NOTHING; SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_stream_tables', ''); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_dependencies', ''); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_change_buffers', ''); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_source_gates', ''); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_watermarks', ''); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_watermark_groups', ''); -- CIT-2: Cross-node advisory lock table for distributed stream table refresh. -- Uses INSERT … ON CONFLICT DO NOTHING for atomic acquisition; timestamp-based -- lease expiry handles crashed holders without requiring heartbeats. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_st_locks ( lock_key TEXT NOT NULL, holder TEXT NOT NULL, acquired_at TIMESTAMPTZ NOT NULL DEFAULT now(), expires_at TIMESTAMPTZ NOT NULL, CONSTRAINT pgt_st_locks_pkey PRIMARY KEY (lock_key) ); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_st_locks', ''); -- CIT-3: Per-worker logical replication slot tracking for distributed sources. -- Each row represents a WAL slot on one Citus worker node that feeds changes -- for a given (stream_table, source) pair into the coordinator's change buffer. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_worker_slots ( pgt_id BIGINT NOT NULL REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, source_relid OID NOT NULL, worker_name TEXT NOT NULL, worker_port INT NOT NULL DEFAULT 5432, slot_name TEXT NOT NULL, last_frontier TEXT, CONSTRAINT pgt_worker_slots_pkey PRIMARY KEY (pgt_id, source_relid, worker_name, worker_port) ); SELECT pg_catalog.pg_extension_config_dump('pgtrickle.pgt_worker_slots', ''); -- G14-SHC: Shared template cache (catalog-backed, UNLOGGED) CREATE UNLOGGED TABLE IF NOT EXISTS pgtrickle.pgt_template_cache ( pgt_id BIGINT PRIMARY KEY REFERENCES pgtrickle.pgt_stream_tables(pgt_id) ON DELETE CASCADE, query_hash BIGINT NOT NULL, delta_sql TEXT NOT NULL, columns TEXT[] NOT NULL, source_oids INTEGER[] NOT NULL, is_dedup BOOLEAN NOT NULL DEFAULT FALSE, key_changed BOOLEAN NOT NULL DEFAULT FALSE, all_algebraic BOOLEAN NOT NULL DEFAULT FALSE, cached_at TIMESTAMPTZ NOT NULL DEFAULT now() ); /* */ /* */ -- src/lib.rs:1471 CREATE OR REPLACE FUNCTION pgtrickle."pause_all"() RETURNS void LANGUAGE plpgsql AS $$ BEGIN UPDATE pgtrickle.pgt_stream_tables SET status = 'PAUSED' WHERE status = 'ACTIVE'; RAISE NOTICE 'pg_trickle: all stream tables paused.'; END; $$; COMMENT ON FUNCTION pgtrickle."pause_all"() IS 'Pause automatic refreshes for every ACTIVE stream table. ' 'Use pgtrickle.resume_all() to re-activate them.'; CREATE OR REPLACE FUNCTION pgtrickle."resume_all"() RETURNS void LANGUAGE plpgsql AS $$ BEGIN UPDATE pgtrickle.pgt_stream_tables SET status = 'ACTIVE' WHERE status = 'PAUSED'; RAISE NOTICE 'pg_trickle: all paused stream tables resumed.'; END; $$; COMMENT ON FUNCTION pgtrickle."resume_all"() IS 'Re-activate all stream tables that were paused with pgtrickle.pause_all().'; /* */ /* */ -- src/lib.rs:1275 -- Create event trigger functions with correct RETURNS event_trigger type. -- pgrx's #[pg_extern] generates RETURNS void, which PostgreSQL rejects for -- event triggers. We create them manually here with the correct return type. CREATE FUNCTION pgtrickle."_on_ddl_end"() RETURNS event_trigger SECURITY DEFINER -- nosemgrep: sql.security-definer.present — event trigger runs with pinned pgtrickle catalog search_path below. SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c AS 'MODULE_PATHNAME', 'pg_trickle_on_ddl_end_wrapper'; CREATE FUNCTION pgtrickle."_on_sql_drop"() RETURNS event_trigger SECURITY DEFINER -- nosemgrep: sql.security-definer.present — event trigger runs with pinned pgtrickle catalog search_path below. SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c AS 'MODULE_PATHNAME', 'pg_trickle_on_sql_drop_wrapper'; -- Event trigger: track ALTER TABLE on upstream sources CREATE EVENT TRIGGER pg_trickle_ddl_tracker ON ddl_command_end EXECUTE FUNCTION pgtrickle._on_ddl_end(); -- Event trigger: track DROP TABLE on upstream sources / ST storage tables CREATE EVENT TRIGGER pg_trickle_drop_tracker ON sql_drop EXECUTE FUNCTION pgtrickle._on_sql_drop(); /* */ /* */ -- src/lib.rs:1734 -- v0.46.0: Slim outbox attachment catalog. -- Maps stream tables to their pg_tide outbox names. -- The full outbox/inbox/relay stack lives in pg_tide (trickle-labs/pg-tide). CREATE TABLE IF NOT EXISTS pgtrickle.pgt_outbox_config ( stream_table_oid OID NOT NULL PRIMARY KEY, stream_table_name TEXT NOT NULL, tide_outbox_name TEXT NOT NULL, -- v0.87.13: immutable pg_tide identity used to reject stale name reuse. pg_tide_extension_oid OID NOT NULL, pg_tide_version TEXT NOT NULL, tide_outbox_created_at TIMESTAMPTZ NOT NULL, -- VA-4 (v0.48.0): optional vector column name for embedding outbox events. embedding_vector_column TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now() ); CREATE INDEX IF NOT EXISTS idx_pgt_outbox_config_name ON pgtrickle.pgt_outbox_config (stream_table_name); COMMENT ON TABLE pgtrickle.pgt_outbox_config IS 'v0.46.0: Catalog of stream tables with a pg_tide outbox attached via attach_outbox(). ' 'v0.48.0: embedding_vector_column set when attached via attach_embedding_outbox(). ' 'v0.87.13: pg_tide extension OID, version, and outbox creation time are ' 'stored to reject stale name reuse.'; /* */ /* */ -- src/lib.rs:756 -- The V2 upgrade is intentionally a rebuild, not an in-place conversion. -- These private catalogs make the destructive window explicit and auditable. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_row_identity_v2_inventory ( inventory_id BOOLEAN PRIMARY KEY DEFAULT TRUE CHECK (inventory_id), recorded_by TEXT NOT NULL, recorded_at TIMESTAMPTZ NOT NULL DEFAULT now(), notes TEXT, acknowledged BOOLEAN NOT NULL DEFAULT FALSE, acknowledged_by TEXT, acknowledged_at TIMESTAMPTZ, CHECK (NOT acknowledged OR acknowledged_at IS NOT NULL) ); CREATE TABLE IF NOT EXISTS pgtrickle.pgt_row_identity_v2_consumers ( consumer_id BIGSERIAL PRIMARY KEY, consumer_name TEXT NOT NULL, consumer_owner TEXT NOT NULL, affected_stream_tables TEXT[] NOT NULL CHECK (cardinality(affected_stream_tables) > 0), consumes_row_id BOOLEAN NOT NULL DEFAULT FALSE, consumes_storage_layout BOOLEAN NOT NULL DEFAULT FALSE, required_schema_change TEXT NOT NULL, resnapshot_status TEXT NOT NULL DEFAULT 'PENDING' CHECK (resnapshot_status IN ('PENDING', 'IN_PROGRESS', 'COMPLETE', 'SKIPPED')), acknowledged BOOLEAN NOT NULL DEFAULT FALSE, acknowledged_at TIMESTAMPTZ, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), updated_at TIMESTAMPTZ NOT NULL DEFAULT now(), CHECK (NOT acknowledged OR acknowledged_at IS NOT NULL) ); SELECT pg_catalog.pg_extension_config_dump( 'pgtrickle.pgt_row_identity_v2_inventory', '' ); SELECT pg_catalog.pg_extension_config_dump( 'pgtrickle.pgt_row_identity_v2_consumers', '' ); CREATE OR REPLACE FUNCTION pgtrickle._row_identity_v2_admin() RETURNS BOOLEAN LANGUAGE sql STABLE SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_admin$ SELECT EXISTS ( SELECT 1 FROM pg_roles r JOIN pg_extension e ON e.extname = 'pg_trickle' WHERE r.rolname = current_user AND (r.rolsuper OR r.oid = e.extowner) ) $row_identity_v2_admin$; CREATE OR REPLACE FUNCTION pgtrickle.row_identity_v2_record_inventory( p_notes TEXT DEFAULT NULL ) RETURNS VOID LANGUAGE plpgsql SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_record_inventory$ BEGIN IF NOT pgtrickle._row_identity_v2_admin() THEN RAISE EXCEPTION 'pg_trickle: row-identity inventory requires the extension owner or superuser' USING ERRCODE = '42501'; END IF; INSERT INTO pgtrickle.pgt_row_identity_v2_inventory ( inventory_id, recorded_by, recorded_at, notes, acknowledged, acknowledged_by, acknowledged_at ) VALUES (TRUE, current_user, now(), p_notes, FALSE, NULL, NULL) ON CONFLICT (inventory_id) DO UPDATE SET recorded_by = EXCLUDED.recorded_by, recorded_at = EXCLUDED.recorded_at, notes = EXCLUDED.notes, acknowledged = FALSE, acknowledged_by = NULL, acknowledged_at = NULL; END $row_identity_v2_record_inventory$; CREATE OR REPLACE FUNCTION pgtrickle.row_identity_v2_register_consumer( p_consumer_name TEXT, p_consumer_owner TEXT, p_affected_stream_tables TEXT[], p_consumes_row_id BOOLEAN, p_consumes_storage_layout BOOLEAN, p_required_schema_change TEXT ) RETURNS BIGINT LANGUAGE plpgsql SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_register_consumer$ DECLARE v_consumer_id BIGINT; BEGIN IF NOT pgtrickle._row_identity_v2_admin() THEN RAISE EXCEPTION 'pg_trickle: row-identity inventory requires the extension owner or superuser' USING ERRCODE = '42501'; END IF; IF NULLIF(btrim(p_consumer_name), '') IS NULL OR NULLIF(btrim(p_consumer_owner), '') IS NULL OR NULLIF(btrim(p_required_schema_change), '') IS NULL OR p_affected_stream_tables IS NULL OR cardinality(p_affected_stream_tables) = 0 OR EXISTS ( SELECT 1 FROM unnest(p_affected_stream_tables) AS affected(name) WHERE NULLIF(btrim(affected.name), '') IS NULL OR position('.' IN btrim(affected.name)) = 0 OR to_regclass(btrim(affected.name)) IS NULL ) THEN RAISE EXCEPTION 'pg_trickle: consumer inventory requires a name, owner, schema-qualified existing stream tables, and a schema-change plan' USING ERRCODE = '22023'; END IF; IF NOT EXISTS ( SELECT 1 FROM pgtrickle.pgt_row_identity_v2_inventory WHERE inventory_id ) THEN RAISE EXCEPTION 'pg_trickle: record the external-consumer inventory before registering consumers' USING ERRCODE = '55000'; END IF; INSERT INTO pgtrickle.pgt_row_identity_v2_consumers ( consumer_name, consumer_owner, affected_stream_tables, consumes_row_id, consumes_storage_layout, required_schema_change ) VALUES ( btrim(p_consumer_name), btrim(p_consumer_owner), p_affected_stream_tables, p_consumes_row_id, p_consumes_storage_layout, btrim(p_required_schema_change) ) RETURNING consumer_id INTO v_consumer_id; UPDATE pgtrickle.pgt_row_identity_v2_inventory SET acknowledged = FALSE, acknowledged_by = NULL, acknowledged_at = NULL; RETURN v_consumer_id; END $row_identity_v2_register_consumer$; CREATE OR REPLACE FUNCTION pgtrickle.row_identity_v2_acknowledge_inventory() RETURNS VOID LANGUAGE plpgsql SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_acknowledge_inventory$ BEGIN IF NOT pgtrickle._row_identity_v2_admin() THEN RAISE EXCEPTION 'pg_trickle: row-identity inventory requires the extension owner or superuser' USING ERRCODE = '42501'; END IF; UPDATE pgtrickle.pgt_row_identity_v2_inventory SET acknowledged = TRUE, acknowledged_by = current_user, acknowledged_at = now(); IF NOT FOUND THEN RAISE EXCEPTION 'pg_trickle: record the external-consumer inventory before acknowledging it' USING ERRCODE = '55000'; END IF; END $row_identity_v2_acknowledge_inventory$; CREATE OR REPLACE FUNCTION pgtrickle.row_identity_v2_acknowledge_consumer( p_consumer_id BIGINT, p_resnapshot_status TEXT DEFAULT 'PENDING' ) RETURNS VOID LANGUAGE plpgsql SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_acknowledge_consumer$ DECLARE v_status TEXT := upper(btrim(p_resnapshot_status)); BEGIN IF NOT pgtrickle._row_identity_v2_admin() THEN RAISE EXCEPTION 'pg_trickle: row-identity inventory requires the extension owner or superuser' USING ERRCODE = '42501'; END IF; IF v_status NOT IN ('PENDING', 'IN_PROGRESS', 'COMPLETE', 'SKIPPED') THEN RAISE EXCEPTION 'pg_trickle: invalid resnapshot status %', p_resnapshot_status USING ERRCODE = '22023'; END IF; UPDATE pgtrickle.pgt_row_identity_v2_consumers SET resnapshot_status = v_status, acknowledged = TRUE, acknowledged_at = now(), updated_at = now() WHERE consumer_id = p_consumer_id; IF NOT FOUND THEN RAISE EXCEPTION 'pg_trickle: external consumer % is not in the inventory', p_consumer_id USING ERRCODE = '22023'; END IF; END $row_identity_v2_acknowledge_consumer$; CREATE OR REPLACE FUNCTION pgtrickle.row_identity_v2_consumer_inventory() RETURNS TABLE ( consumer_id BIGINT, consumer_name TEXT, consumer_owner TEXT, affected_stream_tables TEXT[], consumes_row_id BOOLEAN, consumes_storage_layout BOOLEAN, required_schema_change TEXT, resnapshot_status TEXT, acknowledged BOOLEAN, acknowledged_at TIMESTAMPTZ ) LANGUAGE plpgsql STABLE SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_consumer_inventory$ BEGIN IF NOT pgtrickle._row_identity_v2_admin() THEN RAISE EXCEPTION 'pg_trickle: row-identity inventory requires the extension owner or superuser' USING ERRCODE = '42501'; END IF; RETURN QUERY SELECT c.consumer_id, c.consumer_name, c.consumer_owner, c.affected_stream_tables, c.consumes_row_id, c.consumes_storage_layout, c.required_schema_change, c.resnapshot_status, c.acknowledged, c.acknowledged_at FROM pgtrickle.pgt_row_identity_v2_consumers c ORDER BY c.consumer_id; END $row_identity_v2_consumer_inventory$; CREATE OR REPLACE FUNCTION pgtrickle.row_identity_v2_recreation_preflight() RETURNS JSONB LANGUAGE plpgsql SET search_path TO pgtrickle, pg_catalog, pg_temp AS $row_identity_v2_recreation_preflight$ DECLARE v_checks JSONB[] := ARRAY[]::JSONB[]; v_check JSONB; v_check_list JSONB; v_ok BOOLEAN; v_server_major INTEGER; v_bad_metadata BIGINT; v_bad_storage BIGINT; v_bad_buffers BIGINT; v_bad_sources BIGINT; v_bad_identity_types BIGINT; v_bad_indexes BIGINT; v_bad_index_sizes BIGINT; v_max_index_bytes BIGINT; v_scheduler_paused BOOLEAN; v_inventory_ack BOOLEAN; v_unacknowledged_consumers BIGINT; v_change_schema TEXT := COALESCE( current_setting('pg_trickle.change_buffer_schema', TRUE), 'pgtrickle_changes' ); BEGIN IF NOT pgtrickle._row_identity_v2_admin() THEN RAISE EXCEPTION 'pg_trickle: recreation preflight requires the extension owner or superuser' USING ERRCODE = '42501'; END IF; v_server_major := COALESCE( current_setting('server_version_num', TRUE), '0' )::INTEGER / 10000; SELECT count(*) INTO v_bad_metadata FROM ( SELECT 1 FROM pgtrickle.pgt_stream_tables WHERE row_identity_version IS DISTINCT FROM 2 OR row_probe_version IS DISTINCT FROM 1 UNION ALL SELECT 1 FROM pgtrickle.pgt_change_buffers WHERE row_identity_version IS DISTINCT FROM 2 OR row_probe_version IS DISTINCT FROM 1 ) invalid_metadata; SELECT count(*) INTO v_bad_storage FROM pgtrickle.pgt_stream_tables st LEFT JOIN pg_attribute a ON a.attrelid = st.pgt_relid AND a.attname = '__pgt_row_id' AND a.attnum > 0 AND NOT a.attisdropped WHERE a.atttypid IS DISTINCT FROM 'bytea'::regtype OR a.attnotnull IS DISTINCT FROM TRUE; SELECT count(*) INTO v_bad_buffers FROM pgtrickle.pgt_change_buffers cb LEFT JOIN pg_attribute a ON a.attrelid = to_regclass(format('%I.%I', v_change_schema, cb.buffer_key)) AND a.attname = '__pgt_row_id' AND a.attnum > 0 AND NOT a.attisdropped WHERE a.atttypid IS DISTINCT FROM 'bytea'::regtype OR a.attnotnull IS DISTINCT FROM TRUE; SELECT count(*) INTO v_bad_sources FROM pgtrickle.pgt_dependencies d LEFT JOIN pg_class c ON c.oid = d.source_relid WHERE d.source_type IN ('TABLE', 'STREAM_TABLE', 'VIEW', 'MATVIEW', 'FOREIGN_TABLE') AND c.oid IS NULL; SELECT count(*) INTO v_bad_identity_types FROM pgtrickle.pgt_dependencies d CROSS JOIN LATERAL unnest(COALESCE(d.columns_used, ARRAY[]::TEXT[])) AS u(column_name) LEFT JOIN pg_attribute a ON a.attrelid = d.source_relid AND a.attname = u.column_name AND a.attnum > 0 AND NOT a.attisdropped LEFT JOIN pg_type t ON t.oid = a.atttypid LEFT JOIN pg_collation co ON co.oid = a.attcollation WHERE d.source_type IN ('TABLE', 'STREAM_TABLE', 'VIEW', 'MATVIEW', 'FOREIGN_TABLE') AND ( a.attrelid IS NULL OR ( t.typtype <> 'e' AND t.typname <> ALL (ARRAY[ 'bool', 'int2', 'int4', 'int8', 'oid', 'float4', 'float8', 'numeric', 'text', 'varchar', 'bpchar', 'bytea', 'uuid', 'date', 'time', 'timestamp', 'timestamptz', 'timetz', 'interval', 'inet', 'cidr', 'macaddr', 'macaddr8', 'bit', 'varbit' ]::NAME[]) ) OR a.atttypmod IS NULL OR (co.oid IS NOT NULL AND co.collisdeterministic IS DISTINCT FROM TRUE) ); SELECT count(*) INTO v_bad_indexes FROM pgtrickle.pgt_stream_tables st WHERE NOT EXISTS ( SELECT 1 FROM pg_index i JOIN pg_class ic ON ic.oid = i.indexrelid JOIN pg_am am ON am.oid = ic.relam WHERE i.indrelid = st.pgt_relid AND i.indisvalid AND i.indisready AND am.amname = 'btree' AND NOT (0 = ANY (i.indclass)) AND pg_get_indexdef(i.indexrelid) LIKE '%__pgt_row_id%' ); SELECT count(*) INTO v_bad_index_sizes FROM pgtrickle.pgt_stream_tables st JOIN pg_index i ON i.indrelid = st.pgt_relid WHERE i.indisvalid AND pg_get_indexdef(i.indexrelid) LIKE '%__pgt_row_id%' AND pg_relation_size(i.indexrelid) IS NULL; SELECT COALESCE(max(pg_relation_size(i.indexrelid)), 0) INTO v_max_index_bytes FROM pgtrickle.pgt_stream_tables st JOIN pg_index i ON i.indrelid = st.pgt_relid WHERE i.indisvalid AND pg_get_indexdef(i.indexrelid) LIKE '%__pgt_row_id%'; v_scheduler_paused := COALESCE( current_setting('pg_trickle.enabled', TRUE), 'on' )::BOOLEAN = FALSE; SELECT COALESCE(bool_and(acknowledged), FALSE) INTO v_inventory_ack FROM pgtrickle.pgt_row_identity_v2_inventory WHERE inventory_id; SELECT count(*) INTO v_unacknowledged_consumers FROM pgtrickle.pgt_row_identity_v2_consumers WHERE NOT acknowledged; v_check := jsonb_build_object( 'check', 'SUPPORTED_POSTGRESQL_MAJOR', 'ok', v_server_major = 18, 'detail', format('server major=%s; V2 registry supports major 18', v_server_major) ); v_checks := array_append(v_checks, v_check); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'IDENTITY_METADATA', 'ok', v_bad_metadata = 0, 'detail', format('%s stream-table/buffer rows have stale or unknown version markers', v_bad_metadata) )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'STORAGE_SCHEMA', 'ok', v_bad_storage = 0, 'detail', format('%s stream tables lack __pgt_row_id BYTEA NOT NULL', v_bad_storage) )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'BUFFER_SCHEMA', 'ok', v_bad_buffers = 0, 'detail', format('%s change buffers lack __pgt_row_id BYTEA NOT NULL', v_bad_buffers) )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'IDENTITY_TYPES_AND_COLLATIONS', 'ok', v_bad_identity_types = 0, 'detail', format('%s source identity columns are missing, unsupported, or nondeterministically collated', v_bad_identity_types) )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'SOURCE_KEY_ELIGIBILITY', 'ok', v_bad_sources = 0, 'detail', format('%s source relations are missing; keyless sources remain valid and use exact full-ID matching', v_bad_sources) )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'ROW_ID_INDEXES', 'ok', v_bad_indexes = 0 AND v_bad_index_sizes = 0, 'detail', format('%s stream tables lack a valid row-id index; %s index sizes could not be read; largest=%s bytes', v_bad_indexes, v_bad_index_sizes, v_max_index_bytes) )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'SCHEDULER_PAUSED', 'ok', v_scheduler_paused, 'detail', CASE WHEN v_scheduler_paused THEN 'pg_trickle.enabled is off for this database' ELSE 'set pg_trickle.enabled=off and reload configuration before teardown' END )); v_checks := array_append(v_checks, jsonb_build_object( 'check', 'EXTERNAL_CONSUMER_INVENTORY', 'ok', v_inventory_ack AND v_unacknowledged_consumers = 0, 'detail', format('inventory acknowledged=%s; unacknowledged consumers=%s', v_inventory_ack, v_unacknowledged_consumers) )); SELECT jsonb_agg(value) INTO v_check_list FROM unnest(v_checks) AS checks(value); SELECT COALESCE(bool_and((value ->> 'ok')::BOOLEAN), TRUE) INTO v_ok FROM unnest(v_checks) AS checks(value); RETURN jsonb_build_object( 'ok', v_ok, 'release', '0.87.17', 'checks', COALESCE(v_check_list, '[]'::JSONB), 'instructions', jsonb_build_array( 'Pause the scheduler and wait for in-flight refreshes to drain.', 'Export stream-table definitions and metadata before dropping anything.', 'Drop V1 stream tables in reverse dependency order and clean only unused V1 buffers.', 'Install and restart the v0.87.17 binary, then run ALTER EXTENSION pg_trickle UPDATE.', 'Recreate stream tables in dependency order and perform fresh initial refreshes.', 'Update every external consumer to BYTEA, resnapshot after refresh, and acknowledge completion.', 'Writes made during the recreation window are not replayed; resume writes only after resnapshot.' ) ); END $row_identity_v2_recreation_preflight$; REVOKE ALL ON TABLE pgtrickle.pgt_row_identity_v2_inventory FROM PUBLIC; REVOKE ALL ON TABLE pgtrickle.pgt_row_identity_v2_consumers FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle._row_identity_v2_admin() FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle.row_identity_v2_record_inventory(TEXT) FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle.row_identity_v2_register_consumer(TEXT, TEXT, TEXT[], BOOLEAN, BOOLEAN, TEXT) FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle.row_identity_v2_acknowledge_inventory() FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle.row_identity_v2_acknowledge_consumer(BIGINT, TEXT) FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle.row_identity_v2_consumer_inventory() FROM PUBLIC; REVOKE ALL ON FUNCTION pgtrickle.row_identity_v2_recreation_preflight() FROM PUBLIC; /* */ /* */ -- src/lib.rs:1766 -- VH-2 (v0.48.0): Distance-predicate subscription catalog. -- Stores per-(stream_table, channel) vector distance subscriptions. -- After each non-empty refresh the background worker evaluates the predicate -- and emits pg_notify(channel, payload) when matched_rows > 0. CREATE TABLE IF NOT EXISTS pgtrickle.pgt_distance_subscriptions ( stream_table TEXT NOT NULL, channel TEXT NOT NULL, vector_column TEXT NOT NULL, query_vector TEXT NOT NULL, op TEXT NOT NULL CHECK (op IN ('<->', '<=>', '<#>', '<+>', '<<->>', '<<%>>')), threshold DOUBLE PRECISION NOT NULL CHECK (threshold > 0), created_at TIMESTAMPTZ NOT NULL DEFAULT now(), PRIMARY KEY (stream_table, channel) ); COMMENT ON TABLE pgtrickle.pgt_distance_subscriptions IS 'VH-2 (v0.48.0): Distance-predicate NOTIFY subscriptions per stream table. ' 'Populated via pgtrickle.subscribe_distance() / pgtrickle.unsubscribe_distance().'; /* */ /* */ -- src/lib.rs:1242 -- requires: -- pg_trickle_catalog -- CIT-4: Per-worker replication slot tracking view. -- Safe to query on non-Citus deployments (pgt_change_tracking and -- pgt_worker_slots are local catalog tables; no pg_dist_node reference). -- Returns one row per (stream_table, source, worker) combination. CREATE OR REPLACE VIEW pgtrickle.citus_status AS SELECT st.pgt_id, st.pgt_schema, st.pgt_name, ct.source_relid, ct.source_stable_name, ct.slot_name AS coordinator_slot, ct.source_placement, ct.frontier_per_node, ws.worker_name, ws.worker_port, ws.slot_name AS worker_slot, ws.last_frontier AS worker_frontier FROM pgtrickle.pgt_change_tracking ct JOIN pgtrickle.pgt_stream_tables st ON st.pgt_id = ANY(ct.tracked_by_pgt_ids) LEFT JOIN pgtrickle.pgt_worker_slots ws ON ws.pgt_id = st.pgt_id AND ws.source_relid = ct.source_relid WHERE ct.source_placement = 'distributed'; /* */ /* */ -- src/lib.rs:84 CREATE SCHEMA IF NOT EXISTS pgtrickle; /* pg_trickle::pgtrickle */ /* */ /* */ -- src/api/diagnostics.rs:2434 -- pg_trickle::api::diagnostics::preview_stream_table CREATE FUNCTION pgtrickle."preview_stream_table"( "query" TEXT, /* & str */ "schedule" TEXT DEFAULT 'calculated', /* Option < & str > */ "refresh_mode" TEXT DEFAULT 'AUTO', /* & str */ "target_freshness" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS TABLE ( "property" TEXT, /* String */ "value" TEXT /* String */ ) LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'preview_stream_table_wrapper'; /* */ /* */ -- src/api/mod.rs:2599 -- pg_trickle::api::explain_stream_table CREATE FUNCTION pgtrickle."explain_stream_table"( "name" TEXT /* & str */ ) RETURNS TEXT /* Result < String, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_stream_table_wrapper'; /* */ /* */ -- src/api/create.rs:759 -- pg_trickle::api::create::create_stream_table_realtime CREATE FUNCTION pgtrickle."create_stream_table_realtime"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT false, /* bool */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL /* Option < f64 > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_stream_table_realtime_wrapper'; /* */ /* */ -- src/monitor/mod.rs:2155 -- pg_trickle::monitor::history_prune_status CREATE FUNCTION pgtrickle."history_prune_status"() RETURNS TABLE ( "prune_error_count" bigint, /* i64 */ "last_prune_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "last_rows_deleted" bigint /* Option < i64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'history_prune_status_wrapper'; /* */ /* */ -- src/api/mod.rs:2108 -- pg_trickle::api::list_distance_subscriptions CREATE FUNCTION pgtrickle."list_distance_subscriptions"( "p_stream_table" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS TABLE ( "stream_table" TEXT, /* Option < String > */ "channel" TEXT, /* Option < String > */ "vector_column" TEXT, /* Option < String > */ "op" TEXT, /* Option < String > */ "threshold" double precision, /* Option < f64 > */ "created_at" timestamp with time zone /* Option < pgrx :: datum :: TimestampWithTimeZone > */ ) LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'list_distance_subscriptions_wrapper'; /* */ /* */ -- src/api/create.rs:717 -- pg_trickle::api::create::create_stream_table_fast_append_only CREATE FUNCTION pgtrickle."create_stream_table_fast_append_only"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "schedule" TEXT DEFAULT 'calculated', /* Option < & str > */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL /* Option < f64 > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_stream_table_fast_append_only_wrapper'; /* */ /* */ -- src/api/snapshot.rs:1210 -- pg_trickle::api::snapshot::drop_snapshot CREATE FUNCTION pgtrickle."drop_snapshot"( "p_snapshot_table" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'drop_snapshot_wrapper'; /* */ /* */ -- src/api/mod.rs:3219 -- pg_trickle::api::bulk_drop_stream_tables CREATE FUNCTION pgtrickle."bulk_drop_stream_tables"( "names" TEXT[] /* Vec < Option < String > > */ ) RETURNS INT /* i32 */ STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'bulk_drop_stream_tables_wrapper'; /* */ /* */ -- src/api/alter.rs:3331 -- pg_trickle::api::alter::pause_stream_table CREATE FUNCTION pgtrickle."pause_stream_table"( "name" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pause_stream_table_wrapper'; /* */ /* */ -- src/api/alter.rs:3069 -- pg_trickle::api::alter::repair_stream_table CREATE FUNCTION pgtrickle."repair_stream_table"( "name" TEXT /* & str */ ) RETURNS TEXT /* String */ STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'repair_stream_table_wrapper'; /* */ /* */ -- src/api/mod.rs:3330 -- pg_trickle::api::stream_table_lineage CREATE FUNCTION pgtrickle."stream_table_lineage"( "name" TEXT /* & str */ ) RETURNS TABLE ( "output_col" TEXT, /* Option < String > */ "source_table" TEXT, /* Option < String > */ "source_col" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'stream_table_lineage_wrapper'; /* */ /* */ -- src/api/spec.rs:74 -- pg_trickle::api::spec::stream_table_spec CREATE FUNCTION pgtrickle."stream_table_spec"( "qualified_name" TEXT /* & str */ ) RETURNS jsonb /* Option < pgrx :: JsonB > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'stream_table_spec_by_name_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1159 -- pg_trickle::api::diagnostics::ungate_source CREATE FUNCTION pgtrickle."ungate_source"( "source" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'ungate_source_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:516 -- pg_trickle::api::diagnostics::explain_refresh_mode CREATE FUNCTION pgtrickle."explain_refresh_mode"( "name" TEXT /* & str */ ) RETURNS TABLE ( "configured_mode" TEXT, /* String */ "effective_mode" TEXT, /* Option < String > */ "downgrade_reason" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_refresh_mode_wrapper'; /* */ /* */ -- src/api/self_monitoring.rs:566 -- pg_trickle::api::self_monitoring::explain_dag CREATE FUNCTION pgtrickle."explain_dag"( "format" TEXT DEFAULT 'mermaid' /* Option < & str > */ ) RETURNS TEXT /* Option < String > */ LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_dag_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1386 -- pg_trickle::api::diagnostics::drop_watermark_group CREATE FUNCTION pgtrickle."drop_watermark_group"( "group_name" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'drop_watermark_group_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2076 -- pg_trickle::api::diagnostics::clear_caches CREATE FUNCTION pgtrickle."clear_caches"() RETURNS bigint /* i64 */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'clear_caches_wrapper'; /* */ /* */ -- src/monitor/health.rs:890 -- pg_trickle::monitor::health::worker_pool_status CREATE FUNCTION pgtrickle."worker_pool_status"() RETURNS TABLE ( "active_workers" INT, /* i32 */ "max_workers" INT, /* i32 */ "per_db_cap" INT, /* i32 */ "parallel_mode" TEXT, /* String */ "idle_workers" INT, /* i32 */ "last_scheduler_tick_unix" bigint, /* i64 */ "ring_overflow_count" bigint, /* i64 */ "citus_failure_total" bigint /* i64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'worker_pool_status_wrapper'; /* */ /* */ -- src/api/publication.rs:920 -- pg_trickle::api::publication::set_stream_table_sla CREATE FUNCTION pgtrickle."set_stream_table_sla"( "name" TEXT, /* & str */ "sla" interval /* Interval */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'set_stream_table_sla_wrapper'; /* */ /* */ -- src/monitor/mod.rs:2075 -- pg_trickle::monitor::list_sources CREATE FUNCTION pgtrickle."list_sources"( "name" TEXT /* & str */ ) RETURNS TABLE ( "source_table" TEXT, /* String */ "source_oid" bigint, /* i64 */ "source_type" TEXT, /* String */ "cdc_mode" TEXT, /* String */ "columns_used" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'list_sources_wrapper'; /* */ /* */ -- src/api/refresh_ops.rs:9 -- pg_trickle::api::refresh_ops::refresh_stream_table CREATE FUNCTION pgtrickle."refresh_stream_table"( "name" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'refresh_stream_table_wrapper'; /* */ /* */ -- src/lib.rs:1509 -- requires: -- refresh_stream_table CREATE OR REPLACE FUNCTION pgtrickle."refresh_if_stale"( p_name text, p_max_age interval DEFAULT '5 minutes'::interval ) RETURNS boolean LANGUAGE plpgsql AS $$ DECLARE v_last_end timestamp with time zone; v_refreshed boolean := false; BEGIN SELECT MAX(end_time) INTO v_last_end FROM pgtrickle.pgt_refresh_history h JOIN pgtrickle.pgt_stream_tables s USING (pgt_id) WHERE s.pgt_name = p_name AND h.status = 'COMPLETED'; IF v_last_end IS NULL OR (now() - v_last_end) > p_max_age THEN PERFORM pgtrickle.refresh_stream_table(p_name); v_refreshed := true; END IF; RETURN v_refreshed; END; $$; COMMENT ON FUNCTION pgtrickle."refresh_if_stale"(text, interval) IS 'Refresh the named stream table only when the most recent completed ' 'refresh is older than max_age. Returns TRUE when a refresh was ' 'triggered, FALSE when the table was fresh enough.'; /* */ /* */ -- src/api/mod.rs:2899 -- pg_trickle::api::resume_after_drain CREATE FUNCTION pgtrickle."resume_after_drain"() RETURNS bool /* bool */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'resume_after_drain_wrapper'; /* */ /* */ -- src/api/create.rs:178 -- pg_trickle::api::create::bulk_create CREATE FUNCTION pgtrickle."bulk_create"( "definitions" jsonb /* pgrx :: JsonB */ ) RETURNS jsonb /* pgrx :: JsonB */ STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'bulk_create_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1482 -- pg_trickle::api::diagnostics::watermark_status CREATE FUNCTION pgtrickle."watermark_status"() RETURNS TABLE ( "group_name" TEXT, /* String */ "min_watermark" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "max_watermark" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "lag_secs" double precision, /* Option < f64 > */ "aligned" bool, /* bool */ "sources_with_watermark" INT, /* i32 */ "sources_total" INT /* i32 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'watermark_status_fn_wrapper'; /* */ /* */ -- src/api/alter.rs:1911 -- pg_trickle::api::alter::alter_stream_table CREATE FUNCTION pgtrickle."alter_stream_table"( "name" TEXT, /* & str */ "query" TEXT DEFAULT NULL, /* Option < & str > */ "schedule" TEXT DEFAULT NULL, /* Option < & str > */ "refresh_mode" TEXT DEFAULT NULL, /* Option < & str > */ "status" TEXT DEFAULT NULL, /* Option < & str > */ "diamond_consistency" TEXT DEFAULT NULL, /* Option < & str > */ "diamond_schedule_policy" TEXT DEFAULT NULL, /* Option < & str > */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT NULL, /* Option < bool > */ "pooler_compatibility_mode" bool DEFAULT NULL, /* Option < bool > */ "tier" TEXT DEFAULT NULL, /* Option < & str > */ "fuse" TEXT DEFAULT NULL, /* Option < & str > */ "fuse_ceiling" bigint DEFAULT NULL, /* Option < i64 > */ "fuse_sensitivity" INT DEFAULT NULL, /* Option < i32 > */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL, /* Option < f64 > */ "post_refresh_action" TEXT DEFAULT NULL, /* Option < & str > */ "reindex_drift_threshold" double precision DEFAULT NULL, /* Option < f64 > */ "target_freshness" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'alter_stream_table_wrapper'; /* */ /* */ -- src/api/self_monitoring.rs:240 -- pg_trickle::api::self_monitoring::setup_self_monitoring CREATE FUNCTION pgtrickle."setup_self_monitoring"() RETURNS void STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'setup_self_monitoring_wrapper'; /* */ /* */ -- src/api/helpers.rs:3652 -- pg_trickle::api::helpers::export_definition CREATE FUNCTION pgtrickle."export_definition"( "st_name" TEXT /* & str */ ) RETURNS TEXT /* Result < String, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'export_definition_wrapper'; /* */ /* */ -- src/lib.rs:1549 -- requires: -- export_definition CREATE OR REPLACE FUNCTION pgtrickle."stream_table_definition"( p_name text ) RETURNS text LANGUAGE sql STABLE AS $$ SELECT pgtrickle.export_definition(p_name); $$; COMMENT ON FUNCTION pgtrickle."stream_table_definition"(text) IS 'Return the CREATE STREAM TABLE DDL for the named stream table. ' 'Equivalent to pgtrickle.export_definition(name) — provided as a ' 'more discoverable alias.'; /* */ /* */ -- src/monitor/mod.rs:580 -- pg_trickle::monitor::get_refresh_history CREATE FUNCTION pgtrickle."get_refresh_history"( "name" TEXT, /* & str */ "max_rows" INT DEFAULT 20 /* i32 */ ) RETURNS TABLE ( "refresh_id" bigint, /* i64 */ "data_timestamp" timestamp with time zone, /* TimestampWithTimeZone */ "start_time" timestamp with time zone, /* TimestampWithTimeZone */ "end_time" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "action" TEXT, /* String */ "status" TEXT, /* String */ "rows_inserted" bigint, /* i64 */ "rows_updated" bigint, /* i64 */ "rows_deleted" bigint, /* i64 */ "duration_ms" double precision, /* Option < f64 > */ "error_message" TEXT, /* Option < String > */ "refresh_reason" TEXT, /* Option < String > */ "refresh_reason_detail" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'get_refresh_history_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2607 -- pg_trickle::api::diagnostics::stat_reset_all CREATE FUNCTION pgtrickle."stat_reset_all"() RETURNS void STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'stat_reset_all_wrapper'; /* */ /* */ -- src/ivm.rs:954 -- pg_trickle::ivm::pgt_ivm_apply_delta_enr CREATE FUNCTION pgtrickle."pgt_ivm_apply_delta_enr"( "pgt_id" bigint, /* i64 */ "source_oid" INT, /* i32 */ "has_new" bool, /* bool */ "has_old" bool /* bool */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pgt_ivm_apply_delta_enr_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:235 -- pg_trickle::api::diagnostics::version CREATE FUNCTION pgtrickle."version"() RETURNS TEXT /* & '_ str */ IMMUTABLE STRICT PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'version_wrapper'; /* */ /* */ -- src/citus.rs:103 -- pg_trickle::citus::source_stable_name CREATE FUNCTION pgtrickle."source_stable_name"( "source_oid" oid /* pg_sys :: Oid */ ) RETURNS TEXT /* Option < String > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'sql_stable_name_for_oid_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2409 -- pg_trickle::api::diagnostics::explain_json CREATE FUNCTION pgtrickle."explain_json"( "name" TEXT /* & str */ ) RETURNS jsonb /* Result < pgrx :: JsonB, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_json_wrapper'; /* */ /* */ -- src/api/mod.rs:2916 -- pg_trickle::api::cdc_pause_status CREATE FUNCTION pgtrickle."cdc_pause_status"() RETURNS TABLE ( "paused" bool, /* bool */ "capture_mode" TEXT, /* String */ "note" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'cdc_pause_status_wrapper'; /* */ /* */ -- src/api/create.rs:798 -- pg_trickle::api::create::create_stream_table_batch CREATE FUNCTION pgtrickle."create_stream_table_batch"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT false, /* bool */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL /* Option < f64 > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_stream_table_batch_wrapper'; /* */ /* */ -- src/api/alter.rs:3297 -- pg_trickle::api::alter::set_stream_table_storage_policy CREATE FUNCTION pgtrickle."set_stream_table_storage_policy"( "name" TEXT, /* & str */ "append_only" bool DEFAULT NULL, /* Option < bool > */ "tier" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'set_stream_table_storage_policy_wrapper'; /* */ /* */ -- src/api/mod.rs:1971 -- pg_trickle::api::subscribe CREATE FUNCTION pgtrickle."subscribe"( "stream_table" TEXT, /* & str */ "channel" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'subscribe_wrapper'; /* */ /* */ -- src/api/cluster.rs:18 -- pg_trickle::api::cluster::cluster_worker_summary CREATE FUNCTION pgtrickle."cluster_worker_summary"() RETURNS TABLE ( "db_oid" bigint, /* Option < i64 > */ "db_name" TEXT, /* Option < String > */ "active_workers" INT, /* Option < i32 > */ "scheduler_pid" INT, /* Option < i32 > */ "scheduler_running" bool, /* Option < bool > */ "total_active_workers" INT /* Option < i32 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'cluster_worker_summary_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:329 -- pg_trickle::api::diagnostics::rebuild_cdc_triggers CREATE FUNCTION pgtrickle."rebuild_cdc_triggers"() RETURNS TEXT /* & '_ str */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'rebuild_cdc_triggers_wrapper'; /* */ /* */ -- src/api/create.rs:837 -- pg_trickle::api::create::create_stream_table_cost_optimized CREATE FUNCTION pgtrickle."create_stream_table_cost_optimized"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT false, /* bool */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL /* Option < f64 > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_stream_table_cost_optimized_wrapper'; /* */ /* */ -- src/api/create.rs:404 -- pg_trickle::api::create::create_or_replace_stream_table CREATE FUNCTION pgtrickle."create_or_replace_stream_table"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "schedule" TEXT DEFAULT 'calculated', /* Option < & str > */ "refresh_mode" TEXT DEFAULT 'AUTO', /* & str */ "initialize" bool DEFAULT true, /* bool */ "diamond_consistency" TEXT DEFAULT NULL, /* Option < & str > */ "diamond_schedule_policy" TEXT DEFAULT NULL, /* Option < & str > */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT false, /* bool */ "pooler_compatibility_mode" bool DEFAULT false, /* bool */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL, /* Option < f64 > */ "output_distribution_column" TEXT DEFAULT NULL, /* Option < & str > */ "temporal" bool DEFAULT false, /* bool */ "storage_backend" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_or_replace_stream_table_wrapper'; /* */ /* */ -- src/monitor/health.rs:691 -- pg_trickle::monitor::health::refresh_timeline CREATE FUNCTION pgtrickle."refresh_timeline"( "max_rows" INT DEFAULT 50 /* i32 */ ) RETURNS TABLE ( "start_time" timestamp with time zone, /* TimestampWithTimeZone */ "stream_table" TEXT, /* String */ "action" TEXT, /* String */ "status" TEXT, /* String */ "rows_inserted" bigint, /* i64 */ "rows_updated" bigint, /* i64 */ "rows_deleted" bigint, /* i64 */ "duration_ms" double precision, /* Option < f64 > */ "error_message" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'refresh_timeline_wrapper'; /* */ /* */ -- src/api/planner.rs:117 -- pg_trickle::api::planner::recommend_schedule CREATE FUNCTION pgtrickle."recommend_schedule"( "p_name" TEXT /* & str */ ) RETURNS jsonb /* pgrx :: JsonB */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'recommend_schedule_wrapper'; /* */ /* */ -- src/api/create.rs:29 -- pg_trickle::api::create::create_stream_table CREATE FUNCTION pgtrickle."create_stream_table"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "schedule" TEXT DEFAULT 'calculated', /* Option < & str > */ "refresh_mode" TEXT DEFAULT 'AUTO', /* & str */ "initialize" bool DEFAULT true, /* bool */ "diamond_consistency" TEXT DEFAULT NULL, /* Option < & str > */ "diamond_schedule_policy" TEXT DEFAULT NULL, /* Option < & str > */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT false, /* bool */ "pooler_compatibility_mode" bool DEFAULT false, /* bool */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL, /* Option < f64 > */ "output_distribution_column" TEXT DEFAULT NULL, /* Option < & str > */ "temporal" bool DEFAULT false, /* bool */ "storage_backend" TEXT DEFAULT NULL, /* Option < & str > */ "fillfactor" INT DEFAULT NULL, /* Option < i32 > */ "target_freshness" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_stream_table_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1316 -- pg_trickle::api::diagnostics::advance_watermark CREATE FUNCTION pgtrickle."advance_watermark"( "source" TEXT, /* & str */ "watermark" timestamp with time zone /* TimestampWithTimeZone */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'advance_watermark_wrapper'; /* */ /* */ -- src/api/self_monitoring.rs:453 -- pg_trickle::api::self_monitoring::scheduler_overhead CREATE FUNCTION pgtrickle."scheduler_overhead"() RETURNS TABLE ( "total_refreshes_1h" bigint, /* i64 */ "df_refreshes_1h" bigint, /* i64 */ "df_refresh_fraction" double precision, /* Option < f64 > */ "avg_refresh_ms" double precision, /* Option < f64 > */ "avg_df_refresh_ms" double precision, /* Option < f64 > */ "total_refresh_time_s" double precision, /* Option < f64 > */ "df_refresh_time_s" double precision /* Option < f64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'scheduler_overhead_wrapper'; /* */ /* */ -- src/monitor/mod.rs:751 -- pg_trickle::monitor::slot_health CREATE FUNCTION pgtrickle."slot_health"() RETURNS TABLE ( "slot_name" TEXT, /* String */ "source_relid" bigint, /* i64 */ "active" bool, /* bool */ "retained_wal_bytes" bigint, /* i64 */ "wal_status" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'slot_health_wrapper'; /* */ /* */ -- src/hooks.rs:116 -- pg_trickle::hooks::pg_trickle_on_ddl_end -- Skipped due to `#[pgrx(sql = false)]` /* */ /* */ -- src/diagnostics.rs:980 -- pg_trickle::diagnostics::validate_query CREATE FUNCTION pgtrickle."validate_query"( "query" TEXT /* & str */ ) RETURNS TABLE ( "check_name" TEXT, /* String */ "result" TEXT, /* String */ "severity" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'validate_query_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:443 -- pg_trickle::api::diagnostics::pgt_scc_status CREATE FUNCTION pgtrickle."pgt_scc_status"() RETURNS TABLE ( "scc_id" INT, /* i32 */ "member_count" INT, /* i32 */ "members" TEXT[], /* Vec < String > */ "last_iterations" INT, /* Option < i32 > */ "last_converged_at" timestamp with time zone /* Option < TimestampWithTimeZone > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pgt_scc_status_wrapper'; /* */ /* */ -- src/api/spec.rs:57 -- pg_trickle::api::spec::stream_table_spec CREATE FUNCTION pgtrickle."stream_table_spec"( "relid" oid /* pg_sys :: Oid */ ) RETURNS jsonb /* Option < pgrx :: JsonB > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'stream_table_spec_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1673 -- pg_trickle::api::diagnostics::drop_refresh_group CREATE FUNCTION pgtrickle."drop_refresh_group"( "group_name" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'drop_refresh_group_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2112 -- pg_trickle::api::diagnostics::vector_status CREATE FUNCTION pgtrickle."vector_status"() RETURNS TABLE ( "name" TEXT, /* String */ "post_refresh_action" TEXT, /* String */ "reindex_drift_threshold" double precision, /* Option < f64 > */ "rows_changed_since_last_reindex" bigint, /* i64 */ "last_reindex_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "data_timestamp" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "embedding_lag" interval, /* Option < pgrx :: datum :: Interval > */ "estimated_rows" bigint, /* Option < i64 > */ "drift_pct" double precision /* Option < f64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'vector_status_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:896 -- pg_trickle::api::diagnostics::reset_fuse CREATE FUNCTION pgtrickle."reset_fuse"( "name" TEXT, /* & str */ "action" TEXT DEFAULT 'apply' /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'reset_fuse_wrapper'; /* */ /* */ -- src/citus.rs:806 -- pg_trickle::citus::handle_vp_promoted CREATE FUNCTION pgtrickle."handle_vp_promoted"( "payload" TEXT /* & str */ ) RETURNS bool /* bool */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'sql_handle_vp_promoted_wrapper'; /* */ /* */ -- src/monitor/mod.rs:854 -- pg_trickle::monitor::pgtrickle_refresh_stats CREATE FUNCTION pgtrickle."pgtrickle_refresh_stats"() RETURNS TABLE ( "stream_table" TEXT, /* String */ "mode" TEXT, /* String */ "avg_ms" double precision, /* f64 */ "p95_ms" double precision, /* f64 */ "p99_ms" double precision, /* f64 */ "refresh_count" bigint, /* i64 */ "last_refresh_at" timestamp with time zone /* Option < TimestampWithTimeZone > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pgtrickle_refresh_stats_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1447 -- pg_trickle::api::diagnostics::watermark_groups CREATE FUNCTION pgtrickle."watermark_groups"() RETURNS TABLE ( "group_name" TEXT, /* String */ "source_count" INT, /* i32 */ "tolerance_secs" double precision, /* f64 */ "created_at" timestamp with time zone /* TimestampWithTimeZone */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'watermark_groups_fn_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2214 -- pg_trickle::api::diagnostics::commit_latency_stats CREATE FUNCTION pgtrickle."commit_latency_stats"() RETURNS TABLE ( "pgt_schema" TEXT, /* String */ "pgt_name" TEXT, /* String */ "samples" bigint, /* i64 */ "min_ms" double precision, /* f64 */ "p50_ms" double precision, /* f64 */ "p95_ms" double precision, /* f64 */ "max_ms" double precision, /* f64 */ "tracking_mode" TEXT /* String */ ) STRICT STABLE PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'commit_latency_stats_wrapper'; /* */ /* */ -- src/api/mod.rs:2274 -- pg_trickle::api::embedding_stream_table CREATE FUNCTION pgtrickle."embedding_stream_table"( "name" TEXT, /* & str */ "source_table" TEXT, /* & str */ "vector_column" TEXT, /* & str */ "extra_columns" TEXT DEFAULT NULL, /* Option < & str > */ "refresh_interval" TEXT DEFAULT '1m', /* & str */ "index_type" TEXT DEFAULT 'hnsw', /* & str */ "dry_run" bool DEFAULT false /* bool */ ) RETURNS TABLE ( "action" TEXT /* Option < String > */ ) LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'embedding_stream_table_wrapper'; /* */ /* */ -- src/api/planner.rs:186 -- pg_trickle::api::planner::schedule_recommendations CREATE FUNCTION pgtrickle."schedule_recommendations"() RETURNS TABLE ( "name" TEXT, /* Option < String > */ "current_interval_seconds" double precision, /* Option < f64 > */ "recommended_interval_seconds" double precision, /* Option < f64 > */ "delta_pct" double precision, /* Option < f64 > */ "confidence" double precision, /* Option < f64 > */ "reasoning" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'schedule_recommendations_wrapper'; /* */ /* */ -- src/monitor/mod.rs:1220 -- pg_trickle::monitor::explain_diff_sql CREATE FUNCTION pgtrickle."explain_diff_sql"( "name" TEXT /* & str */ ) RETURNS TEXT /* Option < String > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_diff_sql_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:627 -- pg_trickle::api::diagnostics::explain_delta CREATE FUNCTION pgtrickle."explain_delta"( "name" TEXT, /* & str */ "format" TEXT DEFAULT 'text' /* & str */ ) RETURNS SETOF TEXT /* String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_delta_text_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1394 -- pg_trickle::api::diagnostics::watermarks CREATE FUNCTION pgtrickle."watermarks"() RETURNS TABLE ( "source_table" TEXT, /* String */ "schema_name" TEXT, /* String */ "watermark" timestamp with time zone, /* TimestampWithTimeZone */ "updated_at" timestamp with time zone, /* TimestampWithTimeZone */ "advanced_by" TEXT, /* Option < String > */ "wal_lsn" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'watermarks_fn_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:765 -- pg_trickle::api::diagnostics::shared_buffer_stats CREATE FUNCTION pgtrickle."shared_buffer_stats"() RETURNS TABLE ( "source_oid" bigint, /* i64 */ "source_table" TEXT, /* String */ "consumer_count" INT, /* i32 */ "consumers" TEXT, /* String */ "columns_tracked" INT, /* i32 */ "safe_frontier_lsn" TEXT, /* Option < String > */ "buffer_rows" bigint, /* i64 */ "is_partitioned" bool /* bool */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'shared_buffer_stats_fn_wrapper'; /* */ /* */ -- src/api/outbox.rs:436 -- pg_trickle::api::outbox::detach_outbox CREATE FUNCTION pgtrickle."detach_outbox"( "p_name" TEXT, /* & str */ "p_if_exists" bool DEFAULT false /* bool */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'detach_outbox_wrapper'; /* */ /* */ -- src/api/publication.rs:876 -- pg_trickle::api::publication::drop_stream_table_publication CREATE FUNCTION pgtrickle."drop_stream_table_publication"( "name" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'drop_stream_table_publication_wrapper'; /* */ /* */ -- src/monitor/health.rs:25 -- pg_trickle::monitor::health::health_check CREATE FUNCTION pgtrickle."health_check"() RETURNS TABLE ( "check_name" TEXT, /* String */ "severity" TEXT, /* String */ "detail" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'health_check_wrapper'; /* */ /* */ -- src/monitor/health.rs:932 -- pg_trickle::monitor::health::parallel_job_status CREATE FUNCTION pgtrickle."parallel_job_status"( "max_age_seconds" INT DEFAULT 300 /* i32 */ ) RETURNS TABLE ( "job_id" bigint, /* i64 */ "unit_key" TEXT, /* String */ "unit_kind" TEXT, /* String */ "status" TEXT, /* String */ "member_count" INT, /* i32 */ "attempt_no" INT, /* i32 */ "scheduler_pid" INT, /* i32 */ "worker_pid" INT, /* Option < i32 > */ "enqueued_at" timestamp with time zone, /* TimestampWithTimeZone */ "started_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "finished_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "duration_ms" double precision /* Option < f64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'parallel_job_status_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:314 -- pg_trickle::api::diagnostics::_signal_launcher_rescan CREATE FUNCTION pgtrickle."_signal_launcher_rescan"() RETURNS void STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', '_signal_launcher_rescan_wrapper'; /* */ /* */ -- src/api/helpers.rs:3770 -- pg_trickle::api::helpers::restore_stream_tables CREATE FUNCTION pgtrickle."restore_stream_tables"() RETURNS VOID /* Result < (), crate :: error :: PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'restore_stream_tables_wrapper'; /* */ /* */ -- src/diagnostics.rs:783 -- pg_trickle::diagnostics::list_auxiliary_columns CREATE FUNCTION pgtrickle."list_auxiliary_columns"( "name" TEXT /* & str */ ) RETURNS TABLE ( "column_name" TEXT, /* String */ "data_type" TEXT, /* String */ "purpose" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'list_auxiliary_columns_wrapper'; /* */ /* */ -- src/api/mod.rs:2791 -- pg_trickle::api::view_evolution_status CREATE FUNCTION pgtrickle."view_evolution_status"() RETURNS TABLE ( "stream_table" TEXT, /* Option < String > */ "in_shadow_build" bool, /* Option < bool > */ "shadow_table_name" TEXT, /* Option < String > */ "status" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'view_evolution_status_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1183 -- pg_trickle::api::diagnostics::source_gates CREATE FUNCTION pgtrickle."source_gates"() RETURNS TABLE ( "source_table" TEXT, /* String */ "schema_name" TEXT, /* String */ "gated" bool, /* bool */ "gated_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "ungated_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "gated_by" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'source_gates_fn_wrapper'; /* */ /* */ -- src/api/scheduler_control.rs:70 -- pg_trickle::api::scheduler_control::pause_scheduler CREATE FUNCTION pgtrickle."pause_scheduler"( "nodes" TEXT[] /* pgrx :: Array < & str > */ ) RETURNS TEXT /* & '_ str */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pause_scheduler_wrapper'; /* */ /* */ -- src/hooks.rs:1203 -- pg_trickle::hooks::pg_trickle_on_sql_drop -- Skipped due to `#[pgrx(sql = false)]` /* */ /* */ -- src/api/alter.rs:2637 -- pg_trickle::api::alter::drop_stream_table CREATE FUNCTION pgtrickle."drop_stream_table"( "name" TEXT, /* & str */ "cascade" bool DEFAULT false /* bool */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'drop_stream_table_wrapper'; /* */ /* */ -- src/lib.rs:1572 -- requires: -- create_stream_table -- drop_stream_table -- alter_stream_table CREATE OR REPLACE FUNCTION pgtrickle."canary_begin"( p_name text, p_new_query text ) RETURNS text LANGUAGE plpgsql AS $$ DECLARE v_schema text; v_table text; v_canary text; v_dot int; BEGIN v_dot := strpos(p_name, '.'); IF v_dot > 0 THEN v_schema := substr(p_name, 1, v_dot - 1); v_table := substr(p_name, v_dot + 1); ELSE v_schema := current_schema(); v_table := p_name; END IF; v_canary := '__pgt_canary_' || v_table; -- Drop any stale canary table from a previous run. BEGIN PERFORM pgtrickle.drop_stream_table(v_schema || '.' || v_canary); EXCEPTION WHEN OTHERS THEN NULL; -- ignore if it does not exist END; -- Create the canary stream table with the new query. PERFORM pgtrickle.create_stream_table( v_schema || '.' || v_canary, p_new_query ); RETURN format( 'Canary stream table %I.%I created. Run pgtrickle.canary_diff(%L) to compare.', v_schema, v_canary, p_name ); END; $$; COMMENT ON FUNCTION pgtrickle."canary_begin"(text, text) IS 'Start a shadow/canary test for the named stream table. ' 'Creates __pgt_canary_ with p_new_query and starts refreshing it. ' 'Use canary_diff(name) to inspect differences and canary_promote(name) to ' 'swap canary into production.'; CREATE OR REPLACE FUNCTION pgtrickle."canary_diff"( p_name text ) RETURNS TABLE( row_source text, diff_row text ) LANGUAGE plpgsql AS $$ DECLARE v_schema text; v_table text; v_canary text; v_dot int; v_sql text; BEGIN v_dot := strpos(p_name, '.'); IF v_dot > 0 THEN v_schema := substr(p_name, 1, v_dot - 1); v_table := substr(p_name, v_dot + 1); ELSE v_schema := current_schema(); v_table := p_name; END IF; v_canary := '__pgt_canary_' || v_table; -- Return rows in live-only vs canary-only using EXCEPT (symmetric difference). v_sql := format( '(SELECT %L AS row_source, t::text AS diff_row FROM %I.%I t EXCEPT SELECT %L, c::text FROM %I.%I c) UNION ALL (SELECT %L, c::text FROM %I.%I c EXCEPT SELECT %L, t::text FROM %I.%I t)', 'live_only', v_schema, v_table, 'canary_only', v_schema, v_canary, 'canary_only', v_schema, v_canary, 'live_only', v_schema, v_table ); RETURN QUERY EXECUTE v_sql; END; $$; COMMENT ON FUNCTION pgtrickle."canary_diff"(text) IS 'Compare the live stream table with its canary counterpart. ' 'Returns rows that exist in only one of the two tables. ' 'An empty result set indicates the new query produces the same output.'; CREATE OR REPLACE FUNCTION pgtrickle."canary_promote"( p_name text ) RETURNS text LANGUAGE plpgsql AS $$ DECLARE v_schema text; v_table text; v_canary text; v_dot int; v_new_query text; BEGIN v_dot := strpos(p_name, '.'); IF v_dot > 0 THEN v_schema := substr(p_name, 1, v_dot - 1); v_table := substr(p_name, v_dot + 1); ELSE v_schema := current_schema(); v_table := p_name; END IF; v_canary := '__pgt_canary_' || v_table; -- Read the defining query from the canary table. SELECT defining_query INTO v_new_query FROM pgtrickle.pgt_stream_tables WHERE pgt_schema = v_schema AND pgt_name = v_canary; IF v_new_query IS NULL THEN RAISE EXCEPTION 'No canary found for %. Run pgtrickle.canary_begin() first.', p_name; END IF; -- Promote: alter the live table to use the new query, then drop the canary. PERFORM pgtrickle.alter_stream_table(v_schema || '.' || v_table, query => v_new_query); BEGIN PERFORM pgtrickle.drop_stream_table(v_schema || '.' || v_canary); EXCEPTION WHEN OTHERS THEN NULL; END; RETURN format( 'Canary promoted: %I.%I now uses the canary query. Canary table dropped.', v_schema, v_table ); END; $$; COMMENT ON FUNCTION pgtrickle."canary_promote"(text) IS 'Promote the canary stream table to production. ' 'Calls ALTER STREAM TABLE with the canary query, then drops the canary table. ' 'Run pgtrickle.canary_diff(name) first to confirm the result set matches.'; /* */ /* */ -- src/api/mod.rs:2893 -- pg_trickle::api::is_drained CREATE FUNCTION pgtrickle."is_drained"() RETURNS bool /* Option < bool > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'is_drained_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1603 -- pg_trickle::api::diagnostics::create_refresh_group CREATE FUNCTION pgtrickle."create_refresh_group"( "group_name" TEXT, /* & str */ "members" TEXT[], /* Vec < String > */ "isolation" TEXT DEFAULT 'read_committed' /* & str */ ) RETURNS INT /* Result < i32, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_refresh_group_wrapper'; /* */ /* */ -- src/monitor/mod.rs:709 -- pg_trickle::monitor::st_auto_threshold CREATE FUNCTION pgtrickle."st_auto_threshold"( "name" TEXT /* & str */ ) RETURNS double precision /* Option < f64 > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'st_auto_threshold_wrapper'; /* */ /* */ -- src/monitor/mod.rs:820 -- pg_trickle::monitor::reliability_counters CREATE FUNCTION pgtrickle."reliability_counters"() RETURNS TABLE ( "invalidation_ring_overflows" bigint, /* i64 */ "dag_cycles_detected" bigint, /* i64 */ "template_cache_stale_evictions" bigint /* i64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'reliability_counters_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2572 -- pg_trickle::api::diagnostics::stat_reset CREATE FUNCTION pgtrickle."stat_reset"( "pgt_id" bigint /* i64 */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'stat_reset_wrapper'; /* */ /* */ -- src/api/mod.rs:2048 -- pg_trickle::api::subscribe_distance CREATE FUNCTION pgtrickle."subscribe_distance"( "stream_table" TEXT, /* & str */ "channel" TEXT, /* & str */ "vector_column" TEXT, /* & str */ "query_vector" TEXT, /* & str */ "op" TEXT, /* & str */ "threshold" double precision /* f64 */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'subscribe_distance_wrapper'; /* */ /* */ -- src/api/self_monitoring.rs:358 -- pg_trickle::api::self_monitoring::self_monitoring_status CREATE FUNCTION pgtrickle."self_monitoring_status"() RETURNS TABLE ( "st_name" TEXT, /* String */ "exists" bool, /* bool */ "status" TEXT, /* Option < String > */ "refresh_mode" TEXT, /* Option < String > */ "last_refresh_at" TEXT, /* Option < String > */ "total_refreshes" bigint /* Option < i64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'self_monitoring_status_wrapper'; /* */ /* */ -- src/api/mod.rs:1982 -- pg_trickle::api::unsubscribe CREATE FUNCTION pgtrickle."unsubscribe"( "stream_table" TEXT, /* & str */ "channel" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'unsubscribe_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:288 -- pg_trickle::api::diagnostics::migrate CREATE FUNCTION pgtrickle."migrate"() RETURNS TEXT /* String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'migrate_wrapper'; /* */ /* */ -- src/api/outbox.rs:547 -- pg_trickle::api::outbox::attach_embedding_outbox CREATE FUNCTION pgtrickle."attach_embedding_outbox"( "p_name" TEXT, /* & str */ "p_vector_column" TEXT, /* & str */ "p_retention_hours" INT DEFAULT 24, /* i32 */ "p_inline_threshold_rows" INT DEFAULT 10000 /* i32 */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'attach_embedding_outbox_wrapper'; /* */ /* */ -- src/dvm/row_id_v2.rs:1488 -- pg_trickle::dvm::row_id_v2::row_probe_v1 CREATE FUNCTION pgtrickle."row_probe_v1"( "input" bytea /* Vec < u8 > */ ) RETURNS bytea /* Vec < u8 > */ IMMUTABLE STRICT PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'row_probe_v1_wrapper'; /* */ /* */ -- src/monitor/tree.rs:18 -- pg_trickle::monitor::tree::dependency_tree CREATE FUNCTION pgtrickle."dependency_tree"() RETURNS TABLE ( "tree_line" TEXT, /* String */ "node" TEXT, /* String */ "node_type" TEXT, /* String */ "depth" INT, /* i32 */ "status" TEXT, /* Option < String > */ "refresh_mode" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'dependency_tree_wrapper'; /* */ /* */ -- src/api/helpers.rs:3488 -- pg_trickle::api::helpers::recommend_refresh_mode CREATE FUNCTION pgtrickle."recommend_refresh_mode"( "st_name" TEXT DEFAULT NULL /* Option < String > */ ) RETURNS TABLE ( "pgt_schema" TEXT, /* String */ "pgt_name" TEXT, /* String */ "current_mode" TEXT, /* String */ "effective_mode" TEXT, /* Option < String > */ "recommended_mode" TEXT, /* String */ "confidence" TEXT, /* String */ "reason" TEXT, /* String */ "signals" jsonb /* pgrx :: JsonB */ ) LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'recommend_refresh_mode_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1244 -- pg_trickle::api::diagnostics::bootstrap_gate_status CREATE FUNCTION pgtrickle."bootstrap_gate_status"() RETURNS TABLE ( "source_table" TEXT, /* String */ "schema_name" TEXT, /* String */ "gated" bool, /* bool */ "gated_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "ungated_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "gated_by" TEXT, /* Option < String > */ "gate_duration" interval, /* Option < pgrx :: datum :: Interval > */ "affected_stream_tables" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'bootstrap_gate_status_fn_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1685 -- pg_trickle::api::diagnostics::refresh_groups CREATE FUNCTION pgtrickle."refresh_groups"() RETURNS TABLE ( "group_id" INT, /* i32 */ "group_name" TEXT, /* String */ "member_count" INT, /* i32 */ "isolation" TEXT, /* String */ "created_at" timestamp with time zone /* TimestampWithTimeZone */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'refresh_groups_fn_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1135 -- pg_trickle::api::diagnostics::gate_source CREATE FUNCTION pgtrickle."gate_source"( "source" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'gate_source_wrapper'; /* */ /* */ -- src/hash.rs:46 -- pg_trickle::hash::pg_trickle_hash CREATE FUNCTION pgtrickle."pg_trickle_hash"( "input" TEXT /* Option < & str > */ ) RETURNS bigint /* i64 */ IMMUTABLE PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pg_trickle_hash_wrapper'; /* */ /* */ -- src/diagnostics.rs:608 -- pg_trickle::diagnostics::diagnose_errors CREATE FUNCTION pgtrickle."diagnose_errors"( "name" TEXT /* & str */ ) RETURNS TABLE ( "event_time" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "error_type" TEXT, /* String */ "error_message" TEXT, /* String */ "remediation" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'diagnose_errors_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1734 -- pg_trickle::api::diagnostics::worker_allocation_status CREATE FUNCTION pgtrickle."worker_allocation_status"() RETURNS TABLE ( "db_name" TEXT, /* String */ "workers_used" bigint, /* i64 */ "workers_quota" bigint, /* i64 */ "workers_queued" bigint, /* i64 */ "cluster_active" bigint, /* i64 */ "cluster_max" bigint /* i64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'worker_allocation_status_fn_wrapper'; /* */ /* */ -- src/api/scheduler_control.rs:177 -- pg_trickle::api::scheduler_control::resume_scheduler CREATE FUNCTION pgtrickle."resume_scheduler"( "nodes" TEXT[] /* pgrx :: Array < & str > */ ) RETURNS TEXT /* & '_ str */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'resume_scheduler_wrapper'; /* */ /* */ -- src/api/alter.rs:3268 -- pg_trickle::api::alter::set_stream_table_refresh_policy CREATE FUNCTION pgtrickle."set_stream_table_refresh_policy"( "name" TEXT, /* & str */ "refresh_mode" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'set_stream_table_refresh_policy_wrapper'; /* */ /* */ -- src/diagnostics.rs:262 -- pg_trickle::diagnostics::explain_query_rewrite CREATE FUNCTION pgtrickle."explain_query_rewrite"( "query" TEXT /* & str */ ) RETURNS TABLE ( "pass_name" TEXT, /* String */ "changed" bool, /* bool */ "sql_after" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_query_rewrite_wrapper'; /* */ /* */ -- src/monitor/health.rs:544 -- pg_trickle::monitor::health::health_summary CREATE FUNCTION pgtrickle."health_summary"() RETURNS TABLE ( "total_stream_tables" INT, /* i32 */ "active_count" INT, /* i32 */ "error_count" INT, /* i32 */ "suspended_count" INT, /* i32 */ "stale_count" INT, /* i32 */ "reinit_pending" INT, /* i32 */ "max_staleness_seconds" double precision, /* Option < f64 > */ "scheduler_status" TEXT, /* String */ "cache_hit_rate" double precision /* Option < f64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'health_summary_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:246 -- pg_trickle::api::diagnostics::version_check CREATE FUNCTION pgtrickle."version_check"() RETURNS TEXT /* String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'version_check_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:971 -- pg_trickle::api::diagnostics::fuse_status CREATE FUNCTION pgtrickle."fuse_status"() RETURNS TABLE ( "stream_table" TEXT, /* String */ "fuse_mode" TEXT, /* String */ "fuse_state" TEXT, /* String */ "fuse_ceiling" bigint, /* Option < i64 > */ "effective_ceiling" bigint, /* Option < i64 > */ "fuse_sensitivity" INT, /* Option < i32 > */ "blown_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "blow_reason" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'fuse_status_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2400 -- pg_trickle::api::diagnostics::explain CREATE FUNCTION pgtrickle."explain"( "name" TEXT /* & str */ ) RETURNS TEXT /* Result < String, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_wrapper'; /* */ /* */ -- src/ivm.rs:849 -- pg_trickle::ivm::pgt_ivm_apply_delta CREATE FUNCTION pgtrickle."pgt_ivm_apply_delta"( "pgt_id" bigint, /* i64 */ "source_oid" INT, /* i32 */ "has_new" bool, /* bool */ "has_old" bool /* bool */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pgt_ivm_apply_delta_wrapper'; /* */ /* */ -- src/monitor/health.rs:780 -- pg_trickle::monitor::health::trigger_inventory CREATE FUNCTION pgtrickle."trigger_inventory"() RETURNS TABLE ( "source_table" TEXT, /* String */ "source_oid" bigint, /* i64 */ "trigger_name" TEXT, /* String */ "trigger_type" TEXT, /* String */ "present" bool, /* bool */ "enabled" bool /* bool */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'trigger_inventory_wrapper'; /* */ /* */ -- src/dvm/row_id_v2.rs:1800 -- pg_trickle::dvm::row_id_v2::encode_row_id_v2 CREATE FUNCTION pgtrickle."encode_row_id_v2"( "domain" TEXT, /* & str */ "record" anyelement /* pgrx :: AnyElement */ ) RETURNS bytea /* Vec < u8 > */ STRICT STABLE PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'encode_row_id_v2_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2021 -- pg_trickle::api::diagnostics::lifecycle_preflight CREATE FUNCTION pgtrickle."lifecycle_preflight"() RETURNS TEXT /* String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'lifecycle_preflight_wrapper'; /* */ /* */ -- src/api/mod.rs:3195 -- pg_trickle::api::bulk_alter_stream_tables CREATE FUNCTION pgtrickle."bulk_alter_stream_tables"( "names" TEXT[], /* Vec < Option < String > > */ "params" json /* pgrx :: Json */ ) RETURNS INT /* i32 */ STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'bulk_alter_stream_tables_wrapper'; /* */ /* */ -- src/api/helpers.rs:3593 -- pg_trickle::api::helpers::refresh_efficiency CREATE FUNCTION pgtrickle."refresh_efficiency"() RETURNS TABLE ( "pgt_schema" TEXT, /* String */ "pgt_name" TEXT, /* String */ "refresh_mode" TEXT, /* String */ "total_refreshes" bigint, /* i64 */ "diff_count" bigint, /* i64 */ "full_count" bigint, /* i64 */ "avg_diff_ms" double precision, /* Option < f64 > */ "avg_full_ms" double precision, /* Option < f64 > */ "avg_change_ratio" double precision, /* Option < f64 > */ "diff_speedup" TEXT, /* Option < String > */ "last_refresh_at" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'refresh_efficiency_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:383 -- pg_trickle::api::diagnostics::pgt_status CREATE FUNCTION pgtrickle."pgt_status"() RETURNS TABLE ( "name" TEXT, /* String */ "status" TEXT, /* String */ "refresh_mode" TEXT, /* String */ "is_populated" bool, /* bool */ "consecutive_errors" INT, /* i32 */ "schedule" TEXT, /* Option < String > */ "data_timestamp" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "staleness" interval, /* Option < pgrx :: datum :: Interval > */ "scc_id" INT /* Option < i32 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pgt_status_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:2288 -- pg_trickle::api::diagnostics::tune_recommendations CREATE FUNCTION pgtrickle."tune_recommendations"() RETURNS TABLE ( "guc_name" TEXT, /* String */ "current_value" TEXT, /* String */ "recommended_value" TEXT, /* String */ "reason" TEXT /* String */ ) STRICT STABLE PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'tune_recommendations_wrapper'; /* */ /* */ -- src/api/self_monitoring.rs:305 -- pg_trickle::api::self_monitoring::teardown_self_monitoring CREATE FUNCTION pgtrickle."teardown_self_monitoring"() RETURNS void STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'teardown_self_monitoring_wrapper'; /* */ /* */ -- src/api/outbox.rs:342 -- pg_trickle::api::outbox::attach_outbox CREATE FUNCTION pgtrickle."attach_outbox"( "p_name" TEXT, /* & str */ "p_retention_hours" INT DEFAULT 24, /* i32 */ "p_inline_threshold_rows" INT DEFAULT 10000 /* i32 */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'attach_outbox_wrapper'; /* */ /* */ -- src/api/snapshot.rs:757 -- pg_trickle::api::snapshot::snapshot_stream_table CREATE FUNCTION pgtrickle."snapshot_stream_table"( "p_name" TEXT, /* & str */ "p_target" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS TEXT /* String */ SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'snapshot_stream_table_wrapper'; /* */ /* */ -- src/monitor/mod.rs:1994 -- pg_trickle::monitor::change_buffer_sizes CREATE FUNCTION pgtrickle."change_buffer_sizes"() RETURNS TABLE ( "stream_table" TEXT, /* String */ "source_table" TEXT, /* String */ "source_oid" bigint, /* i64 */ "cdc_mode" TEXT, /* String */ "pending_rows" bigint, /* i64 */ "buffer_bytes" bigint /* i64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'change_buffer_sizes_wrapper'; /* */ /* */ -- src/api/create.rs:89 -- pg_trickle::api::create::create_stream_table_if_not_exists CREATE FUNCTION pgtrickle."create_stream_table_if_not_exists"( "name" TEXT, /* & str */ "query" TEXT, /* & str */ "schedule" TEXT DEFAULT 'calculated', /* Option < & str > */ "refresh_mode" TEXT DEFAULT 'AUTO', /* & str */ "initialize" bool DEFAULT true, /* bool */ "diamond_consistency" TEXT DEFAULT NULL, /* Option < & str > */ "diamond_schedule_policy" TEXT DEFAULT NULL, /* Option < & str > */ "cdc_mode" TEXT DEFAULT NULL, /* Option < & str > */ "append_only" bool DEFAULT false, /* bool */ "pooler_compatibility_mode" bool DEFAULT false, /* bool */ "partition_by" TEXT DEFAULT NULL, /* Option < & str > */ "max_differential_joins" INT DEFAULT NULL, /* Option < i32 > */ "max_delta_fraction" double precision DEFAULT NULL, /* Option < f64 > */ "output_distribution_column" TEXT DEFAULT NULL, /* Option < & str > */ "temporal" bool DEFAULT false, /* bool */ "storage_backend" TEXT DEFAULT NULL, /* Option < & str > */ "fillfactor" INT DEFAULT NULL, /* Option < i32 > */ "target_freshness" TEXT DEFAULT NULL /* Option < & str > */ ) RETURNS void SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_stream_table_if_not_exists_wrapper'; /* */ /* */ -- src/monitor/mod.rs:774 -- pg_trickle::monitor::cache_stats CREATE FUNCTION pgtrickle."cache_stats"() RETURNS TABLE ( "l1_hits" bigint, /* i64 */ "l2_hits" bigint, /* i64 */ "misses" bigint, /* i64 */ "evictions" bigint, /* i64 */ "l1_size" INT, /* i32 */ "l1_bytes" bigint /* i64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'cache_stats_wrapper'; /* */ /* */ -- src/monitor/mod.rs:46 -- pg_trickle::monitor::st_refresh_stats CREATE FUNCTION pgtrickle."st_refresh_stats"() RETURNS TABLE ( "pgt_name" TEXT, /* String */ "pgt_schema" TEXT, /* String */ "status" TEXT, /* String */ "refresh_mode" TEXT, /* String */ "is_populated" bool, /* bool */ "total_refreshes" bigint, /* i64 */ "successful_refreshes" bigint, /* i64 */ "failed_refreshes" bigint, /* i64 */ "total_rows_inserted" bigint, /* i64 */ "total_rows_updated" bigint, /* i64 */ "total_rows_deleted" bigint, /* i64 */ "avg_duration_ms" double precision, /* f64 */ "last_refresh_action" TEXT, /* Option < String > */ "last_refresh_status" TEXT, /* Option < String > */ "last_refresh_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "staleness_secs" double precision, /* Option < f64 > */ "stale" bool, /* bool */ "consecutive_errors" INT, /* i32 */ "schedule" TEXT, /* Option < String > */ "refresh_tier" TEXT, /* String */ "last_error_message" TEXT, /* Option < String > */ "downstream_publication" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'st_refresh_stats_wrapper'; /* */ /* */ -- src/api/snapshot.rs:1127 -- pg_trickle::api::snapshot::list_snapshots CREATE FUNCTION pgtrickle."list_snapshots"( "p_name" TEXT /* & str */ ) RETURNS TABLE ( "snapshot_table" TEXT, /* Option < String > */ "created_at" timestamp with time zone, /* Option < TimestampWithTimeZone > */ "row_count" bigint, /* Option < i64 > */ "frontier" jsonb, /* Option < pgrx :: JsonB > */ "size_bytes" bigint /* Option < i64 > */ ) STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'list_snapshots_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:374 -- pg_trickle::api::diagnostics::parse_duration_seconds CREATE FUNCTION pgtrickle."parse_duration_seconds"( "input" TEXT /* & str */ ) RETURNS bigint /* Option < i64 > */ IMMUTABLE STRICT PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'parse_duration_seconds_wrapper'; /* */ /* */ -- src/lib.rs:1432 -- requires: -- parse_duration_seconds -- ERG-E: One-row health summary for dashboards and alerting. CREATE OR REPLACE VIEW pgtrickle.quick_health AS SELECT (SELECT count(*) FROM pgtrickle.pgt_stream_tables)::bigint AS total_stream_tables, (SELECT count(*) FROM pgtrickle.pgt_stream_tables WHERE status = 'ERROR' OR consecutive_errors > 0)::bigint AS error_tables, (SELECT count(*) FROM pgtrickle.pgt_stream_tables WHERE schedule IS NOT NULL AND schedule !~ '[\s@]' AND last_refresh_at IS NOT NULL AND EXTRACT(EPOCH FROM (now() - last_refresh_at)) > pgtrickle.parse_duration_seconds(schedule))::bigint AS stale_tables, (SELECT count(*) > 0 FROM pg_stat_activity WHERE backend_type = 'pg_trickle scheduler') AS scheduler_running, CASE WHEN (SELECT count(*) FROM pgtrickle.pgt_stream_tables) = 0 THEN 'EMPTY' WHEN (SELECT count(*) FROM pgtrickle.pgt_stream_tables WHERE status = 'SUSPENDED') > 0 THEN 'CRITICAL' WHEN (SELECT count(*) FROM pgtrickle.pgt_stream_tables WHERE status = 'ERROR' OR consecutive_errors > 0) > 0 THEN 'WARNING' WHEN (SELECT count(*) FROM pgtrickle.pgt_stream_tables WHERE schedule IS NOT NULL AND schedule !~ '[\s@]' AND last_refresh_at IS NOT NULL AND EXTRACT(EPOCH FROM (now() - last_refresh_at)) > pgtrickle.parse_duration_seconds(schedule)) > 0 THEN 'WARNING' ELSE 'OK' END AS status; /* */ /* */ -- src/lib.rs:1203 -- requires: -- parse_duration_seconds -- Status overview view (ERR-1d: last_error_message and last_error_at are -- included via st.* from pgt_stream_tables) CREATE OR REPLACE VIEW pgtrickle.stream_tables_info AS SELECT st.*, now() - st.last_refresh_at AS staleness, CASE WHEN st.schedule IS NOT NULL AND st.schedule !~ '[\s@]' THEN EXTRACT(EPOCH FROM (now() - st.last_refresh_at)) > pgtrickle.parse_duration_seconds(st.schedule) ELSE NULL::boolean END AS stale, CASE WHEN st.topk_limit IS NOT NULL THEN TRUE ELSE FALSE END AS is_topk FROM pgtrickle.pgt_stream_tables st; /* */ /* */ -- src/lib.rs:1309 -- requires: -- parse_duration_seconds -- Convenience view: pg_stat_stream_tables -- Combines catalog metadata with aggregate refresh statistics. CREATE OR REPLACE VIEW pgtrickle.pg_stat_stream_tables AS SELECT st.pgt_id, st.pgt_schema, st.pgt_name, st.status, st.refresh_mode, st.is_populated, st.data_timestamp, st.schedule, now() - st.last_refresh_at AS staleness, CASE WHEN st.schedule IS NOT NULL AND st.last_refresh_at IS NOT NULL AND st.schedule !~ '[\s@]' THEN EXTRACT(EPOCH FROM (now() - st.last_refresh_at)) > pgtrickle.parse_duration_seconds(st.schedule) ELSE NULL::boolean END AS stale, st.consecutive_errors, st.needs_reinit, st.last_refresh_at, COALESCE(stats.total_refreshes, 0) AS total_refreshes, COALESCE(stats.successful_refreshes, 0) AS successful_refreshes, COALESCE(stats.failed_refreshes, 0) AS failed_refreshes, COALESCE(stats.total_rows_inserted, 0) AS total_rows_inserted, COALESCE(stats.total_rows_updated, 0) AS total_rows_updated, COALESCE(stats.total_rows_deleted, 0) AS total_rows_deleted, stats.avg_duration_ms, stats.last_action, stats.last_status, (SELECT array_agg(DISTINCT d.cdc_mode ORDER BY d.cdc_mode) FROM pgtrickle.pgt_dependencies d WHERE d.pgt_id = st.pgt_id AND d.source_type = 'TABLE') AS cdc_modes, st.scc_id, st.last_fixpoint_iterations, st.refresh_tier FROM pgtrickle.pgt_stream_tables st LEFT JOIN LATERAL ( SELECT count(*)::bigint AS total_refreshes, count(*) FILTER (WHERE h.status = 'COMPLETED')::bigint AS successful_refreshes, count(*) FILTER (WHERE h.status = 'FAILED')::bigint AS failed_refreshes, COALESCE(sum(h.rows_inserted), 0)::bigint AS total_rows_inserted, COALESCE(sum(h.rows_updated), 0)::bigint AS total_rows_updated, COALESCE(sum(h.rows_deleted), 0)::bigint AS total_rows_deleted, CASE WHEN count(*) FILTER (WHERE h.end_time IS NOT NULL) > 0 THEN avg(EXTRACT(EPOCH FROM (h.end_time - h.start_time)) * 1000) FILTER (WHERE h.end_time IS NOT NULL) ELSE NULL END::float8 AS avg_duration_ms, (SELECT h2.action FROM pgtrickle.pgt_refresh_history h2 WHERE h2.pgt_id = st.pgt_id ORDER BY h2.refresh_id DESC LIMIT 1) AS last_action, (SELECT h2.status FROM pgtrickle.pgt_refresh_history h2 WHERE h2.pgt_id = st.pgt_id ORDER BY h2.refresh_id DESC LIMIT 1) AS last_status, (SELECT h2.initiated_by FROM pgtrickle.pgt_refresh_history h2 WHERE h2.pgt_id = st.pgt_id ORDER BY h2.refresh_id DESC LIMIT 1) AS last_initiated_by, (SELECT h2.freshness_deadline FROM pgtrickle.pgt_refresh_history h2 WHERE h2.pgt_id = st.pgt_id ORDER BY h2.refresh_id DESC LIMIT 1) AS freshness_deadline FROM pgtrickle.pgt_refresh_history h WHERE h.pgt_id = st.pgt_id ) stats ON true; -- v0.86.0: Bounded per-stream-table diagnostic statistics. -- All cumulative values come from summary tables; scraping this view never -- scans refresh history. CREATE OR REPLACE VIEW pgtrickle.pg_stat_pgtrickle AS SELECT st.pgt_id, st.pgt_schema AS schema_name, st.pgt_name AS table_name, COALESCE(s.total_refreshes, 0)::bigint AS total_refreshes, COALESCE(s.total_full_refreshes, 0)::bigint AS total_full_refreshes, COALESCE(s.total_diff_refreshes, 0)::bigint AS total_diff_refreshes, COALESCE(s.total_delta_rows_processed, 0)::bigint AS total_delta_rows_processed, CASE WHEN COALESCE(s.total_refreshes, 0) > 0 THEN s.total_duration_ms::double precision / s.total_refreshes END AS avg_refresh_duration_ms, c.p95_ms AS p95_refresh_duration_ms, c.p99_ms AS p99_refresh_duration_ms, s.last_refresh_at, CASE WHEN COALESCE(s.last_refresh_at, st.created_at) IS NOT NULL THEN EXTRACT(EPOCH FROM (now() - COALESCE(s.last_refresh_at, st.created_at))) * 1000 END AS current_lag_ms, COALESCE(st.requested_refresh_mode, st.refresh_mode) AS requested_refresh_mode, st.target_freshness_mode, st.freshness_deadline_ms AS target_freshness_ms, s.last_full_reason, s.last_full_reason_detail, st.last_error_message AS last_error, st.last_error_at, s.stats_reset_at FROM pgtrickle.pgt_stream_tables st LEFT JOIN pgtrickle.pgt_refresh_summary s ON s.pgt_id = st.pgt_id LEFT JOIN pgtrickle.pgt_cost_model_summary c ON c.pgt_id = st.pgt_id; -- Per-source CDC status view (G5): exposes cdc_mode, slot names, and -- transition timestamps for every TABLE dependency of a stream table. CREATE OR REPLACE VIEW pgtrickle.pgt_cdc_status AS SELECT st.pgt_schema, st.pgt_name, d.source_relid, c.relname AS source_name, n.nspname AS source_schema, d.cdc_mode, d.slot_name, d.decoder_confirmed_lsn, d.transition_started_at FROM pgtrickle.pgt_dependencies d JOIN pgtrickle.pgt_stream_tables st ON st.pgt_id = d.pgt_id JOIN pg_class c ON c.oid = d.source_relid JOIN pg_namespace n ON n.oid = c.relnamespace WHERE d.source_type = 'TABLE'; /* */ /* */ -- src/ivm.rs:1327 -- pg_trickle::ivm::pgt_ivm_handle_truncate CREATE FUNCTION pgtrickle."pgt_ivm_handle_truncate"( "pgt_id" bigint /* i64 */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pgt_ivm_handle_truncate_wrapper'; /* */ /* */ -- src/api/snapshot.rs:996 -- pg_trickle::api::snapshot::restore_from_snapshot CREATE FUNCTION pgtrickle."restore_from_snapshot"( "p_name" TEXT, /* & str */ "p_source" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'restore_from_snapshot_wrapper'; /* */ /* */ -- src/api/metrics_ext.rs:28 -- pg_trickle::api::metrics_ext::metrics_summary CREATE FUNCTION pgtrickle."metrics_summary"() RETURNS TABLE ( "db_name" TEXT, /* Option < String > */ "total_stream_tables" bigint, /* Option < i64 > */ "active_stream_tables" bigint, /* Option < i64 > */ "suspended_stream_tables" bigint, /* Option < i64 > */ "total_refreshes" bigint, /* Option < i64 > */ "successful_refreshes" bigint, /* Option < i64 > */ "failed_refreshes" bigint, /* Option < i64 > */ "total_rows_processed" bigint, /* Option < i64 > */ "active_workers" INT, /* Option < i32 > */ "ivm_lock_parse_error_count" bigint, /* Option < i64 > */ "holdback_probe_calls" bigint, /* Option < i64 > */ "holdback_probe_cache_hits" bigint, /* Option < i64 > */ "holdback_probe_avg_ms" double precision, /* Option < f64 > */ "cleanup_backlog_count" bigint, /* Option < i64 > */ "cleanup_blocked_count" bigint /* Option < i64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'metrics_summary_wrapper'; /* */ /* */ -- src/api/alter.rs:2997 -- pg_trickle::api::alter::resume_stream_table CREATE FUNCTION pgtrickle."resume_stream_table"( "name" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'resume_stream_table_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1825 -- pg_trickle::api::diagnostics::preflight CREATE FUNCTION pgtrickle."preflight"() RETURNS TEXT /* String */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'preflight_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1349 -- pg_trickle::api::diagnostics::create_watermark_group CREATE FUNCTION pgtrickle."create_watermark_group"( "group_name" TEXT, /* & str */ "sources" TEXT[], /* Vec < String > */ "tolerance_secs" double precision DEFAULT 0.0 /* f64 */ ) RETURNS INT /* Result < i32, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'create_watermark_group_wrapper'; /* */ /* */ -- src/api/mod.rs:3252 -- pg_trickle::api::exec_stream_ddl CREATE FUNCTION pgtrickle."exec_stream_ddl"( "cmd" TEXT /* & str */ ) RETURNS bool /* bool */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'exec_stream_ddl_wrapper'; /* */ /* */ -- src/monitor/mod.rs:728 -- pg_trickle::monitor::get_staleness CREATE FUNCTION pgtrickle."get_staleness"( "name" TEXT /* & str */ ) RETURNS double precision /* Option < f64 > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'get_staleness_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:736 -- pg_trickle::api::diagnostics::dedup_stats CREATE FUNCTION pgtrickle."dedup_stats"() RETURNS TABLE ( "total_diff_refreshes" bigint, /* i64 */ "dedup_needed" bigint, /* i64 */ "dedup_ratio_pct" double precision /* f64 */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'dedup_stats_fn_wrapper'; /* */ /* */ -- src/api/mod.rs:1997 -- pg_trickle::api::list_subscriptions CREATE FUNCTION pgtrickle."list_subscriptions"() RETURNS TABLE ( "stream_table" TEXT, /* Option < String > */ "channel" TEXT, /* Option < String > */ "created_at" timestamp with time zone /* Option < pgrx :: datum :: TimestampWithTimeZone > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'list_subscriptions_wrapper'; /* */ /* */ -- src/monitor/mod.rs:931 -- pg_trickle::monitor::explain_st CREATE FUNCTION pgtrickle."explain_st"( "name" TEXT, /* & str */ "with_analyze" bool DEFAULT false /* bool */ ) RETURNS TABLE ( "property" TEXT, /* String */ "value" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'explain_st_wrapper'; /* */ /* */ -- src/api/mod.rs:2850 -- pg_trickle::api::drain CREATE FUNCTION pgtrickle."drain"( "timeout_s" INT DEFAULT 60 /* i32 */ ) RETURNS bool /* bool */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'drain_wrapper'; /* */ /* */ -- src/api/publication.rs:773 -- pg_trickle::api::publication::stream_table_to_publication CREATE FUNCTION pgtrickle."stream_table_to_publication"( "name" TEXT /* & str */ ) RETURNS void STRICT SECURITY DEFINER SET search_path TO pgtrickle, pg_catalog, pg_temp LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'stream_table_to_publication_wrapper'; /* */ /* */ -- src/api/mod.rs:2093 -- pg_trickle::api::unsubscribe_distance CREATE FUNCTION pgtrickle."unsubscribe_distance"( "stream_table" TEXT, /* & str */ "channel" TEXT /* & str */ ) RETURNS VOID /* Result < (), PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'unsubscribe_distance_wrapper'; /* */ /* */ -- src/monitor/health.rs:1107 -- pg_trickle::monitor::health::wal_source_status CREATE FUNCTION pgtrickle."wal_source_status"() RETURNS TABLE ( "source_relid" bigint, /* i64 */ "source_name" TEXT, /* String */ "cdc_mode" TEXT, /* String */ "slot_name" TEXT, /* Option < String > */ "slot_lag_bytes" bigint, /* i64 */ "publication_name" TEXT, /* Option < String > */ "blocked_reason" TEXT, /* Option < String > */ "transition_started_at" TEXT, /* Option < String > */ "decoder_confirmed_lsn" TEXT /* Option < String > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'wal_source_status_wrapper'; /* */ /* */ -- src/api/diagnostics.rs:1042 -- pg_trickle::api::diagnostics::diamond_groups CREATE FUNCTION pgtrickle."diamond_groups"() RETURNS TABLE ( "group_id" INT, /* i32 */ "member_name" TEXT, /* String */ "member_schema" TEXT, /* String */ "is_convergence" bool, /* bool */ "epoch" bigint, /* i64 */ "schedule_policy" TEXT /* String */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'diamond_groups_wrapper'; /* */ /* */ -- src/api/helpers.rs:3408 -- pg_trickle::api::helpers::convert_buffers_to_unlogged CREATE FUNCTION pgtrickle."convert_buffers_to_unlogged"() RETURNS bigint /* Result < i64, PgTrickleError > */ STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'convert_buffers_to_unlogged_wrapper'; /* */ /* */ -- src/monitor/mod.rs:1481 -- pg_trickle::monitor::check_cdc_health CREATE FUNCTION pgtrickle."check_cdc_health"() RETURNS TABLE ( "source_relid" bigint, /* i64 */ "source_table" TEXT, /* String */ "cdc_mode" TEXT, /* String */ "slot_name" TEXT, /* Option < String > */ "lag_bytes" bigint, /* Option < i64 > */ "confirmed_lsn" TEXT, /* Option < String > */ "alert" TEXT, /* Option < String > */ "selective_capture" bool /* bool */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'check_cdc_health_wrapper'; /* */ /* */ -- src/hash.rs:58 -- pg_trickle::hash::pg_trickle_hash_multi CREATE FUNCTION pgtrickle."pg_trickle_hash_multi"( "inputs" TEXT[] /* Vec < Option < String > > */ ) RETURNS bigint /* i64 */ IMMUTABLE STRICT PARALLEL SAFE LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'pg_trickle_hash_multi_wrapper'; /* */ /* */ -- src/api/refresh_ops.rs:38 -- pg_trickle::api::refresh_ops::write_and_refresh CREATE FUNCTION pgtrickle."write_and_refresh"( "sql" TEXT, /* & str */ "stream_table_name" TEXT /* & str */ ) RETURNS void STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'write_and_refresh_wrapper'; /* */ /* */ -- src/api/mod.rs:2520 -- pg_trickle::api::sla_summary CREATE FUNCTION pgtrickle."sla_summary"() RETURNS TABLE ( "stream_table" TEXT, /* Option < String > */ "p50_ms" double precision, /* Option < f64 > */ "p99_ms" double precision, /* Option < f64 > */ "freshness_lag_s" double precision, /* Option < f64 > */ "error_rate" double precision, /* Option < f64 > */ "error_budget_remaining" double precision /* Option < f64 > */ ) STRICT LANGUAGE c /* Rust */ AS 'MODULE_PATHNAME', 'sla_summary_wrapper'; /* */ /* */ -- src/lib.rs:1225 -- requires: -- pg_trickle_acl_policy ALTER FUNCTION pgtrickle.stream_table_to_publication(text) SECURITY DEFINER; -- nosemgrep: sql.security-definer.present — search_path is pinned immediately below. ALTER FUNCTION pgtrickle.stream_table_to_publication(text) SET search_path = pgtrickle, pg_catalog, pg_temp; ALTER FUNCTION pgtrickle.drop_stream_table_publication(text) SECURITY DEFINER; -- nosemgrep: sql.security-definer.present — search_path is pinned immediately below. ALTER FUNCTION pgtrickle.drop_stream_table_publication(text) SET search_path = pgtrickle, pg_catalog, pg_temp; -- src/lib.rs:1794 -- requires: -- pg_trickle_acl_policy ALTER FUNCTION pgtrickle.attach_outbox(text, integer, integer) SECURITY DEFINER; -- nosemgrep: sql.security-definer.present — external pg_tide calls run as the captured caller. ALTER FUNCTION pgtrickle.attach_outbox(text, integer, integer) SET search_path = pgtrickle, pg_catalog, pg_temp; ALTER FUNCTION pgtrickle.detach_outbox(text, boolean) SECURITY DEFINER; -- nosemgrep: sql.security-definer.present — private mapping cleanup is definer-only. ALTER FUNCTION pgtrickle.detach_outbox(text, boolean) SET search_path = pgtrickle, pg_catalog, pg_temp; ALTER FUNCTION pgtrickle.attach_embedding_outbox(text, text, integer, integer) SECURITY DEFINER; -- nosemgrep: sql.security-definer.present — external pg_tide calls run as the captured caller. ALTER FUNCTION pgtrickle.attach_embedding_outbox(text, text, integer, integer) SET search_path = pgtrickle, pg_catalog, pg_temp; -- src/lib.rs:1815 -- requires: -- pg_trickle_event_triggers -- pg_trickle_pause_resume -- pg_trickle_refresh_if_stale -- pg_trickle_stream_table_definition -- pg_trickle_canary -- pg_trickle_distance_subscriptions_catalog -- _signal_launcher_rescan -- preview_stream_table -- create_stream_table_realtime -- create_stream_table_batch -- create_stream_table_cost_optimized -- finalize -- Generated by scripts/check_sql_api_policy.py emit-acl-sql REVOKE EXECUTE ON ALL FUNCTIONS IN SCHEMA pgtrickle FROM PUBLIC; -- Explicit overload policy follows. -- admin_global REVOKE EXECUTE ON FUNCTION pgtrickle.advance_watermark(text, timestamp with time zone) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.clear_caches() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.convert_buffers_to_unlogged() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_refresh_group(text, text[], text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_watermark_group(text, text[], double precision) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.drain(integer) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.drop_refresh_group(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.drop_watermark_group(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.gate_source(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.lifecycle_preflight() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.migrate() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.pause_all() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.pause_scheduler(text[]) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.rebuild_cdc_triggers() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.restore_stream_tables() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.resume_all() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.resume_scheduler(text[]) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.resume_after_drain() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.setup_self_monitoring() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.teardown_self_monitoring() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.ungate_source(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.stat_reset_all() FROM PUBLIC; -- arbitrary_sql REVOKE EXECUTE ON FUNCTION pgtrickle.write_and_refresh(text, text) FROM PUBLIC; -- internal REVOKE EXECUTE ON FUNCTION pgtrickle._signal_launcher_rescan() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.handle_vp_promoted(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.pgt_ivm_apply_delta(bigint, integer, boolean, boolean) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.pgt_ivm_apply_delta_enr(bigint, integer, boolean, boolean) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.pgt_ivm_handle_truncate(bigint) FROM PUBLIC; -- owner_lifecycle REVOKE EXECUTE ON FUNCTION pgtrickle.alter_stream_table(text, text, text, text, text, text, text, text, boolean, boolean, text, text, bigint, integer, text, integer, double precision, text, double precision, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.attach_embedding_outbox(text, text, integer, integer) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.attach_outbox(text, integer, integer) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.bulk_alter_stream_tables(text[], json) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.bulk_create(jsonb) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.bulk_drop_stream_tables(text[]) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.canary_begin(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.canary_diff(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.canary_promote(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_or_replace_stream_table(text, text, text, text, boolean, text, text, text, boolean, boolean, text, integer, double precision, text, boolean, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_stream_table(text, text, text, text, boolean, text, text, text, boolean, boolean, text, integer, double precision, text, boolean, text, integer, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_stream_table_batch(text, text, text, boolean, text, integer, double precision) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_stream_table_cost_optimized(text, text, text, boolean, text, integer, double precision) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_stream_table_fast_append_only(text, text, text, text, text, integer, double precision) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_stream_table_if_not_exists(text, text, text, text, boolean, text, text, text, boolean, boolean, text, integer, double precision, text, boolean, text, integer, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.create_stream_table_realtime(text, text, text, boolean, text, integer, double precision) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.detach_outbox(text, boolean) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.drop_snapshot(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.drop_stream_table(text, boolean) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.drop_stream_table_publication(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.embedding_stream_table(text, text, text, text, text, text, boolean) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.exec_stream_ddl(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.pause_stream_table(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.refresh_efficiency() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.refresh_groups() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.refresh_if_stale(text, interval) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.refresh_stream_table(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.refresh_timeline(integer) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.repair_stream_table(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.reset_fuse(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.restore_from_snapshot(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.resume_stream_table(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.set_stream_table_refresh_policy(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.set_stream_table_sla(text, interval) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.set_stream_table_storage_policy(text, boolean, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.snapshot_stream_table(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.stream_table_to_publication(text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.subscribe(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.subscribe_distance(text, text, text, text, text, double precision) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.unsubscribe(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.unsubscribe_distance(text, text) FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle.stat_reset(bigint) FROM PUBLIC; -- public_read GRANT EXECUTE ON FUNCTION pgtrickle.bootstrap_gate_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.cache_stats() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.cdc_pause_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.change_buffer_sizes() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.check_cdc_health() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.cluster_worker_summary() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.commit_latency_stats() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.dedup_stats() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.dependency_tree() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.diagnose_errors(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.diamond_groups() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_dag(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_delta(text, text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_diff_sql(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_json(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_query_rewrite(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_refresh_mode(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_st(text, boolean) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.explain_stream_table(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.export_definition(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.fuse_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.get_refresh_history(text, integer) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.get_staleness(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.health_check() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.health_summary() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.history_prune_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.is_drained() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.list_auxiliary_columns(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.list_distance_subscriptions(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.list_snapshots(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.list_sources(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.list_subscriptions() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.metrics_summary() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.parallel_job_status(integer) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.parse_duration_seconds(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.encode_row_id_v2(text, anyelement) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.row_probe_v1(bytea) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.pg_trickle_hash(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.pg_trickle_hash_multi(text[]) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.pgt_scc_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.pgt_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.pgtrickle_refresh_stats() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.preflight() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.preview_stream_table(text, text, text, text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.recommend_refresh_mode(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.recommend_schedule(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.reliability_counters() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.schedule_recommendations() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.scheduler_overhead() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.self_monitoring_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.shared_buffer_stats() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.sla_summary() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.slot_health() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.source_gates() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.source_stable_name(oid) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.st_auto_threshold(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.st_refresh_stats() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.stream_table_definition(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.stream_table_lineage(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.stream_table_spec(oid) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.stream_table_spec(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.trigger_inventory() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.tune_recommendations() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.validate_query(text) TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.vector_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.version() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.version_check() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.view_evolution_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.wal_source_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.watermark_groups() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.watermark_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.watermarks() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.worker_allocation_status() TO PUBLIC; GRANT EXECUTE ON FUNCTION pgtrickle.worker_pool_status() TO PUBLIC; -- trigger_entry REVOKE EXECUTE ON FUNCTION pgtrickle._on_ddl_end() FROM PUBLIC; REVOKE EXECUTE ON FUNCTION pgtrickle._on_sql_drop() FROM PUBLIC; SELECT pgtrickle._signal_launcher_rescan(); /* */