LCOV - code coverage report
Current view: top level - src/backend/executor - execReplication.c (source / functions) Coverage Total Hit
Test: PostgreSQL 20devel Lines: 85.8 % 345 296
Test Date: 2026-10-02 20:15:47 Functions: 100.0 % 16 16
Legend: Lines:     hit not hit
Branches: + taken - not taken # not executed
Branches: 69.2 % 263 182

             Branch data     Line data    Source code
       1                 :             : /*-------------------------------------------------------------------------
       2                 :             :  *
       3                 :             :  * execReplication.c
       4                 :             :  *    miscellaneous executor routines for logical replication
       5                 :             :  *
       6                 :             :  * Portions Copyright (c) 1996-2026, PostgreSQL Global Development Group
       7                 :             :  * Portions Copyright (c) 1994, Regents of the University of California
       8                 :             :  *
       9                 :             :  * IDENTIFICATION
      10                 :             :  *    src/backend/executor/execReplication.c
      11                 :             :  *
      12                 :             :  *-------------------------------------------------------------------------
      13                 :             :  */
      14                 :             : 
      15                 :             : #include "postgres.h"
      16                 :             : 
      17                 :             : #include "access/amapi.h"
      18                 :             : #include "access/commit_ts.h"
      19                 :             : #include "access/genam.h"
      20                 :             : #include "access/gist.h"
      21                 :             : #include "access/relscan.h"
      22                 :             : #include "access/tableam.h"
      23                 :             : #include "access/transam.h"
      24                 :             : #include "access/xact.h"
      25                 :             : #include "access/heapam.h"
      26                 :             : #include "catalog/pg_am_d.h"
      27                 :             : #include "commands/trigger.h"
      28                 :             : #include "executor/executor.h"
      29                 :             : #include "executor/nodeModifyTable.h"
      30                 :             : #include "replication/conflict.h"
      31                 :             : #include "replication/logicalrelation.h"
      32                 :             : #include "storage/lmgr.h"
      33                 :             : #include "utils/builtins.h"
      34                 :             : #include "utils/lsyscache.h"
      35                 :             : #include "utils/rel.h"
      36                 :             : #include "utils/snapmgr.h"
      37                 :             : #include "utils/syscache.h"
      38                 :             : #include "utils/typcache.h"
      39                 :             : 
      40                 :             : 
      41                 :             : static bool tuples_equal(TupleTableSlot *slot1, TupleTableSlot *slot2,
      42                 :             :                          TypeCacheEntry **eq, Bitmapset *columns);
      43                 :             : 
      44                 :             : /*
      45                 :             :  * Setup a ScanKey for a search in the relation 'rel' for a tuple 'key' that
      46                 :             :  * is setup to match 'rel' (*NOT* idxrel!).
      47                 :             :  *
      48                 :             :  * Returns how many columns to use for the index scan.
      49                 :             :  *
      50                 :             :  * This is not a generic routine, idxrel must be PK, RI, or an index that can be
      51                 :             :  * used for a REPLICA IDENTITY FULL table. See FindUsableIndexForReplicaIdentityFull()
      52                 :             :  * for details.
      53                 :             :  *
      54                 :             :  * By definition, replication identity of a rel meets all limitations associated
      55                 :             :  * with that. Note that any other index could also meet these limitations.
      56                 :             :  */
      57                 :             : static int
      58                 :       72154 : build_replindex_scan_key(ScanKey skey, Relation rel, Relation idxrel,
      59                 :             :                          TupleTableSlot *searchslot)
      60                 :             : {
      61                 :             :     int         index_attoff;
      62                 :       72154 :     int         skey_attoff = 0;
      63                 :             :     Datum       indclassDatum;
      64                 :             :     oidvector  *opclass;
      65                 :       72154 :     int2vector *indkey = &idxrel->rd_index->indkey;
      66                 :             : 
      67                 :       72154 :     indclassDatum = SysCacheGetAttrNotNull(INDEXRELID, idxrel->rd_indextuple,
      68                 :             :                                            Anum_pg_index_indclass);
      69                 :       72154 :     opclass = (oidvector *) DatumGetPointer(indclassDatum);
      70                 :             : 
      71                 :             :     /* Build scankey for every non-expression attribute in the index. */
      72         [ +  + ]:      144331 :     for (index_attoff = 0; index_attoff < IndexRelationGetNumberOfKeyAttributes(idxrel);
      73                 :       72177 :          index_attoff++)
      74                 :             :     {
      75                 :             :         Oid         operator;
      76                 :             :         Oid         optype;
      77                 :             :         Oid         opfamily;
      78                 :             :         RegProcedure regop;
      79                 :       72177 :         int         table_attno = indkey->values[index_attoff];
      80                 :             :         StrategyNumber eq_strategy;
      81                 :             : 
      82         [ +  + ]:       72177 :         if (!AttributeNumberIsValid(table_attno))
      83                 :             :         {
      84                 :             :             /*
      85                 :             :              * XXX: Currently, we don't support expressions in the scan key,
      86                 :             :              * see code below.
      87                 :             :              */
      88                 :           2 :             continue;
      89                 :             :         }
      90                 :             : 
      91                 :             :         /*
      92                 :             :          * Load the operator info.  We need this to get the equality operator
      93                 :             :          * function for the scan key.
      94                 :             :          */
      95                 :       72175 :         optype = get_opclass_input_type(opclass->values[index_attoff]);
      96                 :       72175 :         opfamily = get_opclass_family(opclass->values[index_attoff]);
      97                 :       72175 :         eq_strategy = IndexAmTranslateCompareType(COMPARE_EQ, idxrel->rd_rel->relam, opfamily, false);
      98                 :       72175 :         operator = get_opfamily_member(opfamily, optype,
      99                 :             :                                        optype,
     100                 :             :                                        eq_strategy);
     101                 :             : 
     102         [ -  + ]:       72175 :         if (!OidIsValid(operator))
     103         [ #  # ]:           0 :             elog(ERROR, "missing operator %d(%u,%u) in opfamily %u",
     104                 :             :                  eq_strategy, optype, optype, opfamily);
     105                 :             : 
     106                 :       72175 :         regop = get_opcode(operator);
     107                 :             : 
     108                 :             :         /* Initialize the scankey. */
     109                 :       72175 :         ScanKeyInit(&skey[skey_attoff],
     110                 :       72175 :                     index_attoff + 1,
     111                 :             :                     eq_strategy,
     112                 :             :                     regop,
     113                 :       72175 :                     searchslot->tts_values[table_attno - 1]);
     114                 :             : 
     115                 :       72175 :         skey[skey_attoff].sk_collation = idxrel->rd_indcollation[index_attoff];
     116                 :             : 
     117                 :             :         /* Check for null value. */
     118         [ +  + ]:       72175 :         if (searchslot->tts_isnull[table_attno - 1])
     119                 :           1 :             skey[skey_attoff].sk_flags |= (SK_ISNULL | SK_SEARCHNULL);
     120                 :             : 
     121                 :       72175 :         skey_attoff++;
     122                 :             :     }
     123                 :             : 
     124                 :             :     /* There must always be at least one attribute for the index scan. */
     125                 :             :     Assert(skey_attoff > 0);
     126                 :             : 
     127                 :       72154 :     return skey_attoff;
     128                 :             : }
     129                 :             : 
     130                 :             : 
     131                 :             : /*
     132                 :             :  * Helper function to check if it is necessary to re-fetch and lock the tuple
     133                 :             :  * due to concurrent modifications. This function should be called after
     134                 :             :  * invoking table_tuple_lock.
     135                 :             :  */
     136                 :             : static bool
     137                 :       72351 : should_refetch_tuple(TM_Result res, TM_FailureData *tmfd)
     138                 :             : {
     139                 :       72351 :     bool        refetch = false;
     140                 :             : 
     141   [ +  -  -  -  :       72351 :     switch (res)
                      - ]
     142                 :             :     {
     143                 :       72351 :         case TM_Ok:
     144                 :       72351 :             break;
     145                 :           0 :         case TM_Updated:
     146                 :             :             /* XXX: Improve handling here */
     147         [ #  # ]:           0 :             if (ItemPointerIndicatesMovedPartitions(&tmfd->ctid))
     148         [ #  # ]:           0 :                 ereport(LOG,
     149                 :             :                         (errcode(ERRCODE_T_R_SERIALIZATION_FAILURE),
     150                 :             :                          errmsg("tuple to be locked was already moved to another partition due to concurrent update, retrying")));
     151                 :             :             else
     152         [ #  # ]:           0 :                 ereport(LOG,
     153                 :             :                         (errcode(ERRCODE_T_R_SERIALIZATION_FAILURE),
     154                 :             :                          errmsg("concurrent update, retrying")));
     155                 :           0 :             refetch = true;
     156                 :           0 :             break;
     157                 :           0 :         case TM_Deleted:
     158                 :             :             /* XXX: Improve handling here */
     159         [ #  # ]:           0 :             ereport(LOG,
     160                 :             :                     (errcode(ERRCODE_T_R_SERIALIZATION_FAILURE),
     161                 :             :                      errmsg("concurrent delete, retrying")));
     162                 :           0 :             refetch = true;
     163                 :           0 :             break;
     164                 :           0 :         case TM_Invisible:
     165         [ #  # ]:           0 :             elog(ERROR, "attempted to lock invisible tuple");
     166                 :             :             break;
     167                 :           0 :         default:
     168         [ #  # ]:           0 :             elog(ERROR, "unexpected table_tuple_lock status: %u", res);
     169                 :             :             break;
     170                 :             :     }
     171                 :             : 
     172                 :       72351 :     return refetch;
     173                 :             : }
     174                 :             : 
     175                 :             : /*
     176                 :             :  * Search the relation 'rel' for tuple using the index.
     177                 :             :  *
     178                 :             :  * If a matching tuple is found, lock it with lockmode, fill the slot with its
     179                 :             :  * contents, and return true.  Return false otherwise.
     180                 :             :  *
     181                 :             :  * 'skipduplicates' specifies whether the first matching tuple can be used
     182                 :             :  * without comparing it against 'searchslot'. If false, all matching tuples are
     183                 :             :  * compared against 'searchslot', which must contain a complete row.
     184                 :             :  */
     185                 :             : bool
     186                 :       72153 : RelationFindReplTupleByIndex(Relation rel, Oid idxoid,
     187                 :             :                              bool skipduplicates,
     188                 :             :                              LockTupleMode lockmode,
     189                 :             :                              TupleTableSlot *searchslot,
     190                 :             :                              TupleTableSlot *outslot)
     191                 :             : {
     192                 :             :     ScanKeyData skey[INDEX_MAX_KEYS];
     193                 :             :     int         skey_attoff;
     194                 :             :     IndexScanDesc scan;
     195                 :             :     SnapshotData snap;
     196                 :             :     TransactionId xwait;
     197                 :             :     Relation    idxrel;
     198                 :             :     bool        found;
     199                 :       72153 :     TypeCacheEntry **eq = NULL;
     200                 :             : 
     201                 :             :     /* Open the index. */
     202                 :       72153 :     idxrel = index_open(idxoid, RowExclusiveLock);
     203                 :             : 
     204                 :       72153 :     InitDirtySnapshot(snap);
     205                 :             : 
     206                 :             :     /* Build scan key. */
     207                 :       72153 :     skey_attoff = build_replindex_scan_key(skey, rel, idxrel, searchslot);
     208                 :             : 
     209                 :             :     /* Start an index scan. */
     210                 :       72153 :     scan = index_beginscan(rel, idxrel, false,
     211                 :             :                            &snap, NULL, skey_attoff, 0, SO_NONE);
     212                 :             : 
     213                 :           0 : retry:
     214                 :       72153 :     found = false;
     215                 :             : 
     216                 :       72153 :     index_rescan(scan, skey, skey_attoff, NULL, 0);
     217                 :             : 
     218                 :             :     /* Try to find the tuple */
     219         [ +  + ]:       72153 :     while (table_index_getnext_slot(scan, ForwardScanDirection, outslot))
     220                 :             :     {
     221                 :             :         /*
     222                 :             :          * Avoid expensive equality check if the index is primary key or
     223                 :             :          * replica identity index.
     224                 :             :          */
     225         [ +  + ]:       72091 :         if (!skipduplicates)
     226                 :             :         {
     227         [ +  - ]:          17 :             if (eq == NULL)
     228                 :          17 :                 eq = palloc0_array(TypeCacheEntry *, outslot->tts_tupleDescriptor->natts);
     229                 :             : 
     230         [ -  + ]:          17 :             if (!tuples_equal(outslot, searchslot, eq, NULL))
     231                 :           0 :                 continue;
     232                 :             :         }
     233                 :             : 
     234                 :       72091 :         ExecMaterializeSlot(outslot);
     235                 :             : 
     236                 :      144182 :         xwait = TransactionIdIsValid(snap.xmin) ?
     237         [ -  + ]:       72091 :             snap.xmin : snap.xmax;
     238                 :             : 
     239                 :             :         /*
     240                 :             :          * If the tuple is locked, wait for locking transaction to finish and
     241                 :             :          * retry.
     242                 :             :          */
     243         [ -  + ]:       72091 :         if (TransactionIdIsValid(xwait))
     244                 :             :         {
     245                 :           0 :             XactLockTableWait(xwait, NULL, NULL, XLTW_None);
     246                 :           0 :             goto retry;
     247                 :             :         }
     248                 :             : 
     249                 :             :         /* Found our tuple and it's not locked */
     250                 :       72091 :         found = true;
     251                 :       72091 :         break;
     252                 :             :     }
     253                 :             : 
     254                 :             :     /* Found tuple, try to lock it in the lockmode. */
     255         [ +  + ]:       72153 :     if (found)
     256                 :             :     {
     257                 :             :         TM_FailureData tmfd;
     258                 :             :         TM_Result   res;
     259                 :             : 
     260                 :       72091 :         PushActiveSnapshot(GetLatestSnapshot());
     261                 :             : 
     262                 :       72091 :         res = table_tuple_lock(rel, &(outslot->tts_tid), GetActiveSnapshot(),
     263                 :             :                                outslot,
     264                 :             :                                GetCurrentCommandId(false),
     265                 :             :                                lockmode,
     266                 :             :                                LockWaitBlock,
     267                 :             :                                0 /* don't follow updates */ ,
     268                 :             :                                &tmfd);
     269                 :             : 
     270                 :       72091 :         PopActiveSnapshot();
     271                 :             : 
     272         [ -  + ]:       72091 :         if (should_refetch_tuple(res, &tmfd))
     273                 :           0 :             goto retry;
     274                 :             :     }
     275                 :             : 
     276                 :       72153 :     index_endscan(scan);
     277                 :             : 
     278                 :             :     /* Don't release lock until commit. */
     279                 :       72153 :     index_close(idxrel, NoLock);
     280                 :             : 
     281                 :       72153 :     return found;
     282                 :             : }
     283                 :             : 
     284                 :             : /*
     285                 :             :  * Compare the tuples in the slots by checking if they have equal values.
     286                 :             :  *
     287                 :             :  * If 'columns' is not null, only the columns specified within it will be
     288                 :             :  * considered for the equality check, ignoring all other columns.
     289                 :             :  */
     290                 :             : static bool
     291                 :      105334 : tuples_equal(TupleTableSlot *slot1, TupleTableSlot *slot2,
     292                 :             :              TypeCacheEntry **eq, Bitmapset *columns)
     293                 :             : {
     294                 :             :     int         attrnum;
     295                 :             : 
     296                 :             :     Assert(slot1->tts_tupleDescriptor->natts ==
     297                 :             :            slot2->tts_tupleDescriptor->natts);
     298                 :             : 
     299                 :      105334 :     slot_getallattrs(slot1);
     300                 :      105334 :     slot_getallattrs(slot2);
     301                 :             : 
     302                 :             :     /* Check equality of the attributes. */
     303         [ +  + ]:      105541 :     for (attrnum = 0; attrnum < slot1->tts_tupleDescriptor->natts; attrnum++)
     304                 :             :     {
     305                 :             :         Form_pg_attribute att;
     306                 :             :         TypeCacheEntry *typentry;
     307                 :             : 
     308                 :      105374 :         att = TupleDescAttr(slot1->tts_tupleDescriptor, attrnum);
     309                 :             : 
     310                 :             :         /*
     311                 :             :          * Ignore dropped and generated columns as the publisher doesn't send
     312                 :             :          * those
     313                 :             :          */
     314   [ +  +  -  + ]:      105374 :         if (att->attisdropped || att->attgenerated)
     315                 :           1 :             continue;
     316                 :             : 
     317                 :             :         /*
     318                 :             :          * Ignore columns that are not listed for checking.
     319                 :             :          */
     320         [ -  + ]:      105373 :         if (columns &&
     321         [ #  # ]:           0 :             !bms_is_member(att->attnum - FirstLowInvalidHeapAttributeNumber,
     322                 :             :                            columns))
     323                 :           0 :             continue;
     324                 :             : 
     325                 :             :         /*
     326                 :             :          * If one value is NULL and other is not, then they are certainly not
     327                 :             :          * equal
     328                 :             :          */
     329         [ -  + ]:      105373 :         if (slot1->tts_isnull[attrnum] != slot2->tts_isnull[attrnum])
     330                 :           0 :             return false;
     331                 :             : 
     332                 :             :         /*
     333                 :             :          * If both are NULL, they can be considered equal.
     334                 :             :          */
     335   [ +  +  -  + ]:      105373 :         if (slot1->tts_isnull[attrnum] || slot2->tts_isnull[attrnum])
     336                 :           1 :             continue;
     337                 :             : 
     338                 :      105372 :         typentry = eq[attrnum];
     339         [ +  + ]:      105372 :         if (typentry == NULL)
     340                 :             :         {
     341                 :         208 :             typentry = lookup_type_cache(att->atttypid,
     342                 :             :                                          TYPECACHE_EQ_OPR_FINFO);
     343         [ -  + ]:         208 :             if (!OidIsValid(typentry->eq_opr_finfo.fn_oid))
     344         [ #  # ]:           0 :                 ereport(ERROR,
     345                 :             :                         (errcode(ERRCODE_UNDEFINED_FUNCTION),
     346                 :             :                          errmsg("could not identify an equality operator for type %s",
     347                 :             :                                 format_type_be(att->atttypid))));
     348                 :         208 :             eq[attrnum] = typentry;
     349                 :             :         }
     350                 :             : 
     351         [ +  + ]:      105372 :         if (!DatumGetBool(FunctionCall2Coll(&typentry->eq_opr_finfo,
     352                 :             :                                             att->attcollation,
     353                 :      105372 :                                             slot1->tts_values[attrnum],
     354                 :      105372 :                                             slot2->tts_values[attrnum])))
     355                 :      105167 :             return false;
     356                 :             :     }
     357                 :             : 
     358                 :         167 :     return true;
     359                 :             : }
     360                 :             : 
     361                 :             : /*
     362                 :             :  * Search the relation 'rel' for tuple using the sequential scan.
     363                 :             :  *
     364                 :             :  * If a matching tuple is found, lock it with lockmode, fill the slot with its
     365                 :             :  * contents, and return true.  Return false otherwise.
     366                 :             :  *
     367                 :             :  * Note that this stops on the first matching tuple.
     368                 :             :  *
     369                 :             :  * This can obviously be quite slow on tables that have more than few rows.
     370                 :             :  */
     371                 :             : bool
     372                 :         152 : RelationFindReplTupleSeq(Relation rel, LockTupleMode lockmode,
     373                 :             :                          TupleTableSlot *searchslot, TupleTableSlot *outslot)
     374                 :             : {
     375                 :             :     TupleTableSlot *scanslot;
     376                 :             :     TableScanDesc scan;
     377                 :             :     SnapshotData snap;
     378                 :             :     TypeCacheEntry **eq;
     379                 :             :     TransactionId xwait;
     380                 :             :     bool        found;
     381                 :         152 :     TupleDesc   desc PG_USED_FOR_ASSERTS_ONLY = RelationGetDescr(rel);
     382                 :             : 
     383                 :             :     Assert(equalTupleDescs(desc, outslot->tts_tupleDescriptor));
     384                 :             : 
     385                 :         152 :     eq = palloc0_array(TypeCacheEntry *, outslot->tts_tupleDescriptor->natts);
     386                 :             : 
     387                 :             :     /* Start a heap scan. */
     388                 :         152 :     InitDirtySnapshot(snap);
     389                 :         152 :     scan = table_beginscan(rel, &snap, 0, NULL,
     390                 :             :                            SO_NONE);
     391                 :         152 :     scanslot = table_slot_create(rel, NULL);
     392                 :             : 
     393                 :           0 : retry:
     394                 :         152 :     found = false;
     395                 :             : 
     396                 :         152 :     table_rescan(scan, NULL);
     397                 :             : 
     398                 :             :     /* Try to find the tuple */
     399         [ +  + ]:      105317 :     while (table_scan_getnextslot(scan, ForwardScanDirection, scanslot))
     400                 :             :     {
     401         [ +  + ]:      105313 :         if (!tuples_equal(scanslot, searchslot, eq, NULL))
     402                 :      105165 :             continue;
     403                 :             : 
     404                 :         148 :         found = true;
     405                 :         148 :         ExecCopySlot(outslot, scanslot);
     406                 :             : 
     407                 :         296 :         xwait = TransactionIdIsValid(snap.xmin) ?
     408         [ -  + ]:         148 :             snap.xmin : snap.xmax;
     409                 :             : 
     410                 :             :         /*
     411                 :             :          * If the tuple is locked, wait for locking transaction to finish and
     412                 :             :          * retry.
     413                 :             :          */
     414         [ -  + ]:         148 :         if (TransactionIdIsValid(xwait))
     415                 :             :         {
     416                 :           0 :             XactLockTableWait(xwait, NULL, NULL, XLTW_None);
     417                 :           0 :             goto retry;
     418                 :             :         }
     419                 :             : 
     420                 :             :         /* Found our tuple and it's not locked */
     421                 :         148 :         break;
     422                 :             :     }
     423                 :             : 
     424                 :             :     /* Found tuple, try to lock it in the lockmode. */
     425         [ +  + ]:         152 :     if (found)
     426                 :             :     {
     427                 :             :         TM_FailureData tmfd;
     428                 :             :         TM_Result   res;
     429                 :             : 
     430                 :         148 :         PushActiveSnapshot(GetLatestSnapshot());
     431                 :             : 
     432                 :         148 :         res = table_tuple_lock(rel, &(outslot->tts_tid), GetActiveSnapshot(),
     433                 :             :                                outslot,
     434                 :             :                                GetCurrentCommandId(false),
     435                 :             :                                lockmode,
     436                 :             :                                LockWaitBlock,
     437                 :             :                                0 /* don't follow updates */ ,
     438                 :             :                                &tmfd);
     439                 :             : 
     440                 :         148 :         PopActiveSnapshot();
     441                 :             : 
     442         [ -  + ]:         148 :         if (should_refetch_tuple(res, &tmfd))
     443                 :           0 :             goto retry;
     444                 :             :     }
     445                 :             : 
     446                 :         152 :     table_endscan(scan);
     447                 :         152 :     ExecDropSingleTupleTableSlot(scanslot);
     448                 :             : 
     449                 :         152 :     return found;
     450                 :             : }
     451                 :             : 
     452                 :             : /*
     453                 :             :  * Build additional index information necessary for conflict detection.
     454                 :             :  */
     455                 :             : static void
     456                 :         115 : BuildConflictIndexInfo(ResultRelInfo *resultRelInfo, Oid conflictindex)
     457                 :             : {
     458         [ +  + ]:         344 :     for (int i = 0; i < resultRelInfo->ri_NumIndices; i++)
     459                 :             :     {
     460                 :         229 :         Relation    indexRelation = resultRelInfo->ri_IndexRelationDescs[i];
     461                 :         229 :         IndexInfo  *indexRelationInfo = resultRelInfo->ri_IndexRelationInfo[i];
     462                 :             : 
     463         [ +  + ]:         229 :         if (conflictindex != RelationGetRelid(indexRelation))
     464                 :         114 :             continue;
     465                 :             : 
     466                 :             :         /*
     467                 :             :          * This Assert will fail if BuildSpeculativeIndexInfo() is called
     468                 :             :          * twice for the given index.
     469                 :             :          */
     470                 :             :         Assert(indexRelationInfo->ii_UniqueOps == NULL);
     471                 :             : 
     472                 :         115 :         BuildSpeculativeIndexInfo(indexRelation, indexRelationInfo);
     473                 :             :     }
     474                 :         115 : }
     475                 :             : 
     476                 :             : /*
     477                 :             :  * If the tuple is recently dead and was deleted by a transaction with a newer
     478                 :             :  * commit timestamp than previously recorded, update the associated transaction
     479                 :             :  * ID, commit time, and origin. This helps ensure that conflict detection uses
     480                 :             :  * the most recent and relevant deletion metadata.
     481                 :             :  */
     482                 :             : static void
     483                 :           3 : update_most_recent_deletion_info(TupleTableSlot *scanslot,
     484                 :             :                                  TransactionId oldestxmin,
     485                 :             :                                  TransactionId *delete_xid,
     486                 :             :                                  TimestampTz *delete_time,
     487                 :             :                                  ReplOriginId *delete_origin)
     488                 :             : {
     489                 :             :     BufferHeapTupleTableSlot *hslot;
     490                 :             :     HeapTuple   tuple;
     491                 :             :     Buffer      buf;
     492                 :           3 :     bool        recently_dead = false;
     493                 :             :     TransactionId xmax;
     494                 :             :     TimestampTz localts;
     495                 :             :     ReplOriginId localorigin;
     496                 :             : 
     497                 :           3 :     hslot = (BufferHeapTupleTableSlot *) scanslot;
     498                 :             : 
     499                 :           3 :     tuple = ExecFetchSlotHeapTuple(scanslot, false, NULL);
     500                 :           3 :     buf = hslot->buffer;
     501                 :             : 
     502                 :           3 :     LockBuffer(buf, BUFFER_LOCK_SHARE);
     503                 :             : 
     504                 :             :     /*
     505                 :             :      * We do not consider HEAPTUPLE_DEAD status because it indicates either
     506                 :             :      * tuples whose inserting transaction was aborted (meaning there is no
     507                 :             :      * commit timestamp or origin), or tuples deleted by a transaction older
     508                 :             :      * than oldestxmin, making it safe to ignore them during conflict
     509                 :             :      * detection (See comments atop worker.c for details).
     510                 :             :      */
     511         [ +  - ]:           3 :     if (HeapTupleSatisfiesVacuum(tuple, oldestxmin, buf) == HEAPTUPLE_RECENTLY_DEAD)
     512                 :           3 :         recently_dead = true;
     513                 :             : 
     514                 :           3 :     LockBuffer(buf, BUFFER_LOCK_UNLOCK);
     515                 :             : 
     516         [ -  + ]:           3 :     if (!recently_dead)
     517                 :           0 :         return;
     518                 :             : 
     519                 :           3 :     xmax = HeapTupleHeaderGetUpdateXid(tuple->t_data);
     520         [ -  + ]:           3 :     if (!TransactionIdIsValid(xmax))
     521                 :           0 :         return;
     522                 :             : 
     523                 :             :     /* Select the dead tuple with the most recent commit timestamp */
     524   [ +  -  +  - ]:           6 :     if (TransactionIdGetCommitTsData(xmax, &localts, &localorigin) &&
     525                 :           3 :         TimestampDifferenceExceeds(*delete_time, localts, 0))
     526                 :             :     {
     527                 :           3 :         *delete_xid = xmax;
     528                 :           3 :         *delete_time = localts;
     529                 :           3 :         *delete_origin = localorigin;
     530                 :             :     }
     531                 :             : }
     532                 :             : 
     533                 :             : /*
     534                 :             :  * Searches the relation 'rel' for the most recently deleted tuple that matches
     535                 :             :  * the values in 'searchslot' and is not yet removable by VACUUM. The function
     536                 :             :  * returns the transaction ID, origin, and commit timestamp of the transaction
     537                 :             :  * that deleted this tuple.
     538                 :             :  *
     539                 :             :  * If 'identidxoid' is valid, it is the replica identity or primary key
     540                 :             :  * index, and only its key columns are compared. Otherwise, all columns are
     541                 :             :  * compared.
     542                 :             :  *
     543                 :             :  * 'oldestxmin' acts as a cutoff transaction ID. Tuples deleted by transactions
     544                 :             :  * with IDs >= 'oldestxmin' are considered recently dead and are eligible for
     545                 :             :  * conflict detection.
     546                 :             :  *
     547                 :             :  * Instead of stopping at the first match, we scan all matching dead tuples to
     548                 :             :  * identify most recent deletion. This is crucial because only the latest
     549                 :             :  * deletion is relevant for resolving conflicts.
     550                 :             :  *
     551                 :             :  * For example, consider a scenario on the subscriber where a row is deleted,
     552                 :             :  * re-inserted, and then deleted again only on the subscriber:
     553                 :             :  *
     554                 :             :  *   - (pk, 1) - deleted at 9:00,
     555                 :             :  *   - (pk, 1) - deleted at 9:02,
     556                 :             :  *
     557                 :             :  * Now, a remote update arrives: (pk, 1) -> (pk, 2), timestamped at 9:01.
     558                 :             :  *
     559                 :             :  * If we mistakenly return the older deletion (9:00), the system may wrongly
     560                 :             :  * apply the remote update using a last-update-wins strategy. Instead, we must
     561                 :             :  * recognize the more recent deletion at 9:02 and skip the update. See
     562                 :             :  * comments atop worker.c for details. Note, as of now, conflict resolution
     563                 :             :  * is not implemented. Consequently, the system may incorrectly report the
     564                 :             :  * older tuple as the conflicted one, leading to misleading results.
     565                 :             :  *
     566                 :             :  * The commit timestamp of the deleting transaction is used to determine which
     567                 :             :  * tuple was deleted most recently.
     568                 :             :  */
     569                 :             : bool
     570                 :           3 : RelationFindDeletedTupleInfoSeq(Relation rel, Oid identidxoid,
     571                 :             :                                 TupleTableSlot *searchslot,
     572                 :             :                                 TransactionId oldestxmin,
     573                 :             :                                 TransactionId *delete_xid,
     574                 :             :                                 ReplOriginId *delete_origin,
     575                 :             :                                 TimestampTz *delete_time)
     576                 :             : {
     577                 :             :     TupleTableSlot *scanslot;
     578                 :             :     TableScanDesc scan;
     579                 :             :     TypeCacheEntry **eq;
     580                 :           3 :     Bitmapset  *indexbitmap = NULL;
     581                 :           3 :     TupleDesc   desc PG_USED_FOR_ASSERTS_ONLY = RelationGetDescr(rel);
     582                 :             : 
     583                 :             :     Assert(equalTupleDescs(desc, searchslot->tts_tupleDescriptor));
     584                 :             : 
     585                 :           3 :     *delete_xid = InvalidTransactionId;
     586                 :           3 :     *delete_origin = InvalidReplOriginId;
     587                 :           3 :     *delete_time = 0;
     588                 :             : 
     589                 :             :     /*
     590                 :             :      * If the caller's replica identity key or primary key is unusable for
     591                 :             :      * locating deleted tuples (see IsIndexUsableForFindingDeletedTuple), a
     592                 :             :      * full table scan becomes necessary. In such cases, comparing the entire
     593                 :             :      * tuple is not required, since the remote tuple might not include all
     594                 :             :      * column values. Instead, the indexed columns alone are sufficient to
     595                 :             :      * identify the target tuple (see logicalrep_rel_mark_updatable).
     596                 :             :      */
     597         [ -  + ]:           3 :     if (OidIsValid(identidxoid))
     598                 :             :     {
     599                 :             :         /* The index must have been locked already */
     600                 :           0 :         Relation    idxrel = index_open(identidxoid, NoLock);
     601                 :             : 
     602                 :             :         /*
     603                 :             :          * The index may no longer be the replica identity if DROP INDEX
     604                 :             :          * CONCURRENTLY or REINDEX CONCURRENTLY ran meanwhile, but it stays
     605                 :             :          * unique and non-partial, which is all we rely on. See
     606                 :             :          * FindReplTupleInLocalRel().
     607                 :             :          */
     608                 :             :         Assert(idxrel->rd_index->indisunique);
     609                 :             :         Assert(heap_attisnull(idxrel->rd_indextuple, Anum_pg_index_indpred,
     610                 :             :                               NULL));
     611                 :             : 
     612         [ #  # ]:           0 :         for (int i = 0; i < idxrel->rd_index->indnkeyatts; i++)
     613                 :           0 :             indexbitmap = bms_add_member(indexbitmap,
     614                 :           0 :                                          idxrel->rd_index->indkey.values[i] -
     615                 :             :                                          FirstLowInvalidHeapAttributeNumber);
     616                 :             : 
     617                 :           0 :         index_close(idxrel, NoLock);
     618                 :             :     }
     619                 :             : 
     620                 :           3 :     eq = palloc0_array(TypeCacheEntry *, searchslot->tts_tupleDescriptor->natts);
     621                 :             : 
     622                 :             :     /*
     623                 :             :      * Start a heap scan using SnapshotAny to identify dead tuples that are
     624                 :             :      * not visible under a standard MVCC snapshot. Tuples from transactions
     625                 :             :      * not yet committed or those just committed prior to the scan are
     626                 :             :      * excluded in update_most_recent_deletion_info().
     627                 :             :      */
     628                 :           3 :     scan = table_beginscan(rel, SnapshotAny, 0, NULL,
     629                 :             :                            SO_NONE);
     630                 :           3 :     scanslot = table_slot_create(rel, NULL);
     631                 :             : 
     632                 :           3 :     table_rescan(scan, NULL);
     633                 :             : 
     634                 :             :     /* Try to find the tuple */
     635         [ +  + ]:           7 :     while (table_scan_getnextslot(scan, ForwardScanDirection, scanslot))
     636                 :             :     {
     637         [ +  + ]:           4 :         if (!tuples_equal(scanslot, searchslot, eq, indexbitmap))
     638                 :           2 :             continue;
     639                 :             : 
     640                 :           2 :         update_most_recent_deletion_info(scanslot, oldestxmin, delete_xid,
     641                 :             :                                          delete_time, delete_origin);
     642                 :             :     }
     643                 :             : 
     644                 :           3 :     table_endscan(scan);
     645                 :           3 :     ExecDropSingleTupleTableSlot(scanslot);
     646                 :             : 
     647                 :           3 :     return *delete_time != 0;
     648                 :             : }
     649                 :             : 
     650                 :             : /*
     651                 :             :  * Similar to RelationFindDeletedTupleInfoSeq() but using index scan to locate
     652                 :             :  * the deleted tuple.
     653                 :             :  *
     654                 :             :  * 'skipduplicates' works as in RelationFindReplTupleByIndex().
     655                 :             :  */
     656                 :             : bool
     657                 :           1 : RelationFindDeletedTupleInfoByIndex(Relation rel, Oid idxoid,
     658                 :             :                                     bool skipduplicates,
     659                 :             :                                     TupleTableSlot *searchslot,
     660                 :             :                                     TransactionId oldestxmin,
     661                 :             :                                     TransactionId *delete_xid,
     662                 :             :                                     ReplOriginId *delete_origin,
     663                 :             :                                     TimestampTz *delete_time)
     664                 :             : {
     665                 :             :     Relation    idxrel;
     666                 :             :     ScanKeyData skey[INDEX_MAX_KEYS];
     667                 :             :     int         skey_attoff;
     668                 :             :     IndexScanDesc scan;
     669                 :             :     TupleTableSlot *scanslot;
     670                 :           1 :     TypeCacheEntry **eq = NULL;
     671                 :           1 :     TupleDesc   desc PG_USED_FOR_ASSERTS_ONLY = RelationGetDescr(rel);
     672                 :             : 
     673                 :             :     Assert(equalTupleDescs(desc, searchslot->tts_tupleDescriptor));
     674                 :             :     Assert(OidIsValid(idxoid));
     675                 :             : 
     676                 :           1 :     *delete_xid = InvalidTransactionId;
     677                 :           1 :     *delete_time = 0;
     678                 :           1 :     *delete_origin = InvalidReplOriginId;
     679                 :             : 
     680                 :           1 :     scanslot = table_slot_create(rel, NULL);
     681                 :             : 
     682                 :           1 :     idxrel = index_open(idxoid, RowExclusiveLock);
     683                 :             : 
     684                 :             :     /* Build scan key. */
     685                 :           1 :     skey_attoff = build_replindex_scan_key(skey, rel, idxrel, searchslot);
     686                 :             : 
     687                 :             :     /*
     688                 :             :      * Start an index scan using SnapshotAny to identify dead tuples that are
     689                 :             :      * not visible under a standard MVCC snapshot. Tuples from transactions
     690                 :             :      * not yet committed or those just committed prior to the scan are
     691                 :             :      * excluded in update_most_recent_deletion_info().
     692                 :             :      */
     693                 :           1 :     scan = index_beginscan(rel, idxrel, false,
     694                 :             :                            SnapshotAny, NULL, skey_attoff, 0, SO_NONE);
     695                 :             : 
     696                 :           1 :     index_rescan(scan, skey, skey_attoff, NULL, 0);
     697                 :             : 
     698                 :             :     /* Try to find the tuple */
     699         [ +  + ]:           2 :     while (table_index_getnext_slot(scan, ForwardScanDirection, scanslot))
     700                 :             :     {
     701                 :             :         /*
     702                 :             :          * Avoid expensive equality check if the index is primary key or
     703                 :             :          * replica identity index.
     704                 :             :          */
     705         [ -  + ]:           1 :         if (!skipduplicates)
     706                 :             :         {
     707         [ #  # ]:           0 :             if (eq == NULL)
     708                 :           0 :                 eq = palloc0_array(TypeCacheEntry *, scanslot->tts_tupleDescriptor->natts);
     709                 :             : 
     710         [ #  # ]:           0 :             if (!tuples_equal(scanslot, searchslot, eq, NULL))
     711                 :           0 :                 continue;
     712                 :             :         }
     713                 :             : 
     714                 :           1 :         update_most_recent_deletion_info(scanslot, oldestxmin, delete_xid,
     715                 :             :                                          delete_time, delete_origin);
     716                 :             :     }
     717                 :             : 
     718                 :           1 :     index_endscan(scan);
     719                 :             : 
     720                 :           1 :     index_close(idxrel, NoLock);
     721                 :             : 
     722                 :           1 :     ExecDropSingleTupleTableSlot(scanslot);
     723                 :             : 
     724                 :           1 :     return *delete_time != 0;
     725                 :             : }
     726                 :             : 
     727                 :             : /*
     728                 :             :  * Find the tuple that violates the passed unique index (conflictindex).
     729                 :             :  *
     730                 :             :  * If the conflicting tuple is found return true, otherwise false.
     731                 :             :  *
     732                 :             :  * We lock the tuple to avoid getting it deleted before the caller can fetch
     733                 :             :  * the required information. Note that if the tuple is deleted before a lock
     734                 :             :  * is acquired, we will retry to find the conflicting tuple again.
     735                 :             :  */
     736                 :             : static bool
     737                 :         115 : FindConflictTuple(ResultRelInfo *resultRelInfo, EState *estate,
     738                 :             :                   Oid conflictindex, TupleTableSlot *slot,
     739                 :             :                   TupleTableSlot **conflictslot)
     740                 :             : {
     741                 :         115 :     Relation    rel = resultRelInfo->ri_RelationDesc;
     742                 :             :     ItemPointerData conflictTid;
     743                 :             :     TM_FailureData tmfd;
     744                 :             :     TM_Result   res;
     745                 :             : 
     746                 :         115 :     *conflictslot = NULL;
     747                 :             : 
     748                 :             :     /*
     749                 :             :      * Build additional information required to check constraints violations.
     750                 :             :      * See check_exclusion_or_unique_constraint().
     751                 :             :      */
     752                 :         115 :     BuildConflictIndexInfo(resultRelInfo, conflictindex);
     753                 :             : 
     754                 :         115 : retry:
     755         [ +  + ]:         229 :     if (ExecCheckIndexConstraints(resultRelInfo, slot, estate,
     756                 :         115 :                                   &conflictTid, &slot->tts_tid,
     757                 :             :                                   list_make1_oid(conflictindex)))
     758                 :             :     {
     759         [ -  + ]:           2 :         if (*conflictslot)
     760                 :           0 :             ExecDropSingleTupleTableSlot(*conflictslot);
     761                 :             : 
     762                 :           2 :         *conflictslot = NULL;
     763                 :           2 :         return false;
     764                 :             :     }
     765                 :             : 
     766                 :         112 :     *conflictslot = table_slot_create(rel, NULL);
     767                 :             : 
     768                 :         112 :     PushActiveSnapshot(GetLatestSnapshot());
     769                 :             : 
     770                 :         112 :     res = table_tuple_lock(rel, &conflictTid, GetActiveSnapshot(),
     771                 :             :                            *conflictslot,
     772                 :             :                            GetCurrentCommandId(false),
     773                 :             :                            LockTupleShare,
     774                 :             :                            LockWaitBlock,
     775                 :             :                            0 /* don't follow updates */ ,
     776                 :             :                            &tmfd);
     777                 :             : 
     778                 :         112 :     PopActiveSnapshot();
     779                 :             : 
     780         [ -  + ]:         112 :     if (should_refetch_tuple(res, &tmfd))
     781                 :           0 :         goto retry;
     782                 :             : 
     783                 :         112 :     return true;
     784                 :             : }
     785                 :             : 
     786                 :             : /*
     787                 :             :  * Check all the unique indexes in 'recheckIndexes' for conflict with the
     788                 :             :  * tuple in 'remoteslot' and report if found.
     789                 :             :  */
     790                 :             : static void
     791                 :          61 : CheckAndReportConflict(ResultRelInfo *resultRelInfo, EState *estate,
     792                 :             :                        ConflictType type, List *recheckIndexes,
     793                 :             :                        TupleTableSlot *searchslot, TupleTableSlot *remoteslot)
     794                 :             : {
     795                 :          61 :     List       *conflicttuples = NIL;
     796                 :             :     TupleTableSlot *conflictslot;
     797                 :             : 
     798                 :             :     /* Check all the unique indexes for conflicts */
     799   [ +  -  +  +  :         235 :     foreach_oid(uniqueidx, resultRelInfo->ri_onConflictArbiterIndexes)
                   +  + ]
     800                 :             :     {
     801   [ +  -  +  + ]:         229 :         if (list_member_oid(recheckIndexes, uniqueidx) &&
     802                 :         115 :             FindConflictTuple(resultRelInfo, estate, uniqueidx, remoteslot,
     803                 :             :                               &conflictslot))
     804                 :             :         {
     805                 :         112 :             ConflictTupleInfo *conflicttuple = palloc0_object(ConflictTupleInfo);
     806                 :             : 
     807                 :         112 :             conflicttuple->slot = conflictslot;
     808                 :         112 :             conflicttuple->indexoid = uniqueidx;
     809                 :             : 
     810                 :         112 :             GetTupleTransactionInfo(conflictslot, &conflicttuple->xmin,
     811                 :             :                                     &conflicttuple->origin, &conflicttuple->ts);
     812                 :             : 
     813                 :         112 :             conflicttuples = lappend(conflicttuples, conflicttuple);
     814                 :             :         }
     815                 :             :     }
     816                 :             : 
     817                 :             :     /* Report the conflict, if found */
     818         [ +  + ]:          60 :     if (conflicttuples)
     819         [ +  + ]:          58 :         ReportApplyConflict(estate, resultRelInfo, ERROR,
     820                 :          58 :                             list_length(conflicttuples) > 1 ? CT_MULTIPLE_UNIQUE_CONFLICTS : type,
     821                 :             :                             searchslot, remoteslot, conflicttuples);
     822                 :           2 : }
     823                 :             : 
     824                 :             : /*
     825                 :             :  * Insert tuple represented in the slot to the relation, update the indexes,
     826                 :             :  * and execute any constraints and per-row triggers.
     827                 :             :  *
     828                 :             :  * Caller is responsible for opening the indexes.
     829                 :             :  */
     830                 :             : void
     831                 :       99563 : ExecSimpleRelationInsert(ResultRelInfo *resultRelInfo,
     832                 :             :                          EState *estate, TupleTableSlot *slot)
     833                 :             : {
     834                 :       99563 :     bool        skip_tuple = false;
     835                 :       99563 :     Relation    rel = resultRelInfo->ri_RelationDesc;
     836                 :             : 
     837                 :             :     /* For now we support only tables. */
     838                 :             :     Assert(rel->rd_rel->relkind == RELKIND_RELATION);
     839                 :             : 
     840                 :       99563 :     CheckCmdReplicaIdentity(rel, CMD_INSERT);
     841                 :             : 
     842                 :             :     /* BEFORE ROW INSERT Triggers */
     843         [ +  + ]:       99563 :     if (resultRelInfo->ri_TrigDesc &&
     844         [ +  + ]:          20 :         resultRelInfo->ri_TrigDesc->trig_insert_before_row)
     845                 :             :     {
     846         [ +  + ]:           3 :         if (!ExecBRInsertTriggers(estate, resultRelInfo, slot))
     847                 :           1 :             skip_tuple = true;  /* "do nothing" */
     848                 :             :     }
     849                 :             : 
     850         [ +  + ]:       99563 :     if (!skip_tuple)
     851                 :             :     {
     852                 :       99562 :         List       *recheckIndexes = NIL;
     853                 :             :         List       *conflictindexes;
     854                 :       99562 :         bool        conflict = false;
     855                 :             : 
     856                 :             :         /* Compute stored generated columns */
     857         [ +  + ]:       99562 :         if (rel->rd_att->constr &&
     858         [ +  + ]:       66089 :             rel->rd_att->constr->has_generated_stored)
     859                 :           4 :             ExecComputeStoredGenerated(resultRelInfo, estate, slot,
     860                 :             :                                        CMD_INSERT);
     861                 :             : 
     862                 :             :         /* Check the constraints of the tuple */
     863         [ +  + ]:       99562 :         if (rel->rd_att->constr)
     864                 :       66089 :             ExecConstraints(resultRelInfo, slot, estate);
     865         [ +  + ]:       99562 :         if (rel->rd_rel->relispartition)
     866                 :         108 :             ExecPartitionCheck(resultRelInfo, slot, estate, true);
     867                 :             : 
     868                 :             :         /* OK, store the tuple and create index entries for it */
     869                 :       99562 :         simple_table_tuple_insert(resultRelInfo->ri_RelationDesc, slot);
     870                 :             : 
     871                 :       99562 :         conflictindexes = resultRelInfo->ri_onConflictArbiterIndexes;
     872                 :             : 
     873         [ +  + ]:       99562 :         if (resultRelInfo->ri_NumIndices > 0)
     874                 :             :         {
     875                 :             :             uint32      flags;
     876                 :             : 
     877         [ +  + ]:       76136 :             if (conflictindexes != NIL)
     878                 :       76132 :                 flags = EIIT_NO_DUPE_ERROR;
     879                 :             :             else
     880                 :           4 :                 flags = 0;
     881                 :       76136 :             recheckIndexes = ExecInsertIndexTuples(resultRelInfo,
     882                 :             :                                                    estate, flags,
     883                 :             :                                                    slot, conflictindexes,
     884                 :             :                                                    &conflict);
     885                 :             :         }
     886                 :             : 
     887                 :             :         /*
     888                 :             :          * Checks the conflict indexes to fetch the conflicting local row and
     889                 :             :          * reports the conflict. We perform this check here, instead of
     890                 :             :          * performing an additional index scan before the actual insertion and
     891                 :             :          * reporting the conflict if any conflicting rows are found. This is
     892                 :             :          * to avoid the overhead of executing the extra scan for each INSERT
     893                 :             :          * operation, even when no conflict arises, which could introduce
     894                 :             :          * significant overhead to replication, particularly in cases where
     895                 :             :          * conflicts are rare.
     896                 :             :          *
     897                 :             :          * XXX OTOH, this could lead to clean-up effort for dead tuples added
     898                 :             :          * in heap and index in case of conflicts. But as conflicts shouldn't
     899                 :             :          * be a frequent thing so we preferred to save the performance
     900                 :             :          * overhead of extra scan before each insertion.
     901                 :             :          */
     902         [ +  + ]:       99562 :         if (conflict)
     903                 :          59 :             CheckAndReportConflict(resultRelInfo, estate, CT_INSERT_EXISTS,
     904                 :             :                                    recheckIndexes, NULL, slot);
     905                 :             : 
     906                 :             :         /* AFTER ROW INSERT Triggers */
     907                 :       99505 :         ExecARInsertTriggers(estate, resultRelInfo, slot,
     908                 :             :                              recheckIndexes, NULL);
     909                 :             : 
     910                 :             :         /*
     911                 :             :          * XXX we should in theory pass a TransitionCaptureState object to the
     912                 :             :          * above to capture transition tuples, but after statement triggers
     913                 :             :          * don't actually get fired by replication yet anyway
     914                 :             :          */
     915                 :             : 
     916                 :       99505 :         list_free(recheckIndexes);
     917                 :             :     }
     918                 :       99506 : }
     919                 :             : 
     920                 :             : /*
     921                 :             :  * Find the searchslot tuple and update it with data in the slot,
     922                 :             :  * update the indexes, and execute any constraints and per-row triggers.
     923                 :             :  *
     924                 :             :  * Caller is responsible for opening the indexes.
     925                 :             :  */
     926                 :             : void
     927                 :       31926 : ExecSimpleRelationUpdate(ResultRelInfo *resultRelInfo,
     928                 :             :                          EState *estate, EPQState *epqstate,
     929                 :             :                          TupleTableSlot *searchslot, TupleTableSlot *slot)
     930                 :             : {
     931                 :       31926 :     bool        skip_tuple = false;
     932                 :       31926 :     Relation    rel = resultRelInfo->ri_RelationDesc;
     933                 :       31926 :     ItemPointer tid = &(searchslot->tts_tid);
     934                 :             : 
     935                 :             :     /*
     936                 :             :      * We support only non-system tables, with
     937                 :             :      * check_publication_add_relation() accountable.
     938                 :             :      */
     939                 :             :     Assert(rel->rd_rel->relkind == RELKIND_RELATION);
     940                 :             :     Assert(!IsCatalogRelation(rel));
     941                 :             : 
     942                 :       31926 :     CheckCmdReplicaIdentity(rel, CMD_UPDATE);
     943                 :             : 
     944                 :             :     /* BEFORE ROW UPDATE Triggers */
     945         [ +  + ]:       31926 :     if (resultRelInfo->ri_TrigDesc &&
     946         [ +  + ]:          10 :         resultRelInfo->ri_TrigDesc->trig_update_before_row)
     947                 :             :     {
     948         [ +  + ]:           3 :         if (!ExecBRUpdateTriggers(estate, epqstate, resultRelInfo,
     949                 :             :                                   tid, NULL, slot, NULL, NULL, false))
     950                 :           2 :             skip_tuple = true;  /* "do nothing" */
     951                 :             :     }
     952                 :             : 
     953         [ +  + ]:       31926 :     if (!skip_tuple)
     954                 :             :     {
     955                 :       31924 :         List       *recheckIndexes = NIL;
     956                 :             :         TU_UpdateIndexes update_indexes;
     957                 :             :         List       *conflictindexes;
     958                 :       31924 :         bool        conflict = false;
     959                 :             : 
     960                 :             :         /* Compute stored generated columns */
     961         [ +  + ]:       31924 :         if (rel->rd_att->constr &&
     962         [ +  + ]:       31878 :             rel->rd_att->constr->has_generated_stored)
     963                 :           2 :             ExecComputeStoredGenerated(resultRelInfo, estate, slot,
     964                 :             :                                        CMD_UPDATE);
     965                 :             : 
     966                 :             :         /* Check the constraints of the tuple */
     967         [ +  + ]:       31924 :         if (rel->rd_att->constr)
     968                 :       31878 :             ExecConstraints(resultRelInfo, slot, estate);
     969         [ +  + ]:       31924 :         if (rel->rd_rel->relispartition)
     970                 :          12 :             ExecPartitionCheck(resultRelInfo, slot, estate, true);
     971                 :             : 
     972                 :       31924 :         simple_table_tuple_update(rel, tid, slot, estate->es_snapshot,
     973                 :             :                                   &update_indexes);
     974                 :             : 
     975                 :       31924 :         conflictindexes = resultRelInfo->ri_onConflictArbiterIndexes;
     976                 :             : 
     977   [ +  +  +  + ]:       31924 :         if (resultRelInfo->ri_NumIndices > 0 && (update_indexes != TU_None))
     978                 :             :         {
     979                 :       20204 :             uint32      flags = EIIT_IS_UPDATE;
     980                 :             : 
     981         [ +  + ]:       20204 :             if (conflictindexes != NIL)
     982                 :       20194 :                 flags |= EIIT_NO_DUPE_ERROR;
     983         [ -  + ]:       20204 :             if (update_indexes == TU_Summarizing)
     984                 :           0 :                 flags |= EIIT_ONLY_SUMMARIZING;
     985                 :       20204 :             recheckIndexes = ExecInsertIndexTuples(resultRelInfo,
     986                 :             :                                                    estate, flags,
     987                 :             :                                                    slot, conflictindexes,
     988                 :             :                                                    &conflict);
     989                 :             :         }
     990                 :             : 
     991                 :             :         /*
     992                 :             :          * Refer to the comments above the call to CheckAndReportConflict() in
     993                 :             :          * ExecSimpleRelationInsert to understand why this check is done at
     994                 :             :          * this point.
     995                 :             :          */
     996         [ +  + ]:       31924 :         if (conflict)
     997                 :           2 :             CheckAndReportConflict(resultRelInfo, estate, CT_UPDATE_EXISTS,
     998                 :             :                                    recheckIndexes, searchslot, slot);
     999                 :             : 
    1000                 :             :         /* AFTER ROW UPDATE Triggers */
    1001                 :       31922 :         ExecARUpdateTriggers(estate, resultRelInfo,
    1002                 :             :                              NULL, NULL,
    1003                 :             :                              tid, NULL, slot,
    1004                 :             :                              recheckIndexes, NULL, false);
    1005                 :             : 
    1006                 :       31922 :         list_free(recheckIndexes);
    1007                 :             :     }
    1008                 :       31924 : }
    1009                 :             : 
    1010                 :             : /*
    1011                 :             :  * Find the searchslot tuple and delete it, and execute any constraints
    1012                 :             :  * and per-row triggers.
    1013                 :             :  *
    1014                 :             :  * Caller is responsible for opening the indexes.
    1015                 :             :  */
    1016                 :             : void
    1017                 :       40313 : ExecSimpleRelationDelete(ResultRelInfo *resultRelInfo,
    1018                 :             :                          EState *estate, EPQState *epqstate,
    1019                 :             :                          TupleTableSlot *searchslot)
    1020                 :             : {
    1021                 :       40313 :     bool        skip_tuple = false;
    1022                 :       40313 :     Relation    rel = resultRelInfo->ri_RelationDesc;
    1023                 :       40313 :     ItemPointer tid = &searchslot->tts_tid;
    1024                 :             : 
    1025                 :       40313 :     CheckCmdReplicaIdentity(rel, CMD_DELETE);
    1026                 :             : 
    1027                 :             :     /* BEFORE ROW DELETE Triggers */
    1028         [ +  + ]:       40313 :     if (resultRelInfo->ri_TrigDesc &&
    1029         [ -  + ]:          10 :         resultRelInfo->ri_TrigDesc->trig_delete_before_row)
    1030                 :             :     {
    1031                 :           0 :         skip_tuple = !ExecBRDeleteTriggers(estate, epqstate, resultRelInfo,
    1032                 :           0 :                                            tid, NULL, NULL, NULL, NULL, false);
    1033                 :             :     }
    1034                 :             : 
    1035         [ +  - ]:       40313 :     if (!skip_tuple)
    1036                 :             :     {
    1037                 :             :         /* OK, delete the tuple */
    1038                 :       40313 :         simple_table_tuple_delete(rel, tid, estate->es_snapshot);
    1039                 :             : 
    1040                 :             :         /* AFTER ROW DELETE Triggers */
    1041                 :       40313 :         ExecARDeleteTriggers(estate, resultRelInfo,
    1042                 :             :                              tid, NULL, NULL, false);
    1043                 :             :     }
    1044                 :       40313 : }
    1045                 :             : 
    1046                 :             : /*
    1047                 :             :  * Check if command can be executed with current replica identity.
    1048                 :             :  */
    1049                 :             : void
    1050                 :      268579 : CheckCmdReplicaIdentity(Relation rel, CmdType cmd)
    1051                 :             : {
    1052                 :             :     PublicationDesc pubdesc;
    1053                 :             : 
    1054                 :             :     /*
    1055                 :             :      * Skip checking the replica identity for partitioned tables, because the
    1056                 :             :      * operations are actually performed on the leaf partitions.
    1057                 :             :      */
    1058         [ +  + ]:      268579 :     if (rel->rd_rel->relkind == RELKIND_PARTITIONED_TABLE)
    1059                 :      252956 :         return;
    1060                 :             : 
    1061                 :             :     /* We only need to do checks for UPDATE and DELETE. */
    1062   [ +  +  +  + ]:      264701 :     if (cmd != CMD_UPDATE && cmd != CMD_DELETE)
    1063                 :      172170 :         return;
    1064                 :             : 
    1065                 :             :     /*
    1066                 :             :      * It is only safe to execute UPDATE/DELETE if the relation does not
    1067                 :             :      * publish UPDATEs or DELETEs, or all the following conditions are
    1068                 :             :      * satisfied:
    1069                 :             :      *
    1070                 :             :      * 1. All columns, referenced in the row filters from publications which
    1071                 :             :      * the relation is in, are valid - i.e. when all referenced columns are
    1072                 :             :      * part of REPLICA IDENTITY.
    1073                 :             :      *
    1074                 :             :      * 2. All columns, referenced in the column lists are valid - i.e. when
    1075                 :             :      * all columns referenced in the REPLICA IDENTITY are covered by the
    1076                 :             :      * column list.
    1077                 :             :      *
    1078                 :             :      * 3. All generated columns in REPLICA IDENTITY of the relation, are valid
    1079                 :             :      * - i.e. when all these generated columns are published.
    1080                 :             :      *
    1081                 :             :      * XXX We could optimize it by first checking whether any of the
    1082                 :             :      * publications have a row filter or column list for this relation, or if
    1083                 :             :      * the relation contains a generated column. If none of these exist and
    1084                 :             :      * the relation has replica identity then we can avoid building the
    1085                 :             :      * descriptor but as this happens only one time it doesn't seem worth the
    1086                 :             :      * additional complexity.
    1087                 :             :      */
    1088                 :       92531 :     RelationBuildPublicationDesc(rel, &pubdesc);
    1089   [ +  +  +  + ]:       92531 :     if (cmd == CMD_UPDATE && !pubdesc.rf_valid_for_update)
    1090         [ +  - ]:          40 :         ereport(ERROR,
    1091                 :             :                 (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
    1092                 :             :                  errmsg("cannot update table \"%s\"",
    1093                 :             :                         RelationGetRelationName(rel)),
    1094                 :             :                  errdetail("Column used in the publication WHERE expression is not part of the replica identity.")));
    1095   [ +  +  +  + ]:       92491 :     else if (cmd == CMD_UPDATE && !pubdesc.cols_valid_for_update)
    1096         [ +  - ]:          72 :         ereport(ERROR,
    1097                 :             :                 (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
    1098                 :             :                  errmsg("cannot update table \"%s\"",
    1099                 :             :                         RelationGetRelationName(rel)),
    1100                 :             :                  errdetail("Column list used by the publication does not cover the replica identity.")));
    1101   [ +  +  +  + ]:       92419 :     else if (cmd == CMD_UPDATE && !pubdesc.gencols_valid_for_update)
    1102         [ +  - ]:          16 :         ereport(ERROR,
    1103                 :             :                 (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
    1104                 :             :                  errmsg("cannot update table \"%s\"",
    1105                 :             :                         RelationGetRelationName(rel)),
    1106                 :             :                  errdetail("Replica identity must not contain unpublished generated columns.")));
    1107   [ +  +  -  + ]:       92403 :     else if (cmd == CMD_DELETE && !pubdesc.rf_valid_for_delete)
    1108         [ #  # ]:           0 :         ereport(ERROR,
    1109                 :             :                 (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
    1110                 :             :                  errmsg("cannot delete from table \"%s\"",
    1111                 :             :                         RelationGetRelationName(rel)),
    1112                 :             :                  errdetail("Column used in the publication WHERE expression is not part of the replica identity.")));
    1113   [ +  +  -  + ]:       92403 :     else if (cmd == CMD_DELETE && !pubdesc.cols_valid_for_delete)
    1114         [ #  # ]:           0 :         ereport(ERROR,
    1115                 :             :                 (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
    1116                 :             :                  errmsg("cannot delete from table \"%s\"",
    1117                 :             :                         RelationGetRelationName(rel)),
    1118                 :             :                  errdetail("Column list used by the publication does not cover the replica identity.")));
    1119   [ +  +  -  + ]:       92403 :     else if (cmd == CMD_DELETE && !pubdesc.gencols_valid_for_delete)
    1120         [ #  # ]:           0 :         ereport(ERROR,
    1121                 :             :                 (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
    1122                 :             :                  errmsg("cannot delete from table \"%s\"",
    1123                 :             :                         RelationGetRelationName(rel)),
    1124                 :             :                  errdetail("Replica identity must not contain unpublished generated columns.")));
    1125                 :             : 
    1126                 :             :     /* If relation has replica identity we are always good. */
    1127         [ +  + ]:       92403 :     if (OidIsValid(RelationGetReplicaIndex(rel)))
    1128                 :       76670 :         return;
    1129                 :             : 
    1130                 :             :     /* REPLICA IDENTITY FULL is also good for UPDATE/DELETE. */
    1131         [ +  + ]:       15733 :     if (rel->rd_rel->relreplident == REPLICA_IDENTITY_FULL)
    1132                 :         238 :         return;
    1133                 :             : 
    1134                 :             :     /*
    1135                 :             :      * This is UPDATE/DELETE and there is no replica identity.
    1136                 :             :      *
    1137                 :             :      * Check if the table publishes UPDATES or DELETES.
    1138                 :             :      */
    1139   [ +  +  +  + ]:       15495 :     if (cmd == CMD_UPDATE && pubdesc.pubactions.pubupdate)
    1140         [ +  - ]:          81 :         ereport(ERROR,
    1141                 :             :                 (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
    1142                 :             :                  errmsg("cannot update table \"%s\" because it does not have a replica identity and publishes updates",
    1143                 :             :                         RelationGetRelationName(rel)),
    1144                 :             :                  errhint("To enable updating the table, set REPLICA IDENTITY using ALTER TABLE.")));
    1145   [ +  +  +  + ]:       15414 :     else if (cmd == CMD_DELETE && pubdesc.pubactions.pubdelete)
    1146         [ +  - ]:           9 :         ereport(ERROR,
    1147                 :             :                 (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
    1148                 :             :                  errmsg("cannot delete from table \"%s\" because it does not have a replica identity and publishes deletes",
    1149                 :             :                         RelationGetRelationName(rel)),
    1150                 :             :                  errhint("To enable deleting from the table, set REPLICA IDENTITY using ALTER TABLE.")));
    1151                 :             : }
    1152                 :             : 
    1153                 :             : 
    1154                 :             : /*
    1155                 :             :  * Check if we support writing into specific relkind of local relation and check
    1156                 :             :  * if it aligns with the relkind of the relation on the publisher.
    1157                 :             :  *
    1158                 :             :  * The nspname and relname are only needed for error reporting.
    1159                 :             :  */
    1160                 :             : void
    1161                 :        1123 : CheckSubscriptionRelkind(char localrelkind, char remoterelkind,
    1162                 :             :                          const char *nspname, const char *relname)
    1163                 :             : {
    1164   [ +  +  +  + ]:        1123 :     if (localrelkind != RELKIND_RELATION &&
    1165         [ -  + ]:          18 :         localrelkind != RELKIND_PARTITIONED_TABLE &&
    1166                 :             :         localrelkind != RELKIND_SEQUENCE)
    1167         [ #  # ]:           0 :         ereport(ERROR,
    1168                 :             :                 (errcode(ERRCODE_WRONG_OBJECT_TYPE),
    1169                 :             :                  errmsg("cannot use relation \"%s.%s\" as logical replication target",
    1170                 :             :                         nspname, relname),
    1171                 :             :                  errdetail_relkind_not_supported(localrelkind)));
    1172                 :             : 
    1173                 :             :     /*
    1174                 :             :      * Allow RELKIND_RELATION and RELKIND_PARTITIONED_TABLE to be treated
    1175                 :             :      * interchangeably, but ensure that sequences (RELKIND_SEQUENCE) match
    1176                 :             :      * exactly on both publisher and subscriber.
    1177                 :             :      */
    1178   [ +  +  +  -  :        1123 :     if ((localrelkind == RELKIND_SEQUENCE && remoterelkind != RELKIND_SEQUENCE) ||
                   +  + ]
    1179         [ -  + ]:        1105 :         (localrelkind != RELKIND_SEQUENCE && remoterelkind == RELKIND_SEQUENCE))
    1180   [ #  #  #  #  :           0 :         ereport(ERROR,
                   #  # ]
    1181                 :             :                 errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
    1182                 :             :         /* translator: 3rd and 4th %s are "sequence" or "table" */
    1183                 :             :                 errmsg("relation \"%s.%s\" type mismatch: source \"%s\", target \"%s\"",
    1184                 :             :                        nspname, relname,
    1185                 :             :                        remoterelkind == RELKIND_SEQUENCE ? "sequence" : "table",
    1186                 :             :                        localrelkind == RELKIND_SEQUENCE ? "sequence" : "table"));
    1187                 :        1123 : }
        

Generated by: LCOV version 2.0-1