Explorar el Código

提交更新sse

yaojianjun hace 9 meses
padre
commit
27975bc6a9

BIN
Tech/__debug_bin.exe


+ 6 - 0
Tech/common/sse_init.go

@@ -0,0 +1,6 @@
+package common
+
+// 初始化SSE管理器
+func init() {
+	GlobalSSEManager.StartCleanupTask()
+}

+ 130 - 0
Tech/common/sse_manager.go

@@ -0,0 +1,130 @@
+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()
+	}
+}

+ 89 - 0
Tech/controllers/sse.go

@@ -0,0 +1,89 @@
+package controllers
+
+import (
+	"Tech/common"
+	"encoding/json"
+	"fmt"
+	"time"
+
+	beego "github.com/beego/beego/v2/server/web"
+)
+
+// SSE处理器
+type SSEController struct {
+	beego.Controller
+}
+
+// 注册SSE连接
+func (c *SSEController) Register() {
+	// 设置SSE响应头
+	c.Ctx.Output.Header("Content-Type", "text/event-stream")
+	c.Ctx.Output.Header("Cache-Control", "no-cache")
+	c.Ctx.Output.Header("Connection", "keep-alive")
+	c.Ctx.Output.Header("Access-Control-Allow-Origin", "*")
+
+	// 生成客户端ID
+	clientID := c.Ctx.Input.Param(":clientid")
+	if clientID == "" {
+		// 如果没有提供客户端ID,生成一个
+		clientID = fmt.Sprintf("client_%d", time.Now().UnixNano())
+	}
+
+	// 创建消息通道
+	messageChan := make(chan common.SSEMessage, 10)
+
+	// 注册连接
+	common.GlobalSSEManager.AddConnection(clientID, messageChan)
+
+	// 确保在函数退出时清理连接
+	defer common.GlobalSSEManager.RemoveConnection(clientID)
+
+	// 发送连接成功消息
+	c.sendSSEMessage(common.SSEMessage{
+		Event: "connected",
+		Data: map[string]string{
+			"clientID": clientID,
+			"message":  "SSE连接已建立",
+		},
+		Time: time.Now().Unix(),
+	})
+
+	// 定期发送心跳
+	heartbeat := time.NewTicker(30 * time.Second)
+	defer heartbeat.Stop()
+
+	// 创建一个context用于监听请求取消
+	ctx := c.Ctx.Request.Context()
+
+	// 监听消息并发送给客户端
+	for {
+		select {
+		case message := <-messageChan:
+			c.sendSSEMessage(message)
+		case <-heartbeat.C:
+			c.sendSSEMessage(common.SSEMessage{
+				Event: "ping",
+				Data:  map[string]string{"status": "alive"},
+				Time:  time.Now().Unix(),
+			})
+		case <-ctx.Done():
+			// 客户端断开连接
+			return
+		}
+	}
+}
+
+// 发送SSE消息
+func (c *SSEController) sendSSEMessage(message common.SSEMessage) {
+	data, err := json.Marshal(message.Data)
+	if err != nil {
+		return
+	}
+	
+	// 格式化SSE消息
+	sseData := fmt.Sprintf("event: %s\ndata: %s\n\n", message.Event, string(data))
+	
+	// 在Beego v2中,使用ResponseController来写入响应
+	c.Ctx.ResponseWriter.Write([]byte(sseData))
+	c.Ctx.ResponseWriter.Flush()
+}

+ 16 - 0
Tech/controllers/电调.go

@@ -2,6 +2,7 @@
 package controllers
 
 import (
+	"Tech/common"
 	"Tech/models"
 	"bytes"
 	"encoding/json"
@@ -370,6 +371,21 @@ func (c *EDispatchController) DoAutoEva(req Addctrllog_req) EvaResult {
 			}
 		}
 
+		// 创建通知数据
+		notificationData := map[string]interface{}{
+			"type":        "evaluation_update",
+			"runLogID":    foundEvaData.RunLogID,
+			"seqNo":       foundEvaData.SeqNo,
+			"evaName":     foundEvaData.EvaName,
+			"isOperated":  1,
+			"operScore":   operScore,
+			"courseScore": totalScore,
+			"timestamp":   time.Now().Format("2006-01-02 15:04:05"),
+		}
+		
+		// 通过SSE广播评价更新消息
+		common.GlobalSSEManager.BroadcastMessage("evaluation_update", notificationData)
+		
 		// 返回成功结果
 		return EvaResult{
 			Code:        200,

+ 3 - 0
Tech/routers/router.go

@@ -114,4 +114,7 @@ func init() {
 
 	// /evalog/setscore
 	beego.Router("/evalog/setscore", &controllers.EvaLogController{}, "post:SetScore")
+	
+	// SSE路由
+	beego.Router("/sse/register/:clientid", &controllers.SSEController{}, "get:Register")
 }

+ 219 - 0
Tech/sse_client_example.html

@@ -0,0 +1,219 @@
+<!DOCTYPE html>
+<html lang="zh-CN">
+<head>
+    <meta charset="UTF-8">
+    <meta name="viewport" content="width=device-width, initial-scale=1.0">
+    <title>自动评价SSE通知示例</title>
+    <style>
+        body {
+            font-family: 'Arial', sans-serif;
+            max-width: 800px;
+            margin: 0 auto;
+            padding: 20px;
+            background-color: #f5f5f5;
+        }
+        .container {
+            background-color: white;
+            border-radius: 8px;
+            box-shadow: 0 2px 10px rgba(0, 0, 0, 0.1);
+            padding: 20px;
+        }
+        h1 {
+            color: #333;
+            text-align: center;
+        }
+        .status {
+            margin: 20px 0;
+            padding: 10px;
+            border-radius: 4px;
+            font-weight: bold;
+        }
+        .connected {
+            background-color: #e8f5e8;
+            color: #2e7d2e;
+            border: 1px solid #a5d6a7;
+        }
+        .disconnected {
+            background-color: #ffebee;
+            color: #c62828;
+            border: 1px solid #ef9a9a;
+        }
+        .updates {
+            margin-top: 20px;
+        }
+        .update-item {
+            margin-bottom: 15px;
+            padding: 15px;
+            border-radius: 4px;
+            border-left: 4px solid #2196f3;
+            background-color: #e3f2fd;
+        }
+        .update-item .timestamp {
+            font-size: 0.9em;
+            color: #555;
+            margin-bottom: 8px;
+        }
+        .update-item .content {
+            font-weight: bold;
+            color: #333;
+        }
+        .score-update {
+            border-left-color: #4caf50;
+            background-color: #e8f5e8;
+        }
+        .notification {
+            border-left-color: #ff9800;
+            background-color: #fff8e1;
+        }
+        button {
+            background-color: #2196f3;
+            color: white;
+            border: none;
+            padding: 10px 15px;
+            border-radius: 4px;
+            cursor: pointer;
+            margin-right: 10px;
+        }
+        button:hover {
+            background-color: #0d8aee;
+        }
+        button:disabled {
+            background-color: #cccccc;
+            cursor: not-allowed;
+        }
+        .controls {
+            margin-bottom: 20px;
+            text-align: center;
+        }
+    </style>
+</head>
+<body>
+    <div class="container">
+        <h1>自动评价SSE通知</h1>
+        
+        <div class="controls">
+            <button id="connectBtn">连接SSE</button>
+            <button id="disconnectBtn" disabled>断开连接</button>
+        </div>
+        
+        <div id="status" class="status disconnected">未连接</div>
+        
+        <div class="updates">
+            <h2>评价更新记录</h2>
+            <div id="updatesList">
+                <p style="color: #666;">暂无更新记录</p>
+            </div>
+        </div>
+    </div>
+
+    <script>
+        let eventSource = null;
+        let clientID = 'web_client_' + Math.floor(Math.random() * 10000);
+        
+        const connectBtn = document.getElementById('connectBtn');
+        const disconnectBtn = document.getElementById('disconnectBtn');
+        const statusDiv = document.getElementById('status');
+        const updatesList = document.getElementById('updatesList');
+        
+        // 连接SSE
+        function connectSSE() {
+            eventSource = new EventSource(`/sse/register/${clientID}`);
+            
+            eventSource.onopen = function(event) {
+                updateStatus('已连接', 'connected');
+                connectBtn.disabled = true;
+                disconnectBtn.disabled = false;
+                console.log('SSE连接已建立');
+            };
+            
+            eventSource.onerror = function(event) {
+                updateStatus('连接错误', 'disconnected');
+                connectBtn.disabled = false;
+                disconnectBtn.disabled = true;
+                console.error('SSE连接错误:', event);
+            };
+            
+            eventSource.addEventListener('connected', function(event) {
+                const data = JSON.parse(event.data);
+                addUpdateItem('系统通知', `SSE连接已建立 (ID: ${data.clientID})`, 'notification');
+            });
+            
+            eventSource.addEventListener('ping', function(event) {
+                const data = JSON.parse(event.data);
+                console.log('心跳:', data);
+            });
+            
+            eventSource.addEventListener('evaluation_update', function(event) {
+                const data = JSON.parse(event.data);
+                handleEvaluationUpdate(data);
+            });
+        }
+        
+        // 断开SSE连接
+        function disconnectSSE() {
+            if (eventSource) {
+                eventSource.close();
+                eventSource = null;
+                updateStatus('已断开', 'disconnected');
+                connectBtn.disabled = false;
+                disconnectBtn.disabled = true;
+                console.log('SSE连接已断开');
+            }
+        }
+        
+        // 更新连接状态
+        function updateStatus(text, className) {
+            statusDiv.textContent = text;
+            statusDiv.className = 'status ' + className;
+        }
+        
+        // 处理评价更新
+        function handleEvaluationUpdate(data) {
+            if (data.type === 'evaluation_update') {
+                const message = `
+                    <strong>评价项目:</strong> ${data.evaName}<br>
+                    <strong>序号:</strong> ${data.seqNo}<br>
+                    <strong>操作得分:</strong> ${data.operScore}<br>
+                    <strong>课程总分:</strong> ${data.courseScore}<br>
+                    <strong>更新时间:</strong> ${data.timestamp}
+                `;
+                
+                addUpdateItem('评价得分更新', message, 'score-update');
+            }
+        }
+        
+        // 添加更新记录
+        function addUpdateItem(title, content, className) {
+            const item = document.createElement('div');
+            item.className = `update-item ${className}`;
+            
+            const timestamp = new Date().toLocaleString();
+            
+            item.innerHTML = `
+                <div class="timestamp">${timestamp}</div>
+                <div class="content">
+                    <strong>${title}</strong><br>
+                    ${content}
+                </div>
+            `;
+            
+            // 如果是第一条记录,清空"暂无更新记录"的提示
+            if (updatesList.children.length === 1 && updatesList.children[0].textContent.includes('暂无更新记录')) {
+                updatesList.innerHTML = '';
+            }
+            
+            // 将新记录插入到顶部
+            updatesList.insertBefore(item, updatesList.firstChild);
+            
+            // 限制显示的记录数量
+            while (updatesList.children.length > 10) {
+                updatesList.removeChild(updatesList.lastChild);
+            }
+        }
+        
+        // 绑定按钮事件
+        connectBtn.addEventListener('click', connectSSE);
+        disconnectBtn.addEventListener('click', disconnectSSE);
+    </script>
+</body>
+</html>