170 lines
4.7 KiB
Go
170 lines
4.7 KiB
Go
package notify
|
||
|
||
import (
|
||
"bytes"
|
||
"encoding/json"
|
||
"fmt"
|
||
"net/http"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
)
|
||
|
||
const (
|
||
// 获取 access_token 的 API 地址
|
||
getTokenURL = "https://qyapi.weixin.qq.com/cgi-bin/gettoken"
|
||
// 发送应用消息的 API 地址
|
||
sendMessageURL = "https://qyapi.weixin.qq.com/cgi-bin/message/send"
|
||
)
|
||
|
||
// wechatNotifier 实现了 Notifier 接口,用于通过企业微信自建应用发送私聊消息。
|
||
type wechatNotifier struct {
|
||
corpID string // 企业ID (CorpID)
|
||
agentID string // 应用ID (AgentID)
|
||
secret string // 应用密钥 (Secret)
|
||
|
||
// 用于线程安全地管理 access_token
|
||
mu sync.Mutex
|
||
accessToken string
|
||
tokenExpiresAt time.Time
|
||
}
|
||
|
||
// NewWechatNotifier 创建一个新的 wechatNotifier 实例。
|
||
// 调用者需要注入企业微信应用的 CorpID, AgentID 和 Secret。
|
||
func NewWechatNotifier(corpID, agentID, secret string) Notifier {
|
||
return &wechatNotifier{
|
||
corpID: corpID,
|
||
agentID: agentID,
|
||
secret: secret,
|
||
}
|
||
}
|
||
|
||
// Send 向指定用户发送一条 markdown 格式的私聊消息。
|
||
// toAddr 参数是接收者的 UserID 列表,用逗号或竖线分隔。
|
||
func (w *wechatNotifier) Send(content AlarmContent, toAddr string) error {
|
||
// 1. 获取有效的 access_token
|
||
token, err := w.getAccessToken()
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
// 2. 构建 markdown 内容
|
||
markdownContent := fmt.Sprintf("## %s\n> 级别: <font color=\"warning\">%s</font>\n> 时间: %s\n\n%s",
|
||
content.Title,
|
||
content.Level.String(),
|
||
content.Timestamp.Format(DefaultTimeFormat),
|
||
content.Message,
|
||
)
|
||
|
||
// 3. 构建请求的 JSON Body
|
||
// 将逗号分隔的 toAddr 转换为竖线分隔,以符合 API 要求
|
||
userList := strings.ReplaceAll(toAddr, ",", "|")
|
||
payload := wechatMessagePayload{
|
||
ToUser: userList,
|
||
MsgType: "markdown",
|
||
AgentID: w.agentID,
|
||
Markdown: struct {
|
||
Content string `json:"content"`
|
||
}{
|
||
Content: markdownContent,
|
||
},
|
||
}
|
||
|
||
jsonBytes, err := json.Marshal(payload)
|
||
if err != nil {
|
||
return fmt.Errorf("序列化企业微信消息失败: %w", err)
|
||
}
|
||
|
||
// 4. 发送 HTTP POST 请求
|
||
url := fmt.Sprintf("%s?access_token=%s", sendMessageURL, token)
|
||
req, err := http.NewRequest("POST", url, bytes.NewBuffer(jsonBytes))
|
||
if err != nil {
|
||
return fmt.Errorf("创建企业微信请求失败: %w", err)
|
||
}
|
||
req.Header.Set("Content-Type", "application/json")
|
||
|
||
client := &http.Client{}
|
||
resp, err := client.Do(req)
|
||
if err != nil {
|
||
return fmt.Errorf("发送企业微信通知失败: %w", err)
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
// 5. 检查响应
|
||
var response wechatResponse
|
||
if err := json.NewDecoder(resp.Body).Decode(&response); err != nil {
|
||
return fmt.Errorf("解析企业微信响应失败: %w", err)
|
||
}
|
||
|
||
if response.ErrCode != 0 {
|
||
return fmt.Errorf("企业微信API返回错误: code=%d, msg=%s", response.ErrCode, response.ErrMsg)
|
||
}
|
||
|
||
return nil
|
||
}
|
||
|
||
// getAccessToken 获取并缓存 access_token,处理了线程安全和自动刷新。
|
||
func (w *wechatNotifier) getAccessToken() (string, error) {
|
||
w.mu.Lock()
|
||
defer w.mu.Unlock()
|
||
|
||
// 如果 token 存在且有效期还有5分钟以上,则直接返回缓存的 token
|
||
if w.accessToken != "" && time.Now().Before(w.tokenExpiresAt.Add(-5*time.Minute)) {
|
||
return w.accessToken, nil
|
||
}
|
||
|
||
// 否则,重新获取 token
|
||
url := fmt.Sprintf("%s?corpid=%s&corpsecret=%s", getTokenURL, w.corpID, w.secret)
|
||
resp, err := http.Get(url)
|
||
if err != nil {
|
||
return "", fmt.Errorf("获取 access_token 请求失败: %w", err)
|
||
}
|
||
defer resp.Body.Close()
|
||
|
||
var tokenResp tokenResponse
|
||
if err := json.NewDecoder(resp.Body).Decode(&tokenResp); err != nil {
|
||
return "", fmt.Errorf("解析 access_token 响应失败: %w", err)
|
||
}
|
||
|
||
if tokenResp.ErrCode != 0 {
|
||
return "", fmt.Errorf("获取 access_token API 返回错误: code=%d, msg=%s", tokenResp.ErrCode, tokenResp.ErrMsg)
|
||
}
|
||
|
||
// 缓存新的 token 和过期时间
|
||
w.accessToken = tokenResp.AccessToken
|
||
w.tokenExpiresAt = time.Now().Add(time.Duration(tokenResp.ExpiresIn) * time.Second)
|
||
|
||
return w.accessToken, nil
|
||
}
|
||
|
||
// Type 返回通知器的类型
|
||
func (w *wechatNotifier) Type() NotifierType {
|
||
return NotifierTypeWeChat
|
||
}
|
||
|
||
// --- API 数据结构 ---
|
||
|
||
// tokenResponse 是获取 access_token API 的响应结构体
|
||
type tokenResponse struct {
|
||
ErrCode int `json:"errcode"`
|
||
ErrMsg string `json:"errmsg"`
|
||
AccessToken string `json:"access_token"`
|
||
ExpiresIn int `json:"expires_in"`
|
||
}
|
||
|
||
// wechatMessagePayload 是发送应用消息 API 的请求体结构
|
||
type wechatMessagePayload struct {
|
||
ToUser string `json:"touser"`
|
||
MsgType string `json:"msgtype"`
|
||
AgentID string `json:"agentid"`
|
||
Markdown struct {
|
||
Content string `json:"content"`
|
||
} `json:"markdown"`
|
||
}
|
||
|
||
// wechatResponse 是企业微信 API 的通用响应结构体
|
||
type wechatResponse struct {
|
||
ErrCode int `json:"errcode"`
|
||
ErrMsg string `json:"errmsg"`
|
||
}
|