codekingpro/portable-devtools
114k
1/*-------------------------------------------------------------------------2 *3 * logicalproto.h4 * logical replication protocol5 *6 * Copyright (c) 2015-2023, PostgreSQL Global Development Group7 *8 * IDENTIFICATION9 * src/include/replication/logicalproto.h10 *11 *-------------------------------------------------------------------------12 */13#ifndef LOGICAL_PROTO_H14#define LOGICAL_PROTO_H15 16#include "access/xact.h"17#include "executor/tuptable.h"18#include "replication/reorderbuffer.h"19#include "utils/rel.h"20 21/*22 * Protocol capabilities23 *24 * LOGICALREP_PROTO_VERSION_NUM is our native protocol.25 * LOGICALREP_PROTO_MAX_VERSION_NUM is the greatest version we can support.26 * LOGICALREP_PROTO_MIN_VERSION_NUM is the oldest version we27 * have backwards compatibility for. The client requests protocol version at28 * connect time.29 *30 * LOGICALREP_PROTO_STREAM_VERSION_NUM is the minimum protocol version with31 * support for streaming large transactions. Introduced in PG14.32 *33 * LOGICALREP_PROTO_TWOPHASE_VERSION_NUM is the minimum protocol version with34 * support for two-phase commit decoding (at prepare time). Introduced in PG15.35 *36 * LOGICALREP_PROTO_STREAM_PARALLEL_VERSION_NUM is the minimum protocol version37 * where we support applying large streaming transactions in parallel.38 * Introduced in PG16.39 */40#define LOGICALREP_PROTO_MIN_VERSION_NUM 141#define LOGICALREP_PROTO_VERSION_NUM 142#define LOGICALREP_PROTO_STREAM_VERSION_NUM 243#define LOGICALREP_PROTO_TWOPHASE_VERSION_NUM 344#define LOGICALREP_PROTO_STREAM_PARALLEL_VERSION_NUM 445#define LOGICALREP_PROTO_MAX_VERSION_NUM LOGICALREP_PROTO_STREAM_PARALLEL_VERSION_NUM46 47/*48 * Logical message types49 *50 * Used by logical replication wire protocol.51 *52 * Note: though this is an enum, the values are used to identify message types53 * in logical replication protocol, which uses a single byte to identify a54 * message type. Hence the values should be single-byte wide and preferably55 * human-readable characters.56 */57typedef enum LogicalRepMsgType58{59 LOGICAL_REP_MSG_BEGIN = 'B',60 LOGICAL_REP_MSG_COMMIT = 'C',61 LOGICAL_REP_MSG_ORIGIN = 'O',62 LOGICAL_REP_MSG_INSERT = 'I',63 LOGICAL_REP_MSG_UPDATE = 'U',64 LOGICAL_REP_MSG_DELETE = 'D',65 LOGICAL_REP_MSG_TRUNCATE = 'T',66 LOGICAL_REP_MSG_RELATION = 'R',67 LOGICAL_REP_MSG_TYPE = 'Y',68 LOGICAL_REP_MSG_MESSAGE = 'M',69 LOGICAL_REP_MSG_BEGIN_PREPARE = 'b',70 LOGICAL_REP_MSG_PREPARE = 'P',71 LOGICAL_REP_MSG_COMMIT_PREPARED = 'K',72 LOGICAL_REP_MSG_ROLLBACK_PREPARED = 'r',73 LOGICAL_REP_MSG_STREAM_START = 'S',74 LOGICAL_REP_MSG_STREAM_STOP = 'E',75 LOGICAL_REP_MSG_STREAM_COMMIT = 'c',76 LOGICAL_REP_MSG_STREAM_ABORT = 'A',77 LOGICAL_REP_MSG_STREAM_PREPARE = 'p'78} LogicalRepMsgType;79 80/*81 * This struct stores a tuple received via logical replication.82 * Keep in mind that the columns correspond to the *remote* table.83 */84typedef struct LogicalRepTupleData85{86 /* Array of StringInfos, one per column; some may be unused */87 StringInfoData *colvalues;88 /* Array of markers for null/unchanged/text/binary, one per column */89 char *colstatus;90 /* Length of above arrays */91 int ncols;92} LogicalRepTupleData;93 94/* Possible values for LogicalRepTupleData.colstatus[colnum] */95/* These values are also used in the on-the-wire protocol */96#define LOGICALREP_COLUMN_NULL 'n'97#define LOGICALREP_COLUMN_UNCHANGED 'u'98#define LOGICALREP_COLUMN_TEXT 't'99#define LOGICALREP_COLUMN_BINARY 'b' /* added in PG14 */100 101typedef uint32 LogicalRepRelId;102 103/* Relation information */104typedef struct LogicalRepRelation105{106 /* Info coming from the remote side. */107 LogicalRepRelId remoteid; /* unique id of the relation */108 char *nspname; /* schema name */109 char *relname; /* relation name */110 int natts; /* number of columns */111 char **attnames; /* column names */112 Oid *atttyps; /* column types */113 char replident; /* replica identity */114 char relkind; /* remote relation kind */115 Bitmapset *attkeys; /* Bitmap of key columns */116} LogicalRepRelation;117 118/* Type mapping info */119typedef struct LogicalRepTyp120{121 Oid remoteid; /* unique id of the remote type */122 char *nspname; /* schema name of remote type */123 char *typname; /* name of the remote type */124} LogicalRepTyp;125 126/* Transaction info */127typedef struct LogicalRepBeginData128{129 XLogRecPtr final_lsn;130 TimestampTz committime;131 TransactionId xid;132} LogicalRepBeginData;133 134typedef struct LogicalRepCommitData135{136 XLogRecPtr commit_lsn;137 XLogRecPtr end_lsn;138 TimestampTz committime;139} LogicalRepCommitData;140 141/*142 * Prepared transaction protocol information for begin_prepare, and prepare.143 */144typedef struct LogicalRepPreparedTxnData145{146 XLogRecPtr prepare_lsn;147 XLogRecPtr end_lsn;148 TimestampTz prepare_time;149 TransactionId xid;150 char gid[GIDSIZE];151} LogicalRepPreparedTxnData;152 153/*154 * Prepared transaction protocol information for commit prepared.155 */156typedef struct LogicalRepCommitPreparedTxnData157{158 XLogRecPtr commit_lsn;159 XLogRecPtr end_lsn;160 TimestampTz commit_time;161 TransactionId xid;162 char gid[GIDSIZE];163} LogicalRepCommitPreparedTxnData;164 165/*166 * Rollback Prepared transaction protocol information. The prepare information167 * prepare_end_lsn and prepare_time are used to check if the downstream has168 * received this prepared transaction in which case it can apply the rollback,169 * otherwise, it can skip the rollback operation. The gid alone is not170 * sufficient because the downstream node can have a prepared transaction with171 * same identifier.172 */173typedef struct LogicalRepRollbackPreparedTxnData174{175 XLogRecPtr prepare_end_lsn;176 XLogRecPtr rollback_end_lsn;177 TimestampTz prepare_time;178 TimestampTz rollback_time;179 TransactionId xid;180 char gid[GIDSIZE];181} LogicalRepRollbackPreparedTxnData;182 183/*184 * Transaction protocol information for stream abort.185 */186typedef struct LogicalRepStreamAbortData187{188 TransactionId xid;189 TransactionId subxid;190 XLogRecPtr abort_lsn;191 TimestampTz abort_time;192} LogicalRepStreamAbortData;193 194extern void logicalrep_write_begin(StringInfo out, ReorderBufferTXN *txn);195extern void logicalrep_read_begin(StringInfo in,196 LogicalRepBeginData *begin_data);197extern void logicalrep_write_commit(StringInfo out, ReorderBufferTXN *txn,198 XLogRecPtr commit_lsn);199extern void logicalrep_read_commit(StringInfo in,200 LogicalRepCommitData *commit_data);201extern void logicalrep_write_begin_prepare(StringInfo out, ReorderBufferTXN *txn);202extern void logicalrep_read_begin_prepare(StringInfo in,203 LogicalRepPreparedTxnData *begin_data);204extern void logicalrep_write_prepare(StringInfo out, ReorderBufferTXN *txn,205 XLogRecPtr prepare_lsn);206extern void logicalrep_read_prepare(StringInfo in,207 LogicalRepPreparedTxnData *prepare_data);208extern void logicalrep_write_commit_prepared(StringInfo out, ReorderBufferTXN *txn,209 XLogRecPtr commit_lsn);210extern void logicalrep_read_commit_prepared(StringInfo in,211 LogicalRepCommitPreparedTxnData *prepare_data);212extern void logicalrep_write_rollback_prepared(StringInfo out, ReorderBufferTXN *txn,213 XLogRecPtr prepare_end_lsn,214 TimestampTz prepare_time);215extern void logicalrep_read_rollback_prepared(StringInfo in,216 LogicalRepRollbackPreparedTxnData *rollback_data);217extern void logicalrep_write_stream_prepare(StringInfo out, ReorderBufferTXN *txn,218 XLogRecPtr prepare_lsn);219extern void logicalrep_read_stream_prepare(StringInfo in,220 LogicalRepPreparedTxnData *prepare_data);221 222extern void logicalrep_write_origin(StringInfo out, const char *origin,223 XLogRecPtr origin_lsn);224extern char *logicalrep_read_origin(StringInfo in, XLogRecPtr *origin_lsn);225extern void logicalrep_write_insert(StringInfo out, TransactionId xid,226 Relation rel,227 TupleTableSlot *newslot,228 bool binary, Bitmapset *columns);229extern LogicalRepRelId logicalrep_read_insert(StringInfo in, LogicalRepTupleData *newtup);230extern void logicalrep_write_update(StringInfo out, TransactionId xid,231 Relation rel,232 TupleTableSlot *oldslot,233 TupleTableSlot *newslot, bool binary, Bitmapset *columns);234extern LogicalRepRelId logicalrep_read_update(StringInfo in,235 bool *has_oldtuple, LogicalRepTupleData *oldtup,236 LogicalRepTupleData *newtup);237extern void logicalrep_write_delete(StringInfo out, TransactionId xid,238 Relation rel, TupleTableSlot *oldslot,239 bool binary, Bitmapset *columns);240extern LogicalRepRelId logicalrep_read_delete(StringInfo in,241 LogicalRepTupleData *oldtup);242extern void logicalrep_write_truncate(StringInfo out, TransactionId xid,243 int nrelids, Oid relids[],244 bool cascade, bool restart_seqs);245extern List *logicalrep_read_truncate(StringInfo in,246 bool *cascade, bool *restart_seqs);247extern void logicalrep_write_message(StringInfo out, TransactionId xid, XLogRecPtr lsn,248 bool transactional, const char *prefix, Size sz, const char *message);249extern void logicalrep_write_rel(StringInfo out, TransactionId xid,250 Relation rel, Bitmapset *columns);251extern LogicalRepRelation *logicalrep_read_rel(StringInfo in);252extern void logicalrep_write_typ(StringInfo out, TransactionId xid,253 Oid typoid);254extern void logicalrep_read_typ(StringInfo in, LogicalRepTyp *ltyp);255extern void logicalrep_write_stream_start(StringInfo out, TransactionId xid,256 bool first_segment);257extern TransactionId logicalrep_read_stream_start(StringInfo in,258 bool *first_segment);259extern void logicalrep_write_stream_stop(StringInfo out);260extern void logicalrep_write_stream_commit(StringInfo out, ReorderBufferTXN *txn,261 XLogRecPtr commit_lsn);262extern TransactionId logicalrep_read_stream_commit(StringInfo in,263 LogicalRepCommitData *commit_data);264extern void logicalrep_write_stream_abort(StringInfo out, TransactionId xid,265 TransactionId subxid,266 XLogRecPtr abort_lsn,267 TimestampTz abort_time,268 bool write_abort_info);269extern void logicalrep_read_stream_abort(StringInfo in,270 LogicalRepStreamAbortData *abort_data,271 bool read_abort_info);272extern const char *logicalrep_message_type(LogicalRepMsgType action);273 274#endif /* LOGICAL_PROTO_H */275 