#!/usr/bin/env python3 import argparse import json import os from pathlib import Path import random import socket import struct import tempfile import time from fs import new_pane, session from ninep import Client BODY_LIMIT = 4096 SCREEN_LIMIT = 4 * 1024 * 1024 FRAGMENTS = [ b'const value = 42;\n', b'fn check() void { if (true) return; }\n', b'\tspaces and tabs\n', 'λ 界 e\u0301 👩\u200d🚀\n'.encode(), b'x', ] OPERATIONS = ['append', 'range', 'replace', 'read', 'save', 'reload', 'cycle', 'reconnect', 'screen'] def emit(**record): print(json.dumps(record, separators=(',', ':')), flush=True) def payload(rng, limit=256): result = bytearray() for _ in range(rng.randint(1, 12)): fragment = rng.choice(FRAGMENTS) if len(result) + len(fragment) <= limit: result.extend(fragment) return bytes(result) or b'x'[:limit] def check_bytes(actual, expected, label): if actual == expected: return at = next((i for i, pair in enumerate(zip(actual, expected)) if pair[0] != pair[1]), min(len(actual), len(expected))) raise AssertionError(f'{label}: mismatch at byte {at}; ' f'length {len(actual)} != {len(expected)}; ' f'{actual[at:at + 32]!r} != {expected[at:at + 32]!r}') def read_fid(client, fid, limit, count): result = bytearray() while chunk := client.read_fid(fid, len(result), count): result.extend(chunk) if len(result) > limit: raise AssertionError(f'read exceeded fixture limit {limit}') return bytes(result) def check_body(client, pane, count=4096): fid = client.open(f'/pane/{pane["serial"]}/body') try: actual = read_fid(client, fid, BODY_LIMIT, count) finally: client.close(fid) check_bytes(actual, pane['body'], f'pane {pane["serial"]}') def check_screen(data): screen = json.loads(data) cols, rows = screen['cols'], screen['rows'] assert cols > 0 and rows > 0 and cols * rows <= 65536, (cols, rows) assert len(screen['cells']) == cols * rows styles = screen['styles'] assert styles for cell in screen['cells']: assert len(cell) == 2 and isinstance(cell[0], str), cell assert isinstance(cell[1], int) and 0 <= cell[1] < len(styles), cell def session_pid(client): if not hasattr(socket, 'SO_PEERCRED'): return None pid, uid, _ = struct.unpack('3i', client.socket.getsockopt( socket.SOL_SOCKET, socket.SO_PEERCRED, struct.calcsize('3i'))) assert pid > 0 and uid == os.getuid(), (pid, uid) return pid def check_session(client, control, pid): assert session_pid(control) == pid if pid is not None: actual = (session_pid(client) if client.socket.family == socket.AF_UNIX else int(client.read('/os/proc/self/stat').split()[0])) assert actual == pid, (actual, pid) def memory(pid): result = {'rss_bytes': None, 'cumulative_peak_rss_bytes': None} if pid is None: return result try: for line in Path(f'/proc/{pid}/status').read_text().splitlines(): fields = line.split() if fields[0] == 'VmRSS:': result['rss_bytes'] = int(fields[1]) * 1024 elif fields[0] == 'VmHWM:': result['cumulative_peak_rss_bytes'] = int(fields[1]) * 1024 except OSError: pass return result def run(args): rng = random.Random(args.seed) started = time.monotonic() deadline = None batch = 0 operation = 'start' operation_count = 0 pid = None try: with tempfile.TemporaryDirectory(prefix='pardes-soak-') as directory: root = Path(directory) options = ['--9p-tcp=tcp!127.0.0.1!0'] if args.transport == 'tcp' else [] with session(str(args.binary.resolve()), root, 'soak', *options, tty=args.tty) as (control, unix_address): pid = session_pid(control) address = unix_address if args.transport == 'tcp': listeners = [line.split('!') for line in control.read('/listeners').decode().splitlines() if line.startswith('tcp!')] assert len(listeners) == 1, listeners transport, host, port = listeners[0] assert transport == 'tcp' and host == '127.0.0.1' and 0 < int(port) < 65536, listeners address = (host, int(port)) client = Client(address) panes = [] try: check_session(client, control, pid) for slot in range(4): body = payload(rng) serial = 1 if slot == 0 else new_pane(client, body) base = f'/pane/{serial}' if slot == 0: client.write(base + '/body', body, truncate=True) path = root / f'pane-{slot}-0.zig' client.write(base + '/name', f'{path}\n'.encode()) client.write(base + '/exec', b'Save\n') panes.append({'serial': serial, 'body': body, 'path': path, 'variant': 0}) emit(kind='start', seed=args.seed, session_pid=pid, transport=args.transport, host='tty' if args.tty else 'detached', batches=args.batches, duration_seconds=args.duration, interval_seconds=args.interval, body_limit_bytes=BODY_LIMIT, live_panes=len(panes), **memory(pid)) if args.duration is not None: deadline = time.monotonic() + args.duration while (args.batches is None or batch < args.batches) and ( deadline is None or time.monotonic() < deadline): batch_started = time.monotonic() operations = OPERATIONS * 2 rng.shuffle(operations) for operation in operations: slot = rng.randrange(len(panes)) pane = panes[slot] base = f'/pane/{pane["serial"]}' if operation == 'append': addition = payload(rng, min(256, BODY_LIMIT - len(pane['body']))) client.write(base + '/body', addition) pane['body'] += addition elif operation == 'range': body = pane['body'] # Keep range endpoints outside combining and ZWJ clusters. boundaries = [0] + [i for i in range(1, len(body)) if body[i - 1] < 128 and body[i] < 128] + [len(body)] lo, hi = sorted((rng.choice(boundaries), rng.choice(boundaries))) replacement = payload(rng, min(256, BODY_LIMIT - len(body) + hi - lo)) if replacement: client.write(base + '/addr', f'#{lo},#{hi}'.encode()) client.write(base + '/data', replacement) pane['body'] = body[:lo] + replacement + body[hi:] elif operation == 'replace': pane['body'] = payload(rng, BODY_LIMIT) client.write(base + '/body', pane['body'], truncate=True) elif operation == 'read': check_body(client, pane, rng.choice([97, 251, 1024, 4096])) elif operation == 'save': pane['variant'] ^= 1 pane['path'] = root / f'pane-{slot}-{pane["variant"]}.zig' client.write(base + '/name', f'{pane["path"]}\n'.encode()) client.write(base + '/exec', b'Save\n') check_bytes(pane['path'].read_bytes(), pane['body'], 'saved file') check_bytes(client.read('/os' + str(pane['path'])), pane['body'], 'OS mount') elif operation == 'reload': pane['body'] = payload(rng) pane['path'].write_bytes(pane['body']) client.write(base + '/ctl', b'get\n') elif operation == 'cycle': old_serial = pane['serial'] client.remove(base) pane['body'] = payload(rng) pane['serial'] = new_pane(client, pane['body']) assert pane['serial'] != old_serial base = f'/pane/{pane["serial"]}' client.write(base + '/name', f'{pane["path"]}\n'.encode()) client.write(base + '/exec', b'Save\n') elif operation == 'reconnect': client.socket.close() client = Client(address) check_session(client, control, pid) check_body(control, pane) elif operation == 'screen': frozen = client.open('/screen') try: before = read_fid(client, frozen, SCREEN_LIMIT, 4096) check_screen(before) addition = payload(rng, min(64, BODY_LIMIT - len(pane['body']))) client.write(base + '/body', addition) pane['body'] += addition after = read_fid(client, frozen, SCREEN_LIMIT, rng.choice([251, 1024])) check_bytes(after, before, 'frozen screen') finally: client.close(frozen) fresh = client.open('/screen') try: check_screen(read_fid(client, fresh, SCREEN_LIMIT, 4096)) finally: client.close(fresh) operation_count += 1 check_body(client, pane) expected = {str(pane['serial']) for pane in panes} assert set(client.list('/pane')) == expected assert set(control.list('/pane')) == expected for pane in panes: check_body(client, pane) check_body(control, pane) batch += 1 emit(kind='batch', seed=args.seed, batch=batch, duration_seconds=time.monotonic() - batch_started, elapsed_seconds=time.monotonic() - started, operations=len(operations), total_operations=operation_count, model_bytes=sum(len(pane['body']) for pane in panes), **memory(pid)) delay = args.interval - (time.monotonic() - batch_started) if deadline is not None: delay = min(delay, deadline - time.monotonic()) if delay > 0: time.sleep(delay) finally: client.socket.close() emit(kind='summary', seed=args.seed, batches=batch, transport=args.transport, host='tty' if args.tty else 'detached', operations=operation_count, duration_seconds=time.monotonic() - started, status='passed') except BaseException as error: emit(kind='failure', seed=args.seed, batch=batch + 1, operation=operation, transport=args.transport, host='tty' if args.tty else 'detached', operations=operation_count, session_pid=pid, elapsed_seconds=time.monotonic() - started, error=str(error)) raise def main(): parser = argparse.ArgumentParser(description='Deterministic valid-operation editor/9P soak; writes JSONL.') parser.add_argument('binary', type=Path) parser.add_argument('--seed', type=int, default=4200) parser.add_argument('--transport', choices=['unix', 'tcp'], default='unix') parser.add_argument('--tty', action='store_true', help='run an owned TTY host instead of a detached host') limit = parser.add_mutually_exclusive_group() limit.add_argument('--batches', type=int) limit.add_argument('--duration', type=float, help='seconds, stopping after the current batch') parser.add_argument('--interval', type=float, default=0, help='minimum seconds between batch starts') args = parser.parse_args() if args.batches is None and args.duration is None: args.batches = 100 if args.batches is not None and args.batches < 1: parser.error('--batches must be positive') if args.duration is not None and not 0 < args.duration < float('inf'): parser.error('--duration must be finite and positive') if not 0 <= args.interval < float('inf'): parser.error('--interval must be finite and nonnegative') run(args) if __name__ == '__main__': main()