mirror of
https://github.com/firestar5683/StarPilot.git
synced 2026-10-04 13:24:13 +08:00
Replace stale preparation status when supervised transport exits
This commit is contained in:
@@ -0,0 +1,40 @@
|
||||
"""Keep preparation status truthful when supervised processes fail."""
|
||||
import json
|
||||
import time
|
||||
|
||||
|
||||
class LaunchFailure(RuntimeError):
|
||||
def __init__(self, component, code, connection=False):
|
||||
self.component, self.code, self.connection = component, code, connection
|
||||
super().__init__(f'{component} exited with status {code}; see {component}.log')
|
||||
|
||||
|
||||
def check_children(children, *, remote, include_receiver=False, allow_clean_receiver=False):
|
||||
names = ['semantic_planner', 'semantic_tunnel'] + (['receiver'] if include_receiver else [])
|
||||
for name in names:
|
||||
child = children.get(name)
|
||||
if child is None:
|
||||
continue
|
||||
code = child.poll()
|
||||
if code is not None and not (name == 'receiver' and code == 0 and allow_clean_receiver):
|
||||
raise LaunchFailure(name, code, name == 'semantic_tunnel' or (name == 'receiver' and remote and code == 255))
|
||||
|
||||
|
||||
def describe_failure(error):
|
||||
connection = isinstance(error, LaunchFailure) and error.connection
|
||||
return {'readiness': 'DEGRADED', 'job_inflight': False,
|
||||
'failure_kind': 'connection_lost' if connection else 'launch_failed',
|
||||
'failure_component': getattr(error, 'component', 'launcher'),
|
||||
'failure_exit_code': getattr(error, 'code', None),
|
||||
'failure_cause': str(error)[-2000:], 'failure_time': time.time()}
|
||||
|
||||
|
||||
def write_failure(path, failure):
|
||||
try:
|
||||
state = json.loads(path.read_text())
|
||||
except (OSError, ValueError):
|
||||
state = {}
|
||||
state.update(failure)
|
||||
temporary = path.with_name(path.name + '.failure.tmp')
|
||||
temporary.write_text(json.dumps(state))
|
||||
temporary.replace(path)
|
||||
@@ -9,6 +9,7 @@ from clock_sync import measure
|
||||
from receiver_environment import assignments as receiver_assignments,resolve_compute
|
||||
from presentation_policy import select_launch
|
||||
from hook_launch import enabled as hook_enabled, start_planner
|
||||
from launch_health import check_children, describe_failure, write_failure
|
||||
from session_seed import select_session, seed_argument, seed_environment, remote_assignments
|
||||
R=Path(__file__).resolve().parents[1]
|
||||
native=Path('/TICI').exists()
|
||||
@@ -74,10 +75,10 @@ if a.replay:
|
||||
from score_archive import latest
|
||||
try:score=latest(a.routeid)
|
||||
except FileNotFoundError:raise SystemExit('No stored RoadScore exists for this route. Generate a score first.')
|
||||
display=None;children=[];logs=[];launch_started=time.monotonic()
|
||||
display=None;children=[];named_children={};failure=None;logs=[];launch_started=time.monotonic()
|
||||
print(('Preparing stored score replay; no generation. ' if a.replay else ('Preparing ACE replay; first preparation may take 10–15 minutes. ' if a.composer=='ace' else 'Preparing RoadScore replay; cold preparation can take 2–3 minutes. '))+('Host speaker enabled.' if a.audible else 'Muted host capture.'),flush=True)
|
||||
def launch(cmd,name,**kw):
|
||||
f=(out/(name+'.log')).open('wb');logs.append(f);c=subprocess.Popen(cmd,stdout=f,stderr=f,env=env,cwd=rt,start_new_session=True,**kw);children.append(c);return c
|
||||
f=(out/(name+'.log')).open('wb');logs.append(f);c=subprocess.Popen(cmd,stdout=f,stderr=f,env=env,cwd=rt,start_new_session=True,**kw);children.append(c);named_children[name]=c;return c
|
||||
try:
|
||||
if not native and Path('/usr/bin/caffeinate').exists():launch(['/usr/bin/caffeinate','-i'],'wake_assertion')
|
||||
if composition_policy=='hook-v2' and not a.transport_only:start_planner(launch,env,out,R,a.bench,native=native)
|
||||
@@ -104,6 +105,7 @@ try:
|
||||
receiver=launch(receiver_command,'receiver',stdin=subprocess.PIPE)
|
||||
deadline=time.monotonic()+(1560 if a.composer=='ace' else 420)
|
||||
while b'BRIDGE_READY' not in (out/'receiver.log').read_bytes():
|
||||
check_children(named_children,remote=not native,include_receiver=True)
|
||||
if receiver.poll() is not None or time.monotonic()>deadline:raise RuntimeError('Bench receiver failed: '+(out/'receiver.log').read_text())
|
||||
time.sleep(.2)
|
||||
audio_host=None
|
||||
@@ -112,12 +114,14 @@ try:
|
||||
audio_host=launch([str(R/'.analysis-venv/bin/python'),str(R/'prototype/host_audio.py'),'--bench',a.bench,'--out',str(out),'--clock',str(out/'clock_sync.json')]+(['--audible'] if a.audible else [])+(['--device',a.audio_device] if a.audio_device else []),'host_audio')
|
||||
deadline=time.monotonic()+20
|
||||
while not (out/'host_audio_ready').exists():
|
||||
check_children(named_children,remote=not native,include_receiver=True)
|
||||
if audio_host.poll() is not None or time.monotonic()>deadline:raise RuntimeError('Host audio failed; see host_audio.log')
|
||||
time.sleep(.1)
|
||||
f=(out/'sender.log').open('wb');logs.append(f)
|
||||
sender=subprocess.Popen([str(py),str(R/'prototype/replay_bridge.py'),'send','--seconds',str(a.duration)],stdout=receiver.stdin,stderr=f,env=env,cwd=rt,start_new_session=True);children.append(sender);receiver.stdin.close()
|
||||
deadline=time.monotonic()+15
|
||||
while b'BRIDGE_SUBSCRIBED' not in (out/'sender.log').read_bytes():
|
||||
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 a.replay:(out/'roadscore_status.json').write_text(json.dumps({'readiness':'READY','style':'Stored score','section':'ARCHIVED SCORE','compute':'none'}))
|
||||
@@ -126,6 +130,7 @@ try:
|
||||
end_watch=ReplayEnd();end_reason='requested duration'
|
||||
state_path=Path('/tmp/replay_state_'+env.get('OPENPILOT_PREFIX','default')+'.json')
|
||||
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())
|
||||
@@ -163,12 +168,21 @@ try:
|
||||
if a.composer=='ace':subprocess.run(['scp',a.bench+':/data/roadscore/generated/ace_link.jsonl',str(out/'ace_link.jsonl')],stdout=subprocess.DEVNULL,stderr=subprocess.DEVNULL)
|
||||
from score_archive import archive
|
||||
archived=archive(a.routeid,out,a.start);print('Score archived:',archived,flush=True)
|
||||
except Exception as error:
|
||||
try:
|
||||
check_children(named_children,remote=not native,include_receiver=True,allow_clean_receiver=True)
|
||||
except Exception as child_error:
|
||||
error=child_error
|
||||
failure=describe_failure(error)
|
||||
write_failure(out/'roadscore_status.json',failure)
|
||||
raise
|
||||
finally:
|
||||
for c in reversed(children):
|
||||
if c.poll() is None:
|
||||
os.killpg(c.pid,signal.SIGTERM)
|
||||
try:c.wait(timeout=10)
|
||||
except subprocess.TimeoutExpired:os.killpg(c.pid,signal.SIGKILL);c.wait()
|
||||
if failure:write_failure(out/'roadscore_status.json',failure)
|
||||
if display:display.close()
|
||||
for f in logs:f.close()
|
||||
print('Replay processes stopped. Worker ownership remains with its supervisor.',flush=True)
|
||||
|
||||
@@ -61,7 +61,11 @@ def overlay_view(state):
|
||||
section = 'INTENT: ' + section
|
||||
elapsed = seconds(state.get('generation_elapsed_seconds'))
|
||||
note = ''
|
||||
if state.get('holding_accepted_music'):
|
||||
if state.get('failure_kind') in ('connection_lost', 'transport_unreachable'):
|
||||
note = 'Connection lost'
|
||||
elif state.get('failure_kind') == 'launch_failed':
|
||||
note = 'Preparation stopped'
|
||||
elif state.get('holding_accepted_music'):
|
||||
note = 'Holding accepted music'
|
||||
elif state.get('worker_failed'):
|
||||
note = 'Composer unavailable'
|
||||
@@ -138,7 +142,9 @@ def draw_panel(rl, font, state, screen_width, screen_height, emphasis_font=None,
|
||||
subtitle = identity + (' / ' + section if section else '')
|
||||
title = 'RoadScore'
|
||||
if view['activity'] == 'DEGRADED':
|
||||
subtitle = ('Music on hold' if state.get('holding_accepted_music') else
|
||||
subtitle = ('Connection lost' if state.get('failure_kind') in ('connection_lost', 'transport_unreachable') else
|
||||
'Preparation stopped' if state.get('failure_kind') == 'launch_failed' else
|
||||
'Music on hold' if state.get('holding_accepted_music') else
|
||||
'Composer unavailable' if state.get('worker_failed') else 'Reserve in use')
|
||||
elif view['event']:
|
||||
subtitle = view['event'].split(' / ')[0]
|
||||
|
||||
@@ -0,0 +1,56 @@
|
||||
import json
|
||||
import sys
|
||||
import tempfile
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from unittest.mock import Mock
|
||||
sys.path.insert(0, str(Path(__file__).resolve().parents[1] / 'prototype'))
|
||||
from launch_health import LaunchFailure, check_children, describe_failure, write_failure
|
||||
from overlay_view import overlay_view
|
||||
|
||||
|
||||
class LaunchHealthTests(unittest.TestCase):
|
||||
def test_dead_tunnel_detected_while_receiver_still_compiles(self):
|
||||
children = {'receiver': Mock(poll=lambda: None), 'semantic_tunnel': Mock(poll=lambda: 255)}
|
||||
with self.assertRaises(LaunchFailure) as caught:
|
||||
check_children(children, remote=True, include_receiver=True)
|
||||
status = describe_failure(caught.exception)
|
||||
self.assertEqual(status['failure_kind'], 'connection_lost')
|
||||
self.assertEqual(status['failure_component'], 'semantic_tunnel')
|
||||
self.assertEqual(overlay_view(status)['note'], 'Connection lost')
|
||||
|
||||
def test_receiver_ssh_exit_distinct_from_remote_application_error(self):
|
||||
for code, kind in [(255, 'connection_lost'), (1, 'launch_failed'), (0, 'launch_failed')]:
|
||||
with self.subTest(code=code):
|
||||
with self.assertRaises(LaunchFailure) as caught:
|
||||
check_children({'receiver': Mock(poll=lambda: code)}, remote=True, include_receiver=True)
|
||||
self.assertEqual(describe_failure(caught.exception)['failure_kind'], kind)
|
||||
check_children({'receiver': Mock(poll=lambda: 0)}, remote=True, include_receiver=True, allow_clean_receiver=True)
|
||||
|
||||
def test_planner_failure_and_timeout_are_not_gpu_diagnoses(self):
|
||||
for error in (LaunchFailure('semantic_planner', -9), TimeoutError('Receiver readiness deadline elapsed')):
|
||||
status = describe_failure(error)
|
||||
self.assertEqual(status['readiness'], 'DEGRADED')
|
||||
self.assertNotIn('worker_failed', status)
|
||||
self.assertIn('failure_cause', status)
|
||||
|
||||
def test_persisted_failure_replaces_preparing_without_losing_profile(self):
|
||||
with tempfile.TemporaryDirectory() as root:
|
||||
path = Path(root)/'roadscore_status.json'
|
||||
path.write_text(json.dumps({'readiness': 'PREPARING', 'profile': 'prism', 'job_inflight': True}))
|
||||
failure = describe_failure(LaunchFailure('receiver', 255, True))
|
||||
write_failure(path, failure)
|
||||
state = json.loads(path.read_text())
|
||||
self.assertEqual(state['profile'], 'prism')
|
||||
self.assertEqual(state['readiness'], 'DEGRADED')
|
||||
self.assertFalse(state['job_inflight'])
|
||||
self.assertNotIn('worker_failed', state)
|
||||
self.assertEqual(state['failure_exit_code'], 255)
|
||||
self.assertEqual(overlay_view(state)['note'], 'Connection lost')
|
||||
path.write_text('{}')
|
||||
write_failure(path, failure)
|
||||
self.assertEqual(json.loads(path.read_text())['readiness'], 'DEGRADED')
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
unittest.main()
|
||||
Reference in New Issue
Block a user