Branch data Line data Source code
1 : : /*-------------------------------------------------------------------------
2 : : * sequencesync.c
3 : : * PostgreSQL logical replication: sequence synchronization
4 : : *
5 : : * Copyright (c) 2025-2026, PostgreSQL Global Development Group
6 : : *
7 : : * IDENTIFICATION
8 : : * src/backend/replication/logical/sequencesync.c
9 : : *
10 : : * NOTES
11 : : * This file contains code for sequence synchronization for
12 : : * logical replication.
13 : : *
14 : : * Sequences requiring synchronization are tracked in the pg_subscription_rel
15 : : * catalog.
16 : : *
17 : : * Sequences to be synchronized will be added with state INIT when either of
18 : : * the following commands is executed:
19 : : * CREATE SUBSCRIPTION
20 : : * ALTER SUBSCRIPTION ... REFRESH PUBLICATION
21 : : *
22 : : * Executing the following command resets all sequences in the subscription to
23 : : * state INIT, triggering re-synchronization:
24 : : * ALTER SUBSCRIPTION ... REFRESH SEQUENCES
25 : : *
26 : : * The apply worker periodically scans pg_subscription_rel for sequences in
27 : : * INIT state. When such sequences are found, it spawns a sequencesync worker
28 : : * to handle synchronization.
29 : : *
30 : : * A single sequencesync worker is responsible for synchronizing all sequences.
31 : : * It begins by retrieving the list of sequences that are flagged for
32 : : * synchronization, i.e., those in the INIT state. These sequences are then
33 : : * processed in batches, allowing multiple entries to be synchronized within a
34 : : * single transaction. The worker fetches the current sequence values and page
35 : : * LSNs from the remote publisher, updates the corresponding sequences on the
36 : : * local subscriber, and finally marks each sequence as READY upon successful
37 : : * synchronization.
38 : : *
39 : : * Sequence state transitions follow this pattern:
40 : : * INIT -> READY
41 : : *
42 : : * To avoid creating too many transactions, up to MAX_SEQUENCES_SYNC_PER_BATCH
43 : : * sequences are synchronized per transaction. The locks on the sequence
44 : : * relation will be periodically released at each transaction commit.
45 : : *
46 : : * XXX: We didn't choose launcher process to maintain the launch of sequencesync
47 : : * worker as it didn't have database connection to access the sequences from the
48 : : * pg_subscription_rel system catalog that need to be synchronized.
49 : : *-------------------------------------------------------------------------
50 : : */
51 : :
52 : : #include "postgres.h"
53 : :
54 : : #include "access/genam.h"
55 : : #include "access/table.h"
56 : : #include "catalog/pg_sequence.h"
57 : : #include "catalog/pg_subscription_rel.h"
58 : : #include "commands/sequence.h"
59 : : #include "pgstat.h"
60 : : #include "postmaster/interrupt.h"
61 : : #include "replication/logicalworker.h"
62 : : #include "replication/worker_internal.h"
63 : : #include "storage/lmgr.h"
64 : : #include "storage/lwlock.h"
65 : : #include "utils/acl.h"
66 : : #include "utils/builtins.h"
67 : : #include "utils/fmgroids.h"
68 : : #include "utils/guc.h"
69 : : #include "utils/inval.h"
70 : : #include "utils/lsyscache.h"
71 : : #include "utils/memutils.h"
72 : : #include "utils/pg_lsn.h"
73 : : #include "utils/syscache.h"
74 : : #include "utils/usercontext.h"
75 : :
76 : : #define REMOTE_SEQ_COL_COUNT 11
77 : :
78 : : typedef enum CopySeqResult
79 : : {
80 : : COPYSEQ_SUCCESS,
81 : : COPYSEQ_MISMATCH,
82 : : COPYSEQ_SUBSCRIBER_INSUFFICIENT_PERM,
83 : : COPYSEQ_PUBLISHER_INSUFFICIENT_PERM,
84 : : COPYSEQ_SKIPPED,
85 : : COPYSEQ_NOT_SUBSCRIBED
86 : : } CopySeqResult;
87 : :
88 : : static List *seqinfos = NIL;
89 : :
90 : : /*
91 : : * Apply worker determines if sequence synchronization is needed.
92 : : *
93 : : * Start a sequencesync worker if one is not already running. The active
94 : : * sequencesync worker will handle all pending sequence synchronization. If any
95 : : * sequences remain unsynchronized after it exits, a new worker can be started
96 : : * in the next iteration.
97 : : */
98 : : void
99 : 8358 : ProcessSequencesForSync(void)
100 : : {
101 : : LogicalRepWorker *sequencesync_worker;
102 : : int nsyncworkers;
103 : : bool has_pending_sequences;
104 : : bool started_tx;
105 : :
106 : 8358 : FetchRelationStates(NULL, &has_pending_sequences, &started_tx);
107 : :
108 [ + + ]: 8358 : if (started_tx)
109 : : {
110 : 196 : CommitTransactionCommand();
111 : 196 : pgstat_report_stat(true);
112 : : }
113 : :
114 [ + + ]: 8358 : if (!has_pending_sequences)
115 : 8331 : return;
116 : :
117 : 48 : LWLockAcquire(LogicalRepWorkerLock, LW_SHARED);
118 : :
119 : : /* Check if there is a sequencesync worker already running? */
120 : 48 : sequencesync_worker = logicalrep_worker_find(WORKERTYPE_SEQUENCESYNC,
121 : 48 : MyLogicalRepWorker->subid,
122 : : InvalidOid, true);
123 [ + + ]: 48 : if (sequencesync_worker)
124 : : {
125 : 21 : LWLockRelease(LogicalRepWorkerLock);
126 : 21 : return;
127 : : }
128 : :
129 : : /*
130 : : * Count running sync workers for this subscription, while we have the
131 : : * lock.
132 : : */
133 : 27 : nsyncworkers = logicalrep_sync_worker_count(MyLogicalRepWorker->subid);
134 : 27 : LWLockRelease(LogicalRepWorkerLock);
135 : :
136 : : /*
137 : : * It is okay to read/update last_seqsync_start_time here in apply worker
138 : : * as we have already ensured that sync worker doesn't exist.
139 : : */
140 : 27 : launch_sync_worker(WORKERTYPE_SEQUENCESYNC, nsyncworkers, InvalidOid,
141 : 27 : &MyLogicalRepWorker->last_seqsync_start_time);
142 : : }
143 : :
144 : : /*
145 : : * get_sequences_string
146 : : *
147 : : * Build a comma-separated string of schema-qualified sequence names
148 : : * for the given list of sequence indexes.
149 : : */
150 : : static void
151 : 8 : get_sequences_string(List *seqindexes, StringInfo buf)
152 : : {
153 : 8 : resetStringInfo(buf);
154 [ + - + + : 24 : foreach_int(seqidx, seqindexes)
+ + ]
155 : : {
156 : : LogicalRepSequenceInfo *seqinfo =
157 : 8 : (LogicalRepSequenceInfo *) list_nth(seqinfos, seqidx);
158 : :
159 [ - + ]: 8 : if (buf->len > 0)
160 : 0 : appendStringInfoString(buf, ", ");
161 : :
162 : 8 : appendStringInfo(buf, "\"%s.%s\"", seqinfo->nspname, seqinfo->seqname);
163 : : }
164 : 8 : }
165 : :
166 : : /*
167 : : * report_sequence_errors
168 : : *
169 : : * Report discrepancies found during sequence synchronization between
170 : : * the publisher and subscriber. Emits warnings for:
171 : : * a) mismatched definitions or concurrent rename
172 : : * b) insufficient privileges on the subscriber
173 : : * c) insufficient privileges on the publisher
174 : : * d) missing sequences on the publisher
175 : : * Then raises an ERROR to indicate synchronization failure.
176 : : */
177 : : static void
178 : 16 : report_sequence_errors(List *mismatched_seqs_idx,
179 : : List *sub_insuffperm_seqs_idx,
180 : : List *pub_insuffperm_seqs_idx,
181 : : List *missing_seqs_idx)
182 : : {
183 : : StringInfoData seqstr;
184 : :
185 : : /* Quick exit if there are no errors to report */
186 [ + + + - : 16 : if (!mismatched_seqs_idx && !sub_insuffperm_seqs_idx &&
+ + ]
187 [ + + ]: 12 : !pub_insuffperm_seqs_idx && !missing_seqs_idx)
188 : 8 : return;
189 : :
190 : 8 : initStringInfo(&seqstr);
191 : :
192 [ + + ]: 8 : if (mismatched_seqs_idx)
193 : : {
194 : 3 : get_sequences_string(mismatched_seqs_idx, &seqstr);
195 [ + - ]: 3 : ereport(WARNING,
196 : : errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
197 : : errmsg_plural("mismatched or renamed sequence on subscriber (%s)",
198 : : "mismatched or renamed sequences on subscriber (%s)",
199 : : list_length(mismatched_seqs_idx),
200 : : seqstr.data));
201 : : }
202 : :
203 [ - + ]: 8 : if (sub_insuffperm_seqs_idx)
204 : : {
205 : 0 : get_sequences_string(sub_insuffperm_seqs_idx, &seqstr);
206 : :
207 : : /*
208 : : * With run_as_owner enabled, sequence synchronization runs as the
209 : : * subscription owner, so a missing UPDATE privilege should be granted
210 : : * to that role. Otherwise, the worker switches to the sequence owner
211 : : * before checking privileges, so no useful GRANT hint can be
212 : : * provided.
213 : : */
214 [ # # # # ]: 0 : ereport(WARNING,
215 : : errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
216 : : errmsg_plural("insufficient privileges on subscriber sequence (%s)",
217 : : "insufficient privileges on subscriber sequences (%s)",
218 : : list_length(sub_insuffperm_seqs_idx),
219 : : seqstr.data),
220 : : MySubscription->runasowner ?
221 : : errhint_plural("Grant UPDATE on the sequence to the subscription "
222 : : "owner on the subscriber.",
223 : : "Grant UPDATE on the sequences to the subscription "
224 : : "owner on the subscriber.",
225 : : list_length(sub_insuffperm_seqs_idx)) : 0);
226 : : }
227 : :
228 [ + + ]: 8 : if (pub_insuffperm_seqs_idx)
229 : : {
230 : 1 : get_sequences_string(pub_insuffperm_seqs_idx, &seqstr);
231 [ + - ]: 1 : ereport(WARNING,
232 : : errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
233 : : errmsg_plural("insufficient privileges on publisher sequence (%s)",
234 : : "insufficient privileges on publisher sequences (%s)",
235 : : list_length(pub_insuffperm_seqs_idx),
236 : : seqstr.data),
237 : : errhint_plural("Grant SELECT on the sequence to the role used for "
238 : : "the replication connection on the publisher.",
239 : : "Grant SELECT on the sequences to the role used for "
240 : : "the replication connection on the publisher.",
241 : : list_length(pub_insuffperm_seqs_idx)));
242 : : }
243 : :
244 [ + + ]: 8 : if (missing_seqs_idx)
245 : : {
246 : 4 : get_sequences_string(missing_seqs_idx, &seqstr);
247 [ + - ]: 4 : ereport(WARNING,
248 : : errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
249 : : errmsg_plural("missing sequence on publisher (%s)",
250 : : "missing sequences on publisher (%s)",
251 : : list_length(missing_seqs_idx),
252 : : seqstr.data));
253 : : }
254 : :
255 [ + - ]: 8 : ereport(ERROR,
256 : : errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
257 : : errmsg("logical replication sequence synchronization failed for subscription \"%s\"",
258 : : MySubscription->name));
259 : : }
260 : :
261 : : /*
262 : : * get_and_validate_seq_info
263 : : *
264 : : * Extracts remote sequence information from the tuple slot received from the
265 : : * publisher, and validates it against the corresponding local sequence
266 : : * definition.
267 : : */
268 : : static CopySeqResult
269 : 31 : get_and_validate_seq_info(TupleTableSlot *slot, Relation *sequence_rel,
270 : : LogicalRepSequenceInfo **seqinfo, int *seqidx)
271 : : {
272 : : bool isnull;
273 : 31 : int col = 0;
274 : : Datum datum;
275 : : bool remote_has_select_priv;
276 : : Oid remote_typid;
277 : : int64 remote_start;
278 : : int64 remote_increment;
279 : : int64 remote_min;
280 : : int64 remote_max;
281 : : bool remote_cycle;
282 : 31 : CopySeqResult result = COPYSEQ_SUCCESS;
283 : : HeapTuple tup;
284 : : Form_pg_sequence local_seq;
285 : : LogicalRepSequenceInfo *seqinfo_local;
286 : :
287 : 31 : *seqidx = DatumGetInt32(slot_getattr(slot, ++col, &isnull));
288 : : Assert(!isnull);
289 : :
290 : : /*
291 : : * The publisher only echoes back an index that we put in the VALUES list,
292 : : * so this should always identify an entry of seqinfos. Check it anyway
293 : : * before using it as a list subscript, since list_nth() does not
294 : : * bounds-check on non-assert builds and we would then write the remote
295 : : * sequence state through a pointer fetched from beyond the list.
296 : : *
297 : : * This only keeps the subscript inside the list. An index that is wrong
298 : : * but still in range is not detected, and cannot be; the sequence it
299 : : * points at then receives another sequence's data. That is the same kind
300 : : * of damage as the publisher reporting a wrong value in any other column,
301 : : * and is likewise beyond what we can check.
302 : : */
303 [ + - - + ]: 31 : if (*seqidx < 0 || *seqidx >= list_length(seqinfos))
304 [ # # ]: 0 : elog(ERROR, "invalid sequence index %d received from the publisher",
305 : : *seqidx);
306 : :
307 : : /* Identify the corresponding local sequence for the given index. */
308 : 31 : *seqinfo = seqinfo_local =
309 : 31 : (LogicalRepSequenceInfo *) list_nth(seqinfos, *seqidx);
310 : :
311 : : /*
312 : : * has_sequence_privilege() itself returns NULL, rather than false, when
313 : : * the sequence has been dropped concurrently after it was identified in
314 : : * the catalog snapshot (see has_sequence_privilege_id()). Treat that as a
315 : : * missing sequence on the publisher.
316 : : */
317 : 31 : datum = slot_getattr(slot, ++col, &isnull);
318 [ + + ]: 31 : if (isnull)
319 : 1 : return COPYSEQ_SKIPPED;
320 : :
321 : 30 : remote_has_select_priv = DatumGetBool(datum);
322 : :
323 : : /*
324 : : * The remote sequence state can be NULL if the publisher lacks the
325 : : * required privileges or if the sequence was dropped concurrently after
326 : : * it was identified in the catalog snapshot (see pg_get_sequence_data()).
327 : : */
328 : 30 : datum = slot_getattr(slot, ++col, &isnull);
329 [ + + ]: 30 : if (isnull)
330 : : {
331 : : /*
332 : : * The sequence was dropped concurrently after it was identified in
333 : : * the catalog snapshot. Treat it as skipped (and, since it no longer
334 : : * exists on the publisher, ultimately missing).
335 : : */
336 [ - + ]: 1 : if (remote_has_select_priv)
337 : 0 : return COPYSEQ_SKIPPED;
338 : :
339 : : /*
340 : : * The publisher lacks the SELECT privilege required by
341 : : * pg_get_sequence_data(). Since has_sequence_privilege() returned
342 : : * false, not NULL, do not classify this sequence as missing on the
343 : : * publisher.
344 : : */
345 : 1 : seqinfo_local->found_on_pub = true;
346 : 1 : return COPYSEQ_PUBLISHER_INSUFFICIENT_PERM;
347 : : }
348 : :
349 : 29 : seqinfo_local->last_value = DatumGetInt64(datum);
350 : :
351 : 29 : seqinfo_local->is_called = DatumGetBool(slot_getattr(slot, ++col, &isnull));
352 : : Assert(!isnull);
353 : :
354 : 29 : seqinfo_local->page_lsn = DatumGetLSN(slot_getattr(slot, ++col, &isnull));
355 : : Assert(!isnull);
356 : :
357 : 29 : remote_typid = DatumGetObjectId(slot_getattr(slot, ++col, &isnull));
358 : : Assert(!isnull);
359 : :
360 : 29 : remote_start = DatumGetInt64(slot_getattr(slot, ++col, &isnull));
361 : : Assert(!isnull);
362 : :
363 : 29 : remote_increment = DatumGetInt64(slot_getattr(slot, ++col, &isnull));
364 : : Assert(!isnull);
365 : :
366 : 29 : remote_min = DatumGetInt64(slot_getattr(slot, ++col, &isnull));
367 : : Assert(!isnull);
368 : :
369 : 29 : remote_max = DatumGetInt64(slot_getattr(slot, ++col, &isnull));
370 : : Assert(!isnull);
371 : :
372 : 29 : remote_cycle = DatumGetBool(slot_getattr(slot, ++col, &isnull));
373 : : Assert(!isnull);
374 : :
375 : : /* Sanity check */
376 : : Assert(col == REMOTE_SEQ_COL_COUNT);
377 : :
378 : 29 : seqinfo_local->found_on_pub = true;
379 : :
380 : 29 : *sequence_rel = try_table_open(seqinfo_local->localrelid, RowExclusiveLock);
381 : :
382 : : /* Sequence was concurrently dropped? */
383 [ - + ]: 29 : if (!*sequence_rel)
384 : 0 : return COPYSEQ_SKIPPED;
385 : :
386 : 29 : tup = SearchSysCache1(SEQRELID, ObjectIdGetDatum(seqinfo_local->localrelid));
387 : :
388 : : /* Sequence was concurrently dropped? */
389 [ - + ]: 29 : if (!HeapTupleIsValid(tup))
390 [ # # ]: 0 : elog(ERROR, "cache lookup failed for sequence %u",
391 : : seqinfo_local->localrelid);
392 : :
393 : 29 : local_seq = (Form_pg_sequence) GETSTRUCT(tup);
394 : :
395 : : /* Sequence parameters for remote/local are the same? */
396 [ + - ]: 29 : if (local_seq->seqtypid != remote_typid ||
397 [ + + ]: 29 : local_seq->seqstart != remote_start ||
398 [ + + ]: 28 : local_seq->seqincrement != remote_increment ||
399 [ + - ]: 26 : local_seq->seqmin != remote_min ||
400 [ + - ]: 26 : local_seq->seqmax != remote_max ||
401 [ - + ]: 26 : local_seq->seqcycle != remote_cycle)
402 : 3 : result = COPYSEQ_MISMATCH;
403 : :
404 : : /* Sequence was concurrently renamed? */
405 [ + - ]: 29 : if (strcmp(seqinfo_local->nspname,
406 : 29 : get_namespace_name(RelationGetNamespace(*sequence_rel))) ||
407 [ - + ]: 29 : strcmp(seqinfo_local->seqname, RelationGetRelationName(*sequence_rel)))
408 : 0 : result = COPYSEQ_MISMATCH;
409 : :
410 : 29 : ReleaseSysCache(tup);
411 : 29 : return result;
412 : : }
413 : :
414 : : /*
415 : : * Apply remote sequence state to local sequence and mark it as
416 : : * synchronized (READY).
417 : : */
418 : : static CopySeqResult
419 : 26 : copy_sequence(LogicalRepSequenceInfo *seqinfo, Oid seqowner)
420 : : {
421 : : UserContext ucxt;
422 : : AclResult aclresult;
423 : 26 : bool run_as_owner = MySubscription->runasowner;
424 : 26 : Oid seqoid = seqinfo->localrelid;
425 : : Relation rel;
426 : :
427 : : /*
428 : : * Take the subscription object lock before checking whether this sequence
429 : : * is still part of the subscription. The lock is held until the end of
430 : : * the transaction, so the check and the state update below are protected
431 : : * from a concurrent ALTER SUBSCRIPTION ... REFRESH PUBLICATION.
432 : : *
433 : : * AlterSubscription() takes this lock in AccessExclusiveLock mode while
434 : : * removing pg_subscription_rel rows, so the row cannot be removed between
435 : : * the check and the state update.
436 : : */
437 : 26 : LockSharedObject(SubscriptionRelationId, MySubscription->oid, 0,
438 : : AccessShareLock);
439 : :
440 : : /*
441 : : * The sequence may no longer be part of the subscription, in which case
442 : : * there is nothing to synchronize and the caller just skips it.
443 : : */
444 [ - + ]: 26 : if (!SearchSysCacheExists2(SUBSCRIPTIONRELMAP,
445 : : ObjectIdGetDatum(seqoid),
446 : : ObjectIdGetDatum(MySubscription->oid)))
447 : 0 : return COPYSEQ_NOT_SUBSCRIBED;
448 : :
449 : : /*
450 : : * If the user did not opt to run as the owner of the subscription
451 : : * ('run_as_owner'), then copy the sequence as the owner of the sequence.
452 : : */
453 [ + - ]: 26 : if (!run_as_owner)
454 : 26 : SwitchToUntrustedUser(seqowner, &ucxt);
455 : :
456 : 26 : aclresult = pg_class_aclcheck(seqoid, GetUserId(), ACL_UPDATE);
457 : :
458 [ - + ]: 26 : if (aclresult != ACLCHECK_OK)
459 : : {
460 [ # # ]: 0 : if (!run_as_owner)
461 : 0 : RestoreUserContext(&ucxt);
462 : :
463 : 0 : return COPYSEQ_SUBSCRIBER_INSUFFICIENT_PERM;
464 : : }
465 : :
466 : : /*
467 : : * The log counter (log_cnt) tracks how many sequence values are still
468 : : * unused locally. It is only relevant to the local node and managed
469 : : * internally by nextval() when allocating new ranges. Since log_cnt does
470 : : * not affect the visible sequence state (like last_value or is_called)
471 : : * and is only used for local caching, it need not be copied to the
472 : : * subscriber during synchronization.
473 : : */
474 : 26 : SetSequence(seqoid, seqinfo->last_value, seqinfo->is_called);
475 : :
476 [ + - ]: 26 : if (!run_as_owner)
477 : 26 : RestoreUserContext(&ucxt);
478 : :
479 : 26 : rel = table_open(SubscriptionRelRelationId, RowExclusiveLock);
480 : :
481 : : /*
482 : : * Record the remote sequence's LSN in pg_subscription_rel and mark the
483 : : * sequence as READY. Both locks it needs are held already, the object
484 : : * lock from further up and the relation lock just taken, so say so rather
485 : : * than have it take and release them again.
486 : : */
487 : 26 : UpdateSubscriptionRelState(MySubscription->oid, seqoid, SUBREL_STATE_READY,
488 : : seqinfo->page_lsn, true);
489 : :
490 : 26 : table_close(rel, NoLock);
491 : :
492 : 26 : return COPYSEQ_SUCCESS;
493 : : }
494 : :
495 : : /*
496 : : * Copy existing data of sequences from the publisher.
497 : : */
498 : : static void
499 : 16 : copy_sequences(WalReceiverConn *conn)
500 : : {
501 : 16 : int cur_batch_base_index = 0;
502 : 16 : int n_seqinfos = list_length(seqinfos);
503 : 16 : List *mismatched_seqs_idx = NIL;
504 : 16 : List *missing_seqs_idx = NIL;
505 : 16 : List *sub_insuffperm_seqs_idx = NIL;
506 : 16 : List *pub_insuffperm_seqs_idx = NIL;
507 : : StringInfoData seqstr;
508 : : StringInfoData cmd;
509 : : MemoryContext oldctx;
510 : :
511 : : /*
512 : : * Sequence synchronization depends on publisher-side functionality
513 : : * introduced in PostgreSQL 19, so it cannot work against an older
514 : : * publisher.
515 : : */
516 [ - + ]: 16 : if (walrcv_server_version(conn) < 190000)
517 [ # # ]: 0 : ereport(ERROR,
518 : : errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
519 : : errmsg("cannot synchronize sequences if the publisher is running a version earlier than PostgreSQL 19"));
520 : :
521 : 16 : initStringInfo(&seqstr);
522 : 16 : initStringInfo(&cmd);
523 : :
524 : : #define MAX_SEQUENCES_SYNC_PER_BATCH 100
525 : :
526 [ - + ]: 16 : elog(DEBUG1,
527 : : "logical replication sequence synchronization for subscription \"%s\" - total unsynchronized: %d",
528 : : MySubscription->name, n_seqinfos);
529 : :
530 [ + + ]: 32 : while (cur_batch_base_index < n_seqinfos)
531 : : {
532 : 16 : Oid seqRow[REMOTE_SEQ_COL_COUNT] = {INT4OID, BOOLOID, INT8OID,
533 : : BOOLOID, LSNOID, OIDOID, INT8OID, INT8OID, INT8OID, INT8OID, BOOLOID};
534 : 16 : int batch_size = 0;
535 : 16 : int batch_succeeded_count = 0;
536 : 16 : int batch_mismatched_count = 0;
537 : 16 : int batch_skipped_count = 0;
538 : 16 : int batch_sub_insuffperm_count = 0;
539 : 16 : int batch_pub_insuffperm_count = 0;
540 : : int batch_missing_count;
541 : :
542 : : WalRcvExecResult *res;
543 : : TupleTableSlot *slot;
544 : :
545 : 16 : StartTransactionCommand();
546 : 16 : maybe_reread_subscription();
547 : :
548 [ + + ]: 50 : for (int idx = cur_batch_base_index; idx < n_seqinfos; idx++)
549 : : {
550 : : char *nspname_literal;
551 : : char *seqname_literal;
552 : :
553 : : LogicalRepSequenceInfo *seqinfo =
554 : 34 : (LogicalRepSequenceInfo *) list_nth(seqinfos, idx);
555 : :
556 [ + + ]: 34 : if (seqstr.len > 0)
557 : 18 : appendStringInfoString(&seqstr, ", ");
558 : :
559 : 34 : nspname_literal = quote_literal_cstr(seqinfo->nspname);
560 : 34 : seqname_literal = quote_literal_cstr(seqinfo->seqname);
561 : :
562 : 34 : appendStringInfo(&seqstr, "(%s, %s, %d)",
563 : : nspname_literal, seqname_literal, idx);
564 : :
565 [ - + ]: 34 : if (++batch_size == MAX_SEQUENCES_SYNC_PER_BATCH)
566 : 0 : break;
567 : : }
568 : :
569 : : /*
570 : : * We deliberately avoid acquiring a local lock on the sequence before
571 : : * querying the publisher to prevent potential distributed deadlocks
572 : : * in bi-directional replication setups.
573 : : *
574 : : * Example scenario:
575 : : *
576 : : * - On each node, a background worker acquires a lock on a sequence
577 : : * as part of a sync operation.
578 : : *
579 : : * - Concurrently, a user transaction attempts to alter the same
580 : : * sequence, waiting on the background worker's lock.
581 : : *
582 : : * - Meanwhile, a query from the other node tries to access metadata
583 : : * that depends on the completion of the alter operation.
584 : : *
585 : : * - This creates a circular wait across nodes:
586 : : *
587 : : * Node-1: Query -> waits on Alter -> waits on Sync Worker
588 : : *
589 : : * Node-2: Query -> waits on Alter -> waits on Sync Worker
590 : : *
591 : : * Since each node only sees part of the wait graph, the deadlock may
592 : : * go undetected, leading to indefinite blocking.
593 : : *
594 : : * Note: Each entry in VALUES includes an index 'seqidx' that
595 : : * represents the sequence's position in the local 'seqinfos' list.
596 : : * This index is propagated to the query results and later used to
597 : : * directly map the fetched publisher sequence rows back to their
598 : : * corresponding local entries without relying on result order or name
599 : : * matching.
600 : : */
601 : 16 : appendStringInfo(&cmd,
602 : : "SELECT s.seqidx, has_sequence_privilege(c.oid, 'SELECT'),\n"
603 : : " ps.*, seq.seqtypid,\n"
604 : : " seq.seqstart, seq.seqincrement, seq.seqmin,\n"
605 : : " seq.seqmax, seq.seqcycle\n"
606 : : "FROM ( VALUES %s ) AS s (schname, seqname, seqidx)\n"
607 : : "JOIN pg_namespace n ON n.nspname = s.schname\n"
608 : : "JOIN pg_class c ON c.relnamespace = n.oid AND c.relname = s.seqname\n"
609 : : "JOIN pg_sequence seq ON seq.seqrelid = c.oid\n"
610 : : "JOIN LATERAL pg_get_sequence_data(seq.seqrelid) AS ps ON true\n",
611 : : seqstr.data);
612 : :
613 : 16 : res = walrcv_exec(conn, cmd.data, lengthof(seqRow), seqRow);
614 [ - + ]: 16 : if (res->status != WALRCV_OK_TUPLES)
615 [ # # ]: 0 : ereport(ERROR,
616 : : errcode(ERRCODE_CONNECTION_FAILURE),
617 : : errmsg("could not fetch sequence information from the publisher: %s",
618 : : res->err));
619 : :
620 : 16 : slot = MakeSingleTupleTableSlot(res->tupledesc, &TTSOpsMinimalTuple);
621 [ + + ]: 47 : while (tuplestore_gettupleslot(res->tuplestore, true, false, slot))
622 : : {
623 : : CopySeqResult sync_status;
624 : : LogicalRepSequenceInfo *seqinfo;
625 : 31 : Relation sequence_rel = NULL;
626 : : int seqidx;
627 : :
628 [ - + ]: 31 : CHECK_FOR_INTERRUPTS();
629 : :
630 [ - + ]: 31 : if (ConfigReloadPending)
631 : : {
632 : 0 : ConfigReloadPending = false;
633 : 0 : ProcessConfigFile(PGC_SIGHUP);
634 : : }
635 : :
636 : 31 : sync_status = get_and_validate_seq_info(slot, &sequence_rel,
637 : : &seqinfo, &seqidx);
638 [ + + ]: 31 : if (sync_status == COPYSEQ_SUCCESS)
639 : 26 : sync_status = copy_sequence(seqinfo,
640 : 26 : sequence_rel->rd_rel->relowner);
641 : :
642 [ + + - + : 31 : switch (sync_status)
+ - - ]
643 : : {
644 : 26 : case COPYSEQ_SUCCESS:
645 [ - + ]: 26 : elog(DEBUG1,
646 : : "logical replication synchronization for subscription \"%s\", sequence \"%s.%s\" has finished",
647 : : MySubscription->name, seqinfo->nspname,
648 : : seqinfo->seqname);
649 : 26 : batch_succeeded_count++;
650 : 26 : break;
651 : 3 : case COPYSEQ_MISMATCH:
652 : :
653 : : /*
654 : : * Remember mismatched sequences in a long-lived memory
655 : : * context since these will be used after the transaction
656 : : * is committed.
657 : : */
658 : 3 : oldctx = MemoryContextSwitchTo(ApplyContext);
659 : 3 : mismatched_seqs_idx = lappend_int(mismatched_seqs_idx,
660 : : seqidx);
661 : 3 : MemoryContextSwitchTo(oldctx);
662 : 3 : batch_mismatched_count++;
663 : 3 : break;
664 : 0 : case COPYSEQ_SUBSCRIBER_INSUFFICIENT_PERM:
665 : :
666 : : /*
667 : : * Remember sequences with insufficient privileges in a
668 : : * long-lived memory context since these will be used
669 : : * after the transaction is committed.
670 : : */
671 : 0 : oldctx = MemoryContextSwitchTo(ApplyContext);
672 : 0 : sub_insuffperm_seqs_idx = lappend_int(sub_insuffperm_seqs_idx,
673 : : seqidx);
674 : 0 : MemoryContextSwitchTo(oldctx);
675 : 0 : batch_sub_insuffperm_count++;
676 : 0 : break;
677 : 1 : case COPYSEQ_PUBLISHER_INSUFFICIENT_PERM:
678 : :
679 : : /*
680 : : * Remember sequences for which the publisher lacks the
681 : : * privileges required by pg_get_sequence_data().
682 : : */
683 : 1 : oldctx = MemoryContextSwitchTo(ApplyContext);
684 : 1 : pub_insuffperm_seqs_idx = lappend_int(pub_insuffperm_seqs_idx,
685 : : seqidx);
686 : 1 : MemoryContextSwitchTo(oldctx);
687 : 1 : batch_pub_insuffperm_count++;
688 : 1 : break;
689 : 1 : case COPYSEQ_SKIPPED:
690 : :
691 : : /*
692 : : * Concurrent removal of a sequence on the subscriber is
693 : : * treated as success, since the only viable action is to
694 : : * skip the corresponding sequence data. Missing sequences
695 : : * on the publisher are treated as ERROR.
696 : : */
697 [ - + ]: 1 : if (seqinfo->found_on_pub)
698 : : {
699 [ # # ]: 0 : ereport(LOG,
700 : : errmsg("skip synchronization of sequence \"%s.%s\" because it has been dropped concurrently",
701 : : seqinfo->nspname,
702 : : seqinfo->seqname));
703 : 0 : batch_skipped_count++;
704 : : }
705 : 1 : break;
706 : 0 : case COPYSEQ_NOT_SUBSCRIBED:
707 : :
708 : : /*
709 : : * A concurrent refresh removed this sequence from the
710 : : * subscription. Skipping it is the only sensible action,
711 : : * and it must not be treated as an error.
712 : : */
713 [ # # ]: 0 : ereport(LOG,
714 : : errmsg("skip synchronization of sequence \"%s.%s\" because it is no longer part of subscription \"%s\"",
715 : : seqinfo->nspname, seqinfo->seqname,
716 : : MySubscription->name));
717 : 0 : batch_skipped_count++;
718 : 0 : break;
719 : : }
720 : :
721 [ + + ]: 31 : if (sequence_rel)
722 : 29 : table_close(sequence_rel, NoLock);
723 : : }
724 : :
725 : 16 : ExecDropSingleTupleTableSlot(slot);
726 : 16 : walrcv_clear_result(res);
727 : 16 : resetStringInfo(&seqstr);
728 : 16 : resetStringInfo(&cmd);
729 : :
730 : 16 : batch_missing_count = batch_size - (batch_succeeded_count +
731 : 16 : batch_mismatched_count +
732 : 16 : batch_sub_insuffperm_count +
733 : 16 : batch_pub_insuffperm_count +
734 : : batch_skipped_count);
735 : :
736 [ - + ]: 16 : elog(DEBUG1,
737 : : "logical replication sequence synchronization for subscription \"%s\" - batch #%d = %d attempted, %d succeeded, %d mismatched, %d subscriber insufficient permission, %d publisher insufficient permission, %d missing from publisher, %d skipped",
738 : : MySubscription->name,
739 : : (cur_batch_base_index / MAX_SEQUENCES_SYNC_PER_BATCH) + 1,
740 : : batch_size, batch_succeeded_count, batch_mismatched_count,
741 : : batch_sub_insuffperm_count, batch_pub_insuffperm_count, batch_missing_count, batch_skipped_count);
742 : :
743 : : /* Commit this batch, and prepare for next batch */
744 : 16 : CommitTransactionCommand();
745 : :
746 [ + + ]: 16 : if (batch_missing_count)
747 : : {
748 [ + + ]: 15 : for (int idx = cur_batch_base_index; idx < cur_batch_base_index + batch_size; idx++)
749 : : {
750 : : LogicalRepSequenceInfo *seqinfo =
751 : 11 : (LogicalRepSequenceInfo *) list_nth(seqinfos, idx);
752 : :
753 : : /* If the sequence was not found on publisher, record it */
754 [ + + ]: 11 : if (!seqinfo->found_on_pub)
755 : 4 : missing_seqs_idx = lappend_int(missing_seqs_idx, idx);
756 : : }
757 : : }
758 : :
759 : : /*
760 : : * cur_batch_base_index is not incremented sequentially because some
761 : : * sequences may be missing, and the number of fetched rows may not
762 : : * match the batch size.
763 : : */
764 : 16 : cur_batch_base_index += batch_size;
765 : : }
766 : :
767 : : /* Report mismatches, permission issues, or missing sequences */
768 : 16 : report_sequence_errors(mismatched_seqs_idx, sub_insuffperm_seqs_idx,
769 : : pub_insuffperm_seqs_idx, missing_seqs_idx);
770 : 8 : }
771 : :
772 : : /*
773 : : * Identifies sequences that require synchronization and initiates the
774 : : * synchronization process.
775 : : */
776 : : static void
777 : 16 : LogicalRepSyncSequences(void)
778 : : {
779 : : char *err;
780 : : bool must_use_password;
781 : : Relation rel;
782 : : HeapTuple tup;
783 : : ScanKeyData skey[2];
784 : : SysScanDesc scan;
785 : 16 : Oid subid = MyLogicalRepWorker->subid;
786 : : StringInfoData app_name;
787 : :
788 : 16 : StartTransactionCommand();
789 : 16 : maybe_reread_subscription();
790 : :
791 : 16 : rel = table_open(SubscriptionRelRelationId, AccessShareLock);
792 : :
793 : 16 : ScanKeyInit(&skey[0],
794 : : Anum_pg_subscription_rel_srsubid,
795 : : BTEqualStrategyNumber, F_OIDEQ,
796 : : ObjectIdGetDatum(subid));
797 : :
798 : 16 : ScanKeyInit(&skey[1],
799 : : Anum_pg_subscription_rel_srsubstate,
800 : : BTEqualStrategyNumber, F_CHAREQ,
801 : : CharGetDatum(SUBREL_STATE_INIT));
802 : :
803 : 16 : scan = systable_beginscan(rel, InvalidOid, false,
804 : : NULL, 2, skey);
805 [ + + ]: 51 : while (HeapTupleIsValid(tup = systable_getnext(scan)))
806 : : {
807 : : Form_pg_subscription_rel subrel;
808 : : LogicalRepSequenceInfo *seq;
809 : : Relation sequence_rel;
810 : : MemoryContext oldctx;
811 : :
812 [ - + ]: 35 : CHECK_FOR_INTERRUPTS();
813 : :
814 : 35 : subrel = (Form_pg_subscription_rel) GETSTRUCT(tup);
815 : :
816 : : /*
817 : : * Lock the sequence so its identity (namespace and name) cannot
818 : : * change under us via a concurrent DROP, RENAME or SET SCHEMA. The
819 : : * lock is released immediately rather than at the transaction end.
820 : : * The later synchronization does not depend on this captured identity
821 : : * remaining valid, as it re-opens the sequence and tolerates
822 : : * concurrent changes. Releasing early also avoids holding one lock
823 : : * per sequence, which could exhaust the lock table.
824 : : */
825 : 35 : sequence_rel = try_table_open(subrel->srrelid, AccessShareLock);
826 : :
827 : : /* Skip if sequence was dropped concurrently */
828 [ - + ]: 35 : if (!sequence_rel)
829 : 0 : continue;
830 : :
831 : : /* Skip if the relation is not a sequence */
832 [ + + ]: 35 : if (sequence_rel->rd_rel->relkind != RELKIND_SEQUENCE)
833 : : {
834 : 1 : table_close(sequence_rel, AccessShareLock);
835 : 1 : continue;
836 : : }
837 : :
838 : : /*
839 : : * Worker needs to process sequences across transaction boundary, so
840 : : * allocate them under long-lived context.
841 : : */
842 : 34 : oldctx = MemoryContextSwitchTo(ApplyContext);
843 : :
844 : 34 : seq = palloc0_object(LogicalRepSequenceInfo);
845 : 34 : seq->localrelid = subrel->srrelid;
846 : 34 : seq->nspname = get_namespace_name(RelationGetNamespace(sequence_rel));
847 : 34 : seq->seqname = pstrdup(RelationGetRelationName(sequence_rel));
848 : 34 : seqinfos = lappend(seqinfos, seq);
849 : :
850 : 34 : MemoryContextSwitchTo(oldctx);
851 : :
852 : 34 : table_close(sequence_rel, AccessShareLock);
853 : : }
854 : :
855 : : /* Cleanup */
856 : 16 : systable_endscan(scan);
857 : 16 : table_close(rel, AccessShareLock);
858 : :
859 : 16 : CommitTransactionCommand();
860 : :
861 : : /*
862 : : * Exit early if no catalog entries found, likely due to concurrent drops.
863 : : */
864 [ - + ]: 16 : if (!seqinfos)
865 : 0 : return;
866 : :
867 : : /* Is the use of a password mandatory? */
868 [ + - ]: 32 : must_use_password = MySubscription->passwordrequired &&
869 [ - + ]: 16 : !MySubscription->ownersuperuser;
870 : :
871 : 16 : initStringInfo(&app_name);
872 : 16 : appendStringInfo(&app_name, "pg_%u_sequence_sync_" UINT64_FORMAT,
873 : 16 : MySubscription->oid, GetSystemIdentifier());
874 : :
875 : : /*
876 : : * Establish the connection to the publisher for sequence synchronization.
877 : : */
878 : 16 : LogRepWorkerWalRcvConn =
879 : 16 : walrcv_connect(MySubscriptionConninfo, true, true,
880 : : must_use_password,
881 : : app_name.data, &err);
882 [ - + ]: 16 : if (LogRepWorkerWalRcvConn == NULL)
883 [ # # ]: 0 : ereport(ERROR,
884 : : errcode(ERRCODE_CONNECTION_FAILURE),
885 : : errmsg("sequencesync worker for subscription \"%s\" could not connect to the publisher: %s",
886 : : MySubscription->name, err));
887 : :
888 : 16 : pfree(app_name.data);
889 : :
890 : 16 : copy_sequences(LogRepWorkerWalRcvConn);
891 : : }
892 : :
893 : : /*
894 : : * Execute the initial sync with error handling. Disable the subscription,
895 : : * if required.
896 : : *
897 : : * Note that we don't handle FATAL errors which are probably because of system
898 : : * resource error and are not repeatable.
899 : : */
900 : : static void
901 : 16 : start_sequence_sync(void)
902 : : {
903 : : Assert(am_sequencesync_worker());
904 : :
905 [ + + ]: 16 : PG_TRY();
906 : : {
907 : : /* Call initial sync. */
908 : 16 : LogicalRepSyncSequences();
909 : : }
910 : 8 : PG_CATCH();
911 : : {
912 [ - + ]: 8 : if (MySubscription->disableonerr)
913 : 0 : DisableSubscriptionAndExit();
914 : : else
915 : : {
916 : : /*
917 : : * Report the worker failed during sequence synchronization. Abort
918 : : * the current transaction so that the stats message is sent in an
919 : : * idle state.
920 : : */
921 : 8 : AbortOutOfAnyTransaction();
922 : 8 : pgstat_report_subscription_error(MySubscription->oid);
923 : :
924 : 8 : PG_RE_THROW();
925 : : }
926 : : }
927 [ - + ]: 8 : PG_END_TRY();
928 : 8 : }
929 : :
930 : : /* Logical Replication sequencesync worker entry point */
931 : : void
932 : 16 : SequenceSyncWorkerMain(Datum main_arg)
933 : : {
934 : 16 : int worker_slot = DatumGetInt32(main_arg);
935 : :
936 : 16 : SetupApplyOrSyncWorker(worker_slot);
937 : :
938 : 16 : start_sequence_sync();
939 : :
940 : 8 : FinishSyncWorker();
941 : : }
|