test_camera_reconnect_budget.py 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266
  1. """The camera reconnect budget counts failures in a row, not every drop.
  2. A stock X1C ends each RTSP session after about a minute. The built-in stream
  3. respawned ffmpeg transparently but counted every respawn against a lifetime
  4. budget of 30, so a live view stopped for good after about half an hour. The
  5. external-camera path allowed three reconnects for the life of the stream.
  6. """
  7. import asyncio
  8. from contextlib import suppress
  9. import pytest
  10. from backend.app.api.routes import camera
  11. from backend.app.services import external_camera
  12. from backend.app.services.camera_profiles import CameraProfile
  13. FRAME = b"\xff\xd8frame\xff\xd9"
  14. class _FakeServer:
  15. def close(self) -> None:
  16. pass
  17. async def wait_closed(self) -> None:
  18. pass
  19. class _Stdout:
  20. def __init__(self, frames: int) -> None:
  21. self._frames = frames
  22. async def read(self, _size: int = -1) -> bytes:
  23. if self._frames <= 0:
  24. return b""
  25. self._frames -= 1
  26. return FRAME
  27. class _Proc:
  28. _next_pid = 78000
  29. def __init__(self, frames: int) -> None:
  30. _Proc._next_pid += 1
  31. self.pid = _Proc._next_pid
  32. self.returncode = None
  33. self.stdout = _Stdout(frames)
  34. self.stderr = None
  35. def terminate(self) -> None:
  36. self.returncode = 0
  37. def kill(self) -> None:
  38. self.returncode = -9
  39. async def wait(self) -> int:
  40. if self.returncode is None:
  41. self.returncode = 0
  42. return self.returncode
  43. @pytest.fixture
  44. def rtsp(monkeypatch):
  45. """Fake ffmpeg sessions: each call spawns the next entry's frame count."""
  46. sessions: list[int] = []
  47. spawned: list[int] = []
  48. delays: list[float] = []
  49. real_sleep = asyncio.sleep
  50. async def _fake_exec(*_args, **_kwargs):
  51. frames = sessions.pop(0) if sessions else 0
  52. spawned.append(frames)
  53. return _Proc(frames)
  54. async def _fake_proxy(_ip: str, _port: int):
  55. return 48998, _FakeServer()
  56. async def _recording_sleep(seconds, *args, **kwargs):
  57. delays.append(seconds)
  58. await real_sleep(0)
  59. monkeypatch.setattr(camera, "get_ffmpeg_path", lambda: "/fake/ffmpeg")
  60. monkeypatch.setattr(camera, "create_tls_proxy", _fake_proxy)
  61. monkeypatch.setattr(camera.asyncio, "create_subprocess_exec", _fake_exec)
  62. monkeypatch.setattr(camera.asyncio, "sleep", _recording_sleep)
  63. return sessions, spawned, delays
  64. def _use_profile(monkeypatch, **kwargs):
  65. monkeypatch.setattr(camera, "get_camera_profile", lambda _model: CameraProfile(**kwargs))
  66. async def _drain(stream, limit: int = 200) -> int:
  67. frames = 0
  68. async for chunk in stream:
  69. if b"image/jpeg" in chunk:
  70. frames += 1
  71. if frames >= limit:
  72. break
  73. return frames
  74. def _stream():
  75. return camera.generate_rtsp_mjpeg_stream(
  76. ip_address="192.0.2.40",
  77. access_code="test-code",
  78. model="X1C",
  79. fps=10,
  80. stream_id="99-fanout-budget",
  81. disconnect_event=asyncio.Event(),
  82. )
  83. async def test_routine_drops_never_use_up_the_budget(rtsp, monkeypatch):
  84. """Ten sessions that each deliver video and end, with a budget of two."""
  85. sessions, spawned, _ = rtsp
  86. _use_profile(monkeypatch, rtsp_reconnect_max=2, rtsp_reconnect_delay=0.2)
  87. sessions.extend([3] * 10)
  88. stream = _stream()
  89. frames = await asyncio.wait_for(_drain(stream, limit=30), timeout=10)
  90. with suppress(Exception):
  91. await stream.aclose()
  92. assert frames == 30
  93. assert len(spawned) == 10
  94. async def test_failures_in_a_row_still_give_up(rtsp, monkeypatch):
  95. sessions, spawned, _ = rtsp
  96. _use_profile(monkeypatch, rtsp_reconnect_max=3, rtsp_reconnect_delay=0.2)
  97. sessions.extend([2]) # then every session fails without a frame
  98. frames = await asyncio.wait_for(_drain(_stream()), timeout=10)
  99. assert frames == 2
  100. # The good session, then three reconnects that each fail -- the budget of
  101. # three consecutive reconnects -- and the stream gives up.
  102. assert spawned == [2, 0, 0, 0]
  103. async def test_failures_back_off_and_a_good_session_resets_the_delay(rtsp, monkeypatch):
  104. sessions, _, delays = rtsp
  105. _use_profile(monkeypatch, rtsp_reconnect_max=6, rtsp_reconnect_delay=0.2, rtsp_reconnect_backoff_max=1.0)
  106. sessions.extend([1, 0, 0, 0, 1])
  107. await asyncio.wait_for(_drain(_stream()), timeout=10)
  108. # 0.1 is the post-spawn startup check, not a reconnect delay.
  109. reconnect_delays = [d for d in delays if d != 0.1]
  110. assert reconnect_delays == [0.2, 0.4, 0.8, 1.0, 0.2, 0.4, 0.8, 1.0, 1.0, 1.0]
  111. # ---------------------------------------------------------------------------
  112. # External cameras
  113. # ---------------------------------------------------------------------------
  114. @pytest.fixture
  115. def external(monkeypatch):
  116. sessions: list[int] = []
  117. opened: list[int] = []
  118. closed: list[int] = []
  119. def _fake_stream_rtsp(_url, _fps, on_process=None):
  120. frames = sessions.pop(0) if sessions else 0
  121. index = len(opened)
  122. opened.append(frames)
  123. async def _gen():
  124. try:
  125. for _ in range(frames):
  126. yield FRAME
  127. finally:
  128. closed.append(index)
  129. return _gen()
  130. async def _no_sleep(_seconds, *args, **kwargs):
  131. return None
  132. monkeypatch.setattr(external_camera, "_stream_rtsp", _fake_stream_rtsp)
  133. monkeypatch.setattr(external_camera.asyncio, "sleep", _no_sleep)
  134. return sessions, opened, closed
  135. async def test_external_routine_drops_keep_the_stream_going(external):
  136. """Used to end for good on the fourth drop."""
  137. sessions, opened, _ = external
  138. sessions.extend([2, 2, 2, 2, 2, 2])
  139. frames = [f async for f in external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10)]
  140. assert len(frames) == 12
  141. assert opened == [2, 2, 2, 2, 2, 2, 0], "a session with no frame ends the stream"
  142. async def test_external_stops_when_asked(external):
  143. sessions, opened, _ = external
  144. sessions.extend([1] * 10)
  145. stop = asyncio.Event()
  146. frames = []
  147. async for frame in external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10, stop_event=stop):
  148. frames.append(frame)
  149. if len(frames) == 3:
  150. stop.set()
  151. assert len(frames) == 3
  152. assert len(opened) == 3
  153. async def test_closing_the_external_stream_closes_the_open_session(external):
  154. """The session owns the ffmpeg process; it must stop when the viewer goes,
  155. not whenever the abandoned iterator happens to be collected."""
  156. sessions, _, closed = external
  157. sessions.extend([50])
  158. stream = external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10)
  159. await anext(stream)
  160. await stream.aclose()
  161. assert closed == [0]
  162. async def test_a_cancelled_viewer_is_not_redialled(monkeypatch):
  163. """The real sessions swallow CancelledError and just end. Without a check,
  164. the now-unlimited reconnect loop would read that as a routine drop."""
  165. opened: list[int] = []
  166. reading = asyncio.Event()
  167. swallow = [True]
  168. def _fake_stream_rtsp(_url, _fps, on_process=None):
  169. opened.append(len(opened))
  170. async def _gen():
  171. yield FRAME
  172. try:
  173. reading.set()
  174. await asyncio.Event().wait() # blocked on the camera
  175. except asyncio.CancelledError:
  176. if not swallow[0]:
  177. raise
  178. return # swallowed, exactly like _stream_rtsp
  179. return _gen()
  180. monkeypatch.setattr(external_camera, "_stream_rtsp", _fake_stream_rtsp)
  181. async def _viewer():
  182. async for _ in external_camera.generate_mjpeg_stream("rtsp://cam.test/live", "rtsp", 10):
  183. pass
  184. task = asyncio.create_task(_viewer())
  185. await asyncio.wait_for(reading.wait(), timeout=5)
  186. task.cancel()
  187. # asyncio.wait, not wait_for: on a regression the loop redials forever,
  188. # and wait_for would hang waiting for the cancelled task to finish.
  189. done, _ = await asyncio.wait({task}, timeout=5)
  190. if not done:
  191. swallow[0] = False
  192. task.cancel()
  193. await asyncio.wait({task}, timeout=5)
  194. assert done, "the cancelled viewer's stream kept running"
  195. assert opened == [0], "a cancelled viewer's stream must not open another session"