feat(p5): messaging, push gateway and AI assistant services
P5 阶段交付物: - services/msg: 消息通知服务(NestJS) - notifications: 发送通知 + ES 全文检索 + search - config/elasticsearch.ts: ES Client 单例 - package.json: 补充 @opentelemetry/sdk-node + exporter-trace-otlp-http - services/push-gateway: WebSocket 推送网关(Go Gin) - internal/hub/hub.go: WebSocket 连接池管理(Register/Unregister/SendToUser) - internal/ws/handler.go: JWT 鉴权 + WebSocket 升级 + 内部推送 API - services/ai: AI 辅助服务(Python FastAPI) - /chat + /chat/stream(SSE 流式) - /generate/question + /optimize/expression - config.py: OpenAI 兼容 API 配置 - packages/shared-proto/proto/msg.proto: NotificationService 契约(send/search) - packages/shared-proto/proto/ai.proto: AiService 契约(含 stream 方法)
This commit is contained in:
13
services/push-gateway/Dockerfile
Normal file
13
services/push-gateway/Dockerfile
Normal file
@@ -0,0 +1,13 @@
|
||||
FROM golang:1.22-alpine AS builder
|
||||
WORKDIR /app
|
||||
COPY go.mod go.sum ./
|
||||
RUN go mod download
|
||||
COPY . .
|
||||
RUN CGO_ENABLED=0 go build -o push-gateway .
|
||||
|
||||
FROM alpine:3.20
|
||||
RUN apk --no-cache add ca-certificates
|
||||
WORKDIR /app
|
||||
COPY --from=builder /app/push-gateway .
|
||||
EXPOSE 8081
|
||||
CMD ["./push-gateway"]
|
||||
42
services/push-gateway/README.md
Normal file
42
services/push-gateway/README.md
Normal file
@@ -0,0 +1,42 @@
|
||||
# Push Gateway 推送网关服务
|
||||
|
||||
> 版本:0.1(P5 骨架)
|
||||
> 端口:8081
|
||||
|
||||
## 职责
|
||||
|
||||
实时推送基础设施服务(Go 实现),管理 WebSocket 长连接。
|
||||
接收 Msg 服务的推送请求,维护用户在线连接池,将消息实时投递到浏览器/移动端。
|
||||
|
||||
## 技术栈
|
||||
|
||||
- Go 1.22 + Gin 1.10
|
||||
- gorilla/websocket 1.5
|
||||
- golang-jwt/jwt/v5(JWT 鉴权)
|
||||
- zap(结构化日志,骨架)
|
||||
|
||||
## 开发
|
||||
|
||||
```bash
|
||||
go mod tidy
|
||||
go run main.go # :8081
|
||||
```
|
||||
|
||||
## API
|
||||
|
||||
| 方法 | 路径 | 说明 |
|
||||
|------|------|------|
|
||||
| GET | /healthz | 健康检查 |
|
||||
| GET | /ws?token=JWT | WebSocket 升级端点(JWT 鉴权) |
|
||||
| POST | /internal/push | 内部推送 API(Msg 服务调用) |
|
||||
|
||||
## 环境变量
|
||||
|
||||
| 变量 | 默认值 | 说明 |
|
||||
|------|--------|------|
|
||||
| PUSH_GATEWAY_PORT | 8081 | 服务端口 |
|
||||
| JWT_SECRET | p1-dev-secret-change-in-production | JWT 密钥 |
|
||||
|
||||
## WebSocket 心跳
|
||||
|
||||
客户端发送 `ping` 文本帧,服务端回复 `pong`。
|
||||
11
services/push-gateway/go.mod
Normal file
11
services/push-gateway/go.mod
Normal file
@@ -0,0 +1,11 @@
|
||||
module github.com/edu-cloud/push-gateway
|
||||
|
||||
go 1.22
|
||||
|
||||
require (
|
||||
github.com/gin-gonic/gin v1.10.0
|
||||
github.com/golang-jwt/jwt/v5 v5.2.1
|
||||
github.com/google/uuid v1.6.0
|
||||
github.com/gorilla/websocket v1.5.3
|
||||
go.uber.org/zap v1.27.0
|
||||
)
|
||||
22
services/push-gateway/internal/config/config.go
Normal file
22
services/push-gateway/internal/config/config.go
Normal file
@@ -0,0 +1,22 @@
|
||||
package config
|
||||
|
||||
import "os"
|
||||
|
||||
type Config struct {
|
||||
Port string
|
||||
JWTSecret string
|
||||
}
|
||||
|
||||
func Load() *Config {
|
||||
return &Config{
|
||||
Port: getEnv("PUSH_GATEWAY_PORT", "8081"),
|
||||
JWTSecret: getEnv("JWT_SECRET", "p1-dev-secret-change-in-production"),
|
||||
}
|
||||
}
|
||||
|
||||
func getEnv(key, fallback string) string {
|
||||
if v := os.Getenv(key); v != "" {
|
||||
return v
|
||||
}
|
||||
return fallback
|
||||
}
|
||||
59
services/push-gateway/internal/hub/hub.go
Normal file
59
services/push-gateway/internal/hub/hub.go
Normal file
@@ -0,0 +1,59 @@
|
||||
package hub
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
)
|
||||
|
||||
// Hub 管理 WebSocket 客户端连接
|
||||
type Hub struct {
|
||||
mu sync.RWMutex
|
||||
clients map[string]map[*websocket.Conn]bool // userID -> connections
|
||||
}
|
||||
|
||||
func NewHub() *Hub {
|
||||
return &Hub{
|
||||
clients: make(map[string]map[*websocket.Conn]bool),
|
||||
}
|
||||
}
|
||||
|
||||
func (h *Hub) Register(userID string, conn *websocket.Conn) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
if h.clients[userID] == nil {
|
||||
h.clients[userID] = make(map[*websocket.Conn]bool)
|
||||
}
|
||||
h.clients[userID][conn] = true
|
||||
}
|
||||
|
||||
func (h *Hub) Unregister(userID string, conn *websocket.Conn) {
|
||||
h.mu.Lock()
|
||||
defer h.mu.Unlock()
|
||||
|
||||
if conns, ok := h.clients[userID]; ok {
|
||||
delete(conns, conn)
|
||||
if len(conns) == 0 {
|
||||
delete(h.clients, userID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// SendToUser 向指定用户的所有连接推送消息
|
||||
func (h *Hub) SendToUser(userID string, message []byte) error {
|
||||
h.mu.RLock()
|
||||
defer h.mu.RUnlock()
|
||||
|
||||
conns, ok := h.clients[userID]
|
||||
if !ok {
|
||||
return nil
|
||||
}
|
||||
|
||||
for conn := range conns {
|
||||
if err := conn.WriteMessage(websocket.TextMessage, message); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
99
services/push-gateway/internal/ws/handler.go
Normal file
99
services/push-gateway/internal/ws/handler.go
Normal file
@@ -0,0 +1,99 @@
|
||||
package ws
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"strings"
|
||||
|
||||
"github.com/edu-cloud/push-gateway/internal/hub"
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/golang-jwt/jwt/v5"
|
||||
"github.com/gorilla/websocket"
|
||||
)
|
||||
|
||||
var upgrader = websocket.Upgrader{
|
||||
CheckOrigin: func(r *http.Request) bool {
|
||||
return true // P5 骨架,生产环境需校验 origin
|
||||
},
|
||||
}
|
||||
|
||||
type Handler struct {
|
||||
hub *hub.Hub
|
||||
jwtSecret string
|
||||
}
|
||||
|
||||
func NewHandler(h *hub.Hub, jwtSecret string) *Handler {
|
||||
return &Handler{hub: h, jwtSecret: jwtSecret}
|
||||
}
|
||||
|
||||
func (h *Handler) HandleWebSocket(c *gin.Context) {
|
||||
// 从 query 参数获取 token(WebSocket 无法设置 Authorization 头)
|
||||
tokenStr := c.Query("token")
|
||||
if tokenStr == "" {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "missing token"})
|
||||
return
|
||||
}
|
||||
|
||||
token, err := jwt.Parse(tokenStr, func(t *jwt.Token) (interface{}, error) {
|
||||
if _, ok := t.Method.(*jwt.SigningMethodHMAC); !ok {
|
||||
return nil, jwt.ErrSignatureInvalid
|
||||
}
|
||||
return []byte(h.jwtSecret), nil
|
||||
})
|
||||
|
||||
if err != nil || !token.Valid {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "invalid token"})
|
||||
return
|
||||
}
|
||||
|
||||
claims, ok := token.Claims.(jwt.MapClaims)
|
||||
if !ok {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "invalid claims"})
|
||||
return
|
||||
}
|
||||
|
||||
userID, ok := claims["sub"].(string)
|
||||
if !ok {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "missing user id"})
|
||||
return
|
||||
}
|
||||
|
||||
conn, err := upgrader.Upgrade(c.Writer, c.Request, nil)
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
defer conn.Close()
|
||||
|
||||
h.hub.Register(userID, conn)
|
||||
defer h.hub.Unregister(userID, conn)
|
||||
|
||||
// 读取循环(保持连接,处理心跳)
|
||||
for {
|
||||
_, msg, err := conn.ReadMessage()
|
||||
if err != nil {
|
||||
break
|
||||
}
|
||||
// 处理心跳 ping
|
||||
if strings.ToLower(string(msg)) == "ping" {
|
||||
conn.WriteMessage(websocket.TextMessage, []byte("pong"))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// PushHandler 接收来自 Msg 服务的推送请求
|
||||
func (h *Handler) PushHandler(c *gin.Context) {
|
||||
var req struct {
|
||||
UserID string `json:"user_id"`
|
||||
Message string `json:"message"`
|
||||
}
|
||||
if err := c.ShouldBindJSON(&req); err != nil {
|
||||
c.JSON(http.StatusBadRequest, gin.H{"success": false, "error": gin.H{"code": "INVALID_REQUEST", "message": err.Error()}})
|
||||
return
|
||||
}
|
||||
|
||||
if err := h.hub.SendToUser(req.UserID, []byte(req.Message)); err != nil {
|
||||
c.JSON(http.StatusInternalServerError, gin.H{"success": false, "error": gin.H{"code": "PUSH_FAILED", "message": err.Error()}})
|
||||
return
|
||||
}
|
||||
|
||||
c.JSON(http.StatusOK, gin.H{"success": true})
|
||||
}
|
||||
63
services/push-gateway/main.go
Normal file
63
services/push-gateway/main.go
Normal file
@@ -0,0 +1,63 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/edu-cloud/push-gateway/internal/config"
|
||||
"github.com/edu-cloud/push-gateway/internal/hub"
|
||||
"github.com/edu-cloud/push-gateway/internal/ws"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
func main() {
|
||||
cfg := config.Load()
|
||||
gin.SetMode(gin.ReleaseMode)
|
||||
|
||||
h := hub.NewHub()
|
||||
wsHandler := ws.NewHandler(h, cfg.JWTSecret)
|
||||
|
||||
r := gin.New()
|
||||
r.Use(gin.Recovery())
|
||||
|
||||
r.GET("/healthz", func(c *gin.Context) {
|
||||
c.JSON(200, gin.H{"status": "ok", "service": "push-gateway"})
|
||||
})
|
||||
|
||||
// WebSocket 升级端点
|
||||
r.GET("/ws", wsHandler.HandleWebSocket)
|
||||
|
||||
// 内部推送 API(Msg 服务调用)
|
||||
api := r.Group("/internal")
|
||||
api.POST("/push", wsHandler.PushHandler)
|
||||
|
||||
srv := &http.Server{
|
||||
Addr: ":" + cfg.Port,
|
||||
Handler: r,
|
||||
ReadTimeout: 10 * time.Second,
|
||||
WriteTimeout: 60 * time.Second, // WebSocket 长连接
|
||||
}
|
||||
|
||||
go func() {
|
||||
log.Printf("Push Gateway listening on :%s", cfg.Port)
|
||||
if err := srv.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
log.Fatalf("listen: %s\n", err)
|
||||
}
|
||||
}()
|
||||
|
||||
quit := make(chan os.Signal, 1)
|
||||
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
|
||||
<-quit
|
||||
log.Println("Shutting down Push Gateway...")
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
if err := srv.Shutdown(ctx); err != nil {
|
||||
log.Fatal("Server forced to shutdown:", err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user