sse.go 1.9 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485
  1. package controllers
  2. import (
  3. "Tech/common"
  4. "encoding/json"
  5. "fmt"
  6. "time"
  7. beego "github.com/beego/beego/v2/server/web"
  8. )
  9. // SSE处理器
  10. type SSEController struct {
  11. beego.Controller
  12. }
  13. // 注册SSE连接
  14. func (c *SSEController) Register() {
  15. // 设置SSE响应头
  16. c.Ctx.Output.Header("Content-Type", "text/event-stream")
  17. c.Ctx.Output.Header("Cache-Control", "no-cache")
  18. c.Ctx.Output.Header("Connection", "keep-alive")
  19. c.Ctx.Output.Header("Access-Control-Allow-Origin", "*")
  20. clientID := fmt.Sprintf("client_%d", time.Now().UnixNano())
  21. // 创建消息通道
  22. messageChan := make(chan common.SSEMessage, 10)
  23. // 注册连接
  24. common.GlobalSSEManager.AddConnection(clientID, messageChan)
  25. // 确保在函数退出时清理连接
  26. defer common.GlobalSSEManager.RemoveConnection(clientID)
  27. // 发送连接成功消息
  28. c.sendSSEMessage(common.SSEMessage{
  29. Event: "connected",
  30. Data: map[string]string{
  31. "clientID": clientID,
  32. "message": "SSE连接已建立",
  33. },
  34. Time: time.Now().Unix(),
  35. })
  36. // 定期发送心跳
  37. heartbeat := time.NewTicker(30 * time.Second)
  38. defer heartbeat.Stop()
  39. // 创建一个context用于监听请求取消
  40. ctx := c.Ctx.Request.Context()
  41. // 监听消息并发送给客户端
  42. for {
  43. select {
  44. case message := <-messageChan:
  45. c.sendSSEMessage(message)
  46. case <-heartbeat.C:
  47. c.sendSSEMessage(common.SSEMessage{
  48. Event: "ping",
  49. Data: map[string]string{"status": "alive"},
  50. Time: time.Now().Unix(),
  51. })
  52. case <-ctx.Done():
  53. // 客户端断开连接
  54. return
  55. }
  56. }
  57. }
  58. // 发送SSE消息
  59. func (c *SSEController) sendSSEMessage(message common.SSEMessage) {
  60. data, err := json.Marshal(message.Data)
  61. if err != nil {
  62. return
  63. }
  64. // 格式化SSE消息
  65. sseData := fmt.Sprintf("event: %s\ndata: %s\n\n", message.Event, string(data))
  66. // 在Beego v2中,使用ResponseController来写入响应
  67. c.Ctx.ResponseWriter.Write([]byte(sseData))
  68. c.Ctx.ResponseWriter.Flush()
  69. }