mirror of
https://github.com/postgres/postgres.git
synced 2025-06-25 01:02:05 +03:00
Previously, the memory used by the logical replication apply worker for processing messages would never be freed, so that could end up using a lot of memory. To improve that, change the existing ApplyContext memory context to ApplyMessageContext and reset that after every message (similar to MessageContext used elsewhere). For consistency of naming, rename the ApplyCacheContext to ApplyContext. Author: Stas Kelvich <s.kelvich@postgrespro.ru>
98 lines
2.7 KiB
C
98 lines
2.7 KiB
C
/*-------------------------------------------------------------------------
|
|
*
|
|
* worker_internal.h
|
|
* Internal headers shared by logical replication workers.
|
|
*
|
|
* Portions Copyright (c) 2016-2017, PostgreSQL Global Development Group
|
|
*
|
|
* src/include/replication/worker_internal.h
|
|
*
|
|
*-------------------------------------------------------------------------
|
|
*/
|
|
#ifndef WORKER_INTERNAL_H
|
|
#define WORKER_INTERNAL_H
|
|
|
|
#include <signal.h>
|
|
|
|
#include "access/xlogdefs.h"
|
|
#include "catalog/pg_subscription.h"
|
|
#include "datatype/timestamp.h"
|
|
#include "storage/lock.h"
|
|
|
|
typedef struct LogicalRepWorker
|
|
{
|
|
/* Time at which this worker was launched. */
|
|
TimestampTz launch_time;
|
|
|
|
/* Indicates if this slot is used or free. */
|
|
bool in_use;
|
|
|
|
/* Increased everytime the slot is taken by new worker. */
|
|
uint16 generation;
|
|
|
|
/* Pointer to proc array. NULL if not running. */
|
|
PGPROC *proc;
|
|
|
|
/* Database id to connect to. */
|
|
Oid dbid;
|
|
|
|
/* User to use for connection (will be same as owner of subscription). */
|
|
Oid userid;
|
|
|
|
/* Subscription id for the worker. */
|
|
Oid subid;
|
|
|
|
/* Used for initial table synchronization. */
|
|
Oid relid;
|
|
char relstate;
|
|
XLogRecPtr relstate_lsn;
|
|
slock_t relmutex;
|
|
|
|
/* Stats. */
|
|
XLogRecPtr last_lsn;
|
|
TimestampTz last_send_time;
|
|
TimestampTz last_recv_time;
|
|
XLogRecPtr reply_lsn;
|
|
TimestampTz reply_time;
|
|
} LogicalRepWorker;
|
|
|
|
/* Main memory context for apply worker. Permanent during worker lifetime. */
|
|
extern MemoryContext ApplyContext;
|
|
|
|
/* libpqreceiver connection */
|
|
extern struct WalReceiverConn *wrconn;
|
|
|
|
/* Worker and subscription objects. */
|
|
extern Subscription *MySubscription;
|
|
extern LogicalRepWorker *MyLogicalRepWorker;
|
|
|
|
extern bool in_remote_transaction;
|
|
extern volatile sig_atomic_t got_SIGHUP;
|
|
extern volatile sig_atomic_t got_SIGTERM;
|
|
|
|
extern void logicalrep_worker_attach(int slot);
|
|
extern LogicalRepWorker *logicalrep_worker_find(Oid subid, Oid relid,
|
|
bool only_running);
|
|
extern void logicalrep_worker_launch(Oid dbid, Oid subid, const char *subname,
|
|
Oid userid, Oid relid);
|
|
extern void logicalrep_worker_stop(Oid subid, Oid relid);
|
|
extern void logicalrep_worker_wakeup(Oid subid, Oid relid);
|
|
extern void logicalrep_worker_wakeup_ptr(LogicalRepWorker *worker);
|
|
|
|
extern int logicalrep_sync_worker_count(Oid subid);
|
|
|
|
extern void logicalrep_worker_sighup(SIGNAL_ARGS);
|
|
extern void logicalrep_worker_sigterm(SIGNAL_ARGS);
|
|
extern char *LogicalRepSyncTableStart(XLogRecPtr *origin_startpos);
|
|
void process_syncing_tables(XLogRecPtr current_lsn);
|
|
void invalidate_syncing_table_states(Datum arg, int cacheid,
|
|
uint32 hashvalue);
|
|
|
|
static inline bool
|
|
am_tablesync_worker(void)
|
|
{
|
|
return OidIsValid(MyLogicalRepWorker->relid);
|
|
}
|
|
|
|
#endif /* WORKER_INTERNAL_H */
|