"""Lab22 fixed offline microscope: prepare, verify, demonstrate, analyze, export.

Only demonstrate may load the model. Every command has a parent watchdog. A
finalized raw run is immutable; neither repairs nor timeouts permit reentry.
"""
import argparse,contextlib,fcntl,importlib.metadata,json,os,platform,selectors,signal,subprocess,sys,time,uuid
from pathlib import Path
from microscope_io import (SCHEMA,RAW_CAP,PAYLOAD_CAP,Store,safe,relative,separate,bounded,read_json,
    digest,encode,stamp,inventory,size,hashes,check_hashes,finalize,finalized,npz,events)
from microscope_imports import (ARTIFACTS,VERSIONS,REPOSITORY,REVISION,DEMO_TEXTS,CAPTURE_KEYS,
    SOURCE17,SOURCE19,address,tokens,lab17,lab19,snapshot)
os.environ['HF_HUB_OFFLINE']='1';os.environ['TRANSFORMERS_OFFLINE']='1'

SCHEDULE=[{'prompt':p,'arm':a} for p in DEMO_TEXTS for a in ('baseline','observe','zero','positive','negative','restore')]+[
    {'prompt':'H1+','arm':'abort'},{'prompt':'H1+','arm':'recovery'}]
LIMITS={'processing_seconds':900,'forwards':14,'tokens':64,'batch_size':1,'raw_bytes':RAW_CAP,
        'numeric_payload_bytes':PAYLOAD_CAP,'member_bytes':20_000_000,'header_bytes':10_000,
        'analysis_seconds':60,'analysis_bytes':5_000_000,'export_seconds':60,'export_bytes':RAW_CAP}
TOLERANCES={'replay_atol':1e-6,'replay_rtol':1e-6,'structure_atol':1e-6,'structure_rtol':1e-5,
            'direction_atol':1e-10,'direction_rtol':1e-10}
PLAN_FIXED={'schema':SCHEMA,'schedule':SCHEDULE,'limits':LIMITS,'tolerances':TOLERANCES,
 'boundary':'raw post-block residual after additions before block 3','block_index':2,
 'metric':{'name':'good-minus-bad final-position logit','candidate_texts':[' good',' bad'],'candidate_ids':[1175,3076]},
 'execution':{'device':'cpu','dtype':'float32','batch_size':1,'seed':19,'deterministic':True,'intra_op':1,'inter_op':1,
    'attention_backend':'eager','use_cache':False,'output_attentions':False,'output_hidden_states':False,'logits_to_keep':1,
    'attention_mask':'explicit all-ones saved mask','padding':False,'truncation':False},
 'expected_checks':['observer','zero','restore','replay_baseline','replay_positive','replay_negative',
    'parallel_addition','final_normalization','normalized_head','raw_head','expected_exception_cleanup','recovery'],
 'prior_evidence_access':'Previously inspected Lab19 prompts/results and corrected post-exposure validation; fresh run is integration replay, not untouched evaluation.'}
SOURCE_NAMES=('microscope.py','microscope_io.py','microscope_imports.py','microscope_analysis.py',
 'microscope-plan.example.json','inspect_activations.py','contrast_experiment.py','inspection_common.py',
 'inspect_residual_stream.py','lab19-fixtures.json')


ENTRYPOINT=Path(__file__).resolve()
VERIFICATION_KIND='canonical_same_environment'

def configure_parser(parser):
    """The canonical CLI has no additional validation modes."""
    pass

def verify_environment(args,root,historical,current):
    for name,old in historical.items():
        if old['python']!=current['python'] or old['packages']!=current['packages'] or old['architecture']!=current['architecture']:
            raise ValueError('Claimed same-runtime replay incompatible: '+name)
    return None

def plan(value):
    if set(value)!=set(PLAN_FIXED)|{'question','predictions'}:raise ValueError('Unknown/missing plan fields')
    if any(encode(value.get(k))!=encode(v) for k,v in PLAN_FIXED.items()):raise ValueError('Only the exact fixed core contract is supported')
    if not isinstance(value['question'],str) or not value['question'].strip() or value['question'].startswith('REPLACE') or len(value['question'])>4000:
        raise ValueError('Write a bounded question before preparation')
    expected={'observer','zero','cleanup','replay','positive','negative'}
    if not isinstance(value['predictions'],dict) or set(value['predictions'])!=expected or any(not isinstance(s,str) or not s.strip() or s.startswith('REPLACE') or len(s)>4000 for s in value['predictions'].values()):
        raise ValueError('Freeze all integration predictions with prior exposure disclosed')
    return value

def sources():return {n:digest(bounded(Path(__file__).with_name(n),1_000_000)) for n in SOURCE_NAMES}
def state(root):
    value=read_json(relative(root,'status.json'))
    if value.get('schema')!=SCHEMA or not isinstance(value.get('run_id'),str):raise ValueError('Unsupported run schema')
    if value.get('provenance')!='measured':raise ValueError('Fixture cannot satisfy measured-run contract')
    return value

def frozen_map(root):
    names=['plan.json','environment.json','artifact-manifest.json','source-manifest.json','imports.json','inputs.json','predictions.json']
    if relative(root,'portability.json').exists():names.append('portability.json')
    names += [n for n in inventory(root) if n.startswith(('imports/','source/'))]
    return {n:digest(bounded(relative(root,n))) for n in names}

def protect(root):
    if state(root).get('verification_kind','canonical_same_environment')!=VERIFICATION_KIND:raise ValueError('Run validation mode differs from entrypoint')
    expected=read_json(relative(root,'frozen.sha256.json'))
    if frozen_map(root)!=expected:raise ValueError('Frozen inputs/source/imports changed')
    if read_json(relative(root,'source-manifest.json'))['current_sources']!=sources():raise ValueError('Current executable source changed; new run required')
    imported=read_json(relative(root,'imports.json'))
    if any(imported[k].get('provenance')!='measured' for k in ('lab17','lab19')):raise ValueError('Imported fixture cannot satisfy measured prerequisite')
    return imported

def validate_imports(root):
    a=lab17(relative(root,'imports/lab17'));b=lab19(relative(root,'imports/lab19'))
    if b['payload_bytes']+sum(v.nbytes for v in a['vectors'].values())>PAYLOAD_CAP:raise ValueError('Raw numeric payload cap')
    original=read_json(relative(root,'imports.json'))
    for name,result in [('lab17',a),('lab19',b)]:
        if result['origin_run_id']!=original[name]['origin_run_id'] or result['file_hashes']!=original[name]['file_hashes']:
            raise ValueError('Import mapping/hash differs')
    return a,b

def actual_environment(torch=None):
    import numpy as np
    packages={k:importlib.metadata.version(k) for k in VERSIONS}
    result={'schema':SCHEMA,'python':platform.python_version(),'implementation':platform.python_implementation(),
        'packages':packages,'numpy':np.__version__,'os':platform.system(),'os_release':platform.release(),
        'architecture':platform.machine(),'cpu_count':os.cpu_count(),'device':'cpu','dtype':'float32',
        'seed':19,'attention_backend':'eager','use_cache':False,'deterministic_requested':True,
        'HF_HUB_OFFLINE':'1','TRANSFORMERS_OFFLINE':'1','routing':'not_applicable: dense specimen',
        'intra_op':torch.get_num_threads() if torch is not None else 'not_verified',
        'inter_op':torch.get_num_interop_threads() if torch is not None else 'not_verified',
        'deterministic_effective':torch.are_deterministic_algorithms_enabled() if torch is not None else 'not_verified'}
    return result

def runtime():
    if platform.python_implementation()!='CPython' or platform.python_version()!='3.13.13':raise ValueError('Pinned CPython3.13.13 required')
    for k,v in VERSIONS.items():
        if importlib.metadata.version(k).split('+')[0]!=v:raise ValueError('Pinned package differs: '+k)
    os.environ['HF_HUB_OFFLINE']='1';os.environ['TRANSFORMERS_OFFLINE']='1'
    import torch
    torch.set_num_threads(1);torch.set_num_interop_threads(1)
    torch.manual_seed(19);torch.use_deterministic_algorithms(True)
    return torch

def local_artifacts(directory):
    directory=safe(directory)
    if not directory.is_dir():raise ValueError('Existing local cache required; no acquisition fallback')
    for n,record in ARTIFACTS.items():
        data=bounded(relative(directory,n),31_000_000)
        if len(data)!=record['bytes'] or digest(data)!=record['sha256']:raise ValueError('Local artifact hash differs: '+n)
    config=read_json(relative(directory,'config.json'))
    for k,v in {'model_type':'gpt_neox','num_hidden_layers':6,'hidden_size':128,'vocab_size':50304,
                'use_parallel_residual':True,'tie_word_embeddings':False}.items():
        if config.get(k)!=v:raise ValueError('Local configuration differs: '+k)

def retokenize(root,directory):
    from transformers import AutoTokenizer
    tokenizer=AutoTokenizer.from_pretrained(str(safe(directory)),local_files_only=True,trust_remote_code=False,use_fast=True)
    if not tokenizer.is_fast:raise ValueError('Fast tokenizer required')
    frozen=read_json(relative(root,'inputs.json'))
    if frozen.get('attention_mask_convention')!='explicit saved all-ones mask':raise ValueError('Unsupported computational convention')
    if [r['id'] for r in frozen['rows']]!=list(DEMO_TEXTS):raise ValueError('Exactly two declared demo inputs')
    for row in frozen['rows']:
        row=tokens(row)
        if row['text']!=DEMO_TEXTS[row['id']]:raise ValueError('Exact frozen text differs')
        encoded=tokenizer(row['text'],add_special_tokens=False,padding=False,truncation=False)
        pieces=[tokenizer.decode([i],clean_up_tokenization_spaces=False) for i in encoded['input_ids']]
        if encoded['input_ids']!=row['token_ids'] or encoded['attention_mask']!=row['attention_mask'] or pieces!=row['decoded_pieces']:
            raise ValueError('Exact local tokenization differs')
    ids=[tokenizer.encode(s,add_special_tokens=False) for s in (' good',' bad')]
    if ids!=[[1175],[3076]]:raise ValueError('Single-token candidate metric differs')
    return tokenizer

def prepare(args,store):
    value=plan(read_json(safe(args.plan),100_000));store.json('plan.json',value)
    a,b=snapshot(store,safe(args.lab17),safe(args.lab19))
    validate_imports(store.root)
    store.json('inputs.json',{'schema':SCHEMA,'rows':b['demo'],'attention_mask_convention':'explicit saved all-ones mask',
        'provenance':'copied actual Lab19 token fields; absent length separately derived'})
    store.json('predictions.json',{'schema':SCHEMA,'frozen_utc':stamp(),'prior_evidence_access':value['prior_evidence_access'],
        'integration_expectations':value['predictions'],'original_lab19_predictions_sha256':b['predictions_sha256'],
        'historical_prediction_status':b['prior_exposure']})
    current=sources()
    for n in SOURCE_NAMES:
        data=bounded(Path(__file__).with_name(n),1_000_000)
        if digest(data)!=current[n]:raise ValueError('Source changed during preparation')
        store.put('source/'+n,data)
    store.json('source-manifest.json',{'schema':SCHEMA,'current_sources':current,'historical_sources':{'lab17':SOURCE17,'lab19':SOURCE19},
        'compatibility_notes':a['compatibility_notes'],
        'inspected_numeric_reuse':'Lab17 capture/run_forward/checked_logits; corrected Lab19 hook/replacement. Imported code is never executed.'})
    store.json('artifact-manifest.json',{'repository':REPOSITORY,'revision':REVISION,'artifacts':ARTIFACTS})
    store.json('environment.json',actual_environment())
    store.json('frozen.sha256.json',frozen_map(store.root))

def verify(args,store):
    protect(store.root);a,b=validate_imports(store.root);local_artifacts(args.model_dir)
    torch=runtime();tokenizer=retokenize(store.root,args.model_dir)
    env=actual_environment(torch)
    historical={'lab17':read_json(relative(store.root,'imports/lab17/manifest.json')),
                'lab19':read_json(relative(store.root,'imports/lab19/environment.json'))}
    comparison=verify_environment(args,store.root,historical,env)
    store.json('environment.json',env,replace=True)
    store.json('verification.json',{'schema':SCHEMA,'verified_utc':stamp(),'passed':True,'model_loaded':False,'forwards':0,
        'verification_kind':state(store.root).get('verification_kind','canonical_same_environment'),'portability_comparison':comparison,
        'historical_numpy':'not recorded in either historical run; actual current NumPy recorded; native numeric execution uses pinned PyTorch',
        'historical_settings':'Historical seed0, current declared seed19; no stochastic sampling in eval. New deterministic flag explicit.',
        'tokenizer_class':type(tokenizer).__name__,'source_sha256':sources(),'checks':['imports','arrays','source','runtime','artifacts','configuration','exact_tokens']})
    store.json('frozen.sha256.json',frozen_map(store.root),replace=True)

class Ledger:
    def __init__(self,store):self.store=store;self.active=None
    def check(self):
        current=state(self.store.root)
        if time.monotonic()>=current['active_deadline_monotonic']:raise TimeoutError('Cumulative processing allowance exhausted')
    def forward(self):
        self.check();current=state(self.store.root);n=current['attempted']
        if n>=14:raise ValueError('Fifteenth forward refused before model call')
        expected=SCHEDULE[n]
        if self.active is None or self.active!={**expected,'attempt':n+1}:raise ValueError('Exact scheduled prompt/arm required')
        current['attempted']=n+1;self.store.json('status.json',current,replace=True)
        self.store.log({**self.active,'timestamp':stamp(),'outcome':'started','cleanup':'pending'})

def measurement_address(run_id,row,arm,attempt,key,source_shape,provenance='measured',parents=None):
    if key=='logits':module,block,normalization,boundary='embed_out',None,'not_applicable: output logits','final-position logits'
    elif key in ('original','replacement','positive_addition','negative_addition'):
        module,block,normalization,boundary='gpt_neox.layers.2/output',2,'raw','raw post-block residual after additions before block 3'
    else:module,block,normalization=address(key);boundary=key
    return {'run_id':run_id,'origin_run_id':run_id,'prompt':row['id'],'arm':arm,'attempt':attempt,
        'module_path':module,'boundary':boundary,'normalization':normalization,'block_index':block,
        'batch_row':0,'token_id':row['token_ids'][-1],'token_position':row['final_position'],
        'source_shape':source_shape,'stored_shape':[50304] if key=='logits' else [128],
        'dtype':'float32','provenance':provenance,'parents':parents or []}

def demonstrate(args,store):
    from inspect_activations import capture,run_forward
    from contrast_experiment import hook,CleanupCheck
    protect(store.root);a,b=validate_imports(store.root);local_artifacts(args.model_dir)
    torch=runtime();retokenize(store.root,args.model_dir)
    if actual_environment(torch)!=read_json(relative(store.root,'environment.json')):raise ValueError('Environment differs from verification')
    if state(store.root).get('verification_kind','canonical_same_environment')!=VERIFICATION_KIND:raise ValueError('Run validation mode differs from entrypoint')
    from inspection_common import load_specimen
    model,_,_=load_specimen(args.model_dir,torch)
    # The pinned implementation returns tuples at blocks; capture/hook assert
    # this runtime boundary on every observed/edited pass before accepting data.
    rows={r['id']:r for r in read_json(relative(store.root,'inputs.json'))['rows']}
    u=torch.tensor(b['direction']['u'],dtype=torch.float64);dose=b['direction']['dose']
    ledger=Ledger(store);logits={};captures={};metadata={};checks=[];baselines={};metrics=[]
    run_state=state(store.root);run_id=run_state['run_id'];value_provenance=run_state['provenance']
    diagnostic_counts={'final_layer_norm':0,'output_head':0}
    def save():
        store.arrays('captures.npz',captures);store.arrays('logits.npz',logits)
        store.json('array-metadata.json',metadata,replace=True)
        store.json('checks.json',checks,replace=True)
    def agree(name,left,right,structural=False):
        record={'name':name,'max_abs_error':float((left-right).abs().max()),'atol':1e-6,
            'rtol':1e-5 if structural else 1e-6,'eligible':True,'passed':bool(torch.allclose(left,right,atol=1e-6,rtol=1e-5 if structural else 1e-6))}
        checks.append(record);store.json('checks.json',checks,replace=True)
        if not record['passed']:raise ValueError('Numerical invariant failed: '+name)
    def diagnostic(kind,module,value):
        diagnostic_counts[kind]+=1
        store.json('diagnostic-module-calls.json',diagnostic_counts,replace=True)
        current=state(store.root);current['diagnostic_module_calls']=diagnostic_counts
        store.json('status.json',current,replace=True)
        event={'module':kind,'call':diagnostic_counts[kind],'attempt':ledger.active['attempt'],'outcome':'started','timestamp':stamp()}
        # A separate fsynced module ledger distinguishes attempted readouts from
        # returned calls if the worker is terminated during the module itself.
        def log(record):store.log(record,'diagnostic-module-events.jsonl')
        log(event);outcome='failed'
        try:
            result=module(value);outcome='returned';return result
        finally:log({**event,'timestamp':stamp(),'outcome':outcome})
    def record_array(number,key,value,row,arm,parents=None):
        name=f'attempt_{number:03d}_'+key
        if value.dtype!=torch.float32 or not torch.isfinite(value).all():raise ValueError('Native float32 finite captures required')
        target=logits if key=='logits' else captures;target[name]=value.detach().cpu().numpy().copy()
        metadata[name]=measurement_address(run_id,row,arm,number,key,[1,1,50304] if key=='logits' else [1,row['length'],128],
            'derived' if parents else value_provenance,parents)
    # Preserve the native float32 additions, including sign/cast order, with
    # named frozen parents. No random direction is produced in this integration.
    for arm,sign in [('positive',1),('negative',-1)]:
        addition=(u*dose*sign*1).float()
        key=arm+'_addition';captures[key]=addition.numpy().copy()
        metadata[key]=measurement_address(run_id,rows['H1+'],arm,0,key,[128],'derived',
            ['imports/lab19/discovery-vectors.npz','imports/lab19/direction.json'])
        metadata[key]['actual_addition_norm_float64']=float(addition.double().norm())
        metadata[key]['source_shape']=[128]
    with torch.inference_mode():
        for number,item in enumerate(SCHEDULE,1):
            ledger.check();row=rows[item['prompt']];arm=item['arm'];holder={'calls':0,'cleanup':True}
            ledger.active={**item,'attempt':number};before={id(m):(len(m._forward_hooks),len(m._forward_pre_hooks)) for m in model.modules()}
            result=None;values={};outcome='failed';error=None
            try:
                if arm=='observe':values,result=capture(model,row,ledger,torch)
                elif arm in ('zero','positive','negative','abort'):
                    source_arm={'positive':'positive_small','negative':'negative_small'}.get(arm,arm)
                    with hook(model,row,source_arm,holder,torch,u=u,dose=dose):result=run_forward(model,row,ledger,torch)
                    if holder['calls']!=1:raise ValueError('Intervention invocation count')
                    values={k:holder[k] for k in ('original','replacement') if k in holder}
                else:result=run_forward(model,row,ledger,torch)
                if arm=='abort':raise ValueError('Expected named cleanup exception did not occur')
                for key,value in values.items():record_array(number,key,value,row,arm)
                record_array(number,'logits',result,row,arm);save();outcome='completed_logits'
            except CleanupCheck as exc:
                if arm!='abort' or str(exc)!='deliberate_cleanup_check':raise
                outcome='expected_exception';error='CleanupCheck: '+str(exc)
            except BaseException as exc:
                error=type(exc).__name__+': '+str(exc);raise
            finally:
                cleanup=all(before[id(m)]==(len(m._forward_hooks),len(m._forward_pre_hooks)) for m in model.modules())
                store.log({**ledger.active,'timestamp':stamp(),'outcome':outcome,'exception':error,
                    'hook_calls':20 if arm=='observe' and result is not None else holder['calls'],'cleanup':'verified' if cleanup else 'failed'})
                current=state(store.root)
                current['finished']+=int(outcome in ('completed_logits','expected_exception'))
                current['diagnostic_module_calls']=diagnostic_counts;store.json('status.json',current,replace=True)
                if not cleanup:raise ValueError('Registered hooks not restored')
            if arm=='abort':
                if holder['calls']!=1 or holder['cleanup'] is not True:raise ValueError('Expected-exception hook cleanup')
                checks.append({'name':'expected_exception_cleanup','eligible':True,'passed':True});save();continue
            baseline=baselines.get(row['id'])
            if arm=='baseline':baselines[row['id']]=result;baseline=result
            if arm in ('observe','zero','restore','recovery'):agree(row['id']+' '+arm,result,baseline)
            if arm in ('baseline','positive','negative'):
                source_arm={'positive':'positive_small','negative':'negative_small'}.get(arm,arm)
                old=next(e for e in b['attempts'] if e['prompt']==row['id'] and e['arm']==source_arm and e['attempt']>=23)
                agree(row['id']+' replay '+arm,result,torch.from_numpy(b['logits'][f"attempt_{old['attempt']:03d}"]))
            if arm=='observe':
                for i in range(6):agree(row['id']+f' parallel_addition_{i}',(values[f'm{i}']+values[f'a{i}'])+values[f'r{i}'],values[f'r{i+1}'],True)
                normalized=diagnostic('final_layer_norm',model.gpt_neox.final_layer_norm,values['r6'])
                agree(row['id']+' final_normalization',normalized,values['h_f'],True)
                normalized_logits=diagnostic('output_head',model.embed_out,values['h_f'])
                agree(row['id']+' normalized_head',normalized_logits,result,True)
                normalized_again=diagnostic('final_layer_norm',model.gpt_neox.final_layer_norm,values['r6'])
                raw_logits=diagnostic('output_head',model.embed_out,normalized_again)
                agree(row['id']+' raw_head',raw_logits,result,True)
                store.json('diagnostic-module-calls.json',diagnostic_counts,replace=True)
            ledger.check();save()
    if state(store.root)['attempted']!=14 or state(store.root)['finished']!=14 or len(logits)!=13 or diagnostic_counts!={'final_layer_norm':4,'output_head':4}:
        raise ValueError('Incomplete exact integration schedule')
    # Analysis is independently regenerated later; this raw metric table remains
    # a complete native-logit readout of all thirteen successful attempts.
    import csv,io
    stream=io.StringIO();writer=csv.writer(stream);writer.writerow(['attempt','prompt','arm','y_good_minus_bad','eligible','provenance'])
    for number,item in enumerate(SCHEDULE,1):
        value=logits.get(f'attempt_{number:03d}_logits')
        writer.writerow([number,item['prompt'],item['arm'],float(value[1175].astype('float64')-value[3076].astype('float64')) if value is not None else '',value is not None,value_provenance if value is not None else 'missing: expected_exception'])
    store.put('metrics.csv',stream.getvalue().encode(),replace=True)


def worker(args):
    store=Store(args.out if args.command in ('prepare','analyze','export') else args.run,
                5_000_000 if args.command=='analyze' else RAW_CAP)
    if args.command=='prepare':prepare(args,store)
    elif args.command=='verify':verify(args,store)
    elif args.command=='demonstrate':demonstrate(args,store)
    else:
        from microscope_analysis import analyze,export
        with relative(args.run,'run.lock').open('rb') as lock:
            fcntl.flock(lock.fileno(),fcntl.LOCK_SH|fcntl.LOCK_NB)
            (analyze if args.command=='analyze' else export)(args,store)

def watch(process,log,root,deadline,cap,reserve=200_000):
    """Supervise stalled native calls and bounded worker output/storage."""
    total_log=0;selector=selectors.DefaultSelector()
    selector.register(process.stdout,selectors.EVENT_READ)
    try:
        while selector.get_map():
            close=relative(root,'closing-deadline.json')
            active_deadline=min(deadline,read_json(close,2000)['deadline_monotonic']) if close.exists() else deadline
            if time.monotonic()>=active_deadline:raise TimeoutError('Parent watchdog processing/terminal-close deadline')
            if size(root)>cap-reserve:raise TimeoutError('Parent watchdog storage reserve')
            for key,_ in selector.select(timeout=min(.1,max(.001,deadline-time.monotonic()))):
                data=os.read(key.fileobj.fileno(),65536)
                if not data:selector.unregister(key.fileobj);continue
                total_log+=len(data)
                if total_log>500_000:raise TimeoutError('Worker log cap')
                Store(root,cap,terminal=reserve==0).append(log,data,active_deadline)
        return process.wait(timeout=max(.001,deadline-time.monotonic()))
    finally:
        selector.close();process.stdout.close()


def finish(args):
    """Supervised terminal bookkeeping; last bounded close is conservatively charged."""
    offline=args.command in ('analyze','export');root=args.out if args.command in ('prepare','analyze','export') else args.run
    store=Store(root,5_000_000 if args.command=='analyze' else RAW_CAP,terminal=True)
    request=read_json(relative(root,'finish-request.json'));started=request['started_monotonic']
    failure=request['failure'];partial=request['interrupted'];rescue=request['failure_finalization']
    terminal_state='partial' if partial else 'failed' if failure else 'completed'
    if offline:
        record=read_json(relative(root,'manifest.json'));record.update(state=terminal_state,reason=failure,finished_utc=stamp())
        record.pop('supervisor_pid',None);record.pop('worker_nonce_sha256',None)
        store.json('manifest.json',record,replace=True)
    else:
        record=state(root);record[args.command]=terminal_state
        if partial and args.command=='demonstrate':
            log=events(relative(root,'attempts.jsonl'));finished={e['attempt'] for e in log if e.get('outcome')!='started'}
            starts={e['attempt']:e for e in log if e.get('outcome')=='started'}
            for number in range(1,record['attempted']+1):
                if number not in finished:
                    event=starts.get(number,{'attempt':number,**SCHEDULE[number-1],'start_event_missing':True})
                    store.log({**event,'timestamp':stamp(),'outcome':'partial','cleanup':'cleanup_unverified','exception':failure})
        # Prepare/verify have no final hash map: keep them blocked until the
        # parent publishes its nonce-bound receipt. Terminal runs instead
        # require the integrity manifest's parent observation for all reads.
        if not failure and args.command!='demonstrate':
            record['pending_close_nonce']=request['close_nonce'];record['active_command']=args.command
        else:
            record.pop('pending_close_nonce',None);record['active_command']=None
        record.pop('active_deadline_monotonic',None)
        record.pop('supervisor_pid',None);record.pop('worker_nonce_sha256',None)
        record['reason']=failure or ''
        record['scientific_outcome']=('supplemental portability replay; not canonical same-environment acceptance' if record.get('verification_kind')=='supplemental_portability' else 'descriptive integration replay; no generalization claim') if args.command=='demonstrate' and not failure else 'unavailable: fresh integration not completed'
        store.json('status.json',record,replace=True)
        if failure or args.command=='demonstrate':
            checks=read_json(relative(root,'checks.json'));checks.append({'name':args.command+' terminal','passed':failure is None,'eligible':not bool(failure),'reason':failure})
            store.json('checks.json',checks,replace=True)
    terminal=offline or bool(failure) or args.command=='demonstrate'
    # Hashing every saved byte is part of processing, before the measured close.
    mapping=hashes(root,('files.sha256.json',)) if terminal else None
    now=time.monotonic();measured=now-started;upper=measured+1.0
    if not rescue and now+1.0>request['deadline_monotonic']:raise TimeoutError('No time for bounded terminal close')
    closing={'deadline_monotonic':now+1.0,'seconds':1.0,'purpose':'final status/manifest serialization and parent receipt publication'}
    store.json('closing-deadline.json',closing,replace=True)
    receipt={'command':args.command,'measured_seconds_through_terminal_hashing':measured,'charged_seconds_upper_bound':upper,
             'close_allowance_seconds':1.0,'failure_finalization_allowance_seconds':10 if rescue else 0,
             'processing_limit_exceeded':rescue,'actual_parent_observation':'pending within charged close'}
    if offline:
        record.update(elapsed_processing_seconds=measured,charged_processing_seconds=upper,limit_seconds=60,
                      timing_boundary='terminal metadata and hashing included; last one-second close charged and parent supervised',
                      storage_bytes_before_final_map=size(root),processing_limit_exceeded=rescue)
        store.json('manifest.json',record,replace=True)
    else:
        record['processing_seconds']=request['previous_processing_seconds']+upper
        record['command_durations'][args.command]=measured
        record.setdefault('command_time_bounds',{})[args.command]=[measured,upper]
        record['timing_boundary']='terminal metadata and hashing included; last one-second close charged and parent supervised'
        record['processing_limit_exceeded']=rescue or record['processing_seconds']>900
        record['storage_bytes_before_final_map']=size(root)
        store.json('status.json',record,replace=True)
    if terminal:
        for name in ('closing-deadline.json','manifest.json' if offline else 'status.json'):
            mapping[name]=digest(bounded(relative(root,name)))
        store.json('files.sha256.json',{'schema':'llm-microscope-integrity/1','files':mapping,'processing_receipt':receipt},replace=True)


def reap(process):
    if process is None:return
    # Cancellation must not break the stop/wait critical section itself.
    with defer_spawn_cancellation(include_alarm=True):
        if process.poll() is None:
            try:os.killpg(process.pid,signal.SIGKILL)
            except ProcessLookupError:pass
        try:process.wait(timeout=5)
        except subprocess.TimeoutExpired:raise RuntimeError('Worker not reaped; evidence remains unfinalized')



@contextlib.contextmanager
def defer_spawn_cancellation(include_alarm=False):
    # Deliver cancellation only after Popen's child is assigned to a reference
    # that the surrounding finally block can stop and reap.
    blocked={signal.SIGTERM,signal.SIGINT}
    if include_alarm:blocked.add(signal.SIGALRM)
    previous=signal.pthread_sigmask(signal.SIG_BLOCK,blocked)
    try:yield
    finally:signal.pthread_sigmask(signal.SIG_SETMASK,previous)


def supervise(args):
    # SIGTERM is catchable cancellation; SIGKILL/power loss cannot run cleanup.
    def cancelled(signum,frame):
        signal.signal(signal.SIGTERM,signal.SIG_IGN)
        raise KeyboardInterrupt('Supervisor cancelled by SIGTERM')
    previous=signal.signal(signal.SIGTERM,cancelled)
    try:return _supervise(args)
    finally:signal.signal(signal.SIGTERM,previous)


def _supervise(args):
    started=getattr(args,'scope_started',time.monotonic());offline=args.command in ('analyze','export');is_prepare=args.command=='prepare'
    if not is_prepare and state(args.run).get('verification_kind','canonical_same_environment')!=VERIFICATION_KIND:
        raise ValueError('Run validation mode differs from entrypoint')
    nonce=uuid.uuid4().hex
    root=safe(args.out if is_prepare or offline else args.run)
    if is_prepare or offline:
        separate(root,*(safe(p) for p in ([args.plan,args.lab17,args.lab19] if is_prepare else [args.run]+([args.analysis] if args.command=='export' else []))))
        root.mkdir(parents=True,mode=0o700,exist_ok=False)
        if offline:
            store=Store(root,5_000_000 if args.command=='analyze' else RAW_CAP)
            store.json('manifest.json',{'schema':SCHEMA,'run_id':str(uuid.uuid4()),'stage':args.command,'state':'running','started_utc':stamp(),
                'supervisor_pid':os.getpid(),'worker_nonce_sha256':digest(nonce.encode())})
        else:
            store=Store(root);store.put('run.lock',b'')
            store.json('status.json',{'schema':SCHEMA,'run_id':str(uuid.uuid4()),'provenance':'measured',
                'prepare':'not_run','verify':'not_run','demonstrate':'not_run','attempted':0,'finished':0,
                'processing_seconds':0,'command_durations':{},'scientific_outcome':'unavailable: not_run','reason':'',
                'active_command':None,'created_utc':stamp(),
                'verification_kind':VERIFICATION_KIND})
            store.json('checks.json',[]);store.put('attempts.jsonl',b'')
    else:
        if relative(root,'files.sha256.json').exists():raise ValueError('Finalized raw run rejects mutation')
        if not relative(root,'run.lock').is_file():raise ValueError('Owned prepared run required')
        store=Store(root)
    lock=contextlib.nullcontext() if offline else relative(root,'run.lock').open('rb')
    with lock as handle:
        if not offline:
            fcntl.flock(handle.fileno(),fcntl.LOCK_EX|fcntl.LOCK_NB)
            if relative(root,'files.sha256.json').exists():raise ValueError('Finalized raw run rejects mutation')
            current=state(root)
            if current.get('pending_close_nonce') is not None or current['active_command'] is not None or any(current[k] in ('running','failed','partial') for k in ('prepare','verify','demonstrate')):
                raise ValueError('Crashed/failed/partial run cannot reenter; preserve it and create a reviewed new run')
            if not is_prepare:
                if current['prepare']!='completed' or current['attempted']!=0 or current['demonstrate']!='not_run':raise ValueError('Unstarted prepared demonstration required')
                if args.command=='verify' and current['verify']!='not_run':raise ValueError('Verification cannot be repeated')
                if args.command=='demonstrate' and current['verify']!='completed':raise ValueError('Successful verification required')
            allowance=min(900-current['processing_seconds'],getattr(args,'scope_allowance',900))
            current[args.command]='running';current['active_command']=args.command
            current['supervisor_pid']=os.getpid();current['worker_nonce_sha256']=digest(nonce.encode())
            current['active_deadline_monotonic']=started+allowance-5;store.json('status.json',current,replace=True)
            store.json('closing-deadline.json',{'deadline_monotonic':started+allowance,
                'purpose':'active command; replaced by bounded terminal close'},replace=True)
        else:allowance=min(60,getattr(args,'scope_allowance',60))
        failure=None;interrupted=False;process=None
        try:
            if time.monotonic()-started>=allowance:raise TimeoutError('Remaining command allowance exhausted')
            with relative(root,args.command+'-worker.log').open('xb') as log:
                with defer_spawn_cancellation(include_alarm=True):
                    process=subprocess.Popen([sys.executable,str(ENTRYPOINT),*sys.argv[1:],'--internal-worker'],
                        stdout=subprocess.PIPE,stderr=subprocess.STDOUT,start_new_session=True,
                        env={**os.environ,'LLMCOURSE_MICROSCOPE_PARENT_PID':str(os.getpid()),
                             'LLMCOURSE_MICROSCOPE_NONCE':nonce})
                code=watch(process,log,root,started+allowance-5,store.cap)
            if code:
                interrupted=code<0
                failure=('Worker terminated by signal; cleanup unverified' if code<0 else 'Worker failed')+'; inspect '+args.command+'-worker.log'
        except (subprocess.TimeoutExpired,TimeoutError,KeyboardInterrupt) as exc:
            interrupted=True;failure=type(exc).__name__+': processing interrupted'
        except BaseException as exc:failure=type(exc).__name__+': '+str(exc)
        finally:
            if process is not None and process.poll() is None:interrupted=True
            reap(process)
        error_file=relative(root,'worker-error.json')
        if error_file.exists() and read_json(error_file).get('type')=='TimeoutError':interrupted=True
        previous=0 if offline else current['processing_seconds']
        request={'command':args.command,'close_nonce':nonce,'started_monotonic':started,'previous_processing_seconds':previous,
            'failure':failure,'interrupted':interrupted,'failure_finalization':False,'deadline_monotonic':started+allowance}
        terminal_store=Store(root,store.cap,terminal=True)
        for rescue in (False,True):
            if rescue:
                request.update(failure='Terminal processing limit/finalization failure: '+str(failure),interrupted=True,
                               failure_finalization=True,deadline_monotonic=time.monotonic()+10)
                # Retain original terminal-close marker as evidence; the new
                # rescue close will replace it under the same owned lock.
                old_close=relative(root,'closing-deadline.json')
                if old_close.exists():terminal_store.put('previous-closing-deadline.json',bounded(old_close),replace=True)
                terminal_store.json('closing-deadline.json',{'deadline_monotonic':request['deadline_monotonic']},replace=True)
                failure=request['failure']
            terminal_store.json('finish-request.json',request,replace=True)
            if not offline:
                current=read_json(relative(root,'status.json'));current['supervisor_pid']=os.getpid();current['worker_nonce_sha256']=digest(nonce.encode())
                terminal_store.json('status.json',current,replace=True)
            else:
                manifest=read_json(relative(root,'manifest.json'));manifest['supervisor_pid']=os.getpid();manifest['worker_nonce_sha256']=digest(nonce.encode())
                terminal_store.json('manifest.json',manifest,replace=True)
            finisher=None
            try:
                with relative(root,args.command+'-finish-worker'+('-rescue' if rescue else '')+'.log').open('xb') as log:
                    argv=[x for x in sys.argv[1:] if x!='--internal-worker']
                    with defer_spawn_cancellation(include_alarm=True):
                        finisher=subprocess.Popen([sys.executable,str(ENTRYPOINT),*argv,'--internal-finish'],
                            stdout=subprocess.PIPE,stderr=subprocess.STDOUT,start_new_session=True,
                            env={**os.environ,'LLMCOURSE_MICROSCOPE_PARENT_PID':str(os.getpid()),'LLMCOURSE_MICROSCOPE_NONCE':nonce})
                    code=watch(finisher,log,root,request['deadline_monotonic'],store.cap,reserve=0)
                if code:raise RuntimeError('Terminal worker failed; inspect finish-worker log')
                break
            except BaseException as exc:
                failure=type(exc).__name__+': '+str(exc)
                if rescue:raise RuntimeError('Bounded failure finalization failed; preserve record without reentry: '+failure)
            finally:reap(finisher)
        elapsed=time.monotonic()-started
        # The final small receipt is within the already charged close. Immutable
        # raw-file hashes exclude the integrity manifest itself by contract.
        closing=read_json(relative(root,'closing-deadline.json'))['deadline_monotonic']
        if time.monotonic()>=closing:raise RuntimeError('Terminal receipt close expired; preserve record for audit')
        import signal
        def receipt_timeout(signum,frame):raise TimeoutError('Bounded receipt publication')
        old_handler=signal.signal(signal.SIGALRM,receipt_timeout)
        signal.setitimer(signal.ITIMER_REAL,max(.001,closing-time.monotonic()))
        try:
            final_map=relative(root,'files.sha256.json')
            if final_map.exists():
                value=read_json(final_map);value['processing_receipt']['actual_parent_observation']=elapsed
                value['processing_receipt']['measurement_boundary']='through supervised terminal worker exit; receipt serialization bounded by charged close'
                terminal_store.json('files.sha256.json',value,replace=True)
            elif not offline:
                current=state(root)
                if current.get('pending_close_nonce')!=nonce or current.get('active_command')!=args.command:
                    raise ValueError('Terminal close nonce/stage mismatch')
                current['command_durations'][args.command]=elapsed
                current['parent_close_receipt']={'command':args.command,'nonce':nonce,'observed_seconds':elapsed}
                current.pop('pending_close_nonce');current['active_command']=None
                terminal_store.json('status.json',current,replace=True)
        finally:
            signal.setitimer(signal.ITIMER_REAL,0);signal.signal(signal.SIGALRM,old_handler)
        if failure:raise RuntimeError(failure)
        print(json.dumps({'command':args.command,'directory':str(root),'elapsed_processing_seconds':elapsed,'state':'completed'}))


def main():
    signal.pthread_sigmask(signal.SIG_UNBLOCK,{signal.SIGTERM,signal.SIGINT,signal.SIGALRM})
    parser=argparse.ArgumentParser(description=__doc__);sub=parser.add_subparsers(dest='command',required=True)
    for command in ('prepare','verify','demonstrate','analyze','export'):
        p=sub.add_parser(command);configure_parser(p);p.add_argument('--internal-worker',action='store_true',help=argparse.SUPPRESS);p.add_argument('--internal-finish',action='store_true',help=argparse.SUPPRESS)
        if command=='prepare':
            for n in ('plan','lab17','lab19','out'):p.add_argument('--'+n,required=True,type=Path)
        else:
            p.add_argument('--run',required=True,type=Path)
            if command in ('verify','demonstrate'):p.add_argument('--model-dir',required=True,type=Path)
            else:p.add_argument('--out',required=True,type=Path)
            if command=='export':
                p.add_argument('--analysis',required=True,type=Path);p.add_argument('--allowlist',required=True,type=Path)
    args=parser.parse_args()
    if args.internal_worker or args.internal_finish:
        # A user invoking this flag cannot bypass the owned running-stage gate.
        current=(state(args.out if args.command=='prepare' else args.run) if args.command in ('prepare','verify','demonstrate')
                 else read_json(relative(args.out,'manifest.json')))
        phase_ok=(current.get('active_command')==args.command and current.get(args.command)=='running') if args.command in ('prepare','verify','demonstrate') else (current.get('stage')==args.command and current.get('state')=='running')
        if ((not phase_ok and not args.internal_finish) or current.get('supervisor_pid')!=os.getppid()
                or os.environ.get('LLMCOURSE_MICROSCOPE_PARENT_PID')!=str(os.getppid())
                or current.get('worker_nonce_sha256')!=digest(os.environ.get('LLMCOURSE_MICROSCOPE_NONCE','').encode())):
            raise ValueError('Live supervisor-owned active stage required')
        if args.internal_finish:finish(args)
        else:
            try:worker(args)
            except BaseException as exc:
                root=args.out if args.command in ('prepare','analyze','export') else args.run
                Store(root,5_000_000 if args.command=='analyze' else RAW_CAP,terminal=True).json('worker-error.json',{'type':type(exc).__name__,'message':str(exc)},replace=True)
                raise
    else:supervise(args)

if __name__=='__main__':main()
