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

574 lines
24 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 模块架构设计文档 — push-gateway
> AIai01Go 网关层)
> 阶段:阶段 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<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. 领域模型
无业务聚合。仅连接管理值对象:
```go
// 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 子协议(应用层消息格式)**
```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/<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 Group`push-gateway`
- 至少一次at-least-once消费失败重试 3 次后入死信队列
- 幂等性:基于 `event_id` 去重Redis SETNXTTL 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: 通道 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
```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 控制帧
- 服务端自动回 Ponggorilla/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:"<uuid>", seq:0}`
- 后续每条推送消息带 `seq` 递增
- 客户端断线重连时,`/ws?token=&session_id=<id>&last_seq=<n>`
- 服务端从 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<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>` 存储用户所在实例集合
- 实例启动时生成唯一 `instanceID`uuid
- 连接建立时 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。