mirror of
https://github.com/firestar5683/StarPilot.git
synced 2026-10-04 13:24:13 +08:00
Capture opt-in native replay screen on bounded JPEG worker
This commit is contained in:
@@ -8,6 +8,8 @@ def install():
|
||||
if os.environ.get('ROADSCORE_OVERLAY')!='1':return
|
||||
import pyray as rl
|
||||
from openpilot.system.ui.lib.application import gui_app,FontWeight
|
||||
from screen_mirror import ScreenMirror,read_ui_rgba
|
||||
mirror=ScreenMirror.from_environ(os.environ)
|
||||
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
|
||||
@@ -100,8 +102,15 @@ def install():
|
||||
if capture_target is not None:capture_image(capture_target)
|
||||
def render(*args,**kwargs):
|
||||
nonlocal native_nav_visible,home_footer_right,home_rect
|
||||
for should_render in original(*args,**kwargs):
|
||||
yield should_render
|
||||
if should_render:draw()
|
||||
native_nav_visible=False;home_footer_right=None;home_rect=None
|
||||
try:
|
||||
for should_render in original(*args,**kwargs):
|
||||
yield should_render
|
||||
if should_render:
|
||||
draw()
|
||||
if mirror is not None:
|
||||
from openpilot.selfdrive.ui.ui_state import ui_state
|
||||
mirror.capture(lambda:read_ui_rgba(rl,gui_app),started=bool(ui_state.started),session_id=raw_state.get('presentation_session_id'))
|
||||
native_nav_visible=False;home_footer_right=None;home_rect=None
|
||||
finally:
|
||||
if mirror is not None:mirror.close()
|
||||
gui_app.render=render
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
"""Opt-in native replay screen capture with one asynchronous JPEG slot."""
|
||||
import io
|
||||
import json
|
||||
import os
|
||||
from pathlib import Path
|
||||
import queue
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
|
||||
from replay_ui_controls import isolated_replay
|
||||
|
||||
|
||||
class ScreenMirror:
|
||||
"""The UI owns readback; a bounded worker owns resize, encoding and disk I/O."""
|
||||
def __init__(self, directory, session_id, *, fps=15, quality=78, max_width=1072):
|
||||
self.directory = Path(directory)
|
||||
if not self.directory.is_absolute() or self.directory == Path('/') or '..' in self.directory.parts:
|
||||
raise ValueError('Mirror directory must be an explicit absolute run directory')
|
||||
if not re.fullmatch(r'[a-zA-Z0-9-]{8,80}', session_id):
|
||||
raise ValueError('Mirror requires a valid showcase session')
|
||||
self.session_id = session_id
|
||||
self.interval = 1 / min(15, max(1, fps))
|
||||
self.quality = min(80, max(75, quality))
|
||||
self.max_width = min(1072, max(1, max_width))
|
||||
self._pending = queue.Queue(maxsize=1)
|
||||
self._busy = threading.Lock()
|
||||
self._closed = threading.Event()
|
||||
self._next_capture = 0.
|
||||
self._frame_id = 0
|
||||
self.last_error = None
|
||||
self._thread = threading.Thread(target=self._encode_loop, name='roadscore-screen-mirror', daemon=True)
|
||||
self._thread.start()
|
||||
|
||||
@classmethod
|
||||
def from_environ(cls, environ):
|
||||
directory = environ.get('ROADSCORE_MIRROR_DIR')
|
||||
if (not directory or not isolated_replay(environ)
|
||||
or environ.get('ROADSCORE_PRESENTATION_POLICY', 'frozen') == 'frozen'
|
||||
or environ.get('ROADSCORE_SEED_ORIGIN') == 'judging-route'):
|
||||
return None
|
||||
try:
|
||||
return cls(directory, environ.get('ROADSCORE_SHOWCASE_SESSION', ''))
|
||||
except (ValueError, TypeError):
|
||||
return None
|
||||
|
||||
def capture(self, readback, *, started, session_id, now=None):
|
||||
"""Drop offroad, stale-session, rate-limited and busy frames before readback."""
|
||||
now = time.monotonic() if now is None else now
|
||||
if (self._closed.is_set() or not started or session_id != self.session_id
|
||||
or now < self._next_capture or not self._busy.acquire(blocking=False)):
|
||||
return False
|
||||
self._next_capture = now + self.interval
|
||||
begin = time.perf_counter()
|
||||
try:
|
||||
data, width, height, flipped = readback()
|
||||
if not 0 < width <= 4096 or not 0 < height <= 4096 or len(data) != width * height * 4:
|
||||
raise ValueError('Mirror readback must be a bounded RGBA frame')
|
||||
self._frame_id += 1
|
||||
job = dict(data=data, width=width, height=height, flipped=bool(flipped),
|
||||
frame_id=self._frame_id, captured_wall=now,
|
||||
capture_ms=(time.perf_counter()-begin)*1000)
|
||||
self._pending.put_nowait(job)
|
||||
return True
|
||||
except Exception as error:
|
||||
self.last_error = f'{type(error).__name__}: {error}'[:240]
|
||||
self._busy.release()
|
||||
return False
|
||||
|
||||
def _encode(self, job):
|
||||
from PIL import Image
|
||||
begin = time.perf_counter()
|
||||
picture = Image.frombytes('RGBA', (job['width'], job['height']), job['data'])
|
||||
if job['flipped']:
|
||||
picture = picture.transpose(Image.Transpose.FLIP_TOP_BOTTOM)
|
||||
if picture.width > self.max_width:
|
||||
picture = picture.resize((self.max_width, max(1, round(picture.height*self.max_width/picture.width))), Image.Resampling.BILINEAR)
|
||||
picture = picture.convert('RGB')
|
||||
encoded = io.BytesIO()
|
||||
picture.save(encoded, format='JPEG', quality=self.quality, optimize=False)
|
||||
metadata = dict(session_id=self.session_id, captured_wall=job['captured_wall'],
|
||||
encoded_wall=time.monotonic(), frame_id=job['frame_id'], width=picture.width,
|
||||
height=picture.height, encode_ms=(time.perf_counter()-begin)*1000,
|
||||
capture_ms=job['capture_ms'])
|
||||
return encoded.getvalue(), metadata
|
||||
|
||||
def _publish(self, data, metadata):
|
||||
self.directory.mkdir(parents=True, exist_ok=True)
|
||||
image_tmp = self.directory / ('.latest-' + self.session_id + '.jpg.tmp')
|
||||
metadata_tmp = self.directory / ('.frame-' + self.session_id + '.json.tmp')
|
||||
try:
|
||||
image_tmp.write_bytes(data)
|
||||
metadata_tmp.write_text(json.dumps(metadata, separators=(',', ':')))
|
||||
if not self._closed.is_set():
|
||||
os.replace(image_tmp, self.directory/'latest.jpg')
|
||||
os.replace(metadata_tmp, self.directory/'frame.json')
|
||||
finally:
|
||||
image_tmp.unlink(missing_ok=True)
|
||||
metadata_tmp.unlink(missing_ok=True)
|
||||
|
||||
def _encode_loop(self):
|
||||
while not self._closed.is_set():
|
||||
try:
|
||||
job = self._pending.get(timeout=.1)
|
||||
except queue.Empty:
|
||||
continue
|
||||
try:
|
||||
if not self._closed.is_set():
|
||||
data, metadata = self._encode(job)
|
||||
if not self._closed.is_set():
|
||||
self._publish(data, metadata)
|
||||
except Exception as error:
|
||||
self.last_error = f'{type(error).__name__}: {error}'[:240]
|
||||
finally:
|
||||
self._busy.release()
|
||||
|
||||
def close(self):
|
||||
"""Do not block UI teardown on encoding or a slow filesystem."""
|
||||
self._closed.set()
|
||||
|
||||
|
||||
def read_ui_rgba(rl, gui_app):
|
||||
"""Call only from the native UI thread while its render target is active."""
|
||||
rl.rl_draw_render_batch_active()
|
||||
image = rl.load_image_from_texture(gui_app._render_texture.texture) if gui_app._render_texture else rl.load_image_from_screen()
|
||||
try:
|
||||
if not 0 < image.width <= 4096 or not 0 < image.height <= 4096:
|
||||
raise ValueError('Mirror render target is too large')
|
||||
data = bytes(rl.ffi.buffer(image.data, image.width*image.height*4))
|
||||
return data, image.width, image.height, bool(gui_app._render_texture)
|
||||
finally:
|
||||
rl.unload_image(image)
|
||||
@@ -0,0 +1,113 @@
|
||||
import json
|
||||
from pathlib import Path
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
import unittest
|
||||
from unittest.mock import Mock, patch
|
||||
from types import SimpleNamespace as NS
|
||||
|
||||
from screen_mirror import ScreenMirror, read_ui_rgba
|
||||
|
||||
|
||||
SESSION = 'abcd1234-session'
|
||||
|
||||
|
||||
class ScreenMirrorTests(unittest.TestCase):
|
||||
def make(self, folder):
|
||||
mirror = ScreenMirror(Path(folder)/'mirror', SESSION)
|
||||
self.addCleanup(mirror.close)
|
||||
return mirror
|
||||
|
||||
def test_opt_in_requires_isolated_non_frozen_session(self):
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
env = dict(ROADSCORE_MIRROR_DIR=folder, ROADSCORE_SHOWCASE_SESSION=SESSION,
|
||||
ROADSCORE_REPLAY_UI_CONTROLS='1', SIMULATION='1',
|
||||
OPENPILOT_PREFIX='roadscore_replay', ROADSCORE_PRESENTATION_POLICY='conservative-v4')
|
||||
mirror = ScreenMirror.from_environ(env)
|
||||
self.assertIsNotNone(mirror)
|
||||
mirror.close()
|
||||
for changed in ({'ROADSCORE_MIRROR_DIR':''}, {'ROADSCORE_SHOWCASE_SESSION':''},
|
||||
{'SIMULATION':'0'}, {'OPENPILOT_PREFIX':'real-car'},
|
||||
{'ROADSCORE_PRESENTATION_POLICY':'frozen'}, {'ROADSCORE_SEED_ORIGIN':'judging-route'}):
|
||||
self.assertIsNone(ScreenMirror.from_environ(dict(env, **changed)))
|
||||
self.assertFalse((Path(folder)/'latest.jpg').exists())
|
||||
|
||||
def test_directory_and_session_are_explicit(self):
|
||||
for directory, session in [('relative/path', SESSION), ('/', SESSION), ('/tmp/../unsafe',SESSION), ('/tmp/safe','')]:
|
||||
with self.assertRaises(ValueError):
|
||||
ScreenMirror(directory,session)
|
||||
|
||||
def test_offroad_wrong_session_and_closed_never_read_back(self):
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
mirror=self.make(folder);read=Mock()
|
||||
self.assertFalse(mirror.capture(read,started=False,session_id=SESSION))
|
||||
self.assertFalse(mirror.capture(read,started=True,session_id='old-session'))
|
||||
mirror.close()
|
||||
self.assertFalse(mirror.capture(read,started=True,session_id=SESSION))
|
||||
read.assert_not_called()
|
||||
|
||||
def test_one_slot_drops_busy_and_rate_limited_frames(self):
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
mirror=self.make(folder)
|
||||
encoding=threading.Event();release=threading.Event()
|
||||
def encode(job):
|
||||
encoding.set();release.wait(2)
|
||||
return b'jpeg', {'frame_id':job['frame_id']}
|
||||
with patch.object(mirror,'_encode',side_effect=encode),patch.object(mirror,'_publish'):
|
||||
try:
|
||||
read=Mock(return_value=(b'\xff'*16,2,2,False))
|
||||
self.assertTrue(mirror.capture(read,started=True,session_id=SESSION,now=1))
|
||||
self.assertTrue(encoding.wait(1))
|
||||
self.assertFalse(mirror.capture(read,started=True,session_id=SESSION,now=1.01))
|
||||
self.assertFalse(mirror.capture(read,started=True,session_id=SESSION,now=2))
|
||||
self.assertEqual(read.call_count,1)
|
||||
self.assertEqual(mirror._pending.qsize(),0)
|
||||
finally:
|
||||
mirror.close();release.set();mirror._thread.join(1)
|
||||
|
||||
def test_async_jpeg_atomic_outputs_flip_resize_and_metadata(self):
|
||||
from PIL import Image
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
mirror=self.make(folder)
|
||||
# Texture readback is upside down; final top must be red.
|
||||
original=Image.new('RGBA',(1200,40),'red')
|
||||
original.paste('blue',(0,0,1200,20))
|
||||
thread_ids=[];done=threading.Event();real=mirror._publish
|
||||
def publish(*args):
|
||||
thread_ids.append(threading.get_ident());real(*args);done.set()
|
||||
with patch.object(mirror,'_publish',side_effect=publish):
|
||||
captured=time.monotonic()
|
||||
self.assertTrue(mirror.capture(lambda:(original.tobytes(),1200,40,True),started=True,session_id=SESSION,now=captured))
|
||||
self.assertTrue(done.wait(2))
|
||||
data=json.loads((mirror.directory/'frame.json').read_text())
|
||||
self.assertEqual(data['session_id'],SESSION)
|
||||
self.assertEqual(data['captured_wall'],captured)
|
||||
self.assertGreaterEqual(data['encoded_wall'],captured)
|
||||
self.assertGreaterEqual(data['encode_ms'],0)
|
||||
self.assertGreaterEqual(data['capture_ms'],0)
|
||||
with Image.open(mirror.directory/'latest.jpg') as result:
|
||||
self.assertEqual(result.size,(1072,36))
|
||||
self.assertGreater(result.getpixel((10,2))[0],200)
|
||||
self.assertGreater(result.getpixel((10,33))[2],200)
|
||||
self.assertNotEqual(thread_ids[0],threading.get_ident())
|
||||
self.assertEqual(sorted(p.name for p in mirror.directory.iterdir()),['frame.json','latest.jpg'])
|
||||
|
||||
def test_readback_error_does_not_break_ui_or_hold_slot(self):
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
mirror=self.make(folder)
|
||||
self.assertFalse(mirror.capture(Mock(side_effect=RuntimeError('GPU read failed')),started=True,session_id=SESSION))
|
||||
self.assertIn('GPU read failed',mirror.last_error)
|
||||
self.assertTrue(mirror._busy.acquire(blocking=False));mirror._busy.release()
|
||||
|
||||
def test_native_texture_orientation_and_release(self):
|
||||
image=NS(data=b'abcd',width=1,height=1)
|
||||
rl=NS(rl_draw_render_batch_active=Mock(),load_image_from_texture=Mock(return_value=image),
|
||||
load_image_from_screen=Mock(return_value=image),unload_image=Mock(),ffi=NS(buffer=lambda data,count:data[:count]))
|
||||
self.assertEqual(read_ui_rgba(rl,NS(_render_texture=NS(texture='texture'))),(b'abcd',1,1,True))
|
||||
rl.load_image_from_texture.assert_called_once_with('texture')
|
||||
rl.unload_image.assert_called_once_with(image)
|
||||
self.assertEqual(read_ui_rgba(rl,NS(_render_texture=None)),(b'abcd',1,1,False))
|
||||
|
||||
|
||||
if __name__=='__main__':unittest.main()
|
||||
@@ -7,7 +7,7 @@ import time
|
||||
import unittest
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace as NS
|
||||
from unittest.mock import patch
|
||||
from unittest.mock import Mock, patch
|
||||
|
||||
ROOT = Path(__file__).resolve().parents[2]
|
||||
sys.path.insert(0, str(ROOT/'roadscore/prototype'))
|
||||
@@ -57,6 +57,8 @@ class OverlayRenderOrderTests(unittest.TestCase):
|
||||
home._render(home.rect)
|
||||
nav._render(None)
|
||||
yield True
|
||||
state.roadscore_replay_prompt_active=True
|
||||
yield True
|
||||
gui=NS(render=native_frames,width=536,height=240,font=lambda _:None)
|
||||
state=NS(started=True,sm={s:NS(alertSize=NS(raw=0)) for s in ('selfdriveState','starpilotSelfdriveState')})
|
||||
modules={'pyray':NS(), 'openpilot.system.ui.lib.application':NS(gui_app=gui,FontWeight=NS(NORMAL=0,SEMI_BOLD=1)),
|
||||
@@ -64,9 +66,10 @@ class OverlayRenderOrderTests(unittest.TestCase):
|
||||
'openpilot.selfdrive.ui.onroad.starpilot.navigation_card':NS(NavigationCardRenderer=Nav),
|
||||
'openpilot.selfdrive.ui.ui_state':NS(ui_state=state)}
|
||||
with tempfile.TemporaryDirectory() as folder:
|
||||
path=Path(folder)/'roadscore_status.json';path.write_text(json.dumps({'readiness':'PREPARING','profile':'prism'}))
|
||||
path=Path(folder)/'roadscore_status.json';path.write_text(json.dumps({'readiness':'PREPARING','profile':'prism','presentation_session_id':'current-session'}))
|
||||
os.utime(path,(time.time()-300,time.time()-300))
|
||||
with patch.dict(sys.modules,modules),patch.dict(os.environ,{'ROADSCORE_OVERLAY':'1','ROADSCORE_STATUS_FILE':str(path),'ROADSCORE_CAPTURE_TIMING':'1'}),patch.object(overlay,'draw_panel',return_value=overlay_view({'readiness':'PREPARING'})) as draw:
|
||||
mirror=NS(capture=Mock(),close=Mock())
|
||||
with patch('screen_mirror.ScreenMirror.from_environ',return_value=mirror),patch.dict(sys.modules,modules),patch.dict(os.environ,{'ROADSCORE_OVERLAY':'1','ROADSCORE_STATUS_FILE':str(path),'ROADSCORE_CAPTURE_TIMING':'1'}),patch.object(overlay,'draw_panel',return_value=overlay_view({'readiness':'PREPARING'})) as draw:
|
||||
overlay.install();list(gui.render())
|
||||
rows=[json.loads(line) for line in path.with_name('ui_frames.jsonl').read_text().splitlines()]
|
||||
self.assertTrue(rows[0]['overlay_visible'])
|
||||
@@ -76,6 +79,12 @@ class OverlayRenderOrderTests(unittest.TestCase):
|
||||
self.assertFalse(rows[1]['overlay_visible'])
|
||||
self.assertEqual(rows[1]['overlay_hidden_reason'],'native_navigation')
|
||||
self.assertEqual(draw.call_count,1)
|
||||
self.assertEqual(rows[2]['overlay_hidden_reason'],'native_alert')
|
||||
self.assertEqual(mirror.capture.call_count,3)
|
||||
for call in mirror.capture.call_args_list:
|
||||
self.assertTrue(call.kwargs['started'])
|
||||
self.assertEqual(call.kwargs['session_id'],'current-session')
|
||||
mirror.close.assert_called_once()
|
||||
self.assertTrue(path.with_name('overlay_status.json').exists())
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user