| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456 |
- """Persistent Notify! Lock Screen widgets, separate from Live Activities.
- These widgets hold stored values: iOS chooses when to refresh them (usually about
- 15 minutes). Bambuddy updates changed content at most once per minute, never
- promises a ticking countdown, and addresses only its own saved WG identifiers.
- """
- import asyncio
- import json
- import logging
- from datetime import timedelta
- import httpx
- from sqlalchemy import and_, case, delete, or_, select, update
- from backend.app.core.database import async_session
- from backend.app.models.notification import NotificationProvider
- from backend.app.models.notification_lock_screen_widget import NotificationLockScreenWidget
- from backend.app.models.printer import Printer
- from backend.app.services.notify_client import NotifyClient, NotifyError
- from backend.app.services.notify_live_activities import PrintSnapshot, _credential_key as _fingerprint, _now
- from backend.app.services.notify_worker import NotifyWorkerLifecycle
- logger = logging.getLogger(__name__)
- _ERROR_PREFIX = "Notify! Lock Screen widget: "
- _CAPACITY_ERROR_PREFIX = "Notify! widget capacity: "
- _INTERVAL = 60
- def _credential_key(config):
- return _fingerprint(
- {
- "device_id": str(config.get("device_id", "")).strip(),
- "token": str(config.get("token", "")).strip(),
- }
- )
- def widget_content(printer_name: str, state) -> dict:
- """A static, compact status card. Faults replace the gauge with a diagnosis."""
- content = {
- "title": printer_name.replace("\x00", "")[:120] or "Bambuddy",
- "value": "Offline",
- "unit": None,
- "detail": "Printer is not connected",
- "symbol": "printer.fill",
- "tint": "#00AE42",
- "progress": None,
- }
- if state is not None and state.connected:
- # Shared parsing keeps HMS advisory filtering and preparation-stage
- # progress consistent with Bambuddy's Live Activity and printer card.
- snapshot = PrintSnapshot.from_state(state)
- tile = snapshot.content(printer_name)
- phase = tile["status"]
- if snapshot.fault:
- content.update(value="Error", detail=snapshot.fault[:120], tint="#E5484D")
- elif phase == "Paused":
- content.update(value="Paused", detail=f"Print paused at {snapshot.progress:.0f}%")
- elif snapshot.state in {"RUNNING", "PRINTING", "PREPARE", "SLICING"}:
- progress = tile["progress"]
- parts = [phase]
- if snapshot.layers:
- parts.append(f"Layer {snapshot.layer}/{snapshot.layers}")
- if snapshot.remaining_seconds > 0:
- minutes = snapshot.remaining_seconds // 60
- eta = f"{minutes // 60}h {minutes % 60}m" if minutes >= 60 else f"{minutes}m"
- parts.append(f"About {eta} left")
- content.update(value=f"{progress:.0f}", unit="%", detail=" · ".join(parts)[:120], progress=progress)
- elif snapshot.state in {"FINISH", "COMPLETED"}:
- content.update(value="Complete", detail="Print completed", progress=100)
- elif snapshot.state == "FAILED":
- content.update(value="Failed", detail="Print failed", tint="#E5484D")
- elif snapshot.state in {"CANCELLED", "ABORTED", "STOPPED"}:
- content.update(value="Stopped", detail="Print stopped")
- elif snapshot.state == "IDLE":
- content.update(value="Idle", detail="Ready to print")
- else:
- content.update(value="Connecting", detail="Waiting for printer status")
- # Notify caps the serialized merged content at 1024 UTF-8 bytes. Reserve
- # space rather than letting emoji printer names make every update fail.
- while len(json.dumps(content, ensure_ascii=False, separators=(",", ":")).encode()) > 950:
- field = "detail" if len(content["detail"]) > 20 else "title"
- content[field] = content[field][:-5]
- return content
- def _status(printer_id):
- from backend.app.services.printer_manager import printer_manager
- return printer_manager.get_status(printer_id)
- class NotifyWidgetService(NotifyWorkerLifecycle):
- ownership_model = NotificationLockScreenWidget
- feature_field = "lock_screen_widgets"
- worker_name = "notify-lock-screen-widgets"
- _cleanup_key = staticmethod(_credential_key)
- def __init__(self, session_factory=None, client=None, state_getter=None):
- self._session = session_factory or async_session
- self._client = client
- self._http: httpx.AsyncClient | None = None
- self._state = state_getter or _status
- self._task: asyncio.Task | None = None
- self._lock = asyncio.Lock()
- self._init_worker()
- self._capacity_devices: set[str] = set()
- self._retry_capacity = False
- self._over_capacity_providers: set[int] = set()
- def providers_changed(self):
- self._retry_capacity = True
- self._capacity_devices.clear()
- super().providers_changed()
- async def _api(self):
- if self._client is None:
- self._http = httpx.AsyncClient(
- timeout=httpx.Timeout(30, connect=5), follow_redirects=False, headers={"User-Agent": "Bambuddy/1.0"}
- )
- self._client = NotifyClient(self._http)
- return self._client
- async def _save(self, row):
- async with self._session() as db:
- values = {column.name: getattr(row, column.name) for column in row.__table__.columns if column.name != "id"}
- # A provider edit may request deletion while this worker is awaiting
- # HTTP. Preserve that intent while capturing a just-returned WG ID,
- # otherwise a rapid off/on would silently cancel the user's cleanup.
- values["state"] = case((NotificationLockScreenWidget.state == "deleting", "deleting"), else_=row.state)
- await db.execute(
- update(NotificationLockScreenWidget).where(NotificationLockScreenWidget.id == row.id).values(**values)
- )
- await db.commit()
- async def _remove(self, row):
- async with self._session() as db:
- await db.execute(delete(NotificationLockScreenWidget).where(NotificationLockScreenWidget.id == row.id))
- await db.commit()
- row.state = "deleted"
- async def tick(self):
- async with self._lock:
- async with self._session() as db:
- providers = (
- await db.scalars(select(NotificationProvider).where(NotificationProvider.provider_type == "notify"))
- ).all()
- enabled = any(self._eligible(provider) for provider in providers)
- previously_enabled = self._provider_enabled
- self._provider_enabled = enabled
- if not enabled and not previously_enabled and not self._cleanup_once and not self._cleanup_pending:
- return
- rows = (await db.scalars(select(NotificationLockScreenWidget))).all()
- self._working_rows = rows
- printers = (
- dict((await db.execute(select(Printer.id, Printer.name).where(Printer.is_active.is_(True)))).all())
- if enabled
- else {}
- )
- retry_capacity, self._retry_capacity = self._retry_capacity, False
- if not retry_capacity:
- self._capacity_devices.update(row.credential_key for row in rows if row.state == "capacity")
- for provider in providers:
- config = json.loads(provider.config) if isinstance(provider.config, str) else provider.config
- enabled = (
- provider.enabled
- and config.get("lock_screen_widgets") is True
- and not str(config.get("device_id", "")).strip().upper().startswith(("GRP", "WB", "MC"))
- )
- credential = _credential_key(config)
- owned = [r for r in rows if r.provider_id == provider.id]
- selected = {pid: name for pid, name in printers.items() if provider.printer_id in (None, pid)}
- if len(selected) > 10:
- # Preserve already-owned widgets first; adding an eleventh
- # printer must not delete one merely because query order changed.
- retained = [row.printer_id for row in owned if row.printer_id in selected]
- chosen = list(dict.fromkeys([*retained, *sorted(selected)]))[:10]
- selected = {pid: selected[pid] for pid in chosen}
- if provider.id not in self._over_capacity_providers:
- logger.warning(
- "Notify widget provider %s exceeds ten active printers; select a printer", provider.id
- )
- self._over_capacity_providers.add(provider.id)
- else:
- self._over_capacity_providers.discard(provider.id)
- for row in owned:
- if retry_capacity and row.state == "capacity":
- row.state, row.failures, row.next_attempt_at = "pending", 0, None
- await self._save(row)
- if self._retired_row(row):
- continue
- if row.credential_key != credential:
- continue # Credential-edit hook cleans up with the previous token.
- if row.state == "deleting" or not enabled or row.printer_id not in selected:
- await self._delete(row, config)
- elif row.state == "uncertain" and not row.failures:
- await self._failure(
- row, NotifyError("Previous creation was interrupted", delivery_state="unknown")
- )
- elif row.state not in {"uncertain", "suppressed", "capacity"}:
- await self._sync(
- row, config, widget_content(selected[row.printer_id], self._state(row.printer_id))
- )
- if not enabled:
- continue
- for printer_id, name in selected.items():
- if credential in self._capacity_devices:
- break
- if any(r.printer_id == printer_id for r in owned):
- continue
- row = NotificationLockScreenWidget(
- provider_id=provider.id,
- printer_id=printer_id,
- credential_key=credential,
- state="pending",
- content="{}",
- created_at=_now(),
- failures=0,
- )
- if self._retired_row(row):
- continue
- self._working_rows.append(row)
- async with self._session() as db:
- db.add(row)
- await db.commit()
- await self._sync(row, config, widget_content(name, self._state(printer_id)))
- self._cleanup_pending = any(
- row.state == "deleting" and row.provider_id in {p.id for p in providers} for row in rows
- )
- async def _sync(self, row, config, content):
- if self._retired_row(row):
- return
- if row.next_attempt_at and row.next_attempt_at > _now():
- return
- encoded = json.dumps(content, ensure_ascii=False, sort_keys=True)
- if row.widget_id and row.content == encoded:
- return
- api = await self._api()
- create_attempted = False
- try:
- if row.widget_id:
- await api.update_widget(row.widget_id, config.get("token", ""), content)
- else:
- # Device and widget IDs can both be eight characters. Prove
- # the configured target is a device before POST: new=1 on a WG
- # URL still addresses that existing (possibly unrelated) widget.
- await api.list_widgets(config.get("device_id", ""), config.get("token", ""))
- # If the process dies after this commit, the create may have
- # happened. The next process must never repeat new=1 blindly.
- row.state = "uncertain"
- await self._save(row)
- if self._retired_row(row):
- row.state = "pending" # No create was sent; cleanup need not warn of an unknown widget.
- return
- create_attempted = True
- result = await api.create_widget(config.get("device_id", ""), config.get("token", ""), content)
- row.widget_id = result["widgetId"]
- row.state = "active"
- row.content = encoded
- row.failures = 0
- row.last_sent_at = _now()
- row.next_attempt_at = _now() + timedelta(seconds=_INTERVAL)
- await self._save(row)
- async with self._session() as db:
- # A successful widget update must not clear a push/Live Activity error.
- await db.execute(
- update(NotificationProvider)
- .where(
- NotificationProvider.id == row.provider_id,
- or_(
- NotificationProvider.last_error.startswith(_ERROR_PREFIX),
- and_(
- NotificationProvider.last_error.startswith(_CAPACITY_ERROR_PREFIX),
- ~select(NotificationLockScreenWidget.id)
- .where(
- NotificationLockScreenWidget.provider_id == row.provider_id,
- NotificationLockScreenWidget.state == "capacity",
- )
- .exists(),
- ),
- ),
- )
- .values(last_error=None, last_error_at=None)
- )
- await db.commit()
- except NotifyError as error:
- if not row.widget_id and not create_attempted:
- # A failed read cannot have created anything. Do not adopt any
- # WG that might appear in its response, and retry safely.
- error = NotifyError(
- str(error),
- status_code=error.status_code,
- delivery_state="not-delivered",
- retry_after_seconds=error.retry_after_seconds or _INTERVAL,
- )
- await self._failure(row, error)
- async def _failure(self, row, error):
- deleting = row.state == "deleting"
- widget_id = getattr(error, "widget_id", None)
- if widget_id:
- row.widget_id = widget_id
- row.state = "active"
- elif not row.widget_id and error.delivery_state == "unknown":
- row.state = "uncertain"
- elif error.status_code in (401, 403, 404):
- row.state = "suppressed"
- elif not row.widget_id:
- capacity = error.status_code == 400 and "10" in str(error.payload.get("message", ""))
- if capacity:
- row.state = "capacity"
- self._capacity_devices.add(row.credential_key)
- elif error.status_code in (429, 503) or error.retry_after_seconds is not None:
- row.state = "pending"
- else:
- row.state = "suppressed"
- if deleting:
- row.state = "deleting"
- row.failures = min((row.failures or 0) + 1, 5)
- delay = max(_INTERVAL * 2 ** (row.failures - 1), error.retry_after_seconds or 0)
- row.next_attempt_at = _now() + timedelta(seconds=delay)
- await self._save(row)
- message = (_CAPACITY_ERROR_PREFIX if row.state == "capacity" else _ERROR_PREFIX) + str(error)
- if row.state == "uncertain":
- message += (
- " Creation could not be confirmed. Check Notify! and remove any duplicate or unwanted widget, "
- "then turn Lock Screen widgets off and on to retry."
- )
- elif row.state == "capacity":
- message += " Free a widget slot in Notify! and save this provider to retry."
- elif error.retry_after_seconds:
- message += f" Retry after {error.retry_after_seconds} seconds."
- async with self._session() as db:
- await db.execute(
- update(NotificationProvider)
- .where(NotificationProvider.id == row.provider_id)
- .values(
- last_error=message,
- last_error_at=_now(),
- )
- )
- await db.commit()
- logger.warning("Notify widget request failed for provider %s (HTTP %s)", row.provider_id, error.status_code)
- async def _delete(self, row, config):
- if row.state != "deleting":
- row.state = "deleting"
- await self._save(row)
- if not row.widget_id:
- await self._remove(row)
- return
- if row.failures and row.next_attempt_at and row.next_attempt_at > _now():
- return
- if row.widget_id:
- try:
- await (await self._api()).delete_widget(row.widget_id, config.get("token", ""))
- except NotifyError as error:
- if error.status_code == 403:
- # The API deliberately conflates a removed WG and invalid
- # credentials. Only an authenticated device list can prove
- # our exact widget is gone; never create/adopt from that list.
- try:
- existing = await (await self._api()).list_widgets(
- config.get("device_id", ""), config.get("token", "")
- )
- except NotifyError as listing_error:
- await self._failure(row, listing_error)
- return
- if any(widget.get("widgetId") == row.widget_id for widget in existing.get("widgets", [])):
- await self._failure(row, error)
- return
- elif error.status_code != 404:
- await self._failure(row, error)
- return
- await self._remove(row)
- async def cleanup_provider(self, provider_id: int, old_config: dict, *, captured_rows=None):
- """Transfer a token rotation, or clean up before old credentials vanish.
- Widgets have no expiry. If a different-device/delete cleanup fails,
- Notify's app must remove the saved widget manually. Never retain or log
- its credential-bearing updateUrl.
- """
- async with self._lock:
- async with self._session() as db:
- current = await db.get(NotificationProvider, provider_id)
- rows = (
- await db.scalars(
- select(NotificationLockScreenWidget).where(
- NotificationLockScreenWidget.provider_id == provider_id,
- NotificationLockScreenWidget.credential_key == _credential_key(old_config),
- )
- )
- ).all()
- if captured_rows is not None:
- merged = {row.id: row for row in rows}
- merged.update(
- {row.id: row for row in captured_rows if row.credential_key == _credential_key(old_config)}
- )
- rows = list(merged.values())
- config = (
- (json.loads(current.config) if isinstance(current.config, str) else current.config) if current else {}
- )
- same_device = (
- current
- and current.provider_type == "notify"
- and str(config.get("device_id", "")).strip() == str(old_config.get("device_id", "")).strip()
- and _credential_key(config) != _credential_key(old_config)
- )
- if same_device:
- # A rotated token still addresses the same device and WG. Never
- # delete/recreate its permanent widget, or revive an uncertain
- # create whose identifier we never received.
- for row in rows:
- row.credential_key = _credential_key(config)
- row.failures, row.next_attempt_at = 0, None
- if row.widget_id and row.state == "suppressed":
- row.state = "active"
- row.content = "{}" # Verify the new token even if printer content is unchanged.
- await self._save(row)
- return
- cleanup_failed = False
- for row in rows:
- if row.state == "deleted":
- continue
- if row.widget_id:
- try:
- await (await self._api()).delete_widget(row.widget_id, old_config.get("token", ""))
- except NotifyError:
- cleanup_failed = True
- elif row.state == "uncertain":
- cleanup_failed = True
- await self._remove(row)
- if cleanup_failed:
- logger.warning(
- "Notify widget cleanup failed for provider %s; remove old widgets in Notify!", provider_id
- )
- # Keep this distinct from transient widget-send errors: success
- # on the new device does not prove an old permanent widget gone.
- async with self._session() as db:
- await db.execute(
- update(NotificationProvider)
- .where(NotificationProvider.id == provider_id)
- .values(
- last_error="Notify! widget cleanup: Remove the previous device's Bambuddy widgets in Notify!; automatic cleanup could not be confirmed.",
- last_error_at=_now(),
- )
- )
- await db.commit()
- notify_widgets = NotifyWidgetService()
|