/* 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 已经把真库装回去了,于是判据在错误的理由上通过。 // 纯函数没有这个�赛跑面。 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) } } } }