Skip to content

Commit 3b53938

Browse files
committed
core: protect ro nodes on disconnect
Adds protection for on_disconnect trigger. After leader change it may return ERR_READONLY while deleting consumers. Now this trigger performs only local operations on ro node.
1 parent a8f5e6b commit 3b53938

3 files changed

Lines changed: 175 additions & 14 deletions

File tree

CHANGELOG.md

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,9 @@ and this project adheres to [Semantic Versioning](http://semver.org/spec/v2.0.0.
1313

1414
### Fixed
1515

16+
- `ERR_READONLY` errors on previous leader after leader change in cluster
17+
while processing `on_disconnect` triggers (#248).
18+
1619
## [1.4.4] - 2025-05-26
1720

1821
The patch release fixes incorrect behavior of the utubettl driver with enabled

queue/abstract.lua

Lines changed: 36 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -479,6 +479,10 @@ end
479479

480480
--- Release all session tasks.
481481
local function release_session_tasks(session_uuid)
482+
if box.info.ro then
483+
return
484+
end
485+
482486
local taken_tasks = box.space._queue_taken_2.index.uuid:select{session_uuid}
483487

484488
for _, task in pairs(taken_tasks) do
@@ -501,26 +505,25 @@ end
501505

502506
function method._on_consumer_disconnect()
503507
local conn_id = connection.id()
508+
local consumers = box.space._queue_consumers
504509

505510
-- wakeup all waiters
506-
while true do
507-
local waiter = box.space._queue_consumers.index.pk:min{conn_id}
508-
if waiter == nil then
509-
break
510-
end
511-
-- Don't touch the other consumers
512-
if waiter[1] ~= conn_id then
513-
break
514-
end
515-
box.space._queue_consumers:delete{waiter[1], waiter[2]}
516-
local cond = conds[waiter[2]]
511+
for _, waiter in consumers.index.pk:pairs(conn_id, { iterator = 'EQ' }) do
512+
local fid = waiter[2]
513+
local cond = conds[fid]
517514
if cond then
518-
releasing_connections[waiter[2]] = true
519-
cond:signal(waiter[2])
515+
releasing_connections[fid] = true
516+
cond:signal(fid)
517+
end
518+
519+
if not box.info.ro then
520+
consumers:delete{waiter[1], waiter[2]}
520521
end
521522
end
522523

523-
session.disconnect(conn_id)
524+
if not box.info.ro then
525+
session.disconnect(conn_id)
526+
end
524527
end
525528

526529
-- function takes tuples and recreates tube
@@ -539,10 +542,29 @@ local function recreate_tube(tube_tuple)
539542
return make_self(driver, space, name, tube_type, id, opts)
540543
end
541544

545+
-- function cleans local temporary spaces on startup to avoid
546+
-- storing old data
547+
local function cleanup_temp_spaces()
548+
if box.info.ro then
549+
return
550+
end
551+
552+
local s = box.space._queue_consumers
553+
if s ~= nil then
554+
s:truncate()
555+
end
556+
557+
s = box.space._queue_session_ids
558+
if s ~= nil then
559+
s:truncate()
560+
end
561+
end
562+
542563
-- Function takes new queue state.
543564
-- The "RUNNING" and "WAITING" states do not require additional actions.
544565
local function on_state_change(state)
545566
if state == queue_state.states.STARTUP then
567+
cleanup_temp_spaces()
546568
local replicaset_mode = queue.cfg['in_replicaset'] or false
547569
-- gh-202: In replicaset mode, tubes can be created and deleted on different nodes.
548570
-- Accordingly, it is necessary to rebuild the queue.tube index.

t/240-ro-on-disconnect.t

Lines changed: 136 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,136 @@
1+
#!/usr/bin/env tarantool
2+
3+
local log = require('log')
4+
local tnt = require('t.tnt')
5+
local test = require('tap').test('')
6+
local fiber = require('fiber')
7+
local queue = require('queue')
8+
9+
local qc = require('queue.compat')
10+
if not qc.check_version({2, 4, 1}) then
11+
log.info('Tests skipped, tarantool version < 2.4.1')
12+
return
13+
end
14+
15+
rawset(_G, 'queue', require('queue'))
16+
17+
local session = require('queue.abstract.queue_session')
18+
local queue_state = require('queue.abstract.queue_state')
19+
20+
test:plan(3)
21+
22+
test:test('on_disconnect handler must be RO-safe', function(test)
23+
test:plan(6)
24+
25+
tnt.cluster.cfg{}
26+
test:ok(tnt.cluster.wait_replica(), 'wait for replica to connect')
27+
28+
queue.cfg{ttr = 0.5, in_replicaset = true}
29+
local tube = queue.create_tube('test_ro_disc', 'fifo', {if_not_exists = true})
30+
test:ok(tube, 'tube created')
31+
32+
local f = fiber.new(function()
33+
queue.tube.test_ro_disc:take(3600)
34+
end)
35+
f:name('queue_waiter_fiber')
36+
37+
local ok = false
38+
for _ = 1, 300 do
39+
if box.space._queue_consumers:count() > 0 then
40+
ok = true
41+
break
42+
end
43+
fiber.sleep(0.01)
44+
end
45+
test:ok(ok, 'waiter registered in _queue_consumers')
46+
47+
box.cfg{read_only = true}
48+
test:ok(box.info.ro, 'instance is RO')
49+
50+
local ok_call, err = pcall(queue._on_consumer_disconnect)
51+
test:ok(ok_call, ('_on_consumer_disconnect() must not fail on RO, err = %s'):format(tostring(err)))
52+
53+
box.cfg{read_only = false}
54+
test:ok(not box.info.ro, 'instance back to RW')
55+
end)
56+
57+
test:test('release_session_tasks: RO-safe', function(test)
58+
test:plan(10)
59+
60+
local tube = queue.create_tube('test_rel', 'fifo', {if_not_exists = true})
61+
test:ok(tube, 'tube created')
62+
63+
-- Create a session + take a task so _queue_taken_2 has a record.
64+
local client = tnt.cluster.connect_master()
65+
test:ok(client.error == nil, 'client connected')
66+
67+
local session_uuid = client:call('queue.identify')
68+
test:ok(session_uuid ~= nil, 'got session_uuid')
69+
70+
test:ok(queue.tube.test_rel:put('data'), 'put task')
71+
local task = client:call('queue.tube.test_rel:take')
72+
test:ok(task ~= nil, 'task taken')
73+
74+
local taken_before = box.space._queue_taken_2.index.uuid:select{session_uuid}
75+
test:is(#taken_before, 1, '_queue_taken_2 has 1 record before')
76+
77+
box.cfg{read_only = true}
78+
test:ok(queue_state.poll(queue_state.states.WAITING, 10), 'state WAITING')
79+
test:ok(box.info.ro, 'instance is RO')
80+
81+
local ok_call, err = pcall(session._on_session_remove, session_uuid)
82+
test:ok(ok_call, ('on_session_remove does not fail on RO, err=%s'):format(tostring(err)))
83+
84+
local taken_after = box.space._queue_taken_2.index.uuid:select{session_uuid}
85+
test:is(#taken_after, 1, '_queue_taken_2 unchanged on RO')
86+
87+
-- Cleanup: back to RW.
88+
box.cfg{read_only = false}
89+
queue_state.poll(queue_state.states.RUNNING, 10)
90+
client:close()
91+
end)
92+
93+
test:test('release_session_tasks: works on RW', function(test)
94+
test:plan(11)
95+
96+
box.cfg{read_only = false}
97+
queue_state.poll(queue_state.states.RUNNING, 10)
98+
test:ok(not box.info.ro, 'instance is RW')
99+
100+
queue.cfg{ttr = 0.5, in_replicaset = true}
101+
local tube = queue.create_tube('test_rel2', 'fifo', {if_not_exists = true})
102+
test:ok(tube, 'tube created')
103+
104+
local client = tnt.cluster.connect_master()
105+
test:ok(client.error == nil, 'client connected')
106+
107+
local session_uuid = client:call('queue.identify')
108+
test:ok(session_uuid ~= nil, 'got session_uuid')
109+
110+
test:ok(queue.tube.test_rel2:put('data2'), 'put task')
111+
local task = client:call('queue.tube.test_rel2:take')
112+
test:ok(task ~= nil, 'task taken')
113+
114+
local taken_before = box.space._queue_taken_2.index.uuid:select{session_uuid}
115+
test:is(#taken_before, 1, '_queue_taken_2 has 1 record before')
116+
117+
-- Call on_session_remove callback directly on RW: must release task.
118+
local ok_call, err = pcall(session._on_session_remove, session_uuid)
119+
test:ok(ok_call, ('on_session_remove ok on RW, err=%s'):format(tostring(err)))
120+
121+
-- Taken record must be removed.
122+
local taken_after = box.space._queue_taken_2.index.uuid:select{session_uuid}
123+
test:is(#taken_after, 0, '_queue_taken_2 record removed')
124+
125+
-- Task must become READY again and be takeable.
126+
local task2 = client:call('queue.tube.test_rel2:take', {0})
127+
test:ok(task2 ~= nil, 'task is takeable again after release')
128+
test:is(task2[3], 'data2', 'task data preserved')
129+
130+
client:close()
131+
end)
132+
133+
rawset(_G, 'queue', nil)
134+
tnt.finish()
135+
os.exit(test:check() and 0 or 1)
136+
-- vim: set ft=lua :

0 commit comments

Comments
 (0)