mirror of
https://github.com/blakeblackshear/frigate.git
synced 2026-10-11 09:12:48 +03:00
CI / AMD64 Build (push) Canceled after 0s
CI / ARM Build (push) Canceled after 0s
CI / Jetson Jetpack 6 (push) Canceled after 0s
CI / AMD64 Extra Build (push) Canceled after 0s
CI / ARM Extra Build (push) Canceled after 0s
CI / Synaptics Build (push) Canceled after 0s
CI / Assemble and push default build (push) Canceled after 0s
* lock zmq sockets shared across threads Several zmq sockets were used from more than one thread at a time, which zmq does not allow. In the embeddings process, GenAI review and object description threads shared one requestor, so overlapping calls silently dropped the second write and could eventually leave a thread waiting forever on a reply that had already arrived, after which nothing was saved until restart. The API's embeddings requestor had the same problem, concurrent config publishes could mix the topic of one update with the payload of another, and the webpush config subscribers were drained from every publishing thread, which could abort the process. Each socket now has a lock. The embeddings requestor lock does not wait, so an overlapping call still returns empty right away instead of stalling the API behind a slow reply. * add test
63 lines
1.9 KiB
Python
63 lines
1.9 KiB
Python
"""Facilitates communication between processes."""
|
|
|
|
import multiprocessing as mp
|
|
import threading
|
|
from _pickle import UnpicklingError
|
|
from multiprocessing.synchronize import Event as MpEvent
|
|
from typing import Any
|
|
|
|
import zmq
|
|
|
|
SOCKET_PUB_SUB = "ipc:///tmp/cache/config"
|
|
|
|
|
|
class ConfigPublisher:
|
|
"""Publishes config changes to different processes."""
|
|
|
|
def __init__(self) -> None:
|
|
self.context = zmq.Context()
|
|
self.socket = self.context.socket(zmq.PUB)
|
|
self.socket.bind(SOCKET_PUB_SUB)
|
|
self.stop_event: MpEvent = mp.Event()
|
|
self.lock = threading.Lock()
|
|
|
|
def publish(self, topic: str, payload: Any) -> None:
|
|
"""There is no communication back to the processes."""
|
|
with self.lock:
|
|
self.socket.send_string(topic, flags=zmq.SNDMORE)
|
|
self.socket.send_pyobj(payload)
|
|
|
|
def stop(self) -> None:
|
|
self.stop_event.set()
|
|
self.socket.close(linger=0)
|
|
self.context.destroy(linger=0)
|
|
|
|
|
|
class ConfigSubscriber:
|
|
"""Simplifies receiving an updated config."""
|
|
|
|
def __init__(self, topic: str, exact: bool = False) -> None:
|
|
self.topic = topic
|
|
self.exact = exact
|
|
self.context = zmq.Context()
|
|
self.socket = self.context.socket(zmq.SUB)
|
|
self.socket.setsockopt_string(zmq.SUBSCRIBE, topic)
|
|
self.socket.connect(SOCKET_PUB_SUB)
|
|
|
|
def check_for_update(self) -> tuple[str, Any] | tuple[None, None]:
|
|
"""Returns updated config or None if no update."""
|
|
try:
|
|
topic = self.socket.recv_string(flags=zmq.NOBLOCK)
|
|
obj = self.socket.recv_pyobj()
|
|
|
|
if not self.exact or self.topic == topic:
|
|
return (topic, obj)
|
|
else:
|
|
return (None, None)
|
|
except (zmq.ZMQError, UnicodeDecodeError, UnpicklingError):
|
|
return (None, None)
|
|
|
|
def stop(self) -> None:
|
|
self.socket.close(linger=0)
|
|
self.context.destroy(linger=0)
|