forked from erp-dev/erp
135 lines
5.5 KiB
Python
135 lines
5.5 KiB
Python
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
|