增加微信通知功能

This commit is contained in:
2026-09-16 18:06:18 +08:00
parent b48688cbab
commit f386899625
23 changed files with 3539 additions and 18 deletions
+282
View File
@@ -0,0 +1,282 @@
package wechatmp
import (
"bytes"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"time"
)
// =============================================================
// 微信公众平台服务端 API 调用(access_token / 带参二维码 / 模板消息 / 订阅通知 / 用户信息)
// =============================================================
const apiBase = "https://api.weixin.qq.com"
var httpClient = &http.Client{Timeout: 10 * time.Second}
// APIError 微信接口业务错误(HTTP 200 但 errcode 非 0)。
// Error() 输出保留 errcode / errmsg 原始信息,便于日志检索与前端展示。
type APIError struct {
API string // 调用的接口(如「获取 access_token」)
Code int
Msg string
}
func (e *APIError) Error() string {
return fmt.Sprintf("%s 失败:errcode=%d errmsg=%s", e.API, e.Code, e.Msg)
}
// Explain 返回常见 errcode 的中文排查建议;未知错误码返回空串。
func (e *APIError) Explain() string {
switch e.Code {
case 40013:
return "AppID 无效:请核对微信后台「设置与开发 → 基本配置」中的 AppID"
case 40125:
return "AppSecret 错误:请在微信后台重置后复制新的 AppSecret 填入本页"
case 40164, 61004:
return "服务器出口 IP 不在白名单:请在微信后台「基本配置 → IP白名单」加入本服务器出口 IP(服务器执行 curl ifconfig.me 可查看)"
case 45009:
return "接口调用频次超限:请稍后重试"
case 48001:
return "接口未获授权:需认证服务号,或该账号未开通此接口能力"
}
return ""
}
// ExplainError 提取微信接口错误的中文排查建议;无可用建议时返回 ("", false)。
func ExplainError(err error) (string, bool) {
var ae *APIError
if errors.As(err, &ae) {
if hint := ae.Explain(); hint != "" {
return hint, true
}
}
return "", false
}
// GetAccessToken 获取 access_token(内存缓存,提前 5 分钟刷新)
func GetAccessToken(cfg *Config) (string, error) {
if cfg == nil || cfg.AppID == "" || cfg.AppSecret == "" {
return "", fmt.Errorf("AppID / AppSecret 未配置")
}
tokenMu.Lock()
if tokenCache.appID == cfg.AppID && tokenCache.token != "" && time.Now().Before(tokenCache.expiresAt) {
t := tokenCache.token
tokenMu.Unlock()
return t, nil
}
tokenMu.Unlock()
api := fmt.Sprintf("%s/cgi-bin/token?grant_type=client_credential&appid=%s&secret=%s",
apiBase, url.QueryEscape(cfg.AppID), url.QueryEscape(cfg.AppSecret))
var resp struct {
AccessToken string `json:"access_token"`
ExpiresIn int `json:"expires_in"`
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
}
if err := httpGetJSON(api, &resp); err != nil {
return "", err
}
if resp.AccessToken == "" {
return "", &APIError{API: "获取 access_token", Code: resp.ErrCode, Msg: resp.ErrMsg}
}
ttl := time.Duration(resp.ExpiresIn) * time.Second
if ttl <= tokenMinTTL {
ttl = time.Hour
}
tokenMu.Lock()
tokenCache = tokenEntry{appID: cfg.AppID, token: resp.AccessToken, expiresAt: time.Now().Add(ttl - tokenMinTTL)}
tokenMu.Unlock()
return resp.AccessToken, nil
}
// CreateQRCode 创建带参数临时二维码(字符串场景值),返回二维码图片地址。
// scene 会通过「关注事件的 EventKey(qrscene_<scene>)/ 扫码事件的 EventKey(<scene>)」回传。
func CreateQRCode(cfg *Config, scene string, expireSeconds int) (string, error) {
token, err := GetAccessToken(cfg)
if err != nil {
return "", err
}
if expireSeconds <= 0 || expireSeconds > 2592000 {
expireSeconds = 300
}
payload := map[string]interface{}{
"expire_seconds": expireSeconds,
"action_name": "QR_STR_SCENE",
"action_info": map[string]interface{}{
"scene": map[string]string{"scene_str": scene},
},
}
api := fmt.Sprintf("%s/cgi-bin/qrcode/create?access_token=%s", apiBase, url.QueryEscape(token))
var resp struct {
Ticket string `json:"ticket"`
ExpireSeconds int `json:"expire_seconds"`
URL string `json:"url"`
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
}
if err := httpPostJSON(api, payload, &resp); err != nil {
return "", err
}
if resp.Ticket == "" {
return "", &APIError{API: "创建带参二维码", Code: resp.ErrCode, Msg: resp.ErrMsg}
}
return "https://mp.weixin.qq.com/cgi-bin/showqrcode?ticket=" + url.QueryEscape(resp.Ticket), nil
}
// SendTemplateMessage 发送模板消息(需认证服务号 + 模板ID)
// data 键为模板字段名(first / keyword1 / remark ...),值为文本内容。
func SendTemplateMessage(cfg *Config, openid, jumpURL string, data map[string]string) error {
if cfg.TemplateID == "" {
return fmt.Errorf("未配置模板消息 ID(认证后在公众号后台申请)")
}
token, err := GetAccessToken(cfg)
if err != nil {
return err
}
fields := make(map[string]map[string]string, len(data))
for k, v := range data {
fields[k] = map[string]string{"value": v}
}
payload := map[string]interface{}{
"touser": openid,
"template_id": cfg.TemplateID,
"data": fields,
}
if jumpURL != "" {
payload["url"] = jumpURL
}
api := fmt.Sprintf("%s/cgi-bin/message/template/send?access_token=%s", apiBase, url.QueryEscape(token))
var resp struct {
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
MsgID int64 `json:"msgid"`
}
if err := httpPostJSON(api, payload, &resp); err != nil {
return err
}
if resp.ErrCode != 0 {
return &APIError{API: "发送模板消息", Code: resp.ErrCode, Msg: resp.ErrMsg}
}
return nil
}
// SendSubscribeMessage 发送订阅通知(一次性订阅,需用户先在页面完成订阅授权)
func SendSubscribeMessage(cfg *Config, openid, templateID, page string, data map[string]string) error {
if templateID == "" {
return fmt.Errorf("未配置订阅通知模板 ID")
}
token, err := GetAccessToken(cfg)
if err != nil {
return err
}
fields := make(map[string]map[string]string, len(data))
for k, v := range data {
fields[k] = map[string]string{"value": v}
}
payload := map[string]interface{}{
"touser": openid,
"template_id": templateID,
"data": fields,
}
if page != "" {
payload["page"] = page
}
api := fmt.Sprintf("%s/cgi-bin/message/subscribe/bizsend?access_token=%s", apiBase, url.QueryEscape(token))
var resp struct {
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
}
if err := httpPostJSON(api, payload, &resp); err != nil {
return err
}
if resp.ErrCode != 0 {
return &APIError{API: "发送订阅通知", Code: resp.ErrCode, Msg: resp.ErrMsg}
}
return nil
}
// UserInfo 用户基本信息(昵称/头像等;未认证公众号可能接口受限,调用失败由上层忽略)
type UserInfo struct {
OpenID string `json:"openid"`
Nickname string `json:"nickname"`
Sex int `json:"sex"`
City string `json:"city"`
Province string `json:"province"`
Country string `json:"country"`
HeadImgURL string `json:"headimgurl"`
Subscribe int `json:"subscribe"`
}
// FetchUserInfo 拉取用户基本信息(失败不算错,由调用方容错)
func FetchUserInfo(cfg *Config, openid string) (*UserInfo, error) {
token, err := GetAccessToken(cfg)
if err != nil {
return nil, err
}
api := fmt.Sprintf("%s/cgi-bin/user/info?access_token=%s&openid=%s&lang=zh_CN",
apiBase, url.QueryEscape(token), url.QueryEscape(openid))
var resp struct {
UserInfo
ErrCode int `json:"errcode"`
ErrMsg string `json:"errmsg"`
}
if err := httpGetJSON(api, &resp); err != nil {
return nil, err
}
if resp.ErrCode != 0 {
return nil, &APIError{API: "获取用户信息", Code: resp.ErrCode, Msg: resp.ErrMsg}
}
return &resp.UserInfo, nil
}
// ============================ HTTP 辅助 ============================
func httpGetJSON(api string, out interface{}) error {
resp, err := httpClient.Get(api)
if err != nil {
return fmt.Errorf("请求微信接口失败: %w", err)
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return fmt.Errorf("读取微信响应失败: %w", err)
}
if err := json.Unmarshal(body, out); err != nil {
return fmt.Errorf("解析微信响应失败: %w(body: %s)", err, truncate(string(body), 200))
}
return nil
}
func httpPostJSON(api string, payload interface{}, out interface{}) error {
bs, err := json.Marshal(payload)
if err != nil {
return err
}
resp, err := httpClient.Post(api, "application/json; charset=utf-8", bytes.NewReader(bs))
if err != nil {
return fmt.Errorf("请求微信接口失败: %w", err)
}
defer resp.Body.Close()
body, err := io.ReadAll(resp.Body)
if err != nil {
return fmt.Errorf("读取微信响应失败: %w", err)
}
if err := json.Unmarshal(body, out); err != nil {
return fmt.Errorf("解析微信响应失败: %w(body: %s)", err, truncate(string(body), 200))
}
return nil
}
func truncate(s string, n int) string {
if len(s) <= n {
return s
}
return s[:n] + "..."
}
+351
View File
@@ -0,0 +1,351 @@
package wechatmp
import (
"crypto/rand"
"encoding/hex"
"errors"
"fmt"
"math/big"
"strings"
"time"
"server/models"
"github.com/beego/beego/v2/client/orm"
)
// =============================================================
// 关注绑定:扫码关注 → 公众号被动回复验证码 → 平台核销完成绑定
// 一个绑定会话 = 一条 yz_platform_wechat_mp_verify_code(scene + 6位验证码 + 过期时间)
// =============================================================
// DefaultVerifyTTL 验证码默认有效期
const DefaultVerifyTTL = 5 * time.Minute
var (
ErrSceneNotScanned = errors.New("尚未检测到扫码关注,请先使用微信扫码并关注公众号")
ErrCodeUsed = errors.New("该验证码已使用,请重新获取二维码")
ErrCodeExpired = errors.New("二维码已过期,请重新获取")
ErrCodeMismatch = errors.New("验证码不正确")
ErrSceneInvalid = errors.New("二维码无效或已失效")
)
// CreateVerifyCode 创建绑定会话并生成带参二维码。
// bindType 为 platform_user / tenant_user,bindID 为发起方用户 ID。
func CreateVerifyCode(bindType string, bindID, bindTid uint64, ttl time.Duration) (scene, code, qrURL string, expireAt time.Time, err error) {
cfg, err := LoadEnabledConfig()
if err != nil {
return "", "", "", time.Time{}, err
}
if ttl <= 0 {
ttl = DefaultVerifyTTL
}
scene = randomScene()
code = randomCode6()
now := time.Now()
expireAt = now.Add(ttl)
row := &models.WechatMpVerifyCode{
Scene: scene,
Code: code,
Status: models.WechatVerifyStatusWaiting,
BindType: bindType,
BindID: bindID,
BindTid: bindTid,
ExpireAt: &expireAt,
CreateTime: now,
UpdateTime: &now,
}
if _, err = models.Orm.Insert(row); err != nil {
return "", "", "", time.Time{}, fmt.Errorf("创建绑定会话失败: %w", err)
}
qrURL, err = CreateQRCode(cfg, scene, int(ttl.Seconds())+120)
if err != nil {
return "", "", "", time.Time{}, err
}
return scene, code, qrURL, expireAt, nil
}
// GetVerifyCode 按 scene 查询绑定会话(过期时惰性标记)
func GetVerifyCode(scene string) (*models.WechatMpVerifyCode, error) {
scene = strings.TrimSpace(scene)
if scene == "" {
return nil, ErrSceneInvalid
}
var row models.WechatMpVerifyCode
if err := models.Orm.QueryTable(new(models.WechatMpVerifyCode)).Filter("scene", scene).One(&row); err != nil {
return nil, ErrSceneInvalid
}
if row.Status != models.WechatVerifyStatusUsed && row.ExpireAt != nil && row.ExpireAt.Before(time.Now()) {
row.Status = models.WechatVerifyStatusExpired
now := time.Now()
_, _ = models.Orm.QueryTable(new(models.WechatMpVerifyCode)).
Filter("id", row.ID).
Update(map[string]interface{}{"status": row.Status, "update_time": now})
}
return &row, nil
}
// ConfirmVerifyCode 核销验证码并完成绑定,返回粉丝记录
func ConfirmVerifyCode(bindType string, bindID, bindTid uint64, scene, code string) (*models.WechatMpFollower, error) {
row, err := GetVerifyCode(scene)
if err != nil {
return nil, err
}
if row.BindType != bindType || row.BindID != bindID {
return nil, errors.New("该二维码不是当前账号发起,请重新获取")
}
switch row.Status {
case models.WechatVerifyStatusWaiting:
return nil, ErrSceneNotScanned
case models.WechatVerifyStatusUsed:
return nil, ErrCodeUsed
case models.WechatVerifyStatusExpired:
return nil, ErrCodeExpired
}
if strings.TrimSpace(code) != row.Code {
return nil, ErrCodeMismatch
}
now := time.Now()
if _, err := models.Orm.QueryTable(new(models.WechatMpVerifyCode)).
Filter("id", row.ID).
Update(map[string]interface{}{"status": models.WechatVerifyStatusUsed, "update_time": now}); err != nil {
return nil, err
}
return BindFollower(row.OpenID, bindType, bindID, bindTid)
}
// BindFollower 将 openid 绑定到指定账号。
// 同一账号只保留一个微信:旧的 openid 绑定关系会被清除。
func BindFollower(openid, bindType string, bindID, bindTid uint64) (*models.WechatMpFollower, error) {
openid = strings.TrimSpace(openid)
if openid == "" {
return nil, errors.New("openid 为空")
}
now := time.Now()
_, _ = models.Orm.QueryTable(new(models.WechatMpFollower)).
Filter("bind_type", bindType).
Filter("bind_id", bindID).
Filter("openid__ne", openid).
Update(map[string]interface{}{
"bind_type": "", "bind_id": 0, "bind_tid": 0, "update_time": now,
})
var row models.WechatMpFollower
err := models.Orm.QueryTable(new(models.WechatMpFollower)).Filter("openid", openid).One(&row)
if err == orm.ErrNoRows {
row = models.WechatMpFollower{
OpenID: openid,
Subscribe: 1,
SubscribeTime: &now,
CreateTime: now,
}
}
row.BindType = bindType
row.BindID = bindID
row.BindTid = bindTid
row.BindTime = &now
row.UpdateTime = &now
if row.ID == 0 {
id, ierr := models.Orm.Insert(&row)
if ierr != nil {
return nil, ierr
}
row.ID = uint64(id)
return &row, nil
}
if _, uerr := models.Orm.Update(&row, "BindType", "BindID", "BindTid", "BindTime", "UpdateTime"); uerr != nil {
return nil, uerr
}
return &row, nil
}
// UnbindByUser 清除某账号的微信绑定(按绑定对象维度)
func UnbindByUser(bindType string, bindID uint64) error {
if bindID == 0 {
return nil
}
now := time.Now()
_, err := models.Orm.QueryTable(new(models.WechatMpFollower)).
Filter("bind_type", bindType).
Filter("bind_id", bindID).
Update(map[string]interface{}{
"bind_type": "", "bind_id": 0, "bind_tid": 0, "update_time": now,
})
return err
}
// GetFollowerByUser 查询某账号已绑定的粉丝记录(未绑定返回 nil, nil)
func GetFollowerByUser(bindType string, bindID uint64) (*models.WechatMpFollower, error) {
if bindID == 0 {
return nil, nil
}
var row models.WechatMpFollower
err := models.Orm.QueryTable(new(models.WechatMpFollower)).
Filter("bind_type", bindType).
Filter("bind_id", bindID).
Filter("delete_time__isnull", true).
OrderBy("-id").
One(&row)
if err == orm.ErrNoRows {
return nil, nil
}
if err != nil {
return nil, err
}
return &row, nil
}
// UpsertSubscribe 关注事件:创建/恢复粉丝记录
func UpsertSubscribe(openid string, at time.Time) (*models.WechatMpFollower, error) {
openid = strings.TrimSpace(openid)
if openid == "" {
return nil, errors.New("openid 为空")
}
var row models.WechatMpFollower
err := models.Orm.QueryTable(new(models.WechatMpFollower)).Filter("openid", openid).One(&row)
if err == orm.ErrNoRows {
row = models.WechatMpFollower{OpenID: openid, Subscribe: 1, SubscribeTime: &at, CreateTime: at, UpdateTime: &at}
id, ierr := models.Orm.Insert(&row)
if ierr != nil {
return nil, ierr
}
row.ID = uint64(id)
return &row, nil
}
if err != nil {
return nil, err
}
row.Subscribe = 1
row.SubscribeTime = &at
row.UnsubscribeTime = nil
row.UpdateTime = &at
if _, err := models.Orm.Update(&row, "Subscribe", "SubscribeTime", "UnsubscribeTime", "UpdateTime"); err != nil {
return nil, err
}
return &row, nil
}
// MarkUnsubscribe 取关事件
func MarkUnsubscribe(openid string, at time.Time) error {
openid = strings.TrimSpace(openid)
if openid == "" {
return nil
}
_, err := models.Orm.QueryTable(new(models.WechatMpFollower)).
Filter("openid", openid).
Update(map[string]interface{}{
"subscribe": 0, "unsubscribe_time": at, "update_time": at,
})
return err
}
// RefreshFollowerProfile 尝试拉取粉丝昵称/头像(未认证公众号可能受限,失败静默忽略)
func RefreshFollowerProfile(cfg *Config, openid string) {
if cfg == nil || openid == "" {
return
}
info, err := FetchUserInfo(cfg, openid)
if err != nil || info == nil {
return
}
now := time.Now()
up := map[string]interface{}{
"sex": info.Sex,
"city": info.City,
"province": info.Province,
"country": info.Country,
"update_time": now,
}
if info.Nickname != "" {
up["nickname"] = info.Nickname
}
if info.HeadImgURL != "" {
up["avatar"] = info.HeadImgURL
}
_, _ = models.Orm.QueryTable(new(models.WechatMpFollower)).Filter("openid", openid).Update(up)
}
// HandleInbound 处理入站消息/事件,返回需要被动回复的文本(空字符串表示不回复)
func HandleInbound(cfg *Config, msg *InboundMessage) string {
if msg == nil {
return ""
}
openid := strings.TrimSpace(msg.FromUserName)
now := time.Now()
switch strings.ToLower(strings.TrimSpace(msg.MsgType)) {
case "event":
switch strings.ToUpper(strings.TrimSpace(msg.Event)) {
case "subscribe":
_, _ = UpsertSubscribe(openid, now)
scene := strings.TrimPrefix(strings.TrimSpace(msg.EventKey), "qrscene_")
if scene != "" {
return handleScanScene(cfg, scene, openid)
}
return "欢迎关注云泽平台!\n如需接收系统提醒、公告推送,请回到系统页面点击「绑定微信」,使用微信扫码获取专属验证码。"
case "scan":
return handleScanScene(cfg, strings.TrimSpace(msg.EventKey), openid)
case "unsubscribe":
_ = MarkUnsubscribe(openid, now)
return ""
}
case "text":
return "如需绑定账号:请回到系统页面点击「绑定微信」,展示专属二维码后使用微信扫码,公众号会自动回复验证码。\n绑定成功后即可接收系统提醒、公告推送。"
}
return ""
}
// handleScanScene 扫码(关注/已关注)后的统一处理:写入验证码会话 + 回复验证码
func handleScanScene(cfg *Config, scene, openid string) string {
if strings.TrimSpace(scene) == "" {
return "未能识别二维码信息,请在系统页面重新获取二维码后再扫码。"
}
row, err := GetVerifyCode(scene)
if err != nil {
return "二维码无效或已失效,请在系统页面重新获取。"
}
switch row.Status {
case models.WechatVerifyStatusUsed:
return "该验证码已使用,请在系统页面重新获取二维码。"
case models.WechatVerifyStatusExpired:
return "二维码已过期,请在系统页面重新获取。"
}
now := time.Now()
if _, uerr := models.Orm.QueryTable(new(models.WechatMpVerifyCode)).
Filter("id", row.ID).
Update(map[string]interface{}{
"openid": openid, "status": models.WechatVerifyStatusScanned, "update_time": now,
}); uerr != nil {
return "系统繁忙,请稍后重新扫码。"
}
// 拉取昵称头像(失败不影响绑定流程)
RefreshFollowerProfile(cfg, openid)
minutes := int(time.Until(*row.ExpireAt).Minutes())
if minutes < 1 {
minutes = 1
}
return fmt.Sprintf("您的验证码是:%s\n请回到系统页面输入完成绑定(%d 分钟内有效)。", row.Code, minutes)
}
// ============================ 随机值 ============================
func randomScene() string {
b := make([]byte, 16)
if _, err := rand.Read(b); err != nil {
return fmt.Sprintf("s%d", time.Now().UnixNano())
}
return hex.EncodeToString(b)
}
func randomCode6() string {
n, err := rand.Int(rand.Reader, big.NewInt(900000))
if err != nil {
return "123456"
}
return fmt.Sprintf("%06d", n.Int64()+100000)
}
+207
View File
@@ -0,0 +1,207 @@
package wechatmp
import (
"errors"
"strings"
"sync"
"time"
"server/models"
"github.com/beego/beego/v2/client/orm"
beego "github.com/beego/beego/v2/server/web"
)
// =============================================================
// 微信公众号(服务号)配置与 access_token
// =============================================================
// 消息加解密方式
const (
EncryptModePlain = "plain" // 明文模式
EncryptModeCompatible = "compatible" // 兼容模式
EncryptModeSafe = "safe" // 安全模式
)
// Config 服务号运行配置(敏感字段已解密)
type Config struct {
ID uint64
AppID string
AppSecret string // 明文
Token string
AESKey string // 明文(43 位 EncodingAESKey)
EncryptMode string
TemplateID string
Verified bool
Enabled bool
Remark string
}
// ErrNotConfigured 未配置或未启用
var ErrNotConfigured = errors.New("微信公众号未配置或未启用")
// LoadConfig 读取服务号配置(单行,敏感字段解密);未配置返回 ErrNotConfigured
func LoadConfig() (*Config, error) {
var row models.WechatMpConfig
err := models.Orm.QueryTable(new(models.WechatMpConfig)).OrderBy("id").One(&row)
if err != nil {
if err == orm.ErrNoRows {
return nil, ErrNotConfigured
}
return nil, err
}
cfg := &Config{
ID: row.ID,
AppID: strings.TrimSpace(row.AppID),
Token: strings.TrimSpace(row.Token),
EncryptMode: strings.TrimSpace(row.EncryptMode),
TemplateID: strings.TrimSpace(row.TemplateID),
Verified: row.Verified == 1,
Enabled: row.Enabled == 1,
Remark: row.Remark,
}
if cfg.EncryptMode == "" {
cfg.EncryptMode = EncryptModePlain
}
if row.AppSecret != "" {
if plain, derr := DecryptSecret(row.AppSecret); derr == nil {
cfg.AppSecret = plain
} else {
return nil, derr
}
}
if row.AESKey != "" {
if plain, derr := DecryptSecret(row.AESKey); derr == nil {
cfg.AESKey = plain
} else {
return nil, derr
}
}
return cfg, nil
}
// LoadEnabledConfig 读取并校验启用状态
func LoadEnabledConfig() (*Config, error) {
cfg, err := LoadConfig()
if err != nil {
return nil, err
}
if !cfg.Enabled || cfg.AppID == "" || cfg.AppSecret == "" {
return nil, ErrNotConfigured
}
return cfg, nil
}
// SaveInput 保存配置入参(敏感字段支持掩码表示不改)
type SaveInput struct {
AppID string
AppSecret string // 掩码 / 空串 = 保持不变
Token string
AESKey string // 掩码 / 空串 = 保持不变
EncryptMode string
TemplateID string
Verified bool
Enabled bool
Remark string
}
// SaveConfig 保存服务号配置(不存在则创建)
func SaveConfig(in SaveInput) error {
if strings.TrimSpace(in.AppID) == "" {
return errors.New("AppID 不能为空")
}
mode := strings.TrimSpace(in.EncryptMode)
if mode == "" {
mode = EncryptModePlain
}
var row models.WechatMpConfig
err := models.Orm.QueryTable(new(models.WechatMpConfig)).OrderBy("id").One(&row)
now := time.Now()
isNew := false
if err != nil {
if err != orm.ErrNoRows {
return err
}
isNew = true
row = models.WechatMpConfig{}
}
// 敏感字段:掩码/空串表示保持原值
secret := strings.TrimSpace(in.AppSecret)
if secret != "" && !IsMasked(secret) {
enc, eerr := EncryptSecret(secret)
if eerr != nil {
return eerr
}
row.AppSecret = enc
}
aesKey := strings.TrimSpace(in.AESKey)
if aesKey != "" && !IsMasked(aesKey) {
enc, eerr := EncryptSecret(aesKey)
if eerr != nil {
return eerr
}
row.AESKey = enc
}
row.AppID = strings.TrimSpace(in.AppID)
row.Token = strings.TrimSpace(in.Token)
row.EncryptMode = mode
row.TemplateID = strings.TrimSpace(in.TemplateID)
row.Remark = strings.TrimSpace(in.Remark)
if in.Verified {
row.Verified = 1
} else {
row.Verified = 0
}
if in.Enabled {
row.Enabled = 1
} else {
row.Enabled = 0
}
row.UpdateTime = &now
if isNew {
row.CreateTime = now
_, err = models.Orm.Insert(&row)
return err
}
_, err = models.Orm.Update(&row)
return err
}
// CallbackURL 微信服务器回调地址(供公众号后台「服务器配置」填写)。
// 优先读取 app.conf 的 wechat_mp_callback_base,其次回落 payment_callback_base(同为对外 API 域名)。
func CallbackURL() string {
base, _ := beego.AppConfig.String("wechat_mp_callback_base")
if strings.TrimSpace(base) == "" {
base, _ = beego.AppConfig.String("payment_callback_base")
}
base = strings.TrimRight(strings.TrimSpace(base), "/")
if base == "" {
return "/api/wechat/mp/callback"
}
return base + "/api/wechat/mp/callback"
}
// ============================ access_token 缓存 ============================
type tokenEntry struct {
appID string
token string
expiresAt time.Time
}
var (
tokenMu sync.Mutex
tokenCache tokenEntry
tokenMinTTL = 5 * time.Minute
)
// resetTokenCacheForTest 仅供测试清理缓存
func resetTokenCacheForTest() {
tokenMu.Lock()
defer tokenMu.Unlock()
tokenCache = tokenEntry{}
}
+106
View File
@@ -0,0 +1,106 @@
package wechatmp
import (
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"errors"
"io"
"strings"
beego "github.com/beego/beego/v2/server/web"
)
// =============================================================
// 微信公众号敏感字段(AppSecret / EncodingAESKey)的对称加密与掩码
//
// yz_platform_wechat_mp_config.app_secret / aes_key 以 AES-256-GCM 加密后存储,
// 格式:base64( "YZW1:" + nonce + ciphertext )。
// 密钥来源:app.conf 的 wechat_mp_secret_key(任意长度字符串,内部做 SHA-256 派生);
// 未配置时退回内置兜底密钥(保证开箱可用),上线前建议在 app.conf 配置该值。
// 注意:更换密钥后,历史密文将无法解密,需要重新保存一次配置。
// =============================================================
const (
cipherPrefix = "YZW1:"
fallbackSecretKey = "yunzer_wechat_mp_default_secret_key_v1"
)
func configKey() []byte {
raw, _ := beego.AppConfig.String("wechat_mp_secret_key")
if strings.TrimSpace(raw) == "" {
raw = fallbackSecretKey
}
sum := sha256.Sum256([]byte(raw))
return sum[:]
}
// EncryptSecret 加密敏感字段;空串原样返回
func EncryptSecret(plain string) (string, error) {
if plain == "" {
return "", nil
}
block, err := aes.NewCipher(configKey())
if err != nil {
return "", err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return "", err
}
nonce := make([]byte, gcm.NonceSize())
if _, err = io.ReadFull(rand.Reader, nonce); err != nil {
return "", err
}
ciphertext := gcm.Seal(nil, nonce, []byte(plain), nil)
buf := append([]byte(cipherPrefix), nonce...)
buf = append(buf, ciphertext...)
return base64.StdEncoding.EncodeToString(buf), nil
}
// DecryptSecret 解密敏感字段;空串原样返回
func DecryptSecret(encoded string) (string, error) {
if encoded == "" {
return "", nil
}
raw, err := base64.StdEncoding.DecodeString(encoded)
if err != nil {
return "", errors.New("密文不是有效的 base64")
}
if len(raw) < len(cipherPrefix) || string(raw[:len(cipherPrefix)]) != cipherPrefix {
return "", errors.New("密文格式不正确(缺少加密前缀)")
}
block, err := aes.NewCipher(configKey())
if err != nil {
return "", err
}
gcm, err := cipher.NewGCM(block)
if err != nil {
return "", err
}
nonce := raw[len(cipherPrefix) : len(cipherPrefix)+gcm.NonceSize()]
ciphertext := raw[len(cipherPrefix)+gcm.NonceSize():]
plain, err := gcm.Open(nil, nonce, ciphertext, nil)
if err != nil {
return "", errors.New("解密失败(wechat_mp_secret_key 是否被更换过?)")
}
return string(plain), nil
}
// MaskSecret 生成掩码:保留末 4 位
func MaskSecret(v string) string {
if v == "" {
return ""
}
if len(v) <= 4 {
return "******"
}
return "******" + v[len(v)-4:]
}
// IsMasked 判断前端回传的值是否仍是掩码(表示「不修改原值」)
func IsMasked(v string) bool {
return strings.HasPrefix(v, "******")
}
+214
View File
@@ -0,0 +1,214 @@
package wechatmp
import (
"crypto/aes"
"crypto/cipher"
"crypto/rand"
"crypto/sha1"
"encoding/base64"
"encoding/binary"
"encoding/hex"
"encoding/xml"
"errors"
"fmt"
"io"
"sort"
"strings"
)
// =============================================================
// 微信被动消息:签名校验 / 安全模式 AES 加解密 / XML 解析与回复构造
// 规范参考:https://developers.weixin.qq.com/doc/offiaccount/Message_Management/Message_encryption_and_decryption_instructions.html
// =============================================================
// pkcs7BlockSize 微信官方填充块大小(注意不是 16,是 32)
const pkcs7BlockSize = 32
// InboundMessage 接收到的消息/事件(已解密后的 XML 内容)
type InboundMessage struct {
XMLName xml.Name `xml:"xml"`
ToUserName string `xml:"ToUserName"`
FromUserName string `xml:"FromUserName"`
CreateTime int64 `xml:"CreateTime"`
MsgType string `xml:"MsgType"` // text / image / event ...
Content string `xml:"Content"`
Event string `xml:"Event"` // subscribe / unsubscribe / SCAN / CLICK ...
EventKey string `xml:"EventKey"` // qrscene_<scene> 或 <scene>
Ticket string `xml:"Ticket"`
MsgID int64 `xml:"MsgId"`
Encrypt string `xml:"Encrypt"` // 安全模式下外层 XML 携带的密文
}
// ParseInboundXML 解析(明文/已解密的)消息 XML
func ParseInboundXML(raw string) (*InboundMessage, error) {
var msg InboundMessage
if err := xml.Unmarshal([]byte(raw), &msg); err != nil {
return nil, fmt.Errorf("解析微信消息失败: %w", err)
}
return &msg, nil
}
// CheckSignature 明文模式签名校验:sha1(sort(token, timestamp, nonce))
func CheckSignature(token, timestamp, nonce, signature string) bool {
if token == "" || signature == "" {
return false
}
arr := []string{token, timestamp, nonce}
sort.Strings(arr)
sum := sha1.Sum([]byte(strings.Join(arr, "")))
return hex.EncodeToString(sum[:]) == signature
}
// CheckMsgSignature 安全/兼容模式签名校验:sha1(sort(token, timestamp, nonce, encrypt))
func CheckMsgSignature(token, timestamp, nonce, encrypt, msgSignature string) bool {
if token == "" || msgSignature == "" {
return false
}
arr := []string{token, timestamp, nonce, encrypt}
sort.Strings(arr)
sum := sha1.Sum([]byte(strings.Join(arr, "")))
return hex.EncodeToString(sum[:]) == msgSignature
}
// MsgSignature 计算安全模式消息签名(回复加密时使用)
func MsgSignature(token, timestamp, nonce, encrypt string) string {
arr := []string{token, timestamp, nonce, encrypt}
sort.Strings(arr)
sum := sha1.Sum([]byte(strings.Join(arr, "")))
return hex.EncodeToString(sum[:])
}
// ============================ AES 加解密 ============================
func decodeAESKey(encodingAESKey string) ([]byte, error) {
key := strings.TrimSpace(encodingAESKey)
if len(key) != 43 {
return nil, fmt.Errorf("EncodingAESKey 长度应为 43 位(当前 %d)", len(key))
}
return base64.StdEncoding.DecodeString(key + "=")
}
// DecryptMessage 安全模式消息解密,返回明文 XML;appID 非空时会校验一致性
func DecryptMessage(aesKey, appID, encryptBase64 string) (string, error) {
key, err := decodeAESKey(aesKey)
if err != nil {
return "", err
}
raw, err := base64.StdEncoding.DecodeString(encryptBase64)
if err != nil {
return "", errors.New("密文不是有效的 base64")
}
block, err := aes.NewCipher(key)
if err != nil {
return "", err
}
if len(raw) < aes.BlockSize || len(raw)%aes.BlockSize != 0 {
return "", errors.New("密文长度不合法")
}
iv := key[:aes.BlockSize]
plain := make([]byte, len(raw))
cipher.NewCBCDecrypter(block, iv).CryptBlocks(plain, raw)
plain, err = pkcs7Unpad(plain)
if err != nil {
return "", err
}
if len(plain) < 20 {
return "", errors.New("解密内容长度不合法")
}
msgLen := int(binary.BigEndian.Uint32(plain[16:20]))
if msgLen < 0 || 20+msgLen > len(plain) {
return "", errors.New("解密内容长度不合法")
}
msg := string(plain[20 : 20+msgLen])
gotAppID := strings.TrimSpace(string(plain[20+msgLen:]))
if appID != "" && gotAppID != appID {
return "", fmt.Errorf("AppID 校验失败(收到 %s)", gotAppID)
}
return msg, nil
}
// EncryptMessage 安全模式消息加密,返回 base64 密文
func EncryptMessage(aesKey, appID, msg string) (string, error) {
key, err := decodeAESKey(aesKey)
if err != nil {
return "", err
}
random := make([]byte, 16)
if _, err := io.ReadFull(rand.Reader, random); err != nil {
return "", err
}
msgBytes := []byte(msg)
lenBytes := make([]byte, 4)
binary.BigEndian.PutUint32(lenBytes, uint32(len(msgBytes)))
buf := make([]byte, 0, 16+4+len(msgBytes)+len(appID))
buf = append(buf, random...)
buf = append(buf, lenBytes...)
buf = append(buf, msgBytes...)
buf = append(buf, []byte(appID)...)
buf = pkcs7Pad(buf)
block, err := aes.NewCipher(key)
if err != nil {
return "", err
}
iv := key[:aes.BlockSize]
ciphertext := make([]byte, len(buf))
cipher.NewCBCEncrypter(block, iv).CryptBlocks(ciphertext, buf)
return base64.StdEncoding.EncodeToString(ciphertext), nil
}
func pkcs7Pad(data []byte) []byte {
pad := pkcs7BlockSize - len(data)%pkcs7BlockSize
if pad <= 0 {
pad = pkcs7BlockSize
}
out := make([]byte, 0, len(data)+pad)
out = append(out, data...)
for i := 0; i < pad; i++ {
out = append(out, byte(pad))
}
return out
}
func pkcs7Unpad(data []byte) ([]byte, error) {
if len(data) == 0 {
return nil, errors.New("空数据")
}
pad := int(data[len(data)-1])
if pad < 1 || pad > pkcs7BlockSize || pad > len(data) {
return nil, errors.New("PKCS7 填充不合法")
}
for i := len(data) - pad; i < len(data); i++ {
if int(data[i]) != pad {
return nil, errors.New("PKCS7 填充不合法")
}
}
return data[:len(data)-pad], nil
}
// ============================ 被动回复构造 ============================
// BuildTextReply 构造文本消息回复 XML(明文模式直接返回)
func BuildTextReply(toUser, fromUser, content string, timestamp int64) string {
esc := func(s string) string {
return strings.NewReplacer(
"]]>", "]] >",
).Replace(s)
}
return fmt.Sprintf(
`<xml><ToUserName><![CDATA[%s]]></ToUserName><FromUserName><![CDATA[%s]]></FromUserName><CreateTime>%d</CreateTime><MsgType><![CDATA[text]]></MsgType><Content><![CDATA[%s]]></Content></xml>`,
esc(toUser), esc(fromUser), timestamp, esc(content))
}
// BuildEncryptedReply 安全模式回复:对明文回复加密并组装外层 XML
func BuildEncryptedReply(cfg *Config, plainReply, timestamp, nonce string) (string, error) {
encrypt, err := EncryptMessage(cfg.AESKey, cfg.AppID, plainReply)
if err != nil {
return "", err
}
sig := MsgSignature(cfg.Token, timestamp, nonce, encrypt)
return fmt.Sprintf(
`<xml><Encrypt><![CDATA[%s]]></Encrypt><MsgSignature><![CDATA[%s]]></MsgSignature><TimeStamp>%s</TimeStamp><Nonce><![CDATA[%s]]></Nonce></xml>`,
encrypt, sig, timestamp, nonce), nil
}
+113
View File
@@ -0,0 +1,113 @@
package wechatmp
import (
"errors"
"time"
"server/models"
)
// =============================================================
// 消息推送(认证后可用):
// 未认证服务号微信不允许主动推送(模板消息/订阅通知/客服消息均受限),
// 本层只负责在认证后通过模板消息把提醒/公告发给已关注粉丝;
// 未认证时返回 ErrNotVerified,由上层给出明确提示。
// =============================================================
// ErrNotVerified 公众号未认证,无法主动推送
var ErrNotVerified = errors.New("公众号未认证:模板消息推送需认证后使用(当前仅支持关注与被动回复)")
// PushMessage 推送内容(对应「标题/内容/备注」通用模板字段)
type PushMessage struct {
Title string // first:首行说明
Content string // keyword1:主体内容
Remark string // remark:备注
URL string // 点击跳转地址(可选)
}
// PushToOpenID 向指定粉丝推送一条模板消息
func PushToOpenID(cfg *Config, openid string, msg PushMessage) error {
if cfg == nil {
return ErrNotConfigured
}
if !cfg.Verified {
return ErrNotVerified
}
if openid == "" {
return errors.New("openid 为空")
}
if msg.Title == "" {
msg.Title = "您有一条新的消息"
}
if len([]rune(msg.Title)) > 100 {
msg.Title = string([]rune(msg.Title)[:100])
}
if len([]rune(msg.Content)) > 200 {
msg.Content = string([]rune(msg.Content)[:200])
}
if len([]rune(msg.Remark)) > 100 {
msg.Remark = string([]rune(msg.Remark)[:100])
}
return SendTemplateMessage(cfg, openid, msg.URL, map[string]string{
"first": msg.Title,
"keyword1": msg.Content,
"remark": msg.Remark,
})
}
// PushScope 推送范围
type PushScope struct {
// BindType 为空表示全部已关注粉丝;platform_user / tenant_user 只推绑定对应类型账号的粉丝
BindType string
// BindID 非 0 时只推给该绑定账号(与 BindType 组合使用)
BindID uint64
// Tid 非 0 时只推给绑定该租户的用户(tenant_user 场景)
Tid uint64
}
// MaxBatchPush 单次批量推送上限
const MaxBatchPush = 500
// PushToFollowers 按范围批量推送,返回成功/失败数量。
// 仅推送给「已关注(subscribe=1)且未删除」的粉丝。
func PushToFollowers(cfg *Config, scope PushScope, msg PushMessage) (sent, failed int, err error) {
if cfg == nil {
return 0, 0, ErrNotConfigured
}
if !cfg.Verified {
return 0, 0, ErrNotVerified
}
qs := models.Orm.QueryTable(new(models.WechatMpFollower)).
Filter("subscribe", 1).
Filter("delete_time__isnull", true)
if scope.BindType != "" {
qs = qs.Filter("bind_type", scope.BindType)
}
if scope.BindID > 0 {
qs = qs.Filter("bind_id", scope.BindID)
}
if scope.Tid > 0 {
qs = qs.Filter("bind_tid", scope.Tid)
}
var rows []models.WechatMpFollower
if _, qerr := qs.OrderBy("-id").Limit(MaxBatchPush).All(&rows); qerr != nil {
return 0, 0, qerr
}
if len(rows) == 0 {
return 0, 0, errors.New("没有匹配的已关注粉丝")
}
for i := range rows {
if perr := PushToOpenID(cfg, rows[i].OpenID, msg); perr != nil {
failed++
continue
}
sent++
// 简单节流,避免触发微信频率限制
if i%20 == 19 {
time.Sleep(200 * time.Millisecond)
}
}
return sent, failed, nil
}