| 1 | #!/usr/bin/env python3 |
| 2 | """Audit completed offline/native runs, including the independent SQLite receipt ledger.""" |
| 3 | import argparse |
| 4 | import hashlib |
| 5 | import json |
| 6 | from pathlib import Path |
| 7 | import cache_images |
| 8 | import sqlite3 |
| 9 | |
| 10 | from native_stress import edit_history, native_history |
| 11 | |
| 12 | |
| 13 | def 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 | |
| 64 | if __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)) |