stop creating a config subscriber per capture thread (#24002)

This commit is contained in:
Josh Hawkins
2026-09-12 07:30:04 -06:00
committed by Nicolas Mowen
parent d78f6a7a98
commit f7c5500ea8
+32 -43
View File
@@ -54,59 +54,48 @@ def capture_frames(
skipped_eps = EventsPerSecond() skipped_eps = EventsPerSecond()
skipped_eps.start() skipped_eps.start()
config_subscriber = CameraConfigUpdateSubscriber( while not stop_event.is_set():
None, {config.name: config}, [CameraConfigUpdateEnum.enabled] # CameraWatchdog applies enabled updates onto this same CameraConfig
) # before it stops ffmpeg. Do not subscribe here: it would be rebuilt per
# ffmpeg restart and strand a pipe in the idle main process config PUB.
if not config.enabled:
logger.debug(f"Stopping capture thread for disabled {config.name}")
break
def get_enabled_state(): fps.value = frame_rate.eps()
"""Fetch the latest enabled state from ZMQ.""" skipped_fps.value = skipped_eps.eps()
config_subscriber.check_for_updates() current_frame.value = datetime.now().timestamp()
return config.enabled frame_name = f"{config.name}_frame{frame_index}"
frame_buffer = frame_manager.write(frame_name)
try: try:
while not stop_event.is_set(): frame_buffer[:] = ffmpeg_process.stdout.read(frame_size)
if not get_enabled_state(): except Exception:
logger.debug(f"Stopping capture thread for disabled {config.name}") # shutdown has been initiated
if stop_event.is_set():
break break
fps.value = frame_rate.eps() logger.error(f"{config.name}: Unable to read frames from ffmpeg process.")
skipped_fps.value = skipped_eps.eps()
current_frame.value = datetime.now().timestamp()
frame_name = f"{config.name}_frame{frame_index}"
frame_buffer = frame_manager.write(frame_name)
try:
frame_buffer[:] = ffmpeg_process.stdout.read(frame_size)
except Exception:
# shutdown has been initiated
if stop_event.is_set():
break
if ffmpeg_process.poll() is not None:
logger.error( logger.error(
f"{config.name}: Unable to read frames from ffmpeg process." f"{config.name}: ffmpeg process is not running. exiting capture thread..."
) )
break
if ffmpeg_process.poll() is not None: continue
logger.error(
f"{config.name}: ffmpeg process is not running. exiting capture thread..."
)
break
continue frame_rate.update()
frame_rate.update() # don't lock the queue to check, just try since it should rarely be full
try:
# add to the queue
frame_queue.put((frame_name, current_frame.value), False)
frame_manager.close(frame_name)
except queue.Full:
# if the queue is full, skip this frame
skipped_eps.update()
# don't lock the queue to check, just try since it should rarely be full frame_index = 0 if frame_index == shm_frame_count - 1 else frame_index + 1
try:
# add to the queue
frame_queue.put((frame_name, current_frame.value), False)
frame_manager.close(frame_name)
except queue.Full:
# if the queue is full, skip this frame
skipped_eps.update()
frame_index = 0 if frame_index == shm_frame_count - 1 else frame_index + 1
finally:
config_subscriber.stop()
class CameraWatchdog(threading.Thread): class CameraWatchdog(threading.Thread):