merchant

Merchant backend to process payments, run by merchants
Log | Files | Refs | Submodules | README | LICENSE

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())