sse_manager.go 2.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131
  1. package common
  2. import (
  3. "sync"
  4. "time"
  5. )
  6. // SSE消息结构体
  7. type SSEMessage struct {
  8. Event string `json:"event"` // 事件类型
  9. Data interface{} `json:"data"` // 数据内容
  10. Time int64 `json:"time"` // 时间戳
  11. }
  12. // SSE连接信息
  13. type SSEConnection struct {
  14. ID string
  15. Channel chan SSEMessage
  16. LastPingTime time.Time
  17. }
  18. // SSE管理器
  19. type SSEManager struct {
  20. connections map[string]*SSEConnection
  21. mutex sync.RWMutex
  22. }
  23. // 全局SSE管理器实例
  24. var GlobalSSEManager = &SSEManager{
  25. connections: make(map[string]*SSEConnection),
  26. }
  27. // 添加SSE连接
  28. func (m *SSEManager) AddConnection(id string, channel chan SSEMessage) {
  29. m.mutex.Lock()
  30. defer m.mutex.Unlock()
  31. m.connections[id] = &SSEConnection{
  32. ID: id,
  33. Channel: channel,
  34. LastPingTime: time.Now(),
  35. }
  36. }
  37. // 移除SSE连接
  38. func (m *SSEManager) RemoveConnection(id string) {
  39. m.mutex.Lock()
  40. defer m.mutex.Unlock()
  41. if conn, exists := m.connections[id]; exists {
  42. close(conn.Channel)
  43. delete(m.connections, id)
  44. }
  45. }
  46. // 发送消息到指定连接
  47. func (m *SSEManager) SendMessageToConnection(id string, event string, data interface{}) bool {
  48. m.mutex.RLock()
  49. defer m.mutex.RUnlock()
  50. if conn, exists := m.connections[id]; exists {
  51. select {
  52. case conn.Channel <- SSEMessage{
  53. Event: event,
  54. Data: data,
  55. Time: time.Now().Unix(),
  56. }:
  57. return true
  58. default:
  59. // 通道已满,可能是客户端断开连接
  60. return false
  61. }
  62. }
  63. return false
  64. }
  65. // 广播消息到所有连接
  66. func (m *SSEManager) BroadcastMessage(event string, data interface{}) {
  67. m.mutex.RLock()
  68. defer m.mutex.RUnlock()
  69. for id, conn := range m.connections {
  70. select {
  71. case conn.Channel <- SSEMessage{
  72. Event: event,
  73. Data: data,
  74. Time: time.Now().Unix(),
  75. }:
  76. default:
  77. // 通道已满,可能是客户端断开连接
  78. close(conn.Channel)
  79. delete(m.connections, id)
  80. }
  81. }
  82. }
  83. // 清理超时连接
  84. func (m *SSEManager) CleanupStaleConnections() {
  85. m.mutex.Lock()
  86. defer m.mutex.Unlock()
  87. now := time.Now()
  88. for id, conn := range m.connections {
  89. // 如果连接超过60秒没有活动,则关闭
  90. if now.Sub(conn.LastPingTime) > 60*time.Second {
  91. close(conn.Channel)
  92. delete(m.connections, id)
  93. }
  94. }
  95. }
  96. // 启动定期清理任务
  97. func (m *SSEManager) StartCleanupTask() {
  98. ticker := time.NewTicker(30 * time.Second)
  99. go func() {
  100. for range ticker.C {
  101. m.CleanupStaleConnections()
  102. }
  103. }()
  104. }
  105. // 更新连接的最后活动时间
  106. func (m *SSEManager) UpdateLastPingTime(id string) {
  107. m.mutex.Lock()
  108. defer m.mutex.Unlock()
  109. if conn, exists := m.connections[id]; exists {
  110. conn.LastPingTime = time.Now()
  111. }
  112. }