package service import ( "bytes" "context" "encoding/json" "fmt" "io" "net/http" "net/url" "sync" "time" "silk-server-go/internal/model" "gorm.io/gorm" "gorm.io/gorm/clause" ) const wechatAPIBase = "https://api.weixin.qq.com" // WechatTemplateKey 风险等级 → 订阅消息场景键(green 不推送) func WechatTemplateKey(level string) string { switch level { case "yellow", "orange", "red": return "inspection" default: return "" } } // IsAuthorized 判断授权模板列表是否包含场景键 func IsAuthorized(authorized []string, key string) bool { for _, k := range authorized { if k == key { return true } } return false } // levelZh 风险等级中文 func levelZh(level string) string { switch level { case "green": return "绿" case "yellow": return "黄" case "orange": return "橙" case "red": return "红" default: return level } } // BuildSubscribeData 构造订阅消息 data(thing1=风险等级 thing2=评分;字段名需与微信模板字段一致) func BuildSubscribeData(level string, score float64) map[string]map[string]string { return map[string]map[string]string{ "thing1": {"value": levelZh(level)}, "thing2": {"value": fmt.Sprintf("%.0f分", score)}, } } // WechatService 微信小程序订阅消息服务(骨架;未配置时 Configured()=false,调用返回明确错误) type WechatService struct { appID string secret string baseURL string accessToken string tokenExpire time.Time mu sync.Mutex httpClient *http.Client } // NewWechatService 创建微信服务 func NewWechatService(appID, secret string) *WechatService { return &WechatService{ appID: appID, secret: secret, baseURL: wechatAPIBase, httpClient: &http.Client{Timeout: 15 * time.Second}, } } // Configured 是否已配置 AppID/Secret func (s *WechatService) Configured() bool { return s.appID != "" && s.secret != "" } // Code2Session 用 wx.login 的 code 换取 openid func (s *WechatService) Code2Session(ctx context.Context, code string) (string, error) { if !s.Configured() { return "", fmt.Errorf("微信未配置(WECHAT_APPID/WECHAT_SECRET)") } u := s.baseURL + "/sns/jscode2session?appid=" + url.QueryEscape(s.appID) + "&secret=" + url.QueryEscape(s.secret) + "&js_code=" + url.QueryEscape(code) + "&grant_type=authorization_code" req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) if err != nil { return "", err } resp, err := s.httpClient.Do(req) if err != nil { return "", err } defer resp.Body.Close() body, err := io.ReadAll(resp.Body) if err != nil { return "", err } var out struct { OpenID string `json:"openid"` ErrCode int `json:"errcode"` ErrMsg string `json:"errmsg"` } if err := json.Unmarshal(body, &out); err != nil { return "", err } if out.ErrCode != 0 { return "", fmt.Errorf("微信 code2session 失败 (%d): %s", out.ErrCode, out.ErrMsg) } return out.OpenID, nil } // getAccessToken 获取并缓存 access_token(提前 60 秒过期) func (s *WechatService) getAccessToken(ctx context.Context) (string, error) { s.mu.Lock() if s.accessToken != "" && time.Now().Before(s.tokenExpire) { tok := s.accessToken s.mu.Unlock() return tok, nil } s.mu.Unlock() if !s.Configured() { return "", fmt.Errorf("微信未配置(WECHAT_APPID/WECHAT_SECRET)") } u := s.baseURL + "/cgi-bin/token?grant_type=client_credential&appid=" + url.QueryEscape(s.appID) + "&secret=" + url.QueryEscape(s.secret) req, err := http.NewRequestWithContext(ctx, http.MethodGet, u, nil) if err != nil { return "", err } resp, err := s.httpClient.Do(req) if err != nil { return "", err } defer resp.Body.Close() body, err := io.ReadAll(resp.Body) if err != nil { return "", err } var out struct { AccessToken string `json:"access_token"` ExpiresIn int `json:"expires_in"` ErrCode int `json:"errcode"` ErrMsg string `json:"errmsg"` } if err := json.Unmarshal(body, &out); err != nil { return "", err } if out.AccessToken == "" { return "", fmt.Errorf("微信 access_token 获取失败 (%d): %s", out.ErrCode, out.ErrMsg) } s.mu.Lock() s.accessToken = out.AccessToken s.tokenExpire = time.Now().Add(time.Duration(out.ExpiresIn-60) * time.Second) s.mu.Unlock() return out.AccessToken, nil } // SendSubscribe 发送订阅消息 func (s *WechatService) SendSubscribe(ctx context.Context, openid, templateID string, data map[string]map[string]string, page string) error { _, err := s.SendSubscribeResult(ctx, openid, templateID, data, page) return err } // WechatSendResult 微信订阅消息响应。 type WechatSendResult struct { ErrCode int `json:"errcode"` ErrMsg string `json:"errmsg"` } // SendSubscribeResult 发送订阅消息并返回微信响应码。 func (s *WechatService) SendSubscribeResult(ctx context.Context, openid, templateID string, data map[string]map[string]string, page string) (WechatSendResult, error) { token, err := s.getAccessToken(ctx) if err != nil { return WechatSendResult{ErrCode: -1, ErrMsg: err.Error()}, err } payload := map[string]any{ "touser": openid, "template_id": templateID, "data": data, } if page != "" { payload["page"] = page } raw, err := json.Marshal(payload) if err != nil { return WechatSendResult{ErrCode: -1, ErrMsg: err.Error()}, err } u := s.baseURL + "/cgi-bin/message/subscribe/send?access_token=" + url.QueryEscape(token) req, err := http.NewRequestWithContext(ctx, http.MethodPost, u, bytes.NewReader(raw)) if err != nil { return WechatSendResult{ErrCode: -1, ErrMsg: err.Error()}, err } req.Header.Set("Content-Type", "application/json") resp, err := s.httpClient.Do(req) if err != nil { return WechatSendResult{ErrCode: -1, ErrMsg: err.Error()}, err } defer resp.Body.Close() body, err := io.ReadAll(resp.Body) if err != nil { return WechatSendResult{ErrCode: -1, ErrMsg: err.Error()}, err } var out WechatSendResult if err := json.Unmarshal(body, &out); err != nil { return WechatSendResult{ErrCode: -1, ErrMsg: err.Error()}, err } if out.ErrCode != 0 { return out, fmt.Errorf("微信订阅消息发送失败 (%d): %s", out.ErrCode, out.ErrMsg) } return out, nil } // WechatSubscribePayload outbox 微信订阅事件载荷。 type WechatSubscribePayload struct { UserID string `json:"userId"` OpenID string `json:"openId"` TemplateKey string `json:"templateKey"` TemplateID string `json:"templateId"` Page string `json:"page"` Title string `json:"title"` Body string `json:"body"` Data map[string]map[string]string `json:"data"` BusinessType string `json:"businessType"` BusinessID string `json:"businessId"` } // NewWechatOutboxHandler 创建微信订阅 outbox 处理器。 func NewWechatOutboxHandler(db *gorm.DB, wechat *WechatService) EventHandler { return func(ctx context.Context, event model.OutboxEvent) error { if event.EventType != OutboxEventWechatSubscribe { return nil } var payload WechatSubscribePayload if err := json.Unmarshal(event.Payload, &payload); err != nil { return err } return handleWechatSubscribe(ctx, db, wechat, event, payload) } } func handleWechatSubscribe(ctx context.Context, db *gorm.DB, wechat *WechatService, event model.OutboxEvent, payload WechatSubscribePayload) error { var binding model.WechatBinding if err := db.Where("user_id = ?", payload.UserID).First(&binding).Error; err != nil { message := "微信未绑定" return saveWechatNotification(ctx, db, event, payload, "", false, WechatSendResult{ErrCode: 40001, ErrMsg: message}, "cancelled", &message) } var authorized []string if len(binding.AuthorizedTemplates) > 0 { _ = json.Unmarshal(binding.AuthorizedTemplates, &authorized) } if !IsAuthorized(authorized, payload.TemplateKey) { message := "用户未授权订阅模板" return saveWechatNotification(ctx, db, event, payload, binding.OpenID, false, WechatSendResult{ErrCode: 43101, ErrMsg: message}, "cancelled", &message) } if payload.TemplateID == "" { message := "订阅消息模板 ID 未配置" if err := saveWechatNotification(ctx, db, event, payload, binding.OpenID, true, WechatSendResult{ErrCode: -1, ErrMsg: message}, "retry", &message); err != nil { return err } return fmt.Errorf("订阅消息模板 ID 未配置") } result, sendErr := wechat.SendSubscribeResult(ctx, binding.OpenID, payload.TemplateID, payload.Data, payload.Page) status := "sent" var lastErr *string maxAttempts := event.MaxAttempts if maxAttempts <= 0 { maxAttempts = 5 } if sendErr != nil { message := sendErr.Error() lastErr = &message status = "retry" if event.Attempts+1 >= maxAttempts { status = "failed" } } if err := saveWechatNotification(ctx, db, event, payload, binding.OpenID, true, result, status, lastErr); err != nil { return err } return sendErr } func saveWechatNotification(ctx context.Context, db *gorm.DB, event model.OutboxEvent, payload WechatSubscribePayload, openID string, authorized bool, response WechatSendResult, status string, lastErr *string) error { data, _ := json.Marshal(payload.Data) providerResponse, _ := json.Marshal(response) userID := payload.UserID eventID := event.EventID notification := model.Notification{ UserID: &userID, Channel: "wechat", Target: openID, Title: payload.Title, Body: payload.Body, Status: status, Attempts: event.Attempts + 1, NextAttemptAt: time.Now(), LastError: lastErr, TemplateID: strPtrOrNil(payload.TemplateID), OpenID: strPtrOrNil(openID), Authorized: authorized, ProviderResponse: providerResponse, Data: data, BusinessType: strPtrOrNil(payload.BusinessType), BusinessID: strPtrOrNil(payload.BusinessID), EventID: &eventID, } return db.WithContext(ctx). Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "event_id"}}, DoUpdates: clause.Assignments(map[string]interface{}{ "status": status, "attempts": event.Attempts + 1, "next_attempt_at": time.Now(), "last_error": lastErr, "authorized": authorized, "provider_response": providerResponse, "updated_at": time.Now(), }), }). Create(¬ification).Error }