From dc3fbd9661baaf99f68c5e92b3a8f236edd3e73c Mon Sep 17 00:00:00 2001 From: Vignesh C Date: Wed, 7 Oct 2026 12:14:38 +0530 Subject: [PATCH] Test to show a slow worker brings down logical replication. Test to show a slow worker brings down logical replication. --- src/backend/replication/logical/worker.c | 2 + .../t/039_parallel_apply_send_timeout.pl | 151 ++++++++++++++++++ 2 files changed, 153 insertions(+) create mode 100644 src/test/subscription/t/039_parallel_apply_send_timeout.pl diff --git a/src/backend/replication/logical/worker.c b/src/backend/replication/logical/worker.c index 08f9b1b514b..5e9fb44e475 100644 --- a/src/backend/replication/logical/worker.c +++ b/src/backend/replication/logical/worker.c @@ -2789,6 +2789,8 @@ apply_handle_begin(StringInfo s) break; case TRANS_PARALLEL_APPLY: + /* For testing worker death before it is tracked as STARTED. */ + INJECTION_POINT("parallel-worker-before-xact-start", NULL); /* Hold the lock until the end of the transaction. */ pa_lock_transaction(MyParallelShared->xid, AccessExclusiveLock); pa_set_xact_state(MyParallelShared, PARALLEL_TRANS_STARTED); diff --git a/src/test/subscription/t/039_parallel_apply_send_timeout.pl b/src/test/subscription/t/039_parallel_apply_send_timeout.pl new file mode 100644 index 00000000000..16ed634cbf0 --- /dev/null +++ b/src/test/subscription/t/039_parallel_apply_send_timeout.pl @@ -0,0 +1,151 @@ + +# Copyright (c) 2026, PostgreSQL Global Development Group + +# Test that a parallel apply worker which is simply slow to start draining +# its input queue -- not crashed, not erroring, just not yet scheduled or +# stuck briefly on something else -- does not make the LEADER itself fail +# with a hard error. This currently fails. +# +# pa_send_data() in applyparallelworker.c retries a non-blocking +# shm_mq_send() for close to 9 seconds (SHM_SEND_TIMEOUT_MS) before giving +# up and returning false. Every call site in worker.c that checks this +# return value responds with an unconditional +# ereport(ERROR, ..., "could not send data to the logical replication +# parallel apply worker"), right next to a comment reading "TODO: Support +# switching to PARTIAL_SERIALIZE mode when the send buffer becomes full" -- +# a fallback that demonstrably does not exist for this path, even though +# PARTIAL_SERIALIZE mode itself is already implemented and used elsewhere in +# this same file for a different trigger. +# +# A parallel worker that is legitimately waiting on a dependency (this +# patch set's own core mechanism) for more than a few seconds is exactly +# the kind of thing that stops it draining its queue. This test does not +# go that far; it demonstrates the simpler, sufficient precondition +# directly: freeze a parallel worker right after it is assigned a +# transaction (before it drains anything beyond the BEGIN message) and push +# more than one queue's worth (16MB) of change data at it in a single +# transaction. The leader should fall back to serializing the overflow to +# disk, not hard-error within a bounded window. +# +# The subscription is also created with disable_on_error, to show just how +# far this reaches: because there is nothing to distinguish this transient, +# self-inflicted send timeout from an actual, unrecoverable data problem, +# the standard advice for operators ("set disable_on_error so a bad apply +# doesn't retry forever") would end up disabling the whole subscription +# over a momentarily-slow worker. +use strict; +use warnings FATAL => 'all'; +use PostgreSQL::Test::Cluster; +use PostgreSQL::Test::Utils; +use Test::More; + +if ($ENV{enable_injection_points} ne 'yes') +{ + plan skip_all => 'Injection points not supported by this build'; +} + +my $node_publisher = PostgreSQL::Test::Cluster->new('publisher'); +$node_publisher->init(allows_streaming => 'logical'); +$node_publisher->start; + +my $node_subscriber = PostgreSQL::Test::Cluster->new('subscriber'); +$node_subscriber->init; +$node_subscriber->append_conf('postgresql.conf', "log_min_messages = debug1"); +$node_subscriber->start; +$node_subscriber->safe_psql('postgres', 'CREATE EXTENSION injection_points'); + +foreach my $node ($node_publisher, $node_subscriber) +{ + $node->safe_psql('postgres', 'CREATE TABLE tab_big (a int PRIMARY KEY, data text)'); +} +$node_publisher->safe_psql('postgres', + "INSERT INTO tab_big SELECT i, 'x' FROM generate_series(1, 500) i"); + +$node_publisher->safe_psql('postgres', 'CREATE PUBLICATION pub FOR TABLE tab_big'); + +my $publisher_connstr = $node_publisher->connstr . ' dbname=postgres'; +$node_subscriber->safe_psql('postgres', + "CREATE SUBSCRIPTION sub CONNECTION '$publisher_connstr' PUBLICATION pub WITH (disable_on_error = true)" +); +$node_subscriber->wait_for_subscription_sync($node_publisher, 'sub'); + +# Freeze the next parallel apply worker right after it is assigned a +# transaction, before it drains anything beyond the BEGIN message. +$node_subscriber->safe_psql('postgres', + "SELECT injection_points_attach('parallel-worker-before-xact-start', 'wait')" +); + +# One non-streamed transaction, updating all 500 rows to carry a 40kB +# payload each (20MB total): comfortably more than the 16MB DSM_QUEUE_SIZE +# the leader is sending into, and all of it goes to the single worker that +# this whole transaction is dispatched to. +my $log_offset = -s $node_subscriber->logfile; +$node_publisher->safe_psql('postgres', + "UPDATE tab_big SET data = repeat('x', 40000)"); + +$node_subscriber->wait_for_event('logical replication parallel worker', + 'parallel-worker-before-xact-start'); + +# The leader keeps trying to push the remaining changes into the frozen +# worker's queue. pa_send_data() retries for close to 9 seconds +# (SHM_SEND_TIMEOUT_MS) before giving up and erroring; wait comfortably +# longer than that. If a PARTIAL_SERIALIZE fallback existed for this path, +# the leader would instead quietly serialize the overflow to disk and keep +# going, with nothing to show for it in the log and the subscription still +# enabled. +sleep(12); + +my $log_contents = slurp_file($node_subscriber->logfile, $log_offset); +unlike( + $log_contents, + qr/could not send data to the logical replication parallel apply worker/, + 'a merely-frozen (not crashed, not erroring) parallel worker that is ' + . 'slow to drain its queue must not make the leader hard-error; it ' + . 'should fall back to PARTIAL_SERIALIZE mode instead, as the code\'s ' + . 'own TODO comments say it should' +); + +# With disable_on_error, start_apply()'s PG_CATCH block (worker.c) would +# treat a real apply error here the same as any other: it would call +# DisableSubscriptionAndExit(), disabling the subscription +# (pg_subscription.subenabled = false) rather than letting the launcher +# just restart the worker to try again. A momentarily-slow parallel worker +# is not an unrecoverable data problem, so the subscription must still be +# enabled. +my $subenabled = $node_subscriber->safe_psql('postgres', + "SELECT subenabled FROM pg_subscription WHERE subname = 'sub'"); +is( $subenabled, 't', + 'disable_on_error must not have disabled the subscription over a ' + . 'merely-slow parallel worker' +); + +if ($log_contents =~ qr/could not send data to the logical replication parallel apply worker/) +{ + # The bug reproduced: the leader's exit already SIGTERMed the frozen + # worker (logicalrep_worker_detach() tears down every parallel worker + # for the subscription whenever the leader exits) and, with + # disable_on_error, nothing will restart it, so there's nothing left + # to wake up and nothing further to check. + $node_subscriber->safe_psql('postgres', + "SELECT injection_points_detach('parallel-worker-before-xact-start')"); +} +else +{ + # No error: the worker should still be genuinely frozen and waiting. + # Release it and confirm the update eventually reaches the subscriber + # correctly. + $node_subscriber->safe_psql('postgres', + "SELECT injection_points_wakeup('parallel-worker-before-xact-start')"); + $node_subscriber->safe_psql('postgres', + "SELECT injection_points_detach('parallel-worker-before-xact-start')"); + $node_publisher->wait_for_catchup('sub'); + + my $result = $node_subscriber->safe_psql('postgres', + "SELECT count(*), count(data = repeat('x', 40000)) FROM tab_big"); + is($result, '500|500', 'the update was eventually applied correctly'); +} + +$node_subscriber->stop('immediate'); +$node_publisher->stop('immediate'); + +done_testing(); -- 2.55.0