Skip to content

Worker 架构 ​

统一 Worker 模型 ​

Herald 中所有 Worker 都是同一种概念,通过 local/remote 区分部署方式:

mermaid
graph TB
    subgraph Scheduler["Herald Scheduler"]
        API["HTTP API"]
    end

    Queue["Queue<br/>memory / redis"]

    subgraph Local["Local Workers"]
        LW1["goroutine"]
        LW2["goroutine"]
    end

    subgraph Remote["Remote Workers"]
        RW1["独立进程"]
    end

    API -->|Push| Queue
    Queue -->|Pop + Ack| LW1
    Queue -->|Pop + Ack| LW2
    Queue -->|Pop + Ack| RW1
    RW1 -.->|WebSocket<br/>注册/心跳| Scheduler
属性Local WorkerRemote Worker
部署进程内 goroutine独立进程(heraldd worker)
任务来源从 Queue Pop从共享 Queue Pop(Redis)
通信直接调用 ProviderWebSocket 注册 + Queue 消费
适用内置 Provider(Telegram、Email 等)需要特殊环境的 Provider

启动方式 ​

bash
# 调度器模式(自动启动 local workers)
heraldd serve --config config.yaml

# 远程 Worker 模式
heraldd worker --config worker.yaml

Worker 注册表 ​

所有 Worker(local 和 remote)都注册到统一的 Registry:

go
type Info struct {
    ID            string    // worker ID
    Mode          Mode      // "local" 或 "remote"
    Capabilities  []string  // 支持的 Provider 类型
    Status        string    // online / offline
    ConnectedAt   time.Time
    LastHeartbeat time.Time
}
  • Local Worker:启动时自动注册,进程退出时自动注销
  • Remote Worker:启动时通过 WebSocket 注册,定期心跳保活,超时自动注销

Queue 接口 ​

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
}
实现适用场景Remote Worker 支持
memory单机部署不支持
redis分布式部署支持

WebSocket 管理通道 ​

WebSocket 不用于任务分发,仅作为远程 Worker 的管理通道:启动时上报 ID 与能力完成注册,之后按固定间隔心跳保活、超时自动注销;Worker 状态变化(session 过期、需要扫码等)也走这条通道上报。

注册消息(Worker → Core) ​

json
{
  "type": "register",
  "worker_id": "wechat-worker-01",
  "mode": "remote",
  "platform": "linux",
  "version": "1.0.0",
  "capabilities": ["wechatmp", "wechat"]
}

心跳消息(Worker → Core) ​

json
{
  "type": "heartbeat",
  "worker_id": "wechat-worker-01",
  "timestamp": 1716780000
}

status 是协议里的可选自由映射(map[string]interface{},protocol/message.go),当前 SDK 与内置 worker 都不填充,示例只发 worker_id + timestamp。

协议消息 ​

注册确认(Core → Worker) ​

json
{
  "type": "register_ack",
  "worker_id": "wechat-worker-01",
  "success": true,
  "server_id": "herald-core-01",
  "timestamp": 1716780000
}

事件消息(Worker → Core) ​

json
{
  "type": "event",
  "worker_id": "wechat-worker-01",
  "event_type": "offline",
  "timestamp": 1716780000,
  "data": {
    "reason": "session_expired"
  }
}

适用场景 ​

Remote Worker 适用于:浏览器自动化(微信公众号、网页版工具)、移动端桥接(Android 通知桥接)、需要 GUI Session 或特殊系统权限的场景,以及想把渠道崩溃隔出调度器进程的部署。

最佳实践 ​

开发和小规模部署用 memory 队列即可;要上分布式再切 redis,代码零改动。Remote Worker 独立进程运行,崩溃不牵连调度器;断线后自动重连并重新注册,心跳掉线能及时暴露连接问题。