obico_detection.py 23 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531
  1. """Obico AI print-failure detection service.
  2. Polls a self-hosted Obico ML API with snapshots from each monitored printer
  3. while a print is running, smooths scores over time, and dispatches a configured
  4. action (notify / pause / pause_and_off) when a sustained failure is detected.
  5. See `obico_smoothing.py` for the per-print EWM + rolling-mean math.
  6. """
  7. import asyncio
  8. import json
  9. import logging
  10. import secrets
  11. import time
  12. from collections import deque
  13. from datetime import datetime, timezone
  14. import httpx
  15. from sqlalchemy import select
  16. from backend.app.core.database import async_session
  17. from backend.app.models.printer import Printer
  18. from backend.app.models.settings import Settings
  19. from backend.app.services.obico_smoothing import (
  20. PrintState,
  21. classify,
  22. score_from_detections,
  23. thresholds,
  24. )
  25. logger = logging.getLogger(__name__)
  26. HISTORY_MAX = 50
  27. HEALTH_TIMEOUT = 5.0
  28. DETECTION_TIMEOUT = 30.0
  29. SNAPSHOT_CAPTURE_TIMEOUT = 20 # seconds — we control this, not Obico
  30. FRAME_CACHE_TTL = 30.0 # seconds — Obico usually fetches within 1s of receiving the URL
  31. # Module-level one-shot frame cache. Obico's ML API is GET-only (/p/?img=URL) and
  32. # fetches the URL itself with a hardcoded 5s read timeout. We capture locally first,
  33. # stash the JPEG under a random nonce, and hand Obico a URL that serves the cached
  34. # bytes instantly — so the 5s ceiling never races RTSP keyframe wait.
  35. _frame_cache: dict[str, tuple[bytes, float]] = {}
  36. _frame_cache_lock = asyncio.Lock()
  37. def auth_headers(token: str | None) -> dict[str, str]:
  38. """Bearer header for the ML API, or nothing when no token is configured.
  39. Obico's ML API gates ``/p/`` behind ``ML_API_TOKEN`` (``ml_api/auth.py``):
  40. with the variable set it answers a bare 401 to any request whose
  41. ``Authorization`` header isn't ``Bearer <token>``, and with it unset it
  42. ignores the header entirely. Sending nothing when unconfigured keeps the
  43. request byte-identical to what shipped before the setting existed.
  44. """
  45. token = (token or "").strip()
  46. return {"Authorization": f"Bearer {token}"} if token else {}
  47. def _prune_frame_cache() -> None:
  48. """Drop entries older than FRAME_CACHE_TTL. Called under the cache lock."""
  49. now = time.monotonic()
  50. expired = [k for k, (_b, ts) in _frame_cache.items() if now - ts > FRAME_CACHE_TTL]
  51. for k in expired:
  52. _frame_cache.pop(k, None)
  53. async def stash_frame(data: bytes) -> str:
  54. """Store JPEG bytes and return a URL-safe nonce that serves them once."""
  55. nonce = secrets.token_urlsafe(32)
  56. async with _frame_cache_lock:
  57. _prune_frame_cache()
  58. _frame_cache[nonce] = (data, time.monotonic())
  59. return nonce
  60. async def pop_frame(nonce: str) -> bytes | None:
  61. """Return and remove a cached frame by nonce; None if missing or expired."""
  62. async with _frame_cache_lock:
  63. _prune_frame_cache()
  64. entry = _frame_cache.pop(nonce, None)
  65. if entry is None:
  66. return None
  67. data, ts = entry
  68. if time.monotonic() - ts > FRAME_CACHE_TTL:
  69. return None
  70. return data
  71. class ObicoDetectionService:
  72. """Singleton service that polls the ML API and acts on sustained failures."""
  73. def __init__(self):
  74. self._task: asyncio.Task | None = None
  75. # printer_id -> PrintState (reset when a new print starts)
  76. self._states: dict[int, PrintState] = {}
  77. # printer_id -> task_name active when state was created (used to detect new prints)
  78. self._state_keys: dict[int, str] = {}
  79. # printer_id -> last classification ("safe"/"warning"/"failure").
  80. # Only written after an inference actually came back, so a missing entry
  81. # means "we have no verdict", which is not the same as "safe" (#2952).
  82. self._last_class: dict[int, str] = {}
  83. # printer_id -> why the most recent poll produced no verdict, or absent
  84. # when the last poll succeeded. Per-printer rather than global so a card
  85. # can say what went wrong for *that* printer.
  86. self._errors: dict[int, str] = {}
  87. # printer_id -> whether an action has already been fired for the current print
  88. self._action_fired: dict[int, bool] = {}
  89. # Global detection event log (most-recent-first)
  90. self._history: deque = deque(maxlen=HISTORY_MAX)
  91. self._last_error: str | None = None
  92. # ---- lifecycle ----
  93. async def start(self):
  94. if self._task is not None:
  95. return
  96. logger.info("Starting Obico detection service")
  97. self._task = asyncio.create_task(self._loop())
  98. def stop(self):
  99. if self._task:
  100. self._task.cancel()
  101. self._task = None
  102. logger.info("Stopped Obico detection service")
  103. # ---- settings ----
  104. async def _load_settings(self) -> dict:
  105. keys = [
  106. "obico_enabled",
  107. "obico_ml_url",
  108. "bambuddy_internal_url",
  109. "obico_ml_token",
  110. "obico_sensitivity",
  111. "obico_action",
  112. "obico_poll_interval",
  113. "obico_enabled_printers",
  114. "external_url",
  115. ]
  116. async with async_session() as db:
  117. result = await db.execute(select(Settings).where(Settings.key.in_(keys)))
  118. rows = {r.key: r.value for r in result.scalars().all()}
  119. enabled_printers_raw = rows.get("obico_enabled_printers", "")
  120. if enabled_printers_raw:
  121. try:
  122. enabled_printers = set(json.loads(enabled_printers_raw))
  123. except json.JSONDecodeError:
  124. enabled_printers = set()
  125. else:
  126. enabled_printers = None # None = all printers
  127. return {
  128. "enabled": rows.get("obico_enabled", "false").lower() == "true",
  129. "ml_url": (rows.get("obico_ml_url") or "").rstrip("/"),
  130. "ml_token": (rows.get("obico_ml_token") or "").strip(),
  131. "sensitivity": rows.get("obico_sensitivity", "medium"),
  132. "action": rows.get("obico_action", "notify"),
  133. "poll_interval": int(rows.get("obico_poll_interval", "10")),
  134. "enabled_printers": enabled_printers,
  135. # Where Obico's ML server fetches snapshots: the Internal URL when
  136. # set, otherwise the public External URL.
  137. "snapshot_base_url": (
  138. (rows.get("bambuddy_internal_url") or "").strip() or (rows.get("external_url") or "").strip()
  139. ).rstrip("/"),
  140. }
  141. # ---- main loop ----
  142. async def _loop(self):
  143. """Poll active printers while enabled. Adjusts interval from settings each cycle."""
  144. while True:
  145. try:
  146. settings = await self._load_settings()
  147. interval = max(5, settings.get("poll_interval", 10))
  148. if not settings["enabled"] or not settings["ml_url"]:
  149. await asyncio.sleep(interval)
  150. continue
  151. await self._poll_once(settings)
  152. await asyncio.sleep(interval)
  153. except asyncio.CancelledError:
  154. break
  155. except Exception as e:
  156. logger.error("Obico detection loop error: %s", e)
  157. self._last_error = str(e) or type(e).__name__
  158. await asyncio.sleep(30)
  159. async def _poll_once(self, settings: dict):
  160. # Late import to avoid cycles at module load time
  161. from backend.app.services.printer_manager import printer_manager
  162. statuses = printer_manager.get_all_statuses()
  163. for printer_id, status in list(statuses.items()):
  164. if settings["enabled_printers"] is not None and printer_id not in settings["enabled_printers"]:
  165. continue
  166. if not printer_manager.is_connected(printer_id):
  167. continue
  168. if not status or getattr(status, "state", None) != "RUNNING":
  169. # Reset state when not printing so the next print starts fresh
  170. self._states.pop(printer_id, None)
  171. self._state_keys.pop(printer_id, None)
  172. self._action_fired.pop(printer_id, None)
  173. self._last_class.pop(printer_id, None)
  174. self._errors.pop(printer_id, None)
  175. continue
  176. await self._check_printer(printer_id, status, settings)
  177. async def _capture_frame(self, printer_id: int) -> bytes | None:
  178. """Capture one JPEG frame from the printer camera. Returns None on failure."""
  179. # Late import to avoid cycles at module load time
  180. from backend.app.services.camera import capture_camera_frame_bytes
  181. from backend.app.services.external_camera import capture_frame as capture_external_frame
  182. async with async_session() as db:
  183. printer = await db.get(Printer, printer_id)
  184. if printer is None:
  185. self._last_error = f"Printer {printer_id} not found"
  186. return None
  187. if printer.external_camera_enabled and printer.external_camera_url:
  188. # Same rule as the built-in branch below, which this used to skip:
  189. # an external camera is single-reader too, so polling while a viewer
  190. # is attached just fails (#2707).
  191. from backend.app.api.routes.camera import live_frame_for_capture
  192. defer, buffered = live_frame_for_capture(printer_id)
  193. if defer:
  194. if buffered:
  195. return buffered
  196. logger.info(
  197. "Obico: viewer attached for printer %s but buffer empty; "
  198. "skipping this poll to avoid competing camera handle (#2707)",
  199. printer_id,
  200. )
  201. return None
  202. return await capture_external_frame(
  203. printer.external_camera_url,
  204. printer.external_camera_type,
  205. timeout=SNAPSHOT_CAPTURE_TIMEOUT,
  206. snapshot_url=printer.external_camera_snapshot_url,
  207. )
  208. # Reuse the fan-out broadcaster's buffered frame when a viewer is
  209. # already watching — avoids opening a second concurrent RTSP socket
  210. # on printers that allow only one camera connection (e.g. X2D
  211. # firmware 01.01.00.00; see #1271). Buffered frame is <1s old while
  212. # a viewer is connected.
  213. #
  214. # When a viewer is attached but no frame is buffered yet (startup
  215. # race, mid-reconnect), we DELIBERATELY skip this poll cycle instead
  216. # of falling through to capture_camera_frame_bytes. Opening a fresh
  217. # RTSP/chamber socket would compete with the live viewer and kick
  218. # the fan-out connection on most firmwares — exactly the freeze
  219. # reported in #1348. The poll loop retries in ~10s.
  220. from backend.app.api.routes.camera import is_stream_active, try_get_active_buffered_frame
  221. if is_stream_active(printer_id):
  222. buffered = try_get_active_buffered_frame(printer_id)
  223. if buffered:
  224. return buffered
  225. logger.info(
  226. "Obico: viewer attached for printer %s but buffer empty; skipping this poll to avoid competing camera socket (#1348)",
  227. printer_id,
  228. )
  229. return None
  230. return await capture_camera_frame_bytes(
  231. ip_address=printer.ip_address,
  232. access_code=printer.access_code,
  233. model=printer.model,
  234. timeout=SNAPSHOT_CAPTURE_TIMEOUT,
  235. )
  236. def _no_verdict(self, printer_id: int, reason: str) -> None:
  237. """Record that this poll produced no verdict for ``printer_id``.
  238. Kept separate from the classification so the status surface can say
  239. "not checking" instead of inheriting the previous verdict — or, worse,
  240. the default "safe" a printer used to get before its first inference.
  241. """
  242. self._errors[printer_id] = reason
  243. self._last_error = reason
  244. logger.warning(reason)
  245. async def _check_printer(self, printer_id: int, status, settings: dict):
  246. task_name = getattr(status, "task_name", None) or getattr(status, "subtask_name", "") or ""
  247. key = f"{task_name}"
  248. if self._state_keys.get(printer_id) != key:
  249. self._states[printer_id] = PrintState()
  250. self._state_keys[printer_id] = key
  251. self._action_fired[printer_id] = False
  252. # Capture locally first, then hand Obico a nonce URL that returns the
  253. # cached bytes instantly. Obico's ML API is GET-only (/p/?img=URL) with a
  254. # hardcoded 5s read timeout which would otherwise race our /camera/snapshot
  255. # keyframe wait.
  256. frame = await self._capture_frame(printer_id)
  257. if not frame:
  258. self._no_verdict(printer_id, f"Failed to capture snapshot for printer {printer_id}")
  259. return
  260. snapshot_base_url = settings.get("snapshot_base_url") or ""
  261. if not snapshot_base_url:
  262. self._no_verdict(
  263. printer_id,
  264. "bambuddy_internal_url and external_url settings are empty — Obico's ML API needs a reachable URL to fetch the snapshot from. "
  265. "Set Settings → Failure Detection → Bambuddy Internal URL or Settings → Network → External URL.",
  266. )
  267. return
  268. nonce = await stash_frame(frame)
  269. snapshot_url = f"{snapshot_base_url}/api/v1/obico/cached-frame/{nonce}"
  270. ml_url = f"{settings['ml_url']}/p/"
  271. try:
  272. async with httpx.AsyncClient(timeout=DETECTION_TIMEOUT) as client:
  273. resp = await client.get(
  274. ml_url,
  275. params={"img": snapshot_url},
  276. headers=auth_headers(settings.get("ml_token")),
  277. )
  278. if resp.status_code == 401:
  279. # The server runs with ML_API_TOKEN set and rejected ours.
  280. # Say so plainly: the health endpoint is ungated, so "Test
  281. # Connection" passes against exactly this configuration and
  282. # a raw 401 gives the user nothing to act on (#2733).
  283. #
  284. # Obico's auth decorator runs before the handler, so a call
  285. # rejected here leaves no trace in the ML API's own log —
  286. # which is how #2952 came to be reported as "the loop never
  287. # calls the ML API" while it was calling it every 10s.
  288. self._no_verdict(
  289. printer_id,
  290. "Obico ML API rejected the token (401). Set Settings → Failure Detection → "
  291. "ML API Token to the ML_API_TOKEN the server runs with, or clear ML_API_TOKEN "
  292. "on the server.",
  293. )
  294. return
  295. resp.raise_for_status()
  296. payload = resp.json()
  297. except Exception as e:
  298. detail = str(e) or type(e).__name__
  299. self._no_verdict(printer_id, f"ML API call failed for printer {printer_id}: {detail}")
  300. return
  301. detections = payload.get("detections", []) if isinstance(payload, dict) else []
  302. current_p = score_from_detections(detections)
  303. state = self._states[printer_id]
  304. score = state.update(current_p)
  305. verdict = classify(score, settings["sensitivity"])
  306. self._last_class[printer_id] = verdict
  307. # A successful capture + ML call clears any transient error from previous
  308. # polls (typical case: cold-start RTSP timeout on first frame after startup,
  309. # followed by healthy polls that otherwise leave the banner stuck in the UI).
  310. self._errors.pop(printer_id, None)
  311. self._last_error = None
  312. # Log every non-safe sample — safe samples would flood history
  313. if verdict != "safe" or detections:
  314. self._history.appendleft(
  315. {
  316. "printer_id": printer_id,
  317. "task_name": task_name,
  318. "timestamp": datetime.now(timezone.utc).isoformat(),
  319. "current_p": round(current_p, 4),
  320. "score": round(score, 4),
  321. "class": verdict,
  322. "detections": len(detections),
  323. }
  324. )
  325. if verdict == "failure" and not self._action_fired.get(printer_id):
  326. self._action_fired[printer_id] = True
  327. await self._dispatch_action(printer_id, settings["action"], task_name, score, frame)
  328. async def _dispatch_action(
  329. self, printer_id: int, action: str, task_name: str, score: float, frame: bytes | None = None
  330. ):
  331. from backend.app.services.obico_actions import execute_action
  332. logger.warning(
  333. "Obico: failure detected on printer %s (task=%r score=%.3f) — action=%s",
  334. printer_id,
  335. task_name,
  336. score,
  337. action,
  338. )
  339. try:
  340. # Same frame the ML model flagged, not a fresh capture — the
  341. # printer's moved on by the time this fires.
  342. await execute_action(printer_id, action, task_name, score, frame)
  343. except Exception as e:
  344. self._last_error = f"Action dispatch failed: {e or type(e).__name__}"
  345. logger.error(self._last_error)
  346. # ---- queries ----
  347. def get_per_printer(self) -> dict:
  348. """Live classification per actively monitored printer.
  349. Only printers with a running, monitored print have a state entry, so
  350. consumers get "show nothing" for idle printers for free.
  351. Four classes, and the two non-verdict ones matter as much as the rest:
  352. ``error`` the most recent poll produced no verdict. ``error`` carries
  353. the reason — a rejected token, an unreachable ML API, a
  354. camera that would not yield a frame, a missing Bambuddy address.
  355. ``unknown`` monitored, but no inference has come back yet. The state
  356. entry is created when the print is first seen, which is
  357. before the first capture, so this is the honest answer for
  358. that window.
  359. ``safe`` / ``warning`` / ``failure``
  360. an actual verdict from an actual inference.
  361. This used to default to ``safe`` whenever no verdict had been recorded,
  362. so a printer whose detection had never once succeeded rendered exactly
  363. like a healthy one: a green badge reading "Safe" at score 0.000. That is
  364. the worst possible failure mode for a safety feature — it asserts the
  365. print is being watched precisely when it is not (#2952).
  366. """
  367. result = {}
  368. for pid, state in self._states.items():
  369. error = self._errors.get(pid)
  370. if error:
  371. verdict = "error"
  372. else:
  373. verdict = self._last_class.get(pid) or "unknown"
  374. result[pid] = {
  375. "class": verdict,
  376. "frame_count": state.frame_count,
  377. "score": round(state.ewm_mean, 4),
  378. "error": error,
  379. }
  380. return result
  381. def get_status(self, sensitivity: str = "medium") -> dict:
  382. # Report the thresholds for the configured sensitivity, not a hardcoded
  383. # "medium" — otherwise the Status panel always shows the medium row
  384. # regardless of the user's selection (#1469). thresholds() falls back
  385. # to the medium multiplier for any unrecognized value.
  386. low, high = thresholds(sensitivity)
  387. return {
  388. "is_running": self._task is not None and not self._task.done(),
  389. "last_error": self._last_error,
  390. "per_printer": self.get_per_printer(),
  391. "thresholds": {"low": low, "high": high},
  392. "history": list(self._history),
  393. }
  394. async def test_connection(self, url: str, token: str = "") -> dict:
  395. """Ping the ML API and check the token. Returns {ok, status_code, body, error, auth_ok}.
  396. The stored ``obico_ml_url`` setting is validated at the schema layer,
  397. but this route takes its URL from the request body, so the same
  398. LAN-service policy has to be applied here or the guard is trivially
  399. sidestepped by testing a URL instead of saving it. The response body
  400. is returned to the caller (it is the health signal — the endpoint
  401. answers "ok"), which is exactly why the destination must be inside
  402. policy before the request is made.
  403. ``token`` is used verbatim — resolving "not supplied" to the saved
  404. setting is the route's job, so this stays a pure outbound call.
  405. Health alone cannot answer whether the token works, because Obico
  406. gates ``/p/`` but leaves ``/hc/`` open — which is how a token-protected
  407. server passed this test while every detection call came back 401
  408. (#2733). So a second, side-effect-free probe follows: ``/p/`` with no
  409. ``img`` parameter. The auth decorator runs before the handler, so 401
  410. means the token was rejected and 422 ("Invalid request params") means
  411. it was accepted. No inference work is done either way.
  412. """
  413. from backend.app.api.routes._url_safety import assert_safe_lan_service_url
  414. try:
  415. assert_safe_lan_service_url(url, label="Obico ML URL")
  416. except ValueError as exc:
  417. return {"ok": False, "status_code": None, "body": None, "error": str(exc), "auth_ok": None}
  418. headers = auth_headers(token)
  419. base = url.rstrip("/")
  420. try:
  421. async with httpx.AsyncClient(timeout=HEALTH_TIMEOUT) as client:
  422. resp = await client.get(f"{base}/hc/", headers=headers)
  423. body = resp.text.strip()
  424. healthy = resp.status_code == 200 and body.lower() == "ok"
  425. if not healthy:
  426. return {
  427. "ok": False,
  428. "status_code": resp.status_code,
  429. "body": body,
  430. "error": None,
  431. "auth_ok": None,
  432. }
  433. auth_ok: bool | None
  434. try:
  435. probe = await client.get(f"{base}/p/", headers=headers)
  436. auth_ok = probe.status_code != 401
  437. except Exception:
  438. # The health check already succeeded, so don't fail the
  439. # whole test on the probe — report the token as unknown.
  440. auth_ok = None
  441. except Exception as e:
  442. return {
  443. "ok": False,
  444. "status_code": None,
  445. "body": None,
  446. "error": str(e) or type(e).__name__,
  447. "auth_ok": None,
  448. }
  449. if auth_ok is False:
  450. return {
  451. "ok": False,
  452. "status_code": 401,
  453. "body": body,
  454. "error": (
  455. "The ML API is reachable but rejected the token. It runs with ML_API_TOKEN set — "
  456. "enter that value as the ML API Token, or clear ML_API_TOKEN on the server."
  457. ),
  458. "auth_ok": False,
  459. }
  460. return {"ok": True, "status_code": resp.status_code, "body": body, "error": None, "auth_ok": auth_ok}
  461. obico_detection_service = ObicoDetectionService()