Schedule native music cues on audible timestamps with fast atomic status relay

This commit is contained in:
firestar5683
2026-09-19 20:29:36 -07:00
parent c2da7800fa
commit 0db08c6c33
5 changed files with 103 additions and 8 deletions
+25
View File
@@ -0,0 +1,25 @@
"""Choose already-rendered musical cues at their scheduled audible time."""
import math
FIELDS = ('phase','kind','amount','activation','strength','section','next_section',
'gesture_active','gesture_queued','turn_signal_music','lead','predicted_peak','scheduled',
'signal_shaker','core_apex','alert_accent','engagement_presentation','motion_presentation')
REFERENCE = 'portaudio-dac-plus-residual-v1'
def audible_state(snapshot, now):
if snapshot.get('presentation_timing_reference') != REFERENCE:
return snapshot
result = {key:value for key,value in snapshot.items() if key not in FIELDS}
due = []
for item in snapshot.get('presentation_timeline', []):
if not isinstance(item,dict):continue
wall=item.get('audible_wall')
if type(wall) in (int,float) and math.isfinite(wall) and wall<=now and isinstance(item.get('cues'),dict):
due.append(item)
if due:
selected=max(due,key=lambda entry:entry['audible_wall'])
result.update({key:value for key,value in selected['cues'].items() if key in FIELDS})
result['presentation_display_lateness_ms']=max(0,(now-selected['audible_wall'])*1000)
result['presentation_display_sequence']=selected.get('sequence')
return result
+9 -4
View File
@@ -60,7 +60,9 @@ env['ROADSCORE_PRESENTATION_POLICY']=presentation['policy']
env['ROADSCORE_COMPOSITION_POLICY']=composition_policy
env['ROADSCORE_COMPOSER']=a.composer
env['ROADSCORE_ACE_PROFILE']=a.profile
if native:env.setdefault('ROADSCORE_REPLAY_PRIME','1')
if native:
env.setdefault('ROADSCORE_REPLAY_PRIME','1')
env['ROADSCORE_NATIVE_CUE_CLOCK']='1'
if native and not a.replay:env['ROADSCORE_AUDIO_DRAIN_FILE']=str(R/'results/current/audio_drained.json')
env['ROADSCORE_OVERLAY_CAPTURE']=str(out/'overlay.png')
env['ROADSCORE_STATUS_FILE']=str(out/'roadscore_status.json')
@@ -92,6 +94,7 @@ if a.replay:
args=[a.routeid,'--allow',services,'--start',str(a.start),'--no-loop','--headless','--cache','2']
if local:args+=['--data_dir',str(local)]
if native:args+=['--no-hw-decoder']
status_relay=None
display=None;children=[];named_children={};failure=None;logs=[];launch_started=time.monotonic()
print(('Preparing stored score replay; no generation. ' if a.replay else 'Preparing RoadScore; waiting for accepted audio. ')+('Speaker output enabled.' if a.audible else 'Muted capture.'),flush=True)
def launch(cmd,name,**kw):
@@ -148,6 +151,9 @@ try:
check_children(named_children,remote=not native,include_receiver=True)
if sender.poll() is not None or time.monotonic()>deadline:raise RuntimeError('Replay subscriber failed')
time.sleep(.1)
if native and not a.replay:
from status_relay import StatusRelay
status_relay=StatusRelay(R/'results/current/status.json',out/'roadscore_status.json');status_relay.start()
if a.replay:(out/'roadscore_status.json').write_text(json.dumps({'readiness':'READY','style':'Stored score','section':'ARCHIVED SCORE','compute':'none'}))
player=launch([str(replay),*args],'replay');started=time.monotonic();print('Normal replay running; '+('host speaker enabled' if a.audible else 'bench muted')+'. Output:',out,flush=True)
from replay_end import ReplayEnd
@@ -156,9 +162,6 @@ try:
while sender.poll() is None:
check_children(named_children,remote=not native,include_receiver=True,allow_clean_receiver=True)
if receiver is not None and receiver.poll() not in (None,0):raise RuntimeError('RoadScore receiver failed; see receiver.log')
if native and not a.replay:
try:(out/'roadscore_status.json').write_text((R/'results/current/status.json').read_text())
except OSError:pass
try:
native_state=json.loads(state_path.read_text())
if end_watch.observe(native_state,(out/'replay.log').read_text(),time.monotonic()):
@@ -196,6 +199,7 @@ try:
from score_archive import archive
archived=archive(a.routeid,out,a.start);print('Score archived:',archived,flush=True)
except Exception as error:
if status_relay:status_relay.close()
try:
check_children(named_children,remote=not native,include_receiver=True,allow_clean_receiver=True)
except Exception as child_error:
@@ -205,6 +209,7 @@ except Exception as error:
write_failure(out/'roadscore_status.json',failure)
raise error from None
finally:
if status_relay:status_relay.close()
for c in reversed(children):
if c.poll() is None:
os.killpg(c.pid,signal.SIGTERM)
+7 -4
View File
@@ -2,12 +2,13 @@
import json,time,os
from overlay_view import draw_panel, overlay_view, EventPresentation, hud_bounds, startup_bounds
from pathlib import Path
from cue_timing import audible_state
def install():
if os.environ.get('ROADSCORE_OVERLAY')!='1':return
import pyray as rl
from openpilot.system.ui.lib.application import gui_app,FontWeight
original=gui_app.render;last=0.;state={};frames=0;captured=False;capture_ready_since=None;captured_events=set();alert_clear_after=0.;native_nav_visible=False;home_footer_right=None;home_rect=None;status_mtime=None;status_read_error=None;event_presentation=EventPresentation();alert_seen_at=None;alert_captured=False
original=gui_app.render;last=0.;state={};raw_state={};frames=0;captured=False;capture_ready_since=None;captured_events=set();alert_clear_after=0.;native_nav_visible=False;home_footer_right=None;home_rect=None;status_mtime=None;status_read_error=None;event_presentation=EventPresentation();alert_seen_at=None;alert_captured=False
# Observe the actual native nav card; it keeps priority over this accessory.
from openpilot.selfdrive.ui.onroad.starpilot.navigation_card import NavigationCardRenderer
original_nav_render=NavigationCardRenderer._render
@@ -39,13 +40,14 @@ def install():
picture.save(target)
finally:rl.unload_image(image)
def draw():
nonlocal last,state,frames,captured,capture_ready_since,alert_clear_after,alert_seen_at,alert_captured,status_mtime,status_read_error
nonlocal last,state,raw_state,frames,captured,capture_ready_since,alert_clear_after,alert_seen_at,alert_captured,status_mtime,status_read_error
now=time.monotonic()
if now-last>.2:
if now-last>.04:
try:
state=json.loads(path.read_text());status_mtime=path.stat().st_mtime;status_read_error=None
raw_state=json.loads(path.read_text());status_mtime=path.stat().st_mtime;status_read_error=None
except (OSError,ValueError) as error:status_read_error=type(error).__name__
last=now
state=audible_state(raw_state,now) if os.environ.get('ROADSCORE_NATIVE_CUE_CLOCK')=='1' else raw_state
from openpilot.selfdrive.ui.ui_state import ui_state
# The native alert owns the display. Leave room for its existing fade-out too.
for service in ('selfdriveState','starpilotSelfdriveState'):
@@ -83,6 +85,7 @@ def install():
view=draw_panel(rl,gui_app.font(FontWeight.NORMAL),state,gui_app.width,gui_app.height,gui_app.font(FontWeight.SEMI_BOLD),presentation,startup=startup,footer_right=home_footer_right or 0)
ready=view['ready']
audit(True,None,view)
state=audible_state(raw_state,now) if os.environ.get('ROADSCORE_NATIVE_CUE_CLOCK')=='1' else raw_state
from openpilot.selfdrive.ui.ui_state import ui_state
if ready and ui_state.started and capture_ready_since is None:capture_ready_since=now
capture_target=None
+29
View File
@@ -0,0 +1,29 @@
"""Atomically relay current native status without waiting on replay supervision."""
import json
import threading
class StatusRelay:
def __init__(self, source, destination, interval=.05):
self.source,self.destination,self.interval=source,destination,interval
self.stop=threading.Event()
self.thread=threading.Thread(target=self.run,daemon=True)
def start(self):
self.thread.start()
def run(self):
previous=None
while not self.stop.is_set():
try:
payload=self.source.read_bytes()
if payload!=previous:
json.loads(payload)
temporary=self.destination.with_suffix('.relay.tmp')
temporary.write_bytes(payload);temporary.replace(self.destination)
previous=payload
except (OSError,ValueError):pass
self.stop.wait(self.interval)
def close(self):
self.stop.set();self.thread.join(timeout=1)
+33
View File
@@ -0,0 +1,33 @@
import json,tempfile,time,unittest
from pathlib import Path
from cue_timing import audible_state,REFERENCE
from status_relay import StatusRelay
class CueTests(unittest.TestCase):
def test_dac_plus_residual_and_live_health(self):
state={'presentation_timing_reference':REFERENCE,'phase':'wrong-early','buffered':44,'route_t':10,'readiness':'DEGRADED','presentation_timeline':[{'audible_wall':10.323,'sequence':1,'cues':{'phase':'stopped','buffered':0}},{'audible_wall':10.45,'sequence':2,'cues':{'phase':'moving'}}]}
self.assertNotIn('phase',audible_state(state,10.322))
shown=audible_state(state,10.324);self.assertEqual(shown['phase'],'stopped');self.assertEqual(shown['buffered'],44);self.assertEqual(shown['readiness'],'DEGRADED')
self.assertEqual(audible_state(state,10.451)['phase'],'moving')
self.assertEqual(state['phase'],'wrong-early')
def test_new_session_does_not_retain_previous_cue(self):
for session in ('old','new'):
state={'presentation_session_id':session,'presentation_timing_reference':REFERENCE,'presentation_timeline':[]}
self.assertNotIn('phase',audible_state(state,100))
def test_legacy_and_invalid_entries(self):
self.assertEqual(audible_state({'phase':'legacy'},3),{'phase':'legacy'})
state={'presentation_timing_reference':REFERENCE,'presentation_timeline':[None,{'audible_wall':float('nan'),'cues':{}},{'audible_wall':5,'cues':{}}]}
self.assertNotIn('phase',audible_state(state,3))
def test_relay_is_atomic_and_keeps_last_valid_status(self):
with tempfile.TemporaryDirectory() as folder:
folder=Path(folder);source=folder/'in.json';target=folder/'out.json';source.write_text('{"phase":"ready"}')
relay=StatusRelay(source,target,interval=.005);relay.start()
try:
deadline=time.monotonic()+1
while not target.exists() and time.monotonic()<deadline:time.sleep(.005)
self.assertEqual(json.loads(target.read_text())['phase'],'ready')
source.write_text('{broken');time.sleep(.03)
self.assertEqual(json.loads(target.read_text())['phase'],'ready')
finally:relay.close()
self.assertFalse(relay.thread.is_alive())
if __name__=='__main__':unittest.main()