1#!/usr/bin/env python3
2"""Kill owned processes across local-cache/remote-file publication boundaries."""
3import argparse
4import cache_images
5import hashlib
6import json
7import os
8from pathlib import Path
9import queue
10import shutil
11import signal
12import sqlite3
13import subprocess
14import threading
15import time
16
17ROOT = Path(__file__).resolve().parent.parent
18BINARY = ROOT / 'target/debug/examples/recovery_probe'
19TOKEN = ' [offline-recovery]'
20
21
22def run_command(binary, args, output):
23 result = subprocess.run([str(binary), *map(str, args)], capture_output=True, text=True, timeout=60,
24 env={**os.environ, 'ONESTORE_RECOVERY_PAUSE': ''})
25 output.with_suffix('.jsonl').write_text(result.stdout)
26 output.with_suffix('.stderr').write_text(result.stderr)
27 result.check_returncode()
28 return [json.loads(line) for line in result.stdout.splitlines()]
29
30
31def kill_at(binary, args, phase, output, receipt_window=None):
32 events = queue.Queue()
33 database = Path(args[1]) / 'cache.sqlite'
34 with output.with_suffix('.jsonl').open('w') as log, output.with_suffix('.stderr').open('w') as error:
35 process = subprocess.Popen([str(binary), *map(str, args)], stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=error, text=True,
36 env={**os.environ, 'ONESTORE_RECOVERY_PAUSE': phase})
37 def collect():
38 try:
39 for line in process.stdout:
40 log.write(line)
41 log.flush()
42 events.put(json.loads(line))
43 finally:
44 events.put(None)
45 reader = threading.Thread(target=collect)
46 reader.start()
47 try:
48 deadline = time.monotonic() + 60
49 while True:
50 event = events.get(timeout=max(0, deadline-time.monotonic()))
51 assert event is not None, f'Process exited before {phase}'
52 if event.get('event') == 'phase' and event['name'] == phase:
53 if receipt_window is not None:
54 assert args[0] == 'sync' and phase == 'publish-after'
55 assert receipt_window == 'wal'
56 # The receipt transaction is cut once its frames are in the log and
57 # its commit frame is not.
58 committed, _ = cache_images.wal_commits(database)
59 process.stdin.write('\n')
60 process.stdin.flush()
61 deadline = time.monotonic() + 5
62 while True:
63 commits, pending = cache_images.wal_commits(database)
64 if commits == committed and pending: break
65 assert commits == committed and process.poll() is None, 'Receipt transaction finished before the requested cut'
66 assert time.monotonic() < deadline, 'Receipt transaction window was not observed'
67 process.kill()
68 break
69 assert process.wait(timeout=10) == -signal.SIGKILL, 'Expected an actual process kill'
70 proof = {'pid': process.pid, 'exit': process.returncode, 'phase': phase}
71 if receipt_window is not None:
72 commits, pending = cache_images.wal_commits(database)
73 proof.update(wal_commits_before=committed, wal_commits_after_kill=commits, uncommitted_frames=pending, receipt_window=receipt_window)
74 assert commits == committed, 'Receipt transaction committed before the process died'
75 output.with_suffix('.cut.json').write_text(json.dumps(proof, indent=2))
76 finally:
77 if process.poll() is None: process.kill()
78 process.wait(timeout=10)
79 process.stdin.close()
80 reader.join(timeout=10)
81 assert not reader.is_alive(), 'Trace reader did not finish'
82 process.stdout.close()
83
84
85def state(rows, original):
86 found = [row for row in rows if row['event'] == 'state']
87 assert len(found) == 1, 'Missing independent post-reopen state'
88 result, = found
89 assert result['local_text'] == original + TOKEN, 'Locally acknowledged text was lost or duplicated'
90 assert result['remote_text'] in [original, original+TOKEN], 'Remote current text is partial, duplicated or invented'
91 assert result['status'] in ('pending', 'uncertain', 'published'), 'Intent disappeared or became an unexplained conflict'
92 if result['status'] == 'pending':
93 assert result['revision'] is None, 'Unattempted intent acquired a publication identity'
94 else:
95 assert isinstance(result['revision'], str) and result['revision'], 'Attempted identity was lost'
96 assert (result['revision'] == result['remote_revision']) == (result['remote_text'] == original+TOKEN), 'Publication identity disagrees with the visible effect'
97 if result['status'] == 'published':
98 assert result['pending'] == [] and result['remote_text'] == original+TOKEN
99 assert result['revision'] == result['remote_revision'], 'Receipt identifies another remote revision'
100 else:
101 assert len(result['pending']) == 1, 'Unacknowledged intent was lost or duplicated'
102 pending, = result['pending']
103 assert pending['id'] == 1 and pending['replacement'] == TOKEN
104 at = len(original.encode('utf-16-le')) // 2
105 assert pending['range'] == [at, at], 'Durable intent range changed'
106 return result
107
108
109def save_image(output, data):
110 digest = hashlib.sha256(data).hexdigest()
111 path = output / 'images' / f'{digest}.one'
112 if not path.exists(): path.write_bytes(data)
113 return digest
114
115
116def confirmation_only(events, before, after):
117 assert not any(row['event'] == 'phase' and row['name'] == 'publish-before' for row in events)
118 assert all(row['offset'] == 212 and row['bytes'] == 40 for row in events if row['event'] == 'write'), 'Recovery republished the remote edit'
119 assert before[:212] == after[:212] and before[252:] == after[252:], 'Confirmation changed content or the transaction count'
120
121
122def prepare_run(source, output):
123 output.mkdir(parents=True, exist_ok=False)
124 (output / 'images').mkdir()
125 binary = output / 'recovery_probe'
126 shutil.copyfile(BINARY, binary)
127 binary.chmod(0o755)
128 shutil.copyfile(__file__, output / Path(__file__).name)
129 source_hash = hashlib.sha256(source.read_bytes()).hexdigest()
130 (output / 'run.json').write_text(json.dumps({'source': str(source), 'source_sha256': source_hash,
131 'binary_sha256': hashlib.sha256(binary.read_bytes()).hexdigest(), 'controller_sha256': hashlib.sha256(Path(__file__).read_bytes()).hexdigest()}, indent=2))
132 return binary, source_hash
133
134
135def run(source, output):
136 binary, source_hash = prepare_run(source, output)
137 baseline = output / 'baseline'
138 initialized = run_command(binary, ['init', baseline, source], output / 'baseline-init')
139 original = next(row['remote_text'] for row in initialized if row['event'] == 'state')
140 state(initialized, original)
141 published = run_command(binary, ['sync', baseline], output / 'baseline-sync')
142 assert state(published, original)['status'] == 'published'
143 phases = [row['name'] for row in published if row['event'] == 'phase']
144 assert len(phases) == len(set(phases)), 'Baseline has ambiguous phase names'
145 assert 'publish-after' in phases and any(name.startswith('write-') for name in phases)
146 confirmation = output / 'confirmation-baseline'
147 run_command(binary, ['init', confirmation, source], output / 'confirmation-init')
148 kill_at(binary, ['sync', confirmation], 'publish-after', output / 'confirmation-setup')
149 confirmed = run_command(binary, ['sync', confirmation], output / 'confirmation-sync')
150 confirmation_phases = [row['name'] for row in confirmed if row['event'] == 'phase']
151 assert len(confirmation_phases) == len(set(confirmation_phases)) and 'confirm-after' in confirmation_phases
152 assert state(confirmed, original)['status'] == 'published'
153 cases = [('local-after', False), *((phase, False) for phase in phases), *((phase, True) for phase in confirmation_phases)]
154 results = []
155 for index, (phase, confirmation) in enumerate(cases):
156 folder = output / f'case-{index:02}'
157 trace = output / f'case-{index:02}-kill'
158 if phase == 'local-after':
159 kill_at(binary, ['init', folder, source], phase, trace)
160 else:
161 run_command(binary, ['init', folder, source], output / f'case-{index:02}-init')
162 if confirmation:
163 kill_at(binary, ['sync', folder], 'publish-after', output / f'case-{index:02}-setup')
164 kill_at(binary, ['sync', folder], phase, trace)
165 before = state(run_command(binary, ['inspect', folder], output / f'case-{index:02}-inspect'), original)
166 before_bytes = (folder / 'remote.one').read_bytes()
167 first = run_command(binary, ['sync', folder], output / f'case-{index:02}-recover')
168 after = state(first, original)
169 if before['status'] == 'uncertain' and before['remote_text'] == original:
170 assert after['status'] == 'uncertain' and after['revision'] == before['revision'], 'Absent uncertain attempt was replayed'
171 assert not any(row['event'] == 'write' for row in first)
172 else:
173 assert after['status'] == 'published'
174 if before['status'] != 'pending':
175 assert after['revision'] == before['revision'], 'Recovery published another revision'
176 confirmation_only(first, before_bytes, (folder / 'remote.one').read_bytes())
177 remote_hash = hashlib.sha256((folder / 'remote.one').read_bytes()).hexdigest()
178 repeated = run_command(binary, ['sync', folder], output / f'case-{index:02}-repeat')
179 assert state(repeated, original) == after, 'Repeated recovery changed durable intent state'
180 assert not any(row['event'] == 'write' for row in repeated), 'Repeated recovery published again'
181 assert hashlib.sha256((folder / 'remote.one').read_bytes()).hexdigest() == remote_hash
182 remote_image = save_image(output, (folder / 'remote.one').read_bytes())
183 connection = sqlite3.connect(folder / 'cache.sqlite')
184 try:
185 assert connection.execute('PRAGMA quick_check').fetchall() == [('ok',)]
186 assert connection.execute('PRAGMA foreign_key_check').fetchall() == []
187 # The image the queue applies to; the local text is the probe's, checked above.
188 base_image = save_image(output, cache_images.image(connection))
189 finally:
190 connection.close()
191 result = {'case': index, 'phase': phase, 'during_confirmation': confirmation, 'before_status': before['status'], 'after_status': after['status'],
192 'visible_before_recovery': before['remote_text'] != original, 'remote_image': remote_image, 'base_image': base_image,
193 'remote_text': after['remote_text'], 'local_text': after['local_text']}
194 results.append(result)
195 (output / 'results.json').write_text(json.dumps(results, indent=2))
196 print(json.dumps({key: value for key, value in result.items() if not key.endswith('_text')}), flush=True)
197 assert hashlib.sha256(source.read_bytes()).hexdigest() == source_hash, 'The frozen source changed'
198 summary = {'cases': len(results), 'process_kills': 1 + len(results) + sum(confirmation for _, confirmation in cases), 'uncertain_absent_preserved': sum(row['after_status']=='uncertain' for row in results),
199 'durable_receipts': sum(row['after_status']=='published' for row in results), 'unique_images': len(list((output/'images').glob('*.one')))}
200 (output/'summary.json').write_text(json.dumps(summary, indent=2))
201 print(json.dumps(summary), flush=True)
202
203
204def receipt_windows(source, output):
205 binary, source_hash = prepare_run(source, output)
206 results = []
207 for window in ('wal',):
208 folder = output / window
209 initialized = run_command(binary, ['init', folder, source], output / (window+'-init'))
210 original = next(row['remote_text'] for row in initialized if row['event'] == 'state')
211 state(initialized, original)
212 kill_at(binary, ['sync', folder], 'publish-after', output / (window+'-kill'), receipt_window=window)
213 before = state(run_command(binary, ['inspect', folder], output / (window+'-inspect')), original)
214 assert before['status'] == 'uncertain' and before['remote_text'] == original+TOKEN, 'Log recovery lost the pending confirmation'
215 before_bytes = (folder / 'remote.one').read_bytes()
216 recovered = run_command(binary, ['sync', folder], output / (window+'-recover'))
217 after = state(recovered, original)
218 assert after['status'] == 'published' and before['revision'] == after['revision']
219 confirmation_only(recovered, before_bytes, (folder / 'remote.one').read_bytes())
220 repeated = run_command(binary, ['sync', folder], output / (window+'-repeat'))
221 assert state(repeated, original) == after and not any(row['event']=='write' for row in repeated)
222 digest = save_image(output, (folder/'remote.one').read_bytes())
223 results.append({'receipt_window':window, 'remote_image':digest, 'remote_text':after['remote_text'], 'revision':after['revision']})
224 print(json.dumps({'receipt_window':window, 'status':after['status'], 'remote_image':digest}), flush=True)
225 assert hashlib.sha256(source.read_bytes()).hexdigest() == source_hash
226 (output/'results.json').write_text(json.dumps(results, indent=2))
227 (output/'summary.json').write_text(json.dumps({'process_kills':len(results), 'recovered_receipts':len(results), 'republished_edits':0}, indent=2))
228
229
230if __name__ == '__main__':
231 parser = argparse.ArgumentParser(description=__doc__)
232 parser.add_argument('source', type=Path)
233 parser.add_argument('output', type=Path)
234 parser.add_argument('--receipt-windows', action='store_true')
235 args = parser.parse_args()
236 (receipt_windows if args.receipt_windows else run)(args.source.resolve(), args.output.resolve())