1#!/usr/bin/env python3
2"""Race independent Rust clients and verify every observed append against committed operations."""
3import argparse
4from collections import Counter
5from contextlib import contextmanager
6import hashlib
7import json
8import os
9from pathlib import Path
10import subprocess
11import time
12
13ROOT = Path(__file__).resolve().parent.parent
14CLIENT = ROOT / 'target/debug/examples/concurrent_client'
15
16
17def 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
73def 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
127def 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
162if __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)