From 462da94d1734a04cb55fc6f2792dfb87d8b8e9e7 Mon Sep 17 00:00:00 2001 From: Alexander Korotkov Date: Tue, 5 Mar 2024 13:35:03 +0200 Subject: [PATCH v11] Implement AFTER clause for BEGIN command The new clause is to be used on standby and specifies waiting for the specific WAL location to be replayed before starting the transaction. This option is useful when the user makes some data changes on primary and needs a guarantee to see these changes on standby. The queue of waiters is stored in the shared memory array sorted by LSN. During replay of WAL waiters whose LSNs are already replayed are deleted from the shared memory array and woken up by setting of their latches. Discussion: https://postgr.es/m/eb12f9b03851bb2583adab5df9579b4b%40postgrespro.ru Author: Kartyshov Ivan, Alexander Korotkov Reviewed-by: Michael Paquier, Peter Eisentraut, Dilip Kumar, Amit Kapila Reviewed-by: Alexander Lakhin --- doc/src/sgml/ref/begin.sgml | 44 ++- doc/src/sgml/ref/start_transaction.sgml | 5 +- src/backend/access/transam/xlogrecovery.c | 7 + src/backend/commands/Makefile | 3 +- src/backend/commands/meson.build | 1 + src/backend/commands/waitlsn.c | 327 ++++++++++++++++++++++ src/backend/storage/ipc/ipci.c | 7 + src/backend/storage/lmgr/proc.c | 6 + src/include/catalog/pg_proc.dat | 5 + src/include/commands/waitlsn.h | 24 ++ src/test/recovery/meson.build | 1 + src/test/recovery/t/043_wait_lsn.pl | 62 ++++ src/tools/pgindent/typedefs.list | 2 + 13 files changed, 491 insertions(+), 3 deletions(-) create mode 100644 src/backend/commands/waitlsn.c create mode 100644 src/include/commands/waitlsn.h create mode 100644 src/test/recovery/t/043_wait_lsn.pl diff --git a/doc/src/sgml/ref/begin.sgml b/doc/src/sgml/ref/begin.sgml index 016b0214874..4ed7002683e 100644 --- a/doc/src/sgml/ref/begin.sgml +++ b/doc/src/sgml/ref/begin.sgml @@ -21,13 +21,16 @@ PostgreSQL documentation -BEGIN [ WORK | TRANSACTION ] [ transaction_mode [, ...] ] +BEGIN [ WORK | TRANSACTION ] [ transaction_mode [, ...] ] wait_event where transaction_mode is one of: ISOLATION LEVEL { SERIALIZABLE | REPEATABLE READ | READ COMMITTED | READ UNCOMMITTED } READ WRITE | READ ONLY [ NOT ] DEFERRABLE + +where wait_event is: + AFTER lsn_value [ WITHIN number_of_milliseconds ] @@ -78,6 +81,36 @@ BEGIN [ WORK | TRANSACTION ] [ transaction_mode + + + AFTER lsn_value + + + AFTER clause is used on standby in + physical streaming replication + and specifies waiting for the specific WAL location (LSN) + to be replayed before starting the transaction. + + + This option is useful when the user makes some data changes on primary + and needs a guarantee to see these changes on standby. The LSN to wait + could be obtained on the primary using + pg_current_wal_insert_lsn + function after committing the relevant changes. + + + + + + WITHIN number_of_milliseconds + + + Provides the timeout for the AFTER clause. + Especially helpful to prevent freezing on streaming replication + connection failures. + + + @@ -123,6 +156,15 @@ BEGIN [ WORK | TRANSACTION ] [ transaction_mode BEGIN; + + + To begin a transaction block after replaying the given LSN. + The command will be canceled if the given LSN is not + reached within the timeout of one second. + +BEGIN AFTER '0/3F0FF791' WITHIN 1000; + + diff --git a/doc/src/sgml/ref/start_transaction.sgml b/doc/src/sgml/ref/start_transaction.sgml index 74ccd7e3456..46a3bcf1a80 100644 --- a/doc/src/sgml/ref/start_transaction.sgml +++ b/doc/src/sgml/ref/start_transaction.sgml @@ -21,13 +21,16 @@ PostgreSQL documentation -START TRANSACTION [ transaction_mode [, ...] ] +START TRANSACTION [ transaction_mode [, ...] ] wait_event where transaction_mode is one of: ISOLATION LEVEL { SERIALIZABLE | REPEATABLE READ | READ COMMITTED | READ UNCOMMITTED } READ WRITE | READ ONLY [ NOT ] DEFERRABLE + +where wait_event is: + AFTER lsn_value [ WITHIN number_of_milliseconds ] diff --git a/src/backend/access/transam/xlogrecovery.c b/src/backend/access/transam/xlogrecovery.c index 29c5bec0847..c211230f6b5 100644 --- a/src/backend/access/transam/xlogrecovery.c +++ b/src/backend/access/transam/xlogrecovery.c @@ -43,6 +43,7 @@ #include "backup/basebackup.h" #include "catalog/pg_control.h" #include "commands/tablespace.h" +#include "commands/waitlsn.h" #include "common/file_utils.h" #include "miscadmin.h" #include "pgstat.h" @@ -1828,6 +1829,12 @@ PerformWalRecovery(void) break; } + /* + * If we replayed an LSN that someone was waiting for, set latches + * in shared memory array to notify the waiter. + */ + WaitLSNSetLatches(XLogRecoveryCtl->lastReplayedEndRecPtr); + /* Else, try to fetch the next WAL record */ record = ReadRecord(xlogprefetcher, LOG, false, replayTLI); } while (record != NULL); diff --git a/src/backend/commands/Makefile b/src/backend/commands/Makefile index 48f7348f91c..cede90c3b98 100644 --- a/src/backend/commands/Makefile +++ b/src/backend/commands/Makefile @@ -61,6 +61,7 @@ OBJS = \ vacuum.o \ vacuumparallel.o \ variable.o \ - view.o + view.o \ + waitlsn.o include $(top_srcdir)/src/backend/common.mk diff --git a/src/backend/commands/meson.build b/src/backend/commands/meson.build index 6dd00a4abde..7549be5dc3b 100644 --- a/src/backend/commands/meson.build +++ b/src/backend/commands/meson.build @@ -50,4 +50,5 @@ backend_sources += files( 'vacuumparallel.c', 'variable.c', 'view.c', + 'waitlsn.c', ) diff --git a/src/backend/commands/waitlsn.c b/src/backend/commands/waitlsn.c new file mode 100644 index 00000000000..691e96439cc --- /dev/null +++ b/src/backend/commands/waitlsn.c @@ -0,0 +1,327 @@ +/*------------------------------------------------------------------------- + * + * waitlsn.c + * Implements waiting for the given LSN, which is used in + * BEGIN AFTER ... [ WITHIN ... ] clause. + * + * Copyright (c) 2024, PostgreSQL Global Development Group + * + * IDENTIFICATION + * src/backend/commands/waitlsn.c + * + *------------------------------------------------------------------------- + */ + +#include +#include +#include "postgres.h" +#include "pgstat.h" +#include "fmgr.h" +#include "access/transam.h" +#include "access/xact.h" +#include "access/xlog.h" +#include "access/xlogdefs.h" +#include "access/xlogrecovery.h" +#include "catalog/pg_type.h" +#include "commands/waitlsn.h" +#include "funcapi.h" +#include "miscadmin.h" +#include "storage/ipc.h" +#include "storage/latch.h" +#include "storage/pmsignal.h" +#include "storage/proc.h" +#include "storage/shmem.h" +#include "storage/spin.h" +#include "storage/sinvaladt.h" +#include "utils/builtins.h" +#include "utils/pg_lsn.h" +#include "utils/snapmgr.h" +#include "utils/timestamp.h" +#include "executor/spi.h" +#include "utils/fmgrprotos.h" + +/* Add to / delete from shared memory array */ +static void addLSNWaiter(XLogRecPtr lsn); +static void deleteLSNWaiter(void); + +/* Shared memory structure */ +typedef struct +{ + int procnum; + XLogRecPtr waitLSN; +} WaitLSNProcInfo; + +typedef struct +{ + pg_atomic_uint64 minLSN; + slock_t mutex; + int numWaitedProcs; + WaitLSNProcInfo procInfos[FLEXIBLE_ARRAY_MEMBER]; +} WaitLSNState; + +static WaitLSNState *state = NULL; +static volatile sig_atomic_t haveShmemItem = false; + +/* + * Report the amount of shared memory space needed for WaitLSNState + */ +Size +WaitLSNShmemSize(void) +{ + Size size; + + size = offsetof(WaitLSNState, procInfos); + size = add_size(size, mul_size(MaxBackends, sizeof(WaitLSNProcInfo))); + return size; +} + +/* Initialize the WaitLSNState in the shared memory */ +void +WaitLSNShmemInit(void) +{ + bool found; + + state = (WaitLSNState *) ShmemInitStruct("WaitLSNState", + WaitLSNShmemSize(), + &found); + if (!found) + { + SpinLockInit(&state->mutex); + state->numWaitedProcs = 0; + pg_atomic_init_u64(&state->minLSN, PG_UINT64_MAX); + } +} + +/* + * Add the information about the LSN waiter backend to the shared memory + * array. + */ +static void +addLSNWaiter(XLogRecPtr lsn) +{ + WaitLSNProcInfo cur; + int i; + + SpinLockAcquire(&state->mutex); + + cur.procnum = MyProcNumber; + cur.waitLSN = lsn; + + for (i = 0; i < state->numWaitedProcs; i++) + { + if (state->procInfos[i].waitLSN >= cur.waitLSN) + { + WaitLSNProcInfo tmp; + + tmp = state->procInfos[i]; + state->procInfos[i] = cur; + cur = tmp; + } + } + state->procInfos[i] = cur; + state->numWaitedProcs++; + + pg_atomic_write_u64(&state->minLSN, state->procInfos[i].waitLSN); + SpinLockRelease(&state->mutex); +} + +/* + * Delete the information about the LSN waiter backend from the shared memory + * array. + */ +static void +deleteLSNWaiter(void) +{ + int i; + bool found = false; + + SpinLockAcquire(&state->mutex); + + for (i = 0; i < state->numWaitedProcs; i++) + { + if (state->procInfos[i].procnum == MyProcNumber) + found = true; + + if (found && i < state->numWaitedProcs - 1) + { + state->procInfos[i] = state->procInfos[i + 1]; + } + } + + if (!found) + { + SpinLockRelease(&state->mutex); + return; + } + state->numWaitedProcs--; + + if (state->numWaitedProcs != 0) + pg_atomic_write_u64(&state->minLSN, state->procInfos[i].waitLSN); + else + pg_atomic_write_u64(&state->minLSN, PG_UINT64_MAX); + + SpinLockRelease(&state->mutex); +} + +/* Set all latches in shared memory to signal that new LSN has been replayed */ +void +WaitLSNSetLatches(XLogRecPtr curLSN) +{ + uint32 i, + numWakeUpProcs; + + if (!state) + return; + + /* + * Fast check if we reached the LSN at least one process is waiting for. + * This is important because saves us from acquiring the spinlock for + * every WAL record replayed. + */ + if (pg_atomic_read_u64(&state->minLSN) > curLSN) + return; + + SpinLockAcquire(&state->mutex); + + /* + * Set latches for process, whose waited LSNs are already replayed. + */ + for (i = 0; i < state->numWaitedProcs; i++) + { + PGPROC *backend; + + if (state->procInfos[i].waitLSN > curLSN) + break; + + backend = GetPGProcByNumber(state->procInfos[i].procnum); + SetLatch(&backend->procLatch); + } + + /* + * Immediately remove those processes from the shmem array. Otherwise, + * shmem array items will be here till corresponding processes wake up and + * delete themselves. + */ + numWakeUpProcs = i; + for (i = 0; i < state->numWaitedProcs - numWakeUpProcs; i++) + state->procInfos[i] = state->procInfos[i + numWakeUpProcs]; + state->numWaitedProcs -= numWakeUpProcs; + + if (state->numWaitedProcs != 0) + pg_atomic_write_u64(&state->minLSN, state->procInfos[i].waitLSN); + else + pg_atomic_write_u64(&state->minLSN, PG_UINT64_MAX); + + SpinLockRelease(&state->mutex); +} + +/* + * Delete our item from shmem array if any. + */ +void +WaitLSNCleanup(void) +{ + if (haveShmemItem) + deleteLSNWaiter(); +} + +/* + * Wait using MyLatch to wait till the given LSN is replayed, the postmaster dies or + * timeout happens. + */ +void +WaitForLSN(XLogRecPtr lsn, int millisecs) +{ + XLogRecPtr curLSN; + int latch_events; + TimestampTz endtime; + + /* Shouldn't be called when shmem isn't initialized */ + Assert(state); + + /* Should be only called by a backend */ + Assert(MyBackendType == B_BACKEND); + + if (!RecoveryInProgress()) + ereport(ERROR, + (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), + errmsg("recovery is not in progress"), + errhint("Waiting for LSN can only be executed during recovery."))); + + endtime = GetCurrentTimestamp() + millisecs * 1000; + + latch_events = WL_TIMEOUT | WL_LATCH_SET | WL_EXIT_ON_PM_DEATH; + addLSNWaiter(lsn); + haveShmemItem = true; + + for (;;) + { + int rc; + long delay_ms; + + /* Check if the waited LSN has been replayed */ + curLSN = GetXLogReplayRecPtr(NULL); + if (lsn <= curLSN) + break; + + if (millisecs > 0) + delay_ms = (endtime - GetCurrentTimestamp()) / 1000; + else + + /* If no timeout is set then wake up in 1 minute for interrupts */ + delay_ms = 60000; + + if (delay_ms <= 0) + break; + + /* + * If received an interruption from CHECK_FOR_INTERRUPTS, then delete + * the current event from array. + */ + CHECK_FOR_INTERRUPTS(); + + /* If postmaster dies, finish immediately */ + if (!PostmasterIsAlive()) + break; + + rc = WaitLatch(MyLatch, latch_events, delay_ms, + WAIT_EVENT_CLIENT_READ); + + if (rc & WL_LATCH_SET) + ResetLatch(MyLatch); + } + + if (lsn > curLSN) + { + deleteLSNWaiter(); + haveShmemItem = false; + ereport(ERROR, + (errcode(ERRCODE_QUERY_CANCELED), + errmsg("canceling waiting for LSN due to timeout"))); + } + else + { + haveShmemItem = false; + } +} + +Datum +pg_wait_lsn(PG_FUNCTION_ARGS) +{ + XLogRecPtr trg_lsn = PG_GETARG_LSN(0); + uint64_t delay = PG_GETARG_INT32(1); + CallContext *context = (CallContext *) fcinfo->context; + + if (context->atomic) + ereport(ERROR, + (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), + errmsg("pg_wait_lsn() must be only called in non-atomic context"))); + + if (ActiveSnapshotSet()) + PopActiveSnapshot(); + Assert(!ActiveSnapshotSet()); + + (void) WaitForLSN(trg_lsn, delay); + + PG_RETURN_VOID(); +} diff --git a/src/backend/storage/ipc/ipci.c b/src/backend/storage/ipc/ipci.c index 521ed5418cc..5aed90c9355 100644 --- a/src/backend/storage/ipc/ipci.c +++ b/src/backend/storage/ipc/ipci.c @@ -25,6 +25,7 @@ #include "access/xlogprefetcher.h" #include "access/xlogrecovery.h" #include "commands/async.h" +#include "commands/waitlsn.h" #include "miscadmin.h" #include "pgstat.h" #include "postmaster/autovacuum.h" @@ -152,6 +153,7 @@ CalculateShmemSize(int *num_semaphores) size = add_size(size, WaitEventExtensionShmemSize()); size = add_size(size, InjectionPointShmemSize()); size = add_size(size, SlotSyncShmemSize()); + size = add_size(size, WaitLSNShmemSize()); #ifdef EXEC_BACKEND size = add_size(size, ShmemBackendArraySize()); #endif @@ -244,6 +246,11 @@ CreateSharedMemoryAndSemaphores(void) /* Initialize subsystems */ CreateOrAttachShmemStructs(); + /* + * Init array of Latches in shared memory for wait lsn + */ + WaitLSNShmemInit(); + #ifdef EXEC_BACKEND /* diff --git a/src/backend/storage/lmgr/proc.c b/src/backend/storage/lmgr/proc.c index 162b1f919db..4b830dc3c85 100644 --- a/src/backend/storage/lmgr/proc.c +++ b/src/backend/storage/lmgr/proc.c @@ -36,6 +36,7 @@ #include "access/transam.h" #include "access/twophase.h" #include "access/xlogutils.h" +#include "commands/waitlsn.h" #include "miscadmin.h" #include "pgstat.h" #include "postmaster/autovacuum.h" @@ -862,6 +863,11 @@ ProcKill(int code, Datum arg) */ LWLockReleaseAll(); + /* + * Cleanup waiting for LSN if any. + */ + WaitLSNCleanup(); + /* Cancel any pending condition variable sleep, too */ ConditionVariableCancelSleep(); diff --git a/src/include/catalog/pg_proc.dat b/src/include/catalog/pg_proc.dat index 700f7daf7b2..4cd85da1f9b 100644 --- a/src/include/catalog/pg_proc.dat +++ b/src/include/catalog/pg_proc.dat @@ -12128,6 +12128,11 @@ prorettype => 'bytea', proargtypes => 'pg_brin_minmax_multi_summary', prosrc => 'brin_minmax_multi_summary_send' }, +{ oid => '16387', descr => 'wait for LSN until timeout', + proname => 'pg_wait_lsn', prokind => 'p', prorettype => 'void', + proargtypes => 'pg_lsn int8', proargnames => '{trg_lsn,delay}', + prosrc => 'pg_wait_lsn' }, + { oid => '6291', descr => 'arbitrary value from among input values', proname => 'any_value', prokind => 'a', proisstrict => 'f', prorettype => 'anyelement', proargtypes => 'anyelement', diff --git a/src/include/commands/waitlsn.h b/src/include/commands/waitlsn.h new file mode 100644 index 00000000000..97f50bc4c77 --- /dev/null +++ b/src/include/commands/waitlsn.h @@ -0,0 +1,24 @@ +/*------------------------------------------------------------------------- + * + * waitlsn.h + * Declarations for LSN waiting routines. + * + * Copyright (c) 2024, PostgreSQL Global Development Group + * + * src/include/commands/waitlsn.h + * + *------------------------------------------------------------------------- + */ +#ifndef WAIT_LSN_H +#define WAIT_LSN_H + +#include "postgres.h" +#include "tcop/dest.h" + +extern void WaitForLSN(XLogRecPtr lsn, int millisecs); +extern Size WaitLSNShmemSize(void); +extern void WaitLSNShmemInit(void); +extern void WaitLSNSetLatches(XLogRecPtr curLSN); +extern void WaitLSNCleanup(void); + +#endif /* WAIT_LSN_H */ diff --git a/src/test/recovery/meson.build b/src/test/recovery/meson.build index b1eb77b1ec1..bc47c93902c 100644 --- a/src/test/recovery/meson.build +++ b/src/test/recovery/meson.build @@ -51,6 +51,7 @@ tests += { 't/040_standby_failover_slots_sync.pl', 't/041_checkpoint_at_promote.pl', 't/042_low_level_backup.pl', + 't/043_wait_lsn.pl', ], }, } diff --git a/src/test/recovery/t/043_wait_lsn.pl b/src/test/recovery/t/043_wait_lsn.pl new file mode 100644 index 00000000000..04b45eefec4 --- /dev/null +++ b/src/test/recovery/t/043_wait_lsn.pl @@ -0,0 +1,62 @@ +# Checks waiting for lsn on standby AFTER +use strict; +use warnings; + +use PostgreSQL::Test::Cluster; +use PostgreSQL::Test::Utils; +use Test::More; + +# Initialize primary node +my $node_primary = PostgreSQL::Test::Cluster->new('primary'); +$node_primary->init(allows_streaming => 1); +$node_primary->start; + +# And some content and take a backup +$node_primary->safe_psql('postgres', + "CREATE TABLE wait_test AS SELECT generate_series(1,10) AS a"); +my $backup_name = 'my_backup'; +$node_primary->backup($backup_name); + +# Create a streaming standby with a 1 second delay from the backup +my $node_standby = PostgreSQL::Test::Cluster->new('standby'); +my $delay = 1; +$node_standby->init_from_backup($node_primary, $backup_name, + has_streaming => 1); +$node_standby->append_conf( + 'postgresql.conf', qq[ + recovery_min_apply_delay = '${delay}s' +]); +$node_standby->start; + + +# Make sure that AFTER works: add new content to primary and memorize +# primary's new LSN, then wait for primary's LSN on standby. Prove that AFTER is +# able to setup an infinite waiting loop and exit it if given no wait timeout. +$node_primary->safe_psql('postgres', + "INSERT INTO wait_test VALUES (generate_series(11, 20))"); +my $lsn1 = + $node_primary->safe_psql('postgres', "SELECT pg_current_wal_lsn()"); +my $output = $node_standby->safe_psql( + 'postgres', qq[ + CALL pg_wait_lsn('${lsn1}', 1000000); + SELECT pg_lsn_cmp(pg_last_wal_replay_lsn(), '${lsn1}'::pg_lsn); +]); + +# Get the current LSN on standby and make sure it's the same as primary's LSN +ok($output eq 0, "standby reached the same LSN as primary AFTER"); + +my $lsn2 = + $node_primary->safe_psql('postgres', "SELECT pg_current_wal_lsn() + 1"); +my $stderr; +$node_standby->safe_psql('postgres', "CALL pg_wait_lsn('${lsn1}', 1);"); +$node_standby->psql( + 'postgres', + "CALL pg_wait_lsn('${lsn2}', 1);", + stderr => \$stderr); +ok( $stderr =~ /canceling waiting for LSN due to timeout/, + "get timeout on waiting for unreachable LSN"); + + +$node_standby->stop; +$node_primary->stop; +done_testing(); diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list index aa7a25b8f8c..7a2ac178bc3 100644 --- a/src/tools/pgindent/typedefs.list +++ b/src/tools/pgindent/typedefs.list @@ -3035,6 +3035,8 @@ WaitEventIO WaitEventIPC WaitEventSet WaitEventTimeout +WaitLSNProcInfo +WaitLSNState WaitPMResult WalCloseMethod WalCompression -- 2.39.3 (Apple Git-145)