Coverage Report

Created: 2026-09-28 06:55

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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, &copybuf->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, &copybuf->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
}