1#!/usr/bin/env python3
2"""Audit completed offline/native runs, including the independent SQLite receipt ledger."""
3import argparse
4import hashlib
5import json
6from pathlib import Path
7import cache_images
8import sqlite3
9
10from native_stress import edit_history, native_history
11
12
13def verify(output, max_gap=120):
14 config = json.loads((output / 'run.json').read_text())
15 assert config['offline'] and config['embedded_smb'] and not config['edit']
16 assert config['stress_clients'] + config['rust_writers'] + config['rust_readers'] >= 12
17 actors = [*(f'w{i}' for i in range(config['rust_writers'])), *(f'r{i}' for i in range(config['rust_readers']))]
18 logs = {actor: [json.loads(line) for line in (output / 'rust' / f'{actor}.jsonl').read_text().splitlines()] for actor in actors}
19 commits, text = edit_history(logs, config['stress_operations'], offline=True)
20 documents = {}
21 if config.get('document_operations'):
22 from offline_document_history import document_history
23 documents = document_history(logs, config['stress_operations'])
24 started = (output / 'rust/start').stat().st_mtime_ns // 1000
25 stopped = (output / 'rust/stop').stat().st_mtime_ns // 1000
26 progress, queues, caches = {}, {}, {}
27 for actor, events in logs.items():
28 times = [event['at_us'] for event in events if event['event'] in ('remote_receipt', 'document_receipt')] if actor.startswith('w') else [event['finished_us'] for event in events if event['event'] == 'read']
29 assert times and times == sorted(times), 'Client progress is absent or went backwards'
30 if actor.startswith('r'): times.append(max(times[-1], stopped))
31 progress[actor] = max(b-a for a, b in zip([started, *times], times)) / 1_000_000
32 if not actor.startswith('w'): continue
33 edits = {event['id']: event for event in events if event['event'] in ('local_commit', 'local_document_commit')}
34 receipts = {event['id']: event for event in events if event['event'] in ('remote_receipt', 'document_receipt')}
35 attempts = {event['revision']: event for event in events if event['event'] == 'remote_attempt'}
36 publication_starts = {id: documents[edit.get('document', edit['object'])]['states'][edit['kind']]['attempt']['started_us'] if edit['event'] == 'local_document_commit'
37 else attempts[receipts[id]['revision']]['started_us'] for id, edit in edits.items()}
38 queues[actor] = max(sum(edit['finished_us'] <= at < publication_starts[id] for id, edit in edits.items()) for at in [edit['finished_us'] for edit in edits.values()])
39 assert queues[actor] >= 2, 'Writer did not establish a durable local queue before publication'
40 path = output / 'rust' / f'{actor}.sqlite'
41 connection = sqlite3.connect(path.resolve().as_uri() + '?mode=ro', uri=True)
42 try:
43 assert connection.execute('PRAGMA quick_check').fetchall() == [('ok',)], 'Cache integrity failed'
44 assert connection.execute('PRAGMA foreign_key_check').fetchall() == [], 'Cache foreign keys failed'
45 for table in ('edits', 'batches', 'payloads', 'remote'):
46 assert connection.execute(f'SELECT count(*) FROM {table}').fetchone() == (0,), 'Completed cache retains unresolved state'
47 persisted = dict(connection.execute('SELECT edit_id, revision FROM receipts'))
48 assert persisted == {id: event['revision'] for id, event in receipts.items()}, 'SQLite receipts differ from observed acknowledgements'
49 base = cache_images.image(connection)
50 caches[actor] = {'receipts': len(persisted), 'image_bytes': len(base), 'image_sha256': hashlib.sha256(base).hexdigest(), 'database_sha256': hashlib.sha256(path.read_bytes()).hexdigest()}
51 finally:
52 connection.close()
53 for i in range(config['stress_clients']):
54 events = [json.loads(line) for line in (output / f'n{i}/stress-events.jsonl').read_text(encoding='utf-8-sig').splitlines()]
55 native_history(events, i, config['stress_operations'], False)
56 times = [event['updated_ticks'] for event in events]
57 assert times == sorted(times)
58 progress[f'n{i}'] = max((b-a for a, b in zip(times, times[1:])), default=0) / 10_000_000
59 assert all(gap <= max_gap for gap in progress.values()), f'Client progress exceeded {max_gap}s: {progress}'
60 return {'remote_publications': len(commits), 'document_publications': sum(len(document['states']) for document in documents.values()), 'remote_text_sha256': hashlib.sha256(text.encode()).hexdigest(), 'maximum_progress_gap_seconds': progress,
61 'queued_before_publication_lower_bound': queues, 'reviewed_placements': sum(event['event'] == 'reviewed_append' for events in logs.values() for event in events), 'caches': caches}
62
63
64if __name__ == '__main__':
65 parser = argparse.ArgumentParser(description=__doc__)
66 parser.add_argument('output', type=Path)
67 parser.add_argument('--max-gap', type=float, default=120)
68 args = parser.parse_args()
69 result = verify(args.output, args.max_gap)
70 (args.output / 'offline-verification.json').write_text(json.dumps(result, indent=2))
71 print(json.dumps(result, indent=2))