| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990 |
- 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()
- }
|