# 模块理解确认书 — push-gateway > AI:ai01(Go 网关层) > 阶段:阶段 1 交付物 > 日期:2026-07-09 > 关联:[004 架构影响地图](../../../docs/architecture/004_architecture_impact_map.md)、[AI 分配方案](../../../docs/architecture/ai-allocation.md) --- ## 1. 我在架构中的位置 - **层级**:L3 网关层(004 §3.1 六层架构,与 api-gateway 并列) - **上游**: - 4 个微前端(浏览器/移动端通过 WebSocket 长连接接入) - msg 服务(通过 HTTP `/internal/push` `/internal/broadcast` 调用) - **下游**: - msg 服务(P5 设计重点之一是与 msg 的 gRPC 推送通道协议;当前骨架走 HTTP) - iam 服务(拉取 JWKS 公钥校验 JWT,P2 起) - Redis(跨实例 Pub/Sub 广播,P5 待实现) - **通信方式**: - 入口:WebSocket(`gorilla/websocket`)+ HTTP(内部 API) - 出口:HTTP(调 msg)/ Redis Pub/Sub(跨实例广播) - **不持有业务状态**:仅持有 WebSocket 连接池(in-memory `map[userID]map[*Connection]bool`),无 DB ## 2. 我的限界上下文 - **我负责**: - WebSocket 长连接生命周期(认证 → 心跳 → 断线清理) - 在线用户 session 映射(userID → 多个 WebSocket 连接) - 接收 msg 服务的推送请求,投递到目标用户的在线连接 - 跨实例广播(P5 待实现,Redis Pub/Sub) - **聚合/实体**:无业务聚合,仅 `Connection` 值对象(UserID + conn + send chan) - **业务领域**:不属于 D1-D6,属于基础设施层 - **我不负责**: - 通知业务逻辑(msg 服务负责:通知模板、已读/未读、通知偏好) - 推送内容持久化(msg 服务落 MySQL) - 离线推送(msg 服务通过 SMS/邮件渠道补推) - JWT 签发(iam 服务签发,本服务只校验) ## 3. 我与外部的契约 ### 3.1 我消费的 proto message | proto | message | 用途 | | --------- | -------------- | ----------------------------------------------------------------------------------- | | msg.proto | `Notification` | 接收 msg 服务的推送请求时,结构对齐(虽走 HTTP JSON,但字段与 `Notification` 一致) | | iam.proto | `UserInfo` | P2 起 RS256 校验 JWT(通过 JWKS 端点,非 gRPC) | > 当前骨架用 HTTP JSON 推送(`{userId, event, data}`),**ai-allocation §5 设计重点要求**与 msg 的 gRPC 推送通道协议。阶段 2 设计文档须明确 gRPC 契约(`PushService.PushToUser` / `PushService.Broadcast`)。 ### 3.2 我暴露的 API 端点 | 方法 | 路径 | 鉴权 | 说明 | | ---- | --------------------- | ------------------------------------------- | ------------------------------ | | GET | `/healthz` | 无 | liveness | | GET | `/readyz` | 无(待实现) | readiness(应检查 Redis 连接) | | GET | `/ws?token=JWT` | JWT(query 参数或 `Authorization: Bearer`) | WebSocket 升级端点 | | POST | `/internal/push` | 待补(内部 API,应校验调用方身份) | msg 服务定向推送 | | POST | `/internal/broadcast` | 待补 | msg 服务广播 | | GET | `/metrics` | 无(待实现) | Prometheus 指标 | ### 3.3 我发布/消费的 Kafka 事件 **当前**:无(骨架用 HTTP `/internal/*`)。 **P5 设计目标**(ai-allocation §5): - **消费**:`edu.notification.events`(msg 服务发布的 `NotificationRequested` 事件),消费后投递到对应用户的在线连接 - **不发布**任何 Kafka 事件(推送结果通过 gRPC response 返 msg) > 阶段 2 设计文档须明确:是 HTTP `/internal/push` 还是 Kafka 消费?建议**双通道**: > > - HTTP `/internal/push`:msg 同步调用,需立即知道推送结果(用户在线/离线) > - Kafka 消费:广播类通知(全校广播),解耦 + 削峰 ### 3.4 错误码前缀 | 错误码 | HTTP | 触发条件 | | ----------------- | ---- | --------------------------------------------------- | | `UNAUTHORIZED` | 401 | 缺失 token / token 无效 / claims 解析失败 | | `INVALID_REQUEST` | 400 | 请求体 JSON 解析失败 | | `INVALID_PAYLOAD` | 400 | event/data 字段缺失或格式错误 | | `PUSH_FAILED` | 500 | 投递失败(当前不会触发,`SendToUser` 静默返回 nil) | | `INTERNAL_ERROR` | 500 | panic 兜底 | ## 4. 我的技术栈 - **语言**:Go 1.22+(go.mod 声明 1.25.0,需与 Dockerfile 对齐) - **框架**:Gin v1.12.0 - **核心依赖**: - `github.com/gorilla/websocket` v1.5.3(WebSocket) - `github.com/golang-jwt/jwt/v5` v5.2.1(JWT 校验) - `go.opentelemetry.io/otel` v1.44.0 + `otelgin` v0.69.0 - **存储**:无 DB;Redis(P5 待引入,用于跨实例 Pub/Sub) - **构建**:Dockerfile 单阶段(builder→runner 但未做非 root 用户、healthcheck、ldflags 优化) ## 5. 我的阶段归属 - **阶段**:P5(沟通与 AI 阶段) - **当前阶段目标**(pending-features §P5): - WebSocket 长连接管理 + 消费 Kafka 广播 + Redis PubSub 跨实例同步 - 单节点支撑 10w+ 连接 - 全校广播推送实时到达 - **依赖上游阶段产出**: - P2 iam 服务(JWT RS256 签发 + JWKS 端点) - P5 msg 服务(推送请求发起方) - P1 Redis 基础设施(已就绪) ## 6. 我需要对齐的黄金模板项(对照 classes 服务) | 项 | classes | push-gateway 现状 | 差距 | | ----------------- | ------------------------------ | ----------------------------------------- | ------------------------------------ | | 权限装饰器 | `@RequirePermission()` | N/A(基础设施,JWT 校验 + 内部 API 鉴权) | ⚠️ 内部 API `/internal/*` 当前无鉴权 | | 错误码前缀 | `CLASSES_*` | 无前缀 | ⚠️ 待统一为 `PUSH_*` | | logger | pino | ❌ 标准 `log` | ⚠️ 待补 `log/slog` | | metrics | `/metrics` | ❌ 无 | ⚠️ 待补:连接数/推送 QPS/心跳 | | tracer | OTel SDK | ✅ | ✅ 对齐 | | `/healthz` | ✅ | ✅ | ✅ 对齐 | | `/readyz` | ✅ 检查 DB | ❌ 无 | ⚠️ 待补:检查 Redis 连接 | | 优雅关闭 | SIGTERM → app.close() | ✅ 5s 超时 | ⚠️ 但未关闭所有 WebSocket 连接 | | 测试覆盖率 | ≥ 80% | 0% | ⚠️ 待补 hub/ws/config 测试 | | Dockerfile | 多阶段 + 非 root + healthcheck | ❌ 单阶段 + root + 无 healthcheck | ⚠️ 待重构 | | Zod 输入验证 | `safeParse` | ✅ `ShouldBindJSON` 等价 | ✅ 对齐 | | GlobalErrorFilter | `GlobalErrorFilter` | ⚠️ `gin.Recovery()` 默认 | ⚠️ 待自定义 JSON 信封错误 | ## 7. 服务审计表 | 服务 | 权限装饰器 | 错误码前缀 | logger | metrics | tracer | /healthz | /readyz | 优雅关闭 | 测试覆盖率 | Dockerfile | | ------------ | ------------------ | ---------- | ----------- | ------- | ------- | -------- | ------- | ---------- | ---------- | ---------- | | push-gateway | ⚠️ 内部 API 无鉴权 | ❌ 无前缀 | ❌ 标准 log | ❌ 无 | ✅ OTel | ✅ | ❌ 无 | ⚠️ 未关 WS | 0% | ❌ 单阶段 | ### 7.1 详细问题清单(按严重度排序) | # | 严重度 | 文件 | 问题 | 修复建议 | | --- | ------ | ------------------------------- | ------------------------------------------------------------------------------------ | -------------------------------------------------------------------------------------------- | | 1 | 高 | 全局 | 无 Redis Pub/Sub,多实例部署时用户连接分散在各自实例,跨实例推送失败 | 引入 Redis Pub/Sub:本实例订阅 `edu.push.user.` channel,跨实例推送通过 PUBLISH 投递 | | 2 | 高 | `internal/ws/handler.go` L22-26 | `CheckOrigin` 永远返回 true,CSRF/Origin 校验缺失 | 从环境变量 `WS_ALLOWED_ORIGINS` 读取白名单 | | 3 | 高 | `internal/ws/handler.go` L67-75 | 心跳用文本消息 `ping/pong`,不符合 RFC 6455 控制帧规范;无空闲超时,僵尸连接永久占用 | 改用 `conn.SetReadDeadline` + WebSocket Ping/Pong 控制帧;60s 无心跳关闭连接 | | 4 | 高 | `internal/ws/handler.go` L19 | `Send` 通道满时静默丢弃,无指标无日志 | 增加 `messages_dropped_total` 指标 + slog.Warn | | 5 | 高 | `internal/hub/hub.go` L44-57 | 无单用户最大连接数限制,DoS 风险 | 限制每用户 ≤ 5 连接,超限拒绝新连接或踢掉最旧 | | 6 | 高 | `main.go` L44-46 | `/internal/*` 无鉴权,内网任何服务可调用 | 引入 `X-Internal-Token` 头校验(与 msg 服务共享密钥) | | 7 | 中 | `internal/ws/handler.go` L75 | 无 reconnect 协议,客户端断线后无法恢复未送达消息 | 协议层支持 `session_id` + `last_seq`,重连后 msg 服务补推(与 msg 协调) | | 8 | 中 | 全局 | HS256 而非 RS256 | P2 起改 RS256,与 api-gateway 共用 JWKS 拉取逻辑 | | 9 | 中 | `internal/config/config.go` | `getEnvBool` 用 `=="true"` 字符串比较,不健壮 | 用 `strconv.ParseBool` | | 10 | 中 | `Dockerfile` | 单阶段构建、root 用户、无 healthcheck、无 ldflags | 重构为多阶段 + 非 root + `wget /healthz` healthcheck | | 11 | 中 | `main.go` L66-71 | 优雅关闭未关闭所有 WebSocket 连接 | Hub 新增 `CloseAll()` 方法,Shutdown 前调用 | | 12 | 中 | 全局 | 无 `/metrics` 端点 | 新增 `metrics.go`:连接数 Gauge / 推送 QPS Counter / 心跳 Counter / 丢弃 Counter | | 13 | 低 | `go.mod` L3 vs `Dockerfile` L1 | go.mod `1.25.0` vs Dockerfile `golang:1.22-alpine` | 统一 | | 14 | 低 | `internal/hub/hub.go` L17-23 | `Send` 非阻塞丢弃,但 `Outgoing` 暴露 chan 给外部读,封装不彻底 | 改为 `Connection` 内部启动写协程,外部只调 `Send` | | 15 | 低 | `README.md` L17 | 写 "Gin 1.10",实际 go.mod 是 v1.12.0 | 修正 | | 16 | 低 | 测试 | 0% 覆盖率 | 补 hub_test.go(注册/注销/推送/广播)、ws_test.go(鉴权/心跳) | ## 8. 风险与假设 - **假设**:P5 msg 服务会通过 gRPC 或 HTTP 调用本服务推送(双通道待设计文档定) - **假设**:Redis 在 P5 部署时可用(P1 已就绪) - **风险**:单实例 in-memory 连接池在水平扩展时失效,必须 Redis Pub/Sub - **风险**:僵尸连接耗尽资源(无空闲超时) - **风险**:内部 API 无鉴权,内网穿透后可被滥用 - **未决**:是否在 push-gateway 消费 Kafka 广播?建议是,避免 msg 服务同步调用造成推送延迟 - **未决**:reconnect 协议是否在本期实现?建议本期仅做基础协议,reconnect 协议作为 P6 硬化项