1import copy
2import unittest
3from native_stress import edit_history
4
5
6class 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
106class 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()