1"""Require durable local progress during a confirmed outage of the owned SMB proxy."""
2import hashlib
3import json
4from pathlib import Path
5import shlex
6import time
7
8from native_runner import windows
9import linux_vm
10from verify_smb_overlap import verify
11from offline_document_history import operation_kinds
12
13
14def 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
134def 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
153def 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
209def 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')}