| 1 | #!/usr/bin/env python3 |
| 2 | """Kill owned processes across local-cache/remote-file publication boundaries.""" |
| 3 | import argparse |
| 4 | import cache_images |
| 5 | import hashlib |
| 6 | import json |
| 7 | import os |
| 8 | from pathlib import Path |
| 9 | import queue |
| 10 | import shutil |
| 11 | import signal |
| 12 | import sqlite3 |
| 13 | import subprocess |
| 14 | import threading |
| 15 | import time |
| 16 | |
| 17 | ROOT = Path(__file__).resolve().parent.parent |
| 18 | BINARY = ROOT / 'target/debug/examples/recovery_probe' |
| 19 | TOKEN = ' [offline-recovery]' |
| 20 | |
| 21 | |
| 22 | def 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 | |
| 31 | def 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 | |
| 85 | def 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 | |
| 109 | def 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 | |
| 116 | def 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 | |
| 122 | def 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 | |
| 135 | def 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 | |
| 204 | def 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 | |
| 230 | if __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()) |