Ver código fonte

Count camera reconnects in a row, not for the life of the stream

maziggy 3 dias atrás
pai
commit
2315b8fb5d

+ 1 - 0
CHANGELOG.md

@@ -46,6 +46,7 @@ All notable changes to Bambuddy will be documented in this file.
 - **The frontend build no longer warns about `path` and `crypto` being externalized for the STEP previewer (#2976)** — `occt-import-js`, the Emscripten build behind STEP previews, requires both modules, but only inside its `ENVIRONMENT_IS_NODE` branches; in the browser it loads its `.wasm` from the URL the preview worker passes and draws randomness from `crypto.getRandomValues`. Vite still externalized both and printed two warnings on every build. `vite.config.ts` now drops exactly those two warnings for that one package through `build.rolldownOptions.onLog`, so an externalization anywhere else, or of any other module, still shows.
 
 ### Fixed
+- **A live camera view no longer stops for good after about half an hour on X1, H2 and P2 printers** — These printers' camera streams come over RTSP, and the printer ends each session after a while; a stock X1 Carbon ends every one after exactly a minute. Bambuddy reconnects straight away, so viewers never notice, but every reconnect counted against a limit of 30 for the life of the stream. Half an hour into a print the camera stopped and did not come back until the page was reopened. The limit now counts failed attempts in a row, and a session that delivered video resets it. A printer that refuses the camera for a while, for example because another app is on it, used to be given up on after 30 attempts in nine seconds. Bambuddy now waits longer between failed attempts, up to five seconds each, so it gets about two minutes. External cameras had the same problem with a limit of three: a camera server that ends its sessions now and then stopped the stream on the fourth drop. External streams now always reconnect after a session that delivered video. Found while investigating #3189.
 - **A slot that reads empty for a moment no longer loses its spool assignment (#3186, reported by @Sawtaytoes)** — An idle X1 Carbon sent one status update that showed a whole AMS unit as empty, with no colour or material in any slot, and Bambuddy deleted all four spool assignments on it at once. The spools never moved, and the next update reported them again. Nothing brought the assignments back, and for a non-RFID spool the assignment is the only record of which spool is in the slot. The #3100 fix covered a blank slot the AMS still reported as occupied, but not one briefly reported as empty. A slot that looks empty, or that drops out of the AMS data, now keeps its assignment for two minutes and loses it only if it is still empty then. Bambuddy checks again by itself when the two minutes are up, so a spool that was really taken out does not wait for the next AMS change. A different spool the AMS can identify, by its RFID tag or its colour and material, still releases the old assignment immediately. A spool it cannot read that goes in during those two minutes releases it once they are up, rather than inheriting the old spool's assignment. Assigning a spool to the slot yourself cancels the wait. Spoolman mode's slot links follow the same rules.
 - **Number fields can be cleared and retyped (#3182, reported by @Carter3DP)** — Every number field corrected its value on each keystroke, so erasing the "1" in the print dialog's Quantity snapped straight back to 1, and getting to 6 meant typing 16 and deleting the 1. It was worse where the minimum is above 1: typing the "6" of 60 into the AMS drying temperature turned it into 45, so the field could only be set with the arrows. Fields now keep what you type while you edit, use it once it is a number in range, and settle on leaving the field: out-of-range numbers are pulled into range, and an empty field goes back to its default. This covers the 38 number fields across the print dialog, per-plate quantities, scheduling, drying, Settings, smart plugs, backups, spools, projects and SpoolBuddy. Two fields could not be set to 0 even though 0 is allowed, the smart plug's delay after drying and AMS humidity; they can now. The drying presets now also keep to their own limits, 30–65 °C on an AMS 2 Pro and 30–85 °C on an AMS HT.
 - **Installing or updating no longer runs out of memory on 2 GB machines (#3181, reported by @PhilippeP62)** — Every installer and updater builds the frontend with `npm run build`, which also ran the TypeScript type check. That check alone needs about 1 GB of Node memory, the whole default allowance on a 2 GB machine such as the standard Proxmox LXC or a 2 GB Raspberry Pi, so the build crashed with "JavaScript heap out of memory". On the Proxmox helper script that left the install without its database and nobody could sign in; Bambuddy's own update script rolled back but could never finish an update. The build now only bundles the frontend, and the type check runs on its own in development and CI (`npm run typecheck`). That also fixes CI's type-check job and the pre-commit hook, which ran `tsc --noEmit` against a config that lists no files and so had been passing without checking anything.

+ 18 - 2
backend/app/api/routes/camera.py

@@ -650,7 +650,14 @@ async def generate_rtsp_mjpeg_stream(
                     ip_address,
                     stream_id,
                 )
-                await asyncio.sleep(profile.rtsp_reconnect_delay)
+                # Fast after a session that delivered frames (reconnect_count
+                # is 1 then), backing off while the printer keeps refusing.
+                await asyncio.sleep(
+                    min(
+                        profile.rtsp_reconnect_delay * 2 ** (reconnect_count - 1),
+                        profile.rtsp_reconnect_backoff_max,
+                    )
+                )
                 if disconnect_event and disconnect_event.is_set():
                     break
 
@@ -698,6 +705,7 @@ async def generate_rtsp_mjpeg_stream(
             buffer = b""
             stream_ended = False
             client_gone = False
+            session_got_frames = False
 
             while True:
                 if disconnect_event and disconnect_event.is_set():
@@ -735,6 +743,7 @@ async def generate_rtsp_mjpeg_stream(
                         frame = buffer[: end_idx + 2]
                         buffer = buffer[end_idx + 2 :]
                         got_any_frames = True
+                        session_got_frames = True
 
                         if printer_id is not None:
                             _last_frames[printer_id] = frame
@@ -784,6 +793,13 @@ async def generate_rtsp_mjpeg_stream(
                 break
 
             if stream_ended:
+                # The budget is for failures in a row. A session that
+                # delivered video was a success, however it ended -- a stock
+                # X1C ends every session after about a minute, which under a
+                # lifetime count stopped the live view for good after half an
+                # hour.
+                if session_got_frames:
+                    reconnect_count = 0
                 reconnect_count += 1
                 continue
 
@@ -792,7 +808,7 @@ async def generate_rtsp_mjpeg_stream(
 
         if reconnect_count > profile.rtsp_reconnect_max:
             logger.error(
-                "RTSP max reconnects (%d) reached for %s (stream_id=%s)",
+                "RTSP max consecutive reconnects (%d) reached for %s (stream_id=%s)",
                 profile.rtsp_reconnect_max,
                 ip_address,
                 stream_id,

+ 11 - 2
backend/app/services/camera_profiles.py

@@ -53,10 +53,19 @@ class CameraProfile:
     # Max consecutive ffmpeg respawns when the printer drops the RTSP
     # session mid-stream. Some firmwares cut the stream after a few
     # seconds (originally noted on P2S), so we transparently respawn
-    # to keep the MJPEG client alive.
+    # to keep the MJPEG client alive. Consecutive means failed sessions
+    # in a row: a session that delivered frames resets the count. A stock
+    # X1C ends every session after about a minute, so a lifetime count
+    # stopped a live view for good after half an hour.
     rtsp_reconnect_max: int = 30
-    # Seconds between ffmpeg respawn attempts.
+    # Seconds before respawning after a session that delivered frames --
+    # the routine drop, reconnected fast enough that viewers don't notice.
     rtsp_reconnect_delay: float = 0.2
+    # Each further failed attempt in a row doubles the delay up to this cap,
+    # so a printer that is refusing sessions for a while (another client on
+    # its camera, a busy camera service) gets minutes rather than 30 dials in
+    # nine seconds before the stream is given up.
+    rtsp_reconnect_backoff_max: float = 5.0
 
     # --- Extra ffmpeg input args ---------------------------------------------
     # Hook for future per-model knobs (e.g. `-fflags` overrides) without

+ 42 - 28
backend/app/services/external_camera.py

@@ -8,6 +8,7 @@ to ensure they are well-formed before use.
 """
 
 import asyncio
+import contextlib
 import functools
 import ipaddress
 import logging
@@ -921,42 +922,55 @@ async def generate_mjpeg_stream(
                 logger.exception("on_frame callback raised")
         return _format_mjpeg_frame(frame)
 
-    if camera_type == "mjpeg":
-        # Proxy MJPEG stream directly, with reconnect on timeout
-        max_retries = 3
-        for attempt in range(max_retries + 1):
+    async def _reconnecting(open_session, label: str):
+        """Yield frames across sessions, reconnecting after each that delivered.
+
+        A session that ends without a single frame stops the stream: the
+        source is down, and retrying is the viewer's call. One that delivered
+        frames and then ended is a routine drop and is always reconnected.
+        This used to be three reconnects for the life of the stream, so a
+        server that closes its sessions periodically ended the stream for good
+        on the fourth drop, however long each session had run -- the external
+        twin of the built-in RTSP path's lifetime reconnect budget.
+        """
+        drops = 0
+        while True:
             frame_yielded = False
-            async for frame in _stream_mjpeg(url):
-                frame_yielded = True
+            # aclosing: stop the session's ffmpeg the moment this generator is
+            # closed, rather than whenever the abandoned iterator is collected.
+            async with contextlib.aclosing(open_session()) as session:
+                async for frame in session:
+                    frame_yielded = True
+                    yield frame
+            if not frame_yielded or (stop_event is not None and stop_event.is_set()):
+                break
+            # The sessions swallow CancelledError and simply end, so a viewer
+            # whose task was cancelled mid-read looks like a routine drop here.
+            # Never redial for a task that is being cancelled (Task.cancelling
+            # is 3.11+; without it this falls back to the next await raising).
+            task = asyncio.current_task()
+            if task is not None and getattr(task, "cancelling", lambda: 0)():
+                break
+            drops += 1
+            logger.warning("External %s stream ended, reconnecting (drop %d)...", label, drops)
+            await asyncio.sleep(2)
+
+    if camera_type == "mjpeg":
+        # Proxy MJPEG stream directly, reconnecting after routine drops.
+        async with contextlib.aclosing(_reconnecting(lambda: _stream_mjpeg(url), "MJPEG")) as frames:
+            async for frame in frames:
                 current_time = asyncio.get_event_loop().time()
                 if current_time - last_frame_time >= frame_interval:
                     last_frame_time = current_time
                     yield _publish(frame)
-            if not frame_yielded or attempt == max_retries or (stop_event is not None and stop_event.is_set()):
-                break
-            logger.warning(
-                "External MJPEG stream ended, reconnecting (attempt %d/%d)...",
-                attempt + 1,
-                max_retries,
-            )
-            await asyncio.sleep(2)
 
     elif camera_type == "rtsp":
-        # Use ffmpeg to convert RTSP to MJPEG, with reconnect on timeout
-        max_retries = 3
-        for attempt in range(max_retries + 1):
-            frame_yielded = False
-            async for frame in _stream_rtsp(url, fps, on_process=on_process):
-                frame_yielded = True
+        # Use ffmpeg to convert RTSP to MJPEG, reconnecting after routine drops.
+        async with contextlib.aclosing(
+            _reconnecting(lambda: _stream_rtsp(url, fps, on_process=on_process), "RTSP")
+        ) as frames:
+            async for frame in frames:
                 yield _publish(frame)
-            if not frame_yielded or attempt == max_retries or (stop_event is not None and stop_event.is_set()):
-                break
-            logger.warning(
-                "External RTSP stream ended, reconnecting (attempt %d/%d)...",
-                attempt + 1,
-                max_retries,
-            )
-            await asyncio.sleep(2)
 
     elif camera_type == "usb":
         # Use ffmpeg to stream from USB camera

+ 266 - 0
backend/tests/unit/test_camera_reconnect_budget.py

@@ -0,0 +1,266 @@
+"""The camera reconnect budget counts failures in a row, not every drop.
+
+A stock X1C ends each RTSP session after about a minute. The built-in stream
+respawned ffmpeg transparently but counted every respawn against a lifetime
+budget of 30, so a live view stopped for good after about half an hour. The
+external-camera path allowed three reconnects for the life of the stream.
+"""
+
+import asyncio
+from contextlib import suppress
+
+import pytest
+
+from backend.app.api.routes import camera
+from backend.app.services import external_camera
+from backend.app.services.camera_profiles import CameraProfile
+
+FRAME = b"\xff\xd8frame\xff\xd9"
+
+
+class _FakeServer:
+    def close(self) -> None:
+        pass
+
+    async def wait_closed(self) -> None:
+        pass
+
+
+class _Stdout:
+    def __init__(self, frames: int) -> None:
+        self._frames = frames
+
+    async def read(self, _size: int = -1) -> bytes:
+        if self._frames <= 0:
+            return b""
+        self._frames -= 1
+        return FRAME
+
+
+class _Proc:
+    _next_pid = 78000
+
+    def __init__(self, frames: int) -> None:
+        _Proc._next_pid += 1
+        self.pid = _Proc._next_pid
+        self.returncode = None
+        self.stdout = _Stdout(frames)
+        self.stderr = None
+
+    def terminate(self) -> None:
+        self.returncode = 0
+
+    def kill(self) -> None:
+        self.returncode = -9
+
+    async def wait(self) -> int:
+        if self.returncode is None:
+            self.returncode = 0
+        return self.returncode
+
+
+@pytest.fixture
+def rtsp(monkeypatch):
+    """Fake ffmpeg sessions: each call spawns the next entry's frame count."""
+    sessions: list[int] = []
+    spawned: list[int] = []
+    delays: list[float] = []
+    real_sleep = asyncio.sleep
+
+    async def _fake_exec(*_args, **_kwargs):
+        frames = sessions.pop(0) if sessions else 0
+        spawned.append(frames)
+        return _Proc(frames)
+
+    async def _fake_proxy(_ip: str, _port: int):
+        return 48998, _FakeServer()
+
+    async def _recording_sleep(seconds, *args, **kwargs):
+        delays.append(seconds)
+        await real_sleep(0)
+
+    monkeypatch.setattr(camera, "get_ffmpeg_path", lambda: "/fake/ffmpeg")
+    monkeypatch.setattr(camera, "create_tls_proxy", _fake_proxy)
+    monkeypatch.setattr(camera.asyncio, "create_subprocess_exec", _fake_exec)
+    monkeypatch.setattr(camera.asyncio, "sleep", _recording_sleep)
+    return sessions, spawned, delays
+
+
+def _use_profile(monkeypatch, **kwargs):
+    monkeypatch.setattr(camera, "get_camera_profile", lambda _model: CameraProfile(**kwargs))
+
+
+async def _drain(stream, limit: int = 200) -> int:
+    frames = 0
+    async for chunk in stream:
+        if b"image/jpeg" in chunk:
+            frames += 1
+        if frames >= limit:
+            break
+    return frames
+
+
+def _stream():
+    return camera.generate_rtsp_mjpeg_stream(
+        ip_address="192.0.2.40",
+        access_code="test-code",
+        model="X1C",
+        fps=10,
+        stream_id="99-fanout-budget",
+        disconnect_event=asyncio.Event(),
+    )
+
+
+async def test_routine_drops_never_use_up_the_budget(rtsp, monkeypatch):
+    """Ten sessions that each deliver video and end, with a budget of two."""
+    sessions, spawned, _ = rtsp
+    _use_profile(monkeypatch, rtsp_reconnect_max=2, rtsp_reconnect_delay=0.2)
+    sessions.extend([3] * 10)
+
+    stream = _stream()
+    frames = await asyncio.wait_for(_drain(stream, limit=30), timeout=10)
+    with suppress(Exception):
+        await stream.aclose()
+
+    assert frames == 30
+    assert len(spawned) == 10
+
+
+async def test_failures_in_a_row_still_give_up(rtsp, monkeypatch):
+    sessions, spawned, _ = rtsp
+    _use_profile(monkeypatch, rtsp_reconnect_max=3, rtsp_reconnect_delay=0.2)
+    sessions.extend([2])  # then every session fails without a frame
+
+    frames = await asyncio.wait_for(_drain(_stream()), timeout=10)
+
+    assert frames == 2
+    # The good session, then three reconnects that each fail -- the budget of
+    # three consecutive reconnects -- and the stream gives up.
+    assert spawned == [2, 0, 0, 0]
+
+
+async def test_failures_back_off_and_a_good_session_resets_the_delay(rtsp, monkeypatch):
+    sessions, _, delays = rtsp
+    _use_profile(monkeypatch, rtsp_reconnect_max=6, rtsp_reconnect_delay=0.2, rtsp_reconnect_backoff_max=1.0)
+    sessions.extend([1, 0, 0, 0, 1])
+
+    await asyncio.wait_for(_drain(_stream()), timeout=10)
+
+    # 0.1 is the post-spawn startup check, not a reconnect delay.
+    reconnect_delays = [d for d in delays if d != 0.1]
+    assert reconnect_delays == [0.2, 0.4, 0.8, 1.0, 0.2, 0.4, 0.8, 1.0, 1.0, 1.0]
+
+
+# ---------------------------------------------------------------------------
+# External cameras
+# ---------------------------------------------------------------------------
+
+
+@pytest.fixture
+def external(monkeypatch):
+    sessions: list[int] = []
+    opened: list[int] = []
+    closed: list[int] = []
+
+    def _fake_stream_rtsp(_url, _fps, on_process=None):
+        frames = sessions.pop(0) if sessions else 0
+        index = len(opened)
+        opened.append(frames)
+
+        async def _gen():
+            try:
+                for _ in range(frames):
+                    yield FRAME
+            finally:
+                closed.append(index)
+
+        return _gen()
+
+    async def _no_sleep(_seconds, *args, **kwargs):
+        return None
+
+    monkeypatch.setattr(external_camera, "_stream_rtsp", _fake_stream_rtsp)
+    monkeypatch.setattr(external_camera.asyncio, "sleep", _no_sleep)
+    return sessions, opened, closed
+
+
+async def test_external_routine_drops_keep_the_stream_going(external):
+    """Used to end for good on the fourth drop."""
+    sessions, opened, _ = external
+    sessions.extend([2, 2, 2, 2, 2, 2])
+
+    frames = [f async for f in external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10)]
+
+    assert len(frames) == 12
+    assert opened == [2, 2, 2, 2, 2, 2, 0], "a session with no frame ends the stream"
+
+
+async def test_external_stops_when_asked(external):
+    sessions, opened, _ = external
+    sessions.extend([1] * 10)
+    stop = asyncio.Event()
+
+    frames = []
+    async for frame in external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10, stop_event=stop):
+        frames.append(frame)
+        if len(frames) == 3:
+            stop.set()
+
+    assert len(frames) == 3
+    assert len(opened) == 3
+
+
+async def test_closing_the_external_stream_closes_the_open_session(external):
+    """The session owns the ffmpeg process; it must stop when the viewer goes,
+    not whenever the abandoned iterator happens to be collected."""
+    sessions, _, closed = external
+    sessions.extend([50])
+
+    stream = external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10)
+    await anext(stream)
+    await stream.aclose()
+
+    assert closed == [0]
+
+
+async def test_a_cancelled_viewer_is_not_redialled(monkeypatch):
+    """The real sessions swallow CancelledError and just end. Without a check,
+    the now-unlimited reconnect loop would read that as a routine drop."""
+    opened: list[int] = []
+    reading = asyncio.Event()
+    swallow = [True]
+
+    def _fake_stream_rtsp(_url, _fps, on_process=None):
+        opened.append(len(opened))
+
+        async def _gen():
+            yield FRAME
+            try:
+                reading.set()
+                await asyncio.Event().wait()  # blocked on the camera
+            except asyncio.CancelledError:
+                if not swallow[0]:
+                    raise
+                return  # swallowed, exactly like _stream_rtsp
+
+        return _gen()
+
+    monkeypatch.setattr(external_camera, "_stream_rtsp", _fake_stream_rtsp)
+
+    async def _viewer():
+        async for _ in external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10):
+            pass
+
+    task = asyncio.create_task(_viewer())
+    await asyncio.wait_for(reading.wait(), timeout=5)
+    task.cancel()
+    # asyncio.wait, not wait_for: on a regression the loop redials forever,
+    # and wait_for would hang waiting for the cancelled task to finish.
+    done, _ = await asyncio.wait({task}, timeout=5)
+    if not done:
+        swallow[0] = False
+        task.cancel()
+        await asyncio.wait({task}, timeout=5)
+
+    assert done, "the cancelled viewer's stream kept running"
+    assert opened == [0], "a cancelled viewer's stream must not open another session"