mirror of
https://github.com/firestar5683/StarPilot.git
synced 2026-10-04 13:24:13 +08:00
Add session-owned saved demo launch and playhead status in Galaxy
This commit is contained in:
@@ -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
|
||||
@@ -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']}
|
||||
@@ -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()
|
||||
@@ -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')
|
||||
|
||||
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user