/src/postgres/src/backend/replication/logical/tablesync.c
Line | Count | Source |
1 | | /*------------------------------------------------------------------------- |
2 | | * tablesync.c |
3 | | * PostgreSQL logical replication: initial table data synchronization |
4 | | * |
5 | | * Copyright (c) 2012-2026, PostgreSQL Global Development Group |
6 | | * |
7 | | * IDENTIFICATION |
8 | | * src/backend/replication/logical/tablesync.c |
9 | | * |
10 | | * NOTES |
11 | | * This file contains code for initial table data synchronization for |
12 | | * logical replication. |
13 | | * |
14 | | * The initial data synchronization is done separately for each table, |
15 | | * in a separate apply worker that only fetches the initial snapshot data |
16 | | * from the publisher and then synchronizes the position in the stream with |
17 | | * the leader apply worker. |
18 | | * |
19 | | * There are several reasons for doing the synchronization this way: |
20 | | * - It allows us to parallelize the initial data synchronization |
21 | | * which lowers the time needed for it to happen. |
22 | | * - The initial synchronization does not have to hold the xid and LSN |
23 | | * for the time it takes to copy data of all tables, causing less |
24 | | * bloat and lower disk consumption compared to doing the |
25 | | * synchronization in a single process for the whole database. |
26 | | * - It allows us to synchronize any tables added after the initial |
27 | | * synchronization has finished. |
28 | | * |
29 | | * The stream position synchronization works in multiple steps: |
30 | | * - Apply worker requests a tablesync worker to start, setting the new |
31 | | * table state to INIT. |
32 | | * - Tablesync worker starts; changes table state from INIT to DATASYNC while |
33 | | * copying. |
34 | | * - Tablesync worker does initial table copy; there is a FINISHEDCOPY (sync |
35 | | * worker specific) state to indicate when the copy phase has completed, so |
36 | | * if the worker crashes with this (non-memory) state then the copy will not |
37 | | * be re-attempted. |
38 | | * - Tablesync worker then sets table state to SYNCWAIT; waits for state change. |
39 | | * - Apply worker periodically checks for tables in SYNCWAIT state. When |
40 | | * any appear, it sets the table state to CATCHUP and starts loop-waiting |
41 | | * until either the table state is set to SYNCDONE or the sync worker |
42 | | * exits. |
43 | | * - After the sync worker has seen the state change to CATCHUP, it will |
44 | | * read the stream and apply changes (acting like an apply worker) until |
45 | | * it catches up to the specified stream position. Then it sets the |
46 | | * state to SYNCDONE. There might be zero changes applied between |
47 | | * CATCHUP and SYNCDONE, because the sync worker might be ahead of the |
48 | | * apply worker. |
49 | | * - Once the state is set to SYNCDONE, the apply will continue tracking |
50 | | * the table until it reaches the SYNCDONE stream position, at which |
51 | | * point it sets state to READY and stops tracking. Again, there might |
52 | | * be zero changes in between. |
53 | | * |
54 | | * So the state progression is always: INIT -> DATASYNC -> FINISHEDCOPY |
55 | | * -> SYNCWAIT -> CATCHUP -> SYNCDONE -> READY. |
56 | | * |
57 | | * The catalog pg_subscription_rel is used to keep information about |
58 | | * subscribed tables and their state. The catalog holds all states |
59 | | * except SYNCWAIT and CATCHUP which are only in shared memory. |
60 | | * |
61 | | * Example flows look like this: |
62 | | * - Apply is in front: |
63 | | * sync:8 |
64 | | * -> set in catalog FINISHEDCOPY |
65 | | * -> set in memory SYNCWAIT |
66 | | * apply:10 |
67 | | * -> set in memory CATCHUP |
68 | | * -> enter wait-loop |
69 | | * sync:10 |
70 | | * -> set in catalog SYNCDONE |
71 | | * -> exit |
72 | | * apply:10 |
73 | | * -> exit wait-loop |
74 | | * -> continue rep |
75 | | * apply:11 |
76 | | * -> set in catalog READY |
77 | | * |
78 | | * - Sync is in front: |
79 | | * sync:10 |
80 | | * -> set in catalog FINISHEDCOPY |
81 | | * -> set in memory SYNCWAIT |
82 | | * apply:8 |
83 | | * -> set in memory CATCHUP |
84 | | * -> continue per-table filtering |
85 | | * sync:10 |
86 | | * -> set in catalog SYNCDONE |
87 | | * -> exit |
88 | | * apply:10 |
89 | | * -> set in catalog READY |
90 | | * -> stop per-table filtering |
91 | | * -> continue rep |
92 | | *------------------------------------------------------------------------- |
93 | | */ |
94 | | |
95 | | #include "postgres.h" |
96 | | |
97 | | #include "access/table.h" |
98 | | #include "access/xact.h" |
99 | | #include "catalog/indexing.h" |
100 | | #include "catalog/pg_subscription_rel.h" |
101 | | #include "catalog/pg_type.h" |
102 | | #include "commands/copy.h" |
103 | | #include "miscadmin.h" |
104 | | #include "nodes/makefuncs.h" |
105 | | #include "parser/parse_relation.h" |
106 | | #include "pgstat.h" |
107 | | #include "replication/logicallauncher.h" |
108 | | #include "replication/logicalrelation.h" |
109 | | #include "replication/logicalworker.h" |
110 | | #include "replication/origin.h" |
111 | | #include "replication/slot.h" |
112 | | #include "replication/walreceiver.h" |
113 | | #include "replication/worker_internal.h" |
114 | | #include "storage/ipc.h" |
115 | | #include "storage/latch.h" |
116 | | #include "storage/lmgr.h" |
117 | | #include "utils/acl.h" |
118 | | #include "utils/array.h" |
119 | | #include "utils/builtins.h" |
120 | | #include "utils/lsyscache.h" |
121 | | #include "utils/rls.h" |
122 | | #include "utils/snapmgr.h" |
123 | | #include "utils/syscache.h" |
124 | | #include "utils/usercontext.h" |
125 | | #include "utils/wait_event.h" |
126 | | |
127 | | List *table_states_not_ready = NIL; |
128 | | |
129 | | static StringInfo copybuf = NULL; |
130 | | |
131 | | /* |
132 | | * Wait until the relation sync state is set in the catalog to the expected |
133 | | * one; return true when it happens. |
134 | | * |
135 | | * Returns false if the table sync worker or the table itself have |
136 | | * disappeared, or the table state has been reset. |
137 | | * |
138 | | * Currently, this is used in the apply worker when transitioning from |
139 | | * CATCHUP state to SYNCDONE. |
140 | | */ |
141 | | static bool |
142 | | wait_for_table_state_change(Oid relid, char expected_state) |
143 | 0 | { |
144 | 0 | char state; |
145 | |
|
146 | 0 | for (;;) |
147 | 0 | { |
148 | 0 | LogicalRepWorker *worker; |
149 | 0 | XLogRecPtr statelsn; |
150 | |
|
151 | 0 | CHECK_FOR_INTERRUPTS(); |
152 | |
|
153 | 0 | InvalidateCatalogSnapshot(); |
154 | 0 | state = GetSubscriptionRelState(MyLogicalRepWorker->subid, |
155 | 0 | relid, &statelsn); |
156 | |
|
157 | 0 | if (state == SUBREL_STATE_UNKNOWN) |
158 | 0 | break; |
159 | | |
160 | 0 | if (state == expected_state) |
161 | 0 | return true; |
162 | | |
163 | | /* Check if the sync worker is still running and bail if not. */ |
164 | 0 | LWLockAcquire(LogicalRepWorkerLock, LW_SHARED); |
165 | 0 | worker = logicalrep_worker_find(WORKERTYPE_TABLESYNC, |
166 | 0 | MyLogicalRepWorker->subid, relid, |
167 | 0 | false); |
168 | 0 | LWLockRelease(LogicalRepWorkerLock); |
169 | 0 | if (!worker) |
170 | 0 | break; |
171 | | |
172 | 0 | (void) WaitLatch(MyLatch, |
173 | 0 | WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, |
174 | 0 | 1000L, WAIT_EVENT_LOGICAL_SYNC_STATE_CHANGE); |
175 | |
|
176 | 0 | ResetLatch(MyLatch); |
177 | 0 | } |
178 | | |
179 | 0 | return false; |
180 | 0 | } |
181 | | |
182 | | /* |
183 | | * Wait until the apply worker changes the state of our synchronization |
184 | | * worker to the expected one. |
185 | | * |
186 | | * Used when transitioning from SYNCWAIT state to CATCHUP. |
187 | | * |
188 | | * Returns false if the apply worker has disappeared. |
189 | | */ |
190 | | static bool |
191 | | wait_for_worker_state_change(char expected_state) |
192 | 0 | { |
193 | 0 | int rc; |
194 | |
|
195 | 0 | for (;;) |
196 | 0 | { |
197 | 0 | LogicalRepWorker *worker; |
198 | |
|
199 | 0 | CHECK_FOR_INTERRUPTS(); |
200 | | |
201 | | /* |
202 | | * Done if already in correct state. (We assume this fetch is atomic |
203 | | * enough to not give a misleading answer if we do it with no lock.) |
204 | | */ |
205 | 0 | if (MyLogicalRepWorker->relstate == expected_state) |
206 | 0 | return true; |
207 | | |
208 | | /* |
209 | | * Bail out if the apply worker has died, else signal it we're |
210 | | * waiting. |
211 | | */ |
212 | 0 | LWLockAcquire(LogicalRepWorkerLock, LW_SHARED); |
213 | 0 | worker = logicalrep_worker_find(WORKERTYPE_APPLY, |
214 | 0 | MyLogicalRepWorker->subid, InvalidOid, |
215 | 0 | false); |
216 | 0 | if (worker && worker->proc) |
217 | 0 | logicalrep_worker_wakeup_ptr(worker); |
218 | 0 | LWLockRelease(LogicalRepWorkerLock); |
219 | 0 | if (!worker) |
220 | 0 | break; |
221 | | |
222 | | /* |
223 | | * Wait. We expect to get a latch signal back from the apply worker, |
224 | | * but use a timeout in case it dies without sending one. |
225 | | */ |
226 | 0 | rc = WaitLatch(MyLatch, |
227 | 0 | WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, |
228 | 0 | 1000L, WAIT_EVENT_LOGICAL_SYNC_STATE_CHANGE); |
229 | |
|
230 | 0 | if (rc & WL_LATCH_SET) |
231 | 0 | ResetLatch(MyLatch); |
232 | 0 | } |
233 | | |
234 | 0 | return false; |
235 | 0 | } |
236 | | |
237 | | /* |
238 | | * Handle table synchronization cooperation from the synchronization |
239 | | * worker. |
240 | | * |
241 | | * If the sync worker is in CATCHUP state and reached (or passed) the |
242 | | * predetermined synchronization point in the WAL stream, mark the table as |
243 | | * SYNCDONE and finish. |
244 | | */ |
245 | | void |
246 | | ProcessSyncingTablesForSync(XLogRecPtr current_lsn) |
247 | 0 | { |
248 | 0 | SpinLockAcquire(&MyLogicalRepWorker->relmutex); |
249 | |
|
250 | 0 | if (MyLogicalRepWorker->relstate == SUBREL_STATE_CATCHUP && |
251 | 0 | current_lsn >= MyLogicalRepWorker->relstate_lsn) |
252 | 0 | { |
253 | 0 | TimeLineID tli; |
254 | 0 | char syncslotname[NAMEDATALEN] = {0}; |
255 | 0 | char originname[NAMEDATALEN] = {0}; |
256 | |
|
257 | 0 | MyLogicalRepWorker->relstate = SUBREL_STATE_SYNCDONE; |
258 | 0 | MyLogicalRepWorker->relstate_lsn = current_lsn; |
259 | |
|
260 | 0 | SpinLockRelease(&MyLogicalRepWorker->relmutex); |
261 | | |
262 | | /* |
263 | | * UpdateSubscriptionRelState must be called within a transaction. |
264 | | */ |
265 | 0 | if (!IsTransactionState()) |
266 | 0 | StartTransactionCommand(); |
267 | |
|
268 | 0 | UpdateSubscriptionRelState(MyLogicalRepWorker->subid, |
269 | 0 | MyLogicalRepWorker->relid, |
270 | 0 | MyLogicalRepWorker->relstate, |
271 | 0 | MyLogicalRepWorker->relstate_lsn, |
272 | 0 | false); |
273 | | |
274 | | /* |
275 | | * End streaming so that LogRepWorkerWalRcvConn can be used to drop |
276 | | * the slot. |
277 | | */ |
278 | 0 | walrcv_endstreaming(LogRepWorkerWalRcvConn, &tli); |
279 | | |
280 | | /* |
281 | | * Cleanup the tablesync slot. |
282 | | * |
283 | | * This has to be done after updating the state because otherwise if |
284 | | * there is an error while doing the database operations we won't be |
285 | | * able to rollback dropped slot. |
286 | | */ |
287 | 0 | ReplicationSlotNameForTablesync(MyLogicalRepWorker->subid, |
288 | 0 | MyLogicalRepWorker->relid, |
289 | 0 | syncslotname, |
290 | 0 | sizeof(syncslotname)); |
291 | | |
292 | | /* |
293 | | * It is important to give an error if we are unable to drop the slot, |
294 | | * otherwise, it won't be dropped till the corresponding subscription |
295 | | * is dropped. So passing missing_ok = false. |
296 | | */ |
297 | 0 | ReplicationSlotDropAtPubNode(LogRepWorkerWalRcvConn, syncslotname, false); |
298 | |
|
299 | 0 | CommitTransactionCommand(); |
300 | 0 | pgstat_report_stat(false); |
301 | | |
302 | | /* |
303 | | * Start a new transaction to clean up the tablesync origin tracking. |
304 | | * This transaction will be ended within the FinishSyncWorker(). Now, |
305 | | * even, if we fail to remove this here, the apply worker will ensure |
306 | | * to clean it up afterward. |
307 | | * |
308 | | * We need to do this after the table state is set to SYNCDONE. |
309 | | * Otherwise, if an error occurs while performing the database |
310 | | * operation, the worker will be restarted and the in-memory state of |
311 | | * replication progress (remote_lsn) won't be rolled-back which would |
312 | | * have been cleared before restart. So, the restarted worker will use |
313 | | * invalid replication progress state resulting in replay of |
314 | | * transactions that have already been applied. |
315 | | */ |
316 | 0 | StartTransactionCommand(); |
317 | |
|
318 | 0 | ReplicationOriginNameForLogicalRep(MyLogicalRepWorker->subid, |
319 | 0 | MyLogicalRepWorker->relid, |
320 | 0 | originname, |
321 | 0 | sizeof(originname)); |
322 | | |
323 | | /* |
324 | | * Resetting the origin session removes the ownership of the slot. |
325 | | * This is needed to allow the origin to be dropped. |
326 | | */ |
327 | 0 | replorigin_session_reset(); |
328 | 0 | replorigin_xact_clear(true); |
329 | | |
330 | | /* |
331 | | * Drop the tablesync's origin tracking if exists. |
332 | | * |
333 | | * There is a chance that the user is concurrently performing refresh |
334 | | * for the subscription where we remove the table state and its origin |
335 | | * or the apply worker would have removed this origin. So passing |
336 | | * missing_ok = true. |
337 | | */ |
338 | 0 | replorigin_drop_by_name(originname, true, false); |
339 | |
|
340 | 0 | FinishSyncWorker(); |
341 | 0 | } |
342 | 0 | else |
343 | 0 | SpinLockRelease(&MyLogicalRepWorker->relmutex); |
344 | 0 | } |
345 | | |
346 | | /* |
347 | | * Handle table synchronization cooperation from the apply worker. |
348 | | * |
349 | | * Walk over all subscription tables that are individually tracked by the |
350 | | * apply process (currently, all that have state other than |
351 | | * SUBREL_STATE_READY) and manage synchronization for them. |
352 | | * |
353 | | * If there are tables that need synchronizing and are not being synchronized |
354 | | * yet, start sync workers for them (if there are free slots for sync |
355 | | * workers). To prevent starting the sync worker for the same relation at a |
356 | | * high frequency after a failure, we store its last start time with each sync |
357 | | * state info. We start the sync worker for the same relation after waiting |
358 | | * at least wal_retrieve_retry_interval. |
359 | | * |
360 | | * For tables that are being synchronized already, check if sync workers |
361 | | * either need action from the apply worker or have finished. This is the |
362 | | * SYNCWAIT to CATCHUP transition. |
363 | | * |
364 | | * If the synchronization position is reached (SYNCDONE), then the table can |
365 | | * be marked as READY and is no longer tracked. |
366 | | */ |
367 | | void |
368 | | ProcessSyncingTablesForApply(XLogRecPtr current_lsn) |
369 | 0 | { |
370 | 0 | struct tablesync_start_time_mapping |
371 | 0 | { |
372 | 0 | Oid relid; |
373 | 0 | TimestampTz last_start_time; |
374 | 0 | }; |
375 | 0 | static HTAB *last_start_times = NULL; |
376 | 0 | ListCell *lc; |
377 | 0 | bool started_tx; |
378 | 0 | bool should_exit = false; |
379 | 0 | Relation rel = NULL; |
380 | |
|
381 | 0 | Assert(!IsTransactionState()); |
382 | | |
383 | | /* We need up-to-date sync state info for subscription tables here. */ |
384 | 0 | FetchRelationStates(NULL, NULL, &started_tx); |
385 | | |
386 | | /* |
387 | | * Prepare a hash table for tracking last start times of workers, to avoid |
388 | | * immediate restarts. We don't need it if there are no tables that need |
389 | | * syncing. |
390 | | */ |
391 | 0 | if (table_states_not_ready != NIL && !last_start_times) |
392 | 0 | { |
393 | 0 | HASHCTL ctl; |
394 | |
|
395 | 0 | ctl.keysize = sizeof(Oid); |
396 | 0 | ctl.entrysize = sizeof(struct tablesync_start_time_mapping); |
397 | 0 | last_start_times = hash_create("Logical replication table sync worker start times", |
398 | 0 | 256, &ctl, HASH_ELEM | HASH_BLOBS); |
399 | 0 | } |
400 | | |
401 | | /* |
402 | | * Clean up the hash table when we're done with all tables (just to |
403 | | * release the bit of memory). |
404 | | */ |
405 | 0 | else if (table_states_not_ready == NIL && last_start_times) |
406 | 0 | { |
407 | 0 | hash_destroy(last_start_times); |
408 | 0 | last_start_times = NULL; |
409 | 0 | } |
410 | | |
411 | | /* |
412 | | * Process all tables that are being synchronized. |
413 | | */ |
414 | 0 | foreach(lc, table_states_not_ready) |
415 | 0 | { |
416 | 0 | SubscriptionRelState *rstate = (SubscriptionRelState *) lfirst(lc); |
417 | |
|
418 | 0 | if (!started_tx) |
419 | 0 | { |
420 | 0 | StartTransactionCommand(); |
421 | 0 | started_tx = true; |
422 | 0 | } |
423 | |
|
424 | 0 | Assert(get_rel_relkind(rstate->relid) != RELKIND_SEQUENCE); |
425 | |
|
426 | 0 | if (rstate->state == SUBREL_STATE_SYNCDONE) |
427 | 0 | { |
428 | | /* |
429 | | * Apply has caught up to the position where the table sync has |
430 | | * finished. Mark the table as ready so that the apply will just |
431 | | * continue to replicate it normally. |
432 | | */ |
433 | 0 | if (current_lsn >= rstate->lsn) |
434 | 0 | { |
435 | 0 | char originname[NAMEDATALEN]; |
436 | |
|
437 | 0 | rstate->state = SUBREL_STATE_READY; |
438 | 0 | rstate->lsn = current_lsn; |
439 | | |
440 | | /* |
441 | | * Remove the tablesync origin tracking if exists. |
442 | | * |
443 | | * There is a chance that the user is concurrently performing |
444 | | * refresh for the subscription where we remove the table |
445 | | * state and its origin or the tablesync worker would have |
446 | | * already removed this origin. We can't rely on tablesync |
447 | | * worker to remove the origin tracking as if there is any |
448 | | * error while dropping we won't restart it to drop the |
449 | | * origin. So passing missing_ok = true. |
450 | | * |
451 | | * Lock the subscription and origin in the same order as we |
452 | | * are doing during DDL commands to avoid deadlocks. See |
453 | | * AlterSubscription_refresh. |
454 | | */ |
455 | 0 | LockSharedObject(SubscriptionRelationId, MyLogicalRepWorker->subid, |
456 | 0 | 0, AccessShareLock); |
457 | |
|
458 | 0 | if (!rel) |
459 | 0 | rel = table_open(SubscriptionRelRelationId, RowExclusiveLock); |
460 | |
|
461 | 0 | ReplicationOriginNameForLogicalRep(MyLogicalRepWorker->subid, |
462 | 0 | rstate->relid, |
463 | 0 | originname, |
464 | 0 | sizeof(originname)); |
465 | 0 | replorigin_drop_by_name(originname, true, false); |
466 | | |
467 | | /* |
468 | | * Update the state to READY only after the origin cleanup. |
469 | | */ |
470 | 0 | UpdateSubscriptionRelState(MyLogicalRepWorker->subid, |
471 | 0 | rstate->relid, rstate->state, |
472 | 0 | rstate->lsn, true); |
473 | 0 | } |
474 | 0 | } |
475 | 0 | else |
476 | 0 | { |
477 | 0 | LogicalRepWorker *syncworker; |
478 | | |
479 | | /* |
480 | | * Look for a sync worker for this relation. |
481 | | */ |
482 | 0 | LWLockAcquire(LogicalRepWorkerLock, LW_SHARED); |
483 | |
|
484 | 0 | syncworker = logicalrep_worker_find(WORKERTYPE_TABLESYNC, |
485 | 0 | MyLogicalRepWorker->subid, |
486 | 0 | rstate->relid, false); |
487 | |
|
488 | 0 | if (syncworker) |
489 | 0 | { |
490 | | /* Found one, update our copy of its state */ |
491 | 0 | SpinLockAcquire(&syncworker->relmutex); |
492 | 0 | rstate->state = syncworker->relstate; |
493 | 0 | rstate->lsn = syncworker->relstate_lsn; |
494 | 0 | if (rstate->state == SUBREL_STATE_SYNCWAIT) |
495 | 0 | { |
496 | | /* |
497 | | * Sync worker is waiting for apply. Tell sync worker it |
498 | | * can catchup now. |
499 | | */ |
500 | 0 | syncworker->relstate = SUBREL_STATE_CATCHUP; |
501 | 0 | syncworker->relstate_lsn = |
502 | 0 | Max(syncworker->relstate_lsn, current_lsn); |
503 | 0 | } |
504 | 0 | SpinLockRelease(&syncworker->relmutex); |
505 | | |
506 | | /* If we told worker to catch up, wait for it. */ |
507 | 0 | if (rstate->state == SUBREL_STATE_SYNCWAIT) |
508 | 0 | { |
509 | | /* Signal the sync worker, as it may be waiting for us. */ |
510 | 0 | if (syncworker->proc) |
511 | 0 | logicalrep_worker_wakeup_ptr(syncworker); |
512 | | |
513 | | /* Now safe to release the LWLock */ |
514 | 0 | LWLockRelease(LogicalRepWorkerLock); |
515 | |
|
516 | 0 | if (started_tx) |
517 | 0 | { |
518 | | /* |
519 | | * We must commit the existing transaction to release |
520 | | * the existing locks before entering a busy loop. |
521 | | * This is required to avoid any undetected deadlocks |
522 | | * due to any existing lock as deadlock detector won't |
523 | | * be able to detect the waits on the latch. |
524 | | * |
525 | | * Also close any tables prior to the commit. |
526 | | */ |
527 | 0 | if (rel) |
528 | 0 | { |
529 | 0 | table_close(rel, NoLock); |
530 | 0 | rel = NULL; |
531 | 0 | } |
532 | 0 | CommitTransactionCommand(); |
533 | 0 | pgstat_report_stat(false); |
534 | 0 | } |
535 | | |
536 | | /* |
537 | | * Enter busy loop and wait for synchronization worker to |
538 | | * reach expected state (or die trying). |
539 | | */ |
540 | 0 | StartTransactionCommand(); |
541 | 0 | started_tx = true; |
542 | |
|
543 | 0 | wait_for_table_state_change(rstate->relid, |
544 | 0 | SUBREL_STATE_SYNCDONE); |
545 | 0 | } |
546 | 0 | else |
547 | 0 | LWLockRelease(LogicalRepWorkerLock); |
548 | 0 | } |
549 | 0 | else |
550 | 0 | { |
551 | | /* |
552 | | * If there is no sync worker for this table yet, count |
553 | | * running sync workers for this subscription, while we have |
554 | | * the lock. |
555 | | */ |
556 | 0 | int nsyncworkers = |
557 | 0 | logicalrep_sync_worker_count(MyLogicalRepWorker->subid); |
558 | 0 | struct tablesync_start_time_mapping *hentry; |
559 | 0 | bool found; |
560 | | |
561 | | /* Now safe to release the LWLock */ |
562 | 0 | LWLockRelease(LogicalRepWorkerLock); |
563 | |
|
564 | 0 | hentry = hash_search(last_start_times, &rstate->relid, |
565 | 0 | HASH_ENTER, &found); |
566 | 0 | if (!found) |
567 | 0 | hentry->last_start_time = 0; |
568 | |
|
569 | 0 | launch_sync_worker(WORKERTYPE_TABLESYNC, nsyncworkers, |
570 | 0 | rstate->relid, &hentry->last_start_time); |
571 | 0 | } |
572 | 0 | } |
573 | 0 | } |
574 | | |
575 | | /* Close table if opened */ |
576 | 0 | if (rel) |
577 | 0 | table_close(rel, NoLock); |
578 | | |
579 | |
|
580 | 0 | if (started_tx) |
581 | 0 | { |
582 | | /* |
583 | | * Even when the two_phase mode is requested by the user, it remains |
584 | | * as 'pending' until all tablesyncs have reached READY state. |
585 | | * |
586 | | * When this happens, we restart the apply worker and (if the |
587 | | * conditions are still ok) then the two_phase tri-state will become |
588 | | * 'enabled' at that time. |
589 | | * |
590 | | * Note: If the subscription has no tables then leave the state as |
591 | | * PENDING, which allows ALTER SUBSCRIPTION ... REFRESH PUBLICATION to |
592 | | * work. |
593 | | */ |
594 | 0 | if (MySubscription->twophasestate == LOGICALREP_TWOPHASE_STATE_PENDING) |
595 | 0 | { |
596 | 0 | CommandCounterIncrement(); /* make updates visible */ |
597 | 0 | if (AllTablesyncsReady()) |
598 | 0 | { |
599 | 0 | ereport(LOG, |
600 | 0 | (errmsg("logical replication apply worker for subscription \"%s\" will restart so that two_phase can be enabled", |
601 | 0 | MySubscription->name))); |
602 | 0 | should_exit = true; |
603 | 0 | } |
604 | 0 | } |
605 | | |
606 | 0 | CommitTransactionCommand(); |
607 | 0 | pgstat_report_stat(true); |
608 | 0 | } |
609 | | |
610 | 0 | if (should_exit) |
611 | 0 | { |
612 | | /* |
613 | | * Reset the last-start time for this worker so that the launcher will |
614 | | * restart it without waiting for wal_retrieve_retry_interval. |
615 | | */ |
616 | 0 | ApplyLauncherForgetWorkerStartTime(MySubscription->oid); |
617 | |
|
618 | 0 | proc_exit(0); |
619 | 0 | } |
620 | 0 | } |
621 | | |
622 | | /* |
623 | | * Create list of columns for COPY based on logical relation mapping. |
624 | | */ |
625 | | static List * |
626 | | make_copy_attnamelist(LogicalRepRelMapEntry *rel) |
627 | 0 | { |
628 | 0 | List *attnamelist = NIL; |
629 | 0 | int i; |
630 | |
|
631 | 0 | for (i = 0; i < rel->remoterel.natts; i++) |
632 | 0 | { |
633 | 0 | attnamelist = lappend(attnamelist, |
634 | 0 | makeString(rel->remoterel.attnames[i])); |
635 | 0 | } |
636 | | |
637 | |
|
638 | 0 | return attnamelist; |
639 | 0 | } |
640 | | |
641 | | /* |
642 | | * Data source callback for the COPY FROM, which reads from the remote |
643 | | * connection and passes the data back to our local COPY. |
644 | | */ |
645 | | static int |
646 | | copy_read_data(void *outbuf, int minread, int maxread) |
647 | 0 | { |
648 | 0 | int bytesread = 0; |
649 | 0 | int avail; |
650 | | |
651 | | /* If there are some leftover data from previous read, use it. */ |
652 | 0 | avail = copybuf->len - copybuf->cursor; |
653 | 0 | if (avail) |
654 | 0 | { |
655 | 0 | if (avail > maxread) |
656 | 0 | avail = maxread; |
657 | 0 | memcpy(outbuf, ©buf->data[copybuf->cursor], avail); |
658 | 0 | copybuf->cursor += avail; |
659 | 0 | maxread -= avail; |
660 | 0 | bytesread += avail; |
661 | 0 | } |
662 | |
|
663 | 0 | while (maxread > 0 && bytesread < minread) |
664 | 0 | { |
665 | 0 | pgsocket fd = PGINVALID_SOCKET; |
666 | 0 | int len; |
667 | 0 | char *buf = NULL; |
668 | |
|
669 | 0 | for (;;) |
670 | 0 | { |
671 | | /* Try read the data. */ |
672 | 0 | len = walrcv_receive(LogRepWorkerWalRcvConn, &buf, &fd); |
673 | |
|
674 | 0 | CHECK_FOR_INTERRUPTS(); |
675 | |
|
676 | 0 | if (len == 0) |
677 | 0 | break; |
678 | 0 | else if (len < 0) |
679 | 0 | return bytesread; |
680 | 0 | else |
681 | 0 | { |
682 | | /* Process the data */ |
683 | 0 | copybuf->data = buf; |
684 | 0 | copybuf->len = len; |
685 | 0 | copybuf->cursor = 0; |
686 | |
|
687 | 0 | avail = copybuf->len - copybuf->cursor; |
688 | 0 | if (avail > maxread) |
689 | 0 | avail = maxread; |
690 | 0 | memcpy(outbuf, ©buf->data[copybuf->cursor], avail); |
691 | 0 | outbuf = (char *) outbuf + avail; |
692 | 0 | copybuf->cursor += avail; |
693 | 0 | maxread -= avail; |
694 | 0 | bytesread += avail; |
695 | 0 | } |
696 | | |
697 | 0 | if (maxread <= 0 || bytesread >= minread) |
698 | 0 | return bytesread; |
699 | 0 | } |
700 | | |
701 | | /* |
702 | | * Wait for more data or latch. |
703 | | */ |
704 | 0 | (void) WaitLatchOrSocket(MyLatch, |
705 | 0 | WL_SOCKET_READABLE | WL_LATCH_SET | |
706 | 0 | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, |
707 | 0 | fd, 1000L, WAIT_EVENT_LOGICAL_SYNC_DATA); |
708 | |
|
709 | 0 | ResetLatch(MyLatch); |
710 | 0 | } |
711 | | |
712 | 0 | return bytesread; |
713 | 0 | } |
714 | | |
715 | | |
716 | | /* |
717 | | * Get information about remote relation in similar fashion the RELATION |
718 | | * message provides during replication. |
719 | | * |
720 | | * This function also returns (a) the relation qualifications to be used in |
721 | | * the COPY command, and (b) whether the remote relation has published any |
722 | | * generated column. |
723 | | */ |
724 | | static void |
725 | | fetch_remote_table_info(char *nspname, char *relname, LogicalRepRelation *lrel, |
726 | | List **qual, bool *gencol_published) |
727 | 0 | { |
728 | 0 | WalRcvExecResult *res; |
729 | 0 | StringInfoData cmd; |
730 | 0 | TupleTableSlot *slot; |
731 | 0 | Oid tableRow[] = {OIDOID, CHAROID, CHAROID}; |
732 | 0 | Oid attrRow[] = {INT2OID, TEXTOID, OIDOID, BOOLOID, BOOLOID}; |
733 | 0 | Oid qualRow[] = {TEXTOID}; |
734 | 0 | bool isnull; |
735 | 0 | int natt; |
736 | 0 | StringInfo pub_names = NULL; |
737 | 0 | Bitmapset *included_cols = NULL; |
738 | 0 | int server_version = walrcv_server_version(LogRepWorkerWalRcvConn); |
739 | |
|
740 | 0 | lrel->nspname = nspname; |
741 | 0 | lrel->relname = relname; |
742 | | |
743 | | /* First fetch Oid and replica identity. */ |
744 | 0 | initStringInfo(&cmd); |
745 | 0 | appendStringInfo(&cmd, "SELECT c.oid, c.relreplident, c.relkind" |
746 | 0 | " FROM pg_catalog.pg_class c" |
747 | 0 | " INNER JOIN pg_catalog.pg_namespace n" |
748 | 0 | " ON (c.relnamespace = n.oid)" |
749 | 0 | " WHERE n.nspname = %s" |
750 | 0 | " AND c.relname = %s", |
751 | 0 | quote_literal_cstr(nspname), |
752 | 0 | quote_literal_cstr(relname)); |
753 | 0 | res = walrcv_exec(LogRepWorkerWalRcvConn, cmd.data, |
754 | 0 | lengthof(tableRow), tableRow); |
755 | |
|
756 | 0 | if (res->status != WALRCV_OK_TUPLES) |
757 | 0 | ereport(ERROR, |
758 | 0 | (errcode(ERRCODE_CONNECTION_FAILURE), |
759 | 0 | errmsg("could not fetch table info for table \"%s.%s\" from publisher: %s", |
760 | 0 | nspname, relname, res->err))); |
761 | | |
762 | 0 | slot = MakeSingleTupleTableSlot(res->tupledesc, &TTSOpsMinimalTuple); |
763 | 0 | if (!tuplestore_gettupleslot(res->tuplestore, true, false, slot)) |
764 | 0 | ereport(ERROR, |
765 | 0 | (errcode(ERRCODE_UNDEFINED_OBJECT), |
766 | 0 | errmsg("table \"%s.%s\" not found on publisher", |
767 | 0 | nspname, relname))); |
768 | | |
769 | 0 | lrel->remoteid = DatumGetObjectId(slot_getattr(slot, 1, &isnull)); |
770 | 0 | Assert(!isnull); |
771 | 0 | lrel->replident = DatumGetChar(slot_getattr(slot, 2, &isnull)); |
772 | 0 | Assert(!isnull); |
773 | 0 | lrel->relkind = DatumGetChar(slot_getattr(slot, 3, &isnull)); |
774 | 0 | Assert(!isnull); |
775 | |
|
776 | 0 | ExecDropSingleTupleTableSlot(slot); |
777 | 0 | walrcv_clear_result(res); |
778 | | |
779 | | |
780 | | /* |
781 | | * Get column lists for each relation. |
782 | | * |
783 | | * We need to do this before fetching info about column names and types, |
784 | | * so that we can skip columns that should not be replicated. |
785 | | */ |
786 | 0 | if (server_version >= 150000) |
787 | 0 | { |
788 | 0 | WalRcvExecResult *pubres; |
789 | 0 | TupleTableSlot *tslot; |
790 | 0 | Oid attrsRow[] = {INT2VECTOROID}; |
791 | | |
792 | | /* Build the pub_names comma-separated string. */ |
793 | 0 | pub_names = makeStringInfo(); |
794 | 0 | GetPublicationsStr(MySubscription->publications, pub_names, true); |
795 | | |
796 | | /* |
797 | | * Fetch info about column lists for the relation (from all the |
798 | | * publications). |
799 | | */ |
800 | 0 | resetStringInfo(&cmd); |
801 | |
|
802 | 0 | if (server_version >= 190000) |
803 | 0 | { |
804 | | /* |
805 | | * We can pass both publication names and relid to |
806 | | * pg_get_publication_tables() since version 19. |
807 | | */ |
808 | 0 | appendStringInfo(&cmd, |
809 | 0 | "SELECT DISTINCT" |
810 | 0 | " (CASE WHEN (array_length(gpt.attrs, 1) = c.relnatts)" |
811 | 0 | " THEN NULL ELSE gpt.attrs END)" |
812 | 0 | " FROM pg_get_publication_tables(ARRAY[%s], %u) gpt," |
813 | 0 | " pg_class c" |
814 | 0 | " WHERE c.oid = gpt.relid", |
815 | 0 | pub_names->data, |
816 | 0 | lrel->remoteid); |
817 | 0 | } |
818 | 0 | else |
819 | 0 | appendStringInfo(&cmd, |
820 | 0 | "SELECT DISTINCT" |
821 | 0 | " (CASE WHEN (array_length(gpt.attrs, 1) = c.relnatts)" |
822 | 0 | " THEN NULL ELSE gpt.attrs END)" |
823 | 0 | " FROM pg_publication p," |
824 | 0 | " LATERAL pg_get_publication_tables(p.pubname) gpt," |
825 | 0 | " pg_class c" |
826 | 0 | " WHERE gpt.relid = %u AND c.oid = gpt.relid" |
827 | 0 | " AND p.pubname IN ( %s )", |
828 | 0 | lrel->remoteid, |
829 | 0 | pub_names->data); |
830 | |
|
831 | 0 | pubres = walrcv_exec(LogRepWorkerWalRcvConn, cmd.data, |
832 | 0 | lengthof(attrsRow), attrsRow); |
833 | |
|
834 | 0 | if (pubres->status != WALRCV_OK_TUPLES) |
835 | 0 | ereport(ERROR, |
836 | 0 | (errcode(ERRCODE_CONNECTION_FAILURE), |
837 | 0 | errmsg("could not fetch column list info for table \"%s.%s\" from publisher: %s", |
838 | 0 | nspname, relname, pubres->err))); |
839 | | |
840 | | /* |
841 | | * We don't support the case where the column list is different for |
842 | | * the same table when combining publications. See comments atop |
843 | | * fetch_relation_list. So there should be only one row returned. |
844 | | * Although we already checked this when creating the subscription, we |
845 | | * still need to check here in case the column list was changed after |
846 | | * creating the subscription and before the sync worker is started. |
847 | | */ |
848 | 0 | if (tuplestore_tuple_count(pubres->tuplestore) > 1) |
849 | 0 | ereport(ERROR, |
850 | 0 | errcode(ERRCODE_FEATURE_NOT_SUPPORTED), |
851 | 0 | errmsg("cannot use different column lists for table \"%s.%s\" in different publications", |
852 | 0 | nspname, relname)); |
853 | | |
854 | | /* |
855 | | * Get the column list and build a single bitmap with the attnums. |
856 | | * |
857 | | * If we find a NULL value, it means all the columns should be |
858 | | * replicated. |
859 | | */ |
860 | 0 | tslot = MakeSingleTupleTableSlot(pubres->tupledesc, &TTSOpsMinimalTuple); |
861 | 0 | if (tuplestore_gettupleslot(pubres->tuplestore, true, false, tslot)) |
862 | 0 | { |
863 | 0 | Datum cfval = slot_getattr(tslot, 1, &isnull); |
864 | |
|
865 | 0 | if (!isnull) |
866 | 0 | { |
867 | 0 | ArrayType *arr; |
868 | 0 | int nelems; |
869 | 0 | int16 *elems; |
870 | |
|
871 | 0 | arr = DatumGetArrayTypeP(cfval); |
872 | 0 | nelems = ARR_DIMS(arr)[0]; |
873 | 0 | elems = (int16 *) ARR_DATA_PTR(arr); |
874 | |
|
875 | 0 | for (natt = 0; natt < nelems; natt++) |
876 | 0 | included_cols = bms_add_member(included_cols, elems[natt]); |
877 | 0 | } |
878 | |
|
879 | 0 | ExecClearTuple(tslot); |
880 | 0 | } |
881 | 0 | ExecDropSingleTupleTableSlot(tslot); |
882 | |
|
883 | 0 | walrcv_clear_result(pubres); |
884 | 0 | } |
885 | | |
886 | | /* |
887 | | * Now fetch column names and types. |
888 | | */ |
889 | 0 | resetStringInfo(&cmd); |
890 | 0 | appendStringInfoString(&cmd, |
891 | 0 | "SELECT a.attnum," |
892 | 0 | " a.attname," |
893 | 0 | " a.atttypid," |
894 | 0 | " a.attnum = ANY(i.indkey)"); |
895 | | |
896 | | /* Generated columns can be replicated since version 18. */ |
897 | 0 | if (server_version >= 180000) |
898 | 0 | appendStringInfoString(&cmd, ", a.attgenerated != ''"); |
899 | |
|
900 | 0 | appendStringInfo(&cmd, |
901 | 0 | " FROM pg_catalog.pg_attribute a" |
902 | 0 | " LEFT JOIN pg_catalog.pg_index i" |
903 | 0 | " ON (i.indexrelid = pg_get_replica_identity_index(%u))" |
904 | 0 | " WHERE a.attnum > 0::pg_catalog.int2" |
905 | 0 | " AND NOT a.attisdropped %s" |
906 | 0 | " AND a.attrelid = %u" |
907 | 0 | " ORDER BY a.attnum", |
908 | 0 | lrel->remoteid, |
909 | 0 | (server_version >= 120000 && server_version < 180000 ? |
910 | 0 | "AND a.attgenerated = ''" : ""), |
911 | 0 | lrel->remoteid); |
912 | 0 | res = walrcv_exec(LogRepWorkerWalRcvConn, cmd.data, |
913 | 0 | server_version >= 180000 ? lengthof(attrRow) : lengthof(attrRow) - 1, attrRow); |
914 | |
|
915 | 0 | if (res->status != WALRCV_OK_TUPLES) |
916 | 0 | ereport(ERROR, |
917 | 0 | (errcode(ERRCODE_CONNECTION_FAILURE), |
918 | 0 | errmsg("could not fetch table info for table \"%s.%s\" from publisher: %s", |
919 | 0 | nspname, relname, res->err))); |
920 | | |
921 | | /* We don't know the number of rows coming, so allocate enough space. */ |
922 | 0 | lrel->attnames = palloc0_array(char *, MaxTupleAttributeNumber); |
923 | 0 | lrel->atttyps = palloc0_array(Oid, MaxTupleAttributeNumber); |
924 | 0 | lrel->attkeys = NULL; |
925 | | |
926 | | /* |
927 | | * Store the columns as a list of names. Ignore those that are not |
928 | | * present in the column list, if there is one. |
929 | | */ |
930 | 0 | natt = 0; |
931 | 0 | slot = MakeSingleTupleTableSlot(res->tupledesc, &TTSOpsMinimalTuple); |
932 | 0 | while (tuplestore_gettupleslot(res->tuplestore, true, false, slot)) |
933 | 0 | { |
934 | 0 | char *rel_colname; |
935 | 0 | AttrNumber attnum; |
936 | |
|
937 | 0 | attnum = DatumGetInt16(slot_getattr(slot, 1, &isnull)); |
938 | 0 | Assert(!isnull); |
939 | | |
940 | | /* If the column is not in the column list, skip it. */ |
941 | 0 | if (included_cols != NULL && !bms_is_member(attnum, included_cols)) |
942 | 0 | { |
943 | 0 | ExecClearTuple(slot); |
944 | 0 | continue; |
945 | 0 | } |
946 | | |
947 | 0 | rel_colname = TextDatumGetCString(slot_getattr(slot, 2, &isnull)); |
948 | 0 | Assert(!isnull); |
949 | |
|
950 | 0 | lrel->attnames[natt] = rel_colname; |
951 | 0 | lrel->atttyps[natt] = DatumGetObjectId(slot_getattr(slot, 3, &isnull)); |
952 | 0 | Assert(!isnull); |
953 | |
|
954 | 0 | if (DatumGetBool(slot_getattr(slot, 4, &isnull))) |
955 | 0 | lrel->attkeys = bms_add_member(lrel->attkeys, natt); |
956 | | |
957 | | /* Remember if the remote table has published any generated column. */ |
958 | 0 | if (server_version >= 180000 && !(*gencol_published)) |
959 | 0 | { |
960 | 0 | *gencol_published = DatumGetBool(slot_getattr(slot, 5, &isnull)); |
961 | 0 | Assert(!isnull); |
962 | 0 | } |
963 | | |
964 | | /* Should never happen. */ |
965 | 0 | if (++natt >= MaxTupleAttributeNumber) |
966 | 0 | elog(ERROR, "too many columns in remote table \"%s.%s\"", |
967 | 0 | nspname, relname); |
968 | | |
969 | 0 | ExecClearTuple(slot); |
970 | 0 | } |
971 | 0 | ExecDropSingleTupleTableSlot(slot); |
972 | |
|
973 | 0 | lrel->natts = natt; |
974 | |
|
975 | 0 | walrcv_clear_result(res); |
976 | | |
977 | | /* |
978 | | * Get relation's row filter expressions. DISTINCT avoids the same |
979 | | * expression of a table in multiple publications from being included |
980 | | * multiple times in the final expression. |
981 | | * |
982 | | * We need to copy the row even if it matches just one of the |
983 | | * publications, so we later combine all the quals with OR. |
984 | | * |
985 | | * For initial synchronization, row filtering can be ignored in following |
986 | | * cases: |
987 | | * |
988 | | * 1) one of the subscribed publications for the table hasn't specified |
989 | | * any row filter |
990 | | * |
991 | | * 2) one of the subscribed publications has puballtables set to true |
992 | | * |
993 | | * 3) one of the subscribed publications is declared as TABLES IN SCHEMA |
994 | | * that includes this relation |
995 | | */ |
996 | 0 | if (server_version >= 150000) |
997 | 0 | { |
998 | | /* Reuse the already-built pub_names. */ |
999 | 0 | Assert(pub_names != NULL); |
1000 | | |
1001 | | /* Check for row filters. */ |
1002 | 0 | resetStringInfo(&cmd); |
1003 | |
|
1004 | 0 | if (server_version >= 190000) |
1005 | 0 | { |
1006 | | /* |
1007 | | * We can pass both publication names and relid to |
1008 | | * pg_get_publication_tables() since version 19. |
1009 | | */ |
1010 | 0 | appendStringInfo(&cmd, |
1011 | 0 | "SELECT DISTINCT pg_get_expr(gpt.qual, gpt.relid)" |
1012 | 0 | " FROM pg_get_publication_tables(ARRAY[%s], %u) gpt", |
1013 | 0 | pub_names->data, |
1014 | 0 | lrel->remoteid); |
1015 | 0 | } |
1016 | 0 | else |
1017 | 0 | appendStringInfo(&cmd, |
1018 | 0 | "SELECT DISTINCT pg_get_expr(gpt.qual, gpt.relid)" |
1019 | 0 | " FROM pg_publication p," |
1020 | 0 | " LATERAL pg_get_publication_tables(p.pubname) gpt" |
1021 | 0 | " WHERE gpt.relid = %u" |
1022 | 0 | " AND p.pubname IN ( %s )", |
1023 | 0 | lrel->remoteid, |
1024 | 0 | pub_names->data); |
1025 | |
|
1026 | 0 | res = walrcv_exec(LogRepWorkerWalRcvConn, cmd.data, 1, qualRow); |
1027 | |
|
1028 | 0 | if (res->status != WALRCV_OK_TUPLES) |
1029 | 0 | ereport(ERROR, |
1030 | 0 | (errmsg("could not fetch table WHERE clause info for table \"%s.%s\" from publisher: %s", |
1031 | 0 | nspname, relname, res->err))); |
1032 | | |
1033 | | /* |
1034 | | * Multiple row filter expressions for the same table will be combined |
1035 | | * by COPY using OR. If any of the filter expressions for this table |
1036 | | * are null, it means the whole table will be copied. In this case it |
1037 | | * is not necessary to construct a unified row filter expression at |
1038 | | * all. |
1039 | | */ |
1040 | 0 | slot = MakeSingleTupleTableSlot(res->tupledesc, &TTSOpsMinimalTuple); |
1041 | 0 | while (tuplestore_gettupleslot(res->tuplestore, true, false, slot)) |
1042 | 0 | { |
1043 | 0 | Datum rf = slot_getattr(slot, 1, &isnull); |
1044 | |
|
1045 | 0 | if (!isnull) |
1046 | 0 | *qual = lappend(*qual, makeString(TextDatumGetCString(rf))); |
1047 | 0 | else |
1048 | 0 | { |
1049 | | /* Ignore filters and cleanup as necessary. */ |
1050 | 0 | if (*qual) |
1051 | 0 | { |
1052 | 0 | list_free_deep(*qual); |
1053 | 0 | *qual = NIL; |
1054 | 0 | } |
1055 | 0 | break; |
1056 | 0 | } |
1057 | | |
1058 | 0 | ExecClearTuple(slot); |
1059 | 0 | } |
1060 | 0 | ExecDropSingleTupleTableSlot(slot); |
1061 | |
|
1062 | 0 | walrcv_clear_result(res); |
1063 | 0 | destroyStringInfo(pub_names); |
1064 | 0 | } |
1065 | | |
1066 | 0 | pfree(cmd.data); |
1067 | 0 | } |
1068 | | |
1069 | | /* |
1070 | | * Copy existing data of a table from publisher. |
1071 | | * |
1072 | | * Caller is responsible for locking the local relation. |
1073 | | */ |
1074 | | static void |
1075 | | copy_table(Relation rel) |
1076 | 0 | { |
1077 | 0 | LogicalRepRelMapEntry *relmapentry; |
1078 | 0 | LogicalRepRelation lrel; |
1079 | 0 | List *qual = NIL; |
1080 | 0 | WalRcvExecResult *res; |
1081 | 0 | StringInfoData cmd; |
1082 | 0 | CopyFromState cstate; |
1083 | 0 | List *attnamelist; |
1084 | 0 | ParseState *pstate; |
1085 | 0 | List *options = NIL; |
1086 | 0 | bool gencol_published = false; |
1087 | 0 | int server_version = walrcv_server_version(LogRepWorkerWalRcvConn); |
1088 | | |
1089 | | /* Get the publisher relation info. */ |
1090 | 0 | fetch_remote_table_info(get_namespace_name(RelationGetNamespace(rel)), |
1091 | 0 | RelationGetRelationName(rel), &lrel, &qual, |
1092 | 0 | &gencol_published); |
1093 | | |
1094 | | /* Put the relation into relmap. */ |
1095 | 0 | logicalrep_relmap_update(&lrel); |
1096 | | |
1097 | | /* Map the publisher relation to local one. */ |
1098 | 0 | relmapentry = logicalrep_rel_open(lrel.remoteid, NoLock); |
1099 | 0 | Assert(rel == relmapentry->localrel); |
1100 | | |
1101 | | /* Start copy on the publisher. */ |
1102 | 0 | initStringInfo(&cmd); |
1103 | | |
1104 | | /* |
1105 | | * Regular or partitioned table with no row filter or generated columns. |
1106 | | * |
1107 | | * "COPY table TO" on a partitioned table is supported since v19. |
1108 | | */ |
1109 | 0 | if ((lrel.relkind == RELKIND_RELATION || |
1110 | 0 | (lrel.relkind == RELKIND_PARTITIONED_TABLE && server_version >= 190000)) && |
1111 | 0 | qual == NIL && !gencol_published) |
1112 | 0 | { |
1113 | 0 | appendStringInfo(&cmd, "COPY %s", |
1114 | 0 | quote_qualified_identifier(lrel.nspname, lrel.relname)); |
1115 | | |
1116 | | /* If the table has columns, then specify the columns */ |
1117 | 0 | if (lrel.natts) |
1118 | 0 | { |
1119 | 0 | appendStringInfoString(&cmd, " ("); |
1120 | | |
1121 | | /* |
1122 | | * XXX Do we need to list the columns in all cases? Maybe we're |
1123 | | * replicating all columns? |
1124 | | */ |
1125 | 0 | for (int i = 0; i < lrel.natts; i++) |
1126 | 0 | { |
1127 | 0 | if (i > 0) |
1128 | 0 | appendStringInfoString(&cmd, ", "); |
1129 | |
|
1130 | 0 | appendStringInfoString(&cmd, quote_identifier(lrel.attnames[i])); |
1131 | 0 | } |
1132 | |
|
1133 | 0 | appendStringInfoChar(&cmd, ')'); |
1134 | 0 | } |
1135 | |
|
1136 | 0 | appendStringInfoString(&cmd, " TO STDOUT"); |
1137 | 0 | } |
1138 | 0 | else |
1139 | 0 | { |
1140 | | /* |
1141 | | * For non-tables and tables with row filters, we need to do COPY |
1142 | | * (SELECT ...), but we can't just do SELECT * because we may need to |
1143 | | * copy only subset of columns including generated columns. For tables |
1144 | | * with any row filters, build a SELECT query with OR'ed row filters |
1145 | | * for COPY. |
1146 | | * |
1147 | | * We also need to use this same COPY (SELECT ...) syntax when |
1148 | | * generated columns are published, because copy of generated columns |
1149 | | * is not supported by the normal COPY. |
1150 | | */ |
1151 | 0 | appendStringInfoString(&cmd, "COPY (SELECT "); |
1152 | 0 | for (int i = 0; i < lrel.natts; i++) |
1153 | 0 | { |
1154 | 0 | appendStringInfoString(&cmd, quote_identifier(lrel.attnames[i])); |
1155 | 0 | if (i < lrel.natts - 1) |
1156 | 0 | appendStringInfoString(&cmd, ", "); |
1157 | 0 | } |
1158 | |
|
1159 | 0 | appendStringInfoString(&cmd, " FROM "); |
1160 | | |
1161 | | /* |
1162 | | * For regular tables, make sure we don't copy data from a child that |
1163 | | * inherits the named table as those will be copied separately. |
1164 | | */ |
1165 | 0 | if (lrel.relkind == RELKIND_RELATION) |
1166 | 0 | appendStringInfoString(&cmd, "ONLY "); |
1167 | |
|
1168 | 0 | appendStringInfoString(&cmd, quote_qualified_identifier(lrel.nspname, lrel.relname)); |
1169 | | /* list of OR'ed filters */ |
1170 | 0 | if (qual != NIL) |
1171 | 0 | { |
1172 | 0 | ListCell *lc; |
1173 | 0 | char *q = strVal(linitial(qual)); |
1174 | |
|
1175 | 0 | appendStringInfo(&cmd, " WHERE %s", q); |
1176 | 0 | for_each_from(lc, qual, 1) |
1177 | 0 | { |
1178 | 0 | q = strVal(lfirst(lc)); |
1179 | 0 | appendStringInfo(&cmd, " OR %s", q); |
1180 | 0 | } |
1181 | 0 | list_free_deep(qual); |
1182 | 0 | } |
1183 | |
|
1184 | 0 | appendStringInfoString(&cmd, ") TO STDOUT"); |
1185 | 0 | } |
1186 | | |
1187 | | /* |
1188 | | * Prior to v16, initial table synchronization will use text format even |
1189 | | * if the binary option is enabled for a subscription. |
1190 | | */ |
1191 | 0 | if (server_version >= 160000 && MySubscription->binary) |
1192 | 0 | { |
1193 | 0 | appendStringInfoString(&cmd, " WITH (FORMAT binary)"); |
1194 | 0 | options = list_make1(makeDefElem("format", |
1195 | 0 | (Node *) makeString("binary"), -1)); |
1196 | 0 | } |
1197 | |
|
1198 | 0 | res = walrcv_exec(LogRepWorkerWalRcvConn, cmd.data, 0, NULL); |
1199 | 0 | pfree(cmd.data); |
1200 | 0 | if (res->status != WALRCV_OK_COPY_OUT) |
1201 | 0 | ereport(ERROR, |
1202 | 0 | (errcode(ERRCODE_CONNECTION_FAILURE), |
1203 | 0 | errmsg("could not start initial contents copy for table \"%s.%s\": %s", |
1204 | 0 | lrel.nspname, lrel.relname, res->err))); |
1205 | 0 | walrcv_clear_result(res); |
1206 | |
|
1207 | 0 | copybuf = makeStringInfo(); |
1208 | |
|
1209 | 0 | pstate = make_parsestate(NULL); |
1210 | 0 | (void) addRangeTableEntryForRelation(pstate, rel, AccessShareLock, |
1211 | 0 | NULL, false, false); |
1212 | |
|
1213 | 0 | attnamelist = make_copy_attnamelist(relmapentry); |
1214 | 0 | cstate = BeginCopyFrom(pstate, rel, NULL, NULL, false, copy_read_data, attnamelist, options); |
1215 | | |
1216 | | /* Do the copy */ |
1217 | 0 | (void) CopyFrom(cstate); |
1218 | 0 | EndCopyFrom(cstate); |
1219 | |
|
1220 | 0 | logicalrep_rel_close(relmapentry, NoLock); |
1221 | 0 | } |
1222 | | |
1223 | | /* |
1224 | | * Determine the tablesync slot name. |
1225 | | * |
1226 | | * The name must not exceed NAMEDATALEN - 1 because of remote node constraints |
1227 | | * on slot name length. We append system_identifier to avoid slot_name |
1228 | | * collision with subscriptions in other clusters. With the current scheme |
1229 | | * pg_%u_sync_%u_UINT64_FORMAT (3 + 10 + 6 + 10 + 20 + '\0'), the maximum |
1230 | | * length of slot_name will be 50. |
1231 | | * |
1232 | | * The returned slot name is stored in the supplied buffer (syncslotname) with |
1233 | | * the given size. |
1234 | | * |
1235 | | * Note: We don't use the subscription slot name as part of tablesync slot name |
1236 | | * because we are responsible for cleaning up these slots and it could become |
1237 | | * impossible to recalculate what name to cleanup if the subscription slot name |
1238 | | * had changed. |
1239 | | */ |
1240 | | void |
1241 | | ReplicationSlotNameForTablesync(Oid suboid, Oid relid, |
1242 | | char *syncslotname, Size szslot) |
1243 | 0 | { |
1244 | 0 | snprintf(syncslotname, szslot, "pg_%u_sync_%u_" UINT64_FORMAT, suboid, |
1245 | 0 | relid, GetSystemIdentifier()); |
1246 | 0 | } |
1247 | | |
1248 | | /* |
1249 | | * Start syncing the table in the sync worker. |
1250 | | * |
1251 | | * If nothing needs to be done to sync the table, we exit the worker without |
1252 | | * any further action. |
1253 | | * |
1254 | | * The returned slot name is palloc'ed in current memory context. |
1255 | | */ |
1256 | | static char * |
1257 | | LogicalRepSyncTableStart(XLogRecPtr *origin_startpos) |
1258 | | { |
1259 | | char *slotname; |
1260 | | char *err; |
1261 | | char relstate; |
1262 | | XLogRecPtr relstate_lsn; |
1263 | | Relation rel; |
1264 | | AclResult aclresult; |
1265 | | WalRcvExecResult *res; |
1266 | | char originname[NAMEDATALEN]; |
1267 | | ReplOriginId originid; |
1268 | | UserContext ucxt; |
1269 | | bool must_use_password; |
1270 | | bool run_as_owner; |
1271 | | |
1272 | | /* Check the state of the table synchronization. */ |
1273 | | StartTransactionCommand(); |
1274 | | relstate = GetSubscriptionRelState(MyLogicalRepWorker->subid, |
1275 | | MyLogicalRepWorker->relid, |
1276 | | &relstate_lsn); |
1277 | | CommitTransactionCommand(); |
1278 | | |
1279 | | /* Is the use of a password mandatory? */ |
1280 | | must_use_password = MySubscription->passwordrequired && |
1281 | | !MySubscription->ownersuperuser; |
1282 | | |
1283 | | SpinLockAcquire(&MyLogicalRepWorker->relmutex); |
1284 | | MyLogicalRepWorker->relstate = relstate; |
1285 | | MyLogicalRepWorker->relstate_lsn = relstate_lsn; |
1286 | | SpinLockRelease(&MyLogicalRepWorker->relmutex); |
1287 | | |
1288 | | /* |
1289 | | * If synchronization is already done or no longer necessary, exit now |
1290 | | * that we've updated shared memory state. |
1291 | | */ |
1292 | | switch (relstate) |
1293 | | { |
1294 | | case SUBREL_STATE_SYNCDONE: |
1295 | | case SUBREL_STATE_READY: |
1296 | | case SUBREL_STATE_UNKNOWN: |
1297 | | FinishSyncWorker(); /* doesn't return */ |
1298 | | } |
1299 | | |
1300 | | /* Calculate the name of the tablesync slot. */ |
1301 | | slotname = (char *) palloc(NAMEDATALEN); |
1302 | | ReplicationSlotNameForTablesync(MySubscription->oid, |
1303 | | MyLogicalRepWorker->relid, |
1304 | | slotname, |
1305 | | NAMEDATALEN); |
1306 | | |
1307 | | /* |
1308 | | * Here we use the slot name instead of the subscription name as the |
1309 | | * application_name, so that it is different from the leader apply worker, |
1310 | | * so that synchronous replication can distinguish them. |
1311 | | */ |
1312 | | LogRepWorkerWalRcvConn = |
1313 | | walrcv_connect(MySubscriptionConninfo, true, true, |
1314 | | must_use_password, |
1315 | | slotname, &err); |
1316 | | if (LogRepWorkerWalRcvConn == NULL) |
1317 | | ereport(ERROR, |
1318 | | (errcode(ERRCODE_CONNECTION_FAILURE), |
1319 | | errmsg("table synchronization worker for subscription \"%s\" could not connect to the publisher: %s", |
1320 | | MySubscription->name, err))); |
1321 | | |
1322 | | Assert(MyLogicalRepWorker->relstate == SUBREL_STATE_INIT || |
1323 | | MyLogicalRepWorker->relstate == SUBREL_STATE_DATASYNC || |
1324 | | MyLogicalRepWorker->relstate == SUBREL_STATE_FINISHEDCOPY); |
1325 | | |
1326 | | /* Assign the origin tracking record name. */ |
1327 | | ReplicationOriginNameForLogicalRep(MySubscription->oid, |
1328 | | MyLogicalRepWorker->relid, |
1329 | | originname, |
1330 | | sizeof(originname)); |
1331 | | |
1332 | | if (MyLogicalRepWorker->relstate == SUBREL_STATE_DATASYNC) |
1333 | | { |
1334 | | /* |
1335 | | * We have previously errored out before finishing the copy so the |
1336 | | * replication slot might exist. We want to remove the slot if it |
1337 | | * already exists and proceed. |
1338 | | * |
1339 | | * XXX We could also instead try to drop the slot, last time we failed |
1340 | | * but for that, we might need to clean up the copy state as it might |
1341 | | * be in the middle of fetching the rows. Also, if there is a network |
1342 | | * breakdown then it wouldn't have succeeded so trying it next time |
1343 | | * seems like a better bet. |
1344 | | */ |
1345 | | ReplicationSlotDropAtPubNode(LogRepWorkerWalRcvConn, slotname, true); |
1346 | | } |
1347 | | else if (MyLogicalRepWorker->relstate == SUBREL_STATE_FINISHEDCOPY) |
1348 | | { |
1349 | | /* |
1350 | | * The COPY phase was previously done, but tablesync then crashed |
1351 | | * before it was able to finish normally. |
1352 | | */ |
1353 | | StartTransactionCommand(); |
1354 | | |
1355 | | /* |
1356 | | * The origin tracking name must already exist. It was created first |
1357 | | * time this tablesync was launched. |
1358 | | */ |
1359 | | originid = replorigin_by_name(originname, false); |
1360 | | replorigin_session_setup(originid, 0); |
1361 | | replorigin_xact_state.origin = originid; |
1362 | | *origin_startpos = replorigin_session_get_progress(false); |
1363 | | |
1364 | | CommitTransactionCommand(); |
1365 | | |
1366 | | goto copy_table_done; |
1367 | | } |
1368 | | |
1369 | | SpinLockAcquire(&MyLogicalRepWorker->relmutex); |
1370 | | MyLogicalRepWorker->relstate = SUBREL_STATE_DATASYNC; |
1371 | | MyLogicalRepWorker->relstate_lsn = InvalidXLogRecPtr; |
1372 | | SpinLockRelease(&MyLogicalRepWorker->relmutex); |
1373 | | |
1374 | | /* |
1375 | | * Update the state, create the replication origin, and make them visible |
1376 | | * to others. |
1377 | | */ |
1378 | | StartTransactionCommand(); |
1379 | | UpdateSubscriptionRelState(MyLogicalRepWorker->subid, |
1380 | | MyLogicalRepWorker->relid, |
1381 | | MyLogicalRepWorker->relstate, |
1382 | | MyLogicalRepWorker->relstate_lsn, |
1383 | | false); |
1384 | | |
1385 | | /* |
1386 | | * Create the replication origin in a separate transaction from the one |
1387 | | * that sets up the origin in shared memory. This prevents the risk that |
1388 | | * changes to the origin in shared memory cannot be rolled back if the |
1389 | | * transaction aborts. |
1390 | | */ |
1391 | | originid = replorigin_by_name(originname, true); |
1392 | | if (!OidIsValid(originid)) |
1393 | | originid = replorigin_create(originname); |
1394 | | |
1395 | | CommitTransactionCommand(); |
1396 | | pgstat_report_stat(true); |
1397 | | |
1398 | | StartTransactionCommand(); |
1399 | | |
1400 | | /* |
1401 | | * Use a standard write lock here. It might be better to disallow access |
1402 | | * to the table while it's being synchronized. But we don't want to block |
1403 | | * the main apply process from working and it has to open the relation in |
1404 | | * RowExclusiveLock when remapping remote relation id to local one. |
1405 | | */ |
1406 | | rel = table_open(MyLogicalRepWorker->relid, RowExclusiveLock); |
1407 | | |
1408 | | /* |
1409 | | * Start a transaction in the remote node in REPEATABLE READ mode. This |
1410 | | * ensures that both the replication slot we create (see below) and the |
1411 | | * COPY are consistent with each other. |
1412 | | */ |
1413 | | res = walrcv_exec(LogRepWorkerWalRcvConn, |
1414 | | "BEGIN READ ONLY ISOLATION LEVEL REPEATABLE READ", |
1415 | | 0, NULL); |
1416 | | if (res->status != WALRCV_OK_COMMAND) |
1417 | | ereport(ERROR, |
1418 | | (errcode(ERRCODE_CONNECTION_FAILURE), |
1419 | | errmsg("table copy could not start transaction on publisher: %s", |
1420 | | res->err))); |
1421 | | walrcv_clear_result(res); |
1422 | | |
1423 | | /* |
1424 | | * Create a new permanent logical decoding slot. This slot will be used |
1425 | | * for the catchup phase after COPY is done, so tell it to use the |
1426 | | * snapshot to make the final data consistent. |
1427 | | */ |
1428 | | walrcv_create_slot(LogRepWorkerWalRcvConn, |
1429 | | slotname, false /* permanent */ , false /* two_phase */ , |
1430 | | MySubscription->failover, |
1431 | | CRS_USE_SNAPSHOT, origin_startpos); |
1432 | | |
1433 | | /* |
1434 | | * Advance the origin to the LSN got from walrcv_create_slot and then set |
1435 | | * up the origin. The advancement is WAL logged for the purpose of |
1436 | | * recovery. Locks are to prevent the replication origin from vanishing |
1437 | | * while advancing. |
1438 | | * |
1439 | | * The purpose of doing these before the copy is to avoid doing the copy |
1440 | | * again due to any error in advancing or setting up origin tracking. |
1441 | | */ |
1442 | | LockRelationOid(ReplicationOriginRelationId, RowExclusiveLock); |
1443 | | replorigin_advance(originid, *origin_startpos, InvalidXLogRecPtr, |
1444 | | true /* go backward */ , true /* WAL log */ ); |
1445 | | UnlockRelationOid(ReplicationOriginRelationId, RowExclusiveLock); |
1446 | | |
1447 | | replorigin_session_setup(originid, 0); |
1448 | | replorigin_xact_state.origin = originid; |
1449 | | |
1450 | | /* |
1451 | | * If the user did not opt to run as the owner of the subscription |
1452 | | * ('run_as_owner'), then copy the table as the owner of the table. |
1453 | | */ |
1454 | | run_as_owner = MySubscription->runasowner; |
1455 | | if (!run_as_owner) |
1456 | | SwitchToUntrustedUser(rel->rd_rel->relowner, &ucxt); |
1457 | | |
1458 | | /* |
1459 | | * Check that our table sync worker has permission to insert into the |
1460 | | * target table. |
1461 | | */ |
1462 | | aclresult = pg_class_aclcheck(RelationGetRelid(rel), GetUserId(), |
1463 | | ACL_INSERT); |
1464 | | if (aclresult != ACLCHECK_OK) |
1465 | | aclcheck_error(aclresult, |
1466 | | get_relkind_objtype(rel->rd_rel->relkind), |
1467 | | RelationGetRelationName(rel)); |
1468 | | |
1469 | | /* |
1470 | | * COPY FROM does not honor RLS policies. That is not a problem for |
1471 | | * subscriptions owned by roles with BYPASSRLS privilege (or superuser, |
1472 | | * who has it implicitly), but other roles should not be able to |
1473 | | * circumvent RLS. Disallow logical replication into RLS enabled |
1474 | | * relations for such roles. |
1475 | | */ |
1476 | | if (check_enable_rls(RelationGetRelid(rel), InvalidOid, false) == RLS_ENABLED) |
1477 | | ereport(ERROR, |
1478 | | (errcode(ERRCODE_FEATURE_NOT_SUPPORTED), |
1479 | | errmsg("user \"%s\" cannot replicate into relation with row-level security enabled: \"%s\"", |
1480 | | GetUserNameFromId(GetUserId(), true), |
1481 | | RelationGetRelationName(rel)))); |
1482 | | |
1483 | | /* Now do the initial data copy */ |
1484 | | PushActiveSnapshot(GetTransactionSnapshot()); |
1485 | | copy_table(rel); |
1486 | | PopActiveSnapshot(); |
1487 | | |
1488 | | res = walrcv_exec(LogRepWorkerWalRcvConn, "COMMIT", 0, NULL); |
1489 | | if (res->status != WALRCV_OK_COMMAND) |
1490 | | ereport(ERROR, |
1491 | | (errcode(ERRCODE_CONNECTION_FAILURE), |
1492 | | errmsg("table copy could not finish transaction on publisher: %s", |
1493 | | res->err))); |
1494 | | walrcv_clear_result(res); |
1495 | | |
1496 | | if (!run_as_owner) |
1497 | | RestoreUserContext(&ucxt); |
1498 | | |
1499 | | table_close(rel, NoLock); |
1500 | | |
1501 | | /* Make the copy visible. */ |
1502 | | CommandCounterIncrement(); |
1503 | | |
1504 | | /* |
1505 | | * Update the persisted state to indicate the COPY phase is done; make it |
1506 | | * visible to others. |
1507 | | */ |
1508 | | UpdateSubscriptionRelState(MyLogicalRepWorker->subid, |
1509 | | MyLogicalRepWorker->relid, |
1510 | | SUBREL_STATE_FINISHEDCOPY, |
1511 | | MyLogicalRepWorker->relstate_lsn, |
1512 | | false); |
1513 | | |
1514 | | CommitTransactionCommand(); |
1515 | | |
1516 | | copy_table_done: |
1517 | | |
1518 | | elog(DEBUG1, |
1519 | | "LogicalRepSyncTableStart: '%s' origin_startpos lsn %X/%08X", |
1520 | | originname, LSN_FORMAT_ARGS(*origin_startpos)); |
1521 | | |
1522 | | /* |
1523 | | * We are done with the initial data synchronization, update the state. |
1524 | | */ |
1525 | | SpinLockAcquire(&MyLogicalRepWorker->relmutex); |
1526 | | MyLogicalRepWorker->relstate = SUBREL_STATE_SYNCWAIT; |
1527 | | MyLogicalRepWorker->relstate_lsn = *origin_startpos; |
1528 | | SpinLockRelease(&MyLogicalRepWorker->relmutex); |
1529 | | |
1530 | | /* |
1531 | | * Finally, wait until the leader apply worker tells us to catch up and |
1532 | | * then return to let LogicalRepApplyLoop do it. |
1533 | | */ |
1534 | | wait_for_worker_state_change(SUBREL_STATE_CATCHUP); |
1535 | | return slotname; |
1536 | | } |
1537 | | |
1538 | | /* |
1539 | | * Execute the initial sync with error handling. Disable the subscription, |
1540 | | * if it's required. |
1541 | | * |
1542 | | * Allocate the slot name in long-lived context on return. Note that we don't |
1543 | | * handle FATAL errors which are probably because of system resource error and |
1544 | | * are not repeatable. |
1545 | | */ |
1546 | | static void |
1547 | | start_table_sync(XLogRecPtr *origin_startpos, char **slotname) |
1548 | 0 | { |
1549 | 0 | char *sync_slotname = NULL; |
1550 | |
|
1551 | 0 | Assert(am_tablesync_worker()); |
1552 | |
|
1553 | 0 | PG_TRY(); |
1554 | 0 | { |
1555 | | /* Call initial sync. */ |
1556 | 0 | sync_slotname = LogicalRepSyncTableStart(origin_startpos); |
1557 | 0 | } |
1558 | 0 | PG_CATCH(); |
1559 | 0 | { |
1560 | 0 | if (MySubscription->disableonerr) |
1561 | 0 | DisableSubscriptionAndExit(); |
1562 | 0 | else |
1563 | 0 | { |
1564 | | /* |
1565 | | * Report the worker failed during table synchronization. Abort |
1566 | | * the current transaction so that the stats message is sent in an |
1567 | | * idle state. |
1568 | | */ |
1569 | 0 | AbortOutOfAnyTransaction(); |
1570 | 0 | pgstat_report_subscription_error(MySubscription->oid); |
1571 | |
|
1572 | 0 | PG_RE_THROW(); |
1573 | 0 | } |
1574 | 0 | } |
1575 | 0 | PG_END_TRY(); |
1576 | | |
1577 | | /* allocate slot name in long-lived context */ |
1578 | 0 | *slotname = MemoryContextStrdup(ApplyContext, sync_slotname); |
1579 | 0 | pfree(sync_slotname); |
1580 | 0 | } |
1581 | | |
1582 | | /* |
1583 | | * Runs the tablesync worker. |
1584 | | * |
1585 | | * It starts syncing tables. After a successful sync, sets streaming options |
1586 | | * and starts streaming to catchup with apply worker. |
1587 | | */ |
1588 | | static void |
1589 | | run_tablesync_worker(void) |
1590 | 0 | { |
1591 | 0 | char originname[NAMEDATALEN]; |
1592 | 0 | XLogRecPtr origin_startpos = InvalidXLogRecPtr; |
1593 | 0 | char *slotname = NULL; |
1594 | 0 | WalRcvStreamOptions options; |
1595 | |
|
1596 | 0 | start_table_sync(&origin_startpos, &slotname); |
1597 | |
|
1598 | 0 | ReplicationOriginNameForLogicalRep(MySubscription->oid, |
1599 | 0 | MyLogicalRepWorker->relid, |
1600 | 0 | originname, |
1601 | 0 | sizeof(originname)); |
1602 | |
|
1603 | 0 | set_apply_error_context_origin(originname); |
1604 | |
|
1605 | 0 | set_stream_options(&options, slotname, &origin_startpos); |
1606 | |
|
1607 | 0 | walrcv_startstreaming(LogRepWorkerWalRcvConn, &options); |
1608 | | |
1609 | | /* Apply the changes till we catchup with the apply worker. */ |
1610 | 0 | start_apply(origin_startpos); |
1611 | 0 | } |
1612 | | |
1613 | | /* Logical Replication Tablesync worker entry point */ |
1614 | | void |
1615 | | TableSyncWorkerMain(Datum main_arg) |
1616 | 0 | { |
1617 | 0 | int worker_slot = DatumGetInt32(main_arg); |
1618 | |
|
1619 | 0 | SetupApplyOrSyncWorker(worker_slot); |
1620 | |
|
1621 | 0 | run_tablesync_worker(); |
1622 | |
|
1623 | 0 | FinishSyncWorker(); |
1624 | 0 | } |
1625 | | |
1626 | | /* |
1627 | | * If the subscription has no tables then return false. |
1628 | | * |
1629 | | * Otherwise, are all tablesyncs READY? |
1630 | | * |
1631 | | * Note: This function is not suitable to be called from outside of apply or |
1632 | | * tablesync workers because MySubscription needs to be already initialized. |
1633 | | */ |
1634 | | bool |
1635 | | AllTablesyncsReady(void) |
1636 | 0 | { |
1637 | 0 | bool started_tx; |
1638 | 0 | bool has_tables; |
1639 | | |
1640 | | /* We need up-to-date sync state info for subscription tables here. */ |
1641 | 0 | FetchRelationStates(&has_tables, NULL, &started_tx); |
1642 | |
|
1643 | 0 | if (started_tx) |
1644 | 0 | { |
1645 | 0 | CommitTransactionCommand(); |
1646 | 0 | pgstat_report_stat(true); |
1647 | 0 | } |
1648 | | |
1649 | | /* |
1650 | | * Return false when there are no tables in subscription or not all tables |
1651 | | * are in ready state; true otherwise. |
1652 | | */ |
1653 | 0 | return has_tables && (table_states_not_ready == NIL); |
1654 | 0 | } |
1655 | | |
1656 | | /* |
1657 | | * Return whether the subscription currently has any tables. |
1658 | | * |
1659 | | * Note: Unlike HasSubscriptionTables(), this function relies on cached |
1660 | | * information for subscription tables. Additionally, it should not be |
1661 | | * invoked outside of apply or tablesync workers, as MySubscription must be |
1662 | | * initialized first. |
1663 | | */ |
1664 | | bool |
1665 | | HasSubscriptionTablesCached(void) |
1666 | 0 | { |
1667 | 0 | bool started_tx; |
1668 | 0 | bool has_tables; |
1669 | | |
1670 | | /* We need up-to-date subscription tables info here */ |
1671 | 0 | FetchRelationStates(&has_tables, NULL, &started_tx); |
1672 | |
|
1673 | 0 | if (started_tx) |
1674 | 0 | { |
1675 | 0 | CommitTransactionCommand(); |
1676 | 0 | pgstat_report_stat(true); |
1677 | 0 | } |
1678 | |
|
1679 | 0 | return has_tables; |
1680 | 0 | } |
1681 | | |
1682 | | /* |
1683 | | * Update the two_phase state of the specified subscription in pg_subscription. |
1684 | | */ |
1685 | | void |
1686 | | UpdateTwoPhaseState(Oid suboid, char new_state) |
1687 | 0 | { |
1688 | 0 | Relation rel; |
1689 | 0 | HeapTuple tup; |
1690 | 0 | bool nulls[Natts_pg_subscription]; |
1691 | 0 | bool replaces[Natts_pg_subscription]; |
1692 | 0 | Datum values[Natts_pg_subscription]; |
1693 | |
|
1694 | 0 | Assert(new_state == LOGICALREP_TWOPHASE_STATE_DISABLED || |
1695 | 0 | new_state == LOGICALREP_TWOPHASE_STATE_PENDING || |
1696 | 0 | new_state == LOGICALREP_TWOPHASE_STATE_ENABLED); |
1697 | |
|
1698 | 0 | rel = table_open(SubscriptionRelationId, RowExclusiveLock); |
1699 | 0 | tup = SearchSysCacheCopy1(SUBSCRIPTIONOID, ObjectIdGetDatum(suboid)); |
1700 | 0 | if (!HeapTupleIsValid(tup)) |
1701 | 0 | elog(ERROR, |
1702 | 0 | "cache lookup failed for subscription oid %u", |
1703 | 0 | suboid); |
1704 | | |
1705 | | /* Form a new tuple. */ |
1706 | 0 | memset(values, 0, sizeof(values)); |
1707 | 0 | memset(nulls, false, sizeof(nulls)); |
1708 | 0 | memset(replaces, false, sizeof(replaces)); |
1709 | | |
1710 | | /* And update/set two_phase state */ |
1711 | 0 | values[Anum_pg_subscription_subtwophasestate - 1] = CharGetDatum(new_state); |
1712 | 0 | replaces[Anum_pg_subscription_subtwophasestate - 1] = true; |
1713 | |
|
1714 | 0 | tup = heap_modify_tuple(tup, RelationGetDescr(rel), |
1715 | 0 | values, nulls, replaces); |
1716 | 0 | CatalogTupleUpdate(rel, &tup->t_self, tup); |
1717 | |
|
1718 | 0 | heap_freetuple(tup); |
1719 | 0 | table_close(rel, RowExclusiveLock); |
1720 | 0 | } |