| 1 | import copy |
| 2 | import unittest |
| 3 | from native_stress import edit_history |
| 4 | |
| 5 | |
| 6 | class OfflineHistoryTests(unittest.TestCase): |
| 7 | def setUp(self): |
| 8 | base = 'Concurrent edits:' |
| 9 | self.logs = {} |
| 10 | for actor, started, before in [('w1', 5, base), ('w0', 10, base + ' [w1:0]')]: |
| 11 | token = f' [{actor}:0]' |
| 12 | self.logs[actor] = [ |
| 13 | {'event': 'ready', 'offline': True}, |
| 14 | {'event': 'local_commit', 'id': 1, 'operation': 0, 'before': base, 'token': token, 'started_us': 1, 'finished_us': 2}, |
| 15 | {'event': 'read', 'text': before, 'started_us': started - 1, 'finished_us': started}, |
| 16 | {'event': 'remote_attempt', 'revision': actor, 'before': before, 'after': before + token, 'state': 'Committed', 'started_us': started, 'finished_us': started + 1}, |
| 17 | {'event': 'remote_receipt', 'id': 1, 'revision': actor, 'at_us': started + 2}, |
| 18 | {'event': 'reopened_receipt', 'id': 1, 'revision': actor}, |
| 19 | {'event': 'done'}, |
| 20 | ] |
| 21 | self.logs['r0'] = [{'event': 'ready'}, {'event': 'read', 'text': base + ' [w1:0] [w0:0]', 'started_us': 14, 'finished_us': 15}, {'event': 'done'}] |
| 22 | |
| 23 | def test_private_local_branches_require_one_separate_remote_history(self): |
| 24 | commits, text = edit_history(self.logs, 1, offline=True) |
| 25 | self.assertEqual([event['token'] for event in commits], [' [w1:0]', ' [w0:0]']) |
| 26 | self.assertEqual(text, self.logs['r0'][1]['text']) |
| 27 | |
| 28 | def test_a_save_that_replaced_a_pending_save_shares_its_intent_and_publication(self): |
| 29 | base = 'Concurrent edits:' |
| 30 | logs = {'w0': [ |
| 31 | {'event': 'ready', 'offline': True}, |
| 32 | {'event': 'local_commit', 'id': 1, 'operation': 0, 'before': base, 'token': ' [w0:0]', 'started_us': 1, 'finished_us': 2}, |
| 33 | {'event': 'local_commit', 'id': 1, 'operation': 1, 'before': base + ' [w0:0]', 'token': ' [w0:1]', 'started_us': 3, 'finished_us': 4}, |
| 34 | {'event': 'remote_attempt', 'revision': 'a', 'before': base, 'after': base + ' [w0:0] [w0:1]', 'state': 'Committed', 'started_us': 5, 'finished_us': 6}, |
| 35 | {'event': 'remote_receipt', 'id': 1, 'revision': 'a', 'at_us': 7}, |
| 36 | {'event': 'remote_receipt', 'id': 1, 'revision': 'a', 'at_us': 7}, |
| 37 | {'event': 'reopened_receipt', 'id': 1, 'revision': 'a'}, |
| 38 | {'event': 'reopened_receipt', 'id': 1, 'revision': 'a'}, |
| 39 | {'event': 'done'}, |
| 40 | ], 'r0': [{'event': 'ready'}, {'event': 'done'}]} |
| 41 | commits, text = edit_history(logs, 2, offline=True) |
| 42 | self.assertEqual([(event['token'], event['operations']) for event in commits], [(' [w0:0] [w0:1]', 2)]) |
| 43 | self.assertEqual(text, base + ' [w0:0] [w0:1]') |
| 44 | logs['w0'][3]['after'] = base + ' [w0:1]' |
| 45 | with self.assertRaises(AssertionError): edit_history(logs, 2, offline=True) |
| 46 | |
| 47 | def test_false_receipts_lost_local_intents_and_changed_reopen_state_fail(self): |
| 48 | for index, field, value in [(1, 'id', 2), (1, 'token', ' [w0:9]'), (3, 'state', 'Unknown'), |
| 49 | (3, 'after', 'Concurrent edits: [w1:0]'), (4, 'at_us', 0), |
| 50 | (5, 'revision', 'changed'), (3, 'revision', 'missing')]: |
| 51 | with self.subTest(index=index, field=field): |
| 52 | logs = copy.deepcopy(self.logs) |
| 53 | logs['w0'][index][field] = value |
| 54 | with self.assertRaises(AssertionError): edit_history(logs, 1, offline=True) |
| 55 | for event in ('local_commit', 'remote_receipt', 'remote_attempt', 'reopened_receipt'): |
| 56 | logs = copy.deepcopy(self.logs) |
| 57 | logs['w0'] = [item for item in logs['w0'] if item['event'] != event] |
| 58 | with self.assertRaises(AssertionError): edit_history(logs, 1, offline=True) |
| 59 | |
| 60 | def test_stale_reads_future_local_tokens_and_duplicate_acknowledgements_fail(self): |
| 61 | logs = copy.deepcopy(self.logs) |
| 62 | logs['r0'][1]['text'] = 'Concurrent edits:' |
| 63 | with self.assertRaises(AssertionError): edit_history(logs, 1, offline=True) |
| 64 | logs = copy.deepcopy(self.logs) |
| 65 | logs['w0'][1]['before'] += ' [w2:0]' |
| 66 | with self.assertRaises(AssertionError): edit_history(logs, 1, offline=True) |
| 67 | for event in ('local_commit', 'remote_receipt', 'remote_attempt', 'reopened_receipt'): |
| 68 | logs = copy.deepcopy(self.logs) |
| 69 | copied = next(item for item in logs['w0'] if item['event'] == event) |
| 70 | logs['w0'].insert(-1, copy.deepcopy(copied)) |
| 71 | with self.assertRaises(AssertionError): edit_history(logs, 1, offline=True) |
| 72 | |
| 73 | def test_uncertain_receipt_requires_the_original_target_revision_and_successful_confirmation(self): |
| 74 | rows = self.logs['w0'] |
| 75 | intent = next(row for row in rows if row['event'] == 'local_commit') |
| 76 | attempt = next(row for row in rows if row['event'] == 'remote_attempt') |
| 77 | receipt = next(row for row in rows if row['event'] == 'remote_receipt') |
| 78 | intent.update(space='page-space', object='paragraph') |
| 79 | attempt.update(space='page-space', object='paragraph', state='Unknown') |
| 80 | receipt['at_us'] = 14 |
| 81 | confirmation = dict(event='remote_confirm', state='Committed', started_us=12, finished_us=13, |
| 82 | revisions={'page-space': [attempt['revision']]}, text=attempt['after']) |
| 83 | rows.insert(4, confirmation) |
| 84 | self.assertEqual(len(edit_history(self.logs, 1, offline=True)[0]), 2) |
| 85 | original = copy.deepcopy(self.logs) |
| 86 | for field, value in [('state', 'NotCommitted'), ('state', 'Unknown'), ('finished_us', 15), |
| 87 | ('revisions', {'other-space': ['w0']}), ('revisions', {'page-space': ['wrong']}), ('text', 'Concurrent edits:')]: |
| 88 | self.logs = copy.deepcopy(original) |
| 89 | next(row for row in self.logs['w0'] if row['event'] == 'remote_confirm')[field] = value |
| 90 | with self.subTest(field=field, value=value), self.assertRaises(AssertionError): edit_history(self.logs, 1, offline=True) |
| 91 | self.logs = copy.deepcopy(original) |
| 92 | replay = dict(attempt, revision='another-attempt', state='NotCommitted', started_us=11, finished_us=12) |
| 93 | self.logs['w0'].insert(4, replay) |
| 94 | with self.assertRaisesRegex(AssertionError, 'replayed'): edit_history(self.logs, 1, offline=True) |
| 95 | |
| 96 | def test_partial_capture_does_not_promote_pending_local_success(self): |
| 97 | logs = copy.deepcopy(self.logs) |
| 98 | logs['w0'] = logs['w0'][:2] |
| 99 | logs['r0'] = [{'event': 'ready'}] |
| 100 | commits, text = edit_history(logs, 1, offline=True, partial=True) |
| 101 | self.assertEqual(len(commits), 1) |
| 102 | self.assertEqual(text, 'Concurrent edits: [w1:0]') |
| 103 | with self.assertRaises(AssertionError): edit_history(logs, 1, offline=True) |
| 104 | |
| 105 | |
| 106 | class OfflineLedgerTests(unittest.TestCase): |
| 107 | def test_independent_ledger_queue_depth_and_progress_checks(self): |
| 108 | import json |
| 109 | import os |
| 110 | from pathlib import Path |
| 111 | import sqlite3 |
| 112 | import tempfile |
| 113 | from verify_offline import verify |
| 114 | with tempfile.TemporaryDirectory() as folder: |
| 115 | output = Path(folder) |
| 116 | (output / 'rust').mkdir() |
| 117 | (output / 'run.json').write_text(json.dumps({'offline': True, 'embedded_smb': True, 'edit': False, 'stress_clients': 4, 'rust_writers': 4, 'rust_readers': 4, 'stress_operations': 2})) |
| 118 | (output / 'rust/stop').touch() |
| 119 | os.utime(output / 'rust/stop', ns=(0, 0)) |
| 120 | (output / 'rust/start').touch() |
| 121 | os.utime(output / 'rust/start', ns=(0, 0)) |
| 122 | logs = {f'w{i}': [{'event': 'ready', 'offline': True}] for i in range(4)} |
| 123 | for actor, events in logs.items(): |
| 124 | before = 'Concurrent edits:' |
| 125 | for op in range(2): |
| 126 | token = f' [{actor}:{op}]' |
| 127 | events.append({'event': 'local_commit', 'id': op+1, 'operation': op, 'token': token, 'before': before, 'started_us': 1+op*2, 'finished_us': 2+op*2}) |
| 128 | before += token |
| 129 | before = 'Concurrent edits:' |
| 130 | timestamp = 100 |
| 131 | for op in range(2): |
| 132 | for actor, events in logs.items(): |
| 133 | token, revision = f' [{actor}:{op}]', f'{actor}-{op}' |
| 134 | events.append({'event': 'remote_attempt', 'before': before, 'after': before+token, 'revision': revision, 'state': 'Committed', 'started_us': timestamp, 'finished_us': timestamp+1}) |
| 135 | events.append({'event': 'remote_receipt', 'id': op+1, 'revision': revision, 'at_us': timestamp+2}) |
| 136 | before += token |
| 137 | timestamp += 10 |
| 138 | for actor, events in logs.items(): |
| 139 | events.extend({'event': 'reopened_receipt', 'id': op+1, 'revision': f'{actor}-{op}'} for op in range(2)) |
| 140 | events.append({'event': 'done'}) |
| 141 | connection = sqlite3.connect(output / 'rust' / f'{actor}.sqlite') |
| 142 | # Schema 15 in WAL mode, as a drained cache leaves it. |
| 143 | connection.executescript('PRAGMA journal_mode=WAL; CREATE TABLE receipts(edit_id INTEGER, revision TEXT); CREATE TABLE edits(id INTEGER); CREATE TABLE batches(id INTEGER); CREATE TABLE payloads(sha256 BLOB); CREATE TABLE base(chunk INTEGER, bytes BLOB); CREATE TABLE remote(chunk INTEGER, bytes BLOB);') |
| 144 | connection.executemany('INSERT INTO receipts VALUES (?,?)', [(op+1, f'{actor}-{op}') for op in range(2)]) |
| 145 | connection.execute('INSERT INTO base VALUES (0, ?)', (b'opaque image',)) |
| 146 | connection.commit() |
| 147 | connection.close() |
| 148 | for i in range(4): |
| 149 | logs[f'r{i}'] = [{'event': 'ready'}, {'event': 'read', 'text': before, 'started_us': timestamp, 'finished_us': timestamp+1}, {'event': 'done'}] |
| 150 | (output / f'n{i}').mkdir() |
| 151 | prefix = f'Native {i}:' |
| 152 | native = [] |
| 153 | for op in range(2): |
| 154 | token = f' [n{i}:{op}]' |
| 155 | native.append({'operation': op, 'token': token, 'before': prefix, 'updated_ticks': op*10_000_000}) |
| 156 | prefix += token |
| 157 | (output / f'n{i}/stress-events.jsonl').write_text('\n'.join(json.dumps(event) for event in native)) |
| 158 | for actor, events in logs.items(): |
| 159 | (output / 'rust' / f'{actor}.jsonl').write_text('\n'.join(json.dumps(event) for event in events)) |
| 160 | result = verify(output, max_gap=2) |
| 161 | self.assertEqual(result['remote_publications'], 8) |
| 162 | self.assertEqual(result['queued_before_publication_lower_bound'], {f'w{i}': 2 for i in range(4)}) |
| 163 | with self.assertRaisesRegex(AssertionError, 'progress exceeded'): verify(output, max_gap=.5) |
| 164 | os.utime(output / 'rust/stop', ns=(3_000_000_000, 3_000_000_000)) |
| 165 | with self.assertRaisesRegex(AssertionError, 'progress exceeded'): verify(output, max_gap=2) |
| 166 | os.utime(output / 'rust/stop', ns=(0, 0)) |
| 167 | connection = sqlite3.connect(output / 'rust/w0.sqlite') |
| 168 | for sql, undo in [("UPDATE receipts SET revision='wrong' WHERE edit_id=1", "UPDATE receipts SET revision='w0-0' WHERE edit_id=1"), |
| 169 | ('INSERT INTO batches VALUES (1)', 'DELETE FROM batches'), |
| 170 | ("INSERT INTO remote VALUES (0, X'00')", 'DELETE FROM remote')]: |
| 171 | connection.execute(sql) |
| 172 | connection.commit() |
| 173 | with self.assertRaises(AssertionError): verify(output) |
| 174 | connection.execute(undo) |
| 175 | connection.commit() |
| 176 | connection.close() |