Add opt-in asynchronous replay control forwarding helper

This commit is contained in:
firestar5683
2026-09-20 10:19:46 -07:00
parent 4de3e304c0
commit f7f9cf047c
2 changed files with 556 additions and 0 deletions
+254
View File
@@ -0,0 +1,254 @@
"""Optional replay-control forwarding; no video, audio, or transport synchronization.
Integration contract: apply the local action independently, then call submit from
the HTTP/control thread and poll its Future. Never wait on it in an audio callback.
Nothing is contacted at import/construction, or unless enabled=True and an
explicit Galaxy peer were supplied. There is no discovery, pairing,
credential lookup, retry of writes, or launch/vehicle-control endpoint.
Galaxy's demo.available attests to a <=2 second fresh replay app status; the POST
checks that status and its own presentation session again. Its requested_* reply
only acknowledges command receipt. A subsequent status read confirms adoption.
Timeout after a POST can mean the peer applied it: report uncertainty, never
retry or roll back either target. Local playback remains the caller's concern.
"""
from concurrent.futures import Future
from dataclasses import dataclass, field
import json
import ipaddress
import math
import re
import threading
import time
from urllib.error import HTTPError, URLError
from urllib.parse import quote, unquote, urlsplit
from urllib.request import HTTPRedirectHandler, ProxyHandler, Request, build_opener
DISCLOSURE = 'Controls only; video and music are not synchronized.'
ACTIONS = {'demo_engagement': ('mode', ('recorded', 'engaged', 'disengaged')),
'demo_signal': ('signal_mode', ('recorded', 'left', 'right', 'off'))}
@dataclass(frozen=True)
class GalaxyPeer:
"""Explicit Galaxy tunnel with cookie, or separately opted-in private LAN URL.
Cookie format comes from Galaxy's /api/galaxy/session UI. This helper never
calls that endpoint or reads credential files. The existing LAN service has no
login check: it authorizes replay commands by offroad state and fresh session.
allow_lan_http=True permits only a literal RFC1918 IPv4 address on port 8082,
without credentials. It is not an authenticated/encrypted transport. The caller
must explicitly select that local device. URLs/secrets stay out of repr.
"""
base_url: str = field(repr=False)
session_cookie: str | None = field(default=None, repr=False)
name: str = 'comma'
allow_lan_http: bool = False
def __post_init__(self):
try:
url = urlsplit(self.base_url)
plain = (not any(char.isspace() for char in self.base_url)
and url.hostname and not url.username and not url.password
and not url.query and not url.fragment)
tunnel = plain and url.scheme == 'https' and url.port in (None, 443) and re.fullmatch(r'/[A-Za-z0-9]{16}/?', url.path)
lan = False
if plain and url.scheme == 'http' and self.allow_lan_http is True and url.port == 8082 and url.path in ('', '/'):
address = ipaddress.IPv4Address(url.hostname)
lan = any(address in ipaddress.IPv4Network(network) for network in ('10.0.0.0/8', '172.16.0.0/12', '192.168.0.0/16'))
except (TypeError, ValueError):
tunnel = lan = False
if not (tunnel or lan) or type(self.allow_lan_http) is not bool:
raise ValueError('Explicit HTTPS Galaxy tunnel or opted-in private IPv4 LAN address on port 8082 required')
if not isinstance(self.name, str) or not re.fullmatch(r'[A-Za-z0-9_-]{1,40}', self.name):
raise ValueError('A short peer name is required')
if lan:
if self.session_cookie is not None:
raise ValueError('Credentials are not sent over LAN HTTP')
else:
credential = unquote(self.session_cookie) if isinstance(self.session_cookie, str) else ''
if not re.fullmatch(re.escape(url.path.strip('/')) + r':[a-fA-F0-9]{64}', credential):
raise ValueError('Explicit Galaxy session cookie must match the configured tunnel slug')
object.__setattr__(self, 'session_cookie', quote(credential, safe=''))
object.__setattr__(self, 'base_url', self.base_url.rstrip('/'))
class _NoRedirect(HTTPRedirectHandler):
def redirect_request(self, request, fp, code, msg, headers, newurl):
raise ValueError('Peer redirect rejected')
def _http_json(method, url, payload, headers, timeout):
"""One bounded-size response, no redirect, ambient proxy, or credential cache."""
body = None if payload is None else json.dumps(payload).encode()
request = Request(url, data=body, headers=headers, method=method)
with build_opener(ProxyHandler({}), _NoRedirect()).open(request, timeout=timeout) as response:
if response.status != 200:
raise ValueError('Unexpected peer response')
if response.headers.get_content_type() != 'application/json':
raise ValueError('Expected peer JSON')
content = response.read(65537)
if len(content) > 65536:
raise ValueError('Peer response too large')
result = json.loads(content)
if not isinstance(result, dict):
raise ValueError('Expected peer status object')
return result
def _session(status):
if not isinstance(status, dict):
raise ValueError('invalid_peer_status')
demo = status.get('demo')
live = status.get('live', {})
if (status.get('available') is not True or status.get('offroad') is not True
or status.get('state') not in ('READY', 'GENERATING')
or not isinstance(demo, dict) or demo.get('available') is not True
or (isinstance(live, dict) and live.get('enabled') is True)):
raise ValueError('peer_not_ready_for_replay')
# Current Galaxy omits these fields; reject contradictory richer replies too.
if (status.get('input_mode', 'replay') != 'replay'
or status.get('mode') in ('stored', 'stored-score', 'live')
or status.get('route') == 'live' or status.get('judging_locked') is True):
raise ValueError('peer_not_ready_for_replay')
session = demo.get('session_id')
if (not isinstance(session, str) or not 1 <= len(session) <= 256
or any(ord(char) < 32 for char in session)
or demo.get('mode') not in ACTIONS['demo_engagement'][1]
or demo.get('signal_mode') not in ACTIONS['demo_signal'][1]):
raise ValueError('invalid_peer_session')
return session
class PairedDemoControls:
"""At most one background action; overloaded/disabled submissions finish now.
submit(action, value) returns a Future of {targets: {name: result}, disclosure,
music_video_synchronized: False}. Caller combines its own local target outcome.
Results distinguish disabled/rejected/busy/failed/timed_out/acknowledged/applied.
There is no queue of stale controls, and each action fetches the peer session.
"""
def __init__(self, peer=None, *, enabled=False, timeout=1.5, transport=_http_json):
if type(enabled) is not bool or (peer is not None and not isinstance(peer, GalaxyPeer)):
raise ValueError('Explicit peer and boolean enable required')
if type(timeout) not in (int, float) or not math.isfinite(timeout) or not .1 <= timeout <= 5:
raise ValueError('Timeout must be between 0.1 and 5 seconds')
if enabled and peer is None:
raise ValueError('Enabling paired controls requires an explicit peer')
self.peer, self.enabled, self.timeout, self.transport = peer, enabled, timeout, transport
self._lock = threading.Lock()
self._active = None
self._closed = False
def _result(self, status, action, value, *, acknowledged=False, applied=False,
error=None, session_id=None, delivery_unknown=False):
name = self.peer.name if self.peer else 'peer'
return {'targets': {name: {'status': status, 'action': action, 'value': value,
'acknowledged': acknowledged, 'applied': applied,
'error': error, 'session_id': session_id,
'delivery_unknown': delivery_unknown}},
'disclosure': DISCLOSURE, 'music_video_synchronized': False}
def submit(self, action, value):
"""Nonblocking submission for a control/HTTP thread, never an audio callback."""
future = Future()
# A caller cancellation cannot undo a possibly delivered peer command.
future.set_running_or_notify_cancel()
if type(action) is not str or type(value) is not str or action not in ACTIONS or value not in ACTIONS[action][1]:
future.set_result(self._result('rejected', action, value, error='unsupported_replay_action'))
return future
with self._lock:
if self._closed or not self.enabled:
future.set_result(self._result('disabled', action, value))
return future
if self._active is not None:
future.set_result(self._result('busy', action, value, error='peer_action_in_progress'))
return future
operation = {'future': future, 'action': action, 'value': value,
'sent': False, 'acknowledged': False, 'session': None, 'settled': False}
self._active = operation
deadline = time.monotonic() + self.timeout
timer = threading.Timer(self.timeout, self._expire, args=(operation,))
timer.daemon = True
timer.start()
threading.Thread(target=self._run, args=(operation, deadline, timer),
name='paired-replay-control', daemon=True).start()
return future
def _finish(self, operation, status, error=None, applied=False):
with self._lock:
if operation['settled']:
return
operation['settled'] = True
result = self._result(status, operation['action'], operation['value'], error=error, applied=applied,
acknowledged=operation['acknowledged'], session_id=operation['session'],
delivery_unknown=operation['sent'] and not operation['acknowledged'])
# Future callbacks may call submit/close; do not invoke them under our lock.
operation['future'].set_result(result)
def _expire(self, operation):
self._finish(operation, 'acknowledged' if operation['acknowledged'] else 'timed_out',
'peer_application_unconfirmed' if operation['acknowledged'] else 'peer_timeout')
def _run(self, operation, deadline, timer):
def request(method, endpoint, payload=None):
remaining = deadline - time.monotonic()
with self._lock:
if self._closed or operation['settled'] or remaining <= 0:
raise TimeoutError()
if method == 'POST':
operation['sent'] = True
headers = {'Accept': 'application/json', 'Content-Type': 'application/json',
'Cache-Control': 'no-store'}
if self.peer.session_cookie is not None:
headers['Cookie'] = 'galaxy_session=' + self.peer.session_cookie
result = self.transport(method, self.peer.base_url + '/api/roadscore/' + endpoint,
payload, headers, min(remaining, .75))
if time.monotonic() >= deadline:
raise TimeoutError()
return result
try:
action, value = operation['action'], operation['value']
field_name = ACTIONS[action][0]
session = _session(request('GET', 'status'))
operation['session'] = session
reply = request('POST', action, {'session_id': session, field_name: value})
if (not isinstance(reply, dict) or reply.get('requested_' + field_name) != value
or not isinstance(reply.get('demo'), dict) or reply['demo'].get('available') is not True
or reply['demo'].get('session_id') != session):
raise ValueError('invalid_peer_acknowledgement')
operation['acknowledged'] = True
while True:
status = request('GET', 'status')
if _session(status) != session:
raise ValueError('peer_session_changed')
if status['demo'][field_name] == value:
self._finish(operation, 'applied', applied=True)
break
# Only status reads are repeated. Never resend a possibly applied write.
time.sleep(min(.05, max(0., deadline - time.monotonic())))
except TimeoutError:
self._expire(operation)
except HTTPError as error:
self._finish(operation, 'failed', 'peer_http_' + str(error.code))
except (OSError, URLError):
self._finish(operation, 'failed', 'peer_transport_error')
except Exception as error:
known = {'peer_not_ready_for_replay', 'invalid_peer_status', 'invalid_peer_session',
'invalid_peer_acknowledgement', 'peer_session_changed'}
self._finish(operation, 'failed', str(error) if str(error) in known else 'invalid_peer_response')
finally:
timer.cancel()
with self._lock:
if self._active is operation:
self._active = None
def close(self):
"""Stop accepting commands without waiting for a blocked peer transport."""
with self._lock:
self._closed = True
active = self._active
if active is not None:
self._finish(active, 'failed', 'forwarder_closed')
@@ -0,0 +1,302 @@
"""Local injected transport only; no sockets, hardware, audio, or server changes."""
import copy
import importlib.util
import json
from pathlib import Path
import tempfile
import threading
import time
import unittest
from unittest.mock import patch
from urllib.error import HTTPError
from paired_demo_controls import GalaxyPeer, PairedDemoControls, _NoRedirect, _http_json
BASE = 'https://galaxy.firestar.link/ABCDEFGHIJKLMNOP'
COOKIE = 'ABCDEFGHIJKLMNOP%3A' + 'a' * 64
READY = {'available': True, 'offroad': True, 'state': 'READY', 'live': {'enabled': False},
'demo': {'available': True, 'session_id': 'independent-peer-session',
'mode': 'recorded', 'signal_mode': 'recorded'}}
def target(result):
return result['targets'].get('comma', result['targets'].get('peer'))
class MockGalaxy:
def __init__(self):
self.status = copy.deepcopy(READY)
self.calls = []
self.adopt = True
def __call__(self, method, url, payload, headers, timeout):
self.calls.append((method, url, payload, headers, timeout))
if method == 'GET':
return copy.deepcopy(self.status)
field = 'mode' if url.endswith('/demo_engagement') else 'signal_mode'
if self.adopt:
self.status['demo'][field] = payload[field]
return {'requested_' + field: payload[field], 'demo': copy.deepcopy(self.status['demo'])}
class PairedTests(unittest.TestCase):
def setUp(self):
self.transport = MockGalaxy()
self.peer = GalaxyPeer(BASE, COOKIE)
def forwarder(self, **options):
result = PairedDemoControls(self.peer, enabled=True, transport=self.transport, **options)
self.addCleanup(result.close)
return result
def test_default_disabled_even_with_explicit_peer(self):
for peer in (None, self.peer):
controls = PairedDemoControls(peer, transport=self.transport)
result = controls.submit('demo_signal', 'left').result(0)
self.assertEqual(target(result)['status'], 'disabled')
self.assertFalse(result['music_video_synchronized'])
self.assertEqual(self.transport.calls, [])
def test_explicit_enable_and_authenticated_https_required(self):
with self.assertRaises(ValueError):
PairedDemoControls(enabled=True)
for url in ('http://192.0.2.1:8082', 'https://host/', BASE + '?token=secret',
BASE.replace('https://', 'https://user:secret@'), BASE + '/other', BASE + '\n'):
with self.subTest(url=url), self.assertRaises(ValueError):
GalaxyPeer(url, COOKIE)
for cookie in ('', 'wrong', 'wrongslug%3A' + 'a' * 64, COOKIE + '\r\nInjected: value'):
with self.assertRaises(ValueError):
GalaxyPeer(BASE, cookie)
self.assertNotIn(COOKIE, repr(self.peer))
self.assertNotIn(BASE, repr(self.peer))
def test_no_import_or_constructor_network(self):
with patch('paired_demo_controls.build_opener', side_effect=AssertionError('No network')):
controls = PairedDemoControls(self.peer, enabled=True)
controls.close()
self.assertEqual(self.transport.calls, [])
def test_lan_requires_separate_opt_in_private_literal_and_exact_port(self):
for url in ('http://192.168.8.156:8082', 'http://10.1.2.3:8082', 'http://172.16.0.1:8082'):
with self.assertRaises(ValueError):
GalaxyPeer(url)
peer = GalaxyPeer(url, allow_lan_http=True)
controls = PairedDemoControls(peer, enabled=True, transport=self.transport)
self.addCleanup(controls.close)
self.assertTrue(target(controls.submit('demo_signal', 'left').result(1))['applied'])
self.assertNotIn('Cookie', self.transport.calls[-1][3])
self.assertTrue(self.transport.calls[-1][1].startswith(url + '/api/roadscore/'))
for url in ('http://8.8.8.8:8082', 'http://127.0.0.1:8082', 'http://169.254.1.1:8082',
'http://192.0.2.1:8082', 'http://192.168.8.156', 'http://192.168.8.156:80',
'http://device.local:8082', 'http://192.168.8.156:8082/mobile/',
'http://192.168.8.156:8082/#/roadscore', 'http://[::1]:8082'):
with self.subTest(url=url), self.assertRaises(ValueError):
GalaxyPeer(url, allow_lan_http=True)
with self.assertRaises(ValueError):
GalaxyPeer('http://192.168.8.156:8082', COOKIE, allow_lan_http=True)
def test_only_replay_actions_and_values(self):
controls = self.forwarder()
for action, value in (('power', 'off'), ('play', 'on'), ('start', 'on'), ('live', True),
('demo_signal', 'hazards'), ('demo_engagement', 'controlsAllowed'),
('demo_signal', {'signal_mode': 'left'}), ({}, 'left')):
result = controls.submit(action, value).result(0)
self.assertEqual(target(result)['status'], 'rejected')
self.assertEqual(self.transport.calls, [])
def test_independent_session_exact_post_and_per_target_ack(self):
result = self.forwarder().submit('demo_signal', 'right').result(1)
outcome = target(result)
self.assertEqual(outcome['status'], 'applied')
self.assertTrue(outcome['acknowledged'])
self.assertTrue(outcome['applied'])
self.assertFalse(outcome['delivery_unknown'])
self.assertFalse(result['music_video_synchronized'])
self.assertIn('not synchronized', result['disclosure'])
self.assertEqual([call[0] for call in self.transport.calls], ['GET', 'POST', 'GET'])
post = self.transport.calls[1]
self.assertEqual(post[1], BASE + '/api/roadscore/demo_signal')
self.assertEqual(post[2], {'session_id': 'independent-peer-session', 'signal_mode': 'right'})
self.assertEqual(post[3]['Cookie'], 'galaxy_session=' + COOKIE)
self.assertTrue(all(0 < call[4] <= .75 for call in self.transport.calls))
def test_engagement_uses_its_distinct_payload(self):
result = self.forwarder().submit('demo_engagement', 'disengaged').result(1)
self.assertTrue(target(result)['applied'])
self.assertEqual(self.transport.calls[1][2], {'session_id': 'independent-peer-session', 'mode': 'disengaged'})
def test_current_galaxy_operator_wire_contract_without_network(self):
source = Path(__file__).resolve().parents[2] / 'starpilot/system/the_galaxy/roadscore.py'
spec = importlib.util.spec_from_file_location('paired_test_galaxy', source)
galaxy = importlib.util.module_from_spec(spec)
spec.loader.exec_module(galaxy)
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
run = root / 'results/current'
run.mkdir(parents=True)
state = {'route': 'fixture-route', 'input_mode': 'replay', 'command_wall': time.monotonic(),
'presentation_session_id': 'galaxy-independent-session',
'engagement_presentation': {'enabled': True}}
(run / 'status.json').write_text(json.dumps(state))
operator = galaxy.Operator(root=root, device=True, offroad=lambda: True)
def local_transport(method, url, payload, headers, timeout):
if method == 'GET':
return {**READY, 'demo': operator.demo_status()}
result = operator.operate(url.rsplit('/', 1)[1], payload, True)
# Mock the replay app consuming the command after Galaxy acknowledged it.
command = json.loads((run / 'demo_engagement.json').read_text())
state.update(demo_engagement_mode=command['mode'], demo_signal_mode=command['signal_mode'])
(run / 'status.json').write_text(json.dumps(state))
return result
self.transport = local_transport
outcome = target(self.forwarder().submit('demo_signal', 'right').result(1))
self.assertEqual(outcome['session_id'], 'galaxy-independent-session')
self.assertTrue(outcome['applied'])
self.assertEqual({p.name for p in run.iterdir()}, {'status.json', 'demo_engagement.json'})
def test_unready_stale_live_stored_and_unknown_rejected_before_post(self):
patches = [{'available': False}, {'offroad': False}, {'state': 'PREPARING'}, {'state': 'DEGRADED'},
{'demo': {'available': False}}, {'demo': {}}, {'input_mode': 'live'},
{'mode': 'stored-score'}, {'live': {'enabled': True}}, {'judging_locked': True}]
for changes in patches:
with self.subTest(changes=changes):
self.transport = MockGalaxy()
self.transport.status.update(changes)
outcome = target(self.forwarder().submit('demo_signal', 'left').result(1))
self.assertEqual(outcome['status'], 'failed')
self.assertFalse(outcome['acknowledged'])
self.assertEqual([c[0] for c in self.transport.calls], ['GET'])
def test_missing_invalid_peer_session_cannot_use_local_session(self):
for session in (None, '', 123, '\nbad', 'x' * 257):
self.transport = MockGalaxy()
self.transport.status['demo']['session_id'] = session
outcome = target(self.forwarder().submit('demo_signal', 'left').result(1))
self.assertEqual(outcome['error'], 'invalid_peer_session')
self.assertEqual(len(self.transport.calls), 1)
def test_http_auth_error_is_sanitized_and_never_retried(self):
self.transport = lambda *args: (_ for _ in ()).throw(HTTPError(BASE, 401, COOKIE, {}, None))
result = self.forwarder().submit('demo_signal', 'left').result(1)
self.assertEqual(target(result)['error'], 'peer_http_401')
self.assertNotIn(COOKIE, json.dumps(result))
self.assertNotIn(BASE, json.dumps(result))
def test_application_not_claimed_from_post_receipt_alone(self):
self.transport.adopt = False
outcome = target(self.forwarder(timeout=.1).submit('demo_signal', 'left').result(1))
self.assertEqual(outcome['status'], 'acknowledged')
self.assertTrue(outcome['acknowledged'])
self.assertFalse(outcome['applied'])
self.assertEqual(outcome['error'], 'peer_application_unconfirmed')
self.assertEqual(sum(c[0] == 'POST' for c in self.transport.calls), 1)
def test_session_change_after_post_does_not_retry_into_new_session(self):
original = self.transport
def changed(method, *args):
value = original(method, *args)
if method == 'POST':
original.status['demo']['session_id'] = 'new-session'
return value
self.transport = changed
outcome = target(self.forwarder().submit('demo_signal', 'left').result(1))
self.assertEqual(outcome['error'], 'peer_session_changed')
self.assertTrue(outcome['acknowledged'])
self.assertFalse(outcome['applied'])
self.assertEqual(sum(c[0] == 'POST' for c in original.calls), 1)
def test_slow_status_has_bounded_future_no_queue_and_no_late_post(self):
entered, release, finished = threading.Event(), threading.Event(), threading.Event()
calls = []
def blocked(method, *args):
calls.append(method)
entered.set()
release.wait(2)
finished.set()
return copy.deepcopy(READY)
self.transport = blocked
controls = self.forwarder(timeout=.1)
started = time.monotonic()
future = controls.submit('demo_signal', 'left')
self.assertLess(time.monotonic() - started, .1)
self.assertTrue(entered.wait(1))
self.assertEqual(target(controls.submit('demo_signal', 'right').result(0))['status'], 'busy')
try:
outcome = target(future.result(1))
self.assertEqual(outcome['status'], 'timed_out')
self.assertFalse(outcome['delivery_unknown'])
self.assertEqual(target(controls.submit('demo_signal', 'off').result(0))['status'], 'busy')
finally:
release.set()
self.assertTrue(finished.wait(1))
self.assertEqual(calls, ['GET'])
def test_timeout_during_write_reports_unknown_without_rollback(self):
entered, release = threading.Event(), threading.Event()
original = self.transport
def blocked(method, *args):
value = original(method, *args)
if method == 'POST':
entered.set()
release.wait(2)
return value
self.transport = blocked
controls = self.forwarder(timeout=.1)
future = controls.submit('demo_signal', 'left')
self.assertTrue(entered.wait(1))
try:
outcome = target(future.result(1))
self.assertTrue(outcome['delivery_unknown'])
self.assertFalse(outcome['applied'])
self.assertFalse(outcome['acknowledged'])
self.assertEqual(original.status['demo']['signal_mode'], 'left')
self.assertEqual(sum(c[0] == 'POST' for c in original.calls), 1)
finally:
release.set()
def test_future_callback_can_close_without_deadlock(self):
entered, release = threading.Event(), threading.Event()
original = self.transport
def blocked(method, *args):
if method == 'POST':
entered.set()
release.wait(1)
return original(method, *args)
self.transport = blocked
controls = self.forwarder()
future = controls.submit('demo_signal', 'left')
self.assertTrue(entered.wait(1))
callback_done = threading.Event()
future.add_done_callback(lambda _: (controls.close(), callback_done.set()))
release.set()
future.result(1)
self.assertTrue(callback_done.wait(1))
self.assertEqual(target(controls.submit('demo_signal', 'off').result(0))['status'], 'disabled')
def test_redirects_are_rejected_instead_of_forwarding_cookie(self):
with self.assertRaises(ValueError):
_NoRedirect().redirect_request(None, None, 302, '', {}, 'https://other.example/')
def test_response_reader_is_bounded_and_does_not_honor_proxies(self):
class Response:
status = 200
class headers:
@staticmethod
def get_content_type(): return 'application/json'
def __enter__(self): return self
def __exit__(self, *args): pass
def read(self, count):
self.count = count
return b'x' * count
response = Response()
with patch('paired_demo_controls.build_opener') as build:
build.return_value.open.return_value = response
with self.assertRaises(ValueError):
_http_json('GET', BASE, None, {}, .2)
self.assertEqual(build.call_args.args[0].proxies, {})
self.assertEqual(response.count, 65537)
if __name__ == '__main__':
unittest.main()