From be168d14e9e402237a01178766af0be8c9b8d4e4 Mon Sep 17 00:00:00 2001 From: DevTekVE Date: Sat, 26 Apr 2025 19:33:35 +0200 Subject: [PATCH] SL: updating localproxy implementation (#841) * Adding capabilities to route localProxy via sunnylink * Undo * Thx lint * get api token * cert is not valid when it's an IP. Still use cert, but don't validate --- sunnypilot/sunnylink/athena/sunnylinkd.py | 17 ++++++++++++++-- sunnypilot/sunnylink/utils.py | 9 +++++++++ system/athena/athenad.py | 24 +++++++++++++---------- 3 files changed, 38 insertions(+), 12 deletions(-) diff --git a/sunnypilot/sunnylink/athena/sunnylinkd.py b/sunnypilot/sunnylink/athena/sunnylinkd.py index 2c4b95fa4..a67582062 100755 --- a/sunnypilot/sunnylink/athena/sunnylinkd.py +++ b/sunnypilot/sunnylink/athena/sunnylinkd.py @@ -10,11 +10,12 @@ import threading import time from jsonrpc import dispatcher +from functools import partial from openpilot.common.params import Params from openpilot.common.realtime import set_core_affinity from openpilot.common.swaglog import cloudlog from openpilot.system.athena.athenad import ws_send, jsonrpc_handler, \ - recv_queue, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler + recv_queue, UploadQueueCache, upload_queue, cur_upload_items, backoff, ws_manage, log_handler, start_local_proxy_shim from websocket import (ABNF, WebSocket, WebSocketException, WebSocketTimeoutException, create_connection) @@ -50,7 +51,7 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: # threading.Thread(target=sunny_log_handler, args=(end_event, comma_prime_cellular_end_event), name='log_handler'), # threading.Thread(target=stat_handler, args=(end_event,), name='stat_handler'), ] + [ - threading.Thread(target=jsonrpc_handler, args=(end_event,), name=f'worker_{x}') + threading.Thread(target=jsonrpc_handler, args=(end_event, partial(startLocalProxy, end_event),), name=f'worker_{x}') for x in range(HANDLER_THREADS) ] @@ -201,6 +202,18 @@ def saveParams(params_to_update: dict[str, str], compression: bool = False) -> N cloudlog.error(f"sunnylinkd.saveParams.exception {e}") +def startLocalProxy(global_end_event: threading.Event, remote_ws_uri: str, local_port: int) -> dict[str, int]: + cloudlog.debug("athena.startLocalProxy.starting") + ws = create_connection( + remote_ws_uri, + header={"Authorization": f"Bearer {sunnylink_api.get_token()}"}, + enable_multithread=True, + sslopt={"cert_reqs": ssl.CERT_NONE} + ) + + return start_local_proxy_shim(global_end_event, local_port, ws) + + def main(exit_event: threading.Event = None): try: set_core_affinity([0, 1, 2, 3]) diff --git a/sunnypilot/sunnylink/utils.py b/sunnypilot/sunnylink/utils.py index 57523d739..55b1f6f28 100644 --- a/sunnypilot/sunnylink/utils.py +++ b/sunnypilot/sunnylink/utils.py @@ -46,3 +46,12 @@ def register_sunnylink(): sunnylink_id = SunnylinkApi(None).register_device(None, **extra_args) print(f"SunnyLinkId: {sunnylink_id}") + + +def get_api_token(): + """Get the API token for the device.""" + params = Params() + sunnylink_dongle_id = params.get("SunnylinkDongleId", encoding='utf-8') + sunnylink_api = SunnylinkApi(sunnylink_dongle_id) + token = sunnylink_api.get_token() + print(f"API Token: {token}") diff --git a/system/athena/athenad.py b/system/athena/athenad.py index 4bcdc9b83..3f3fc3d70 100755 --- a/system/athena/athenad.py +++ b/system/athena/athenad.py @@ -184,8 +184,8 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: thread.join() -def jsonrpc_handler(end_event: threading.Event) -> None: - dispatcher["startLocalProxy"] = partial(startLocalProxy, end_event) +def jsonrpc_handler(end_event: threading.Event, localProxyHandler = None) -> None: + dispatcher["startLocalProxy"] = localProxyHandler or partial(startLocalProxy, end_event) while not end_event.is_set(): try: data = recv_queue.get(timeout=1) @@ -480,7 +480,19 @@ def setRouteViewed(route: str) -> dict[str, int | str]: def startLocalProxy(global_end_event: threading.Event, remote_ws_uri: str, local_port: int) -> dict[str, int]: + cloudlog.debug("athena.startLocalProxy.starting") + dongle_id = Params().get("DongleId").decode('utf8') + identity_token = Api(dongle_id).get_token() + ws = create_connection(remote_ws_uri, cookie="jwt=" + identity_token, enable_multithread=True) + + return start_local_proxy_shim(global_end_event, local_port, ws) + + +def start_local_proxy_shim(global_end_event: threading.Event, local_port: int, ws: WebSocket) -> dict[str, int]: try: + if ws.sock is None: + raise Exception("WebSocket is not connected") + # migration, can be removed once 0.9.8 is out for a while if local_port == 8022: local_port = 22 @@ -488,14 +500,6 @@ def startLocalProxy(global_end_event: threading.Event, remote_ws_uri: str, local if local_port not in LOCAL_PORT_WHITELIST: raise Exception("Requested local port not whitelisted") - cloudlog.debug("athena.startLocalProxy.starting") - - dongle_id = Params().get("DongleId").decode('utf8') - identity_token = Api(dongle_id).get_token() - ws = create_connection(remote_ws_uri, - cookie="jwt=" + identity_token, - enable_multithread=True) - # Set TOS to keep connection responsive while under load. # DSCP of 36/HDD_LINUX_AC_VI with the minimum delay flag ws.sock.setsockopt(socket.IPPROTO_IP, socket.IP_TOS, 0x90)