Files
Edu/services/push-gateway/docs/02-architecture-design.md
SpecialX a0f6c228af docs(push-gateway): ai01 阶段1+2 模块理解确认书与架构设计文档
阶段1交付:services/push-gateway/docs/01-understanding.md
- 8节:架构位置/限界上下文/契约/技术栈/阶段归属P5骨架/审计
- 审计16项差距(6高/5中/5低)

阶段2交付:services/push-gateway/docs/02-architecture-design.md
- 13节:分层图/Connection+Hub模型/Redis key/WS端点
- 双通道协议/生命周期/心跳RFC6455/多实例PubSub
- 容量目标10w+连接

AI Agent: ai01 (api-gateway/push-gateway)
Branch: main
Coordinator: coord
2026-07-09 17:22:23 +08:00

24 KiB
Raw Blame History

模块架构设计文档 — push-gateway

AIai01Go 网关层) 阶段:阶段 2 交付物 日期2026-07-09 关联:01 理解确认书004 架构影响地图

本文档覆盖 ai-allocation §5 设计重点WebSocket 连接生命周期、与 msg 的 gRPC 推送通道协议、用户 session 映射、水平扩展方案Redis Pub/Sub 跨实例广播)。


1. 模块内部分层图

graph TB
    subgraph Client["客户端"]
        B[浏览器/移动端]
    end

    subgraph External["外部入口层"]
        WS[WebSocket Upgrader<br/>/ws?token=JWT]
    end

    subgraph MW["中间件链"]
        M1[Recovery]
        M2[OTelgin]
        M3[RequestID]
        M4[OriginGuard<br/>CheckOrigin 白名单]
        M5[InternalAuth<br/>/internal/* 校验]
    end

    subgraph Hub["Hub 连接池(核心)"]
        H[Hub<br/>map userID -> map connID -> Connection]
        Conn[Connection<br/>userID + ws + send chan + heartbeat]
    end

    subgraph Writer["写协程池"]
        W[每 Connection 一个写协程<br/>消费 send chan 写 ws]
    end

    subgraph Reader["读协程"]
        R[每 Connection 一个读协程<br/>处理 Ping 控制帧 + 消息]
    end

    subgraph PushAPI["内部推送 API"]
        P1[POST /internal/push]
        P2[POST /internal/broadcast]
    end

    subgraph Redis["Redis Pub/Sub横向扩展"]
        Sub[Subscriber<br/>订阅 edu.push.user.* + edu.push.broadcast]
        Pub[Publisher<br/>本实例无此用户时 PUBLISH]
    end

    subgraph Kafka["Kafka 消费者(可选)"]
        K[Consumer Group<br/>topic: edu.notification.events]
    end

    subgraph Msg["msg 服务"]
        MSG[msg :3007]
    end

    B -->|WebSocket| WS
    WS --> M1 --> M2 --> M3 --> M4 --> H
    H --> Conn
    Conn --> W
    Conn --> R
    R -->|Ping| Conn
    MSG -->|HTTP /internal/push| P1
    MSG -->|HTTP /internal/broadcast| P2
    P1 --> H
    P2 --> H
    P1 -.本地无此用户.-> Pub --> Sub
    Sub -.广播到本实例.-> H
    MSG -.NotificationRequested 事件.-> K
    K --> H

    style Redis fill:#fff4e6
    style Kafka fill:#fff4e6

关键设计决策

  • Hub 是单例,进程内唯一,持有所有连接
  • 每 Connection 两个 goroutine:读协程(消费 ws+ 写协程(消费 send chan
  • Redis Pub/Sub 解决跨实例广播:本实例无目标用户时 PUBLISH 到 edu.push.user.<userID> channel所有实例订阅此 channel持有该用户的实例收到后投递
  • Kafka 消费可选:广播类通知走 Kafka 解耦;定向推送走 HTTP /internal/push(同步返结果)

2. 领域模型

无业务聚合。仅连接管理值对象:

// Connection 包装单个 WebSocket 连接
type Connection struct {
    UserID    string
    ConnID    string                 // 连接唯一 IDuuid
    conn      *websocket.Conn
    send      chan []byte             // 异步发送通道cap 64
    closeOnce sync.Once
    closed    bool
    mu        sync.Mutex
    lastPing  time.Time              // 最后心跳时间
}

// Hub 管理所有在线连接
type Hub struct {
    mu       sync.RWMutex
    clients  map[string]map[string]*Connection  // userID -> connID -> Connection
    counters map[string]int                       // userID -> 连接计数(快速判断超限)
}

不变式

  • Hub.clients[userID][connID] 一一对应 Connection
  • Connection.send channel 关闭后不可再写
  • Connection.closed 为 true 后所有操作 no-op

3. 数据模型

无数据库。所有状态在内存 + Redis。

3.1 Redis 数据结构

Key 模式 类型 TTL 用途
edu:push:online:<userID> SET成员是 instanceID 60s心跳续期 在线用户所在实例集合(用于判断用户是否在线)
edu:push:session:<userID>:<connID> HASH 60s 单连接元数据instanceID, device, lastPing
edu:push:channel:user:<userID> Pub/Sub channel 跨实例定向推送 channel
edu:push:channel:broadcast Pub/Sub channel 跨实例广播 channel

3.2 在线状态判定

用户在线 := Redis SET edu:push:online:<userID> 非空
用户离线 := SET 为空或不存在

续期机制:每 30s 客户端心跳时,本实例 SADD 自己 + EXPIRE 60s。

4. API 设计

4.1 WebSocket 端点

方法 路径 鉴权 说明
GET /ws JWTquery ?token=Authorization: Bearer WebSocket 升级

WebSocket 子协议(应用层消息格式)

// 服务端 → 客户端:推送消息
{
  "type": "message",
  "event": "notification.created",
  "data": { ... },
  "seq": 12345,                  // 全局递增序列号(用于 reconnect
  "timestamp": "2026-07-09T..."
}

// 服务端 → 客户端:确认
{ "type": "ack", "seq": 12345 }

// 客户端 → 服务端心跳WebSocket Ping 控制帧,非文本消息)
// 客户端 → 服务端:消息确认
{ "type": "ack", "seq": 12345 }

4.2 内部推送 API

方法 路径 鉴权 请求体 响应
POST /internal/push X-Internal-Token {user_id, event, data, ttl?} {success, delivered, online}
POST /internal/broadcast X-Internal-Token {event, data, filter?} {success, reached}
GET /internal/online/<userID> X-Internal-Token {online: bool, instances: []}

X-Internal-Token 校验

  • 启动时从 INTERNAL_API_TOKEN 环境变量加载
  • 与 msg 服务共享同一密钥K8s Secret 注入)
  • 缺失/不匹配 → 401

响应语义

  • delivered: true:本实例或跨实例成功投递到至少一个连接
  • online: false用户离线msg 服务应走离线推送SMS/邮件)
  • delivered: false, online: true:投递失败(连接满/异常msg 服务应重试或落库

4.3 健康检查

端点 检查逻辑
/healthz 进程存活,返回 200
/readyz 检查 Redis 连接PING+ 在线连接数 > 0 时认为就绪Redis 不可达返回 503

4.4 指标端点

端点 暴露指标
/metrics 见 §6.4

5. 事件设计

5.1 我消费的 Kafka 事件

Topic Message 触发场景 消费动作
edu.notification.events NotificationRequested msg 服务收到通知请求后发布 消费 → 调用 Hub.SendToUser 投递 → 若离线则 ACK 不重投msg 服务落库)

消费语义

  • Consumer Grouppush-gateway
  • 至少一次at-least-once消费失败重试 3 次后入死信队列
  • 幂等性:基于 event_id 去重Redis SETNXTTL 24h

5.2 我发布的 Kafka 事件

。推送结果通过 HTTP /internal/push 同步响应返 msg 服务。

5.3 与 msg 服务的双通道协议

sequenceDiagram
    participant Msg as msg 服务
    participant GW as push-gateway
    participant Hub as Hub
    participant Redis as Redis Pub/Sub
    participant Client as 客户端

    rect rgb(240, 248, 255)
        Note over Msg,Client: 通道 1HTTP 同步推送(定向)
        Msg->>GW: POST /internal/push {user_id, event, data}
        GW->>Hub: SendToUser(user_id, msg)
        alt 本实例有此用户
            Hub->>Client: WebSocket 推送
            Client-->>Hub: ack
            Hub-->>GW: delivered=true
        else 本实例无此用户
            GW->>Redis: PUBLISH edu.push.channel.user.<userID> msg
            Note over Redis: 持有该用户的实例订阅后投递
            Redis-->>GW: 投递结果(异步,通过另一个 ack channel
            GW-->>Msg: delivered=true, online=true
        end
    end

    rect rgb(240, 255, 240)
        Note over Msg,Client: 通道 2Kafka 异步消费(广播)
        Msg->>Msg: 发布 NotificationRequested 事件到 Kafka
        Note over GW: Consumer Group 异步消费
        GW->>Hub: Broadcast(msg)
        Hub->>Client: WebSocket 推送(所有在线连接)
    end

双通道选择规则msg 服务决策):

  • 定向推送(单用户/少数用户)→ HTTP /internal/push(需立即知道在线/离线)
  • 广播(全校/全年级)→ Kafka 事件(解耦 + 削峰)
  • 紧急通知(考试开始)→ HTTP /internal/push(实时性优先)

6. 横切关注点对齐清单

6.1 权限校验

端点 鉴权方式
/ws JWT RS256P2 起,与 api-gateway 共用 JWKS 拉取逻辑)
/internal/* X-Internal-Token 头校验
/healthz /readyz /metrics

6.2 错误码清单

错误码 HTTP 触发条件
PUSH_UNAUTHORIZED 401 缺失/无效 token
PUSH_INVALID_REQUEST 400 JSON 解析失败
PUSH_INVALID_PAYLOAD 400 event/data 字段缺失
PUSH_TOO_MANY_CONNECTIONS 429 单用户连接数超限(>5
PUSH_INTERNAL_ERROR 500 panic / Redis 不可达

6.3 Logger

// internal/observability/logger.go待新增
var Logger *slog.Logger

func InitLogger(level string) {
    Logger = slog.New(slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{Level: parseLevel(level)}))
    slog.SetDefault(Logger)
}

日志字段timestamp level service=push-gateway request_id trace_id user_id conn_id event

6.4 Metrics 指标清单

指标名 类型 标签 描述
push_gateway_active_connections Gauge 当前在线连接数
push_gateway_connections_per_user Gauge 每用户连接数分布(仅暴露均值/最大)
push_gateway_messages_pushed_total Counter event, result 推送消息总数result=delivered/dropped/offline
push_gateway_messages_dropped_total Counter reason 丢弃消息数reason=channel_full/buffer_overflow
push_gateway_heartbeat_total Counter 心跳接收总数
push_gateway_disconnect_total Counter reason 断开连接数reason=idle/error/closed
push_gateway_redis_pubsub_latency_seconds Histogram direction Pub/Sub 延迟publish/subscribe
push_gateway_kafka_consumed_total Counter topic, partition Kafka 消费数
go_* prom-client 默认 Go runtime 指标

6.5 Tracer

已实现 OTel tracer。待补

  • WebSocket 升级 spanws.upgrade
  • 消息推送 spanpush.message,标签 user_id/event
  • Redis Pub/Sub spanredis.publish / redis.subscribe

6.6 /healthz 检查逻辑

func Healthz(c *gin.Context) {
    c.JSON(200, gin.H{
        "status":    "ok",
        "service":   "push-gateway",
        "version":   Version,
        "timestamp": time.Now().UTC().Format(time.RFC3339),
        "connections": hub.ActiveConnections(),
    })
}

6.7 /readyz 检查逻辑

func Readyz(redisClient *redis.Client) gin.HandlerFunc {
    return func(c *gin.Context) {
        ctx, cancel := context.WithTimeout(c.Request.Context(), 1*time.Second)
        defer cancel()
        if err := redisClient.Ping(ctx).Err(); err != nil {
            c.JSON(503, gin.H{"status":"error","error":"redis unreachable"})
            return
        }
        c.JSON(200, gin.H{"status":"ok"})
    }
}

6.8 优雅关闭顺序

// 1. 收到 SIGTERM
// 2. 标记 Hub 为 closing拒绝新连接
// 3. 向所有在线连接发送 close 帧code=1001 going away
// 4. 等待所有连接关闭10s 超时)
// 5. srv.Shutdown(ctx) 停止 HTTP 服务
// 6. 关闭 Kafka consumer提交 offset
// 7. 关闭 Redis subscriber
// 8. tracerShutdown() flush span
// 9. os.Exit(0)

7. WebSocket 连接生命周期(核心设计)

stateDiagram-v2
    [*] --> Connecting: 客户端 GET /ws?token=
    Connecting --> Authenticating: Upgrade 成功
    Authenticating --> Rejected: JWT 无效
    Rejected --> [*]: 返回 401
    Authenticating --> Connected: JWT 有效
    Connected --> Heartbeating: 注册到 Hub
    Heartbeating --> Heartbeating: 收到 Ping → 回 Pong
    Heartbeating --> Reconnecting: 客户端断线
    Heartbeating --> Closing: 服务端 SIGTERM
    Heartbeating --> IdleTimeout: 60s 无心跳
    IdleTimeout --> Closing: 服务端主动关闭
    Reconnecting --> Connected: 客户端重连(带 session_id + last_seq
    Closing --> [*]: 发送 close 帧 + 注销 Hub

7.1 心跳协议RFC 6455 控制帧)

规则

  • 客户端每 30s 发送 WebSocket Ping 控制帧
  • 服务端自动回 Ponggorilla/websocket 默认行为)
  • 服务端 conn.SetReadDeadline(60s)60s 未收到任何消息(含 Ping则关闭连接
  • 服务端每 30s 检查 lastPing,超过 60s 主动关闭

不使用文本消息 ping/pong(当前骨架的实现,需重构)。

7.2 单用户连接数限制

const MaxConnectionsPerUser = 5

func (h *Hub) Register(userID string, conn *websocket.Conn) (*Connection, error) {
    h.mu.Lock()
    defer h.mu.Unlock()
    if h.counters[userID] >= MaxConnectionsPerUser {
        return nil, ErrTooManyConnections
    }
    // ... 注册逻辑
}

超限策略:拒绝新连接,返回 close 帧code=1008 policy violation客户端应处理为"已达到最大连接数"。

7.3 Reconnect 协议P6 硬化项,本期仅预留)

协议

  • 客户端首次连接时,服务端返回 {type:"hello", session_id:"<uuid>", seq:0}
  • 后续每条推送消息带 seq 递增
  • 客户端断线重连时,/ws?token=&session_id=<id>&last_seq=<n>
  • 服务端从 msg 服务拉取 last_seq+1 到当前的所有未送达消息补推

本期不实现,仅协议预留。

7.4 Origin 校验

var upgrader = websocket.Upgrader{
    CheckOrigin: func(r *http.Request) bool {
        origin := r.Header.Get("Origin")
        if origin == "" {
            return false // 无 Origin 拒绝
        }
        for _, allowed := range allowedOrigins {
            if origin == allowed {
                return true
            }
        }
        return false
    },
}

白名单来源:环境变量 WS_ALLOWED_ORIGINS(逗号分隔)。

8. 横向扩展方案Redis Pub/Sub

8.1 多实例部署架构

graph TB
    subgraph LB["负载均衡器"]
        L[Sticky Session 或 IP Hash<br/>保证同一用户路由到同一实例]
    end

    subgraph Instances["push-gateway 实例集群"]
        I1[Instance 1<br/>持有用户 A,B,C]
        I2[Instance 2<br/>持有用户 D,E,F]
        I3[Instance 3<br/>持有用户 G,H,I]
    end

    subgraph Redis["Redis Pub/Sub"]
        R[(Redis<br/>edu.push.channel.user.*<br/>edu.push.channel.broadcast)]
    end

    subgraph Msg["msg 服务"]
        M[msg]
    end

    M -->|HTTP /internal/push user=A| I1
    M -->|HTTP /internal/push user=D| I2
    M -->|HTTP /internal/push user=Z 不知在哪| R
    R --> I1
    R --> I2
    R --> I3
    M -->|Kafka broadcast| K[Kafka]
    K --> I1
    K --> I2
    K --> I3

8.2 跨实例推送流程

  1. msg 服务调本实例 /internal/push {user_id: A}
  2. 本实例 Hub 查 clients[A]
    • 命中 → 直接投递
    • 未命中 → 查 Redis edu:push:online:A
      • 不存在 → 用户离线,返 delivered:false, online:false
      • 存在 → PUBLISH 到 edu.push.channel.user.A,所有实例订阅此 channel
  3. 持有用户 A 的实例收到 Pub/Sub 消息 → 投递到本地连接 → 通过另一 ack channel 返结果
  4. 本实例收到 ack → 返 msg 服务 delivered:true

8.3 广播优化

  • 不遍历所有用户:直接 PUBLISH 到 edu.push.channel.broadcast
  • 所有实例订阅此 channel收到后遍历本地连接投递
  • 单实例广播成本 O(local_connections),集群总成本 O(sum_connections)

8.4 在线状态一致性

  • Redis SET edu:push:online:<userID> 存储用户所在实例集合
  • 实例启动时生成唯一 instanceIDuuid
  • 连接建立时 SADD 自己到 SETEXPIRE 60s
  • 心跳续期时 EXPIRE 60s
  • 连接断开时 SREM 自己;若 SET 空DEL key
  • 实例崩溃时 SET 自然过期60s

9. 与其他模块的交互点(契约清单)

方向 对方服务 协议 接口/事件 用途
被调用 msg HTTP POST /internal/push 定向推送
被调用 msg HTTP POST /internal/broadcast 广播
被调用 msg HTTP GET /internal/online/<userID> 查在线状态
消费 msg Kafka edu.notification.events / NotificationRequested 异步广播通知
调用 iam HTTP GET /.well-known/jwks.json RS256 公钥P2
调用 Redis Redis Pub/Sub + SET 跨实例广播 + 在线状态
被调用 微前端 WebSocket /ws 长连接
被调用 K8s/Docker HTTP /healthz /readyz 探针
被调用 Prometheus HTTP GET /metrics 指标采集

10. 风险与假设

  • 假设Redis 在 P5 部署时可用P1 已就绪)
  • 假设msg 服务 P5 实现 NotificationRequested 事件发布
  • 假设iam P2 暴露 JWKS 端点
  • 风险:单实例内存连接池上限受 Go runtime 内存限制10w 连接约 1-2GB
  • 风险Redis Pub/Sub 是 fire-and-forget实例崩溃时未投递的消息丢失msg 服务落库兜底)
  • 风险:跨实例 ack 链路复杂,可能引入延迟(建议 P6 用 Redis Stream 替代 Pub/Sub 实现持久化)
  • 未决reconnect 协议是否本期实现建议本期仅协议预留P6 实现
  • 未决:是否在 push-gateway 做 SSE 支持建议不支持WebSocket 已够用,避免协议膨胀

11. 容量与性能目标

指标 目标 实现方式
单实例最大连接数 10w+ goroutine per connection + 内存优化
推送延迟(本实例) < 50ms 异步 send chan
推送延迟(跨实例) < 200ms Redis Pub/Sub
心跳间隔 30s 客户端 Ping
空闲超时 60s SetReadDeadline
单用户最大连接数 5 Hub.counters 限制
消息体上限 64KB WebSocket 消息大小限制

12. 实施优先级

优先级 任务 阶段
P0 Origin 校验CheckOrigin 白名单) P5 立即
P0 /internal/* 鉴权X-Internal-Token P5 立即
P0 心跳改用 WebSocket 控制帧 + SetReadDeadline P5 立即
P0 单用户连接数限制 P5 立即
P0 Send 通道满时记录指标+日志 P5 立即
P0 Dockerfile 重构(多阶段 + 非 root + healthcheck P5 立即
P0 引入 log/slog 替换标准 log P5 立即
P1 Redis Pub/Sub 跨实例广播 P5
P1 /metrics 端点 + prom-client P5
P1 /readyz 真实健康检查Redis PING P5
P1 优雅关闭所有 WebSocket 连接 P5
P2 Kafka 消费 edu.notification.events P5
P2 JWT RS256 升级(与 api-gateway 共用 JWKS P5依赖 iam
P2 在线状态 Redis SET 维护 P5
P3 Reconnect 协议session_id + last_seq P6
P3 Redis Stream 替代 Pub/Sub持久化 P6
P3 测试覆盖率 ≥ 80% P6

13. 与 api-gateway 共用的横切代码

为避免重复实现,建议以下模块提取到 packages/shared-go/(由 coord 维护ai01 提交需求):

模块 用途 当前位置
observability/tracer.go OTel tracer 初始化 两服务各有一份,完全相同
observability/logger.go slog 结构化日志 待新增(两服务共用)
auth/jwks.go JWKS 公钥拉取与缓存 待新增两服务共用P2
config/env.go 环境变量加载工具 两服务各有 getEnv可统一

提交方式:通过 # proto-change 渠道向 coord 声明需求coord 在 shared-go 中建立ai01 在本服务中改为 import。