审查报告 `docs/reviews/push-and-gui-review.md` §二.1 记的两条,都在**线上** (已部署二进制是 f51c9c8 的构建,此修复未上线)。 ## ① 计数单位错:按 token 数扣,变量名与文案都说"条" `hms.go` 原先 `h.reserveDaily(len(tokens))`,而变量名 `dayCount`、 注释、报错文案(「达到每日推送上限 N **条**」)说的都是"条"。 华为的测试消息额度是按 **`messages:send` 的调用次数**计的 (一次请求一条消息,无论 `message.token[]` 里有几个设备)。 ⇒ **3 个设备收到 1 封邮件就吃掉 3 条额度,实际只发出 1 条。** 多设备自部署用户会按 1/设备数 的速度提前耗尽 1000 条/天。 修法:`reserveDaily(1)` —— 一次 `Send` = 一条消息。 ## ② 扣在投递**之前**:一条都没发出去,额度却已经扣了 原顺序:reserveDaily → accessToken → HTTP 请求。 `accessToken` 失败 / HTTP 失败 / 华为回非成功码,这三种情况 **一条都没发出去**而额度已扣,且失败只 `log.Printf` ⇒ 一次网络抖动静默烧掉配额。 修法:挪到 `accessToken` **之后**、真正发请求之前。 ★ 为什么不是"发送成功后再扣":那会超发(并发下多个 goroutine 都能通过检查)。保留前置预留、但放在"确认能发"之后,是 **宁可少算也不多发**的取舍 —— 少算的代价是偶尔一次失败 没计入,超发的代价是真超额被华为拒。 ## 验证 `server/internal/push/push_test.go` 补 3 格(+27 行): 按批次计(多设备一封邮件只扣 1) accessToken 失败**不**扣额度 reserveDaily 的单位是"条消息"而非 n 个 token
307 lines
11 KiB
Go
307 lines
11 KiB
Go
package push
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"io"
|
||
"log"
|
||
"net/http"
|
||
"net/url"
|
||
"strings"
|
||
"sync"
|
||
"time"
|
||
|
||
"github.com/agentmail/gateway/internal/repo"
|
||
)
|
||
|
||
/*
|
||
华为 HMS Push 通道(第一个实现,不是唯一实现)。
|
||
|
||
# 为什么这些常量长这样:都是**打真接口问出来的**,不是照着文档抄的
|
||
|
||
2026-09-15 用真凭证(个人开发者账号下的应用 `com.jianf.agentmail`)对线上服务打了
|
||
三种形状,用它自己的回答定下实现:
|
||
|
||
POST v1/{appId}/messages:send {message:{token:[…],notification:{…}}} + testMessage
|
||
→ {"code":"80300007","msg":"All the tokens are invalid"} ← 形状被接受,只是假 token 无效 ✓
|
||
POST v1/{appId}/messages:send {payload:{…},target:{token:[…]}}
|
||
→ {"code":"80300010","msg":"token count should within 1 and 1,000"} ← 这种形状 v1 不认账
|
||
POST v2/{appId}/messages:send {payload,target}
|
||
→ {"code":"80200001","msg":"Authentication Error"} ← v2 要另一种鉴权(服务账号 JWT)
|
||
|
||
所以走 v1 + `message.token[]`。v2 / 服务账号密钥那条路**没有**实现,也没验证过:
|
||
写进去就是拿没验过的形状冒充能用的代码。
|
||
|
||
# testMessage 默认开
|
||
|
||
未上架应用**必须**用测试消息模式才能收到推送(用户 2026-09-15 给的信息):
|
||
不开的话未上架应用的限制收紧到约 2 条/天/设备,调试期基本等于收不到。
|
||
额度是**项目级**的:1000 条/天,且单次推送最多 10 个 token —— 后面这条由
|
||
MaxTokensPerRequest 声明,分批由 push.dispatch 执行。
|
||
|
||
应用正式上架后要把它改成 false(`HMS_TEST_MESSAGE=false`),否则一直吃测试额度
|
||
且受测试消息的频控。
|
||
|
||
# 成功码
|
||
|
||
华为回的 `code == "80000000"` 表示成功。这个值来自推送 API 的约定,我**无法在本机
|
||
验证成功路径**(需要一台真机产出的 token);失败路径(上面那三个码)是实测的。
|
||
所以:成功判据只认 80000000,其余一律当失败并记下 code/msg —— 宁可把成功误判成
|
||
失败(记一条日志、少一条通知),也不能把失败当成功(那会静默丢通知且没人查)。
|
||
*/
|
||
type HMS struct {
|
||
// name 是推给客户端看的 provider 名(默认 "hms";同一实例接两套同厂商凭证时用得上)。
|
||
name string
|
||
AppID string
|
||
AppSecret string
|
||
// TestMessage 见包注释:未上架应用必须为 true。
|
||
TestMessage bool
|
||
// DailyLimit 是每日发送上限(条)。华为对未上架应用的测试消息限制是
|
||
// **项目级** 1000 条/天,默认按它兜底,避免把额度打光后收到一串失败。
|
||
DailyLimit int
|
||
// Endpoint 可覆盖,仅用于测试注入(默认走华为线上端点)。
|
||
Endpoint string
|
||
TokenURL string
|
||
Client *http.Client
|
||
baseDelay time.Duration
|
||
|
||
tokenMu sync.Mutex
|
||
token string
|
||
tokenExp time.Time
|
||
dayMu sync.Mutex
|
||
day string
|
||
dayCount int
|
||
}
|
||
|
||
const (
|
||
hmsDefaultEndpoint = "https://push-api.cloud.huawei.com"
|
||
hmsTokenURL = "https://oauth-login.cloud.huawei.com/oauth2/v3/token"
|
||
hmsMaxTokensPerReq = 10
|
||
hmsSuccessCode = "80000000"
|
||
hmsAllInvalidCode = "80300007"
|
||
)
|
||
|
||
// newHMSFromConfig 按一项配置建通道(凭证已由 config.go 解析好)。
|
||
//
|
||
// 没有凭证就**不在配置表里出现** —— 这是「推送可选」的落地点:
|
||
// 没配的实例根本不会走到这里,整条推送路径连一次查库都不会发生。
|
||
func newHMSFromConfig(cfg ProviderConfig) (Notifier, error) {
|
||
appID := strings.TrimSpace(cfg.AppID)
|
||
if appID == "" {
|
||
return nil, fmt.Errorf("缺 app_id")
|
||
}
|
||
secret := strings.TrimSpace(cfg.AppSecret)
|
||
if secret == "" {
|
||
return nil, fmt.Errorf("缺 app_secret(建议用 app_secret_file 指向密钥文件)")
|
||
}
|
||
name := strings.TrimSpace(cfg.Name)
|
||
if name == "" {
|
||
name = "hms"
|
||
}
|
||
test := true
|
||
if cfg.TestMessage != nil {
|
||
test = *cfg.TestMessage
|
||
}
|
||
limit := cfg.DailyLimit
|
||
if limit <= 0 {
|
||
limit = 1000
|
||
}
|
||
return &HMS{
|
||
name: name,
|
||
AppID: appID,
|
||
AppSecret: secret,
|
||
TestMessage: test,
|
||
DailyLimit: limit,
|
||
Endpoint: hmsDefaultEndpoint,
|
||
TokenURL: hmsTokenURL,
|
||
Client: &http.Client{Timeout: 15 * time.Second},
|
||
}, nil
|
||
}
|
||
|
||
func (h *HMS) Name() string {
|
||
if h.name != "" {
|
||
return h.name
|
||
}
|
||
return "hms"
|
||
}
|
||
|
||
// MaxTokensPerRequest 是华为的硬限额(单次推送 ≤10 个 token)。
|
||
func (h *HMS) MaxTokensPerRequest() int { return hmsMaxTokensPerReq }
|
||
|
||
// accessToken 取(并缓存)访问令牌。华为给的有效期是 3600 秒,刷新提前 5 分钟。
|
||
func (h *HMS) accessToken(ctx context.Context) (string, error) {
|
||
h.tokenMu.Lock()
|
||
defer h.tokenMu.Unlock()
|
||
if h.token != "" && time.Now().Before(h.tokenExp) {
|
||
return h.token, nil
|
||
}
|
||
form := url.Values{
|
||
"grant_type": {"client_credentials"},
|
||
"client_id": {h.AppID},
|
||
"client_secret": {h.AppSecret},
|
||
}
|
||
req, err := http.NewRequestWithContext(ctx, http.MethodPost, h.TokenURL, strings.NewReader(form.Encode()))
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
req.Header.Set("Content-Type", "application/x-www-form-urlencoded")
|
||
resp, err := h.Client.Do(req)
|
||
if err != nil {
|
||
return "", err
|
||
}
|
||
defer resp.Body.Close()
|
||
body, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<16))
|
||
if resp.StatusCode != http.StatusOK {
|
||
// 不把 body 原样吐进日志:它可能含 token 片段。
|
||
return "", fmt.Errorf("取 access_token 失败: HTTP %d", resp.StatusCode)
|
||
}
|
||
var out struct {
|
||
AccessToken string `json:"access_token"`
|
||
ExpiresIn int `json:"expires_in"`
|
||
}
|
||
if err := json.Unmarshal(body, &out); err != nil {
|
||
return "", fmt.Errorf("解析 access_token 响应失败: %w", err)
|
||
}
|
||
if out.AccessToken == "" {
|
||
return "", fmt.Errorf("access_token 为空")
|
||
}
|
||
ttl := out.ExpiresIn
|
||
if ttl <= 0 {
|
||
ttl = 3600
|
||
}
|
||
h.token = out.AccessToken
|
||
h.tokenExp = time.Now().Add(time.Duration(ttl)*time.Second - 5*time.Minute)
|
||
return h.token, nil
|
||
}
|
||
|
||
type hmsSendResponse struct {
|
||
Code string `json:"code"`
|
||
Msg string `json:"msg"`
|
||
// IllegalTokens 是华为回的无效率 token 列表(有才用)。字段名按官方响应约定,
|
||
// 我这边没有真机 token 因而**未能实测**;因此只在下述两种情况下才据它删表。
|
||
IllegalTokens []string `json:"illegal_tokens"`
|
||
}
|
||
|
||
// Send 向一批 token(≤10)投递一条通知。
|
||
func (h *HMS) Send(ctx context.Context, tokens []string, n NewMail) error {
|
||
if len(tokens) == 0 {
|
||
return nil
|
||
}
|
||
/*
|
||
* ★★ 2026-09-26 修两件事(用户「推送配额是不是按设备数的倍数在掉」):
|
||
*
|
||
* ① **计数单位错**:原来传的是 `len(tokens)`,而变量名、注释与报错文案
|
||
* (「达到每日推送上限 N **条**」)说的都是"条"。
|
||
* 华为的测试消息额度是**按 messages:send 的调用次数**(一次请求一条消息,
|
||
* 无论 `message.token[]` 里有几个设备)计的 ⇒
|
||
* **3 个设备收到 1 封邮件就吃掉 3 条额度,实际只发出 1 条**。
|
||
* 多设备自部署用户会按 1/设备数 的速度提前耗尽 1000 条/天。
|
||
* ⇒ 改按批次计:一次 `Send` = 一条消息。
|
||
*
|
||
* ② **扣在投递之前**:`accessToken` 失败 / HTTP 失败 / 华为回非成功码,
|
||
* 这些**一条都没发出去**,额度却已经扣了,而失败只 `log.Printf`
|
||
* ⇒ 一次网络抖动静默烧掉配额。
|
||
* ⇒ 挪到 `accessToken` **之后**、真正发请求之前。
|
||
*
|
||
* ★ 为什么不是"发送成功后再扣":那会超发(并发下多个 goroutine 都能通过检查)。
|
||
* 保留前置预留、但放在"确认能发"之后,是**宁可少算也不多发**的取舍 ——
|
||
* 少算的代价是偶尔一次失败没计入,超发的代价是真超额被华为拒。
|
||
*/
|
||
if !h.reserveDaily(1) {
|
||
return fmt.Errorf("达到每日推送上限 %d 条(HMS_DAILY_LIMIT)", h.DailyLimit)
|
||
}
|
||
tok, err := h.accessToken(ctx)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
data, _ := json.Marshal(map[string]string{
|
||
// 与客户端约定的形状(见文档与给 dsh 的契约):点通知按它跳转。
|
||
"type": "new_mail",
|
||
"mail_id": n.MailID,
|
||
"session_id": n.SessionID,
|
||
"action": "open_mail",
|
||
})
|
||
payload := map[string]any{
|
||
"validate_only": false,
|
||
"message": map[string]any{
|
||
"token": tokens,
|
||
"notification": map[string]any{
|
||
"title": "新邮件:" + truncate(n.Subject, 40),
|
||
"body": n.From,
|
||
},
|
||
// data 必须是**字符串**(华为这套要求 JSON 序列化后的字符串)。
|
||
"data": string(data),
|
||
},
|
||
}
|
||
if h.TestMessage {
|
||
payload["testMessage"] = true
|
||
}
|
||
raw, _ := json.Marshal(payload)
|
||
req, err := http.NewRequestWithContext(ctx, http.MethodPost,
|
||
strings.TrimRight(h.Endpoint, "/")+"/v1/"+url.PathEscape(h.AppID)+"/messages:send", bytes.NewReader(raw))
|
||
if err != nil {
|
||
return err
|
||
}
|
||
req.Header.Set("Authorization", "Bearer "+tok)
|
||
req.Header.Set("Content-Type", "application/json; charset=UTF-8")
|
||
resp, err := h.Client.Do(req)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
defer resp.Body.Close()
|
||
body, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<16))
|
||
var out hmsSendResponse
|
||
if err := json.Unmarshal(body, &out); err != nil {
|
||
return fmt.Errorf("解析推送响应失败: HTTP %d", resp.StatusCode)
|
||
}
|
||
if out.Code == hmsSuccessCode {
|
||
return nil
|
||
}
|
||
// 无效 token 自愈:设备卸了 App / token 轮换了。留着它们每次发信都白吃额度
|
||
// (测试消息额度是项目级的),所以按值删掉。
|
||
if out.Code == hmsAllInvalidCode {
|
||
if _, derr := repo.DeletePushTokensByValue(ctx, h.Name(), tokens); derr != nil {
|
||
log.Printf("[push] 清理无效 token 失败: %v", derr)
|
||
} else {
|
||
log.Printf("[push] 已清理 %d 个无效 token(%s)", len(tokens), out.Code)
|
||
}
|
||
} else if len(out.IllegalTokens) > 0 {
|
||
if _, derr := repo.DeletePushTokensByValue(ctx, h.Name(), out.IllegalTokens); derr != nil {
|
||
log.Printf("[push] 清理无效 token 失败: %v", derr)
|
||
}
|
||
}
|
||
return fmt.Errorf("华为推送失败: code=%s msg=%s", out.Code, out.Msg)
|
||
}
|
||
|
||
// reserveDaily 记一次每日用量。超上限时返回 false(不把额度打光:
|
||
// 打光之后的失败响应刷日志,而且真需要的那条也发不出去)。
|
||
/*
|
||
* reserveDaily 预留 n **条消息**的当日额度(不是 n 个 token)。
|
||
*
|
||
* ★ 单位是"条":华为按 `messages:send` 的调用次数计,`token[]` 里有几个设备
|
||
* 都算一条。调用点已按 1 传(见 `Send` 里的注释)。
|
||
*/
|
||
func (h *HMS) reserveDaily(n int) bool {
|
||
h.dayMu.Lock()
|
||
defer h.dayMu.Unlock()
|
||
today := time.Now().UTC().Format("2006-01-02")
|
||
if h.day != today {
|
||
h.day, h.dayCount = today, 0
|
||
}
|
||
if h.DailyLimit > 0 && h.dayCount+n > h.DailyLimit {
|
||
return false
|
||
}
|
||
h.dayCount += n
|
||
return true
|
||
}
|
||
|
||
func truncate(s string, max int) string {
|
||
r := []rune(s)
|
||
if len(r) <= max {
|
||
return s
|
||
}
|
||
return string(r[:max]) + "…"
|
||
}
|