|
| 1 | +#!/usr/bin/env tarantool |
| 2 | + |
| 3 | +local fio = require('fio') |
| 4 | +local tnt = require('t.tnt') |
| 5 | +local test = require('tap').test('') |
| 6 | +local uuid = require('uuid') |
| 7 | +local queue = require('queue') |
| 8 | +local fiber = require('fiber') |
| 9 | + |
| 10 | +local queue_state = require('queue.abstract.queue_state') |
| 11 | +rawset(_G, 'queue', require('queue')) |
| 12 | + |
| 13 | +-- Replica connection handler |
| 14 | +local conn = {} |
| 15 | + |
| 16 | +test:plan(5) |
| 17 | + |
| 18 | +test:test('Check master-replica setup', function(test) |
| 19 | + test:plan(8) |
| 20 | + local engine = os.getenv('ENGINE') or 'memtx' |
| 21 | + tnt.cluster.cfg{} |
| 22 | + |
| 23 | + test:ok(rawget(box, 'space'), 'box started') |
| 24 | + test:ok(queue, 'queue is loaded') |
| 25 | + |
| 26 | + test:ok(tnt.cluster.wait_replica(), 'wait for replica to connect') |
| 27 | + conn = tnt.cluster.connect_replica() |
| 28 | + test:ok(conn.error == nil, 'no errors on connect to replica') |
| 29 | + test:ok(conn:ping(), 'ping replica') |
| 30 | + test:is(queue.state(), 'RUNNING', 'check master queue state') |
| 31 | + conn:eval('rawset(_G, "queue", require("queue"))') |
| 32 | + test:is(conn:call('queue.state'), 'INIT', 'check replica queue state') |
| 33 | + |
| 34 | + -- Setup tube. Set ttr = 0.5 for sessions expire testing. |
| 35 | + conn:call('queue.cfg', {{ttr = 0.5}}) |
| 36 | + queue.cfg{ttr = 0.5} |
| 37 | + local tube = queue.create_tube('test', 'fifo', {engine = engine}) |
| 38 | + test:ok(tube, 'test tube created') |
| 39 | +end) |
| 40 | + |
| 41 | +test:test('Check queue state switching', function(test) |
| 42 | + test:plan(2) |
| 43 | + box.cfg{read_only = true} |
| 44 | + test:ok(queue_state.poll(queue_state.states.WAITING, 10), |
| 45 | + "queue state changed to waiting") |
| 46 | + box.cfg{read_only = false} |
| 47 | + test:ok(queue_state.poll(queue_state.states.RUNNING, 10), |
| 48 | + "queue state changed to running") |
| 49 | +end) |
| 50 | + |
| 51 | +test:test('Check session resuming', function(test) |
| 52 | + test:plan(17) |
| 53 | + local client = tnt.cluster.connect_master() |
| 54 | + test:ok(client.error == nil, 'no errors on client connect to master') |
| 55 | + local session_uuid = client:call('queue.identify') |
| 56 | + local uuid_obj = uuid.frombin(session_uuid) |
| 57 | + |
| 58 | + test:ok(queue.tube.test:put('testdata'), 'put task') |
| 59 | + local task_master = client:call('queue.tube.test:take') |
| 60 | + test:ok(task_master, 'task was taken') |
| 61 | + test:is(task_master[3], 'testdata', 'task.data') |
| 62 | + client:close() |
| 63 | + |
| 64 | + local qt = box.space._queue_taken_2:select() |
| 65 | + test:is(uuid.frombin(qt[1][4]):str(), uuid_obj:str(), |
| 66 | + 'task taken by actual uuid') |
| 67 | + |
| 68 | + -- wait for disconnect collback |
| 69 | + local attempts = 0 |
| 70 | + while true do |
| 71 | + local is = box.space._queue_inactive_sessions:select() |
| 72 | + |
| 73 | + if is[1] then |
| 74 | + test:is(uuid.frombin(is[1][1]):str(), uuid_obj:str(), |
| 75 | + 'check inactive sessions') |
| 76 | + break |
| 77 | + end |
| 78 | + |
| 79 | + attempts = attempts + 1 |
| 80 | + if attempts == 10 then |
| 81 | + test:ok(false, 'check inactive sessions') |
| 82 | + return false |
| 83 | + end |
| 84 | + fiber.sleep(0.01) |
| 85 | + end |
| 86 | + |
| 87 | + -- switch roles |
| 88 | + box.cfg{read_only = true} |
| 89 | + queue_state.poll(queue_state.states.WAITING, 10) |
| 90 | + test:is(queue.state(), 'WAITING', 'master state is waiting') |
| 91 | + conn:eval('box.cfg{read_only=false}') |
| 92 | + conn:eval([[ |
| 93 | + queue_state = require('queue.abstract.queue_state') |
| 94 | + queue_state.poll(queue_state.states.RUNNING, 10) |
| 95 | + ]]) |
| 96 | + test:is(conn:call('queue.state'), 'RUNNING', 'replica state is running') |
| 97 | + |
| 98 | + local cfg = conn:eval('return queue.cfg') |
| 99 | + test:is(cfg.ttr, 0.5, 'check cfg applied after lazy start') |
| 100 | + |
| 101 | + test:ok(conn:call('queue.identify', {session_uuid}), 'identify old session') |
| 102 | + local stat = conn:call('queue.statistics') |
| 103 | + test:is(stat.test.tasks.taken, 1, 'taken tasks count') |
| 104 | + test:is(stat.test.tasks.done, 0, 'done tasks count') |
| 105 | + local task_replica = conn:call('queue.tube.test:ack', {task_master[1]}) |
| 106 | + test:is(task_replica[3], 'testdata', 'check task data') |
| 107 | + local stat = conn:call('queue.statistics') |
| 108 | + test:is(stat.test.tasks.taken, 0, 'taken tasks count after ack()') |
| 109 | + test:is(stat.test.tasks.done, 1, 'done tasks count after ack()') |
| 110 | + |
| 111 | + -- switch roles back |
| 112 | + conn:eval('box.cfg{read_only=true}') |
| 113 | + conn:eval([[ |
| 114 | + queue_state = require('queue.abstract.queue_state') |
| 115 | + queue_state.poll(queue_state.states.WAITING, 10) |
| 116 | + ]]) |
| 117 | + box.cfg{read_only = false} |
| 118 | + queue_state.poll(queue_state.states.RUNNING, 10) |
| 119 | + test:is(queue.state(), 'RUNNING', 'master state is running') |
| 120 | + test:is(conn:call('queue.state'), 'WAITING', 'replica state is waiting') |
| 121 | +end) |
| 122 | + |
| 123 | +test:test('Check task is cleaned after migrate', function(test) |
| 124 | + test:plan(9) |
| 125 | + local client = tnt.cluster.connect_master() |
| 126 | + local session_uuid = client:call('queue.identify') |
| 127 | + local uuid_obj = uuid.frombin(session_uuid) |
| 128 | + test:ok(queue.tube.test:put('testdata'), 'put task') |
| 129 | + test:ok(client:call('queue.tube.test:take'), 'take task from master') |
| 130 | + client:close() |
| 131 | + |
| 132 | + -- wait for disconnect collback |
| 133 | + local attempts = 0 |
| 134 | + while true do |
| 135 | + local is = box.space._queue_inactive_sessions:select() |
| 136 | + |
| 137 | + if is[1] then |
| 138 | + test:is(uuid.frombin(is[1][1]):str(), uuid_obj:str(), |
| 139 | + 'check inactive sessions') |
| 140 | + break |
| 141 | + end |
| 142 | + |
| 143 | + attempts = attempts + 1 |
| 144 | + if attempts == 10 then |
| 145 | + test:ok(false, 'check inactive sessions') |
| 146 | + return false |
| 147 | + end |
| 148 | + fiber.sleep(0.01) |
| 149 | + end |
| 150 | + |
| 151 | + -- switch roles |
| 152 | + box.cfg{read_only = true} |
| 153 | + |
| 154 | + queue_state.poll(queue_state.states.WAITING, 10) |
| 155 | + test:is(queue.state(), 'WAITING', 'master state is waiting') |
| 156 | + conn:eval('box.cfg{read_only=false}') |
| 157 | + conn:eval([[ |
| 158 | + queue_state = require('queue.abstract.queue_state') |
| 159 | + queue_state.poll(queue_state.states.RUNNING, 10) |
| 160 | + ]]) |
| 161 | + test:is(conn:call('queue.state'), 'RUNNING', 'replica state is running') |
| 162 | + |
| 163 | + -- check task |
| 164 | + local stat = conn:call('queue.statistics') |
| 165 | + test:is(stat.test.tasks.taken, 1, 'taken tasks count before timeout') |
| 166 | + fiber.sleep(1) |
| 167 | + local stat = conn:call('queue.statistics') |
| 168 | + test:is(stat.test.tasks.taken, 0, 'taken tasks count after timeout') |
| 169 | + |
| 170 | + -- switch roles back |
| 171 | + conn:eval('box.cfg{read_only=true}') |
| 172 | + conn:eval([[ |
| 173 | + queue_state = require('queue.abstract.queue_state') |
| 174 | + queue_state.poll(queue_state.states.WAITING, 10) |
| 175 | + ]]) |
| 176 | + box.cfg{read_only = false} |
| 177 | + queue_state.poll(queue_state.states.RUNNING, 10) |
| 178 | + test:is(queue.state(), 'RUNNING', 'master state is running') |
| 179 | + test:is(conn:call('queue.state'), 'WAITING', 'replica state is waiting') |
| 180 | +end) |
| 181 | + |
| 182 | +test:test('Check release_all method', function(test) |
| 183 | + test:plan(6) |
| 184 | + test:ok(queue.tube.test:put('testdata'), 'put task #0') |
| 185 | + test:ok(queue.tube.test:put('testdata'), 'put task #1') |
| 186 | + test:ok(queue.tube.test:take(), 'take task #0') |
| 187 | + test:ok(queue.tube.test:take(), 'take task #1') |
| 188 | + test:is(queue.statistics().test.tasks.taken, 2, |
| 189 | + 'taken tasks count before release_all') |
| 190 | + queue.tube.test:release_all() |
| 191 | + test:is(queue.statistics().test.tasks.taken, 0, |
| 192 | + 'taken tasks count after release_all') |
| 193 | +end) |
| 194 | + |
| 195 | +rawset(_G, 'queue', nil) |
| 196 | +conn:eval('rawset(_G, "queue", nil)') |
| 197 | +conn:close() |
| 198 | +tnt.finish() |
| 199 | +os.exit(test:check() and 0 or 1) |
| 200 | +-- vim: set ft=lua : |
0 commit comments