deerflow-code/offline-backend-20260512/backend/packages/harness/deerflow/persistence/notifications/sender.py
2026-09-07 18:24:55 +08:00

73 lines
2.1 KiB
Python

"""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)