Add opt-in isolated cached-plan resident session reset

This commit is contained in:
firestar5683
2026-09-19 16:50:23 -07:00
parent 19b8cc33e3
commit a23492f627
6 changed files with 177 additions and 26 deletions
@@ -26,6 +26,11 @@ if composition_policy in ('hook-v2','hook-cache-v1') and not windowed:raise Valu
if composition_policy=='hook-cache-v1':
from cached_composition import CachedComposition,validate_bank
validate_bank(Path(os.environ['ROADSCORE_PLAN_BANK']),profile=selected())
resident=os.environ.get('ROADSCORE_RESIDENT')=='1'
if resident and composition_policy!='hook-cache-v1':raise ValueError('Resident mode requires local current conditioning')
if resident and (G/'session_request.json').exists():raise RuntimeError('Unacknowledged resident request requires owner inspection before restart')
from contextlib import nullcontext
from resident_session import session_lease,validate_session
from tinygrad import Device
profile=selected();preparation_id=f'{profile}_{time.time_ns()}'
lock=open(G/'gpu.lock','w');fcntl.flock(lock,fcntl.LOCK_EX|fcntl.LOCK_NB)
@@ -66,29 +71,56 @@ try:
save_wave(stem.with_suffix('.wav'),wave);np.save(stem.with_suffix('.npy'),latent);write_json(stem.with_suffix('.json'),row)
return {'directory':str(folder),'stem':stem.name,'output_gain':POLICY.output_gain}
return record
print('ACE_PREPARING',profile,flush=True)
initial=None;last=None;preparation=[];slot=0;first_accepted_audio_seconds=None
while initial is None or len(initial)/48000<POLICY.initial_buffer_seconds:
role=('initial' if windowed else 'verse') if initial is None else ('verse' if windowed else 'repaint_verse')
if planned:role=planned.begin(last)
wave,last_new,stats=qualified.run(role,sample_seed(base_seed,"prepare",slot),last,record=record_for('prepare_'+preparation_id+'_'+str(slot)))
preparation.append(stats)
if wave is None:raise RuntimeError('Preparation rejected after bounded quality retries; inspect generated/quality')
if planned:planned.accept(wave,last_new)
if initial is None:
first_accepted_audio_seconds=time.monotonic()-BOOT;initial=wave.copy()
else:
overlap=2*48000;prefix=round(stats['prefix_seconds']*48000);alpha=np.linspace(0,1,overlap)[:,None]
initial[-overlap:]=initial[-overlap:]*(1-alpha)+wave[prefix-overlap:prefix]*alpha
initial=np.concatenate([initial,wave[prefix:]])
last=last_new;slot+=1
write_json(G/'ace_worker_state.json',{'pid':os.getpid(),'generation_seed':base_seed,'composition_policy':composition_policy,'profile':profile,'phase':'preparing','accepted_chunks':slot,'accepted_buffer_seconds':len(initial)/48000,'first_accepted_audio_seconds':first_accepted_audio_seconds,'elapsed_seconds':time.monotonic()-BOOT})
if slot>8:raise RuntimeError('Initial buffer did not fill within bounded preparation')
save_wave(G/'ace_initial.wav',initial);np.save(G/'ace_initial.npy',last)
write_json(G/'ace_initial.json',{'generation_seed':base_seed,'composition_policy':composition_policy,'composer':'ace','startup_seconds':time.monotonic()-BOOT,'first_accepted_audio_seconds':first_accepted_audio_seconds,'model_load_seconds':model_load_seconds,'host_peak_rss_kib':resource.getrusage(resource.RUSAGE_SELF).ru_maxrss,'duration':len(initial)/48000,'source_identity':'kpop_control','prepared_profile':profile if windowed else 'legacy','preparation_id':preparation_id,'continuation_policy':'quality-gated fixed lookahead','generation':preparation,'prepared_identity':True,'output_gain':POLICY.output_gain,'created_wall':time.time()})
write_json(G/'ace_worker_state.json',{'pid':os.getpid(),'generation_seed':base_seed,'composition_policy':composition_policy,'profile':profile,'phase':'READY','initial_buffer_seconds':len(initial)/48000})
(G/'request.json').unlink(missing_ok=True);ready.write_text('ace');print('ACE_READY',flush=True)
def prepare_session(selection=None):
global profile,base_seed,preparation_id,planned,qualified
if selection is not None:
new_profile,new_seed,new_policy,bank_hash=validate_session(selection)
if new_policy!=composition_policy or bank_hash!=planned.bank_hash:raise ValueError('Resident conditioning identity changed')
if (G/'request.json').exists() or (G/'busy').exists():raise RuntimeError('Pending continuation prevents session reset')
next_id=f'{new_profile}_{time.time_ns()}'
next_plan=CachedComposition(Path(os.environ['ROADSCORE_PLAN_BANK']),G/'hook_sessions'/next_id,new_seed,new_profile)
if next_plan.bank_hash!=bank_hash:raise ValueError('Conditioning bank changed on disk')
ready.unlink(missing_ok=True)
profile,base_seed,preparation_id,planned=new_profile,new_seed,next_id,next_plan
qualified=QualifiedGenerator(generate,policy=HOOK_POLICY)
session_started=time.monotonic() if selection else BOOT
print('ACE_PREPARING',profile,flush=True)
initial=None;last=None;preparation=[];slot=0;first_accepted_audio_seconds=None
while initial is None or len(initial)/48000<POLICY.initial_buffer_seconds:
role=('initial' if windowed else 'verse') if initial is None else ('verse' if windowed else 'repaint_verse')
if planned:role=planned.begin(last)
wave,last_new,stats=qualified.run(role,sample_seed(base_seed,"prepare",slot),last,record=record_for('prepare_'+preparation_id+'_'+str(slot)))
preparation.append(stats)
if wave is None:raise RuntimeError('Preparation rejected after bounded quality retries; inspect generated/quality')
if planned:planned.accept(wave,last_new)
if initial is None:
first_accepted_audio_seconds=time.monotonic()-session_started;initial=wave.copy()
else:
overlap=2*48000;prefix=round(stats['prefix_seconds']*48000);alpha=np.linspace(0,1,overlap)[:,None]
initial[-overlap:]=initial[-overlap:]*(1-alpha)+wave[prefix-overlap:prefix]*alpha
initial=np.concatenate([initial,wave[prefix:]])
last=last_new;slot+=1
write_json(G/'ace_worker_state.json',{'pid':os.getpid(),'generation_seed':base_seed,'composition_policy':composition_policy,'profile':profile,'phase':'preparing','accepted_chunks':slot,'accepted_buffer_seconds':len(initial)/48000,'first_accepted_audio_seconds':first_accepted_audio_seconds,'elapsed_seconds':time.monotonic()-session_started})
if slot>8:raise RuntimeError('Initial buffer did not fill within bounded preparation')
save_wave(G/'ace_initial.wav',initial);np.save(G/'ace_initial.npy',last)
write_json(G/'ace_initial.json',{'generation_seed':base_seed,'composition_policy':composition_policy,'composer':'ace','startup_seconds':time.monotonic()-session_started,'first_accepted_audio_seconds':first_accepted_audio_seconds,'model_load_seconds':0. if selection else model_load_seconds,'resident_model_load_seconds':model_load_seconds,'worker_uptime_seconds':time.monotonic()-BOOT,'resident_reused':bool(selection),'resident_capable':resident,'resident_model_identity':str(id(c)),'conditioning_bank_sha256':planned.bank_hash if resident else None,'host_peak_rss_kib':resource.getrusage(resource.RUSAGE_SELF).ru_maxrss,'duration':len(initial)/48000,'source_identity':'kpop_control','prepared_profile':profile if windowed else 'legacy','preparation_id':preparation_id,'continuation_policy':'quality-gated fixed lookahead','generation':preparation,'prepared_identity':True,'output_gain':POLICY.output_gain,'created_wall':time.time()})
write_json(G/'ace_worker_state.json',{'pid':os.getpid(),'generation_seed':base_seed,'composition_policy':composition_policy,'profile':profile,'phase':'READY','initial_buffer_seconds':len(initial)/48000})
(G/'request.json').unlink(missing_ok=True);ready.write_text('ace');print('ACE_READY',flush=True)
with session_lease(G) if resident else nullcontext():prepare_session()
while True:
session_request=G/'session_request.json'
if resident and session_request.exists():
selection=json.loads(session_request.read_text());session_request.unlink()
try:
with session_lease(G):prepare_session(selection)
write_json(G/'session_result.json',{**selection,'phase':'READY','preparation_id':preparation_id})
except BlockingIOError:
write_json(G/'session_result.json',{'id':selection.get('id'),'error':'Playback owns the resident session'})
except Exception as error:
write_json(G/'session_result.json',{'id':selection.get('id'),'error':str(error)})
ready.unlink(missing_ok=True)
raise
continue
request=G/'request.json'
if not request.exists():time.sleep(.1);continue
req=json.loads(request.read_text());request.unlink()
@@ -0,0 +1,11 @@
# Opt-in current-bank resident service
This candidate is OFF by default and has not passed native cold/warm audio equivalence. The normal baseline still owns and stops its worker.
After baseline evidence is protected, the sole hardware owner can validate with `ROADSCORE_RESIDENT=1 ./onroad --roadscore route1 --muted`. The first launch still loads models and compiles. After a successful preparation the power supervisor remains alive; its PID and process-start ticks are recorded in `generated/resident_owner.json`. Stop that exact verified owner with SIGTERM to restore its CPU settings and stop its GPU child. Do not kill unrelated workers.
Each subsequent launch chooses its normal fresh seed and requests a token-bound idle preparation. The worker retains DiT/VAE/decoder objects and compiled graphs, but recreates CachedComposition, retry estimator, accepted-history map, and initial audio/latent state. Seed, profile, policy and conditioning-bank hash must all match the request acknowledgment. A lease prevents reset while audio playback owns the files. Same-seed reproduction still rebuilds the session; it never replays a consumed prepared initial automatically. No preparation runs alongside playback.
Measure actual first accepted audio, startup-to-READY, and replay start; the initial buffer remains 112 seconds until measured native throughput supports a smaller value. Use identical explicit seeds for the cold/warm equivalence check, comparing initial PCM/latent and quality decisions. Then verify a normal launch logs a different seed. Preserve the preceding archive before switching sessions. No performance or audible equivalence claim follows from the CPU tests.
Do not hot-switch conditioning banks or reuse a timed-out request. Bank changes or outstanding requests fail closed and require the hardware owner to inspect preserved state. A worker exception ends its power supervisor, which restores CPU settings. This is a local opt-in event service, not an installed boot daemon.
+5
View File
@@ -25,6 +25,11 @@ p=argparse.ArgumentParser();p.add_argument('--root',default='/data/roadscore');p
def terminate(sig,frame):raise KeyboardInterrupt
signal.signal(signal.SIGTERM,terminate)
root=Path(a.root)
if os.environ.get('ROADSCORE_RESIDENT')=='1':
import atexit
from resident_session import session_lease
playback_lease=session_lease(root/'generated');playback_lease.__enter__()
atexit.register(playback_lease.__exit__,None,None,None)
from privacy_guard import install
install(root)
run=root/'results/current';run.mkdir(parents=True,exist_ok=True)
+16 -3
View File
@@ -9,10 +9,12 @@ flock -n 9 || { echo "Another native RoadScore session owns this bench"; exit 1;
power_pid=""
audio_pid=""
bridge_pid=""
worker_reused=1
resident_keep=0
cleanup() {
if [ -n "$bridge_pid" ]; then kill "$bridge_pid" 2>/dev/null || true; wait "$bridge_pid" 2>/dev/null || true; fi
if [ -n "$audio_pid" ]; then kill "$audio_pid" 2>/dev/null || true; wait "$audio_pid" 2>/dev/null || true; fi
if [ -n "$power_pid" ]; then kill "$power_pid" 2>/dev/null || true; wait "$power_pid" 2>/dev/null || true; fi
if [ -n "$power_pid" ] && [ "$resident_keep" != 1 ]; then kill "$power_pid" 2>/dev/null || true; wait "$power_pid" 2>/dev/null || true; fi
}
trap cleanup EXIT
trap 'exit 130' HUP INT TERM
@@ -22,9 +24,13 @@ export ROADSCORE_WORKER="$worker_script"
preparation_wait=360
if [ "${ROADSCORE_COMPOSER:-ace}" = ace ]; then preparation_wait=1500; fi
if ! pgrep -f "^/data/sa3-feasibility/venv/bin/python -u ${worker_script}$" >/dev/null; then
worker_reused=0
rm -f generated/worker_ready
env -u OPENPILOT_PREFIX /usr/local/venv/bin/python -u prototype/power_worker.py > results/native_worker.log 2>&1 < /dev/null &
setsid env -u OPENPILOT_PREFIX /usr/local/venv/bin/python -u prototype/power_worker.py > results/native_worker.log 2>&1 < /dev/null &
power_pid=$!
if [ "${ROADSCORE_RESIDENT:-0}" = 1 ]; then
/usr/local/venv/bin/python -c 'import json,sys,time; from pathlib import Path; p=int(sys.argv[1]); Path("generated/resident_owner.json").write_text(json.dumps({"power_worker_pid":p,"process_start_ticks":Path(f"/proc/{p}/stat").read_text().split()[21],"created_wall":time.time(),"stop":"SIGTERM power_worker_pid; it restores CPU and stops its child"}))' "$power_pid"
fi
fi
for attempt in $(seq 1 "$preparation_wait"); do
[ -f generated/worker_ready ] && break
@@ -33,7 +39,14 @@ for attempt in $(seq 1 "$preparation_wait"); do
sleep 1
done
[ -f generated/worker_ready ] || { echo 'Worker preparation timed out'; exit 1; }
if [ "${ROADSCORE_COMPOSER:-ace}" = ace ]; then /usr/local/venv/bin/python prototype/prepared_session.py; fi
if [ "${ROADSCORE_COMPOSER:-ace}" = ace ]; then
if [ "${ROADSCORE_RESIDENT:-0}" = 1 ] && [ "$worker_reused" = 1 ]; then
ROADSCORE_PREPARE_RESIDENT=1 /usr/local/venv/bin/python prototype/prepared_session.py
else
/usr/local/venv/bin/python prototype/prepared_session.py
fi
fi
if [ "${ROADSCORE_RESIDENT:-0}" = 1 ]; then resident_keep=1; fi
export OPENPILOT_PREFIX=roadscore_native
export PYTHONPATH=/data/openpilot:/data/roadscore/prototype:/data/roadscore-feasibility/venv/lib/python3.12/site-packages
mkdir -p /dev/shm/msgq_roadscore_native
+12 -1
View File
@@ -9,4 +9,15 @@ def verify(metadata,environ):
if metadata.get('prepared_profile')!=environ.get('ROADSCORE_ACE_PROFILE','prism'):raise ValueError('Resident composer profile differs from this launch')
if metadata.get('composition_policy','prepared-v1')!=environ.get('ROADSCORE_COMPOSITION_POLICY','prepared-v1'):raise ValueError('Resident composer policy differs from this launch')
if __name__=='__main__':verify(json.loads(Path('/data/roadscore/generated/ace_initial.json').read_text()),os.environ)
if __name__=='__main__':
generated=Path('/data/roadscore/generated')
metadata=json.loads((generated/'ace_initial.json').read_text())
if os.environ.get('ROADSCORE_PREPARE_RESIDENT')=='1':
if not metadata.get('resident_capable'):raise RuntimeError('Existing worker does not support resident handoff')
from cached_composition import validate_bank,digest
from resident_session import request_preparation
bank=Path(os.environ['ROADSCORE_PLAN_BANK']);profile=os.environ.get('ROADSCORE_ACE_PROFILE','prism')
validate_bank(bank,profile)
request_preparation(generated,profile,int(os.environ['ROADSCORE_GENERATION_SEED']),os.environ['ROADSCORE_COMPOSITION_POLICY'],digest(bank/'bank.json'))
metadata=json.loads((generated/'ace_initial.json').read_text())
verify(metadata,os.environ)
@@ -0,0 +1,79 @@
"""Exercise the worker's real reset function without importing or initializing a GPU."""
import ast
import os
from pathlib import Path
import resource
import time
from types import SimpleNamespace
import numpy as np
import pytest
from generation_seed import sample_seed
from resident_session import validate_session
BANK = 'a' * 64
class Plan:
def __init__(self, bank, root, seed, profile):
self.bank_hash = BANK
self.index = 0
def begin(self, previous):
assert previous is None and self.index == 0
return 'initial'
def accept(self, wave, latent):
self.index += 1
class Generator:
def __init__(self, generate, policy):
self.policy = policy
def run(self, role, seed, previous, record):
assert previous is None
rng = np.random.default_rng(seed)
return rng.standard_normal((480, 2)), rng.standard_normal((1, 4, 64)), {'prefix_seconds': 0}
@pytest.fixture
def worker(tmp_path, monkeypatch):
path=Path(__file__).resolve().parents[1]/'experiments/ace_chestnut_20260916/ace_worker.py'
function=next(n for n in ast.walk(ast.parse(path.read_text())) if isinstance(n, ast.FunctionDef) and n.name=='prepare_session')
code=compile(ast.fix_missing_locations(ast.Module(body=[function],type_ignores=[])),str(path),'exec')
saved={};audio=[];model=object();policy=SimpleNamespace(initial_buffer_seconds=.01,output_gain=.65)
monkeypatch.setenv('ROADSCORE_PLAN_BANK',str(tmp_path/'bank'))
state=dict(profile='prism',base_seed=123,preparation_id='cold',planned=Plan(None,None,123,'prism'),
qualified=Generator(None,policy),composition_policy='hook-cache-v1',validate_session=validate_session,
G=tmp_path,Path=Path,os=os,time=time,np=np,resource=resource,BOOT=time.monotonic(),
ready=tmp_path/'worker_ready',CachedComposition=Plan,QualifiedGenerator=Generator,
generate=None,HOOK_POLICY=policy,POLICY=policy,windowed=True,resident=True,c=model,
model_load_seconds=12.,record_for=lambda job:None,sample_seed=sample_seed,
write_json=lambda path,value:saved.update({path.name:value}),
save_wave=lambda path,wave:audio.append(wave.copy()))
exec(code,state)
return state,saved,audio,model
def selection(seed=123,bank=BANK):
return dict(profile='prism',generation_seed=seed,composition_policy='hook-cache-v1',bank_sha256=bank)
def test_cold_then_warm_same_seed_equal_and_new_seed_varies(worker):
state,saved,audio,model=worker
state['prepare_session']()
old_plan=state['planned'];old_gate=state['qualified']
state['prepare_session'](selection())
np.testing.assert_array_equal(audio[0],audio[1])
assert state['planned'] is not old_plan and state['qualified'] is not old_gate
assert state['c'] is model
assert saved['ace_initial.json']['resident_reused'] is True
assert saved['ace_initial.json']['model_load_seconds']==0
state['prepare_session'](selection(124))
assert not np.array_equal(audio[0],audio[2])
assert saved['ace_initial.json']['generation_seed']==124
def test_pending_continuation_cannot_be_erased_by_reset(worker):
state,saved,audio,model=worker
state['prepare_session']()
request=state['G']/'request.json';request.write_text('preserve')
old=state['planned']
with pytest.raises(RuntimeError,match='Pending continuation'):state['prepare_session'](selection(124))
assert request.read_text()=='preserve' and state['planned'] is old
def test_wrong_bank_does_not_change_existing_session(worker):
state,saved,audio,model=worker
with pytest.raises(ValueError,match='identity changed'):state['prepare_session'](selection(124,'b'*64))
assert state['base_seed']==123 and audio==[]