forked from huawei/openGauss-server
Compare commits
17 Commits
| Author | SHA1 | Date |
|---|---|---|
|
|
a548f5c3c6 | |
|
|
eb4e54d4eb | |
|
|
3f7a909887 | |
|
|
d8b89ceea5 | |
|
|
2250adfd4b | |
|
|
fd34b5da2a | |
|
|
7f97d633f1 | |
|
|
6ba4c95f5f | |
|
|
5755feee5d | |
|
|
cd5d44d66c | |
|
|
2d657cddbe | |
|
|
ccd674c159 | |
|
|
8292238381 | |
|
|
5054ddd002 | |
|
|
48a5033c18 | |
|
|
811a9afcd9 | |
|
|
bb0aa02eb0 |
|
|
@ -5923,7 +5923,7 @@ int main(int argc, char** argv)
|
|||
&option_index)) != -1)
|
||||
#endif
|
||||
#else
|
||||
while ((c = getopt_long(argc, argv, "b:cD:e:fi:G:l:m:M:N:o:O:p:P:r:R:v:x:sS:t:u:U:wWZ:dqL:T:Q:", long_options,
|
||||
while ((c = getopt_long(argc, argv, "b:cD:e:fi:G:l:m:M:N:o:O:p:P:r:R:v:x:sS:t:u:U:wWZ:C:dqL:T:Q:", long_options,
|
||||
&option_index)) != -1)
|
||||
#endif
|
||||
#endif
|
||||
|
|
|
|||
|
|
@ -233,6 +233,7 @@ char* all_data_nodename_list = NULL;
|
|||
const uint32 USTORE_UPGRADE_VERSION = 92368;
|
||||
const uint32 PACKAGE_ENHANCEMENT = 92444;
|
||||
const uint32 SUBSCRIPTION_VERSION = 92580;
|
||||
const uint32 SUBSCRIPTION_BINARY_VERSION_NUM = 92606;
|
||||
|
||||
#ifdef DUMPSYSLOG
|
||||
char* syslogpath = NULL;
|
||||
|
|
@ -4444,7 +4445,9 @@ void getSubscriptions(Archive *fout)
|
|||
int i_subslotname;
|
||||
int i_subsynccommit;
|
||||
int i_subpublications;
|
||||
int i, ntups;
|
||||
int i_subbinary;
|
||||
int i;
|
||||
int ntups;
|
||||
|
||||
if (no_subscriptions || GetVersionNum(fout) < SUBSCRIPTION_VERSION) {
|
||||
return;
|
||||
|
|
@ -4469,14 +4472,20 @@ void getSubscriptions(Archive *fout)
|
|||
resetPQExpBuffer(query);
|
||||
|
||||
/* Get the subscriptions in current database. */
|
||||
appendPQExpBuffer(query,
|
||||
"SELECT s.tableoid, s.oid, s.subname,"
|
||||
"(%s s.subowner) AS rolname, "
|
||||
" s.subconninfo, s.subslotname, s.subsynccommit, s.subpublications "
|
||||
"FROM pg_catalog.pg_subscription s "
|
||||
appendPQExpBuffer(query, "SELECT s.tableoid, s.oid, s.subname,"
|
||||
"(%s s.subowner) AS rolname, s.subconninfo, s.subslotname, "
|
||||
"s.subsynccommit, s.subpublications, \n", username_subquery);
|
||||
|
||||
if (GetVersionNum(fout) >= SUBSCRIPTION_BINARY_VERSION_NUM) {
|
||||
appendPQExpBuffer(query, " s.subbinary\n");
|
||||
} else {
|
||||
appendPQExpBuffer(query, " false AS subbinary\n");
|
||||
}
|
||||
|
||||
appendPQExpBuffer(query, "FROM pg_catalog.pg_subscription s "
|
||||
"WHERE s.subdbid = (SELECT oid FROM pg_catalog.pg_database"
|
||||
" WHERE datname = current_database())",
|
||||
username_subquery);
|
||||
" WHERE datname = current_database())");
|
||||
|
||||
res = ExecuteSqlQuery(fout, query->data, PGRES_TUPLES_OK);
|
||||
|
||||
ntups = PQntuples(res);
|
||||
|
|
@ -4494,6 +4503,7 @@ void getSubscriptions(Archive *fout)
|
|||
i_subslotname = PQfnumber(res, "subslotname");
|
||||
i_subsynccommit = PQfnumber(res, "subsynccommit");
|
||||
i_subpublications = PQfnumber(res, "subpublications");
|
||||
i_subbinary = PQfnumber(res, "subbinary");
|
||||
|
||||
subinfo = (SubscriptionInfo *)pg_malloc(ntups * sizeof(SubscriptionInfo));
|
||||
|
||||
|
|
@ -4512,6 +4522,7 @@ void getSubscriptions(Archive *fout)
|
|||
}
|
||||
subinfo[i].subsynccommit = gs_strdup(PQgetvalue(res, i, i_subsynccommit));
|
||||
subinfo[i].subpublications = gs_strdup(PQgetvalue(res, i, i_subpublications));
|
||||
subinfo[i].subbinary = gs_strdup(PQgetvalue(res, i, i_subbinary));
|
||||
|
||||
if (strlen(subinfo[i].rolname) == 0) {
|
||||
write_msg(NULL, "WARNING: owner of subscription \"%s\" appears to be invalid\n", subinfo[i].dobj.name);
|
||||
|
|
@ -4578,6 +4589,10 @@ static void dumpSubscription(Archive *fout, const SubscriptionInfo *subinfo)
|
|||
appendPQExpBufferStr(query, "NONE");
|
||||
}
|
||||
|
||||
if (strcmp(subinfo->subbinary, "t") == 0) {
|
||||
appendPQExpBuffer(query, ", binary = true");
|
||||
}
|
||||
|
||||
if (strcmp(subinfo->subsynccommit, "off") != 0) {
|
||||
appendPQExpBuffer(query, ", synchronous_commit = %s", fmtId(subinfo->subsynccommit));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -498,6 +498,7 @@ typedef struct _SubscriptionInfo {
|
|||
char *subslotname;
|
||||
char *subsynccommit;
|
||||
char *subpublications;
|
||||
char *subbinary;
|
||||
} SubscriptionInfo;
|
||||
|
||||
/* global decls */
|
||||
|
|
|
|||
|
|
@ -31,6 +31,9 @@
|
|||
it will be backuped up in external dirs */
|
||||
parray *pgdata_nobackup_dir = NULL;
|
||||
|
||||
/* list of logical replication slots */
|
||||
parray *logical_replslot = NULL;
|
||||
|
||||
static int standby_message_timeout_local = 10 ; /* 10 sec = default */
|
||||
static XLogRecPtr stop_backup_lsn = InvalidXLogRecPtr;
|
||||
static XLogRecPtr stop_stream_lsn = InvalidXLogRecPtr;
|
||||
|
|
@ -89,10 +92,11 @@ static void backup_cleanup(bool fatal, void *userdata);
|
|||
|
||||
static void *backup_files(void *arg);
|
||||
|
||||
static void do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool backup_logs);
|
||||
static void do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool backup_logs,
|
||||
bool backup_replslots);
|
||||
|
||||
static void pg_start_backup(const char *label, bool smooth, pgBackup *backup,
|
||||
PGNodeInfo *nodeInfo, PGconn *conn);
|
||||
PGNodeInfo *nodeInfo, PGconn *conn, bool backup_replslots);
|
||||
static void pg_stop_backup(pgBackup *backup, PGconn *pg_startbackup_conn, PGNodeInfo *nodeInfo);
|
||||
static int checkpoint_timeout(PGconn *backup_conn);
|
||||
|
||||
|
|
@ -558,7 +562,7 @@ static void sync_files(parray *database_map, const char *database_path, parray *
|
|||
* Move files from 'pgdata' to a subdirectory in 'backup_path'.
|
||||
*/
|
||||
static void
|
||||
do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool backup_logs)
|
||||
do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool backup_logs, bool backup_replslots)
|
||||
{
|
||||
int i;
|
||||
char database_path[MAXPGPATH];
|
||||
|
|
@ -591,7 +595,7 @@ do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool
|
|||
securec_check_c(rc, "\0", "\0");
|
||||
|
||||
/* Call pg_start_backup function in openGauss connect */
|
||||
pg_start_backup(label, smooth_checkpoint, ¤t, nodeInfo, backup_conn);
|
||||
pg_start_backup(label, smooth_checkpoint, ¤t, nodeInfo, backup_conn, backup_replslots);
|
||||
|
||||
/* Obtain current timeline */
|
||||
#if PG_VERSION_NUM >= 90600
|
||||
|
|
@ -624,10 +628,10 @@ do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool
|
|||
/* list files with the logical path. omit $PGDATA */
|
||||
if (fio_is_remote(FIO_DB_HOST))
|
||||
fio_list_dir(backup_files_list, instance_config.pgdata,
|
||||
true, true, false, backup_logs, true, 0);
|
||||
true, true, false, backup_logs, true, 0, backup_replslots);
|
||||
else
|
||||
dir_list_file(backup_files_list, instance_config.pgdata,
|
||||
true, true, false, backup_logs, true, 0, FIO_LOCAL_HOST);
|
||||
true, true, false, backup_logs, true, 0, FIO_LOCAL_HOST, backup_replslots);
|
||||
|
||||
/*
|
||||
* Get database_map (name to oid) for use in partial restore feature.
|
||||
|
|
@ -749,6 +753,11 @@ do_backup_instance(PGconn *backup_conn, PGNodeInfo *nodeInfo, bool no_sync, bool
|
|||
}
|
||||
pgdata_nobackup_dir = NULL;
|
||||
|
||||
if (logical_replslot) {
|
||||
free_dir_list(logical_replslot);
|
||||
}
|
||||
logical_replslot = NULL;
|
||||
|
||||
/* Cleanup */
|
||||
if (backup_list)
|
||||
{
|
||||
|
|
@ -849,7 +858,7 @@ static void do_after_backup()
|
|||
*/
|
||||
int
|
||||
do_backup(time_t start_time, pgSetBackupParams *set_backup_params,
|
||||
bool no_validate, bool no_sync, bool backup_logs)
|
||||
bool no_validate, bool no_sync, bool backup_logs, bool backup_replslots)
|
||||
{
|
||||
PGconn *backup_conn = NULL;
|
||||
PGNodeInfo nodeInfo;
|
||||
|
|
@ -925,7 +934,7 @@ do_backup(time_t start_time, pgSetBackupParams *set_backup_params,
|
|||
add_note(¤t, set_backup_params->note);
|
||||
|
||||
/* backup data */
|
||||
do_backup_instance(backup_conn, &nodeInfo, no_sync, backup_logs);
|
||||
do_backup_instance(backup_conn, &nodeInfo, no_sync, backup_logs, backup_replslots);
|
||||
pgut_atexit_pop(backup_cleanup, NULL);
|
||||
|
||||
/* compute size of wal files of this backup stored in the archive */
|
||||
|
|
@ -1034,13 +1043,15 @@ confirm_block_size(PGconn *conn, const char *name, int blcksz)
|
|||
*/
|
||||
static void
|
||||
pg_start_backup(const char *label, bool smooth, pgBackup *backup,
|
||||
PGNodeInfo *nodeInfo, PGconn *conn)
|
||||
PGNodeInfo *nodeInfo, PGconn *conn, bool backup_replslots)
|
||||
{
|
||||
PGresult *res;
|
||||
const char *params[2];
|
||||
uint32 lsn_hi;
|
||||
uint32 lsn_lo;
|
||||
int ret;
|
||||
int i;
|
||||
XLogRecPtr startLsn;
|
||||
|
||||
params[0] = label;
|
||||
|
||||
|
|
@ -1068,7 +1079,33 @@ pg_start_backup(const char *label, bool smooth, pgBackup *backup,
|
|||
XLogDataFromLSN(ret, PQgetvalue(res, 0, 0), &lsn_hi, &lsn_lo);
|
||||
securec_check_for_sscanf_s(ret, 2, "\0", "\0");
|
||||
/* Calculate LSN */
|
||||
backup->start_lsn = ((uint64) lsn_hi )<< 32 | lsn_lo;
|
||||
startLsn = ((uint64) lsn_hi )<< 32 | lsn_lo;
|
||||
|
||||
if (backup_replslots) {
|
||||
logical_replslot = parray_new();
|
||||
/* query for logical replication slots of subscriptions */
|
||||
res = pgut_execute(conn,
|
||||
"SELECT slot_name, restart_lsn FROM pg_catalog.pg_get_replication_slots()"
|
||||
"WHERE slot_type = 'logical' AND plugin = 'pgoutput'", 0, NULL);
|
||||
if (PQntuples(res) == 0) {
|
||||
elog(LOG, "logical replication slots for subscriptions not found");
|
||||
} else {
|
||||
XLogRecPtr repslotLsn;
|
||||
|
||||
for (i = 0; i < PQntuples(res); i++) {
|
||||
XLogDataFromLSN(ret, PQgetvalue(res, i, 1), &lsn_hi, &lsn_lo);
|
||||
securec_check_for_sscanf_s(ret, 2, "\0", "\0");
|
||||
repslotLsn = ((uint64) lsn_hi )<< 32 | lsn_lo;
|
||||
startLsn = Min(startLsn, repslotLsn);
|
||||
|
||||
char* slotname = pg_strdup(PQgetvalue(res, i, 0));
|
||||
parray_append(logical_replslot, slotname);
|
||||
}
|
||||
elog(WARNING, "logical replication slots for subscriptions will be backed up. "
|
||||
"If don't use them after restoring, please drop them to avoid affecting xlog recycling.");
|
||||
}
|
||||
}
|
||||
backup->start_lsn = startLsn;
|
||||
|
||||
PQclear(res);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -42,13 +42,6 @@ const char *pgdata_exclude_dir[] =
|
|||
(const char *)"pg_stat_tmp",
|
||||
(const char *)"pgsql_tmp",
|
||||
|
||||
/*
|
||||
* It is generally not useful to backup the contents of this directory even
|
||||
* if the intention is to restore to another master. See backup.sgml for a
|
||||
* more detailed description.
|
||||
*/
|
||||
(const char *)"pg_replslot",
|
||||
|
||||
/* Contents removed on startup, see dsm_cleanup_for_mmap(). */
|
||||
(const char *)"pg_dynshmem",
|
||||
|
||||
|
|
@ -68,7 +61,7 @@ const char *pgdata_exclude_dir[] =
|
|||
(const char *)"pg_subtrans",
|
||||
|
||||
/* end of list */
|
||||
NULL, /* pg_log will be set later */
|
||||
NULL, /* pg_log and pg_replslot will be set later */
|
||||
NULL
|
||||
};
|
||||
|
||||
|
|
@ -128,17 +121,20 @@ may be removed int the future */
|
|||
|
||||
static int pgCompareString(const void *str1, const void *str2);
|
||||
|
||||
static char dir_check_file(pgFile *file, bool backup_logs);
|
||||
static char dir_check_file(pgFile *file, bool backup_logs, bool backup_replslots);
|
||||
static char check_in_tablespace(pgFile *file, bool in_tablespace);
|
||||
static char check_db_dir(pgFile *file);
|
||||
static char check_digit_file(pgFile *file);
|
||||
static char check_nobackup_dir(pgFile *file);
|
||||
static void dir_list_file_internal(parray *files, pgFile *parent, const char *parent_dir,
|
||||
bool exclude, bool follow_symlink, bool backup_logs,
|
||||
bool skip_hidden, int external_dir_num, fio_location location);
|
||||
bool skip_hidden, int external_dir_num, fio_location location,
|
||||
bool backup_replslots);
|
||||
static void opt_path_map(ConfigOption *opt, const char *arg,
|
||||
TablespaceList *list, const char *type);
|
||||
|
||||
char check_logical_replslot_dir(const char *rel_path);
|
||||
|
||||
/* Tablespace mapping */
|
||||
static TablespaceList tablespace_dirs = {NULL, NULL};
|
||||
/* Extra directories mapping */
|
||||
|
|
@ -538,7 +534,7 @@ db_map_entry_free(void *entry)
|
|||
void
|
||||
dir_list_file(parray *files, const char *root, bool exclude, bool follow_symlink,
|
||||
bool add_root, bool backup_logs, bool skip_hidden, int external_dir_num,
|
||||
fio_location location)
|
||||
fio_location location, bool backup_replslots)
|
||||
{
|
||||
pgFile *file;
|
||||
|
||||
|
|
@ -565,7 +561,7 @@ dir_list_file(parray *files, const char *root, bool exclude, bool follow_symlink
|
|||
parray_append(files, file);
|
||||
|
||||
dir_list_file_internal(files, file, root, exclude, follow_symlink,
|
||||
backup_logs, skip_hidden, external_dir_num, location);
|
||||
backup_logs, skip_hidden, external_dir_num, location, backup_replslots);
|
||||
|
||||
if (!add_root)
|
||||
pgFileFree(file);
|
||||
|
|
@ -589,7 +585,7 @@ dir_list_file(parray *files, const char *root, bool exclude, bool follow_symlink
|
|||
* - datafiles
|
||||
*/
|
||||
static char
|
||||
dir_check_file(pgFile *file, bool backup_logs)
|
||||
dir_check_file(pgFile *file, bool backup_logs, bool backup_replslots)
|
||||
{
|
||||
int i;
|
||||
int sscanf_res;
|
||||
|
|
@ -652,6 +648,29 @@ dir_check_file(pgFile *file, bool backup_logs)
|
|||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* Backup pg_replslot if it is specified.
|
||||
* It is generally not useful to backup the contents of this directory even
|
||||
* if the intention is to restore to another master. See backup.sgml for a
|
||||
* more detailed description.
|
||||
*/
|
||||
if (!backup_replslots) {
|
||||
if (strcmp(file->rel_path, PG_REPLSLOT_DIR) == 0) {
|
||||
/* Skip */
|
||||
elog(VERBOSE, "Excluding directory content: %s", file->rel_path);
|
||||
return CHECK_EXCLUDE_FALSE;
|
||||
}
|
||||
} else {
|
||||
/*
|
||||
* Check file that under pg_replslot and judge whether it
|
||||
* belonged to logical replication slots for subscriptions.
|
||||
*/
|
||||
if (strcmp(file->rel_path, PG_REPLSLOT_DIR) != 0 &&
|
||||
path_is_prefix_of_path(PG_REPLSLOT_DIR, file->rel_path)) {
|
||||
return check_logical_replslot_dir(file->rel_path);
|
||||
}
|
||||
}
|
||||
|
||||
ret = check_nobackup_dir(file);
|
||||
if (ret != -1) { /* -1 means need backup */
|
||||
return ret;
|
||||
|
|
@ -749,6 +768,35 @@ static char check_nobackup_dir(pgFile *file)
|
|||
return ret;
|
||||
}
|
||||
|
||||
char check_logical_replslot_dir(const char *rel_path)
|
||||
{
|
||||
char ret = CHECK_FALSE;
|
||||
int i = 0;
|
||||
char *tmp = pg_strdup(rel_path);
|
||||
char *p;
|
||||
#define DIRECTORY_DELIMITER "/"
|
||||
|
||||
if (logical_replslot) {
|
||||
/* extract slot name from rel_path, such as sub1 from pg_replslot/sub1/snap */
|
||||
p = strtok(tmp, DIRECTORY_DELIMITER);
|
||||
if (p != NULL) {
|
||||
p = strtok(NULL, DIRECTORY_DELIMITER);
|
||||
}
|
||||
|
||||
for (i = 0; p != NULL && i < (int)parray_num(logical_replslot); i++) {
|
||||
char *slotName = (char *)parray_get(logical_replslot, i);
|
||||
if (strcmp(p, slotName) == 0) {
|
||||
pfree(tmp);
|
||||
return CHECK_TRUE;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
ret = CHECK_TRUE;
|
||||
}
|
||||
pfree(tmp);
|
||||
return ret;
|
||||
}
|
||||
|
||||
static char check_db_dir(pgFile *file)
|
||||
{
|
||||
char ret = -1;
|
||||
|
|
@ -889,7 +937,8 @@ bool SkipSomeDirFile(pgFile *file, struct dirent *dent, bool skipHidden)
|
|||
static void
|
||||
dir_list_file_internal(parray *files, pgFile *parent, const char *parent_dir,
|
||||
bool exclude, bool follow_symlink, bool backup_logs,
|
||||
bool skip_hidden, int external_dir_num, fio_location location)
|
||||
bool skip_hidden, int external_dir_num, fio_location location,
|
||||
bool backup_replslots)
|
||||
{
|
||||
DIR *dir;
|
||||
struct dirent *dent;
|
||||
|
|
@ -937,7 +986,7 @@ dir_list_file_internal(parray *files, pgFile *parent, const char *parent_dir,
|
|||
|
||||
if (exclude)
|
||||
{
|
||||
check_res = dir_check_file(file, backup_logs);
|
||||
check_res = dir_check_file(file, backup_logs, backup_replslots);
|
||||
if (check_res == CHECK_FALSE)
|
||||
{
|
||||
/* Skip */
|
||||
|
|
@ -963,7 +1012,7 @@ dir_list_file_internal(parray *files, pgFile *parent, const char *parent_dir,
|
|||
*/
|
||||
if (S_ISDIR(file->mode))
|
||||
dir_list_file_internal(files, file, child, exclude, follow_symlink,
|
||||
backup_logs, skip_hidden, external_dir_num, location);
|
||||
backup_logs, skip_hidden, external_dir_num, location, backup_replslots);
|
||||
}
|
||||
|
||||
if (errno && errno != ENOENT)
|
||||
|
|
|
|||
|
|
@ -51,6 +51,7 @@ typedef struct
|
|||
bool exclusive_backup;
|
||||
bool skip_hidden;
|
||||
int external_dir_num;
|
||||
bool backup_replslots;
|
||||
} fio_list_dir_request;
|
||||
|
||||
typedef struct
|
||||
|
|
@ -1794,7 +1795,7 @@ cleanup:
|
|||
/* Compile the array of files located on remote machine in directory root */
|
||||
void fio_list_dir(parray *files, const char *root, bool exclude,
|
||||
bool follow_symlink, bool add_root, bool backup_logs,
|
||||
bool skip_hidden, int external_dir_num)
|
||||
bool skip_hidden, int external_dir_num, bool backup_replslots)
|
||||
{
|
||||
fio_header hdr;
|
||||
fio_list_dir_request req;
|
||||
|
|
@ -1811,6 +1812,7 @@ void fio_list_dir(parray *files, const char *root, bool exclude,
|
|||
req.exclusive_backup = exclusive_backup;
|
||||
req.skip_hidden = skip_hidden;
|
||||
req.external_dir_num = external_dir_num;
|
||||
req.backup_replslots = backup_replslots;
|
||||
|
||||
hdr.cop = FIO_LIST_DIR;
|
||||
hdr.size = sizeof(req);
|
||||
|
|
@ -1870,7 +1872,14 @@ void fio_list_dir(parray *files, const char *root, bool exclude,
|
|||
securec_check_ss_c(nRet, "\0", "\0");
|
||||
}
|
||||
|
||||
|
||||
/*
|
||||
* Check file that under pg_replslot and judge whether it
|
||||
* belonged to logical replication slots for subscriptions.
|
||||
*/
|
||||
if (backup_replslots && strcmp(buf, PG_REPLSLOT_DIR) != 0 &&
|
||||
path_is_prefix_of_path(PG_REPLSLOT_DIR, buf) && check_logical_replslot_dir(file->rel_path) != 1) {
|
||||
continue;
|
||||
}
|
||||
|
||||
parray_append(files, file);
|
||||
}
|
||||
|
|
@ -1914,7 +1923,7 @@ static void fio_list_dir_impl(int out, char* buf)
|
|||
|
||||
dir_list_file(file_files, req->path, req->exclude, req->follow_symlink,
|
||||
req->add_root, req->backup_logs, req->skip_hidden,
|
||||
req->external_dir_num, FIO_LOCAL_HOST);
|
||||
req->external_dir_num, FIO_LOCAL_HOST, req->backup_replslots);
|
||||
|
||||
/* send information about files to the main process */
|
||||
for (i = 0; i < (int)parray_num(file_files); i++)
|
||||
|
|
|
|||
|
|
@ -163,5 +163,7 @@ extern z_off_t fio_gzseek(gzFile f, z_off_t offset, int whence);
|
|||
extern const char* fio_gzerror(gzFile file, int *errnum);
|
||||
#endif
|
||||
|
||||
extern char check_logical_replslot_dir(const char *rel_path);
|
||||
|
||||
#endif
|
||||
|
||||
|
|
|
|||
|
|
@ -154,6 +154,7 @@ void help_pg_probackup(void)
|
|||
printf(_(" [--remote-port=port] [--ssh-options=ssh_options]\n"));
|
||||
printf(_(" [--remote-libpath=libpath]\n"));
|
||||
printf(_(" [--ttl=interval] [--expire-time=time]\n"));
|
||||
printf(_(" [--backup-pg-replslot]\n"));
|
||||
printf(_(" [--help]\n"));
|
||||
|
||||
printf(_("\n %s restore -B backup-path --instance=instance_name\n"), PROGRAM_NAME);
|
||||
|
|
@ -420,6 +421,7 @@ static void help_backup(void)
|
|||
printf(_(" [--remote-port=port] [--ssh-options=ssh_options]\n"));
|
||||
printf(_(" [--remote-libpath=libpath]\n"));
|
||||
printf(_(" [--ttl=interval] [--expire-time=time]\n\n"));
|
||||
printf(_(" [--backup-pg-replslot]\n"));
|
||||
|
||||
printf(_(" -B, --backup-path=backup-path location of the backup storage area\n"));
|
||||
printf(_(" --instance=instance_name name of the instance\n"));
|
||||
|
|
@ -441,6 +443,7 @@ static void help_backup(void)
|
|||
printf(_(" --note=text add note to backup\n"));
|
||||
printf(_(" (example: --note='backup before app update to v13.1')\n"));
|
||||
printf(_(" --archive-timeout=timeout wait timeout for WAL segment archiving (default: 5min)\n"));
|
||||
printf(_(" --backup-pg-replslot] backup of '%s' directory\n"), PG_REPLSLOT_DIR);
|
||||
|
||||
printf(_("\n Logging options:\n"));
|
||||
printf(_(" --log-level-console=log-level-console\n"));
|
||||
|
|
|
|||
|
|
@ -77,6 +77,7 @@ int rw_timeout = 0;
|
|||
|
||||
/* backup options */
|
||||
bool backup_logs = false;
|
||||
bool backup_replslots = false;
|
||||
bool smooth_checkpoint;
|
||||
char *remote_agent;
|
||||
static char *backup_note = NULL;
|
||||
|
|
@ -186,6 +187,7 @@ static ConfigOption cmd_options[] =
|
|||
{ 'b', 145, "wal", &delete_wal, SOURCE_CMD_STRICT },
|
||||
{ 'b', 146, "expired", &delete_expired, SOURCE_CMD_STRICT },
|
||||
{ 's', 172, "status", &delete_status, SOURCE_CMD_STRICT },
|
||||
{ 'b', 186, "backup-pg-replslot", &backup_replslots, SOURCE_CMD_STRICT},
|
||||
|
||||
{ 'b', 147, "force", &force, SOURCE_CMD_STRICT },
|
||||
{ 'b', 148, "compress", &compress_shortcut, SOURCE_CMD_STRICT },
|
||||
|
|
@ -550,7 +552,7 @@ static int do_actual_operate()
|
|||
elog(ERROR, "required parameter not specified: BACKUP_MODE "
|
||||
"(-b, --backup-mode)");
|
||||
|
||||
return do_backup(start_time, set_backup_params, no_validate, no_sync, backup_logs);
|
||||
return do_backup(start_time, set_backup_params, no_validate, no_sync, backup_logs, backup_replslots);
|
||||
}
|
||||
case RESTORE_CMD:
|
||||
return do_restore_or_validate(current.backup_id,
|
||||
|
|
|
|||
|
|
@ -69,6 +69,7 @@ extern const char *PROGRAM_FULL_PATH;
|
|||
#define HEADER_MAP "page_header_map"
|
||||
#define HEADER_MAP_TMP "page_header_map_tmp"
|
||||
#define PG_RELATIVE_TBLSPC_DIR "pg_location"
|
||||
#define PG_REPLSLOT_DIR "pg_replslot"
|
||||
|
||||
/* Timeout defaults */
|
||||
#define ARCHIVE_TIMEOUT_DEFAULT 300
|
||||
|
|
|
|||
|
|
@ -54,6 +54,9 @@ extern bool smooth_checkpoint;
|
|||
it will be backuped up in external dirs */
|
||||
extern parray *pgdata_nobackup_dir;
|
||||
|
||||
/* list of logical replication slots */
|
||||
extern parray *logical_replslot;
|
||||
|
||||
/* remote probackup options */
|
||||
extern char* remote_agent;
|
||||
|
||||
|
|
@ -89,7 +92,7 @@ extern const char *pgdata_exclude_dir[];
|
|||
|
||||
/* in backup.c */
|
||||
extern int do_backup(time_t start_time, pgSetBackupParams *set_backup_params,
|
||||
bool no_validate, bool no_sync, bool backup_logs);
|
||||
bool no_validate, bool no_sync, bool backup_logs, bool backup_replslots);
|
||||
extern BackupMode parse_backup_mode(const char *value);
|
||||
extern const char *deparse_backup_mode(BackupMode mode);
|
||||
extern void process_block_change(ForkNumber forknum, const RelFileNode rnode,
|
||||
|
|
@ -239,7 +242,8 @@ extern const char* deparse_compress_alg(int alg);
|
|||
/* in dir.c */
|
||||
extern void dir_list_file(parray *files, const char *root, bool exclude,
|
||||
bool follow_symlink, bool add_root, bool backup_logs,
|
||||
bool skip_hidden, int external_dir_num, fio_location location);
|
||||
bool skip_hidden, int external_dir_num, fio_location location,
|
||||
bool backup_replslots = false);
|
||||
|
||||
extern void create_data_directories(parray *dest_files,
|
||||
const char *data_dir,
|
||||
|
|
@ -432,7 +436,8 @@ extern int fio_send_file(const char *from_fullpath, const char *to_fullpath, FIL
|
|||
pgFile *file, char **errormsg);
|
||||
|
||||
extern void fio_list_dir(parray *files, const char *root, bool exclude, bool follow_symlink,
|
||||
bool add_root, bool backup_logs, bool skip_hidden, int external_dir_num);
|
||||
bool add_root, bool backup_logs, bool skip_hidden, int external_dir_num,
|
||||
bool backup_replslots = false);
|
||||
|
||||
extern bool pgut_rmtree(const char *path, bool rmtopdir, bool strict);
|
||||
|
||||
|
|
|
|||
|
|
@ -28,7 +28,6 @@
|
|||
#include "utils/builtins.h"
|
||||
#include "utils/fmgroids.h"
|
||||
#include "utils/syscache.h"
|
||||
#include "replication/worker_internal.h"
|
||||
|
||||
static List *textarray_to_stringlist(ArrayType *textarray);
|
||||
|
||||
|
|
@ -91,6 +90,13 @@ Subscription *GetSubscription(Oid subid, bool missing_ok)
|
|||
}
|
||||
sub->publications = textarray_to_stringlist(DatumGetArrayTypeP(datum));
|
||||
|
||||
datum = SysCacheGetAttr(SUBSCRIPTIONOID, tup, Anum_pg_subscription_subbinary, &isnull);
|
||||
if (unlikely(isnull)) {
|
||||
ereport(ERROR, (errcode(ERRCODE_UNEXPECTED_NULL_VALUE),
|
||||
errmsg("null binary for subscription %u", subid)));
|
||||
}
|
||||
sub->binary = DatumGetBool(datum);
|
||||
|
||||
ReleaseSysCache(tup);
|
||||
|
||||
return sub;
|
||||
|
|
@ -183,7 +189,7 @@ char *get_subscription_name(Oid subid, bool missing_ok)
|
|||
}
|
||||
|
||||
/* Clear the list content, only deal with DefElem and string content */
|
||||
static void ClearListContent(List *list)
|
||||
void ClearListContent(List *list)
|
||||
{
|
||||
ListCell *cell = NULL;
|
||||
foreach(cell, list) {
|
||||
|
|
@ -203,25 +209,6 @@ static void ClearListContent(List *list)
|
|||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* Decrypt conninfo for subscription.
|
||||
* IMPORTANT: caller should clear and free the memory after using it immediately
|
||||
*/
|
||||
char *DecryptConninfo(char *encryptConninfo)
|
||||
{
|
||||
const char* sensitiveOptionsArray[] = {"password"};
|
||||
const int sensitiveArrayLength = lengthof(sensitiveOptionsArray);
|
||||
List *defList = ConninfoToDefList(encryptConninfo);
|
||||
DecryptOptions(defList, sensitiveOptionsArray, sensitiveArrayLength, SUBSCRIPTION_MODE);
|
||||
char *decryptConninfo = DefListToString(defList);
|
||||
|
||||
/* defList has plain content, clear it before free */
|
||||
ClearListContent(defList);
|
||||
list_free_ext(defList);
|
||||
/* IMPORTANT: caller should clear and free the memory after using it immediately */
|
||||
return decryptConninfo;
|
||||
}
|
||||
|
||||
/*
|
||||
* Convert text array to list of strings.
|
||||
*
|
||||
|
|
|
|||
|
|
@ -59,7 +59,7 @@ bool open_join_children = true;
|
|||
bool will_shutdown = false;
|
||||
|
||||
/* hard-wired binary version number */
|
||||
const uint32 GRAND_VERSION_NUM = 92605;
|
||||
const uint32 GRAND_VERSION_NUM = 92606;
|
||||
|
||||
const uint32 PREDPUSH_SAME_LEVEL_VERSION_NUM = 92522;
|
||||
const uint32 UPSERT_WHERE_VERSION_NUM = 92514;
|
||||
|
|
@ -101,6 +101,7 @@ const uint32 PRIVS_DIRECTORY_VERSION_NUM = 92460;
|
|||
const uint32 COMMENT_RECORD_PARAM_VERSION_NUM = 92484;
|
||||
const uint32 SCAN_BATCH_MODE_VERSION_NUM = 92568;
|
||||
const uint32 PUBLICATION_VERSION_NUM = 92580;
|
||||
const uint32 SUBSCRIPTION_BINARY_VERSION_NUM = 92606;
|
||||
|
||||
/* Version number of the guc parameter backend_version added in V500R001C20 */
|
||||
const uint32 V5R1C20_BACKEND_VERSION_NUM = 92305;
|
||||
|
|
|
|||
|
|
@ -26,9 +26,11 @@ import logging
|
|||
try:
|
||||
from .dao.gsql_execute import GSqlExecute
|
||||
from .dao.execute_factory import ExecuteFactory
|
||||
from .mcts import MCTS
|
||||
except ImportError:
|
||||
from dao.gsql_execute import GSqlExecute
|
||||
from dao.execute_factory import ExecuteFactory
|
||||
from mcts import MCTS
|
||||
|
||||
ENABLE_MULTI_NODE = False
|
||||
SAMPLE_NUM = 5
|
||||
|
|
@ -192,9 +194,12 @@ class IndexAdvisor:
|
|||
self.workload_used_index))
|
||||
if DRIVER:
|
||||
self.db.close_conn()
|
||||
|
||||
opt_config = greedy_determine_opt_config(self.workload_info[0], atomic_config_total,
|
||||
candidate_indexes, self.index_cost_total[0])
|
||||
if MAX_INDEX_STORAGE:
|
||||
opt_config = MCTS(self.workload_info[0], atomic_config_total, candidate_indexes,
|
||||
MAX_INDEX_STORAGE, MAX_INDEX_NUM)
|
||||
else:
|
||||
opt_config = greedy_determine_opt_config(self.workload_info[0], atomic_config_total,
|
||||
candidate_indexes, self.index_cost_total[0])
|
||||
self.retain_lower_cost_index(candidate_indexes)
|
||||
if len(opt_config) == 0:
|
||||
print("No optimal indexes generated!")
|
||||
|
|
@ -943,7 +948,7 @@ def check_parameter(args):
|
|||
raise argparse.ArgumentTypeError("%s is an invalid positive int value" %
|
||||
args.max_index_num)
|
||||
if args.max_index_storage is not None and args.max_index_storage <= 0:
|
||||
raise argparse.ArgumentTypeError("%s is an invalid positive int value" %
|
||||
raise argparse.ArgumentTypeError("%s is an invalid positive float value" %
|
||||
args.max_index_storage)
|
||||
JSON_TYPE = args.json
|
||||
MAX_INDEX_NUM = args.max_index_num
|
||||
|
|
@ -971,7 +976,7 @@ def main(argv):
|
|||
arg_parser.add_argument(
|
||||
"--max_index_num", help="Maximum number of suggested indexes", type=int)
|
||||
arg_parser.add_argument("--max_index_storage",
|
||||
help="Maximum storage of suggested indexes/MB", type=int)
|
||||
help="Maximum storage of suggested indexes/MB", type=float)
|
||||
arg_parser.add_argument("--multi_iter_mode", action='store_true',
|
||||
help="Whether to use multi-iteration algorithm", default=False)
|
||||
arg_parser.add_argument("--multi_node", action='store_true',
|
||||
|
|
|
|||
|
|
@ -0,0 +1,397 @@
|
|||
import sys
|
||||
import math
|
||||
import random
|
||||
import copy
|
||||
|
||||
STORAGE_THRESHOLD = 0
|
||||
AVAILABLE_CHOICES = None
|
||||
ATOMIC_CHOICES = None
|
||||
WORKLOAD_INFO = None
|
||||
MAX_INDEX_NUM = 0
|
||||
|
||||
|
||||
def is_same_index(index, compared_index):
|
||||
return index.table == compared_index.table and \
|
||||
index.columns == compared_index.columns and \
|
||||
index.index_type == compared_index.index_type
|
||||
|
||||
|
||||
def atomic_config_is_valid(atomic_config, config):
|
||||
# if candidate indexes contains all atomic index of current config1, then record it
|
||||
for atomic_index in atomic_config:
|
||||
is_exist = False
|
||||
for index in config:
|
||||
if is_same_index(index, atomic_index):
|
||||
index.storage = atomic_index.storage
|
||||
is_exist = True
|
||||
break
|
||||
if not is_exist:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
def find_subsets_num(choice):
|
||||
atomic_subsets_num = []
|
||||
for pos, atomic in enumerate(ATOMIC_CHOICES):
|
||||
if not atomic or len(atomic) > len(choice):
|
||||
continue
|
||||
# find valid atomic index
|
||||
if atomic_config_is_valid(atomic, choice):
|
||||
atomic_subsets_num.append(pos)
|
||||
# find the same atomic index as the candidate index
|
||||
if len(atomic) == 1 and (is_same_index(choice[-1], atomic[0])):
|
||||
choice[-1].atomic_pos = pos
|
||||
return atomic_subsets_num
|
||||
|
||||
|
||||
def find_best_benefit(choice):
|
||||
atomic_subsets_num = find_subsets_num(choice)
|
||||
total_benefit = 0
|
||||
for ind, obj in enumerate(WORKLOAD_INFO):
|
||||
# calculate the best benefit for the current sql
|
||||
max_benefit = 0
|
||||
for pos in atomic_subsets_num:
|
||||
if (obj.cost_list[0] - obj.cost_list[pos]) > max_benefit:
|
||||
max_benefit = obj.cost_list[0] - obj.cost_list[pos]
|
||||
total_benefit += max_benefit
|
||||
return total_benefit
|
||||
|
||||
|
||||
def get_diff(available_choices, choices):
|
||||
except_choices = copy.copy(available_choices)
|
||||
for i in available_choices:
|
||||
for j in choices:
|
||||
if is_same_index(i, j):
|
||||
except_choices.remove(i)
|
||||
return except_choices
|
||||
|
||||
|
||||
class State(object):
|
||||
"""
|
||||
The game state of the Monte Carlo tree search,
|
||||
the state data recorded under a certain Node node,
|
||||
including the current game score, the current number of game rounds,
|
||||
and the execution record from the beginning to the current.
|
||||
|
||||
It is necessary to realize whether the current state has reached the end of the game state,
|
||||
and support the operation of randomly fetching from the Action collection.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.current_storage = 0.0
|
||||
self.current_benefit = 0.0
|
||||
# record the sum of choices up to the current state
|
||||
self.accumulation_choices = []
|
||||
# record available choices of current state
|
||||
self.available_choices = []
|
||||
self.displayable_choices = []
|
||||
|
||||
def get_available_choices(self):
|
||||
return self.available_choices
|
||||
|
||||
def set_available_choices(self, choices):
|
||||
self.available_choices = choices
|
||||
|
||||
def get_current_storage(self):
|
||||
return self.current_storage
|
||||
|
||||
def set_current_storage(self, value):
|
||||
self.current_storage = value
|
||||
|
||||
def get_current_benefit(self):
|
||||
return self.current_benefit
|
||||
|
||||
def set_current_benefit(self, value):
|
||||
self.current_benefit = value
|
||||
|
||||
def get_accumulation_choices(self):
|
||||
return self.accumulation_choices
|
||||
|
||||
def set_accumulation_choices(self, choices):
|
||||
self.accumulation_choices = choices
|
||||
|
||||
def is_terminal(self):
|
||||
# the current node is a leaf node
|
||||
return len(self.accumulation_choices) == MAX_INDEX_NUM
|
||||
|
||||
def compute_benefit(self):
|
||||
return self.current_benefit
|
||||
|
||||
def get_next_state_with_random_choice(self):
|
||||
# ensure that the choices taken are not repeated
|
||||
if not self.available_choices:
|
||||
return None
|
||||
random_choice = random.choice([choice for choice in self.available_choices])
|
||||
self.available_choices.remove(random_choice)
|
||||
choice = copy.copy(self.accumulation_choices)
|
||||
choice.append(random_choice)
|
||||
benefit = find_best_benefit(choice)
|
||||
# if current choice not satisfy restrictions, then continue get next choice
|
||||
if benefit <= self.current_benefit or \
|
||||
self.current_storage + random_choice.storage > STORAGE_THRESHOLD:
|
||||
return self.get_next_state_with_random_choice()
|
||||
|
||||
next_state = State()
|
||||
# initialize the properties of the new state
|
||||
next_state.set_accumulation_choices(choice)
|
||||
next_state.set_current_benefit(benefit)
|
||||
next_state.set_current_storage(self.current_storage + random_choice.storage)
|
||||
next_state.set_available_choices(get_diff(AVAILABLE_CHOICES, choice))
|
||||
return next_state
|
||||
|
||||
def __repr__(self):
|
||||
self.displayable_choices = ['{}: {}'.format(choice.table, choice.columns)
|
||||
for choice in self.accumulation_choices]
|
||||
return "reward: {}, storage :{}, choices: {}".format(
|
||||
self.current_benefit, self.current_storage, self.displayable_choices)
|
||||
|
||||
|
||||
class Node(object):
|
||||
"""
|
||||
The Node of the Monte Carlo tree search tree contains the parent node and
|
||||
current point information,
|
||||
which is used to calculate the traversal times and quality value of the UCB,
|
||||
and the State of the Node selected by the game.
|
||||
"""
|
||||
def __init__(self):
|
||||
self.visit_number = 0
|
||||
self.quality = 0.0
|
||||
|
||||
self.parent = None
|
||||
self.children = []
|
||||
self.state = None
|
||||
|
||||
def get_parent(self):
|
||||
return self.parent
|
||||
|
||||
def set_parent(self, parent):
|
||||
self.parent = parent
|
||||
|
||||
def get_children(self):
|
||||
return self.children
|
||||
|
||||
def expand_child(self, node):
|
||||
node.set_parent(self)
|
||||
self.children.append(node)
|
||||
|
||||
def set_state(self, state):
|
||||
self.state = state
|
||||
|
||||
def get_state(self):
|
||||
return self.state
|
||||
|
||||
def get_visit_number(self):
|
||||
return self.visit_number
|
||||
|
||||
def set_visit_number(self, number):
|
||||
self.visit_number = number
|
||||
|
||||
def update_visit_number(self):
|
||||
self.visit_number += 1
|
||||
|
||||
def get_quality_value(self):
|
||||
return self.quality
|
||||
|
||||
def set_quality_value(self, value):
|
||||
self.quality = value
|
||||
|
||||
def update_quality_value(self, reward):
|
||||
self.quality += reward
|
||||
|
||||
def is_all_expand(self):
|
||||
return len(self.children) == \
|
||||
len(AVAILABLE_CHOICES) - len(self.get_state().get_accumulation_choices())
|
||||
|
||||
def __repr__(self):
|
||||
return "Node: {}, Q/N: {}/{}, State: {}".format(
|
||||
hash(self), self.quality, self.visit_number, self.state)
|
||||
|
||||
|
||||
def tree_policy(node):
|
||||
"""
|
||||
In the Selection and Expansion stages of Monte Carlo tree search,
|
||||
the node that needs to be searched (such as the root node) is passed in,
|
||||
and the best node that needs to be expanded is returned
|
||||
according to the exploration/exploitation algorithm.
|
||||
Note that if the node is a leaf node, it will be returned directly.
|
||||
|
||||
The basic strategy is to first find the child nodes that have not been selected at present,
|
||||
and select them randomly if there are more than one. If both are selected,
|
||||
find the one with the largest UCB value that has weighed exploration/exploitation,
|
||||
and randomly select if the UCB values are equal.
|
||||
"""
|
||||
|
||||
# check if the current node is leaf node
|
||||
while node and not node.get_state().is_terminal():
|
||||
|
||||
if node.is_all_expand():
|
||||
node = best_child(node, True)
|
||||
else:
|
||||
# return the new sub node
|
||||
sub_node = expand(node)
|
||||
# when there is no node that satisfies the condition in the remaining nodes,
|
||||
# this state is empty
|
||||
if sub_node.get_state():
|
||||
return sub_node
|
||||
|
||||
# return the leaf node
|
||||
return node
|
||||
|
||||
|
||||
def default_policy(node):
|
||||
"""
|
||||
In the Simulation stage of Monte Carlo tree search, input a node that needs to be expanded,
|
||||
create a new node after random operation, and return the reward of the new node.
|
||||
Note that the input node should not be a child node,
|
||||
and there are unexecuted Actions that can be expendable.
|
||||
|
||||
The basic strategy is to choose the Action at random.
|
||||
"""
|
||||
|
||||
# get the state of the game
|
||||
current_state = copy.deepcopy(node.get_state())
|
||||
|
||||
# run until the game over
|
||||
while not current_state.is_terminal():
|
||||
# pick one random action to play and get next state
|
||||
next_state = current_state.get_next_state_with_random_choice()
|
||||
if not next_state:
|
||||
break
|
||||
current_state = next_state
|
||||
|
||||
final_state_reward = current_state.compute_benefit()
|
||||
return final_state_reward
|
||||
|
||||
|
||||
def expand(node):
|
||||
"""
|
||||
Enter a node, expand a new node on the node, use the random method to execute the Action,
|
||||
and return the new node. Note that it is necessary to ensure that the newly
|
||||
added nodes are different from other node Action
|
||||
"""
|
||||
|
||||
new_state = node.get_state().get_next_state_with_random_choice()
|
||||
sub_node = Node()
|
||||
sub_node.set_state(new_state)
|
||||
node.expand_child(sub_node)
|
||||
|
||||
return sub_node
|
||||
|
||||
|
||||
def best_child(node, is_exploration):
|
||||
"""
|
||||
Using the UCB algorithm,
|
||||
select the child node with the highest score after weighing the exploration and exploitation.
|
||||
Note that if it is the prediction stage,
|
||||
the current Q-value score with the highest score is directly selected.
|
||||
"""
|
||||
|
||||
best_score = -sys.maxsize
|
||||
best_sub_node = None
|
||||
|
||||
# travel all sub nodes to find the best one
|
||||
for sub_node in node.get_children():
|
||||
# The children nodes of the node contains the children node whose state is empty,
|
||||
# this kind of node comes from the node that does not meet the conditions.
|
||||
if not sub_node.get_state():
|
||||
continue
|
||||
# ignore exploration for inference
|
||||
if is_exploration:
|
||||
C = 1 / math.sqrt(2.0)
|
||||
else:
|
||||
C = 0.0
|
||||
|
||||
# UCB = quality / times + C * sqrt(2 * ln(total_times) / times)
|
||||
left = sub_node.get_quality_value() / sub_node.get_visit_number()
|
||||
right = 2.0 * math.log(node.get_visit_number()) / sub_node.get_visit_number()
|
||||
score = left + C * math.sqrt(right)
|
||||
# get the maximum score, while filtering nodes that do not meet the space constraints and
|
||||
# nodes that have no revenue
|
||||
if score > best_score \
|
||||
and sub_node.get_state().get_current_storage() <= STORAGE_THRESHOLD \
|
||||
and sub_node.get_state().get_current_benefit() > 0:
|
||||
best_sub_node = sub_node
|
||||
best_score = score
|
||||
|
||||
return best_sub_node
|
||||
|
||||
|
||||
def backpropagate(node, reward):
|
||||
"""
|
||||
In the Backpropagation stage of Monte Carlo tree search,
|
||||
input the node that needs to be expended and the reward of the newly executed Action,
|
||||
feed it back to the expend node and all upstream nodes,
|
||||
and update the corresponding data.
|
||||
"""
|
||||
|
||||
# update util the root node
|
||||
while node is not None:
|
||||
# update the visit number
|
||||
node.update_visit_number()
|
||||
|
||||
# update the quality value
|
||||
node.update_quality_value(reward)
|
||||
|
||||
# change the node to the parent node
|
||||
node = node.parent
|
||||
|
||||
|
||||
def monte_carlo_tree_search(node):
|
||||
"""
|
||||
Implement the Monte Carlo tree search algorithm, pass in a root node,
|
||||
expand new nodes and update data according to the
|
||||
tree structure that has been explored before in a limited time,
|
||||
and then return as long as the child node with the highest exploitation.
|
||||
|
||||
When making predictions,
|
||||
you only need to select the node with the largest exploitation according to the Q value,
|
||||
and find the next optimal node.
|
||||
"""
|
||||
|
||||
computation_budget = len(AVAILABLE_CHOICES) * 3
|
||||
|
||||
# run as much as possible under the computation budget
|
||||
for i in range(computation_budget):
|
||||
# 1. find the best node to expand
|
||||
expand_node = tree_policy(node)
|
||||
if not expand_node:
|
||||
# when it is None, it means that all nodes are added but no nodes meet the space limit
|
||||
break
|
||||
# 2. random get next action and get reward
|
||||
reward = default_policy(expand_node)
|
||||
|
||||
# 3. update all passing nodes with reward
|
||||
backpropagate(expand_node, reward)
|
||||
|
||||
# get the best next node
|
||||
best_next_node = best_child(node, False)
|
||||
|
||||
return best_next_node
|
||||
|
||||
|
||||
def MCTS(workload_info, atomic_choices, available_choices, storage_threshold, max_index_num):
|
||||
global ATOMIC_CHOICES, STORAGE_THRESHOLD, WORKLOAD_INFO, AVAILABLE_CHOICES, MAX_INDEX_NUM
|
||||
WORKLOAD_INFO = workload_info
|
||||
AVAILABLE_CHOICES = available_choices
|
||||
ATOMIC_CHOICES = atomic_choices
|
||||
STORAGE_THRESHOLD = storage_threshold
|
||||
MAX_INDEX_NUM = max_index_num if max_index_num else len(available_choices)
|
||||
|
||||
# create the initialized state and initialized node
|
||||
init_state = State()
|
||||
choices = copy.copy(available_choices)
|
||||
init_state.set_available_choices(choices)
|
||||
init_node = Node()
|
||||
init_node.set_state(init_state)
|
||||
current_node = init_node
|
||||
|
||||
opt_config = []
|
||||
# set the rounds to play
|
||||
for i in range(len(AVAILABLE_CHOICES)):
|
||||
if current_node:
|
||||
current_node = monte_carlo_tree_search(current_node)
|
||||
if current_node:
|
||||
opt_config = current_node.state.accumulation_choices
|
||||
else:
|
||||
break
|
||||
return opt_config
|
||||
|
|
@ -27,6 +27,7 @@ from collections.abc import Iterable
|
|||
from collections import defaultdict
|
||||
|
||||
import index_advisor_workload as iaw
|
||||
import mcts
|
||||
|
||||
|
||||
def hash_any(obj):
|
||||
|
|
@ -227,6 +228,32 @@ select * from student_range_part1 where credit=1;
|
|||
|
||||
class IndexAdvisorTester(unittest.TestCase):
|
||||
|
||||
def test_mcts(self):
|
||||
storage_threshold = 12
|
||||
index1 = iaw.IndexItem('public.a', 'col1', index_type='global')
|
||||
index2 = iaw.IndexItem('public.b', 'col1', index_type='global')
|
||||
index3 = iaw.IndexItem('public.c', 'col1', index_type='global')
|
||||
index4 = iaw.IndexItem('public.d', 'col1', index_type='global')
|
||||
|
||||
atomic_index1 = iaw.IndexItem('public.a', 'col1', index_type='global')
|
||||
atomic_index2 = iaw.IndexItem('public.b', 'col1', index_type='global')
|
||||
atomic_index3 = iaw.IndexItem('public.c', 'col1', index_type='global')
|
||||
atomic_index4 = iaw.IndexItem('public.d', 'col1', index_type='global')
|
||||
|
||||
atomic_index1.storage = 10
|
||||
atomic_index2.storage = 4
|
||||
atomic_index3.storage = 7
|
||||
available_choices = [index1, index2, index3, index4]
|
||||
atomic_choices = [[], [atomic_index2], [atomic_index1], [atomic_index3],
|
||||
[atomic_index2, atomic_index3], [atomic_index4]]
|
||||
query = iaw.QueryItem('select * from gia_01', 1)
|
||||
query.cost_list = [10, 7, 5, 9, 4, 11]
|
||||
workload_info = [query]
|
||||
|
||||
results = mcts.MCTS(workload_info, atomic_choices, available_choices, storage_threshold, 2)
|
||||
self.assertLessEqual([index1.atomic_pos, index2.atomic_pos, index3.atomic_pos], [2, 1, 3])
|
||||
self.assertSetEqual({results[0].table, results[1].table}, {'public.b', 'public.c'})
|
||||
|
||||
def test_get_indexable_columns(self):
|
||||
tables = 'table1 table2 table2 table3 table3 table3'.split()
|
||||
columns = 'col1,col2 col2 col3 col1,col2 col2,col3 col2,col5'.split()
|
||||
|
|
|
|||
|
|
@ -44,7 +44,7 @@
|
|||
#include "utils/array.h"
|
||||
#include "utils/acl.h"
|
||||
|
||||
static void ConnectPublisher(char *conninfo, char* slotname);
|
||||
static bool ConnectPublisher(char* conninfo, char* slotname);
|
||||
static void CreateSlotInPublisher(char *slotname);
|
||||
static void ValidateReplicationSlot(char *slotname, List *publications);
|
||||
|
||||
|
|
@ -56,7 +56,7 @@ static void ValidateReplicationSlot(char *slotname, List *publications);
|
|||
* accommodate that.
|
||||
*/
|
||||
static void parse_subscription_options(const List *options, char **conninfo, List **publications, bool *enabled_given,
|
||||
bool *enabled, bool *slot_name_given, char **slot_name, char **synchronous_commit)
|
||||
bool *enabled, bool *slot_name_given, char **slot_name, char **synchronous_commit, bool *binary_given, bool *binary)
|
||||
{
|
||||
ListCell *lc;
|
||||
|
||||
|
|
@ -76,6 +76,10 @@ static void parse_subscription_options(const List *options, char **conninfo, Lis
|
|||
if (synchronous_commit) {
|
||||
*synchronous_commit = NULL;
|
||||
}
|
||||
if (binary) {
|
||||
*binary_given = false;
|
||||
*binary = false;
|
||||
}
|
||||
|
||||
/* Parse options */
|
||||
foreach (lc, options) {
|
||||
|
|
@ -124,6 +128,15 @@ static void parse_subscription_options(const List *options, char **conninfo, Lis
|
|||
/* Test if the given value is valid for synchronous_commit GUC. */
|
||||
(void)set_config_option("synchronous_commit", *synchronous_commit, PGC_BACKEND, PGC_S_TEST, GUC_ACTION_SET,
|
||||
false, 0, false);
|
||||
} else if (strcmp(defel->defname, "binary") == 0 && binary) {
|
||||
if (*binary_given) {
|
||||
ereport(ERROR,
|
||||
(errcode(ERRCODE_SYNTAX_ERROR),
|
||||
errmsg("conflicting or redundant options")));
|
||||
}
|
||||
|
||||
*binary_given = true;
|
||||
*binary = defGetBoolean(defel);
|
||||
} else {
|
||||
ereport(ERROR,
|
||||
(errcode(ERRCODE_SYNTAX_ERROR), errmsg("unrecognized subscription parameter: %s", defel->defname)));
|
||||
|
|
@ -197,26 +210,82 @@ static Datum publicationListToArray(List *publist)
|
|||
}
|
||||
|
||||
/*
|
||||
* connect publisher and create slot.
|
||||
* the input conninfo should be encrypt, we will decrypt password inside
|
||||
* Parse the original connection string which is encrypted, poll all hosts and ports,
|
||||
* and try to connect to the publisher.
|
||||
* When checkRemoteMode is true, the remotemode must be normal or primary.
|
||||
* Return true to indicate successful connection.
|
||||
*/
|
||||
static void ConnectPublisher(char *conninfo, char *slotname)
|
||||
bool AttemptConnectPublisher(const char *conninfoOriginal, char* slotname, bool checkRemoteMode)
|
||||
{
|
||||
size_t conninfoLen = strlen(conninfoOriginal) + 1;
|
||||
|
||||
char* conninfo = NULL;
|
||||
StringInfoData conninfoWithoutHostport;
|
||||
initStringInfo(&conninfoWithoutHostport);
|
||||
HostPort* hostPortList[MAX_REPLNODE_NUM] = {NULL};
|
||||
ParseConninfo(conninfoOriginal, &conninfoWithoutHostport, hostPortList);
|
||||
if (hostPortList[0] == NULL) {
|
||||
ereport(ERROR, (errcode(ERRCODE_SYNTAX_ERROR), errmsg(
|
||||
"invalid connection string syntax, missing host and port")));
|
||||
}
|
||||
bool connectSuccess = false;
|
||||
conninfo = (char*)palloc(conninfoLen * sizeof(char));
|
||||
for (int i = 0; i < MAX_REPLNODE_NUM; ++i) {
|
||||
if (hostPortList[i] == NULL) {
|
||||
break;
|
||||
}
|
||||
int ret = snprintf_s(conninfo, conninfoLen, conninfoLen - 1,
|
||||
"%s host=%s port=%s", conninfoWithoutHostport.data,
|
||||
hostPortList[i]->host, hostPortList[i]->port);
|
||||
securec_check_ss(ret, "\0", "\0");
|
||||
|
||||
connectSuccess = ConnectPublisher(conninfo, slotname);
|
||||
if (!connectSuccess) {
|
||||
/* try next host */
|
||||
continue;
|
||||
}
|
||||
if (!checkRemoteMode) {
|
||||
break;
|
||||
}
|
||||
ServerMode publisherServerMde = IdentifyRemoteMode();
|
||||
if (publisherServerMde == NORMAL_MODE || publisherServerMde == PRIMARY_MODE) {
|
||||
break;
|
||||
}
|
||||
/* it's a standby, try next host */
|
||||
(WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_disconnect();
|
||||
connectSuccess = false;
|
||||
}
|
||||
pfree_ext(conninfo);
|
||||
|
||||
/* clean up */
|
||||
FreeStringInfo(&conninfoWithoutHostport);
|
||||
for (int i = 0; i < MAX_REPLNODE_NUM; ++i) {
|
||||
if (hostPortList[i] == NULL) {
|
||||
break;
|
||||
}
|
||||
pfree_ext(hostPortList[i]->host);
|
||||
pfree_ext(hostPortList[i]->port);
|
||||
pfree_ext(hostPortList[i]);
|
||||
}
|
||||
return connectSuccess;
|
||||
}
|
||||
|
||||
/*
|
||||
* connect to publisher with conninfo
|
||||
*/
|
||||
static bool ConnectPublisher(char* conninfo, char* slotname)
|
||||
{
|
||||
/* Try to connect to the publisher. */
|
||||
volatile WalRcvData *walrcv = t_thrd.walreceiverfuncs_cxt.WalRcv;
|
||||
SpinLockAcquire(&walrcv->mutex);
|
||||
walrcv->conn_target = REPCONNTARGET_PUBLICATION;
|
||||
SpinLockRelease(&walrcv->mutex);
|
||||
|
||||
char *decryptConninfo = DecryptConninfo(conninfo);
|
||||
char* decryptConninfo = EncryptOrDecryptConninfo(conninfo, 'D');
|
||||
bool connectSuccess = (WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_connect(decryptConninfo, NULL, slotname, -1);
|
||||
int rc = memset_s(decryptConninfo, strlen(decryptConninfo), 0, strlen(decryptConninfo));
|
||||
securec_check(rc, "", "");
|
||||
pfree_ext(decryptConninfo);
|
||||
|
||||
if (!connectSuccess) {
|
||||
ereport(ERROR, (errcode(ERRCODE_CONNECTION_FAILURE), errmsg("could not connect to the publisher")));
|
||||
}
|
||||
return connectSuccess;
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
@ -293,9 +362,10 @@ ObjectAddress CreateSubscription(CreateSubscriptionStmt *stmt, bool isTopLevel)
|
|||
bool enabled_given = false;
|
||||
bool enabled = true;
|
||||
char *synchronous_commit;
|
||||
char *conninfo;
|
||||
char *slotname;
|
||||
bool slotname_given;
|
||||
bool binary;
|
||||
bool binary_given;
|
||||
char originname[NAMEDATALEN];
|
||||
List *publications;
|
||||
int rc;
|
||||
|
|
@ -305,7 +375,7 @@ ObjectAddress CreateSubscription(CreateSubscriptionStmt *stmt, bool isTopLevel)
|
|||
* Connection and publication should not be specified here.
|
||||
*/
|
||||
parse_subscription_options(stmt->options, NULL, NULL, &enabled_given, &enabled, &slotname_given, &slotname,
|
||||
&synchronous_commit);
|
||||
&synchronous_commit, &binary_given, &binary);
|
||||
|
||||
/*
|
||||
* Since creating a replication slot is not transactional, rolling back
|
||||
|
|
@ -333,11 +403,10 @@ ObjectAddress CreateSubscription(CreateSubscriptionStmt *stmt, bool isTopLevel)
|
|||
synchronous_commit = "off";
|
||||
}
|
||||
|
||||
conninfo = stmt->conninfo;
|
||||
publications = stmt->publication;
|
||||
|
||||
/* Check the connection info string. */
|
||||
libpqrcv_check_conninfo(conninfo);
|
||||
libpqrcv_check_conninfo(stmt->conninfo);
|
||||
|
||||
/* Everything ok, form a new tuple. */
|
||||
rc = memset_s(values, sizeof(values), 0, sizeof(values));
|
||||
|
|
@ -349,18 +418,12 @@ ObjectAddress CreateSubscription(CreateSubscriptionStmt *stmt, bool isTopLevel)
|
|||
values[Anum_pg_subscription_subname - 1] = DirectFunctionCall1(namein, CStringGetDatum(stmt->subname));
|
||||
values[Anum_pg_subscription_subowner - 1] = ObjectIdGetDatum(owner);
|
||||
values[Anum_pg_subscription_subenabled - 1] = BoolGetDatum(enabled);
|
||||
values[Anum_pg_subscription_subbinary - 1] = BoolGetDatum(binary);
|
||||
|
||||
/* encrypt conninfo */
|
||||
List *conninfoList = ConninfoToDefList(stmt->conninfo);
|
||||
/* Sensitive options for subscription, will be encrypted when saved to catalog. */
|
||||
const char* sensitiveOptionsArray[] = {"password"};
|
||||
const int sensitiveArrayLength = lengthof(sensitiveOptionsArray);
|
||||
EncryptGenericOptions(conninfoList, sensitiveOptionsArray, sensitiveArrayLength, SUBSCRIPTION_MODE);
|
||||
char *encryptConninfo = DefListToString(conninfoList);
|
||||
|
||||
char *encryptConninfo = EncryptOrDecryptConninfo(stmt->conninfo, 'E');
|
||||
values[Anum_pg_subscription_subconninfo - 1] = CStringGetTextDatum(encryptConninfo);
|
||||
|
||||
pfree_ext(conninfoList);
|
||||
if (enabled) {
|
||||
if (!slotname_given) {
|
||||
slotname = stmt->subname;
|
||||
|
|
@ -396,11 +459,14 @@ ObjectAddress CreateSubscription(CreateSubscriptionStmt *stmt, bool isTopLevel)
|
|||
*/
|
||||
if (enabled) {
|
||||
Assert(slotname);
|
||||
ConnectPublisher(encryptConninfo, slotname);
|
||||
|
||||
if (!AttemptConnectPublisher(encryptConninfo, slotname, true)) {
|
||||
ereport(ERROR, (errcode(ERRCODE_CONNECTION_FAILURE), errmsg("Failed to connect to publisher.")));
|
||||
}
|
||||
|
||||
CreateSlotInPublisher(slotname);
|
||||
(WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_disconnect();
|
||||
}
|
||||
|
||||
pfree_ext(encryptConninfo);
|
||||
heap_close(rel, RowExclusiveLock);
|
||||
rc = memset_s(stmt->conninfo, strlen(stmt->conninfo), 0, strlen(stmt->conninfo));
|
||||
|
|
@ -439,6 +505,8 @@ ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt)
|
|||
Oid subid;
|
||||
bool enabled_given = false;
|
||||
bool enabled;
|
||||
bool binary_given;
|
||||
bool binary;
|
||||
char *synchronous_commit;
|
||||
char *conninfo;
|
||||
char *slot_name;
|
||||
|
|
@ -473,7 +541,7 @@ ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt)
|
|||
|
||||
/* Parse options. */
|
||||
parse_subscription_options(stmt->options, &conninfo, &publications, &enabled_given, &enabled, &slotname_given,
|
||||
&slot_name, &synchronous_commit);
|
||||
&slot_name, &synchronous_commit, &binary_given, &binary);
|
||||
|
||||
/* Form a new tuple. */
|
||||
rc = memset_s(nulls, sizeof(nulls), false, sizeof(nulls));
|
||||
|
|
@ -490,23 +558,15 @@ ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt)
|
|||
if (conninfo) {
|
||||
/* Check the connection info string. */
|
||||
libpqrcv_check_conninfo(conninfo);
|
||||
|
||||
/* encrypt conninfo */
|
||||
List *conninfoList = ConninfoToDefList(conninfo);
|
||||
/* Sensitive options for subscription, will be encrypted when saved to catalog. */
|
||||
const char* sensitiveOptionsArray[] = {"password"};
|
||||
const int sensitiveArrayLength = lengthof(sensitiveOptionsArray);
|
||||
EncryptGenericOptions(conninfoList, sensitiveOptionsArray, sensitiveArrayLength, SUBSCRIPTION_MODE);
|
||||
encryptConninfo = DefListToString(conninfoList);
|
||||
needFreeConninfo = true;
|
||||
|
||||
encryptConninfo = EncryptOrDecryptConninfo(conninfo, 'E');
|
||||
rc = memset_s(conninfo, strlen(conninfo), 0, strlen(conninfo));
|
||||
securec_check(rc, "\0", "\0");
|
||||
values[Anum_pg_subscription_subconninfo - 1] = CStringGetTextDatum(encryptConninfo);
|
||||
replaces[Anum_pg_subscription_subconninfo - 1] = true;
|
||||
needFreeConninfo = true;
|
||||
|
||||
pfree_ext(conninfoList);
|
||||
|
||||
/* need to check whether new conninfo can be used to connect to new publisher */
|
||||
if (sub->enabled || (enabled_given && enabled)) {
|
||||
/* we need to check whether new conninfo can be used to connect to new publisher */
|
||||
checkConn = true;
|
||||
}
|
||||
}
|
||||
|
|
@ -548,6 +608,10 @@ ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt)
|
|||
values[Anum_pg_subscription_subsynccommit - 1] = CStringGetTextDatum(synchronous_commit);
|
||||
replaces[Anum_pg_subscription_subsynccommit - 1] = true;
|
||||
}
|
||||
if (binary_given) {
|
||||
values[Anum_pg_subscription_subbinary - 1] = BoolGetDatum(binary);
|
||||
replaces[Anum_pg_subscription_subbinary - 1] = true;
|
||||
}
|
||||
if (publications != NIL) {
|
||||
values[Anum_pg_subscription_subpublications - 1] = publicationListToArray(publications);
|
||||
replaces[Anum_pg_subscription_subpublications - 1] = true;
|
||||
|
|
@ -570,16 +634,18 @@ ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt)
|
|||
if (sub->enabled && !enabled) {
|
||||
ereport(ERROR, (errmsg("If you want to deactivate this subscription, use DROP SUBSCRIPTION.")));
|
||||
}
|
||||
/* enable subscription */
|
||||
if (!sub->enabled && enabled) {
|
||||
/* if slot hasn't been created, then create it */
|
||||
if (!sub->slotname || !*(sub->slotname)) {
|
||||
/* enabling subscription, but slot hasn't been created,
|
||||
* then mark createSlot to true.
|
||||
*/
|
||||
if (!sub->enabled && enabled && (!sub->slotname || !*(sub->slotname))) {
|
||||
createSlot = true;
|
||||
}
|
||||
}
|
||||
|
||||
if (checkConn || createSlot || validateSlot) {
|
||||
ConnectPublisher(encryptConninfo, finalSlotName);
|
||||
if (!AttemptConnectPublisher(encryptConninfo, finalSlotName, true)) {
|
||||
ereport(ERROR, (errcode(ERRCODE_CONNECTION_FAILURE), errmsg(
|
||||
checkConn ? "The new conninfo cannot connect to new publisher." : "Failed to connect to publisher.")));
|
||||
}
|
||||
|
||||
if (createSlot) {
|
||||
CreateSlotInPublisher(finalSlotName);
|
||||
|
|
@ -597,12 +663,6 @@ ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt)
|
|||
if (needFreeConninfo) {
|
||||
pfree_ext(encryptConninfo);
|
||||
}
|
||||
|
||||
if (conninfo) {
|
||||
rc = memset_s(conninfo, strlen(conninfo), 0, strlen(conninfo));
|
||||
securec_check(rc, "", "");
|
||||
}
|
||||
|
||||
return myself;
|
||||
}
|
||||
|
||||
|
|
@ -753,7 +813,11 @@ void DropSubscription(DropSubscriptionStmt *stmt, bool isTopLevel)
|
|||
initStringInfo(&cmd);
|
||||
appendStringInfo(&cmd, "DROP_REPLICATION_SLOT %s", quote_identifier(slotname));
|
||||
|
||||
ConnectPublisher(conninfo, slotname);
|
||||
if (!AttemptConnectPublisher(conninfo, slotname, true)) {
|
||||
ereport(ERROR, (errcode(ERRCODE_CONNECTION_FAILURE), errmsg(
|
||||
"could not connect to publisher.")));
|
||||
}
|
||||
|
||||
PG_TRY();
|
||||
{
|
||||
int sqlstate = 0;
|
||||
|
|
@ -779,6 +843,7 @@ void DropSubscription(DropSubscriptionStmt *stmt, bool isTopLevel)
|
|||
|
||||
(WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_disconnect();
|
||||
|
||||
pfree_ext(conninfo);
|
||||
pfree(cmd.data);
|
||||
heap_close(rel, NoLock);
|
||||
}
|
||||
|
|
@ -908,3 +973,149 @@ void RenameSubscription(List *oldname, const char *newname)
|
|||
|
||||
return;
|
||||
}
|
||||
|
||||
/*
|
||||
* Parse the host or port string into a string array,
|
||||
* where host and port are separated by ",".
|
||||
* input: conn --- host or port string separated by ","
|
||||
* output: connArray --- host or port string array
|
||||
* return: the length of connArray
|
||||
* for example:
|
||||
* (1):
|
||||
* conn = 1.1.1.1,2.2.2.2,...,9.9.9.9
|
||||
* connArray = {
|
||||
* 1,.1.1.1,
|
||||
* 2.2.2.2,
|
||||
* ...,
|
||||
* 9.9.9.9
|
||||
* }
|
||||
* return 9
|
||||
* (2):
|
||||
* conn = 1,2,...,9
|
||||
* connArray = {1,2,...,9}
|
||||
* return 9
|
||||
*/
|
||||
static int HostsPortsToArray(const char* conn, char** connArray)
|
||||
{
|
||||
if (conn == NULL) {
|
||||
return 0;
|
||||
}
|
||||
char* cp = NULL;
|
||||
char* cur = NULL;
|
||||
char *buf = pstrdup(conn);
|
||||
|
||||
cp = buf;
|
||||
int i = 0;
|
||||
while (*cp) {
|
||||
cur = cp;
|
||||
while (*cp && *cp != ',') {
|
||||
++cp;
|
||||
}
|
||||
if (*cp == ',') {
|
||||
*cp = '\0';
|
||||
++cp;
|
||||
}
|
||||
if (i >= MAX_REPLNODE_NUM) {
|
||||
ereport(ERROR, (errmsg("Currently, a maximum of %d servers are "
|
||||
"supported.", MAX_REPLNODE_NUM)));
|
||||
}
|
||||
connArray[i++] = pstrdup(cur);
|
||||
|
||||
if (*cp == 0) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
pfree(buf);
|
||||
return i;
|
||||
}
|
||||
|
||||
/*
|
||||
* parse host and port
|
||||
*/
|
||||
static void ParseHostPort(char* hoststr, char* portstr, HostPort** hostPortList)
|
||||
{
|
||||
char* hosts[MAX_REPLNODE_NUM] = {NULL};
|
||||
char* ports[MAX_REPLNODE_NUM] = {NULL};
|
||||
int hostNum = HostsPortsToArray(hoststr, hosts);
|
||||
int portNum = HostsPortsToArray(portstr, ports);
|
||||
if (hostNum != portNum) {
|
||||
ereport(ERROR, (errcode(ERRCODE_SYNTAX_ERROR), errmsg("The number of host and port are inconsistent.")));
|
||||
}
|
||||
|
||||
for (int i = 0; i < hostNum; ++i) {
|
||||
hostPortList[i] = (HostPort*)palloc(sizeof(HostPort));
|
||||
hostPortList[i]->host = hosts[i];
|
||||
hostPortList[i]->port = ports[i];
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* Parse conninfo
|
||||
* conninfo format:
|
||||
* 'dbname=abc user=username password=xxxx host=ip1,ip2,...,ip9 port=p1,p2,...,p9'
|
||||
* after parsing:
|
||||
* conninfoWithoutHostPort:
|
||||
* 'dbname=abc user=username password=xxxx'
|
||||
* hostPortList:
|
||||
* {
|
||||
* {host=ip1, port=p1},
|
||||
* {host=ip2, port=p2},
|
||||
* ...
|
||||
* {host=ip9, port=p9}
|
||||
* }
|
||||
*/
|
||||
void ParseConninfo(const char* conninfo, StringInfoData* conninfoWithoutHostPort, HostPort** hostPortList)
|
||||
{
|
||||
List* conninfoList = ConninfoToDefList(conninfo);
|
||||
ListCell* l = NULL;
|
||||
|
||||
char* hostStr = NULL;
|
||||
char* portStr = NULL;
|
||||
foreach (l, conninfoList) {
|
||||
DefElem* defel = (DefElem*)lfirst(l);
|
||||
if (pg_strcasecmp(defel->defname, "host") == 0) {
|
||||
hostStr = defGetString(defel);
|
||||
} else if (pg_strcasecmp(defel->defname, "port") == 0) {
|
||||
portStr = defGetString(defel);
|
||||
} else {
|
||||
appendStringInfo(conninfoWithoutHostPort, "%s=%s ", defel->defname, defGetString(defel));
|
||||
}
|
||||
}
|
||||
if (hostPortList != NULL) {
|
||||
ParseHostPort(hostStr, portStr, hostPortList);
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* encrypt conninfo when action = 'E'
|
||||
* decrypt conninfo when action = 'D'
|
||||
* conninfoNew: encrypted or decrypted conninfo
|
||||
*/
|
||||
char* EncryptOrDecryptConninfo(const char* conninfo, const char action)
|
||||
{
|
||||
/* parse conninfo to list */
|
||||
List *conninfoList = ConninfoToDefList(conninfo);
|
||||
/* Sensitive options for subscription */
|
||||
const char* sensitiveOptionsArray[] = {"password"};
|
||||
const int sensitiveArrayLength = lengthof(sensitiveOptionsArray);
|
||||
switch (action) {
|
||||
/* Encrypt */
|
||||
case 'E':
|
||||
EncryptGenericOptions(conninfoList, sensitiveOptionsArray, sensitiveArrayLength, SUBSCRIPTION_MODE);
|
||||
break;
|
||||
|
||||
/* Decrypt */
|
||||
case 'D':
|
||||
DecryptOptions(conninfoList, sensitiveOptionsArray, sensitiveArrayLength, SUBSCRIPTION_MODE);
|
||||
break;
|
||||
|
||||
default:
|
||||
break;
|
||||
}
|
||||
|
||||
char* conninfoNew = DefListToString(conninfoList);
|
||||
ClearListContent(conninfoList);
|
||||
list_free_ext(conninfoList);
|
||||
|
||||
return conninfoNew;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -827,9 +827,6 @@ void InitBSqlPluginHookIfNeeded()
|
|||
{
|
||||
const char* b_sql_plugin = "b_sql_plugin";
|
||||
CFunInfo tmpCF;
|
||||
if (!CheckIfExtensionExists(b_sql_plugin)) {
|
||||
return;
|
||||
}
|
||||
|
||||
tmpCF = load_external_function(b_sql_plugin, INIT_PLUGIN_OBJECT, false, false);
|
||||
if (tmpCF.user_fn != NULL) {
|
||||
|
|
@ -865,9 +862,11 @@ List* pg_parse_query(const char* query_string, List** query_string_locationlist)
|
|||
|
||||
List* (*parser_hook)(const char*, List**) = raw_parser;
|
||||
#ifndef ENABLE_MULTIPLE_NODES
|
||||
int id = GetCustomParserId();
|
||||
if (id >= 0 && g_instance.raw_parser_hook[id] != NULL) {
|
||||
parser_hook = (List* (*)(const char*, List**))g_instance.raw_parser_hook[id];
|
||||
if (u_sess->attr.attr_sql.b_sql_plugin) {
|
||||
int id = GetCustomParserId();
|
||||
if (id >= 0 && g_instance.raw_parser_hook[id] != NULL) {
|
||||
parser_hook = (List* (*)(const char*, List**))g_instance.raw_parser_hook[id];
|
||||
}
|
||||
}
|
||||
#endif
|
||||
raw_parsetree_list = parser_hook(query_string, query_string_locationlist);
|
||||
|
|
@ -6115,6 +6114,9 @@ void ProcessInterrupts(void)
|
|||
|
||||
/* The logical replication launcher can be stopped at any time. */
|
||||
proc_exit(0);
|
||||
} else if (IsLogicalWorker()) {
|
||||
ereport(FATAL, (errcode(ERRCODE_ADMIN_SHUTDOWN),
|
||||
errmsg("terminating logical replication worker due to administrator command")));
|
||||
#endif
|
||||
} else if (IsTxnSnapCapturerProcess()) {
|
||||
ereport(FATAL,
|
||||
|
|
@ -7572,7 +7574,7 @@ int PostgresMain(int argc, char* argv[], const char* dbname, const char* usernam
|
|||
init_set_params_htab();
|
||||
|
||||
#ifndef ENABLE_MULTIPLE_NODES
|
||||
if (u_sess->proc_cxt.MyDatabaseId != InvalidOid && DB_IS_CMPT(B_FORMAT)) {
|
||||
if (u_sess->proc_cxt.MyDatabaseId != InvalidOid && DB_IS_CMPT(B_FORMAT) && u_sess->attr.attr_sql.b_sql_plugin) {
|
||||
InitBSqlPluginHookIfNeeded();
|
||||
}
|
||||
#endif
|
||||
|
|
|
|||
|
|
@ -848,7 +848,7 @@ static bool InitSession(knl_session_context* session)
|
|||
t_thrd.proc_cxt.PostInit->InitSession();
|
||||
|
||||
#ifndef ENABLE_MULTIPLE_NODES
|
||||
if (u_sess->proc_cxt.MyDatabaseId != InvalidOid && DB_IS_CMPT(B_FORMAT)) {
|
||||
if (u_sess->proc_cxt.MyDatabaseId != InvalidOid && DB_IS_CMPT(B_FORMAT) && u_sess->attr.attr_sql.b_sql_plugin) {
|
||||
InitBSqlPluginHookIfNeeded();
|
||||
}
|
||||
#endif
|
||||
|
|
|
|||
|
|
@ -255,8 +255,14 @@ void StartRemoteStreaming(const LibpqrcvConnectParam *options)
|
|||
stringlist_to_identifierstr(t_thrd.libwalreceiver_cxt.streamConn, options->publicationNames);
|
||||
appendStringInfo(&cmd, ", publication_names %s",
|
||||
PQescapeLiteral(t_thrd.libwalreceiver_cxt.streamConn, pubnames_str, strlen(pubnames_str)));
|
||||
appendStringInfoChar(&cmd, ')');
|
||||
pfree(pubnames_str);
|
||||
|
||||
if (options->binary && PQserverVersion(t_thrd.libwalreceiver_cxt.streamConn) >= 90204) {
|
||||
appendStringInfoString(&cmd, ", binary 'true'");
|
||||
ereport(DEBUG5, (errmsg("append binary true")));
|
||||
}
|
||||
|
||||
appendStringInfoChar(&cmd, ')');
|
||||
}
|
||||
|
||||
PGresult *res = libpqrcv_PQexec(cmd.data);
|
||||
|
|
@ -420,7 +426,7 @@ void IdentifyRemoteSystem(bool checkRemote)
|
|||
}
|
||||
|
||||
/* identify remote mode, should do this after connect success. */
|
||||
static ServerMode IdentifyRemoteMode()
|
||||
ServerMode IdentifyRemoteMode()
|
||||
{
|
||||
Assert(t_thrd.libwalreceiver_cxt.streamConn != NULL);
|
||||
volatile WalRcvData *walrcv = t_thrd.walreceiverfuncs_cxt.WalRcv;
|
||||
|
|
@ -443,7 +449,9 @@ static ServerMode IdentifyRemoteMode()
|
|||
num_fields)));
|
||||
}
|
||||
remoteMode = (ServerMode)pg_strtoint32(PQgetvalue(res, 0, 0));
|
||||
if (!t_thrd.walreceiver_cxt.AmWalReceiverForFailover && (!IS_PRIMARY_NORMAL(remoteMode)) &&
|
||||
if (walrcv->conn_target != REPCONNTARGET_PUBLICATION &&
|
||||
!t_thrd.walreceiver_cxt.AmWalReceiverForFailover &&
|
||||
(!IS_PRIMARY_NORMAL(remoteMode)) &&
|
||||
/* remoteMode of cascade standby is a standby */
|
||||
!t_thrd.xlog_cxt.is_cascade_standby && !IS_SHARED_STORAGE_MODE) {
|
||||
PQclear(res);
|
||||
|
|
|
|||
|
|
@ -18,7 +18,6 @@
|
|||
#include "catalog/pg_type.h"
|
||||
#include "libpq/pqformat.h"
|
||||
#include "replication/logicalproto.h"
|
||||
#include "utils/builtins.h"
|
||||
#include "utils/lsyscache.h"
|
||||
#include "utils/syscache.h"
|
||||
|
||||
|
|
@ -28,7 +27,7 @@
|
|||
static const int LOGICALREP_IS_REPLICA_IDENTITY = 1;
|
||||
|
||||
static void logicalrep_write_attrs(StringInfo out, Relation rel);
|
||||
static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple);
|
||||
static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple, bool binary);
|
||||
|
||||
static void logicalrep_read_attrs(StringInfo in, LogicalRepRelation *rel);
|
||||
static void logicalrep_read_tuple(StringInfo in, LogicalRepTupleData *tuple);
|
||||
|
|
@ -115,7 +114,7 @@ void logicalrep_write_origin(StringInfo out, const char *origin, XLogRecPtr orig
|
|||
/*
|
||||
* Write INSERT to the output stream.
|
||||
*/
|
||||
void logicalrep_write_insert(StringInfo out, Relation rel, HeapTuple newtuple)
|
||||
void logicalrep_write_insert(StringInfo out, Relation rel, HeapTuple newtuple, bool binary)
|
||||
{
|
||||
pq_sendbyte(out, 'I'); /* action INSERT */
|
||||
|
||||
|
|
@ -123,7 +122,7 @@ void logicalrep_write_insert(StringInfo out, Relation rel, HeapTuple newtuple)
|
|||
pq_sendint32(out, RelationGetRelid(rel));
|
||||
|
||||
pq_sendbyte(out, 'N'); /* new tuple follows */
|
||||
logicalrep_write_tuple(out, rel, newtuple);
|
||||
logicalrep_write_tuple(out, rel, newtuple, binary);
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
@ -151,7 +150,7 @@ LogicalRepRelId logicalrep_read_insert(StringInfo in, LogicalRepTupleData *newtu
|
|||
/*
|
||||
* Write UPDATE to the output stream.
|
||||
*/
|
||||
void logicalrep_write_update(StringInfo out, Relation rel, HeapTuple oldtuple, HeapTuple newtuple)
|
||||
void logicalrep_write_update(StringInfo out, Relation rel, HeapTuple oldtuple, HeapTuple newtuple, bool binary)
|
||||
{
|
||||
pq_sendbyte(out, 'U'); /* action UPDATE */
|
||||
|
||||
|
|
@ -166,11 +165,11 @@ void logicalrep_write_update(StringInfo out, Relation rel, HeapTuple oldtuple, H
|
|||
pq_sendbyte(out, 'O'); /* old tuple follows */
|
||||
else
|
||||
pq_sendbyte(out, 'K'); /* old key follows */
|
||||
logicalrep_write_tuple(out, rel, oldtuple);
|
||||
logicalrep_write_tuple(out, rel, oldtuple, binary);
|
||||
}
|
||||
|
||||
pq_sendbyte(out, 'N'); /* new tuple follows */
|
||||
logicalrep_write_tuple(out, rel, newtuple);
|
||||
logicalrep_write_tuple(out, rel, newtuple, binary);
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
@ -213,7 +212,7 @@ LogicalRepRelId logicalrep_read_update(StringInfo in, bool *has_oldtuple, Logica
|
|||
/*
|
||||
* Write DELETE to the output stream.
|
||||
*/
|
||||
void logicalrep_write_delete(StringInfo out, Relation rel, HeapTuple oldtuple)
|
||||
void logicalrep_write_delete(StringInfo out, Relation rel, HeapTuple oldtuple, bool binary)
|
||||
{
|
||||
char relreplident = RelationGetRelReplident(rel);
|
||||
Assert(relreplident == REPLICA_IDENTITY_DEFAULT ||
|
||||
|
|
@ -229,7 +228,7 @@ void logicalrep_write_delete(StringInfo out, Relation rel, HeapTuple oldtuple)
|
|||
else
|
||||
pq_sendbyte(out, 'K'); /* old key follows */
|
||||
|
||||
logicalrep_write_tuple(out, rel, oldtuple);
|
||||
logicalrep_write_tuple(out, rel, oldtuple, binary);
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
@ -344,7 +343,7 @@ void logicalrep_read_typ(StringInfo in, LogicalRepTyp *ltyp)
|
|||
/*
|
||||
* Write a tuple to the outputstream, in the most efficient format possible.
|
||||
*/
|
||||
static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple)
|
||||
static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple, bool binary)
|
||||
{
|
||||
TupleDesc desc;
|
||||
Datum values[MaxTupleAttributeNumber];
|
||||
|
|
@ -371,7 +370,6 @@ static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple
|
|||
HeapTuple typtup;
|
||||
Form_pg_type typclass;
|
||||
Form_pg_attribute att = desc->attrs[i];
|
||||
char *outputstr;
|
||||
|
||||
/* skip dropped columns */
|
||||
if (att->attisdropped || GetGeneratedCol(desc, i)) {
|
||||
|
|
@ -379,7 +377,7 @@ static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple
|
|||
}
|
||||
|
||||
if (isnull[i]) {
|
||||
pq_sendbyte(out, 'n'); /* null column */
|
||||
pq_sendbyte(out, LOGICALREP_COLUMN_NULL); /* null column */
|
||||
continue;
|
||||
}
|
||||
|
||||
|
|
@ -388,61 +386,91 @@ static void logicalrep_write_tuple(StringInfo out, Relation rel, HeapTuple tuple
|
|||
elog(ERROR, "cache lookup failed for type %u", att->atttypid);
|
||||
typclass = (Form_pg_type)GETSTRUCT(typtup);
|
||||
|
||||
pq_sendbyte(out, 't'); /* 'text' data follows */
|
||||
if (!typclass->typbyval && typclass->typlen == -1) {
|
||||
/* definitely detoasted Datum */
|
||||
Datum val = PointerGetDatum(PG_DETOAST_DATUM(values[i]));
|
||||
outputstr = OidOutputFunctionCall(typclass->typoutput, val);
|
||||
/*
|
||||
* Send in binary if requested and type has suitable send function.
|
||||
*/
|
||||
if (binary && OidIsValid(typclass->typsend)) {
|
||||
bytea* outputbytes = NULL;
|
||||
pq_sendbyte(out, LOGICALREP_COLUMN_BINARY);
|
||||
if (!typclass->typbyval && typclass->typlen == -1) {
|
||||
/* definitely detoasted Datum */
|
||||
Datum val = PointerGetDatum(PG_DETOAST_DATUM(values[i]));
|
||||
outputbytes = OidSendFunctionCall(typclass->typsend, val);
|
||||
} else {
|
||||
outputbytes = OidSendFunctionCall(typclass->typsend, values[i]);
|
||||
}
|
||||
int len = VARSIZE(outputbytes) - VARHDRSZ;
|
||||
pq_sendint(out, len, 4); /* length */
|
||||
pq_sendbytes(out, VARDATA(outputbytes), len); /* data */
|
||||
if (outputbytes != NULL) {
|
||||
pfree(outputbytes);
|
||||
}
|
||||
} else {
|
||||
outputstr = OidOutputFunctionCall(typclass->typoutput, values[i]);
|
||||
char* outputstr = NULL;
|
||||
pq_sendbyte(out, LOGICALREP_COLUMN_TEXT);
|
||||
if (!typclass->typbyval && typclass->typlen == -1) {
|
||||
/* definitely detoasted Datum */
|
||||
Datum val = PointerGetDatum(PG_DETOAST_DATUM(values[i]));
|
||||
outputstr = OidOutputFunctionCall(typclass->typoutput, val);
|
||||
} else {
|
||||
outputstr = OidOutputFunctionCall(typclass->typoutput, values[i]);
|
||||
}
|
||||
pq_sendcountedtext(out, outputstr, strlen(outputstr), false);
|
||||
if (outputstr != NULL) {
|
||||
pfree(outputstr);
|
||||
}
|
||||
}
|
||||
pq_sendcountedtext(out, outputstr, strlen(outputstr), false);
|
||||
pfree(outputstr);
|
||||
ReleaseSysCache(typtup);
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
* Read tuple in remote format from stream.
|
||||
*
|
||||
* The returned tuple points into the input stringinfo.
|
||||
* Read tuple in logical replication format from stream.
|
||||
*/
|
||||
static void logicalrep_read_tuple(StringInfo in, LogicalRepTupleData *tuple)
|
||||
{
|
||||
uint16 i;
|
||||
uint16 natts;
|
||||
int rc;
|
||||
|
||||
/* Get number of attributes. */
|
||||
natts = pq_getmsgint(in, sizeof(uint16));
|
||||
uint16 natts = pq_getmsgint(in, sizeof(uint16));
|
||||
|
||||
rc = memset_s(tuple->changed, sizeof(tuple->changed), 0, sizeof(tuple->changed));
|
||||
securec_check(rc, "", "");
|
||||
/* Allocate space for per-column values; zero out unused StringInfoDatas */
|
||||
tuple->colvalues = (StringInfoData *) palloc0(natts * sizeof(StringInfoData));
|
||||
tuple->colstatus = (char *) palloc(natts * sizeof(char));
|
||||
tuple->ncols = natts;
|
||||
|
||||
/* Read the data */
|
||||
for (i = 0; i < natts; i++) {
|
||||
char kind;
|
||||
|
||||
kind = pq_getmsgbyte(in);
|
||||
for (uint16 i = 0; i < natts; i++) {
|
||||
char kind = pq_getmsgbyte(in);
|
||||
tuple->colstatus[i] = kind;
|
||||
uint32 len;
|
||||
StringInfo value = &tuple->colvalues[i];
|
||||
|
||||
switch (kind) {
|
||||
case 'n': /* null */
|
||||
tuple->values[i] = NULL;
|
||||
tuple->changed[i] = true;
|
||||
case LOGICALREP_COLUMN_NULL: /* null */
|
||||
/* nothing more to do */
|
||||
break;
|
||||
case 't': { /* text formatted value */
|
||||
uint32 len;
|
||||
tuple->changed[i] = true;
|
||||
|
||||
len = pq_getmsgint(in, sizeof(uint32)); /* read length */
|
||||
case LOGICALREP_COLUMN_TEXT:
|
||||
len = pq_getmsgint(in, sizeof(uint32)); /* read length */
|
||||
|
||||
/* and data */
|
||||
tuple->values[i] = (char *)palloc(len + 1);
|
||||
pq_copymsgbytes(in, tuple->values[i], len);
|
||||
tuple->values[i][len] = '\0';
|
||||
value->data = (char *) palloc((len + 1) * sizeof(char));
|
||||
pq_copymsgbytes(in, value->data, len);
|
||||
value->data[len] = '\0';
|
||||
/* make StringInfo fully valid */
|
||||
value->len = len;
|
||||
value->cursor = 0;
|
||||
value->maxlen = len;
|
||||
break;
|
||||
case LOGICALREP_COLUMN_BINARY:
|
||||
len = pq_getmsgint(in, sizeof(uint32)); /* read length */
|
||||
|
||||
/* and data */
|
||||
value->data = (char *)palloc0((len + 1) * sizeof(char));
|
||||
pq_copymsgbytes(in, value->data, len);
|
||||
/* make StringInfo fully valid */
|
||||
value->len = len;
|
||||
value->cursor = 0;
|
||||
value->maxlen = len;
|
||||
break;
|
||||
}
|
||||
default:
|
||||
elog(ERROR, "unrecognized data representation type '%c'", kind);
|
||||
break;
|
||||
|
|
@ -451,7 +479,7 @@ static void logicalrep_read_tuple(StringInfo in, LogicalRepTupleData *tuple)
|
|||
}
|
||||
|
||||
/*
|
||||
* Write relation attributes to the stream.
|
||||
* Write relation attributes metadata to the stream.
|
||||
*/
|
||||
static void logicalrep_write_attrs(StringInfo out, Relation rel)
|
||||
{
|
||||
|
|
@ -504,7 +532,7 @@ static void logicalrep_write_attrs(StringInfo out, Relation rel)
|
|||
}
|
||||
|
||||
/*
|
||||
* Read relation attribute names from the stream.
|
||||
* Read relation attribute metadata from the stream.
|
||||
*/
|
||||
static void logicalrep_read_attrs(StringInfo in, LogicalRepRelation *rel)
|
||||
{
|
||||
|
|
@ -573,3 +601,25 @@ static const char *logicalrep_read_namespace(StringInfo in)
|
|||
|
||||
return nspname;
|
||||
}
|
||||
|
||||
/*
|
||||
* Write conninfo to the output stream.
|
||||
*/
|
||||
void logicalrep_write_conninfo(StringInfo out, char* conninfo)
|
||||
{
|
||||
pq_sendbyte(out, 'S'); /* action */
|
||||
|
||||
pq_writestring(out, conninfo); /* conninfo follows */
|
||||
}
|
||||
|
||||
/*
|
||||
* Read conninfo from stream.
|
||||
*/
|
||||
void logicalrep_read_conninfo(StringInfo in, char** conninfo)
|
||||
{
|
||||
const char* conninfoTemp = pq_getmsgstring(in);
|
||||
size_t conninfoLen = strlen(conninfoTemp) + 1;
|
||||
*conninfo = (char*)palloc(conninfoLen);
|
||||
int rc = strcpy_s(*conninfo, conninfoLen, conninfoTemp);
|
||||
securec_check(rc, "", "");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -41,6 +41,7 @@
|
|||
#include "catalog/pg_partition_fn.h"
|
||||
|
||||
#include "commands/trigger.h"
|
||||
#include "commands/subscriptioncmds.h"
|
||||
|
||||
#include "executor/executor.h"
|
||||
#include "executor/node/nodeModifyTable.h"
|
||||
|
|
@ -111,6 +112,8 @@ static void store_flush_position(XLogRecPtr remote_lsn);
|
|||
static void reread_subscription(void);
|
||||
static void ApplyWorkerProcessMsg(char type, StringInfo s, XLogRecPtr *lastRcv);
|
||||
static void apply_dispatch(StringInfo s);
|
||||
static void apply_handle_conninfo(StringInfo s);
|
||||
static void UpdateConninfo(char* standbysInfo);
|
||||
|
||||
/* SIGHUP: set flag to re-read config file at next convenient time */
|
||||
static void LogicalrepWorkerSighub(SIGNAL_ARGS)
|
||||
|
|
@ -118,18 +121,6 @@ static void LogicalrepWorkerSighub(SIGNAL_ARGS)
|
|||
t_thrd.applyworker_cxt.got_SIGHUP = true;
|
||||
}
|
||||
|
||||
/* SIGTERM: time to die */
|
||||
static void LogicalrepWorkerSigterm(SIGNAL_ARGS)
|
||||
{
|
||||
int saveErrno = errno;
|
||||
|
||||
t_thrd.applyworker_cxt.got_SIGTERM = true;
|
||||
if (t_thrd.proc)
|
||||
SetLatch(&t_thrd.proc->procLatch);
|
||||
|
||||
errno = saveErrno;
|
||||
}
|
||||
|
||||
/*
|
||||
* Make sure that we started local transaction.
|
||||
*
|
||||
|
|
@ -265,11 +256,11 @@ static void slot_store_error_callback(void *arg)
|
|||
}
|
||||
|
||||
/*
|
||||
* Store data in C string form into slot.
|
||||
* This is similar to BuildTupleFromCStrings but TupleTableSlot fits our
|
||||
* use better.
|
||||
* Store tuple data into slot.
|
||||
*
|
||||
* Incoming data can be either text or binary format.
|
||||
*/
|
||||
static void slot_store_cstrings(TupleTableSlot *slot, LogicalRepRelMapEntry *rel, char **values)
|
||||
static void slot_store_data(TupleTableSlot *slot, LogicalRepRelMapEntry *rel, LogicalRepTupleData *tupleData)
|
||||
{
|
||||
int natts = slot->tts_tupleDescriptor->natts;
|
||||
int i;
|
||||
|
|
@ -286,19 +277,52 @@ static void slot_store_cstrings(TupleTableSlot *slot, LogicalRepRelMapEntry *rel
|
|||
errcallback.previous = t_thrd.log_cxt.error_context_stack;
|
||||
t_thrd.log_cxt.error_context_stack = &errcallback;
|
||||
|
||||
/* Call the "in" function for each non-dropped attribute */
|
||||
/* Call the "in" function for each non-dropped, non-null attribute */
|
||||
for (i = 0; i < natts; i++) {
|
||||
Form_pg_attribute att = slot->tts_tupleDescriptor->attrs[i];
|
||||
int remoteattnum = rel->attrmap[i];
|
||||
Oid typinput;
|
||||
Oid typioparam;
|
||||
|
||||
if (!att->attisdropped && remoteattnum >= 0 && values[remoteattnum] != NULL) {
|
||||
if (!att->attisdropped && remoteattnum >= 0) {
|
||||
StringInfo colvalue = &tupleData->colvalues[remoteattnum];
|
||||
errarg.remote_attnum = remoteattnum;
|
||||
|
||||
getTypeInputInfo(att->atttypid, &typinput, &typioparam);
|
||||
slot->tts_values[i] = OidInputFunctionCall(typinput, values[remoteattnum], typioparam, att->atttypmod);
|
||||
slot->tts_isnull[i] = false;
|
||||
if (tupleData->colstatus[remoteattnum] == LOGICALREP_COLUMN_TEXT) {
|
||||
Oid typinput;
|
||||
Oid typioparam;
|
||||
|
||||
getTypeInputInfo(att->atttypid, &typinput, &typioparam);
|
||||
slot->tts_values[i] = OidInputFunctionCall(typinput, colvalue->data, typioparam, att->atttypmod);
|
||||
slot->tts_isnull[i] = false;
|
||||
} else if (tupleData->colstatus[remoteattnum] == LOGICALREP_COLUMN_BINARY) {
|
||||
Oid typreceive;
|
||||
Oid typioparam;
|
||||
|
||||
/*
|
||||
* In some code paths we may be asked to re-parse the same
|
||||
* tuple data. Reset the StringInfo's cursor so that works.
|
||||
*/
|
||||
colvalue->cursor = 0;
|
||||
|
||||
getTypeBinaryInputInfo(att->atttypid, &typreceive, &typioparam);
|
||||
slot->tts_values[i] = OidReceiveFunctionCall(typreceive, colvalue, typioparam, att->atttypmod);
|
||||
|
||||
/* Trouble if it didn't eat the whole buffer */
|
||||
if (colvalue->cursor != colvalue->len) {
|
||||
ereport(ERROR, (errcode(ERRCODE_INVALID_BINARY_REPRESENTATION),
|
||||
errmsg("incorrect binary data format in logical replication column %d",
|
||||
remoteattnum + 1)));
|
||||
}
|
||||
slot->tts_isnull[i] = false;
|
||||
} else {
|
||||
/*
|
||||
* NULL value from remote. (We don't expect to see
|
||||
* LOGICALREP_COLUMN_UNCHANGED here, but if we do, treat it as
|
||||
* NULL.)
|
||||
*/
|
||||
slot->tts_values[i] = (Datum) 0;
|
||||
slot->tts_isnull[i] = true;
|
||||
}
|
||||
/* Reset attnum for error callback */
|
||||
errarg.remote_attnum = -1;
|
||||
} else {
|
||||
/*
|
||||
|
|
@ -318,18 +342,19 @@ static void slot_store_cstrings(TupleTableSlot *slot, LogicalRepRelMapEntry *rel
|
|||
}
|
||||
|
||||
/*
|
||||
* Replace selected columns with user data provided as C strings.
|
||||
* Replace updated columns with data from the LogicalRepTupleData struct.
|
||||
* This is somewhat similar to heap_modify_tuple but also calls the type
|
||||
* input functions on the user data.
|
||||
* "slot" is filled with a copy of the tuple in "srcslot", with
|
||||
* columns selected by the "replaces" array replaced with data values
|
||||
* from "values".
|
||||
*
|
||||
* "slot" is filled with a copy of the tuple in "srcslot", replacing
|
||||
* columns provided in "tupleData" and leaving others as-is.
|
||||
*
|
||||
* Caution: unreplaced pass-by-ref columns in "slot" will point into the
|
||||
* storage for "srcslot". This is OK for current usage, but someday we may
|
||||
* need to materialize "slot" at the end to make it independent of "srcslot".
|
||||
*/
|
||||
static void slot_modify_cstrings(TupleTableSlot *slot, TupleTableSlot *srcslot, LogicalRepRelMapEntry *rel,
|
||||
char **values, const bool *replaces)
|
||||
static void slot_modify_data(TupleTableSlot *slot, TupleTableSlot *srcslot, LogicalRepRelMapEntry *rel,
|
||||
LogicalRepTupleData *tupleData)
|
||||
{
|
||||
int natts = slot->tts_tupleDescriptor->natts;
|
||||
int i;
|
||||
|
|
@ -364,23 +389,47 @@ static void slot_modify_cstrings(TupleTableSlot *slot, TupleTableSlot *srcslot,
|
|||
Form_pg_attribute att = slot->tts_tupleDescriptor->attrs[i];
|
||||
int remoteattnum = rel->attrmap[i];
|
||||
|
||||
if (remoteattnum < 0 || !replaces[remoteattnum]) {
|
||||
if (remoteattnum < 0) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (values[remoteattnum] != NULL) {
|
||||
Oid typinput;
|
||||
Oid typioparam;
|
||||
|
||||
if (tupleData->colstatus[remoteattnum] != LOGICALREP_COLUMN_UNCHANGED) {
|
||||
StringInfo colvalue = &tupleData->colvalues[remoteattnum];
|
||||
errarg.remote_attnum = remoteattnum;
|
||||
|
||||
getTypeInputInfo(att->atttypid, &typinput, &typioparam);
|
||||
slot->tts_values[i] = OidInputFunctionCall(typinput, values[remoteattnum], typioparam, att->atttypmod);
|
||||
slot->tts_isnull[i] = false;
|
||||
if (tupleData->colstatus[remoteattnum] == LOGICALREP_COLUMN_TEXT) {
|
||||
Oid typinput;
|
||||
Oid typioparam;
|
||||
|
||||
getTypeInputInfo(att->atttypid, &typinput, &typioparam);
|
||||
slot->tts_values[i] = OidInputFunctionCall(typinput, colvalue->data, typioparam, att->atttypmod);
|
||||
slot->tts_isnull[i] = false;
|
||||
} else if (tupleData->colstatus[remoteattnum] == LOGICALREP_COLUMN_BINARY) {
|
||||
Oid typreceive;
|
||||
Oid typioparam;
|
||||
|
||||
/*
|
||||
* In some code paths we may be asked to re-parse the same
|
||||
* tuple data. Reset the StringInfo's cursor so that works.
|
||||
*/
|
||||
colvalue->cursor = 0;
|
||||
|
||||
getTypeBinaryInputInfo(att->atttypid, &typreceive, &typioparam);
|
||||
slot->tts_values[i] = OidReceiveFunctionCall(typreceive, colvalue, typioparam, att->atttypmod);
|
||||
|
||||
/* Trouble if it didn't eat the whole buffer */
|
||||
if (colvalue->cursor != colvalue->len) {
|
||||
ereport(ERROR, (errcode(ERRCODE_INVALID_BINARY_REPRESENTATION),
|
||||
errmsg("incorrect binary data format in logical replication column %d", remoteattnum + 1)));
|
||||
}
|
||||
slot->tts_isnull[i] = false;
|
||||
} else {
|
||||
/* must be LOGICALREP_COLUMN_NULL */
|
||||
slot->tts_values[i] = (Datum) 0;
|
||||
slot->tts_isnull[i] = true;
|
||||
}
|
||||
|
||||
errarg.remote_attnum = -1;
|
||||
} else {
|
||||
slot->tts_values[i] = (Datum)0;
|
||||
slot->tts_isnull[i] = true;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -516,7 +565,7 @@ static void apply_handle_insert(StringInfo s)
|
|||
PushActiveSnapshot(GetTransactionSnapshot());
|
||||
/* Process and store remote tuple in the slot */
|
||||
oldctx = MemoryContextSwitchTo(GetPerTupleMemoryContext(estate));
|
||||
slot_store_cstrings(remoteslot, rel, newtup.values);
|
||||
slot_store_data(remoteslot, rel, &newtup);
|
||||
slot_fill_defaults(rel, estate, remoteslot);
|
||||
MemoryContextSwitchTo(oldctx);
|
||||
|
||||
|
|
@ -646,7 +695,7 @@ static void apply_handle_update(StringInfo s)
|
|||
int remoteattnum = rel->attrmap[i];
|
||||
if (!att->attisdropped && remoteattnum >= 0) {
|
||||
Assert(remoteattnum < newtup.ncols);
|
||||
if (newtup.changed[i]) {
|
||||
if (newtup.colstatus[i] != LOGICALREP_COLUMN_UNCHANGED) {
|
||||
target_rte->updatedCols = bms_add_member(target_rte->updatedCols,
|
||||
i + 1 - FirstLowInvalidHeapAttributeNumber);
|
||||
}
|
||||
|
|
@ -660,7 +709,7 @@ static void apply_handle_update(StringInfo s)
|
|||
|
||||
/* Build the search tuple. */
|
||||
oldctx = MemoryContextSwitchTo(GetPerTupleMemoryContext(estate));
|
||||
slot_store_cstrings(remoteslot, rel, has_oldtup ? oldtup.values : newtup.values);
|
||||
slot_store_data(remoteslot, rel, has_oldtup ? &oldtup : &newtup);
|
||||
MemoryContextSwitchTo(oldctx);
|
||||
|
||||
/*
|
||||
|
|
@ -683,7 +732,7 @@ static void apply_handle_update(StringInfo s)
|
|||
if (found) {
|
||||
/* Process and store remote tuple in the slot */
|
||||
oldctx = MemoryContextSwitchTo(GetPerTupleMemoryContext(estate));
|
||||
slot_modify_cstrings(remoteslot, localslot, rel, newtup.values, newtup.changed);
|
||||
slot_modify_data(remoteslot, localslot, rel, &newtup);
|
||||
MemoryContextSwitchTo(oldctx);
|
||||
|
||||
EvalPlanQualSetSlot(&epqstate, remoteslot);
|
||||
|
|
@ -748,7 +797,7 @@ static void apply_handle_delete(StringInfo s)
|
|||
|
||||
/* Find the tuple using the replica identity index. */
|
||||
oldctx = MemoryContextSwitchTo(GetPerTupleMemoryContext(estate));
|
||||
slot_store_cstrings(remoteslot, rel, oldtup.values);
|
||||
slot_store_data(remoteslot, rel, &oldtup);
|
||||
MemoryContextSwitchTo(oldctx);
|
||||
|
||||
/*
|
||||
|
|
@ -824,6 +873,9 @@ static void apply_dispatch(StringInfo s)
|
|||
case 'O':
|
||||
apply_handle_origin(s);
|
||||
break;
|
||||
case 'S':
|
||||
apply_handle_conninfo(s);
|
||||
break;
|
||||
default:
|
||||
ereport(ERROR, (errcode(ERRCODE_PROTOCOL_VIOLATION),
|
||||
errmsg("invalid logical replication message type \"%c\"", action)));
|
||||
|
|
@ -1004,13 +1056,15 @@ static void ApplyLoop(void)
|
|||
/* mark as idle, before starting to loop */
|
||||
pgstat_report_activity(STATE_IDLE, NULL);
|
||||
|
||||
while (!t_thrd.applyworker_cxt.got_SIGTERM) {
|
||||
for (;;) {
|
||||
MemoryContextSwitchTo(t_thrd.applyworker_cxt.messageContext);
|
||||
|
||||
int len;
|
||||
char *buf = NULL;
|
||||
unsigned char type;
|
||||
|
||||
CHECK_FOR_INTERRUPTS();
|
||||
|
||||
/* Wait a while for data to arrive */
|
||||
if ((WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_receive(NAPTIME_PER_CYCLE, &type, &buf, &len)) {
|
||||
StringInfoData s;
|
||||
|
|
@ -1224,6 +1278,17 @@ static void reread_subscription(void)
|
|||
proc_exit(0);
|
||||
}
|
||||
|
||||
/*
|
||||
* Exit if any parameter that affects the remote connection was changed.
|
||||
* The launcher will start a new worker.
|
||||
*/
|
||||
if (strcmp(newsub->name, t_thrd.applyworker_cxt.mySubscription->name) != 0 ||
|
||||
newsub->binary != t_thrd.applyworker_cxt.mySubscription->binary) {
|
||||
ereport(LOG, (errmsg("logical replication apply worker for subscription \"%s\" "
|
||||
"will restart because of a parameter change", t_thrd.applyworker_cxt.mySubscription->name)));
|
||||
proc_exit(0);
|
||||
}
|
||||
|
||||
/* !slotname should never happen when enabled is true. */
|
||||
Assert(newsub->slotname);
|
||||
|
||||
|
|
@ -1296,7 +1361,7 @@ void ApplyWorkerMain()
|
|||
*/
|
||||
gspqsignal(SIGHUP, LogicalrepWorkerSighub);
|
||||
gspqsignal(SIGINT, StatementCancelHandler);
|
||||
gspqsignal(SIGTERM, LogicalrepWorkerSigterm);
|
||||
gspqsignal(SIGTERM, die);
|
||||
|
||||
gspqsignal(SIGQUIT, quickdie);
|
||||
gspqsignal(SIGALRM, handle_sig_alarm);
|
||||
|
|
@ -1422,14 +1487,9 @@ void ApplyWorkerMain()
|
|||
|
||||
CommitTransactionCommand();
|
||||
|
||||
char *decryptConnInfo = DecryptConninfo(t_thrd.applyworker_cxt.mySubscription->conninfo);
|
||||
bool connectSuccess = (WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_connect(decryptConnInfo, NULL,
|
||||
t_thrd.applyworker_cxt.mySubscription->name, -1);
|
||||
rc = memset_s(decryptConnInfo, strlen(decryptConnInfo), 0, strlen(decryptConnInfo));
|
||||
securec_check(rc, "", "");
|
||||
pfree_ext(decryptConnInfo);
|
||||
if (!connectSuccess) {
|
||||
ereport(ERROR, (errcode(ERRCODE_CONNECTION_FAILURE), errmsg("could not connect to the publisher")));
|
||||
if (!AttemptConnectPublisher(t_thrd.applyworker_cxt.mySubscription->conninfo,
|
||||
t_thrd.applyworker_cxt.mySubscription->name, true)) {
|
||||
ereport(ERROR, (errcode(ERRCODE_CONNECTION_FAILURE), errmsg("Failed to connect to publisher.")));
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
@ -1447,6 +1507,7 @@ void ApplyWorkerMain()
|
|||
options.slotname = t_thrd.applyworker_cxt.mySubscription->slotname;
|
||||
options.protoVersion = LOGICALREP_PROTO_VERSION_NUM;
|
||||
options.publicationNames = t_thrd.applyworker_cxt.mySubscription->publications;
|
||||
options.binary = t_thrd.applyworker_cxt.mySubscription->binary;
|
||||
|
||||
/* Start streaming from the slot. */
|
||||
(WalReceiverFuncTable[GET_FUNC_IDX]).walrcv_startstreaming(&options);
|
||||
|
|
@ -1595,3 +1656,72 @@ char* DefListToString(const List *defList)
|
|||
return buf.data;
|
||||
}
|
||||
|
||||
/*
|
||||
* Handle conninfo update message.
|
||||
*/
|
||||
static void apply_handle_conninfo(StringInfo s)
|
||||
{
|
||||
char* standbysInfo = NULL;
|
||||
logicalrep_read_conninfo(s, &standbysInfo);
|
||||
UpdateConninfo(standbysInfo);
|
||||
pfree_ext(standbysInfo);
|
||||
}
|
||||
|
||||
static void UpdateConninfo(char* standbysInfo)
|
||||
{
|
||||
Relation rel;
|
||||
bool nulls[Natts_pg_subscription];
|
||||
bool replaces[Natts_pg_subscription];
|
||||
Datum values[Natts_pg_subscription];
|
||||
HeapTuple tup;
|
||||
Subscription* sub = t_thrd.applyworker_cxt.mySubscription;
|
||||
Oid subid = sub->oid;
|
||||
|
||||
StartTransactionCommand();
|
||||
rel = heap_open(SubscriptionRelationId, RowExclusiveLock);
|
||||
/* Fetch the existing tuple. */
|
||||
tup = SearchSysCacheCopy2(SUBSCRIPTIONNAME, u_sess->proc_cxt.MyDatabaseId,
|
||||
CStringGetDatum(t_thrd.applyworker_cxt.mySubscription->name));
|
||||
if (!HeapTupleIsValid(tup)) {
|
||||
ereport(ERROR, (errcode(ERRCODE_UNDEFINED_OBJECT), errmsg("subscription \"%s\" does not exist",
|
||||
t_thrd.applyworker_cxt.mySubscription->name)));
|
||||
}
|
||||
subid = HeapTupleGetOid(tup);
|
||||
|
||||
/* Form a new tuple. */
|
||||
int rc = memset_s(nulls, sizeof(nulls), false, sizeof(nulls));
|
||||
securec_check(rc, "", "");
|
||||
rc = memset_s(values, sizeof(values), 0, sizeof(values));
|
||||
securec_check(rc, "", "");
|
||||
rc = memset_s(replaces, sizeof(replaces), false, sizeof(replaces));
|
||||
securec_check(rc, "", "");
|
||||
|
||||
/* get conninfoWithoutHostport */
|
||||
StringInfoData conninfoWithoutHostport;
|
||||
initStringInfo(&conninfoWithoutHostport);
|
||||
ParseConninfo(sub->conninfo, &conninfoWithoutHostport, (HostPort**)NULL);
|
||||
|
||||
/* join conninfoWithoutHostport together with standbysinfo */
|
||||
appendStringInfo(&conninfoWithoutHostport, " %s", standbysInfo);
|
||||
/* Replace connection information */
|
||||
values[Anum_pg_subscription_subconninfo - 1] = CStringGetTextDatum(conninfoWithoutHostport.data);
|
||||
replaces[Anum_pg_subscription_subconninfo - 1] = true;
|
||||
tup = heap_modify_tuple(tup, RelationGetDescr(rel), values, nulls, replaces);
|
||||
|
||||
/* Update the catalog. */
|
||||
simple_heap_update(rel, &tup->t_self, tup);
|
||||
CatalogUpdateIndexes(rel, tup);
|
||||
|
||||
heap_close(rel, RowExclusiveLock);
|
||||
CommitTransactionCommand();
|
||||
|
||||
ereport(LOG, (errmsg("Update conninfo successfully, new conninfo %s.", standbysInfo)));
|
||||
}
|
||||
|
||||
/*
|
||||
* Is current process a logical replication worker?
|
||||
*/
|
||||
bool IsLogicalWorker(void)
|
||||
{
|
||||
return t_thrd.applyworker_cxt.curWorker != NULL;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -15,6 +15,8 @@
|
|||
|
||||
#include "catalog/pg_publication.h"
|
||||
|
||||
#include "commands/defrem.h"
|
||||
|
||||
#include "replication/logical.h"
|
||||
#include "replication/logicalproto.h"
|
||||
#include "replication/origin.h"
|
||||
|
|
@ -44,6 +46,8 @@ static bool pgoutput_origin_filter(LogicalDecodingContext *ctx, RepOriginId orig
|
|||
|
||||
static List *LoadPublications(List *pubnames);
|
||||
static void publication_invalidation_cb(Datum arg, int cacheid, uint32 hashvalue);
|
||||
static bool ReplconninfoChanged();
|
||||
static void GetConninfo(StringInfoData* standbysInfo);
|
||||
|
||||
/* Entry in the map used to remember which relation schemas we sent. */
|
||||
typedef struct RelationSyncEntry {
|
||||
|
|
@ -79,6 +83,9 @@ static void parse_output_parameters(List *options, PGOutputData *data)
|
|||
ListCell *lc;
|
||||
bool protocol_version_given = false;
|
||||
bool publication_names_given = false;
|
||||
bool binary_option_given = false;
|
||||
|
||||
data->binary = false;
|
||||
|
||||
foreach (lc, options) {
|
||||
DefElem *defel = (DefElem *)lfirst(lc);
|
||||
|
|
@ -108,6 +115,12 @@ static void parse_output_parameters(List *options, PGOutputData *data)
|
|||
|
||||
if (!SplitIdentifierString(strVal(defel->arg), ',', &(data->publication_names)))
|
||||
ereport(ERROR, (errcode(ERRCODE_INVALID_NAME), errmsg("invalid publication_names syntax")));
|
||||
} else if (strcmp(defel->defname, "binary") == 0) {
|
||||
if (binary_option_given)
|
||||
ereport(ERROR, (errcode(ERRCODE_SYNTAX_ERROR), errmsg("conflicting or redundant options")));
|
||||
binary_option_given = true;
|
||||
|
||||
data->binary = defGetBoolean(defel);
|
||||
} else
|
||||
elog(ERROR, "unrecognized pgoutput option: %s", defel->defname);
|
||||
}
|
||||
|
|
@ -206,6 +219,22 @@ static void pgoutput_commit_txn(LogicalDecodingContext *ctx, ReorderBufferTXN *t
|
|||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_commit(ctx->out, txn, commit_lsn);
|
||||
OutputPluginWrite(ctx, true);
|
||||
|
||||
/*
|
||||
* Send the newest connecttion information to the subscriber,
|
||||
* when the connection information about the standby changes.
|
||||
*/
|
||||
if (ReplconninfoChanged()) {
|
||||
StringInfoData standbysInfo;
|
||||
initStringInfo(&standbysInfo);
|
||||
|
||||
GetConninfo(&standbysInfo);
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_conninfo(ctx->out, standbysInfo.data);
|
||||
OutputPluginWrite(ctx, true);
|
||||
|
||||
FreeStringInfo(&standbysInfo);
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
|
|
@ -300,19 +329,19 @@ static void pgoutput_change(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
|
|||
switch (change->action) {
|
||||
case REORDER_BUFFER_CHANGE_INSERT:
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_insert(ctx->out, relation, &change->data.tp.newtuple->tuple);
|
||||
logicalrep_write_insert(ctx->out, relation, &change->data.tp.newtuple->tuple, data->binary);
|
||||
OutputPluginWrite(ctx, true);
|
||||
break;
|
||||
case REORDER_BUFFER_CHANGE_UINSERT:
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_insert(ctx->out, relation, (HeapTuple)(&change->data.utp.newtuple->tuple));
|
||||
logicalrep_write_insert(ctx->out, relation, (HeapTuple)(&change->data.utp.newtuple->tuple), data->binary);
|
||||
OutputPluginWrite(ctx, true);
|
||||
break;
|
||||
case REORDER_BUFFER_CHANGE_UPDATE: {
|
||||
HeapTuple oldtuple = change->data.tp.oldtuple ? &change->data.tp.oldtuple->tuple : NULL;
|
||||
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_update(ctx->out, relation, oldtuple, &change->data.tp.newtuple->tuple);
|
||||
logicalrep_write_update(ctx->out, relation, oldtuple, &change->data.tp.newtuple->tuple, data->binary);
|
||||
OutputPluginWrite(ctx, true);
|
||||
break;
|
||||
}
|
||||
|
|
@ -320,14 +349,14 @@ static void pgoutput_change(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
|
|||
HeapTuple oldtuple = change->data.utp.oldtuple ? ((HeapTuple)(&change->data.utp.oldtuple->tuple)) : NULL;
|
||||
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_update(ctx->out, relation, oldtuple, (HeapTuple)(&change->data.utp.newtuple->tuple));
|
||||
logicalrep_write_update(ctx->out, relation, oldtuple, (HeapTuple)(&change->data.utp.newtuple->tuple), data->binary);
|
||||
OutputPluginWrite(ctx, true);
|
||||
break;
|
||||
}
|
||||
case REORDER_BUFFER_CHANGE_DELETE:
|
||||
if (change->data.tp.oldtuple) {
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_delete(ctx->out, relation, &change->data.tp.oldtuple->tuple);
|
||||
logicalrep_write_delete(ctx->out, relation, &change->data.tp.oldtuple->tuple, data->binary);
|
||||
OutputPluginWrite(ctx, true);
|
||||
} else
|
||||
elog(DEBUG1, "didn't send DELETE change because of missing oldtuple");
|
||||
|
|
@ -335,7 +364,7 @@ static void pgoutput_change(LogicalDecodingContext *ctx, ReorderBufferTXN *txn,
|
|||
case REORDER_BUFFER_CHANGE_UDELETE:
|
||||
if (change->data.utp.oldtuple) {
|
||||
OutputPluginPrepareWrite(ctx, true);
|
||||
logicalrep_write_delete(ctx->out, relation, (HeapTuple)(&change->data.utp.oldtuple->tuple));
|
||||
logicalrep_write_delete(ctx->out, relation, (HeapTuple)(&change->data.utp.oldtuple->tuple), data->binary);
|
||||
OutputPluginWrite(ctx, true);
|
||||
} else
|
||||
elog(DEBUG1, "didn't send DELETE change because of missing oldtuple");
|
||||
|
|
@ -605,3 +634,43 @@ static void rel_sync_cache_publication_cb(Datum arg, int cacheid, uint32 hashval
|
|||
entry->pubactions.pubdelete = false;
|
||||
}
|
||||
}
|
||||
|
||||
static void GetConninfo(StringInfoData* standbysInfo)
|
||||
{
|
||||
bool primaryJoined = false;
|
||||
StringInfoData hosts;
|
||||
StringInfoData ports;
|
||||
initStringInfo(&hosts);
|
||||
initStringInfo(&ports);
|
||||
for (int i = 1; i < MAX_REPLNODE_NUM + 1; ++i) {
|
||||
t_thrd.postmaster_cxt.ReplConnChangeType[i] = 0;
|
||||
if (t_thrd.postmaster_cxt.ReplConnArray[i] == NULL) {
|
||||
continue;
|
||||
}
|
||||
if (!primaryJoined) {
|
||||
appendStringInfo(&hosts, "%s,%s",
|
||||
t_thrd.postmaster_cxt.ReplConnArray[i]->localhost,
|
||||
t_thrd.postmaster_cxt.ReplConnArray[i]->remotehost);
|
||||
appendStringInfo(&ports, "%d,%d",
|
||||
t_thrd.postmaster_cxt.ReplConnArray[i]->localport,
|
||||
t_thrd.postmaster_cxt.ReplConnArray[i]->remoteport);
|
||||
primaryJoined = true;
|
||||
} else {
|
||||
appendStringInfo(&hosts, ",%s",
|
||||
t_thrd.postmaster_cxt.ReplConnArray[i]->remotehost);
|
||||
appendStringInfo(&ports, ",%d",
|
||||
t_thrd.postmaster_cxt.ReplConnArray[i]->remoteport);
|
||||
}
|
||||
}
|
||||
appendStringInfo(standbysInfo, "host=%s port=%s", hosts.data, ports.data);
|
||||
}
|
||||
|
||||
static inline bool ReplconninfoChanged()
|
||||
{
|
||||
for (int i = 1; i < MAX_REPLNODE_NUM; ++i) {
|
||||
if (t_thrd.postmaster_cxt.ReplConnChangeType[i]) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -50,13 +50,15 @@ CATALOG(pg_subscription,6126) BKI_SHARED_RELATION BKI_ROWTYPE_OID(6128) BKI_SCHE
|
|||
NameData subslotname; /* Slot name on publisher */
|
||||
text subsynccommit; /* Synchronous commit setting for worker */
|
||||
text subpublications[1]; /* List of publications subscribed to */
|
||||
bool subbinary; /* True if the subscription wants the
|
||||
* publisher to send data in binary */
|
||||
#endif
|
||||
}
|
||||
FormData_pg_subscription;
|
||||
|
||||
typedef FormData_pg_subscription *Form_pg_subscription;
|
||||
|
||||
#define Natts_pg_subscription 8
|
||||
#define Natts_pg_subscription 9
|
||||
#define Anum_pg_subscription_subdbid 1
|
||||
#define Anum_pg_subscription_subname 2
|
||||
#define Anum_pg_subscription_subowner 3
|
||||
|
|
@ -65,6 +67,7 @@ typedef FormData_pg_subscription *Form_pg_subscription;
|
|||
#define Anum_pg_subscription_subslotname 6
|
||||
#define Anum_pg_subscription_subsynccommit 7
|
||||
#define Anum_pg_subscription_subpublications 8
|
||||
#define Anum_pg_subscription_subbinary 9
|
||||
|
||||
|
||||
typedef struct Subscription {
|
||||
|
|
@ -77,6 +80,7 @@ typedef struct Subscription {
|
|||
char *slotname; /* Name of the replication slot */
|
||||
char *synccommit; /* Synchronous commit setting for worker */
|
||||
List *publications; /* List of publication names to subscribe to */
|
||||
bool binary; /* Indicates if the subscription wants data in binary format */
|
||||
} Subscription;
|
||||
|
||||
|
||||
|
|
@ -86,6 +90,7 @@ extern Oid get_subscription_oid(const char *subname, bool missing_ok);
|
|||
extern char *get_subscription_name(Oid subid, bool missing_ok);
|
||||
|
||||
extern int CountDBSubscriptions(Oid dbid);
|
||||
extern char *DecryptConninfo(char *encryptConninfo);
|
||||
extern void ClearListContent(List *list);
|
||||
|
||||
|
||||
#endif /* PG_SUBSCRIPTION_H */
|
||||
|
|
|
|||
|
|
@ -17,6 +17,11 @@
|
|||
|
||||
#include "nodes/parsenodes.h"
|
||||
|
||||
typedef struct HostPort {
|
||||
char* host;
|
||||
char* port;
|
||||
} HostPort;
|
||||
|
||||
extern ObjectAddress CreateSubscription(CreateSubscriptionStmt *stmt, bool isTopLevel);
|
||||
extern ObjectAddress AlterSubscription(AlterSubscriptionStmt *stmt);
|
||||
extern void DropSubscription(DropSubscriptionStmt *stmt, bool isTopLevel);
|
||||
|
|
@ -24,6 +29,12 @@ extern void DropSubscription(DropSubscriptionStmt *stmt, bool isTopLevel);
|
|||
extern ObjectAddress AlterSubscriptionOwner(const char *name, Oid newOwnerId);
|
||||
extern void AlterSubscriptionOwner_oid(Oid subid, Oid newOwnerId);
|
||||
extern void RenameSubscription(List* oldname, const char* newname);
|
||||
extern void AddStandbysInfo(char* standbysInfo);
|
||||
extern void DropStandbysInfo(char* standbysInfo);
|
||||
|
||||
extern void ParseConninfo(const char* conninfo, StringInfoData* conninfoWithoutHostPort, HostPort** hostPortList);
|
||||
extern char* EncryptOrDecryptConninfo(const char* conninfo, const char action);
|
||||
extern bool AttemptConnectPublisher(const char *conninfoOriginal, char* slotname, bool checkRemoteMode);
|
||||
|
||||
#endif /* SUBSCRIPTIONCMDS_H */
|
||||
|
||||
|
|
|
|||
|
|
@ -90,6 +90,7 @@ extern const uint32 SUPPORT_DATA_REPAIR;
|
|||
extern const uint32 SCAN_BATCH_MODE_VERSION_NUM;
|
||||
extern const uint32 RELMAP_4K_VERSION_NUM;
|
||||
extern const uint32 PUBLICATION_VERSION_NUM;
|
||||
extern const uint32 SUBSCRIPTION_BINARY_VERSION_NUM;
|
||||
extern const uint32 ANALYZER_HOOK_VERSION_NUM;
|
||||
extern const uint32 SUPPORT_HASH_XLOG_VERSION_NUM;
|
||||
extern const uint32 PITR_INIT_VERSION_NUM;
|
||||
|
|
|
|||
|
|
@ -35,6 +35,7 @@ typedef struct LibpqrcvConnectParam {
|
|||
bool logical;
|
||||
uint32 protoVersion; /* Logical protocol version */
|
||||
List *publicationNames; /* String list of publications */
|
||||
bool binary; /* Ask publisher to use binary */
|
||||
}LibpqrcvConnectParam;
|
||||
|
||||
extern int32 pg_atoi(char* s, int size, int c);
|
||||
|
|
@ -53,5 +54,6 @@ extern bool libpqrcv_command(const char *cmd, char **err, int *sqlstate);
|
|||
extern void IdentifyRemoteSystem(bool checkRemote);
|
||||
extern void CreateRemoteReplicationSlot(XLogRecPtr startpoint, const char* slotname, bool isLogical);
|
||||
extern void StartRemoteStreaming(const LibpqrcvConnectParam *options);
|
||||
extern ServerMode IdentifyRemoteMode();
|
||||
|
||||
#endif
|
||||
|
|
|
|||
|
|
@ -35,12 +35,21 @@
|
|||
* Keep in mind that the columns correspond to the *remote* table.
|
||||
*/
|
||||
typedef struct LogicalRepTupleData {
|
||||
char *values[MaxTupleAttributeNumber]; /* value in out function format or NULL if values is NULL */
|
||||
bool changed[MaxTupleAttributeNumber]; /* marker for changed/unchanged values */
|
||||
/* Array of StringInfos, one per column; some may be unused */
|
||||
StringInfoData *colvalues;
|
||||
/* Array of markers for null/unchanged/text/binary, one per column */
|
||||
char *colstatus;
|
||||
/* Length of above arrays */
|
||||
int ncols;
|
||||
} LogicalRepTupleData;
|
||||
|
||||
/* Possible values for LogicalRepTupleData.colstatus[colnum] */
|
||||
/* These values are also used in the on-the-wire protocol */
|
||||
#define LOGICALREP_COLUMN_NULL 'n'
|
||||
#define LOGICALREP_COLUMN_UNCHANGED 'u'
|
||||
#define LOGICALREP_COLUMN_TEXT 't'
|
||||
#define LOGICALREP_COLUMN_BINARY 'b' /* added in PG14 */
|
||||
|
||||
typedef uint32 LogicalRepRelId;
|
||||
|
||||
/* Relation information */
|
||||
|
|
@ -81,16 +90,18 @@ extern void logicalrep_read_begin(StringInfo in, LogicalRepBeginData *begin_data
|
|||
extern void logicalrep_write_commit(StringInfo out, ReorderBufferTXN *txn, XLogRecPtr commit_lsn);
|
||||
extern void logicalrep_read_commit(StringInfo in, LogicalRepCommitData *commit_data);
|
||||
extern void logicalrep_write_origin(StringInfo out, const char *origin, XLogRecPtr origin_lsn);
|
||||
extern void logicalrep_write_insert(StringInfo out, Relation rel, HeapTuple newtuple);
|
||||
extern void logicalrep_write_insert(StringInfo out, Relation rel, HeapTuple newtuple, bool binary);
|
||||
extern LogicalRepRelId logicalrep_read_insert(StringInfo in, LogicalRepTupleData *newtup);
|
||||
extern void logicalrep_write_update(StringInfo out, Relation rel, HeapTuple oldtuple, HeapTuple newtuple);
|
||||
extern void logicalrep_write_update(StringInfo out, Relation rel, HeapTuple oldtuple, HeapTuple newtuple, bool binary);
|
||||
extern LogicalRepRelId logicalrep_read_update(StringInfo in, bool *has_oldtuple, LogicalRepTupleData *oldtup,
|
||||
LogicalRepTupleData *newtup);
|
||||
extern void logicalrep_write_delete(StringInfo out, Relation rel, HeapTuple oldtuple);
|
||||
extern void logicalrep_write_delete(StringInfo out, Relation rel, HeapTuple oldtuple, bool binary);
|
||||
extern LogicalRepRelId logicalrep_read_delete(StringInfo in, LogicalRepTupleData *oldtup);
|
||||
extern void logicalrep_write_rel(StringInfo out, Relation rel);
|
||||
extern LogicalRepRelation *logicalrep_read_rel(StringInfo in);
|
||||
extern void logicalrep_write_typ(StringInfo out, Oid typoid);
|
||||
extern void logicalrep_read_typ(StringInfo out, LogicalRepTyp *ltyp);
|
||||
extern void logicalrep_write_conninfo(StringInfo out, char* conninfo);
|
||||
extern void logicalrep_read_conninfo(StringInfo in, char** conninfo);
|
||||
|
||||
#endif /* LOGICALREP_PROTO_H */
|
||||
|
|
|
|||
|
|
@ -13,5 +13,6 @@
|
|||
#define LOGICALWORKER_H
|
||||
|
||||
extern void ApplyWorkerMain();
|
||||
extern bool IsLogicalWorker(void);
|
||||
|
||||
#endif /* LOGICALWORKER_H */
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ typedef struct PGOutputData {
|
|||
|
||||
List *publication_names;
|
||||
List *publications;
|
||||
bool binary;
|
||||
} PGOutputData;
|
||||
|
||||
#endif /* PGOUTPUT_H */
|
||||
|
|
|
|||
|
|
@ -24,6 +24,7 @@ CREATE SUBSCRIPTION testsub CONNECTION 'foo';
|
|||
CREATE SUBSCRIPTION testsub PUBLICATION foo;
|
||||
-- fail - could not connect to the publisher
|
||||
create subscription testsub2 connection 'host=abc' publication pub;
|
||||
create subscription testsub2 connection 'host=abc port=12345' publication pub;
|
||||
set client_min_messages to error;
|
||||
-- fail - syntax error, invalid connection string syntax: missing "="
|
||||
CREATE SUBSCRIPTION testsub CONNECTION 'testconn' PUBLICATION testpub;
|
||||
|
|
@ -32,12 +33,15 @@ CREATE SUBSCRIPTION testsub CONNECTION 'dbname=doesnotexist' PUBLICATION testpub
|
|||
CREATE SUBSCRIPTION testsub CONNECTION 'dbname=doesnotexist' PUBLICATION testpub WITH (ENABLED=false, slot_name='testsub', synchronous_commit=off);
|
||||
-- create SUBSCRIPTION with conninfo in two single quote, used to check mask string bug
|
||||
CREATE SUBSCRIPTION testsub_maskconninfo CONNECTION 'host=''1.2.3.4'' port=''12345'' user=''username'' dbname=''postgres'' password=''password_1234''' PUBLICATION testpub WITH (ENABLED=false, slot_name='testsub', synchronous_commit=off);
|
||||
|
||||
-- fail - The number of host and port are inconsistent
|
||||
create subscription sub1 connection 'dbname=postgres user=pubusr password=Huawei@123 host=192.168.0.38,192.168.0.38,192.168.0.38 port=14001,14501' publication pub1;
|
||||
-- fail - a maximum of 9 servers are supported
|
||||
create subscription sub1 connection 'dbname=postgres user=pubusr password=Huawei@123 host=192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38 port=14001,14501' publication pub1;
|
||||
-- alter connection
|
||||
ALTER SUBSCRIPTION testsub CONNECTION 'host=''1.2.3.4'' port=''12345'' user=''username'' dbname=''postgres'' password=''password_1234''';
|
||||
ALTER SUBSCRIPTION testsub CONNECTION 'dbname=does_not_exist';
|
||||
reset client_min_messages;
|
||||
select subname, pg_get_userbyid(subowner) as Owner, subenabled, subconninfo, subpublications from pg_subscription where subname='testsub';
|
||||
select subname, pg_get_userbyid(subowner) as Owner, subenabled, subconninfo, subpublications, subbinary from pg_subscription where subname='testsub';
|
||||
--- alter subscription
|
||||
------ set publication
|
||||
ALTER SUBSCRIPTION testsub SET PUBLICATION testpub2, testpub3;
|
||||
|
|
@ -57,6 +61,9 @@ select subname, subenabled, subsynccommit from pg_subscription where subname='t
|
|||
ALTER SUBSCRIPTION testsub SET (slot_name='testsub');
|
||||
-- alter owner
|
||||
ALTER SUBSCRIPTION testsub owner to regress_subscription_user2;
|
||||
-- alter subbinary to true
|
||||
ALTER SUBSCRIPTION testsub SET (binary=true);
|
||||
select subname, subbinary from pg_subscription where subname='testsub';
|
||||
--rename
|
||||
ALTER SUBSCRIPTION testsub rename to testsub_rename;
|
||||
--- inside a transaction block
|
||||
|
|
|
|||
|
|
@ -60,8 +60,10 @@ LINE 1: CREATE SUBSCRIPTION testsub PUBLICATION foo;
|
|||
^
|
||||
-- fail - could not connect to the publisher
|
||||
create subscription testsub2 connection 'host=abc' publication pub;
|
||||
ERROR: The number of host and port are inconsistent.
|
||||
create subscription testsub2 connection 'host=abc port=12345' publication pub;
|
||||
WARNING: apply worker could not connect to the remote server
|
||||
ERROR: could not connect to the publisher
|
||||
ERROR: Failed to connect to publisher.
|
||||
set client_min_messages to error;
|
||||
-- fail - syntax error, invalid connection string syntax: missing "="
|
||||
CREATE SUBSCRIPTION testsub CONNECTION 'testconn' PUBLICATION testpub;
|
||||
|
|
@ -72,14 +74,20 @@ ERROR: unrecognized subscription parameter: create_slot
|
|||
CREATE SUBSCRIPTION testsub CONNECTION 'dbname=doesnotexist' PUBLICATION testpub WITH (ENABLED=false, slot_name='testsub', synchronous_commit=off);
|
||||
-- create SUBSCRIPTION with conninfo in two single quote, used to check mask string bug
|
||||
CREATE SUBSCRIPTION testsub_maskconninfo CONNECTION 'host=''1.2.3.4'' port=''12345'' user=''username'' dbname=''postgres'' password=''password_1234''' PUBLICATION testpub WITH (ENABLED=false, slot_name='testsub', synchronous_commit=off);
|
||||
-- fail - The number of host and port are inconsistent
|
||||
create subscription sub1 connection 'dbname=postgres user=pubusr password=Huawei@123 host=192.168.0.38,192.168.0.38,192.168.0.38 port=14001,14501' publication pub1;
|
||||
ERROR: The number of host and port are inconsistent.
|
||||
-- fail - a maximum of 9 servers are supported
|
||||
create subscription sub1 connection 'dbname=postgres user=pubusr password=Huawei@123 host=192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38,192.168.0.38 port=14001,14501' publication pub1;
|
||||
ERROR: Currently, a maximum of 9 servers are supported.
|
||||
-- alter connection
|
||||
ALTER SUBSCRIPTION testsub CONNECTION 'host=''1.2.3.4'' port=''12345'' user=''username'' dbname=''postgres'' password=''password_1234''';
|
||||
ALTER SUBSCRIPTION testsub CONNECTION 'dbname=does_not_exist';
|
||||
reset client_min_messages;
|
||||
select subname, pg_get_userbyid(subowner) as Owner, subenabled, subconninfo, subpublications from pg_subscription where subname='testsub';
|
||||
subname | owner | subenabled | subconninfo | subpublications
|
||||
---------+---------------------------+------------+------------------------+-----------------
|
||||
testsub | regress_subscription_user | f | dbname=does_not_exist | {testpub}
|
||||
select subname, pg_get_userbyid(subowner) as Owner, subenabled, subconninfo, subpublications, subbinary from pg_subscription where subname='testsub';
|
||||
subname | owner | subenabled | subconninfo | subpublications | subbinary
|
||||
---------+---------------------------+------------+------------------------+-----------------+-----------
|
||||
testsub | regress_subscription_user | f | dbname=does_not_exist | {testpub} | f
|
||||
(1 row)
|
||||
|
||||
--- alter subscription
|
||||
|
|
@ -122,6 +130,14 @@ ALTER SUBSCRIPTION testsub SET (slot_name='testsub');
|
|||
ERROR: Currently enabled=false, cannot change slot_name to a non-null value.
|
||||
-- alter owner
|
||||
ALTER SUBSCRIPTION testsub owner to regress_subscription_user2;
|
||||
-- alter subbinary to true
|
||||
ALTER SUBSCRIPTION testsub SET (binary=true);
|
||||
select subname, subbinary from pg_subscription where subname='testsub';
|
||||
subname | subbinary
|
||||
---------+-----------
|
||||
testsub | t
|
||||
(1 row)
|
||||
|
||||
--rename
|
||||
ALTER SUBSCRIPTION testsub rename to testsub_rename;
|
||||
--- inside a transaction block
|
||||
|
|
@ -134,8 +150,7 @@ COMMIT;
|
|||
-- -- active SUBSCRIPTION
|
||||
BEGIN;
|
||||
ALTER SUBSCRIPTION testsub_rename ENABLE;
|
||||
WARNING: apply worker could not connect to the remote server
|
||||
ERROR: could not connect to the publisher
|
||||
ERROR: invalid connection string syntax, missing host and port
|
||||
select subname, subenabled from pg_subscription where subname='testsub_rename';
|
||||
ERROR: current transaction is aborted, commands ignored until end of transaction block, firstChar[Q]
|
||||
ALTER SUBSCRIPTION testsub_rename SET (ENABLED=false);
|
||||
|
|
@ -281,10 +296,11 @@ SELECT object_name,detail_info FROM pg_query_audit('2022-01-13 9:30:00', '2031-1
|
|||
testsub_maskconninfo | ALTER SUBSCRIPTION testsub_maskconninfo SET (conninfo='*************************************************************************************************************************;
|
||||
testsub | ALTER SUBSCRIPTION testsub SET (synchronous_commit=on);
|
||||
testsub | ALTER SUBSCRIPTION testsub owner to regress_subscription_user2;
|
||||
testsub | ALTER SUBSCRIPTION testsub SET (binary=true);
|
||||
testsub | ALTER SUBSCRIPTION testsub rename to testsub_rename;
|
||||
testsub_rename | DROP SUBSCRIPTION IF EXISTS testsub_rename;
|
||||
testsub_maskconninfo | DROP SUBSCRIPTION IF EXISTS testsub_maskconninfo;
|
||||
(14 rows)
|
||||
(15 rows)
|
||||
|
||||
--clear audit log
|
||||
SELECT pg_delete_audit('1012-11-10', '3012-11-11');
|
||||
|
|
|
|||
Loading…
Reference in New Issue