use strict;
use warnings FATAL => 'all';
use PostgreSQL::Test::Cluster;
use PostgreSQL::Test::Utils;
use Test::More;

# Reproduces the race from "Table sync race with REFRESH PUBLICATION"
# (Ayush Tiwari, pgsql-hackers, Sep 2026): the apply worker can cache a
# table's SYNCDONE state, then block on AlterSubscription_refresh()'s
# AccessExclusiveLock on the subscription (subscriptioncmds.c) while a
# concurrent REFRESH PUBLICATION removes that table. Confirmed fix
# (v2-0001-Recheck-table-sync-state-after-refresh.patch) re-reads the
# state via GetSubscriptionRelState() under an AccessShareLock on the
# same shared object (tablesync.c, ProcessSyncingTablesForApply) and
# skips the READY transition if it's no longer SYNCDONE. Without the
# fix, the apply worker tries to update a pg_subscription_rel row that
# REFRESH already deleted, errors out, and disable_on_error disables
# the whole subscription.

my $publisher = PostgreSQL::Test::Cluster->new('publisher');
$publisher->init(allows_streaming => 'logical');
my $subscriber = PostgreSQL::Test::Cluster->new('subscriber');
$subscriber->init;
$publisher->start;
$subscriber->start;

# injection_points pauses a backend at a named point in the C code and
# resumes it on command from another session -- needed because this is
# a lock-acquisition-order race, not something reproducible just by
# timing separate psql calls against each other.
plan skip_all => 'injection points not supported by this build'
  unless $subscriber->check_extension('injection_points');
$subscriber->safe_psql('postgres', 'CREATE EXTENSION injection_points');

my $connstr = $publisher->connstr . ' dbname=postgres';
my $schema = q{
CREATE TABLE tab_first (a int PRIMARY KEY);
CREATE TABLE tab_second (a int PRIMARY KEY);
CREATE TABLE tab_sync (a int PRIMARY KEY);
};

# Matches pg_locks rows for a "shared object" lock where classid is
# pg_subscription's own OID and objid is this subscription's OID.
# Confirmed: this is exactly what both AlterSubscription_refresh()'s
# LockSharedObject(SubscriptionRelationId, subid, 0, AccessExclusiveLock)
# and ProcessSyncingTablesForApply()'s
# LockSharedObject(SubscriptionRelationId, subid, 0, AccessShareLock)
# target -- same classid/objid, different lock modes, which is what
# makes one block on the other.
my $subscription_lock = q{
locks.classid = 'pg_subscription'::regclass
AND locks.objid = (SELECT oid FROM pg_subscription WHERE subname = 'sub_sync')};

# Polls until a backend of the given type is shown WAITING (not granted)
# on a lock matching $target. This is the test's core primitive for
# proving "worker X is now genuinely stuck trying to acquire lock Y",
# used repeatedly below to pin down exact interleavings.
sub wait_for_lock
{
	my ($backend, $type, $target) = @_;
	$subscriber->poll_query_until('postgres', qq{
SELECT count(*) = 1 FROM pg_locks locks
JOIN pg_stat_activity activity USING (pid)
WHERE activity.backend_type = '$backend' AND NOT locks.granted
  AND locks.locktype = '$type' AND $target;
}) or die "$backend did not wait for $type lock";
}

# tab_sync starts OUTSIDE the publication on purpose. It's added via
# REFRESH partway through the test, so its tablesync worker launches
# fresh mid-test rather than during CREATE SUBSCRIPTION's initial sync.
$publisher->safe_psql('postgres', $schema . q{
INSERT INTO tab_sync VALUES (1);
CREATE PUBLICATION pub_sync FOR TABLE tab_first, tab_second;
});
$subscriber->safe_psql('postgres', $schema);
$subscriber->safe_psql('postgres',
"CREATE SUBSCRIPTION sub_sync CONNECTION '$connstr' PUBLICATION pub_sync WITH (disable_on_error = true)");
# disable_on_error is the test's real assertion target: it's what turns
# the unfixed bug's ERROR into an observable, poll-able state change
# (subenabled flips to false) rather than just a log line to grep for.
$subscriber->wait_for_subscription_sync($publisher, 'sub_sync');

# --- Phase 1: add tab_sync, stall its initial copy so we control timing --

# A SHARE lock on tab_sync blocks the new tablesync worker's TRUNCATE
# (which needs AccessExclusiveLock) right at the start of its initial
# copy, so we get a stable point to synchronize on before letting it run.
my $sync_blocker = $subscriber->background_psql('postgres');
$sync_blocker->query_safe('BEGIN; LOCK TABLE tab_sync IN SHARE MODE');
$publisher->safe_psql('postgres', 'ALTER PUBLICATION pub_sync ADD TABLE tab_sync');
$subscriber->safe_psql('postgres', 'ALTER SUBSCRIPTION sub_sync REFRESH PUBLICATION');
wait_for_lock('logical replication tablesync worker', 'relation',
"locks.relation = 'tab_sync'::regclass");

# --- Phase 2: stall the apply worker on two already-READY tables --------

# tab_first/tab_second are already READY, so their INSERTs are applied
# directly by the apply worker rather than by a tablesync worker.
# Taking SHARE locks on both from the subscriber side lets the test
# stall apply's DML application one table at a time in later phases.
my $first_blocker = $subscriber->background_psql('postgres');
$first_blocker->query_safe('BEGIN; LOCK TABLE tab_first IN SHARE MODE');
my $second_blocker = $subscriber->background_psql('postgres');
$second_blocker->query_safe('BEGIN; LOCK TABLE tab_second IN SHARE MODE');
$publisher->safe_psql('postgres', 'INSERT INTO tab_first VALUES (1)');
$publisher->safe_psql('postgres', 'INSERT INTO tab_second VALUES (1)');
wait_for_lock('logical replication apply worker', 'relation',
"locks.relation = 'tab_first'::regclass");
my $apply_pid = $subscriber->safe_psql('postgres', q{
SELECT pid FROM pg_stat_subscription WHERE subname = 'sub_sync' AND worker_type = 'apply';
});

# --- Phase 3: drive tab_sync's tablesync all the way to SYNCDONE --------

# Releasing tab_sync lets its tablesync worker finish COPYing, record
# srsubstate='f' (FINISHEDCOPY), then block in-memory in SYNCWAIT
# (wait_event = LogicalSyncStateChange) until the apply worker promotes
# it to CATCHUP.
$sync_blocker->query_safe('COMMIT');
$subscriber->poll_query_until('postgres', q{
SELECT count(*) = 1 FROM pg_stat_activity activity, pg_subscription_rel rel
WHERE activity.backend_type = 'logical replication tablesync worker'
  AND activity.wait_event = 'LogicalSyncStateChange'
  AND rel.srrelid = 'tab_sync'::regclass AND rel.srsubstate = 'f';
}) or die 'tablesync did not reach SYNCWAIT';

# Releasing tab_first lets apply finish/commit that INSERT, which is
# when it calls ProcessSyncingTablesForApply() and notices tab_sync
# waiting: SYNCWAIT -> CATCHUP. With no further changes queued for
# tab_sync, the tablesync worker immediately reaches SYNCDONE and
# exits. Apply then moves on to the tab_second INSERT and stalls there.
$first_blocker->query_safe('COMMIT');
wait_for_lock('logical replication apply worker', 'relation',
"locks.relation = 'tab_second'::regclass");

# Confirms the setup landed exactly as intended: tab_sync's row is
# SYNCDONE *before* REFRESH runs below, and apply is currently stuck
# elsewhere (tab_second) rather than having already consumed that
# SYNCDONE state. This is the precondition the whole race depends on --
# without it, the rest of the test wouldn't actually be exercising the
# bug.
$subscriber->safe_psql('postgres', q{
SELECT srsubstate FROM pg_subscription_rel WHERE srrelid = 'tab_sync'::regclass;
}) eq 's' or die 'table sync did not complete before refresh';

# --- Phase 4: concurrently remove tab_sync from the publication --------

$publisher->safe_psql('postgres', 'ALTER PUBLICATION pub_sync DROP TABLE tab_sync');

# Pauses REFRESH's own backend at a point inside its removal logic for
# tab_sync -- confirmed to be a PRE-EXISTING injection point (reused
# here, not added by this fix), consistent with the fix's commit
# message noting "this table-side issue predates the sequence-side fix
# in dca73f7dd03" (an analogous REFRESH SEQUENCES race). I have NOT
# independently verified the exact source line this call sits on in
# subscriptioncmds.c; I'm relying on the patch author's own framing
# ("the existing refresh injection point") plus the commit-message
# cross-reference. What IS independently confirmed: REFRESH holds its
# AccessExclusiveLock on the subscription for the full removal
# sequence, so pausing anywhere inside that removal necessarily also
# holds the lock the apply worker needs next (see $subscription_lock
# comment above).
$subscriber->safe_psql('postgres', q{
SELECT injection_points_attach('subscription-refresh-before-origin-check', 'wait');
});
my $refresh = $subscriber->background_psql('postgres');
$refresh->query_until(qr/starting_refresh/, q{
\echo starting_refresh
ALTER SUBSCRIPTION sub_sync REFRESH PUBLICATION;
});
$subscriber->wait_for_event('client backend',
'subscription-refresh-before-origin-check');

# --- Phase 5: force the apply worker into the race window --------------

# Releasing tab_second lets apply finish/commit that INSERT. A fresh
# publisher-side change then forces apply's next loop iteration to call
# ProcessSyncingTablesForApply() again -- this time reaching for the
# AccessShareLock on the subscription, which REFRESH (paused above)
# is holding as AccessExclusiveLock. The wait_for_lock call below is
# the test's actual proof that the race window is open: it confirms
# apply is genuinely blocked on the identical shared-object lock
# REFRESH holds mid-removal of tab_sync -- not just "probably" blocked
# on it by timing coincidence.
$second_blocker->query_safe('COMMIT');
$publisher->safe_psql('postgres', 'INSERT INTO tab_first VALUES (2)');
wait_for_lock('logical replication apply worker', 'object',
"locks.pid = $apply_pid AND $subscription_lock");

# --- Phase 6: release both, then check the actual outcome --------------

# Releasing REFRESH lets it finish removing tab_sync's row/origin and
# drop its AccessExclusiveLock. Apply then acquires its AccessShareLock
# and re-reads state via GetSubscriptionRelState() -- confirmed
# (v2-0001) to return SUBREL_STATE_UNKNOWN here, since the row is gone.
# The fix's "if (current_relstate != SUBREL_STATE_SYNCDONE) continue;"
# means apply logs DEBUG1 and moves on, instead of trying to write a
# READY state into a row that no longer exists.
$subscriber->safe_psql('postgres', q{
SELECT injection_points_wakeup('subscription-refresh-before-origin-check');
});
$refresh->quit or die 'refresh failed';

# One more change proves apply is still alive and functioning normally
# after the race, not merely that it avoided an immediate crash.
$publisher->safe_psql('postgres', 'INSERT INTO tab_first VALUES (3)');

# The actual A/B assertion: either replication caught up normally
# (count=3, fix works) or the subscription got disabled (bug
# reproduced). Either branch ends the poll, so this can't hang --
# the `is()` check below is what actually fails the test.
$subscriber->poll_query_until('postgres', q{
SELECT (SELECT count(*) = 3 FROM tab_first) OR NOT subenabled
FROM pg_subscription WHERE subname = 'sub_sync';
}) or die 'apply neither caught up nor disabled the subscription';

# With the fix: stays enabled. Without it: apply's ERROR plus
# disable_on_error would have flipped subenabled to false here.
is($subscriber->safe_psql('postgres', q{
SELECT subenabled FROM pg_subscription WHERE subname = 'sub_sync';
}), 't', 'refresh does not disable the subscription');

$subscriber->safe_psql('postgres', q{
SELECT injection_points_detach('subscription-refresh-before-origin-check');
DROP SUBSCRIPTION sub_sync;
DROP TABLE tab_first, tab_second, tab_sync;
});
$publisher->safe_psql('postgres', q{
DROP PUBLICATION pub_sync;
DROP TABLE tab_first, tab_second, tab_sync;
});
$sync_blocker->quit;
$first_blocker->quit;
$second_blocker->quit;

$subscriber->stop('fast');
$publisher->stop('fast');
done_testing();