1"""Overlapping native and Rust editing histories on one shared section."""
2import json
3import os
4import time
5
6from native_runner import windows
7from concurrent_rust import running_clients
8from document_model import ordered_pages, walk
9
10
11def edit_history(logs, operations, edit=False, *, partial=False, offline=False):
12 if offline:
13 assert not edit, 'Offline acceptance currently uses append intents'
14 from offline_history import publication_links
15 links = publication_links(logs, operations, partial)
16 else:
17 links = {}
18 for actor, events in logs.items():
19 assert events and events[0]['event'] == 'ready', 'Missing client start'
20 assert partial or events[-1]['event'] == 'done', 'Incomplete client log'
21 if not actor.startswith('w'): continue
22 intents = {event['attempt']: event for event in events if event['event'] == 'intent'}
23 commits = [event for event in events if event['event'] == 'commit']
24 assert len(commits) <= operations, 'Unexpected Rust acknowledgement'
25 expected = len(commits) if partial else operations
26 assert sorted(event['operation'] for event in commits) == list(range(expected)), 'Missing or duplicate Rust acknowledgements'
27 for event in commits:
28 intent = intents[event['attempt']]
29 token = f' [{actor}:{event["operation"]}]'
30 offset = len(intent['before'].encode('utf-16-le')) // 2
31 assert event['token'] == intent['token'] == token, 'Acknowledgement differs from intended edit'
32 assert intent['replacement'] == (' café 🦀' if edit else '') + token, 'Unexpected replacement text'
33 assert intent['operation'] == event['operation'] and intent['source_transaction'] == event['source_transaction'], 'Acknowledgement used another intent'
34 start, end = intent['range']
35 assert len('Concurrent edits:') <= start <= end <= offset, 'Invalid edit range'
36 if not edit: assert start == end == offset, 'Expected an append intent'
37 units = intent['before'].encode('utf-16-le')
38 after = units[:start * 2].decode('utf-16-le') + intent['replacement'] + units[end * 2:].decode('utf-16-le')
39 assert intent['before'] not in links, 'Acknowledged Rust edits branched from the same content'
40 links[intent['before']] = event, after
41 text = 'Concurrent edits:'
42 versions = {text: 0}
43 ordered = []
44 while text in links:
45 event, text = links.pop(text)
46 ordered.append(event)
47 versions[text] = len(ordered)
48 assert not links, 'Acknowledged edits do not form one complete content history'
49 for events in logs.values():
50 previous = 0
51 for event in events:
52 if event['event'] != 'read': continue
53 assert event['text'] in versions, 'A reader observed partial, lost or invented content'
54 version = versions[event['text']]
55 assert version >= previous, 'A reader went backwards in content history'
56 previous = version
57 for i, commit in enumerate(ordered):
58 if commit['finished_us'] < event['started_us']:
59 assert version >= i + 1, 'A read missed an acknowledged edit'
60 if commit['started_us'] > event['finished_us']:
61 assert version <= i, 'A read observed a future edit'
62 return ordered, text
63
64
65def native_history(events, actor, operations, edit):
66 assert len(events) == operations, 'Missing native acknowledgements'
67 prefix = f'Native {actor}:'
68 expected = prefix
69 for j, event in enumerate(events):
70 assert event['operation'] == j and event['token'] == f' [n{actor}:{j}]' and event['before'] == expected, 'Native history lost or invented an edit'
71 if edit:
72 start, end = event['range']
73 units = expected.encode('utf-16-le')
74 assert len(prefix) <= start <= end <= len(units) // 2, 'Invalid native edit range'
75 assert event['replacement'] == ' café 🦀' + event['token'], 'Unexpected native replacement'
76 expected = units[:start * 2].decode('utf-16-le') + event['replacement'] + units[end * 2:].decode('utf-16-le')
77 else:
78 expected += event['token']
79 return expected
80
81
82def verify_capture(output, capture):
83 import xml.etree.ElementTree as ET
84 from native_format import native_characters
85 from native_xml import ns
86
87 config = json.loads((output / 'run.json').read_text())
88 operations, edit = config['stress_operations'], config['edit']
89 actors = [f'w{i}' for i in range(config['rust_writers'])] + [f'r{i}' for i in range(config['rust_readers'])]
90 logs = {actor: [json.loads(line) for line in (output / 'rust' / f'{actor}.jsonl').read_text().splitlines()] for actor in actors}
91 commits, rust_text = edit_history(logs, operations, edit, offline=config.get('offline', False))
92 native = [[json.loads(line) for line in (output / f'n{i}' / 'stress-events.jsonl').read_text(encoding='utf-8-sig').splitlines()]
93 for i in range(config['stress_clients'])]
94 expected = [rust_text, *(native_history(events, i, operations, edit) for i, events in enumerate(native))]
95 documents = {}
96 if config.get('document_operations'):
97 from offline_document_history import document_history
98 documents = document_history(logs, operations)
99 expected.extend(''.join(char for char, *_ in list(document['states'].values())[-1]['characters']) for document in documents.values())
100 page_file, = capture.glob('page-*.xml')
101 page = ET.parse(page_file).getroot()
102 paragraphs = native_characters(page, page.findall('one:Outline', ns))
103 actual = [''.join(char for char, _ in paragraph) for paragraph in paragraphs]
104 assert sorted(actual) == sorted(expected), 'Fresh native content differs from the complete recorded editing history'
105 checks = 0
106 if edit:
107 for text, events in zip(expected[1:], native, strict=True):
108 for _, style in paragraphs[actual.index(text)]:
109 for field in ('bold', 'italic'):
110 assert bool(style.get(field)) == events[-1][field], 'Fresh native formatting differs from its last recorded editing intent'
111 checks += 1
112 if documents:
113 from offline_document_history import verify_native
114 checks += verify_native(paragraphs, (list(document['states'].values())[-1]['characters'] for document in documents.values()))
115 return {'rust_intents': len(commits), 'document_intents': sum(len(document['states']) for document in documents.values()), 'native_intents': sum(map(len, native)),
116 'exact_paragraphs': len(expected), 'native_intended_format_checks': checks}
117
118
119def exercise(output, shared, clients, action, wait_action, wait_text, checkpoint, operations, sync_every, rust_writers, rust_readers, edit=False, seed=710, embedded_smb=False):
120 wait_text(clients[0], ['Concurrent edits:'])
121 action(clients[0], 'prepare-stress', clients=len(clients))
122 prefixes = [f'Native {i}:' for i in range(len(clients))]
123 for client in clients:
124 wait_text(client, ['Concurrent edits:', *prefixes])
125 checkpoint('stress-initial')
126 config = json.loads((output / 'run.json').read_text())
127 maintenance = config.get('maintenance', False)
128 disconnect = config.get('disconnect', False)
129 offline_outage = config.get('offline_outage', False)
130 offline_lost_reply = config.get('offline_lost_reply', False)
131 if maintenance:
132 (output / 'maintenance-ready').touch()
133 deadline = time.monotonic() + 300
134 while not (output / 'maintenance-start').exists():
135 if time.monotonic() > deadline: raise TimeoutError('Maintenance controller did not start the workload')
136 time.sleep(.1)
137 clocks = []
138 for client in clients:
139 before = time.time_ns() // 1000
140 result = windows.do_health(client['name'])
141 after = time.time_ns() // 1000
142 native = result['utc_us']
143 clocks.append({'native_minus_host_us': [native - after - 15625, native - before + 15625]})
144 (output / 'clocks.json').write_text(json.dumps(clocks, indent=2))
145 sequences = [action(client, 'stress', wait=False, actor=i, prefix=prefixes[i], operations=operations, seed=seed + 10000 + i, sync_every=sync_every, edit=edit, maintenance=maintenance)
146 for i, client in enumerate(clients)]
147 for client, sequence in zip(clients, sequences):
148 deadline = time.monotonic() + 60
149 while time.monotonic() < deadline:
150 result = windows.do_cmd(f'if exist C:\\one-tests\\runs\\capture\\outbox\\{sequence}\\ready echo ready', target=client['name'])
151 if 'ready' in result.get('stdout', ''): break
152 time.sleep(.1)
153 else: raise TimeoutError('Native stress client did not reach the barrier')
154 folder = output / 'rust'
155 folder.mkdir()
156 start = folder / 'start'
157 from concurrent_rust import ROOT
158 binaries = ROOT / 'target' / config.get('client_profile', 'debug') / 'examples'
159 source = 'm6-collaboration\\synthetic.one' if embedded_smb else shared / 'synthetic.one'
160 executable = binaries / ('smb_concurrent_client' if embedded_smb else 'concurrent_client')
161 if disconnect: executable = binaries / 'smb_reconnect_client'
162 if config.get('offline'): executable = binaries / 'smb_offline_client'
163 environment = {**os.environ, 'ONESTORE_MAINTENANCE_DIR': str(folder)} if maintenance else None
164 reader_executable = None
165 if offline_outage:
166 environment = {**os.environ, 'ONESTORE_OFFLINE_OUTAGE_DIR': str(folder)}
167 reader_executable = binaries / 'smb_reconnect_client'
168 if offline_lost_reply:
169 captures = folder / 'confirmations'
170 captures.mkdir()
171 environment = {**os.environ, 'ONESTORE_OFFLINE_CONFIRM_DIR': str(captures)}
172 reader_executable = binaries / 'smb_reconnect_client'
173 if config.get('document_operations'):
174 environment = {**(environment or os.environ), 'ONESTORE_OFFLINE_DOCUMENTS': '1'}
175 if offline_lost_reply:
176 environment['ONESTORE_OFFLINE_FORMAT_REPLY_DIR'] = str(folder)
177 try:
178 with running_clients(folder, source, rust_writers, rust_readers, operations, seed, timeout=config.get('client_timeout', 600), edit=edit, executable=executable, environment=environment, reader_executable=reader_executable) as processes:
179 (shared / 'stress-start').write_text('start')
180 for client, sequence in zip(clients, sequences):
181 deadline = time.monotonic() + 60
182 while time.monotonic() < deadline:
183 result = windows.do_cmd(f'if exist C:\\one-tests\\runs\\capture\\outbox\\{sequence}\\editing echo editing', target=client['name'])
184 if 'editing' in result.get('stdout', ''): break
185 time.sleep(.1)
186 else: raise TimeoutError('Native stress client did not acknowledge its first edit')
187 start.touch()
188 if disconnect:
189 from native_disconnect import interrupt
190 interrupt(output, clients, sequences, processes)
191 if offline_outage or offline_lost_reply:
192 from offline_outage import interrupt
193 interrupt(output, clients, sequences, processes)
194 if embedded_smb:
195 import linux_vm
196 server = json.loads((output / 'run.json').read_text())['server']
197 with (output / 'server-locks.jsonl').open('w') as trace:
198 for _ in range(100):
199 captured = linux_vm.run_ssh(server, 'sudo smbstatus --byterange --json', timeout=5)
200 captured.check_returncode()
201 trace.write(json.dumps(json.loads(captured.stdout)) + '\n')
202 trace.flush()
203 time.sleep(.1)
204 for client, sequence in zip(clients, sequences):
205 wait_action(client, sequence, 'stress')
206 finally:
207 for client, sequence in zip(clients, sequences):
208 local = client['folder'] / 'stress-events.jsonl'
209 try:
210 result = windows.do_get(f'C:\\one-tests\\runs\\capture\\outbox\\{sequence}\\events.jsonl', local, client['name'])
211 except Exception as error:
212 result = {'error': str(error)}
213 (client['folder'] / 'stress-capture.json').write_text(json.dumps(result, indent=2))
214 for client, clock in zip(clients, clocks):
215 before = time.time_ns() // 1000
216 native = windows.do_health(client['name'])['utc_us']
217 after = time.time_ns() // 1000
218 clock['after_native_minus_host_us'] = [native - after - 15625, native - before + 15625]
219 (output / 'clocks.json').write_text(json.dumps(clocks, indent=2))
220 rust = {actor: [json.loads(line) for line in (folder / f'{actor}.jsonl').read_text().splitlines()] for actor in processes}
221 commits, expected_rust = edit_history(rust, operations, edit, offline=config.get('offline', False))
222 assert sum(event.get('operations', 1) for event in commits) == rust_writers * operations, 'Missing Rust acknowledgements'
223 native_events = []
224 for client in clients:
225 local = client['folder'] / 'stress-events.jsonl'
226 events = [json.loads(line) for line in local.read_text(encoding='utf-8-sig').splitlines()]
227 native_events.append(events)
228 prefixes = [native_history(events, i, operations, edit) for i, events in enumerate(native_events)]
229 expected = [expected_rust, *prefixes]
230 documents = {}
231 if config.get('document_operations'):
232 from offline_document_history import document_history
233 documents = document_history(rust, operations)
234 expected.extend(''.join(char for char, *_ in list(document['states'].values())[-1]['characters']) for document in documents.values())
235 for client in clients: wait_text(client, expected)
236 checkpoint('stress-final')
237 model = json.loads((output / 'stress-final/model/document.json').read_text())
238 for _, _, revision, _ in ordered_pages(model):
239 assert not revision['nodes'][revision['roots']['1']]['spaces'], 'Disjoint edits created conflict pages'
240 paragraphs = [n['kind']['text'] for _, _, revision, page in ordered_pages(model)
241 for _, n in walk(revision, page) if n['kind']['type'] == 'RichText']
242 assert sorted(paragraphs) == sorted(expected), 'Final shared state lost, duplicated or added unrecorded content'
243 if documents:
244 from offline_document_history import verify_model
245 verify_model(model, documents)
246 if edit:
247 resolved = json.loads((output / 'stress-final/model/text.json').read_text())
248 for i, events in enumerate(native_events):
249 runs, = [resolved[sid][rid][oid] for sid, rid, revision, page in ordered_pages(model)
250 for oid, node in walk(revision, page)
251 if node['kind']['type'] == 'RichText' and node['kind']['text'] == prefixes[i]]
252 for run in runs:
253 for field in ('bold', 'italic'):
254 assert bool(run['format'][field]) == events[-1][field], 'Final native formatting differs from its last edit'
255 reads = [e for events in rust.values() for e in events if e['event'] == 'read']
256 overlap = 0
257 for events, clock in zip(native_events, clocks):
258 low = min(clock['native_minus_host_us'][0], clock['after_native_minus_host_us'][0])
259 high = max(clock['native_minus_host_us'][1], clock['after_native_minus_host_us'][1])
260 for native in events:
261 latest_start = (native['update_started_ticks'] - 621355968000000000) // 10 - low
262 earliest_end = (native['updated_ticks'] - 621355968000000000) // 10 - high
263 overlap += sum(max(latest_start, rust['started_us']) < min(earliest_end, rust['finished_us']) for rust in commits)
264 result = {'native_writers': len(clients), 'rust_writers': rust_writers, 'rust_readers': rust_readers,
265 'sync_every': sync_every, 'edit': edit, 'seed': seed, 'offline': config.get('offline', False), 'native_edits': operations * len(clients), 'rust_commits': len(commits), 'document_commits': sum(len(document['states']) for document in documents.values()), 'rust_reads': len(reads),
266 'clock_bounded_native_rust_call_overlaps': overlap, 'converged_paragraphs': len(expected)}
267 (output / 'result.json').write_text(json.dumps(result, indent=2))
268 assert overlap, 'No native/Rust call overlap established within clock uncertainty'
269 print(json.dumps(result), flush=True)
270
271
272if __name__ == '__main__':
273 import argparse
274 from pathlib import Path
275
276 parser = argparse.ArgumentParser(description='Check a fresh native capture against concurrent editing intents.')
277 parser.add_argument('run', type=Path)
278 parser.add_argument('capture', type=Path)
279 args = parser.parse_args()
280 print(json.dumps(verify_capture(args.run, args.capture), indent=2))