notify_widgets.py 22 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456
  1. """Persistent Notify! Lock Screen widgets, separate from Live Activities.
  2. These widgets hold stored values: iOS chooses when to refresh them (usually about
  3. 15 minutes). Bambuddy updates changed content at most once per minute, never
  4. promises a ticking countdown, and addresses only its own saved WG identifiers.
  5. """
  6. import asyncio
  7. import json
  8. import logging
  9. from datetime import timedelta
  10. import httpx
  11. from sqlalchemy import and_, case, delete, or_, select, update
  12. from backend.app.core.database import async_session
  13. from backend.app.models.notification import NotificationProvider
  14. from backend.app.models.notification_lock_screen_widget import NotificationLockScreenWidget
  15. from backend.app.models.printer import Printer
  16. from backend.app.services.notify_client import NotifyClient, NotifyError
  17. from backend.app.services.notify_live_activities import PrintSnapshot, _credential_key as _fingerprint, _now
  18. from backend.app.services.notify_worker import NotifyWorkerLifecycle
  19. logger = logging.getLogger(__name__)
  20. _ERROR_PREFIX = "Notify! Lock Screen widget: "
  21. _CAPACITY_ERROR_PREFIX = "Notify! widget capacity: "
  22. _INTERVAL = 60
  23. def _credential_key(config):
  24. return _fingerprint(
  25. {
  26. "device_id": str(config.get("device_id", "")).strip(),
  27. "token": str(config.get("token", "")).strip(),
  28. }
  29. )
  30. def widget_content(printer_name: str, state) -> dict:
  31. """A static, compact status card. Faults replace the gauge with a diagnosis."""
  32. content = {
  33. "title": printer_name.replace("\x00", "")[:120] or "Bambuddy",
  34. "value": "Offline",
  35. "unit": None,
  36. "detail": "Printer is not connected",
  37. "symbol": "printer.fill",
  38. "tint": "#00AE42",
  39. "progress": None,
  40. }
  41. if state is not None and state.connected:
  42. # Shared parsing keeps HMS advisory filtering and preparation-stage
  43. # progress consistent with Bambuddy's Live Activity and printer card.
  44. snapshot = PrintSnapshot.from_state(state)
  45. tile = snapshot.content(printer_name)
  46. phase = tile["status"]
  47. if snapshot.fault:
  48. content.update(value="Error", detail=snapshot.fault[:120], tint="#E5484D")
  49. elif phase == "Paused":
  50. content.update(value="Paused", detail=f"Print paused at {snapshot.progress:.0f}%")
  51. elif snapshot.state in {"RUNNING", "PRINTING", "PREPARE", "SLICING"}:
  52. progress = tile["progress"]
  53. parts = [phase]
  54. if snapshot.layers:
  55. parts.append(f"Layer {snapshot.layer}/{snapshot.layers}")
  56. if snapshot.remaining_seconds > 0:
  57. minutes = snapshot.remaining_seconds // 60
  58. eta = f"{minutes // 60}h {minutes % 60}m" if minutes >= 60 else f"{minutes}m"
  59. parts.append(f"About {eta} left")
  60. content.update(value=f"{progress:.0f}", unit="%", detail=" · ".join(parts)[:120], progress=progress)
  61. elif snapshot.state in {"FINISH", "COMPLETED"}:
  62. content.update(value="Complete", detail="Print completed", progress=100)
  63. elif snapshot.state == "FAILED":
  64. content.update(value="Failed", detail="Print failed", tint="#E5484D")
  65. elif snapshot.state in {"CANCELLED", "ABORTED", "STOPPED"}:
  66. content.update(value="Stopped", detail="Print stopped")
  67. elif snapshot.state == "IDLE":
  68. content.update(value="Idle", detail="Ready to print")
  69. else:
  70. content.update(value="Connecting", detail="Waiting for printer status")
  71. # Notify caps the serialized merged content at 1024 UTF-8 bytes. Reserve
  72. # space rather than letting emoji printer names make every update fail.
  73. while len(json.dumps(content, ensure_ascii=False, separators=(",", ":")).encode()) > 950:
  74. field = "detail" if len(content["detail"]) > 20 else "title"
  75. content[field] = content[field][:-5]
  76. return content
  77. def _status(printer_id):
  78. from backend.app.services.printer_manager import printer_manager
  79. return printer_manager.get_status(printer_id)
  80. class NotifyWidgetService(NotifyWorkerLifecycle):
  81. ownership_model = NotificationLockScreenWidget
  82. feature_field = "lock_screen_widgets"
  83. worker_name = "notify-lock-screen-widgets"
  84. _cleanup_key = staticmethod(_credential_key)
  85. def __init__(self, session_factory=None, client=None, state_getter=None):
  86. self._session = session_factory or async_session
  87. self._client = client
  88. self._http: httpx.AsyncClient | None = None
  89. self._state = state_getter or _status
  90. self._task: asyncio.Task | None = None
  91. self._lock = asyncio.Lock()
  92. self._init_worker()
  93. self._capacity_devices: set[str] = set()
  94. self._retry_capacity = False
  95. self._over_capacity_providers: set[int] = set()
  96. def providers_changed(self):
  97. self._retry_capacity = True
  98. self._capacity_devices.clear()
  99. super().providers_changed()
  100. async def _api(self):
  101. if self._client is None:
  102. self._http = httpx.AsyncClient(
  103. timeout=httpx.Timeout(30, connect=5), follow_redirects=False, headers={"User-Agent": "Bambuddy/1.0"}
  104. )
  105. self._client = NotifyClient(self._http)
  106. return self._client
  107. async def _save(self, row):
  108. async with self._session() as db:
  109. values = {column.name: getattr(row, column.name) for column in row.__table__.columns if column.name != "id"}
  110. # A provider edit may request deletion while this worker is awaiting
  111. # HTTP. Preserve that intent while capturing a just-returned WG ID,
  112. # otherwise a rapid off/on would silently cancel the user's cleanup.
  113. values["state"] = case((NotificationLockScreenWidget.state == "deleting", "deleting"), else_=row.state)
  114. await db.execute(
  115. update(NotificationLockScreenWidget).where(NotificationLockScreenWidget.id == row.id).values(**values)
  116. )
  117. await db.commit()
  118. async def _remove(self, row):
  119. async with self._session() as db:
  120. await db.execute(delete(NotificationLockScreenWidget).where(NotificationLockScreenWidget.id == row.id))
  121. await db.commit()
  122. row.state = "deleted"
  123. async def tick(self):
  124. async with self._lock:
  125. async with self._session() as db:
  126. providers = (
  127. await db.scalars(select(NotificationProvider).where(NotificationProvider.provider_type == "notify"))
  128. ).all()
  129. enabled = any(self._eligible(provider) for provider in providers)
  130. previously_enabled = self._provider_enabled
  131. self._provider_enabled = enabled
  132. if not enabled and not previously_enabled and not self._cleanup_once and not self._cleanup_pending:
  133. return
  134. rows = (await db.scalars(select(NotificationLockScreenWidget))).all()
  135. self._working_rows = rows
  136. printers = (
  137. dict((await db.execute(select(Printer.id, Printer.name).where(Printer.is_active.is_(True)))).all())
  138. if enabled
  139. else {}
  140. )
  141. retry_capacity, self._retry_capacity = self._retry_capacity, False
  142. if not retry_capacity:
  143. self._capacity_devices.update(row.credential_key for row in rows if row.state == "capacity")
  144. for provider in providers:
  145. config = json.loads(provider.config) if isinstance(provider.config, str) else provider.config
  146. enabled = (
  147. provider.enabled
  148. and config.get("lock_screen_widgets") is True
  149. and not str(config.get("device_id", "")).strip().upper().startswith(("GRP", "WB", "MC"))
  150. )
  151. credential = _credential_key(config)
  152. owned = [r for r in rows if r.provider_id == provider.id]
  153. selected = {pid: name for pid, name in printers.items() if provider.printer_id in (None, pid)}
  154. if len(selected) > 10:
  155. # Preserve already-owned widgets first; adding an eleventh
  156. # printer must not delete one merely because query order changed.
  157. retained = [row.printer_id for row in owned if row.printer_id in selected]
  158. chosen = list(dict.fromkeys([*retained, *sorted(selected)]))[:10]
  159. selected = {pid: selected[pid] for pid in chosen}
  160. if provider.id not in self._over_capacity_providers:
  161. logger.warning(
  162. "Notify widget provider %s exceeds ten active printers; select a printer", provider.id
  163. )
  164. self._over_capacity_providers.add(provider.id)
  165. else:
  166. self._over_capacity_providers.discard(provider.id)
  167. for row in owned:
  168. if retry_capacity and row.state == "capacity":
  169. row.state, row.failures, row.next_attempt_at = "pending", 0, None
  170. await self._save(row)
  171. if self._retired_row(row):
  172. continue
  173. if row.credential_key != credential:
  174. continue # Credential-edit hook cleans up with the previous token.
  175. if row.state == "deleting" or not enabled or row.printer_id not in selected:
  176. await self._delete(row, config)
  177. elif row.state == "uncertain" and not row.failures:
  178. await self._failure(
  179. row, NotifyError("Previous creation was interrupted", delivery_state="unknown")
  180. )
  181. elif row.state not in {"uncertain", "suppressed", "capacity"}:
  182. await self._sync(
  183. row, config, widget_content(selected[row.printer_id], self._state(row.printer_id))
  184. )
  185. if not enabled:
  186. continue
  187. for printer_id, name in selected.items():
  188. if credential in self._capacity_devices:
  189. break
  190. if any(r.printer_id == printer_id for r in owned):
  191. continue
  192. row = NotificationLockScreenWidget(
  193. provider_id=provider.id,
  194. printer_id=printer_id,
  195. credential_key=credential,
  196. state="pending",
  197. content="{}",
  198. created_at=_now(),
  199. failures=0,
  200. )
  201. if self._retired_row(row):
  202. continue
  203. self._working_rows.append(row)
  204. async with self._session() as db:
  205. db.add(row)
  206. await db.commit()
  207. await self._sync(row, config, widget_content(name, self._state(printer_id)))
  208. self._cleanup_pending = any(
  209. row.state == "deleting" and row.provider_id in {p.id for p in providers} for row in rows
  210. )
  211. async def _sync(self, row, config, content):
  212. if self._retired_row(row):
  213. return
  214. if row.next_attempt_at and row.next_attempt_at > _now():
  215. return
  216. encoded = json.dumps(content, ensure_ascii=False, sort_keys=True)
  217. if row.widget_id and row.content == encoded:
  218. return
  219. api = await self._api()
  220. create_attempted = False
  221. try:
  222. if row.widget_id:
  223. await api.update_widget(row.widget_id, config.get("token", ""), content)
  224. else:
  225. # Device and widget IDs can both be eight characters. Prove
  226. # the configured target is a device before POST: new=1 on a WG
  227. # URL still addresses that existing (possibly unrelated) widget.
  228. await api.list_widgets(config.get("device_id", ""), config.get("token", ""))
  229. # If the process dies after this commit, the create may have
  230. # happened. The next process must never repeat new=1 blindly.
  231. row.state = "uncertain"
  232. await self._save(row)
  233. if self._retired_row(row):
  234. row.state = "pending" # No create was sent; cleanup need not warn of an unknown widget.
  235. return
  236. create_attempted = True
  237. result = await api.create_widget(config.get("device_id", ""), config.get("token", ""), content)
  238. row.widget_id = result["widgetId"]
  239. row.state = "active"
  240. row.content = encoded
  241. row.failures = 0
  242. row.last_sent_at = _now()
  243. row.next_attempt_at = _now() + timedelta(seconds=_INTERVAL)
  244. await self._save(row)
  245. async with self._session() as db:
  246. # A successful widget update must not clear a push/Live Activity error.
  247. await db.execute(
  248. update(NotificationProvider)
  249. .where(
  250. NotificationProvider.id == row.provider_id,
  251. or_(
  252. NotificationProvider.last_error.startswith(_ERROR_PREFIX),
  253. and_(
  254. NotificationProvider.last_error.startswith(_CAPACITY_ERROR_PREFIX),
  255. ~select(NotificationLockScreenWidget.id)
  256. .where(
  257. NotificationLockScreenWidget.provider_id == row.provider_id,
  258. NotificationLockScreenWidget.state == "capacity",
  259. )
  260. .exists(),
  261. ),
  262. ),
  263. )
  264. .values(last_error=None, last_error_at=None)
  265. )
  266. await db.commit()
  267. except NotifyError as error:
  268. if not row.widget_id and not create_attempted:
  269. # A failed read cannot have created anything. Do not adopt any
  270. # WG that might appear in its response, and retry safely.
  271. error = NotifyError(
  272. str(error),
  273. status_code=error.status_code,
  274. delivery_state="not-delivered",
  275. retry_after_seconds=error.retry_after_seconds or _INTERVAL,
  276. )
  277. await self._failure(row, error)
  278. async def _failure(self, row, error):
  279. deleting = row.state == "deleting"
  280. widget_id = getattr(error, "widget_id", None)
  281. if widget_id:
  282. row.widget_id = widget_id
  283. row.state = "active"
  284. elif not row.widget_id and error.delivery_state == "unknown":
  285. row.state = "uncertain"
  286. elif error.status_code in (401, 403, 404):
  287. row.state = "suppressed"
  288. elif not row.widget_id:
  289. capacity = error.status_code == 400 and "10" in str(error.payload.get("message", ""))
  290. if capacity:
  291. row.state = "capacity"
  292. self._capacity_devices.add(row.credential_key)
  293. elif error.status_code in (429, 503) or error.retry_after_seconds is not None:
  294. row.state = "pending"
  295. else:
  296. row.state = "suppressed"
  297. if deleting:
  298. row.state = "deleting"
  299. row.failures = min((row.failures or 0) + 1, 5)
  300. delay = max(_INTERVAL * 2 ** (row.failures - 1), error.retry_after_seconds or 0)
  301. row.next_attempt_at = _now() + timedelta(seconds=delay)
  302. await self._save(row)
  303. message = (_CAPACITY_ERROR_PREFIX if row.state == "capacity" else _ERROR_PREFIX) + str(error)
  304. if row.state == "uncertain":
  305. message += (
  306. " Creation could not be confirmed. Check Notify! and remove any duplicate or unwanted widget, "
  307. "then turn Lock Screen widgets off and on to retry."
  308. )
  309. elif row.state == "capacity":
  310. message += " Free a widget slot in Notify! and save this provider to retry."
  311. elif error.retry_after_seconds:
  312. message += f" Retry after {error.retry_after_seconds} seconds."
  313. async with self._session() as db:
  314. await db.execute(
  315. update(NotificationProvider)
  316. .where(NotificationProvider.id == row.provider_id)
  317. .values(
  318. last_error=message,
  319. last_error_at=_now(),
  320. )
  321. )
  322. await db.commit()
  323. logger.warning("Notify widget request failed for provider %s (HTTP %s)", row.provider_id, error.status_code)
  324. async def _delete(self, row, config):
  325. if row.state != "deleting":
  326. row.state = "deleting"
  327. await self._save(row)
  328. if not row.widget_id:
  329. await self._remove(row)
  330. return
  331. if row.failures and row.next_attempt_at and row.next_attempt_at > _now():
  332. return
  333. if row.widget_id:
  334. try:
  335. await (await self._api()).delete_widget(row.widget_id, config.get("token", ""))
  336. except NotifyError as error:
  337. if error.status_code == 403:
  338. # The API deliberately conflates a removed WG and invalid
  339. # credentials. Only an authenticated device list can prove
  340. # our exact widget is gone; never create/adopt from that list.
  341. try:
  342. existing = await (await self._api()).list_widgets(
  343. config.get("device_id", ""), config.get("token", "")
  344. )
  345. except NotifyError as listing_error:
  346. await self._failure(row, listing_error)
  347. return
  348. if any(widget.get("widgetId") == row.widget_id for widget in existing.get("widgets", [])):
  349. await self._failure(row, error)
  350. return
  351. elif error.status_code != 404:
  352. await self._failure(row, error)
  353. return
  354. await self._remove(row)
  355. async def cleanup_provider(self, provider_id: int, old_config: dict, *, captured_rows=None):
  356. """Transfer a token rotation, or clean up before old credentials vanish.
  357. Widgets have no expiry. If a different-device/delete cleanup fails,
  358. Notify's app must remove the saved widget manually. Never retain or log
  359. its credential-bearing updateUrl.
  360. """
  361. async with self._lock:
  362. async with self._session() as db:
  363. current = await db.get(NotificationProvider, provider_id)
  364. rows = (
  365. await db.scalars(
  366. select(NotificationLockScreenWidget).where(
  367. NotificationLockScreenWidget.provider_id == provider_id,
  368. NotificationLockScreenWidget.credential_key == _credential_key(old_config),
  369. )
  370. )
  371. ).all()
  372. if captured_rows is not None:
  373. merged = {row.id: row for row in rows}
  374. merged.update(
  375. {row.id: row for row in captured_rows if row.credential_key == _credential_key(old_config)}
  376. )
  377. rows = list(merged.values())
  378. config = (
  379. (json.loads(current.config) if isinstance(current.config, str) else current.config) if current else {}
  380. )
  381. same_device = (
  382. current
  383. and current.provider_type == "notify"
  384. and str(config.get("device_id", "")).strip() == str(old_config.get("device_id", "")).strip()
  385. and _credential_key(config) != _credential_key(old_config)
  386. )
  387. if same_device:
  388. # A rotated token still addresses the same device and WG. Never
  389. # delete/recreate its permanent widget, or revive an uncertain
  390. # create whose identifier we never received.
  391. for row in rows:
  392. row.credential_key = _credential_key(config)
  393. row.failures, row.next_attempt_at = 0, None
  394. if row.widget_id and row.state == "suppressed":
  395. row.state = "active"
  396. row.content = "{}" # Verify the new token even if printer content is unchanged.
  397. await self._save(row)
  398. return
  399. cleanup_failed = False
  400. for row in rows:
  401. if row.state == "deleted":
  402. continue
  403. if row.widget_id:
  404. try:
  405. await (await self._api()).delete_widget(row.widget_id, old_config.get("token", ""))
  406. except NotifyError:
  407. cleanup_failed = True
  408. elif row.state == "uncertain":
  409. cleanup_failed = True
  410. await self._remove(row)
  411. if cleanup_failed:
  412. logger.warning(
  413. "Notify widget cleanup failed for provider %s; remove old widgets in Notify!", provider_id
  414. )
  415. # Keep this distinct from transient widget-send errors: success
  416. # on the new device does not prove an old permanent widget gone.
  417. async with self._session() as db:
  418. await db.execute(
  419. update(NotificationProvider)
  420. .where(NotificationProvider.id == provider_id)
  421. .values(
  422. last_error="Notify! widget cleanup: Remove the previous device's Bambuddy widgets in Notify!; automatic cleanup could not be confirmed.",
  423. last_error_at=_now(),
  424. )
  425. )
  426. await db.commit()
  427. notify_widgets = NotifyWidgetService()