| 1 | """Overlapping native and Rust editing histories on one shared section.""" |
| 2 | import json |
| 3 | import os |
| 4 | import time |
| 5 | |
| 6 | from native_runner import windows |
| 7 | from concurrent_rust import running_clients |
| 8 | from document_model import ordered_pages, walk |
| 9 | |
| 10 | |
| 11 | def 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 | |
| 65 | def 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 | |
| 82 | def 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 | |
| 119 | def 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 | |
| 272 | if __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)) |