# 模块架构设计文档 — push-gateway
> AI:ai01(Go 网关层)
> 阶段:阶段 2 交付物
> 日期:2026-07-09
> 关联:[01 理解确认书](./01-understanding.md)、[004 架构影响地图](../../../docs/architecture/004_architecture_impact_map.md)
>
> 本文档覆盖 ai-allocation §5 设计重点:WebSocket 连接生命周期、与 msg 的 gRPC 推送通道协议、用户 session 映射、水平扩展方案(Redis Pub/Sub 跨实例广播)。
---
## 1. 模块内部分层图
```mermaid
graph TB
subgraph Client["客户端"]
B[浏览器/移动端]
end
subgraph External["外部入口层"]
WS[WebSocket Upgrader
/ws?token=JWT]
end
subgraph MW["中间件链"]
M1[Recovery]
M2[OTelgin]
M3[RequestID]
M4[OriginGuard
CheckOrigin 白名单]
M5[InternalAuth
/internal/* 校验]
end
subgraph Hub["Hub 连接池(核心)"]
H[Hub
map userID -> map connID -> Connection]
Conn[Connection
userID + ws + send chan + heartbeat]
end
subgraph Writer["写协程池"]
W[每 Connection 一个写协程
消费 send chan 写 ws]
end
subgraph Reader["读协程"]
R[每 Connection 一个读协程
处理 Ping 控制帧 + 消息]
end
subgraph PushAPI["内部推送 API"]
P1[POST /internal/push]
P2[POST /internal/broadcast]
end
subgraph Redis["Redis Pub/Sub(横向扩展)"]
Sub[Subscriber
订阅 edu.push.user.* + edu.push.broadcast]
Pub[Publisher
本实例无此用户时 PUBLISH]
end
subgraph Kafka["Kafka 消费者(可选)"]
K[Consumer Group
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.` channel,所有实例订阅此 channel,持有该用户的实例收到后投递
- **Kafka 消费可选**:广播类通知走 Kafka 解耦;定向推送走 HTTP `/internal/push`(同步返结果)
## 2. 领域模型
无业务聚合。仅连接管理值对象:
```go
// Connection 包装单个 WebSocket 连接
type Connection struct {
UserID string
ConnID string // 连接唯一 ID(uuid)
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:` | SET(成员是 `instanceID`) | 60s(心跳续期) | 在线用户所在实例集合(用于判断用户是否在线) |
| `edu:push:session::` | HASH | 60s | 单连接元数据(instanceID, device, lastPing) |
| `edu:push:channel:user:` | Pub/Sub channel | — | 跨实例定向推送 channel |
| `edu:push:channel:broadcast` | Pub/Sub channel | — | 跨实例广播 channel |
### 3.2 在线状态判定
```
用户在线 := Redis SET edu:push:online: 非空
用户离线 := SET 为空或不存在
```
**续期机制**:每 30s 客户端心跳时,本实例 SADD 自己 + EXPIRE 60s。
## 4. API 设计
### 4.1 WebSocket 端点
| 方法 | 路径 | 鉴权 | 说明 |
| ---- | ----- | ------------------------------------------------- | -------------- |
| GET | `/ws` | JWT(query `?token=` 或 `Authorization: Bearer`) | WebSocket 升级 |
**WebSocket 子协议(应用层消息格式)**:
```json
// 服务端 → 客户端:推送消息
{
"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/` | `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 Group:`push-gateway`
- 至少一次(at-least-once),消费失败重试 3 次后入死信队列
- 幂等性:基于 `event_id` 去重(Redis SETNX,TTL 24h)
### 5.2 我发布的 Kafka 事件
**无**。推送结果通过 HTTP `/internal/push` 同步响应返 msg 服务。
### 5.3 与 msg 服务的双通道协议
```mermaid
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: 通道 1:HTTP 同步推送(定向)
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. 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: 通道 2:Kafka 异步消费(广播)
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 RS256(P2 起,与 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
```go
// 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 升级 span(`ws.upgrade`)
- 消息推送 span(`push.message`,标签 user_id/event)
- Redis Pub/Sub span(`redis.publish` / `redis.subscribe`)
### 6.6 /healthz 检查逻辑
```go
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 检查逻辑
```go
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 优雅关闭顺序
```go
// 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 连接生命周期(核心设计)
```mermaid
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 控制帧
- 服务端自动回 Pong(gorilla/websocket 默认行为)
- 服务端 `conn.SetReadDeadline(60s)`:60s 未收到任何消息(含 Ping)则关闭连接
- 服务端每 30s 检查 `lastPing`,超过 60s 主动关闭
**不使用文本消息 `ping/pong`**(当前骨架的实现,需重构)。
### 7.2 单用户连接数限制
```go
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:"", seq:0}`
- 后续每条推送消息带 `seq` 递增
- 客户端断线重连时,`/ws?token=&session_id=&last_seq=`
- 服务端从 msg 服务拉取 `last_seq+1` 到当前的所有未送达消息补推
**本期不实现**,仅协议预留。
### 7.4 Origin 校验
```go
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 多实例部署架构
```mermaid
graph TB
subgraph LB["负载均衡器"]
L[Sticky Session 或 IP Hash
保证同一用户路由到同一实例]
end
subgraph Instances["push-gateway 实例集群"]
I1[Instance 1
持有用户 A,B,C]
I2[Instance 2
持有用户 D,E,F]
I3[Instance 3
持有用户 G,H,I]
end
subgraph Redis["Redis Pub/Sub"]
R[(Redis
edu.push.channel.user.*
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:` 存储用户所在实例集合
- 实例启动时生成唯一 `instanceID`(uuid)
- 连接建立时 SADD 自己到 SET,EXPIRE 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/` | 查在线状态 |
| 消费 | 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。