Skip to content

Commit dcfd4cf

Browse files
committed
fix dbconsumer_next
Signed-off-by: Emelia Lei <wlei29@bloomberg.net>
1 parent ddab29b commit dcfd4cf

4 files changed

Lines changed: 336 additions & 1 deletion

File tree

lua/sp.c

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -933,11 +933,15 @@ static int dbconsumer_next(Lua L)
933933
struct sqlclntstate *clnt = sp->clnt;
934934
if (!clnt->intrans) {
935935
/* First write done by this txn */
936+
rc = start_new_transaction(clnt);
937+
if (rc) {
938+
luaL_error(L, "%s: start_new_transaction intrans:%d rc:%d\n",
939+
__func__, clnt->intrans, rc);
940+
}
936941
rc = osql_sock_start_no_reorder(clnt, OSQL_SOCK_REQ, 0, 0);
937942
if (rc) {
938943
luaL_error(L, "%s osql_sock_start rc:%d", __func__, rc);
939944
}
940-
clnt->intrans = 1;
941945
}
942946
Q4SP(qname, q->info.spname);
943947
++clnt->osql_max_trans;

tests/consumer.test/runit

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@ if [ "$TESTCASE" == "consumer" ]; then
77
./t06.sh
88
./t07.sh
99
./t08.sh
10+
./t10.sh
1011
./cdb2api_drain.sh
1112
fi
1213
${TESTSROOTDIR}/tools/compare_results.sh -s -d $1

tests/consumer.test/t10.expected

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
test1 updated rows:
2+
13 1
3+
14 0
4+
test1 queue depth:
5+
0
6+
test2 inserted rows:
7+
99
8+
test2 queue depth:
9+
0
10+
test3 remaining rows:
11+
1
12+
3
13+
test3 queue depth:
14+
0
15+
test4 updated rows:
16+
1 1
17+
2 0
18+
3 1
19+
4 0
20+
5 1
21+
test4 queue depth:
22+
0
23+
test5 inserted rows:
24+
10
25+
20
26+
30
27+
test5 queue depth:
28+
0
29+
test6 remaining rows:
30+
1
31+
3
32+
5
33+
test6 queue depth:
34+
0

tests/consumer.test/t10.sh

Lines changed: 296 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,296 @@
1+
#!/bin/bash
2+
3+
# Test: consumer:next() followed by DML in an explicit transaction.
4+
# Reproduces bug where consumer:next() did not initialize shadow_tran,
5+
# causing "wrong sql transaction" errors on subsequent DML.
6+
7+
cdb2sql="${CDB2SQL_EXE} -tabs -s ${CDB2_OPTIONS} ${DBNAME} default"
8+
9+
(
10+
# ---- Test 1: consumer:next() + UPDATE ----
11+
12+
$cdb2sql "DROP TABLE IF EXISTS bug_src"
13+
$cdb2sql "DROP TABLE IF EXISTS bug_tbl"
14+
15+
$cdb2sql <<'EOF' >/dev/null 2>&1
16+
CREATE TABLE bug_src(tktnum int)$$
17+
CREATE TABLE bug_tbl(tktnum int, status int)$$
18+
CREATE PROCEDURE next_upd_test VERSION 'bar' {
19+
local function main()
20+
local consumer = db:consumer()
21+
local event = consumer:get()
22+
db:begin()
23+
consumer:next()
24+
local stmt = db:prepare("UPDATE bug_tbl SET status=1 WHERE tktnum=@tktnum AND status=0")
25+
stmt:bind("tktnum", event.new.tktnum)
26+
local rc = stmt:exec()
27+
if rc ~= 0 then
28+
db:rollback()
29+
return -201, db:error()
30+
else
31+
db:commit()
32+
end
33+
end
34+
}$$
35+
CREATE LUA CONSUMER next_upd_test ON (TABLE bug_src FOR INSERT);
36+
EOF
37+
38+
$cdb2sql "INSERT INTO bug_tbl VALUES(13, 0), (14, 0)" >/dev/null 2>&1
39+
$cdb2sql "INSERT INTO bug_src VALUES(13)" >/dev/null 2>&1
40+
41+
sleep 2
42+
43+
$cdb2sql "EXEC PROCEDURE next_upd_test()" >/dev/null 2>&1
44+
45+
echo "test1 updated rows:"
46+
$cdb2sql "SELECT tktnum, status FROM bug_tbl ORDER BY tktnum"
47+
48+
echo "test1 queue depth:"
49+
$cdb2sql "SELECT depth FROM comdb2_queues WHERE spname='next_upd_test'"
50+
51+
$cdb2sql "DROP LUA CONSUMER next_upd_test" >/dev/null 2>&1
52+
$cdb2sql "DROP TABLE bug_src" >/dev/null 2>&1
53+
$cdb2sql "DROP TABLE bug_tbl" >/dev/null 2>&1
54+
55+
# ---- Test 2: consumer:next() + INSERT ----
56+
57+
$cdb2sql <<'EOF' >/dev/null 2>&1
58+
CREATE TABLE ins_src(i int)$$
59+
CREATE TABLE ins_dst(i int)$$
60+
CREATE PROCEDURE next_ins_test VERSION 'bar' {
61+
local function main()
62+
local consumer = db:consumer()
63+
local event = consumer:get()
64+
db:begin()
65+
consumer:next()
66+
local stmt = db:prepare("INSERT INTO ins_dst VALUES(@i)")
67+
stmt:bind("i", event.new.i)
68+
local rc = stmt:exec()
69+
if rc ~= 0 then
70+
db:rollback()
71+
return -201, db:error()
72+
else
73+
db:commit()
74+
end
75+
end
76+
}$$
77+
CREATE LUA CONSUMER next_ins_test ON (TABLE ins_src FOR INSERT);
78+
EOF
79+
80+
$cdb2sql "INSERT INTO ins_src VALUES(99)" >/dev/null 2>&1
81+
82+
sleep 2
83+
84+
$cdb2sql "EXEC PROCEDURE next_ins_test()" >/dev/null 2>&1
85+
86+
echo "test2 inserted rows:"
87+
$cdb2sql "SELECT i FROM ins_dst ORDER BY i"
88+
89+
echo "test2 queue depth:"
90+
$cdb2sql "SELECT depth FROM comdb2_queues WHERE spname='next_ins_test'"
91+
92+
$cdb2sql "DROP LUA CONSUMER next_ins_test" >/dev/null 2>&1
93+
$cdb2sql "DROP TABLE ins_src" >/dev/null 2>&1
94+
$cdb2sql "DROP TABLE ins_dst" >/dev/null 2>&1
95+
96+
# ---- Test 3: consumer:next() + DELETE ----
97+
98+
$cdb2sql <<'EOF' >/dev/null 2>&1
99+
CREATE TABLE del_src(i int)$$
100+
CREATE TABLE del_tbl(i int)$$
101+
CREATE PROCEDURE next_del_test VERSION 'bar' {
102+
local function main()
103+
local consumer = db:consumer()
104+
local event = consumer:get()
105+
db:begin()
106+
consumer:next()
107+
local stmt = db:prepare("DELETE FROM del_tbl WHERE i=@i")
108+
stmt:bind("i", event.new.i)
109+
local rc = stmt:exec()
110+
if rc ~= 0 then
111+
db:rollback()
112+
return -201, db:error()
113+
else
114+
db:commit()
115+
end
116+
end
117+
}$$
118+
CREATE LUA CONSUMER next_del_test ON (TABLE del_src FOR INSERT);
119+
EOF
120+
121+
$cdb2sql "INSERT INTO del_tbl VALUES(1), (2), (3)" >/dev/null 2>&1
122+
$cdb2sql "INSERT INTO del_src VALUES(2)" >/dev/null 2>&1
123+
124+
sleep 2
125+
126+
$cdb2sql "EXEC PROCEDURE next_del_test()" >/dev/null 2>&1
127+
128+
echo "test3 remaining rows:"
129+
$cdb2sql "SELECT i FROM del_tbl ORDER BY i"
130+
131+
echo "test3 queue depth:"
132+
$cdb2sql "SELECT depth FROM comdb2_queues WHERE spname='next_del_test'"
133+
134+
$cdb2sql "DROP LUA CONSUMER next_del_test" >/dev/null 2>&1
135+
$cdb2sql "DROP TABLE del_src" >/dev/null 2>&1
136+
$cdb2sql "DROP TABLE del_tbl" >/dev/null 2>&1
137+
138+
# ---- Test 4: consumer:next() + UPDATE in a loop (multiple events) ----
139+
140+
$cdb2sql <<'EOF' >/dev/null 2>&1
141+
CREATE TABLE loop_src(tktnum int)$$
142+
CREATE TABLE loop_tbl(tktnum int, status int)$$
143+
CREATE PROCEDURE next_loop_test VERSION 'bar' {
144+
local function main()
145+
local consumer = db:consumer()
146+
while true do
147+
local event = consumer:get()
148+
db:begin()
149+
consumer:next()
150+
local stmt = db:prepare("UPDATE loop_tbl SET status=1 WHERE tktnum=@tktnum AND status=0")
151+
stmt:bind("tktnum", event.new.tktnum)
152+
local rc = stmt:exec()
153+
if rc ~= 0 then
154+
db:rollback()
155+
return -201, db:error()
156+
else
157+
db:commit()
158+
end
159+
end
160+
end
161+
}$$
162+
CREATE LUA CONSUMER next_loop_test ON (TABLE loop_src FOR INSERT);
163+
EOF
164+
165+
$cdb2sql "INSERT INTO loop_tbl VALUES(1,0),(2,0),(3,0),(4,0),(5,0)" >/dev/null 2>&1
166+
$cdb2sql "INSERT INTO loop_src VALUES(1),(3),(5)" >/dev/null 2>&1
167+
168+
sleep 2
169+
170+
$cdb2sql "EXEC PROCEDURE next_loop_test()" >/dev/null 2>&1 &
171+
CONSUMER_PID=$!
172+
173+
sleep 4
174+
175+
kill $CONSUMER_PID 2>/dev/null
176+
wait $CONSUMER_PID 2>/dev/null
177+
178+
echo "test4 updated rows:"
179+
$cdb2sql "SELECT tktnum, status FROM loop_tbl ORDER BY tktnum"
180+
181+
echo "test4 queue depth:"
182+
$cdb2sql "SELECT depth FROM comdb2_queues WHERE spname='next_loop_test'"
183+
184+
$cdb2sql "DROP LUA CONSUMER next_loop_test" >/dev/null 2>&1
185+
$cdb2sql "DROP TABLE loop_src" >/dev/null 2>&1
186+
$cdb2sql "DROP TABLE loop_tbl" >/dev/null 2>&1
187+
188+
# ---- Test 5: consumer:next() + INSERT in a loop (multiple events) ----
189+
190+
$cdb2sql <<'EOF' >/dev/null 2>&1
191+
CREATE TABLE loop_ins_src(i int)$$
192+
CREATE TABLE loop_ins_dst(i int)$$
193+
CREATE PROCEDURE next_loop_ins VERSION 'bar' {
194+
local function main()
195+
local consumer = db:consumer()
196+
while true do
197+
local event = consumer:get()
198+
db:begin()
199+
consumer:next()
200+
local stmt = db:prepare("INSERT INTO loop_ins_dst VALUES(@i)")
201+
stmt:bind("i", event.new.i)
202+
local rc = stmt:exec()
203+
if rc ~= 0 then
204+
db:rollback()
205+
return -201, db:error()
206+
else
207+
db:commit()
208+
end
209+
end
210+
end
211+
}$$
212+
CREATE LUA CONSUMER next_loop_ins ON (TABLE loop_ins_src FOR INSERT);
213+
EOF
214+
215+
$cdb2sql "INSERT INTO loop_ins_src VALUES(10),(20),(30)" >/dev/null 2>&1
216+
217+
sleep 2
218+
219+
$cdb2sql "EXEC PROCEDURE next_loop_ins()" >/dev/null 2>&1 &
220+
CONSUMER_PID=$!
221+
222+
sleep 4
223+
224+
kill $CONSUMER_PID 2>/dev/null
225+
wait $CONSUMER_PID 2>/dev/null
226+
227+
echo "test5 inserted rows:"
228+
$cdb2sql "SELECT i FROM loop_ins_dst ORDER BY i"
229+
230+
echo "test5 queue depth:"
231+
$cdb2sql "SELECT depth FROM comdb2_queues WHERE spname='next_loop_ins'"
232+
233+
$cdb2sql "DROP LUA CONSUMER next_loop_ins" >/dev/null 2>&1
234+
$cdb2sql "DROP TABLE loop_ins_src" >/dev/null 2>&1
235+
$cdb2sql "DROP TABLE loop_ins_dst" >/dev/null 2>&1
236+
237+
# ---- Test 6: consumer:next() + DELETE in a loop (multiple events) ----
238+
239+
$cdb2sql <<'EOF' >/dev/null 2>&1
240+
CREATE TABLE loop_del_src(i int)$$
241+
CREATE TABLE loop_del_tbl(i int)$$
242+
CREATE PROCEDURE next_loop_del VERSION 'bar' {
243+
local function main()
244+
local consumer = db:consumer()
245+
while true do
246+
local event = consumer:get()
247+
db:begin()
248+
consumer:next()
249+
local stmt = db:prepare("DELETE FROM loop_del_tbl WHERE i=@i")
250+
stmt:bind("i", event.new.i)
251+
local rc = stmt:exec()
252+
if rc ~= 0 then
253+
db:rollback()
254+
return -201, db:error()
255+
else
256+
db:commit()
257+
end
258+
end
259+
end
260+
}$$
261+
CREATE LUA CONSUMER next_loop_del ON (TABLE loop_del_src FOR INSERT);
262+
EOF
263+
264+
$cdb2sql "INSERT INTO loop_del_tbl VALUES(1),(2),(3),(4),(5)" >/dev/null 2>&1
265+
$cdb2sql "INSERT INTO loop_del_src VALUES(2),(4)" >/dev/null 2>&1
266+
267+
sleep 2
268+
269+
$cdb2sql "EXEC PROCEDURE next_loop_del()" >/dev/null 2>&1 &
270+
CONSUMER_PID=$!
271+
272+
sleep 4
273+
274+
kill $CONSUMER_PID 2>/dev/null
275+
wait $CONSUMER_PID 2>/dev/null
276+
277+
echo "test6 remaining rows:"
278+
$cdb2sql "SELECT i FROM loop_del_tbl ORDER BY i"
279+
280+
echo "test6 queue depth:"
281+
$cdb2sql "SELECT depth FROM comdb2_queues WHERE spname='next_loop_del'"
282+
283+
$cdb2sql "DROP LUA CONSUMER next_loop_del" >/dev/null 2>&1
284+
$cdb2sql "DROP TABLE loop_del_src" >/dev/null 2>&1
285+
$cdb2sql "DROP TABLE loop_del_tbl" >/dev/null 2>&1
286+
287+
) > t10.output 2>&1
288+
289+
set -x
290+
diff -q t10.output t10.expected
291+
if [[ $? -ne 0 ]]; then
292+
diff t10.output t10.expected | head -30
293+
exit 1
294+
fi
295+
echo "passed t10"
296+
exit 0

0 commit comments

Comments
 (0)