1#!/usr/bin/env python3
2"""Replay shared notebook edits with disposable OneNote clients and a Linux server."""
3import argparse
4import errno
5from contextlib import ExitStack, contextmanager
6from concurrent.futures import ThreadPoolExecutor
7from threading import Lock
8import hashlib
9import json
10import os
11from pathlib import Path
12import shutil
13import signal
14import subprocess
15import tarfile
16import time
17import xml.etree.ElementTree as ET
18
19from native_runner import ROOT, clone, command, windows, vm, collect_artifacts
20from native_xml import texts
21from document_model import EXPORTER, ordered_pages, view, walk
22import linux_vm
23
24
25def 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
41def 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
50def 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
574if __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)