Files
Edu/services/push-gateway

Push Gateway 推送网关服务

版本:1.0(P5 仲裁回写) 端口:8081 关联:02 架构设计、modules/push-gateway/README

职责

实时推送基础设施服务(Go 实现),管理 WebSocket 长连接(仅 WebSocket,不支持 SSE — ISSUE-001)。 接收 msg 服务的推送请求(HTTP /internal/* + Kafka edu.notification.requested),维护用户在线连接池, 将消息实时投递到浏览器/移动端。无 DB(ISSUE-005),审计/落库由 msg 负责。

技术栈

  • Go 1.25+ / Gin 1.12+
  • gorilla/websocket 1.5(RFC 6455 控制帧心跳)
  • segmentio/kafka-go 0.4(纯 Go,无 cgo)
  • redis/go-redis/v9(在线 SET + Pub/Sub 跨实例 fanout)
  • golang-jwt/jwt/v5 + shared-go/jwks(RS256 校验,5min 缓存)
  • log/slog(结构化日志,JSON prod / text dev)
  • prometheus/client_golang(9 个指标)
  • shared-go(tracer/logger/env/jwks 复用 — ARB-015 §17.5)

开发

go mod tidy
go run main.go   # :8081

# DevMode(跳过内部鉴权,接受 dev-token)
DEV_MODE=true go run main.go

# 生产模式(需设置 JWT_SECRET + INTERNAL_API_TOKEN)
JWT_SECRET=xxx INTERNAL_API_TOKEN=yyy go run main.go

API

方法 路径 鉴权 说明
GET /ws?token=JWT JWT RS256 WebSocket 升级(亦支持 Authorization: Bearer)
POST /internal/push X-Internal-Token 定向推送(msg 服务调用)
POST /internal/broadcast X-Internal-Token 广播
GET /internal/online/<userID> X-Internal-Token 查在线状态
GET /healthz 无 liveness 探针
GET /readyz 无 readiness 探针(软失败 — ISSUE-058)
GET /metrics 无 Prometheus 指标

环境变量

变量 默认值 说明
PUSH_GATEWAY_PORT 8081 HTTP 监听端口
DEV_MODE false 开发模式(跳过内部鉴权,接受 dev-token)
JWT_SECRET (DevMode 有默认) JWT 密钥(生产必填)
JWKS_URL http://localhost:50052/.well-known/jwks.json iam JWKS 端点
INTERNAL_API_TOKEN (生产必填) /internal/* 鉴权 token(ISSUE-002)
REDIS_URL redis://localhost:6379/0 Redis 连接 URL
KAFKA_BROKERS localhost:9092 Kafka broker 列表(逗号分隔)
KAFKA_NOTIFICATION_TOPIC edu.notification.requested 消费 topic(ISSUE-053)
KAFKA_CONSUMER_GROUP push-gateway consumer group
WS_ALLOWED_ORIGINS (空) WebSocket Origin 白名单(逗号分隔)
MAX_CONNS_PER_USER 5 单用户最大连接数
HEARTBEAT_INTERVAL_SECONDS 30 心跳间隔
INSTANCE_ID (hostname) 实例标识(Redis SET 成员)
OTEL_EXPORTER_OTLP_ENDPOINT localhost:4318 OTLP 端点

WebSocket 心跳(RFC 6455 控制帧)

  • 客户端每 30s 发送 WebSocket Ping 控制帧(非文本消息)
  • 服务端自动回 Pong(gorilla/websocket 默认行为)
  • 服务端 SetReadDeadline(60s):60s 未收到任何帧则断开
  • 每次 Pong 刷新 Redis 在线 SET 的 TTL(60s)

内部包结构

internal/
├─ config/           # 环境变量配置(shared-go/env)
├─ hub/              # 连接池(userID -> connID -> Connection)
├─ ws/               # WebSocket 升级 + /internal/* HTTP API
├─ redisclient/      # Redis Pub/Sub + 在线 SET(ISSUE-058 启动重建)
├─ kafkaconsumer/    # Kafka 消费(edu.notification.requested)
├─ observability/    # slog logger + OTel tracer + Prometheus metrics
└─ health/           # /healthz + /readyz(软失败)

优雅关闭顺序

  1. 收到 SIGTERM → 标记 Hub closing(拒绝新连接,/readyz 返 503)
  2. Hub.CloseAll():向所有在线连接发送 close 帧(code=1001 going away)
  3. srv.Shutdown(10s):停止 HTTP 服务,drain 在途请求
  4. Kafka consumer cancel + reader.Close
  5. Redis Pub/Sub cancel + client.Close
  6. Tracer shutdown(flush span)

关键设计决策

  • 软失败(ISSUE-058/006):Redis/Kafka 故障时 /readyz 返 200 + degraded:true,不断本地 WebSocket
  • Redis SET 启动重建(ISSUE-058):启动时 SREM 旧 instanceID,基于内存重建
  • 无 gRPC(ISSUE-007):仅 HTTP + WebSocket
  • shared-go 复用(ARB-015 §17.5):tracer/logger/env/jwks 不重复实现

完整 ADR 见 02 架构设计 §14。