LCOV - code coverage report
Current view: top level - contrib/test_decoding - test_decoding.c (source / functions) Coverage Total Hit
Test: PostgreSQL 20devel Lines: 86.1 % 410 353
Test Date: 2026-09-05 06:15:58 Functions: 96.4 % 28 27
Legend: Lines:     hit not hit
Branches: + taken - not taken # not executed
Branches: 68.6 % 236 162

             Branch data     Line data    Source code
       1                 :             : /*-------------------------------------------------------------------------
       2                 :             :  *
       3                 :             :  * test_decoding.c
       4                 :             :  *        example logical decoding output plugin
       5                 :             :  *
       6                 :             :  * Copyright (c) 2012-2026, PostgreSQL Global Development Group
       7                 :             :  *
       8                 :             :  * IDENTIFICATION
       9                 :             :  *        contrib/test_decoding/test_decoding.c
      10                 :             :  *
      11                 :             :  *-------------------------------------------------------------------------
      12                 :             :  */
      13                 :             : #include "postgres.h"
      14                 :             : 
      15                 :             : #include "catalog/pg_type.h"
      16                 :             : 
      17                 :             : #include "replication/logical.h"
      18                 :             : #include "replication/origin.h"
      19                 :             : 
      20                 :             : #include "utils/builtins.h"
      21                 :             : #include "utils/lsyscache.h"
      22                 :             : #include "utils/memutils.h"
      23                 :             : #include "utils/rel.h"
      24                 :             : 
      25                 :         140 : PG_MODULE_MAGIC_EXT(
      26                 :             :                     .name = "test_decoding",
      27                 :             :                     .version = PG_VERSION
      28                 :             : );
      29                 :             : 
      30                 :             : typedef struct
      31                 :             : {
      32                 :             :     MemoryContext context;
      33                 :             :     bool        include_xids;
      34                 :             :     bool        include_timestamp;
      35                 :             :     bool        skip_empty_xacts;
      36                 :             :     bool        only_local;
      37                 :             : } TestDecodingData;
      38                 :             : 
      39                 :             : /*
      40                 :             :  * Maintain the per-transaction level variables to track whether the
      41                 :             :  * transaction and or streams have written any changes. In streaming mode the
      42                 :             :  * transaction can be decoded in streams so along with maintaining whether the
      43                 :             :  * transaction has written any changes, we also need to track whether the
      44                 :             :  * current stream has written any changes. This is required so that if user
      45                 :             :  * has requested to skip the empty transactions we can skip the empty streams
      46                 :             :  * even though the transaction has written some changes.
      47                 :             :  */
      48                 :             : typedef struct
      49                 :             : {
      50                 :             :     bool        xact_wrote_changes;
      51                 :             :     bool        stream_wrote_changes;
      52                 :             : } TestDecodingTxnData;
      53                 :             : 
      54                 :             : static void pg_decode_startup(LogicalDecodingContext *ctx, OutputPluginOptions *opt,
      55                 :             :                               bool is_init);
      56                 :             : static void pg_decode_shutdown(LogicalDecodingContext *ctx);
      57                 :             : static void pg_decode_begin_txn(LogicalDecodingContext *ctx,
      58                 :             :                                 ReorderBufferTXN *txn);
      59                 :             : static void pg_output_begin(LogicalDecodingContext *ctx,
      60                 :             :                             TestDecodingData *data,
      61                 :             :                             ReorderBufferTXN *txn,
      62                 :             :                             bool last_write);
      63                 :             : static void pg_decode_commit_txn(LogicalDecodingContext *ctx,
      64                 :             :                                  ReorderBufferTXN *txn, XLogRecPtr commit_lsn);
      65                 :             : static void pg_decode_change(LogicalDecodingContext *ctx,
      66                 :             :                              ReorderBufferTXN *txn, Relation relation,
      67                 :             :                              ReorderBufferChange *change);
      68                 :             : static void pg_decode_truncate(LogicalDecodingContext *ctx,
      69                 :             :                                ReorderBufferTXN *txn,
      70                 :             :                                int nrelations, Relation relations[],
      71                 :             :                                ReorderBufferChange *change);
      72                 :             : static bool pg_decode_filter(LogicalDecodingContext *ctx,
      73                 :             :                              ReplOriginId origin_id);
      74                 :             : static void pg_decode_message(LogicalDecodingContext *ctx,
      75                 :             :                               ReorderBufferTXN *txn, XLogRecPtr lsn,
      76                 :             :                               bool transactional, const char *prefix,
      77                 :             :                               Size sz, const char *message);
      78                 :             : static bool pg_decode_filter_prepare(LogicalDecodingContext *ctx,
      79                 :             :                                      TransactionId xid,
      80                 :             :                                      const char *gid);
      81                 :             : static void pg_decode_begin_prepare_txn(LogicalDecodingContext *ctx,
      82                 :             :                                         ReorderBufferTXN *txn);
      83                 :             : static void pg_decode_prepare_txn(LogicalDecodingContext *ctx,
      84                 :             :                                   ReorderBufferTXN *txn,
      85                 :             :                                   XLogRecPtr prepare_lsn);
      86                 :             : static void pg_decode_commit_prepared_txn(LogicalDecodingContext *ctx,
      87                 :             :                                           ReorderBufferTXN *txn,
      88                 :             :                                           XLogRecPtr commit_lsn);
      89                 :             : static void pg_decode_rollback_prepared_txn(LogicalDecodingContext *ctx,
      90                 :             :                                             ReorderBufferTXN *txn,
      91                 :             :                                             XLogRecPtr prepare_end_lsn,
      92                 :             :                                             TimestampTz prepare_time);
      93                 :             : static void pg_decode_stream_start(LogicalDecodingContext *ctx,
      94                 :             :                                    ReorderBufferTXN *txn);
      95                 :             : static void pg_output_stream_start(LogicalDecodingContext *ctx,
      96                 :             :                                    TestDecodingData *data,
      97                 :             :                                    ReorderBufferTXN *txn,
      98                 :             :                                    bool last_write);
      99                 :             : static void pg_decode_stream_stop(LogicalDecodingContext *ctx,
     100                 :             :                                   ReorderBufferTXN *txn);
     101                 :             : static void pg_decode_stream_abort(LogicalDecodingContext *ctx,
     102                 :             :                                    ReorderBufferTXN *txn,
     103                 :             :                                    XLogRecPtr abort_lsn);
     104                 :             : static void pg_decode_stream_prepare(LogicalDecodingContext *ctx,
     105                 :             :                                      ReorderBufferTXN *txn,
     106                 :             :                                      XLogRecPtr prepare_lsn);
     107                 :             : static void pg_decode_stream_commit(LogicalDecodingContext *ctx,
     108                 :             :                                     ReorderBufferTXN *txn,
     109                 :             :                                     XLogRecPtr commit_lsn);
     110                 :             : static void pg_decode_stream_change(LogicalDecodingContext *ctx,
     111                 :             :                                     ReorderBufferTXN *txn,
     112                 :             :                                     Relation relation,
     113                 :             :                                     ReorderBufferChange *change);
     114                 :             : static void pg_decode_stream_message(LogicalDecodingContext *ctx,
     115                 :             :                                      ReorderBufferTXN *txn, XLogRecPtr lsn,
     116                 :             :                                      bool transactional, const char *prefix,
     117                 :             :                                      Size sz, const char *message);
     118                 :             : static void pg_decode_stream_truncate(LogicalDecodingContext *ctx,
     119                 :             :                                       ReorderBufferTXN *txn,
     120                 :             :                                       int nrelations, Relation relations[],
     121                 :             :                                       ReorderBufferChange *change);
     122                 :             : 
     123                 :             : void
     124                 :         140 : _PG_init(void)
     125                 :             : {
     126                 :             :     /* other plugins can perform things here */
     127                 :         140 : }
     128                 :             : 
     129                 :             : /* specify output plugin callbacks */
     130                 :             : void
     131                 :         374 : _PG_output_plugin_init(OutputPluginCallbacks *cb)
     132                 :             : {
     133                 :         374 :     cb->startup_cb = pg_decode_startup;
     134                 :         374 :     cb->begin_cb = pg_decode_begin_txn;
     135                 :         374 :     cb->change_cb = pg_decode_change;
     136                 :         374 :     cb->truncate_cb = pg_decode_truncate;
     137                 :         374 :     cb->commit_cb = pg_decode_commit_txn;
     138                 :         374 :     cb->filter_by_origin_cb = pg_decode_filter;
     139                 :         374 :     cb->shutdown_cb = pg_decode_shutdown;
     140                 :         374 :     cb->message_cb = pg_decode_message;
     141                 :         374 :     cb->filter_prepare_cb = pg_decode_filter_prepare;
     142                 :         374 :     cb->begin_prepare_cb = pg_decode_begin_prepare_txn;
     143                 :         374 :     cb->prepare_cb = pg_decode_prepare_txn;
     144                 :         374 :     cb->commit_prepared_cb = pg_decode_commit_prepared_txn;
     145                 :         374 :     cb->rollback_prepared_cb = pg_decode_rollback_prepared_txn;
     146                 :         374 :     cb->stream_start_cb = pg_decode_stream_start;
     147                 :         374 :     cb->stream_stop_cb = pg_decode_stream_stop;
     148                 :         374 :     cb->stream_abort_cb = pg_decode_stream_abort;
     149                 :         374 :     cb->stream_prepare_cb = pg_decode_stream_prepare;
     150                 :         374 :     cb->stream_commit_cb = pg_decode_stream_commit;
     151                 :         374 :     cb->stream_change_cb = pg_decode_stream_change;
     152                 :         374 :     cb->stream_message_cb = pg_decode_stream_message;
     153                 :         374 :     cb->stream_truncate_cb = pg_decode_stream_truncate;
     154                 :         374 : }
     155                 :             : 
     156                 :             : 
     157                 :             : /* initialize this plugin */
     158                 :             : static void
     159                 :         374 : pg_decode_startup(LogicalDecodingContext *ctx, OutputPluginOptions *opt,
     160                 :             :                   bool is_init)
     161                 :             : {
     162                 :             :     ListCell   *option;
     163                 :             :     TestDecodingData *data;
     164                 :         374 :     bool        enable_streaming = false;
     165                 :             : 
     166                 :         374 :     data = palloc0_object(TestDecodingData);
     167                 :         374 :     data->context = AllocSetContextCreate(ctx->context,
     168                 :             :                                           "text conversion context",
     169                 :             :                                           ALLOCSET_DEFAULT_SIZES);
     170                 :         374 :     data->include_xids = true;
     171                 :         374 :     data->include_timestamp = false;
     172                 :         374 :     data->skip_empty_xacts = false;
     173                 :         374 :     data->only_local = false;
     174                 :             : 
     175                 :         374 :     ctx->output_plugin_private = data;
     176                 :             : 
     177                 :         374 :     opt->output_type = OUTPUT_PLUGIN_TEXTUAL_OUTPUT;
     178                 :         374 :     opt->receive_rewrites = false;
     179                 :             : 
     180   [ +  +  +  +  :         757 :     foreach(option, ctx->output_plugin_options)
                   +  + ]
     181                 :             :     {
     182                 :         386 :         DefElem    *elem = lfirst(option);
     183                 :             : 
     184                 :             :         Assert(elem->arg == NULL || IsA(elem->arg, String));
     185                 :             : 
     186         [ +  + ]:         386 :         if (strcmp(elem->defname, "include-xids") == 0)
     187                 :             :         {
     188                 :             :             /* if option does not provide a value, it means its value is true */
     189         [ -  + ]:         182 :             if (elem->arg == NULL)
     190                 :           0 :                 data->include_xids = true;
     191         [ +  + ]:         182 :             else if (!parse_bool(strVal(elem->arg), &data->include_xids))
     192         [ +  - ]:           2 :                 ereport(ERROR,
     193                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     194                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     195                 :             :                                 strVal(elem->arg), elem->defname)));
     196                 :             :         }
     197         [ +  + ]:         204 :         else if (strcmp(elem->defname, "include-timestamp") == 0)
     198                 :             :         {
     199         [ -  + ]:           1 :             if (elem->arg == NULL)
     200                 :           0 :                 data->include_timestamp = true;
     201         [ -  + ]:           1 :             else if (!parse_bool(strVal(elem->arg), &data->include_timestamp))
     202         [ #  # ]:           0 :                 ereport(ERROR,
     203                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     204                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     205                 :             :                                 strVal(elem->arg), elem->defname)));
     206                 :             :         }
     207         [ +  + ]:         203 :         else if (strcmp(elem->defname, "force-binary") == 0)
     208                 :             :         {
     209                 :             :             bool        force_binary;
     210                 :             : 
     211         [ -  + ]:           6 :             if (elem->arg == NULL)
     212                 :           0 :                 continue;
     213         [ -  + ]:           6 :             else if (!parse_bool(strVal(elem->arg), &force_binary))
     214         [ #  # ]:           0 :                 ereport(ERROR,
     215                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     216                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     217                 :             :                                 strVal(elem->arg), elem->defname)));
     218                 :             : 
     219         [ +  + ]:           6 :             if (force_binary)
     220                 :           2 :                 opt->output_type = OUTPUT_PLUGIN_BINARY_OUTPUT;
     221                 :             :         }
     222         [ +  + ]:         197 :         else if (strcmp(elem->defname, "skip-empty-xacts") == 0)
     223                 :             :         {
     224                 :             : 
     225         [ -  + ]:         180 :             if (elem->arg == NULL)
     226                 :           0 :                 data->skip_empty_xacts = true;
     227         [ -  + ]:         180 :             else if (!parse_bool(strVal(elem->arg), &data->skip_empty_xacts))
     228         [ #  # ]:           0 :                 ereport(ERROR,
     229                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     230                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     231                 :             :                                 strVal(elem->arg), elem->defname)));
     232                 :             :         }
     233         [ +  + ]:          17 :         else if (strcmp(elem->defname, "only-local") == 0)
     234                 :             :         {
     235                 :             : 
     236         [ -  + ]:           3 :             if (elem->arg == NULL)
     237                 :           0 :                 data->only_local = true;
     238         [ -  + ]:           3 :             else if (!parse_bool(strVal(elem->arg), &data->only_local))
     239         [ #  # ]:           0 :                 ereport(ERROR,
     240                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     241                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     242                 :             :                                 strVal(elem->arg), elem->defname)));
     243                 :             :         }
     244         [ +  + ]:          14 :         else if (strcmp(elem->defname, "include-rewrites") == 0)
     245                 :             :         {
     246                 :             : 
     247         [ -  + ]:           1 :             if (elem->arg == NULL)
     248                 :           0 :                 continue;
     249         [ -  + ]:           1 :             else if (!parse_bool(strVal(elem->arg), &opt->receive_rewrites))
     250         [ #  # ]:           0 :                 ereport(ERROR,
     251                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     252                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     253                 :             :                                 strVal(elem->arg), elem->defname)));
     254                 :             :         }
     255         [ +  + ]:          13 :         else if (strcmp(elem->defname, "stream-changes") == 0)
     256                 :             :         {
     257         [ -  + ]:          12 :             if (elem->arg == NULL)
     258                 :           0 :                 continue;
     259         [ -  + ]:          12 :             else if (!parse_bool(strVal(elem->arg), &enable_streaming))
     260         [ #  # ]:           0 :                 ereport(ERROR,
     261                 :             :                         (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     262                 :             :                          errmsg("could not parse value \"%s\" for parameter \"%s\"",
     263                 :             :                                 strVal(elem->arg), elem->defname)));
     264                 :             :         }
     265                 :             :         else
     266                 :             :         {
     267   [ +  -  +  - ]:           1 :             ereport(ERROR,
     268                 :             :                     (errcode(ERRCODE_INVALID_PARAMETER_VALUE),
     269                 :             :                      errmsg("option \"%s\" = \"%s\" is unknown",
     270                 :             :                             elem->defname,
     271                 :             :                             elem->arg ? strVal(elem->arg) : "(null)")));
     272                 :             :         }
     273                 :             :     }
     274                 :             : 
     275                 :         371 :     ctx->streaming &= enable_streaming;
     276                 :         371 : }
     277                 :             : 
     278                 :             : /* cleanup this plugin's resources */
     279                 :             : static void
     280                 :         360 : pg_decode_shutdown(LogicalDecodingContext *ctx)
     281                 :             : {
     282                 :         360 :     TestDecodingData *data = ctx->output_plugin_private;
     283                 :             : 
     284                 :             :     /* cleanup our own resources via memory context reset */
     285                 :         360 :     MemoryContextDelete(data->context);
     286                 :         360 : }
     287                 :             : 
     288                 :             : /* BEGIN callback */
     289                 :             : static void
     290                 :         463 : pg_decode_begin_txn(LogicalDecodingContext *ctx, ReorderBufferTXN *txn)
     291                 :             : {
     292                 :         463 :     TestDecodingData *data = ctx->output_plugin_private;
     293                 :             :     TestDecodingTxnData *txndata =
     294                 :         463 :         MemoryContextAllocZero(ctx->context, sizeof(TestDecodingTxnData));
     295                 :             : 
     296                 :         463 :     txndata->xact_wrote_changes = false;
     297                 :         463 :     txn->output_plugin_private = txndata;
     298                 :             : 
     299                 :             :     /*
     300                 :             :      * If asked to skip empty transactions, we'll emit BEGIN at the point
     301                 :             :      * where the first operation is received for this transaction.
     302                 :             :      */
     303         [ +  + ]:         463 :     if (data->skip_empty_xacts)
     304                 :         413 :         return;
     305                 :             : 
     306                 :          50 :     pg_output_begin(ctx, data, txn, true);
     307                 :             : }
     308                 :             : 
     309                 :             : static void
     310                 :         290 : pg_output_begin(LogicalDecodingContext *ctx, TestDecodingData *data, ReorderBufferTXN *txn, bool last_write)
     311                 :             : {
     312                 :         290 :     OutputPluginPrepareWrite(ctx, last_write);
     313         [ +  + ]:         290 :     if (data->include_xids)
     314                 :          46 :         appendStringInfo(ctx->out, "BEGIN %u", txn->xid);
     315                 :             :     else
     316                 :         244 :         appendStringInfoString(ctx->out, "BEGIN");
     317                 :         290 :     OutputPluginWrite(ctx, last_write);
     318                 :         290 : }
     319                 :             : 
     320                 :             : /* COMMIT callback */
     321                 :             : static void
     322                 :         463 : pg_decode_commit_txn(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
     323                 :             :                      XLogRecPtr commit_lsn)
     324                 :             : {
     325                 :         463 :     TestDecodingData *data = ctx->output_plugin_private;
     326                 :         463 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     327                 :         463 :     bool        xact_wrote_changes = txndata->xact_wrote_changes;
     328                 :             : 
     329                 :         463 :     pfree(txndata);
     330                 :         463 :     txn->output_plugin_private = NULL;
     331                 :             : 
     332   [ +  +  +  + ]:         463 :     if (data->skip_empty_xacts && !xact_wrote_changes)
     333                 :         181 :         return;
     334                 :             : 
     335                 :         282 :     OutputPluginPrepareWrite(ctx, true);
     336         [ +  + ]:         282 :     if (data->include_xids)
     337                 :          45 :         appendStringInfo(ctx->out, "COMMIT %u", txn->xid);
     338                 :             :     else
     339                 :         237 :         appendStringInfoString(ctx->out, "COMMIT");
     340                 :             : 
     341         [ +  + ]:         282 :     if (data->include_timestamp)
     342                 :           1 :         appendStringInfo(ctx->out, " (at %s)",
     343                 :             :                          timestamptz_to_str(txn->commit_time));
     344                 :             : 
     345                 :         282 :     OutputPluginWrite(ctx, true);
     346                 :             : }
     347                 :             : 
     348                 :             : /* BEGIN PREPARE callback */
     349                 :             : static void
     350                 :           9 : pg_decode_begin_prepare_txn(LogicalDecodingContext *ctx, ReorderBufferTXN *txn)
     351                 :             : {
     352                 :           9 :     TestDecodingData *data = ctx->output_plugin_private;
     353                 :             :     TestDecodingTxnData *txndata =
     354                 :           9 :         MemoryContextAllocZero(ctx->context, sizeof(TestDecodingTxnData));
     355                 :             : 
     356                 :           9 :     txndata->xact_wrote_changes = false;
     357                 :           9 :     txn->output_plugin_private = txndata;
     358                 :             : 
     359                 :             :     /*
     360                 :             :      * If asked to skip empty transactions, we'll emit BEGIN at the point
     361                 :             :      * where the first operation is received for this transaction.
     362                 :             :      */
     363         [ +  + ]:           9 :     if (data->skip_empty_xacts)
     364                 :           8 :         return;
     365                 :             : 
     366                 :           1 :     pg_output_begin(ctx, data, txn, true);
     367                 :             : }
     368                 :             : 
     369                 :             : /* PREPARE callback */
     370                 :             : static void
     371                 :           9 : pg_decode_prepare_txn(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
     372                 :             :                       XLogRecPtr prepare_lsn)
     373                 :             : {
     374                 :           9 :     TestDecodingData *data = ctx->output_plugin_private;
     375                 :           9 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     376                 :             : 
     377                 :             :     /*
     378                 :             :      * If asked to skip empty transactions, we'll emit PREPARE at the point
     379                 :             :      * where the first operation is received for this transaction.
     380                 :             :      */
     381   [ +  +  +  + ]:           9 :     if (data->skip_empty_xacts && !txndata->xact_wrote_changes)
     382                 :           1 :         return;
     383                 :             : 
     384                 :           8 :     OutputPluginPrepareWrite(ctx, true);
     385                 :             : 
     386                 :           8 :     appendStringInfo(ctx->out, "PREPARE TRANSACTION %s",
     387                 :           8 :                      quote_literal_cstr(txn->gid));
     388                 :             : 
     389         [ +  + ]:           8 :     if (data->include_xids)
     390                 :           1 :         appendStringInfo(ctx->out, ", txid %u", txn->xid);
     391                 :             : 
     392         [ -  + ]:           8 :     if (data->include_timestamp)
     393                 :           0 :         appendStringInfo(ctx->out, " (at %s)",
     394                 :             :                          timestamptz_to_str(txn->prepare_time));
     395                 :             : 
     396                 :           8 :     OutputPluginWrite(ctx, true);
     397                 :             : }
     398                 :             : 
     399                 :             : /* COMMIT PREPARED callback */
     400                 :             : static void
     401                 :           8 : pg_decode_commit_prepared_txn(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
     402                 :             :                               XLogRecPtr commit_lsn)
     403                 :             : {
     404                 :           8 :     TestDecodingData *data = ctx->output_plugin_private;
     405                 :             : 
     406                 :           8 :     OutputPluginPrepareWrite(ctx, true);
     407                 :             : 
     408                 :           8 :     appendStringInfo(ctx->out, "COMMIT PREPARED %s",
     409                 :           8 :                      quote_literal_cstr(txn->gid));
     410                 :             : 
     411         [ +  + ]:           8 :     if (data->include_xids)
     412                 :           1 :         appendStringInfo(ctx->out, ", txid %u", txn->xid);
     413                 :             : 
     414         [ -  + ]:           8 :     if (data->include_timestamp)
     415                 :           0 :         appendStringInfo(ctx->out, " (at %s)",
     416                 :             :                          timestamptz_to_str(txn->commit_time));
     417                 :             : 
     418                 :           8 :     OutputPluginWrite(ctx, true);
     419                 :           8 : }
     420                 :             : 
     421                 :             : /* ROLLBACK PREPARED callback */
     422                 :             : static void
     423                 :           2 : pg_decode_rollback_prepared_txn(LogicalDecodingContext *ctx,
     424                 :             :                                 ReorderBufferTXN *txn,
     425                 :             :                                 XLogRecPtr prepare_end_lsn,
     426                 :             :                                 TimestampTz prepare_time)
     427                 :             : {
     428                 :           2 :     TestDecodingData *data = ctx->output_plugin_private;
     429                 :             : 
     430                 :           2 :     OutputPluginPrepareWrite(ctx, true);
     431                 :             : 
     432                 :           2 :     appendStringInfo(ctx->out, "ROLLBACK PREPARED %s",
     433                 :           2 :                      quote_literal_cstr(txn->gid));
     434                 :             : 
     435         [ -  + ]:           2 :     if (data->include_xids)
     436                 :           0 :         appendStringInfo(ctx->out, ", txid %u", txn->xid);
     437                 :             : 
     438         [ -  + ]:           2 :     if (data->include_timestamp)
     439                 :           0 :         appendStringInfo(ctx->out, " (at %s)",
     440                 :             :                          timestamptz_to_str(txn->commit_time));
     441                 :             : 
     442                 :           2 :     OutputPluginWrite(ctx, true);
     443                 :           2 : }
     444                 :             : 
     445                 :             : /*
     446                 :             :  * Filter out two-phase transactions.
     447                 :             :  *
     448                 :             :  * Each plugin can implement its own filtering logic. Here we demonstrate a
     449                 :             :  * simple logic by checking the GID. If the GID contains the "_nodecode"
     450                 :             :  * substring, then we filter it out.
     451                 :             :  */
     452                 :             : static bool
     453                 :         170 : pg_decode_filter_prepare(LogicalDecodingContext *ctx, TransactionId xid,
     454                 :             :                          const char *gid)
     455                 :             : {
     456         [ +  + ]:         170 :     if (strstr(gid, "_nodecode") != NULL)
     457                 :          14 :         return true;
     458                 :             : 
     459                 :         156 :     return false;
     460                 :             : }
     461                 :             : 
     462                 :             : static bool
     463                 :      943229 : pg_decode_filter(LogicalDecodingContext *ctx,
     464                 :             :                  ReplOriginId origin_id)
     465                 :             : {
     466                 :      943229 :     TestDecodingData *data = ctx->output_plugin_private;
     467                 :             : 
     468   [ +  +  +  + ]:      943229 :     if (data->only_local && origin_id != InvalidReplOriginId)
     469                 :           9 :         return true;
     470                 :      943220 :     return false;
     471                 :             : }
     472                 :             : 
     473                 :             : /*
     474                 :             :  * Print literal `outputstr' already represented as string of type `typid'
     475                 :             :  * into stringbuf `s'.
     476                 :             :  *
     477                 :             :  * Some builtin types aren't quoted, the rest is quoted. Escaping is done
     478                 :             :  * per standard SQL rules.
     479                 :             :  */
     480                 :             : static void
     481                 :      176280 : print_literal(StringInfo s, Oid typid, char *outputstr)
     482                 :             : {
     483                 :             :     const char *valptr;
     484                 :             : 
     485   [ +  -  -  + ]:      176280 :     switch (typid)
     486                 :             :     {
     487                 :       60417 :         case INT2OID:
     488                 :             :         case INT4OID:
     489                 :             :         case INT8OID:
     490                 :             :         case OIDOID:
     491                 :             :         case FLOAT4OID:
     492                 :             :         case FLOAT8OID:
     493                 :             :         case NUMERICOID:
     494                 :             :             /* NB: We don't care about Inf, NaN et al. */
     495                 :       60417 :             appendStringInfoString(s, outputstr);
     496                 :       60417 :             break;
     497                 :             : 
     498                 :           0 :         case BITOID:
     499                 :             :         case VARBITOID:
     500                 :           0 :             appendStringInfo(s, "B'%s'", outputstr);
     501                 :           0 :             break;
     502                 :             : 
     503                 :           0 :         case BOOLOID:
     504         [ #  # ]:           0 :             if (strcmp(outputstr, "t") == 0)
     505                 :           0 :                 appendStringInfoString(s, "true");
     506                 :             :             else
     507                 :           0 :                 appendStringInfoString(s, "false");
     508                 :           0 :             break;
     509                 :             : 
     510                 :      115863 :         default:
     511                 :      115863 :             appendStringInfoChar(s, '\'');
     512         [ +  + ]:     5427925 :             for (valptr = outputstr; *valptr; valptr++)
     513                 :             :             {
     514                 :     5312062 :                 char        ch = *valptr;
     515                 :             : 
     516         [ +  + ]:     5312062 :                 if (SQL_STR_DOUBLE(ch, false))
     517                 :          64 :                     appendStringInfoChar(s, ch);
     518                 :     5312062 :                 appendStringInfoChar(s, ch);
     519                 :             :             }
     520                 :      115863 :             appendStringInfoChar(s, '\'');
     521                 :      115863 :             break;
     522                 :             :     }
     523                 :      176280 : }
     524                 :             : 
     525                 :             : /* print the tuple 'tuple' into the StringInfo s */
     526                 :             : static void
     527                 :      145730 : tuple_to_stringinfo(StringInfo s, TupleDesc tupdesc, HeapTuple tuple, bool skip_nulls)
     528                 :             : {
     529                 :             :     int         natt;
     530                 :             : 
     531                 :             :     /* print all columns individually */
     532         [ +  + ]:      347774 :     for (natt = 0; natt < tupdesc->natts; natt++)
     533                 :             :     {
     534                 :             :         Form_pg_attribute attr; /* the attribute itself */
     535                 :             :         Oid         typid;      /* type of current attribute */
     536                 :             :         Oid         typoutput;  /* output function */
     537                 :             :         bool        typisvarlena;
     538                 :             :         Datum       origval;    /* possibly toasted Datum */
     539                 :             :         bool        isnull;     /* column is null? */
     540                 :             : 
     541                 :      202044 :         attr = TupleDescAttr(tupdesc, natt);
     542                 :             : 
     543                 :             :         /*
     544                 :             :          * don't print dropped columns, we can't be sure everything is
     545                 :             :          * available for them
     546                 :             :          */
     547         [ +  + ]:      202044 :         if (attr->attisdropped)
     548                 :        5135 :             continue;
     549                 :             : 
     550                 :             :         /*
     551                 :             :          * Don't print system columns, oid will already have been printed if
     552                 :             :          * present.
     553                 :             :          */
     554         [ -  + ]:      201972 :         if (attr->attnum < 0)
     555                 :           0 :             continue;
     556                 :             : 
     557                 :             :         /*
     558                 :             :          * Virtual generated columns are always stored as null in the tuple,
     559                 :             :          * so don't print them at all. A printed null would not be
     560                 :             :          * distinguishable from a column that really contains a null. Stored
     561                 :             :          * generated columns are printed as usual since their values are
     562                 :             :          * actually on disk.
     563                 :             :          */
     564         [ +  + ]:      201972 :         if (attr->attgenerated == ATTRIBUTE_GENERATED_VIRTUAL)
     565                 :           1 :             continue;
     566                 :             : 
     567                 :      201971 :         typid = attr->atttypid;
     568                 :             : 
     569                 :             :         /* get Datum from tuple */
     570                 :      201971 :         origval = heap_getattr(tuple, natt + 1, tupdesc, &isnull);
     571                 :             : 
     572   [ +  +  +  + ]:      201971 :         if (isnull && skip_nulls)
     573                 :        5062 :             continue;
     574                 :             : 
     575                 :             :         /* print attribute name */
     576                 :      196909 :         appendStringInfoChar(s, ' ');
     577                 :      196909 :         appendStringInfoString(s, quote_identifier(NameStr(attr->attname)));
     578                 :             : 
     579                 :             :         /* print attribute type */
     580                 :      196909 :         appendStringInfoChar(s, '[');
     581                 :      196909 :         appendStringInfoString(s, format_type_be(typid));
     582                 :      196909 :         appendStringInfoChar(s, ']');
     583                 :             : 
     584                 :             :         /* query output function */
     585                 :      196909 :         getTypeOutputInfo(typid,
     586                 :             :                           &typoutput, &typisvarlena);
     587                 :             : 
     588                 :             :         /* print separator */
     589                 :      196909 :         appendStringInfoChar(s, ':');
     590                 :             : 
     591                 :             :         /* print data */
     592         [ +  + ]:      196909 :         if (isnull)
     593                 :       20617 :             appendStringInfoString(s, "null");
     594   [ +  +  +  + ]:      176292 :         else if (typisvarlena && VARATT_IS_EXTERNAL_ONDISK(DatumGetPointer(origval)))
     595                 :          12 :             appendStringInfoString(s, "unchanged-toast-datum");
     596         [ +  + ]:      176280 :         else if (!typisvarlena)
     597                 :       60421 :             print_literal(s, typid,
     598                 :             :                           OidOutputFunctionCall(typoutput, origval));
     599                 :             :         else
     600                 :             :         {
     601                 :             :             Datum       val;    /* definitely detoasted Datum */
     602                 :             : 
     603                 :      115859 :             val = PointerGetDatum(PG_DETOAST_DATUM(origval));
     604                 :      115859 :             print_literal(s, typid, OidOutputFunctionCall(typoutput, val));
     605                 :             :         }
     606                 :             :     }
     607                 :      145730 : }
     608                 :             : 
     609                 :             : /*
     610                 :             :  * callback for individual changed tuples
     611                 :             :  */
     612                 :             : static void
     613                 :      150718 : pg_decode_change(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
     614                 :             :                  Relation relation, ReorderBufferChange *change)
     615                 :             : {
     616                 :             :     TestDecodingData *data;
     617                 :             :     TestDecodingTxnData *txndata;
     618                 :             :     Form_pg_class class_form;
     619                 :             :     TupleDesc   tupdesc;
     620                 :             :     MemoryContext old;
     621                 :             : 
     622                 :      150718 :     data = ctx->output_plugin_private;
     623                 :      150718 :     txndata = txn->output_plugin_private;
     624                 :             : 
     625                 :             :     /* output BEGIN if we haven't yet */
     626   [ +  +  +  + ]:      150718 :     if (data->skip_empty_xacts && !txndata->xact_wrote_changes)
     627                 :             :     {
     628                 :         229 :         pg_output_begin(ctx, data, txn, false);
     629                 :             :     }
     630                 :      150718 :     txndata->xact_wrote_changes = true;
     631                 :             : 
     632                 :      150718 :     class_form = RelationGetForm(relation);
     633                 :      150718 :     tupdesc = RelationGetDescr(relation);
     634                 :             : 
     635                 :             :     /* Avoid leaking memory by using and resetting our own context */
     636                 :      150718 :     old = MemoryContextSwitchTo(data->context);
     637                 :             : 
     638                 :      150718 :     OutputPluginPrepareWrite(ctx, true);
     639                 :             : 
     640                 :      150718 :     appendStringInfoString(ctx->out, "table ");
     641                 :      150718 :     appendStringInfoString(ctx->out,
     642                 :      150718 :                            quote_qualified_identifier(get_namespace_name(get_rel_namespace(RelationGetRelid(relation))),
     643         [ +  + ]:      150718 :                                                       class_form->relrewrite ?
     644                 :           1 :                                                       get_rel_name(class_form->relrewrite) :
     645                 :             :                                                       NameStr(class_form->relname)));
     646                 :      150718 :     appendStringInfoChar(ctx->out, ':');
     647                 :             : 
     648   [ +  +  +  - ]:      150718 :     switch (change->action)
     649                 :             :     {
     650                 :      133151 :         case REORDER_BUFFER_CHANGE_INSERT:
     651                 :      133151 :             appendStringInfoString(ctx->out, " INSERT:");
     652         [ -  + ]:      133151 :             if (change->data.tp.newtuple == NULL)
     653                 :           0 :                 appendStringInfoString(ctx->out, " (no-tuple-data)");
     654                 :             :             else
     655                 :      133151 :                 tuple_to_stringinfo(ctx->out, tupdesc,
     656                 :             :                                     change->data.tp.newtuple,
     657                 :             :                                     false);
     658                 :      133151 :             break;
     659                 :        7544 :         case REORDER_BUFFER_CHANGE_UPDATE:
     660                 :        7544 :             appendStringInfoString(ctx->out, " UPDATE:");
     661         [ +  + ]:        7544 :             if (change->data.tp.oldtuple != NULL)
     662                 :             :             {
     663                 :          19 :                 appendStringInfoString(ctx->out, " old-key:");
     664                 :          19 :                 tuple_to_stringinfo(ctx->out, tupdesc,
     665                 :             :                                     change->data.tp.oldtuple,
     666                 :             :                                     true);
     667                 :          19 :                 appendStringInfoString(ctx->out, " new-tuple:");
     668                 :             :             }
     669                 :             : 
     670         [ -  + ]:        7544 :             if (change->data.tp.newtuple == NULL)
     671                 :           0 :                 appendStringInfoString(ctx->out, " (no-tuple-data)");
     672                 :             :             else
     673                 :        7544 :                 tuple_to_stringinfo(ctx->out, tupdesc,
     674                 :             :                                     change->data.tp.newtuple,
     675                 :             :                                     false);
     676                 :        7544 :             break;
     677                 :       10023 :         case REORDER_BUFFER_CHANGE_DELETE:
     678                 :       10023 :             appendStringInfoString(ctx->out, " DELETE:");
     679                 :             : 
     680                 :             :             /* if there was no PK, we only know that a delete happened */
     681         [ +  + ]:       10023 :             if (change->data.tp.oldtuple == NULL)
     682                 :        5007 :                 appendStringInfoString(ctx->out, " (no-tuple-data)");
     683                 :             :             /* In DELETE, only the replica identity is present; display that */
     684                 :             :             else
     685                 :        5016 :                 tuple_to_stringinfo(ctx->out, tupdesc,
     686                 :             :                                     change->data.tp.oldtuple,
     687                 :             :                                     true);
     688                 :       10023 :             break;
     689                 :      150718 :         default:
     690                 :             :             Assert(false);
     691                 :             :     }
     692                 :             : 
     693                 :      150718 :     MemoryContextSwitchTo(old);
     694                 :      150718 :     MemoryContextReset(data->context);
     695                 :             : 
     696                 :      150718 :     OutputPluginWrite(ctx, true);
     697                 :      150718 : }
     698                 :             : 
     699                 :             : static void
     700                 :           8 : pg_decode_truncate(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
     701                 :             :                    int nrelations, Relation relations[], ReorderBufferChange *change)
     702                 :             : {
     703                 :             :     TestDecodingData *data;
     704                 :             :     TestDecodingTxnData *txndata;
     705                 :             :     MemoryContext old;
     706                 :             :     int         i;
     707                 :             : 
     708                 :           8 :     data = ctx->output_plugin_private;
     709                 :           8 :     txndata = txn->output_plugin_private;
     710                 :             : 
     711                 :             :     /* output BEGIN if we haven't yet */
     712   [ +  +  +  - ]:           8 :     if (data->skip_empty_xacts && !txndata->xact_wrote_changes)
     713                 :             :     {
     714                 :           7 :         pg_output_begin(ctx, data, txn, false);
     715                 :             :     }
     716                 :           8 :     txndata->xact_wrote_changes = true;
     717                 :             : 
     718                 :             :     /* Avoid leaking memory by using and resetting our own context */
     719                 :           8 :     old = MemoryContextSwitchTo(data->context);
     720                 :             : 
     721                 :           8 :     OutputPluginPrepareWrite(ctx, true);
     722                 :             : 
     723                 :           8 :     appendStringInfoString(ctx->out, "table ");
     724                 :             : 
     725         [ +  + ]:          17 :     for (i = 0; i < nrelations; i++)
     726                 :             :     {
     727         [ +  + ]:           9 :         if (i > 0)
     728                 :           1 :             appendStringInfoString(ctx->out, ", ");
     729                 :             : 
     730                 :           9 :         appendStringInfoString(ctx->out,
     731                 :           9 :                                quote_qualified_identifier(get_namespace_name(relations[i]->rd_rel->relnamespace),
     732                 :           9 :                                                           NameStr(relations[i]->rd_rel->relname)));
     733                 :             :     }
     734                 :             : 
     735                 :           8 :     appendStringInfoString(ctx->out, ": TRUNCATE:");
     736                 :             : 
     737         [ +  + ]:           8 :     if (change->data.truncate.restart_seqs
     738         [ -  + ]:           7 :         || change->data.truncate.cascade)
     739                 :             :     {
     740         [ +  - ]:           1 :         if (change->data.truncate.restart_seqs)
     741                 :           1 :             appendStringInfoString(ctx->out, " restart_seqs");
     742         [ +  - ]:           1 :         if (change->data.truncate.cascade)
     743                 :           1 :             appendStringInfoString(ctx->out, " cascade");
     744                 :             :     }
     745                 :             :     else
     746                 :           7 :         appendStringInfoString(ctx->out, " (no-flags)");
     747                 :             : 
     748                 :           8 :     MemoryContextSwitchTo(old);
     749                 :           8 :     MemoryContextReset(data->context);
     750                 :             : 
     751                 :           8 :     OutputPluginWrite(ctx, true);
     752                 :           8 : }
     753                 :             : 
     754                 :             : static void
     755                 :          10 : pg_decode_message(LogicalDecodingContext *ctx,
     756                 :             :                   ReorderBufferTXN *txn, XLogRecPtr lsn, bool transactional,
     757                 :             :                   const char *prefix, Size sz, const char *message)
     758                 :             : {
     759                 :          10 :     TestDecodingData *data = ctx->output_plugin_private;
     760                 :             :     TestDecodingTxnData *txndata;
     761                 :             : 
     762         [ +  + ]:          10 :     txndata = transactional ? txn->output_plugin_private : NULL;
     763                 :             : 
     764                 :             :     /* output BEGIN if we haven't yet for transactional messages */
     765   [ +  +  +  -  :          10 :     if (transactional && data->skip_empty_xacts && !txndata->xact_wrote_changes)
                   +  + ]
     766                 :           3 :         pg_output_begin(ctx, data, txn, false);
     767                 :             : 
     768         [ +  + ]:          10 :     if (transactional)
     769                 :           5 :         txndata->xact_wrote_changes = true;
     770                 :             : 
     771                 :          10 :     OutputPluginPrepareWrite(ctx, true);
     772                 :          10 :     appendStringInfo(ctx->out, "message: transactional: %d prefix: %s, sz: %zu content:",
     773                 :             :                      transactional, prefix, sz);
     774                 :          10 :     appendBinaryStringInfo(ctx->out, message, sz);
     775                 :          10 :     OutputPluginWrite(ctx, true);
     776                 :          10 : }
     777                 :             : 
     778                 :             : static void
     779                 :          40 : pg_decode_stream_start(LogicalDecodingContext *ctx,
     780                 :             :                        ReorderBufferTXN *txn)
     781                 :             : {
     782                 :          40 :     TestDecodingData *data = ctx->output_plugin_private;
     783                 :          40 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     784                 :             : 
     785                 :             :     /*
     786                 :             :      * Allocate the txn plugin data for the first stream in the transaction.
     787                 :             :      */
     788         [ +  + ]:          40 :     if (txndata == NULL)
     789                 :             :     {
     790                 :             :         txndata =
     791                 :           9 :             MemoryContextAllocZero(ctx->context, sizeof(TestDecodingTxnData));
     792                 :           9 :         txndata->xact_wrote_changes = false;
     793                 :           9 :         txn->output_plugin_private = txndata;
     794                 :             :     }
     795                 :             : 
     796                 :          40 :     txndata->stream_wrote_changes = false;
     797         [ +  - ]:          40 :     if (data->skip_empty_xacts)
     798                 :          40 :         return;
     799                 :           0 :     pg_output_stream_start(ctx, data, txn, true);
     800                 :             : }
     801                 :             : 
     802                 :             : static void
     803                 :          11 : pg_output_stream_start(LogicalDecodingContext *ctx, TestDecodingData *data, ReorderBufferTXN *txn, bool last_write)
     804                 :             : {
     805                 :          11 :     OutputPluginPrepareWrite(ctx, last_write);
     806         [ -  + ]:          11 :     if (data->include_xids)
     807                 :           0 :         appendStringInfo(ctx->out, "opening a streamed block for transaction TXN %u", txn->xid);
     808                 :             :     else
     809                 :          11 :         appendStringInfoString(ctx->out, "opening a streamed block for transaction");
     810                 :          11 :     OutputPluginWrite(ctx, last_write);
     811                 :          11 : }
     812                 :             : 
     813                 :             : static void
     814                 :          40 : pg_decode_stream_stop(LogicalDecodingContext *ctx,
     815                 :             :                       ReorderBufferTXN *txn)
     816                 :             : {
     817                 :          40 :     TestDecodingData *data = ctx->output_plugin_private;
     818                 :          40 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     819                 :             : 
     820   [ +  -  +  + ]:          40 :     if (data->skip_empty_xacts && !txndata->stream_wrote_changes)
     821                 :          29 :         return;
     822                 :             : 
     823                 :          11 :     OutputPluginPrepareWrite(ctx, true);
     824         [ -  + ]:          11 :     if (data->include_xids)
     825                 :           0 :         appendStringInfo(ctx->out, "closing a streamed block for transaction TXN %u", txn->xid);
     826                 :             :     else
     827                 :          11 :         appendStringInfoString(ctx->out, "closing a streamed block for transaction");
     828                 :          11 :     OutputPluginWrite(ctx, true);
     829                 :             : }
     830                 :             : 
     831                 :             : static void
     832                 :           4 : pg_decode_stream_abort(LogicalDecodingContext *ctx,
     833                 :             :                        ReorderBufferTXN *txn,
     834                 :             :                        XLogRecPtr abort_lsn)
     835                 :             : {
     836                 :           4 :     TestDecodingData *data = ctx->output_plugin_private;
     837                 :             : 
     838                 :             :     /*
     839                 :             :      * stream abort can be sent for an individual subtransaction but we
     840                 :             :      * maintain the output_plugin_private only under the toptxn so if this is
     841                 :             :      * not the toptxn then fetch the toptxn.
     842                 :             :      */
     843         [ +  - ]:           4 :     ReorderBufferTXN *toptxn = rbtxn_get_toptxn(txn);
     844                 :           4 :     TestDecodingTxnData *txndata = toptxn->output_plugin_private;
     845                 :           4 :     bool        xact_wrote_changes = txndata->xact_wrote_changes;
     846                 :             : 
     847         [ -  + ]:           4 :     if (rbtxn_is_toptxn(txn))
     848                 :             :     {
     849                 :             :         Assert(txn->output_plugin_private != NULL);
     850                 :           0 :         pfree(txndata);
     851                 :           0 :         txn->output_plugin_private = NULL;
     852                 :             :     }
     853                 :             : 
     854   [ +  -  -  + ]:           4 :     if (data->skip_empty_xacts && !xact_wrote_changes)
     855                 :           0 :         return;
     856                 :             : 
     857                 :           4 :     OutputPluginPrepareWrite(ctx, true);
     858         [ -  + ]:           4 :     if (data->include_xids)
     859                 :           0 :         appendStringInfo(ctx->out, "aborting streamed (sub)transaction TXN %u", txn->xid);
     860                 :             :     else
     861                 :           4 :         appendStringInfoString(ctx->out, "aborting streamed (sub)transaction");
     862                 :           4 :     OutputPluginWrite(ctx, true);
     863                 :             : }
     864                 :             : 
     865                 :             : static void
     866                 :           1 : pg_decode_stream_prepare(LogicalDecodingContext *ctx,
     867                 :             :                          ReorderBufferTXN *txn,
     868                 :             :                          XLogRecPtr prepare_lsn)
     869                 :             : {
     870                 :           1 :     TestDecodingData *data = ctx->output_plugin_private;
     871                 :           1 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     872                 :             : 
     873   [ +  -  -  + ]:           1 :     if (data->skip_empty_xacts && !txndata->xact_wrote_changes)
     874                 :           0 :         return;
     875                 :             : 
     876                 :           1 :     OutputPluginPrepareWrite(ctx, true);
     877                 :             : 
     878         [ -  + ]:           1 :     if (data->include_xids)
     879                 :           0 :         appendStringInfo(ctx->out, "preparing streamed transaction TXN %s, txid %u",
     880                 :           0 :                          quote_literal_cstr(txn->gid), txn->xid);
     881                 :             :     else
     882                 :           1 :         appendStringInfo(ctx->out, "preparing streamed transaction %s",
     883                 :           1 :                          quote_literal_cstr(txn->gid));
     884                 :             : 
     885         [ -  + ]:           1 :     if (data->include_timestamp)
     886                 :           0 :         appendStringInfo(ctx->out, " (at %s)",
     887                 :             :                          timestamptz_to_str(txn->prepare_time));
     888                 :             : 
     889                 :           1 :     OutputPluginWrite(ctx, true);
     890                 :             : }
     891                 :             : 
     892                 :             : static void
     893                 :           5 : pg_decode_stream_commit(LogicalDecodingContext *ctx,
     894                 :             :                         ReorderBufferTXN *txn,
     895                 :             :                         XLogRecPtr commit_lsn)
     896                 :             : {
     897                 :           5 :     TestDecodingData *data = ctx->output_plugin_private;
     898                 :           5 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     899                 :           5 :     bool        xact_wrote_changes = txndata->xact_wrote_changes;
     900                 :             : 
     901                 :           5 :     pfree(txndata);
     902                 :           5 :     txn->output_plugin_private = NULL;
     903                 :             : 
     904   [ +  -  -  + ]:           5 :     if (data->skip_empty_xacts && !xact_wrote_changes)
     905                 :           0 :         return;
     906                 :             : 
     907                 :           5 :     OutputPluginPrepareWrite(ctx, true);
     908                 :             : 
     909         [ -  + ]:           5 :     if (data->include_xids)
     910                 :           0 :         appendStringInfo(ctx->out, "committing streamed transaction TXN %u", txn->xid);
     911                 :             :     else
     912                 :           5 :         appendStringInfoString(ctx->out, "committing streamed transaction");
     913                 :             : 
     914         [ -  + ]:           5 :     if (data->include_timestamp)
     915                 :           0 :         appendStringInfo(ctx->out, " (at %s)",
     916                 :             :                          timestamptz_to_str(txn->commit_time));
     917                 :             : 
     918                 :           5 :     OutputPluginWrite(ctx, true);
     919                 :             : }
     920                 :             : 
     921                 :             : /*
     922                 :             :  * In streaming mode, we don't display the changes as the transaction can abort
     923                 :             :  * at a later point in time.  We don't want users to see the changes until the
     924                 :             :  * transaction is committed.
     925                 :             :  */
     926                 :             : static void
     927                 :          65 : pg_decode_stream_change(LogicalDecodingContext *ctx,
     928                 :             :                         ReorderBufferTXN *txn,
     929                 :             :                         Relation relation,
     930                 :             :                         ReorderBufferChange *change)
     931                 :             : {
     932                 :          65 :     TestDecodingData *data = ctx->output_plugin_private;
     933                 :          65 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
     934                 :             : 
     935                 :             :     /* output stream start if we haven't yet */
     936   [ +  -  +  + ]:          65 :     if (data->skip_empty_xacts && !txndata->stream_wrote_changes)
     937                 :             :     {
     938                 :           8 :         pg_output_stream_start(ctx, data, txn, false);
     939                 :             :     }
     940                 :          65 :     txndata->xact_wrote_changes = txndata->stream_wrote_changes = true;
     941                 :             : 
     942                 :          65 :     OutputPluginPrepareWrite(ctx, true);
     943         [ -  + ]:          65 :     if (data->include_xids)
     944                 :           0 :         appendStringInfo(ctx->out, "streaming change for TXN %u", txn->xid);
     945                 :             :     else
     946                 :          65 :         appendStringInfoString(ctx->out, "streaming change for transaction");
     947                 :          65 :     OutputPluginWrite(ctx, true);
     948                 :          65 : }
     949                 :             : 
     950                 :             : /*
     951                 :             :  * In streaming mode, we don't display the contents for transactional messages
     952                 :             :  * as the transaction can abort at a later point in time.  We don't want users to
     953                 :             :  * see the message contents until the transaction is committed.
     954                 :             :  */
     955                 :             : static void
     956                 :           3 : pg_decode_stream_message(LogicalDecodingContext *ctx,
     957                 :             :                          ReorderBufferTXN *txn, XLogRecPtr lsn, bool transactional,
     958                 :             :                          const char *prefix, Size sz, const char *message)
     959                 :             : {
     960                 :             :     /* Output stream start if we haven't yet for transactional messages. */
     961         [ +  - ]:           3 :     if (transactional)
     962                 :             :     {
     963                 :           3 :         TestDecodingData *data = ctx->output_plugin_private;
     964                 :           3 :         TestDecodingTxnData *txndata = txn->output_plugin_private;
     965                 :             : 
     966   [ +  -  +  - ]:           3 :         if (data->skip_empty_xacts && !txndata->stream_wrote_changes)
     967                 :             :         {
     968                 :           3 :             pg_output_stream_start(ctx, data, txn, false);
     969                 :             :         }
     970                 :           3 :         txndata->xact_wrote_changes = txndata->stream_wrote_changes = true;
     971                 :             :     }
     972                 :             : 
     973                 :           3 :     OutputPluginPrepareWrite(ctx, true);
     974                 :             : 
     975         [ +  - ]:           3 :     if (transactional)
     976                 :             :     {
     977                 :           3 :         appendStringInfo(ctx->out, "streaming message: transactional: %d prefix: %s, sz: %zu",
     978                 :             :                          transactional, prefix, sz);
     979                 :             :     }
     980                 :             :     else
     981                 :             :     {
     982                 :           0 :         appendStringInfo(ctx->out, "streaming message: transactional: %d prefix: %s, sz: %zu content:",
     983                 :             :                          transactional, prefix, sz);
     984                 :           0 :         appendBinaryStringInfo(ctx->out, message, sz);
     985                 :             :     }
     986                 :             : 
     987                 :           3 :     OutputPluginWrite(ctx, true);
     988                 :           3 : }
     989                 :             : 
     990                 :             : /*
     991                 :             :  * In streaming mode, we don't display the detailed information of Truncate.
     992                 :             :  * See pg_decode_stream_change.
     993                 :             :  */
     994                 :             : static void
     995                 :           0 : pg_decode_stream_truncate(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
     996                 :             :                           int nrelations, Relation relations[],
     997                 :             :                           ReorderBufferChange *change)
     998                 :             : {
     999                 :           0 :     TestDecodingData *data = ctx->output_plugin_private;
    1000                 :           0 :     TestDecodingTxnData *txndata = txn->output_plugin_private;
    1001                 :             : 
    1002   [ #  #  #  # ]:           0 :     if (data->skip_empty_xacts && !txndata->stream_wrote_changes)
    1003                 :             :     {
    1004                 :           0 :         pg_output_stream_start(ctx, data, txn, false);
    1005                 :             :     }
    1006                 :           0 :     txndata->xact_wrote_changes = txndata->stream_wrote_changes = true;
    1007                 :             : 
    1008                 :           0 :     OutputPluginPrepareWrite(ctx, true);
    1009         [ #  # ]:           0 :     if (data->include_xids)
    1010                 :           0 :         appendStringInfo(ctx->out, "streaming truncate for TXN %u", txn->xid);
    1011                 :             :     else
    1012                 :           0 :         appendStringInfoString(ctx->out, "streaming truncate for transaction");
    1013                 :           0 :     OutputPluginWrite(ctx, true);
    1014                 :           0 : }
        

Generated by: LCOV version 2.0-1