Skip to content

模块设计 ​

核心模块 ​

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"]

核心职责:

  1. 启动 N 个 local worker goroutine
  2. 每个 worker 循环:Pop → Deliver → Ack/Nack
  3. 管理远程 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 ​

核心编排层,流程:

  1. 去重 — 基于内容的 SHA256 稳定 key
  2. 路由 — 根据 type/level 解析目标渠道
  3. 模板渲染 — 替换模板变量
  4. 投递规划 — 为每个渠道生成 DeliveryTask
  5. 入队 — 推送到 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
}