diff --git a/roadscore/experiments/ace_chestnut_20260916/ace_worker.py b/roadscore/experiments/ace_chestnut_20260916/ace_worker.py index e3e38e57c3..0c46312b1e 100644 --- a/roadscore/experiments/ace_chestnut_20260916/ace_worker.py +++ b/roadscore/experiments/ace_chestnut_20260916/ace_worker.py @@ -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)/480008: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)/480008: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() diff --git a/roadscore/prototype/RESIDENT_CACHED_STARTUP.md b/roadscore/prototype/RESIDENT_CACHED_STARTUP.md new file mode 100644 index 0000000000..367813132b --- /dev/null +++ b/roadscore/prototype/RESIDENT_CACHED_STARTUP.md @@ -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. diff --git a/roadscore/prototype/app.py b/roadscore/prototype/app.py index b99de6a413..1ef9b01d89 100644 --- a/roadscore/prototype/app.py +++ b/roadscore/prototype/app.py @@ -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) diff --git a/roadscore/prototype/native_receiver.sh b/roadscore/prototype/native_receiver.sh index 7ec7cc8971..3eaa79a608 100755 --- a/roadscore/prototype/native_receiver.sh +++ b/roadscore/prototype/native_receiver.sh @@ -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 diff --git a/roadscore/prototype/prepared_session.py b/roadscore/prototype/prepared_session.py index 27aa782712..97a2882498 100644 --- a/roadscore/prototype/prepared_session.py +++ b/roadscore/prototype/prepared_session.py @@ -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) diff --git a/roadscore/prototype/test_resident_worker_reset.py b/roadscore/prototype/test_resident_worker_reset.py new file mode 100644 index 0000000000..37afa4a41f --- /dev/null +++ b/roadscore/prototype/test_resident_worker_reset.py @@ -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==[]