revert: drop the ZMQ video change from this branch
Committed by mistake in daabf83 - the video stream work was explicitly
put on hold and does not belong in this branch.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01JHoUkj66jxByS2ypY5h9Mn
This commit is contained in:
@@ -1,6 +1,5 @@
|
||||
import json
|
||||
import time
|
||||
from threading import Lock
|
||||
|
||||
import cv2
|
||||
import numpy as np
|
||||
@@ -15,17 +14,6 @@ from aare.gui.constants import LOGGER_NAME
|
||||
|
||||
logger = setup_logger(LOGGER_NAME)
|
||||
|
||||
# Frames handed to the GUI thread but not painted yet. The image signal is a
|
||||
# queued connection, so without a cap the subscriber hands over frames faster
|
||||
# than they can be painted and the Qt event queue grows without bound - lag
|
||||
# that no amount of draining on the socket can fix.
|
||||
MAX_FRAMES_IN_FLIGHT = 1
|
||||
|
||||
# The socket's own queue. Only a memory bound: _recv_latest() drops whatever
|
||||
# piled up behind the newest frame anyway. Has to stay above the number of
|
||||
# parts in a frame, or the pipe cannot assemble a multipart message.
|
||||
RCV_QUEUE_MESSAGES = 10
|
||||
|
||||
|
||||
class PredictionSubscriber(QThread):
|
||||
prediction = Signal(dict)
|
||||
@@ -42,14 +30,6 @@ class PredictionSubscriber(QThread):
|
||||
self._sock = self._ctx.socket(zmq.SUB)
|
||||
self._sock.setsockopt(zmq.RCVTIMEO, 500)
|
||||
self._sock.setsockopt(zmq.LINGER, 0)
|
||||
self._sock.setsockopt(zmq.RCVHWM, RCV_QUEUE_MESSAGES)
|
||||
|
||||
# Delivery of an image frees the slot for the next one. Queued back to
|
||||
# the thread this object lives in (the GUI thread), like the painting
|
||||
# slots, so it runs once the GUI has worked through the frame.
|
||||
self._frames_in_flight = 0
|
||||
self._frames_in_flight_lock = Lock()
|
||||
self.image.connect(self._on_image_delivered)
|
||||
|
||||
self._emit_images = True
|
||||
self.running = True
|
||||
@@ -82,37 +62,6 @@ class PredictionSubscriber(QThread):
|
||||
def set_emit_images(self, enabled: bool) -> None:
|
||||
self._emit_images = enabled
|
||||
|
||||
def _recv_latest(self) -> tuple[list[bytes], int]:
|
||||
"""The newest frame on the socket, plus how many older ones were
|
||||
dropped to get to it. This is what CONFLATE would do, except that
|
||||
CONFLATE keeps a single message *part* and every frame here is
|
||||
multipart (header, image, detections), so it cannot be used."""
|
||||
sock = self._sock
|
||||
if sock is None: # run() has already torn the socket down
|
||||
raise zmq.Again
|
||||
parts = sock.recv_multipart()
|
||||
dropped = 0
|
||||
while True:
|
||||
try:
|
||||
parts = sock.recv_multipart(zmq.NOBLOCK)
|
||||
except zmq.Again:
|
||||
return parts, dropped
|
||||
dropped += 1
|
||||
|
||||
@Slot(QPixmap)
|
||||
def _on_image_delivered(self, _pixmap: QPixmap) -> None:
|
||||
with self._frames_in_flight_lock:
|
||||
self._frames_in_flight = max(0, self._frames_in_flight - 1)
|
||||
|
||||
def _gui_ready_for_frame(self) -> bool:
|
||||
with self._frames_in_flight_lock:
|
||||
return self._frames_in_flight < MAX_FRAMES_IN_FLIGHT
|
||||
|
||||
def _emit_image(self, pixmap: QPixmap) -> None:
|
||||
with self._frames_in_flight_lock:
|
||||
self._frames_in_flight += 1
|
||||
self.image.emit(pixmap)
|
||||
|
||||
def _set_camera_available(self, available: bool, error: str | None = None) -> None:
|
||||
if available != self._camera_available:
|
||||
self._camera_available = available
|
||||
@@ -203,10 +152,12 @@ class PredictionSubscriber(QThread):
|
||||
self.focus_measure.emit(sharpness)
|
||||
|
||||
def run(self):
|
||||
self._debug_last_log_ts = time.perf_counter()
|
||||
self._debug_msg_count = 0
|
||||
try:
|
||||
while self.running:
|
||||
try:
|
||||
parts, dropped = self._recv_latest()
|
||||
parts = self._sock.recv_multipart()
|
||||
except zmq.Again:
|
||||
now = time.perf_counter()
|
||||
elapsed = now - self._fps_window_start
|
||||
@@ -228,8 +179,7 @@ class PredictionSubscriber(QThread):
|
||||
|
||||
now = time.perf_counter()
|
||||
self._last_frame_time = now
|
||||
# dropped frames count too: this is the camera's rate, not ours
|
||||
self._fps_frame_count += 1 + dropped
|
||||
self._fps_frame_count += 1
|
||||
|
||||
elapsed = now - self._fps_window_start
|
||||
if elapsed >= self._fps_emit_period_s:
|
||||
@@ -261,12 +211,7 @@ class PredictionSubscriber(QThread):
|
||||
target = next((d for d in json_dicts if "target_point" in d), None)
|
||||
image_bytes = max(non_json_parts, key=len) if non_json_parts else None
|
||||
|
||||
# The decode is skipped along with the emit: a frame the GUI
|
||||
# is too busy to paint is not worth decoding. The sharpness
|
||||
# read-out is the exception, it works on the decoded image.
|
||||
emit_image = self._emit_images and self._gui_ready_for_frame()
|
||||
|
||||
if header and image_bytes and (emit_image or self._measure_focus):
|
||||
if self._emit_images and header and image_bytes:
|
||||
rgb = self._decode_rgb_image(header, image_bytes)
|
||||
if rgb is None:
|
||||
self._set_camera_available(
|
||||
@@ -275,9 +220,9 @@ class PredictionSubscriber(QThread):
|
||||
else:
|
||||
self._set_camera_available(True)
|
||||
self._emit_focus_measure_if_enabled(rgb)
|
||||
if emit_image and self.running:
|
||||
self._emit_image(self._rgb_to_pixmap(rgb))
|
||||
elif self._emit_images and not (header and image_bytes):
|
||||
if self.running:
|
||||
self.image.emit(self._rgb_to_pixmap(rgb))
|
||||
elif self._emit_images:
|
||||
self._set_camera_available(
|
||||
False, "Sample camera feed unavailable: no frame header in zmq stream"
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user