| 1 | #!/usr/bin/env python3 |
| 2 | """Race independent Rust clients and verify every observed append against committed operations.""" |
| 3 | import argparse |
| 4 | from collections import Counter |
| 5 | from contextlib import contextmanager |
| 6 | import hashlib |
| 7 | import json |
| 8 | import os |
| 9 | from pathlib import Path |
| 10 | import subprocess |
| 11 | import time |
| 12 | |
| 13 | ROOT = Path(__file__).resolve().parent.parent |
| 14 | CLIENT = ROOT / 'target/debug/examples/concurrent_client' |
| 15 | |
| 16 | |
| 17 | def verify(logs, initial_transaction, writers, operations, edit=False): |
| 18 | commits = [(actor, event) for actor, events in logs.items() for event in events if event['event'] == 'commit'] |
| 19 | assert Counter(actor for actor, _ in commits) == {f'w{i}': operations for i in range(writers)}, 'Missing or extra acknowledged operations' |
| 20 | ordered = sorted(commits, key=lambda pair: pair[1]['source_transaction']) |
| 21 | versions = {initial_transaction: 'Concurrent edits:'} |
| 22 | seen = set() |
| 23 | for i, (actor, event) in enumerate(ordered): |
| 24 | transaction = initial_transaction + i |
| 25 | assert event['source_transaction'] == transaction, 'Two writers committed the same snapshot or skipped a transaction' |
| 26 | assert event['token'] == f' [{actor}:{event["operation"]}]', 'Committed token differs from intended operation' |
| 27 | assert (actor, event['operation']) not in seen, 'An operation committed twice' |
| 28 | seen.add((actor, event['operation'])) |
| 29 | if edit: |
| 30 | intent, = [e for e in logs[actor] if e['event'] == 'intent' and e['attempt'] == event['attempt']] |
| 31 | assert intent['before'] == versions[transaction] and intent['source_transaction'] == transaction, 'Intent used another snapshot' |
| 32 | assert intent['operation'] == event['operation'] and intent['token'] == event['token'], 'Acknowledgement differs from intent' |
| 33 | assert intent['replacement'] == ' café 🦀' + event['token'], 'Unexpected replacement text' |
| 34 | start, end = intent['range'] |
| 35 | units = intent['before'].encode('utf-16-le') |
| 36 | assert len(versions[initial_transaction]) <= start <= end <= len(units) // 2, 'Invalid edit range' |
| 37 | units[:start * 2].decode('utf-16-le') |
| 38 | units[end * 2:].decode('utf-16-le') |
| 39 | versions[transaction + 1] = (units[:start * 2] + intent['replacement'].encode('utf-16-le') + units[end * 2:]).decode('utf-16-le') |
| 40 | else: |
| 41 | versions[transaction + 1] = versions[transaction] + event['token'] |
| 42 | observations = 0 |
| 43 | intervals = [] |
| 44 | for actor, events in logs.items(): |
| 45 | assert events and events[0]['event'] == 'ready' and events[-1]['event'] == 'done', 'Client did not finish' |
| 46 | previous = initial_transaction |
| 47 | for event in events: |
| 48 | if event['event'] in ('commit', 'retry'): |
| 49 | intervals.append((actor, event['started_us'], event['finished_us'])) |
| 50 | if event['event'] != 'read': |
| 51 | continue |
| 52 | transaction = event['transaction'] |
| 53 | assert transaction >= previous, 'A reader went backwards' |
| 54 | previous = transaction |
| 55 | assert transaction in versions and event['text'] == versions[transaction], 'A reader observed a lost, partial, duplicated or invented edit' |
| 56 | for _, commit in ordered: |
| 57 | published = commit['source_transaction'] + 1 |
| 58 | if commit['finished_us'] < event['started_us']: |
| 59 | assert transaction >= published, 'Read missed an already acknowledged commit' |
| 60 | if commit['started_us'] > event['finished_us']: |
| 61 | assert transaction < published, 'Read observed a future commit' |
| 62 | observations += 1 |
| 63 | overlap = sum(a != b and max(start, other_start) < min(end, other_end) |
| 64 | for i, (a, start, end) in enumerate(intervals) |
| 65 | for b, other_start, other_end in intervals[i + 1:]) |
| 66 | assert overlap, 'Writer calls never overlapped' |
| 67 | return {'commits': len(commits), 'observations': observations, 'overlapping_writer_calls': overlap, |
| 68 | 'retries': sum(e['event'] == 'retry' for events in logs.values() for e in events), |
| 69 | 'final_transaction': initial_transaction + len(commits), 'final_text': versions[max(versions)]} |
| 70 | |
| 71 | |
| 72 | @contextmanager |
| 73 | def running_clients(output, source, writers, readers, operations, seed, timeout=600, edit=False, executable=CLIENT, environment=None, reader_executable=None): |
| 74 | if not 0 < timeout < 2**64 / 1000: |
| 75 | raise ValueError('Choose a finite positive client timeout.') |
| 76 | timeout_ms = int(timeout * 1000) |
| 77 | if not 0 < timeout_ms < 2**64: |
| 78 | raise ValueError('Client timeout does not fit the subprocess clock.') |
| 79 | environment = {**(os.environ if environment is None else environment), 'ONESTORE_CLIENT_TIMEOUT_MS': str(timeout_ms)} |
| 80 | start, stop = output / 'start', output / 'stop' |
| 81 | manifest = {'writers': writers, 'readers': readers, 'operations': operations, |
| 82 | 'seed': seed, 'edit': edit, 'timeout_ms': timeout_ms, 'client_sha256': hashlib.sha256(executable.read_bytes()).hexdigest(), 'executable': str(executable)} |
| 83 | if reader_executable is not None: |
| 84 | manifest.update(reader_executable=str(reader_executable), reader_sha256=hashlib.sha256(reader_executable.read_bytes()).hexdigest()) |
| 85 | (output / 'clients.json').write_text(json.dumps(manifest, indent=2)) |
| 86 | processes = {} |
| 87 | streams = [] |
| 88 | try: |
| 89 | for mode, count in [('write', writers), ('read', readers)]: |
| 90 | client = reader_executable if mode == 'read' and reader_executable is not None else executable |
| 91 | for i in range(count): |
| 92 | actor = mode[0] + str(i) |
| 93 | out = (output / f'{actor}.jsonl').open('w') |
| 94 | err = (output / f'{actor}.stderr').open('w') |
| 95 | streams.extend([out, err]) |
| 96 | processes[actor] = subprocess.Popen([client, 'edit' if edit and mode == 'write' else mode, source, actor, str(operations), start, stop, |
| 97 | str(seed + i + (10000 if mode == 'read' else 0))], stdout=out, stderr=err, env=environment) |
| 98 | deadline = time.monotonic() + timeout |
| 99 | while not all((output / f'{actor}.jsonl').stat().st_size for actor in processes): |
| 100 | if any(p.poll() is not None for p in processes.values()) or time.monotonic() > deadline: |
| 101 | raise RuntimeError('A client failed before the start barrier.') |
| 102 | time.sleep(.01) |
| 103 | yield processes |
| 104 | while any(p.poll() is None for p in processes.values()): |
| 105 | failed = {actor: p.returncode for actor, p in processes.items() if p.returncode not in (None, 0)} |
| 106 | if failed or time.monotonic() > deadline: |
| 107 | raise RuntimeError(f'Concurrent clients failed or timed out: {failed}') |
| 108 | if all(p.poll() == 0 for actor, p in processes.items() if actor.startswith('w')): |
| 109 | stop.touch(exist_ok=True) |
| 110 | time.sleep(.05) |
| 111 | assert all(p.returncode == 0 for p in processes.values()), 'A client exited unsuccessfully' |
| 112 | finally: |
| 113 | for process in processes.values(): |
| 114 | if process.poll() is None: |
| 115 | process.terminate() |
| 116 | for process in processes.values(): |
| 117 | try: |
| 118 | process.wait(timeout=5) |
| 119 | except subprocess.TimeoutExpired: |
| 120 | process.kill() |
| 121 | process.wait() |
| 122 | for stream in streams: |
| 123 | stream.close() |
| 124 | (output / 'teardown.json').write_text(json.dumps({actor: p.returncode for actor, p in processes.items()}, indent=2)) |
| 125 | |
| 126 | |
| 127 | def run(output, writers, readers, operations, seed, notebook=None, timeout=600, edit=False): |
| 128 | if min(writers, readers, operations) < 1 or writers < 2: |
| 129 | raise ValueError('Choose at least two writers, one reader and one operation.') |
| 130 | output = output.resolve() |
| 131 | output.mkdir(parents=True, exist_ok=False) |
| 132 | notebook = output / 'notebook' if notebook is None else notebook.resolve() |
| 133 | notebook.mkdir() |
| 134 | source = notebook / 'synthetic.one' |
| 135 | seed_file = output / 'initial.one' |
| 136 | subprocess.run([CLIENT, 'init', seed_file], check=True) |
| 137 | initial = seed_file.read_bytes() |
| 138 | with source.open('xb') as stream: |
| 139 | stream.write(initial) |
| 140 | stream.flush() |
| 141 | os.fsync(stream.fileno()) |
| 142 | start, stop = output / 'start', output / 'stop' |
| 143 | (output / 'run.json').write_text(json.dumps({'initial_sha256': hashlib.sha256(initial).hexdigest(), |
| 144 | 'generator_sha256': hashlib.sha256(CLIENT.read_bytes()).hexdigest()}, indent=2)) |
| 145 | started = time.monotonic() |
| 146 | with running_clients(output, source, writers, readers, operations, seed, timeout, edit) as processes: |
| 147 | start.touch() |
| 148 | logs = {actor: [json.loads(line) for line in (output / f'{actor}.jsonl').read_text().splitlines()] for actor in processes} |
| 149 | result = verify(logs, int.from_bytes(initial[96:100], 'little'), writers, operations, edit) |
| 150 | final = subprocess.run([CLIENT, 'read', source, 'final', '1', start, stop, str(seed)], check=True, capture_output=True, text=True) |
| 151 | (output / 'final.jsonl').write_text(final.stdout) |
| 152 | final_read, = [json.loads(line) for line in final.stdout.splitlines() if json.loads(line)['event'] == 'read'] |
| 153 | assert final_read['text'] == result['final_text'] and final_read['transaction'] == result['final_transaction'], 'Final file lost an acknowledged edit' |
| 154 | result['elapsed_seconds'] = time.monotonic() - started |
| 155 | final_bytes = source.read_bytes() |
| 156 | result['final_sha256'] = hashlib.sha256(final_bytes).hexdigest() |
| 157 | (output / 'final.one').write_bytes(final_bytes) |
| 158 | (output / 'result.json').write_text(json.dumps(result, indent=2)) |
| 159 | print(json.dumps({k: v for k, v in result.items() if k != 'final_text'}), flush=True) |
| 160 | |
| 161 | |
| 162 | if __name__ == '__main__': |
| 163 | parser = argparse.ArgumentParser(description=__doc__) |
| 164 | parser.add_argument('output', type=Path) |
| 165 | parser.add_argument('--writers', type=int, default=4) |
| 166 | parser.add_argument('--readers', type=int, default=3) |
| 167 | parser.add_argument('--operations', type=int, default=30) |
| 168 | parser.add_argument('--seed', type=int, default=1) |
| 169 | parser.add_argument('--notebook-dir', type=Path) |
| 170 | parser.add_argument('--timeout', type=int, default=600) |
| 171 | parser.add_argument('--edit', action='store_true') |
| 172 | args = parser.parse_args() |
| 173 | run(args.output, args.writers, args.readers, args.operations, args.seed, args.notebook_dir, args.timeout, args.edit) |