package item import ( "errors" "fmt" "go.uber.org/zap" "gorm.io/gorm" "math" "server/dao" "server/global" "server/model" "server/model/item/request" "time" ) type BluetoothService struct { } func HandleTempScanDeviceData(bv model.BeaconValue) { // 查询网关 var gw dao.Gateway global.GVA_DB.Where("device_name = ?", bv.DeviceName).First(&gw) // 组装临时记录 temp := dao.TempScanDevice{ GatewayID: gw.ID, GatewaySN: bv.DeviceName, AreaNo: bv.Area, ClusterNo: bv.Cluster, Number: bv.Number, UUID: bv.UUID, FirmwareVer: bv.Version, DeviceType: "lamp", IsSelected: 0, ScanTime: time.Now(), } // 同网关+UUID自动覆盖旧扫描记录 global.GVA_DB.Where("gateway_id = ? AND uuid = ?", gw.ID, temp.UUID).Assign(temp).FirstOrCreate(&temp) } func HandleHeartBeat(deviceName, uuid string) { now := time.Now() updateGatewayData := map[string]interface{}{ "online": 1, "online_time": now, } // 网关deviceName匹配,更新网关在线时间 err := global.GVA_DB.Model(&dao.Gateway{}).Where("device_name = ?", deviceName). Updates(updateGatewayData).Error if err != nil { // 网关不存在/更新失败仅打日志,不阻断灯具更新逻辑 global.GVA_LOG.Error("心跳:", zap.Any("err", err)) return } // 以UUID为唯一主键匹配灯具 updateData := map[string]interface{}{ "online": 1, "online_time": now, } // 仅更新,不存在则不操作 global.GVA_DB.Model(&dao.Bluetooth{}). Where("uuid = ?", uuid). Updates(updateData) } func (bs *BluetoothService) QueryAllTempScanDevices() ([]dao.TempScanDevice, error) { return dao.QueryAllTempScanDevices() } func (bs *BluetoothService) DeviceSendCmd(info request.DeviceSendCmdData) error { gateway, err := dao.QueryGatewayByDeviceName(info.DeviceName) if err != nil { return err } return DeviceSendCmd(gateway.ProductKey, gateway.DeviceName, gateway.NetPwd, "00 00", "00 00", info.Action, info.Params) } func (bs *BluetoothService) GetDeviceData(info request.GetDeviceOrderData) error { gateway, err := dao.QueryGatewayByDeviceName(info.DeviceName) if err != nil { return err } return GetSetting(gateway.ProductKey, gateway.DeviceName, gateway.NetPwd, "00 00", "00 00") } func GetDeviceConsumption() error { gateways, err := dao.QueryAllGateways() if err != nil { return err } for _, gateway := range gateways { go DeviceSendCmd(gateway.ProductKey, gateway.DeviceName, gateway.NetPwd, "00 00", "00 00", "getConsumption", "") } return nil } func (bs *BluetoothService) ScanAllDevice(info request.GetDeviceOrderData) error { gateway, err := dao.QueryGatewayByDeviceName(info.DeviceName) if err != nil { return err } return ScanAll(gateway.ProductKey, gateway.DeviceName, gateway.NetPwd, "00 00") } func (bs *BluetoothService) CreateBluetooth(info request.CreateBluetoothData) error { // 1、插入正式灯具表 blu := dao.Bluetooth{ GatewayID: uint(info.GatewayId), UUID: info.Device.UUID, AreaNo: info.Device.AreaNo, ClusterNo: info.Device.ClusterNo, Number: info.Device.Number, DeviceType: info.Device.DeviceType, Online: 0, } err := blu.CreateBluetooth() if err != nil { return err } // 2、更新对应临时设备 IsSelected=1(已入库) tempModel := dao.TempScanDevice{} err = global.GVA_DB.Model(&tempModel). Where("gateway_id = ? AND uuid = ?", info.Device.GatewayID, info.Device.UUID). Update("is_selected", 1).Error return err } func HandleSceneData(sv model.SceneValue) { now := time.Now() // 更新字段映射 updateMap := map[string]interface{}{ "scene_no": sv.SceneNo, "high_bright": sv.HighBright, "standby_bright": sv.StandbyBright, "cct_bright": sv.CctBright, "delay_time": sv.DelayTime, "delay_time2": sv.DelayTime2, "light_mode": sv.LightMode, "delay_mode": sv.DelayMode, "als_control_mode": sv.AlsControlMode, "scene_validity": sv.SceneValidity, "online": 1, "online_time": now, } // 根据唯一UUID更新灯具记录 err := global.GVA_DB.Model(&dao.Bluetooth{}). Where("uuid = ?", sv.UUID). Updates(updateMap).Error if err != nil { global.GVA_LOG.Error("处理情景参数上报失败", zap.String("uuid", sv.UUID), zap.Error(err)) return } // 顺带刷新网关在线时间(DeviceName为网关SN) err = global.GVA_DB.Model(&dao.Gateway{}). Where("device_name = ?", sv.DeviceName). Updates(map[string]interface{}{ "online_time": now, "online": 1, }).Error if err != nil { global.GVA_LOG.Error("情景上报刷新网关在线失败", zap.String("gatewaySn", sv.DeviceName), zap.Error(err)) } } func (bs *BluetoothService) QueryBluetoothsByGatewayID(gatewayId int) ([]dao.Bluetooth, error) { return dao.QueryBluetoothsByGatewayID(gatewayId) } func CreateDeviceConsumption(val model.ConsumptionValue) error { // 1. 获取网关信息 gateway, err := dao.QueryGatewayByDeviceName(val.DeviceName) if err != nil { return err } // 2. 查询该灯具(UUID)上一次上报的记录 var lastRecord dao.DeviceConsumption // 根据 UUID 倒序查询最新的一条 err = global.GVA_DB.Where("uuid = ?", val.UUID). Order("acquisition_ts DESC"). First(&lastRecord).Error // 3. 初始化本次存储的“增量”数据(默认等于硬件上报值,以防第一次上报) deltaTimeDur := val.TimeDur deltaEnergyDur := val.EnergyDur // 4. 如果历史记录存在,计算出差值 if err == nil { // 说明查到了上次记录 // 防御性编程:防止硬件重启导致 TimeDur 数值变小变成负数 if val.TimeDur > lastRecord.TimeDur { deltaTimeDur = val.TimeDur - lastRecord.TimeDur } if val.EnergyDur > lastRecord.EnergyDur { deltaEnergyDur = val.EnergyDur - lastRecord.EnergyDur } } // 5. 根据差值,重新计算本次的实际和理论能耗 (度) deltaActualKwh := float64(deltaEnergyDur) * val.Power / 3600 / 1000 deltaTheoryKwh := float64(deltaTimeDur) * val.Power / 3600 / 1000 // 6. 构建要落库的结构体,此时存入的都是 【差值】 dc := dao.DeviceConsumption{ GatewayID: gateway.ID, UUID: val.UUID, AreaNo: val.Area, Number: val.Number, DeviceName: val.DeviceName, AcquisitionTs: val.AcquisitionTime, TimeDur: deltaTimeDur, // 【关键修改】存差值,而不是源值 SensorDur: val.SensorDur, // 如果 sensor_dur 也是总累积,也需要按此逻辑做减法 EnergyDur: deltaEnergyDur, // 【关键修改】存差值 Power: val.Power, ActualKwh: deltaActualKwh, // 【关键修改】存差值对应的实际耗电 TheoryKwh: deltaTheoryKwh, // 【关键修改】存差值对应的理论耗电 } // 7. 调用你原有的 dao 方法,直接落库 return dc.CreateDeviceConsumption() } // SummaryWeekConsumption 按网关维度聚合每周总能耗 func SummaryWeekConsumption() error { db := global.GVA_DB.Begin() defer func() { if err := recover(); err != nil { global.GVA_LOG.Error("能耗汇总panic", zap.Any("err", err)) db.Rollback() } }() now := time.Now() curYear, curWeek := now.ISOWeek() // ========== 第一步:先统计【每一盏灯具】每周能耗 ========== // 先删除本周已存在的灯具维度汇总数据,避免重复累加 err := db.Where("year = ? AND week = ? AND uuid <> ''", curYear, curWeek). Delete(&dao.WeekConsumptionSummary{}).Error if err != nil { db.Rollback() global.GVA_LOG.Error("删除本周灯具历史汇总失败", zap.Error(err)) return err } type LampSummary struct { GatewayID uint UUID string Year int Week int ActualTotal float64 TheoryTotal float64 } var lampSummaryList []LampSummary // 按网关、灯具UUID、年、周分组求和 err = db.Model(&dao.DeviceConsumption{}). Select(`gateway_id, uuid, YEAR(FROM_UNIXTIME(acquisition_ts)) AS year, WEEK(FROM_UNIXTIME(acquisition_ts), 1) AS week, SUM(actual_kwh) AS actual_total, SUM(theory_kwh) AS theory_total`). Group("gateway_id, uuid, year, week"). Scan(&lampSummaryList).Error if err != nil { db.Rollback() global.GVA_LOG.Error("灯具周能耗聚合失败", zap.Error(err)) return err } // 批量插入灯具维度周汇总 var lampInsertList []dao.WeekConsumptionSummary for _, item := range lampSummaryList { lampInsertList = append(lampInsertList, dao.WeekConsumptionSummary{ GatewayID: item.GatewayID, UUID: item.UUID, Year: item.Year, Week: item.Week, ActualTotal: item.ActualTotal, TheoryTotal: item.TheoryTotal, }) } if len(lampInsertList) > 0 { if err = db.CreateInBatches(lampInsertList, 500).Error; err != nil { db.Rollback() return err } } // ========== 第二步:统计【单个网关下所有灯具】每周总能耗(前端图表用) ========== // UUID为空代表网关全局汇总 err = db.Where("year = ? AND week = ? AND uuid = ''", curYear, curWeek). Delete(&dao.WeekConsumptionSummary{}).Error if err != nil { db.Rollback() return err } type GatewaySummary struct { GatewayID uint Year int Week int ActualTotal float64 TheoryTotal float64 } var gatewaySummaryList []GatewaySummary // 只按网关、年、周分组,自动累加该网关下所有灯具能耗 err = db.Model(&dao.DeviceConsumption{}). Select(`gateway_id, YEAR(FROM_UNIXTIME(acquisition_ts)) AS year, WEEK(FROM_UNIXTIME(acquisition_ts), 1) AS week, SUM(actual_kwh) AS actual_total, SUM(theory_kwh) AS theory_total`). Group("gateway_id, year, week"). Scan(&gatewaySummaryList).Error if err != nil { db.Rollback() return err } // 网关总汇总:UUID赋值为空字符串 var gatewayInsertList []dao.WeekConsumptionSummary for _, item := range gatewaySummaryList { gatewayInsertList = append(gatewayInsertList, dao.WeekConsumptionSummary{ GatewayID: item.GatewayID, UUID: "", // 空标识网关全局汇总 Year: item.Year, Week: item.Week, ActualTotal: item.ActualTotal, TheoryTotal: item.TheoryTotal, }) } if len(gatewayInsertList) > 0 { if err = db.CreateInBatches(gatewayInsertList, 500).Error; err != nil { db.Rollback() return err } } db.Commit() global.GVA_LOG.Info("本周网关&灯具能耗汇总统计完成") return nil } // QueryDataDashboardParameter 获取数据大屏综合参数 func (bs *BluetoothService) QueryDataDashboardParameter(req model.DashboardRequest) model.DashboardResponse { db := global.GVA_DB // 1. 获取所有网关列表 (下拉框数据) var gateways []dao.Gateway if err := db.Find(&gateways).Error; err != nil { global.GVA_LOG.Error("获取网关列表失败", zap.Error(err)) } // 2. 获取可供选择的历史周列表 (下拉框数据) // 注意:如果选了网关,只展示该网关下的历史周 var rawWeeks []struct { Year int Week int } weekQuery := db.Model(&dao.WeekConsumptionSummary{}).Select("DISTINCT year, week").Where("uuid = ''") if req.GatewayID != 0 { weekQuery = weekQuery.Where("gateway_id = ?", req.GatewayID) } weekQuery.Order("year DESC, week DESC").Scan(&rawWeeks) var availableWeeks []model.AvailableWeek for _, r := range rawWeeks { availableWeeks = append(availableWeeks, model.AvailableWeek{ Year: r.Year, Week: r.Week, Label: fmt.Sprintf("%d 第 %d 周", r.Year, r.Week), }) } // 3. 默认时间处理:若前端没传年份/周数,则使用表中最新的周 if req.Year == 0 || req.Week == 0 { if len(rawWeeks) > 0 { req.Year, req.Week = rawWeeks[0].Year, rawWeeks[0].Week } else { // 如果没有任何记录,则使用当前系统周(兜底) now := time.Now() req.Year, req.Week = now.ISOWeek() } } // ================= 核心数据聚合 ================= // 4. 获取指定网关、指定周的汇总能耗数据 (实际耗电、理论耗电) var actualEnergy, theoryEnergy float64 var weekSummary dao.WeekConsumptionSummary querySummary := db.Where("uuid = ''") if req.GatewayID != 0 { querySummary = querySummary.Where("gateway_id = ?", req.GatewayID) } err := querySummary.Where("year = ? AND week = ?", req.Year, req.Week). First(&weekSummary).Error if err == nil { actualEnergy = weekSummary.ActualTotal theoryEnergy = weekSummary.TheoryTotal } else if !errors.Is(err, gorm.ErrRecordNotFound) { global.GVA_LOG.Error("查询周汇总能耗失败", zap.Error(err)) } // 5. 计算碳排放和占比 carbonReduction := math.Max(0, theoryEnergy-actualEnergy) carbonReductionRatio := 0.0 if theoryEnergy > 0 { carbonReductionRatio = (carbonReduction / theoryEnergy) * 100 } // ================= 统计灯具与预警数 ================= var totalLuminaires int64 lampsQuery := db.Model(&dao.Bluetooth{}) if req.GatewayID != 0 { lampsQuery = lampsQuery.Where("gateway_id = ?", req.GatewayID) } lampsQuery.Count(&totalLuminaires) var warningCount int64 alarmQuery := db.Model(&dao.Bluetooth{}).Where("online = 0") if req.GatewayID != 0 { alarmQuery = alarmQuery.Where("gateway_id = ?", req.GatewayID) } alarmQuery.Count(&warningCount) // ================= 获取回路设备(智能控制)列表 ================= var devices []dao.Device if err := db.Group("device.id").Find(&devices).Error; err != nil { global.GVA_LOG.Error("查询回路设备失败", zap.Error(err)) } loopDeviceCount := len(devices) // ================= 趋势图数据计算 ================= var trends []model.TrendData var yearOnYearTrends []model.TrendData // 只有当传了确切的周数才处理趋势图,否则无法计算时间范围 if req.Year != 0 && req.Week != 0 { // 6. 获取本周每日能耗趋势 startTs, endTs := getWeekRange(req.Year, req.Week) // 修正后的ISO周计算 var dailyData []dailyTrendDB trendQuery := db.Model(&dao.DeviceConsumption{}) if req.GatewayID != 0 { trendQuery = trendQuery.Where("gateway_id = ?", req.GatewayID) } err = trendQuery. Select("DATE_FORMAT(CONVERT_TZ(FROM_UNIXTIME(acquisition_ts), '+00:00', '+08:00'), '%Y-%m-%d') as date_str, "+ "SUM(actual_kwh) as actual_total, SUM(theory_kwh) as theory_total"). Where("acquisition_ts >= ? AND acquisition_ts <= ?", startTs, endTs). Group("date_str"). Order("date_str ASC"). Scan(&dailyData).Error if err != nil { global.GVA_LOG.Error("获取本周每日趋势失败", zap.Error(err)) } // 序列化为 7天的标准 TrendData (补齐缺失的日期为0) trendMap := make(map[string]model.TrendData) for _, d := range dailyData { trendMap[d.DateStr] = model.TrendData{Date: d.DateStr, Actual: d.Actual, Theory: d.Theory} } for i := 0; i < 7; i++ { dayUnix := startTs + int64(i*24*3600) dayStr := time.Unix(dayUnix, 0).Format("2006-01-02") if val, ok := trendMap[dayStr]; ok { trends = append(trends, val) } else { trends = append(trends, model.TrendData{Date: dayStr, Actual: 0, Theory: 0}) } } // 7. 获取去年同周同比数据 (按星期几分组) lastYear, lastWeek := getPrevYearWeek(req.Year, req.Week) lastStartTs, lastEndTs := getWeekRange(lastYear, lastWeek) var lastYearData []weeklyTrendDB yearTrendQuery := db.Model(&dao.DeviceConsumption{}) if req.GatewayID != 0 { yearTrendQuery = yearTrendQuery.Where("gateway_id = ?", req.GatewayID) } err = yearTrendQuery. Select("WEEKDAY(CONVERT_TZ(FROM_UNIXTIME(acquisition_ts), '+00:00', '+08:00')) as weekday, "+ "SUM(actual_kwh) as actual_total, SUM(theory_kwh) as theory_total"). Where("acquisition_ts >= ? AND acquisition_ts <= ?", lastStartTs, lastEndTs). Group("weekday"). Order("weekday ASC"). Scan(&lastYearData).Error if err != nil { global.GVA_LOG.Error("获取去年同比趋势失败", zap.Error(err)) } // 星期名称映射 weekNameMap := map[int]string{0: "星期一", 1: "星期二", 2: "星期三", 3: "星期四", 4: "星期五", 5: "星期六", 6: "星期日"} yearOnYearTrends = make([]model.TrendData, 7) for i := 0; i < 7; i++ { yearOnYearTrends[i] = model.TrendData{Date: weekNameMap[i], Actual: 0, Theory: 0} } for _, d := range lastYearData { if d.Weekday >= 0 && d.Weekday <= 6 { yearOnYearTrends[d.Weekday] = model.TrendData{ Date: weekNameMap[d.Weekday], Actual: d.Actual, Theory: d.Theory, } } } } // ================= 组装返回结构 ================= return model.DashboardResponse{ Gateways: gateways, AvailableWeeks: availableWeeks, CurrentGatewayID: req.GatewayID, TotalLuminaires: int(totalLuminaires), ElectricityConsumption: actualEnergy, CarbonReduction: carbonReduction, CarbonReductionRatio: carbonReductionRatio, WarningCount: int(warningCount), LoopDeviceCount: loopDeviceCount, Devices: devices, GatewayCount: len(gateways), Trends: trends, YearOnYearTrends: yearOnYearTrends, } } // ================= 内部辅助函数 (ISO周严格计算,修正了偏移问题) ================= // getWeekRange 根据年、周数计算该周周一 00:00:00 和周日 23:59:59 func getWeekRange(year, week int) (startTs, endTs int64) { // ISO 8601标准:第一周包含第一个星期四 jan1 := time.Date(year, time.January, 1, 0, 0, 0, 0, time.Local) offsetToThursday := int(time.Thursday - jan1.Weekday()) if offsetToThursday < 0 { offsetToThursday += 7 } firstThursday := jan1.AddDate(0, 0, offsetToThursday) // 第一个星期一 = 第一个星期四 - 3天 firstMonday := firstThursday.AddDate(0, 0, -3) // 目标周的周一 targetMonday := firstMonday.AddDate(0, 0, (week-1)*7) start := time.Date(targetMonday.Year(), targetMonday.Month(), targetMonday.Day(), 0, 0, 0, 0, time.Local) end := start.AddDate(0, 0, 7).Add(-time.Second) return start.Unix(), end.Unix() } // getPrevYearWeek 获取去年的同一周 func getPrevYearWeek(curYear, curWeek int) (prevYear, prevWeek int) { // 基于相同算法稳健减去365天 jan1 := time.Date(curYear, time.January, 1, 0, 0, 0, 0, time.Local) offsetToThursday := int(time.Thursday - jan1.Weekday()) if offsetToThursday < 0 { offsetToThursday += 7 } firstThursday := jan1.AddDate(0, 0, offsetToThursday) firstMonday := firstThursday.AddDate(0, 0, -3) targetMonday := firstMonday.AddDate(0, 0, (curWeek-1)*7) lastYearMonday := targetMonday.AddDate(0, 0, -364) return lastYearMonday.ISOWeek() } // 结构体本地引用 (若全局有定义可移除) type dailyTrendDB struct { DateStr string `gorm:"column:date_str"` Actual float64 `gorm:"column:actual_total"` Theory float64 `gorm:"column:theory_total"` } type weeklyTrendDB struct { Weekday int `gorm:"column:weekday"` Actual float64 `gorm:"column:actual_total"` Theory float64 `gorm:"column:theory_total"` } // 定时删除 2 年前的记录 func CleanupOldDeviceConsumption() { // 计算 2 年的时间戳 twoYearsAgo := time.Now().AddDate(-2, 0, 0).Unix() // 循环分批删除,直到没有数据可删 for { result := global.GVA_DB.Where("acquisition_ts < ?", twoYearsAgo). Limit(5000). Delete(&dao.DeviceConsumption{}) if result.Error != nil { global.GVA_LOG.Error("清理历史明细失败", zap.Error(result.Error)) break // 遇到错误终止循环 } // 如果这一批删了 0 条,说明历史数据已经全部清完,退出循环 if result.RowsAffected == 0 { global.GVA_LOG.Info("历史数据清理完毕") break } // 稍微停顿 200 毫秒,给主库一点喘息的时间,避免 DELETE 压力过大导致主从同步延迟 time.Sleep(200 * time.Millisecond) } } func (bs *BluetoothService) QueryBluetoothList(info request.SearchBluetoothList) ([]dao.Bluetooth, int64, error) { limit := info.PageSize offset := info.PageSize * (info.Page - 1) return dao.QueryBluetoothList(info.GatewayId, limit, offset) }