| 1 | #!/usr/bin/env python3 |
| 2 | """Replay shared notebook edits with disposable OneNote clients and a Linux server.""" |
| 3 | import argparse |
| 4 | import errno |
| 5 | from contextlib import ExitStack, contextmanager |
| 6 | from concurrent.futures import ThreadPoolExecutor |
| 7 | from threading import Lock |
| 8 | import hashlib |
| 9 | import json |
| 10 | import os |
| 11 | from pathlib import Path |
| 12 | import shutil |
| 13 | import signal |
| 14 | import subprocess |
| 15 | import tarfile |
| 16 | import time |
| 17 | import xml.etree.ElementTree as ET |
| 18 | |
| 19 | from native_runner import ROOT, clone, command, windows, vm, collect_artifacts |
| 20 | from native_xml import texts |
| 21 | from document_model import EXPORTER, ordered_pages, view, walk |
| 22 | import linux_vm |
| 23 | |
| 24 | |
| 25 | def reachable_page_text(model): |
| 26 | (sid, _, _, _), = ordered_pages(model) |
| 27 | pending, retained = [sid], {} |
| 28 | while pending: |
| 29 | sid = pending.pop() |
| 30 | if sid in retained: continue |
| 31 | _, revision = view(model, sid) |
| 32 | manifest = revision['nodes'][revision['roots']['1']] |
| 33 | assert manifest['kind']['type'] == 'Manifest' |
| 34 | page, = manifest['content'] |
| 35 | assert revision['nodes'][page]['kind']['type'] == 'Page' |
| 36 | retained[sid] = sorted(node['kind']['text'] for _, node in walk(revision, page) if node['kind']['type'] == 'RichText') |
| 37 | pending.extend(manifest['spaces']) |
| 38 | return retained |
| 39 | |
| 40 | |
| 41 | def verify_final_state(model): |
| 42 | (sid, main), (conflict_sid, competing) = reachable_page_text(model).items() |
| 43 | main, competing = set(main), set(competing) |
| 44 | assert {'Native lock recovery.', 'Client B online cell.'} <= main |
| 45 | edits = {'Client A concurrent edit.', 'Client B concurrent edit.'} |
| 46 | assert len(main & edits) == len(competing & edits) == 1 and edits <= main | competing |
| 47 | return {'main_space': sid, 'conflict_space': conflict_sid, 'competing_edits': sorted(edits)} |
| 48 | |
| 49 | |
| 50 | def replay(output, server, stress_clients=0, stress_operations=30, sync_every=1, rust_writers=4, rust_readers=3, edit=False, seed=710, conflict_clients=0, abrupt=False, embedded_smb=False, maintenance=False, fixture=None, disconnect=False, client_timeout=600, offline=False, offline_outage=False, offline_lost_reply=False, client_profile="debug", document_operations=False, record_writes=False, offline_client_reply=False): |
| 51 | if fixture is not None and not stress_clients: |
| 52 | raise ValueError('Use a fixture with stress mode.') |
| 53 | if linux_vm.instance_path(server).exists(): |
| 54 | raise ValueError('Choose a new Linux VM name; existing machines are not owned by this run.') |
| 55 | output = output.resolve() |
| 56 | output.mkdir(parents=True, exist_ok=False) |
| 57 | for profile in (('debug', 'release') if client_profile == 'release' else ('debug',)): |
| 58 | build = ['cargo', 'build', '--locked', '--workspace', '--all-features', '--examples'] |
| 59 | if profile == 'release': build.append('--release') |
| 60 | with (output / f'build-{profile}.log').open('w') as log: |
| 61 | result = subprocess.run(build, cwd=ROOT, env={**os.environ, 'CARGO_TARGET_DIR': str(ROOT / 'target')}, |
| 62 | stdout=log, stderr=subprocess.STDOUT) |
| 63 | (output / f'build-{profile}.json').write_text(json.dumps({'command': build, 'exit': result.returncode}, indent=2)) |
| 64 | result.check_returncode() |
| 65 | mount = output / 'mount' |
| 66 | mount.mkdir() |
| 67 | scripts = output / 'scripts' |
| 68 | scripts.mkdir() |
| 69 | for name in ('cold.ps1', 'collaborate.ps1', 'network.ps1', 'stress.ps1', 'text.ps1'): |
| 70 | shutil.copyfile(ROOT / 'tools/native' / name, scripts / name) |
| 71 | harness = ('native_collaboration.py', 'native_maintenance.py', 'native_disconnect.py', 'native_stress.py', 'offline_history.py', 'offline_document_history.py', 'offline_outage.py', 'verify_offline.py', 'concurrent_rust.py', 'native_runner.py', |
| 72 | 'crash_recovery.py', 'smb-proxy.py', 'verify_smb_overlap.py', 'w7/crash.py', 'w7/vm.py', 'w7/linux_vm.py') |
| 73 | for name in harness: |
| 74 | saved = output / 'harness' / name |
| 75 | saved.parent.mkdir(parents=True, exist_ok=True) |
| 76 | shutil.copyfile(ROOT / 'tools' / name, saved) |
| 77 | (output / 'run.json').write_text(json.dumps({'server': server, 'stress_clients': stress_clients, 'conflict_clients': conflict_clients, 'stress_operations': stress_operations, 'sync_every': sync_every, 'rust_writers': rust_writers, 'rust_readers': rust_readers, 'edit': edit, 'seed': seed, 'abrupt': abrupt, 'embedded_smb': embedded_smb, 'maintenance': maintenance, 'fixture': str(fixture) if fixture is not None else None, |
| 78 | 'disconnect': disconnect, 'client_timeout': client_timeout, 'client_profile': client_profile, 'offline': offline, 'document_operations': document_operations, 'record_writes': record_writes, 'offline_outage': offline_outage, 'offline_lost_reply': offline_lost_reply, 'offline_client_reply': offline_client_reply, 'harness_sha256': {name: hashlib.sha256((output / 'harness' / name).read_bytes()).hexdigest() for name in harness}, 'scripts': {p.name: hashlib.sha256(p.read_bytes()).hexdigest() for p in scripts.iterdir()}}, indent=2)) |
| 79 | |
| 80 | def ssh(text): |
| 81 | result = linux_vm.run_ssh(server, text, timeout=90) |
| 82 | with (output / 'server.jsonl').open('a') as log: |
| 83 | log.write(json.dumps({'command': text, 'exit': result.returncode, 'stdout': result.stdout, 'stderr': result.stderr}) + '\n') |
| 84 | result.check_returncode() |
| 85 | return result.stdout |
| 86 | |
| 87 | try: |
| 88 | if not linux_vm.instance_path(server).exists(): |
| 89 | linux_vm.create_instance(server) |
| 90 | linux_vm.launch(server) |
| 91 | linux_vm.wait_instance(server, 600) |
| 92 | config = linux_vm.load_instance(server) |
| 93 | (output / 'linux.json').write_text(json.dumps(config, indent=2)) |
| 94 | if embedded_smb: |
| 95 | subprocess.run(linux_vm.ssh_argv(server, 'cat > /tmp/smb-proxy.py'), |
| 96 | input=(output / 'harness/smb-proxy.py').read_bytes(), check=True) |
| 97 | ssh("sudo sed -i '/^\\[global\\]/a smb ports = 1445' /etc/samba/smb.conf && sudo systemctl restart smbd") |
| 98 | subprocess.run(linux_vm.ssh_argv(server, 'cat > /tmp/smb-control.json'), |
| 99 | input=json.dumps({'record_writes': record_writes}).encode(), check=True) |
| 100 | ssh("sudo sh -c 'nohup python3 /tmp/smb-proxy.py /tmp/smb-control.json --port 445 --bind 0.0.0.0 --server 127.0.0.1 --server-port 1445 > /tmp/smb-trace.jsonl 2>&1 < /dev/null &'") |
| 101 | ssh("sleep 1; sudo ss -ltn | grep ':445 '") |
| 102 | os.environ['ONESTORE_SMB_LAB'] = f'127.0.0.1:{config["samba_port"]}' |
| 103 | os.environ['ONESTORE_SMB_SHARE'] = 'agent' |
| 104 | def unmount(force=False): |
| 105 | for options in ([['-f']] if force else [[], ['-f']]): |
| 106 | if not os.path.ismount(mount): return |
| 107 | result = subprocess.run(['/sbin/umount', *options, str(mount)], capture_output=True, text=True, timeout=60) |
| 108 | with (output / 'unmounts.jsonl').open('a') as log: |
| 109 | log.write(json.dumps({'force': bool(options), 'exit': result.returncode, 'stderr': result.stderr}) + '\n') |
| 110 | if os.path.ismount(mount): |
| 111 | raise RuntimeError('The owned SMB mount remains attached; preserve its server until it is unmounted.') |
| 112 | def reconnect_mount(): |
| 113 | unmount() |
| 114 | subprocess.run(['/sbin/mount_smbfs', '-N', f'//guest@127.0.0.1:{config["samba_port"]}/agent', mount], check=True, stdin=subprocess.DEVNULL) |
| 115 | ssh('mkdir /srv/agent/m6-collaboration') |
| 116 | source = ROOT / 'corpus/native-ink/cold-ui-ink/notebook' |
| 117 | if stress_clients or conflict_clients or abrupt: |
| 118 | source = output / 'input' |
| 119 | if fixture is not None: |
| 120 | source.mkdir() |
| 121 | for name in ('synthetic.one', 'Open Notebook.onetoc2'): |
| 122 | shutil.copyfile(Path(fixture) / name, source / name) |
| 123 | elif maintenance: |
| 124 | subprocess.run([ROOT / 'target/debug/examples/maintenance_fixture', source, 'Concurrent edits:'], check=True) |
| 125 | else: |
| 126 | subprocess.run([ROOT / 'target/debug/examples/create_notebook', source, 'Concurrent edits:', 'Concurrency test'], check=True) |
| 127 | with tarfile.open(output / 'input.tar', 'w', dereference=True) as archive: |
| 128 | for name in ('synthetic.one', 'Open Notebook.onetoc2'): |
| 129 | archive.add(source / name, arcname=name) |
| 130 | with (output / 'input.tar').open('rb') as stream: |
| 131 | subprocess.run(linux_vm.ssh_argv(server, 'tar xf - -C /srv/agent/m6-collaboration'), stdin=stream, check=True) |
| 132 | ssh('chmod u+w /srv/agent/m6-collaboration/*') |
| 133 | reconnect_mount() |
| 134 | shared = mount / 'm6-collaboration' |
| 135 | assert (shared / 'synthetic.one').read_bytes() == (source / 'synthetic.one').read_bytes() |
| 136 | for probe in ('sudo smbd --version', 'uname -a', 'testparm -s 2>&1'): |
| 137 | ssh(probe) |
| 138 | with ExitStack() as stack: |
| 139 | def start_controller(client, resume=False): |
| 140 | name, folder = client['name'], client['folder'] |
| 141 | command(name, 'if exist C:\\one-tests\\runs\\capture\\ready del C:\\one-tests\\runs\\capture\\ready', folder) |
| 142 | options = f' -ResumeCache -StartSequence {client["sequence"] + 1}' if resume else '' |
| 143 | result = windows.do_spawn('cmd /c powershell -NoProfile -NonInteractive -ExecutionPolicy Bypass -File C:\\one-tests\\collaborate.ps1 -Root C:\\one-tests\\runs\\capture -SharedPath \\\\192.168.77.1\\agent\\m6-collaboration -CloneHost ONE-' + name.upper() + options + f' > C:\\one-tests\\runs\\capture\\controller-{client["sequence"]}.log 2>&1', name) |
| 144 | (folder / f'spawn-{client["sequence"]}.json').write_text(json.dumps(result, indent=2)) |
| 145 | if result.get('error'): raise RuntimeError(result['error']) |
| 146 | deadline = time.monotonic() + 120 |
| 147 | while time.monotonic() < deadline: |
| 148 | result = windows.do_cmd('if exist C:\\one-tests\\runs\\capture\\ready (echo ready)', target=name) |
| 149 | if result.get('exit') == 0 and 'ready' in result.get('stdout', ''): return |
| 150 | time.sleep(1) |
| 151 | raise RuntimeError('The native collaboration client did not open the shared notebook') |
| 152 | |
| 153 | labels = [f'n{i}' for i in range(stress_clients or conflict_clients)] if stress_clients or conflict_clients else ['a', 'b'] |
| 154 | for label in labels: |
| 155 | (output / label).mkdir() |
| 156 | ownership = Lock() |
| 157 | def start_clone(label): |
| 158 | manager = clone(output / label) |
| 159 | name = manager.__enter__() |
| 160 | with ownership: |
| 161 | stack.push(manager) |
| 162 | return name |
| 163 | # Joining workers before stack exit retains ownership even if another boot fails. |
| 164 | with ThreadPoolExecutor(max_workers=len(labels)) as pool: |
| 165 | names = dict(zip(labels, pool.map(start_clone, labels))) |
| 166 | def prepare_client(label): |
| 167 | folder = output / label |
| 168 | name = names[label] |
| 169 | for local in scripts.iterdir(): |
| 170 | remote = 'cold-current.ps1' if local.name == 'cold.ps1' else local.name |
| 171 | result = windows.do_put(local, 'C:\\one-tests\\' + remote, name) |
| 172 | if result.get('error'): raise RuntimeError(result['error']) |
| 173 | command(name, 'mkdir C:\\one-tests\\runs\\capture', folder) |
| 174 | command(name, 'powershell -NoProfile -ExecutionPolicy Bypass -File C:\\one-tests\\network.ps1 -LabMac ' + vm.lab_mac(name), folder) |
| 175 | command(name, 'ipconfig', folder) |
| 176 | command(name, 'dir \\\\192.168.77.1\\agent\\m6-collaboration', folder) |
| 177 | client = {'name': name, 'folder': folder, 'sequence': 0} |
| 178 | start_controller(client) |
| 179 | command(name, 'ipconfig', folder) |
| 180 | print('Collaboration ready:', label, name, flush=True) |
| 181 | return client |
| 182 | with ThreadPoolExecutor(max_workers=len(labels)) as pool: |
| 183 | clients = list(pool.map(prepare_client, labels)) |
| 184 | |
| 185 | def action(client, action, wait=True, **parameters): |
| 186 | client['sequence'] += 1 |
| 187 | sequence = client['sequence'] |
| 188 | local = client['folder'] / f'command-{sequence:04}.json' |
| 189 | local.write_text(json.dumps({'action': action, **parameters}, ensure_ascii=False)) |
| 190 | inbox = f'C:\\one-tests\\runs\\capture\\inbox\\{sequence:04}' |
| 191 | result = windows.do_put(local, inbox + '.tmp', client['name']) |
| 192 | if result.get('error'): raise RuntimeError(result['error']) |
| 193 | command(client['name'], f'move {inbox}.tmp {inbox}.json', client['folder']) |
| 194 | if not wait: return sequence |
| 195 | return wait_action(client, sequence, action) |
| 196 | |
| 197 | def wait_action(client, sequence, action): |
| 198 | deadline = time.monotonic() + 600 |
| 199 | remote = f'C:\\one-tests\\runs\\capture\\outbox\\{sequence}' |
| 200 | while time.monotonic() < deadline: |
| 201 | result = windows.do_cmd(f'if exist {remote}\\error (type {remote}\\error & exit /b 1) else (if exist {remote}\\done (echo complete))', target=client['name']) |
| 202 | if result.get('exit') != 0: raise RuntimeError(str(result)) |
| 203 | if 'complete' in result.get('stdout', ''): break |
| 204 | time.sleep(.5) |
| 205 | else: raise TimeoutError(f'Native command {sequence} did not complete') |
| 206 | if action == 'snapshot': |
| 207 | hierarchy = client['folder'] / f'hierarchy-{sequence:04}.xml' |
| 208 | result = windows.do_get(remote + '\\hierarchy.xml', hierarchy, client['name']) |
| 209 | if result.get('error'): raise RuntimeError(result['error']) |
| 210 | xml = hierarchy.read_text(encoding='utf-8-sig').strip() |
| 211 | if xml == '<?xml version="1.0"?>': return [] |
| 212 | pages = [node for node in ET.fromstring(xml).iter() if node.tag.endswith('}Page')] |
| 213 | observed = [] |
| 214 | for i in range(len(pages)): |
| 215 | destination = client['folder'] / f'snapshot-{sequence:04}-{i}.xml' |
| 216 | result = windows.do_get(remote + f'\\page-{i}.xml', destination, client['name']) |
| 217 | if result.get('error'): raise RuntimeError(result['error']) |
| 218 | observed.extend(texts(ET.parse(destination).getroot())) |
| 219 | return observed |
| 220 | |
| 221 | def wait_text(client, expected): |
| 222 | # Background sync of a few hundred queued native edits takes minutes on a loaded host. |
| 223 | deadline = time.monotonic() + 600 |
| 224 | while time.monotonic() < deadline: |
| 225 | action(client, 'sync') |
| 226 | observed = action(client, 'snapshot') |
| 227 | if all(text in observed for text in expected): return observed |
| 228 | time.sleep(1) |
| 229 | raise AssertionError(('Native synchronization did not converge', expected, observed)) |
| 230 | |
| 231 | @contextmanager |
| 232 | def locked_file(name): |
| 233 | deadline = time.monotonic() + 120 |
| 234 | while True: |
| 235 | try: |
| 236 | descriptor = os.open(shared / name, os.O_RDONLY | os.O_EXLOCK | os.O_NONBLOCK) |
| 237 | break |
| 238 | except OSError as error: |
| 239 | with (output / 'read-retries.jsonl').open('a') as log: |
| 240 | log.write(json.dumps({'file': name, 'errno': error.errno, 'error': str(error)}) + '\n') |
| 241 | if error.errno not in (errno.ENOENT, errno.EACCES, errno.EAGAIN): raise |
| 242 | if time.monotonic() >= deadline: raise |
| 243 | time.sleep(.1) |
| 244 | with os.fdopen(descriptor, 'rb') as stream: |
| 245 | yield stream |
| 246 | |
| 247 | def snapshot_file(name): |
| 248 | with locked_file(name) as stream: |
| 249 | return stream.read() |
| 250 | |
| 251 | def checkpoint(label): |
| 252 | if embedded_smb: reconnect_mount() |
| 253 | deadline, previous, incomplete = time.monotonic() + 120, None, 0 |
| 254 | while time.monotonic() < deadline: |
| 255 | snapshot = {name: snapshot_file(name) for name in ('synthetic.one', 'Open Notebook.onetoc2')} |
| 256 | if snapshot == previous: |
| 257 | time.sleep(.1) |
| 258 | continue |
| 259 | previous = snapshot |
| 260 | target = output / label / 'notebook' |
| 261 | target.mkdir(parents=True) |
| 262 | for name, data in snapshot.items(): (target / name).write_bytes(data) |
| 263 | result = subprocess.run([EXPORTER, target / 'synthetic.one', target.parent / 'model'], capture_output=True, text=True) |
| 264 | if result.returncode == 0: |
| 265 | (target.parent / 'locks.txt').write_text(ssh('sudo smbstatus --locks')) |
| 266 | print('Captured collaboration state:', label, flush=True) |
| 267 | return |
| 268 | (target.parent / 'parse-error.txt').write_text(result.stderr) |
| 269 | if result.stderr.strip() != 'Error: Error { offset: 0, message: "Document context has no revision" }': |
| 270 | result.check_returncode() |
| 271 | incomplete += 1 |
| 272 | saved = output / f'{label}-incomplete-{incomplete}' |
| 273 | target.parent.rename(saved) |
| 274 | print('Preserved incomplete native document:', saved.name, flush=True) |
| 275 | time.sleep(.1) |
| 276 | raise TimeoutError(f'Native document references did not settle: {label}') |
| 277 | |
| 278 | if abrupt and not conflict_clients: |
| 279 | from concurrent_rust import running_clients |
| 280 | from crash_recovery import active_text, verify_text |
| 281 | import crash |
| 282 | |
| 283 | action(clients[0], 'prepare-stress', clients=len(clients)) |
| 284 | native_text = [f'Native {i}:' for i in range(len(clients))] |
| 285 | for client in clients: wait_text(client, ['Concurrent edits:', *native_text]) |
| 286 | checkpoint('crash-initial') |
| 287 | folder = output / 'rust-interrupted' |
| 288 | folder.mkdir() |
| 289 | try: |
| 290 | with running_clients(folder, shared / 'synthetic.one', rust_writers, rust_readers, 500, seed) as processes: |
| 291 | (folder / 'start').touch() |
| 292 | for i, client in enumerate(clients): |
| 293 | changed = native_text[i] + ' Before server stop.' |
| 294 | action(client, 'edit', expected=native_text[i], text=changed) |
| 295 | native_text[i] = changed |
| 296 | for client in clients: action(client, 'sync') |
| 297 | for client in clients: wait_text(client, native_text) |
| 298 | deadline = time.monotonic() + 60 |
| 299 | while True: |
| 300 | commits = sum(line.count('"event":"commit"') for path in folder.glob('w*.jsonl') for line in path.read_text().splitlines()) |
| 301 | if commits >= 10: break |
| 302 | if any(p.poll() is not None for p in processes.values()) or time.monotonic() >= deadline: |
| 303 | raise AssertionError('Writers did not overlap the server interruption') |
| 304 | time.sleep(.02) |
| 305 | assert all(p.poll() is None for p in processes.values()), 'A client stopped before the planned interruption' |
| 306 | (output / 'server-crash.json').write_text(json.dumps(crash.stop('linux', server), indent=2)) |
| 307 | for client in clients: vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': False}) |
| 308 | for process in processes.values(): process.terminate() |
| 309 | unmount(force=True) |
| 310 | raise InterruptedError('Recorded server interruption') |
| 311 | except InterruptedError: |
| 312 | pass |
| 313 | logs = {actor: [json.loads(line) for line in (folder / f'{actor}.jsonl').read_text().splitlines()] for actor in processes} |
| 314 | assert not list(folder.glob('invalid-*.one')), 'A client captured malformed storage' |
| 315 | linux_vm.launch(server) |
| 316 | linux_vm.wait_instance(server, 600) |
| 317 | recovered = output / 'server-before-native-reconnect' |
| 318 | recovered.mkdir() |
| 319 | with (recovered / 'synthetic.one').open('wb') as stream: |
| 320 | subprocess.run(linux_vm.ssh_argv(server, 'cat /srv/agent/m6-collaboration/synthetic.one'), stdout=stream, check=True, timeout=60) |
| 321 | subprocess.run([EXPORTER, recovered / 'synthetic.one', recovered / 'model'], check=True) |
| 322 | model = json.loads((recovered / 'model/document.json').read_text()) |
| 323 | text = active_text(model) |
| 324 | retained = verify_text('Concurrent edits:', text, logs) |
| 325 | (output / 'server-retention.json').write_text(json.dumps(retained, indent=2)) |
| 326 | reconnect_mount() |
| 327 | for client in clients: vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': True}) |
| 328 | for client in clients: action(client, 'sync') |
| 329 | for client in clients: wait_text(client, [text, *native_text]) |
| 330 | checkpoint('server-recovered') |
| 331 | |
| 332 | client = clients[0] |
| 333 | pending_text = native_text[0] + ' Offline cache edit.' |
| 334 | with locked_file('synthetic.one'): |
| 335 | vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': False}) |
| 336 | action(client, 'edit', expected=native_text[0], text=pending_text) |
| 337 | assert pending_text in action(client, 'snapshot') |
| 338 | (output / 'client-crash.json').write_text(json.dumps(crash.stop('windows', client['name']), indent=2)) |
| 339 | checkpoint('client-stopped') |
| 340 | stopped = json.loads((output / 'client-stopped/model/document.json').read_text()) |
| 341 | assert pending_text not in [value for values in reachable_page_text(stopped).values() for value in values], 'The pending cache edit already reached the server' |
| 342 | vm.start_instance(client['name']) |
| 343 | vm.wait_instance(client['name'], 300) |
| 344 | command(client['name'], 'powershell -NoProfile -ExecutionPolicy Bypass -File C:\\one-tests\\network.ps1 -LabMac ' + vm.lab_mac(client['name']), client['folder']) |
| 345 | start_controller(client, resume=True) |
| 346 | deadline = time.monotonic() + 120 |
| 347 | while True: |
| 348 | action(client, 'sync') |
| 349 | observed = action(client, 'snapshot') |
| 350 | if pending_text in observed or native_text[0] in observed: break |
| 351 | if time.monotonic() >= deadline: |
| 352 | raise AssertionError('Recovered cache has neither acknowledged version') |
| 353 | time.sleep(1) |
| 354 | survived = pending_text in observed |
| 355 | native_text[0] = pending_text if survived else native_text[0] |
| 356 | (output / 'native-cache-result.json').write_text(json.dumps({'cache_acknowledged': pending_text, 'retained_after_abrupt_stop': survived, 'server_durability_acknowledged': False}, indent=2)) |
| 357 | for client in clients: wait_text(client, [text, *native_text]) |
| 358 | checkpoint('client-recovered') |
| 359 | recovered_model = json.loads((output / 'client-recovered/model/document.json').read_text()) |
| 360 | recovered_pages = reachable_page_text(recovered_model) |
| 361 | main_space, *conflict_spaces = recovered_pages |
| 362 | known = {'Concurrent edits:', *[f'Native {i}:' for i in range(len(clients))], *native_text, pending_text} |
| 363 | conflicts = {sid: recovered_pages[sid] for sid in conflict_spaces} |
| 364 | assert all(value in known or (value.endswith(']') and text.startswith(value)) |
| 365 | for values in conflicts.values() for value in values), 'A recovered conflict contains unrecorded content' |
| 366 | (output / 'cache-conflicts.json').write_text(json.dumps(conflicts, indent=2)) |
| 367 | |
| 368 | continued = output / 'rust-continued' |
| 369 | continued.mkdir() |
| 370 | with running_clients(continued, shared / 'synthetic.one', rust_writers, rust_readers, stress_operations, seed + 1) as processes: |
| 371 | (continued / 'start').touch() |
| 372 | continuation = {actor: [json.loads(line) for line in (continued / f'{actor}.jsonl').read_text().splitlines()] for actor in processes} |
| 373 | checkpoint('continued') |
| 374 | model = json.loads((output / 'continued/model/document.json').read_text()) |
| 375 | final_text = active_text(model) |
| 376 | result = verify_text(text, final_text, continuation) |
| 377 | for client in clients: wait_text(client, [final_text, *native_text]) |
| 378 | (output / 'result.json').write_text(json.dumps({'server': retained, 'continuation': result, 'native_cache_retained': survived, 'expected_text': sorted([final_text, *native_text])}, indent=2)) |
| 379 | elif conflict_clients: |
| 380 | from concurrent_rust import running_clients |
| 381 | from native_stress import edit_history |
| 382 | |
| 383 | original = 'Concurrent edits:' |
| 384 | native = [f'Native conflict {i}: café 🦀' for i in range(conflict_clients)] |
| 385 | for client in clients: wait_text(client, [original]) |
| 386 | checkpoint('initial') |
| 387 | try: |
| 388 | with locked_file('synthetic.one'): |
| 389 | for client in clients: vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': False}) |
| 390 | for client, text in zip(clients, native): |
| 391 | action(client, 'edit', expected=original, text=text) |
| 392 | assert text in action(client, 'snapshot'), 'Native cache did not retain its intended edit' |
| 393 | folder = output / 'rust' |
| 394 | folder.mkdir() |
| 395 | with running_clients(folder, shared / 'synthetic.one', rust_writers, rust_readers, stress_operations, seed, edit=edit) as processes: |
| 396 | (folder / 'start').touch() |
| 397 | logs = {actor: [json.loads(line) for line in (folder / f'{actor}.jsonl').read_text().splitlines()] for actor in processes} |
| 398 | commits, rust_text = edit_history(logs, stress_operations, edit) |
| 399 | checkpoint('rust-before-reconnect') |
| 400 | finally: |
| 401 | for client in clients: vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': True}) |
| 402 | expected = {rust_text, *native} |
| 403 | deadline = time.monotonic() + 120 |
| 404 | while True: |
| 405 | for client in clients: action(client, 'sync') |
| 406 | checkpoint_label = f'convergence-{clients[0]["sequence"]}' |
| 407 | checkpoint(checkpoint_label) |
| 408 | model = json.loads((output / checkpoint_label / 'model/document.json').read_text()) |
| 409 | retained = reachable_page_text(model) |
| 410 | if sorted(expected) == sorted(text for texts in retained.values() for text in texts): break |
| 411 | if time.monotonic() >= deadline: |
| 412 | (output / 'lost-edit.json').write_text(json.dumps({'expected': sorted(expected), |
| 413 | 'reachable_page_text': retained, 'checkpoint': checkpoint_label}, indent=2)) |
| 414 | raise AssertionError(('Competing Rust/native edits differ from recorded results', sorted(expected), retained)) |
| 415 | time.sleep(1) |
| 416 | checkpoint('conflict-final') |
| 417 | (output / 'live-retention.json').write_text(json.dumps(retained, indent=2)) |
| 418 | if abrupt: |
| 419 | import crash |
| 420 | with locked_file('synthetic.one') as stream: |
| 421 | os.fsync(stream.fileno()) |
| 422 | (output / 'conflict-server-crash.json').write_text(json.dumps(crash.stop('linux', server), indent=2)) |
| 423 | for client in clients: vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': False}) |
| 424 | linux_vm.launch(server) |
| 425 | linux_vm.wait_instance(server, 600) |
| 426 | reconnect_mount() |
| 427 | checkpoint('conflict-server-recovered') |
| 428 | model = json.loads((output / 'conflict-server-recovered/model/document.json').read_text()) |
| 429 | observed = [text for texts in reachable_page_text(model).values() for text in texts] |
| 430 | assert sorted(observed) == sorted(expected), 'Server interruption lost a flushed competing result' |
| 431 | for client in clients: vm.qmp(client['name'], 'set_link', {'name': 'lab', 'up': True}) |
| 432 | for client in clients: action(client, 'sync') |
| 433 | elif stress_clients: |
| 434 | from native_stress import exercise |
| 435 | exercise(output, shared, clients, action, wait_action, wait_text, checkpoint, stress_operations, sync_every, rust_writers, rust_readers, edit, seed, embedded_smb) |
| 436 | else: |
| 437 | a, b = clients |
| 438 | original = 'Fictitious: café, 東京, مرحبا' |
| 439 | rust_text = 'Rust disjoint edit.'.ljust(len(original.encode('utf-16-le')) // 2, '.') |
| 440 | other = 'Fictitious positioned outline.' |
| 441 | for client in clients: wait_text(client, [original, other]) |
| 442 | checkpoint('initial') |
| 443 | subprocess.run([ROOT / 'target/debug/examples/edit_property', shared / 'synthetic.one', '--in-place', '6000e', '1c001c22', (original + '\0').encode('utf-16-le').hex(), (rust_text + '\0').encode('utf-16-le').hex()], check=True) |
| 444 | action(a, 'edit', expected=other, text='Native disjoint edit.') |
| 445 | for client in clients: wait_text(client, [rust_text, 'Native disjoint edit.']) |
| 446 | checkpoint('disjoint') |
| 447 | ssh('sudo systemctl stop smbd && ! systemctl is-active --quiet smbd') |
| 448 | try: |
| 449 | action(a, 'edit', expected=rust_text, text='Client A concurrent edit.') |
| 450 | action(b, 'edit', expected=rust_text, text='Client B concurrent edit.') |
| 451 | finally: |
| 452 | ssh('sudo systemctl start smbd && systemctl is-active smbd') |
| 453 | reconnect_mount() |
| 454 | for client in clients: action(client, 'sync') |
| 455 | deadline = time.monotonic() + 120 |
| 456 | while time.monotonic() < deadline: |
| 457 | for client in clients: action(client, 'sync') |
| 458 | temporary = output / 'conflict-model' |
| 459 | if temporary.exists(): shutil.rmtree(temporary) |
| 460 | snapshot = output / 'conflict.one' |
| 461 | snapshot.write_bytes(snapshot_file('synthetic.one')) |
| 462 | subprocess.run([EXPORTER, snapshot, temporary], check=True) |
| 463 | model = json.loads((temporary / 'document.json').read_text()) |
| 464 | values = [node['kind'].get('text') for space in model['spaces'].values() for revision in space['revisions'].values() for node in revision['nodes'].values()] |
| 465 | if 'Client A concurrent edit.' in values and 'Client B concurrent edit.' in values: break |
| 466 | time.sleep(1) |
| 467 | else: raise AssertionError('Both concurrent edits were not retained') |
| 468 | checkpoint('conflict') |
| 469 | ssh('sudo systemctl restart smbd && systemctl is-active smbd') |
| 470 | reconnect_mount() |
| 471 | for client in clients: wait_text(client, ['Native disjoint edit.']) |
| 472 | checkpoint('server-restart') |
| 473 | vm.qmp(a['name'], 'set_link', {'name': 'lab', 'up': False}) |
| 474 | try: |
| 475 | action(a, 'edit', expected='Native disjoint edit.', text='Client A offline edit.') |
| 476 | assert 'Client A offline edit.' in action(a, 'snapshot') |
| 477 | action(b, 'edit', expected='Left cell', text='Client B online cell.') |
| 478 | wait_text(b, ['Client B online cell.', 'Native disjoint edit.']) |
| 479 | checkpoint('transport-disconnected') |
| 480 | finally: |
| 481 | vm.qmp(a['name'], 'set_link', {'name': 'lab', 'up': True}) |
| 482 | for client in clients: wait_text(client, ['Client A offline edit.', 'Client B online cell.']) |
| 483 | checkpoint('transport-reconnected') |
| 484 | with os.fdopen(os.open(shared / 'synthetic.one', os.O_RDONLY | os.O_EXLOCK | os.O_NONBLOCK), 'rb') as locked: |
| 485 | before = locked.read() |
| 486 | action(a, 'edit', expected='Client A offline edit.', text='Native lock recovery.') |
| 487 | action(a, 'sync') |
| 488 | assert 'Native lock recovery.' in action(a, 'snapshot') |
| 489 | locked.seek(0) |
| 490 | assert locked.read() == before, 'The shared file changed while exclusively locked' |
| 491 | (output / 'contention.json').write_text(json.dumps({'unchanged_sha256': hashlib.sha256(before).hexdigest(), 'locks': ssh('sudo smbstatus --locks')}, indent=2)) |
| 492 | for client in clients: wait_text(client, ['Native lock recovery.', 'Client B online cell.']) |
| 493 | checkpoint('lock-released') |
| 494 | action(a, 'kill-process') |
| 495 | for client in clients: wait_text(client, ['Native lock recovery.', 'Client B online cell.']) |
| 496 | checkpoint('process-restarted') |
| 497 | vm.qmp(a['name'], 'quit') |
| 498 | deadline = time.monotonic() + 10 |
| 499 | while vm.running(a['name']) and time.monotonic() < deadline: time.sleep(.1) |
| 500 | assert not vm.running(a['name']), 'The owned clone did not terminate' |
| 501 | vm.start_instance(a['name']) |
| 502 | vm.wait_instance(a['name'], 300) |
| 503 | command(a['name'], 'powershell -NoProfile -ExecutionPolicy Bypass -File C:\\one-tests\\network.ps1 -LabMac ' + vm.lab_mac(a['name']), a['folder']) |
| 504 | start_controller(a, resume=True) |
| 505 | for client in clients: wait_text(client, ['Native lock recovery.', 'Client B online cell.']) |
| 506 | checkpoint('vm-restarted') |
| 507 | final = json.loads((output / 'vm-restarted/model/document.json').read_text()) |
| 508 | (output / 'retention.json').write_text(json.dumps(verify_final_state(final), indent=2)) |
| 509 | for client in clients: |
| 510 | action(client, 'close') |
| 511 | collect_artifacts(client['name'], client['folder'], '*') |
| 512 | if abrupt and not conflict_clients: |
| 513 | checkpoint('crash-closed') |
| 514 | model = json.loads((output / 'crash-closed/model/document.json').read_text()) |
| 515 | actual = reachable_page_text(model) |
| 516 | assert actual.pop(main_space) == sorted([final_text, *native_text]), 'Application closure changed recovered edits' |
| 517 | assert actual == conflicts, 'Application closure changed recovered conflict pages' |
| 518 | elif stress_clients: |
| 519 | checkpoint('stress-closed') |
| 520 | elif conflict_clients: |
| 521 | checkpoint('conflict-closed') |
| 522 | model = json.loads((output / 'conflict-closed/model/document.json').read_text()) |
| 523 | retained = reachable_page_text(model) |
| 524 | assert sorted(expected) == sorted(text for texts in retained.values() for text in texts), 'Closing native clients changed the competing results' |
| 525 | (output / 'result.json').write_text(json.dumps({'native_writers': conflict_clients, |
| 526 | 'rust_writers': rust_writers, 'rust_readers': rust_readers, 'rust_commits': len(commits), |
| 527 | 'competing_results': sorted(expected), 'reachable_page_text': retained}, indent=2)) |
| 528 | if not (stress_clients or conflict_clients or abrupt): (output / 'result.json').write_text(json.dumps({'disjoint': True, 'conflict_retained': True, 'server_restart': True, 'transport_reconnect': True, 'lock_contention': True, 'process_restart': True, 'vm_restart': True}, indent=2)) |
| 529 | except BaseException: |
| 530 | if linux_vm.instance_path(server).exists() and linux_vm.running(server): |
| 531 | failure = output / 'failure-after-client-shutdown' |
| 532 | failure.mkdir(exist_ok=True) |
| 533 | try: |
| 534 | for name in ('synthetic.one', 'Open Notebook.onetoc2'): |
| 535 | with (failure / name).open('wb') as stream: |
| 536 | subprocess.run(linux_vm.ssh_argv(server, f'cat /srv/agent/m6-collaboration/"{name}"'), |
| 537 | stdout=stream, stderr=subprocess.PIPE, check=True, timeout=30) |
| 538 | except Exception as error: |
| 539 | (failure / 'capture-error.txt').write_text(str(error)) |
| 540 | raise |
| 541 | finally: |
| 542 | if embedded_smb and linux_vm.instance_path(server).exists() and linux_vm.running(server): |
| 543 | try: |
| 544 | with (output / 'smb-trace.jsonl').open('wb') as trace: |
| 545 | subprocess.run(linux_vm.ssh_argv(server, 'cat /tmp/smb-trace.jsonl'), stdout=trace, check=True, timeout=30) |
| 546 | except Exception as error: |
| 547 | (output / 'trace-error.txt').write_text(str(error)) |
| 548 | if os.path.ismount(mount): unmount() |
| 549 | if linux_vm.instance_path(server).exists(): |
| 550 | try: |
| 551 | if linux_vm.running(server): linux_vm.shutdown(server, 60) |
| 552 | finally: |
| 553 | if linux_vm.running(server): linux_vm.qmp(server, 'quit') |
| 554 | deadline = time.monotonic() + 10 |
| 555 | while linux_vm.running(server) and time.monotonic() < deadline: time.sleep(.1) |
| 556 | linux_vm.delete_instance(server) |
| 557 | (output / 'teardown.json').write_text(json.dumps({'linux_absent': not linux_vm.instance_path(server).exists()}, indent=2)) |
| 558 | if embedded_smb: |
| 559 | from verify_smb_overlap import verify |
| 560 | with (output / 'smb-trace.jsonl').open() as trace: |
| 561 | overlap = verify(json.loads(line) for line in trace) |
| 562 | (output / 'overlap.json').write_text(json.dumps(overlap, indent=2)) |
| 563 | if offline_lost_reply: |
| 564 | from offline_outage import verify_lost_reply |
| 565 | (output / 'offline-lost-reply-verification.json').write_text(json.dumps(verify_lost_reply(output), indent=2)) |
| 566 | if offline_outage: |
| 567 | from offline_outage import verify_outage |
| 568 | (output / 'offline-outage-verification.json').write_text(json.dumps(verify_outage(output), indent=2)) |
| 569 | if disconnect: |
| 570 | from native_disconnect import verify_disconnect |
| 571 | (output / 'disconnect-verification.json').write_text(json.dumps(verify_disconnect(output), indent=2)) |
| 572 | |
| 573 | |
| 574 | if __name__ == '__main__': |
| 575 | parser = argparse.ArgumentParser(description=__doc__) |
| 576 | parser.add_argument('output', type=Path) |
| 577 | parser.add_argument('--linux', required=True, help='Disposable Linux VM owned by this run; it is deleted on exit.') |
| 578 | parser.add_argument('--stress-clients', type=int, default=0) |
| 579 | parser.add_argument('--conflict-clients', type=int, default=0, help='Native clients editing the same paragraph offline while Rust writes the server copy.') |
| 580 | parser.add_argument('--stress-operations', type=int, default=30) |
| 581 | parser.add_argument('--sync-every', type=int, default=1, help='Native edits between explicit sync requests; zero uses background sync.') |
| 582 | parser.add_argument('--rust-writers', type=int, default=4) |
| 583 | parser.add_argument('--rust-readers', type=int, default=3) |
| 584 | parser.add_argument('--edit', action='store_true', help='Use random text replacements in every writer.') |
| 585 | parser.add_argument('--seed', type=int, default=710) |
| 586 | parser.add_argument('--abrupt', action='store_true', help='Abrupt server and client stops with preserved-disk recovery and intent accounting.') |
| 587 | parser.add_argument('--embedded-smb', action='store_true', help='Run Rust stress clients through the embedded SMB adapter.') |
| 588 | parser.add_argument('--maintenance', action='store_true', help='Pause the workload for the owned maintenance controller.') |
| 589 | parser.add_argument('--offline', action='store_true', help='Use durable local queues and traced offline workers for Rust writers.') |
| 590 | parser.add_argument('--document-operations', action='store_true', help='Queue paragraph/outline creation and formatting alongside offline text edits.') |
| 591 | parser.add_argument('--record-writes', action='store_true', help='Retain owned lab write payloads for revision replay.') |
| 592 | parser.add_argument('--offline-lost-reply', action='store_true', help='Drop a publication reply and require confirmation of its original revision.') |
| 593 | parser.add_argument('--offline-client-reply', action='store_true', help='Disconnect only the formatting writer and require peer publication before it reconciles.') |
| 594 | parser.add_argument('--offline-outage', action='store_true', help='Queue local edits during an owned SMB outage, then require recovery.') |
| 595 | parser.add_argument('--disconnect', action='store_true', help='Interrupt and reconnect the embedded append workload twice.') |
| 596 | parser.add_argument('--client-profile', choices=('debug', 'release'), default='debug', help='Cargo build profile for Rust stress clients') |
| 597 | parser.add_argument('--client-timeout', type=float, default=600, help='Maximum seconds for the Rust workload, including its start barrier.') |
| 598 | parser.add_argument('--fixture', type=Path, help='Copy this fixture directory into the disposable stress notebook.') |
| 599 | args = parser.parse_args() |
| 600 | if not 0 < args.client_timeout < 2**64 / 1000: parser.error('Choose a finite positive client timeout.') |
| 601 | if args.maintenance and not (args.embedded_smb and args.stress_clients): parser.error('--maintenance requires embedded SMB stress mode.') |
| 602 | if args.document_operations and not args.offline: parser.error('--document-operations requires --offline.') |
| 603 | if args.record_writes and not args.embedded_smb: parser.error('--record-writes requires --embedded-smb.') |
| 604 | if args.offline_client_reply and not (args.offline_lost_reply and args.document_operations): parser.error('--offline-client-reply requires --offline-lost-reply and --document-operations.') |
| 605 | if args.offline_lost_reply and (not args.offline or args.offline_outage or args.sync_every or args.stress_operations < 8): parser.error('--offline-lost-reply requires --offline, --sync-every 0, at least eight operations and no --offline-outage.') |
| 606 | if args.offline_outage and (not args.offline or args.sync_every or args.stress_operations < 8): parser.error('--offline-outage requires --offline, --sync-every 0 and at least eight operations.') |
| 607 | if args.offline and (not args.embedded_smb or not args.stress_clients or args.edit or args.disconnect or args.maintenance): parser.error('--offline requires embedded append stress without disconnect or maintenance.') |
| 608 | if args.embedded_smb and not args.stress_clients: parser.error('--embedded-smb requires --stress-clients.') |
| 609 | if args.disconnect and (not args.embedded_smb or args.edit or args.maintenance or args.sync_every): parser.error('--disconnect requires embedded append stress with --sync-every 0 and no maintenance.') |
| 610 | if args.rust_writers < 2 or args.rust_readers < 1: parser.error('Use at least two Rust writers and one reader.') |
| 611 | if args.sync_every < 0 or args.stress_operations <= 0: parser.error('Use a nonnegative sync interval and positive operation count.') |
| 612 | if args.stress_clients and args.stress_clients < 3: parser.error('Stress mode requires at least three native clients.') |
| 613 | if args.conflict_clients < 0 or (args.conflict_clients and args.stress_clients): parser.error('Choose a positive conflict client count without stress mode.') |
| 614 | if args.abrupt and (args.stress_clients or args.edit): parser.error('Abrupt recovery uses append intents or the offline-conflict workload.') |
| 615 | def interrupted(_signal, _frame): raise KeyboardInterrupt |
| 616 | signal.signal(signal.SIGTERM, interrupted) |
| 617 | replay(args.output, args.linux, args.stress_clients, args.stress_operations, args.sync_every, args.rust_writers, args.rust_readers, args.edit, args.seed, args.conflict_clients, args.abrupt, args.embedded_smb, args.maintenance, args.fixture, args.disconnect, args.client_timeout, args.offline, args.offline_outage, args.offline_lost_reply, args.client_profile, args.document_operations, args.record_writes, args.offline_client_reply) |