diff --git a/roadscore/MAC_SHOWCASE.md b/roadscore/MAC_SHOWCASE.md index 4d90d5a42e..a455143c9d 100644 --- a/roadscore/MAC_SHOWCASE.md +++ b/roadscore/MAC_SHOWCASE.md @@ -16,4 +16,14 @@ The archive must contain `dry.wav`, `launch.json`, `audio_blocks.jsonl`, `replay `--score-archive PATH` selects another compatible complete recording explicitly. `--check` validates the local prerequisites without starting playback. `--muted`, `--duration SECONDS`, and `--no-browser` support verification. Live generation and the existing stored final-score replay retain their separate launch modes. -The Mac uses its selected system audio output. Bluetooth delay estimates from the comma do not transfer automatically to a different Mac output. Current controls do not synchronize a second device's playback. +The Mac uses its selected system audio output. Bluetooth delay estimates from the comma do not transfer automatically to a different Mac output. + +Optional paired controls require an explicit Galaxy LAN URL each launch. Replace this example with the comma's private IPv4 address: + +```sh +./onroad --roadscore route1 --prepared-showcase --paired-comma http://192.168.1.50:8082 +``` + +The default is unpaired. The comma must already have a fresh, ready RoadScore replay running while offroad. Each button writes the Mac action first, then forwards only that engagement or signal action in the background using the comma's own session ID. The page reports Mac state and whether the comma received or applied the latest command. An unavailable comma leaves the Mac buttons usable. A timeout can mean delivery is unknown; writes are not retried or rolled back automatically. + +The local control page remains bound to `127.0.0.1`. Pairing permits only an explicit private IPv4 address on port 8082, without redirects or discovery. The existing LAN Galaxy service uses offroad and fresh replay-session checks; it has no separate HTTP login or encryption. `--check` does not contact the peer. Pairing neither starts the comma nor synchronizes its video or music with the Mac. diff --git a/roadscore/prototype/mac_showcase.html b/roadscore/prototype/mac_showcase.html index 0bac9accde..885aa9d00b 100644 --- a/roadscore/prototype/mac_showcase.html +++ b/roadscore/prototype/mac_showcase.html @@ -3,5 +3,31 @@
Prepared Prism music · live presentation controls
These controls simulate the replay display and music only.
Waiting for replay…
No model generation, vehicle control, or connection to the comma. - +Waiting for replay…
Mac waiting; comma not paired
No model generation or vehicle control. Optional comma pairing forwards controls only. Video and music are not synchronized. + diff --git a/roadscore/prototype/mac_showcase.py b/roadscore/prototype/mac_showcase.py index a6fd291480..2f2059dba7 100644 --- a/roadscore/prototype/mac_showcase.py +++ b/roadscore/prototype/mac_showcase.py @@ -12,6 +12,7 @@ import threading import time import uuid from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from paired_demo_controls import DISCLOSURE, GalaxyPeer, PairedDemoControls HERE = Path(__file__).resolve().parent @@ -22,7 +23,7 @@ def write_json(path, value): temporary.replace(path) -def control_server(project, out, shared, port): +def control_server(project, out, shared, port, forwarder=None): spec = importlib.util.spec_from_file_location('showcase_galaxy', project/'starpilot/system/the_galaxy/roadscore.py') module = importlib.util.module_from_spec(spec) spec.loader.exec_module(module) @@ -30,6 +31,27 @@ def control_server(project, out, shared, port): (operator_root/'results').mkdir(parents=True) (operator_root/'results/current').symlink_to(out) operator = module.Operator(operator_root, device=False, offroad=lambda: True) + control_lock = threading.Lock() + sequence = 0 + paired = dict(enabled=forwarder is not None, sequence=0, request=None, music_video_synchronized=False, disclosure=DISCLOSURE, + targets={'comma':dict(status='idle' if forwarder else 'disabled',acknowledged=False,applied=False,error=None)}) + def paired_snapshot(demo): + with control_lock:result = {**paired, 'targets':dict(paired['targets'])} + result['targets']['mac'] = dict(status='ready' if demo['available'] else 'unready', + mode=demo['mode'],signal_mode=demo['signal_mode'],acknowledged=False,applied=False) + request = result['request'] + if request: + applied = demo['available'] and demo['session_id'] == request['session_id'] and demo[request['field']] == request['value'] + result['targets']['mac'].update(acknowledged=True,applied=bool(applied),action=request['action'],value=request['value']) + return result + def failed_peer(): + return dict(targets={'comma':dict(status='failed',acknowledged=False,applied=False,error='peer_submission_failed')}) + def received(future, request_sequence): + try:result = future.result() + except Exception:result = failed_peer() + with control_lock: + # A slow result cannot replace the status of a newer button action. + if sequence == request_sequence:paired['targets'] = result['targets'] class Handler(BaseHTTPRequestHandler): def log_message(self, *args):pass def send_json(self, value, status=200): @@ -41,23 +63,46 @@ def control_server(project, out, shared, port): self.end_headers();self.wfile.write(data) def do_GET(self): if self.path == '/status': - self.send_json({'demo': operator.demo_status(), **shared}) + demo = operator.demo_status() + self.send_json({'demo': demo, **shared, 'paired_controls':paired_snapshot(demo)}) elif self.path == '/': data = (HERE/'mac_showcase.html').read_bytes() self.send_response(200);self.send_header('Content-Type', 'text/html; charset=utf-8');self.send_header('Content-Length', str(len(data)));self.end_headers();self.wfile.write(data) else:self.send_error(404) def do_POST(self): + nonlocal sequence if self.path != '/control':self.send_error(404);return - if not module.control_origin_allowed(self.headers.get('Origin'), self.headers.get('Host'), 'http', self.headers.get('Sec-Fetch-Site')): + allowed_hosts = {f'127.0.0.1:{self.server.server_port}',f'localhost:{self.server.server_port}'} + if (self.headers.get('Host') not in allowed_hosts + or not module.control_origin_allowed(self.headers.get('Origin'), self.headers.get('Host'), 'http', self.headers.get('Sec-Fetch-Site'))): self.send_json({'error':'Cross-origin controls are not allowed'},403);return try: length = int(self.headers.get('Content-Length', '0')) if not 0 < length <= 1024:raise ValueError('Invalid control size') data = json.loads(self.rfile.read(length)) field = 'signal_mode' if 'signal_mode' in data else 'mode' - self.send_json(operator.demo_engagement(data, field)) + future = None + with control_lock: + # Serialize local writes with their submissions. The peer is optional; + # waiting for its network/status response never holds this lock. + result = operator.demo_engagement(data, field) + sequence += 1;request_sequence = sequence + action = 'demo_signal' if field == 'signal_mode' else 'demo_engagement' + paired.update(sequence=sequence,request=dict(action=action,field=field,value=data[field],session_id=data['session_id'])) + if forwarder is not None: + try: + future = forwarder.submit(action, data[field]) + paired['targets'] = {'comma':dict(status='pending',action=action,value=data[field], + acknowledged=False,applied=False,error=None)} + except Exception:paired['targets'] = failed_peer()['targets'] + if future is not None:future.add_done_callback(lambda completed:received(completed,request_sequence)) + self.send_json({**result,'control_sequence':request_sequence}) except (ValueError, TypeError) as error:self.send_json({'error':str(error)},400) - server = ThreadingHTTPServer(('127.0.0.1',port),Handler) + class ControlServer(ThreadingHTTPServer): + def server_close(self): + if forwarder is not None:forwarder.close() + super().server_close() + server = ControlServer(('127.0.0.1',port),Handler) threading.Thread(target=server.serve_forever,daemon=True).start() return server @@ -97,8 +142,6 @@ def audio_worker(a): max_drift = 0. first_frame = None shared = dict(state='PREPARING',audio_s=0.,muted=a.muted) - server = control_server(a.project_root,a.out,shared,a.port) - write_json(a.out/'controls.json',{'url':f'http://127.0.0.1:{server.server_port}/'}) def callback(out, n, ti, status): nonlocal position,rendered,done,flags,max_drift,first_frame out.fill(0) @@ -121,7 +164,11 @@ def audio_worker(a): errors.append(repr(error));done=True started=None;last_status=0. trace=(a.out/'presentation.jsonl').open('w',buffering=1) + forwarder = PairedDemoControls(GalaxyPeer(a.paired_comma,allow_lan_http=True),enabled=True) if a.paired_comma else None + server = None try: + server = control_server(a.project_root,a.out,shared,a.port,forwarder) + write_json(a.out/'controls.json',{'url':f'http://127.0.0.1:{server.server_port}/'}) with sd.OutputStream(device=a.audio_device,samplerate=rate,channels=2,blocksize=960,dtype='float32',callback=callback): (a.out/'prepared_ready').write_text('ready') while not done: @@ -154,8 +201,14 @@ def audio_worker(a): if started and now-started>=a.duration:break if started and now-received['modelV2']>2 and (position or 0)/rate < len(audio)/rate-2:raise RuntimeError('Replay model stream stopped before prepared audio ended') finally: - write_json(a.out/'audio_drained.json',dict(wall=time.monotonic(),drained=not errors)) - server.shutdown();trace.close() + try:write_json(a.out/'audio_drained.json',dict(wall=time.monotonic(),drained=not errors)) + finally: + try: + if server is not None: + try:server.shutdown() + finally:server.server_close() + elif forwarder is not None:forwarder.close() + finally:trace.close() write_json(a.out/'prepared_summary.json',dict(generation_invoked=False,source=str(a.score_archive),first_source_frame=first_frame,last_source_frame=position,sample_rate=rate,portaudio_flags=flags,max_clock_error_seconds=max_drift,callback_errors=errors,muted=a.muted,session_id=session,manual_scope='isolated replay display and presentation only')) if errors:raise RuntimeError(errors[0]) if max_drift>.05:raise RuntimeError('Prepared audio clock drift exceeded 50 ms') @@ -172,6 +225,12 @@ def parser(): p.add_argument('--out',type=Path) p.add_argument('--duration',type=float,default=float('inf')) p.add_argument('--port',type=int,default=0) + def paired_url(value): + try:GalaxyPeer(value,allow_lan_http=True) + except ValueError as error:raise argparse.ArgumentTypeError(str(error)) from error + return value + p.add_argument('--paired-comma',type=paired_url,default=None,metavar='URL', + help='Optional explicit Galaxy LAN URL http://PRIVATE_IPV4:8082; forwards replay controls only') p.add_argument('--muted',action='store_true');p.add_argument('--headless',action='store_true') p.add_argument('--no-browser',action='store_true');p.add_argument('--audio-device');p.add_argument('--check',action='store_true');p.add_argument('--audio-worker',action='store_true',help=argparse.SUPPRESS) return p @@ -217,7 +276,7 @@ def main(): args[args.index('--data_dir')+1]=str(local) if '--no-hw-decoder' in args:args.remove('--no-hw-decoder') write_json(out/'status.json',dict(readiness='PREPARING',style='Prism',compute='prepared-core')) - write_json(out/'launch.json',dict(core_sha256=hashlib.sha256((a.score_archive/'dry.wav').read_bytes()).hexdigest(),curve_plan_sha256=hashlib.sha256(a.curve_plan.read_bytes()).hexdigest() if a.curve_plan else None,mode='prepared-interactive-showcase',route=a.route,source=str(a.score_archive),runtime=str(rt),native_replay_args=args,generation_invoked=False,network_required=False,session_id=session,muted=a.muted or a.headless)) + write_json(out/'launch.json',dict(core_sha256=hashlib.sha256((a.score_archive/'dry.wav').read_bytes()).hexdigest(),curve_plan_sha256=hashlib.sha256(a.curve_plan.read_bytes()).hexdigest() if a.curve_plan else None,mode='prepared-interactive-showcase',route=a.route,source=str(a.score_archive),runtime=str(rt),native_replay_args=args,generation_invoked=False,network_required=False,session_id=session,muted=a.muted or a.headless,paired_comma=a.paired_comma,paired_controls_scope=DISCLOSURE)) check=subprocess.run([str(py),'-c','from cereal import messaging; import sounddevice,soundfile; from prepared_core import load_archive; import sys; a,r,m=load_archive(sys.argv[1],sys.argv[2]); print("Prepared core:",len(a)/r,"seconds; local replay ready")',str(a.score_archive),a.route],cwd=rt,env=env) if check.returncode:raise SystemExit(check.returncode) if a.check:return @@ -233,7 +292,7 @@ def main(): with (out/'seed.log').open('w') as log:subprocess.run([str(py),str(rt/'tools/replay/onroad_config.py'),'seed',*args],env=env,cwd=rt,stdout=log,stderr=subprocess.STDOUT,check=True) seed_finished=time.monotonic() if not a.headless:ui=start([str(py),str(HERE/'normal_ui_audit.py')],'ui') - audio=start([str(py),str(__file__),'--audio-worker',a.route,'--score-archive',str(a.score_archive),'--project-root',str(project),'--out',str(out),'--duration',str(a.duration),'--port',str(a.port)]+(['--curve-plan',str(a.curve_plan)] if a.curve_plan else [])+(['--muted'] if a.muted or a.headless else [])+(['--audio-device',a.audio_device] if a.audio_device else []),'audio') + audio=start([str(py),str(__file__),'--audio-worker',a.route,'--score-archive',str(a.score_archive),'--project-root',str(project),'--out',str(out),'--duration',str(a.duration),'--port',str(a.port)]+(['--curve-plan',str(a.curve_plan)] if a.curve_plan else [])+(['--muted'] if a.muted or a.headless else [])+(['--audio-device',a.audio_device] if a.audio_device else [])+(['--paired-comma',a.paired_comma] if a.paired_comma else []),'audio') deadline=time.monotonic()+30 while not (out/'prepared_ready').exists(): if audio.poll() is not None or time.monotonic()>deadline:raise RuntimeError('Prepared audio did not become ready; see '+str(out/'audio.log')) diff --git a/roadscore/prototype/test_mac_showcase_controls.py b/roadscore/prototype/test_mac_showcase_controls.py new file mode 100644 index 0000000000..efe38d658a --- /dev/null +++ b/roadscore/prototype/test_mac_showcase_controls.py @@ -0,0 +1,176 @@ +"""Loopback HTTP integration with mocked peer transport; no audio or peer access.""" +from concurrent.futures import Future +import contextlib +import copy +import http.client +import io +import json +from pathlib import Path +import tempfile +import threading +import time +import unittest + +from mac_showcase import control_server, parser +from paired_demo_controls import GalaxyPeer, PairedDemoControls + + +class ManualForwarder: + def __init__(self, out): + self.out = out + self.calls = [] + self.closed = False + + def submit(self, action, value): + # This mock observes the real local Galaxy command file before forwarding. + command = json.loads((self.out/'demo_engagement.json').read_text()) + field = 'mode' if action == 'demo_engagement' else 'signal_mode' + if command[field] != value:raise AssertionError('Peer submitted before local write') + future = Future() + self.calls.append((action, value, future)) + return future + + def close(self):self.closed = True + + +class ControlTests(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory() + self.addCleanup(self.temp.cleanup) + self.out = Path(self.temp.name) + self.project = Path(__file__).resolve().parents[2] + self.state = dict(input_mode='replay',route='fixture-route',presentation_session_id='mac-session', + command_wall=time.monotonic(),engagement_presentation={'enabled':True}) + self.write_state() + self.server = None + self.addCleanup(self.stop) + + def write_state(self): + self.state['command_wall'] = time.monotonic() + (self.out/'status.json').write_text(json.dumps(self.state)) + + def start(self, forwarder=None): + self.server = control_server(self.project,self.out,dict(state='READY',audio_s=1.,muted=True),0,forwarder) + self.assertEqual(self.server.server_address[0], '127.0.0.1') + + def stop(self): + if self.server is not None: + try:self.server.shutdown() + finally:self.server.server_close();self.server=None + + def request(self, method='GET', path='/status', body=None, headers=None): + connection = http.client.HTTPConnection('127.0.0.1',self.server.server_port,timeout=1) + try: + connection.request(method,path,None if body is None else json.dumps(body),headers=headers or {'Content-Type':'application/json'}) + response = connection.getresponse() + return response.status,json.loads(response.read()) + finally:connection.close() + + def command(self, value='engaged', field='mode', session='mac-session', headers=None): + return self.request('POST','/control',{'session_id':session,field:value},headers) + + def test_default_has_no_peer_and_local_ack_is_not_application(self): + self.start() + status, result = self.command() + self.assertEqual(status,200) + self.assertEqual(result['requested_mode'],'engaged') + status, result = self.request() + paired = result['paired_controls'] + self.assertFalse(paired['enabled']) + self.assertEqual(paired['targets']['comma']['status'],'disabled') + self.assertTrue(paired['targets']['mac']['acknowledged']) + self.assertFalse(paired['targets']['mac']['applied']) + self.state['demo_engagement_mode']='engaged';self.write_state() + self.assertTrue(self.request()[1]['paired_controls']['targets']['mac']['applied']) + + def test_local_write_then_peer_submission_and_session_rejection(self): + peer = ManualForwarder(self.out) + self.start(peer) + self.assertEqual(self.command(session='old-session')[0],400) + self.assertEqual(peer.calls,[]) + self.assertEqual(self.command('left','signal_mode')[0],200) + paired = self.request()[1]['paired_controls'] + self.assertEqual(paired['sequence'],1) + self.assertEqual(paired['request']['session_id'],'mac-session') + self.assertEqual(paired['targets']['comma']['status'],'pending') + self.assertEqual(peer.calls[0][:2],('demo_signal','left')) + self.assertFalse(paired['music_video_synchronized']) + + def test_older_peer_result_cannot_replace_newer_action(self): + peer = ManualForwarder(self.out) + self.start(peer) + self.assertEqual(self.command('left','signal_mode')[1]['control_sequence'],1) + self.assertEqual(self.command('right','signal_mode')[1]['control_sequence'],2) + peer.calls[0][2].set_result({'targets':{'comma':{'status':'applied','value':'left','applied':True}}}) + paired = self.request()[1]['paired_controls'] + self.assertEqual(paired['sequence'],2) + self.assertEqual(paired['targets']['comma']['status'],'pending') + self.assertEqual(paired['targets']['comma']['value'],'right') + peer.calls[1][2].set_result({'targets':{'comma':{'status':'failed','value':'right','error':'peer_not_ready_for_replay'}}}) + status = self.request()[1] + self.assertTrue(status['demo']['available']) + self.assertEqual(status['paired_controls']['targets']['comma']['error'],'peer_not_ready_for_replay') + + def test_cross_origin_and_rebinding_host_rejected_before_both_writes(self): + peer = ManualForwarder(self.out) + self.start(peer) + for headers in ({'Origin':'https://other.example','Sec-Fetch-Site':'cross-site'}, + {'Host':'other.example','Origin':'http://other.example'}): + self.assertEqual(self.command(headers=headers)[0],403) + self.assertEqual(peer.calls,[]) + self.assertFalse((self.out/'demo_engagement.json').exists()) + origin = f'http://127.0.0.1:{self.server.server_port}' + self.assertEqual(self.command(headers={'Origin':origin,'Sec-Fetch-Site':'same-origin'})[0],200) + + def test_peer_future_failure_keeps_local_controls_ready(self): + peer = ManualForwarder(self.out) + self.start(peer) + self.assertEqual(self.command()[0],200) + peer.calls[0][2].set_exception(RuntimeError('Mock peer error')) + result = self.request()[1] + self.assertTrue(result['demo']['available']) + self.assertEqual(result['paired_controls']['targets']['comma']['error'],'peer_submission_failed') + self.assertEqual(self.command('disengaged')[0],200) + self.stop() + self.assertTrue(peer.closed) + + def test_slow_mock_peer_never_delays_local_http_and_uses_own_session(self): + entered, release = threading.Event(), threading.Event() + calls = [] + peer_status = dict(available=True,offroad=True,state='READY',live={'enabled':False}, + demo=dict(available=True,session_id='comma-session',mode='recorded',signal_mode='recorded')) + def transport(method, url, payload, headers, timeout): + calls.append((method,payload)) + if len(calls)==1:entered.set();release.wait(1) + if method=='GET':return copy.deepcopy(peer_status) + field = 'signal_mode' if url.endswith('demo_signal') else 'mode' + peer_status['demo'][field]=payload[field] + return {'requested_'+field:payload[field],'demo':copy.deepcopy(peer_status['demo'])} + forwarder = PairedDemoControls(GalaxyPeer('http://192.168.1.50:8082',allow_lan_http=True),enabled=True,transport=transport) + self.start(forwarder) + try: + before = time.monotonic() + self.assertEqual(self.command('left','signal_mode')[0],200) + self.assertLess(time.monotonic()-before,.25) + self.assertTrue(entered.wait(.5)) + self.assertEqual(self.request()[1]['paired_controls']['targets']['comma']['status'],'pending') + finally:release.set() + deadline = time.monotonic()+1 + while True: + result = self.request()[1]['paired_controls']['targets']['comma'] + if result['status']!='pending' or time.monotonic()>deadline:break + time.sleep(.01) + self.assertTrue(result['applied']) + self.assertEqual(calls[1][1],{'session_id':'comma-session','signal_mode':'left'}) + self.assertEqual(json.loads((self.out/'demo_engagement.json').read_text())['session_id'],'mac-session') + + def test_cli_pairing_is_explicit_and_lan_only_without_credentials(self): + self.assertIsNone(parser().parse_args([]).paired_comma) + value = 'http://192.168.1.50:8082' + self.assertEqual(parser().parse_args(['--paired-comma',value]).paired_comma,value) + for value in ('http://device.local:8082','http://8.8.8.8:8082','http://192.168.1.50:8082/mobile/#/roadscore'): + with self.subTest(url=value),contextlib.redirect_stderr(io.StringIO()),self.assertRaises(SystemExit): + parser().parse_args(['--paired-comma',value]) + + +if __name__=='__main__':unittest.main()