feat(server-go): 微信订阅消息骨架(#11,绑定/订阅/巡检触发,配置占位)
This commit is contained in:
@@ -0,0 +1,213 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
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 {
|
||||
token, err := s.getAccessToken(ctx)
|
||||
if err != nil {
|
||||
return 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 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 err
|
||||
}
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
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 {
|
||||
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("微信订阅消息发送失败 (%d): %s", out.ErrCode, out.ErrMsg)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,112 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestWechatTemplateKey(t *testing.T) {
|
||||
cases := map[string]string{
|
||||
"green": "",
|
||||
"yellow": "inspection",
|
||||
"orange": "inspection",
|
||||
"red": "inspection",
|
||||
"": "",
|
||||
"unknown": "",
|
||||
}
|
||||
for level, want := range cases {
|
||||
if got := WechatTemplateKey(level); got != want {
|
||||
t.Errorf("WechatTemplateKey(%q) = %q, want %q", level, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestIsAuthorized(t *testing.T) {
|
||||
auth := []string{"alarm", "inspection"}
|
||||
if !IsAuthorized(auth, "inspection") {
|
||||
t.Error("inspection 应在授权列表内")
|
||||
}
|
||||
if IsAuthorized(auth, "push") {
|
||||
t.Error("push 不应在授权列表内")
|
||||
}
|
||||
if IsAuthorized(nil, "inspection") {
|
||||
t.Error("空授权列表应返回 false")
|
||||
}
|
||||
}
|
||||
|
||||
func TestBuildSubscribeData(t *testing.T) {
|
||||
data := BuildSubscribeData("orange", 75)
|
||||
if data["thing1"]["value"] != "橙" {
|
||||
t.Errorf("thing1 应为橙色等级,实际 %v", data["thing1"])
|
||||
}
|
||||
if data["thing2"]["value"] != "75分" {
|
||||
t.Errorf("thing2 应为 75分,实际 %v", data["thing2"])
|
||||
}
|
||||
}
|
||||
|
||||
func TestWechatCode2Session(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Path != "/sns/jscode2session" {
|
||||
t.Errorf("path = %s", r.URL.Path)
|
||||
}
|
||||
q := r.URL.Query()
|
||||
if q.Get("appid") != "app1" || q.Get("secret") != "sec1" || q.Get("js_code") != "code123" {
|
||||
t.Errorf("参数不正确: %v", q)
|
||||
}
|
||||
_ = json.NewEncoder(w).Encode(map[string]string{"openid": "oAbC123", "session_key": "sk"})
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
s := NewWechatService("app1", "sec1")
|
||||
s.baseURL = srv.URL
|
||||
openid, err := s.Code2Session(context.Background(), "code123")
|
||||
if err != nil {
|
||||
t.Fatalf("Code2Session 错误: %v", err)
|
||||
}
|
||||
if openid != "oAbC123" {
|
||||
t.Errorf("openid = %s", openid)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWechatSendSubscribe(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.URL.Path != "/cgi-bin/message/subscribe/send" {
|
||||
t.Errorf("path = %s", r.URL.Path)
|
||||
}
|
||||
var body map[string]any
|
||||
_ = json.NewDecoder(r.Body).Decode(&body)
|
||||
if body["touser"] != "oAbC123" || body["template_id"] != "tmpl1" {
|
||||
t.Errorf("body 不正确: %v", body)
|
||||
}
|
||||
_, _ = w.Write([]byte(`{"errcode":0,"errmsg":"ok"}`))
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
s := NewWechatService("app1", "sec1")
|
||||
s.baseURL = srv.URL
|
||||
s.accessToken = "fake-token"
|
||||
s.tokenExpire = time.Now().Add(time.Hour)
|
||||
err := s.SendSubscribe(context.Background(), "oAbC123", "tmpl1", BuildSubscribeData("red", 88), "")
|
||||
if err != nil {
|
||||
t.Fatalf("SendSubscribe 错误: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWechatSendSubscribeErrorCode(t *testing.T) {
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
_, _ = w.Write([]byte(`{"errcode":40003,"errmsg":"invalid openid"}`))
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
s := NewWechatService("app1", "sec1")
|
||||
s.baseURL = srv.URL
|
||||
s.accessToken = "fake-token"
|
||||
s.tokenExpire = time.Now().Add(time.Hour)
|
||||
if err := s.SendSubscribe(context.Background(), "bad", "tmpl1", nil, ""); err == nil {
|
||||
t.Error("errcode!=0 应返回错误")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user