Gate short startup tests with verified resident buffer handoff

This commit is contained in:
firestar5683
2026-09-19 21:03:17 -07:00
parent 9c7582d73d
commit 270169eedb
10 changed files with 166 additions and 11 deletions
@@ -31,8 +31,8 @@ if resident and composition_policy!='hook-cache-v1':raise ValueError('Resident m
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 startup_buffer import initial_target
initial_buffer_target=initial_target(os.environ,composition_policy,POLICY.initial_buffer_seconds)
from startup_buffer import initial_target,session_target,SESSION_PROTOCOL
boot_initial_buffer_target=initial_buffer_target=initial_target(os.environ,composition_policy,POLICY.initial_buffer_seconds)
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)
@@ -74,9 +74,10 @@ try:
return {'directory':str(folder),'stem':stem.name,'output_gain':POLICY.output_gain}
return record
def prepare_session(selection=None):
global profile,base_seed,preparation_id,planned,qualified
global profile,base_seed,preparation_id,planned,qualified,initial_buffer_target
if selection is not None:
new_profile,new_seed,new_policy,bank_hash=validate_session(selection)
next_buffer_target=session_target(selection,new_policy,boot_initial_buffer_target)
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()}'
@@ -84,6 +85,7 @@ try:
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
initial_buffer_target=next_buffer_target
qualified=QualifiedGenerator(generate,policy=HOOK_POLICY)
session_started=time.monotonic() if selection else BOOT
write_json(G/'ace_worker_state.json',{'pid':os.getpid(),'generation_seed':base_seed,'composition_policy':composition_policy,'profile':profile,'phase':'preparing','accepted_chunks':0,'accepted_buffer_seconds':0.,'elapsed_seconds':0.,'resident_reused':bool(selection)})
@@ -106,7 +108,7 @@ try:
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,'initial_buffer_target_seconds':initial_buffer_target,'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_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_session_protocol':SESSION_PROTOCOL if resident else None,'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,'initial_buffer_target_seconds':initial_buffer_target,'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()
@@ -116,7 +118,7 @@ try:
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})
write_json(G/'session_result.json',{**selection,'phase':'READY','preparation_id':preparation_id,'applied_initial_buffer_target_seconds':initial_buffer_target})
except BlockingIOError:
write_json(G/'session_result.json',{'id':selection.get('id'),'error':'Playback owns the resident session'})
except Exception as error:
+22
View File
@@ -0,0 +1,22 @@
# Short-startup controlled candidate
Normal launches retain the 100-second target. Only `ROADSCORE_TEST_SHORT_STARTUP=1` selects target 40 unless explicitly overridden. Whole accepted chunks quantize that to approximately **54 seconds** (26 + 28), not 40. At the reported current hardware timings (first accepted ~17 seconds, continuation 24.5–25.55 seconds), expected preparation is roughly **42–44 seconds**, compared with observed 91.0–91.3 seconds for 110 seconds of accepted audio. That estimate excludes launcher overhead and quality retries; it is not a measured short-startup result.
The per-session target is bounded (35–112 in explicit test mode, 80–112 otherwise), transmitted using resident protocol 2, validated before reset, acknowledged as actually applied, and checked against the completed bundle. No quality threshold, seed rule, reroll limit, deadline guard, or hold behavior is weakened. Old workers cannot apply a new target without **one owner-controlled restart after current playback**. Same-target old-worker launches remain compatible. An old worker requesting a different target fails before enqueue; echoing unknown fields cannot masquerade as applying the change.
`startup_buffer_audit.py` models 20 jobs alternating the supplied 24.5/25.55-second observations, adding 28 seconds accepted music per success. Guard starts at 40 seconds, updates to max(30, observed*1.2)+10, request threshold remains90, accepted-tail hold adds24 seconds. The adverse sequence begins with rejected23-second work followed by a26-second retry if its original deadline permits. This is a **sensitivity simulation**, not a replay of complete hardware trace timestamps.
| Initial actual reserve | All accepted: minimum / holds | Early rejected job: minimum / holds / fresh jobs |
|---|---|---|
|82 seconds|57.50 / 0|33.00 / 0 / 20|
|54 seconds (40 target)|29.50 / 0|29.45 / 1 / 19|
|40 seconds hypothetical|20.00 / 1|20.00 / 2 / 19|
|30 seconds hypothetical|29.50 / 1|29.45 / 2 / 19|
All simulated underflow counts are zero **only under instantaneous viable holds and zero callback/scheduling jitter**. These are not hardware underflow predictions. The54-second case cannot admit the early second quality attempt inside the original deadline: it keeps accepted music through a reported hold, then resumes fresh work.30 seconds is deliberately not an allowed target. The existing adaptive GenerationBudget and inflight20-second hold reserve remain unchanged.
The last supplied real long-run evidence was20 jobs, no retries, minimum64 seconds. A newer observed rejection plus retry took49 seconds. These motivate measuring holds and minimum buffer during the controlled short-startup test; the prior no-retry64-second minimum is not proof the smaller reserve is safe.
## Additional launcher time
Source inspection finds no fixed26–34-second sleep. `normal_onroad.py` sequentially seeds route params/acquires display before starting the receiver; `native_receiver.sh` then runs manifest capture and app initialization after preparation, with at most roughly two1-second polling delays plus bridge subscription. In `app.py`, full-source `analyze_music` and full-source `assess_grid` run before `RhythmTimeline.add` assesses overlapping windows again; the standalone grid is immediately replaced by `timeline.at(0)`. Removing that duplicate while preserving required audit fields is a clear candidate. Actual stage timestamps or a labeled offline profile are required to assign the26–34 seconds quantitatively. No app or launcher edit is included here.
+1 -1
View File
@@ -15,5 +15,5 @@ def configure(session, replay, composer, *, native, transport_only, environ, roo
environ.pop('ROADSCORE_PLANNER_URL', None)
environ.pop('ROADSCORE_PLANNER_TOKEN', None)
environ.setdefault('ROADSCORE_RESIDENT', '1')
environ.setdefault('ROADSCORE_INITIAL_BUFFER_SECONDS', '100')
environ.setdefault('ROADSCORE_INITIAL_BUFFER_SECONDS', '40' if environ.get('ROADSCORE_TEST_SHORT_STARTUP') == '1' else '100')
return 'hook-cache-v1'
+14 -1
View File
@@ -2,14 +2,27 @@
import json
import os
from pathlib import Path
from startup_buffer import initial_target,require_session_protocol,SESSION_PROTOCOL
def verify(metadata,environ):
if 'ROADSCORE_INITIAL_BUFFER_SECONDS' in environ:
target=initial_target(environ,environ.get('ROADSCORE_COMPOSITION_POLICY','prepared-v1'))
if metadata.get('initial_buffer_target_seconds')!=target:raise ValueError('Prepared audio buffer target differs from requested session')
if metadata.get('resident_capable') and environ.get('ROADSCORE_RESIDENT')!='1':raise ValueError('Resident service requires an opted-in fresh handoff and playback lease')
if int(metadata.get('generation_seed',-1))!=int(environ['ROADSCORE_GENERATION_SEED']):raise ValueError('Resident composer belongs to another seed; its owner must finish/release that session')
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')
def requested_buffer(metadata,environ):
if 'ROADSCORE_INITIAL_BUFFER_SECONDS' not in environ:return None
target=initial_target(environ,environ['ROADSCORE_COMPOSITION_POLICY'])
if metadata.get('resident_session_protocol',0)>=SESSION_PROTOCOL:return target
if metadata.get('initial_buffer_target_seconds')==target:return None
require_session_protocol(metadata)
return target
if __name__=='__main__':
generated=Path('/data/roadscore/generated')
metadata=json.loads((generated/'ace_initial.json').read_text())
@@ -19,6 +32,6 @@ if __name__=='__main__':
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'))
request_preparation(generated,profile,int(os.environ['ROADSCORE_GENERATION_SEED']),os.environ['ROADSCORE_COMPOSITION_POLICY'],digest(bank/'bank.json'),initial_buffer_seconds=requested_buffer(metadata,os.environ),short_startup_test=os.environ.get('ROADSCORE_TEST_SHORT_STARTUP')=='1')
metadata=json.loads((generated/'ace_initial.json').read_text())
verify(metadata,os.environ)
+11 -1
View File
@@ -7,6 +7,7 @@ from pathlib import Path
import re
import time
import uuid
from startup_buffer import session_target, require_session_protocol, SESSION_PROTOCOL
@contextmanager
@@ -40,13 +41,20 @@ def validate_session(request):
raise ValueError('Resident handoff requires hook-cache-v1')
if not isinstance(bank, str) or re.fullmatch('[0-9a-f]{64}', bank) is None:
raise ValueError('Resident handoff requires the conditioning bank SHA-256')
session_target(request, policy)
return profile, seed, policy, bank
def request_preparation(generated, profile, seed, composition_policy, bank_sha256, timeout=1560):
def request_preparation(generated, profile, seed, composition_policy, bank_sha256, timeout=1560, *, initial_buffer_seconds=None, short_startup_test=False):
generated = Path(generated)
selection = {'profile': profile, 'generation_seed': seed,
'composition_policy': composition_policy, 'bank_sha256': bank_sha256}
if initial_buffer_seconds is not None:
if short_startup_test:selection['startup_policy']='short-startup-test-v1'
selection.update(session_protocol=SESSION_PROTOCOL, initial_buffer_target_seconds=initial_buffer_seconds)
# Old workers echo unknown request fields: capability must be checked before enqueue.
metadata = json.loads((generated / 'ace_initial.json').read_text())
require_session_protocol(metadata)
validate_session(selection)
if not math.isfinite(timeout) or timeout <= 0:
raise ValueError('Preparation timeout must be finite and positive')
@@ -73,6 +81,8 @@ def request_preparation(generated, profile, seed, composition_policy, bank_sha25
raise RuntimeError(str(result['error']))
if validate_session(result) != validate_session(selection):
raise RuntimeError('Resident acknowledgment belongs to a different session')
if initial_buffer_seconds is not None and result.get('applied_initial_buffer_target_seconds') != initial_buffer_seconds:
raise RuntimeError('Resident worker did not acknowledge the applied buffer target')
if result.get('phase') != 'READY' or not result.get('preparation_id'):
raise RuntimeError('Resident preparation did not report READY with provenance')
return result
+23 -2
View File
@@ -9,6 +9,27 @@ def initial_target(environ, composition_policy, default=112.):
if composition_policy != 'hook-cache-v1':
raise ValueError('Initial buffer override is only supported for local cached planning')
value = float(value)
if not math.isfinite(value) or not 80 <= value <= default:
raise ValueError('Validated candidate range is 80 to 112 seconds')
minimum = 35. if environ.get('ROADSCORE_TEST_SHORT_STARTUP') == '1' else 80.
if not math.isfinite(value) or not minimum <= value <= default:
raise ValueError('Startup reserve outside supported range (35–112 only for explicit short-startup test; otherwise 80–112)')
return value
SESSION_PROTOCOL = 2
def session_target(request, composition_policy, fallback=112.):
"""Bound a requested reserve without changing generation/quality policy."""
if 'initial_buffer_target_seconds' not in request:
return fallback
if request.get('session_protocol') != SESSION_PROTOCOL:
raise ValueError('Per-session buffer target requires resident protocol 2')
value = request['initial_buffer_target_seconds']
if isinstance(value, bool) or not isinstance(value, (int, float)):
raise ValueError('Session buffer target must be numeric seconds')
return initial_target({'ROADSCORE_INITIAL_BUFFER_SECONDS': value, 'ROADSCORE_TEST_SHORT_STARTUP': '1' if request.get('startup_policy') == 'short-startup-test-v1' else '0'}, composition_policy)
def require_session_protocol(metadata):
if metadata.get('resident_session_protocol', 0) < SESSION_PROTOCOL:
raise RuntimeError('Existing worker cannot change its buffer per session; its owner must restart it once with the new code after playback finishes')
@@ -0,0 +1,36 @@
"""Offline reserve sensitivity model; no audio, routes, GPU or device operations.
This models buffer arithmetic/guards, not callback timing or actual hold viability.
"""
import json
def simulate(initial, jobs=20, reject_first=False):
buffer=float(initial);minimum=buffer;holds=0;underflows=0;accepted=0;estimate=30.;elapsed=0.
for job in range(jobs):
if buffer>90:
elapsed+=buffer-90;buffer=90.
required=estimate+10
if buffer<required:
buffer+=24;holds+=1
deadline_buffer=buffer
durations=[23.,26.] if reject_first and job==0 else [25.55 if job%2 else 24.5]
spent=0.
for index,duration in enumerate(durations):
if deadline_buffer-spent<estimate+10:
if buffer<max(40,estimate+12):buffer+=24;holds+=1
break
# Runtime holds during in-flight work once the accepted reserve reaches 20s.
if buffer-duration<20:
before=max(0.,buffer-20);elapsed+=before;spent+=before;duration-=before;buffer-=before
minimum=min(minimum,buffer);buffer+=24;holds+=1
buffer-=duration;spent+=duration;elapsed+=duration;minimum=min(minimum,buffer)
underflows+=int(buffer<0)
estimate=max(30.,durations[index]*1.2)
if index==len(durations)-1:buffer+=28;accepted+=1
return dict(initial_seconds=initial,accepted_jobs=accepted,holds=holds,
underflows_model_only=underflows,minimum_buffer_seconds=round(minimum,2),elapsed_seconds=round(elapsed,2))
if __name__=='__main__':
print(json.dumps({scenario:[simulate(initial,reject_first=retry) for initial in (82,54,40,30)]
for scenario,retry in [('all_accepted',False),('first_job_rejected_retry_49s',True)]},indent=2))
@@ -60,3 +60,20 @@ class ResidentTests(unittest.TestCase):
self.assertEqual(original,(self.root/'session_request.json').read_bytes())
if __name__=='__main__':unittest.main()
class BufferProtocolTests(ResidentTests):
def test_old_worker_fails_before_queue(self):
(self.root/'ace_initial.json').write_text(json.dumps({'resident_capable':True}))
with self.assertRaisesRegex(RuntimeError,'restart it once'):
request_preparation(self.root,'prism',123,'hook-cache-v1','a'*64,initial_buffer_seconds=40,short_startup_test=True)
self.assertFalse((self.root/'session_request.json').exists())
def test_echo_without_applied_target_fails(self):
(self.root/'ace_initial.json').write_text(json.dumps({'resident_session_protocol':2}))
self.responder()
with self.assertRaisesRegex(RuntimeError,'applied buffer'):
request_preparation(self.root,'prism',123,'hook-cache-v1','a'*64,initial_buffer_seconds=40,short_startup_test=True)
def test_exact_applied_target_succeeds(self):
(self.root/'ace_initial.json').write_text(json.dumps({'resident_session_protocol':2}))
self.responder({'applied_initial_buffer_target_seconds':40})
result=request_preparation(self.root,'prism',123,'hook-cache-v1','a'*64,initial_buffer_seconds=40,short_startup_test=True)
self.assertEqual(result['initial_buffer_target_seconds'],40)
@@ -9,6 +9,7 @@ import numpy as np
import pytest
from generation_seed import sample_seed
from resident_session import validate_session
from startup_buffer import session_target,SESSION_PROTOCOL
BANK = 'a' * 64
@@ -42,7 +43,7 @@ def worker(tmp_path, monkeypatch):
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.,initial_buffer_target=.01,record_for=lambda job:None,sample_seed=sample_seed,
model_load_seconds=12.,initial_buffer_target=.01,boot_initial_buffer_target=.01,session_target=session_target,SESSION_PROTOCOL=SESSION_PROTOCOL,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)
@@ -84,3 +85,17 @@ def test_nonresident_launch_cannot_reuse_resident_audio_with_same_seed():
metadata={'generation_seed':123,'prepared_profile':'prism','composition_policy':'hook-cache-v1','resident_capable':True}
env={'ROADSCORE_GENERATION_SEED':'123','ROADSCORE_ACE_PROFILE':'prism','ROADSCORE_COMPOSITION_POLICY':'hook-cache-v1'}
with pytest.raises(ValueError,match='opted-in'):verify(metadata,env)
def test_explicit_session_target_replaces_boot_target(worker):
state,saved,audio,model=worker
class BufferGenerator(Generator):
def run(self,role,seed,previous,record):
assert previous is None
return np.zeros((40*48000,2),dtype=np.float32),np.zeros((1,4,64)),{'prefix_seconds':0}
state['QualifiedGenerator']=BufferGenerator
request={**selection(),'session_protocol':2,'initial_buffer_target_seconds':40,'startup_policy':'short-startup-test-v1'}
state['prepare_session'](request)
assert state['initial_buffer_target']==40
assert saved['ace_initial.json']['initial_buffer_target_seconds']==40
assert saved['ace_initial.json']['duration']==40
assert state['c'] is model
@@ -12,3 +12,22 @@ def test_unsupported_reserve_rejected(value):
def test_frozen_policy_unchanged_and_rejects_override():
assert initial_target({},'prepared-v1')==112
with pytest.raises(ValueError):initial_target({'ROADSCORE_INITIAL_BUFFER_SECONDS':'80'},'prepared-v1')
def test_short_startup_requires_explicit_test_policy():
from startup_buffer import session_target
with pytest.raises(ValueError):initial_target({'ROADSCORE_INITIAL_BUFFER_SECONDS':'40'},'hook-cache-v1')
assert initial_target({'ROADSCORE_INITIAL_BUFFER_SECONDS':'40','ROADSCORE_TEST_SHORT_STARTUP':'1'},'hook-cache-v1')==40
request={'session_protocol':2,'initial_buffer_target_seconds':40,'startup_policy':'short-startup-test-v1'}
assert session_target(request,'hook-cache-v1')==40
with pytest.raises(ValueError):session_target({**request,'initial_buffer_target_seconds':30},'hook-cache-v1')
with pytest.raises(ValueError):session_target({**request,'session_protocol':1},'hook-cache-v1')
def test_old_worker_same_target_compatible_changed_target_requires_restart():
from prepared_session import requested_buffer
metadata={'initial_buffer_target_seconds':100}
env={'ROADSCORE_INITIAL_BUFFER_SECONDS':'100','ROADSCORE_COMPOSITION_POLICY':'hook-cache-v1'}
assert requested_buffer(metadata,env) is None
with pytest.raises(RuntimeError,match='restart it once'):
requested_buffer(metadata,{**env,'ROADSCORE_INITIAL_BUFFER_SECONDS':'40','ROADSCORE_TEST_SHORT_STARTUP':'1'})
assert requested_buffer({**metadata,'resident_session_protocol':2},env)==100