From 735da70382ebcd32480eee1866fbd9de536612f4 Mon Sep 17 00:00:00 2001 From: royjr Date: Sat, 25 Jul 2026 04:57:48 -0400 Subject: [PATCH] close --- openpilot/tools/camerastream/compressed_vipc.py | 16 ++++++++++++---- openpilot/tools/wgpu/zmq.py | 9 +++++++++ 2 files changed, 21 insertions(+), 4 deletions(-) diff --git a/openpilot/tools/camerastream/compressed_vipc.py b/openpilot/tools/camerastream/compressed_vipc.py index 19f8849f20..bf376cb6ff 100755 --- a/openpilot/tools/camerastream/compressed_vipc.py +++ b/openpilot/tools/camerastream/compressed_vipc.py @@ -109,17 +109,25 @@ class CompressedVipc: while min(sm.recv_frame.values()) == 0: sm.update(100) + stream_dimensions = { + vst: (sm[ENCODE_SOCKETS[vst]].width, sm[ENCODE_SOCKETS[vst]].height) + for vst in vision_streams + } + # The metadata subscribers are setup-only. Leaving them connected creates a + # second unread camera subscription whose TCP queues grow for the entire run. + sm.close() + self.vipc_server = VisionIpcServer(server_name) for vst in vision_streams: - ed = sm[ENCODE_SOCKETS[vst]] - self.vipc_server.create_buffers(vst, 4, ed.width, ed.height) + width, height = stream_dimensions[vst] + self.vipc_server.create_buffers(vst, 4, width, height) self.vipc_server.start_listener() self.procs = [] process_context = multiprocessing.get_context("fork") for vst in vision_streams: - ed = sm[ENCODE_SOCKETS[vst]] - p = process_context.Process(target=decoder, args=(addr, self.vipc_server, vst, ed.width, ed.height, debug)) + width, height = stream_dimensions[vst] + p = process_context.Process(target=decoder, args=(addr, self.vipc_server, vst, width, height, debug)) p.start() self.procs.append(p) diff --git a/openpilot/tools/wgpu/zmq.py b/openpilot/tools/wgpu/zmq.py index dcf88fb33f..7d883a24a2 100644 --- a/openpilot/tools/wgpu/zmq.py +++ b/openpilot/tools/wgpu/zmq.py @@ -43,6 +43,10 @@ class ZmqSubSocket: messages.append(message) return messages + def close(self) -> None: + self.socket.close(linger=0) + self.context.term() + class ZmqSubMaster: def __init__(self, services: list[str], address: str): @@ -77,6 +81,11 @@ class ZmqSubMaster: self.updated[service] = True self.recv_frame[service] = self.frame + def close(self) -> None: + for sub in self.sockets.values(): + self.poller.unregister(sub.socket) + sub.close() + class ZmqPubMaster: def __init__(self, services: list[str]):