"""Helper for writing notifications without callers having to thread a store. The harness can NOT import ``app.*`` (firewall), so the Gateway app injects its ``NotificationStore`` instance via ``set_notification_store(...)`` at startup. Until that happens, every call is a no-op (silently logged) so non-Gateway entry points (tests, embedded client) don't crash on missing infrastructure. """ from __future__ import annotations import logging import uuid from typing import Any from deerflow.persistence.notifications.base import NotificationStore logger = logging.getLogger(__name__) _store: NotificationStore | None = None def set_notification_store(store: NotificationStore | None) -> None: """Inject the active ``NotificationStore`` instance for the process. Pass ``None`` to clear (e.g. on shutdown or in tests). """ global _store _store = store def get_notification_store() -> NotificationStore | None: return _store async def notify( user_id: str, *, type: str, title: str = "", body: str = "", payload: dict[str, Any] | None = None, ) -> None: """Write a notification for ``user_id``. Never raises — failures are logged. Args: user_id: Recipient. type: Stable identifier (e.g. ``agent_skill_removed``, ``skill_taken_down``). title: One-line title shown in the inbox. body: Longer description shown when the user expands the entry. payload: Structured context (skill_name, agent_id, actor_user_id, ...). """ try: if _store is None: logger.warning( "notify(user_id=%s, type=%s) dropped — NotificationStore not configured", user_id, type, ) return await _store.create( { "id": str(uuid.uuid4()), "user_id": user_id, "type": type, "title": title, "body": body, "payload": payload or {}, } ) except Exception: logger.exception("Failed to write notification (user_id=%s, type=%s)", user_id, type)