Files
MailUI4Agents/server/internal/push/hms.go
JianFeeeee 186cf53804 fix(push): 额度预留**真的**挪到 accessToken 之后(上一版只写了注释)
## 起因:pi 在邮件驱动的一轮里当场抓出来的

pi 收到那封 `[收尾验证]` 邮件后,自己翻代码核对,
在会话文件里写下(原文):

    The code contradicts its own comment (item ②: reserve should be *after* accessToken)
    The commit only changed the argument (`len(tokens)` → `1`) and the comment — it
    Fix ① (per-batch) is real and tested. Fix ② is claimed but not implemented.
    Confirmed — the bug is real.

它甚至自己造了探针(`zz_probe_test.go`,跑完已删)来实证。

**我独立复核确认它是对的**:
`reserveDaily(1)` 在第 212 行,`accessToken` 在第 215 行 ——
预留仍在**之前**。2026-09-26 那次我只改了 ①(`len(tokens)` → `1`),
把 ② 写进了注释,**代码没动**。

## 为什么当时那批判据没接住

`push_test.go` 原有 3 格只验 ①(按批次计),**造不出「accessToken 失败」这条路** ——
`hmsStub` 的 `/token` 永远返回 200 + 令牌。

⇒ 「注释说修了」与「代码真修了」能分家,而没有任何东西会发现。

## 改法

① `hmsStub` 加 `failToken` 开关(`/token` 可返回 400)。
② `reserveDaily(1)` 挪到 `accessToken` 成功**之后**、真正发请求之前。
   仍保持**前置预留**语义(不是"发成功后再扣")—— 那会超发,
   并发下多个 goroutine 都能通过检查。宁可少算也不多发。
③ 新增 `TestHMSAccessTokenFailureDoesNotBurnQuota`:三次 accessToken 失败后
   断言 `dayCount == 0`、零推送发出、且恢复正常后仍能发(额度没被吃掉)。

## 变异验证(这格判据本该在 2026-09-26 就存在)

把 `reserveDaily` 挪回 `accessToken` 之前(= 还原成 bug)⇒

    ★ accessToken 失败不该扣额度,实际已扣 3 条
    (一次网络抖动静默烧配额就是这么来的)

## 教训(与本仓 python-probe-shadowing / baseline-residue 同族)

**「我写了注释说明怎么修」不等于「我改了代码」。**
审查报告给了两条,我处理了一条,把另一条**誊进了注释**就当做了。
写完注释应当立刻核对行号 —— 那是 5 秒钟的事,而这次是别人替我发现的。

★ 另一层:**别人(或另一个 Agent)独立复核出来的结论,要自己再验一遍再改**。
我逐条查了行号才动手,没有因为"pi 说的"就直接信。
2026-09-28 09:49:22 +08:00

310 lines
11 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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` = 一条消息。
*
* ② **扣在投递之前**:原顺序是 reserveDaily → accessToken → HTTP 请求。
* `accessToken` 失败 / HTTP 失败 / 华为回非成功码,这些**一条都没发出去**,
* 额度却已经扣了,而失败只 `log.Printf`
* ⇒ 一次网络抖动静默烧掉配额。
* ⇒ 预留挪到 `accessToken` **之后**、真正发请求之前(见下方调用点)。
*
* ★ 为什么不是"发送成功后再扣":那会超发(并发下多个 goroutine
* 都能通过检查)。保留前置预留、但放在"确认能发"之后,是**宁可少算也不多发**
* 的取舍 —— 少算的代价是偶尔一次失败没计入,超发的代价是真超额被华为拒。
*/
tok, err := h.accessToken(ctx)
if err != nil {
return err
}
// 必须在 accessToken **之后**预留(②):拿不到 token 就一条都发不出去。
// 挪回上面一行会让 DailyLimit=1 时一次网络抖动烧穿当日配额。
if !h.reserveDaily(1) {
return fmt.Errorf("达到每日推送上限 %d 条(HMS_DAILY_LIMIT)", h.DailyLimit)
}
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]) + "…"
}