From a9e4f913ecf94c3e6fbfd0afbe0eea9a31d3c315 Mon Sep 17 00:00:00 2001 From: nicoboy Date: Wed, 1 Jul 2026 16:18:38 +0200 Subject: [PATCH] fix: reconnexion continue au flux GStreamer au lieu d'un timeout fixe MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit L'ancienne sonde abandonnait définitivement après 15s et restait bloquée en mode synthétique même quand le flux réel arrivait juste après (latence Mycelium variable, attente du premier keyframe H.264). Remplacé par un thread de fond (gstreamer_watcher) qui retente en continu et bascule la source active dès qu'une frame valide est lue, y compris après une perte de flux en cours de route. Démarrage MJPEG immédiat, plus d'attente bloquante. --- receive_cam.py | 111 +++++++++++++++++++++++++++++-------------------- 1 file changed, 65 insertions(+), 46 deletions(-) diff --git a/receive_cam.py b/receive_cam.py index 05330c9..7ec1527 100644 --- a/receive_cam.py +++ b/receive_cam.py @@ -49,6 +49,12 @@ _frame_lock = threading.Lock() _current_frame = None _running = True +# Source active partagée entre le watcher GStreamer et la boucle de capture : +# "synthetic" par défaut, bascule vers "gstreamer" dès qu'une frame réelle arrive. +_source_lock = threading.Lock() +_active_cap: cv2.VideoCapture | None = None +_active_source = "synthetic" + # Dernières stats Pi Zero reçues via UDP 5602 (thread-safe) _stats_lock = threading.Lock() _pi_stats: dict | None = None @@ -119,15 +125,39 @@ def draw_overlay(frame: np.ndarray) -> np.ndarray: return out -def _probe_gstreamer(result: list) -> None: - """Thread : tente une lecture GStreamer et stocke le cap si succès.""" - cap = cv2.VideoCapture(GSTREAMER_PIPELINE, cv2.CAP_GSTREAMER) - if cap.isOpened(): - ret, _ = cap.read() # bloque jusqu'à la première frame - if ret: - result.append(cap) - return - cap.release() +def gstreamer_watcher() -> None: + """ + Thread de fond permanent : tant que la source active n'est pas "gstreamer", + tente d'ouvrir le pipeline et de lire une première frame (le read() bloque + naturellement jusqu'à ce que des paquets valides arrivent, ce qui remplace + l'ancien timeout fixe). Bascule la source active dès qu'une frame arrive, + y compris après un démarrage en mode synthétique ou une perte de flux. + """ + global _active_cap, _active_source + + while _running: + with _source_lock: + already_connected = _active_source == "gstreamer" + if already_connected: + time.sleep(1) + continue + + cap = cv2.VideoCapture(GSTREAMER_PIPELINE, cv2.CAP_GSTREAMER) + if not cap.isOpened(): + cap.release() + time.sleep(2) + continue + + ret, _ = cap.read() # bloque jusqu'à la première frame valide + if not ret: + cap.release() + time.sleep(2) + continue + + print("[CAM] Flux GStreamer UDP reçu.") + with _source_lock: + _active_cap = cap + _active_source = "gstreamer" def _make_waiting_frame() -> np.ndarray: @@ -140,50 +170,34 @@ def _make_waiting_frame() -> np.ndarray: return frame -def open_capture() -> tuple[cv2.VideoCapture | None, str]: - """ - Priorité de source : - 1. Flux GStreamer UDP 5600 (Pi Zero) - 2. Frame synthétique (pas de hardware) - """ - print("[CAM] Sonde flux GStreamer UDP 5600 (timeout 15 s)...") - result: list = [] - t = threading.Thread(target=_probe_gstreamer, args=(result,), daemon=True) - t.start() - t.join(timeout=15) - - if result: - print("[CAM] Flux GStreamer UDP reçu.") - return result[0], "gstreamer" - - print("[CAM] Flux UDP indisponible — mode synthétique (attente Pi Zero)") - return None, "synthetic" - - def capture_loop(model: YOLO) -> None: - """Lit les frames, passe dans YOLOv8, stocke le résultat annoté.""" - global _current_frame, _running + """ + Lit les frames de la source active (gstreamer ou synthétique), passe dans + YOLOv8, stocke le résultat annoté. La source est mise à jour en tâche de + fond par gstreamer_watcher() : ce thread n'a pas besoin d'attendre le flux + réel au démarrage, il bascule dessus dès qu'il devient disponible. + """ + global _current_frame, _running, _active_cap, _active_source - cap, source = open_capture() - print(f"[CAM] Source active : {source}") - - # Frame synthétique réutilisée en boucle (pas d'inférence : frame vide) - synthetic_frame = _make_waiting_frame() if source == "synthetic" else None + print("[CAM] Démarrage en mode synthétique — connexion GStreamer en tâche de fond.") + synthetic_frame = _make_waiting_frame() while _running: + with _source_lock: + cap, source = _active_cap, _active_source + if source == "synthetic": frame = synthetic_frame time.sleep(1 / 10) # 10 fps synthétiques else: ret, frame = cap.read() if not ret: - if source == "gstreamer": - print("[CAM] Flux GStreamer perdu, nouvelle tentative...") - time.sleep(1) - continue - print("[CAM] Webcam perdue, arrêt.") - _running = False - break + print("[CAM] Flux GStreamer perdu, retour en mode synthétique.") + cap.release() + with _source_lock: + _active_cap = None + _active_source = "synthetic" + continue # Inférence YOLOv8 — skip en mode synthétique (frame vide sans objet) if model is not None and source != "synthetic": @@ -198,8 +212,9 @@ def capture_loop(model: YOLO) -> None: with _frame_lock: _current_frame = annotated - if cap is not None: - cap.release() + with _source_lock: + if _active_cap is not None: + _active_cap.release() print("[CAM] Capture terminée.") @@ -261,7 +276,11 @@ def main(): ts = threading.Thread(target=sideband_loop, daemon=True) ts.start() - # Thread de capture + inférence + # Thread de fond : connexion/reconnexion continue au flux GStreamer réel + tg = threading.Thread(target=gstreamer_watcher, daemon=True) + tg.start() + + # Thread de capture + inférence (démarre immédiatement en mode synthétique) t = threading.Thread(target=capture_loop, args=(model,), daemon=True) t.start()