/src/postgres/src/backend/replication/slotfuncs.c
Line | Count | Source |
1 | | /*------------------------------------------------------------------------- |
2 | | * |
3 | | * slotfuncs.c |
4 | | * Support functions for replication slots |
5 | | * |
6 | | * Copyright (c) 2012-2026, PostgreSQL Global Development Group |
7 | | * |
8 | | * IDENTIFICATION |
9 | | * src/backend/replication/slotfuncs.c |
10 | | * |
11 | | *------------------------------------------------------------------------- |
12 | | */ |
13 | | #include "postgres.h" |
14 | | |
15 | | #include "access/htup_details.h" |
16 | | #include "access/xlog_internal.h" |
17 | | #include "access/xlogrecovery.h" |
18 | | #include "access/xlogutils.h" |
19 | | #include "funcapi.h" |
20 | | #include "replication/logical.h" |
21 | | #include "replication/slot.h" |
22 | | #include "replication/slotsync.h" |
23 | | #include "storage/proc.h" |
24 | | #include "utils/builtins.h" |
25 | | #include "utils/guc.h" |
26 | | #include "utils/pg_lsn.h" |
27 | | |
28 | | /* |
29 | | * Map SlotSyncSkipReason enum values to human-readable names. |
30 | | */ |
31 | | static const char *SlotSyncSkipReasonNames[] = { |
32 | | [SS_SKIP_NONE] = "none", |
33 | | [SS_SKIP_WAL_NOT_FLUSHED] = "wal_not_flushed", |
34 | | [SS_SKIP_WAL_OR_ROWS_REMOVED] = "wal_or_rows_removed", |
35 | | [SS_SKIP_NO_CONSISTENT_SNAPSHOT] = "no_consistent_snapshot", |
36 | | [SS_SKIP_INVALID] = "slot_invalidated" |
37 | | }; |
38 | | |
39 | | /* |
40 | | * Helper function for creating a new physical replication slot with |
41 | | * given arguments. Note that this function doesn't release the created |
42 | | * slot. |
43 | | * |
44 | | * If restart_lsn is a valid value, we use it without WAL reservation |
45 | | * routine. So the caller must guarantee that WAL is available. |
46 | | */ |
47 | | static void |
48 | | create_physical_replication_slot(char *name, bool immediately_reserve, |
49 | | bool temporary, XLogRecPtr restart_lsn) |
50 | 0 | { |
51 | 0 | Assert(!MyReplicationSlot); |
52 | | |
53 | | /* acquire replication slot, this will check for conflicting names */ |
54 | 0 | ReplicationSlotCreate(name, false, |
55 | 0 | temporary ? RS_TEMPORARY : RS_PERSISTENT, false, |
56 | 0 | false, false, false); |
57 | |
|
58 | 0 | if (immediately_reserve) |
59 | 0 | { |
60 | | /* Reserve WAL as the user asked for it */ |
61 | 0 | if (!XLogRecPtrIsValid(restart_lsn)) |
62 | 0 | ReplicationSlotReserveWal(); |
63 | 0 | else |
64 | 0 | MyReplicationSlot->data.restart_lsn = restart_lsn; |
65 | | |
66 | | /* Write this slot to disk */ |
67 | 0 | ReplicationSlotMarkDirty(); |
68 | 0 | ReplicationSlotSave(); |
69 | 0 | } |
70 | 0 | } |
71 | | |
72 | | /* |
73 | | * SQL function for creating a new physical (streaming replication) |
74 | | * replication slot. |
75 | | */ |
76 | | Datum |
77 | | pg_create_physical_replication_slot(PG_FUNCTION_ARGS) |
78 | 0 | { |
79 | 0 | Name name = PG_GETARG_NAME(0); |
80 | 0 | bool immediately_reserve = PG_GETARG_BOOL(1); |
81 | 0 | bool temporary = PG_GETARG_BOOL(2); |
82 | 0 | Datum values[2]; |
83 | 0 | bool nulls[2]; |
84 | 0 | TupleDesc tupdesc; |
85 | 0 | HeapTuple tuple; |
86 | 0 | Datum result; |
87 | |
|
88 | 0 | if (get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE) |
89 | 0 | elog(ERROR, "return type must be a row type"); |
90 | | |
91 | 0 | CheckSlotPermissions(); |
92 | |
|
93 | 0 | CheckSlotRequirements(false); |
94 | |
|
95 | 0 | create_physical_replication_slot(NameStr(*name), |
96 | 0 | immediately_reserve, |
97 | 0 | temporary, |
98 | 0 | InvalidXLogRecPtr); |
99 | |
|
100 | 0 | values[0] = NameGetDatum(&MyReplicationSlot->data.name); |
101 | 0 | nulls[0] = false; |
102 | |
|
103 | 0 | if (immediately_reserve) |
104 | 0 | { |
105 | 0 | values[1] = LSNGetDatum(MyReplicationSlot->data.restart_lsn); |
106 | 0 | nulls[1] = false; |
107 | 0 | } |
108 | 0 | else |
109 | 0 | nulls[1] = true; |
110 | |
|
111 | 0 | tuple = heap_form_tuple(tupdesc, values, nulls); |
112 | 0 | result = HeapTupleGetDatum(tuple); |
113 | |
|
114 | 0 | ReplicationSlotRelease(); |
115 | |
|
116 | 0 | PG_RETURN_DATUM(result); |
117 | 0 | } |
118 | | |
119 | | |
120 | | /* |
121 | | * Helper function for creating a new logical replication slot with |
122 | | * given arguments. Note that this function doesn't release the created |
123 | | * slot. |
124 | | * |
125 | | * When find_startpoint is false, the slot's confirmed_flush is not set; it's |
126 | | * caller's responsibility to ensure it's set to something sensible. |
127 | | */ |
128 | | static void |
129 | | create_logical_replication_slot(char *name, char *plugin, |
130 | | bool temporary, bool two_phase, |
131 | | bool failover, |
132 | | XLogRecPtr restart_lsn, |
133 | | bool find_startpoint) |
134 | 0 | { |
135 | 0 | LogicalDecodingContext *ctx = NULL; |
136 | |
|
137 | 0 | Assert(!MyReplicationSlot); |
138 | | |
139 | | /* |
140 | | * Acquire a logical decoding slot, this will check for conflicting names. |
141 | | * Initially create persistent slot as ephemeral - that allows us to |
142 | | * nicely handle errors during initialization because it'll get dropped if |
143 | | * this transaction fails. We'll make it persistent at the end. Temporary |
144 | | * slots can be created as temporary from beginning as they get dropped on |
145 | | * error as well. |
146 | | */ |
147 | 0 | ReplicationSlotCreate(name, true, |
148 | 0 | temporary ? RS_TEMPORARY : RS_EPHEMERAL, two_phase, |
149 | 0 | false, failover, false); |
150 | | |
151 | | /* |
152 | | * Ensure the logical decoding is enabled before initializing the logical |
153 | | * decoding context. |
154 | | */ |
155 | 0 | EnsureLogicalDecodingEnabled(); |
156 | | |
157 | | /* |
158 | | * Outside of recovery, holding a valid logical slot prevents logical |
159 | | * decoding from being disabled. During recovery, however, replaying a |
160 | | * status change record can disable it at any time regardless of slot |
161 | | * existence, so we cannot assert that it is still enabled here. That is |
162 | | * harmless: such a replay invalidates slots, so this slot creation fails |
163 | | * afterwards. |
164 | | */ |
165 | 0 | Assert(RecoveryInProgress() || IsLogicalDecodingEnabled()); |
166 | | |
167 | | /* |
168 | | * Create logical decoding context to find start point or, if we don't |
169 | | * need it, to 1) bump slot's restart_lsn and xmin 2) check plugin sanity. |
170 | | * |
171 | | * Note: when !find_startpoint this is still important, because it's at |
172 | | * this point that the output plugin is validated. |
173 | | */ |
174 | 0 | ctx = CreateInitDecodingContext(plugin, NIL, |
175 | 0 | false, /* just catalogs is OK */ |
176 | 0 | false, /* not repack */ |
177 | 0 | restart_lsn, |
178 | 0 | XL_ROUTINE(.page_read = read_local_xlog_page, |
179 | 0 | .segment_open = wal_segment_open, |
180 | 0 | .segment_close = wal_segment_close), |
181 | 0 | NULL, NULL, NULL); |
182 | | |
183 | | /* |
184 | | * If caller needs us to determine the decoding start point, do so now. |
185 | | * This might take a while. |
186 | | */ |
187 | 0 | if (find_startpoint) |
188 | 0 | DecodingContextFindStartpoint(ctx); |
189 | | |
190 | | /* don't need the decoding context anymore */ |
191 | 0 | FreeDecodingContext(ctx); |
192 | 0 | } |
193 | | |
194 | | /* |
195 | | * SQL function for creating a new logical replication slot. |
196 | | */ |
197 | | Datum |
198 | | pg_create_logical_replication_slot(PG_FUNCTION_ARGS) |
199 | 0 | { |
200 | 0 | Name name = PG_GETARG_NAME(0); |
201 | 0 | Name plugin = PG_GETARG_NAME(1); |
202 | 0 | bool temporary = PG_GETARG_BOOL(2); |
203 | 0 | bool two_phase = PG_GETARG_BOOL(3); |
204 | 0 | bool failover = PG_GETARG_BOOL(4); |
205 | 0 | Datum result; |
206 | 0 | TupleDesc tupdesc; |
207 | 0 | HeapTuple tuple; |
208 | 0 | Datum values[2]; |
209 | 0 | bool nulls[2]; |
210 | |
|
211 | 0 | if (get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE) |
212 | 0 | elog(ERROR, "return type must be a row type"); |
213 | | |
214 | 0 | CheckSlotPermissions(); |
215 | |
|
216 | 0 | CheckLogicalDecodingRequirements(false); |
217 | |
|
218 | 0 | create_logical_replication_slot(NameStr(*name), |
219 | 0 | NameStr(*plugin), |
220 | 0 | temporary, |
221 | 0 | two_phase, |
222 | 0 | failover, |
223 | 0 | InvalidXLogRecPtr, |
224 | 0 | true); |
225 | |
|
226 | 0 | values[0] = NameGetDatum(&MyReplicationSlot->data.name); |
227 | 0 | values[1] = LSNGetDatum(MyReplicationSlot->data.confirmed_flush); |
228 | |
|
229 | 0 | memset(nulls, 0, sizeof(nulls)); |
230 | |
|
231 | 0 | tuple = heap_form_tuple(tupdesc, values, nulls); |
232 | 0 | result = HeapTupleGetDatum(tuple); |
233 | | |
234 | | /* ok, slot is now fully created, mark it as persistent if needed */ |
235 | 0 | if (!temporary) |
236 | 0 | ReplicationSlotPersist(); |
237 | 0 | ReplicationSlotRelease(); |
238 | |
|
239 | 0 | PG_RETURN_DATUM(result); |
240 | 0 | } |
241 | | |
242 | | |
243 | | /* |
244 | | * SQL function for dropping a replication slot. |
245 | | */ |
246 | | Datum |
247 | | pg_drop_replication_slot(PG_FUNCTION_ARGS) |
248 | 0 | { |
249 | 0 | Name name = PG_GETARG_NAME(0); |
250 | |
|
251 | 0 | CheckSlotPermissions(); |
252 | |
|
253 | 0 | CheckSlotRequirements(false); |
254 | |
|
255 | 0 | ReplicationSlotDrop(NameStr(*name), true); |
256 | |
|
257 | 0 | PG_RETURN_VOID(); |
258 | 0 | } |
259 | | |
260 | | /* |
261 | | * pg_get_replication_slots - SQL SRF showing all replication slots |
262 | | * that currently exist on the database cluster. |
263 | | */ |
264 | | Datum |
265 | | pg_get_replication_slots(PG_FUNCTION_ARGS) |
266 | 0 | { |
267 | 0 | #define PG_GET_REPLICATION_SLOTS_COLS 21 |
268 | 0 | ReturnSetInfo *rsinfo = (ReturnSetInfo *) fcinfo->resultinfo; |
269 | 0 | XLogRecPtr currlsn; |
270 | 0 | int slotno; |
271 | | |
272 | | /* |
273 | | * We don't require any special permission to see this function's data |
274 | | * because nothing should be sensitive. The most critical being the slot |
275 | | * name, which shouldn't contain anything particularly sensitive. |
276 | | */ |
277 | |
|
278 | 0 | InitMaterializedSRF(fcinfo, 0); |
279 | |
|
280 | 0 | currlsn = GetXLogWriteRecPtr(); |
281 | |
|
282 | 0 | LWLockAcquire(ReplicationSlotControlLock, LW_SHARED); |
283 | 0 | for (slotno = 0; slotno < max_replication_slots + max_repack_replication_slots; slotno++) |
284 | 0 | { |
285 | 0 | ReplicationSlot *slot = &ReplicationSlotCtl->replication_slots[slotno]; |
286 | 0 | ReplicationSlot slot_contents; |
287 | 0 | Datum values[PG_GET_REPLICATION_SLOTS_COLS]; |
288 | 0 | bool nulls[PG_GET_REPLICATION_SLOTS_COLS]; |
289 | 0 | WALAvailability walstate; |
290 | 0 | int i; |
291 | 0 | ReplicationSlotInvalidationCause cause; |
292 | |
|
293 | 0 | if (!slot->in_use) |
294 | 0 | continue; |
295 | | |
296 | | /* Copy slot contents while holding spinlock, then examine at leisure */ |
297 | 0 | SpinLockAcquire(&slot->mutex); |
298 | 0 | slot_contents = *slot; |
299 | 0 | SpinLockRelease(&slot->mutex); |
300 | |
|
301 | 0 | memset(values, 0, sizeof(values)); |
302 | 0 | memset(nulls, 0, sizeof(nulls)); |
303 | |
|
304 | 0 | i = 0; |
305 | 0 | values[i++] = NameGetDatum(&slot_contents.data.name); |
306 | |
|
307 | 0 | if (slot_contents.data.database == InvalidOid) |
308 | 0 | nulls[i++] = true; |
309 | 0 | else |
310 | 0 | values[i++] = NameGetDatum(&slot_contents.data.plugin); |
311 | |
|
312 | 0 | if (slot_contents.data.database == InvalidOid) |
313 | 0 | values[i++] = CStringGetTextDatum("physical"); |
314 | 0 | else |
315 | 0 | values[i++] = CStringGetTextDatum("logical"); |
316 | |
|
317 | 0 | if (slot_contents.data.database == InvalidOid) |
318 | 0 | nulls[i++] = true; |
319 | 0 | else |
320 | 0 | values[i++] = ObjectIdGetDatum(slot_contents.data.database); |
321 | |
|
322 | 0 | values[i++] = BoolGetDatum(slot_contents.data.persistency == RS_TEMPORARY); |
323 | 0 | values[i++] = BoolGetDatum(slot_contents.active_proc != INVALID_PROC_NUMBER); |
324 | |
|
325 | 0 | if (slot_contents.active_proc != INVALID_PROC_NUMBER) |
326 | 0 | values[i++] = Int32GetDatum(GetPGProcByNumber(slot_contents.active_proc)->pid); |
327 | 0 | else |
328 | 0 | nulls[i++] = true; |
329 | |
|
330 | 0 | if (slot_contents.data.xmin != InvalidTransactionId) |
331 | 0 | values[i++] = TransactionIdGetDatum(slot_contents.data.xmin); |
332 | 0 | else |
333 | 0 | nulls[i++] = true; |
334 | |
|
335 | 0 | if (slot_contents.data.catalog_xmin != InvalidTransactionId) |
336 | 0 | values[i++] = TransactionIdGetDatum(slot_contents.data.catalog_xmin); |
337 | 0 | else |
338 | 0 | nulls[i++] = true; |
339 | |
|
340 | 0 | if (XLogRecPtrIsValid(slot_contents.data.restart_lsn)) |
341 | 0 | values[i++] = LSNGetDatum(slot_contents.data.restart_lsn); |
342 | 0 | else |
343 | 0 | nulls[i++] = true; |
344 | |
|
345 | 0 | if (XLogRecPtrIsValid(slot_contents.data.confirmed_flush)) |
346 | 0 | values[i++] = LSNGetDatum(slot_contents.data.confirmed_flush); |
347 | 0 | else |
348 | 0 | nulls[i++] = true; |
349 | | |
350 | | /* |
351 | | * If the slot has not been invalidated, test availability from |
352 | | * restart_lsn. |
353 | | */ |
354 | 0 | if (slot_contents.data.invalidated != RS_INVAL_NONE) |
355 | 0 | walstate = WALAVAIL_REMOVED; |
356 | 0 | else |
357 | 0 | walstate = GetWALAvailability(slot_contents.data.restart_lsn); |
358 | |
|
359 | 0 | switch (walstate) |
360 | 0 | { |
361 | 0 | case WALAVAIL_INVALID_LSN: |
362 | 0 | nulls[i++] = true; |
363 | 0 | break; |
364 | | |
365 | 0 | case WALAVAIL_RESERVED: |
366 | 0 | values[i++] = CStringGetTextDatum("reserved"); |
367 | 0 | break; |
368 | | |
369 | 0 | case WALAVAIL_EXTENDED: |
370 | 0 | values[i++] = CStringGetTextDatum("extended"); |
371 | 0 | break; |
372 | | |
373 | 0 | case WALAVAIL_UNRESERVED: |
374 | 0 | values[i++] = CStringGetTextDatum("unreserved"); |
375 | 0 | break; |
376 | | |
377 | 0 | case WALAVAIL_REMOVED: |
378 | | |
379 | | /* |
380 | | * If we read the restart_lsn long enough ago, maybe that file |
381 | | * has been removed by now. However, the walsender could have |
382 | | * moved forward enough that it jumped to another file after |
383 | | * we looked. If checkpointer signalled the process to |
384 | | * termination, then it's definitely lost; but if a process is |
385 | | * still alive, then "unreserved" seems more appropriate. |
386 | | * |
387 | | * If we do change it, save the state for safe_wal_size below. |
388 | | */ |
389 | 0 | if (XLogRecPtrIsValid(slot_contents.data.restart_lsn)) |
390 | 0 | { |
391 | 0 | ProcNumber procno; |
392 | |
|
393 | 0 | SpinLockAcquire(&slot->mutex); |
394 | 0 | procno = slot->active_proc; |
395 | 0 | slot_contents.data.restart_lsn = slot->data.restart_lsn; |
396 | 0 | SpinLockRelease(&slot->mutex); |
397 | 0 | if (procno != INVALID_PROC_NUMBER) |
398 | 0 | { |
399 | 0 | values[i++] = CStringGetTextDatum("unreserved"); |
400 | 0 | walstate = WALAVAIL_UNRESERVED; |
401 | 0 | break; |
402 | 0 | } |
403 | 0 | } |
404 | 0 | values[i++] = CStringGetTextDatum("lost"); |
405 | 0 | break; |
406 | 0 | } |
407 | | |
408 | | /* |
409 | | * safe_wal_size is only computed for slots that have not been lost, |
410 | | * and only if there's a configured maximum size. |
411 | | */ |
412 | 0 | if (walstate == WALAVAIL_REMOVED || max_slot_wal_keep_size_mb < 0) |
413 | 0 | nulls[i++] = true; |
414 | 0 | else |
415 | 0 | { |
416 | 0 | XLogSegNo targetSeg; |
417 | 0 | uint64 slotKeepSegs; |
418 | 0 | uint64 keepSegs; |
419 | 0 | XLogSegNo failSeg; |
420 | 0 | XLogRecPtr failLSN; |
421 | |
|
422 | 0 | XLByteToSeg(slot_contents.data.restart_lsn, targetSeg, wal_segment_size); |
423 | | |
424 | | /* determine how many segments can be kept by slots */ |
425 | 0 | slotKeepSegs = XLogMBVarToSegs(max_slot_wal_keep_size_mb, wal_segment_size); |
426 | | /* ditto for wal_keep_size */ |
427 | 0 | keepSegs = XLogMBVarToSegs(wal_keep_size_mb, wal_segment_size); |
428 | | |
429 | | /* if currpos reaches failLSN, we lose our segment */ |
430 | 0 | failSeg = targetSeg + Max(slotKeepSegs, keepSegs) + 1; |
431 | 0 | XLogSegNoOffsetToRecPtr(failSeg, 0, wal_segment_size, failLSN); |
432 | |
|
433 | 0 | values[i++] = Int64GetDatum(failLSN - currlsn); |
434 | 0 | } |
435 | |
|
436 | 0 | values[i++] = BoolGetDatum(slot_contents.data.two_phase); |
437 | |
|
438 | 0 | if (slot_contents.data.two_phase && |
439 | 0 | XLogRecPtrIsValid(slot_contents.data.two_phase_at)) |
440 | 0 | values[i++] = LSNGetDatum(slot_contents.data.two_phase_at); |
441 | 0 | else |
442 | 0 | nulls[i++] = true; |
443 | |
|
444 | 0 | if (slot_contents.inactive_since > 0) |
445 | 0 | values[i++] = TimestampTzGetDatum(slot_contents.inactive_since); |
446 | 0 | else |
447 | 0 | nulls[i++] = true; |
448 | |
|
449 | 0 | cause = slot_contents.data.invalidated; |
450 | |
|
451 | 0 | if (SlotIsPhysical(&slot_contents)) |
452 | 0 | nulls[i++] = true; |
453 | 0 | else |
454 | 0 | { |
455 | | /* |
456 | | * rows_removed and wal_level_insufficient are the only two |
457 | | * reasons for the logical slot's conflict with recovery. |
458 | | */ |
459 | 0 | if (cause == RS_INVAL_HORIZON || |
460 | 0 | cause == RS_INVAL_WAL_LEVEL) |
461 | 0 | values[i++] = BoolGetDatum(true); |
462 | 0 | else |
463 | 0 | values[i++] = BoolGetDatum(false); |
464 | 0 | } |
465 | |
|
466 | 0 | if (cause == RS_INVAL_NONE) |
467 | 0 | nulls[i++] = true; |
468 | 0 | else |
469 | 0 | values[i++] = CStringGetTextDatum(GetSlotInvalidationCauseName(cause)); |
470 | |
|
471 | 0 | values[i++] = BoolGetDatum(slot_contents.data.failover); |
472 | |
|
473 | 0 | values[i++] = BoolGetDatum(slot_contents.data.synced); |
474 | |
|
475 | 0 | if (slot_contents.slotsync_skip_reason == SS_SKIP_NONE) |
476 | 0 | nulls[i++] = true; |
477 | 0 | else |
478 | 0 | values[i++] = CStringGetTextDatum(SlotSyncSkipReasonNames[slot_contents.slotsync_skip_reason]); |
479 | |
|
480 | 0 | Assert(i == PG_GET_REPLICATION_SLOTS_COLS); |
481 | |
|
482 | 0 | tuplestore_putvalues(rsinfo->setResult, rsinfo->setDesc, |
483 | 0 | values, nulls); |
484 | 0 | } |
485 | | |
486 | 0 | LWLockRelease(ReplicationSlotControlLock); |
487 | |
|
488 | 0 | return (Datum) 0; |
489 | 0 | } |
490 | | |
491 | | /* |
492 | | * Helper function for advancing our physical replication slot forward. |
493 | | * |
494 | | * The LSN position to move to is compared simply to the slot's restart_lsn, |
495 | | * knowing that any position older than that would be removed by successive |
496 | | * checkpoints. |
497 | | */ |
498 | | static XLogRecPtr |
499 | | pg_physical_replication_slot_advance(XLogRecPtr moveto) |
500 | 0 | { |
501 | 0 | XLogRecPtr startlsn = MyReplicationSlot->data.restart_lsn; |
502 | 0 | XLogRecPtr retlsn = startlsn; |
503 | |
|
504 | 0 | Assert(XLogRecPtrIsValid(moveto)); |
505 | |
|
506 | 0 | if (startlsn < moveto) |
507 | 0 | { |
508 | 0 | SpinLockAcquire(&MyReplicationSlot->mutex); |
509 | 0 | MyReplicationSlot->data.restart_lsn = moveto; |
510 | 0 | SpinLockRelease(&MyReplicationSlot->mutex); |
511 | 0 | retlsn = moveto; |
512 | | |
513 | | /* |
514 | | * Dirty the slot so as it is written out at the next checkpoint. Note |
515 | | * that the LSN position advanced may still be lost in the event of a |
516 | | * crash, but this makes the data consistent after a clean shutdown. |
517 | | */ |
518 | 0 | ReplicationSlotMarkDirty(); |
519 | | |
520 | | /* |
521 | | * Wake up logical walsenders holding logical failover slots after |
522 | | * updating the restart_lsn of the physical slot. |
523 | | */ |
524 | 0 | PhysicalWakeupLogicalWalSnd(); |
525 | 0 | } |
526 | |
|
527 | 0 | return retlsn; |
528 | 0 | } |
529 | | |
530 | | /* |
531 | | * Advance our logical replication slot forward. See |
532 | | * LogicalSlotAdvanceAndCheckSnapState for details. |
533 | | */ |
534 | | static XLogRecPtr |
535 | | pg_logical_replication_slot_advance(XLogRecPtr moveto) |
536 | 0 | { |
537 | 0 | return LogicalSlotAdvanceAndCheckSnapState(moveto, NULL); |
538 | 0 | } |
539 | | |
540 | | /* |
541 | | * SQL function for moving the position in a replication slot. |
542 | | */ |
543 | | Datum |
544 | | pg_replication_slot_advance(PG_FUNCTION_ARGS) |
545 | 0 | { |
546 | 0 | Name slotname = PG_GETARG_NAME(0); |
547 | 0 | XLogRecPtr moveto = PG_GETARG_LSN(1); |
548 | 0 | XLogRecPtr endlsn; |
549 | 0 | XLogRecPtr minlsn; |
550 | 0 | TupleDesc tupdesc; |
551 | 0 | Datum values[2]; |
552 | 0 | bool nulls[2]; |
553 | 0 | HeapTuple tuple; |
554 | 0 | Datum result; |
555 | |
|
556 | 0 | Assert(!MyReplicationSlot); |
557 | |
|
558 | 0 | CheckSlotPermissions(); |
559 | |
|
560 | 0 | if (!XLogRecPtrIsValid(moveto)) |
561 | 0 | ereport(ERROR, |
562 | 0 | (errcode(ERRCODE_INVALID_PARAMETER_VALUE), |
563 | 0 | errmsg("invalid target WAL LSN"))); |
564 | | |
565 | | /* Build a tuple descriptor for our result type */ |
566 | 0 | if (get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE) |
567 | 0 | elog(ERROR, "return type must be a row type"); |
568 | | |
569 | | /* |
570 | | * We can't move slot past what's been flushed/replayed so clamp the |
571 | | * target position accordingly. |
572 | | */ |
573 | 0 | if (!RecoveryInProgress()) |
574 | 0 | moveto = Min(moveto, GetFlushRecPtr(NULL)); |
575 | 0 | else |
576 | 0 | moveto = Min(moveto, GetXLogReplayRecPtr(NULL)); |
577 | | |
578 | | /* Acquire the slot so we "own" it */ |
579 | 0 | ReplicationSlotAcquire(NameStr(*slotname), true, true); |
580 | | |
581 | | /* A slot whose restart_lsn has never been reserved cannot be advanced */ |
582 | 0 | if (!XLogRecPtrIsValid(MyReplicationSlot->data.restart_lsn)) |
583 | 0 | ereport(ERROR, |
584 | 0 | (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), |
585 | 0 | errmsg("replication slot \"%s\" cannot be advanced", |
586 | 0 | NameStr(*slotname)), |
587 | 0 | errdetail("This slot has never previously reserved WAL, or it has been invalidated."))); |
588 | | |
589 | | /* |
590 | | * Check if the slot is not moving backwards. Physical slots rely simply |
591 | | * on restart_lsn as a minimum point, while logical slots have confirmed |
592 | | * consumption up to confirmed_flush, meaning that in both cases data |
593 | | * older than that is not available anymore. |
594 | | */ |
595 | 0 | if (OidIsValid(MyReplicationSlot->data.database)) |
596 | 0 | minlsn = MyReplicationSlot->data.confirmed_flush; |
597 | 0 | else |
598 | 0 | minlsn = MyReplicationSlot->data.restart_lsn; |
599 | |
|
600 | 0 | if (moveto < minlsn) |
601 | 0 | ereport(ERROR, |
602 | 0 | (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), |
603 | 0 | errmsg("cannot advance replication slot to %X/%08X, minimum is %X/%08X", |
604 | 0 | LSN_FORMAT_ARGS(moveto), LSN_FORMAT_ARGS(minlsn)))); |
605 | | |
606 | | /* Do the actual slot update, depending on the slot type */ |
607 | 0 | if (OidIsValid(MyReplicationSlot->data.database)) |
608 | 0 | endlsn = pg_logical_replication_slot_advance(moveto); |
609 | 0 | else |
610 | 0 | endlsn = pg_physical_replication_slot_advance(moveto); |
611 | |
|
612 | 0 | values[0] = NameGetDatum(&MyReplicationSlot->data.name); |
613 | 0 | nulls[0] = false; |
614 | | |
615 | | /* |
616 | | * Recompute the minimum LSN and xmin across all slots to adjust with the |
617 | | * advancing potentially done. |
618 | | */ |
619 | 0 | ReplicationSlotsComputeRequiredXmin(false); |
620 | 0 | ReplicationSlotsComputeRequiredLSN(); |
621 | |
|
622 | 0 | ReplicationSlotRelease(); |
623 | | |
624 | | /* Return the reached position. */ |
625 | 0 | values[1] = LSNGetDatum(endlsn); |
626 | 0 | nulls[1] = false; |
627 | |
|
628 | 0 | tuple = heap_form_tuple(tupdesc, values, nulls); |
629 | 0 | result = HeapTupleGetDatum(tuple); |
630 | |
|
631 | 0 | PG_RETURN_DATUM(result); |
632 | 0 | } |
633 | | |
634 | | /* |
635 | | * Helper function of copying a replication slot. |
636 | | */ |
637 | | static Datum |
638 | | copy_replication_slot(FunctionCallInfo fcinfo, bool logical_slot) |
639 | 0 | { |
640 | 0 | Name src_name = PG_GETARG_NAME(0); |
641 | 0 | Name dst_name = PG_GETARG_NAME(1); |
642 | 0 | ReplicationSlot *src = NULL; |
643 | 0 | ReplicationSlot first_slot_contents; |
644 | 0 | ReplicationSlot second_slot_contents; |
645 | 0 | XLogRecPtr src_restart_lsn; |
646 | 0 | bool src_islogical; |
647 | 0 | bool temporary; |
648 | 0 | char *plugin; |
649 | 0 | Datum values[2]; |
650 | 0 | bool nulls[2]; |
651 | 0 | Datum result; |
652 | 0 | TupleDesc tupdesc; |
653 | 0 | HeapTuple tuple; |
654 | |
|
655 | 0 | if (get_call_result_type(fcinfo, NULL, &tupdesc) != TYPEFUNC_COMPOSITE) |
656 | 0 | elog(ERROR, "return type must be a row type"); |
657 | | |
658 | 0 | CheckSlotPermissions(); |
659 | |
|
660 | 0 | if (logical_slot) |
661 | 0 | CheckLogicalDecodingRequirements(false); |
662 | 0 | else |
663 | 0 | CheckSlotRequirements(false); |
664 | |
|
665 | 0 | LWLockAcquire(ReplicationSlotControlLock, LW_SHARED); |
666 | | |
667 | | /* |
668 | | * We need to prevent the source slot's reserved WAL from being removed, |
669 | | * but we don't want to lock that slot for very long, and it can advance |
670 | | * in the meantime. So obtain the source slot's data, and create a new |
671 | | * slot using its restart_lsn. Afterwards we lock the source slot again |
672 | | * and verify that the data we copied (name, type) has not changed |
673 | | * incompatibly. No inconvenient WAL removal can occur once the new slot |
674 | | * is created -- but since WAL removal could have occurred before we |
675 | | * managed to create the new slot, we advance the new slot's restart_lsn |
676 | | * to the source slot's updated restart_lsn the second time we lock it. |
677 | | */ |
678 | 0 | for (int i = 0; i < max_replication_slots + max_repack_replication_slots; i++) |
679 | 0 | { |
680 | 0 | ReplicationSlot *s = &ReplicationSlotCtl->replication_slots[i]; |
681 | |
|
682 | 0 | if (s->in_use && strcmp(NameStr(s->data.name), NameStr(*src_name)) == 0) |
683 | 0 | { |
684 | | /* Copy the slot contents while holding spinlock */ |
685 | 0 | SpinLockAcquire(&s->mutex); |
686 | 0 | first_slot_contents = *s; |
687 | 0 | SpinLockRelease(&s->mutex); |
688 | 0 | src = s; |
689 | 0 | break; |
690 | 0 | } |
691 | 0 | } |
692 | |
|
693 | 0 | LWLockRelease(ReplicationSlotControlLock); |
694 | |
|
695 | 0 | if (src == NULL) |
696 | 0 | ereport(ERROR, |
697 | 0 | (errcode(ERRCODE_UNDEFINED_OBJECT), |
698 | 0 | errmsg("replication slot \"%s\" does not exist", NameStr(*src_name)))); |
699 | | |
700 | 0 | src_islogical = SlotIsLogical(&first_slot_contents); |
701 | 0 | src_restart_lsn = first_slot_contents.data.restart_lsn; |
702 | 0 | temporary = (first_slot_contents.data.persistency == RS_TEMPORARY); |
703 | 0 | plugin = logical_slot ? NameStr(first_slot_contents.data.plugin) : NULL; |
704 | | |
705 | | /* Check type of replication slot */ |
706 | 0 | if (src_islogical != logical_slot) |
707 | 0 | ereport(ERROR, |
708 | 0 | (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), |
709 | 0 | src_islogical ? |
710 | 0 | errmsg("cannot copy physical replication slot \"%s\" as a logical replication slot", |
711 | 0 | NameStr(*src_name)) : |
712 | 0 | errmsg("cannot copy logical replication slot \"%s\" as a physical replication slot", |
713 | 0 | NameStr(*src_name)))); |
714 | | |
715 | | /* Copying non-reserved slot doesn't make sense */ |
716 | 0 | if (!XLogRecPtrIsValid(src_restart_lsn)) |
717 | 0 | ereport(ERROR, |
718 | 0 | (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), |
719 | 0 | errmsg("cannot copy a replication slot that doesn't reserve WAL"))); |
720 | | |
721 | | /* Cannot copy an invalidated replication slot */ |
722 | 0 | if (first_slot_contents.data.invalidated != RS_INVAL_NONE) |
723 | 0 | ereport(ERROR, |
724 | 0 | errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), |
725 | 0 | errmsg("cannot copy invalidated replication slot \"%s\"", |
726 | 0 | NameStr(*src_name))); |
727 | | |
728 | | /* Overwrite params from optional arguments */ |
729 | 0 | if (PG_NARGS() >= 3) |
730 | 0 | temporary = PG_GETARG_BOOL(2); |
731 | 0 | if (PG_NARGS() >= 4) |
732 | 0 | { |
733 | 0 | Assert(logical_slot); |
734 | 0 | plugin = NameStr(*(PG_GETARG_NAME(3))); |
735 | 0 | } |
736 | | |
737 | | /* Create new slot and acquire it */ |
738 | 0 | if (logical_slot) |
739 | 0 | { |
740 | | /* |
741 | | * We must not try to read WAL, since we haven't reserved it yet -- |
742 | | * hence pass find_startpoint false. confirmed_flush will be set |
743 | | * below, by copying from the source slot. |
744 | | * |
745 | | * We don't copy the failover option to prevent potential issues with |
746 | | * slot synchronization. For instance, if a slot was synchronized to |
747 | | * the standby, then dropped on the primary, and immediately recreated |
748 | | * by copying from another existing slot with much earlier restart_lsn |
749 | | * and confirmed_flush_lsn, the slot synchronization would only |
750 | | * observe the LSN of the same slot moving backward. As slot |
751 | | * synchronization does not copy the restart_lsn and |
752 | | * confirmed_flush_lsn backward (see update_local_synced_slot() for |
753 | | * details), if a failover happens before the primary's slot catches |
754 | | * up, logical replication cannot continue using the synchronized slot |
755 | | * on the promoted standby because the slot retains the restart_lsn |
756 | | * and confirmed_flush_lsn that are much later than expected. |
757 | | */ |
758 | 0 | create_logical_replication_slot(NameStr(*dst_name), |
759 | 0 | plugin, |
760 | 0 | temporary, |
761 | 0 | false, |
762 | 0 | false, |
763 | 0 | src_restart_lsn, |
764 | 0 | false); |
765 | 0 | } |
766 | 0 | else |
767 | 0 | create_physical_replication_slot(NameStr(*dst_name), |
768 | 0 | true, |
769 | 0 | temporary, |
770 | 0 | src_restart_lsn); |
771 | | |
772 | | /* |
773 | | * Update the destination slot to current values of the source slot; |
774 | | * recheck that the source slot is still the one we saw previously. |
775 | | */ |
776 | 0 | { |
777 | 0 | TransactionId copy_effective_xmin; |
778 | 0 | TransactionId copy_effective_catalog_xmin; |
779 | 0 | TransactionId copy_xmin; |
780 | 0 | TransactionId copy_catalog_xmin; |
781 | 0 | XLogRecPtr copy_restart_lsn; |
782 | 0 | XLogRecPtr copy_confirmed_flush; |
783 | 0 | bool copy_islogical; |
784 | 0 | char *copy_name; |
785 | | |
786 | | /* Copy data of source slot again */ |
787 | 0 | SpinLockAcquire(&src->mutex); |
788 | 0 | second_slot_contents = *src; |
789 | 0 | SpinLockRelease(&src->mutex); |
790 | |
|
791 | 0 | copy_effective_xmin = second_slot_contents.effective_xmin; |
792 | 0 | copy_effective_catalog_xmin = second_slot_contents.effective_catalog_xmin; |
793 | |
|
794 | 0 | copy_xmin = second_slot_contents.data.xmin; |
795 | 0 | copy_catalog_xmin = second_slot_contents.data.catalog_xmin; |
796 | 0 | copy_restart_lsn = second_slot_contents.data.restart_lsn; |
797 | 0 | copy_confirmed_flush = second_slot_contents.data.confirmed_flush; |
798 | | |
799 | | /* for existence check */ |
800 | 0 | copy_name = NameStr(second_slot_contents.data.name); |
801 | 0 | copy_islogical = SlotIsLogical(&second_slot_contents); |
802 | | |
803 | | /* |
804 | | * Check if the source slot still exists and is valid. We regard it as |
805 | | * invalid if the type of replication slot or name has been changed, |
806 | | * or the restart_lsn either is invalid or has gone backward. (The |
807 | | * restart_lsn could go backwards if the source slot is dropped and |
808 | | * copied from an older slot during installation.) |
809 | | * |
810 | | * Since erroring out will release and drop the destination slot we |
811 | | * don't need to release it here. |
812 | | */ |
813 | 0 | if (copy_restart_lsn < src_restart_lsn || |
814 | 0 | src_islogical != copy_islogical || |
815 | 0 | strcmp(copy_name, NameStr(*src_name)) != 0) |
816 | 0 | ereport(ERROR, |
817 | 0 | (errmsg("could not copy replication slot \"%s\"", |
818 | 0 | NameStr(*src_name)), |
819 | 0 | errdetail("The source replication slot was modified incompatibly during the copy operation."))); |
820 | | |
821 | | /* The source slot must have a consistent snapshot */ |
822 | 0 | if (src_islogical && !XLogRecPtrIsValid(copy_confirmed_flush)) |
823 | 0 | ereport(ERROR, |
824 | 0 | (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), |
825 | 0 | errmsg("cannot copy unfinished logical replication slot \"%s\"", |
826 | 0 | NameStr(*src_name)), |
827 | 0 | errhint("Retry when the source replication slot's confirmed_flush_lsn is valid."))); |
828 | | |
829 | | /* |
830 | | * Copying an invalid slot doesn't make sense. Note that the source |
831 | | * slot can become invalid after we create the new slot and copy the |
832 | | * data of source slot. This is possible because the operations in |
833 | | * InvalidateObsoleteReplicationSlots() are not serialized with this |
834 | | * function. Even though we can't detect such a case here, the copied |
835 | | * slot will become invalid in the next checkpoint cycle. |
836 | | */ |
837 | 0 | if (second_slot_contents.data.invalidated != RS_INVAL_NONE) |
838 | 0 | ereport(ERROR, |
839 | 0 | errmsg("cannot copy replication slot \"%s\"", |
840 | 0 | NameStr(*src_name)), |
841 | 0 | errdetail("The source replication slot was invalidated during the copy operation.")); |
842 | | |
843 | | /* Install copied values again */ |
844 | 0 | SpinLockAcquire(&MyReplicationSlot->mutex); |
845 | 0 | MyReplicationSlot->effective_xmin = copy_effective_xmin; |
846 | 0 | MyReplicationSlot->effective_catalog_xmin = copy_effective_catalog_xmin; |
847 | |
|
848 | 0 | MyReplicationSlot->data.xmin = copy_xmin; |
849 | 0 | MyReplicationSlot->data.catalog_xmin = copy_catalog_xmin; |
850 | 0 | MyReplicationSlot->data.restart_lsn = copy_restart_lsn; |
851 | 0 | MyReplicationSlot->data.confirmed_flush = copy_confirmed_flush; |
852 | 0 | SpinLockRelease(&MyReplicationSlot->mutex); |
853 | |
|
854 | 0 | ReplicationSlotMarkDirty(); |
855 | 0 | ReplicationSlotsComputeRequiredXmin(false); |
856 | 0 | ReplicationSlotsComputeRequiredLSN(); |
857 | 0 | ReplicationSlotSave(); |
858 | |
|
859 | | #ifdef USE_ASSERT_CHECKING |
860 | | /* Check that the restart_lsn is available */ |
861 | | { |
862 | | XLogSegNo segno; |
863 | | |
864 | | XLByteToSeg(copy_restart_lsn, segno, wal_segment_size); |
865 | | Assert(XLogGetLastRemovedSegno() < segno); |
866 | | } |
867 | | #endif |
868 | 0 | } |
869 | | |
870 | | /* target slot fully created, mark as persistent if needed */ |
871 | 0 | if (logical_slot && !temporary) |
872 | 0 | ReplicationSlotPersist(); |
873 | | |
874 | | /* All done. Set up the return values */ |
875 | 0 | values[0] = NameGetDatum(dst_name); |
876 | 0 | nulls[0] = false; |
877 | 0 | if (XLogRecPtrIsValid(MyReplicationSlot->data.confirmed_flush)) |
878 | 0 | { |
879 | 0 | values[1] = LSNGetDatum(MyReplicationSlot->data.confirmed_flush); |
880 | 0 | nulls[1] = false; |
881 | 0 | } |
882 | 0 | else |
883 | 0 | nulls[1] = true; |
884 | |
|
885 | 0 | tuple = heap_form_tuple(tupdesc, values, nulls); |
886 | 0 | result = HeapTupleGetDatum(tuple); |
887 | |
|
888 | 0 | ReplicationSlotRelease(); |
889 | |
|
890 | 0 | PG_RETURN_DATUM(result); |
891 | 0 | } |
892 | | |
893 | | /* The wrappers below are all to appease opr_sanity */ |
894 | | Datum |
895 | | pg_copy_logical_replication_slot_a(PG_FUNCTION_ARGS) |
896 | 0 | { |
897 | 0 | return copy_replication_slot(fcinfo, true); |
898 | 0 | } |
899 | | |
900 | | Datum |
901 | | pg_copy_logical_replication_slot_b(PG_FUNCTION_ARGS) |
902 | 0 | { |
903 | 0 | return copy_replication_slot(fcinfo, true); |
904 | 0 | } |
905 | | |
906 | | Datum |
907 | | pg_copy_logical_replication_slot_c(PG_FUNCTION_ARGS) |
908 | 0 | { |
909 | 0 | return copy_replication_slot(fcinfo, true); |
910 | 0 | } |
911 | | |
912 | | Datum |
913 | | pg_copy_physical_replication_slot_a(PG_FUNCTION_ARGS) |
914 | 0 | { |
915 | 0 | return copy_replication_slot(fcinfo, false); |
916 | 0 | } |
917 | | |
918 | | Datum |
919 | | pg_copy_physical_replication_slot_b(PG_FUNCTION_ARGS) |
920 | 0 | { |
921 | 0 | return copy_replication_slot(fcinfo, false); |
922 | 0 | } |
923 | | |
924 | | /* |
925 | | * Synchronize failover enabled replication slots to a standby server |
926 | | * from the primary server. |
927 | | */ |
928 | | Datum |
929 | | pg_sync_replication_slots(PG_FUNCTION_ARGS) |
930 | 0 | { |
931 | 0 | WalReceiverConn *wrconn; |
932 | 0 | char *err; |
933 | 0 | StringInfoData app_name; |
934 | |
|
935 | 0 | CheckSlotPermissions(); |
936 | |
|
937 | 0 | if (!RecoveryInProgress()) |
938 | 0 | ereport(ERROR, |
939 | 0 | errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), |
940 | 0 | errmsg("replication slots can only be synchronized to a standby server")); |
941 | | |
942 | 0 | ValidateSlotSyncParams(ERROR); |
943 | | |
944 | | /* Load the libpq-specific functions */ |
945 | 0 | load_file("libpqwalreceiver", false); |
946 | |
|
947 | 0 | (void) CheckAndGetDbnameFromConninfo(); |
948 | |
|
949 | 0 | initStringInfo(&app_name); |
950 | 0 | if (cluster_name[0]) |
951 | 0 | appendStringInfo(&app_name, "%s_slotsync", cluster_name); |
952 | 0 | else |
953 | 0 | appendStringInfoString(&app_name, "slotsync"); |
954 | | |
955 | | /* Connect to the primary server. */ |
956 | 0 | wrconn = walrcv_connect(PrimaryConnInfo, false, false, false, |
957 | 0 | app_name.data, &err); |
958 | |
|
959 | 0 | if (!wrconn) |
960 | 0 | ereport(ERROR, |
961 | 0 | errcode(ERRCODE_CONNECTION_FAILURE), |
962 | 0 | errmsg("synchronization worker \"%s\" could not connect to the primary server: %s", |
963 | 0 | app_name.data, err)); |
964 | | |
965 | 0 | pfree(app_name.data); |
966 | |
|
967 | 0 | SyncReplicationSlots(wrconn); |
968 | |
|
969 | 0 | walrcv_disconnect(wrconn); |
970 | |
|
971 | 0 | PG_RETURN_VOID(); |
972 | 0 | } |