Files
PaperClipAI/scripts/ci/pi-intel/owned_processes.py
DottaandPaperclip 05f5c9a4e0 fix(ci): bound diagnostic process inspection within phase deadlines
Reduce diagnostic polling pressure while retaining fatal inspection failures and exact process cleanup. Keep shared inspection defaults and original startup assertions unchanged.

Co-Authored-By: Paperclip <noreply@paperclip.ing>
2026-10-01 19:19:52 -05:00

254 lines
16 KiB
Python

"""Private qualification cleanup. Journals identities, never argv/environment.
Postgres is not an ordinary kill-tree leaf. The installed embedded-postgres
stop() uses SIGINT and awaits exit. Give a known wrapper sole graceful-signal ownership before native cleanup.
Unknown database/initializer identities never receive fallback TERM/KILL.
"""
import datetime, hashlib, json, os, pathlib, shlex, signal, stat, subprocess, threading, time
from cleanup_policy import (OWNER_GRACE_SECONDS,POSTGRES_GRACE_SECONDS,RESIDUAL_GRACE_SECONDS,KILL_GRACE_SECONDS,INSPECTION_RETRY_SECONDS,STOP_DEADLINE_SECONDS)
def now(): return datetime.datetime.now(datetime.timezone.utc).isoformat()
def atomic_json(path, value):
path=pathlib.Path(path); tmp=path.with_suffix('.tmp')
tmp.write_text(json.dumps(value,indent=2)+'\n'); tmp.chmod(0o600); tmp.replace(path)
def process_table(*, timeout=2):
result={}
for line in subprocess.check_output(['/bin/ps','-axo','pid=,ppid=,lstart=,stat=,comm='],text=True,timeout=timeout).splitlines():
row=line.split(maxsplit=8)
if len(row)!=9: continue
status={'R':'running','S':'sleeping','I':'idle','T':'stopped','Z':'zombie','U':'uninterruptible'}.get(row[7][:1],'unknown')
result[int(row[0])]={'parent':int(row[1]),'started':' '.join(row[2:7]),'status':status,'executable':row[8]}
return result
class ProcessInspectionUnavailable(Exception):
"""No trustworthy argv snapshot; caller must revalidate identity."""
def command_tokens(pid, *, timeout=2):
# Inspect only an already owned PID; never persist or print this output.
try:value=subprocess.check_output(['/bin/ps','-p',str(pid),'-o','command='],text=True,stderr=subprocess.DEVNULL,timeout=timeout)
except (subprocess.SubprocessError,OSError) as error:raise ProcessInspectionUnavailable() from error
if not value.strip() or len(value)>16384:raise ProcessInspectionUnavailable()
try:return shlex.split(value)
except ValueError as error:raise ProcessInspectionUnavailable() from error
class OwnedProcesses:
def __init__(self, journal, private_root, app_root, sidecar_path, *, table=process_table,
argv=command_tokens, send=os.kill, clock=time.monotonic, sleep=time.sleep,
app_grace=OWNER_GRACE_SECONDS, postgres_grace=POSTGRES_GRACE_SECONDS, root_grace=RESIDUAL_GRACE_SECONDS, kill_grace=KILL_GRACE_SECONDS):
self.journal=pathlib.Path(journal); self.root=pathlib.Path(private_root).resolve()
st=self.root.lstat()
if not stat.S_ISDIR(st.st_mode) or st.st_uid!=os.getuid() or st.st_mode&0o077: raise ValueError('private root ownership invalid')
self.root_identity=(st.st_dev,st.st_ino,st.st_uid)
self.app=str(pathlib.Path(app_root)); self.sidecar=str(pathlib.Path(sidecar_path))
self.table=table;self.argv=argv;self.send=send;self.clock=clock;self.sleep=sleep
self.budgets=(app_grace,postgres_grace,root_grace,kill_grace)
self.owned={};self.actions=[];self.root_pid=None;self.lock=threading.RLock();self.stopping=False;self.stop_result=None
self.postmasters={} # trusted paths are private in-memory only
self.stop_done=threading.Event();self.stop_thread=None;self.stop_deadline=None;self.inspection_failures=0
def _role(self,pid,row):
name=pathlib.Path(row['executable']).name
if name=='paperclip-runnerd': return 'runnerd'
if name=='copilot': return 'copilot'
if name=='initdb': return 'initdb'
if name=='postgres' or name.startswith('postgres:'):
tokens=self.argv(pid)
if '-D' in tokens:
i=tokens.index('-D')
if i+1<len(tokens):
data=pathlib.Path(tokens[i+1]).resolve()
if data.is_relative_to(self.root):
self.postmasters[pid]=data
return 'postgres_postmaster_candidate'
return 'postgres_backend'
if name in ('node','nodejs'):
tokens=self.argv(pid)
if self.sidecar in tokens:return 'acpx_sidecar'
if self.app+'/tests/runner-e2e/server-entry.ts' in tokens:return 'paperclip_server'
if self.app+'/tests/runner-e2e/server.ts' in tokens:return 'product_server_supervisor'
if self.app+'/tests/runner-e2e/launch.ts' in tokens:return 'product_launcher'
return 'node'
return {'sh':'shell','bash':'shell','zsh':'shell','python3':'supervisor','Python':'supervisor','Chromium':'browser'}.get(name,'other')
def observe(self,root_pid=None):
with self.lock:
if root_pid is not None and self.root_pid is None:self.root_pid=root_pid
rows=self.table();live={pid for pid,r in self.owned.items() if rows.get(pid,{}).get('started')==r['started']}
if self.root_pid not in self.owned and self.root_pid in rows: live.add(self.root_pid)
while True:
extra={pid for pid,r in rows.items() if r['parent'] in live}-live
if not extra:break
live|=extra
for pid in list(live):
if self.stop_deadline is not None and self.clock()>=self.stop_deadline:raise ProcessInspectionUnavailable()
row=rows[pid]
if self.stopping or pid not in self.owned or self.owned[pid]['role']=='unresolved_inspection':
try:role=self._role(pid,row)
except ProcessInspectionUnavailable:
current=self.table().get(pid)
if current is None or current['started']!=row['started']:
# Natural exit/reuse between snapshots is not a
# supervisor failure and never creates ownership.
live.discard(pid);rows.pop(pid,None);continue
role='unresolved_inspection'
if pid not in self.owned:self.owned[pid]={'parent':row['parent'],'started':row['started'],'role':role,'firstObservedAt':now(),'firstObservedMonotonic':self.clock()}
else:self.owned[pid]['role']=role
self.owned[pid].update(status=row['status'],lastObservedAt=now(),lastObservedMonotonic=self.clock())
for pid,r in self.owned.items():
if pid not in live:r['status']='gone_or_identity_changed'
self._persist()
return rows
def _persist(self):
with self.lock:atomic_json(self.journal,{'schema':'qualification-owned-processes/v2','recordedAt':now(),'rootPid':self.root_pid,'owned':self.owned,'actions':self.actions,'stopResult':self.stop_result,'rawArgvRetained':False,'rawEnvironmentRetained':False,'inspectionFailures':self.inspection_failures})
def live(self,rows=None):
rows=self.observe() if rows is None else rows
return {pid for pid,r in self.owned.items() if rows.get(pid,{}).get('started')==r['started'] and rows[pid].get('status')!='zombie'}
def _fresh(self):
end=min(self.clock()+INSPECTION_RETRY_SECONDS,self.stop_deadline or float('inf'))
while True:
try:
rows=self.observe()
live=self.live(rows)
if any(self.owned[pid]['role']=='unresolved_inspection' or rows[pid].get('status')=='unknown' for pid in live):raise ProcessInspectionUnavailable()
return rows
except (ProcessInspectionUnavailable,subprocess.SubprocessError,OSError):
self.inspection_failures+=1
if self.clock()>=end:raise ProcessInspectionUnavailable() from None
self.sleep(min(.1,end-self.clock()))
def _signal(self,pid,sig,reason):
if pid<=1 or pid==os.getpid():raise ProcessInspectionUnavailable()
rows=self._fresh();live=self.live(rows)
if pid not in live:return False
if reason=='sole_graceful_owner' and pid not in self._owners(live,rows):return False
role=self.owned[pid]['role']
# Reclassify at every signal boundary, including exec under the same PID.
database=self._postgres(live) or any(self.owned[p]['role']=='initdb' for p in live)
if (sig==signal.SIGKILL or reason=='post_owner_non_database_fallback') and database:raise ProcessInspectionUnavailable()
if role=='initdb' or (role.startswith('postgres_') and (sig!=signal.SIGINT or not self._verified_postmaster(pid))):raise ProcessInspectionUnavailable()
action={'at':now(),'pid':pid,'started':self.owned[pid]['started'],'role':role,'signal':signal.Signals(sig).name,'reason':reason,'delivery':'pending'}
with self.lock:self.actions.append(action);self._persist()
try:self.send(pid,sig)
except ProcessLookupError:
with self.lock:action['delivery']='not_delivered_process_absent';self._persist()
return False
with self.lock:action['delivery']='delivered';self._persist()
return True
def _postgres(self,live):return {pid for pid in live if self.owned[pid]['role'].startswith('postgres_')}
def _verified_postmaster(self,pid):
rows=self.observe(); r=self.owned.get(pid)
if not r or pid not in self.live(rows) or pid not in self.postmasters:return False
root_stat=self.root.lstat()
if (root_stat.st_dev,root_stat.st_ino,root_stat.st_uid)!=self.root_identity:return False
data=self.postmasters[pid];file=data/'postmaster.pid'
try:
fd=os.open(file,os.O_RDONLY|os.O_NOFOLLOW)
try:
st=os.fstat(fd)
if not stat.S_ISREG(st.st_mode) or st.st_uid!=os.getuid() or st.st_nlink!=1 or st.st_size>4096:return False
lines=os.read(fd,4096).decode().splitlines()
finally:os.close(fd)
if int(lines[0])!=pid or pathlib.Path(lines[1]).resolve()!=data:return False
started=time.mktime(time.strptime(r['started'],'%a %b %d %H:%M:%S %Y'))
if abs(int(lines[2])-started)>2:return False
return True
except (OSError,ValueError,IndexError):return False
def _wait(self,seconds,predicate):
end=min(self.clock()+seconds,self.stop_deadline or float('inf'))
while True:
rows=self._fresh();live=self.live(rows)
if predicate(live):return True
if self.stop_deadline is not None and self.clock()>=self.stop_deadline-1e-9:raise ProcessInspectionUnavailable()
if self.clock()>=end-1e-9:return False
self.sleep(min(.1,end-self.clock()))
def stop(self,root_pid=None):
thread=threading.get_ident()
with self.lock:
if self.stopping:owner=self.stop_thread==thread
else:
self.stopping=True;self.stop_thread=thread;owner=None
self.stop_deadline=self.clock()+STOP_DEADLINE_SECONDS
if owner is not None:
if not owner:self.stop_done.wait(timeout=STOP_DEADLINE_SECONDS+2)
return self.stop_result or {'status':'cleanup_in_progress','remainingOwnedPids':None,'postgresKilled':False}
try:
return self._stop(root_pid)
except (ProcessInspectionUnavailable,subprocess.SubprocessError,OSError):
return self._finish('unresolved_process_inspection',None)
finally:self.stop_done.set()
def _finish(self,status,live):
self.stop_result={'status':status,'remainingOwnedPids':sorted(live) if live is not None else None,
'lastObservedOwnedPids':sorted(pid for pid,r in self.owned.items() if r['status']!='gone_or_identity_changed'),
'inspectionComplete':live is not None,'postgresKilled':False,
'cleanupDurationSeconds':round(self.clock()-(self.stop_deadline-STOP_DEADLINE_SECONDS),3),
'cleanupDeadlineSeconds':STOP_DEADLINE_SECONDS}
if self.clock()>self.stop_deadline:self.stop_result['status']='unresolved_cleanup_deadline'
self._persist();return self.stop_result
def _owners(self,live,rows):
# Prefer a wrapper while a server/database is live. Its stopServer path
# owns the single signal to the real server. Never broadcast to providers.
wrappers={pid for pid in live if self.owned[pid]['role']=='product_server_supervisor'}
if len(wrappers)>1:
# tsx may expose both its launcher and inner Node worker as server.ts.
# Accept only one current, identity-validated ancestry chain; signal
# its deepest wrapper, which owns stopServer, never both wrappers.
def ancestors(pid):
seen=set()
while pid in live:
if pid in seen:raise ProcessInspectionUnavailable()
seen.add(pid)
current=rows.get(pid)
if not current or current['started']!=self.owned[pid]['started']:raise ProcessInspectionUnavailable()
pid=current['parent']
return seen
candidates={pid for pid in wrappers if wrappers.issubset(ancestors(pid))}
if len(candidates)!=1:raise ProcessInspectionUnavailable()
owner=next(iter(candidates))
if any(owner not in ancestors(pid) for pid in live if self.owned[pid]['role']=='paperclip_server'):
raise ProcessInspectionUnavailable()
return {owner}
if wrappers:return wrappers
servers={pid for pid in live if self.owned[pid]['role']=='paperclip_server'}
if servers:return servers
if self._postgres(live) or any(self.owned[p]['role']=='initdb' for p in live):return set()
if self.root_pid in live:return {self.root_pid}
return {pid for pid in live if self.owned[pid]['parent'] not in live}
def _stop(self,root_pid=None):
with self.lock:
if self.root_pid is None:self.root_pid=root_pid
app_grace,pg_grace,root_grace,kill_grace=self.budgets
live=self.live(self._fresh())
if not live:return self._finish('stopped',live)
delivered=set();selection_end=min(self.clock()+INSPECTION_RETRY_SECONDS,self.stop_deadline)
while not delivered:
rows=self._fresh();live=self.live(rows)
if not live:return self._finish('stopped',live)
owners=self._owners(live,rows)
if not owners:return self._finish('unresolved_database_owner',live)
for pid in sorted(owners):
if self._signal(pid,signal.SIGTERM,'sole_graceful_owner'):
delivered.add(pid);break
if not delivered:
if self.clock()>=selection_end:return self._finish('unresolved_graceful_owner_delivery',self.live(self._fresh()))
self.sleep(min(.1,selection_end-self.clock()))
owners=delivered # Never signal a nested owner after any successful delivery.
# Even an early root exit does not shorten this window for retained children.
if self._wait(app_grace,lambda live:not live):return self._finish('stopped',set())
live=self.live(self._fresh())
if any(self.owned[p]['role']=='initdb' for p in live):return self._finish('unresolved_database_initialization',live)
pg=self._postgres(live)
for pid in sorted(pg):
if self._verified_postmaster(pid):self._signal(pid,signal.SIGINT,'verified_postmaster_fast_shutdown')
if pg and not self._wait(pg_grace,lambda live:not self._postgres(live)):
return self._finish('unresolved_postgres_shutdown',self.live(self._fresh()))
live=self.live(self._fresh())
if any(self.owned[p]['role']=='initdb' for p in live):return self._finish('unresolved_database_initialization',live)
# Inner owner had its full50s. Only now request remaining non-DB exits.
for pid in sorted(live):
if pid in owners or self.owned[pid]['role'] in ('paperclip_server','product_server_supervisor','product_launcher'):continue
self._signal(pid,signal.SIGTERM,'post_owner_non_database_fallback')
if self._wait(root_grace,lambda live:not live):return self._finish('stopped',set())
live=self.live(self._fresh())
if self._postgres(live) or any(self.owned[p]['role']=='initdb' for p in live):return self._finish('unresolved_late_database',live)
for pid in sorted(live,reverse=True):self._signal(pid,signal.SIGKILL,'revalidated_non_database_fallback')
self._wait(kill_grace,lambda live:not live)
live=self.live(self._fresh())
return self._finish('stopped' if not live else 'unresolved_owned_processes',live)