Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
logicalproto.h275 linesDownload Raw Back to replication
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 
codekingpro/portable-devtools · Team Ai