From c20f6760d8ea5c1e72f6037ba30f5cc4ec44cae7 Mon Sep 17 00:00:00 2001 From: firestar5683 <168790843+firestar5683@users.noreply.github.com> Date: Sun, 20 Sep 2026 10:55:42 -0700 Subject: [PATCH] Add session-owned saved demo launch and playhead status in Galaxy --- roadscore/prototype/demo_catalog.py | 34 +++++ roadscore/prototype/demo_session.py | 140 ++++++++++++++++++ roadscore/prototype/test_demo_session.py | 73 +++++++++ starpilot/system/the_galaxy/roadscore.py | 38 ++++- .../the_galaxy/tests/test_roadscore_demo.py | 10 ++ 5 files changed, 292 insertions(+), 3 deletions(-) create mode 100644 roadscore/prototype/demo_catalog.py create mode 100644 roadscore/prototype/demo_session.py create mode 100644 roadscore/prototype/test_demo_session.py diff --git a/roadscore/prototype/demo_catalog.py b/roadscore/prototype/demo_catalog.py new file mode 100644 index 0000000000..310af9534a --- /dev/null +++ b/roadscore/prototype/demo_catalog.py @@ -0,0 +1,34 @@ +"""Local saved-demo registry. Entries identify original core audio and route assets.""" +import json +import re +from pathlib import Path + + +def entry(root, alias): + root = Path(root) + if not isinstance(alias, str) or not re.fullmatch(r'route[1-9][0-9]*', alias): + raise ValueError('Choose a registered demo route alias') + path = root / 'assets/demo_catalog.json' + try: + catalog = json.loads(path.read_text()) + except (OSError, ValueError) as error: + raise ValueError('Saved demo catalog is unavailable') from error + if catalog.get('version') != 1 or not isinstance(catalog.get('routes'), dict): + raise ValueError('Unsupported saved demo catalog') + item = catalog['routes'].get(alias) + if not isinstance(item, dict) or item.get('ready') is not True: + raise ValueError(f'{alias} has not been prepared as a saved demo') + if not isinstance(item.get('route'), str) or not item['route']: + raise ValueError('Saved demo is missing its route identity') + result = {**item, 'alias': alias} + for key in ('archive', 'curve_plan'): + value = item.get(key) + if key == 'curve_plan' and value is None: + continue + if not isinstance(value, str) or not value: + raise ValueError('Saved demo is missing its core archive') + asset = Path(value) + result[key] = str((asset if asset.is_absolute() else root / asset).resolve()) + if not Path(result[key]).exists(): + raise ValueError(f'Saved demo asset is missing: {key}') + return result diff --git a/roadscore/prototype/demo_session.py b/roadscore/prototype/demo_session.py new file mode 100644 index 0000000000..a12b19719d --- /dev/null +++ b/roadscore/prototype/demo_session.py @@ -0,0 +1,140 @@ +"""Galaxy-owned saved replay lifecycle; never starts a composer or a vehicle service.""" +import fcntl +import json +import math +import os +from pathlib import Path +import re +import signal +import subprocess +import threading +import time + +from demo_catalog import entry + + +def read(path): + try: + data = json.loads(Path(path).read_text()) + return data if isinstance(data, dict) else {} + except (OSError, ValueError): + return {} + + +def write(path, value): + path = Path(path) + temporary = path.with_suffix('.tmp') + temporary.write_text(json.dumps(value)) + temporary.replace(path) + + +def process_ticks(pid): + try: + fields = Path(f'/proc/{pid}/stat').read_text().rsplit(') ', 1)[1].split() + arguments = Path(f'/proc/{pid}/cmdline').read_bytes().split(b'\0') + if fields[0] == 'Z' or not any(argument.endswith(b'/native_prepared_showcase.py') for argument in arguments): + return None + return fields[19] + except (OSError, ValueError, IndexError): + return None + + +class DemoSession: + def __init__(self, root, offroad, *, spawn=subprocess.Popen, ticks=process_ticks, stop_group=None): + self.root = Path(root) + self.offroad, self.spawn, self.ticks = offroad, spawn, ticks + self.stop_group = stop_group or (lambda pid: os.killpg(pid, signal.SIGTERM)) + self.lock = threading.Lock() + self.owner_path = self.root / 'generated/demo_owner.json' + + def status(self): + owner = read(self.owner_path) + pid = owner.get('pid') + running = type(pid) is int and pid > 1 and self.ticks(pid) == owner.get('start_ticks') and owner.get('start_ticks') is not None + result = {key: owner.get(key) for key in ('request_id', 'alias', 'route', 'pid', 'out', 'muted')} + result.update(running=bool(running), generation_invoked=False) + if owner.get('out'): + out = Path(owner['out']) + ready = read(out/'demo_ready.json') + result['prepared'] = bool(running and ready.get('ready') is True and ready.get('session_id')) + result['ready_session_id'] = ready.get('session_id') if result['prepared'] else None + state = read(out/'status.json') + result['presentation_session_id'] = state.get('presentation_session_id') if running else None + result['failure'] = read(out/'failure.json').get('error') + return result + + def start(self, data): + if (not isinstance(data, dict) or set(data) != {'alias', 'request_id', 'muted'} + or not isinstance(data['request_id'], str) or not re.fullmatch(r'[a-zA-Z0-9-]{8,80}', data['request_id']) + or type(data['muted']) is not bool): + raise ValueError('Saved demo requires an alias, request ID and explicit output choice') + if not self.offroad(): + raise ValueError('Saved demo playback requires the vehicle to be offroad') + selected = entry(self.root, data['alias']) + with self.lock: + current = self.status() + if current['running']: + if current['request_id'] == data['request_id'] and current['alias'] == data['alias'] and current['muted'] == data['muted']: + return current + raise ValueError('A saved demo is already running; stop that owned session first') + if read(self.owner_path).get('request_id') == data['request_id']: + raise ValueError('This launch request already finished; use a new request ID') + (self.root/'generated').mkdir(parents=True, exist_ok=True) + with (self.root/'generated/native_session.lock').open('a') as lease: + try: fcntl.flock(lease, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: raise ValueError('Another RoadScore replay is active') from None + out = self.root/'results'/('demo_' + data['request_id']) + out.mkdir(parents=True, exist_ok=False) + command = ['/usr/local/venv/bin/python', str(self.root/'prototype/native_prepared_showcase.py'), + '--roadscore', data['alias'], '--demo', '--hold-start', '--out', str(out), + '--score-archive', selected['archive']] + if selected.get('curve_plan'): command += ['--curve-plan', selected['curve_plan']] + if data['muted']: command.append('--muted') + env = os.environ.copy() + for key in ('ZMQ', 'OPENPILOT_PREFIX', 'OPENPILOT_ZMQ_NAMESPACE', 'PARAMS_ROOT'): + env.pop(key, None) + env['PYTHONPATH'] = ':'.join((str(self.root/'prototype'), '/data/openpilot', '/data/roadscore-feasibility/venv/lib/python3.12/site-packages')) + with (out/'launcher.log').open('wb') as log: + process = self.spawn(command, cwd='/data/openpilot', env=env, stdin=subprocess.DEVNULL, + stdout=log, stderr=subprocess.STDOUT, start_new_session=True, close_fds=True) + stamp = None + for _ in range(20): + stamp = self.ticks(process.pid) + if stamp is not None or process.poll() is not None: break + time.sleep(.01) + if stamp is None: + raise ValueError('Saved demo launcher failed; inspect its launcher.log') + write(self.owner_path, {**data, 'route': selected['route'], 'pid': process.pid, + 'start_ticks': stamp, 'out': str(out), 'created_wall': time.monotonic()}) + return self.status() + + def release(self, data): + if not isinstance(data, dict) or set(data) != {'request_id', 'session_id'}: + raise ValueError('Starting prepared playback requires both session identities') + if not self.offroad(): + raise ValueError('Saved demo playback requires the vehicle to be offroad') + with self.lock: + current = self.status() + if (not current['running'] or not current.get('prepared') + or current.get('request_id') != data['request_id'] + or current.get('ready_session_id') != data['session_id']): + raise ValueError('No matching prepared demo is ready') + path = Path(current['out'])/'start.json' + previous = read(path) + deadline = previous.get('start_at_wall') + if previous.get('session_id') != data['session_id'] or type(deadline) not in (int, float) or not math.isfinite(deadline): + deadline = time.monotonic() + 1. + write(path, {'session_id': data['session_id'], 'play': True, 'start_at_wall': deadline}) + return {'request_id': data['request_id'], 'session_id': data['session_id'], + 'start_at_wall': deadline, 'server_wall': time.monotonic()} + + def stop(self, data): + if not isinstance(data, dict) or set(data) != {'request_id'}: + raise ValueError('Stopping a saved demo requires its request ID') + with self.lock: + current = self.status() + if current.get('request_id') != data['request_id']: + raise ValueError('Saved demo session changed; refusing to stop another session') + if current['running']: + self.stop_group(current['pid']) + return {'stopping': current['running'], 'request_id': data['request_id']} diff --git a/roadscore/prototype/test_demo_session.py b/roadscore/prototype/test_demo_session.py new file mode 100644 index 0000000000..ddf603f7ff --- /dev/null +++ b/roadscore/prototype/test_demo_session.py @@ -0,0 +1,73 @@ +import fcntl +import json +from pathlib import Path +import tempfile +from types import SimpleNamespace +import unittest + +from demo_catalog import entry +from demo_session import DemoSession, write + + +class SavedDemoTests(unittest.TestCase): + def setUp(self): + self.tmp = tempfile.TemporaryDirectory();self.addCleanup(self.tmp.cleanup) + self.root = Path(self.tmp.name) + for folder in ('assets', 'generated', 'archive'): (self.root/folder).mkdir() + write(self.root/'assets/demo_catalog.json', {'version':1,'routes':{'route1':{'ready':True,'route':'private/route','archive':'archive'},'route2':{'ready':False}}}) + self.calls=[];self.stops=[];self.running=True;self.offroad=True + def spawn(command, **kwargs): + self.calls.append((command,kwargs));return SimpleNamespace(pid=123,poll=lambda:None) + self.service=DemoSession(self.root,lambda:self.offroad,spawn=spawn, + ticks=lambda pid:'birth1' if self.running and pid==123 else None, + stop_group=self.stops.append) + self.request={'alias':'route1','request_id':'test-session-1','muted':True} + + def test_only_registered_ready_route_can_launch(self): + for alias in ('route2','route3','../route1','route1;other'): + with self.assertRaises(ValueError):entry(self.root,alias) + self.assertEqual(entry(self.root,'route1')['archive'],str((self.root/'archive').resolve())) + self.offroad=False + with self.assertRaises(ValueError):self.service.start(self.request) + self.assertEqual(self.calls,[]) + + def test_existing_native_replay_blocks_launch(self): + with (self.root/'generated/native_session.lock').open('a') as lease: + fcntl.flock(lease,fcntl.LOCK_EX|fcntl.LOCK_NB) + with self.assertRaises(ValueError):self.service.start(self.request) + self.assertEqual(self.calls,[]) + + def test_launch_is_owned_and_idempotent_without_a_composer(self): + current=self.service.start(self.request) + self.assertTrue(current['running']);self.assertFalse(current['generation_invoked']) + self.assertEqual(self.service.start(self.request)['pid'],123) + self.assertEqual(len(self.calls),1) + command,options=self.calls[0] + self.assertTrue(command[1].endswith('/native_prepared_showcase.py')) + self.assertIn('--hold-start',command);self.assertIn('--muted',command) + self.assertTrue(options['close_fds']);self.assertTrue(options['start_new_session']) + with self.assertRaises(ValueError):self.service.start({**self.request,'request_id':'different-request'}) + with self.assertRaises(ValueError):self.service.stop({'request_id':'different-request'}) + self.service.stop({'request_id':self.request['request_id']}) + self.assertEqual(self.stops,[123]) + + def test_matching_ready_session_is_required_for_release(self): + current=self.service.start(self.request);out=Path(current['out']) + request={'request_id':self.request['request_id'],'session_id':'prepared-session'} + with self.assertRaises(ValueError):self.service.release(request) + write(out/'demo_ready.json',{'ready':True,'session_id':'prepared-session'}) + with self.assertRaises(ValueError):self.service.release({**request,'session_id':'old-session'}) + first=self.service.release(request);again=self.service.release(request) + self.assertEqual(first['start_at_wall'],again['start_at_wall']) + self.assertGreater(first['start_at_wall'],first['server_wall']) + self.assertEqual(json.loads((out/'start.json').read_text())['session_id'],'prepared-session') + + def test_reused_or_exited_pid_is_never_signalled(self): + self.service.start(self.request);self.running=False + self.assertFalse(self.service.status()['running']) + self.service.stop({'request_id':self.request['request_id']}) + self.assertEqual(self.stops,[]) + with self.assertRaises(ValueError):self.service.start(self.request) + + +if __name__=='__main__':unittest.main() diff --git a/starpilot/system/the_galaxy/roadscore.py b/starpilot/system/the_galaxy/roadscore.py index 07ac3dcf21..5bfdedc139 100644 --- a/starpilot/system/the_galaxy/roadscore.py +++ b/starpilot/system/the_galaxy/roadscore.py @@ -104,6 +104,19 @@ class Operator: self.live_owner = None self.live_init_lock = threading.Lock() self.output_init_lock = threading.Lock() + self.demo_owner = None + self.demo_init_lock = threading.Lock() + + def demo_controller(self): + if not self.device: + raise ValueError('Paired demo launch is available on the comma') + with self.demo_init_lock: + if self.demo_owner is None: + path = self.root/'prototype' + if str(path) not in sys.path:sys.path.insert(0,str(path)) + from demo_session import DemoSession + self.demo_owner = DemoSession(self.root, self.offroad) + return self.demo_owner def live_controller(self): if not self.device: @@ -150,8 +163,17 @@ class Operator: available=bool(fresh and state.get('engagement_presentation',{}).get('enabled') is True and state.get('input_mode')=='replay' and state.get('route') not in (None,'','live') and isinstance(session,str) and session) mode=state.get('demo_engagement_mode','recorded') signal_mode=state.get('demo_signal_mode','recorded') - return {'available':available,'session_id':session if available else None,'signal_mode':signal_mode if signal_mode in ('recorded','left','right','off') else 'recorded','mode':mode if mode in ('recorded','engaged','disengaged') else 'recorded', - 'reason':'' if available else 'Start a RoadScore replay to simulate displayed engagement, signals and music.'} + result = {'available':available,'session_id':session if available else None,'signal_mode':signal_mode if signal_mode in ('recorded','left','right','off') else 'recorded','mode':mode if mode in ('recorded','engaged','disengaged') else 'recorded', + 'reason':'' if available else 'Start a RoadScore replay to simulate displayed engagement, signals and music.'} + if available: + result['route'] = state['route'] + result['readiness'] = state.get('readiness') + result['compute'] = state.get('compute') + result['status_age_seconds'] = time.monotonic()-stamp + values = {key:state.get(key) for key in ('route_t','elapsed','source_model_ns')} + if all(type(value) in (int,float) and math.isfinite(value) for value in values.values()): + result['playhead'] = {**values,'sampled_wall':stamp,'server_wall':time.monotonic()} + return result def demo_engagement(self,data,field='mode'): allowed=('recorded','engaged','disengaged') if field=='mode' else ('recorded','left','right','off') @@ -197,10 +219,12 @@ class Operator: state = output['state'] demo=self.demo_status() if not offroad:demo.update(available=False,reason='Replay simulation requires the vehicle to be offroad.') + prepared_demo = demo['available'] and demo.get('compute') == 'prepared-core' + if prepared_demo and demo.get('readiness') in STATES:state = demo['readiness'] return dict(available=self.device and self.root.exists(), state=state, profiles=PROFILES, live=self.live_status(), demo=demo, profile=worker.get('profile') if live else settings.get('profile', 'prism'), selected_profile=settings.get('profile', 'prism'), - composer='ace' if live else None, backend='Chestnut' if live else None, + composer='ace' if live or prepared_demo else None, backend='Prepared local audio' if prepared_demo else 'Chestnut' if live else None, generation_seed=worker.get('generation_seed') if live else None, offroad=bool(offroad), locked=bool(locked), preparing=self.preparing, can_prepare=False, @@ -209,6 +233,14 @@ class Operator: calibrating=output.get('calibrating', False), session_muted=output.get('session_muted', True), output=output.get('output'), latency_ms=output.get('latency_ms'), error=self.error or output.get('error')) def operate(self, action, data, offroad): + if action in ('demo_start','demo_ready','demo_play','demo_stop'): + if not offroad:raise ValueError('Saved demo playback requires the vehicle to be offroad') + owner = self.demo_controller() + if action == 'demo_start':return owner.start(data) + if action == 'demo_play':return owner.release(data) + if action == 'demo_stop':return owner.stop(data) + if data:raise ValueError('Demo readiness takes no arguments') + return owner.status() if action in ('demo_engagement','demo_signal'): if not offroad:raise ValueError('Replay simulation requires the vehicle to be offroad') return self.demo_engagement(data,'signal_mode' if action=='demo_signal' else 'mode') diff --git a/starpilot/system/the_galaxy/tests/test_roadscore_demo.py b/starpilot/system/the_galaxy/tests/test_roadscore_demo.py index 210982bdb4..9d76c03792 100644 --- a/starpilot/system/the_galaxy/tests/test_roadscore_demo.py +++ b/starpilot/system/the_galaxy/tests/test_roadscore_demo.py @@ -60,5 +60,15 @@ class DemoTests(unittest.TestCase): for payload in ({'mode':'engaged'}, {'mode':'engaged','session_id':'session-a','pid':1}): with self.assertRaises(ValueError):self.operator.operate('demo_engagement',payload,True) self.assertFalse((self.run/'demo_engagement.json').exists()) + def test_prepared_demo_readiness_does_not_require_a_warm_worker(self): + from unittest.mock import patch + self.state.update(readiness='READY',compute='prepared-core',route_t=12.,elapsed=12.1,source_model_ns=123456789) + self.write() + with patch.object(module,'current_worker',return_value={}),patch.object(self.operator,'target',return_value={}),patch.object(self.operator,'live_status',return_value={'enabled':False}): + state=self.operator.status(True) + self.assertEqual(state['state'],'READY') + self.assertEqual(state['backend'],'Prepared local audio') + self.assertEqual(state['demo']['playhead']['route_t'],12.) + self.assertGreaterEqual(state['demo']['status_age_seconds'],0.) if __name__=='__main__':unittest.main()