package controllers import ( "Tech/common" "encoding/json" "fmt" "time" beego "github.com/beego/beego/v2/server/web" ) // SSE处理器 type SSEController struct { beego.Controller } // 注册SSE连接 func (c *SSEController) Register() { // 设置SSE响应头 c.Ctx.Output.Header("Content-Type", "text/event-stream") c.Ctx.Output.Header("Cache-Control", "no-cache") c.Ctx.Output.Header("Connection", "keep-alive") c.Ctx.Output.Header("Access-Control-Allow-Origin", "*") // 生成客户端ID clientID := c.Ctx.Input.Param(":clientid") if clientID == "" { // 如果没有提供客户端ID,生成一个 clientID = fmt.Sprintf("client_%d", time.Now().UnixNano()) } // 创建消息通道 messageChan := make(chan common.SSEMessage, 10) // 注册连接 common.GlobalSSEManager.AddConnection(clientID, messageChan) // 确保在函数退出时清理连接 defer common.GlobalSSEManager.RemoveConnection(clientID) // 发送连接成功消息 c.sendSSEMessage(common.SSEMessage{ Event: "connected", Data: map[string]string{ "clientID": clientID, "message": "SSE连接已建立", }, Time: time.Now().Unix(), }) // 定期发送心跳 heartbeat := time.NewTicker(30 * time.Second) defer heartbeat.Stop() // 创建一个context用于监听请求取消 ctx := c.Ctx.Request.Context() // 监听消息并发送给客户端 for { select { case message := <-messageChan: c.sendSSEMessage(message) case <-heartbeat.C: c.sendSSEMessage(common.SSEMessage{ Event: "ping", Data: map[string]string{"status": "alive"}, Time: time.Now().Unix(), }) case <-ctx.Done(): // 客户端断开连接 return } } } // 发送SSE消息 func (c *SSEController) sendSSEMessage(message common.SSEMessage) { data, err := json.Marshal(message.Data) if err != nil { return } // 格式化SSE消息 sseData := fmt.Sprintf("event: %s\ndata: %s\n\n", message.Event, string(data)) // 在Beego v2中,使用ResponseController来写入响应 c.Ctx.ResponseWriter.Write([]byte(sseData)) c.Ctx.ResponseWriter.Flush() }