From 03e01517ec45f1ced329f5cd3aed21b61f4c72c6 Mon Sep 17 00:00:00 2001 From: MoreTore Date: Mon, 17 Mar 2025 02:24:58 -0500 Subject: [PATCH] dont drop connections --- system/athena/athenad.py | 9 ++++++--- system/athena/streamer.py | 6 +++--- 2 files changed, 9 insertions(+), 6 deletions(-) diff --git a/system/athena/athenad.py b/system/athena/athenad.py index f3c3602a1..8cfac869f 100755 --- a/system/athena/athenad.py +++ b/system/athena/athenad.py @@ -149,7 +149,6 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: threading.Thread(target=upload_handler, args=(end_event,), name='upload_handler'), threading.Thread(target=log_handler, args=(end_event,), name='log_handler'), threading.Thread(target=stat_handler, args=(end_event,), name='stat_handler'), - threading.Thread(target=rtc_handler, args=(end_event, sdp_send_queue, sdp_recv_queue, ice_send_queue), name='rtc_handler'), ] + [ threading.Thread(target=jsonrpc_handler, args=(end_event,), name=f'worker_{x}') for x in range(HANDLER_THREADS) @@ -170,12 +169,12 @@ def handle_long_poll(ws: WebSocket, exit_event: threading.Event | None) -> None: thread.join() -def rtc_handler(end_event: threading.Event, sdp_send_queue: queue.Queue, sdp_recv_queue: queue.Queue, ice_recv_queue: queue.Queue) -> None: +def rtc_handler(exit_event: threading.Event, sdp_send_queue: queue.Queue, sdp_recv_queue: queue.Queue, ice_recv_queue: queue.Queue) -> None: loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) try: streamer = Streamer(sdp_send_queue, sdp_recv_queue, ice_recv_queue) - loop.run_until_complete(streamer.event_loop(end_event)) + loop.run_until_complete(streamer.event_loop(exit_event)) finally: loop.close() @@ -841,6 +840,10 @@ def main(exit_event: threading.Event = None): conn_start = None conn_retries = 0 + + #if Params().get_bool("EnableStreamer"): + threading.Thread(target=rtc_handler, args=(exit_event, sdp_send_queue, sdp_recv_queue, ice_send_queue), name='rtc_handler').start() + while exit_event is None or not exit_event.is_set(): try: if conn_start is None: diff --git a/system/athena/streamer.py b/system/athena/streamer.py index 8cf0567cc..e00a045e5 100644 --- a/system/athena/streamer.py +++ b/system/athena/streamer.py @@ -263,17 +263,17 @@ class Streamer: except Exception: logger.exception("Error during stop:", ) - async def event_loop(self, end_event: threading.Event): + async def event_loop(self, exit_event: threading.Event): """ Main event loop that processes signaling messages and maintains the PeerConnection. - Runs until end_event is set. + Runs until exit_event is set. """ self.camera = NativeProcess("camerad", "system/camerad", ["./camerad"], True) self.encoder = NativeProcess("encoderd", "system/loggerd", ["./encoderd", "--stream"], True) logger.info("Native processes for camera and encoder initialized.") stop_states = ['failed', 'closed'] connecting_states = ['connecting', 'new'] - while not end_event.is_set(): + while exit_event is None or not exit_event.is_set(): self.onroad = self.params.get_bool("IsOnroad") # support some functions while onroad try: try: