sse.go 2.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990
  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. // 生成客户端ID
  21. clientID := c.Ctx.Input.Param(":clientid")
  22. if clientID == "" {
  23. // 如果没有提供客户端ID,生成一个
  24. clientID = fmt.Sprintf("client_%d", time.Now().UnixNano())
  25. }
  26. // 创建消息通道
  27. messageChan := make(chan common.SSEMessage, 10)
  28. // 注册连接
  29. common.GlobalSSEManager.AddConnection(clientID, messageChan)
  30. // 确保在函数退出时清理连接
  31. defer common.GlobalSSEManager.RemoveConnection(clientID)
  32. // 发送连接成功消息
  33. c.sendSSEMessage(common.SSEMessage{
  34. Event: "connected",
  35. Data: map[string]string{
  36. "clientID": clientID,
  37. "message": "SSE连接已建立",
  38. },
  39. Time: time.Now().Unix(),
  40. })
  41. // 定期发送心跳
  42. heartbeat := time.NewTicker(30 * time.Second)
  43. defer heartbeat.Stop()
  44. // 创建一个context用于监听请求取消
  45. ctx := c.Ctx.Request.Context()
  46. // 监听消息并发送给客户端
  47. for {
  48. select {
  49. case message := <-messageChan:
  50. c.sendSSEMessage(message)
  51. case <-heartbeat.C:
  52. c.sendSSEMessage(common.SSEMessage{
  53. Event: "ping",
  54. Data: map[string]string{"status": "alive"},
  55. Time: time.Now().Unix(),
  56. })
  57. case <-ctx.Done():
  58. // 客户端断开连接
  59. return
  60. }
  61. }
  62. }
  63. // 发送SSE消息
  64. func (c *SSEController) sendSSEMessage(message common.SSEMessage) {
  65. data, err := json.Marshal(message.Data)
  66. if err != nil {
  67. return
  68. }
  69. // 格式化SSE消息
  70. sseData := fmt.Sprintf("event: %s\ndata: %s\n\n", message.Event, string(data))
  71. // 在Beego v2中,使用ResponseController来写入响应
  72. c.Ctx.ResponseWriter.Write([]byte(sseData))
  73. c.Ctx.ResponseWriter.Flush()
  74. }