# Reproduces the FSM race condition from https://github.com/paradedb/paradedb/issues/4935, # fixed in https://github.com/paradedb/paradedb/pull/5067. Without the fix this suite # crashes the subscriber within minutes (SIGSEGV or SegmentMetaEntryHeader: UnexpectedEnd); # with the fix it runs indefinitely. [[server]] name = "Publisher" [server.style.Automatic] postgresql_conf = "Publisher" [server.setup] sql = """ DROP TABLE IF EXISTS test CASCADE; CREATE TABLE test ( id SERIAL8 NOT NULL PRIMARY KEY, message TEXT, old_message TEXT, unique_id BIGINT UNIQUE ) WITH (autovacuum_enabled = false); DROP SEQUENCE IF EXISTS test_unique_id_seq; CREATE SEQUENCE test_unique_id_seq START 1; DROP PUBLICATION IF EXISTS stressgres_pub; CREATE PUBLICATION stressgres_pub FOR ALL TABLES; INSERT INTO test (message, unique_id) SELECT (ARRAY['beer wine cheese', 'beer wine', 'beer cheese', 'beer', 'wine cheese', 'wine', 'cheese', 'bread butter'])[1 + (i % 8)] || ' ' || i::text, nextval('test_unique_id_seq') FROM generate_series(1, 100) AS s(i); """ [server.teardown] sql = "" [server.monitor] refresh_ms = 10 log_columns = ["replication_lag:MB"] sql = """ SELECT pid, pg_wal_lsn_diff(sent_lsn, replay_lsn) AS replication_lag, application_name::text, state::text FROM pg_stat_replication; """ [[server]] name = "Subscriber" [server.style.Automatic] postgresql_conf = "Subscriber" [server.setup] sql = """ DROP TABLE IF EXISTS test CASCADE; CREATE EXTENSION IF NOT EXISTS pg_search CASCADE; CREATE TABLE test ( id SERIAL8 NOT NULL PRIMARY KEY, message TEXT, old_message TEXT, unique_id BIGINT UNIQUE ) WITH ( autovacuum_enabled = true, autovacuum_vacuum_scale_factor = 0, autovacuum_vacuum_threshold = 50, autovacuum_vacuum_insert_threshold = 50 ); DROP SUBSCRIPTION IF EXISTS stressgres_sub; CREATE SUBSCRIPTION stressgres_sub CONNECTION '@Publisher_CONNSTR@' PUBLICATION stressgres_pub; SELECT pg_sleep(5); CREATE INDEX idxtest ON test USING paradedb (id, message) WITH (layer_sizes = '10kb, 100kb, 1mb, 100mb'); CREATE OR REPLACE FUNCTION assert(a bigint, b bigint) RETURNS bool LANGUAGE plpgsql AS $$ BEGIN IF a <> b THEN RAISE EXCEPTION 'Assertion failed: % <> %', a, b; END IF; RETURN true; END; $$; CREATE OR REPLACE FUNCTION assert_plan_contains( p_query text, p_expected_text text ) RETURNS boolean AS $$ DECLARE plan_line text; full_plan text := ''; BEGIN FOR plan_line IN EXECUTE 'EXPLAIN (VERBOSE) ' || p_query LOOP full_plan := full_plan || plan_line || chr(10); IF plan_line ILIKE '%' || p_expected_text || '%' THEN RETURN true; END IF; END LOOP; RAISE EXCEPTION 'Plan assertion failed: expected "%" not found in plan for query: %. Actual plan:%', p_expected_text, p_query, chr(10) || full_plan; END; $$ LANGUAGE plpgsql; ALTER SYSTEM SET autovacuum_naptime TO '1s'; SELECT pg_reload_conf(); """ [server.teardown] sql = "" [server.monitor] refresh_ms = 10 title = "Index Info Monitor" destination = "Subscriber" sql = """ SELECT segno, visible, recyclable, xmax, num_docs, num_deleted, byte_size FROM paradedb.index_info('idxtest', true) ORDER BY byte_size DESC; """ [[jobs]] refresh_ms = 5 title = "Index Size Info" destination = "Subscriber" log_columns = ["pages", "relation_size:MB", "segment_count"] log_tps = false sql = """ SELECT count(*) FILTER (WHERE visible) AS visible, count(*) FILTER (WHERE recyclable) AS recyclable, count(*) AS segment_count, pg_relation_size('idxtest') / 8192 AS pages, pg_relation_size('idxtest') AS relation_size, pg_size_pretty(pg_relation_size('idxtest')) FROM paradedb.index_info('idxtest', true); """ destinations = ["Subscriber"] [[jobs]] refresh_ms = 5 title = "Unordered Top K Base Scan" on_connect = """ SET max_parallel_workers_per_gather = 0; SELECT assert_plan_contains( $q$ SELECT id, message, old_message FROM test WHERE message &&& 'beer wine' LIMIT 50 $q$, 'TopKScanExecState' ); """ sql = """ SELECT id, message, old_message FROM test WHERE message &&& 'beer wine' LIMIT 50; """ destinations = ["Subscriber"] [[jobs]] refresh_ms = 5 title = "Normal Base Scan" on_connect = """ SET max_parallel_workers_per_gather = 0; SELECT assert_plan_contains( $q$ SELECT id, message, old_message FROM test WHERE message &&& 'beer wine' $q$, 'NormalScanExecState' ); """ sql = """ SELECT id, message, old_message FROM test WHERE message &&& 'beer wine'; """ destinations = ["Subscriber"] [[jobs]] refresh_ms = 5 title = "Aggregate Scan" on_connect = """ SELECT assert_plan_contains( $q$ SELECT count(*) FROM test WHERE message ||| 'beer' $q$, 'ParadeDB Aggregate Scan' ); """ sql = """ SELECT assert(count(*), 51), count(*) FROM test WHERE message ||| 'beer'; """ destinations = ["Subscriber"] [[jobs]] refresh_ms = 5 title = "Postgres Index Scan Fallback" on_connect = """ SET max_parallel_workers = 0; SET paradedb.enable_custom_scan = off; SET enable_seqscan = off; SELECT assert_plan_contains( $q$ SELECT id, message, old_message FROM test WHERE message ||| 'bread' LIMIT 25 $q$, 'Index Scan' ); """ sql = """ SELECT id, message, old_message FROM test WHERE message ||| 'bread' LIMIT 25; """ destinations = ["Subscriber"] [[jobs]] refresh_ms = 50 log_tps = false title = "Update 1..50" sql = """ UPDATE test SET message = message || ' ' || txid_current(), old_message = message WHERE id <= 50; """ destinations = ["Publisher"] [[jobs]] refresh_ms = 50 log_tps = false title = "Update 51..100" sql = """ UPDATE test SET message = message || ' ' || txid_current(), old_message = message WHERE id > 50 AND id <= 100; """ destinations = ["Publisher"] [[jobs]] refresh_ms = 75 log_tps = false title = "Insert value" sql = """ INSERT INTO test (message, unique_id) VALUES ('new ' || txid_current(), nextval('test_unique_id_seq')); """ destinations = ["Publisher"] [[jobs]] refresh_ms = 100 log_tps = false pause_keycode = 'd' cancel_keycode = 'D' title = "Delete values" sql = """ DELETE FROM test WHERE id > 200; """ destinations = ["Publisher"]