"""Offline rpi/Pi runtime comparison. No real model credentials or cloud requests.
Fixtures carry identical logical text but each runtime's native JSONL format.
Local deterministic Anthropic endpoint measures overhead, NOT coding quality.
"""
import argparse, copy, hashlib, http.server, json, os, pathlib, platform, queue
import shutil, statistics, subprocess, tempfile, threading, time, uuid

P = argparse.ArgumentParser()
P.add_argument('--rpi', required=True)
P.add_argument('--pi', required=True, help='native Pi JS entrypoint')
P.add_argument('--node', default='node')
P.add_argument('--output', required=True)
P.add_argument('--runs', type=int, default=5)
P.add_argument('--warmup', type=int, default=2)
P.add_argument('--sizes', default='100,1000,5000')
P.add_argument('--turns', default='20,100')
A = P.parse_args()
ROOT = pathlib.Path(tempfile.mkdtemp(prefix='rpi-pi-runtime-'))
REQUESTS = queue.Queue()
FILE_TEXT = 'deterministic tool result\n' * 8
TASK_TURNS = 0

class Handler(http.server.BaseHTTPRequestHandler):
    def log_message(self, *args): pass
    def do_POST(self):
        received = time.perf_counter()
        req = json.loads(self.rfile.read(int(self.headers['Content-Length'])))
        messages = req.get('messages', [])
        results = sum(1 for m in messages for c in m.get('content', []) if isinstance(c, dict) and c.get('type') == 'tool_result')
        turns = TASK_TURNS
        text = json.dumps(messages, ensure_ascii=False)
        tool_results=[c for m in messages for c in m.get('content', []) if isinstance(c,dict) and c.get('type')=='tool_result']
        tool_ok=all(not c.get('is_error',False) and 'deterministic tool result' in json.dumps(c) for c in tool_results)
        REQUESTS.put({'received':received,'messages':len(messages),'results':results,'has_tail':'BENCH_TAIL' in text,'body_bytes':len(json.dumps(req).encode()),'tool_ok':tool_ok})
        content = {'type':'tool_use','id':f'call_{results}','name':'read','input':{'path':str(ROOT/'work'/'fixture.txt')}} if results < turns else {'type':'text','text':'BENCH_DONE'}
        stop = 'tool_use' if results < turns else 'end_turn'
        msg = {'id':'msg_bench','type':'message','role':'assistant','model':'bench-model','content':[content], 'stop_reason':stop,'stop_sequence':None,'usage':{'input_tokens':1,'output_tokens':1}}
        if not req.get('stream'):
            payload = json.dumps(msg).encode(); self.send_response(200);self.send_header('Content-Type','application/json');self.send_header('Content-Length',str(len(payload)));self.end_headers();self.wfile.write(payload);return
        start = dict(msg,content=[],stop_reason=None)
        if content['type']=='tool_use':
            block=dict(content,input={});delta={'type':'input_json_delta','partial_json':json.dumps(content['input'])}
        else:
            block={'type':'text','text':''};delta={'type':'text_delta','text':content['text']}
        events=[('message_start',{'type':'message_start','message':start}),('content_block_start',{'type':'content_block_start','index':0,'content_block':block}),('content_block_delta',{'type':'content_block_delta','index':0,'delta':delta}),('content_block_stop',{'type':'content_block_stop','index':0}),('message_delta',{'type':'message_delta','delta':{'stop_reason':stop,'stop_sequence':None},'usage':{'output_tokens':1}}),('message_stop',{'type':'message_stop'})]
        data=''.join('event: '+k+'\ndata: '+json.dumps(v)+'\n\n' for k,v in events).encode()
        self.send_response(200);self.send_header('Content-Type','text/event-stream');self.send_header('Content-Length',str(len(data)));self.end_headers();self.wfile.write(data)

SERVER=http.server.ThreadingHTTPServer(('127.0.0.1',0),Handler)
threading.Thread(target=SERVER.serve_forever,daemon=True).start()
BASE=f'http://127.0.0.1:{SERVER.server_port}'
WORK=ROOT/'work';WORK.mkdir();(WORK/'fixture.txt').write_text(FILE_TEXT)

class Client:
    def __init__(self, tool, session=None):
        self.tool=tool; self.q=queue.Queue();self.stderr=[];self.events=[]
        config=ROOT/f'cfg-{uuid.uuid4().hex}';config.mkdir()
        model={'id':'bench-model','name':'Benchmark','reasoning':False,'input':['text'],'cost':{'input':0,'output':0,'cacheRead':0,'cacheWrite':0},'contextWindow':10000000,'maxTokens':1024}
        (config/'models.json').write_text(json.dumps({'providers':{'bench':{'api':'anthropic-messages','baseUrl':BASE,'apiKey':'bench-placeholder','models':[model]}}}))
        (config/'settings.json').write_text(json.dumps({'compaction':{'enabled':False},'retry':{'enabled':False}}))
        env=dict(os.environ)
        for k in list(env):
            if any(x in k.upper() for x in ['API_KEY','AUTH_TOKEN','SECRET']):env.pop(k)
        env.update(PI_CODING_AGENT_DIR=str(config),RPI_CODING_AGENT_DIR=str(config),PI_OFFLINE='1',RPI_OFFLINE='1',PI_SKIP_VERSION_CHECK='1',ANTHROPIC_API_KEY='bench-placeholder')
        cmd=([A.rpi] if tool=='rpi' else [A.node,A.pi])+['--mode','rpc','--model','bench/bench-model']
        if session:
            sid=json.loads(session.read_text(encoding='utf-8').splitlines()[0])['id']
            cmd+=['--session-dir',str(session.parent.parent),'--session',sid if tool=='rpi' else str(session)]
        self.started=time.perf_counter()
        self.p=subprocess.Popen(cmd,cwd=WORK,env=env,stdin=subprocess.PIPE,stdout=subprocess.PIPE,stderr=subprocess.PIPE,text=True,encoding='utf-8',bufsize=1)
        threading.Thread(target=self.read,args=(self.p.stdout,False),daemon=True).start()
        threading.Thread(target=self.read,args=(self.p.stderr,True),daemon=True).start()
    def read(self, stream, err):
        for line in stream:
            if err:self.stderr.append(line);continue
            try:self.q.put((time.perf_counter(),json.loads(line)))
            except json.JSONDecodeError:pass
    def send(self, kind, **extra):
        id=uuid.uuid4().hex;self.p.stdin.write(json.dumps(dict(type=kind,id=id,**extra))+'\n');self.p.stdin.flush();return id
    def response(self,id, timeout=40):
        until=time.perf_counter()+timeout
        while time.perf_counter()<until:
            try:t,m=self.q.get(timeout=max(.01,until-time.perf_counter()))
            except queue.Empty:break
            self.events.append(m)
            if m.get('type')=='response' and m.get('id')==id:
                if m.get('success') is False or m.get('status')=='error':raise RuntimeError(str(m))
                return t,m
        raise RuntimeError(self.tool+' timed out '+''.join(self.stderr)[-1500:])
    def ready(self):return self.response(self.send('get_state'))
    def prompt(self,text):
        t=time.perf_counter();id=self.send('prompt',content=text,message=text);end,m=self.response(id,120)
        if self.tool=='pi' and not any(x.get('type')=='agent_end' for x in self.events):
            until=time.perf_counter()+120
            while time.perf_counter()<until:
                et,e=self.q.get(timeout=until-time.perf_counter());self.events.append(e)
                if e.get('type')=='agent_end':end=et;break
            else:raise RuntimeError('no agent_end')
        encoded=json.dumps(self.events)+json.dumps(m)
        if 'BENCH_DONE' not in encoded:raise RuntimeError('missing completed output '+encoded[-600:])
        return (end-t)*1000
    def rss(self):
        if os.name=='nt':
            x=subprocess.check_output(['powershell','-NoProfile','-Command',f'(Get-Process -Id {self.p.pid}).WorkingSet64'],encoding='utf-8')
            return int(x.strip())/1048576
        return int(subprocess.check_output(['ps','-o','rss=','-p',str(self.p.pid)]))/1024
    def close(self):
        self.p.kill();self.p.wait(timeout=10)
        self.p.stdin.close()

def drain():
    out=[]
    while not REQUESTS.empty():out.append(REQUESTS.get())
    return out

def fixture(tool,n):
    id=str(uuid.uuid4());ts=1791298800000;iso='2026-10-06T15:00:00.000Z';parent=None
    folder=ROOT/f'{tool}-{n}-sessions'/'bucket';folder.mkdir(parents=True)
    f=folder/f'{id}.jsonl'
    head={'kind':'header','version':4,'id':id,'createdAt':ts,'cwd':str(WORK)} if tool=='rpi' else {'type':'session','version':3,'id':id,'timestamp':iso,'cwd':str(WORK)}
    lines=[json.dumps(head)]
    for i in range(n):
        eid=f'{i+1:08x}';text=f'MESSAGE_{i:06d} '+('x'*1000)+(' BENCH_TAIL' if i==n-1 else '')
        msg={'role':'user','content':[{'type':'text','text':text}],'timestamp':ts+i}
        if i%2:
            msg.update(role='assistant',api='anthropic-messages',provider='bench',model='bench-model',stopReason='stop',usage={'input':1,'output':1,'cacheRead':0,'cacheWrite':0,'totalTokens':2,'cost':{'input':0,'output':0,'cacheRead':0,'cacheWrite':0,'total':0}})
        if tool=='rpi':msg['kind']=msg['role']
        entry={'type':'message','id':eid,'parentId':parent,'timestamp':ts+i if tool=='rpi' else iso,'message':msg}
        if tool=='rpi':entry.update(kind='entry',lane='main',seq=i+1)
        lines.append(json.dumps(entry));parent=eid
    f.write_text('\n'.join(lines)+'\n',encoding='utf-8');return f,id,parent

report={'created':time.strftime('%Y-%m-%dT%H:%M:%S%z'),'platform':platform.platform(),'node':subprocess.check_output([A.node,'--version'],encoding='utf-8').strip(),'rpi':subprocess.check_output([A.rpi,'--version'],encoding='utf-8').strip(),'pi':subprocess.check_output([A.node,A.pi,'--version'],encoding='utf-8').strip(),'runs':A.runs,'warmup':A.warmup,'fixture_text_chars':1000,'load':{},'loop':{},'method':'native JSONL completed text history; OS cache warm; local deterministic Anthropic endpoint; sequential real read tool; no cloud; no extensions; memory separately sampled'}
try:
    for n in map(int,A.sizes.split(',')):
        report['load'][n]={}
        for tool in ['rpi','pi']:
            f,sid,leaf=fixture(tool,n);times=[];first=[];payloads=[];original=f.read_bytes()
            for j in range(A.warmup+A.runs):
                f.write_bytes(original)
                drain();c=Client(tool,f)
                try:
                    rt,state=c.ready();elapsed=(rt-c.started)*1000
                    data=state.get('data',state.get('result',{}))
                    if data.get('sessionId')!=sid:raise RuntimeError('session was not loaded '+str(state))
                    if tool=='rpi' and data.get('leafId')!=leaf:raise RuntimeError('wrong leaf '+str(state))
                    pt=time.perf_counter();c.prompt('Verify restored context.');reqs=drain()
                    if len(reqs)!=1 or not reqs[0]['has_tail'] or reqs[0]['messages']!=n+1:raise RuntimeError('context mismatch '+str(reqs))
                    if j>=A.warmup:times.append(elapsed);first.append((reqs[0]['received']-pt)*1000);payloads.append(reqs[0]['body_bytes'])
                finally:c.close()
            f.write_bytes(original)
            c=Client(tool,f)
            try:c.ready();rss=c.rss()
            finally:c.close()
            row={'ready_ms':times,'ready_median_ms':statistics.median(times),'prompt_to_request_ms':first,'prompt_to_request_median_ms':statistics.median(first),'rss_ready_mib':rss,'file_bytes':f.stat().st_size,'request_bytes':payloads}
            report['load'][n][tool]=row;print('LOAD',n,tool,json.dumps(row),flush=True)
    for turns in map(int,A.turns.split(',')):
        TASK_TURNS = turns
        report['loop'][turns]={}
        for tool in ['rpi','pi']:
            times=[]
            for j in range(A.warmup+A.runs):
                drain();c=Client(tool)
                try:
                    c.ready();elapsed=c.prompt(f'BENCH_LOOP={turns} Execute the scripted read calls.');reqs=drain()
                    if len(reqs)!=turns+1 or reqs[-1]['results']!=turns or not all(q['tool_ok'] for q in reqs):raise RuntimeError('tool loop mismatch '+str(reqs))
                    if j>=A.warmup:times.append(elapsed)
                finally:c.close()
            row={'duration_ms':times,'median_ms':statistics.median(times),'tool_calls':turns,'provider_requests':turns+1,'tool_file_bytes':len(FILE_TEXT.encode())}
            report['loop'][turns][tool]=row;print('LOOP',turns,tool,json.dumps(row),flush=True)
except Exception as e:
    report['error']=repr(e);print('ERROR',repr(e),flush=True)
finally:
    report['script_sha256']=hashlib.sha256(pathlib.Path(__file__).read_bytes()).hexdigest()
    pathlib.Path(A.output).parent.mkdir(parents=True,exist_ok=True)
    pathlib.Path(A.output).write_text(json.dumps(report,ensure_ascii=False,indent=2),encoding='utf-8')
    SERVER.shutdown();shutil.rmtree(ROOT,ignore_errors=True)
if 'error' in report:raise SystemExit(1)
