codekingpro/portable-devtools
115k
1/*-------------------------------------------------------------------------2 * logical.h3 * PostgreSQL logical decoding coordination4 *5 * Copyright (c) 2012-2023, PostgreSQL Global Development Group6 *7 *-------------------------------------------------------------------------8 */9#ifndef LOGICAL_H10#define LOGICAL_H11 12#include "access/xlog.h"13#include "access/xlogreader.h"14#include "replication/output_plugin.h"15#include "replication/slot.h"16 17struct LogicalDecodingContext;18 19typedef void (*LogicalOutputPluginWriterWrite) (struct LogicalDecodingContext *lr,20 XLogRecPtr Ptr,21 TransactionId xid,22 bool last_write23);24 25typedef LogicalOutputPluginWriterWrite LogicalOutputPluginWriterPrepareWrite;26 27typedef void (*LogicalOutputPluginWriterUpdateProgress) (struct LogicalDecodingContext *lr,28 XLogRecPtr Ptr,29 TransactionId xid,30 bool skipped_xact31);32 33typedef struct LogicalDecodingContext34{35 /* memory context this is all allocated in */36 MemoryContext context;37 38 /* The associated replication slot */39 ReplicationSlot *slot;40 41 /* infrastructure pieces for decoding */42 XLogReaderState *reader;43 struct ReorderBuffer *reorder;44 struct SnapBuild *snapshot_builder;45 46 /*47 * Marks the logical decoding context as fast forward decoding one. Such a48 * context does not have plugin loaded so most of the following properties49 * are unused.50 */51 bool fast_forward;52 53 OutputPluginCallbacks callbacks;54 OutputPluginOptions options;55 56 /*57 * User specified options58 */59 List *output_plugin_options;60 61 /*62 * User-Provided callback for writing/streaming out data.63 */64 LogicalOutputPluginWriterPrepareWrite prepare_write;65 LogicalOutputPluginWriterWrite write;66 LogicalOutputPluginWriterUpdateProgress update_progress;67 68 /*69 * Output buffer.70 */71 StringInfo out;72 73 /*74 * Private data pointer of the output plugin.75 */76 void *output_plugin_private;77 78 /*79 * Private data pointer for the data writer.80 */81 void *output_writer_private;82 83 /*84 * Does the output plugin support streaming, and is it enabled?85 */86 bool streaming;87 88 /*89 * Does the output plugin support two-phase decoding, and is it enabled?90 */91 bool twophase;92 93 /*94 * Is two-phase option given by output plugin?95 *96 * This flag indicates that the plugin passed in the two-phase option as97 * part of the START_STREAMING command. We can't rely solely on the98 * twophase flag which only tells whether the plugin provided all the99 * necessary two-phase callbacks.100 */101 bool twophase_opt_given;102 103 /*104 * State for writing output.105 */106 bool accept_writes;107 bool prepared_write;108 XLogRecPtr write_location;109 TransactionId write_xid;110 /* Are we processing the end LSN of a transaction? */111 bool end_xact;112} LogicalDecodingContext;113 114 115extern void CheckLogicalDecodingRequirements(void);116 117extern LogicalDecodingContext *CreateInitDecodingContext(const char *plugin,118 List *output_plugin_options,119 bool need_full_snapshot,120 XLogRecPtr restart_lsn,121 XLogReaderRoutine *xl_routine,122 LogicalOutputPluginWriterPrepareWrite prepare_write,123 LogicalOutputPluginWriterWrite do_write,124 LogicalOutputPluginWriterUpdateProgress update_progress);125extern LogicalDecodingContext *CreateDecodingContext(XLogRecPtr start_lsn,126 List *output_plugin_options,127 bool fast_forward,128 XLogReaderRoutine *xl_routine,129 LogicalOutputPluginWriterPrepareWrite prepare_write,130 LogicalOutputPluginWriterWrite do_write,131 LogicalOutputPluginWriterUpdateProgress update_progress);132extern void DecodingContextFindStartpoint(LogicalDecodingContext *ctx);133extern bool DecodingContextReady(LogicalDecodingContext *ctx);134extern void FreeDecodingContext(LogicalDecodingContext *ctx);135 136extern void LogicalIncreaseXminForSlot(XLogRecPtr current_lsn,137 TransactionId xmin);138extern void LogicalIncreaseRestartDecodingForSlot(XLogRecPtr current_lsn,139 XLogRecPtr restart_lsn);140extern void LogicalConfirmReceivedLocation(XLogRecPtr lsn);141 142extern bool filter_prepare_cb_wrapper(LogicalDecodingContext *ctx,143 TransactionId xid, const char *gid);144extern bool filter_by_origin_cb_wrapper(LogicalDecodingContext *ctx, RepOriginId origin_id);145extern void ResetLogicalStreamingState(void);146extern void UpdateDecodingStats(LogicalDecodingContext *ctx);147 148#endif149 