Files
frigate/frigate/comms/inter_process.py
T

94 lines
3.0 KiB
Python
Raw Normal View History

"""Facilitates communication between processes."""
2025-08-08 06:08:37 -06:00
import logging
import multiprocessing as mp
import threading
2026-07-06 09:28:02 -08:00
from collections.abc import Callable
from multiprocessing.synchronize import Event as MpEvent
2026-07-06 09:28:02 -08:00
from typing import Any
import zmq
2025-02-10 20:47:15 -06:00
from frigate.comms.base_communicator import Communicator
2025-08-08 06:08:37 -06:00
logger = logging.getLogger(__name__)
SOCKET_REP_REQ = "ipc:///tmp/cache/comms"
class InterProcessCommunicator(Communicator):
def __init__(self) -> None:
# bound eagerly so subprocesses starting before start_communicators()
# can still connect; their requests queue in zmq until the reader runs
self.context = zmq.Context()
self.socket = self.context.socket(zmq.REP)
self.socket.bind(SOCKET_REP_REQ)
self.stop_event: MpEvent = mp.Event()
self.reader_thread: threading.Thread | None = None
2025-08-08 06:08:37 -06:00
def publish(self, topic: str, payload: Any, retain: bool = False) -> None:
"""There is no communication back to the processes."""
pass
def subscribe(self, receiver: Callable) -> None:
self._dispatcher = receiver
def start(self) -> None:
self.reader_thread = threading.Thread(target=self.read)
self.reader_thread.start()
def read(self) -> None:
while not self.stop_event.is_set():
while True: # load all messages that are queued
has_message, _, _ = zmq.select([self.socket], [], [], 1)
if not has_message:
break
try:
2025-08-08 06:08:37 -06:00
raw = self.socket.recv_json(flags=zmq.NOBLOCK)
2025-08-08 06:08:37 -06:00
if isinstance(raw, list):
(topic, value) = raw
response = self._dispatcher(topic, value)
else:
logging.warning(
f"Received unexpected data type in ZMQ recv_json: {type(raw)}"
)
response = None
if response is not None:
self.socket.send_json(response)
else:
self.socket.send_json([])
except zmq.ZMQError:
break
def stop(self) -> None:
self.stop_event.set()
if self.reader_thread is not None:
self.reader_thread.join()
2026-03-04 10:07:34 -06:00
self.socket.close(linger=0)
self.context.destroy(linger=0)
class InterProcessRequestor:
"""Simplifies sending data to InterProcessCommunicator and getting a reply."""
def __init__(self) -> None:
self.context = zmq.Context()
self.socket = self.context.socket(zmq.REQ)
self.socket.connect(SOCKET_REP_REQ)
2025-05-13 16:27:20 +02:00
def send_data(self, topic: str, data: Any) -> Any:
"""Sends data and then waits for reply."""
2024-10-11 10:47:23 -06:00
try:
self.socket.send_json((topic, data))
return self.socket.recv_json()
except zmq.ZMQError:
return ""
def stop(self) -> None:
2026-03-04 10:07:34 -06:00
self.socket.close(linger=0)
self.context.destroy(linger=0)