test_kyc_refresh.py (28999B)
1 #!/usr/bin/env python3 2 3 # This file is part of TALER 4 # Copyright (C) 2026 Taler Systems SA 5 # 6 # TALER is free software; you can redistribute it and/or modify it under the 7 # terms of the GNU General Public License as published by the Free Software 8 # Foundation; either version 3, or (at your option) any later version. 9 # 10 # TALER is distributed in the hope that it will be useful, but WITHOUT ANY 11 # WARRANTY; without even the implied warranty of MERCHANTABILITY or FITNESS FOR 12 # A PARTICULAR PURPOSE. See the GNU General Public License for more details. 13 # 14 # You should have received a copy of the GNU General Public License along with 15 # TALER; see the file COPYING. If not, see <http://www.gnu.org/licenses/> 16 17 """Exercise KYC refreshes with real PostgreSQL notifications and merchant HTTPD. 18 19 Requires PostgreSQL server tools and Python psycopg2. All database state and 20 HTTP listeners belong to the fixture. No installed database is used. 21 """ 22 23 import base64 24 from concurrent.futures import ThreadPoolExecutor 25 import hashlib 26 import json 27 from pathlib import Path 28 import select 29 import shutil 30 import socket 31 import struct 32 import subprocess 33 import sys 34 import tempfile 35 import time 36 import unittest 37 from urllib.parse import urlencode 38 from urllib.request import Request, urlopen 39 40 from test_order_sequence_migrations import Database, postgres_cluster, run 41 42 43 ALPHABET = str.maketrans("ABCDEFGHIJKLMNOPQRSTUVWXYZ234567", 44 "0123456789ABCDEFGHJKMNPQRSTVWXYZ") 45 EXCHANGE = "http://localhost:1/" 46 OTHER_EXCHANGE = "http://localhost:2/" 47 # Synthetic exchange master key (private seed is 31 zero bytes followed by 1) 48 # and its signature over the account below with no restrictions or gateways. 49 MASTER_PUB = bytes.fromhex( 50 "4cb5abf6ad79fbf5abbccafcc269d85cd2651ed4b885b5869f241aedf0a5ba29") 51 ACCOUNT_SIG = bytes.fromhex( 52 "42063a0cc9152e9896bb23b0446e14d9bf6e072fbbdcbf3097ac640eba3744ee" 53 "817368fe93ae567c98ad8f5c27ab5022c82bf73cd024a52baa4e2920bca01806") 54 FOREVER = 9223372036854775807 55 WIRE = bytes([42]) * 64 56 TOKEN = bytes([43]) * 32 57 58 59 def encode(data): 60 return base64.b32encode(data).decode().rstrip("=").translate(ALPHABET) 61 62 63 def channel(event_type, body=b""): 64 header = struct.pack("!HH", 4 + len(body), event_type) 65 return "x" + encode(hashlib.sha512(header + body).digest()[:32]).lower() 66 67 68 def public_key(instance): 69 return instance.to_bytes(32, "big") 70 71 72 # Use the existing protocol event's documented channel, independent of its enum. 73 REFRESH = "xj40p0cfmz0dt6sfz70vrq19kg1hp6aj1q6vcdzn4n2fgpsag4kdg" 74 75 76 class KycRefresh(unittest.TestCase): 77 def sql(self, statement, params=()): 78 with self.db.cursor() as cur: 79 cur.execute(statement, params or None) 80 return cur.fetchall() if cur.description else [] 81 82 def events(self, wait=0): 83 deadline = time.monotonic() + wait 84 events = [] 85 while True: 86 self.listener.poll() 87 events.extend(self.listener.notifies) 88 self.listener.notifies.clear() 89 remaining = deadline - time.monotonic() 90 if remaining <= 0: 91 return events 92 select.select([self.listener], [], [], remaining) 93 94 def setUp(self): 95 self.db = self.connect() 96 self.listener = self.connect() 97 self.listener.autocommit = True 98 with self.listener.cursor() as cur: 99 for event in [REFRESH] + [channel(t, public_key(i) + suffix) 100 for i in (1, 2) 101 for t, suffix in ((1115, b""), (1113, WIRE))]: 102 cur.execute('LISTEN "' + event + '"') 103 self.sql("SET search_path TO merchant_instance_1") 104 # Direct reset avoids generating status events before each scenario. 105 for i in (1, 2): 106 self.sql(f'SET search_path TO merchant_instance_{i}') 107 self.sql("DELETE FROM merchant_tos_accepted") 108 self.sql(f"UPDATE merchant_instance_{i}.merchant_kyc SET " 109 "kyc_ok=true, exchange_http_status=200, exchange_ec_code=0, " 110 "aml_review=false, jaccount_limits='[]', last_rule_gen=1, " 111 "access_token=%s, next_kyc_poll=%s, kyc_backoff=0", 112 (TOKEN, FOREVER)) 113 self.sql('SET search_path TO merchant_instance_1') 114 self.events() 115 116 def tearDown(self): 117 self.listener.close() 118 self.db.close() 119 120 def read(self, refresh, instance=1, exchange=EXCHANGE, wire=WIRE): 121 self.sql(f"SET search_path TO merchant_instance_{instance}") 122 return self.sql("SELECT * FROM merchant_do_account_kyc_get_status" 123 "(%s,%s,%s,%s,%s)", 124 (time.time_ns() // 1000, exchange, wire, refresh, instance)) 125 126 def store(self, instance=1, exchange=EXCHANGE, **changes): 127 values = dict(http=200, ec=0, token=TOKEN, limits="[]", aml=False, 128 ok=True, generation=1, due=FOREVER, backoff=0) 129 values.update(changes) 130 self.sql(f"SET search_path TO merchant_instance_{instance}") 131 return self.sql("SELECT * FROM merchant_do_account_kyc_set_status" 132 "(%s,%s,%s,%s,%s,%s,%s::jsonb,%s,%s,%s,%s,%s,%s,%s)", 133 (WIRE, exchange, int(time.time()) * 1000000, 134 values['http'], values['ec'], values['token'], 135 values['limits'], values['aml'], values['ok'], 136 channel(1113, public_key(instance) + WIRE), 137 channel(1115, public_key(instance)), values['generation'], 138 values['due'], values['backoff'])) 139 140 def test_read_only_and_transactional_refresh(self): 141 self.assertEqual(1, len(self.read(False))) 142 self.assertEqual([], self.events()) 143 self.assertEqual([(FOREVER,)], self.sql("SELECT next_kyc_poll FROM merchant_kyc")) 144 self.assertIsNone(self.read(True, exchange="http://unmatched.invalid/")[0][2]) 145 self.assertEqual([], self.events()) 146 self.db.autocommit = False 147 self.read(True) 148 self.assertEqual([], self.events()) 149 self.db.rollback() 150 self.assertEqual([], self.events()) 151 self.db.autocommit = True 152 self.read(True) 153 events = self.events(0.1) 154 self.assertEqual([(REFRESH, encode(struct.pack('!Q', 1)))], 155 [(e.channel, e.payload) for e in events]) 156 due = self.sql("SELECT next_kyc_poll FROM merchant_kyc") 157 self.read(False) 158 self.assertEqual(due, self.sql("SELECT next_kyc_poll FROM merchant_kyc")) 159 self.assertEqual([], self.events()) 160 161 def test_status_change_notifications(self): 162 self.store(due=17, backoff=123) 163 self.assertEqual([], self.events()) 164 changes = dict(http=202, ec=123, token=None, limits=None, aml=True, 165 ok=False, generation=2) 166 current = {} 167 for field, value in changes.items(): 168 with self.subTest(field=field): 169 current[field] = value 170 self.store(**current) 171 self.assertEqual({channel(1113, public_key(1) + WIRE), 172 channel(1115, public_key(1))}, 173 {e.channel for e in self.events(0.1)}) 174 self.store(**current) 175 self.assertEqual([], self.events()) 176 self.store() 177 self.assertEqual(2, len(self.events(0.1))) # Includes NULL -> value. 178 self.sql("DELETE FROM merchant_kyc") 179 self.store() 180 self.assertEqual(2, len(self.events(0.1))) 181 args = (WIRE, EXCHANGE, int(time.time()) * 1000000, 502, False, 182 channel(1113, public_key(1) + WIRE), channel(1115, public_key(1))) 183 for expected in (2, 0): 184 self.sql("SELECT * FROM merchant_do_account_kyc_set_failed" 185 "(%s,%s,%s,%s,%s,%s,%s)", args) 186 self.assertEqual(expected, len(self.events(0.1))) 187 188 def test_refresh_only_accesses_target_schema(self): 189 self.read(True, instance=1) 190 self.read(True, instance=2) 191 self.db.autocommit = False 192 rows = self.sql("SELECT * FROM merchant.account_kyc_get_outdated(%s,1)", 193 (time.time_ns() // 1000,)) 194 self.assertEqual([('test-1', WIRE, EXCHANGE)], 195 [(i, bytes(w), e) for i, w, e in rows]) 196 touched = self.sql("SELECT DISTINCT n.nspname FROM pg_locks l " 197 "JOIN pg_class c ON c.oid=l.relation " 198 "JOIN pg_namespace n ON n.oid=c.relnamespace " 199 "WHERE l.pid=pg_backend_pid() " 200 "AND n.nspname LIKE 'merchant_instance_%'") 201 self.assertEqual([('merchant_instance_1',)], touched) 202 self.db.rollback() 203 self.db.autocommit = True 204 self.assertEqual([], self.sql( 205 "SELECT * FROM merchant.account_kyc_get_outdated(%s,999999)", (FOREVER,))) 206 self.sql("SET search_path TO merchant_instance_1") 207 self.sql("UPDATE merchant_accounts SET active=false") 208 try: 209 self.assertEqual([], self.sql( 210 "SELECT * FROM merchant.account_kyc_get_outdated(%s,1)", (FOREVER,))) 211 finally: 212 self.sql("UPDATE merchant_accounts SET active=true") 213 214 def get(self, instance=1, query=""): 215 with urlopen(self.url + f"instances/test-{instance}/private/kyc" + query, 216 timeout=8) as response: 217 return json.load(response), response.headers['ETag'] 218 219 def wait_for_refresh(self): 220 deadline = time.monotonic() + 5 221 while time.monotonic() < deadline: 222 if any(e.channel == REFRESH for e in self.events(0.05)): 223 return 224 self.fail("HTTP request did not prompt its initial refresh") 225 226 def test_long_poll_rereads_do_not_refresh(self): 227 self.check_long_poll() 228 229 def test_account_long_poll_rereads_do_not_refresh(self): 230 self.check_long_poll("&h_wire=" + encode(WIRE)) 231 232 def check_long_poll(self, account_filter=""): 233 body, etag = self.get() 234 self.assertEqual('ready', body['kyc_data'][0]['status']) 235 self.wait_for_refresh() 236 with ThreadPoolExecutor() as pool: 237 pending = pool.submit(self.get, 1, "?timeout_ms=5000&lp_not_etag=" 238 + etag.strip('"') + account_filter) 239 self.wait_for_refresh() 240 # A real status change that does not affect the response ETag forces 241 # an internal reread, which must not request another exchange check. 242 self.store(generation=2) 243 self.assertFalse(any(e.channel == REFRESH for e in self.events(0.25))) 244 self.assertFalse(pending.done()) 245 # Another instance's change must not wake this long poll either. 246 self.store(instance=2, http=202, ok=False) 247 self.assertFalse(any(e.channel == REFRESH for e in self.events(0.25))) 248 self.assertFalse(pending.done()) 249 self.store(http=202, ok=False, generation=2) 250 changed, _ = pending.result(timeout=3) 251 self.assertEqual('kyc-required', changed['kyc_data'][0]['status']) 252 self.assertFalse(any(e.channel == REFRESH for e in self.events(0.1))) 253 self.get() 254 self.wait_for_refresh() # Separate HTTP requests still refresh. 255 256 def install_keys(self, exchange=EXCHANGE, **changes): 257 keys = dict( 258 version='33:0:0', currency='EUR', asset_type='fiat', 259 master_public_key=encode(MASTER_PUB), 260 reserve_closing_delay={'d_us': 3600000000}, 261 list_issue_date={'t_s': int(time.time())}, 262 global_fees=[], signkeys=[], denominations=[], auditors=[], 263 wire_fees={}, wads=[], kyc_enabled=True, kyc_swap_tos_acceptance=False, 264 hard_limits=[], zero_limits=[], stefan_abs='EUR:0', stefan_log='EUR:0', 265 stefan_lin=0.0, currency_specification=dict( 266 name='Euro', num_fractional_input_digits=2, 267 num_fractional_normal_digits=2, num_fractional_trailing_zero_digits=2, 268 alt_unit_names={'0': 'EUR'}), 269 accounts=[dict( 270 payto_uri='payto://x-taler-bank/localhost/exchange?receiver-name=Exchange', 271 credit_restrictions=[], debit_restrictions=[], master_sig=encode(ACCOUNT_SIG))]) 272 keys.update(changes) 273 self.sql("INSERT INTO merchant.merchant_exchange_keys " 274 "(exchange_url,keys_json,first_retry,expiration_time," 275 "exchange_http_status,exchange_ec_code) VALUES (%s,%s,0,0,200,0) " 276 "ON CONFLICT (exchange_url) DO UPDATE SET keys_json=EXCLUDED.keys_json", 277 (exchange, json.dumps(dict(version=0, exchange_url=exchange, 278 expire={'t_s': int(time.time()) + 3600}, keys=keys)))) 279 self.sql("SELECT pg_notify(%s,%s)", 280 (channel(1110), encode((exchange + '\0').encode()))) 281 282 def test_initial_keys_wake_long_poll(self): 283 query = '?' + urlencode({'exchange_url': OTHER_EXCHANGE}) 284 self.store(instance=2, exchange=OTHER_EXCHANGE, 285 http=404, token=None, limits=None, ok=False, generation=0) 286 try: 287 body, etag = self.get(instance=2, query=query) 288 self.assertTrue(body['kyc_data'][0]['no_keys']) 289 self.wait_for_refresh() 290 with ThreadPoolExecutor() as pool: 291 pending = pool.submit(self.get, 2, query + "&timeout_ms=5000&lp_not_etag=" 292 + etag.strip('"')) 293 self.wait_for_refresh() 294 self.install_keys(exchange=OTHER_EXCHANGE) 295 body, new_etag = pending.result(timeout=2) 296 self.assertNotEqual(etag, new_etag) 297 self.assertFalse(body['kyc_data'][0]['no_keys']) 298 self.assertEqual('kyc-wire-required', body['kyc_data'][0]['status']) 299 self.assertEqual([], self.events(0.1)) 300 finally: 301 self.sql("SET search_path TO merchant_instance_2") 302 self.sql("DELETE FROM merchant_kyc WHERE exchange_url=%s", (OTHER_EXCHANGE,)) 303 304 def test_keys_long_poll(self): 305 self.check_keys_long_poll() 306 307 def test_keys_exchange_long_poll(self): 308 self.check_keys_long_poll("&" + urlencode({'exchange_url': EXCHANGE})) 309 310 def test_keys_account_long_poll(self): 311 self.check_keys_long_poll("&h_wire=" + encode(WIRE)) 312 313 def check_keys_long_poll(self, account_filter=""): 314 status = dict(http=404, token=None, limits=None, ok=False, generation=0) 315 self.store(**status) 316 self.install_keys() 317 # Wait until HTTPD has consumed the key notification. The keys go 318 # through the production deserializer, including signature validation. 319 deadline = time.monotonic() + 5 320 while True: 321 body, etag = self.get() 322 self.wait_for_refresh() 323 data = body['kyc_data'][0] 324 if not data['no_keys'] and not data['kyc_swap_tos_acceptance'] and data['limits'] == []: 325 break 326 self.assertLess(time.monotonic(), deadline) 327 time.sleep(0.02) 328 self.assertEqual('kyc-wire-required', data['status']) 329 self.assertTrue(data['payto_kycauths']) 330 cases = [ 331 ('kyc_swap_tos_acceptance', True), 332 ('hard_limits', [dict(operation_type='DEPOSIT', threshold='EUR:10', 333 timeframe={'d_us': 60000000}, soft_limit=False)]), 334 ('zero_limits', [dict(operation_type='WITHDRAW')]), 335 ('accounts', []), 336 ] 337 changes = {} 338 with ThreadPoolExecutor() as pool: 339 for field, value in cases: 340 with self.subTest(field=field): 341 pending = pool.submit(self.get, 1, "?timeout_ms=5000&lp_not_etag=" 342 + etag.strip('"') + account_filter) 343 self.wait_for_refresh() 344 # Repeated or unrelated keys do not finish this poll or 345 # request another exchange KYC check. 346 self.install_keys(**changes) 347 self.install_keys(exchange=OTHER_EXCHANGE, kyc_swap_tos_acceptance=True) 348 self.assertFalse(any(e.channel == REFRESH for e in self.events(0.15))) 349 self.assertFalse(pending.done()) 350 changes[field] = value 351 self.install_keys(**changes) 352 self.store(**status) # The same cached /kyc-check result. 353 body, new_etag = pending.result(timeout=2) 354 self.assertNotEqual(etag, new_etag) 355 etag = new_etag 356 data = body['kyc_data'][0] 357 if field == 'kyc_swap_tos_acceptance': 358 self.assertTrue(data[field]) 359 elif field == 'hard_limits': 360 self.assertEqual('EUR:10', data['limits'][0]['threshold']) 361 elif field == 'zero_limits': 362 self.assertTrue(data['limits'][1]['disallowed']) 363 else: 364 self.assertEqual('kyc-wire-impossible', data['status']) 365 self.assertFalse(data.get('payto_kycauths')) 366 self.assertEqual([], self.events(0.1)) 367 368 def accept_tos(self, version, instance=1): 369 request = Request(self.url + f"instances/test-{instance}/private/accept-tos-early", 370 data=json.dumps(dict(exchange_url=EXCHANGE, 371 tos_version=version)).encode(), 372 headers={'Content-Type': 'application/json'}) 373 with urlopen(request, timeout=3) as response: 374 self.assertEqual(204, response.status) 375 376 def set_tos(self, version): 377 return self.sql("SELECT merchant_do_set_tos_accepted_early(%s,%s,%s,%s)", 378 (EXCHANGE, version, struct.pack('!HH', 100, 1113) + public_key(1), 379 channel(1115, public_key(1)))) 380 381 def test_tos_notifications_are_transactional_and_scoped(self): 382 # Even accounts without a cached KYC row can report early acceptance. 383 extra_wire = bytes([44]) * 64 384 extra_channel = channel(1113, public_key(1) + extra_wire) 385 self.sql("INSERT INTO merchant_accounts(h_wire,salt,payto_uri,active) " 386 "VALUES(%s,%s,'payto://x-taler-bank/localhost/extra?receiver-name=Test',true)", 387 (extra_wire, bytes(16))) 388 with self.listener.cursor() as cur: 389 cur.execute('LISTEN "' + extra_channel + '"') 390 self.events() 391 try: 392 self.db.autocommit = False 393 self.assertEqual([(True,)], self.set_tos('v1')) 394 self.assertEqual([], self.events(0.1)) 395 self.db.rollback() 396 self.db.autocommit = True 397 self.assertEqual([], self.events()) 398 self.assertEqual([], self.sql("SELECT * FROM merchant_tos_accepted")) 399 for version in ('v1', 'v2', None): 400 self.assertEqual([(True,)], self.set_tos(version)) 401 self.assertEqual({channel(1115, public_key(1)), 402 channel(1113, public_key(1) + WIRE), extra_channel}, 403 {e.channel for e in self.events(0.1)}) 404 self.assertEqual([(False,)], self.set_tos(version)) 405 self.assertEqual([], self.events(0.1)) 406 finally: 407 self.sql("DELETE FROM merchant_accounts WHERE h_wire=%s", (extra_wire,)) 408 409 def test_tos_long_poll(self): 410 self.check_tos_long_poll() 411 412 def test_tos_account_long_poll(self): 413 self.check_tos_long_poll("&h_wire=" + encode(WIRE)) 414 415 def check_tos_long_poll(self, account_filter=""): 416 self.store(http=202, ok=False) 417 self.events(0.1) 418 body, etag = self.get() 419 self.assertNotIn('tos_accepted_early', body['kyc_data'][0]) 420 self.wait_for_refresh() 421 with ThreadPoolExecutor() as pool: 422 for version in ('v1', 'v2', None): 423 with self.subTest(version=version): 424 pending = pool.submit(self.get, 1, "?timeout_ms=5000&lp_not_etag=" 425 + etag.strip('"') + account_filter) 426 self.wait_for_refresh() 427 if version == 'v2': 428 self.accept_tos('v1') # An identical acceptance is silent. 429 self.assertEqual([], self.events(0.1)) 430 self.accept_tos(version or 'v3', instance=2) 431 self.assertFalse(any(e.channel == REFRESH for e in self.events(0.1))) 432 self.assertFalse(pending.done()) 433 if version is None: 434 # The actual 409 callback clears acceptance; the following 435 # unchanged KYC write must not be needed to wake the poll. 436 self.transition('tos-conflict') 437 self.store(http=202, ok=False) 438 else: 439 self.accept_tos(version) 440 body, new_etag = pending.result(timeout=2) 441 self.assertNotEqual(etag, new_etag) 442 etag = new_etag 443 self.assertEqual(version, body['kyc_data'][0].get('tos_accepted_early')) 444 self.assertEqual('kyc-required', body['kyc_data'][0]['status']) 445 self.assertEqual({channel(1115, public_key(1)), 446 channel(1113, public_key(1) + WIRE)}, 447 {e.channel for e in self.events(0.1)}) 448 449 def test_eligibility_loss_wakes_long_poll(self): 450 body, etag = self.get() 451 self.assertEqual('ready', body['kyc_data'][0]['status']) 452 self.wait_for_refresh() 453 with ThreadPoolExecutor() as pool: 454 pending = pool.submit(self.get, 1, "?timeout_ms=5000&lp_not_etag=" 455 + etag.strip('"')) 456 self.wait_for_refresh() 457 self.transition('eligibility') 458 body, new_etag = pending.result(timeout=2) 459 self.assertNotEqual(etag, new_etag) 460 # /keys snapshots are injected into kyccheck, while HTTPD has no 461 # keys in this fixture. Verify that it returns the persisted 462 # ineligibility error and drops the stale ready status/token. 463 self.assertEqual('no-exchange-keys', body['kyc_data'][0]['status']) 464 self.assertEqual(2628, body['kyc_data'][0]['exchange_code']) 465 self.assertNotIn('access_token', body['kyc_data'][0]) 466 self.assertEqual({channel(1115, public_key(1)), 467 channel(1113, public_key(1) + WIRE)}, 468 {e.channel for e in self.events(0.1)}) 469 body, _ = self.get(instance=2) 470 self.assertEqual('ready', body['kyc_data'][0]['status']) 471 472 def test_upgrade_existing_procedure_signature(self): 473 # Procedure synchronization must replace the old three-argument API. 474 self.sql("CREATE FUNCTION merchant_do_account_kyc_get_status(bigint,text,bytea) " 475 "RETURNS integer LANGUAGE sql AS 'SELECT 0'") 476 self.sql('CALL merchant.sync_instance_procedures(1)') 477 self.assertEqual([(None,)], self.sql( 478 "SELECT to_regprocedure('merchant_do_account_kyc_get_status(bigint,text,bytea)')")) 479 self.assertEqual(1, len(self.read(False))) 480 481 482 def main(): 483 try: 484 import psycopg2 485 except ImportError: 486 print("KYC refresh integration tests require Python psycopg2") 487 return 77 488 source, build = (Path(arg).resolve() for arg in sys.argv[1:]) 489 if not shutil.which('pg_config'): 490 return 77 491 bindir = Path(run(['pg_config', '--bindir'])) 492 with postgres_cluster(bindir, max_locks=1024) as env, tempfile.TemporaryDirectory( 493 prefix='merchant-kyc-') as tmp: 494 sql_dir = build / 'src/backenddb/sql-schema' 495 db = Database('talercheck', bindir, env, sql_dir) 496 run(['createdb', 'talercheck'], env=env) 497 db.apply_file(source / 'src/backenddb/sql-schema/versioning.sql') 498 for migration in sorted(sql_dir.glob('merchant-????.sql')): 499 db.apply_file(migration) 500 db.apply_file(sql_dir / 'global_procedures.sql') 501 db.apply_file(sql_dir / 'instance_procedures.sql') 502 for i in range(1, 201): 503 db.sql(f""" 504 INSERT INTO merchant.merchant_instances 505 (merchant_serial, merchant_id, merchant_name, merchant_pub, 506 merchant_priv, address, jurisdiction, default_wire_transfer_delay, 507 default_pay_delay, use_stefan) 508 VALUES ({i}, 'test-{i}', 'Test', decode('{public_key(i).hex()}','hex'), 509 decode('{public_key(i).hex()}','hex'), '{{}}', '{{}}', 1, 1, false); 510 SET search_path TO merchant_instance_{i}; 511 INSERT INTO merchant_instance_{i}.merchant_accounts 512 (h_wire, salt, payto_uri, active) 513 VALUES (decode('{WIRE.hex()}','hex'), decode(repeat('00',16),'hex'), 514 'payto://x-taler-bank/localhost/account-{i}?receiver-name=Test', true); 515 INSERT INTO merchant_instance_{i}.merchant_kyc 516 (account_serial,exchange_url,kyc_timestamp,kyc_ok,access_token, 517 exchange_http_status,jaccount_limits,last_rule_gen,next_kyc_poll) 518 VALUES (1,'{EXCHANGE}',0,true,decode('{TOKEN.hex()}','hex'),200,'[]',1,{FOREVER}); 519 """) 520 def connect(): 521 conn = psycopg2.connect(dbname='talercheck', user='postgres', 522 host=env['PGHOST'], port=env['PGPORT']) 523 conn.autocommit = True 524 return conn 525 KycRefresh.connect = staticmethod(connect) 526 with socket.socket() as sock: 527 sock.bind(('127.0.0.1', 0)) 528 port = sock.getsockname()[1] 529 KycRefresh.url = f'http://127.0.0.1:{port}/' 530 config = Path(tmp, 'merchant.conf') 531 config.write_text(f""" 532 @INLINE@ {source / "src/backend/merchant.conf"} 533 @INLINE@ {source / "src/util/currencies.conf"} 534 [merchant] 535 CURRENCY=EUR 536 SERVE=tcp 537 PORT={port} 538 BIND_TO=127.0.0.1 539 BASE_URL={KycRefresh.url} 540 [merchantdb-postgres] 541 CONFIG=postgres:///talercheck 542 SQL_DIR={sql_dir}/ 543 [taler] 544 CURRENCY=EUR 545 [merchant-exchange-kyc-test] 546 EXCHANGE_BASE_URL={EXCHANGE} 547 CURRENCY=EUR 548 MASTER_KEY={encode(MASTER_PUB)} 549 [merchant-exchange-other-test] 550 EXCHANGE_BASE_URL={OTHER_EXCHANGE} 551 CURRENCY=EUR 552 MASTER_KEY={encode(MASTER_PUB)} 553 [merchant-exchange-kudos] 554 DISABLED=YES 555 [merchant-exchange-chf] 556 DISABLED=YES 557 """) 558 env['LD_LIBRARY_PATH'] = ':'.join(str(build / 'src' / d) 559 for d in ('backenddb', 'util', 'bank', 'lib')) \ 560 + ':' + env.get('LD_LIBRARY_PATH', '') 561 prefix = Path(tmp, 'prefix') 562 resources = prefix / 'share/taler-merchant' 563 resources.mkdir(parents=True) 564 (resources / 'templates').symlink_to(source / 'src/frontend') 565 (resources / 'spa').symlink_to(source / 'contrib/spa') 566 env['TALER_MERCHANT_PREFIX'] = str(prefix) 567 KycRefresh.transition = staticmethod(lambda mode: run( 568 [str(build / 'src/backend/test_merchant_kyccheck'), str(config), mode], env=env)) 569 with open(Path(tmp, 'httpd.log'), 'w+') as log: 570 process = subprocess.Popen([str(build / 'src/backend/taler-merchant-httpd'), 571 '-c', str(config), '-L', 'INFO'], 572 env=env, stdout=log, stderr=log) 573 try: 574 deadline = time.monotonic() + 10 575 while time.monotonic() < deadline: 576 if process.poll() is not None: 577 raise RuntimeError('merchant HTTPD exited at startup') 578 try: 579 with urlopen(KycRefresh.url + 'config', timeout=0.2): 580 break 581 except OSError: 582 time.sleep(0.05) 583 else: 584 raise RuntimeError('merchant HTTPD did not start') 585 suite = unittest.defaultTestLoader.loadTestsFromTestCase(KycRefresh) 586 result = unittest.TextTestRunner(verbosity=2).run(suite) 587 return 0 if result.wasSuccessful() else 1 588 finally: 589 process.terminate() 590 try: 591 process.wait(timeout=10) 592 except subprocess.TimeoutExpired: 593 process.kill() 594 process.wait() 595 log.seek(0) 596 Path(build / 'kyc-refresh-httpd.log').write_text(log.read()) 597 598 599 if __name__ == '__main__': 600 sys.exit(main())