#!/usr/bin/env python3
"""Finite EV-2 protocol reader. Python 3.9+, standard library only; no real I/O.
Run: python foundations-event-io-checker.py [--case all|readiness|cancellation|framing]
--quantum and --attempts change the readiness experiment, not the stated protocol.
Checks deliberately use exceptions, so python -O preserves every check.
"""
import argparse
import itertools
import json
from collections import deque
from dataclasses import dataclass, field

CHECKS = 0

def check(condition, message):
    global CHECKS
    CHECKS += 1
    if not condition:
        raise RuntimeError(message)

def rejects(action, message):
    try:
        action()
    except ValueError:
        check(True, message)
    else:
        check(False, message)

def positive(text):
    try:
        n = int(text)
    except ValueError as exc:
        raise argparse.ArgumentTypeError('expected an integer') from exc
    if n <= 0 or n > 10000:
        raise argparse.ArgumentTypeError('must be between 1 and 10000')
    return n

class Loop:
    def __init__(self, payloads, quantum=2, attempts=1):
        self.sources = {key: bytearray(value) for key, value in payloads.items()}
        self.states = {key: 'WAITING' for key in payloads}
        self.hints = {key: False for key in payloads}
        self.outputs = {key: bytearray() for key in payloads}
        self.queue = deque()
        self.quantum, self.attempts = quantum, attempts
        self.trace = []
        self.calls = self.no_progress = self.notifications = self.harvests = 0
        self.interrupts = {key: 0 for key in payloads}
        self.closed = set()
        self.errors = set()

    def invariant(self):
        check(len(set(self.queue)) == len(self.queue), 'duplicate queue member')
        check(set(self.queue) == {k for k, s in self.states.items() if s == 'QUEUED'}, 'ready membership')
        check(all(self.hints[k] for k in self.queue), 'queued without continuation')

    def enqueue(self, key):
        self.states[key] = 'QUEUED'
        self.queue.append(key)

    def notify(self, key):
        self.notifications += 1
        if key not in self.states or self.states[key] in ('FINISHED', 'FAILED'):
            return
        self.hints[key] = True
        if self.states[key] == 'WAITING':
            self.enqueue(key)
        self.invariant()

    def turn(self, notifications=(), broken_drop=False):
        # One bounded event-harvest batch even while queue is nonempty.
        self.harvests += 1
        for key in notifications:
            self.notify(key)
        if not self.queue:
            return False
        key = self.queue.popleft()
        self.states[key] = 'RUNNING'
        budget = self.quantum
        events = []
        for _ in range(self.attempts):
            if budget == 0:
                break
            self.calls += 1
            if key in self.errors:
                result = 'EIO'
                self.no_progress += 1
                self.hints[key] = False
                self.states[key] = 'FAILED'
            elif self.interrupts[key]:
                self.interrupts[key] -= 1
                result = 'EINTR'
                self.no_progress += 1
            elif self.sources[key]:
                n = min(budget, len(self.sources[key]))
                data = bytes(self.sources[key][:n])
                del self.sources[key][:n]
                self.outputs[key].extend(data)
                budget -= n
                result = data.decode('ascii')
            else:
                self.no_progress += 1
                result = 'EOF' if key in self.closed else 'EAGAIN'
                self.hints[key] = False
                self.states[key] = 'FINISHED' if result == 'EOF' else 'WAITING'
            events.append(result)
            if self.states[key] != 'RUNNING':
                break
        if self.states[key] == 'RUNNING':
            if broken_drop:
                self.hints[key] = False
                self.states[key] = 'WAITING'
            else:
                self.enqueue(key)
        self.invariant()
        self.trace.append({'connection': key, 'returns': events, 'queue': list(self.queue)})
        return True

    def drain(self):
        # For these finite fixtures each payload byte and each injected EINTR
        # consumes progress; an empty source incurs at most one final EAGAIN.
        bound = sum(map(len, self.sources.values())) + sum(self.interrupts.values()) + len(self.sources) + 1
        for _ in range(bound):
            if not self.turn():
                return
        raise RuntimeError('finite fixture failed to reach idle')

def readiness(q, attempts):
    loop = Loop({'A': b'abcdef', 'B': b'XY'}, q, attempts)
    loop.notify('A'); loop.notify('B'); loop.drain()
    check(bytes(loop.outputs['A']) == b'abcdef' and bytes(loop.outputs['B']) == b'XY', 'lost bytes')
    if q == 2 and attempts == 1:
        check([(r['connection'], r['returns']) for r in loop.trace] == [
            ('A', ['ab']), ('B', ['XY']), ('A', ['cd']), ('B', ['EAGAIN']), ('A', ['ef']), ('A', ['EAGAIN'])], 'default trace')
    dup = Loop({'A': b'x'})
    dup.notify('A'); dup.notify('A')
    check(list(dup.queue) == ['A'], 'duplicate notification was not merged')
    broken = Loop({'A': b'abcdef', 'B': b'XY'})
    broken.turn(['A', 'B'], broken_drop=True)
    broken.turn(broken_drop=True)
    check(not broken.queue and bytes(broken.sources['A']) == b'cdef', 'missing ET counterexample')
    # A possible LT stream re-notifies still-readable A after every quantum.
    lt = Loop({'A': b'abcdef', 'B': b'XY'})
    lt.turn(['A', 'B'], broken_drop=True)
    lt.turn(['A'], broken_drop=True)
    lt.turn(['A'], broken_drop=True)
    lt.turn(['A'], broken_drop=True)
    check(bytes(lt.outputs['A']) == b'abcdef', 'LT fixture failed')
    arrival = Loop({'A': b'abcdefghij', 'B': b'XY'})
    arrival.turn(['A']); arrival.turn(['B']); arrival.turn()
    check(arrival.trace[2]['connection'] == 'B', 'new notification starved behind busy queue')
    interrupted = Loop({'A': b'z'}, attempts=2)
    interrupted.interrupts['A'] = 3
    interrupted.turn(['A'])
    check(interrupted.trace[0]['returns'] == ['EINTR', 'EINTR'] and list(interrupted.queue) == ['A'], 'attempt budget lost')
    interrupted.drain()
    eof = Loop({'A': b''}); eof.closed.add('A'); eof.turn(['A'])
    check(eof.states['A'] == 'FINISHED', 'EOF not terminal')
    fatal = Loop({'A': b'x'}); fatal.errors.add('A'); fatal.turn(['A'])
    fatal.notify('A')
    check(fatal.states['A'] == 'FAILED' and not fatal.queue and not fatal.hints['A'], 'fatal read resurrected')
    stale = Loop({'slot3-generation8': b'n'})
    stale.notify('slot3-generation7')
    check(not stale.queue and stale.sources['slot3-generation8'] == b'n', 'stale fd event touched new generation')
    return {'trace': loop.trace, 'bytes': sum(map(len, loop.outputs.values())), 'read_calls': loop.calls,
            'no_progress_calls': loop.no_progress, 'initial_notifications': loop.notifications,
            'registrations': 2, 'nonblocking_harvests_including_idle_check': loop.harvests,
            'broken_ET_remaining_A': broken.sources['A'].decode(), 'LT_fixture_A': lt.outputs['A'].decode(),
            'new_arrival_trace': arrival.trace, 'bounded_EINTR_first_turn': interrupted.trace[0]}

@dataclass
class Request:
    generation: int = 7
    logical: str = 'WAITING'
    target: str = 'INFLIGHT'
    cancel: str = 'NONE'
    pin: int = 1
    holders: set = field(default_factory=set)
    buffer_free: bool = False
    retired: bool = False
    frees: int = 0
    pin_releases: int = 0
    trace: list = field(default_factory=list)

    def finish_event(self, event):
        if self.target != 'INFLIGHT' and not self.pin and not self.holders and not self.buffer_free:
            self.buffer_free = True
            self.frees += 1
        if self.buffer_free and self.cancel != 'PENDING':
            self.retired = True
        check((self.target == 'INFLIGHT') == bool(self.pin), 'target/pin conservation')
        check(not self.buffer_free or (not self.pin and not self.holders and self.target != 'INFLIGHT'), 'premature buffer reuse')
        check(self.frees <= 1 and self.pin_releases <= 1, 'double release')
        check(not self.retired or self.cancel != 'PENDING', 'identity retired before cancel CQE')
        self.trace.append({'event': event, 'logical': self.logical, 'target': self.target, 'cancel': self.cancel,
                           'pin': self.pin, 'borrowers': sorted(self.holders), 'buffer_free': self.buffer_free, 'retired': self.retired})

    def deadline(self, accepted=True):
        if self.logical == 'WAITING':
            self.logical = 'TIMEOUT'
            self.cancel = 'PENDING' if accepted else 'DONE:submit-error'
        self.finish_event('deadline')

    def complete_target(self, result='ok', borrower=None, before_deadline=False):
        if self.target != 'INFLIGHT':
            raise ValueError('duplicate target terminal')
        if borrower is not None:
            if borrower in self.holders:
                raise ValueError('duplicate borrower')
            self.holders.add(borrower)  # Acquire before publication or d release.
        self.target = 'TERMINAL:' + result
        self.pin = 0
        self.pin_releases += 1
        if self.logical == 'WAITING':
            self.logical = 'RESULT' if before_deadline else 'TIMEOUT'
        self.finish_event('target:' + result)

    def complete_cancel(self, result):
        if self.cancel != 'PENDING':
            raise ValueError('unexpected cancel terminal')
        self.cancel = 'DONE:' + result
        self.finish_event('cancel:' + result)

    def release(self, token):
        if token not in self.holders:
            raise ValueError('unknown or repeated borrower release')
        self.holders.remove(token)
        self.finish_event('release:' + token)

    def transfer(self, old, new):
        if old not in self.holders or new in self.holders or self.buffer_free:
            raise ValueError('invalid transfer')
        self.holders.add(new)
        self.holders.remove(old)
        self.finish_event('transfer:' + old + '->' + new)

    def matched(self, kind, generation):
        # Stable-registry lookup precedes any access via an event-carried address.
        return not self.retired and kind in ('target', 'cancel') and generation == self.generation

def cancellation():
    r = Request(); r.finish_event('submit'); r.deadline()
    r.complete_target(borrower='callback'); r.release('callback'); r.complete_cancel('not-found')
    check(r.retired and r.logical == 'TIMEOUT' and r.frees == 1, 'main drain')
    six = []
    for order in itertools.permutations('CTL'):
        # Address-retention lease only: no payload access while device writes.
        x = Request(holders={'reserve'}); x.deadline()
        for event in order:
            {'C': lambda: x.complete_cancel('success'), 'T': lambda: x.complete_target('canceled'),
             'L': lambda: x.release('reserve')}[event]()
        check(x.retired and x.frees == 1 and x.pin_releases == 1, 'six-order retirement')
        six.append({'order': ''.join(order), 'trace': x.trace})
    three = []
    for order in itertools.permutations('CTL'):
        if order.index('L') < order.index('T'):
            continue
        x = Request(); x.deadline()
        for event in order:
            {'C': lambda: x.complete_cancel('success'), 'T': lambda: x.complete_target('canceled', 'callback'),
             'L': lambda: x.release('callback')}[event]()
        check(x.retired, 'callback order retirement')
        three.append(''.join(order))
    variants = {}
    x = Request(); x.complete_target(before_deadline=True); x.deadline()
    check(x.logical == 'RESULT' and x.cancel == 'NONE' and x.retired, 'target wins')
    variants['target_before_deadline'] = x.trace
    for code in ('already-running', 'not-found-wrong-id', 'invalid'):
        x = Request(); x.deadline(); x.complete_cancel(code)
        check(not x.buffer_free and x.pin == 1, 'cancel result freed target')
        x.complete_target('io-error')
        check(x.retired, 'eventual target error drain')
        variants[code] = x.trace
    x = Request(); x.deadline(accepted=False); x.complete_target('ok')
    check(x.retired, 'cancel submission error drain')
    variants['cancel_submit_error'] = x.trace
    x = Request(); x.complete_target(borrower='callback', before_deadline=True)
    try:
        raise LookupError('deliberate callback failure')
    except LookupError:
        pass
    finally:
        x.release('callback')
    check(x.retired, 'exception leaked lease')
    rejects(lambda: x.release('callback'), 'duplicate release rejected')
    rejects(lambda: x.complete_target(), 'duplicate completion rejected')
    variants['callback_error'] = x.trace
    x = Request(); x.deadline(); x.complete_target(borrower='callback')
    x.transfer('callback', 'background'); x.complete_cancel('not-found')
    check(not x.buffer_free and not x.retired, 'background not charged')
    check(not x.matched('target', 8), 'wrong generation matched')
    x.release('background'); check(x.retired, 'background failed to drain')
    check(not x.matched('cancel', 7), 'retired registry accepted stale event')
    variants['background_transfer'] = x.trace
    stuck = Request(); stuck.deadline(); stuck.complete_cancel('already-running')
    check(not stuck.retired and stuck.pin == 1, 'missing target must remain retained')
    # Q=1 includes target-complete borrowed buffer; target count alone is wrong.
    borrowed = Request(); borrowed.complete_target(borrower='background', before_deadline=True)
    admitted = [borrowed]
    check(sum(not a.retired for a in admitted) == 1 and all(a.target != 'INFLIGHT' for a in admitted), 'admission counterexample')
    orphan = Request(); orphan.deadline()  # target CQE exists externally, but no drain consumes it
    check(orphan.pin == 1 and not orphan.retired, 'orphan drain counterexample')
    return {'main_trace': r.trace, 'reserve_lease_six_orders': six, 'read_callback_three_orders': three,
            'variants': variants, 'external_failure_retained': stuck.trace[-1],
            'no_drain_retained': orphan.trace[-1], 'Q1_borrowed_request_still_occupies_slot': True}

class Framer:
    def __init__(self, capacity=9, max_payload=3):
        self.capacity, self.maximum = capacity, max_payload + 2
        self.raw = bytearray()  # Models source bytes, outside the application budget.
        self.part = bytearray(); self.expected = None; self.reserved = 0
        self.output = []; self.state = 'WAITING'; self.hint = False
        self.trace = []; self.peak = 0

    def used(self):
        return sum(len(x) + 2 for x in self.output) + self.reserved

    def record(self, event):
        self.peak = max(self.peak, self.used())
        check(self.used() <= self.capacity, 'capacity exceeded')
        check(len(self.part) <= self.reserved, 'read without reservation')
        self.trace.append({'event': event, 'state': self.state, 'hint': self.hint,
                           'partial_hex': self.part.hex(), 'output': [x.decode() for x in self.output],
                           'used_or_reserved': self.used(), 'source_remaining_hex': self.raw.hex()})

    def arrive(self, data):
        if self.state in ('FINISHED', 'FAILED'):
            raise ValueError('input after terminal state')
        self.raw.extend(data); self.hint = True
        if self.state == 'WAITING': self.state = 'QUEUED'
        self.record('arrival:' + data.hex())

    def turn(self, attempts=2):
        if self.state != 'QUEUED': return
        frames = 0
        for _ in range(attempts):
            if frames == 1: break
            if not self.reserved:
                if self.used() + self.maximum > self.capacity:
                    self.state = 'PAUSED'; self.record('capacity'); return
                self.reserved = self.maximum
            if not self.raw:
                self.hint = False; self.state = 'WAITING'; self.record('EAGAIN'); return
            need = (2 if self.expected is None else self.expected) - len(self.part)
            n = min(need, len(self.raw))
            self.part.extend(self.raw[:n]); del self.raw[:n]
            if len(self.part) == 2 and self.expected is None:
                size = int.from_bytes(self.part, 'big')
                if size + 2 > self.maximum:
                    self.fail('oversize frame')
                self.expected = size + 2; self.reserved = self.expected
            if self.expected is not None and len(self.part) == self.expected:
                self.output.append(bytes(self.part[2:])); self.part.clear()
                self.expected = None; self.reserved = 0; frames += 1
            self.record('read:' + str(n))
        # Capacity suspension can happen at the quantum boundary too.
        if not self.reserved and self.used() + self.maximum > self.capacity:
            self.state = 'PAUSED'
        self.record('yield')

    def consume(self):
        item = self.output.pop(0)
        if self.state == 'PAUSED': self.state = 'QUEUED' if self.hint else 'WAITING'
        self.record('consume:' + item.decode())
        return item

    def fail(self, reason):
        self.state = 'FAILED'; self.hint = False
        self.part.clear(); self.expected = None; self.reserved = 0
        self.record('failure:' + reason)
        raise ValueError(reason)

    def eof(self):
        if self.state in ('FINISHED', 'FAILED'):
            raise ValueError('EOF after terminal state')
        if self.raw:
            raise ValueError('EOF before source drained')
        if self.part:
            self.fail('truncated frame')
        self.state = 'FINISHED'; self.hint = False; self.reserved = 0
        self.record('EOF')

def framing():
    f = Framer()
    for chunk in (b'\x00', b'\x03CA', b'T\x00\x02OK'):
        f.arrive(chunk)
        for _ in range(10):
            if f.state != 'QUEUED': break
            f.turn()
    check(f.output == [b'CAT'] and f.state == 'PAUSED' and f.hint and f.raw == b'\x00\x02OK', 'frame continuation')
    results = [f.consume()]
    check(f.state == 'QUEUED', 'credit did not requeue')
    for _ in range(10):
        if f.state != 'QUEUED': break
        f.turn()
    results.append(f.consume())
    check(results == [b'CAT', b'OK'], 'frame outputs')
    f.eof(); check(f.used() == 0 and not f.hint, 'frame final drain')
    oversize = Framer(); oversize.arrive(b'\x00\x04')
    rejects(oversize.turn, 'oversize rejected')
    check(oversize.state == 'FAILED' and not oversize.hint and oversize.used() == 0, 'oversize not terminal')
    rejects(lambda: oversize.arrive(b'x'), 'oversize resumed')
    oversize.turn(); check(not oversize.output, 'failed parser emitted')
    partial = Framer(); partial.arrive(b'\x00'); partial.turn()
    rejects(partial.eof, 'partial EOF rejected')
    check(partial.state == 'FAILED' and not partial.hint, 'truncation not terminal')
    rejects(lambda: partial.arrive(b'\x00'), 'truncated parser resumed')
    partial.turn(); check(not partial.output, 'truncated parser emitted')
    premature = Framer(); premature.arrive(b'\x00\x00')
    rejects(premature.eof, 'EOF before source drained rejected')
    empty = Framer(); empty.arrive(b'\x00\x00'); empty.turn()
    check(empty.output == [b''], 'zero frame')
    return {'trace': f.trace, 'decoded': [x.decode() for x in results], 'peak_used_or_reserved': f.peak,
            'final_used_or_reserved': f.used(), 'capacity': f.capacity}

def main():
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument('--case', choices=('all', 'readiness', 'cancellation', 'framing'), default='all')
    parser.add_argument('--quantum', type=positive, default=2)
    parser.add_argument('--attempts', type=positive, default=1)
    args = parser.parse_args()
    result = {'model': 'EV-2', 'scope': 'finite protocol model, not epoll/io_uring conformance'}
    if args.case in ('all', 'readiness'): result['readiness'] = readiness(args.quantum, args.attempts)
    if args.case in ('all', 'cancellation'): result['cancellation'] = cancellation()
    if args.case in ('all', 'framing'): result['framing'] = framing()
    result['checks'] = CHECKS
    result['status'] = 'passed'
    print(json.dumps(result, ensure_ascii=False, indent=2))

if __name__ == '__main__':
    main()
