#!/usr/bin/env python3
"""Three explicitly bounded teaching interfaces; standard library, stdout only.
Not a full NTP implementation, replicated lease service, or recovery protocol.
All time arithmetic is exact Fraction. HLC uses integer timestamp components.
"""
from fractions import Fraction as F
from itertools import product
import json
import random


def require(condition, message):
    if not condition:
        raise AssertionError(message)


def rational(x):
    if type(x) not in (int, F):
        raise ValueError('exact int or Fraction required')
    return F(x)


def estimate(t1, t2, t3, t4):
    t1, t2, t3, t4 = map(rational, (t1, t2, t3, t4))
    if t4 < t1 or t3 < t2:
        raise ValueError('time order violates constant-offset unit-rate model')
    low, high = t3-t4, t2-t1
    if low > high:
        raise ValueError('negative apparent propagation in this model')
    return dict(stamps=(t1,t2,t3,t4), processing=t3-t2, propagation=high-low,
                interval=(low,high), midpoint=(low+high)/2, radius=(high-low)/2)


def delay_witness(record, offset):
    offset = rational(offset)
    low, high = record['interval']
    if not low <= offset <= high:
        raise ValueError('offset outside feasible interval')
    return dict(offset=offset, forward=high-offset, backward=offset-low,
                processing=record['processing'])


def intersect_intervals(intervals):
    low = high = None
    for a,b in intervals:
        a,b = rational(a),rational(b)
        if a>b:
            raise ValueError('reversed interval')
        low = a if low is None else max(low,a)
        high = b if high is None else min(high,b)
    if low is None:
        return dict(status='NO_OBSERVATION', interval=None)
    return dict(status='COMPATIBLE' if low<=high else 'INCOMPATIBLE',
                interval=(low,high) if low<=high else None,
                lower_constraint=low,upper_constraint=high)


def expand_interval(interval, relative_rate, elapsed_real):
    a,b=map(rational,interval);rate=rational(relative_rate);elapsed=rational(elapsed_real)
    if a>b or rate<0 or elapsed<0:
        raise ValueError('invalid uncertainty propagation')
    return a-rate*elapsed,b+rate*elapsed


def lease_duration(length, rho_client, rho_server):
    length,rho_client,rho_server=map(rational,(length,rho_client,rho_server))
    if length<=0 or not 0<=rho_client<1 or not 0<=rho_server<1:
        raise ValueError('positive length, drift bounds in [0,1) required')
    return length*(1-rho_client)/(1+rho_server)


class LeaseServer:
    """Honest serial grantor, no crash or early revoke. Clock rates are premises.
    Resource lookup is expected O(1) using a dict. Grant counter is unbounded.
    """
    def __init__(self):
        self.resources={}
        self.next_epoch=0

    def grant(self, resource, client, request_id, length, server_now):
        length,server_now=map(rational,(length,server_now))
        if length<=0 or type(request_id) is not int or request_id<1:
            raise ValueError('invalid grant request')
        old=self.resources.get(resource)
        if old is not None and server_now<old['server_expiry']:
            return None
        self.next_epoch+=1
        grant=dict(resource=resource,client=client,request_id=request_id,length=length,
                   epoch=self.next_epoch,server_expiry=server_now+length)
        self.resources[resource]=grant
        return grant.copy()


class LeaseClient:
    """One resource, one outstanding request; no renewal/early release.
    Caller supplies a clock satisfying the declared rate bound. Identity is
    a per-client monotonically increasing request counter, never reused.
    """
    def __init__(self, resource, client, rho_client, rho_server):
        lease_duration(1,rho_client,rho_server)
        self.resource,self.client=resource,client
        self.rho_client,self.rho_server=rational(rho_client),rational(rho_server)
        self.serial=0;self.phase='IDLE';self.deadline=None;self.length=None;self.epoch=None

    def advance(self, now):
        now=rational(now)
        if self.phase!='IDLE' and now>=self.deadline:
            self.phase='IDLE';self.epoch=None
        return now

    def begin(self, length, now):
        now=self.advance(now)
        if self.phase!='IDLE':
            raise ValueError('only one outstanding/active request; no renewal')
        duration=lease_duration(length,self.rho_client,self.rho_server)
        self.serial+=1;self.length=rational(length);self.deadline=now+duration
        self.phase='WAITING';self.epoch=None
        return dict(resource=self.resource,client=self.client,request_id=self.serial,length=self.length)

    def receive(self, grant, now):
        now=self.advance(now)
        if self.phase!='WAITING' or grant is None:
            return False
        expected=(self.resource,self.client,self.serial,self.length)
        actual=tuple(grant.get(k) for k in ('resource','client','request_id','length'))
        if actual!=expected:
            return False
        # advance already retired attempts at or after the original deadline.
        self.phase='ACTIVE';self.epoch=grant['epoch']
        return True

    def valid(self, now):
        self.advance(now)
        return self.phase=='ACTIVE'


class HLC:
    def __init__(self, max_counter=None):
        if max_counter is not None and (type(max_counter) is not int or max_counter<0):
            raise ValueError('nonnegative integer counter bound required')
        self.l=0;self.c=0;self.max_counter=max_counter

    @staticmethod
    def stamp(value):
        if not isinstance(value,(tuple,list)) or len(value)!=2 or any(type(v) is not int or v<0 for v in value):
            raise ValueError('nonnegative integer pair required')
        return tuple(value)

    def event(self, physical, message=None):
        if type(physical) is not int or physical<0:
            raise ValueError('nonnegative integer physical reading required')
        old=(self.l,self.c)
        if message is None:
            new_l=max(self.l,physical)
            new_c=self.c+1 if new_l==self.l else 0
            branch='LOCAL_SAME' if new_l==self.l else 'LOCAL_PHYSICAL'
        else:
            ml,mc=self.stamp(message)
            new_l=max(self.l,ml,physical)
            if new_l==self.l==ml:
                new_c=max(self.c,mc)+1;branch='BOTH'
            elif new_l==self.l:
                new_c=self.c+1;branch='LOCAL'
            elif new_l==ml:
                new_c=mc+1;branch='MESSAGE'
            else:
                new_c=0;branch='PHYSICAL'
        if self.max_counter is not None and new_c>self.max_counter:
            raise OverflowError('counter capacity exhausted; no state published')
        self.l,self.c=new_l,new_c
        return dict(old=old,physical=physical,message=message,branch=branch,new=(new_l,new_c))


class PiecewiseClock:
    """Test oracle only: rates before and after a real-time switch, then affine.
    Real time is available to this offline verifier, never to the lease client.
    """
    def __init__(self, offset, before, after, switch=2):
        self.offset,self.before,self.after,self.switch=map(rational,(offset,before,after,switch))
        if self.before<=0 or self.after<=0 or self.switch<0:
            raise ValueError('test clock must be continuous and strictly advancing')

    def read(self,t):
        t=rational(t)
        if t<0:raise ValueError('nonnegative real time required')
        return self.offset+self.before*min(t,self.switch)+self.after*max(F(0),t-self.switch)

    def hit(self,reading):
        reading=rational(reading)
        if reading<self.offset:raise ValueError('threshold before clock origin')
        at_switch=self.read(self.switch)
        if reading<=at_switch:return (reading-self.offset)/self.before
        return self.switch+(reading-at_switch)/self.after


def main():
    counts={}
    probes=[estimate(*x) for x in [(100,134,136,112),(200,232,233,205),(300,331,332,304),(400,441,442,405)]]
    require([x['interval'] for x in probes]==[(24,34),(28,32),(28,31),(37,41)],'four probe intervals')
    witnesses=[delay_witness(probes[0],v) for v in (24,30,34)]
    trials=offset_witnesses=0
    for offset,forward,backward,processing in product(range(-7,8),range(9),range(9),range(5)):
        t1=100;t2=t1+forward+offset;t3=t2+processing;t4=t1+forward+processing+backward
        r=estimate(t1,t2,t3,t4)
        require(r['propagation']==forward+backward and r['processing']==processing,'elapsed decomposition')
        require(r['interval'][0]<=offset<=r['interval'][1],'true offset feasible')
        require(r['midpoint']-offset==F(forward-backward,2),'asymmetry bias')
        for theta in range(int(r['interval'][0]),int(r['interval'][1])+1):
            w=delay_witness(r,theta)
            require(w['forward']>=0 and w['backward']>=0,'nonnegative witness')
            require(theta+w['forward']==t2-t1 and theta-w['backward']==t3-t4,'reconstructed observations')
            offset_witnesses+=1
        trials+=1
    counts.update(probe_worlds=trials,feasible_integer_witnesses=offset_witnesses)
    require(intersect_intervals(p['interval'] for p in probes[:3])['interval']==(28,31),'three-way intersection')
    require(intersect_intervals(p['interval'] for p in probes)['status']=='INCOMPATIBLE','empty intersection')
    require(expand_interval((28,32),F(1,10),10)==(27,33),'uncertainty aging')
    require(estimate(100,135,137,112)['midpoint']==30,'symmetric migration')
    # Exhaustive two-segment rate corners and interior rates, distinct client/server bounds.
    clock_cases=0;equalities=0
    for rc,rs in product((F(0),F(1,10),F(1,3)),repeat=2):
        for cr,sr in product(product((1-rc,F(1),1+rc),repeat=2),product((1-rs,F(1),1+rs),repeat=2)):
            C=PiecewiseClock(100,*cr);S=PiecewiseClock(700,*sr)
            for t0,gap,L in product((F(0),F(3)),(F(0),F(1,2),F(4)),(F(1),F(11))):
                g=t0+gap;duration=lease_duration(L,rc,rs)
                client_stop=C.hit(C.read(t0)+duration);server_stop=S.hit(S.read(g)+L)
                require(client_stop<=t0+L/(1+rs)<=g+L/(1+rs)<=server_stop,'drift inequalities')
                require(client_stop<=server_stop,'no overlapping valid lease interval')
                clock_cases+=1;equalities+=client_stop==server_stop
    counts.update(piecewise_rate_schedules=clock_cases,tight_deadline_equalities=equalities)
    C=PiecewiseClock(100,F(9,10),F(9,10));S=PiecewiseClock(700,F(11,10),F(11,10))
    server=LeaseServer();client=LeaseClient('R','A',F(1,10),F(1,10))
    req=client.begin(11,C.read(0));require(not client.valid(C.read(1)),'request not yet granted')
    grant=server.grant(**req,server_now=S.read(2));require(grant is not None,'first grant')
    require(client.receive(grant,C.read(4)),'timely matching reply')
    require(not client.receive(grant,C.read(4)),'duplicate does not start new duration')
    require(client.valid(C.read(9)) and not client.valid(C.read(10)),'strict deadline')
    require(server.grant('R','B',1,11,S.read(11)) is None,'no early server regrant')
    second=server.grant('R','B',2,11,S.read(12));require(second is not None,'server equality may regrant')
    late=LeaseClient('R','A',F(1,10),F(1,10));late.begin(11,C.read(0))
    require(not late.receive(grant,C.read(11)) and not late.valid(C.read(11)),'late reply retired')
    late.begin(11,C.read(13));require(not late.receive(grant,C.read(13)),'old nonce cannot satisfy fresh attempt')
    stop=C.hit(109);server_stop=S.hit(grant['server_expiry'])
    wrong_reply=C.hit(C.read(4)+9);wrong_unscaled=C.hit(C.read(0)+11)
    require((stop,server_stop,wrong_reply,wrong_unscaled)==(10,12,14,F(110,9)),'deadline example')
    # Test the existing resource fencing interface, not a replicated grant protocol.
    resource={'highest':40,'value':0}
    def write(epoch,value):
        if epoch<resource['highest']:return False
        resource['highest']=epoch;resource['value']=value;return True
    resource['highest']=42;require(write(42,9),'new holder')
    require(not write(41,8) and resource['value']==9,'late stale effect rejected')
    without_fencing=8
    lease_example=dict(client_deadline=F(109),server_deadline=grant['server_expiry'],
        client_real_stop=stop,server_real_stop=server_stop,wrong_reply_real_stop=wrong_reply,
        wrong_unscaled_real_stop=wrong_unscaled,late_reply_rejected=True,
        restored_clock_at_13=C.read(13),stopped_clock_example=C.read(9),
        new_epoch=second['epoch'],fenced_value=resource['value'],unfenced_value=without_fencing)
    # All receive branch combinations; independent lower-bound formulation.
    branch_counts={}
    for old_l,old_c,p,ml,mc in product(range(5),range(4),range(5),range(5),range(4)):
        h=HLC();h.l,h.c=old_l,old_c;r=h.event(p,(ml,mc));new_l=max(old_l,p,ml)
        constraints=[0]
        for l,c in ((old_l,old_c),(ml,mc)):
            if l==new_l:constraints.append(c+1)
        wanted=(new_l,max(constraints));require(r['new']==wanted,'minimal lexicographic upper bound')
        require(wanted>(old_l,old_c) and wanted>(ml,mc),'strict event and message extension')
        branch_counts[r['branch']]=branch_counts.get(r['branch'],0)+1
    counts['receive_rule_combinations']=sum(branch_counts.values())
    # Actual message DAG oracle tracks independent ancestry and maximum observed physical value.
    rng=random.Random(2026100911);events_checked=causal_pairs=bounded_events=0
    for bounded in (False,True):
        for trial in range(240):
            clocks=[HLC() for _ in range(4)];last=[None]*4;messages=[];rows=[];ancestors=[];epsilon=7
            for t in range(100):
                receiver=rng.randrange(4);available=[x for x in messages if x['to']==receiver]
                msg=rng.choice(available) if available and rng.randrange(3)==0 else None
                old=last[receiver];anc=set()
                for previous in (old,None if msg is None else msg['event']):
                    if previous is not None:anc|=ancestors[previous]|{previous}
                p=max(0,t+rng.randrange(-epsilon,epsilon+1)) if bounded else rng.randrange(180)
                r=clocks[receiver].event(p,None if msg is None else msg['stamp'])
                maximum=max([0,p]+[rows[j]['physical'] for j in anc])
                require(r['new'][0]==maximum,'ancestral physical maximum')
                for j in anc:require(rows[j]['new']<r['new'],'full causal closure');causal_pairs+=1
                if bounded:
                    require(p<=r['new'][0]<=t+epsilon and r['new'][0]-p<=2*epsilon,'absolute-time distance')
                    bounded_events+=1
                rows.append(r);ancestors.append(anc);last[receiver]=t
                if msg is None and rng.randrange(2)==0:
                    targets=[j for j in range(4) if j!=receiver]
                    messages.append(dict(to=rng.choice(targets),event=t,stamp=r['new']))
                events_checked+=1
    counts.update(dag_events=events_checked,causal_pair_checks=causal_pairs,absolute_error_events=bounded_events)
    clocks={x:HLC() for x in 'ABC'};sent={};trace=[]
    operations=[('a','A',100,'send','m1'),('b','B',90,'send','m0'),('c','B',91,'receive','m1'),
       ('d','B',90,'send','m2'),('e','A',99,'receive','m2'),('f','C',120,'send','m3'),
       ('g','A',101,'receive','m3'),('h','A',130,'receive','m0'),('i','A',129,'receive','m2')]
    for id,who,p,action,msg in operations:
        r=clocks[who].event(p,sent[msg] if action=='receive' else None)
        if action=='send':sent[msg]=r['new']
        trace.append(dict(id=id,process=who,action=action,message_id=msg,**r))
    require([x['new'] for x in trace]==[(100,0),(90,0),(100,1),(100,2),(100,3),(120,0),(120,1),(130,0),(130,1)],'nine-event trace')
    trace_edges=set();previous_by_process={};send_by_message={}
    for row in trace:
        id,who=row['id'],row['process']
        if who in previous_by_process:trace_edges.add((previous_by_process[who],id))
        previous_by_process[who]=id
        if row['action']=='send':send_by_message[row['message_id']]=id
        else:trace_edges.add((send_by_message[row['message_id']],id))
    closure=set(trace_edges)
    for pivot in [row['id'] for row in trace]:
        closure|={(a,b) for a,k in list(closure) for j,b in list(closure) if k==j==pivot}
    by_id={row['id']:row for row in trace}
    require(all(by_id[a]['new']<by_id[b]['new'] for a,b in closure),'fixed trace causal closure')
    require(('d','f') not in closure and ('f','d') not in closure and by_id['d']['new']<by_id['f']['new'],'concurrent but ordered')
    limited=HLC(3)
    for _ in range(4):limited.event(100)
    try:limited.event(100)
    except OverflowError:pass
    else:raise AssertionError('counter overflow silently accepted')
    require((limited.l,limited.c)==(100,3),'overflow leaves state unchanged')
    equal=(HLC().event(100)['new'],HLC().event(100)['new']);require(equal[0]==equal[1],'concurrent equal stamps')
    invalid=[lambda:estimate(0,3,2,5),lambda:estimate(5,3,4,2),lambda:estimate(0,0,5,1),
        lambda:estimate(0.1,2,3,4),lambda:lease_duration(0,0,0),lambda:lease_duration(1,1,0),
        lambda:HLC().event(True),lambda:HLC().event(1,(-1,0)),lambda:HLC(-1)]
    for test in invalid:
        try:test()
        except ValueError:pass
        else:raise AssertionError('invalid input accepted')
    counts['invalid_inputs']=len(invalid)
    result=dict(status='PASS',checks=counts,probes=probes,first_probe_witnesses=witnesses,
       compatible=intersect_intervals(p['interval'] for p in probes[:3]),
       incompatible=intersect_intervals(p['interval'] for p in probes),lease=lease_example,
       hlc_trace=trace,hlc_causal_edges=sorted(trace_edges),hlc_causal_pairs=sorted(closure),
       concurrent_pair=dict(events=('d','f'),timestamp_order=True,causal_order=False),
       receive_branch_counts=branch_counts,concurrent_equal_stamps=equal,
       overflow=dict(rejected=True,retained=(limited.l,limited.c)),
       scope='Finite exact arithmetic tests; static assumptions and proofs remain required. No network or crash recovery implementation.')
    def encode(x):
        if isinstance(x,F):return str(x)
        raise TypeError(type(x).__name__)
    print(json.dumps(result,ensure_ascii=False,indent=2,default=encode))


if __name__=='__main__':
    main()
