| 1 | import copy |
| 2 | import json |
| 3 | from pathlib import Path |
| 4 | import tempfile |
| 5 | import unittest |
| 6 | from unittest.mock import patch |
| 7 | |
| 8 | from offline_outage import verify_outage, verify_lost_reply |
| 9 | |
| 10 | |
| 11 | class OutageOracle(unittest.TestCase): |
| 12 | def setUp(self): |
| 13 | temporary = tempfile.TemporaryDirectory() |
| 14 | self.addCleanup(temporary.cleanup) |
| 15 | self.root = Path(temporary.name) |
| 16 | (self.root / 'rust').mkdir() |
| 17 | self.config = dict(offline=True, offline_outage=True, embedded_smb=True, stress_clients=4, rust_writers=4, rust_readers=4, stress_operations=8) |
| 18 | self.sample = dict(down_started_us=1_000_000, up_started_us=5_000_000, native_before=[1]*4, native_during=[4]*4, reader_errors_before={f'r{i}': 0 for i in range(4)}) |
| 19 | (self.root / 'clocks.json').write_text(json.dumps([dict(native_minus_host_us=[-10, 10], after_native_minus_host_us=[-10, 10])]*4)) |
| 20 | for i in range(4): |
| 21 | (self.root / f'n{i}').mkdir() |
| 22 | (self.root / f'n{i}/stress-events.jsonl').write_text(json.dumps(dict(update_started_ticks=621355968020000000, updated_ticks=621355968030000000))) |
| 23 | self.logs = {} |
| 24 | for i in range(4): |
| 25 | self.logs[f'w{i}'] = [dict(event='ready'), dict(event='transport_connected', at_us=0), |
| 26 | dict(event='publication_paused', revision=f'w{i}', at_us=100), |
| 27 | *[dict(event='local_commit', operation=j, started_us=10 if j == 0 else 2_000_000+j, finished_us=20 if j == 0 else 2_000_010+j) for j in range(8)], |
| 28 | dict(event='remote_attempt', started_us=6_000_000, state='NotCommitted', revision=f'w{i}'), |
| 29 | dict(event='transport_connected', at_us=6_000_001), |
| 30 | *[dict(event='remote_receipt', at_us=7_000_000+j) for j in range(8)], dict(event='done')] |
| 31 | self.logs[f'r{i}'] = [dict(event='ready'), dict(event='transport_read_error'), dict(event='transport_connected', at_us=6_000_000), |
| 32 | *[dict(event='read', started_us=7_000_000+j, finished_us=7_000_001+j) for j in range(3)], dict(event='done')] |
| 33 | self.trace = [dict(control=dict(phase='offline-down', mode='down')), dict(control=dict(phase='offline-reconnected'))] |
| 34 | |
| 35 | def verify(self): |
| 36 | (self.root / 'run.json').write_text(json.dumps(self.config)) |
| 37 | (self.root / 'offline-outage-progress.json').write_text(json.dumps(self.sample)) |
| 38 | (self.root / 'smb-trace.jsonl').write_text('\n'.join(map(json.dumps, self.trace))) |
| 39 | for actor, events in self.logs.items(): |
| 40 | (self.root / 'rust' / f'{actor}.jsonl').write_text('\n'.join(map(json.dumps, events))) |
| 41 | with patch('offline_outage.verify', return_value={'guarded_pairs': 1}) as overlap: |
| 42 | result = verify_outage(self.root) |
| 43 | overlap.assert_called_once_with(self.trace, phase='offline-reconnected') |
| 44 | return result |
| 45 | |
| 46 | def test_confirmed_outage_requires_local_and_native_progress_and_fresh_sessions(self): |
| 47 | self.assertEqual(self.verify()['local_edits_while_down'], {f'w{i}': 7 for i in range(4)}) |
| 48 | original = copy.deepcopy(self.logs) |
| 49 | for event, field, value in [('local_commit', 'started_us', 0), ('publication_paused', 'at_us', 3_000_000), |
| 50 | ('remote_attempt', 'state', 'Unknown'), ('remote_attempt', 'started_us', 3_000_000), |
| 51 | ('remote_attempt', 'revision', 'other'), ('remote_receipt', 'at_us', 0)]: |
| 52 | self.logs = copy.deepcopy(original) |
| 53 | row = next(row for row in self.logs['w0'] if row['event'] == event and (event != 'local_commit' or row['operation'] == 1)) |
| 54 | row[field] = value |
| 55 | with self.subTest(event=event, field=field), self.assertRaises(AssertionError): self.verify() |
| 56 | for actor, event in [('r0', 'transport_connected'), ('r0', 'transport_read_error'), ('r0', 'read'), ('w0', 'publication_paused')]: |
| 57 | self.logs = copy.deepcopy(original) |
| 58 | self.logs[actor] = [row for row in self.logs[actor] if row['event'] != event] |
| 59 | with self.subTest(actor=actor, event=event), self.assertRaises(AssertionError): self.verify() |
| 60 | |
| 61 | def test_lost_reply_requires_one_original_revision_receipt_and_captured_confirmation(self): |
| 62 | self.config.update(offline_outage=False, offline_lost_reply=True) |
| 63 | self.sample.update(before={actor: 3 for actor in self.logs}, after={actor: 6 for actor in self.logs}) |
| 64 | attempt = next(row for row in self.logs['w0'] if row['event'] == 'remote_attempt') |
| 65 | attempt.update(state='Unknown', started_us=500_000, finished_us=1_100_000) |
| 66 | receipt = next(row for row in self.logs['w0'] if row['event'] == 'remote_receipt') |
| 67 | receipt['revision'] = attempt['revision'] |
| 68 | for row in self.logs['w0']: |
| 69 | if row['event'] == 'remote_receipt' and row is not receipt: row['revision'] = 'other' |
| 70 | folder = self.root / 'rust/confirmations' |
| 71 | folder.mkdir() |
| 72 | (folder / 'confirmation.one').write_bytes(b'owned captured image') |
| 73 | confirmation = dict(event='remote_confirm', state='Committed', started_us=6_000_000, finished_us=6_000_001, |
| 74 | text='Concurrent edits: [w0:0]', capture='confirmation.one') |
| 75 | self.logs['w0'].insert(-1, confirmation) |
| 76 | self.trace = [dict(control=dict(phase='offline-reply-cut', cut=9, peer='10.0.2.2', offset=96, direction='response')), |
| 77 | dict(direction='request', command=9, offset=96, connection=1, message=2), |
| 78 | dict(cut=dict(direction='response', command=9, status='0x0', connection=1, message=2)), |
| 79 | dict(control=dict(phase='offline-reply-reconnected'))] |
| 80 | def check(): |
| 81 | (self.root / 'run.json').write_text(json.dumps(self.config)) |
| 82 | (self.root / 'offline-lost-reply-progress.json').write_text(json.dumps(self.sample)) |
| 83 | (self.root / 'smb-trace.jsonl').write_text('\n'.join(map(json.dumps, self.trace))) |
| 84 | for actor, rows in self.logs.items(): |
| 85 | (self.root / 'rust' / f'{actor}.jsonl').write_text('\n'.join(map(json.dumps, rows))) |
| 86 | with patch('offline_history.publication_links') as ledger, patch('offline_document_history.document_history') as documents, patch('offline_outage.verify', return_value={'guarded_pairs': 1}) as overlap: |
| 87 | if self.config.get('offline_client_reply'): |
| 88 | documents.return_value = {'target': {'states': {'format': {'attempt': {'receipt_revision': 'effect-revision'}}}}} |
| 89 | result = verify_lost_reply(self.root) |
| 90 | ledger.assert_called_once_with(self.logs, self.config['stress_operations']) |
| 91 | if self.config.get('document_operations'): |
| 92 | documents.assert_called_once_with(self.logs, self.config['stress_operations']) |
| 93 | else: |
| 94 | documents.assert_not_called() |
| 95 | overlap.assert_called_once_with(self.trace, phase='offline-reply-reconnected') |
| 96 | return result |
| 97 | self.assertEqual(check()['confirmed_revision'], 'w0') |
| 98 | self.config['document_operations'] = True |
| 99 | self.sample['format_released_us'] = 200_000 |
| 100 | receipt['event'] = 'document_receipt' |
| 101 | receipt['id'] = 17 |
| 102 | intent = dict(event='local_document_commit', id=17, kind='format') |
| 103 | self.logs['w0'].insert(-1, intent) |
| 104 | for actor, rows in self.logs.items(): |
| 105 | if actor.startswith('w'): |
| 106 | next(row for row in rows if row['event'] == 'publication_paused')['kind'] = 'format' |
| 107 | self.assertEqual(check()['confirmed_revision'], 'w0') |
| 108 | intent['kind'] = 'insert' |
| 109 | with self.assertRaisesRegex(AssertionError, 'not formatting'): check() |
| 110 | intent['kind'] = 'format' |
| 111 | self.sample['format_released_us'] = 900_000 |
| 112 | with self.assertRaises(AssertionError): check() |
| 113 | self.sample['format_released_us'] = 200_000 |
| 114 | paused = next(row for row in self.logs['w0'] if row['event'] == 'publication_paused') |
| 115 | paused['revision'] = 'stale-prepared-revision' |
| 116 | prior = dict(event='remote_attempt', revision=paused['revision'], state='NotCommitted', |
| 117 | started_us=200_001, finished_us=200_002) |
| 118 | self.logs['w0'].insert(-1, prior) |
| 119 | self.assertEqual(check()['confirmed_revision'], 'w0') |
| 120 | prior['finished_us'] = 600_000 |
| 121 | with self.assertRaisesRegex(AssertionError, 'proving it unpublished'): check() |
| 122 | self.logs['w0'].remove(prior) |
| 123 | paused['revision'] = attempt['revision'] |
| 124 | self.config['offline_client_reply'] = True |
| 125 | receipt['revision'] = 'effect-revision' |
| 126 | attempt.update(space='space', document_changes={'target': {}}) |
| 127 | retirement = dict(event='revision_retired', space='space', revision=attempt['revision'], started_us=2_200_000, finished_us=2_300_000) |
| 128 | self.logs['r0'].append(retirement) |
| 129 | (self.root / 'rust/offline-retired.one').write_bytes(b'retired revision snapshot') |
| 130 | self.sample['during'] = {actor: 3 if actor == 'w0' else 6 for actor in self.logs} |
| 131 | self.trace[0]['control']['scope'] = 'connection' |
| 132 | barrier = dict(event='confirmation_paused', revision=attempt['revision'], at_us=1_200_000) |
| 133 | self.logs['w0'].append(barrier) |
| 134 | peer_rows = [] |
| 135 | for actor, rows in self.logs.items(): |
| 136 | if actor == 'w0': continue |
| 137 | for i in range(3): |
| 138 | row = dict(event='remote_attempt' if actor.startswith('w') else 'read', state='Committed', |
| 139 | started_us=2_000_000+i, finished_us=2_000_010+i) |
| 140 | rows.append(row) |
| 141 | peer_rows.append((rows, row)) |
| 142 | wire = [dict(connection=7, opened=True, peer=['192.168.77.12', 445]), |
| 143 | dict(connection=7, message=2, command=5, direction='request', path='owned/synthetic.one'), |
| 144 | dict(connection=7, message=2, command=5, direction='response', status='0x0', file_id='native-file'), |
| 145 | dict(connection=7, message=3, command=9, direction='request', time=2.0, file_id='native-file'), |
| 146 | dict(connection=7, message=3, command=9, direction='response', time=2.1, status='0x0', written=4)] |
| 147 | self.trace.extend(wire) |
| 148 | result = check() |
| 149 | self.assertEqual(result['native_writes_during_client_disconnect'], 1) |
| 150 | self.assertEqual(result['peer_operations_during_client_disconnect'], {actor: 3 for actor in self.logs if actor != 'w0'}) |
| 151 | peer_rows[0][1]['finished_us'] = 6_000_000 |
| 152 | with self.assertRaisesRegex(AssertionError, 'insufficient completed I/O'): check() |
| 153 | peer_rows[0][1]['finished_us'] = 2_000_010 |
| 154 | wire[-1]['time'] = 6.0 |
| 155 | with self.assertRaisesRegex(AssertionError, 'No successful native writes'): check() |
| 156 | wire[-1]['time'] = 2.1 |
| 157 | barrier['revision'] = 'other' |
| 158 | with self.assertRaises(AssertionError): check() |
| 159 | self.logs['w0'].remove(barrier) |
| 160 | self.logs['r0'].remove(retirement) |
| 161 | receipt['revision'] = attempt['revision'] |
| 162 | for rows, row in peer_rows: rows.remove(row) |
| 163 | del self.trace[-len(wire):] |
| 164 | del self.trace[0]['control']['scope'] |
| 165 | self.config['offline_client_reply'] = False |
| 166 | self.config['document_operations'] = False |
| 167 | receipt['event'] = 'remote_receipt' |
| 168 | self.logs['w0'].remove(intent) |
| 169 | attempt['state'] = 'Committed' |
| 170 | with self.assertRaises(AssertionError): check() |
| 171 | attempt['state'] = 'Unknown' |
| 172 | self.trace[1]['offset'] = 100 |
| 173 | with self.assertRaises(AssertionError): check() |
| 174 | self.trace[1]['offset'] = 96 |
| 175 | (folder / 'unrecorded.one').write_bytes(b'extra') |
| 176 | with self.assertRaises(AssertionError): check() |
| 177 | |
| 178 | def test_document_outage_requires_both_local_operations_and_delayed_receipts(self): |
| 179 | self.config['document_operations'] = True |
| 180 | for actor, events in self.logs.items(): |
| 181 | if not actor.startswith('w'): continue |
| 182 | events[-1:-1] = [dict(event='local_document_commit', operation=operation, kind=kind, |
| 183 | started_us=10 if operation == 0 else 2_100_000+operation, |
| 184 | finished_us=20 if operation == 0 else 2_100_010+operation) |
| 185 | for operation in range(8) for kind in ('insert', 'format')] |
| 186 | events[-1:-1] = [dict(event='document_receipt', at_us=7_100_000+i) for i in range(16)] |
| 187 | self.assertEqual(self.verify()['local_edits_while_down'], {f'w{i}': 21 for i in range(4)}) |
| 188 | original = copy.deepcopy(self.logs) |
| 189 | for event, field, value in [('local_document_commit', 'finished_us', 6_000_000), |
| 190 | ('local_document_commit', 'kind', 'insert'), |
| 191 | ('document_receipt', 'at_us', 0)]: |
| 192 | self.logs = copy.deepcopy(original) |
| 193 | row = next(row for row in self.logs['w0'] if row['event'] == event |
| 194 | and (event != 'local_document_commit' or row['operation'] == 1 and row['kind'] == 'format')) |
| 195 | row[field] = value |
| 196 | with self.subTest(event=event, field=field), self.assertRaises(AssertionError): self.verify() |
| 197 | |
| 198 | def test_native_edits_outside_confirmed_window_fail(self): |
| 199 | (self.root / 'n0/stress-events.jsonl').write_text(json.dumps(dict(update_started_ticks=621355968000000000, updated_ticks=621355968000000100))) |
| 200 | with self.assertRaisesRegex(AssertionError, 'timestamps'): self.verify() |
| 201 | |
| 202 | def test_native_stall_short_outage_and_wrong_control_fail(self): |
| 203 | self.sample['native_before'].append(1) |
| 204 | with self.assertRaises(AssertionError): self.verify() |
| 205 | self.sample['native_before'].pop() |
| 206 | self.sample['native_during'][0] = 1 |
| 207 | with self.assertRaises(AssertionError): self.verify() |
| 208 | self.sample['native_during'][0] = 4 |
| 209 | self.sample['up_started_us'] = 1_000_001 |
| 210 | with self.assertRaises(AssertionError): self.verify() |
| 211 | self.sample['up_started_us'] = 5_000_000 |
| 212 | self.trace[0]['control'].pop('mode') |
| 213 | with self.assertRaises(AssertionError): self.verify() |
| 214 | |
| 215 | |
| 216 | if __name__ == '__main__': unittest.main() |