Skip to content

Commit 63df626

Browse files
Publication and Subscription are corrected and refactored (#388)
The minimal changes. We fix and check (via assert) an usage of dbname/username only.
1 parent 27b6827 commit 63df626

1 file changed

Lines changed: 152 additions & 27 deletions

File tree

src/pubsub.py

Lines changed: 152 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,12 @@ def __init__(self, name, node, tables=None, dbname=None, username=None):
6363
dbname: database name used to connect and perform subscription.
6464
username: username used to connect to the database.
6565
"""
66+
assert type(name) is str
67+
assert node is not None
68+
assert node.os_ops is not None
69+
assert dbname is None or type(dbname) is str
70+
assert username is None or type(username) is str
71+
6672
self.name = name
6773
self.node = node
6874
self.dbname = dbname or default_dbname()
@@ -71,15 +77,32 @@ def __init__(self, name, node, tables=None, dbname=None, username=None):
7177
# create publication in database
7278
t = "table " + ", ".join(tables) if tables else "all tables"
7379
query = "create publication {} for {}"
74-
node.execute(query.format(name, t), dbname=dbname, username=username)
80+
self.node.execute(
81+
query.format(name, t),
82+
dbname=self.dbname,
83+
username=self.username,
84+
)
7585

7686
def drop(self, dbname=None, username=None):
7787
"""
7888
Drop publication
7989
"""
80-
self.node.execute("drop publication {}".format(self.name),
81-
dbname=dbname,
82-
username=username)
90+
assert dbname is None or type(dbname) is str
91+
assert username is None or type(username) is str
92+
93+
#
94+
# [2026-07-10] [BUG FIX]
95+
# dbname and username are ignored.
96+
# We will use settings of our object.
97+
#
98+
assert dbname is None or dbname == self.dbname
99+
assert username is None or username == self.username
100+
101+
self.node.execute(
102+
"drop publication {}".format(self.name),
103+
dbname=self.dbname,
104+
username=self.username,
105+
)
83106

84107
def add_tables(self, tables, dbname=None, username=None):
85108
"""
@@ -89,13 +112,26 @@ def add_tables(self, tables, dbname=None, username=None):
89112
Args:
90113
tables: a list of tables to be added to the publication.
91114
"""
115+
assert dbname is None or type(dbname) is str
116+
assert username is None or type(username) is str
117+
118+
#
119+
# [2026-07-10] [BUG FIX]
120+
# dbname and username are ignored.
121+
# We will use settings of our object.
122+
#
123+
assert dbname is None or dbname == self.dbname
124+
assert username is None or username == self.username
125+
92126
if not tables:
93127
raise ValueError("Tables list is empty")
94128

95129
query = "alter publication {} add table {}"
96-
self.node.execute(query.format(self.name, ", ".join(tables)),
97-
dbname=dbname or self.dbname,
98-
username=username or self.username)
130+
self.node.execute(
131+
query.format(self.name, ", ".join(tables)),
132+
dbname=self.dbname,
133+
username=self.username,
134+
)
99135

100136

101137
class Subscription(object):
@@ -121,9 +157,17 @@ def __init__(self,
121157
<https://www.postgresql.org/docs/current/static/sql-createsubscription.html>`_
122158
for details).
123159
"""
160+
assert type(name) is str
161+
assert node is not None
162+
assert node.os_ops is not None
163+
assert dbname is None or type(dbname) is str
164+
assert username is None or type(username) is str
165+
124166
self.name = name
125167
self.node = node
126168
self.pub = publication
169+
self.dbname = dbname or default_dbname()
170+
self.username = username or default_username()
127171

128172
# connection info
129173
conninfo = {
@@ -142,38 +186,99 @@ def __init__(self,
142186
query += " with ({})".format(options_string(**params))
143187

144188
# Note: cannot run 'create subscription' query in transaction mode
145-
node.execute(query, dbname=dbname, username=username)
189+
self.node.execute(
190+
query,
191+
dbname=self.dbname,
192+
username=self.username,
193+
)
146194

147195
def disable(self, dbname=None, username=None):
148196
"""
149197
Disables the running subscription.
150198
"""
199+
assert dbname is None or type(dbname) is str
200+
assert username is None or type(username) is str
201+
202+
#
203+
# [2026-07-10] [BUG FIX]
204+
# dbname and username are ignored.
205+
# We will use settings of our object.
206+
#
207+
assert dbname is None or dbname == self.dbname
208+
assert username is None or username == self.username
209+
151210
query = "alter subscription {} disable"
152-
self.node.execute(query.format(self.name), dbname=None, username=None)
211+
self.node.execute(
212+
query.format(self.name),
213+
dbname=self.dbname,
214+
username=self.username,
215+
)
153216

154217
def enable(self, dbname=None, username=None):
155218
"""
156219
Enables the previously disabled subscription.
157220
"""
221+
assert dbname is None or type(dbname) is str
222+
assert username is None or type(username) is str
223+
224+
#
225+
# [2026-07-10] [BUG FIX]
226+
# dbname and username were and are ignored.
227+
# We will use settings of our object.
228+
#
229+
assert dbname is None or dbname == self.dbname
230+
assert username is None or username == self.username
231+
158232
query = "alter subscription {} enable"
159-
self.node.execute(query.format(self.name), dbname=None, username=None)
233+
234+
self.node.execute(
235+
query.format(self.name),
236+
dbname=self.dbname,
237+
username=self.username,
238+
)
160239

161240
def refresh(self, copy_data=True, dbname=None, username=None):
162241
"""
163242
Disables the running subscription.
164243
"""
244+
assert dbname is None or type(dbname) is str
245+
assert username is None or type(username) is str
246+
247+
#
248+
# [2026-07-10] [BUG FIX]
249+
# dbname and username are ignored.
250+
# We will use settings of our object.
251+
#
252+
assert dbname is None or dbname == self.dbname
253+
assert username is None or username == self.username
254+
165255
query = "alter subscription {} refresh publication with (copy_data={})"
166-
self.node.execute(query.format(self.name, copy_data),
167-
dbname=dbname,
168-
username=username)
256+
self.node.execute(
257+
query.format(self.name, copy_data),
258+
dbname=self.dbname,
259+
username=self.username,
260+
)
169261

170262
def drop(self, dbname=None, username=None):
171263
"""
172264
Drops subscription
173265
"""
174-
self.node.execute("drop subscription {}".format(self.name),
175-
dbname=dbname,
176-
username=username)
266+
assert dbname is None or type(dbname) is str
267+
assert username is None or type(username) is str
268+
269+
#
270+
# [2026-07-10] [BUG FIX]
271+
# dbname and username are ignored.
272+
# We will use settings of our object.
273+
#
274+
assert dbname is None or dbname == self.dbname
275+
assert username is None or username == self.username
276+
277+
self.node.execute(
278+
"drop subscription {}".format(self.name),
279+
dbname=self.dbname,
280+
username=self.username,
281+
)
177282

178283
def catchup(self, username=None):
179284
"""
@@ -182,14 +287,32 @@ def catchup(self, username=None):
182287
Args:
183288
username: remote node's user name.
184289
"""
290+
assert username is None or type(username) is str
291+
292+
#
293+
# [2026-07-10] [BUG FIX]
294+
# username is ignored.
295+
# We will use settings of objects.
296+
#
297+
assert username is None or username == self.username
298+
185299
try:
186-
pub_lsn = self.pub.node.execute(query="select pg_current_wal_lsn()",
187-
dbname=None,
188-
username=None)[0][0] # yapf: disable
300+
#
301+
# [2026-07-10]
302+
# About dbname=None and username=None
303+
# We will try to use self.pub.xxx the next time. OK?
304+
#
305+
pub_lsn = self.pub.node.execute(
306+
query="select pg_current_wal_lsn()",
307+
dbname=None,
308+
username=None,
309+
)[0][0] # yapf: disable
189310
# create dummy xact, as LR replicates only on commit.
190-
self.pub.node.execute(query="select txid_current()",
191-
dbname=None,
192-
username=None)
311+
self.pub.node.execute(
312+
query="select txid_current()",
313+
dbname=None,
314+
username=None,
315+
)
193316
query = """
194317
select '{}'::pg_lsn - replay_lsn <= 0
195318
from pg_catalog.pg_stat_replication where application_name = '{}'
@@ -199,8 +322,9 @@ def catchup(self, username=None):
199322
self.pub.node.poll_query_until(
200323
query=query,
201324
dbname=self.pub.dbname,
202-
username=username or self.pub.username,
203-
max_attempts=LOGICAL_REPL_MAX_CATCHUP_ATTEMPTS)
325+
username=self.pub.username,
326+
max_attempts=LOGICAL_REPL_MAX_CATCHUP_ATTEMPTS,
327+
)
204328

205329
# Now, wait until there are no tablesync workers: probably
206330
# replay_lsn above was sent with changes of new tables just skipped;
@@ -210,8 +334,9 @@ def catchup(self, username=None):
210334
"""
211335
self.node.poll_query_until(
212336
query=query,
213-
dbname=self.pub.dbname,
214-
username=username or self.pub.username,
215-
max_attempts=LOGICAL_REPL_MAX_CATCHUP_ATTEMPTS)
337+
dbname=self.dbname,
338+
username=self.username,
339+
max_attempts=LOGICAL_REPL_MAX_CATCHUP_ATTEMPTS,
340+
)
216341
except Exception as e:
217342
raise_from(CatchUpException("Failed to catch up"), e)

0 commit comments

Comments
 (0)