|
| 1 | +# -*- coding: utf-8 -*- |
| 2 | + |
| 3 | +from __future__ import print_function |
| 4 | + |
| 5 | +import sys |
| 6 | +import time |
| 7 | +import unittest |
| 8 | +import warnings |
| 9 | +from time import sleep |
| 10 | + |
| 11 | +import tarantool |
| 12 | +from tarantool.error import ClusterTolopogyError, DatabaseError, NetworkError |
| 13 | + |
| 14 | +from .lib.skip import skip_or_run_sql_test, skip_or_run_conn_pool_test |
| 15 | +from .lib.tarantool_server import TarantoolServer |
| 16 | + |
| 17 | + |
| 18 | +def create_server(_id): |
| 19 | + srv = TarantoolServer() |
| 20 | + srv.script = 'test/suites/box.lua' |
| 21 | + srv.start() |
| 22 | + srv.admin("box.schema.user.create('test', {password = 'test', " + |
| 23 | + "if_not_exists = true})") |
| 24 | + srv.admin("box.schema.user.grant('test', 'execute', 'universe')") |
| 25 | + srv.admin("box.schema.space.create('test')") |
| 26 | + srv.admin(r"box.space.test:format({" |
| 27 | + +r" { name = 'pk', type = 'string' }," + |
| 28 | + r" { name = 'id', type = 'number', is_nullable = true }" + |
| 29 | + r"})") |
| 30 | + srv.admin(r"box.space.test:create_index('pk'," + |
| 31 | + r"{ unique = true," + |
| 32 | + r" parts = {{field = 1, type = 'string'}}})") |
| 33 | + srv.admin(r"box.space.test:create_index('id'," + |
| 34 | + r"{ unique = true," + |
| 35 | + r" parts = {{field = 2, type = 'number', is_nullable=true}}})") |
| 36 | + srv.admin("box.schema.user.grant('test', 'read,write', 'space', 'test')") |
| 37 | + |
| 38 | + # Create srv_id function (for testing purposes). |
| 39 | + srv.admin("function srv_id() return %s end" % _id) |
| 40 | + return srv |
| 41 | + |
| 42 | + |
| 43 | +@unittest.skipIf(sys.platform.startswith("win"), |
| 44 | + 'Pool tests on windows platform are not supported') |
| 45 | +class TestSuite_Pool(unittest.TestCase): |
| 46 | + def set_ro(self, srv, read_only): |
| 47 | + if read_only: |
| 48 | + req = r'box.cfg{read_only = true}' |
| 49 | + else: |
| 50 | + req = r'box.cfg{read_only = false}' |
| 51 | + |
| 52 | + srv.admin(req) |
| 53 | + |
| 54 | + def set_cluster_ro(self, read_only_list): |
| 55 | + assert len(self.servers) == len(read_only_list) |
| 56 | + |
| 57 | + for i in range(len(self.servers)): |
| 58 | + self.set_ro(self.servers[i], read_only_list[i]) |
| 59 | + |
| 60 | + def retry(self, func, count=5, timeout=0.5): |
| 61 | + for i in range(count): |
| 62 | + try: |
| 63 | + func() |
| 64 | + except Exception as e: |
| 65 | + if i + 1 == count: |
| 66 | + raise e |
| 67 | + |
| 68 | + time.sleep(timeout) |
| 69 | + |
| 70 | + @classmethod |
| 71 | + def setUpClass(self): |
| 72 | + print(' POOL '.center(70, '='), file=sys.stderr) |
| 73 | + print('-' * 70, file=sys.stderr) |
| 74 | + |
| 75 | + @skip_or_run_conn_pool_test |
| 76 | + def setUp(self): |
| 77 | + # Create five servers and extract helpful fields for tests. |
| 78 | + self.servers = [] |
| 79 | + self.addrs = [] |
| 80 | + for i in range(5): |
| 81 | + srv = create_server(i) |
| 82 | + self.servers.append(srv) |
| 83 | + self.addrs.append({'host': srv.host, 'port': srv.args['primary']}) |
| 84 | + |
| 85 | + def test_00_basic(self): |
| 86 | + self.set_cluster_ro([False, False, True, False, True]) |
| 87 | + |
| 88 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 89 | + |
| 90 | + self.assertSequenceEqual(self.pool.eval('return box.info().ro', mode=tarantool.Mode.RW), [False]) |
| 91 | + self.assertSequenceEqual(self.pool.eval('return box.info().ro', mode=tarantool.Mode.PREFER_RW), [False]) |
| 92 | + self.assertSequenceEqual(self.pool.eval('return box.info().ro', mode=tarantool.Mode.PREFER_RO), [True]) |
| 93 | + |
| 94 | + def test_01_roundrobin(self): |
| 95 | + self.set_cluster_ro([False, False, True, False, True]) |
| 96 | + RW_ports = set([str(self.addrs[0]['port']), str(self.addrs[1]['port']), str(self.addrs[3]['port'])]) |
| 97 | + RO_ports = set([str(self.addrs[2]['port']), str(self.addrs[4]['port'])]) |
| 98 | + all_ports = set() |
| 99 | + for addr in self.addrs: |
| 100 | + all_ports.add(str(addr['port'])) |
| 101 | + |
| 102 | + self.pool = tarantool.ConnectionPool( |
| 103 | + addrs=self.addrs, |
| 104 | + user='test', |
| 105 | + password='test', |
| 106 | + refresh_delay=0.2) |
| 107 | + |
| 108 | + # Expect RW iterate through all RW instances. |
| 109 | + RW_ports_result = set() |
| 110 | + for i in range(len(self.servers)): |
| 111 | + resp = self.pool.eval('return box.cfg.listen', mode=tarantool.Mode.RW) |
| 112 | + RW_ports_result.add(resp.data[0]) |
| 113 | + |
| 114 | + self.assertSetEqual(RW_ports_result, RW_ports) |
| 115 | + |
| 116 | + # Expect PREFER_RW iterate through all RW instances if there is at least one. |
| 117 | + PREFER_RW_ports_result = set() |
| 118 | + for i in range(len(self.servers)): |
| 119 | + resp = self.pool.eval('return box.cfg.listen', mode=tarantool.Mode.PREFER_RW) |
| 120 | + PREFER_RW_ports_result.add(resp.data[0]) |
| 121 | + |
| 122 | + self.assertSetEqual(PREFER_RW_ports_result, RW_ports) |
| 123 | + |
| 124 | + # Expect PREFER_RO iterate through all RO instances if there is at least one. |
| 125 | + PREFER_RO_ports_result = set() |
| 126 | + for i in range(len(self.servers)): |
| 127 | + resp = self.pool.eval('return box.cfg.listen', mode=tarantool.Mode.PREFER_RO) |
| 128 | + PREFER_RO_ports_result.add(resp.data[0]) |
| 129 | + |
| 130 | + self.assertSetEqual(PREFER_RO_ports_result, RO_ports) |
| 131 | + |
| 132 | + # Expect PREFER_RW iterate through all instances if there are no RW. |
| 133 | + self.set_cluster_ro([True, True, True, True, True]) |
| 134 | + |
| 135 | + def expect_PREFER_RW_iterate_through_all_instances_if_there_are_no_RW(): |
| 136 | + PREFER_RW_ports_result_all_ro = set() |
| 137 | + for i in range(len(self.servers)): |
| 138 | + resp = self.pool.eval('return box.cfg.listen', mode=tarantool.Mode.PREFER_RW) |
| 139 | + PREFER_RW_ports_result_all_ro.add(resp.data[0]) |
| 140 | + |
| 141 | + self.assertSetEqual(PREFER_RW_ports_result_all_ro, all_ports) |
| 142 | + |
| 143 | + self.retry(func=expect_PREFER_RW_iterate_through_all_instances_if_there_are_no_RW) |
| 144 | + |
| 145 | + # Expect PREFER_RO iterate through all instances if there are no RO. |
| 146 | + self.set_cluster_ro([False, False, False, False, False]) |
| 147 | + |
| 148 | + def expect_PREFER_RO_iterate_through_all_instances_if_there_are_no_RO(): |
| 149 | + PREFER_RO_ports_result_all_rw = set() |
| 150 | + for i in range(len(self.servers)): |
| 151 | + resp = self.pool.eval('return box.cfg.listen', mode=tarantool.Mode.PREFER_RO) |
| 152 | + PREFER_RO_ports_result_all_rw.add(resp.data[0]) |
| 153 | + |
| 154 | + self.assertSetEqual(PREFER_RO_ports_result_all_rw, all_ports) |
| 155 | + |
| 156 | + self.retry(func=expect_PREFER_RO_iterate_through_all_instances_if_there_are_no_RO) |
| 157 | + |
| 158 | + def test_02_exception_raise(self): |
| 159 | + self.set_cluster_ro([False, False, True, False, True]) |
| 160 | + |
| 161 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 162 | + with self.assertRaises(DatabaseError): |
| 163 | + self.pool.call('non_existing_procedure') |
| 164 | + |
| 165 | + def test_03_insert(self): |
| 166 | + self.set_cluster_ro([False, True, False, True, True]) |
| 167 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 168 | + |
| 169 | + self.assertSequenceEqual(self.pool.insert('test', ['test_05_insert_1', 1]), [['test_05_insert_1', 1]]) |
| 170 | + self.assertSequenceEqual(self.pool.insert('test', ['test_05_insert_2', 2]), [['test_05_insert_2', 2]]) |
| 171 | + self.assertSequenceEqual(self.pool.insert('test', ['test_05_insert_3', 3]), [['test_05_insert_3', 3]]) |
| 172 | + |
| 173 | + conn_0 = tarantool.connect( |
| 174 | + host=self.addrs[0]['host'], |
| 175 | + port=self.addrs[0]['port'], |
| 176 | + user='test', |
| 177 | + password='test') |
| 178 | + conn_2 = tarantool.connect( |
| 179 | + host=self.addrs[2]['host'], |
| 180 | + port=self.addrs[2]['port'], |
| 181 | + user='test', |
| 182 | + password='test') |
| 183 | + |
| 184 | + for key in ['test_05_insert_1', 'test_05_insert_2', 'test_05_insert_3']: |
| 185 | + resp_0 = conn_0.select('test', key) |
| 186 | + resp_2 = conn_2.select('test', key) |
| 187 | + self.assertEqual(len(resp_0) + len(resp_2), 1) |
| 188 | + |
| 189 | + def test_04_delete(self): |
| 190 | + self.set_cluster_ro([True, True, True, False, True]) |
| 191 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 192 | + |
| 193 | + conn_3 = tarantool.connect( |
| 194 | + host=self.addrs[3]['host'], |
| 195 | + port=self.addrs[3]['port'], |
| 196 | + user='test', |
| 197 | + password='test') |
| 198 | + |
| 199 | + conn_3.insert('test', ['test_06_delete_1', 1]) |
| 200 | + conn_3.insert('test', ['test_06_delete_2', 2]) |
| 201 | + |
| 202 | + self.assertSequenceEqual(self.pool.delete('test', 'test_06_delete_1'), [['test_06_delete_1', 1]]) |
| 203 | + self.assertSequenceEqual(conn_3.select('test', 'test_06_delete_1'), []) |
| 204 | + |
| 205 | + self.assertSequenceEqual(self.pool.delete('test', 2, index='id'), [['test_06_delete_2', 2]]) |
| 206 | + self.assertSequenceEqual(conn_3.select('test', 'test_06_delete_2'), []) |
| 207 | + |
| 208 | + def test_05_upsert(self): |
| 209 | + self.set_cluster_ro([True, False, True, True, True]) |
| 210 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 211 | + |
| 212 | + conn_1 = tarantool.connect( |
| 213 | + host=self.addrs[1]['host'], |
| 214 | + port=self.addrs[1]['port'], |
| 215 | + user='test', |
| 216 | + password='test') |
| 217 | + |
| 218 | + self.assertSequenceEqual(self.pool.upsert('test', ['test_07_upsert', 3], [('+', 1, 1)]), []) |
| 219 | + self.assertSequenceEqual(conn_1.select('test', 'test_07_upsert'), [['test_07_upsert', 3]]) |
| 220 | + self.assertSequenceEqual(self.pool.upsert('test', ['test_07_upsert', 3], [('+', 1, 1)]), []) |
| 221 | + self.assertSequenceEqual(conn_1.select('test', 'test_07_upsert'), [['test_07_upsert', 4]]) |
| 222 | + |
| 223 | + def test_06_update(self): |
| 224 | + self.set_cluster_ro([True, True, True, True, False]) |
| 225 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 226 | + |
| 227 | + conn_4 = tarantool.connect( |
| 228 | + host=self.addrs[4]['host'], |
| 229 | + port=self.addrs[4]['port'], |
| 230 | + user='test', |
| 231 | + password='test') |
| 232 | + conn_4.insert('test', ['test_08_update', 3]) |
| 233 | + |
| 234 | + self.assertSequenceEqual(self.pool.update('test', ('test_08_update',), [('+', 1, 1)]), [['test_08_update', 4]]) |
| 235 | + self.assertSequenceEqual(conn_4.select('test', 'test_08_update'), [['test_08_update', 4]]) |
| 236 | + |
| 237 | + def test_07_replace(self): |
| 238 | + self.set_cluster_ro([True, True, True, True, False]) |
| 239 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 240 | + |
| 241 | + conn_4 = tarantool.connect( |
| 242 | + host=self.addrs[4]['host'], |
| 243 | + port=self.addrs[4]['port'], |
| 244 | + user='test', |
| 245 | + password='test') |
| 246 | + conn_4.insert('test', ['test_09_replace', 3]) |
| 247 | + |
| 248 | + self.assertSequenceEqual(self.pool.update('test', ('test_09_replace',) , [('+', 1, 1)]), [['test_09_replace', 4]]) |
| 249 | + self.assertSequenceEqual(conn_4.select('test', 'test_09_replace'), [['test_09_replace', 4]]) |
| 250 | + |
| 251 | + def test_08_select(self): |
| 252 | + self.set_cluster_ro([False, False, False, False, False]) |
| 253 | + |
| 254 | + for addr in self.addrs: |
| 255 | + conn = tarantool.connect( |
| 256 | + host=addr['host'], |
| 257 | + port=addr['port'], |
| 258 | + user='test', |
| 259 | + password='test') |
| 260 | + conn.insert('test', ['test_10_select', 3]) |
| 261 | + |
| 262 | + self.set_cluster_ro([False, True, False, True, True]) |
| 263 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 264 | + |
| 265 | + self.assertSequenceEqual(self.pool.select('test', 'test_10_select'), [['test_10_select', 3]]) |
| 266 | + self.assertSequenceEqual(self.pool.select('test', ['test_10_select']), [['test_10_select', 3]]) |
| 267 | + self.assertSequenceEqual(self.pool.select('test', 3, index='id'), [['test_10_select', 3]]) |
| 268 | + self.assertSequenceEqual(self.pool.select('test', [3], index='id'), [['test_10_select', 3]]) |
| 269 | + |
| 270 | + def test_09_ping(self): |
| 271 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 272 | + |
| 273 | + self.assertTrue(self.pool.ping() > 0) |
| 274 | + self.assertEqual(self.pool.ping(notime=True), "Success") |
| 275 | + |
| 276 | + def test_10_call(self): |
| 277 | + self.set_cluster_ro([False, True, False, True, True]) |
| 278 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 279 | + |
| 280 | + self.assertEqual(self.pool.call('box.info')[0]['ro'], False) |
| 281 | + self.assertEqual(self.pool.call('box.info', mode=tarantool.Mode.RW)[0]['ro'], False) |
| 282 | + self.assertEqual(self.pool.call('box.info', mode=tarantool.Mode.PREFER_RO)[0]['ro'], True) |
| 283 | + |
| 284 | + @skip_or_run_sql_test |
| 285 | + def test_11_execute(self): |
| 286 | + self.set_cluster_ro([False, True, True, True, True]) |
| 287 | + self.pool = tarantool.ConnectionPool(addrs=self.addrs, user='test', password='test') |
| 288 | + |
| 289 | + resp = self.pool.execute( |
| 290 | + 'insert into "test" values (:pk, :id)', |
| 291 | + { 'pk': 'test_11_execute', 'id': 1}) |
| 292 | + self.assertEqual(resp.affected_row_count, 1) |
| 293 | + self.assertEqual(resp.data, None) |
| 294 | + |
| 295 | + conn_0 = tarantool.connect( |
| 296 | + host=self.addrs[0]['host'], |
| 297 | + port=self.addrs[0]['port'], |
| 298 | + user='test', |
| 299 | + password='test') |
| 300 | + |
| 301 | + self.assertSequenceEqual(conn_0.select('test', 'test_11_execute'), [['test_11_execute', 1]]) |
| 302 | + |
| 303 | + def test_12_failover(self): |
| 304 | + self.set_cluster_ro([False, True, True, True, True]) |
| 305 | + self.pool = tarantool.ConnectionPool( |
| 306 | + addrs=self.addrs, |
| 307 | + user='test', |
| 308 | + password='test', |
| 309 | + refresh_delay=0.2) |
| 310 | + |
| 311 | + # Simulate failover |
| 312 | + self.servers[0].stop() |
| 313 | + self.set_ro(self.servers[1], False) |
| 314 | + |
| 315 | + def expect_RW_request_execute_on_new_master(): |
| 316 | + self.assertSequenceEqual( |
| 317 | + self.pool.eval('return box.cfg.listen', mode=tarantool.Mode.RW), |
| 318 | + [ str(self.addrs[1]['port']) ]) |
| 319 | + |
| 320 | + self.retry(func=expect_RW_request_execute_on_new_master) |
| 321 | + |
| 322 | + def test_13_cluster_with_instances_dead_in_runtime_is_ok(self): |
| 323 | + self.set_cluster_ro([False, True, False, True, True]) |
| 324 | + self.servers[0].stop() |
| 325 | + |
| 326 | + self.pool = tarantool.ConnectionPool( |
| 327 | + addrs=self.addrs, |
| 328 | + user='test', |
| 329 | + password='test', |
| 330 | + refresh_delay=0.2) |
| 331 | + |
| 332 | + self.pool.ping(mode=tarantool.Mode.RW) |
| 333 | + |
| 334 | + def test_14_cluster_with_dead_instances_on_start_is_ok(self): |
| 335 | + self.set_cluster_ro([False, True, True, True, True]) |
| 336 | + self.servers[0].stop() |
| 337 | + |
| 338 | + self.pool = tarantool.ConnectionPool( |
| 339 | + addrs=self.addrs, |
| 340 | + user='test', |
| 341 | + password='test', |
| 342 | + refresh_delay=0.2) |
| 343 | + |
| 344 | + self.servers[0].start() |
| 345 | + |
| 346 | + def ping_RW(): |
| 347 | + self.pool.ping(mode=tarantool.Mode.RW) |
| 348 | + |
| 349 | + self.retry(func=ping_RW) |
| 350 | + |
| 351 | + def tearDown(self): |
| 352 | + if hasattr(self, 'pool'): |
| 353 | + self.pool.close() |
| 354 | + |
| 355 | + for srv in self.servers: |
| 356 | + srv.stop() |
| 357 | + srv.clean() |
0 commit comments