"""#307 evidence-only isolated-home driver. Run under retained native job launcher."""
import argparse
import ctypes
from ctypes import wintypes
import datetime
import hashlib
import json
import os
from pathlib import Path
import re
import shutil
import subprocess
import sys
import threading
import time
import traceback

import psutil

BASE = Path(__file__).resolve().parent
PERCH_ENV = ('OWL_SESSION_ID', 'SPT_AGENT_ID', 'SPT_ENDPOINT_ID', 'SPT_SESSION_NAME',
             'SPT_SESSION_ID', 'SPT_ADAPTER', 'SPT_HOST_PID', 'SPT_INJECT_VERIFY_ECHO')
k32 = ctypes.WinDLL('kernel32', use_last_error=True)
k32.CreateFileW.argtypes = [wintypes.LPCWSTR, wintypes.DWORD, wintypes.DWORD, ctypes.c_void_p,
                           wintypes.DWORD, wintypes.DWORD, wintypes.HANDLE]
k32.CreateFileW.restype = wintypes.HANDLE
k32.GetNamedPipeServerProcessId.argtypes = [wintypes.HANDLE, ctypes.POINTER(wintypes.ULONG)]
k32.GetNamedPipeServerProcessId.restype = wintypes.BOOL
k32.CloseHandle.argtypes = [wintypes.HANDLE]
k32.GetProcessId.argtypes = [wintypes.HANDLE]
k32.GetProcessId.restype = wintypes.DWORD
INVALID_HANDLE = ctypes.c_void_p(-1).value


def now_ms():
    return time.time_ns() // 1000000


def sha(path):
    with Path(path).open('rb') as f:
        return hashlib.file_digest(f, 'sha256').hexdigest()


def save(path, data):
    Path(path).write_text(json.dumps(data, indent=2) + '\n', encoding='utf-8')


def pipe_pid(name):
    handle = k32.CreateFileW('\\\\.\\pipe\\' + name, 0xC0000000, 0, None, 3, 0x40000000, None)
    if handle == INVALID_HANDLE:
        return {'error': ctypes.get_last_error()}
    try:
        pid = wintypes.ULONG()
        if not k32.GetNamedPipeServerProcessId(handle, ctypes.byref(pid)):
            return {'error': ctypes.get_last_error()}
        return {'pid': pid.value}
    finally:
        k32.CloseHandle(handle)


class Rig:
    def __init__(self, args):
        self.args = args
        self.out = Path(args.output).resolve()
        self.out.mkdir(exist_ok=False)
        self.home = self.out / ('home-' + args.lane)
        self.started = time.monotonic()
        self.deadline = self.started + 300
        self.cleanup_deadline = None
        self.events = (self.out / 'events.jsonl').open('a', encoding='utf-8', buffering=1)
        self.lock = threading.Lock()
        self.broker = None
        self.rc = None
        self.read_allowed = threading.Event()
        self.read_allowed.set()
        self.reader = None
        self.trace_stop = threading.Event()
        self.trace_reader = None
        self.records = []
        self.owned = {}
        self.streams = []
        self.phase = 'admission'
        self.rc_number = 0
        self.cmd_number = 0
        self.backup = None
        self.restored = False
        self.ancestor_ids = {os.getpid(), *(p.pid for p in psutil.Process().parents())}
        self.preexisting_spt = {p.pid for p in psutil.process_iter(['name'])
                               if (p.info['name'] or '').lower() == 'spt.exe'}
        self.resident = Path(os.environ['LOCALAPPDATA']) / 'spt-core'
        assert not self.home.exists()
        assert self.resident.resolve() not in self.home.parents
        self.env = dict(os.environ)
        removed = []
        for key in list(self.env):
            if key.startswith(('SPT_', 'OWL_', 'CLAUDE_')):
                removed.append(key)
                del self.env[key]
        self.env.update(SPT_HOME=str(self.home), SPT_TEST_EPHEMERAL_ADVISORY_PORTS='1',
                        SPT_INSTALL_NO_FIREWALL='1', SPT_PUMP_TRACE='1', PYTHONUNBUFFERED='1')
        self.spt_source = Path(args.binary).resolve()
        self.pin = args.sha256
        assert sha(self.spt_source) == self.pin
        self.bin_dir = self.out / 'bin'
        self.bin_dir.mkdir()
        self.spt = self.bin_dir / 'spt.exe'
        shutil.copy2(self.spt_source, self.spt)
        assert sha(self.spt) == self.pin
        self.sockets = self.socket_names(self.env)
        resident_env = {**self.env, 'SPT_HOME': str(self.resident)}
        self.resident_sockets = self.socket_names(resident_env)
        assert not set(self.sockets) & set(self.resident_sockets)
        before = {name: pipe_pid(name) for name in self.sockets}
        assert all(row.get('error') == 2 for row in before.values()), before
        self.home.mkdir()
        # Valid default config, required so the cadence backup has exact original bytes.
        (self.home / 'daemon.json').write_bytes(b'{}\n')
        self.event('ADMITTED_HOME', home=str(self.home), explicit_env=self.env['SPT_HOME'],
                   sockets=self.sockets, resident_sockets=self.resident_sockets,
                   sockets_before=before, excluded_pids=sorted(self.preexisting_spt | self.ancestor_ids),
                   source=str(self.spt_source), executable=str(self.spt), sha256=self.pin,
                   removed_environment_names=removed, observation_bound_s=300, teardown_bound_s=120)

    def socket_names(self, env):
        r = subprocess.run([self.args.socket_helper], env=env, capture_output=True, text=True,
                           timeout=10)
        assert r.returncode == 0, r.stderr
        lines = r.stdout.splitlines()
        assert len(lines) == 3 and lines[0] == env['SPT_HOME'], lines
        assert lines[1].startswith('spt-daemon-broker-') and lines[2].startswith('spt-daemon-seed-')
        return lines[1:]

    def event(self, kind, **data):
        row = {'event': kind, 'at_ms': now_ms(), 'elapsed_s': time.monotonic() - self.started,
               'phase': self.phase, 'box': os.environ.get('COMPUTERNAME'),
               'route': 'hfenduleam-isolated-native-local', **data}
        with self.lock:
            self.events.write(json.dumps(row) + '\n')
        return row

    def remaining(self):
        return max(0, (self.cleanup_deadline or self.deadline) - time.monotonic())

    def identify(self, pid, role):
        proc = psutil.Process(pid)
        row = {'pid': pid, 'birth': proc.create_time(), 'exe': proc.exe(), 'ppid': proc.ppid(),
               'role': role, 'ancestors': [p.pid for p in proc.parents()]}
        assert pid not in self.preexisting_spt | self.ancestor_ids, row
        assert os.getpid() in row['ancestors'], row
        self.owned[pid] = row
        self.event('OWNED_PROCESS', **row)
        return row

    def guard(self, reason):
        assert self.broker is not None and self.broker.poll() is None
        expected = self.owned[self.broker.pid]
        proc = psutil.Process(self.broker.pid)
        assert proc.create_time() == expected['birth']
        image = Path(proc.exe()).resolve()
        # Same-byte apply may rename the still-running broker image to its owned aside.
        assert image.parent == self.bin_dir and sha(image) == self.pin
        assert k32.GetProcessId(int(self.broker._handle)) == self.broker.pid
        assert self.broker.pid not in self.preexisting_spt | self.ancestor_ids
        assert os.getpid() in [p.pid for p in proc.parents()]
        assert self.env['SPT_HOME'] == str(self.home)
        actual = self.socket_names(self.env)
        assert actual == self.sockets and not set(actual) & set(self.resident_sockets)
        endpoints = {name: pipe_pid(name) for name in actual}
        assert all(row.get('pid') == self.broker.pid for row in endpoints.values()), endpoints
        self.event('MUTATION_GUARD', reason=reason, home=self.env['SPT_HOME'],
                   broker_pid=self.broker.pid, birth=expected['birth'], endpoints=endpoints,
                   owned_handle_pid=k32.GetProcessId(int(self.broker._handle)))

    def command(self, args, mutate=False, timeout=20, env_extra=None, executable=None):
        if mutate:
            self.guard(' '.join(args))
        assert self.remaining() > 0, 'phase deadline exhausted before command'
        env = {**self.env, **(env_extra or {})}
        assert env['SPT_HOME'] == str(self.home)
        argv = [str(executable or self.spt), *args]
        # cmd parses its own command tail, not CRT argv quoting.
        launch = (subprocess.list2cmdline(argv[:3]) + ' ' + argv[3]
                  if Path(argv[0]).name.lower() == 'cmd.exe' else argv)
        self.cmd_number += 1
        num = self.cmd_number
        self.event('COMMAND_START', command=num, argv=argv, home=env['SPT_HOME'])
        began = time.monotonic()
        result = subprocess.run(launch, env=env, cwd=self.out, capture_output=True, text=True,
                                encoding='utf-8', errors='replace',
                                timeout=min(timeout, max(.1, self.remaining())))
        record = {'argv': argv, 'exit': result.returncode, 'stdout': result.stdout,
                  'stderr': result.stderr, 'elapsed_ms': (time.monotonic()-began)*1000,
                  'home': env['SPT_HOME']}
        save(self.out / f'command-{num:03}.json', record)
        self.event('COMMAND_END', command=num, exit=result.returncode,
                   elapsed_ms=record['elapsed_ms'])
        return result

    def start(self):
        self.phase = 'startup'
        assert sha(self.spt) == self.pin
        log = (self.out/'daemon.stderr.log').open('wb')
        self.streams.append(log)
        self.broker = subprocess.Popen([str(self.spt), 'daemon', 'run'], env=self.env,
                                      cwd=self.out, stdin=subprocess.DEVNULL,
                                      stdout=subprocess.DEVNULL, stderr=log,
                                      creationflags=subprocess.CREATE_NO_WINDOW)
        self.identify(self.broker.pid, 'broker')
        def trace_output():
            pending = ''
            sink = self.home/'logs/daemon.stderr.log'
            while not sink.exists() and not self.trace_stop.wait(.02):
                pass
            if self.trace_stop.is_set():
                return
            with sink.open('r',encoding='utf-8',errors='replace') as stream:
                while not self.trace_stop.is_set():
                    pending += stream.read(65536)
                    lines = pending.split('\n')
                    pending = lines.pop()
                    for line in lines:
                        if any(token in line for token in ('BRAIN_','CONTROLLER_SLOT_CLOSED','SUBSCRIBE_DECISION','ENDPOINT_RUN')):
                            self.event('DAEMON_TRACE',line=line)
                    self.trace_stop.wait(.02)
        self.trace_reader = threading.Thread(target=trace_output,daemon=True)
        self.trace_reader.start()
        end = min(self.deadline, time.monotonic() + 45)
        while time.monotonic() < end:
            if self.broker.poll() is not None:
                raise RuntimeError('isolated daemon exited during startup')
            ready = self.home/'brain.ready'
            try:
                state = json.loads(ready.read_text())
                pid = state['pid']
                brain = self.identify(pid, 'brain')
                assert self.broker.pid in brain['ancestors']
                self.guard('startup complete')
                save(self.out/'brain-ready-start.json', state)
                break
            except (OSError, ValueError, KeyError, psutil.Error, AssertionError):
                time.sleep(.2)
        else:
            raise RuntimeError('isolated daemon failed 45s startup admission')
        print('ISOLATED_DAEMON_READY', flush=True)

    def fixture(self):
        self.guard('create disposable adapter files')
        src = self.home/'srcs/evidencerig'
        src.mkdir(parents=True)
        drops = self.home/'drops'
        drops.mkdir()
        shutil.copy2(self.spt, src/'psychebin.exe')
        mode = 'hitch' if self.args.lane == 'w1c' else 'refresh'
        # The command template follows dummy_harness_e2e's existing contract.
        quote = lambda x: '"' + str(x).replace('\\', '/') + '"'
        cmd = (f'{quote(sys.executable)} -u {quote(BASE/"generator.py")} '
               f'--spt {quote(self.spt)} --home {quote(self.home)} '
               f'--id {{id}} --session-id {{session_id}} --mode {mode}')
        manifest = ('[adapter]\nname="evidencerig"\nkind="harness"\nversion="1"\n'
                    'min_spt_core_version="0"\n\n[session.self]\ncommand=\''+cmd+'\'\n\n'
                    '[session]\ncommune_dir="'+drops.as_posix()+'"\nsignoff_dir="'+drops.as_posix()+'"\n\n'
                    '[session.psyche_init]\ncommand=\'psychebin ready {id}\'\ncwd="{psyche_dir}"\nkeys=[]\n')
        (src/'manifest.toml').write_text(manifest, encoding='utf-8')
        for args in (['adapter','add',str(src)],
                     ['endpoint','create',self.args.lane+'rig','--adapter','evidencerig'],
                     ['endpoint','start',self.args.lane+'rig']):
            r = self.command(args, mutate=True, timeout=35)
            if r.returncode:
                raise RuntimeError('fixture command refused: '+r.stderr)
        self.attach()
        end = min(self.deadline, time.monotonic()+15)
        while time.monotonic() < end:
            if any(row['kind']=='PROG' for row in self.records):
                print('RIG_READY', flush=True)
                return
            time.sleep(.1)
        raise RuntimeError('real RC progress not observed within attach bound')

    def attach(self):
        self.guard('rc attach')
        self.rc_number += 1
        err = (self.out/f'rc-{self.rc_number}.stderr.log').open('wb')
        self.streams.append(err)
        self.rc = subprocess.Popen([str(self.spt),'rc',self.args.lane+'rig'], env=self.env,
                                   cwd=self.out, stdin=subprocess.PIPE, stdout=subprocess.PIPE,
                                   stderr=err, creationflags=subprocess.CREATE_NO_WINDOW)
        self.identify(self.rc.pid, 'rc-'+str(self.rc_number))
        self.read_allowed.set()
        child, number = self.rc, self.rc_number
        def read_output():
            acc = ''
            raw = (self.out/f'rc-{number}.chunks.jsonl').open('a',encoding='utf-8',buffering=1)
            parsed = (self.out/f'rc-{number}.progress.jsonl').open('a',encoding='utf-8',buffering=1)
            try:
                while True:
                    self.read_allowed.wait()
                    chunk = os.read(child.stdout.fileno(), 8192)
                    if not chunk:
                        break
                    received = now_ms()
                    text = chunk.decode('utf-8', errors='replace')
                    raw.write(json.dumps({'received_ms':received,'phase':self.phase,'text':text})+'\n')
                    acc += text
                    # Scan terminated semantic markers, preserving untouched raw chunks beside them.
                    pattern = r'(PROG (\d+) (\d{13}) ([0-9a-f]{8})|ACK (\S+) (\d+) (\d{13})|STATE (\d+) ([0-9a-f]{16}) (\d{13}))(?=[\r\n\x1b])'
                    matches = list(re.finditer(pattern,acc))
                    for m in matches:
                        fields = m.group(1).split()
                        row = {'kind':fields[0], 'fields':fields[1:], 'received_ms':received,
                               'phase':self.phase,'rc':number}
                        if fields[0]=='PROG':
                            row.update(counter=int(fields[1]),generated_ms=int(fields[2]),generation=fields[3],
                                       delivery_lag_ms=received-int(fields[2]))
                        elif fields[0]=='ACK':
                            row.update(tag=fields[1],ordinal=int(fields[2]),generated_ms=int(fields[3]))
                        else:
                            row.update(ordinal=int(fields[1]),chain=fields[2],generated_ms=int(fields[3]))
                        self.records.append(row)
                        if len(self.records) > 10000:
                            del self.records[:8000]
                        parsed.write(json.dumps(row)+'\n')
                    if matches:
                        acc = acc[matches[-1].end():]
                    if len(acc)>16384:
                        self.event('PARSER_UNMATCHED_TRUNCATION', characters=len(acc)-8192, rc=number)
                        acc=acc[-8192:]
            except Exception as e:
                self.event('RC_READER_ERROR',error=repr(e),rc=number)
            finally:
                raw.close()
                parsed.close()
        self.reader=threading.Thread(target=read_output,daemon=True)
        self.reader.start()

    def send(self, tag):
        self.guard('rc input '+tag)
        assert self.rc and self.rc.poll() is None
        self.event('INPUT_SENT',tag=tag,rc=self.rc_number)
        self.rc.stdin.write((tag+'\r').encode())
        self.rc.stdin.flush()

    def detach(self):
        if not self.rc or self.rc.poll() is not None:
            return
        self.guard('rc detach')
        self.read_allowed.set()
        self.rc.stdin.write(b'\x02d')
        self.rc.stdin.flush()
        try:
            self.rc.wait(timeout=min(8,max(.1,self.remaining())))
        except subprocess.TimeoutExpired:
            self.event('DETACH_NOT_COMPLETED_WITHIN_WINDOW',pid=self.rc.pid)
            return
        self.reader.join(timeout=2)
        self.event('RC_DETACHED',pid=self.rc.pid,exit=self.rc.returncode)

    def probe(self):
        # Preserve the operator's cmd /d /c diagnostic invocation, with pinned absolute binary.
        command = f'set SPT_RC_HITCH_DIAG=1&& "{self.spt}" node status --json'
        r = self.command(['/d','/c',command],executable=os.environ.get('COMSPEC','cmd.exe'),timeout=7)
        try:
            payload=json.loads(r.stdout)
        except ValueError:
            payload={'probe_error':'non-JSON output','stdout':r.stdout}
        valid=(r.returncode==0 and payload.get('broker_pid')==self.broker.pid
               and payload.get('net_enabled') is True
               and type(payload.get('net_canary_age_ms')) is int
               and type(payload.get('active_dial_tasks')) is int)
        self.event('HITCH_PROBE',valid=valid,exit=r.returncode,**payload)
        return valid,payload

    def observe(self, phase, seconds, probes=True):
        self.phase=phase
        begin=time.monotonic()
        end=min(self.deadline,begin+seconds)
        self.event('PHASE_START',requested_s=seconds)
        tick=0
        while time.monotonic()<end:
            if probes:
                self.probe()
            if tick%2==0 and self.rc and self.rc.poll() is None:
                self.send(f'TAG-{phase}-{tick:03}')
            tick+=1
            time.sleep(min(1,max(0,end-time.monotonic())))
        self.event('PHASE_END',requested_s=seconds,actual_s=time.monotonic()-begin,
                   truncated=end < begin+seconds)

    def hitch(self):
        valid,payload=self.probe()
        if not valid:
            raise RuntimeError('numeric net-enabled hitch instrument admission refused: '+str(payload))
        self.observe('baseline',45)
        self.guard('backup and set isolated cadence')
        config=self.home/'daemon.json'
        self.backup=config.read_bytes()
        (self.out/'daemon.json.original').write_bytes(self.backup)
        modified=json.loads(self.backup)
        for key in ('notif_pump_period_ms','registry_pump_period_ms','sync_pull_period_ms'):
            modified[key]=600000
        config.write_text(json.dumps(modified,indent=2)+'\n',encoding='utf-8')
        r=self.command(['node','refresh'],mutate=True)
        assert r.returncode==0,r.stderr
        self.observe('cadence-priming',10)
        self.observe('cadence-600000',120)
        self.restore()
        self.observe('restored',45)

    def restore(self):
        if self.backup is not None and not self.restored:
            self.guard('restore exact isolated daemon config bytes')
            config=self.home/'daemon.json'
            config.write_bytes(self.backup)
            assert config.read_bytes()==self.backup
            self.restored=True
            r=self.command(['node','refresh'],mutate=True)
            assert r.returncode==0,r.stderr
            self.event('CONFIG_RESTORED',sha256=sha(config),refresh_exit=r.returncode)

    def refresh_trial(self):
        self.observe('refresh-baseline',10,probes=False)
        self.read_allowed.clear()
        self.event('READER_PAUSED',rc=self.rc_number)
        self.observe('blocked-before-refresh',10,probes=False)
        baseline=[r for r in self.records if r['kind']=='PROG']
        save(self.out/'pre-refresh-progress.json',baseline[-1] if baseline else None)
        refresh=self.command(['daemon','refresh'],mutate=True,timeout=30)
        self.event('REAL_REFRESH_RESULT',exit=refresh.returncode,stdout=refresh.stdout,stderr=refresh.stderr)
        if refresh.returncode:
            raise RuntimeError('real refresh refused')
        self.observe('refresh-blocked',20,probes=False)
        self.read_allowed.set()
        self.event('READER_RESUMED',rc=self.rc_number)
        self.observe('refresh-recovery',35,probes=False)
        self.detach()
        if self.rc.poll() is None:
            raise RuntimeError('reattach arm refused: previous RC did not detach')
        self.attach()
        self.observe('reattached',10,probes=False)
        self.phase='signed-trial-admission'
        self.guard('generate isolated signing key')
        key='hertz-307-'+str(now_ms())
        keygen=subprocess.run([self.args.xtask,'debug-keygen',key],env=self.env,cwd=self.out,
                              capture_output=True,text=True,timeout=min(10,self.remaining()))
        if keygen.returncode:
            self.event('Q3_BLOCKED',reason='debug-keygen refused',exit=keygen.returncode)
            return
        values={}
        for name in ('key_id','public_hex','seed_hex'):
            found=re.findall(r'^'+name+r':\s*(\S+)\s*$',keygen.stdout,re.M)
            assert len(found)==1, 'ambiguous keygen output'
            values[name]=found[0]
        assert values['key_id']==key
        assert re.fullmatch('[0-9a-f]{64}',values['public_hex'])
        assert re.fullmatch('[0-9a-f]{64}',values['seed_hex'])
        echoed=re.findall(r'debug-pin --key-id '+re.escape(key)+r' --public-key ([0-9a-f]{64})',keygen.stdout)
        assert echoed==[values['public_hex']]
        self.event('SIGNING_KEY_GENERATED',key_id=key,public_hex=values['public_hex'],
                   secret_recorded=False,source='one native keygen invocation')
        pin=self.command(['debug-pin','--key-id',key,'--public-key',values['public_hex'],
                          '--home',str(self.home)],mutate=True,executable=self.args.xtask)
        if pin.returncode:
            self.event('Q3_BLOCKED',reason='isolated debug-pin refused',exit=pin.returncode)
            return
        assert not (self.home/'releases/applied-state.json').exists()
        assert not (self.home/'releases/applied.json').exists()
        assert sha(self.spt_source)==self.pin and sha(self.spt)==self.pin
        staged=self.command(['debug-rollout','--key-id',key,'--channel','debug','--version','1',
                             '--artifact','x86_64-pc-windows-msvc='+str(self.spt_source),
                             '--stage-dir',str(self.home/'releases'),
                             '--state',str(self.out/'debug-rollout-state.json')],
                            mutate=True,timeout=30,executable=self.args.xtask,
                            env_extra={'SPT_DEBUG_RELEASE_SEED':values.pop('seed_hex')})
        del keygen
        if staged.returncode:
            self.event('Q3_BLOCKED',reason='signed stage refused',exit=staged.returncode)
            return
        signed=json.loads((self.home/'releases/release.json').read_text())
        metadata=json.loads(signed['metadata_json'])
        save(self.out/'signed-update-set.json',signed)
        artifact=metadata['artifacts']['x86_64-pc-windows-msvc']
        assert artifact['artifact_sha256']==self.pin
        assert artifact['brain_ipc_version']==1 and artifact['broker_resource_abi']==1
        assert sha(self.home/'releases/artifacts/x86_64-pc-windows-msvc.bin')==self.pin
        self.event('SIGNED_TRIAL_ADMITTED',metadata=metadata,artifact_sha256=self.pin)
        self.phase='trial-blocked'
        self.read_allowed.clear()
        self.event('READER_PAUSED',rc=self.rc_number)
        self.observe('blocked-before-apply',10,probes=False)
        finished=threading.Event()
        def watch_trial():
            last={}
            with (self.out/'trial-records.jsonl').open('a',encoding='utf-8',buffering=1) as stream:
                while not finished.is_set():
                    for relative in ('releases/applied-state.json','releases/applied.json',
                                     'releases/last-outcome.json','brain.ready'):
                        try:
                            raw=(self.home/relative).read_text()
                            if raw!=last.get(relative):
                                parsed=json.loads(raw)
                                stream.write(json.dumps({'observed_ms':now_ms(),'path':relative,'record':parsed})+'\n')
                                last[relative]=raw
                        except (OSError,ValueError):
                            pass
                    finished.wait(.02)
        watcher=threading.Thread(target=watch_trial,daemon=True)
        watcher.start()
        try:
            apply=self.command(['update','apply'],mutate=True,timeout=30)
            self.event('SIGNED_APPLY_RESULT',exit=apply.returncode,stdout=apply.stdout,stderr=apply.stderr)
            if apply.returncode:
                self.event('Q3_BLOCKED',reason='real signed update apply refused',exit=apply.returncode,
                           stdout=apply.stdout,stderr=apply.stderr)
            self.observe('trial-blocked',20,probes=False)
            self.read_allowed.set()
            self.event('READER_RESUMED',rc=self.rc_number)
            self.observe('trial-recovery',55,probes=False)
        finally:
            finished.set()
            watcher.join(timeout=2)
            self.read_allowed.set()
        save(self.out/'post-trial-binary.json',{'sha256':sha(self.spt),'source_sha256':self.pin})

    def cleanup(self):
        self.phase='cleanup'
        self.cleanup_deadline=time.monotonic()+120
        self.event('CLEANUP_START')
        errors=[]
        self.read_allowed.set()
        sink = self.home/'logs/daemon.stderr.log'
        if sink.exists():
            shutil.copy2(sink,self.out/'daemon-before-cleanup.log')
        if self.broker and self.broker.poll() is None:
            for action in (self.restore,self.detach):
                try:
                    action()
                except Exception as e:
                    errors.append(repr(e))
            try:
                # Capture descendants while the owner is still alive, not after a broken ancestry hop.
                for child in psutil.Process(self.broker.pid).children(recursive=True):
                    try:
                        self.identify(child.pid,'teardown-descendant')
                    except psutil.NoSuchProcess:
                        pass
                r=self.command(['daemon','stop','--force'],mutate=True,timeout=45)
                self.event('STOP_RESULT',exit=r.returncode,stdout=r.stdout,stderr=r.stderr)
                self.broker.wait(timeout=min(15,max(.1,self.remaining())))
            except Exception as e:
                errors.append(repr(e))
        survivors=[]
        unknown=[]
        for pid,old in self.owned.items():
            try:
                p=psutil.Process(pid)
                if p.create_time()==old['birth'] and p.is_running():
                    survivors.append(old)
            except psutil.NoSuchProcess:
                pass
            except psutil.Error as e:
                unknown.append({'pid':pid,'error':str(e)})
        sockets={name:pipe_pid(name) for name in self.sockets}
        confirmed=(not survivors and not unknown and all(x.get('error')==2 for x in sockets.values()))
        row={'confirmed_before_job_close':confirmed,'survivors':survivors,'unknown':unknown,
             'sockets':sockets,'errors':errors,'retained_home':str(self.home)}
        save(self.out/'teardown.json',row)
        self.event('CLEANUP_END',**row)
        self.trace_stop.set()
        if self.trace_reader:
            self.trace_reader.join(timeout=2)
        if sink.exists():
            shutil.copy2(sink,self.out/'daemon-complete.log')
        for f in self.streams:
            f.close()
        self.events.close()
        return confirmed and not errors


def main():
    p=argparse.ArgumentParser()
    p.add_argument('--lane',choices=['w1c','w1b'],required=True)
    p.add_argument('--output',required=True)
    p.add_argument('--binary',required=True)
    p.add_argument('--sha256',required=True)
    p.add_argument('--socket-helper',required=True)
    p.add_argument('--xtask',required=True)
    args=p.parse_args()
    rig=None
    code=0
    try:
        rig=Rig(args)
        rig.start()
        rig.fixture()
        if args.lane=='w1c':
            rig.hitch()
        else:
            rig.refresh_trial()
    except Exception as e:
        code=1
        traceback.print_exc()
        if rig:
            rig.event('EXPERIMENT_ERROR',error=repr(e))
    finally:
        if rig and not rig.cleanup():
            code=2
    return code

if __name__=='__main__':
    raise SystemExit(main())
