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