notify_live_activities.py 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752
  1. """Persistent Notify! Live Activities, independent of push and digest selection.
  2. Only the worker performs network I/O. MQTT callbacks copy their mutable state and
  3. wake it; database sessions are closed before every HTTP call. A durable intent is
  4. written *before* starting a tile: after a crash or an ambiguous network failure we
  5. never blindly repeat a push-to-start and consume another of the five device slots.
  6. """
  7. import asyncio
  8. import hashlib
  9. import json
  10. import logging
  11. from dataclasses import dataclass, replace
  12. from datetime import datetime, timedelta, timezone
  13. from uuid import uuid4
  14. import httpx
  15. from sqlalchemy import delete, select, update
  16. from backend.app.core.database import async_session
  17. from backend.app.models.notification import NotificationProvider
  18. from backend.app.models.notification_live_activity import NotificationLiveActivity
  19. from backend.app.models.printer import Printer
  20. from backend.app.services.notify_client import NotifyClient, NotifyError, _https_url
  21. from backend.app.services.notify_worker import NotifyWorkerLifecycle
  22. logger = logging.getLogger(__name__)
  23. _RUNNING = {"RUNNING", "PRINTING", "PAUSE", "PREPARE", "SLICING"}
  24. _TERMINAL = {"IDLE", "FINISH", "FAILED", "COMPLETED", "CANCELLED", "ABORTED", "STOPPED"}
  25. _STOPPED = {"ended", "suppressed"}
  26. _ROLLOVER = {"expired", "overdue", "abandoned"}
  27. _INTERVAL = 60
  28. _STARTUP_GRACE = timedelta(minutes=2)
  29. def _now() -> datetime:
  30. return datetime.now(timezone.utc).replace(tzinfo=None)
  31. def _date(value) -> datetime | None:
  32. if not value:
  33. return None
  34. try:
  35. return datetime.fromisoformat(value.replace("Z", "+00:00")).astimezone(timezone.utc).replace(tzinfo=None)
  36. except (ValueError, AttributeError):
  37. return None
  38. def _config(provider) -> dict:
  39. return json.loads(provider.config) if isinstance(provider.config, str) else provider.config or {}
  40. def _credential_key(config: dict) -> str:
  41. return hashlib.sha256(f"{config.get('device_id', '')}:{config.get('token', '')}".encode()).hexdigest()
  42. def _identity(state) -> str:
  43. job_id = str(getattr(state, "subtask_id", "") or "")
  44. if job_id and job_id != "0":
  45. return f"job:{job_id}"[:160]
  46. # raw_data contains MQTT deltas, so gcode_start_time may disappear on the
  47. # next push. The stable job ID is preferred; the filename fallback is given
  48. # a fresh generation by real start callbacks and adopted on restart.
  49. filename = getattr(state, "subtask_name", None) or getattr(state, "current_print", None) or ""
  50. digest = hashlib.sha256(str(filename).encode()).hexdigest()[:32]
  51. return f"file:{digest}"
  52. def _printer_fault(state) -> str | None:
  53. """Match the printer card's actionable HMS filtering, including print_error.
  54. Bambu's MQTT parser folds print_error into hms_errors. Level-3 sixteen-digit
  55. HMS notices without actions (for example an open cover) are informational;
  56. eight-digit task-stopping print_error prompts at that level still count.
  57. """
  58. for error in getattr(state, "hms_errors", None) or []:
  59. severity = getattr(error, "severity", 0)
  60. description = str(getattr(error, "description", "") or "")
  61. actions = bool(getattr(error, "actions", None))
  62. notice = severity == 3 and len(str(getattr(error, "full_code", "") or "")) == 16
  63. if severity >= 1 and (actions or (description and not notice)):
  64. code = getattr(error, "full_code", "") or getattr(error, "code", "")
  65. return (description or f"Printer needs attention ({code})").replace("\x00", "")[:170]
  66. return None
  67. @dataclass(frozen=True)
  68. class PrintSnapshot:
  69. key: str
  70. connected: bool
  71. state: str
  72. filename: str
  73. progress: float
  74. remaining_seconds: int
  75. layer: int
  76. layers: int
  77. job_key: str | None = None
  78. fault: str | None = None
  79. temperatures: tuple = ()
  80. stage: int = -1
  81. @classmethod
  82. def from_state(cls, state):
  83. return cls(
  84. key=_identity(state),
  85. job_key=_identity(state) if _identity(state).startswith("job:") else None,
  86. connected=bool(state.connected),
  87. state=str(state.state or "").upper(),
  88. filename=str(getattr(state, "subtask_name", None) or getattr(state, "current_print", None) or "Print"),
  89. progress=max(0, min(100, float(state.progress or 0))),
  90. remaining_seconds=max(0, int(getattr(state, "remaining_time", 0) or 0) * 60),
  91. layer=int(state.layer_num or 0),
  92. layers=int(state.total_layers or 0),
  93. temperatures=tuple((getattr(state, "temperatures", {}) or {}).items()),
  94. stage=getattr(state, "stg_cur", -1),
  95. fault=_printer_fault(state),
  96. )
  97. def content(self, printer_name: str, config: dict | None = None) -> dict:
  98. config = config or {}
  99. phase = "Printing"
  100. if not self.connected:
  101. phase = "Printer offline"
  102. elif self.state == "PAUSE":
  103. phase = "Paused"
  104. elif self.state in ("PREPARE", "SLICING") or (
  105. self.state in {"RUNNING", "PRINTING"}
  106. and self.layer < 1
  107. and (self.layers > 0 or self.stage not in {0, -1, 255})
  108. ):
  109. phase = "Preparing"
  110. elif self.state in _TERMINAL:
  111. phase = {
  112. "FINISH": "Complete",
  113. "COMPLETED": "Complete",
  114. "FAILED": "Failed",
  115. "CANCELLED": "Stopped",
  116. "ABORTED": "Stopped",
  117. "STOPPED": "Stopped",
  118. "IDLE": "Finished",
  119. }[self.state]
  120. if self.connected and self.fault and self.state in _RUNNING:
  121. lower = self.fault.lower()
  122. runout = "filament" in lower and any(
  123. word in lower for word in ("run out", "ran out", "runout", "exhaust", "empty")
  124. )
  125. phase = "Filament runout" if runout else "Printer error"
  126. # Character and UTF-8 limits leave ample room for server-owned fields
  127. # within Notify's 2048-byte merged-content budget, including emoji names.
  128. filename = (
  129. "Print in progress"
  130. if config.get("live_activity_privacy") is True
  131. else self.filename.replace("\x00", "")[:120]
  132. )
  133. layer = f" · Layer {self.layer}/{self.layers}" if self.layers else ""
  134. progress = 0 if phase == "Preparing" else self.progress
  135. content = {
  136. "title": printer_name.replace("\x00", "")[:80] or "Bambuddy",
  137. "body": self.fault
  138. if self.connected and self.fault and self.state in _RUNNING
  139. else f"{filename}{layer}"[:170],
  140. "symbol": str(config.get("live_activity_symbol") or "printer.fill")[:64],
  141. "tint": config.get("live_activity_tint") or "#00AE42",
  142. "progress": 100 if phase == "Complete" else progress,
  143. "status": phase,
  144. "endsIn": self.remaining_seconds
  145. if self.connected
  146. and not self.fault
  147. and self.state in {"RUNNING", "PRINTING"}
  148. and 0 < self.remaining_seconds <= 86400
  149. else None,
  150. }
  151. style = config.get("live_activity_style", "bar")
  152. content.update(
  153. steps=10 if style == "segments" else None, step=int(progress // 10) if style == "segments" else None
  154. )
  155. if style == "none":
  156. content["progress"] = None
  157. if (
  158. config.get("live_activity_stage") is True
  159. and self.connected
  160. and not self.fault
  161. and self.state in {"RUNNING", "PRINTING"}
  162. ):
  163. from backend.app.services.bambu_mqtt import get_stage_name
  164. content["status"] = get_stage_name(self.stage)[:40] if self.stage >= 0 else phase
  165. eta = f"{self.remaining_seconds // 3600}h {(self.remaining_seconds // 60) % 60}m"
  166. content["trailing"] = eta if self.remaining_seconds > 86400 and phase == "Printing" else None
  167. # The default tile uses the native countdown. Metric chips are opt-in;
  168. # send null when disabled so merge-patch clears previously selected chips.
  169. metrics = []
  170. temperatures = dict(self.temperatures)
  171. for field in (config.get("live_activity_metrics") or [])[:6]:
  172. label, value, unit = "", None, ""
  173. if field == "progress":
  174. label, value = "Progress", f"{progress:.0f}%"
  175. elif field == "eta" and phase == "Printing" and self.remaining_seconds > 0:
  176. label, value = "Remaining", eta
  177. elif field == "layers" and self.layers:
  178. label, value = "Layer", f"{self.layer}/{self.layers}"
  179. elif field in {"nozzle", "bed", "chamber"} and temperatures.get(field) is not None:
  180. label, value, unit = field.title(), f"{temperatures[field]:.0f}", "°C"
  181. if value is not None:
  182. metrics.append({"label": label, "value": value[:16], "unit": unit})
  183. content["metrics"] = None if self.fault else metrics or None
  184. button = _https_url(config.get("live_activity_button_url"))
  185. content["button"] = (
  186. {"title": "Open Bambuddy", "url": button, "open": True} if button and len(button) <= 512 else None
  187. )
  188. # Notify adds timestamps when merging. Reserve 200 bytes for those;
  189. # shorten display text first, then drop optional cells/button if needed.
  190. def size():
  191. return len(json.dumps(content, ensure_ascii=False, separators=(",", ":")).encode())
  192. while size() > 1800 and len(content["body"]) > 20:
  193. content["body"] = content["body"][:-10]
  194. if size() > 1800:
  195. content["button"] = None
  196. while size() > 1800 and content["metrics"]:
  197. content["metrics"].pop()
  198. return content
  199. class NotifyLiveActivityService(NotifyWorkerLifecycle):
  200. ownership_model = NotificationLiveActivity
  201. feature_field = "live_activities"
  202. worker_interval = 30
  203. worker_name = "notify-live-activities"
  204. stopped_states = ("ended", "suppressed")
  205. _cleanup_key = staticmethod(_credential_key)
  206. def __init__(self, session_factory=None, client=None):
  207. self._session = session_factory or async_session
  208. self._client = client
  209. self._http: httpx.AsyncClient | None = None
  210. self._snapshots: dict[int, PrintSnapshot] = {}
  211. self._new_fallbacks: dict[int, str] = {}
  212. self._started_at: dict[int, datetime] = {}
  213. self._progress_resets: dict[int, tuple[float, datetime]] = {}
  214. self._finished: list[tuple[int, str | None, str]] = []
  215. self._lock = asyncio.Lock()
  216. self._wake = asyncio.Event()
  217. self._task: asyncio.Task | None = None
  218. self._init_worker()
  219. self._status_grace_until: datetime | None = None
  220. def _on_activated(self):
  221. # No MQTT observations are accepted while dormant. Anything retained
  222. # from before opt-out is stale until a real status arrives after opt-in.
  223. self._snapshots.clear()
  224. self._finished.clear()
  225. self._new_fallbacks.clear()
  226. self._started_at.clear()
  227. self._progress_resets.clear()
  228. self._status_grace_until = None
  229. def observe(self, printer_id: int, state) -> None:
  230. if self._provider_enabled is False:
  231. return
  232. snapshot = PrintSnapshot.from_state(state)
  233. if snapshot.state in {"", "UNKNOWN"}:
  234. previous = self._snapshots.get(printer_id)
  235. if previous is None or snapshot.connected:
  236. return # Broker connection alone is not a real printer status.
  237. snapshot = replace(previous, connected=False)
  238. reset = self._progress_resets.get(printer_id)
  239. if reset:
  240. initial, until = reset
  241. if (
  242. snapshot.state not in _RUNNING
  243. or snapshot.progress < initial
  244. or snapshot.progress <= 5
  245. or _now() >= until
  246. ):
  247. self._progress_resets.pop(printer_id, None)
  248. else:
  249. snapshot = replace(snapshot, progress=0, remaining_seconds=0)
  250. previous = self._snapshots.get(printer_id)
  251. if (
  252. previous
  253. and previous.state in _RUNNING
  254. and snapshot.state in _RUNNING
  255. and previous.filename == snapshot.filename
  256. ):
  257. # A delta may omit or belatedly supply a job ID. Retain the print's
  258. # established identity until an explicit start or a different ID.
  259. if snapshot.key.startswith("file:") or previous.key.startswith("file:"):
  260. snapshot = replace(snapshot, key=previous.key)
  261. self._snapshots[printer_id] = snapshot
  262. if self._provider_enabled and (
  263. previous is None
  264. or (previous.key, previous.state, previous.connected, previous.fault)
  265. != (
  266. snapshot.key,
  267. snapshot.state,
  268. snapshot.connected,
  269. snapshot.fault,
  270. )
  271. ):
  272. self._wake.set()
  273. def print_started(self, printer_id: int, state, data: dict | None = None) -> None:
  274. if state is None or self._provider_enabled is False:
  275. return
  276. snapshot = PrintSnapshot.from_state(state)
  277. job_id = (data or {}).get("subtask_id") or ((data or {}).get("raw_data") or {}).get("subtask_id")
  278. if job_id and str(job_id) != "0":
  279. snapshot = replace(snapshot, key=f"job:{job_id}", job_key=f"job:{job_id}")
  280. self._progress_resets[printer_id] = (snapshot.progress, _now() + timedelta(minutes=2))
  281. snapshot = replace(snapshot, progress=0, remaining_seconds=0)
  282. self._started_at[printer_id] = _now()
  283. self._new_fallbacks.pop(printer_id, None)
  284. if snapshot.key.startswith("file:"):
  285. self._new_fallbacks[printer_id] = f"{snapshot.key}:{uuid4().hex}"
  286. self._snapshots[printer_id] = snapshot
  287. self._wake.set()
  288. def print_finished(self, printer_id: int, data: dict) -> None:
  289. if self._provider_enabled is False:
  290. return
  291. job = data.get("subtask_id") or (data.get("raw_data") or {}).get("subtask_id")
  292. snapshot = self._snapshots.get(printer_id)
  293. filename = data.get("subtask_name") or data.get("filename")
  294. key = f"job:{job}" if job and str(job) != "0" else None
  295. if key is None and snapshot and (not filename or snapshot.filename == filename):
  296. key = (
  297. self._new_fallbacks.get(printer_id, snapshot.key) if snapshot.key.startswith("file:") else snapshot.key
  298. )
  299. if key is None:
  300. # No unambiguous ownership: the real terminal state can reconcile
  301. # it, but a delayed event must never end the printer's next job.
  302. return
  303. self._finished.append((printer_id, key, str(data.get("status", "completed")).upper()))
  304. snapshot = self._snapshots.get(printer_id)
  305. if snapshot and key in (snapshot.key, snapshot.job_key):
  306. self._snapshots[printer_id] = replace(snapshot, state=self._finished[-1][2], remaining_seconds=0)
  307. self._wake.set()
  308. async def _api(self):
  309. if self._client is None:
  310. self._http = httpx.AsyncClient(
  311. timeout=httpx.Timeout(30, connect=5), follow_redirects=False, headers={"User-Agent": "Bambuddy/1.0"}
  312. )
  313. self._client = NotifyClient(self._http)
  314. return self._client
  315. async def _save(self, row) -> None:
  316. async with self._session() as db:
  317. values = {column.name: getattr(row, column.name) for column in row.__table__.columns if column.name != "id"}
  318. # Never resurrect ownership after provider deletion. A queued cleanup
  319. # still holds this object's newly returned remote ID.
  320. await db.execute(
  321. update(NotificationLiveActivity).where(NotificationLiveActivity.id == row.id).values(**values)
  322. )
  323. await db.commit()
  324. async def tick(self) -> None:
  325. async with self._lock:
  326. async with self._session() as db:
  327. providers = (
  328. await db.scalars(select(NotificationProvider).where(NotificationProvider.provider_type == "notify"))
  329. ).all()
  330. any_enabled = any(self._eligible(provider) for provider in providers)
  331. previously_enabled = self._provider_enabled
  332. self._provider_enabled = any_enabled
  333. if any_enabled and previously_enabled is False:
  334. self._on_activated()
  335. if any_enabled and self._status_grace_until is None:
  336. # Start when reconciliation begins, not at module import.
  337. self._status_grace_until = _now() + _STARTUP_GRACE
  338. if not any_enabled and not previously_enabled and not self._cleanup_once and not self._cleanup_pending:
  339. return
  340. rows = (await db.scalars(select(NotificationLiveActivity))).all()
  341. self._working_rows = rows
  342. names = dict((await db.execute(select(Printer.id, Printer.name))).all()) if any_enabled else {}
  343. finished, self._finished = self._finished, []
  344. snapshots = self._snapshots.copy()
  345. for provider in providers:
  346. config = _config(provider)
  347. owned = [r for r in rows if r.provider_id == provider.id]
  348. enabled = bool(
  349. provider.enabled
  350. and config.get("live_activities") is True
  351. and not str(config.get("device_id", "")).upper().startswith(("GRP", "MC", "WB"))
  352. )
  353. for row in owned:
  354. snapshot = snapshots.get(row.printer_id)
  355. selected = provider.printer_id is None or provider.printer_id == row.printer_id
  356. if (
  357. row.state == "ended"
  358. and row.end_reason == "disabled"
  359. and enabled
  360. and selected
  361. and snapshot
  362. and snapshot.connected
  363. and snapshot.state in _RUNNING
  364. and self._matches(row, snapshot, row.printer_id)
  365. ):
  366. # An explicit re-enable may resume this job's tile. A
  367. # user-dismissed/failed-to-appear tile remains suppressed.
  368. row.state, row.activity_id, row.end_reason = "pending", None, None
  369. row.failures, row.next_attempt_at, row.last_sent_at = 0, None, None
  370. await self._save(row)
  371. if (
  372. snapshot
  373. and snapshot.job_key
  374. and not row.job_key
  375. and self._matches(row, snapshot, row.printer_id)
  376. ):
  377. row.job_key = snapshot.job_key
  378. await self._save(row)
  379. if self._retired_row(row):
  380. continue
  381. if row.state in _STOPPED:
  382. continue
  383. if row.credential_key != _credential_key(config):
  384. # The edit hook owns cleanup using the old credentials.
  385. continue
  386. if row.state == "ending":
  387. await self._end(row, config, row.end_reason or "STOPPED")
  388. continue
  389. snapshot = snapshots.get(row.printer_id)
  390. final = next(
  391. (
  392. status
  393. for pid, key, status in finished
  394. if pid == row.printer_id and key in (row.print_key, row.job_key)
  395. ),
  396. None,
  397. )
  398. selected = provider.printer_id is None or provider.printer_id == row.printer_id
  399. same_print = snapshot and self._matches(row, snapshot, row.printer_id)
  400. if final or not enabled or not selected or row.printer_id not in names:
  401. await self._end(
  402. row, config, final or ("DISABLED" if not enabled or not selected else "STOPPED")
  403. )
  404. elif snapshot and snapshot.connected and snapshot.state in _TERMINAL:
  405. await self._end(row, config, snapshot.state)
  406. elif snapshot and snapshot.connected and snapshot.state in _RUNNING and not same_print:
  407. await self._end(row, config, "STOPPED")
  408. elif snapshot is None and self._status_grace_until and _now() < self._status_grace_until:
  409. # Allow reconnect time without leaving a switched-off
  410. # printer's saved countdown running indefinitely.
  411. continue
  412. else:
  413. # A missing/offline printer is not evidence the print ended.
  414. # Clear its countdown and keep its last known progress.
  415. content = (
  416. snapshot.content(names[row.printer_id], config) if snapshot else json.loads(row.content)
  417. )
  418. if snapshot is None or not snapshot.connected:
  419. content.update(status="Printer offline", endsIn=None, trailing=None)
  420. if snapshot is None:
  421. content["metrics"] = None # Saved ETA/temperature chips are stale after a restart.
  422. await self._sync(
  423. row,
  424. config,
  425. content,
  426. can_start=bool(
  427. snapshot
  428. and snapshot.connected
  429. and snapshot.state in _RUNNING
  430. and enabled
  431. and not self._quiet(provider)
  432. ),
  433. )
  434. if not enabled or self._quiet(provider):
  435. continue
  436. for printer_id, snapshot in snapshots.items():
  437. if (
  438. not snapshot.connected
  439. or snapshot.state not in _RUNNING
  440. or printer_id not in names
  441. or provider.printer_id not in (None, printer_id)
  442. ):
  443. continue
  444. # Finish the old tile before using another device slot for
  445. # this printer. A failed DELETE remains durable and retries.
  446. if any(r.printer_id == printer_id and r.state == "ending" and r.activity_id for r in owned):
  447. continue
  448. if any(self._matches(r, snapshot, printer_id) for r in owned):
  449. continue
  450. key = (
  451. self._new_fallbacks.get(printer_id, snapshot.key)
  452. if snapshot.key.startswith("file:")
  453. else snapshot.key
  454. )
  455. row = NotificationLiveActivity(
  456. provider_id=provider.id,
  457. printer_id=printer_id,
  458. print_key=key,
  459. print_name=snapshot.filename[:255],
  460. job_key=snapshot.job_key,
  461. credential_key=_credential_key(config),
  462. state="pending",
  463. content=json.dumps(snapshot.content(names[printer_id], config)),
  464. created_at=_now(),
  465. )
  466. if self._retired_row(row):
  467. continue
  468. self._working_rows.append(row)
  469. async with self._session() as db:
  470. db.add(row)
  471. await db.commit()
  472. await self._sync(row, config, json.loads(row.content), can_start=True)
  473. self._cleanup_pending = any(
  474. row.state == "ending" and row.provider_id in {p.id for p in providers} for row in rows
  475. )
  476. # Tombstones need no maintenance while the integration is dormant.
  477. if not any_enabled:
  478. return
  479. # Tombstones must outlive any plausible print, but need not grow forever.
  480. async with self._session() as db:
  481. await db.execute(
  482. delete(NotificationLiveActivity).where(
  483. NotificationLiveActivity.state.in_(_STOPPED),
  484. NotificationLiveActivity.created_at < _now() - timedelta(days=30),
  485. NotificationLiveActivity.printer_id.not_in(
  486. [pid for pid, snapshot in snapshots.items() if snapshot.state in _RUNNING]
  487. ),
  488. )
  489. )
  490. await db.commit()
  491. def _matches(self, row, snapshot, printer_id):
  492. if row.printer_id != printer_id:
  493. return False
  494. if row.job_key and snapshot.job_key:
  495. return row.job_key == snapshot.job_key
  496. fallback = self._new_fallbacks.get(printer_id)
  497. if fallback and snapshot.key.startswith("file:"):
  498. return row.print_key == fallback
  499. if (
  500. row.print_key.startswith("file:")
  501. and snapshot.key.startswith("job:")
  502. and row.print_name == snapshot.filename[:255]
  503. and row.job_key is None
  504. and self._started_at.get(printer_id, row.created_at) <= row.created_at
  505. ):
  506. return True # Restart after firmware finally supplied its job ID.
  507. return row.print_key == snapshot.key or (
  508. snapshot.key.startswith("file:") and row.print_key.startswith(snapshot.key + ":")
  509. )
  510. @staticmethod
  511. def _quiet(provider) -> bool:
  512. # Share Bambuddy's configured/local-time quiet-hours semantics. Digest
  513. # and event toggles deliberately never participate in this lifecycle.
  514. from backend.app.services.notification_service import notification_service
  515. return notification_service._is_in_quiet_hours(provider)
  516. async def _sync(self, row, config, content, *, can_start):
  517. if self._retired_row(row):
  518. return
  519. now = _now()
  520. urgent = json.loads(row.content).get("status") != content.get("status")
  521. if row.next_attempt_at and row.next_attempt_at > now and not (urgent and row.last_sent_at and not row.failures):
  522. return
  523. # Store an absolute deadline rather than repeatedly resetting a frozen
  524. # MQTT minute estimate. Refresh only when that estimate or phase changes.
  525. seconds = content.get("endsIn")
  526. if seconds is None:
  527. row.eta_seconds, row.eta_deadline = None, None
  528. elif seconds != row.eta_seconds or row.eta_deadline is None:
  529. row.eta_seconds = seconds
  530. row.eta_deadline = now + timedelta(seconds=seconds)
  531. if row.eta_deadline:
  532. remaining = int((row.eta_deadline - now).total_seconds())
  533. content["endsIn"] = remaining if remaining > 0 else None
  534. api = await self._api()
  535. token = config.get("token", "")
  536. try:
  537. if row.activity_id:
  538. remote = await api.get_activity(row.activity_id, token)
  539. state, reason = remote.get("state"), remote.get("endReason")
  540. if state in {"ended", "dismissed"}:
  541. if reason not in _ROLLOVER or not can_start:
  542. row.state = "suppressed" if reason not in _ROLLOVER else "pending"
  543. row.end_reason = reason or state
  544. if row.state == "pending":
  545. row.activity_id = None
  546. await self._save(row)
  547. return
  548. row.activity_id = None
  549. row.state = "pending"
  550. await self._save(row)
  551. else:
  552. row.state = state or "active"
  553. row.expires_at = _date(remote.get("expiresAt")) or row.expires_at
  554. if can_start and row.expires_at and row.expires_at <= now:
  555. # Poll first: a dismissed tile must never be rolled over.
  556. await api.end_activity(row.activity_id, token)
  557. row.activity_id = None
  558. row.state = "pending"
  559. await self._save(row)
  560. else:
  561. if self._retired_row(row):
  562. return
  563. await api.update_activity(row.activity_id, token, content)
  564. row.content = json.dumps(content)
  565. row.failures = 0
  566. row.last_sent_at = now
  567. row.next_attempt_at = now + timedelta(seconds=_INTERVAL)
  568. await self._save(row)
  569. return
  570. if row.state == "uncertain":
  571. # No usable ID after a timeout/crash. The API has no client
  572. # idempotency key; do not adopt an unrelated device tile.
  573. row.state = "suppressed"
  574. row.end_reason = "unconfirmed-start"
  575. await self._save(row)
  576. return
  577. if not can_start or row.state in _STOPPED:
  578. return
  579. row.state = "uncertain"
  580. row.content = json.dumps(content)
  581. await self._save(row)
  582. if self._retired_row(row):
  583. return
  584. result = await api.start_activity(config.get("device_id", ""), token, content)
  585. row.activity_id = result["activityId"]
  586. row.state = "starting"
  587. row.failures = 0
  588. row.expires_at = _date(result.get("expiresAt")) or now + timedelta(hours=8)
  589. row.last_sent_at = now
  590. row.next_attempt_at = now + timedelta(seconds=_INTERVAL)
  591. await self._save(row)
  592. except NotifyError as error:
  593. await self._failure(row, error)
  594. async def _failure(self, row, error):
  595. ending = row.state == "ending"
  596. if error.activity_id:
  597. row.activity_id = error.activity_id
  598. row.state = "starting"
  599. elif not row.activity_id and error.delivery_state == "unknown":
  600. row.state = "suppressed"
  601. row.end_reason = "unconfirmed-start"
  602. elif error.status_code in (401, 403):
  603. row.state = "suppressed"
  604. row.end_reason = "credentials"
  605. elif error.status_code == 410 and row.activity_id:
  606. # GET on the next pass determines whether this was a dismissal or
  607. # the eight-hour ceiling; never infer rollover from HTTP 410 alone.
  608. pass
  609. elif not row.activity_id:
  610. if error.status_code in (409, 429, 503) or error.retry_after_seconds is not None:
  611. row.state = "pending"
  612. elif error.status_code == 400 and "5" in str(error.payload.get("message", "")):
  613. row.state = "pending" # Device cap: another printer may free a slot.
  614. else:
  615. row.state = "suppressed"
  616. row.end_reason = "start-rejected"
  617. if ending:
  618. row.state = "ending"
  619. row.failures = min((row.failures or 0) + 1, 5)
  620. row.next_attempt_at = _now() + timedelta(
  621. seconds=max(60 * 2 ** (row.failures - 1), error.retry_after_seconds or 0)
  622. )
  623. # Do not log upstream bodies, URLs, or credentials.
  624. logger.warning(
  625. "Notify Live Activity request failed for provider %s (HTTP %s)", row.provider_id, error.status_code
  626. )
  627. await self._save(row)
  628. message = str(error)
  629. if error.retry_after_seconds:
  630. message += f" Retry after {error.retry_after_seconds} seconds."
  631. async with self._session() as db:
  632. await db.execute(
  633. update(NotificationProvider)
  634. .where(NotificationProvider.id == row.provider_id)
  635. .values(
  636. last_error=message,
  637. last_error_at=_now(),
  638. )
  639. )
  640. await db.commit()
  641. async def _end(self, row, config, status):
  642. content = json.loads(row.content)
  643. content.update(
  644. status={"COMPLETED": "Complete", "FINISH": "Complete", "FAILED": "Failed"}.get(status, "Stopped"),
  645. endsIn=None,
  646. )
  647. if status in {"COMPLETED", "FINISH"}:
  648. content["progress"] = 100
  649. # Persist the terminal decision before I/O. A failed end must still
  650. # finish this exact job after restart or after the next job begins.
  651. if row.state != "ending":
  652. row.state = "ending"
  653. row.end_reason = status
  654. row.content = json.dumps(content)
  655. await self._save(row)
  656. if row.failures and row.next_attempt_at and row.next_attempt_at > _now():
  657. return
  658. if row.activity_id:
  659. try:
  660. await (await self._api()).end_activity(row.activity_id, config.get("token", ""), content)
  661. except NotifyError as error:
  662. if error.status_code not in (403, 404, 410):
  663. await self._failure(row, error)
  664. return
  665. row.state = "ended"
  666. row.end_reason = status.lower()
  667. row.content = json.dumps(content)
  668. await self._save(row)
  669. async def cleanup_provider(self, provider_id: int, old_config: dict, *, captured_rows=None) -> None:
  670. """Called after saving a credential/scope change or deleting a provider.
  671. Cleanup is best effort (the remote eight-hour ceiling bounds an outage).
  672. No old credential is retained in a second database location.
  673. """
  674. async with self._lock:
  675. async with self._session() as db:
  676. rows = (
  677. await db.scalars(
  678. select(NotificationLiveActivity).where(
  679. NotificationLiveActivity.provider_id == provider_id,
  680. NotificationLiveActivity.credential_key == _credential_key(old_config),
  681. )
  682. )
  683. ).all()
  684. if captured_rows is not None:
  685. merged = {row.id: row for row in rows}
  686. merged.update(
  687. {row.id: row for row in captured_rows if row.credential_key == _credential_key(old_config)}
  688. )
  689. rows = list(merged.values())
  690. for row in rows:
  691. if row.activity_id and row.state not in _STOPPED:
  692. try:
  693. await (await self._api()).end_activity(row.activity_id, old_config.get("token", ""))
  694. except NotifyError:
  695. logger.warning("Notify Live Activity cleanup failed for provider %s", provider_id)
  696. async with self._session() as db:
  697. await db.execute(
  698. delete(NotificationLiveActivity).where(
  699. NotificationLiveActivity.provider_id == provider_id,
  700. NotificationLiveActivity.credential_key == _credential_key(old_config),
  701. )
  702. )
  703. await db.commit()
  704. for row in rows:
  705. row.state = "ended"
  706. self._wake.set()
  707. notify_live_activities = NotifyLiveActivityService()