external_camera.py 54 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339
  1. """External camera service.
  2. Supports MJPEG streams, RTSP streams (via ffmpeg), HTTP snapshot URLs, and USB cameras.
  3. Security Note: This service intentionally makes requests to user-configured camera URLs.
  4. This is necessary functionality for external camera integration. URLs are validated
  5. to ensure they are well-formed before use.
  6. """
  7. import asyncio
  8. import contextlib
  9. import functools
  10. import ipaddress
  11. import logging
  12. import re
  13. import shutil
  14. import socket
  15. from collections.abc import AsyncGenerator, Callable
  16. from pathlib import Path
  17. from urllib.parse import urlparse
  18. import aiohttp
  19. from backend.app.core.logging_filters import redact_url_credentials
  20. from backend.app.utils.ffmpeg_output import NO_FFMPEG_OUTPUT, summarize_ffmpeg_stderr
  21. logger = logging.getLogger(__name__)
  22. # Protocols ffmpeg may use for an RTSP input. RTSP negotiates its media
  23. # transport at runtime, so the transports have to be here alongside rtsp itself;
  24. # tls and crypto cover encrypted variants. Everything ffmpeg would otherwise
  25. # accept behind an -i — file, http, tcp to anywhere, concat — is left out, so a
  26. # stream that references something outside itself cannot pull it in.
  27. _RTSP_PROTOCOL_WHITELIST = "rtsp,rtp,udp,tcp,tls,crypto"
  28. def _blocked_host_reason(hostname: str) -> str | None:
  29. """Describe why *hostname* is a destination we refuse to fetch, or None to allow it.
  30. Camera URLs are user-supplied and reach the network — over aiohttp for the
  31. HTTP types, and as an ``ffmpeg -i`` argument for RTSP — so this is where the
  32. SSRF boundary sits. LAN addresses are deliberately allowed: cameras live on
  33. the same network as Bambuddy, and blocking RFC-1918 would remove the feature
  34. rather than protect it. What is left to refuse is the host talking to
  35. itself, the unspecified address, link-local (which is where the cloud
  36. metadata endpoint lives), and the metadata hostnames.
  37. IP literals are classified with ``ipaddress`` rather than compared against a
  38. list of spellings, because 127.0.0.1, 127.0.0.2, 2130706433, 0177.0.0.1,
  39. 127.1 and ::ffff:127.0.0.1 all arrive at loopback and a list of strings only
  40. ever catches whichever one someone thought to write down. ``inet_aton``
  41. comes first because it accepts the legacy octal, decimal and short forms
  42. that ``ip_address`` rejects — the C resolvers behind aiohttp and ffmpeg
  43. accept them, so refusing to understand them here would only mean not seeing
  44. where the request is actually going.
  45. """
  46. host = hostname.lower()
  47. ip: ipaddress.IPv4Address | ipaddress.IPv6Address | None = None
  48. try:
  49. ip = ipaddress.ip_address(socket.inet_aton(host))
  50. except OSError:
  51. try:
  52. ip = ipaddress.ip_address(host)
  53. except ValueError:
  54. ip = None
  55. if ip is None:
  56. # A name, not an address. It is not resolved here on purpose: aiohttp
  57. # and ffmpeg each resolve independently afterwards, so a check here
  58. # decides nothing about where they end up (DNS rebinding), while a
  59. # lookup on every capture would break LAN cameras behind slow or
  60. # intermittent local DNS.
  61. if host == "localhost" or host.endswith(".localhost"):
  62. return "localhost"
  63. if host in ("metadata.google.internal", "metadata.google"):
  64. return "a cloud metadata service"
  65. return None
  66. # ::ffff:127.0.0.1 is loopback wearing an IPv6 spelling.
  67. mapped = getattr(ip, "ipv4_mapped", None)
  68. if mapped is not None:
  69. ip = mapped
  70. if ip.is_loopback:
  71. return "loopback"
  72. if ip.is_unspecified:
  73. return "the unspecified address"
  74. if ip.is_link_local:
  75. return "a link-local address (the cloud metadata range)"
  76. return None
  77. def _sanitize_camera_url(url: str, allowed_schemes: tuple[str, ...] = ("http", "https", "rtsp")) -> str | None:
  78. """Validate and sanitize camera URL, returning a safe reconstructed URL.
  79. This validates that the URL is well-formed, uses an allowed scheme, does not
  80. target the host itself or a cloud metadata service, and returns a URL
  81. reconstructed from the validated components.
  82. Note: This intentionally allows user-provided URLs as that is the
  83. purpose of external camera configuration. Local network IPs are
  84. allowed since cameras are typically on the same LAN.
  85. Args:
  86. url: URL to validate and sanitize
  87. allowed_schemes: Tuple of allowed URL schemes
  88. Returns:
  89. Sanitized URL string if valid, None otherwise
  90. """
  91. try:
  92. parsed = urlparse(url)
  93. if not parsed.scheme or not parsed.netloc:
  94. return None
  95. # Validate scheme against allowlist
  96. scheme = parsed.scheme.lower()
  97. if scheme not in allowed_schemes:
  98. return None
  99. hostname = parsed.hostname or ""
  100. if not hostname:
  101. return None
  102. blocked = _blocked_host_reason(hostname)
  103. if blocked:
  104. logger.warning("Blocked camera URL targeting %s: %s", blocked, hostname)
  105. return None
  106. # Reconstruct URL from validated components to break taint chain
  107. # This creates a new string from validated parts
  108. #
  109. # The credentials are carried across verbatim from netloc rather than
  110. # via parsed.username/.password, which urlparse has already percent-
  111. # decoded: re-emitting those would corrupt any password containing an
  112. # @ or a :. They have to survive at all because most RTSP cameras — and
  113. # a fair number of MJPEG ones — carry their login in the URL, and
  114. # dropping it turns every one of them into an authentication failure.
  115. netloc = parsed.netloc
  116. userinfo = f"{netloc.rsplit('@', 1)[0]}@" if "@" in netloc else ""
  117. # parsed.hostname has already stripped the brackets off an IPv6 literal;
  118. # without them back the result is not a URL any client can parse.
  119. host_str = f"[{hostname}]" if ":" in hostname else hostname
  120. port_str = f":{parsed.port}" if parsed.port else ""
  121. path = parsed.path or ""
  122. query = f"?{parsed.query}" if parsed.query else ""
  123. fragment = f"#{parsed.fragment}" if parsed.fragment else ""
  124. # Build sanitized URL from validated components
  125. sanitized = f"{scheme}://{userinfo}{host_str}{port_str}{path}{query}{fragment}"
  126. return sanitized
  127. except ValueError:
  128. return None
  129. def _validate_camera_url(url: str, allowed_schemes: tuple[str, ...] = ("http", "https", "rtsp")) -> bool:
  130. """Validate camera URL format (legacy wrapper).
  131. Args:
  132. url: URL to validate
  133. allowed_schemes: Tuple of allowed URL schemes
  134. Returns:
  135. True if URL is valid, False otherwise
  136. """
  137. return _sanitize_camera_url(url, allowed_schemes) is not None
  138. def list_usb_cameras() -> list[dict]:
  139. """List available USB cameras (V4L2 devices on Linux).
  140. Returns:
  141. List of dicts with {device: str, name: str, capabilities: list}
  142. """
  143. cameras = []
  144. video_devices = sorted(Path("/dev").glob("video*"))
  145. for device in video_devices:
  146. device_path = str(device)
  147. info = {"device": device_path, "name": device.name, "capabilities": []}
  148. # Try to get device info via v4l2-ctl
  149. v4l2_ctl = shutil.which("v4l2-ctl")
  150. if v4l2_ctl:
  151. import subprocess
  152. try:
  153. result = subprocess.run(
  154. [v4l2_ctl, "-d", device_path, "--info"],
  155. capture_output=True,
  156. text=True,
  157. timeout=5,
  158. )
  159. if result.returncode == 0:
  160. # Parse device name from output
  161. for line in result.stdout.splitlines():
  162. if "Card type" in line:
  163. info["name"] = line.split(":", 1)[1].strip()
  164. elif "Driver name" in line:
  165. info["driver"] = line.split(":", 1)[1].strip()
  166. # Check if device supports video capture
  167. result = subprocess.run(
  168. [v4l2_ctl, "-d", device_path, "--list-formats"],
  169. capture_output=True,
  170. text=True,
  171. timeout=5,
  172. )
  173. if result.returncode == 0 and result.stdout.strip():
  174. info["capabilities"].append("capture")
  175. # Parse available formats
  176. formats = re.findall(r"'(\w+)'", result.stdout)
  177. info["formats"] = list(set(formats))
  178. except (subprocess.TimeoutExpired, Exception) as e:
  179. logger.debug("v4l2-ctl failed for %s: %s", device_path, e)
  180. # Only include devices that look like video capture devices
  181. # Skip metadata devices (typically odd numbered like video1, video3)
  182. try:
  183. device_num = int(device.name.replace("video", ""))
  184. # Even numbered devices are usually capture, odd are metadata
  185. # But also check if we got capabilities
  186. if info.get("capabilities") or device_num % 2 == 0:
  187. cameras.append(info)
  188. except ValueError:
  189. cameras.append(info)
  190. return cameras
  191. def get_ffmpeg_path() -> str | None:
  192. """Get the path to ffmpeg executable."""
  193. # Try shutil.which first
  194. path = shutil.which("ffmpeg")
  195. if path:
  196. return path
  197. # Check common locations (systemd services may have limited PATH)
  198. for common_path in ["/usr/bin/ffmpeg", "/usr/local/bin/ffmpeg", "/opt/homebrew/bin/ffmpeg"]:
  199. if Path(common_path).exists():
  200. return common_path
  201. return None
  202. # In-flight one-shot captures, keyed by (url, camera_type, snapshot_url) —
  203. # the tuple that actually identifies the physical resource being contended
  204. # (#2707 comment thread, following #2705's shape for the built-in path).
  205. #
  206. # V4L2 USB devices allow exactly one open handle, and is_stream_active() /
  207. # try_get_active_buffered_frame() (#2707) only stop a one-shot capturer from
  208. # competing with the fan-out live view. They do nothing for capturer-vs-
  209. # capturer with no viewer attached, where every consumer correctly concludes
  210. # it isn't competing with a viewer and then collides with the others -
  211. # exactly the #2705 report, just for this module's callers instead of
  212. # capture_camera_frame_bytes()'s (Obico polling, the in-print frame bank,
  213. # the finish-photo moment, plate detection, and the notification snapshot
  214. # all reach capture_frame() independently).
  215. #
  216. # snapshot_url is part of the key (not just url/camera_type) because it
  217. # routes to a completely different endpoint (#1177) - two printers that
  218. # share a camera_url but differ only in snapshot_url must not coalesce.
  219. _inflight_captures: dict[tuple[str, str, str | None], asyncio.Task[bytes | None]] = {}
  220. def capture_in_flight(url: str, camera_type: str, snapshot_url: str | None = None) -> bool:
  221. """Return True iff a one-shot capture for this key is running right now.
  222. Mirrors camera.py's capture_in_flight() for the built-in path - for a
  223. caller that needs to know it will JOIN someone else's capture rather
  224. than open its own connection. Ordinary consumers should ignore this:
  225. they want "a recent frame", and capture_frame() already does the right
  226. thing for them.
  227. """
  228. task = _inflight_captures.get((url, camera_type, snapshot_url))
  229. return task is not None and not task.done()
  230. def _discard_inflight_capture(key: tuple[str, str, str | None], task: asyncio.Task) -> None:
  231. """Done-callback: drop the finished task from the in-flight registry.
  232. Guarded on identity so a slow task that finishes after a newer capture
  233. has registered for the same key can't evict its successor.
  234. Also retrieves the exception, if any: the leader normally awaits the
  235. task and would surface it, but a leader whose own caller was cancelled
  236. leaves nobody to collect it, and an unretrieved task exception is
  237. logged by asyncio as a warning with a traceback at an arbitrary later
  238. point otherwise.
  239. """
  240. if _inflight_captures.get(key) is task:
  241. del _inflight_captures[key]
  242. if not task.cancelled() and task.exception() is not None:
  243. logger.debug("In-flight external-camera capture for %s ended in an exception", _log_key(key))
  244. def _log_key(key: tuple[str, str, str | None]) -> str:
  245. """Render an in-flight key for a log line, with credentials redacted.
  246. Unlike camera.py's coalescing — which is keyed by IP address and so has
  247. nothing to hide — these keys carry the camera URL, and an RTSP camera URL
  248. routinely embeds ``user:pass@``. Redact before truncating: slicing first
  249. can cut the URL short of the ``@`` the pattern anchors on and leave the
  250. password in the log, which is why every other URL log in this module does
  251. it in this order.
  252. """
  253. return redact_url_credentials(key[0])[:50] if key[0] else "None"
  254. async def capture_frame(
  255. url: str,
  256. camera_type: str,
  257. timeout: int = 15,
  258. snapshot_url: str | None = None,
  259. ) -> bytes | None:
  260. """Capture single frame from external camera.
  261. Args:
  262. url: Live-stream URL (MJPEG stream, RTSP URL, HTTP snapshot URL, or USB device path).
  263. camera_type: "mjpeg", "rtsp", "snapshot", or "usb".
  264. timeout: Connection timeout in seconds. Applies to this caller's own
  265. wait, including when it joins another caller's capture - call
  266. sites disagree about the value, and a follower must not silently
  267. inherit the leader's deadline in either direction.
  268. snapshot_url: Optional override for single-frame capture. When set, fetched
  269. via plain HTTP GET regardless of `camera_type`. Bypasses MJPEG warm-up
  270. handling on sources that expose a dedicated frame endpoint (e.g. go2rtc's
  271. `/api/frame.jpeg` reliably returns a clean image while the MJPEG stream's
  272. first frame is often the encoder's stale keyframe). #1177.
  273. Returns:
  274. JPEG bytes or None on failure
  275. Concurrent callers for the same (url, camera_type, snapshot_url) share
  276. one capture (#2705-shape fix, filed for the external-camera path as a
  277. follow-up on #2707): the first opens the connection, everyone arriving
  278. while it's in flight awaits the same result. This coalesces; it does
  279. not cache - a call that arrives after the previous capture finished
  280. always captures fresh, since plate detection and the finish-photo path
  281. judge a running print from these frames and a stale one there is worse
  282. than a slow one (#1397).
  283. """
  284. key = (url, camera_type, snapshot_url)
  285. # A follower whose leader fails takes a turn of its own rather than
  286. # inheriting a failure it never had a chance to avoid - by then the
  287. # leader has finished, so there's no connection left to compete with.
  288. # Bounded at two rounds: if the capture we joined AND its replacement
  289. # both failed, a third attempt won't help, and this caller has already
  290. # spent its patience.
  291. for _ in range(2):
  292. leader = _inflight_captures.get(key)
  293. if leader is None or leader.done():
  294. break
  295. try:
  296. frame = await asyncio.wait_for(asyncio.shield(leader), timeout=timeout)
  297. except TimeoutError:
  298. # shield() keeps the capture running for whoever else is still
  299. # waiting on it - giving up is this caller's decision alone.
  300. logger.warning(
  301. "Gave up waiting %ss on the in-flight external-camera capture for %s", timeout, _log_key(key)
  302. )
  303. return None
  304. except asyncio.CancelledError:
  305. # Distinguish "the capture I joined was cancelled" from "I was
  306. # cancelled". Only the former is ours to recover from.
  307. if not leader.cancelled():
  308. raise
  309. logger.info("In-flight external-camera capture for %s was cancelled; capturing our own", _log_key(key))
  310. continue
  311. if frame is not None:
  312. logger.debug(
  313. "Reusing in-flight external-camera capture for %s: %d bytes (no second connection opened)",
  314. _log_key(key),
  315. len(frame),
  316. )
  317. return frame
  318. logger.debug("In-flight external-camera capture for %s failed; capturing our own", _log_key(key))
  319. else:
  320. return None
  321. task = asyncio.create_task(_capture_frame_uncoalesced(url, camera_type, timeout, snapshot_url))
  322. _inflight_captures[key] = task
  323. task.add_done_callback(functools.partial(_discard_inflight_capture, key))
  324. # No wait_for here: this caller IS the capture, and each dispatched
  325. # _capture_* function already enforces `timeout` internally, where it
  326. # can also kill the ffmpeg process - a second deadline on top would
  327. # abandon the subprocess instead of killing it. shield() so a cancelled
  328. # leader (a client navigating away mid-request is routine) doesn't take
  329. # the capture down with it - followers already waiting on it still get
  330. # their frame.
  331. return await asyncio.shield(task)
  332. async def _capture_frame_uncoalesced(
  333. url: str,
  334. camera_type: str,
  335. timeout: int,
  336. snapshot_url: str | None,
  337. ) -> bytes | None:
  338. """Open a connection and capture one frame. See capture_frame().
  339. Callers want that wrapper, not this: it opens a connection
  340. unconditionally, which is the collision #2705/#2707 are about.
  341. Failure is reported as ``None``, never as an exception. That is load-
  342. bearing now that captures are shared: the coalescing wrapper hands one
  343. task's outcome to every caller waiting on it, and it can only give a
  344. follower its own turn for an outcome it can recognise. An exception
  345. escaping here would instead propagate to every follower at once —
  346. turning one caller's failure into N — and none of them would retry.
  347. The per-type helpers below each catch what they expect and return None,
  348. but they catch narrowly (``aiohttp.ClientError``/``OSError``/timeouts),
  349. so this is the structural guarantee rather than one contingent on their
  350. coverage. Mirrors ``_capture_camera_frame_bytes_uncoalesced`` in
  351. camera.py, which ends in the same blanket catch for the same reason.
  352. """
  353. try:
  354. if snapshot_url:
  355. # Redact before truncating — slicing first can cut the URL short of the
  356. # ``@`` the pattern anchors on and leave the password in the log.
  357. logger.debug("capture_frame using snapshot override url=%s...", redact_url_credentials(snapshot_url)[:50])
  358. return await _capture_snapshot(snapshot_url, timeout)
  359. logger.debug(
  360. "capture_frame called: type=%s, url=%s...",
  361. camera_type,
  362. redact_url_credentials(url)[:50] if url else "None",
  363. )
  364. if camera_type == "mjpeg":
  365. return await _capture_mjpeg_frame(url, timeout)
  366. elif camera_type == "rtsp":
  367. return await _capture_rtsp_frame(url, timeout)
  368. elif camera_type == "snapshot":
  369. return await _capture_snapshot(url, timeout)
  370. elif camera_type == "usb":
  371. return await _capture_usb_frame(url, timeout)
  372. else:
  373. logger.warning("Unknown camera type: %s", camera_type)
  374. return None
  375. except asyncio.CancelledError:
  376. # Cancellation is not a capture failure and must stay distinguishable:
  377. # the wrapper checks ``leader.cancelled()`` to decide whether a
  378. # follower may take its own turn.
  379. raise
  380. except Exception:
  381. logger.exception("External camera capture failed for %s", redact_url_credentials(url)[:50] if url else "None")
  382. return None
  383. def _safe_usb_device_path(device: str) -> str | None:
  384. """Rebuild a /dev/videoN path from a validated device number, or None.
  385. Validate device path - must be /dev/videoN format where N is 0-99. This
  386. prevents path traversal by using a strict allowlist approach: the returned
  387. path is built from an integer, which cannot carry a traversal, rather than
  388. from any part of the caller's string.
  389. Returns None if the device does not exist, so a caller cannot hand ffmpeg a
  390. path to something that is not a device node.
  391. """
  392. device_match = re.match(r"^/dev/video(\d{1,2})$", device)
  393. if not device_match:
  394. logger.error("Invalid USB device path format: %s", device)
  395. return None
  396. # Convert to integer to break taint chain - integers cannot contain path traversal
  397. # lgtm[py/path-injection] - device_num is validated integer 0-99
  398. device_num = int(device_match.group(1)) # Safe: regex guarantees 1-2 digits
  399. # Construct safe path from validated integer (completely untainted)
  400. safe_device_path = Path(f"/dev/video{device_num}") # lgtm[py/path-injection]
  401. if not safe_device_path.exists():
  402. logger.error("USB device does not exist: %s", safe_device_path)
  403. return None
  404. return str(safe_device_path) # lgtm[py/path-injection]
  405. async def _capture_usb_frame(device: str, timeout: int) -> bytes | None:
  406. """Capture frame from USB camera using ffmpeg."""
  407. ffmpeg = get_ffmpeg_path()
  408. if not ffmpeg:
  409. logger.error("ffmpeg not found - required for USB camera capture")
  410. return None
  411. safe_device = _safe_usb_device_path(device)
  412. if not safe_device:
  413. return None
  414. # Use the safe path for ffmpeg - this is a hardcoded /dev/videoN path
  415. device = safe_device # lgtm[py/path-injection]
  416. # Use ffmpeg to grab a single frame from USB camera
  417. cmd = [
  418. ffmpeg,
  419. "-f",
  420. "v4l2",
  421. "-i",
  422. device,
  423. "-frames:v",
  424. "1",
  425. "-f",
  426. "image2pipe",
  427. "-vcodec",
  428. "mjpeg",
  429. "-q:v",
  430. "2",
  431. "-",
  432. ]
  433. try:
  434. logger.debug("Running USB capture: %s", " ".join(cmd))
  435. process = await asyncio.create_subprocess_exec(
  436. *cmd,
  437. stdout=asyncio.subprocess.PIPE,
  438. stderr=asyncio.subprocess.PIPE,
  439. )
  440. stdout, stderr = await asyncio.wait_for(process.communicate(), timeout=timeout)
  441. if process.returncode != 0:
  442. logger.error("ffmpeg USB capture failed: %s", summarize_ffmpeg_stderr(stderr) or NO_FFMPEG_OUTPUT)
  443. return None
  444. if not stdout or len(stdout) < 100:
  445. logger.error("ffmpeg returned empty or too small frame from USB camera")
  446. return None
  447. return stdout
  448. except TimeoutError:
  449. logger.warning("USB frame capture timed out after %ss", timeout)
  450. if process:
  451. process.kill()
  452. return None
  453. except OSError as e:
  454. logger.error("USB frame capture failed: %s", e)
  455. return None
  456. async def _capture_mjpeg_frame(url: str, timeout: int) -> bytes | None:
  457. """Extract a single representative frame from an MJPEG stream.
  458. Many MJPEG sources — go2rtc most notably (#1177), and several IP cameras —
  459. emit a "warm-up" frame on the byte that follows connection accept: usually
  460. the last keyframe held in the encoder, which is often black or stale until
  461. the encoder catches up to live content. To return a frame that's actually
  462. representative of the scene we read past the first frame and return the
  463. second; if the connection closes / times out / hits the buffer cap before
  464. a second frame ever arrives we fall back to the first so callers still
  465. get *something* (better than degrading slow / single-frame streams to None,
  466. which would regress every code path that consumed pre-fix behaviour).
  467. Note: this function intentionally makes requests to user-configured URLs.
  468. External camera support requires connecting to user-specified camera
  469. endpoints. URL is sanitized and dangerous destinations are blocked.
  470. """
  471. safe_url = _sanitize_camera_url(url, ("http", "https"))
  472. if not safe_url:
  473. logger.error("Invalid MJPEG URL format: %s...", redact_url_credentials(url)[:50])
  474. return None
  475. jpeg_start = b"\xff\xd8"
  476. jpeg_end = b"\xff\xd9"
  477. first_frame: bytes | None = None # warm-up frame; fallback if no second arrives
  478. buffer = b""
  479. try:
  480. async with (
  481. aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=timeout)) as session,
  482. session.get(safe_url) as response,
  483. ):
  484. if response.status != 200:
  485. logger.error("MJPEG stream returned status %s", response.status)
  486. return None
  487. async for chunk in response.content.iter_chunked(8192):
  488. buffer += chunk
  489. # A single chunk can carry multiple frames (e.g. high-FPS sources)
  490. # or a partial frame. Drain every complete frame we already have
  491. # before pulling the next chunk.
  492. while True:
  493. start_idx = buffer.find(jpeg_start)
  494. if start_idx == -1:
  495. # No frame start yet — drop trailing garbage, keep waiting.
  496. break
  497. end_idx = buffer.find(jpeg_end, start_idx + 2)
  498. if end_idx == -1:
  499. # Partial frame; trim already-discarded prefix so the
  500. # buffer stays bounded across long-running streams.
  501. if start_idx > 0:
  502. buffer = buffer[start_idx:]
  503. break
  504. frame = buffer[start_idx : end_idx + 2]
  505. buffer = buffer[end_idx + 2 :]
  506. if first_frame is None:
  507. first_frame = frame # warm-up; keep but don't return yet
  508. continue
  509. return frame # representative second frame
  510. if len(buffer) > 5 * 1024 * 1024: # 5MB limit
  511. logger.warning("MJPEG buffer exceeded 5MB without finding frame")
  512. break # exit chunk loop, fall through to first_frame fallback
  513. except TimeoutError:
  514. logger.warning("MJPEG frame capture timed out after %ss", timeout)
  515. except (aiohttp.ClientError, OSError) as e:
  516. logger.error("MJPEG frame capture failed: %s", e)
  517. # Stream ended / timed out / buffer cap before a second frame arrived.
  518. # Return whatever warm-up frame we managed to read; better an iffy frame
  519. # than None for callers that need *some* image (snapshot UX, plate-detect
  520. # CV, finish photo). None only if no frame ever arrived at all.
  521. return first_frame
  522. async def _capture_rtsp_frame(url: str, timeout: int) -> bytes | None:
  523. """Capture frame from RTSP using ffmpeg.
  524. For rtsps:// URLs, a local TLS proxy is used to avoid GnuTLS issues.
  525. Note: this function intentionally connects to user-configured URLs, the same
  526. as the MJPEG and snapshot paths. The URL is sanitized and dangerous
  527. destinations are blocked before it reaches ffmpeg.
  528. """
  529. ffmpeg = get_ffmpeg_path()
  530. if not ffmpeg:
  531. logger.error("ffmpeg not found - required for RTSP capture")
  532. return None
  533. # ffmpeg's -i accepts every protocol it was built with, so an unchecked URL
  534. # here is a request to any host and scheme the caller names, not merely to a
  535. # camera. Restricting the scheme to RTSP is what keeps this a camera fetch.
  536. safe_url = _sanitize_camera_url(url, ("rtsp", "rtsps"))
  537. if not safe_url:
  538. logger.error("Invalid RTSP URL: %s...", redact_url_credentials(url)[:50])
  539. return None
  540. # If rtsps://, use TLS proxy
  541. proxy_server = None
  542. effective_url = safe_url
  543. if safe_url.lower().startswith("rtsps://"):
  544. try:
  545. from urllib.parse import urlparse
  546. from backend.app.services.camera import close_tls_proxy, create_tls_proxy
  547. parsed = urlparse(safe_url)
  548. target_port = parsed.port or 322
  549. proxy_port, proxy_server = await create_tls_proxy(parsed.hostname, target_port)
  550. userinfo = ""
  551. if parsed.username:
  552. userinfo = parsed.username
  553. if parsed.password:
  554. userinfo += f":{parsed.password}"
  555. userinfo += "@"
  556. # Points at loopback deliberately, and is built after the check
  557. # above rather than re-checked: the destination that mattered was
  558. # the one the caller named, and it has already been vetted.
  559. effective_url = f"rtsp://{userinfo}127.0.0.1:{proxy_port}{parsed.path}"
  560. if parsed.query:
  561. effective_url += f"?{parsed.query}"
  562. except Exception as e:
  563. logger.warning("Failed to create TLS proxy for RTSP capture, falling back: %s", e)
  564. effective_url = safe_url
  565. cmd = [
  566. ffmpeg,
  567. "-rtsp_transport",
  568. "tcp",
  569. # Belt and braces on the scheme check above: a demuxer that follows a
  570. # reference out of the stream cannot leave these protocols either.
  571. "-protocol_whitelist",
  572. _RTSP_PROTOCOL_WHITELIST,
  573. "-i",
  574. effective_url,
  575. "-frames:v",
  576. "1",
  577. "-f",
  578. "image2pipe",
  579. "-vcodec",
  580. "mjpeg",
  581. "-q:v",
  582. "2",
  583. "-",
  584. ]
  585. try:
  586. logger.debug("Running ffmpeg RTSP capture...")
  587. process = await asyncio.create_subprocess_exec(
  588. *cmd,
  589. stdout=asyncio.subprocess.PIPE,
  590. stderr=asyncio.subprocess.PIPE,
  591. )
  592. stdout, stderr = await asyncio.wait_for(process.communicate(), timeout=timeout)
  593. logger.debug(
  594. "ffmpeg returned: code=%s, stdout=%s bytes, stderr=%s bytes",
  595. process.returncode,
  596. len(stdout),
  597. len(stderr),
  598. )
  599. if process.returncode != 0:
  600. # The summariser masks the camera password the input URL carries.
  601. logger.error("ffmpeg RTSP capture failed: %s", summarize_ffmpeg_stderr(stderr) or NO_FFMPEG_OUTPUT)
  602. return None
  603. if not stdout or len(stdout) < 100:
  604. logger.error("ffmpeg returned empty or too small frame")
  605. return None
  606. return stdout
  607. except TimeoutError:
  608. logger.warning("RTSP frame capture timed out after %ss", timeout)
  609. if process:
  610. process.kill()
  611. return None
  612. except OSError as e:
  613. logger.error("RTSP frame capture failed: %s", e)
  614. return None
  615. finally:
  616. if proxy_server:
  617. await close_tls_proxy(proxy_server)
  618. def _transcode_to_jpeg(data: bytes) -> bytes | None:
  619. """Decode an arbitrary still image (PNG/WebP/BMP/GIF/...) and re-encode as JPEG.
  620. Some camera/proxy snapshot endpoints serve stills as PNG or WebP rather than
  621. JPEG. A browser opened directly at the URL renders those fine, but our MJPEG
  622. ``multipart/x-mixed-replace`` stream hard-labels every part
  623. ``Content-Type: image/jpeg`` — so a non-JPEG payload makes the browser reject
  624. the frame and drop the whole stream ("connection lost", #1902). Transcoding to
  625. JPEG keeps the stream genuinely MJPEG and also keeps the JPEG-only downstream
  626. (plate detection, Obico, finish photo) working.
  627. Returns None if the bytes are not a decodable image (e.g. an HTML error page)
  628. or if the imaging libraries are unavailable — callers fall back to the raw
  629. bytes so behaviour is never worse than before.
  630. """
  631. try:
  632. import cv2
  633. import numpy as np
  634. except ImportError:
  635. return None
  636. try:
  637. img = cv2.imdecode(np.frombuffer(data, dtype=np.uint8), cv2.IMREAD_COLOR)
  638. if img is None:
  639. return None
  640. ok, buf = cv2.imencode(".jpg", img, [cv2.IMWRITE_JPEG_QUALITY, 85])
  641. if not ok:
  642. return None
  643. return buf.tobytes()
  644. except Exception as e: # cv2 raises cv2.error (a subclass of Exception) on bad input
  645. logger.debug("Snapshot transcode to JPEG failed: %s", e)
  646. return None
  647. async def _capture_snapshot(url: str, timeout: int) -> bytes | None:
  648. """Fetch snapshot from HTTP URL.
  649. Note: This function intentionally makes requests to user-configured URLs.
  650. External camera support requires connecting to user-specified camera endpoints.
  651. URL is sanitized and dangerous destinations are blocked.
  652. """
  653. # Sanitize URL - returns reconstructed URL from validated components
  654. safe_url = _sanitize_camera_url(url, ("http", "https"))
  655. if not safe_url:
  656. logger.error("Invalid snapshot URL format: %s...", redact_url_credentials(url)[:50])
  657. return None
  658. try:
  659. async with (
  660. aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=timeout)) as session,
  661. session.get(safe_url) as response,
  662. ):
  663. if response.status != 200:
  664. logger.error("Snapshot URL returned status %s", response.status)
  665. return None
  666. data = await response.read()
  667. except TimeoutError:
  668. logger.warning("Snapshot capture timed out after %ss", timeout)
  669. return None
  670. except (aiohttp.ClientError, OSError) as e:
  671. logger.error("Snapshot capture failed: %s", e)
  672. return None
  673. # Fast path: already JPEG (SOI marker), stream it as-is (no decode/re-encode).
  674. if data.startswith(b"\xff\xd8"):
  675. return data
  676. # Not JPEG. Many snapshot endpoints serve PNG/WebP/BMP — transcode to JPEG so
  677. # the browser's MJPEG stream (and JPEG-only downstream) keep working instead of
  678. # dropping the connection (#1902). Run off the event loop: cv2 decode/encode is
  679. # CPU-bound and this can be polled at up to 15 fps while a camera view is open.
  680. transcoded = await asyncio.to_thread(_transcode_to_jpeg, data)
  681. if transcoded is not None:
  682. logger.debug(
  683. "Transcoded non-JPEG snapshot (%d bytes, header %s) to JPEG",
  684. len(data),
  685. data[:4].hex(),
  686. )
  687. return transcoded
  688. # Couldn't decode it as an image at all — most likely not an image response
  689. # (HTML error page, auth redirect, wrong URL). Return the raw bytes as a last
  690. # resort (unchanged behaviour) but log enough to debug.
  691. logger.warning(
  692. "External camera snapshot is not a decodable image "
  693. "(%d bytes, header %s) — verify the camera URL returns an image",
  694. len(data),
  695. data[:4].hex(),
  696. )
  697. return data
  698. async def test_connection(url: str, camera_type: str) -> dict:
  699. """Test camera connection.
  700. Returns:
  701. Dict with {success: bool, error?: str, resolution?: str, coalesced: bool}
  702. ``coalesced`` is True when the frame came from a capture that was already
  703. running rather than from a connection this test opened. Captures are shared
  704. (see ``capture_frame``), so a test that lands while Obico is polling — or
  705. while any other one-shot consumer is mid-capture — gets that frame back and
  706. would otherwise report a healthy connection it never made, which is the one
  707. answer a *connection test* must not give silently. Forcing an uncoalesced
  708. capture here would be worse: it would open the second handle to a
  709. single-reader device that this whole mechanism exists to prevent. So the
  710. test still shares, and says so. Mirrors the ``coalesced_capture`` code the
  711. built-in diagnostic reports for the same situation (camera_diagnose.py).
  712. """
  713. logger.info("Testing camera connection: type=%s, url=%s...", camera_type, redact_url_credentials(url)[:50])
  714. # Sampled before the call, while it can still distinguish "someone else is
  715. # mid-capture" from "I am the one capturing".
  716. coalesced = capture_in_flight(url, camera_type)
  717. try:
  718. frame = await capture_frame(url, camera_type, timeout=10)
  719. logger.info("Capture result: %s bytes%s", len(frame) if frame else 0, " (coalesced)" if coalesced else "")
  720. if frame:
  721. # Try to get resolution from JPEG header
  722. resolution = None
  723. try:
  724. # Simple JPEG dimension extraction
  725. # SOF0 marker is FF C0, followed by length, precision, height, width
  726. sof_markers = [b"\xff\xc0", b"\xff\xc1", b"\xff\xc2"]
  727. for marker in sof_markers:
  728. idx = frame.find(marker)
  729. if idx != -1 and idx + 9 <= len(frame):
  730. height = (frame[idx + 5] << 8) | frame[idx + 6]
  731. width = (frame[idx + 7] << 8) | frame[idx + 8]
  732. resolution = f"{width}x{height}"
  733. break
  734. except (IndexError, ValueError):
  735. pass # Resolution detection is optional; fall back to default
  736. return {"success": True, "resolution": resolution, "coalesced": coalesced}
  737. else:
  738. return {"success": False, "error": "Failed to capture frame from camera", "coalesced": coalesced}
  739. except Exception as e:
  740. # Sanitize error message - don't expose internal details
  741. error_type = type(e).__name__
  742. logger.error("Camera connection test failed: %s", e)
  743. return {"success": False, "error": f"Connection failed: {error_type}", "coalesced": coalesced}
  744. async def generate_mjpeg_stream(
  745. url: str,
  746. camera_type: str,
  747. fps: int = 10,
  748. *,
  749. on_process: Callable[[asyncio.subprocess.Process], None] | None = None,
  750. on_frame: Callable[[bytes], None] | None = None,
  751. stop_event: asyncio.Event | None = None,
  752. ) -> AsyncGenerator[bytes, None]:
  753. """Generator yielding MJPEG frames for streaming.
  754. Args:
  755. url: Camera URL or USB device path
  756. camera_type: "mjpeg", "rtsp", "snapshot", or "usb"
  757. fps: Target frames per second
  758. on_process: Called with the spawned ffmpeg process for the ``usb`` and
  759. ``rtsp`` paths so the route layer can register it into the shared
  760. stream registries — that's what lets ``/camera/stop`` and the orphan
  761. janitor find and kill a leaked ffmpeg that's holding a USB device
  762. open (#2675). Without it the process is reachable only from this
  763. generator's own ``finally``, which an abrupt client disconnect can
  764. skip (same cancellation-timing class as #776).
  765. on_frame: Called with each RAW frame, before it is wrapped for the wire,
  766. so the route layer can publish it as the printer's buffered frame
  767. (#2707). It has to be a callback: what this generator yields is
  768. multipart-wrapped, so a consumer of the stream cannot recover the
  769. JPEG, and until now nothing populated the buffer for external
  770. cameras at all — leaving every one-shot consumer (layer timelapse,
  771. finish photo, Obico, plate check) with nothing to reuse and no
  772. option but to open a competing handle on a single-reader device.
  773. Exceptions are logged and swallowed: buffering must never be able
  774. to break the live stream.
  775. stop_event: When set, the reconnect loops stop retrying — so an explicit
  776. stop (which kills the current ffmpeg) doesn't immediately respawn a
  777. new process and reacquire the device.
  778. Yields:
  779. MJPEG frame data with HTTP multipart boundaries
  780. """
  781. frame_interval = 1.0 / max(fps, 1)
  782. last_frame_time = 0.0
  783. def _publish(frame: bytes) -> bytes:
  784. """Hand the raw frame to on_frame, then format it for the wire."""
  785. if on_frame is not None:
  786. try:
  787. on_frame(frame)
  788. except Exception:
  789. logger.exception("on_frame callback raised")
  790. return _format_mjpeg_frame(frame)
  791. async def _reconnecting(open_session, label: str):
  792. """Yield frames across sessions, reconnecting after each that delivered.
  793. A session that ends without a single frame stops the stream: the
  794. source is down, and retrying is the viewer's call. One that delivered
  795. frames and then ended is a routine drop and is always reconnected.
  796. This used to be three reconnects for the life of the stream, so a
  797. server that closes its sessions periodically ended the stream for good
  798. on the fourth drop, however long each session had run -- the external
  799. twin of the built-in RTSP path's lifetime reconnect budget.
  800. """
  801. drops = 0
  802. while True:
  803. frame_yielded = False
  804. # aclosing: stop the session's ffmpeg the moment this generator is
  805. # closed, rather than whenever the abandoned iterator is collected.
  806. async with contextlib.aclosing(open_session()) as session:
  807. async for frame in session:
  808. frame_yielded = True
  809. yield frame
  810. if not frame_yielded or (stop_event is not None and stop_event.is_set()):
  811. break
  812. # The sessions swallow CancelledError and simply end, so a viewer
  813. # whose task was cancelled mid-read looks like a routine drop here.
  814. # Never redial for a task that is being cancelled (Task.cancelling
  815. # is 3.11+; without it this falls back to the next await raising).
  816. task = asyncio.current_task()
  817. if task is not None and getattr(task, "cancelling", lambda: 0)():
  818. break
  819. drops += 1
  820. logger.warning("External %s stream ended, reconnecting (drop %d)...", label, drops)
  821. await asyncio.sleep(2)
  822. if camera_type == "mjpeg":
  823. # Proxy MJPEG stream directly, reconnecting after routine drops.
  824. async with contextlib.aclosing(_reconnecting(lambda: _stream_mjpeg(url), "MJPEG")) as frames:
  825. async for frame in frames:
  826. current_time = asyncio.get_event_loop().time()
  827. if current_time - last_frame_time >= frame_interval:
  828. last_frame_time = current_time
  829. yield _publish(frame)
  830. elif camera_type == "rtsp":
  831. # Use ffmpeg to convert RTSP to MJPEG, reconnecting after routine drops.
  832. async with contextlib.aclosing(
  833. _reconnecting(lambda: _stream_rtsp(url, fps, on_process=on_process), "RTSP")
  834. ) as frames:
  835. async for frame in frames:
  836. yield _publish(frame)
  837. elif camera_type == "usb":
  838. # Use ffmpeg to stream from USB camera
  839. async for frame in _stream_usb(url, fps, on_process=on_process):
  840. yield _publish(frame)
  841. elif camera_type == "snapshot":
  842. # Poll snapshot URL at interval
  843. while True:
  844. try:
  845. frame = await _capture_snapshot(url, timeout=10)
  846. if frame:
  847. yield _publish(frame)
  848. await asyncio.sleep(frame_interval)
  849. except asyncio.CancelledError:
  850. break
  851. except (aiohttp.ClientError, OSError) as e:
  852. logger.warning("Snapshot poll failed: %s", e)
  853. await asyncio.sleep(frame_interval)
  854. def _format_mjpeg_frame(frame: bytes) -> bytes:
  855. """Format frame for MJPEG HTTP response."""
  856. return (
  857. b"--frame\r\n"
  858. b"Content-Type: image/jpeg\r\n"
  859. b"Content-Length: " + str(len(frame)).encode() + b"\r\n"
  860. b"\r\n" + frame + b"\r\n"
  861. )
  862. async def _stream_mjpeg(url: str) -> AsyncGenerator[bytes, None]:
  863. """Stream frames from MJPEG URL.
  864. Note: This function intentionally makes requests to user-configured URLs.
  865. External camera support requires connecting to user-specified camera endpoints.
  866. URL is sanitized and dangerous destinations are blocked.
  867. """
  868. # Sanitize URL - returns reconstructed URL from validated components
  869. safe_url = _sanitize_camera_url(url, ("http", "https"))
  870. if not safe_url:
  871. logger.error("Invalid MJPEG stream URL: %s...", redact_url_credentials(url)[:50])
  872. return
  873. try:
  874. timeout = aiohttp.ClientTimeout(total=None, sock_read=30)
  875. async with aiohttp.ClientSession(timeout=timeout) as session, session.get(safe_url) as response:
  876. if response.status != 200:
  877. logger.error("MJPEG stream returned status %s", response.status)
  878. return
  879. buffer = b""
  880. jpeg_start = b"\xff\xd8"
  881. jpeg_end = b"\xff\xd9"
  882. async for chunk in response.content.iter_chunked(8192):
  883. buffer += chunk
  884. # Extract complete frames from buffer
  885. while True:
  886. start_idx = buffer.find(jpeg_start)
  887. if start_idx == -1:
  888. buffer = buffer[-2:] if len(buffer) > 2 else buffer
  889. break
  890. if start_idx > 0:
  891. buffer = buffer[start_idx:]
  892. end_idx = buffer.find(jpeg_end, 2)
  893. if end_idx == -1:
  894. break
  895. frame = buffer[: end_idx + 2]
  896. buffer = buffer[end_idx + 2 :]
  897. yield frame
  898. except asyncio.CancelledError:
  899. logger.info("MJPEG stream cancelled")
  900. except (aiohttp.ClientError, OSError) as e:
  901. logger.error("MJPEG stream error: %s", e)
  902. async def _stream_rtsp(
  903. url: str,
  904. fps: int,
  905. *,
  906. on_process: Callable[[asyncio.subprocess.Process], None] | None = None,
  907. ) -> AsyncGenerator[bytes, None]:
  908. """Stream frames from RTSP URL via ffmpeg.
  909. For rtsps:// URLs, a local TLS proxy (Python OpenSSL) is used instead
  910. of relying on ffmpeg's GnuTLS backend, which has compatibility issues
  911. with some printer firmwares.
  912. Note: this function intentionally connects to user-configured URLs. The URL
  913. is sanitized and dangerous destinations are blocked before it reaches
  914. ffmpeg — see ``_capture_rtsp_frame``, which guards the one-shot path the
  915. same way.
  916. """
  917. ffmpeg = get_ffmpeg_path()
  918. if not ffmpeg:
  919. logger.error("ffmpeg not found - required for RTSP streaming")
  920. return
  921. from backend.app.services.camera import rtsp_socket_timeout_flag
  922. safe_url = _sanitize_camera_url(url, ("rtsp", "rtsps"))
  923. if not safe_url:
  924. logger.error("Invalid RTSP stream URL: %s...", redact_url_credentials(url)[:50])
  925. return
  926. # If the URL uses rtsps://, set up a TLS proxy so ffmpeg uses plain rtsp://
  927. proxy_server = None
  928. effective_url = safe_url
  929. if safe_url.lower().startswith("rtsps://"):
  930. try:
  931. from urllib.parse import urlparse
  932. from backend.app.services.camera import close_tls_proxy, create_tls_proxy
  933. parsed = urlparse(safe_url)
  934. target_port = parsed.port or 322
  935. proxy_port, proxy_server = await create_tls_proxy(parsed.hostname, target_port)
  936. # Rewrite URL: rtsps://user:pass@host:port/path → rtsp://user:pass@127.0.0.1:proxy/path
  937. userinfo = ""
  938. if parsed.username:
  939. userinfo = parsed.username
  940. if parsed.password:
  941. userinfo += f":{parsed.password}"
  942. userinfo += "@"
  943. # Loopback by design, and built after the check above rather than
  944. # re-checked — see the same rewrite in _capture_rtsp_frame.
  945. effective_url = f"rtsp://{userinfo}127.0.0.1:{proxy_port}{parsed.path}"
  946. if parsed.query:
  947. effective_url += f"?{parsed.query}"
  948. except Exception as e:
  949. logger.warning("Failed to create TLS proxy for RTSP, falling back to direct: %s", e)
  950. effective_url = safe_url
  951. cmd = [
  952. ffmpeg,
  953. "-rtsp_transport",
  954. "tcp",
  955. "-rtsp_flags",
  956. "prefer_tcp",
  957. "-protocol_whitelist",
  958. _RTSP_PROTOCOL_WHITELIST,
  959. # Socket I/O timeout name varies by ffmpeg version (#1504); see
  960. # `rtsp_socket_timeout_flag()` in services.camera.
  961. f"-{rtsp_socket_timeout_flag()}",
  962. "30000000",
  963. "-buffer_size",
  964. "1024000",
  965. "-max_delay",
  966. "500000",
  967. # No probe cap here (#3082). The input is whatever camera the user
  968. # owns, so there is no stream to tune a fast-start probe against: a
  969. # 32-byte probe expires before a source that carries SPS/PPS in-band
  970. # rather than in its SDP has sent them, and ffmpeg then starts no
  971. # H.264 decoder and emits nothing at all. ffmpeg's defaults are a
  972. # ceiling rather than a wait, so a camera that announces itself in the
  973. # first packet still starts as fast as it ever did.
  974. #
  975. # `_capture_rtsp_frame` has always run on those defaults, which is how
  976. # a camera could pass the connection test and still show a black live
  977. # view. The printer path is the opposite case — a known Bambu camera
  978. # per model — and keeps its tuning in `camera_profiles.py`.
  979. "-fflags",
  980. "nobuffer",
  981. "-flags",
  982. "low_delay",
  983. "-i",
  984. effective_url,
  985. "-f",
  986. "mjpeg",
  987. "-q:v",
  988. "5",
  989. "-r",
  990. str(fps),
  991. "-an",
  992. "-",
  993. ]
  994. process = None
  995. try:
  996. process = await asyncio.create_subprocess_exec(
  997. *cmd,
  998. stdout=asyncio.subprocess.PIPE,
  999. stderr=asyncio.subprocess.PIPE,
  1000. )
  1001. # Register immediately — before the startup probe below — so a process
  1002. # that hangs on connect (rather than exiting) is still reachable by the
  1003. # stop endpoint / orphan janitor (#2675).
  1004. if on_process is not None:
  1005. on_process(process)
  1006. # Brief check for immediate startup failures
  1007. await asyncio.sleep(0.1)
  1008. if process.returncode is not None:
  1009. stderr = await process.stderr.read()
  1010. # The summariser masks the camera password the input URL carries.
  1011. logger.error(
  1012. "ffmpeg RTSP stream failed immediately: %s", summarize_ffmpeg_stderr(stderr) or NO_FFMPEG_OUTPUT
  1013. )
  1014. return
  1015. buffer = b""
  1016. jpeg_start = b"\xff\xd8"
  1017. jpeg_end = b"\xff\xd9"
  1018. while True:
  1019. try:
  1020. chunk = await asyncio.wait_for(process.stdout.read(8192), timeout=30.0)
  1021. if not chunk:
  1022. break
  1023. buffer += chunk
  1024. # Extract complete frames
  1025. while True:
  1026. start_idx = buffer.find(jpeg_start)
  1027. if start_idx == -1:
  1028. buffer = buffer[-2:] if len(buffer) > 2 else buffer
  1029. break
  1030. if start_idx > 0:
  1031. buffer = buffer[start_idx:]
  1032. end_idx = buffer.find(jpeg_end, 2)
  1033. if end_idx == -1:
  1034. break
  1035. frame = buffer[: end_idx + 2]
  1036. buffer = buffer[end_idx + 2 :]
  1037. yield frame
  1038. except TimeoutError:
  1039. logger.warning("RTSP stream read timeout")
  1040. break
  1041. except asyncio.CancelledError:
  1042. logger.info("RTSP stream cancelled")
  1043. except OSError as e:
  1044. logger.error("RTSP stream error: %s", e)
  1045. finally:
  1046. if process and process.returncode is None:
  1047. process.terminate()
  1048. try:
  1049. await asyncio.wait_for(process.wait(), timeout=2.0)
  1050. except TimeoutError:
  1051. process.kill()
  1052. await process.wait()
  1053. if proxy_server:
  1054. await close_tls_proxy(proxy_server)
  1055. async def _stream_usb(
  1056. device: str,
  1057. fps: int,
  1058. *,
  1059. on_process: Callable[[asyncio.subprocess.Process], None] | None = None,
  1060. ) -> AsyncGenerator[bytes, None]:
  1061. """Stream frames from USB camera via ffmpeg."""
  1062. ffmpeg = get_ffmpeg_path()
  1063. if not ffmpeg:
  1064. logger.error("ffmpeg not found - required for USB camera streaming")
  1065. return
  1066. # Same validation as the one-shot path: a prefix check accepted
  1067. # /dev/video/../../<anything that exists>, which -f v4l2 would then refuse
  1068. # rather than the check refusing it.
  1069. safe_device = _safe_usb_device_path(device)
  1070. if not safe_device:
  1071. return
  1072. device = safe_device
  1073. # ffmpeg command to stream from USB camera (v4l2)
  1074. cmd = [
  1075. ffmpeg,
  1076. "-f",
  1077. "v4l2",
  1078. "-framerate",
  1079. str(fps),
  1080. "-i",
  1081. device,
  1082. "-f",
  1083. "mjpeg",
  1084. "-q:v",
  1085. "5",
  1086. "-r",
  1087. str(fps),
  1088. "-",
  1089. ]
  1090. process = None
  1091. try:
  1092. logger.info("Starting USB camera stream from %s at %s fps", device, fps)
  1093. process = await asyncio.create_subprocess_exec(
  1094. *cmd,
  1095. stdout=asyncio.subprocess.PIPE,
  1096. stderr=asyncio.subprocess.PIPE,
  1097. )
  1098. # Register immediately — before the startup probe below — so a process
  1099. # that hangs in open()/ioctl on a still-locked device (rather than
  1100. # exiting with a "busy" error) is still reachable by the stop endpoint /
  1101. # orphan janitor (#2675).
  1102. if on_process is not None:
  1103. on_process(process)
  1104. # Give ffmpeg a moment to start and check for immediate failures
  1105. await asyncio.sleep(0.5)
  1106. if process.returncode is not None:
  1107. stderr = await process.stderr.read()
  1108. logger.error(
  1109. "ffmpeg USB stream failed immediately: %s", summarize_ffmpeg_stderr(stderr) or NO_FFMPEG_OUTPUT
  1110. )
  1111. return
  1112. buffer = b""
  1113. jpeg_start = b"\xff\xd8"
  1114. jpeg_end = b"\xff\xd9"
  1115. while True:
  1116. try:
  1117. chunk = await asyncio.wait_for(process.stdout.read(8192), timeout=30.0)
  1118. if not chunk:
  1119. break
  1120. buffer += chunk
  1121. # Extract complete frames
  1122. while True:
  1123. start_idx = buffer.find(jpeg_start)
  1124. if start_idx == -1:
  1125. buffer = buffer[-2:] if len(buffer) > 2 else buffer
  1126. break
  1127. if start_idx > 0:
  1128. buffer = buffer[start_idx:]
  1129. end_idx = buffer.find(jpeg_end, 2)
  1130. if end_idx == -1:
  1131. break
  1132. frame = buffer[: end_idx + 2]
  1133. buffer = buffer[end_idx + 2 :]
  1134. yield frame
  1135. except TimeoutError:
  1136. logger.warning("USB stream read timeout")
  1137. break
  1138. except asyncio.CancelledError:
  1139. logger.info("USB stream cancelled")
  1140. except OSError as e:
  1141. logger.error("USB stream error: %s", e)
  1142. finally:
  1143. if process and process.returncode is None:
  1144. process.terminate()
  1145. try:
  1146. await asyncio.wait_for(process.wait(), timeout=2.0)
  1147. except TimeoutError:
  1148. process.kill()
  1149. await process.wait()