模块设计
核心模块
Worker Pool(统一 Worker 池)
Worker Pool 管理一组 Worker,统一从 Queue 消费任务并投递:
mermaid
graph LR
Queue -->|Pop| WP["Worker Pool"]
WP -->|Ack/Nack| Queue
WP --> RM["Runtime Manager"]
RM --> Provider["Provider"]核心职责:
- 启动 N 个 local worker goroutine
- 每个 worker 循环:Pop → Deliver → Ack/Nack
- 管理远程 Worker 注册信息
源码位置: core/worker/pool.go
Registry(Worker 注册表)
统一管理所有 Worker 的元信息:
mermaid
graph LR
subgraph Registry["Registry"]
L0["local-0 → mode: local"]
L1["local-1 → mode: local"]
W1["worker-01 → mode: remote"]
W2["worker-02 → mode: remote"]
end源码位置: core/worker/registry.go
数据流
mermaid
graph LR
API["API Request"] --> Handler --> NS["NotificationService"]
NS -->|路由| Route["Route Engine"]
NS -->|去重| Dedup["Dedup"]
NS -->|渲染| Template["Template"]
NS --> Planner["DeliveryPlanner"]
Planner -->|Binding 解析| Queue["Queue"]
Queue --> WP["Worker Pool"]
WP -->|Provider.Deliver| Provider["Provider"]核心类型
Notification(通知意图)
API 层的输入,描述"用户想发什么"。
go
type Notification struct {
ID string
Type string
Level string
AudienceID string // 受众维度(空为广播)
Channels []string
Recipients map[string][]string
TemplateRef string
Params map[string]any
Content *DirectContent
CreatedAt time.Time
RelationType string // subscription / enrollment(受众关系上下文)
Source string // 入口来源(bot / preference_center / admin / app:<ns>)
}DeliveryTask(投递任务)
Provider 层的输入,描述"怎么发到具体渠道"。
go
type DeliveryTask struct {
ID string
Provider string
Targets []string
Payload DeliveryPayload
Level string
AlertID string // 确认身份透传到交互卡片按钮
Status DeliveryStatus // 投递状态机(queued/delivered/failed/dead)
RetryCount int // 已重试次数(attempts = RetryCount + 1)
MaxAttempts int // 首次投递时记入的总尝试预算
LastError string // 终态为 failed/dead 时的最后一次错误
CreatedAt time.Time
NextRetryAt *time.Time // 等退避重投期间由队列持留
RelationType string // subscription / enrollment(受众关系上下文)
Source string // 入口来源(bot / preference_center / admin / app:<ns>)
AudienceID string // §13.4 查询维度:投给谁(匿名 notify 为空)
Category string // §13.4 查询维度:按哪个品类触发
EventID string // 集成方事件身份,随 §13.5 回调回带
}
type DeliveryPayload struct {
Kind PayloadKind
Content *RenderedContent
ProviderTemplate *ProviderTemplatePayload
Raw map[string]any
}Queue(任务队列)
go
type Queue interface {
Push(ctx context.Context, task *DeliveryTask) error
Pop(ctx context.Context) (*DeliveryTask, error)
Ack(ctx context.Context, taskID string) error
Nack(ctx context.Context, taskID string, reason error) error
Size() int
Close() error
}模块列表
1. NotificationService
核心编排层,流程:
- 去重 — 基于内容的 SHA256 稳定 key
- 路由 — 根据 type/level 解析目标渠道
- 模板渲染 — 替换模板变量
- 投递规划 — 为每个渠道生成 DeliveryTask
- 入队 — 推送到 Queue
2. DeliveryPlanner
根据 Provider 能力和模板 Binding 生成 DeliveryTask:
- SMS Provider →
ProviderTemplatePayload - 内容 Provider →
RenderedContent
3. Template + Binding
一个业务模板可映射到多个渠道配置:
yaml
templates:
server_alert:
bindings:
email: { format: html }
telegram: { format: markdown }
aliyunsms:
template_code: "SMS_123456"
params: { "主机": "host", "状态": "status" }4. Route Engine
负责:Notification Type/Level → Provider
5. Retry
接入 runtime.Manager.Deliver(),支持指数退避重试。
6. Dedup
基于内容的稳定 key(SHA256),窗口 5 分钟。
7. Queue
队列抽象,当前实现:
| 实现 | 适用场景 | Remote Worker |
|---|---|---|
memory | 单机部署 | 不支持 |
redis | 分布式部署 | 支持 |
8. WebSocket Hub
远程 Worker 的管理通道(非任务分发通道):
- Worker 注册与能力记录
- 心跳保活与超时注销
- 状态事件上报
9. Provider Capability
每个 Provider 声明自己的能力:
go
type ProviderCapability struct {
PayloadKinds []PayloadKind
ContentFormats []string
SupportsBatch bool
SupportsTemplate bool
}