package common import ( "Tech/models" "bytes" "crypto/md5" "encoding/hex" "encoding/json" "fmt" "io" "net/http" "time" "github.com/beego/beego/v2/client/orm" beego "github.com/beego/beego/v2/server/web" "github.com/sirupsen/logrus" ) // EquipmentData 表示设备数据 type EquipmentData struct { DataID string `json:"dataId"` DataName string `json:"dataName"` Value string `json:"value"` ID string `json:"id"` EqpmtTypeID int `json:"eqpmtTypeID"` } // ElectricalDataFetcher 电调数据获取器 type ElectricalDataFetcher struct { manager *models.ElectricModelManager electricalIP string electricalPort int lastDataHash string lastEquipments []map[string]interface{} // 记录上次广播的设备数据 isRunning bool } var fetcher *ElectricalDataFetcher // InitElectricalDataFetcher 初始化电调数据获取器 func InitElectricalDataFetcher() { fetcher = &ElectricalDataFetcher{ manager: models.GetInstance(), isRunning: false, lastEquipments: nil, // 初始化为空,第一次会推送全量数据 } // 从配置文件读取电调服务配置 electricalIP, err := beego.AppConfig.String("ElectricalIP") if err != nil { LogError("读取电调服务IP配置失败", err) return } electricalPort, err := beego.AppConfig.Int("ElectricalPort") if err != nil { LogError("读取电调服务端口配置失败", err) return } fetcher.electricalIP = electricalIP fetcher.electricalPort = electricalPort // 启动WebSocket管理器 GlobalWSManager.Start() // 启动定时获取数据任务 fetcher.Start() } // Start 启动定时获取数据任务 func (f *ElectricalDataFetcher) Start() { if f.isRunning { return } f.isRunning = true // 获取全量数据更新间隔配置,默认5秒 intervalSeconds := 5 if interval, err := beego.AppConfig.Int("EGetAllDataInterval"); err == nil { intervalSeconds = interval } ticker := time.NewTicker(time.Duration(intervalSeconds) * time.Second) go func() { // 初始延迟1秒执行,避免系统启动时的资源竞争 time.Sleep(1 * time.Second) for range ticker.C { f.fetchAndBroadcastData() } }() LogInfo(fmt.Sprintf("电调数据定时获取任务已启动,间隔 %d 秒", intervalSeconds)) } // fetchAndBroadcastData 获取并广播数据 func (f *ElectricalDataFetcher) fetchAndBroadcastData() { defer func() { if r := recover(); r != nil { LogError("获取电调数据时发生异常", nil, logrus.Fields{ "panic": r, }) } }() // 获取最新的课程运行记录以获取线路号 o := orm.NewOrm() var latestCourse models.CourseRunLog err := o.QueryTable(new(models.CourseRunLog)). OrderBy("-RunLogID"). One(&latestCourse) lineID := 0 if err == nil && err != orm.ErrNoRows { // 根据运行课程ID查询课程信息以获取线路号 var courseInfo models.CourseInfo err = o.QueryTable(new(models.CourseInfo)). Filter("CourseID", latestCourse.RunCourseID). One(&courseInfo) if err == nil { lineID = courseInfo.LineNo } } // 构造电调服务地址 electricalAddr := fmt.Sprintf("http://%s:%d/ServerCtrl/AllEqpmtList", f.electricalIP, f.electricalPort) // 准备请求数据 type RequestData struct { LineID int `json:"lineID"` StationID int `json:"StationID"` ViewID int `json:"ViewID"` } reqData := RequestData{ LineID: lineID, StationID: 0, ViewID: 0, } // 将请求数据转换为JSON reqBody, err := json.Marshal(reqData) if err != nil { LogError("构造请求数据失败", err) return } // 调用电调服务获取所有设备列表接口 resp, err := http.Post(electricalAddr, "application/json", bytes.NewBuffer(reqBody)) if err != nil { LogError("调用电调服务获取所有设备列表接口失败", err) return } defer resp.Body.Close() // 检查响应状态 if resp.StatusCode != http.StatusOK { bodyBytes, _ := io.ReadAll(resp.Body) LogError("电调服务返回异常状态码", fmt.Errorf("状态码: %d, 响应: %s", resp.StatusCode, string(bodyBytes))) return } // 读取并解析响应体 bodyBytes, err := io.ReadAll(resp.Body) if err != nil { LogError("读取电调服务响应失败", err) return } // 解析JSON响应 var result map[string]interface{} if err := json.Unmarshal(bodyBytes, &result); err != nil { LogError("解析电调服务JSON响应失败", err) return } // 检查是否有data字段 data, exists := result["data"] if !exists { LogError("电调服务响应中没有data字段", nil) return } // 解析data中的设备列表 var equipments []models.DiandiaoModel if equipBytes, err := json.Marshal(data); err == nil { if err := json.Unmarshal(equipBytes, &equipments); err != nil { LogError("解析设备列表数据失败", err) return } } // 更新模型管理器的总包设备数据 f.manager.UpdateBatchData(equipments) // 转换为标准设备数据格式 equipmentsData := f.manager.ConvertToEquipmentData(equipments) // 加入测试代码:生成随机值修改第一个设备数据 // if len(equipmentsData) > 0 { // // 生成0-100之间的随机值 // randomValue := time.Now().Unix() % 101 // equipmentsData[0]["value"] = fmt.Sprintf("%d", randomValue) // LogDebug(fmt.Sprintf("测试:修改第一个设备数据的值为 %d", randomValue), nil) // } // 比较数据差异,获取变化的设备 changedEquipments := f.compareWithLastData(equipmentsData) // 如果没有变化,则跳过广播 if len(changedEquipments) == 0 { return } // 更新上次广播的数据 f.lastEquipments = equipmentsData // 广播变化的设备数据到所有WebSocket连接 GlobalWSManager.Broadcast(changedEquipments) LogDebug(fmt.Sprintf("成功获取并广播变化设备数据,变化数量: %d", len(changedEquipments)), logrus.Fields{ "lineID": lineID, "totalCount": len(equipmentsData), "changedCount": len(changedEquipments), }) } // compareWithLastData 比较当前数据与上次数据的差异 func (f *ElectricalDataFetcher) compareWithLastData(currentData []map[string]interface{}) []map[string]interface{} { // 如果是第一次(没有上次数据),返回所有数据 if f.lastEquipments == nil { return currentData } // 构建上次数据的映射 lastDataMap := make(map[string]map[string]interface{}) for _, equip := range f.lastEquipments { if id, ok := equip["id"].(string); ok { lastDataMap[id] = equip } } // 构建当前数据的映射 currentDataMap := make(map[string]map[string]interface{}) for _, equip := range currentData { if id, ok := equip["id"].(string); ok { currentDataMap[id] = equip } } // 查找变化的设备 var changedEquipments []map[string]interface{} allIDs := make(map[string]bool) // 检查上次数据中的设备 for id, lastEquip := range lastDataMap { allIDs[id] = true currentEquip, exists := currentDataMap[id] if !exists { // 设备被删除了,不推送删除的设备 continue } // 比较值是否变化 if lastEquip["value"] != currentEquip["value"] { changedEquipments = append(changedEquipments, currentEquip) } } // 检查新增加的设备 for id, currentEquip := range currentDataMap { if !allIDs[id] { changedEquipments = append(changedEquipments, currentEquip) } } return changedEquipments } // calculateDataHash 计算数据哈希 func (f *ElectricalDataFetcher) calculateDataHash(data []map[string]interface{}) string { // 将数据序列化为JSON jsonData, err := json.Marshal(data) if err != nil { return "" } // 计算MD5哈希 hash := md5.Sum(jsonData) return hex.EncodeToString(hash[:]) } // GetCurrentData 获取当前缓存的设备数据 func (f *ElectricalDataFetcher) GetCurrentData() ([]map[string]interface{}, error) { // 获取所有设备数据 equipments := f.manager.GetAllEquipmentData() // 转换为标准设备数据格式 equipmentsData := f.manager.ConvertToEquipmentData(equipments) return equipmentsData, nil } // SendCurrentDataToClient 发送当前数据到指定客户端 func SendCurrentDataToClient(clientID string) { if fetcher == nil { LogError("电调数据获取器未初始化", nil) return } // 获取当前数据 equipments, err := fetcher.GetCurrentData() if err != nil { LogError("获取当前设备数据失败", err) return } // 发送到指定客户端(直接发送设备数组) if err := GlobalWSManager.BroadcastToClient(clientID, equipments); err != nil { LogError("发送设备数据到客户端失败", err, logrus.Fields{ "clientID": clientID, }) return } LogInfo(fmt.Sprintf("成功发送全量设备数据到客户端 %s,设备数量: %d", clientID, len(equipments))) }