| 1 | """Require durable local progress during a confirmed outage of the owned SMB proxy.""" |
| 2 | import hashlib |
| 3 | import json |
| 4 | from pathlib import Path |
| 5 | import shlex |
| 6 | import time |
| 7 | |
| 8 | from native_runner import windows |
| 9 | import linux_vm |
| 10 | from verify_smb_overlap import verify |
| 11 | from offline_document_history import operation_kinds |
| 12 | |
| 13 | |
| 14 | def interrupt(output, clients, sequences, processes): |
| 15 | config = json.loads((output / 'run.json').read_text()) |
| 16 | folder = output / 'rust' |
| 17 | samples = {} |
| 18 | |
| 19 | def logs(): |
| 20 | result = {} |
| 21 | for actor in processes: |
| 22 | text = (folder / f'{actor}.jsonl').read_text() |
| 23 | result[actor] = [json.loads(line) for line in text[:text.rfind('\n') + 1].splitlines()] |
| 24 | return result |
| 25 | |
| 26 | def wait_for(predicate, message, timeout=90): |
| 27 | deadline = time.monotonic() + timeout |
| 28 | while not predicate(): |
| 29 | assert all(p.poll() in (None, 0) for p in processes.values()), 'A client failed during the offline outage' |
| 30 | if time.monotonic() > deadline: raise TimeoutError(message) |
| 31 | time.sleep(.05) |
| 32 | |
| 33 | def ssh(command): |
| 34 | result = linux_vm.run_ssh(config['server'], command, timeout=15) |
| 35 | with (output / 'offline-outage-server.jsonl').open('a') as stream: |
| 36 | stream.write(json.dumps({'command': command, 'exit': result.returncode, 'stdout': result.stdout, 'stderr': result.stderr}) + '\n') |
| 37 | result.check_returncode() |
| 38 | return result.stdout |
| 39 | |
| 40 | def phase(control): |
| 41 | text = json.dumps(control) |
| 42 | ssh("printf '%s' " + shlex.quote(text) + ' > /tmp/smb-control.tmp && mv /tmp/smb-control.tmp /tmp/smb-control.json') |
| 43 | wait_for(lambda: text in ssh("grep -F '\"control\":' /tmp/smb-trace.jsonl | tail -n 1"), 'Proxy did not acknowledge the outage phase') |
| 44 | |
| 45 | def native_counts(label): |
| 46 | counts = [] |
| 47 | for actor, (client, sequence) in enumerate(zip(clients, sequences)): |
| 48 | capture = output / f'offline-{label}-n{actor}.jsonl' |
| 49 | result = windows.do_get(f'C:\\one-tests\\runs\\capture\\outbox\\{sequence}\\events.jsonl', capture, client['name']) |
| 50 | assert not result.get('error'), result |
| 51 | text = capture.read_text(encoding='utf-8-sig') |
| 52 | rows = [json.loads(line) for line in text[:text.rfind('\n') + 1].splitlines()] |
| 53 | assert 0 < len(rows) < config['stress_operations'], 'Native client was inactive during the outage campaign' |
| 54 | counts.append(len(rows)) |
| 55 | return counts |
| 56 | |
| 57 | writers = [actor for actor in processes if actor.startswith('w')] |
| 58 | readers = [actor for actor in processes if actor.startswith('r')] |
| 59 | if config.get('offline_lost_reply'): |
| 60 | isolated = config.get('offline_client_reply', False) |
| 61 | def counts(): |
| 62 | return {actor: sum(row['event'] == ('remote_receipt' if actor in writers else 'read') for row in rows) |
| 63 | for actor, rows in logs().items()} |
| 64 | if config.get('document_operations'): |
| 65 | wait_for(lambda: all((folder / f'offline-paused-{actor}').exists() for actor in writers) |
| 66 | and all(counts()[actor] > 0 for actor in readers), 'Clients did not reach the formatting publication barrier') |
| 67 | else: |
| 68 | wait_for(lambda: all(value >= 3 for value in counts().values()), 'Clients made no progress before the reply cut') |
| 69 | samples['before'] = counts() |
| 70 | samples['native_before'] = native_counts('reply-before') |
| 71 | try: |
| 72 | control = {'phase': 'offline-reply-cut', 'cut': 9, 'peer': '10.0.2.2', 'offset': 96, 'direction': 'response'} |
| 73 | if isolated: |
| 74 | control['scope'] = 'connection' |
| 75 | (folder / 'offline-paused-w0.isolate').touch() |
| 76 | phase(control) |
| 77 | if config.get('document_operations'): |
| 78 | samples['format_released_us'] = time.time_ns() // 1000 |
| 79 | (folder / 'offline-paused-w0.resume').touch() |
| 80 | wait_for(lambda: int(ssh("grep -c '\"cut\": {' /tmp/smb-trace.jsonl || true").strip()) == 1, |
| 81 | 'The publication reply was not interrupted') |
| 82 | samples['down_started_us'] = time.time_ns() // 1000 |
| 83 | if isolated: |
| 84 | wait_for(lambda: any(row['event'] == 'confirmation_paused' for row in logs()['w0']), |
| 85 | 'The disconnected writer did not retain its uncertain publication') |
| 86 | for actor in writers[1:]: (folder / f'offline-paused-{actor}.resume').touch() |
| 87 | wait_for(lambda: all(value >= samples['before'][actor] + 3 for actor, value in counts().items() if actor != 'w0'), |
| 88 | 'Peers did not advance while the writer was disconnected') |
| 89 | wait_for(lambda: (folder / 'offline-retired.one').exists(), |
| 90 | 'Native maintenance did not retire the isolated revision') |
| 91 | time.sleep(3) |
| 92 | if isolated: samples['during'] = counts() |
| 93 | samples['native_during'] = native_counts('reply-during') |
| 94 | finally: |
| 95 | samples['up_started_us'] = time.time_ns() // 1000 |
| 96 | phase({'phase': 'offline-reply-reconnected'}) |
| 97 | if config.get('document_operations'): |
| 98 | for actor in writers: (folder / f'offline-paused-{actor}.resume').touch() |
| 99 | if isolated: (folder / 'offline-paused-w0.confirmation-resume').touch() |
| 100 | (output / 'offline-lost-reply-progress.json').write_text(json.dumps(samples, indent=2)) |
| 101 | wait_for(lambda: all(value >= samples['before'][actor] + 3 for actor, value in counts().items()), |
| 102 | 'A client failed to progress after the lost publication reply', timeout=120) |
| 103 | samples['after'] = counts() |
| 104 | (output / 'offline-lost-reply-progress.json').write_text(json.dumps(samples, indent=2)) |
| 105 | return |
| 106 | wait_for(lambda: all((folder / f'offline-paused-{actor}').exists() for actor in writers) |
| 107 | and (not config.get('document_operations') or all( |
| 108 | sum(row['event'] == 'local_document_commit' for row in logs()[actor]) == len(operation_kinds(logs()[actor])) for actor in writers)) |
| 109 | and all(any(row['event'] == 'read' for row in logs()[actor]) for actor in readers), |
| 110 | 'Clients did not reach the pre-publication outage barrier') |
| 111 | samples['native_before'] = native_counts('before') |
| 112 | samples['reader_errors_before'] = {actor: sum(row['event'] == 'transport_read_error' for row in logs()[actor]) for actor in readers} |
| 113 | try: |
| 114 | phase({'phase': 'offline-down', 'mode': 'down'}) |
| 115 | samples['down_started_us'] = time.time_ns() // 1000 |
| 116 | (folder / 'offline-outage-down').touch() |
| 117 | wait_for(lambda: all(sum(row['event'] == 'local_commit' for row in events) == 8 for actor, events in logs().items() if actor in writers) |
| 118 | and (not config.get('document_operations') or all( |
| 119 | sum(row['event'] == 'local_document_commit' for row in logs()[actor]) == 8 * len(operation_kinds(logs()[actor])) for actor in writers)) |
| 120 | and all(sum(row['event'] == 'transport_read_error' for row in logs()[actor]) > samples['reader_errors_before'][actor] for actor in readers), |
| 121 | 'Local queues or disconnected readers failed to progress during the outage') |
| 122 | time.sleep(3) |
| 123 | samples['native_during'] = native_counts('during') |
| 124 | samples['up_started_us'] = time.time_ns() // 1000 |
| 125 | finally: |
| 126 | phase({'phase': 'offline-reconnected'}) |
| 127 | (folder / 'offline-outage-resumed').touch() |
| 128 | (output / 'offline-outage-progress.json').write_text(json.dumps(samples, indent=2)) |
| 129 | wait_for(lambda: all(sum(row['event'] == 'remote_receipt' for row in logs()[actor]) >= 3 for actor in writers) |
| 130 | and all(sum(row['event'] == 'read' and row['started_us'] > samples['up_started_us'] for row in logs()[actor]) >= 3 for actor in readers), |
| 131 | 'A client failed to progress after the offline outage') |
| 132 | |
| 133 | |
| 134 | def native_progress(output, config, sample): |
| 135 | down, up = sample['down_started_us'], sample['up_started_us'] |
| 136 | assert len(sample['native_before']) == len(sample['native_during']) == config['stress_clients'] |
| 137 | native = list(zip(sample['native_before'], sample['native_during'])) |
| 138 | assert all(0 < before < during < config['stress_operations'] for before, during in native), 'Native local edits did not advance during the outage' |
| 139 | clocks = json.loads((output / 'clocks.json').read_text()) |
| 140 | native_inside = [] |
| 141 | assert len(clocks) == config['stress_clients'] |
| 142 | for index, clock in enumerate(clocks): |
| 143 | low = min(clock['native_minus_host_us'][0], clock['after_native_minus_host_us'][0]) |
| 144 | high = max(clock['native_minus_host_us'][1], clock['after_native_minus_host_us'][1]) |
| 145 | rows = [json.loads(line) for line in (output / f'n{index}/stress-events.jsonl').read_text(encoding='utf-8-sig').splitlines()] |
| 146 | inside = sum((row['update_started_ticks'] - 621355968000000000) // 10 - high > down |
| 147 | and (row['updated_ticks'] - 621355968000000000) // 10 - low < up for row in rows) |
| 148 | assert inside, 'Native timestamps do not establish local edits inside the confirmed outage' |
| 149 | native_inside.append(inside) |
| 150 | return native_inside |
| 151 | |
| 152 | |
| 153 | def verify_outage(output): |
| 154 | config = json.loads((output / 'run.json').read_text()) |
| 155 | assert config['offline'] and config['offline_outage'] and config['embedded_smb'] |
| 156 | assert config['stress_clients'] + config['rust_writers'] + config['rust_readers'] >= 12 |
| 157 | sample = json.loads((output / 'offline-outage-progress.json').read_text()) |
| 158 | down, up = sample['down_started_us'], sample['up_started_us'] |
| 159 | assert up - down >= 3_000_000, 'Confirmed outage lasted less than three seconds' |
| 160 | native_inside = native_progress(output, config, sample) |
| 161 | queues, reconnects = {}, {} |
| 162 | for mode, count in [('w', config['rust_writers']), ('r', config['rust_readers'])]: |
| 163 | for index in range(count): |
| 164 | actor = f'{mode}{index}' |
| 165 | events = [json.loads(line) for line in (output / 'rust' / f'{actor}.jsonl').read_text().splitlines()] |
| 166 | assert events[0]['event'] == 'ready' and events[-1]['event'] == 'done' |
| 167 | connected = [row for row in events if row['event'] == 'transport_connected'] |
| 168 | assert any(row['at_us'] > up for row in connected), 'No fresh transport after outage' |
| 169 | reconnects[actor] = len(connected) |
| 170 | if mode == 'r': |
| 171 | assert sum(row['event'] == 'transport_read_error' for row in events) > sample['reader_errors_before'][actor] |
| 172 | assert sum(row['event'] == 'read' and row['started_us'] > up for row in events) >= 3 |
| 173 | continue |
| 174 | local = [row for row in events if row['event'] == 'local_commit'] |
| 175 | assert len(local) == config['stress_operations'] |
| 176 | assert local[0]['finished_us'] < down |
| 177 | queued = [row for row in local if down < row['started_us'] <= row['finished_us'] < up] |
| 178 | assert [row['operation'] for row in queued] == list(range(1, 8)), 'Seven local edits were not accepted while SMB was down' |
| 179 | queues[actor] = len(queued) |
| 180 | paused = [row for row in events if row['event'] == 'publication_paused'] |
| 181 | assert len(paused) == 1 and paused[0]['at_us'] < down |
| 182 | attempts = [row for row in events if row['event'] == 'remote_attempt'] |
| 183 | assert attempts and all(row['started_us'] > up for row in attempts), 'Publication escaped the outage barrier' |
| 184 | assert attempts[0]['state'] == 'NotCommitted' and attempts[0]['revision'] == paused[0]['revision'], 'Disconnected pre-I/O attempt was not safely rejected' |
| 185 | assert all(row['state'] != 'Unknown' for row in attempts), 'Unexpected uncertain publication requires separate recovery evidence' |
| 186 | receipts = [row for row in events if row['event'] == 'remote_receipt'] |
| 187 | assert len(receipts) == config['stress_operations'] and all(row['at_us'] > up for row in receipts) |
| 188 | if config.get('document_operations'): |
| 189 | kinds = operation_kinds(events) |
| 190 | edits = [row for row in events if row['event'] == 'local_document_commit'] |
| 191 | assert [(row['operation'], row['kind']) for row in edits] == [ |
| 192 | (operation, kind) for operation in range(config['stress_operations']) for kind in kinds] |
| 193 | assert all(row['finished_us'] < down for row in edits[:len(kinds)]), 'Initial document edits missed the outage barrier' |
| 194 | queued = [row for row in edits if down < row['started_us'] <= row['finished_us'] < up] |
| 195 | assert [(row['operation'], row['kind']) for row in queued] == [ |
| 196 | (operation, kind) for operation in range(1, 8) for kind in kinds], 'Document edits did not persist during the outage' |
| 197 | queues[actor] += len(queued) |
| 198 | receipts = [row for row in events if row['event'] == 'document_receipt'] |
| 199 | assert len(receipts) == config['stress_operations'] * len(kinds) and all(row['at_us'] > up for row in receipts) |
| 200 | trace = [json.loads(line) for line in (output / 'smb-trace.jsonl').read_text().splitlines()] |
| 201 | controls = [row['control'] for row in trace if row.get('control', {}).get('phase', '').startswith('offline-')] |
| 202 | assert controls == [{'phase': 'offline-down', 'mode': 'down'}, {'phase': 'offline-reconnected'}], 'Unexpected outage control sequence' |
| 203 | return {'outages': 1, 'confirmed_down_seconds': (up - down) / 1_000_000, |
| 204 | 'local_edits_while_down': queues, 'transport_connection_events': reconnects, |
| 205 | 'native_local_edits_inside_confirmed_outage': native_inside, |
| 206 | 'resumed_guarded_io_overlap': verify(trace, phase='offline-reconnected')} |
| 207 | |
| 208 | |
| 209 | def verify_lost_reply(output): |
| 210 | from offline_history import publication_links, tokens |
| 211 | config = json.loads((output / 'run.json').read_text()) |
| 212 | assert config['offline'] and config['offline_lost_reply'] and config['embedded_smb'] |
| 213 | assert config['stress_clients'] + config['rust_writers'] + config['rust_readers'] >= 12 |
| 214 | sample = json.loads((output / 'offline-lost-reply-progress.json').read_text()) |
| 215 | isolated = config.get('offline_client_reply', False) |
| 216 | assert sample['up_started_us'] - sample['down_started_us'] >= 3_000_000 |
| 217 | actors = [*(f'w{i}' for i in range(config['rust_writers'])), *(f'r{i}' for i in range(config['rust_readers']))] |
| 218 | assert set(sample['before']) == set(sample['after']) == set(actors) |
| 219 | assert all(sample['after'][actor] >= sample['before'][actor] + 3 for actor in actors) |
| 220 | if not config.get('document_operations'): |
| 221 | assert all(sample['before'][actor] >= 3 for actor in actors) |
| 222 | logs = {actor: [json.loads(line) for line in (output / 'rust' / f'{actor}.jsonl').read_text().splitlines()] for actor in actors} |
| 223 | publication_links(logs, config['stress_operations']) |
| 224 | documents = {} |
| 225 | if config.get('document_operations'): |
| 226 | from offline_document_history import document_history |
| 227 | documents = document_history(logs, config['stress_operations']) |
| 228 | for actor, rows in logs.items(): |
| 229 | if isolated and actor != 'w0': continue |
| 230 | assert any(row['event'] == 'transport_connected' and row['at_us'] > sample['up_started_us'] for row in rows), f'{actor} did not reconnect' |
| 231 | unknown = [(actor, row) for actor, rows in logs.items() for row in rows if row['event'] == 'remote_attempt' and row['state'] == 'Unknown'] |
| 232 | assert len(unknown) == 1, 'The reply cut did not establish exactly one uncertain publication' |
| 233 | actor, attempt = unknown[0] |
| 234 | assert attempt['started_us'] < sample['down_started_us'], 'Uncertain attempt started after the reply cut' |
| 235 | receipts = [row for row in logs[actor] if row['event'] in ('remote_receipt', 'document_receipt') and row['revision'] == attempt['revision']] |
| 236 | if not receipts and documents: |
| 237 | target, = attempt['document_changes'] |
| 238 | confirmed_revision = documents[target]['states']['format']['attempt']['receipt_revision'] |
| 239 | receipts = [row for row in logs[actor] if row['event'] == 'document_receipt' and row['revision'] == confirmed_revision] |
| 240 | assert len(receipts) == 1 and receipts[0]['at_us'] > sample['up_started_us'] |
| 241 | peer_progress = {} |
| 242 | if isolated: |
| 243 | assert sample['during']['w0'] == sample['before']['w0'] |
| 244 | retired, = [row for row in logs['r0'] if row['event'] == 'revision_retired'] |
| 245 | assert retired['space'] == attempt['space'] and retired['revision'] == attempt['revision'] |
| 246 | assert sample['down_started_us'] < retired['started_us'] <= retired['finished_us'] < sample['up_started_us'] |
| 247 | assert receipts[0]['revision'] != attempt['revision'], 'The client-disconnect gate did not exercise confirmation after revision retirement' |
| 248 | paused, = [row for row in logs['w0'] if row['event'] == 'confirmation_paused'] |
| 249 | assert paused['revision'] == attempt['revision'] and attempt['finished_us'] <= paused['at_us'] < sample['up_started_us'] |
| 250 | for peer, rows in logs.items(): |
| 251 | if peer == 'w0': continue |
| 252 | assert sample['during'][peer] >= sample['before'][peer] + 3 |
| 253 | if peer.startswith('w'): |
| 254 | progress = [row for row in rows if row['event'] == 'remote_attempt' and row['state'] == 'Committed' |
| 255 | and sample['down_started_us'] < row['started_us'] <= row['finished_us'] < sample['up_started_us']] |
| 256 | else: |
| 257 | progress = [row for row in rows if row['event'] == 'read' |
| 258 | and sample['down_started_us'] < row['started_us'] <= row['finished_us'] < sample['up_started_us']] |
| 259 | assert len(progress) >= 3, f'{peer} has insufficient completed I/O while the writer was disconnected' |
| 260 | peer_progress[peer] = len(progress) |
| 261 | if config.get('document_operations'): |
| 262 | assert actor == 'w0' and sample['format_released_us'] <= attempt['started_us'] |
| 263 | intent, = [row for row in logs[actor] if row['event'] == 'local_document_commit' and row['id'] == receipts[0]['id']] |
| 264 | assert intent['kind'] == 'format', 'The interrupted publication was not formatting' |
| 265 | for writer in (f'w{i}' for i in range(config['rust_writers'])): |
| 266 | paused = [row for row in logs[writer] if row['event'] == 'publication_paused'] |
| 267 | assert len(paused) == 1 and paused[0]['kind'] == 'format' and paused[0]['at_us'] < sample['format_released_us'] |
| 268 | paused, = [row for row in logs[actor] if row['event'] == 'publication_paused'] |
| 269 | if paused['revision'] != attempt['revision']: |
| 270 | prior, = [row for row in logs[actor] if row['event'] == 'remote_attempt' and row['revision'] == paused['revision']] |
| 271 | assert prior['state'] == 'NotCommitted' and sample['format_released_us'] <= prior['started_us'] <= prior['finished_us'] <= attempt['started_us'], 'Paused revision was replaced without proving it unpublished' |
| 272 | captures = {} |
| 273 | for rows in logs.values(): |
| 274 | for row in rows: |
| 275 | if row['event'] != 'remote_confirm': continue |
| 276 | name = row['capture'] |
| 277 | assert Path(name).name == name and name not in captures |
| 278 | tokens(row['text']) |
| 279 | assert row['started_us'] > sample['up_started_us'] |
| 280 | captures[name] = {'sha256': hashlib.sha256((output / 'rust/confirmations' / name).read_bytes()).hexdigest(), 'state': row['state']} |
| 281 | assert captures and set(captures) == {path.name for path in (output / 'rust/confirmations').glob('*.one')} |
| 282 | trace = [json.loads(line) for line in (output / 'smb-trace.jsonl').read_text().splitlines()] |
| 283 | controls = [row['control'] for row in trace if row.get('control', {}).get('phase', '').startswith('offline-')] |
| 284 | expected = {'phase': 'offline-reply-cut', 'cut': 9, 'peer': '10.0.2.2', 'offset': 96, 'direction': 'response'} |
| 285 | if isolated: expected['scope'] = 'connection' |
| 286 | assert controls == [expected, {'phase': 'offline-reply-reconnected'}] |
| 287 | cuts = [row['cut'] for row in trace if 'cut' in row] |
| 288 | assert len(cuts) == 1 and cuts[0]['direction'] == 'response' and cuts[0]['command'] == 9 and cuts[0]['status'] == '0x0' |
| 289 | request, = [row for row in trace if row.get('direction') == 'request' and row['connection'] == cuts[0]['connection'] and row['message'] == cuts[0]['message']] |
| 290 | assert request['offset'] == 96 and request['command'] == 9 |
| 291 | native_writes = 0 |
| 292 | if isolated: |
| 293 | peers = {row['connection']: row['peer'][0] for row in trace if row.get('opened')} |
| 294 | pending, files = {}, {} |
| 295 | for row in trace: |
| 296 | if row.get('command') not in (5, 6, 9): continue |
| 297 | key = row['connection'], row['message'] |
| 298 | if row['direction'] == 'request': |
| 299 | pending[key] = row |
| 300 | if row['command'] == 6: files.pop((key[0], row['file_id']), None) |
| 301 | elif row['status'] == '0x0' and key in pending: |
| 302 | request = pending.pop(key) |
| 303 | if row['command'] == 5: |
| 304 | files[key[0], row['file_id']] = request['path'].lower() |
| 305 | elif (row['command'] == 9 and peers[key[0]].startswith('192.168.77.') and row.get('written', 0) > 0 |
| 306 | and files.get((key[0], request['file_id']), '').endswith('synthetic.one') |
| 307 | and sample['down_started_us'] < request['time'] * 1_000_000 <= row['time'] * 1_000_000 < sample['up_started_us']): |
| 308 | native_writes += 1 |
| 309 | assert native_writes, 'No successful native writes while the isolated writer awaited reconciliation' |
| 310 | return {'reply_cuts': 1, 'uncertain_actor': actor, 'attempted_revision': attempt['revision'], 'confirmed_revision': receipts[0]['revision'], |
| 311 | 'peer_operations_during_client_disconnect': peer_progress, |
| 312 | 'native_writes_during_client_disconnect': native_writes, |
| 313 | 'retired_snapshot_sha256': hashlib.sha256((output / 'rust/offline-retired.one').read_bytes()).hexdigest() if isolated else None, |
| 314 | 'confirmation_snapshots': captures, 'native_local_edits_inside_confirmed_outage': native_progress(output, config, sample), |
| 315 | 'resumed_guarded_io_overlap': verify(trace, phase='offline-reply-reconnected')} |