package common import ( "sync" "time" ) // SSE消息结构体 type SSEMessage struct { Event string `json:"event"` // 事件类型 Data interface{} `json:"data"` // 数据内容 Time int64 `json:"time"` // 时间戳 } // SSE连接信息 type SSEConnection struct { ID string Channel chan SSEMessage LastPingTime time.Time } // SSE管理器 type SSEManager struct { connections map[string]*SSEConnection mutex sync.RWMutex } // 全局SSE管理器实例 var GlobalSSEManager = &SSEManager{ connections: make(map[string]*SSEConnection), } // 添加SSE连接 func (m *SSEManager) AddConnection(id string, channel chan SSEMessage) { m.mutex.Lock() defer m.mutex.Unlock() m.connections[id] = &SSEConnection{ ID: id, Channel: channel, LastPingTime: time.Now(), } } // 移除SSE连接 func (m *SSEManager) RemoveConnection(id string) { m.mutex.Lock() defer m.mutex.Unlock() if conn, exists := m.connections[id]; exists { close(conn.Channel) delete(m.connections, id) } } // 发送消息到指定连接 func (m *SSEManager) SendMessageToConnection(id string, event string, data interface{}) bool { m.mutex.RLock() defer m.mutex.RUnlock() if conn, exists := m.connections[id]; exists { select { case conn.Channel <- SSEMessage{ Event: event, Data: data, Time: time.Now().Unix(), }: return true default: // 通道已满,可能是客户端断开连接 return false } } return false } // 广播消息到所有连接 func (m *SSEManager) BroadcastMessage(event string, data interface{}) { m.mutex.RLock() defer m.mutex.RUnlock() for id, conn := range m.connections { select { case conn.Channel <- SSEMessage{ Event: event, Data: data, Time: time.Now().Unix(), }: default: // 通道已满,可能是客户端断开连接 close(conn.Channel) delete(m.connections, id) } } } // 清理超时连接 func (m *SSEManager) CleanupStaleConnections() { m.mutex.Lock() defer m.mutex.Unlock() now := time.Now() for id, conn := range m.connections { // 如果连接超过60秒没有活动,则关闭 if now.Sub(conn.LastPingTime) > 60*time.Second { close(conn.Channel) delete(m.connections, id) } } } // 启动定期清理任务 func (m *SSEManager) StartCleanupTask() { ticker := time.NewTicker(30 * time.Second) go func() { for range ticker.C { m.CleanupStaleConnections() } }() } // 更新连接的最后活动时间 func (m *SSEManager) UpdateLastPingTime(id string) { m.mutex.Lock() defer m.mutex.Unlock() if conn, exists := m.connections[id]; exists { conn.LastPingTime = time.Now() } }