diff --git a/CHANGELOG.md b/CHANGELOG.md index 1391fe32b..c728ba421 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -45,6 +45,7 @@ All notable changes to Bambuddy will be documented in this file. - **Reopening the camera quickly could leave the new stream invisible to Bambuddy (#2707)** — Closing a camera view and opening it again straight away could leave the newly started stream unregistered, even though it was running and showing frames. The consequences were all indirect, which is what made it hard to spot: Bambuddy believed no viewer was attached, so Obico polling and snapshots would open a second camera connection and fight the live view — precisely what the guards added in #1348 and #1271 exist to prevent; the background cleanup task saw a camera process with no stream attached to it and killed the live stream as an orphan, usually within a minute; and pressing Stop reported that it had stopped nothing while the view was still running. **Root cause.** Each printer's fan-out stream was registered under a key derived from the printer alone, so every successive stream for that printer reused the same key, and the departing stream's cleanup removed whatever was registered under it — including its own replacement. The same cleanup also cleared the printer's most recent camera frame unconditionally, discarding the new stream's frame. It needed the old and new streams to overlap, which the four-second teardown fixed above made easy to hit. **Fix.** Each stream now gets its own registry key, so one stream can only ever clean up after itself — the same approach the external-camera path already uses (#2675) — and the shared per-printer frame is only released when no stream for that printer is left running. Covered by tests for both halves, including one that drives the real cleanup path with a second stream already registered. - **Closing the camera held the printer's camera connection for four more seconds, then logged an error that wasn't true (#2707)** — Every time a camera view closed, the log recorded `ffmpeg didn't terminate gracefully, killing` and then `ffmpeg did not exit within 2.0s of SIGKILL; abandoning wait`. Both waits expired every single time, so each close cost a fixed four seconds — and because Bambu firmware allows exactly one camera connection, that was four seconds in which nothing else could use the camera: reopening the view, a snapshot, Obico, or the diagnostic. **Root cause.** ffmpeg is started with its output and error streams as pipes, and the shutdown path had stopped reading them. A process whose output pipe is full blocks mid-write, and ffmpeg's shutdown signal only sets a flag that its main loop checks on the next pass, so the polite request could never be acted on and the grace period was dead time. The forced kill did work — but Python cannot report an exit while a pipe is still unread, so the second wait expired too and Bambuddy concluded the process was stuck when it had already gone. Measured at 4.00s per close before, ~0.15s after. **Fix.** Both pipes are now drained while the process is being stopped, which makes the polite shutdown effective and the exit observable. The forced-kill path and its time limit remain as backstops, so a genuinely wedged process still can't hang a stream, a Stop request, or the cleanup task. This also corrects the conclusion recorded for #2580: that 12-hour hang was the unbounded form of this same self-inflicted stall rather than a stuck ffmpeg, so bounding the wait had capped the symptom without removing the cause. Covered by tests that drive a real subprocess — the fault lives in Python's pipe bookkeeping, so a stand-in object would pass against the broken code — including one that verifies a process ignoring the polite signal still has its forced exit observed rather than abandoned. - **Two snapshots taken at the same moment opened two competing camera connections (#2705, reporter @gzimbric)** — Bambu firmware allows exactly one camera connection at a time. Bambuddy already knew this: a snapshot taken while somebody is watching the live view reuses the viewer's frame instead of opening a second socket. What nothing covered was two *snapshots* overlapping with no viewer attached at all — an Obico poll and a printer-wall refresh landing 200 ms apart, each correctly concluding it wasn't competing with a viewer, and then colliding with each other. On the reporter's P2S this knocked over the live stream that was feeding the camera wall, which was then reaped for having received no frames for 58 seconds. Eight paths take one-shot frames independently — Obico polling, `/camera/snapshot`, the finish-photo capture and its disk-writing sibling, plate detection, the camera connection test, and the diagnostic — so any pair of them could overlap, and a shorter Obico interval widened the window. **Fix.** Simultaneous captures for the same printer now share one connection: the first opens it, everyone arriving while it is in flight gets the same frame. Every consumer here wants "a recent frame" rather than a frame stamped at its own microsecond, so identical bytes are the right answer. This shares captures, it does not cache them — a request arriving after the previous capture finished still takes a fresh frame, because plate detection and the finish photo judge a running print from these images and a stale frame there is worse than a slow one. Each caller keeps its own deadline (they range from 10 to 30 seconds) rather than inheriting whichever one happened to open the connection, giving up alone leaves the capture running for whoever else is waiting on it, and a capture that fails doesn't hand its failure to callers that never got an attempt of their own — they retry, which by then competes with nothing. One visible consequence: when the **Diagnose** tool shares a capture this way its frame-capture stage is labelled `coalesced_capture`, because the pass is real but the timing shown is mostly time spent waiting, and a diagnostic must not report on a connection it never opened. Wiki updated. Covered by tests for the reported collision, the five-callers-one-connection case the reporter verified on live hardware, per-printer isolation, staying coalescing rather than becoming a cache, registry cleanup, a failed capture not poisoning its followers, bounded retry, a follower abandoning its wait without sabotaging the capture, and cancellation from either side. +- **The same one-shot-capture collision could happen on external cameras too, with no viewer attached (#2707 follow-up, reporter @bitbarista)** — #2705 fixed simultaneous captures colliding on the built-in camera path, keyed by printer IP through `capture_camera_frame_bytes()`. External cameras reach the same kind of collision through a different function — `external_camera.capture_frame()` — that #2705 didn't touch, and a V4L2 USB device allows exactly one open handle just like Bambu's own RTSP limit. Nothing coalesced two one-shot capturers here either: Obico polling, the in-print frame bank, the finish-photo moment, plate detection and the notification snapshot could each open their own connection to the same USB camera and collide, with `is_stream_active()` unable to help since that guard only stops a capturer from competing with an *attached viewer*, not with another capturer. **Fix.** The same shape of fix as #2705, applied to `capture_frame()`: concurrent callers for the same camera (URL, type, and — since #1177's snapshot override routes to a different endpoint entirely — snapshot URL) share one capture rather than opening a second connection. Coalesces, does not cache, so a call after the previous one finishes always captures fresh. Each caller keeps its own timeout, giving up leaves the capture running for whoever else is waiting, and a capture that fails doesn't hand its failure to a caller that never got a turn of its own. Covered by tests mirroring #2705's: the reported-shape collision, five callers sharing one connection, per-camera and per-snapshot-URL isolation, staying coalescing rather than becoming a cache, registry cleanup, a failed capture not poisoning its followers, bounded retry, a follower abandoning its wait without sabotaging the capture, and cancellation from either side. - **Auto-matched filament showed a green tick when the colour was plainly wrong (#2687, reporter @pchulpjoost)** — The Filament Mapping panel reported a slot as matched, with the header reading **(Ready)**, while the swatch beside it showed the slice wanted dark red and the tray it had picked held Dark Green. Manually selecting that very same tray from the dropdown correctly reported the colour mismatch, which is what made the disagreement so visible. **Root cause.** Auto-match ranks candidate trays by filament preset ID (`tray_info_idx`) first, and when exactly one loaded tray carried the preset the slice asked for, that tray was accepted as a *definitive* match on the assumption "same preset means same spool, so the colour must agree too". The preset ID names the **variant**, not the spool — `GFA00` is PLA Basic, `GFA01` PLA Matte, `GFA17` PLA Translucent, in every colour Bambu sells it. So a user with one Matte spool loaded matched every Matte requirement regardless of colour, and the colour comparison was never reached. This is why the report came in for PLA Matte in particular: generic PLA Basic is usually loaded several times over, which sent the match down a different path that did compare colours correctly. **Fix.** The colour verdict is now taken from the tray that was actually selected, never from which rule selected it, and the automatic and manual paths share one comparison so they cannot drift apart again. The preset still decides *selection*, because the Basic/Matte/Silk distinction matters ([#2650](https://github.com/maziggy/bambuddy/issues/2650)) — a wrong-coloured tray of the right variant is still chosen, but it is now reported as an amber **Color mismatch** instead of a green tick, and you can print anyway or pick another slot. A near-enough shade still counts as a match, and a 3MF that specifies no colour for a slot is satisfied by any colour rather than being flagged. Dispatch behaviour is unchanged: **Force color match** already required an exact colour before sending a job, so nothing was ever printed in the wrong colour because of this — the panel was simply telling you it was fine when it wasn't. Frontend-only. Wiki updated. Covered by tests for the unique-preset wrong-colour case, agreement between the auto and manual verdicts, the near-shade and colourless-requirement cases, and the multi-preset path that already worked. - **P1-series archives kept the worse finish photo when the timelapse arrived late (#2704 follow-up)** — When a print records a timelapse, Bambuddy prefers the video's last frame as the finish photo: the firmware stops recording after the toolhead parks but before the end G-code drops the bed, so it frames the finished print properly, where a live camera grab at that moment catches an already-lowered plate. Bambuddy waited 60 seconds for the video and then gave up, because the print-complete notification is waiting on that photo and holding a notification for minutes is worse than sending it with the live grab. On P1-series printers the video usually arrives later than that — they write MJPEG AVI instead of H.264 MP4 and serve it slowly, so across the support bundles their median was 33 seconds but the 90th percentile was 167 and the slowest observed was 546; every other model finished inside 26 seconds. The result was that the printers most in need of the better photo were the ones that never got it. **Fix.** The notification still goes out on the same 60-second bound with the live grab, so nothing gets slower. If the video was still on its way when that bound expired, Bambuddy now keeps waiting in the background and adds the extracted frame to the archive when it lands, at the front of the photo list so opening the gallery shows it first. The live grab is kept rather than replaced — the notification that already went out links to that exact file, and removing it would leave a broken image in Discord or Telegram. Covered by tests for the ordering, the longer budget, idempotency and the cases where the video never arrives. - **Timelapses that never got attached, and a Scan button that could not find them (#2704)** — Timelapse was on for the print, the video never arrived in the archive, and pressing **Scan for Timelapse** afterwards turned up nothing. Measured across 247 support bundles, this was not rare: of 457 automatic scans only 262 ever attached a video. **Root cause, part one.** The scan looked four times, at 5, 10, 20 and 30 seconds, then stopped. The printer writes the video only after the print ends and a long print makes a large file, so it often arrived after the last look — the attempt that found the video was the first one 272 times and then 17 / 13 / 13, a flat tail against the cutoff rather than a decaying one. What ran after those four attempts was a fallback that searched for the print's name inside the video filename; Bambu firmware only ever writes `video_`, so in 247 bundles it fired 159 times and matched exactly zero. **Root cause, part two.** The manual Scan button had no such snapshot to work from and matched by filename timestamp, by FTP modification time, or by there being exactly one video on the printer — all of which read a clock the printer cannot set, because a printer in LAN Only mode never reaches Bambu's time server. The reporter's P1S was six and a half days out, which defeats every one of those. **Fix.** The automatic scan now polls for several minutes instead of giving up after about a minute, and the name-match fallback is gone. The list of videos present when the print started is saved with the archive, so the comparison survives a Bambuddy restart mid-print and the manual Scan button can use it too — same clock-independent comparison, no timestamps anywhere. When a previous print's video lands late and two files look new, the one already attached to another archive is ruled out by name rather than by picking whichever the printer listed first, which could attach the wrong video. **Bambuddy now deletes a timelapse from the printer once it has been archived**, which keeps the printer's folder down to unclaimed videos and stops P1-series cards filling up with AVIs; your copy is in the archive, where you can watch, edit, download or remove it. That delete only happens after the transfer has been checked against the size the printer reported — which also fixes a silent truncation: an FTPS transfer that ended early produced a partial video that was attached as though it were complete. Because the first look happens seconds after the print ends — while the printer may still be writing the video — the file is also re-checked afterwards and only accepted once it has stopped growing, so a partial video is never mistaken for a finished one and the printer's copy is never removed on the strength of one. Wiki updated. Covered by tests for candidate selection, the download check gating the delete, the poll bounds, baseline persistence and the manual scan. diff --git a/backend/app/services/external_camera.py b/backend/app/services/external_camera.py index 283073a42..15a2e07c6 100644 --- a/backend/app/services/external_camera.py +++ b/backend/app/services/external_camera.py @@ -8,6 +8,7 @@ to ensure they are well-formed before use. """ import asyncio +import functools import logging import re import shutil @@ -175,6 +176,57 @@ def get_ffmpeg_path() -> str | None: return None +# In-flight one-shot captures, keyed by (url, camera_type, snapshot_url) — +# the tuple that actually identifies the physical resource being contended +# (#2707 comment thread, following #2705's shape for the built-in path). +# +# V4L2 USB devices allow exactly one open handle, and is_stream_active() / +# try_get_active_buffered_frame() (#2707) only stop a one-shot capturer from +# competing with the fan-out live view. They do nothing for capturer-vs- +# capturer with no viewer attached, where every consumer correctly concludes +# it isn't competing with a viewer and then collides with the others - +# exactly the #2705 report, just for this module's callers instead of +# capture_camera_frame_bytes()'s (Obico polling, the in-print frame bank, +# the finish-photo moment, plate detection, and the notification snapshot +# all reach capture_frame() independently). +# +# snapshot_url is part of the key (not just url/camera_type) because it +# routes to a completely different endpoint (#1177) - two printers that +# share a camera_url but differ only in snapshot_url must not coalesce. +_inflight_captures: dict[tuple[str, str, str | None], asyncio.Task[bytes | None]] = {} + + +def capture_in_flight(url: str, camera_type: str, snapshot_url: str | None = None) -> bool: + """Return True iff a one-shot capture for this key is running right now. + + Mirrors camera.py's capture_in_flight() for the built-in path - for a + caller that needs to know it will JOIN someone else's capture rather + than open its own connection. Ordinary consumers should ignore this: + they want "a recent frame", and capture_frame() already does the right + thing for them. + """ + task = _inflight_captures.get((url, camera_type, snapshot_url)) + return task is not None and not task.done() + + +def _discard_inflight_capture(key: tuple[str, str, str | None], task: asyncio.Task) -> None: + """Done-callback: drop the finished task from the in-flight registry. + + Guarded on identity so a slow task that finishes after a newer capture + has registered for the same key can't evict its successor. + + Also retrieves the exception, if any: the leader normally awaits the + task and would surface it, but a leader whose own caller was cancelled + leaves nobody to collect it, and an unretrieved task exception is + logged by asyncio as a warning with a traceback at an arbitrary later + point otherwise. + """ + if _inflight_captures.get(key) is task: + del _inflight_captures[key] + if not task.cancelled() and task.exception() is not None: + logger.debug("In-flight external-camera capture for %s ended in an exception", key[0]) + + async def capture_frame( url: str, camera_type: str, @@ -186,7 +238,10 @@ async def capture_frame( Args: url: Live-stream URL (MJPEG stream, RTSP URL, HTTP snapshot URL, or USB device path). camera_type: "mjpeg", "rtsp", "snapshot", or "usb". - timeout: Connection timeout in seconds. + timeout: Connection timeout in seconds. Applies to this caller's own + wait, including when it joins another caller's capture - call + sites disagree about the value, and a follower must not silently + inherit the leader's deadline in either direction. snapshot_url: Optional override for single-frame capture. When set, fetched via plain HTTP GET regardless of `camera_type`. Bypasses MJPEG warm-up handling on sources that expose a dedicated frame endpoint (e.g. go2rtc's @@ -195,6 +250,76 @@ async def capture_frame( Returns: JPEG bytes or None on failure + + Concurrent callers for the same (url, camera_type, snapshot_url) share + one capture (#2705-shape fix, filed for the external-camera path as a + follow-up on #2707): the first opens the connection, everyone arriving + while it's in flight awaits the same result. This coalesces; it does + not cache - a call that arrives after the previous capture finished + always captures fresh, since plate detection and the finish-photo path + judge a running print from these frames and a stale one there is worse + than a slow one (#1397). + """ + key = (url, camera_type, snapshot_url) + + # A follower whose leader fails takes a turn of its own rather than + # inheriting a failure it never had a chance to avoid - by then the + # leader has finished, so there's no connection left to compete with. + # Bounded at two rounds: if the capture we joined AND its replacement + # both failed, a third attempt won't help, and this caller has already + # spent its patience. + for _ in range(2): + leader = _inflight_captures.get(key) + if leader is None or leader.done(): + break + try: + frame = await asyncio.wait_for(asyncio.shield(leader), timeout=timeout) + except TimeoutError: + # shield() keeps the capture running for whoever else is still + # waiting on it - giving up is this caller's decision alone. + logger.warning("Gave up waiting %ss on the in-flight external-camera capture for %s", timeout, key[0]) + return None + except asyncio.CancelledError: + # Distinguish "the capture I joined was cancelled" from "I was + # cancelled". Only the former is ours to recover from. + if not leader.cancelled(): + raise + logger.info("In-flight external-camera capture for %s was cancelled; capturing our own", key[0]) + continue + if frame is not None: + logger.debug( + "Reusing in-flight external-camera capture for %s: %d bytes (no second connection opened)", + key[0], + len(frame), + ) + return frame + logger.debug("In-flight external-camera capture for %s failed; capturing our own", key[0]) + else: + return None + + task = asyncio.create_task(_capture_frame_uncoalesced(url, camera_type, timeout, snapshot_url)) + _inflight_captures[key] = task + task.add_done_callback(functools.partial(_discard_inflight_capture, key)) + # No wait_for here: this caller IS the capture, and each dispatched + # _capture_* function already enforces `timeout` internally, where it + # can also kill the ffmpeg process - a second deadline on top would + # abandon the subprocess instead of killing it. shield() so a cancelled + # leader (a client navigating away mid-request is routine) doesn't take + # the capture down with it - followers already waiting on it still get + # their frame. + return await asyncio.shield(task) + + +async def _capture_frame_uncoalesced( + url: str, + camera_type: str, + timeout: int, + snapshot_url: str | None, +) -> bytes | None: + """Open a connection and capture one frame. See capture_frame(). + + Callers want that wrapper, not this: it opens a connection + unconditionally, which is the collision #2705/#2707 are about. """ if snapshot_url: # Redact before truncating — slicing first can cut the URL short of the diff --git a/backend/tests/unit/services/test_external_camera_capture_coalescing.py b/backend/tests/unit/services/test_external_camera_capture_coalescing.py new file mode 100644 index 000000000..c40fcc2d7 --- /dev/null +++ b/backend/tests/unit/services/test_external_camera_capture_coalescing.py @@ -0,0 +1,308 @@ +"""Single-flight coalescing of one-shot external-camera captures (#2705-shape +fix, filed against the external-camera path as a follow-up on #2707). + +V4L2 USB devices allow exactly one open handle - the same one-connection +limit #2705 covers for Bambu firmware. The #2707 guards (``is_stream_active`` +/ ``try_get_active_buffered_frame``) only keep a one-shot capturer from +competing with the fan-out live view; nothing kept the capturers from +competing with EACH OTHER when no viewer is attached, so an Obico poll and +the in-print frame bank (say) could each open their own connection to the +same USB device and collide. + +These tests drive ``capture_frame`` at the public boundary and count how +many times the underlying capture ran, since "how many connections did we +open" is the entire point of the fix. Mirrors +``test_camera_capture_coalescing.py``'s structure for the built-in path. +""" + +import asyncio + +import pytest + +from backend.app.services import external_camera as ec_module +from backend.app.services.external_camera import capture_frame, capture_in_flight + +FRAME_A = b"\xff\xd8" + b"a" * 200 + b"\xff\xd9" +FRAME_B = b"\xff\xd8" + b"b" * 200 + b"\xff\xd9" + + +@pytest.fixture(autouse=True) +def _clear_inflight(): + """The registry is module-global; don't leak tasks between tests.""" + ec_module._inflight_captures.clear() + yield + ec_module._inflight_captures.clear() + + +class RecordingCapture: + """Stand-in for the real capture, recording each call. + + ``gate`` (when set) holds every capture open until released, which is how + these tests create the overlap window that used to produce two + connections. + """ + + def __init__(self, frames=(FRAME_A, FRAME_B), gate: asyncio.Event | None = None): + self.calls: list[tuple[str, str, str | None, int]] = [] + self._frames = list(frames) + self._gate = gate + self.started = asyncio.Event() + + async def __call__(self, url, camera_type, timeout, snapshot_url): + self.calls.append((url, camera_type, snapshot_url, timeout)) + self.started.set() + if self._gate is not None: + await self._gate.wait() + return self._frames.pop(0) if self._frames else None + + @property + def count(self) -> int: + return len(self.calls) + + +@pytest.fixture +def patch_capture(monkeypatch): + def _install(capture): + monkeypatch.setattr(ec_module, "_capture_frame_uncoalesced", capture) + return capture + + return _install + + +async def _let_leader_start(capture: RecordingCapture) -> None: + """Wait until the leader is inside the capture, so the next caller joins it. + + Without this the second caller can reach the registry before the first + has even been scheduled, which tests a different (and uninteresting) race. + """ + await asyncio.wait_for(capture.started.wait(), timeout=1) + + +@pytest.mark.asyncio +async def test_simultaneous_callers_share_one_capture(patch_capture): + """The reported collision: two consumers, one connection, two frames.""" + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=20)) + await _let_leader_start(capture) + follower = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=15)) + await asyncio.sleep(0) + gate.set() + + assert await leader == FRAME_A + assert await follower == FRAME_A + assert capture.count == 1 + + +@pytest.mark.asyncio +async def test_five_callers_one_capture(patch_capture): + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + first = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await _let_leader_start(capture) + rest = [asyncio.create_task(capture_frame("/dev/video1", "usb")) for _ in range(4)] + await asyncio.sleep(0) + gate.set() + + assert await asyncio.gather(first, *rest) == [FRAME_A] * 5 + assert capture.count == 1 + + +@pytest.mark.asyncio +async def test_different_cameras_do_not_coalesce(patch_capture): + """The one-connection limit is per camera, so the key must be too.""" + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + one = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await _let_leader_start(capture) + two = asyncio.create_task(capture_frame("/dev/video2", "usb")) + await asyncio.sleep(0) + gate.set() + + assert {await one, await two} == {FRAME_A, FRAME_B} + assert capture.count == 2 + assert {url for url, *_ in capture.calls} == {"/dev/video1", "/dev/video2"} + + +@pytest.mark.asyncio +async def test_different_snapshot_url_does_not_coalesce(patch_capture): + """#1177's snapshot_url override routes to a different endpoint entirely - + two printers sharing a camera_url but differing only in snapshot_url must + not share a capture.""" + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + one = asyncio.create_task(capture_frame("http://cam/", "mjpeg", snapshot_url="http://cam/frame1.jpg")) + await _let_leader_start(capture) + two = asyncio.create_task(capture_frame("http://cam/", "mjpeg", snapshot_url="http://cam/frame2.jpg")) + await asyncio.sleep(0) + gate.set() + + assert {await one, await two} == {FRAME_A, FRAME_B} + assert capture.count == 2 + + +@pytest.mark.asyncio +async def test_coalescing_is_not_caching(patch_capture): + """Sequential callers each capture fresh. + + Deliberate: plate detection and the finish-photo path decide things about + a running print from these frames, and #1397 was a finish photo a few + seconds stale showing the bed already lowered. + """ + capture = patch_capture(RecordingCapture()) + + assert await capture_frame("/dev/video1", "usb") == FRAME_A + assert await capture_frame("/dev/video1", "usb") == FRAME_B + assert capture.count == 2 + + +@pytest.mark.asyncio +async def test_registry_is_empty_after_a_capture_finishes(patch_capture): + """No leak, and nothing left behind for the next caller to join.""" + patch_capture(RecordingCapture()) + + await capture_frame("/dev/video1", "usb") + await asyncio.sleep(0) # let the done-callback run + + assert ec_module._inflight_captures == {} + assert capture_in_flight("/dev/video1", "usb") is False + + +@pytest.mark.asyncio +async def test_failed_leader_does_not_poison_its_followers(patch_capture): + """A follower that never got its own attempt gets one when the leader fails. + + Safe by then: the leader has finished, so there is no connection to + compete with. This also covers the follower whose timeout is LONGER than + the leader's — it isn't cut short by someone else's deadline. + """ + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(frames=(None, FRAME_B), gate=gate)) + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=10)) + await _let_leader_start(capture) + follower = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=20)) + await asyncio.sleep(0) + gate.set() + + assert await leader is None + assert await follower == FRAME_B + assert capture.count == 2 + + +@pytest.mark.asyncio +async def test_two_consecutive_failures_give_up(patch_capture): + """Bounded retry: a follower doesn't chase failing captures forever. + + Two followers behind a failing leader. The first takes its own turn, the + second joins THAT capture, and when it fails too the second gives up + rather than opening a third connection. + """ + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(frames=(None, None), gate=gate)) + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await _let_leader_start(capture) + first = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await asyncio.sleep(0) + second = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await asyncio.sleep(0) + gate.set() + + assert await leader is None + assert await first is None + assert await second is None + # The leader's capture plus one retry — not one per disappointed caller. + assert capture.count == 2 + + +@pytest.mark.asyncio +async def test_follower_timeout_does_not_sabotage_the_capture(patch_capture): + """A follower giving up leaves the capture running for everyone else. + + Call sites disagree about the timeout, so a follower must be able to + abandon a join without cancelling a capture other callers are still + waiting on. + """ + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=30)) + await _let_leader_start(capture) + impatient = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=0.01)) + patient = asyncio.create_task(capture_frame("/dev/video1", "usb", timeout=30)) + + assert await impatient is None # gave up on its own deadline + gate.set() + + assert await leader == FRAME_A + assert await patient == FRAME_A # unaffected by the one that walked away + assert capture.count == 1 + + +@pytest.mark.asyncio +async def test_cancelled_leader_still_delivers_to_followers(patch_capture): + """Snapshot/capture requests get cancelled routinely (client navigates + away mid-request). The follower must not lose the frame because the + caller that happened to open the connection went away.""" + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await _let_leader_start(capture) + follower = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await asyncio.sleep(0) + + leader.cancel() + with pytest.raises(asyncio.CancelledError): + await leader + gate.set() + + assert await follower == FRAME_A + assert capture.count == 1 + + +@pytest.mark.asyncio +async def test_cancelling_a_follower_leaves_the_leader_alone(patch_capture): + """The mirror case: the follower's cancellation is its own business.""" + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await _let_leader_start(capture) + follower = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await asyncio.sleep(0) + + follower.cancel() + with pytest.raises(asyncio.CancelledError): + await follower + gate.set() + + assert await leader == FRAME_A + assert capture.count == 1 + + +@pytest.mark.asyncio +async def test_capture_in_flight_reports_the_window(patch_capture): + """The predicate a diagnose-style caller would use to know it will join, + not measure its own connection.""" + gate = asyncio.Event() + capture = patch_capture(RecordingCapture(gate=gate)) + + assert capture_in_flight("/dev/video1", "usb") is False + + leader = asyncio.create_task(capture_frame("/dev/video1", "usb")) + await _let_leader_start(capture) + + assert capture_in_flight("/dev/video1", "usb") is True + assert capture_in_flight("/dev/video2", "usb") is False # per camera + + gate.set() + await leader + await asyncio.sleep(0) + + assert capture_in_flight("/dev/video1", "usb") is False