import logging import json from api_v1.utils.wecom_webhook import send_wecom_webhook_message from notifier.models import NotifierChannelEnum from notifier.message_api import send_message_api_news_to_agents, send_message_api_text_message from notifier.serializers import ( MessageNewsRequestSerializer, MessageTemplatePayloadSerializer, MessageTextRequestSerializer, ) logger = logging.getLogger(__name__) class BaseNotifierBackend: channel = "" def notify(self, *, notifier, content: str, context: dict) -> dict: raise NotImplementedError class WeComWebhookNotifierBackend(BaseNotifierBackend): channel = NotifierChannelEnum.WECOM_WEBHOOK def notify(self, *, notifier, content: str, context: dict) -> dict: content = (content or "").strip() if not content: raise ValueError("企业微信机器人消息 content 不能为空") msgtype = str(notifier.get_config_value("msgtype", "markdown") or "markdown").strip().lower() timeout_seconds = float(notifier.get_config_value("timeout_seconds", 10.0) or 10.0) key = str(notifier.get_config_value("key", "") or "").strip() response = send_wecom_webhook_message( content=content, msgtype=msgtype, key=key or None, timeout_seconds=timeout_seconds, ) if not response.ok: raise RuntimeError(f"WeCom webhook 返回失败: errcode={response.errcode}, errmsg={response.errmsg}") result = { "channel": self.channel, "msgtype": msgtype, "errcode": response.errcode, "errmsg": response.errmsg, } logger.info("[notifier.backends] wecom webhook sent: notifier_id=%s result=%s", notifier.id, result) return result class MessageAPINotifierBackend(BaseNotifierBackend): channel = NotifierChannelEnum.MESSAGE_API def _load_rendered_payload(self, *, content: str) -> dict: content = (content or "").strip() if not content: raise ValueError("消息发送 API 模板渲染结果不能为空") try: payload = json.loads(content) except json.JSONDecodeError as exc: raise ValueError(f"消息发送 API 模板渲染结果不是合法 JSON: {exc}") from exc if not isinstance(payload, dict): raise ValueError("消息发送 API 模板渲染结果必须是 JSON 对象") serializer = MessageTemplatePayloadSerializer(data=payload) serializer.is_valid(raise_exception=True) return serializer.validated_data def notify(self, *, notifier, content: str, context: dict) -> dict: rendered_payload = self._load_rendered_payload(content=content) msgtype = rendered_payload["msgtype"] if msgtype == "text": payload = { "agent_id": rendered_payload.get("agent_id", notifier.get_config_value("agent_id")), "content": rendered_payload["content"], } resolved_user_ids = rendered_payload.get("user_ids") or notifier.get_config_value("user_ids") if resolved_user_ids: payload["user_ids"] = resolved_user_ids serializer = MessageTextRequestSerializer(data=payload) serializer.is_valid(raise_exception=True) data = serializer.validated_data response = send_message_api_text_message( agent_id=data["agent_id"], content=data["content"], user_ids=data.get("user_ids"), timeout_seconds=float(notifier.get_config_value("timeout_seconds", 10.0) or 10.0), ) result = { "channel": self.channel, "msgtype": "text", "agent_id": data["agent_id"], "errcode": response.errcode, "errmsg": response.errmsg, } logger.info("[notifier.backends] message api text sent: notifier_id=%s result=%s", notifier.id, result) return result payload = { "agent_ids": rendered_payload.get("agent_ids", notifier.get_config_value("agent_ids")), "title": rendered_payload["title"], "description": rendered_payload["description"], "url": rendered_payload["url"], "image_url": rendered_payload["image_url"], } resolved_user_ids = rendered_payload.get("user_ids") or notifier.get_config_value("user_ids") if resolved_user_ids: payload["user_ids"] = resolved_user_ids serializer = MessageNewsRequestSerializer(data=payload) serializer.is_valid(raise_exception=True) data = serializer.validated_data results = send_message_api_news_to_agents( agent_ids=data["agent_ids"], title=data["title"], description=data["description"], url=data["url"], image_url=data["image_url"], user_ids=data.get("user_ids"), timeout_seconds=float(notifier.get_config_value("timeout_seconds", 10.0) or 10.0), ) if not any(item.get("ok") for item in results): raise RuntimeError("All agent sends failed") result = { "channel": self.channel, "msgtype": "news", "agent_ids": data["agent_ids"], "ok_count": sum(1 for item in results if item.get("ok")), "total_count": len(results), "results": results, } logger.info("[notifier.backends] message api news sent: notifier_id=%s result=%s", notifier.id, result) return result