Files
MailUI4Agents/server/internal/push/push.go
JianFeeeee 46fa7fa729 feat(push): 可选、配置式、多厂商的推送通道(HMS 为首个实现)
用户要求:推送密钥必须是可选项(自部署后端不能写死推送方式),且要支持
多厂商配置式接入 —— 每个用户各自部署服务器、自己选厂商、自己配凭证。
所以落地成:

· internal/push:通道抽象 + 工厂表(RegisterType),加厂商不改配置层与端点形状;
  HMS 只是第一个实现(internal/push/hms.go)
· 配置在 PUSH_CONFIG(默认 <AGENTMAIL_DATA_DIR>/push.json),一项一个厂商,
  凭证走文件(app_secret_file / files.*,建议 600);环境变量只是可选覆盖
· 没配 = 整条推送路径连一次查库都不发生(shouldDispatch 早退);
  单项配错(未知类型/密钥读不到/enabled:false)只跳过那一条,不影响启动
· push_tokens 表带 provider 维度 + 三个 /me/devices/push-token 端点;
  没配推送时端点照存并回 enabled:false(登记成功 != 服务端开了推送)
· notify.Recipients 末尾异步挂钩:收件人名单直接用 SSE 那份 seen(两条通道
  共用同一份"谁该收到"的判据);失败只记日志,绝不拖住收信

HMS 的形状是拿真凭证打线上接口问出来的(v1 + message.token[] + testMessage;
payload/target 形状 v1 不认、v2 要服务账号 JWT)。未上架应用必须 test_message=true,
单批 ≤10 token(MaxTokensPerRequest 声明)、每日 1000 条兜底(项目级额度)。
实测:App ID + App Secret 能换到 access_token(3600s);形状被线上服务接受。

判据:repo 6 条 + push 12 条 + handler 3 组,全部做过**变异验证** ——
过程中抓出两条假判据(异步分发与 t.Cleanup 赛跑而假绿;密钥文件优先级没被覆盖)
并补掉。Go 全量测试与 go vet 干净。

★ 未验:端到端真机送达(需要真机 token + 客户端按 com.jianf.agentmail 重编并签名,
签名指纹还要在 AGC 登记)—— 从未真正发出过一条能到达设备的推送。
详见 docs/HMS-PUSH-PLAN.md 的「实现状态」一节。
2026-09-15 11:21:00 +08:00

180 lines
6.2 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 是「新邮件」的第二条送达通道(第一条是 SSE)。
# 它是可选的,且默认关闭
自部署实例通常**一个推送渠道都没配** —— 那是正常状态,不是配置错误:客户端
在线时 SSE 已经够用,推送只解决「App 不在前台 / 被系统杀掉」这一种情形。
因此本包所有入口在没注册任何 provider 时都是**立即返回**:不查库、不建连接、
不刷日志(用户 2026-09-15 的明确要求:「不能写死推送方式,因为我们是自部署后端」
「即推送密钥应当是可选项」)。
# 为什么不写死华为
表与端点都带 `provider` 维度,Notifier 是接口:加一个通道(web push、别的厂商)
只加一个实现 + 一行注册,不动 schema、不动端点形状、不动调用方。
华为 HMS 只是第一个实现(internal/push/hms.go)。
# 为什么发送是异步且会丢
推送发生在**收信路径**上(notify.Recipients),而它已经在库事务之外的下发阶段:
一个慢的推送 HTTP 请求不能拖住邮件送达 —— 收信是主功能,推送是锦上添花。
因此:有界并发 + 超时 + 失败只记日志。**宁可丢一条通知,不可慢一封邮件**。
*/
package push
import (
"context"
"log"
"sync"
"time"
"github.com/agentmail/gateway/internal/repo"
)
// NewMail 是要推送的一封新邮件。
//
// 字段刻意少:推送只负责「通知栏那一行 + 点进去看哪封信」。正文不进通知,
// 否则锁屏上就会露出邮件内容(SSE 是给已解锁的在线客户端用的,两者隐私模型不同)。
type NewMail struct {
MailID string
SessionID string
// From 是发件方名字,用于通知标题。
From string
// Subject 是邮件主题。
Subject string
// Recipients 是**该收到这封信的人名**(= SSE 的收件判据:主收件人 + 抄送方)。
// 两条通道共用同一份名单,不各算一套。
Recipients []string
}
// Notifier 是一个推送通道。
type Notifier interface {
// Name 是 provider 标识,与 push_tokens.provider 的值一致(如 "hms")。
Name() string
// MaxTokensPerRequest 是单次请求能带的最大 token 数(厂商限额,如华为测试消息 ≤10)。
// 由通道自己声明,而不是调用方写死一个「10」—— 限额是通道的属性。
MaxTokensPerRequest() int
// Send 向一批 token 投递同一条通知。返回错误只用于**记日志**。
Send(ctx context.Context, tokens []string, n NewMail) error
}
const (
// maxInFlight 是在途推送任务上限。超了就丢掉这一轮(记日志),不排队:
// 排队的后果是通知在几十秒后集中弹出来,那比丢掉更糟。
maxInFlight = 4
// sendTimeout 是单个 provider 单批的超时。
sendTimeout = 10 * time.Second
)
var (
mu sync.RWMutex
providers []Notifier
slots = make(chan struct{}, maxInFlight)
)
// Register 注册一个推送通道。由 main 按配置调用 —— 没配就不注册。
func Register(n Notifier) {
if n == nil {
return
}
mu.Lock()
defer mu.Unlock()
providers = append(providers, n)
log.Printf("[push] 通道已启用: %s(单批最多 %d 个 token)", n.Name(), n.MaxTokensPerRequest())
}
// Enabled 报告是否配了任何推送通道。
//
// 端点据此回 `enabled`,客户端据此知道自己「登记了也可能收不到」——
// 而不是以为登记失败。
func Enabled() bool {
mu.RLock()
defer mu.RUnlock()
return len(providers) > 0
}
// Names 返回已启用的通道名(端点回给客户端看)。
func Names() []string {
mu.RLock()
defer mu.RUnlock()
out := make([]string, 0, len(providers))
for _, p := range providers {
out = append(out, p.Name())
}
return out
}
// shouldDispatch 报告这封邮件值不值得进推送管线:没配通道、或没有收件人 → 不值。
//
// 为什么单独成一个**纯函数**(不碰库、不起 goroutine):因为“没配通道时零开销”
// 这条判据必须能**同步**验证。2026-09-15 变异验证实测过:把判据写成“调 NotifyNewMail
// 后用 nil 库不 panic”,去掉本函数里的 Enabled() 后测试**仍然绿** —— 分发在
// goroutine 里跑,而 t.Cleanup 已经把真库装回去了,于是判据在错误的理由上通过。
// 纯函数没有这个<E8BF99>赛跑面。
func shouldDispatch(n NewMail) bool {
return Enabled() && len(n.Recipients) > 0
}
// NotifyNewMail 异步把一封新邮件推给收件方登记的设备。
//
// 调用方(notify.Recipients)**永远不因此拿到错误**:推送失败不该影响收信,
// 也不该让调用方写一半成功一半失败的处理逻辑。
func NotifyNewMail(ctx context.Context, n NewMail) {
if !shouldDispatch(n) {
return
}
select {
case slots <- struct{}{}:
default:
log.Printf("[push] 在途任务已达上限 %d,跳过本轮推送(可选通道,丢一条通知不影响收信)", maxInFlight)
return
}
go func() {
defer func() { <-slots }()
// 用 Background 而不是请求的 ctx:发信请求一旦返回,ctx 就被取消,
// 挂在它上面的推送会被立刻掐断(而收信方恰恰是那个已经离线的人)。
sendCtx, cancel := context.WithTimeout(context.Background(), sendTimeout)
defer cancel()
dispatch(sendCtx, n)
}()
}
// dispatch 按 provider 分组投递,每批不超过该通道声明的上限。
func dispatch(ctx context.Context, n NewMail) {
tokens, err := repo.ListPushTokensOfOwners(ctx, n.Recipients)
if err != nil {
log.Printf("[push] 读推送登记失败(不影响收信): %v", err)
return
}
if len(tokens) == 0 {
return
}
grouped := map[string][]string{}
for _, t := range tokens {
grouped[t.Provider] = append(grouped[t.Provider], t.Token)
}
mu.RLock()
list := append([]Notifier(nil), providers...)
mu.RUnlock()
for _, p := range list {
ts := grouped[p.Name()]
if len(ts) == 0 {
continue
}
batch := p.MaxTokensPerRequest()
if batch <= 0 {
batch = 1
}
for i := 0; i < len(ts); i += batch {
end := i + batch
if end > len(ts) {
end = len(ts)
}
if err := p.Send(ctx, ts[i:end], n); err != nil {
log.Printf("[push] %s 投递失败(不影响收信): %v", p.Name(), err)
}
}
}
}