Files
frigate/frigate/comms/embeddings_updater.py
T
Josh HawkinsandGitHub cc86519603
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 (#24607)
* 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
2026-10-09 13:28:00 -06:00

101 lines
2.9 KiB
Python

"""Facilitates communication between processes."""
import logging
import threading
from collections.abc import Callable
from enum import Enum
from typing import Any
import zmq
logger = logging.getLogger(__name__)
SOCKET_REP_REQ = "ipc:///tmp/cache/embeddings"
class EmbeddingsRequestEnum(Enum):
# audio
transcribe_audio = "transcribe_audio"
# custom classification
reload_classification_model = "reload_classification_model"
# face
clear_face_classifier = "clear_face_classifier"
recognize_face = "recognize_face"
register_face = "register_face"
reprocess_face = "reprocess_face"
# semantic search
embed_description = "embed_description"
embed_thumbnail = "embed_thumbnail"
generate_search = "generate_search"
reindex = "reindex"
# LPR
reprocess_plate = "reprocess_plate"
# Review Descriptions
summarize_review = "summarize_review"
class EmbeddingsResponder:
def __init__(self) -> None:
self.context = zmq.Context()
self.socket = self.context.socket(zmq.REP)
self.socket.bind(SOCKET_REP_REQ)
def check_for_request(self, process: Callable) -> None:
while True: # load all messages that are queued
has_message, _, _ = zmq.select([self.socket], [], [], 0.01)
if not has_message:
break
try:
raw = self.socket.recv_json(flags=zmq.NOBLOCK)
if isinstance(raw, list):
(topic, value) = raw
response = process(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.socket.close()
self.context.destroy()
class EmbeddingsRequestor:
"""Simplifies sending data to EmbeddingsResponder 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)
self.lock = threading.Lock()
def send_data(self, topic: str, data: Any) -> Any:
"""Sends data and then waits for reply."""
# an overlapping call fails fast so a slow reply can't stall the API
if not self.lock.acquire(blocking=False):
return ""
try:
self.socket.send_json((topic, data))
return self.socket.recv_json()
except zmq.ZMQError:
return ""
finally:
self.lock.release()
def stop(self) -> None:
self.socket.close()
self.context.destroy()