"""Bounded 24-attempt resource slice; reuse pinned public-30 single workers.
Run from a repository checkout. Raw DOM, screenshots and driver logs stay private.
Three concurrent workers target three different domains and meet a ready barrier.
"""
import argparse
import asyncio
from collections import Counter
from datetime import datetime, timezone
import hashlib
import json
import os
from pathlib import Path
import re
import shlex
import signal
import subprocess
import sys
import time
from types import SimpleNamespace
from urllib.parse import urlsplit

HERE = Path(__file__).resolve().parent
RUNNERS = {
 'patchright': ('patchright-public-30-20260929', '42dd3b229682a6152bd6fa0158dc38dd42c9e728abd0e26aa1ccc19906ab32bd'),
 'camoufox': ('camoufox-public-30-20260929', '76b6a8ce0dd3f56200484b16a67d6757d87c54c67ed8c1dbcbb9da79718884c7'),
}
POLICY = {'sample_interval_seconds':0.1,'min_host_available_mib':3072,
 'max_attempt_tree_pss_mib':4096,'max_aggregate_tree_pss_mib':6144,
 'ready_timeout_seconds':20,'worker_timeout_seconds':60,'slice_timeout_seconds':480,
 'group_gap_seconds':2,'application_retries':0,'repetitions':2,'max_concurrent_targets':3,
 'pss_scope':'Each isolated Python worker plus its driver/browser descendants, from process spawn through startup, ready barrier, navigation, settle, DOM capture, viewport screenshot and teardown. Coordinator excluded. Aggregate is the sum of unique live worker trees in the same sweep, not the sum of their separate peak readings.',
 'sampling_caveat':'PSS is Linux proportional set size in MiB (1024^2 bytes), shared with all processes on the host. Sweeps read processes sequentially; sweep duration and actual spacing are saved. Readings may miss brief peaks or processes. No RSS fallback. Concurrent shared pages can change per-worker attribution. Host is not an idle lab.'}

def utc():return datetime.now(timezone.utc).isoformat()
def sha(data):return hashlib.sha256(data).hexdigest()
def write(path,data):
 temporary=Path(str(path)+'.tmp');temporary.write_text(json.dumps(data,indent=2)+'\n');temporary.replace(path)
def memory():
 return {line.split(':')[0]:int(line.split()[1])/1024 for line in Path('/proc/meminfo').read_text().splitlines() if line.split(':')[0] in ('MemTotal','MemAvailable','SwapTotal','SwapFree')}
def birth(pid):
 try:return Path(f'/proc/{pid}/stat').read_text().rsplit(')',1)[1].split()[19]
 except (OSError,IndexError):return None
def tree(root):
 seen={};pending=[root]
 while pending:
  pid=pending.pop()
  if pid in seen:continue
  start=birth(pid)
  if start is None:continue
  seen[pid]=start
  try:
   for task in Path(f'/proc/{pid}/task').iterdir():pending.extend(int(p) for p in (task/'children').read_text().split())
  except (OSError,ValueError):pass
 return seen
def pss(members):
 total=read=0
 for pid,start in members.items():
  if birth(pid)!=start:continue
  try:
   match=re.search(r'^Pss:\s+(\d+)',Path(f'/proc/{pid}/smaps_rollup').read_text(),re.M)
   if match:total+=int(match.group(1));read+=1
  except OSError:pass
 return (round(total/1024,3) if read else None),read

def signal_owned(members,sig):
 for pid,start in reversed(list(members.items())):
  if birth(pid)==start:
   try:os.kill(pid,sig)
   except ProcessLookupError:pass

def direct_env(tool,cache):
 env=os.environ.copy()
 for key in list(env):
  if key.lower() in ('http_proxy','https_proxy','all_proxy'):env.pop(key)
 env['PYTHONDONTWRITEBYTECODE']='1'
 if tool=='camoufox':env['XDG_CACHE_HOME']=str(cache)
 return env

async def ready_wait(args,target):
 write(args.output/f"{target['id']}-ready.json",{'ready_at':utc(),'pid':os.getpid(),'state':'Fresh browser and page ready; no target navigation yet'})
 while not (args.output/'go.json').exists():await asyncio.sleep(0.02)

# Insert only a pre-navigation barrier. Browser launch, capture and validators
# remain the exact pinned single() implementation; hashes reject runner drift.
def load_worker(tool):
 folder,digest=RUNNERS[tool];path=HERE.parent/folder/'run.py';source=path.read_bytes()
 assert sha(source)==digest,'Frozen runner changed; stop before target requests'
 text=source.decode();marker='                row["navigation_started_at"] = utc()'
 assert text.count(marker)==1
 text=text.replace(marker,'                await ready_wait(args, target)\n'+marker)
 namespace={'__name__':'frozen_single','__file__':str(path),'ready_wait':ready_wait}
 exec(compile(text,str(path),'exec'),namespace)
 return namespace

def worker(args):
 from importlib.metadata import version
 assert version(args.worker)==('1.63.0' if args.worker=='patchright' else '0.5.6')
 if args.worker=='camoufox':assert version('playwright')=='1.62.0'
 module=load_worker(args.worker)
 target=next(t for t in json.loads(args.manifest.read_text())['targets'] if t['id']==args.target)
 asyncio.run(module['single'](args,target))

def group(args,tool,mode,rep,targets,order,run):
 if time.monotonic()-run['_clock']>POLICY['slice_timeout_seconds']:raise RuntimeError('Whole-slice time guard')
 ram=memory()
 if ram['MemAvailable']<POLICY['min_host_available_mib']:raise RuntimeError('Host has less than 3 GiB available before group')
 group_id=f'{order:02d}-{tool}-{mode}-r{rep}';out=args.output/group_id;out.mkdir(mode=0o700)
 g={'id':group_id,'order':order,'tool':tool,'mode':mode,'repetition':rep,'target_ids':[t['id'] for t in targets],
    'started_at':utc(),'host_ram_before':ram,'barrier_released_at':None,'attempts':[],'samples':[],'blocker':None}
 run['groups'].append(g);save(run,args.output)
 procs={};owned={};started={};logs={};handled=set();clock=time.monotonic();go_released=False
 try:
  for target in targets:
   tid=target['id'];logs[tid]=open(out/f'{tid}-driver.log','wb')
   command=[args.python_paths[tool],str(Path(__file__).resolve()),'--worker',tool,'--target',tid,'--binary',args.binaries[tool],
     '--manifest',str(args.manifest.resolve()),'--output',str(out.resolve())]
   procs[tid]=subprocess.Popen(command,stdout=logs[tid],stderr=logs[tid],env=direct_env(tool,args.camoufox_cache),start_new_session=True)
   owned[tid]={};started[tid]=time.monotonic()
  while len(handled)<len(procs):
   sweep_start=time.monotonic();union={};per={};read_counts={}
   for tid,proc in procs.items():
    current=tree(proc.pid) if proc.poll() is None else {pid:b for pid,b in owned[tid].items() if birth(pid)==b}
    owned[tid].update(current);union.update(current);reading,count=pss(current);per[tid]=reading;read_counts[tid]=count
   # Sum per-worker readings from this sweep; roots are independent and disjoint.
   aggregate=round(sum(v for v in per.values() if v is not None),3) if any(v is not None for v in per.values()) else None
   ram=memory();sample={'utc':utc(),'elapsed_seconds':round(sweep_start-clock,4),'sweep_seconds':round(time.monotonic()-sweep_start,4),
    'tree_pss_mib':per,'read_processes':read_counts,'aggregate_pss_mib':aggregate,'host_available_mib':round(ram['MemAvailable'],3)}
   g['samples'].append(sample)
   if ram['MemAvailable']<POLICY['min_host_available_mib'] or (aggregate or 0)>POLICY['max_aggregate_tree_pss_mib'] or any((v or 0)>POLICY['max_attempt_tree_pss_mib'] for v in per.values()):g['blocker']='unsafe_memory_guard'
   if time.monotonic()-run['_clock']>POLICY['slice_timeout_seconds']:g['blocker']='slice_timeout'
   if not go_released:
    if all((out/f'{tid}-ready.json').exists() for tid in procs):
     g['barrier_released_at']=utc();write(out/'go.json',{'released_at':g['barrier_released_at']});go_released=True
    elif any(proc.poll() is not None for proc in procs.values()):g['blocker']='setup_process_exited_before_barrier'
    elif time.monotonic()-clock>POLICY['ready_timeout_seconds']:g['blocker']='ready_barrier_timeout'
   for tid,proc in procs.items():
    if tid in handled:continue
    if time.monotonic()-started[tid]>POLICY['worker_timeout_seconds']:g['blocker']='worker_timeout'
    if g['blocker'] or proc.poll() is not None:
     owned[tid].update(tree(proc.pid))
     signal_owned(owned[tid],signal.SIGTERM);time.sleep(0.1);signal_owned(owned[tid],signal.SIGKILL)
     proc.wait(timeout=3);handled.add(tid)
   if g['blocker']:break
   time.sleep(POLICY['sample_interval_seconds'])
 finally:
  for tid,proc in procs.items():
   owned[tid].update(tree(proc.pid));signal_owned(owned[tid],signal.SIGTERM);time.sleep(0.05);signal_owned(owned[tid],signal.SIGKILL)
   proc.wait(timeout=3);logs[tid].close()
   row_path=out/f'{tid}.json';row=json.loads(row_path.read_text()) if row_path.exists() else {'target_id':tid,'status':None,'usable_content':False,'navigation_started_at':None}
   row.update(tool=tool,mode=mode,repetition=rep,group_id=group_id,worker_process_exit=proc.returncode,
    worker_ready_at=json.loads((out/f'{tid}-ready.json').read_text())['ready_at'] if (out/f'{tid}-ready.json').exists() else None,
    peak_sampled_tree_pss_mib=max((s['tree_pss_mib'][tid] for s in g['samples'] if s['tree_pss_mib'][tid] is not None),default=None),
    pss_sample_count=sum(s['tree_pss_mib'][tid] is not None for s in g['samples']),
    worker_lifecycle_seconds=round(time.monotonic()-started[tid],3),
    cleanup_surviving_owned_pids=[pid for pid,b in owned[tid].items() if birth(pid)==b])
   if g['blocker'] or proc.returncode!=0:row.update(classification=g['blocker'] or 'worker_process_error',usable_content=False,guard_error=g['blocker'] or f'exit {proc.returncode}')
   write(row_path,row);g['attempts'].append(row)
  g.update(finished_at=utc(),group_wall_seconds=round(time.monotonic()-clock,3),
   peak_sampled_aggregate_pss_mib=max((s['aggregate_pss_mib'] for s in g['samples'] if s['aggregate_pss_mib'] is not None),default=None),host_ram_after=memory())
  # Measure overlap of the actual navigation-to-DOM capture intervals, not launch count.
  events=[]
  for row in g['attempts']:
   if row.get('navigation_started_at') and row.get('latency_seconds') is not None:
    start=datetime.fromisoformat(row['navigation_started_at']).timestamp();events.extend([(start,1),(start+row['latency_seconds'],-1)])
  active=maximum=0
  for _,delta in sorted(events):active+=delta;maximum=max(maximum,active)
  g['max_overlapping_navigation_to_capture_intervals']=maximum
  save(run,args.output)
 print(group_id,' | '.join(f"{a['target_id']}: {a.get('status')} {a.get('classification')} PSS={a['peak_sampled_tree_pss_mib']}" for a in g['attempts']),f"aggregate={g['peak_sampled_aggregate_pss_mib']} overlap={maximum}",flush=True)
 if g['blocker']:raise RuntimeError(g['blocker'])
 if any(a['cleanup_surviving_owned_pids'] for a in g['attempts']):raise RuntimeError('Owned-process cleanup incomplete; stop before next group')

def save(run,output):write(output/'results.json',{k:v for k,v in run.items() if not k.startswith('_')})

def main(args):
 assert args.output.resolve()!=HERE and 'public' not in args.output.resolve().parts,'Raw captures must use a private output directory'
 subset=json.loads(args.manifest.read_text());assert len(subset['targets'])==len({urlsplit(t['url']).hostname for t in subset['targets']})==3
 baselines={tool:json.loads((HERE.parent/folder/'results.json').read_text()) for tool,(folder,_) in RUNNERS.items()}
 versions={}
 for tool,d in baselines.items():
  assert sha(Path(args.binaries[tool]).read_bytes())==d['binary_sha256'],'Installed binary changed'
  assert sha((HERE.parent/RUNNERS[tool][0]/'run.py').read_bytes())==RUNNERS[tool][1]
  env=direct_env(tool,args.camoufox_cache)
  command="from importlib.metadata import version;import json,platform;print(json.dumps({'tool':version('"+tool+"'),'python':platform.python_version()}))"
  v=json.loads(subprocess.check_output([args.python_paths[tool],'-c',command],env=env,text=True));assert v['tool']==d['tool_version']
  versions[tool]={'tool_version':v['tool'],'python_version':v['python'],'binary_sha256':d['binary_sha256'],'executable_version':d['executable_version'],
    'baseline_runner_sha256':RUNNERS[tool][1],'baseline_config':d['config'],'binary':args.binaries[tool],'python':args.python_paths[tool]}
 addon=args.camoufox_cache/'camoufox/addons/UBO/manifest.json';assert sha(addon.read_bytes())==baselines['camoufox']['default_addon']['manifest_sha256'],'Default cached addon changed'
 os.umask(0o077);args.output.mkdir(parents=True,exist_ok=False,mode=0o700)
 run={'run_id':args.output.name,'started_at':utc(),'manifest':subset,'manifest_sha256':sha(args.manifest.read_bytes()),'script_sha256':sha(Path(__file__).read_bytes()),
  'executed_command':shlex.join([sys.executable,*sys.argv]),'policy':POLICY,'versions':versions,'default_addon':baselines['camoufox']['default_addon'],
  'network':'Same local host, direct outbound. Proxy environment variables removed case-insensitively. Chromium --no-proxy-server; Firefox network.proxy.type=0. No login, CAPTCHA solving, paid service or proxy.',
  'host_ram_preflight':memory(),'groups':[],'_clock':time.monotonic()}
 # Mirror the complete block order in repetition 2; both tools lead one round.
 schedule=[(1,'solo','patchright'),(1,'solo','camoufox'),(1,'concurrent','patchright'),(1,'concurrent','camoufox'),
           (2,'concurrent','camoufox'),(2,'concurrent','patchright'),(2,'solo','camoufox'),(2,'solo','patchright')]
 run['schedule']=[{'repetition':r,'mode':m,'tool':t} for r,m,t in schedule];save(run,args.output)
 order=0
 try:
  for rep,mode,tool in schedule:
   targets=subset['targets'] if rep==1 else list(reversed(subset['targets']))
   chunks=[[t] for t in targets] if mode=='solo' else [targets]
   for chunk in chunks:
    if order:time.sleep(POLICY['group_gap_seconds'])
    order+=1;group(args,tool,mode,rep,chunk,order,run)
 except BaseException as exc:
  run['blocker']=f'{type(exc).__name__}: {exc}';run['finished_at']=utc();save(run,args.output);raise
 run['finished_at']=utc();save(run,args.output)
 print('Completed',sum(len(g['attempts']) for g in run['groups']),'attempts',flush=True)

if __name__=='__main__':
 parser=argparse.ArgumentParser(description=__doc__)
 parser.add_argument('--manifest',type=Path,default=HERE/'targets.json');parser.add_argument('--output',type=Path,required=True)
 parser.add_argument('--worker',choices=RUNNERS);parser.add_argument('--target');parser.add_argument('--binary',type=Path)
 parser.add_argument('--patchright-python');parser.add_argument('--camoufox-python');parser.add_argument('--patchright-binary');parser.add_argument('--camoufox-binary')
 parser.add_argument('--camoufox-cache',type=Path)
 args=parser.parse_args()
 if args.worker:worker(args)
 else:
  assert all([args.patchright_python,args.camoufox_python,args.patchright_binary,args.camoufox_binary,args.camoufox_cache]),'Supply both existing Python environments, binaries and Camoufox cache'
  args.python_paths={'patchright':args.patchright_python,'camoufox':args.camoufox_python};args.binaries={'patchright':args.patchright_binary,'camoufox':args.camoufox_binary};main(args)
