Files
bambuddy/backend/tests/unit/test_camera_reconnect_budget.py

267 lines
8.1 KiB
Python

"""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"