Просмотр исходного кода

@author liuqing
@commit 网关配置下发功能和日志功能

lq 2 месяцев назад
Родитель
Сommit
2ad0287533

+ 5 - 4
cloud/ipolesvr/gwhandler.go

@@ -33,8 +33,8 @@ type GwHandler struct {
 }
 
 func (o *GwHandler) SubscribeTopics() {
-	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ONLINE), mqtt.AtMostOnce, o.HandlerData)
-	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_WILL), mqtt.AtMostOnce, o.HandlerData)
+	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ONLINE), mqtt.AtLeastOnce, o.HandlerData)
+	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_WILL), mqtt.AtLeastOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_APP_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SET_APP_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SERIAL_ACK), mqtt.AtMostOnce, o.HandlerData)
@@ -45,6 +45,7 @@ func (o *GwHandler) SubscribeTopics() {
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SET_MODEL_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_LOG_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_REMOVE_LOG_ACK), mqtt.AtMostOnce, o.HandlerData)
+	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_LOG_CFG_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_SYS_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ITS_ACK), mqtt.AtMostOnce, o.HandlerData)
 	GetMQTTMgr().Subscribe(GetVagueTopic(protocol.DT_GATEWAY, protocol.TP_GW_ONVIFDEV_ACK), mqtt.AtMostOnce, o.HandlerData)
@@ -100,7 +101,7 @@ func (o *GwHandler) Handler(args ...interface{}) interface{} {
 				cacheState(obj.Id, obj.Time, 1)
 				GetEventMgr().PushEvent(&EventObject{ID: obj.Id, EventType: models.ET_OFFLINE, Time: util.MlNow()})
 			}
-		case protocol.TP_GW_SET_APP_ACK, protocol.TP_GW_SET_MODEL_ACK, protocol.TP_GW_SET_RTU_ACK, protocol.TP_GW_SET_SERIAL_ACK, protocol.TP_GW_REMOVE_LOG_ACK:
+		case protocol.TP_GW_SET_APP_ACK, protocol.TP_GW_SET_MODEL_ACK, protocol.TP_GW_SET_RTU_ACK, protocol.TP_GW_SET_SERIAL_ACK, protocol.TP_GW_REMOVE_LOG_ACK, protocol.TP_GW_LOG_CFG_ACK:
 			var obj protocol.Pack_Ack
 			if err := obj.DeCode(m.PayloadString()); err == nil {
 				o := models.DeviceCmdRecord{
@@ -321,7 +322,7 @@ func (o *GwHandler) Handler(args ...interface{}) interface{} {
 }
 
 func SaveFile(tenant, GID, strType string, fo []protocol.FileObject) {
-	Dir := tenant + string(filepath.Separator) + GID + string(filepath.Separator) + strType + string(filepath.Separator)
+	Dir := "/opt/ipole_data/gateway_logs" + string(filepath.Separator) + tenant + string(filepath.Separator) + GID + string(filepath.Separator) + strType + string(filepath.Separator)
 	err := os.MkdirAll(Dir, os.ModePerm)
 	if err != nil {
 		return

+ 636 - 0
cloud/websvr/controllers/cgateway.go

@@ -2,11 +2,16 @@ package controllers
 
 import (
 	"fmt"
+	"os"
+	"strconv"
 	"strings"
+	"time"
 
 	"github.com/astaxie/beego"
 
 	"lc/common/models"
+	"lc/common/mqtt"
+	"lc/common/protocol"
 )
 
 type GatewayController struct {
@@ -151,3 +156,634 @@ func (o *GatewayController) ImportGateway() {
 	}
 	o.Response(Success, "成功", resp)
 }
+
+// PushSerialConfig @Title 下发串口配置到网关
+// @Description 下发 serial.json 配置到指定网关
+// @Param   body controllers.ReqGatewayConfig true "数据"
+// @Success 0 {int} BaseResponse.Code "成功"
+// @Failure 1 {int} BaseResponse.Code "失败"
+// @router /v1/serial/push [post]
+func (o *GatewayController) PushSerialConfig() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	if gid == "" || tenant == "" {
+		o.Response(Failure, "gid和tenant不能为空", nil)
+		return
+	}
+
+	body := o.Ctx.Input.RequestBody
+	if len(body) == 0 {
+		o.Response(Failure, "请求body不能为空,请在表单中导入或添加串口配置后再下发", nil)
+		return
+	}
+	sc := protocol.SerialConfig{Serial: make(map[uint8]*protocol.SerialPort)}
+	if err := json.Unmarshal(body, &sc); err != nil || len(sc.Serial) == 0 {
+		o.Response(Failure, "串口配置JSON解析失败或为空", nil)
+		return
+	}
+	content := string(body)
+
+	seq := GetNextUint64()
+	var obj protocol.Pack_SeqFileObject
+	obj.Data.File = "serial.json"
+	obj.Data.Content = content
+	str, err := obj.EnCode(gid, seq)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
+		return
+	}
+
+	topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_SET_SERIAL)
+	if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
+		o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
+		return
+	}
+
+	// 下发成功后,将配置存档到DB
+	pushedCodes := make(map[int]bool)
+	for _, v := range sc.Serial {
+		models.G_db.Save(&models.GatewaySerial{
+			ID: gid, ComID: int(v.Code),
+			Interface: v.Interface, Address: v.Address,
+			BaudRate: v.BaudRate, DataBits: int(v.DataBits),
+			StopBits: int(v.StopBits), Parity: v.Parity,
+			Timeout: int(v.Timeout), ProtocolType: int(v.ProtocolType),
+		})
+		pushedCodes[int(v.Code)] = true
+	}
+	// 标记不在此次推送中的旧串口为删除
+	var oldSerials []models.GatewaySerial
+	models.G_db.Where("id = ?", gid).Find(&oldSerials)
+	for _, s := range oldSerials {
+		if !pushedCodes[s.ComID] {
+			models.G_db.Delete(&s)
+		}
+	}
+
+	// 记录指令
+	dcr := models.DeviceCmdRecord{ID: seq, GID: gid, DID: gid, Topic: topic, Message: str, State: 0}
+	if err := models.G_db.Create(&dcr).Error; err != nil {
+		beego.Error("PushSerialConfig:指令入库失败:", err.Error())
+	}
+
+	o.Response(Success, "串口配置已下发并保存", nil)
+}
+
+// PushDevConfig @Title 下发设备配置到网关
+// @Description 下发 dev/{code}.json 配置到指定网关
+// @Param   body controllers.ReqDevConfig true "数据"
+// @Success 0 {int} BaseResponse.Code "成功"
+// @Failure 1 {int} BaseResponse.Code "失败"
+// @router /v1/dev/push [post]
+func (o *GatewayController) PushDevConfig() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	codeStr := strings.Trim(o.GetString("code"), " ")
+	if gid == "" || tenant == "" || codeStr == "" {
+		o.Response(Failure, "gid、tenant、code不能为空", nil)
+		return
+	}
+	code, err := strconv.Atoi(codeStr)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("code格式错误:%s", err.Error()), nil)
+		return
+	}
+
+	var mdc protocol.MapDevConfig
+	// 优先从请求body解析(前端表单直接提交的配置)
+	body := o.Ctx.Input.RequestBody
+	if len(body) == 0 {
+		o.Response(Failure, "请求body不能为空,请在表单中导入或添加设备配置后再下发", nil)
+		return
+	}
+	if err := json.Unmarshal(body, &mdc); err != nil || len(mdc.Rtu) == 0 {
+		o.Response(Failure, "设备配置JSON解析失败或为空", nil)
+		return
+	}
+	content := string(body)
+
+	seq := GetNextUint64()
+	var obj protocol.Pack_SeqFileObject
+	obj.Data.File = codeStr + ".json"
+	obj.Data.Content = content
+	str, err := obj.EnCode(gid, seq)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
+		return
+	}
+
+	topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_SET_RTU)
+	if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
+		o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
+		return
+	}
+
+	// 下发成功后,将配置存档到DB
+	pushedDevCodes := make(map[string]bool)
+	for _, v := range mdc.Rtu {
+		models.G_db.Save(&models.GatewayDevice{
+			ID: v.DevCode, Name: v.Name, GID: gid,
+			ComID: int(v.Code), RtuID: int(v.DevID), TID: int(v.TID),
+			SendCloud: v.SendCloud, WaitTime: int(v.WaitTime),
+			ProtocolType: int(v.ProtocolType), DevType: int(v.DevType),
+			Tenant: tenant, State: 1,
+		})
+		pushedDevCodes[v.DevCode] = true
+	}
+	// 标记不在此次推送中且同串口下的旧设备为删除
+	var oldDevices []models.GatewayDevice
+	models.G_db.Where("g_id = ? AND com_id = ? AND state = 1", gid, code).Find(&oldDevices)
+	for _, d := range oldDevices {
+		if !pushedDevCodes[d.ID] {
+			models.G_db.Model(&d).Update("state", 0)
+		}
+	}
+
+	// 记录指令
+	dcr := models.DeviceCmdRecord{ID: seq, GID: gid, DID: gid, Topic: topic, Message: str, State: 0}
+	if err := models.G_db.Create(&dcr).Error; err != nil {
+		beego.Error("PushDevConfig:指令入库失败:", err.Error())
+	}
+
+	o.Response(Success, "设备配置已下发并保存", nil)
+}
+
+// PushModelConfig @Title 下发物模型配置到网关
+// @Description 下发 model/{tid}.json 配置到指定网关
+// @Param   body controllers.ReqModelConfig true "数据"
+// @Success 0 {int} BaseResponse.Code "成功"
+// @Failure 1 {int} BaseResponse.Code "失败"
+// @router /v1/model/push [post]
+func (o *GatewayController) PushModelConfig() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	tidStr := strings.Trim(o.GetString("tid"), " ")
+	if gid == "" || tenant == "" || tidStr == "" {
+		o.Response(Failure, "gid、tenant、tid不能为空", nil)
+		return
+	}
+	tid, err := strconv.Atoi(tidStr)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("tid格式错误:%s", err.Error()), nil)
+		return
+	}
+
+	// 必须从请求body取模型JSON(不下发时以表单数据为准)
+	body := o.Ctx.Input.RequestBody
+	if len(body) == 0 {
+		o.Response(Failure, "请求body不能为空,请在表单中选择JSON文件后再下发", nil)
+		return
+	}
+	fileContent := string(body)
+	var iot protocol.IotModel
+	if err := json.Unmarshal(body, &iot); err != nil {
+		o.Response(Failure, fmt.Sprintf("物模型JSON解析失败:%s", err.Error()), nil)
+		return
+	}
+	device := iot.Device
+	modelName := iot.Model
+	protocolName := iot.Protocol
+
+	seq := GetNextUint64()
+	var obj protocol.Pack_SeqFileObject
+	obj.Data.File = tidStr + ".json"
+	obj.Data.Content = fileContent
+	str, err := obj.EnCode(gid, seq)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
+		return
+	}
+
+	topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_SET_MODEL)
+	if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
+		o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
+		return
+	}
+
+	// 记录指令
+	dcr := models.DeviceCmdRecord{ID: seq, GID: gid, DID: gid, Topic: topic, Message: str, State: 0}
+	if err := models.G_db.Create(&dcr).Error; err != nil {
+		beego.Error("PushModelConfig:指令入库失败:", err.Error())
+	}
+
+	// 下发成功后存档到 t_gateway_model
+	models.G_db.Save(&models.GatewayModel{
+		GID: gid, TID: uint16(tid),
+		Device: device, Model: modelName, Protocol: protocolName,
+		File: fileContent, PushedAt: time.Now(),
+	})
+
+	o.Response(Success, "物模型配置已下发并保存", nil)
+}
+
+// SerialConfigList @Title 查询网关串口列表
+// @Description 查询指定网关的所有串口配置
+// @Param   gid query string true "网关ID"
+// @Success 0 {int} BaseResponse.Code "成功"
+// @router /v1/serial/list [get]
+func (o *GatewayController) SerialConfigList() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	if gid == "" {
+		o.Response(Failure, "gid不能为空", nil)
+		return
+	}
+	var serials []models.GatewaySerial
+	if err := models.G_db.Where("id = ?", gid).Find(&serials).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "成功", serials)
+}
+
+// SerialConfigSave @Title 保存串口配置
+// @Description 新增或更新单条串口配置
+// @Param   gid query string true "网关ID"
+// @Param   comid query int true "串口编号"
+// @Success 0 {int} BaseResponse.Code "成功"
+// @router /v1/serial/save [post]
+func (o *GatewayController) SerialConfigSave() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	comIDStr := strings.Trim(o.GetString("comid"), " ")
+	if gid == "" || comIDStr == "" {
+		o.Response(Failure, "gid和comid不能为空", nil)
+		return
+	}
+	comID, err := strconv.Atoi(comIDStr)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("comid格式错误:%s", err.Error()), nil)
+		return
+	}
+	baudRate, _ := strconv.Atoi(strings.Trim(o.GetString("baudrate"), " "))
+	dataBits, _ := strconv.Atoi(strings.Trim(o.GetString("databits"), " "))
+	stopBits, _ := strconv.Atoi(strings.Trim(o.GetString("stopbits"), " "))
+	timeout, _ := strconv.Atoi(strings.Trim(o.GetString("timeout"), " "))
+	protocolType, _ := strconv.Atoi(strings.Trim(o.GetString("protocoltype"), " "))
+
+	s := models.GatewaySerial{
+		ID:           gid,
+		ComID:        comID,
+		Interface:    strings.Trim(o.GetString("interface"), " "),
+		Address:      strings.Trim(o.GetString("address"), " "),
+		BaudRate:     baudRate,
+		DataBits:     dataBits,
+		StopBits:     stopBits,
+		Parity:       strings.Trim(o.GetString("parity"), " "),
+		Timeout:      timeout,
+		ProtocolType: protocolType,
+	}
+	if err := models.G_db.Save(&s).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("保存失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "保存成功", nil)
+}
+
+// SerialConfigDelete @Title 删除串口配置
+// @Description 删除单条串口配置
+// @Param   gid query string true "网关ID"
+// @Param   comid query int true "串口编号"
+// @router /v1/serial/delete [post]
+func (o *GatewayController) SerialConfigDelete() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	comIDStr := strings.Trim(o.GetString("comid"), " ")
+	if gid == "" || comIDStr == "" {
+		o.Response(Failure, "gid和comid不能为空", nil)
+		return
+	}
+	comID, err := strconv.Atoi(comIDStr)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("comid格式错误:%s", err.Error()), nil)
+		return
+	}
+	if err := models.G_db.Where("id = ? AND com_id = ?", gid, comID).Delete(&models.GatewaySerial{}).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("删除失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "删除成功", nil)
+}
+
+// DevConfigList @Title 查询网关设备列表
+// @Description 查询指定网关某串口下的设备配置
+// @Param   gid query string true "网关ID"
+// @Param   code query int false "串口编号,0查全部"
+// @router /v1/dev/list [get]
+func (o *GatewayController) DevConfigList() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	if gid == "" {
+		o.Response(Failure, "gid不能为空", nil)
+		return
+	}
+	codeStr := strings.Trim(o.GetString("code"), " ")
+	var devices []models.GatewayDevice
+	db := models.G_db.Where("g_id = ? AND state = 1", gid)
+	if codeStr != "" {
+		code, _ := strconv.Atoi(codeStr)
+		db = db.Where("com_id = ?", code)
+	}
+	if err := db.Find(&devices).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "成功", devices)
+}
+
+// DevConfigSave @Title 保存设备配置
+// @Description 新增或更新单条设备配置
+// @router /v1/dev/save [post]
+func (o *GatewayController) DevConfigSave() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	devCode := strings.Trim(o.GetString("devcode"), " ")
+	if gid == "" || devCode == "" {
+		o.Response(Failure, "gid和devcode不能为空", nil)
+		return
+	}
+	comID, _ := strconv.Atoi(strings.Trim(o.GetString("comid"), " "))
+	rtuID, _ := strconv.Atoi(strings.Trim(o.GetString("rtuid"), " "))
+	tid, _ := strconv.Atoi(strings.Trim(o.GetString("tid"), " "))
+	protocolType, _ := strconv.Atoi(strings.Trim(o.GetString("protocoltype"), " "))
+	devType, _ := strconv.Atoi(strings.Trim(o.GetString("devtype"), " "))
+	sendCloud, _ := strconv.Atoi(strings.Trim(o.GetString("sendcloud"), " "))
+	waitTime, _ := strconv.Atoi(strings.Trim(o.GetString("waittime"), " "))
+
+	d := models.GatewayDevice{
+		ID:           devCode,
+		Name:         strings.Trim(o.GetString("name"), " "),
+		GID:          gid,
+		ComID:        comID,
+		RtuID:        rtuID,
+		TID:          tid,
+		SendCloud:    sendCloud,
+		WaitTime:     waitTime,
+		ProtocolType: protocolType,
+		DevType:      devType,
+		Tenant:       tenant,
+		State:        1,
+	}
+	if err := models.G_db.Save(&d).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("保存失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "保存成功", nil)
+}
+
+// DevConfigDelete @Title 删除设备配置
+// @Description 删除单条设备配置(软删除)
+// @Param   devcode query string true "设备编码"
+// @router /v1/dev/delete [post]
+func (o *GatewayController) DevConfigDelete() {
+	devCode := strings.Trim(o.GetString("devcode"), " ")
+	if devCode == "" {
+		o.Response(Failure, "devcode不能为空", nil)
+		return
+	}
+	if err := models.G_db.Model(&models.GatewayDevice{}).Where("id = ?", devCode).Update("state", 0).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("删除失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "删除成功", nil)
+}
+
+// ModelList @Title 查询物模型列表
+// @Description 查询网关已下发的物模型(默认查所有,传gid查指定网关)
+// @router /v1/model/list [get]
+func (o *GatewayController) ModelList() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	type ModelSummary struct {
+		TID      uint16 `json:"tid"`
+		Device   string `json:"device"`
+		Model    string `json:"model"`
+		Protocol string `json:"protocol"`
+		File     string `json:"file,omitempty"`
+	}
+	var list []ModelSummary
+	if gid != "" {
+		var gms []models.GatewayModel
+		if err := models.G_db.Where("g_id = ?", gid).Find(&gms).Error; err != nil {
+			o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
+			return
+		}
+		for _, v := range gms {
+			list = append(list, ModelSummary{TID: v.TID, Device: v.Device, Model: v.Model, Protocol: v.Protocol, File: v.File})
+		}
+	} else {
+		var gms []models.GatewayModel
+		if err := models.G_db.Find(&gms).Error; err != nil {
+			o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
+			return
+		}
+		for _, v := range gms {
+			list = append(list, ModelSummary{TID: v.TID, Device: v.Device, Model: v.Model, Protocol: v.Protocol, File: v.File})
+		}
+	}
+	o.Response(Success, "成功", list)
+}
+
+// ModelSave @Title 保存物模型
+// @Description 新增或更新物模型到 t_gateway_model
+// @router /v1/model/save [post]
+func (o *GatewayController) ModelSave() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tidStr := strings.Trim(o.GetString("tid"), " ")
+	if gid == "" || tidStr == "" {
+		o.Response(Failure, "gid和tid不能为空", nil)
+		return
+	}
+	tid, err := strconv.Atoi(tidStr)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("tid格式错误:%s", err.Error()), nil)
+		return
+	}
+	device := strings.Trim(o.GetString("device"), " ")
+	modelName := strings.Trim(o.GetString("model"), " ")
+	protocol := strings.Trim(o.GetString("protocol"), " ")
+	fileContent := strings.Trim(o.GetString("file"), " ")
+	if fileContent == "" {
+		// 尝试从上传的请求body中读取文件内容
+		fileContent = string(o.Ctx.Input.RequestBody)
+	}
+	m := models.GatewayModel{
+		GID: gid, TID: uint16(tid),
+		Device: device, Model: modelName, Protocol: protocol,
+		File: fileContent, PushedAt: time.Now(),
+	}
+	if err := models.G_db.Save(&m).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("保存失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "保存成功", nil)
+}
+
+// ModelDelete @Title 删除物模型
+// @Description 从 t_gateway_model 删除指定物模型
+// @router /v1/model/delete [post]
+func (o *GatewayController) ModelDelete() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tidStr := strings.Trim(o.GetString("tid"), " ")
+	if gid == "" || tidStr == "" {
+		o.Response(Failure, "gid和tid不能为空", nil)
+		return
+	}
+	tid, err := strconv.Atoi(tidStr)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("tid格式错误:%s", err.Error()), nil)
+		return
+	}
+	if err := models.G_db.Where("g_id = ? AND t_id = ?", gid, tid).Delete(&models.GatewayModel{}).Error; err != nil {
+		o.Response(Failure, fmt.Sprintf("删除失败:%s", err.Error()), nil)
+		return
+	}
+	o.Response(Success, "删除成功", nil)
+}
+
+// LogQuery @Title 查询边缘日志
+// @Description 触发边缘上传日志文件
+// @router /v1/log/query [post]
+func (o *GatewayController) LogQuery() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	if gid == "" || tenant == "" {
+		o.Response(Failure, "gid和tenant不能为空", nil)
+		return
+	}
+	seq := GetNextUint64()
+	var obj protocol.Pack_IDObject
+	str, err := obj.EnCode(gid, seq, 0)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
+		return
+	}
+	topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_LOG)
+	beego.Info("LogQuery:发布,topic=", topic, ",payload=", str)
+	if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
+		beego.Error("LogQuery:MQTT发布失败,topic=", topic, ",err=", err)
+		o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
+		return
+	}
+	beego.Info("LogQuery:发布成功,topic=", topic)
+	o.Response(Success, "日志查询指令已发送,请稍后刷新查看", map[string]interface{}{"seq": seq})
+}
+
+// LogList @Title 列出已保存的日志文件
+// @Description 列出网关已上传的日志文件
+// @router /v1/log/list [get]
+func (o *GatewayController) LogList() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	if gid == "" || tenant == "" {
+		o.Response(Failure, "gid和tenant不能为空", nil)
+		return
+	}
+	dir := "/opt/ipole_data/gateway_logs/" + tenant + "/" + gid + "/log/"
+	entries, err := os.ReadDir(dir)
+	if err != nil {
+		o.Response(Failure, "该网关暂无日志文件,请先点击「查询日志」", nil)
+		return
+	}
+	type LogFile struct {
+		Name string `json:"name"`
+		Size int64  `json:"size"`
+		Time string `json:"time"`
+	}
+	var files []LogFile
+	for _, e := range entries {
+		if !e.IsDir() {
+			info, _ := e.Info()
+			f := LogFile{Name: e.Name(), Size: info.Size(), Time: info.ModTime().Format("2006-01-02 15:04:05")}
+			files = append(files, f)
+		}
+	}
+	o.Response(Success, "成功", files)
+}
+
+// LogDelete @Title 删除边缘日志
+// @Description 远程删除边缘网关的所有日志文件
+// @Param   gid query string true "网关ID"
+// @Param   tenant query string true "租户ID"
+// @router /v1/log/delete [post]
+func (o *GatewayController) LogDelete() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	if gid == "" || tenant == "" {
+		o.Response(Failure, "gid和tenant不能为空", nil)
+		return
+	}
+	seq := GetNextUint64()
+	var obj protocol.Pack_IDObject
+	str, err := obj.EnCode(gid, seq, 0)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
+		return
+	}
+	topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_REMOVE_LOG)
+	beego.Info("LogDelete:发布,topic=", topic)
+	if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
+		beego.Error("LogDelete:MQTT发布失败,topic=", topic, ",err=", err)
+		o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
+		return
+	}
+	beego.Info("LogDelete:发布成功,topic=", topic)
+	o.Response(Success, "日志删除指令已发送", nil)
+}
+
+// LogCfg @Title 设置边缘日志等级
+// @Description 云端下发日志等级配置到边缘网关
+// @Param   gid query string true "网关ID"
+// @Param   tenant query string true "租户ID"
+// @Success 0 {int} BaseResponse.Code "成功"
+// @Failure 1 {int} BaseResponse.Code "失败"
+// @router /v1/log/cfg [post]
+func (o *GatewayController) LogCfg() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	if gid == "" || tenant == "" {
+		o.Response(Failure, "gid和tenant不能为空", nil)
+		return
+	}
+
+	body := o.Ctx.Input.RequestBody
+	if len(body) == 0 {
+		o.Response(Failure, "请求body为空,请提供日志等级配置", nil)
+		return
+	}
+
+	topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_LOG_CFG)
+	beego.Info("LogCfg:发布,topic=", topic)
+	if err := GetMqttHandler().PublishString(topic, string(body), mqtt.AtLeastOnce); err != nil {
+		beego.Error("LogCfg:MQTT发布失败,topic=", topic, ",err=", err)
+		o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
+		return
+	}
+	beego.Info("LogCfg:发布成功,topic=", topic)
+	o.Response(Success, "日志等级配置已下发", nil)
+}
+
+// LogRead @Title 读取日志内容
+// @Description 读取指定日志文件内容
+// @router /v1/log/read [get]
+func (o *GatewayController) LogRead() {
+	gid := strings.Trim(o.GetString("gid"), " ")
+	tenant := strings.Trim(o.GetString("tenant"), " ")
+	fname := strings.Trim(o.GetString("file"), " ")
+	if gid == "" || tenant == "" || fname == "" {
+		o.Response(Failure, "gid、tenant、file不能为空", nil)
+		return
+	}
+	// 防止路径遍历
+	if strings.Contains(fname, "..") || strings.Contains(fname, "/") || strings.Contains(fname, "\\") {
+		o.Response(Failure, "非法文件名", nil)
+		return
+	}
+	path := "/opt/ipole_data/gateway_logs/" + tenant + "/" + gid + "/log/" + fname
+	buf, err := os.ReadFile(path)
+	if err != nil {
+		o.Response(Failure, fmt.Sprintf("读取失败:%s", err.Error()), nil)
+		return
+	}
+	// 限制返回最近100KB,避免超大文件
+	content := string(buf)
+	if len(content) > 100*1024 {
+		content = content[len(content)-100*1024:]
+	}
+	o.Response(Success, "成功", content)
+}

+ 153 - 0
cloud/websvr/routers/commentsRouter_.go

@@ -457,6 +457,159 @@ func init() {
 			Filters:          nil,
 			Params:           nil})
 
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "PushSerialConfig",
+			Router:           `/v1/serial/push`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "PushDevConfig",
+			Router:           `/v1/dev/push`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "PushModelConfig",
+			Router:           `/v1/model/push`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "SerialConfigList",
+			Router:           `/v1/serial/list`,
+			AllowHTTPMethods: []string{"get"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "SerialConfigSave",
+			Router:           `/v1/serial/save`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "SerialConfigDelete",
+			Router:           `/v1/serial/delete`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "DevConfigList",
+			Router:           `/v1/dev/list`,
+			AllowHTTPMethods: []string{"get"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "DevConfigSave",
+			Router:           `/v1/dev/save`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "DevConfigDelete",
+			Router:           `/v1/dev/delete`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "ModelList",
+			Router:           `/v1/model/list`,
+			AllowHTTPMethods: []string{"get"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "ModelSave",
+			Router:           `/v1/model/save`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "ModelDelete",
+			Router:           `/v1/model/delete`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "LogQuery",
+			Router:           `/v1/log/query`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "LogList",
+			Router:           `/v1/log/list`,
+			AllowHTTPMethods: []string{"get"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "LogRead",
+			Router:           `/v1/log/read`,
+			AllowHTTPMethods: []string{"get"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "LogDelete",
+			Router:           `/v1/log/delete`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
+	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:GatewayController"],
+		beego.ControllerComments{
+			Method:           "LogCfg",
+			Router:           `/v1/log/cfg`,
+			AllowHTTPMethods: []string{"post"},
+			MethodParams:     param.Make(),
+			Filters:          nil,
+			Params:           nil})
+
 	beego.GlobalControllerRouter["lc/cloud/websvr/controllers:IotModelController"] = append(beego.GlobalControllerRouter["lc/cloud/websvr/controllers:IotModelController"],
 		beego.ControllerComments{
 			Method:           "ModelUpload",

+ 19 - 0
common/models/gatewaymodel.go

@@ -0,0 +1,19 @@
+package models
+
+import "time"
+
+// GatewayModel 网关已下发的物模型快照
+type GatewayModel struct {
+	GID       string    `gorm:"type:varchar(32);primary_key"` // 网关ID
+	TID       uint16    `gorm:"type:int;primary_key"`         // 物模型ID
+	Device    string    `gorm:"type:varchar(40)"`             // 设备名称
+	Model     string    `gorm:"type:varchar(18)"`             // 设备型号
+	Protocol  string    `gorm:"type:varchar(18)"`             // 协议类型
+	File      string    `gorm:"type:text"`                    // 模型文件内容快照
+	PushedAt  time.Time `gorm:"type:datetime"`                // 下发时间
+	CreatedAt time.Time
+}
+
+func (GatewayModel) TableName() string {
+	return "t_gateway_model"
+}

+ 3 - 0
common/models/init.go

@@ -188,6 +188,9 @@ func CreateTable() {
 	if !G_db.HasTable(&CableGuardianStatus{}) {
 		G_db.Set("gorm:table_options", "ENGINE=InnoDB").CreateTable(&CableGuardianStatus{})
 	}
+	if !G_db.HasTable(&GatewayModel{}) {
+		G_db.Set("gorm:table_options", "ENGINE=InnoDB").CreateTable(&GatewayModel{})
+	}
 }
 
 func InitDB(conf *util.MySQLConfig) {

+ 6 - 2
common/mqtt/client_wrapper.go

@@ -69,9 +69,13 @@ func (o *MqttClient) OnConnectHandler() {
 	}
 	topic, str := o.MqttOnline.GetOnlineMsg()
 	if topic != "" {
-		if err := o.PublishString(topic, str, 0); err != nil {
-			logrus.Errorf("发布上线消息失败: %v", err)
+		if err := o.PublishString(topic, str, 1); err != nil {
+			logrus.Errorf("发布上线消息失败: topic=%s, err=%v", topic, err)
+		} else {
+			logrus.Infof("发布上线消息成功: topic=%s", topic)
 		}
+	} else {
+		logrus.Warnln("发布上线消息跳过: GetOnlineMsg返回空topic,请检查appConfig.GID是否已加载")
 	}
 }
 

+ 21 - 4
common/mqtt/mgr.go

@@ -83,25 +83,42 @@ func (o *MQTTMgr) SetRestartFn(fn RestartFunc) {
 
 // Subscribe 订阅主题,支持方向选择
 func (o *MQTTMgr) Subscribe(topic string, qos QOS, handler MessageHandler, tp OptType) {
+	logrus.Infof("MQTTMgr.Subscribe:topic=%s,qos=%d,direction=%d", topic, qos, tp)
 	switch tp {
 	case ToAll:
 		if o.Cloud != nil {
 			o.Cloud.Handle(topic, handler)
-			_ = o.Cloud.Subscribe(topic, qos)
+			if err := o.Cloud.Subscribe(topic, qos); err != nil {
+				logrus.Errorf("MQTTMgr.Subscribe:Cloud订阅失败,topic=%s,err=%v", topic, err)
+			}
+		} else {
+			logrus.Warnf("MQTTMgr.Subscribe:Cloud客户端为nil,无法订阅topic=%s", topic)
 		}
 		if o.Edge != nil {
 			o.Edge.Handle(topic, handler)
-			_ = o.Edge.Subscribe(topic, qos)
+			if err := o.Edge.Subscribe(topic, qos); err != nil {
+				logrus.Errorf("MQTTMgr.Subscribe:Edge订阅失败,topic=%s,err=%v", topic, err)
+			}
+		} else {
+			logrus.Warnf("MQTTMgr.Subscribe:Edge客户端为nil,无法订阅topic=%s", topic)
 		}
 	case ToCloud:
 		if o.Cloud != nil {
 			o.Cloud.Handle(topic, handler)
-			_ = o.Cloud.Subscribe(topic, qos)
+			if err := o.Cloud.Subscribe(topic, qos); err != nil {
+				logrus.Errorf("MQTTMgr.Subscribe:Cloud订阅失败,topic=%s,err=%v", topic, err)
+			}
+		} else {
+			logrus.Warnf("MQTTMgr.Subscribe:Cloud客户端为nil,无法订阅topic=%s", topic)
 		}
 	case ToEdge:
 		if o.Edge != nil {
 			o.Edge.Handle(topic, handler)
-			_ = o.Edge.Subscribe(topic, qos)
+			if err := o.Edge.Subscribe(topic, qos); err != nil {
+				logrus.Errorf("MQTTMgr.Subscribe:Edge订阅失败,topic=%s,err=%v", topic, err)
+			}
+		} else {
+			logrus.Warnf("MQTTMgr.Subscribe:Edge客户端为nil,无法订阅topic=%s", topic)
 		}
 	}
 }

+ 8 - 8
common/mqtt/mqtt.go

@@ -93,17 +93,20 @@ func NewClient(options ClientOptions, connhandler ConnHandler) (*Client, error)
 
 	pahoOptions.SetCleanSession(false)
 
-	var client Client
+	client := &Client{
+		Options:     options,
+		router:      newRouter(),
+		connhandler: connhandler,
+	}
 	pahoOptions.SetConnectionLostHandler(client.ConnectionLostHandler) //断连
 	pahoOptions.SetOnConnectHandler(client.OnConnectHandler)           //连接
 	if t, m := connhandler.GetWill(); t != "" {
-		pahoOptions.SetWill(t, m, 0, false) //遗嘱消息
+		pahoOptions.SetWill(t, m, 1, false) //遗嘱消息,QoS=1确保送达
 	}
 
 	pahoClient := paho.NewClient(pahoOptions)
-	router := newRouter()
 	pahoClient.AddRoute("#", handle(func(message Message) {
-		routes := router.match(&message)
+		routes := client.router.match(&message)
 		for _, route := range routes {
 			m := message
 			m.vars = route.vars(&message)
@@ -112,11 +115,8 @@ func NewClient(options ClientOptions, connhandler ConnHandler) (*Client, error)
 	}))
 
 	client.client = pahoClient
-	client.Options = options
-	client.router = router
-	client.connhandler = connhandler
 
-	return &client, nil
+	return client, nil
 }
 
 // Connect tries to establish a connection with the mqtt servers

+ 2 - 0
common/protocol/topic.go

@@ -39,6 +39,8 @@ var (
 	TP_GW_LOG_ACK        string = "log/ack"
 	TP_GW_REMOVE_LOG     string = "rlog"
 	TP_GW_REMOVE_LOG_ACK string = "rlog/ack"
+	TP_GW_LOG_CFG        string = "logcfg"
+	TP_GW_LOG_CFG_ACK    string = "logcfg/ack"
 	TP_GW_ONLINE         string = "online"
 	TP_GW_WILL           string = "will"
 	TP_GW_SYS            string = "sys"

+ 233 - 0
common/util/taglog.go

@@ -0,0 +1,233 @@
+package util
+
+import (
+	"encoding/json"
+	"os"
+	"sync"
+
+	"github.com/sirupsen/logrus"
+)
+
+// TagLogger 带设备标签的日志包装器,支持运行时级别控制
+type TagLogger struct {
+	mu          sync.RWMutex
+	globalLevel logrus.Level
+	tagLevels   map[string]logrus.Level
+	persistPath string
+}
+
+// LogLevelConfig 持久化配置结构
+type LogLevelConfig struct {
+	Level string            `json:"level"`
+	Tags  map[string]string `json:"tags"`
+}
+
+var _tagLog *TagLogger
+
+// InitTagLog 初始化全局 TagLogger
+func InitTagLog(globalLevel logrus.Level) {
+	_tagLog = &TagLogger{
+		globalLevel: globalLevel,
+		tagLevels:   make(map[string]logrus.Level),
+		persistPath: "conf/logcfg.json",
+	}
+}
+
+// GetTagLog 获取全局 TagLogger 实例
+func GetTagLog() *TagLogger {
+	if _tagLog == nil {
+		InitTagLog(logrus.InfoLevel)
+	}
+	return _tagLog
+}
+
+// SetPersistPath 设置持久化文件路径
+func (o *TagLogger) SetPersistPath(path string) {
+	o.persistPath = path
+}
+
+// SetLevel 设置全局默认级别(线程安全)
+func (o *TagLogger) SetLevel(level logrus.Level) {
+	o.mu.Lock()
+	defer o.mu.Unlock()
+	o.globalLevel = level
+}
+
+// SetTagLevel 设置单个 tag 的覆盖级别
+func (o *TagLogger) SetTagLevel(tag string, level logrus.Level) {
+	o.mu.Lock()
+	defer o.mu.Unlock()
+	o.tagLevels[tag] = level
+}
+
+// RemoveTagLevel 移除 tag 的级别覆盖(恢复跟随全局)
+func (o *TagLogger) RemoveTagLevel(tag string) {
+	o.mu.Lock()
+	defer o.mu.Unlock()
+	delete(o.tagLevels, tag)
+}
+
+// GetLevel 获取指定 tag 的有效级别
+func (o *TagLogger) GetLevel(tag string) logrus.Level {
+	o.mu.RLock()
+	defer o.mu.RUnlock()
+	if lv, ok := o.tagLevels[tag]; ok {
+		return lv
+	}
+	return o.globalLevel
+}
+
+// GetGlobalLevel 获取全局默认级别
+func (o *TagLogger) GetGlobalLevel() logrus.Level {
+	o.mu.RLock()
+	defer o.mu.RUnlock()
+	return o.globalLevel
+}
+
+// GetAllTagLevels 获取所有 tag 级别(用于前端展示)
+func (o *TagLogger) GetAllTagLevels() map[string]string {
+	o.mu.RLock()
+	defer o.mu.RUnlock()
+	result := make(map[string]string)
+	for k, v := range o.tagLevels {
+		result[k] = v.String()
+	}
+	return result
+}
+
+// shouldLog 判断指定 tag 和 level 是否应输出
+func (o *TagLogger) shouldLog(tag string, level logrus.Level) bool {
+	o.mu.RLock()
+	defer o.mu.RUnlock()
+	effectiveLevel := o.globalLevel
+	if lv, ok := o.tagLevels[tag]; ok {
+		effectiveLevel = lv
+	}
+	return level <= effectiveLevel
+}
+
+// Log 返回带 tag 字段的 logrus.Entry。注意:此方法不做级别过滤,
+// 需要级别过滤的调用方应使用 Error/Info/Debug 等带级方法或 LogAt。
+func (o *TagLogger) Log(tag string) *logrus.Entry {
+	return logrus.WithField("tag", tag)
+}
+
+// LogAt 指定级别,若级别不足返回 nil
+func (o *TagLogger) LogAt(tag string, level logrus.Level) *logrus.Entry {
+	if !o.shouldLog(tag, level) {
+		return nil
+	}
+	entry := logrus.WithField("tag", tag)
+	return entry
+}
+
+// Debug ...
+func (o *TagLogger) Debug(tag string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.DebugLevel); e != nil {
+		e.Debug(args...)
+	}
+}
+
+// Info ...
+func (o *TagLogger) Info(tag string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.InfoLevel); e != nil {
+		e.Info(args...)
+	}
+}
+
+// Warn ...
+func (o *TagLogger) Warn(tag string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.WarnLevel); e != nil {
+		e.Warn(args...)
+	}
+}
+
+// Error ...
+func (o *TagLogger) Error(tag string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.ErrorLevel); e != nil {
+		e.Error(args...)
+	}
+}
+
+// Debugf ...
+func (o *TagLogger) Debugf(tag string, format string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.DebugLevel); e != nil {
+		e.Debugf(format, args...)
+	}
+}
+
+// Infof ...
+func (o *TagLogger) Infof(tag string, format string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.InfoLevel); e != nil {
+		e.Infof(format, args...)
+	}
+}
+
+// Warnf ...
+func (o *TagLogger) Warnf(tag string, format string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.WarnLevel); e != nil {
+		e.Warnf(format, args...)
+	}
+}
+
+// Errorf ...
+func (o *TagLogger) Errorf(tag string, format string, args ...interface{}) {
+	if e := o.LogAt(tag, logrus.ErrorLevel); e != nil {
+		e.Errorf(format, args...)
+	}
+}
+
+// SaveToFile 持久化当前配置到文件
+func (o *TagLogger) SaveToFile() error {
+	o.mu.RLock()
+	cfg := LogLevelConfig{
+		Level: o.globalLevel.String(),
+		Tags:  make(map[string]string, len(o.tagLevels)),
+	}
+	for k, v := range o.tagLevels {
+		cfg.Tags[k] = v.String()
+	}
+	o.mu.RUnlock()
+
+	f, err := os.Create(o.persistPath)
+	if err != nil {
+		return err
+	}
+	defer f.Close()
+
+	data, err := json.MarshalIndent(cfg, "", "  ")
+	if err != nil {
+		return err
+	}
+	_, err = f.Write(data)
+	return err
+}
+
+// LoadFromFile 从文件加载配置
+func (o *TagLogger) LoadFromFile() error {
+	data, err := os.ReadFile(o.persistPath)
+	if err != nil {
+		return err
+	}
+	var cfg LogLevelConfig
+	if err := json.Unmarshal(data, &cfg); err != nil {
+		return err
+	}
+
+	o.mu.Lock()
+	defer o.mu.Unlock()
+
+	if lv, err := logrus.ParseLevel(cfg.Level); err == nil {
+		o.globalLevel = lv
+	} else if cfg.Level != "" {
+		logrus.Warnf("TagLogger.LoadFromFile:无效的全局级别'%s',使用默认值", cfg.Level)
+	}
+	for k, v := range cfg.Tags {
+		if lv, err := logrus.ParseLevel(v); err == nil {
+			o.tagLevels[k] = lv
+		} else {
+			logrus.Warnf("TagLogger.LoadFromFile:无效的tag级别'%s=%s',已忽略", k, v)
+		}
+	}
+	return nil
+}

+ 12 - 13
edge/ipole/auto_reload.go

@@ -7,7 +7,6 @@ import (
 	"time"
 
 	"github.com/radovskyb/watcher"
-	"github.com/sirupsen/logrus"
 
 	"lc/common/util"
 )
@@ -29,7 +28,7 @@ func WatchDevConfig(args ...interface{}) interface{} {
 						case watcher.Create, watcher.Write:
 							rtuArr, err := LoadDev(uint8(code))
 							if err != nil {
-								logrus.Errorf("加载RTU文件[%s]失败:%s", strconv.Itoa(code)+".json", err.Error())
+								util.GetTagLog().Errorf("sys", "加载RTU文件[%s]失败:%s", strconv.Itoa(code)+".json", err.Error())
 							} else {
 								GetDeviceMgr().UpdateDevices(uint8(code), rtuArr)
 							}
@@ -40,7 +39,7 @@ func WatchDevConfig(args ...interface{}) interface{} {
 					}
 				}
 			case err := <-w.Error:
-				logrus.Errorf("WatchDevConfig:发生错误:%s", err.Error())
+				util.GetTagLog().Errorf("sys", "WatchDevConfig:发生错误:%s", err.Error())
 			case <-w.Closed:
 				return
 			}
@@ -48,10 +47,10 @@ func WatchDevConfig(args ...interface{}) interface{} {
 	}()
 
 	if err := w.Add(util.GetPath(1)); err != nil {
-		logrus.Errorf("WatchDevConfig:Add(conf)发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "WatchDevConfig:Add(conf)发生错误:%s", err.Error())
 	}
 	if err := w.Start(time.Second); err != nil {
-		logrus.Errorf("WatchDevConfig:Start发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "WatchDevConfig:Start发生错误:%s", err.Error())
 	}
 	return 0
 }
@@ -81,17 +80,17 @@ func WatchModelConfig(args ...interface{}) interface{} {
 					}
 				}
 			case err := <-w.Error:
-				logrus.Errorf("WatchModelConfig:发生错误:%s", err.Error())
+				util.GetTagLog().Errorf("sys", "WatchModelConfig:发生错误:%s", err.Error())
 			case <-w.Closed:
 				return
 			}
 		}
 	}()
 	if err := w.Add(util.GetPath(2)); err != nil {
-		logrus.Errorf("WatchModelConfig:Add(conf)发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "WatchModelConfig:Add(conf)发生错误:%s", err.Error())
 	}
 	if err := w.Start(time.Second); err != nil {
-		logrus.Errorf("WatchModelConfig:Start发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "WatchModelConfig:Start发生错误:%s", err.Error())
 	}
 	return 0
 }
@@ -117,31 +116,31 @@ func WatchConfConfig(args ...interface{}) interface{} {
 					}
 				}
 			case err := <-w.Error:
-				logrus.Errorf("WatchConfConfig:发生错误:%s", err.Error())
+				util.GetTagLog().Errorf("sys", "WatchConfConfig:发生错误:%s", err.Error())
 			case <-w.Closed:
 				return
 			}
 		}
 	}()
 	if err := w.Add(util.GetPath(0)); err != nil {
-		logrus.Errorf("WatchConfConfig:Add(conf)发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "WatchConfConfig:Add(conf)发生错误:%s", err.Error())
 	}
 	if err := w.Start(time.Second); err != nil {
-		logrus.Errorf("WatchConfConfig:Start发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "WatchConfConfig:Start发生错误:%s", err.Error())
 	}
 	return 0
 }
 
 func UpdateAppConfig() {
 	if err := loadAppConfig(); err != nil {
-		logrus.Errorf("UpdateAppConfig:重新加载app.json文件失败:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "UpdateAppConfig:重新加载app.json文件失败:%s", err.Error())
 		return
 	}
 }
 
 func UpdateSerialConfig() {
 	if err := loadSerialConfig(); err != nil {
-		logrus.Errorf("UpdateSerialConfig:重新加载serial.json文件失败:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "UpdateSerialConfig:重新加载serial.json文件失败:%s", err.Error())
 		return
 	}
 	GetSerialMgr().AddSerialPorts(serialConfig.Serial)

+ 114 - 0
edge/ipole/cleanlog.go

@@ -0,0 +1,114 @@
+package main
+
+import (
+	"os"
+	"path/filepath"
+	"sort"
+	"strings"
+	"time"
+
+	"lc/common/util"
+)
+
+const (
+	logRetentionDays = 7
+	logMaxTotalSize  = 100 * 1024 * 1024 // 100MB
+)
+
+type logFileInfo struct {
+	name    string
+	size    int64
+	modTime time.Time
+}
+
+// removeLogFile 删除单个日志文件并返回释放的字节数
+func removeLogFile(dir string, f logFileInfo) (bool, int64) {
+	path := filepath.Join(dir, f.name)
+	if err := os.Remove(path); err != nil {
+		util.GetTagLog().Warnf("sys", "CleanLogs:删除失败:%s,err=%v", f.name, err)
+		return false, 0
+	}
+	return true, f.size
+}
+
+// CleanLogs 定时清理日志文件:保留7天,总大小不超过100MB
+func CleanLogs(args ...interface{}) interface{} {
+	for {
+		time.Sleep(1 * time.Hour)
+
+		dir := util.GetPath(3) // {workdir}/log/
+		entries, err := os.ReadDir(dir)
+		if err != nil {
+			util.GetTagLog().Warnf("sys", "CleanLogs:读取日志目录失败:%s", err)
+			continue
+		}
+
+		var logFiles []logFileInfo
+		var totalSize int64
+		now := time.Now()
+
+		for _, e := range entries {
+			if e.IsDir() || !strings.HasSuffix(e.Name(), ".log") {
+				continue
+			}
+			info, err := e.Info()
+			if err != nil {
+				continue
+			}
+			lf := logFileInfo{
+				name:    e.Name(),
+				size:    info.Size(),
+				modTime: info.ModTime(),
+			}
+			logFiles = append(logFiles, lf)
+			totalSize += lf.size
+		}
+
+		if len(logFiles) == 0 {
+			continue
+		}
+
+		// 按修改时间升序排列(最旧的在前)
+		sort.Slice(logFiles, func(i, j int) bool {
+			return logFiles[i].modTime.Before(logFiles[j].modTime)
+		})
+
+		var deleted int
+		var freed int64
+
+		// Step 1: 删除超过7天的文件
+		cutoff := now.AddDate(0, 0, -logRetentionDays)
+		remaining := make([]logFileInfo, 0, len(logFiles))
+		for _, f := range logFiles {
+			if f.modTime.Before(cutoff) {
+				if ok, sz := removeLogFile(dir, f); ok {
+					deleted++
+					freed += sz
+					totalSize -= sz
+				}
+			} else {
+				remaining = append(remaining, f)
+			}
+		}
+
+		// Step 2: 若剩余总大小超过100MB,从最旧的文件开始删除(索引遍历避免切片复制)
+		idx := 0
+		for totalSize > logMaxTotalSize && idx < len(remaining) {
+			if ok, sz := removeLogFile(dir, remaining[idx]); ok {
+				deleted++
+				freed += sz
+				totalSize -= sz
+			}
+			idx++
+		}
+
+		if deleted > 0 {
+			kept := len(remaining) - idx
+			if kept < 0 {
+				kept = 0
+			}
+			util.GetTagLog().Infof("sys", "CleanLogs:清理完成,删除%d个文件,释放%.1fKB,剩余%d个文件,总大小%.1fKB",
+				deleted, float64(freed)/1024, kept, float64(totalSize)/1024)
+		}
+	}
+}

+ 43 - 44
edge/ipole/concentrator.go

@@ -11,7 +11,6 @@ import (
 	"time"
 
 	"github.com/go-redis/redis/v7"
-	"github.com/sirupsen/logrus"
 	"github.com/valyala/bytebufferpool"
 
 	"lc/common/mqtt"
@@ -185,20 +184,20 @@ func (o *Concentrator) UpdateModel2(mi *ModelInfo) {
 		return
 	}
 	if mi.Flag == 0 {
-		logrus.Errorf("Concentrator.UpdateModel2:设备[%s]的物模型[tid=%d]模型文件被删除,下次启动即将生效。", o.devinfo.DevCode, mi.TID)
+		util.GetTagLog().Errorf("concentrator", "Concentrator.UpdateModel2:设备[%s]的物模型[tid=%d]模型文件被删除,下次启动即将生效。", o.devinfo.DevCode, mi.TID)
 		return
 	}
-	logrus.Debugf("Concentrator.UpdateModel2:更新设备[%s]的物模型[%d]", o.devinfo.DevCode, mi.TID)
+	util.GetTagLog().Debugf("concentrator", "Concentrator.UpdateModel2:更新设备[%s]的物模型[%d]", o.devinfo.DevCode, mi.TID)
 	iot, err := loadModel(mi.TID)
 	if err != nil {
-		logrus.Errorf("Concentrator.UpdateModel2:加载模型[%d]文件错误:%s", mi.TID, err.Error())
+		util.GetTagLog().Errorf("concentrator", "Concentrator.UpdateModel2:加载模型[%d]文件错误:%s", mi.TID, err.Error())
 		return
 	}
 	if iot.Protocol == ConcentratorProtocol { //合法的物模型
 		o.model = iot
-		logrus.Infof("Concentrator.UpdateModel2:更新设备[%s]的物模型[%d]成功", o.devinfo.DevCode, mi.TID)
+		util.GetTagLog().Infof("concentrator", "Concentrator.UpdateModel2:更新设备[%s]的物模型[%d]成功", o.devinfo.DevCode, mi.TID)
 	} else {
-		logrus.Error("Concentrator.UpdateModel2:物模型错误,TID和文件名tid不一致或协议非ModbusRTU协议")
+		util.GetTagLog().Error("concentrator", "Concentrator.UpdateModel2:物模型错误,TID和文件名tid不一致或协议非ModbusRTU协议")
 	}
 }
 
@@ -208,7 +207,7 @@ func (o *Concentrator) ReloadOOTFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("Concentrator.ReloadOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("concentrator", "Concentrator.ReloadOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	for k, v := range mapdata {
@@ -227,7 +226,7 @@ func (o *Concentrator) ReloadSwitchOOTFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("Concentrator.ReloadSwitchOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("concentrator", "Concentrator.ReloadSwitchOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	for k, v := range mapdata {
@@ -246,7 +245,7 @@ func (o *Concentrator) ReloadBroadCastFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("Concentrator.ReloadBroadCastFromRedis设备[%s]从redis加载广播恢复截止时间失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("concentrator", "Concentrator.ReloadBroadCastFromRedis设备[%s]从redis加载广播恢复截止时间失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	if t, err := util.MlParseTime(strTime); err == nil {
@@ -261,7 +260,7 @@ func (o *Concentrator) ReloadLampAlarmFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("Concentrator.ReloadLampAlarmFromRedis设备[%s]从redis加载广播恢复截止时间失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("concentrator", "Concentrator.ReloadLampAlarmFromRedis设备[%s]从redis加载广播恢复截止时间失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	for k, v := range mapAlarm {
@@ -269,7 +268,7 @@ func (o *Concentrator) ReloadLampAlarmFromRedis() error {
 		if err := json.UnmarshalFromString(v, &lai); err == nil {
 			o.mapLampAlarm[k] = &lai
 		} else {
-			logrus.Errorf("从redis获取的告警信息还原失败,原内容:%s,失败原因:%s", v, err.Error())
+			util.GetTagLog().Errorf("concentrator", "从redis获取的告警信息还原失败,原内容:%s,失败原因:%s", v, err.Error())
 		}
 	}
 	return nil
@@ -278,8 +277,8 @@ func (o *Concentrator) ReloadLampAlarmFromRedis() error {
 func (o *Concentrator) Handle() {
 	defer func() {
 		if err := recover(); err != nil {
-			logrus.Errorf("Concentrator.Handle发生异常:%v", err)
-			logrus.Errorf("Concentrator.Handle发生异常,堆栈信息:%s", string(debug.Stack()))
+			util.GetTagLog().Errorf("concentrator", "Concentrator.Handle发生异常:%v", err)
+			util.GetTagLog().Errorf("concentrator", "Concentrator.Handle发生异常,堆栈信息:%s", string(debug.Stack()))
 			time.Sleep(5 * time.Second)
 			go o.Handle()
 		}
@@ -295,7 +294,7 @@ func (o *Concentrator) Handle() {
 	for {
 		select {
 		case <-o.ctx.Done():
-			logrus.Errorf("设备[%s]的HandlePole退出,原因:%v", o.devinfo.DevCode, o.ctx.Err())
+			util.GetTagLog().Errorf("concentrator", "设备[%s]的HandlePole退出,原因:%v", o.devinfo.DevCode, o.ctx.Err())
 			exit = true
 		case devinfo_ := <-o.chanDevInfo:
 			o.devinfo = devinfo_
@@ -308,7 +307,7 @@ func (o *Concentrator) Handle() {
 					if fn, ok := o.mapTopicHandle[mm.Topic()]; ok {
 						fn(mm)
 					} else {
-						logrus.Errorf("Concentrator.Handle:不支持的主题:%s", mm.Topic())
+						util.GetTagLog().Errorf("concentrator", "Concentrator.Handle:不支持的主题:%s", mm.Topic())
 					}
 				}
 			} else {
@@ -350,14 +349,14 @@ func (o *Concentrator) CheckRecoveryAuto(force bool) {
 			delete(o.mapTempLampsOOT, k)
 		}
 		if err := o.BroadcastAuto(); err != nil {
-			logrus.Errorf("广播恢复时控失败:%s", err.Error())
+			util.GetTagLog().Errorf("concentrator", "广播恢复时控失败:%s", err.Error())
 		} else {
 			o.broadcastAutoTime = time.Time{}
 			//删除redis中所有临时开关灯记录
 			if err := redisEdgeData.Del(LampSwitchPrefix + o.devinfo.DevCode).Err(); err != nil {
-				logrus.Errorf("手动广播恢复时控,更新redis失败:%s", err.Error())
+				util.GetTagLog().Errorf("concentrator", "手动广播恢复时控,更新redis失败:%s", err.Error())
 			}
-			logrus.Info("广播恢复时控成功")
+			util.GetTagLog().Info("concentrator", "广播恢复时控成功")
 		}
 	} else {
 		//如果广播控模式未过期,则判断单灯是否手动控过期,过期则恢复
@@ -365,9 +364,9 @@ func (o *Concentrator) CheckRecoveryAuto(force bool) {
 		for k, v := range o.mapTempLampsOOT {
 			if time.Time(v.End).Before(util.MlNow()) {
 				if err := o.SetPoleAuto(k); err != nil {
-					logrus.Errorf("单灯[%d]恢复时控失败:%s", k, err.Error())
+					util.GetTagLog().Errorf("concentrator", "单灯[%d]恢复时控失败:%s", k, err.Error())
 				} else {
-					logrus.Infof("单灯[%d]恢复时控成功", k)
+					util.GetTagLog().Infof("concentrator", "单灯[%d]恢复时控成功", k)
 					strList = append(strList, strconv.Itoa(int(k)))
 					delete(o.mapTempLampsOOT, k)
 				}
@@ -375,7 +374,7 @@ func (o *Concentrator) CheckRecoveryAuto(force bool) {
 		}
 		if len(strList) > 0 {
 			if err := redisEdgeData.HDel(LampSwitchPrefix+o.devinfo.DevCode, strList...).Err(); err != nil {
-				logrus.Errorf("手动恢复灯控[%v]时控模式,更新redis失败:%s", strList, err.Error())
+				util.GetTagLog().Errorf("concentrator", "手动恢复灯控[%v]时控模式,更新redis失败:%s", strList, err.Error())
 			}
 		}
 	}
@@ -472,7 +471,7 @@ func (o *Concentrator) UploadLampAlarm() {
 			mapRedis := make(map[string]interface{})
 			mapRedis[k] = strAlarm
 			if err := redisEdgeData.HSet(LampAlarmPrefix+o.devinfo.DevCode, mapRedis).Err(); err != nil {
-				logrus.Errorf("告警信息缓存入redis失败:%s", err.Error())
+				util.GetTagLog().Errorf("concentrator", "告警信息缓存入redis失败:%s", err.Error())
 			}
 
 		} else if v.Alarm.EndTime != "" { //告警结束上报
@@ -486,7 +485,7 @@ func (o *Concentrator) UploadLampAlarm() {
 	}
 	if len(toDelete) > 0 {
 		if err := redisEdgeData.HDel(LampAlarmPrefix+o.devinfo.DevCode, toDelete...).Err(); err != nil {
-			logrus.Errorf("告警信息从redis删除失败:%s", err.Error())
+			util.GetTagLog().Errorf("concentrator", "告警信息从redis删除失败:%s", err.Error())
 		}
 	}
 }
@@ -544,14 +543,14 @@ func (o *Concentrator) CheckLampAlarm(lnd LampNumberDID, b1, b2 uint8) {
 		if except == protocol.LE_OK { //告警结束
 			a.Alarm.EndTime = now.Format("2006-01-02 15:04:05")
 			a.Alarm.Brightness = b1
-			logrus.Debugf("灯控[%s]状态[%s]", lnd.DID, ConvertExcept(except))
+			util.GetTagLog().Debugf("concentrator", "灯控[%s]状态[%s]", lnd.DID, ConvertExcept(except))
 		}
 	} else {
 		if except != protocol.LE_OK { //告警开始
 			a := protocol.LampAlarm{DID: lnd.DID, AlarmType: except, AlarmBrightness: b1, StartTime: now.Format("2006-01-02 15:04:05")}
 			lai := LampAlarmInfo{Alarm: &a, Send: false}
 			o.mapLampAlarm[lnd.DID] = &lai
-			logrus.Debugf("灯控[%s]状态[%s]", lnd.DID, ConvertExcept(except))
+			util.GetTagLog().Debugf("concentrator", "灯控[%s]状态[%s]", lnd.DID, ConvertExcept(except))
 		}
 	}
 }
@@ -703,7 +702,7 @@ func (o *Concentrator) GetOnOffTime(PoleID uint32, Cmd uint8) ([]zigbee.OnOffTim
 		}
 		return ret, nil
 	} else {
-		logrus.Errorf("读取开关灯时间返回的内容错误:%s", hex.EncodeToString(recvdata))
+		util.GetTagLog().Errorf("concentrator", "读取开关灯时间返回的内容错误:%s", hex.EncodeToString(recvdata))
 	}
 	return nil, errors.New("读取开关灯时间返回的内容错误")
 }
@@ -732,7 +731,7 @@ func (o *Concentrator) ReadPoleTime(PoleID uint32) (uint8, uint8, uint8, error)
 	if len(pgfcresp.Data) >= 5 {
 		return pgfcresp.Data[2], pgfcresp.Data[3], pgfcresp.Data[4], nil
 	} else {
-		logrus.Errorf("读取单灯时间返回的内容错误:%s", hex.EncodeToString(recvdata))
+		util.GetTagLog().Errorf("concentrator", "读取单灯时间返回的内容错误:%s", hex.EncodeToString(recvdata))
 	}
 	return 0, 0, 0, errors.New("读取单灯时间返回的内容错误")
 }
@@ -774,7 +773,7 @@ func (o *Concentrator) ReadElectricalPara(PoleID uint32) (*ElecPara, error) {
 		ep.Degree[1] = binary.BigEndian.Uint16(pgfcresp.Data[11:13])
 		return &ep, nil
 	} else {
-		logrus.Errorf("读取电流电压电度返回的内容错误:%s", hex.EncodeToString(recvdata))
+		util.GetTagLog().Errorf("concentrator", "读取电流电压电度返回的内容错误:%s", hex.EncodeToString(recvdata))
 	}
 	return nil, errors.New("读取电流电压电度返回错误")
 }
@@ -819,7 +818,7 @@ func (o *Concentrator) GetBrightness(PoleID uint32) (uint8, uint8, error) {
 	if pgfcresp.Cmd == zigbee.CmdReadBrightness && len(pgfcresp.Data) >= 4 {
 		return pgfcresp.Data[2], pgfcresp.Data[3], nil
 	} else {
-		logrus.Errorf("查询亮度返回的内容错误:%s", hex.EncodeToString(recvdata))
+		util.GetTagLog().Errorf("concentrator", "查询亮度返回的内容错误:%s", hex.EncodeToString(recvdata))
 	}
 	return 0, 0, errors.New("查询亮度返回错误")
 }
@@ -861,7 +860,7 @@ func (o *Concentrator) HandleTpChzbSetBroadcasttime(m mqtt.Message) {
 	var ret protocol.Pack_Ack
 	var err error
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	err = o.BroadcastTime()
@@ -875,7 +874,7 @@ func (o *Concentrator) HandleTpChzbSetWaittime(m mqtt.Message) {
 	var ret protocol.Pack_Ack
 	var err error
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Data.Waittime < 1000 || obj.Data.Waittime > 15000 {
@@ -894,7 +893,7 @@ func (o *Concentrator) HandleTpChzbSetSwitch(m mqtt.Message) {
 	var err error
 	mapIpole := make(map[uint32]*protocol.StateError)
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
@@ -931,7 +930,7 @@ func (o *Concentrator) HandleTpChzbSetSwitch(m mqtt.Message) {
 		}
 	}
 	if err := redisEdgeData.HSet(LampSwitchPrefix+o.devinfo.DevCode, mapRedisTempLampsOOT).Err(); err != nil {
-		logrus.Errorf("手动开关灯时间设置[内容:%v]缓存到redis失败:%s", mapRedisTempLampsOOT, err.Error())
+		util.GetTagLog().Errorf("concentrator", "手动开关灯时间设置[内容:%v]缓存到redis失败:%s", mapRedisTempLampsOOT, err.Error())
 	}
 	if str, err := ret.EnCode(o.devinfo.DevCode, appConfig.GID, obj.Seq, mapIpole); err == nil {
 		GetMQTTMgr().Publish(GetTopic(protocol.DT_CONCENTRATOR, o.devinfo.DevCode, protocol.TP_CHZB_SET_SWITCH_ACK), str, mqtt.AtMostOnce, ToAll)
@@ -943,7 +942,7 @@ func (o *Concentrator) HandleTpChzbSetRecoveryAuto(m mqtt.Message) {
 	var err error
 	mapIpole := make(map[uint32]*protocol.StateError)
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
@@ -959,12 +958,12 @@ func (o *Concentrator) HandleTpChzbSetRecoveryAuto(m mqtt.Message) {
 				delete(o.mapTempLampsOOT, v)
 				strList = append(strList, strconv.Itoa(int(v)))
 			} else {
-				logrus.Errorf("手动单灯[%d]恢复时控失败:%s", v, err.Error())
+				util.GetTagLog().Errorf("concentrator", "手动单灯[%d]恢复时控失败:%s", v, err.Error())
 			}
 			mapIpole[v] = protocol.NewStateError(err)
 		}
 		if err := redisEdgeData.HDel(LampSwitchPrefix+o.devinfo.DevCode, strList...).Err(); err != nil {
-			logrus.Errorf("手动恢复灯控[%v]时控模式,更新redis失败:%s", strList, err.Error())
+			util.GetTagLog().Errorf("concentrator", "手动恢复灯控[%v]时控模式,更新redis失败:%s", strList, err.Error())
 		}
 	}
 
@@ -979,14 +978,14 @@ func (o *Concentrator) HandleTpChzbSetOnofftime(m mqtt.Message) {
 	var err error
 	mapIpole := make(map[uint32]*protocol.StateError)
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
 		return
 	}
 	if len(obj.Data.LampIDs) == 0 || len(obj.Data.OnOffTime) == 0 {
-		logrus.Errorf("Handle_TP_CHZB_SET_ONOFFTIME:错误,灯控编号[%v],时间段个数:%v",
+		util.GetTagLog().Errorf("concentrator", "Handle_TP_CHZB_SET_ONOFFTIME:错误,灯控编号[%v],时间段个数:%v",
 			obj.Data.LampIDs, obj.Data.OnOffTime)
 		return
 	}
@@ -1008,9 +1007,9 @@ func (o *Concentrator) HandleTpChzbSetOnofftime(m mqtt.Message) {
 			continue
 		}
 		if err = o.SetOnOffTime(v, zigbee.CmdSetOnofftime, data); err != nil {
-			logrus.Errorf("单灯[%d]设置开关灯时间失败:%s", v, err.Error())
+			util.GetTagLog().Errorf("concentrator", "单灯[%d]设置开关灯时间失败:%s", v, err.Error())
 		} else {
-			logrus.Infof("单灯[%d]设置开关灯时间成功", v)
+			util.GetTagLog().Infof("concentrator", "单灯[%d]设置开关灯时间成功", v)
 		}
 		mapIpole[v] = protocol.NewStateError(err)
 		mapRedisOOT[strconv.Itoa(int(v))] = datastr //缓存到redis
@@ -1018,7 +1017,7 @@ func (o *Concentrator) HandleTpChzbSetOnofftime(m mqtt.Message) {
 	}
 	//持久缓存到redis,以便于重启后读取进内存中
 	if err := redisEdgeData.HSet(LampOotPrefix+o.devinfo.DevCode, mapRedisOOT).Err(); err != nil {
-		logrus.Errorf("灯控时间设置[内容:%v]缓存到redis失败:%s", mapRedisOOT, err.Error())
+		util.GetTagLog().Errorf("concentrator", "灯控时间设置[内容:%v]缓存到redis失败:%s", mapRedisOOT, err.Error())
 	}
 	if str, err := ret.EnCode(o.devinfo.DevCode, appConfig.GID, obj.Seq, mapIpole); err == nil {
 		GetMQTTMgr().Publish(GetTopic(protocol.DT_CONCENTRATOR, o.devinfo.DevCode, protocol.TP_CHZB_SET_ONOFFTIME_ACK), str, mqtt.AtMostOnce, ToCloud)
@@ -1030,7 +1029,7 @@ func (o *Concentrator) HandleTpChzbQueryOnofftime(m mqtt.Message) {
 	var oot []zigbee.OnOffTime
 	var err error
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
@@ -1061,7 +1060,7 @@ func (o *Concentrator) HandleTpChzbSetUpdateLamp(m mqtt.Message) {
 	var ret protocol.Pack_Ack
 	var err error
 	if err = obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
@@ -1086,7 +1085,7 @@ func (o *Concentrator) HandleTpChzbQueryTime(m mqtt.Message) {
 	var obj protocol.Pack_CHZB_QueryTime
 	var ret protocol.Pack_CHZB_QueryTimeAck
 	if err := obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("concentrator", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	var pt *protocol.CHZB_LampTime = nil

+ 2 - 3
edge/ipole/config.go

@@ -2,7 +2,6 @@ package main
 
 import (
 	"errors"
-	"github.com/sirupsen/logrus"
 	"strconv"
 
 	"lc/common/configor"
@@ -35,7 +34,7 @@ func loadSerialConfig() error {
 	return err
 }
 
-//加载模型文件
+// 加载模型文件
 func loadModel(tid uint16) (*protocol.IotModel, error) {
 	var o protocol.IotModel
 	var fname = util.GetPath(2) + strconv.Itoa(int(tid)) + ".json"
@@ -54,7 +53,7 @@ func LoadDev(code uint8) ([]protocol.DevInfo, error) {
 	var arrDev protocol.MapDevConfig
 	err := configor.Load(&arrDev, util.GetPath(1)+strconv.Itoa(int(code))+".json")
 	if err != nil {
-		logrus.Errorf("LoadDev err = :%s", err.Error())
+		util.GetTagLog().Errorf("sys", "LoadDev err = :%s", err.Error())
 		return nil, err
 
 	}

+ 7 - 8
edge/ipole/devmgr.go

@@ -4,9 +4,8 @@ import (
 	"fmt"
 	"sync"
 
-	"github.com/sirupsen/logrus"
-
 	"lc/common/protocol"
+	"lc/common/util"
 )
 
 type ModelInfo struct {
@@ -32,7 +31,7 @@ func CreateDevice(t uint8, devinfo *protocol.DevInfo) Device {
 	case 2:
 		device = NewYmLampController(devinfo)
 	default:
-		logrus.Errorf("AddDevices:不支持的协议的设备:mode=%d", t)
+		util.GetTagLog().Errorf("sys", "AddDevices:不支持的协议的设备:mode=%d", t)
 	}
 	return device
 }
@@ -94,18 +93,18 @@ func (o *DeviceMgr) checkDevIDConflict(v *protocol.DevInfo) error {
 // AddDevices 新增
 func (o *DeviceMgr) AddDevices(devinfos []protocol.DevInfo) {
 	if err := validateDevInfos(devinfos); err != nil {
-		logrus.Errorf("AddDevices: 配置校验失败: %v", err)
+		util.GetTagLog().Errorf("sys", "AddDevices: 配置校验失败: %v", err)
 		return
 	}
 	o.mu.Lock()
 	defer o.mu.Unlock()
 	for _, v := range devinfos {
 		if _, ok := o.mapDevice[v.DevCode]; ok {
-			logrus.Errorf("AddDevices:已添加该设备,不可重复添加:devcode=%s", v.DevCode)
+			util.GetTagLog().Errorf("sys", "AddDevices:已添加该设备,不可重复添加:devcode=%s", v.DevCode)
 			continue
 		}
 		if err := o.checkDevIDConflict(&v); err != nil {
-			logrus.Errorf("AddDevices: %v", err)
+			util.GetTagLog().Errorf("sys", "AddDevices: %v", err)
 			continue
 		}
 		di := v
@@ -130,7 +129,7 @@ func (o *DeviceMgr) RemoveDevice(devcode string) {
 func (o *DeviceMgr) UpdateDevices(code uint8, devinfos []protocol.DevInfo) {
 	if len(devinfos) > 0 {
 		if err := validateDevInfos(devinfos); err != nil {
-			logrus.Errorf("UpdateDevices: 配置校验失败: %v", err)
+			util.GetTagLog().Errorf("sys", "UpdateDevices: 配置校验失败: %v", err)
 			return
 		}
 	}
@@ -179,7 +178,7 @@ func (o *DeviceMgr) UpdateDevices(code uint8, devinfos []protocol.DevInfo) {
 			}
 		} else {
 			if err := o.checkDevIDConflict(&v); err != nil {
-				logrus.Errorf("UpdateDevices: %v", err)
+				util.GetTagLog().Errorf("sys", "UpdateDevices: %v", err)
 				continue
 			}
 		}

+ 24 - 6
edge/ipole/handlefile.go

@@ -19,7 +19,7 @@ func HandleFile(topic string, obj *protocol.Pack_SeqFileObject) error {
 		}
 		return SaveFile(obj.Data.File, obj.Data.Content, 0)
 	case protocol.TP_GW_SET_SERIAL:
-		var sc protocol.SerialConfig
+		sc := protocol.SerialConfig{Serial: make(map[uint8]*protocol.SerialPort)}
 		if err := json.UnmarshalFromString(obj.Data.Content, &sc); err != nil {
 			return err
 		}
@@ -50,23 +50,41 @@ func SaveFile(fname, content string, flag uint8) error {
 	case 2:
 		SavePath = util.GetPath(2)
 	}
+	tmpPath := util.GetPath(4)
+	util.GetTagLog().Infof("sys", "SaveFile:开始,fname=%s,flag=%d,tmpPath=%s,savePath=%s", fname, flag, tmpPath, SavePath)
+
 	//存成临时文件
-	if err := os.WriteFile(util.GetPath(4)+fname, []byte(content), os.ModePerm); err != nil {
+	if err := os.MkdirAll(tmpPath, os.ModePerm); err != nil {
+		util.GetTagLog().Errorf("sys", "SaveFile:MkdirAll失败,path=%s,err=%v", tmpPath, err)
 		return err
 	}
-	//备份原文件
-	if err := os.Rename(SavePath+fname, SavePath+fname+".bak"); err != nil {
+	if err := os.WriteFile(tmpPath+fname, []byte(content), os.ModePerm); err != nil {
+		util.GetTagLog().Errorf("sys", "SaveFile:写tmp文件失败,path=%s,err=%v", tmpPath+fname, err)
 		return err
 	}
+	util.GetTagLog().Infof("sys", "SaveFile:tmp文件写入成功,path=%s,len=%d", tmpPath+fname, len(content))
+
+	//备份原文件(不存在则跳过)
+	if _, err := os.Stat(SavePath + fname); err == nil {
+		if err := os.Rename(SavePath+fname, SavePath+fname+".bak"); err != nil {
+			util.GetTagLog().Errorf("sys", "SaveFile:备份失败,src=%s,err=%v", SavePath+fname, err)
+			return err
+		}
+		util.GetTagLog().Infof("sys", "SaveFile:备份成功,src=%s", SavePath+fname+".bak")
+	} else {
+		util.GetTagLog().Infof("sys", "SaveFile:原文件不存在,跳过备份,path=%s", SavePath+fname)
+	}
 
 	//拷贝新文件
-	if err := os.Rename(util.GetPath(4)+fname, SavePath+fname); err != nil {
+	if err := os.Rename(tmpPath+fname, SavePath+fname); err != nil {
+		util.GetTagLog().Errorf("sys", "SaveFile:替换失败,src=%s,dst=%s,err=%v", tmpPath+fname, SavePath+fname, err)
 		//拷贝失败,则还原文件
 		if err := os.Rename(SavePath+fname+".bak", SavePath+fname); err != nil {
-			return err
+			util.GetTagLog().Errorf("sys", "SaveFile:还原失败,src=%s,err=%v", SavePath+fname+".bak", err)
 		}
 		return err
 	}
+	util.GetTagLog().Infof("sys", "SaveFile:替换成功,path=%s", SavePath+fname)
 	return nil
 }
 

+ 69 - 0
edge/ipole/logcfg.go

@@ -0,0 +1,69 @@
+package main
+
+import (
+	"github.com/sirupsen/logrus"
+
+	"lc/common/mqtt"
+	"lc/common/protocol"
+	"lc/common/util"
+)
+
+// InitLogCfg 初始化日志配置:从持久化文件加载,或使用默认值
+func InitLogCfg() {
+	tl := util.GetTagLog()
+	tl.SetPersistPath("conf/logcfg.json")
+	if err := tl.LoadFromFile(); err != nil {
+		tl.Infof("sys", "InitLogCfg:未找到持久化配置,使用默认debug级别,err=%v", err)
+		tl.SetLevel(logrus.DebugLevel)
+	} else {
+		tl.Infof("sys", "InitLogCfg:已加载持久化配置,global=%s,tags=%v",
+			tl.GetGlobalLevel().String(), tl.GetAllTagLevels())
+	}
+}
+
+// HandleTpLogCfg 处理云端下发的日志等级配置
+func HandleTpLogCfg(m mqtt.Message) {
+	tl := util.GetTagLog()
+	tl.Infof("sys", "HandleTpLogCfg:收到日志配置,topic=%s", m.Topic())
+
+	var cfg util.LogLevelConfig
+	cfg.Tags = make(map[string]string)
+	if err := json.Unmarshal([]byte(m.PayloadString()), &cfg); err != nil {
+		tl.Errorf("sys", "HandleTpLogCfg:解析失败,payload=%s,err=%v", m.PayloadString(), err)
+		return
+	}
+
+	// 设置全局级别
+	if cfg.Level != "" {
+		if lv, err := logrus.ParseLevel(cfg.Level); err == nil {
+			tl.SetLevel(lv)
+			tl.Infof("sys", "HandleTpLogCfg:全局级别已更新为%s", cfg.Level)
+		} else {
+			tl.Warnf("sys", "HandleTpLogCfg:无效的全局级别:%s", cfg.Level)
+		}
+	}
+
+	// 设置 tag 级别
+	for tag, levelStr := range cfg.Tags {
+		if lv, err := logrus.ParseLevel(levelStr); err == nil {
+			tl.SetTagLevel(tag, lv)
+			tl.Infof("sys", "HandleTpLogCfg:tag[%s]级别已更新为%s", tag, levelStr)
+		} else {
+			tl.Warnf("sys", "HandleTpLogCfg:无效的tag级别:%s=%s", tag, levelStr)
+		}
+	}
+
+	// 异步持久化,避免阻塞MQTT消息处理
+	go func() {
+		if err := tl.SaveToFile(); err != nil {
+			tl.Errorf("sys", "HandleTpLogCfg:持久化失败,err=%v", err)
+		}
+	}()
+
+	// 发送ACK
+	var ack protocol.Pack_Ack
+	if str, err := ack.EnCode(appConfig.GID, appConfig.GID, 0, nil); err == nil {
+		topic := GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_CFG_ACK)
+		GetMQTTMgr().Publish(topic, str, 0, ToCloud)
+	}
+}

+ 18 - 14
edge/ipole/main.go

@@ -12,7 +12,6 @@ import (
 
 	jsoniter "github.com/json-iterator/go"
 	"github.com/sanbornm/go-selfupdate/selfupdate"
-	"github.com/sirupsen/logrus"
 	"github.com/thinkgos/timing/v3"
 
 	"lc/common/util"
@@ -42,7 +41,7 @@ func Stat(args ...interface{}) interface{} {
 func GetNextUint64() uint64 {
 	u64, err := IDGen.NextId()
 	if err != nil {
-		logrus.Errorf("IDGen.NextId发生错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "IDGen.NextId发生错误:%s", err.Error())
 		u64 = util.MlNow().Unix()
 	}
 	return uint64(u64)
@@ -69,7 +68,7 @@ func WatchGoforever(args ...interface{}) interface{} {
 
 		isRun, err := CheckProRunning("goforever")
 		if err != nil {
-			logrus.Errorf("检查goforever命令失败:%s", err.Error())
+			util.GetTagLog().Errorf("sys", "检查goforever命令失败:%s", err.Error())
 		} else {
 			if !isRun {
 				//这里重启goforever
@@ -77,12 +76,12 @@ func WatchGoforever(args ...interface{}) interface{} {
 				//err := exec.Command("/bin/sh", "-c", "/usr/app/goforever/goforever ./goforever &").Run()
 				err := exec.Command("/bin/sh", "-c", "/usr/app/goforever ./goforever &").Run()
 				if err != nil {
-					logrus.Errorf("重启goforever出错: %s", err.Error())
+					util.GetTagLog().Errorf("sys", "重启goforever出错: %s", err.Error())
 				} else {
-					logrus.Info("重启goforever成功,当前时间:", util.MlNow().String())
+					util.GetTagLog().Info("sys", "重启goforever成功,当前时间:", util.MlNow().String())
 				}
 			} else {
-				logrus.Info("goforever进程正在运行:", util.MlNow().String())
+				util.GetTagLog().Info("sys", "goforever进程正在运行:", util.MlNow().String())
 			}
 		}
 	}
@@ -97,9 +96,9 @@ func SyncTime(args ...interface{}) interface{} {
 		}
 		err := exec.Command("ntpclient", "-h", appConfig.NtpServer, "-s").Run()
 		if err != nil {
-			logrus.Errorf("时间同步失败:%s", err.Error())
+			util.GetTagLog().Errorf("sys", "时间同步失败:%s", err.Error())
 		} else {
-			logrus.Info("时间同步成功,当前时间:", util.MlNow().String())
+			util.GetTagLog().Info("sys", "时间同步成功,当前时间:", util.MlNow().String())
 		}
 		time.Sleep(30 * time.Minute)
 	}
@@ -110,7 +109,7 @@ func SignalProcess(args ...interface{}) interface{} {
 		ch := make(chan os.Signal)
 		signal.Notify(ch, syscall.SIGHUP, syscall.SIGINT, syscall.SIGTERM, syscall.SIGKILL, syscall.SIGQUIT)
 
-		logrus.Info("检测到退出信号:", <-ch)
+		util.GetTagLog().Info("sys", "检测到退出信号:", <-ch)
 
 	}
 	return 0
@@ -131,13 +130,13 @@ func main() {
 	util.InitLogrus("release")
 
 	go func() {
-		logrus.Infoln(http.ListenAndServe(":9999", nil))
+		util.GetTagLog().Info("sys", http.ListenAndServe(":9999", nil))
 	}()
 
-	logrus.Infof("当前程序版本:%s", appname+" "+version)
+	util.GetTagLog().Infof("sys", "当前程序版本:%s", appname+" "+version)
 
 	if err := loadAppConfig(); err != nil {
-		logrus.Errorf("loadAppConfig错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "loadAppConfig错误:%s", err.Error())
 		return
 	}
 	//升级检查
@@ -146,7 +145,7 @@ func main() {
 	}
 
 	if err := loadSerialConfig(); err != nil {
-		logrus.Errorf("loadAppConfig错误:%s", err.Error())
+		util.GetTagLog().Errorf("sys", "loadAppConfig错误:%s", err.Error())
 		return
 	}
 
@@ -159,6 +158,9 @@ func main() {
 
 	InitCloudMqttSubscribeTopics()
 
+	//初始化日志配置(加载持久化级别或使用默认值)
+	InitLogCfg()
+
 	//打开串口
 	GetSerialMgr().AddSerialPorts(serialConfig.Serial)
 
@@ -166,7 +168,7 @@ func main() {
 	for k := range serialConfig.Serial {
 		devinfos, err := LoadDev(k)
 		if err != nil {
-			logrus.Warnf("加载串口[%d]的设备配置文件dev/%d.json失败: %v ,该串口无设备将被跳过", k, k, err)
+			util.GetTagLog().Warnf("sys", "加载串口[%d]的设备配置文件dev/%d.json失败: %v ,该串口无设备将被跳过", k, k, err)
 			continue
 		}
 		GetDeviceMgr().AddDevices(devinfos)
@@ -186,6 +188,8 @@ func main() {
 	gopool.Add(GetMQTTMgr().MQTTConnectMgr, 5)
 	gopool.Add(Stat, 6)
 	gopool.Add(RadarReceive, 7)
+	gopool.Add(CleanLogs, 9)
+	gopool.Add(Heartbeat, 10)
 	gopool.Run()
 	gopool.Wait()
 }

+ 40 - 37
edge/ipole/modbusrtu.go

@@ -9,7 +9,6 @@ import (
 	"runtime/debug"
 	"time"
 
-	"github.com/sirupsen/logrus"
 	"github.com/thinkgos/timing/v3"
 
 	"lc/common/mqtt"
@@ -62,12 +61,12 @@ func NewModbusRtu(info *protocol.DevInfo) Device {
 
 	iot, err := loadModel(info.TID)
 	if err != nil {
-		logrus.Errorf("ReloadModel:加载模型[tid=%d]文件发生错误:%s", info.TID, err.Error())
+		util.GetTagLog().Errorf("modbus", "ReloadModel:加载模型[tid=%d]文件发生错误:%s", info.TID, err.Error())
 	} else {
 		if iot.Protocol == ModbusRtuProtocol {
 			rtu.model = iot
 		} else {
-			logrus.Error("ModbusRtu.UpdateModel:物模型错误,非ModbusRTU协议")
+			util.GetTagLog().Error("modbus", "ModbusRtu.UpdateModel:物模型错误,非ModbusRTU协议")
 		}
 	}
 	mapRtuUploadManager.Store(info.DevCode, NewRtuUploadManager(info))
@@ -115,22 +114,22 @@ func (o *ModbusRtu) UpdateModel2(mi *ModelInfo) {
 		return
 	}
 	if mi.Flag == 0 {
-		logrus.Errorf("设备[%s]的物模型[tid=%d]模型文件被删除,下次启动即将生效。", o.devinfo.DevCode, mi.TID)
+		util.GetTagLog().Errorf("modbus", "设备[%s]的物模型[tid=%d]模型文件被删除,下次启动即将生效。", o.devinfo.DevCode, mi.TID)
 		return
 	}
-	logrus.Debugf("ModbusRtu.UpdateModel2:更新设备[%s]的物模型[%d]", o.devinfo.DevCode, mi.TID)
+	util.GetTagLog().Debugf("modbus", "ModbusRtu.UpdateModel2:更新设备[%s]的物模型[%d]", o.devinfo.DevCode, mi.TID)
 	iot, err := loadModel(mi.TID)
 	if err != nil {
-		logrus.Errorf("ModbusRtu.UpdateModel2:加载模型[%d]文件错误:%s", mi.TID, err.Error())
+		util.GetTagLog().Errorf("modbus", "ModbusRtu.UpdateModel2:加载模型[%d]文件错误:%s", mi.TID, err.Error())
 		return
 	}
 	if iot.Protocol == ModbusRtuProtocol { //合法的物模型
 		o.model = iot
 		o.clearRequest()
 		o.updateRequest()
-		logrus.Infof("ModbusRtu.UpdateModel2:更新设备[%s]的物模型[%d]成功", o.devinfo.DevCode, mi.TID)
+		util.GetTagLog().Infof("modbus", "ModbusRtu.UpdateModel2:更新设备[%s]的物模型[%d]成功", o.devinfo.DevCode, mi.TID)
 	} else {
-		logrus.Error("ModbusRtu.UpdateModel2:物模型错误,TID和文件名tid不一致或协议非ModbusRTU协议")
+		util.GetTagLog().Error("modbus", "ModbusRtu.UpdateModel2:物模型错误,TID和文件名tid不一致或协议非ModbusRTU协议")
 	}
 }
 
@@ -151,8 +150,8 @@ func (o *ModbusRtu) GetDevType() string {
 func (o *ModbusRtu) HandleData() {
 	defer func() {
 		if err := recover(); err != nil {
-			logrus.Error("ModbusRtu.HandleData:panic:", err)
-			logrus.Error("stack:", string(debug.Stack()))
+			util.GetTagLog().Error("modbus", "ModbusRtu.HandleData:panic:", err)
+			util.GetTagLog().Error("modbus", "stack:", string(debug.Stack()))
 			time.Sleep(5 * time.Second)
 			go o.HandleData()
 		}
@@ -162,7 +161,7 @@ func (o *ModbusRtu) HandleData() {
 	for {
 		select {
 		case <-o.ctx.Done():
-			logrus.Errorf("设备[%s]的HandleData退出,原因:%v", o.devinfo.DevCode, o.ctx.Err())
+			util.GetTagLog().Errorf("modbus", "设备[%s]的HandleData退出,原因:%v", o.devinfo.DevCode, o.ctx.Err())
 			return
 		case devinfo_ := <-o.chanDevInfo:
 			o.devinfo = devinfo_
@@ -184,12 +183,16 @@ func (o *ModbusRtu) clearRequest() {
 }
 
 func (o *ModbusRtu) updateRequest() {
+	if o.model == nil {
+		util.GetTagLog().Errorf("modbus", "设备[%s]的物模型为nil,无法更新采集请求,请先下发TID=%d的物模型文件", o.devinfo.DevCode, o.devinfo.TID)
+		return
+	}
 	for k, v := range o.model.Packet {
 		r := Request{CID: k, Rtuinfo: o.devinfo, FuncCode: v.Code, Address: v.Addr,
 			Quantity: v.Quantity, ScanRate: time.Duration(v.Cycle) * time.Millisecond,
 		}
 		if err := o.AddGatherJob(&r); err != nil {
-			logrus.Errorf("给串口[%d]添加采集任务[DevCode=%s,SlaveID=%d,TID=%d,CID=%d]失败:%s",
+			util.GetTagLog().Errorf("modbus", "给串口[%d]添加采集任务[DevCode=%s,SlaveID=%d,TID=%d,CID=%d]失败:%s",
 				o.devinfo.Code, o.devinfo.DevCode, o.devinfo.DevID, o.devinfo.TID, k, err.Error())
 		} else {
 			o.reqList = append(o.reqList, &r)
@@ -232,7 +235,7 @@ func (o *ModbusRtu) ProcReadDiscretes(cid uint8, address, quality uint16, valBuf
 }
 
 func (o *ModbusRtu) ProcReadHoldingRegisters(cid uint8, address, quality uint16, valBuf []byte) {
-	logrus.Debugf("收到自寄存器地址%d;开始的数据:%v", address, hex.EncodeToString(valBuf))
+	util.GetTagLog().Debugf("modbus", "收到自寄存器地址%d;开始的数据:%v", address, hex.EncodeToString(valBuf))
 	dataLen := len(valBuf)
 	if o.model.Packet[cid].Resplen != uint(dataLen) {
 		return
@@ -243,7 +246,7 @@ func (o *ModbusRtu) ProcReadHoldingRegisters(cid uint8, address, quality uint16,
 			continue
 		}
 		if int(v.Start+v.Len) > dataLen { //索引超长,忽略该项目
-			logrus.Errorf("物模型[TID=%d]配置的数据项[sid=%d]start加len超过采集响应的长度", o.model.TID, v.SID)
+			util.GetTagLog().Errorf("modbus", "物模型[TID=%d]配置的数据项[sid=%d]start加len超过采集响应的长度", o.model.TID, v.SID)
 			continue
 		}
 		var fVal float64
@@ -278,7 +281,7 @@ func (o *ModbusRtu) ProcReadHoldingRegisters(cid uint8, address, quality uint16,
 		}
 		fVal = fVal - float64(v.Base)
 
-		logrus.Debugf(v.NameZh, ":", fVal)
+		util.GetTagLog().Debugf("modbus", v.NameZh, ":", fVal)
 
 		if v.Type == 0 || v.Type == 1 { //整数
 			dataMap[v.SID] = Precision(fVal, 0, true)
@@ -295,20 +298,20 @@ func (o *ModbusRtu) ProcReadHoldingRegisters(cid uint8, address, quality uint16,
 }
 
 func (o *ModbusRtu) ProcReadInputRegisters(cid uint8, address, quality uint16, valBuf []byte) {
-	logrus.Debugf("收到自寄存器地址%d, 开始的数据:%v, 长度为:%d", address, hex.EncodeToString(valBuf), len(valBuf))
+	util.GetTagLog().Debugf("modbus", "收到自寄存器地址%d, 开始的数据:%v, 长度为:%d", address, hex.EncodeToString(valBuf), len(valBuf))
 	dataLen := len(valBuf)
 	if o.model.Packet[cid].Resplen != uint(dataLen) {
-		logrus.Errorf("ProcReadInputRegisters len no equal")
+		util.GetTagLog().Errorf("modbus", "ProcReadInputRegisters len no equal")
 		return
 	}
 	dataMap := make(map[uint16]float64)
 	for _, v := range o.model.DataUp {
 		if v.Cid != cid || v.Len == 0 {
-			logrus.Errorf("ProcReadInputRegisters Cid no equal")
+			util.GetTagLog().Errorf("modbus", "ProcReadInputRegisters Cid no equal")
 			continue
 		}
 		if int(v.Start+v.Len) > dataLen { //索引超长,忽略该项目
-			logrus.Errorf("物模型[TID=%d]配置的数据项[sid=%d]start加len超过采集响应的长度", o.model.TID, v.SID)
+			util.GetTagLog().Errorf("modbus", "物模型[TID=%d]配置的数据项[sid=%d]start加len超过采集响应的长度", o.model.TID, v.SID)
 			continue
 		}
 		var fVal float64
@@ -350,7 +353,7 @@ func (o *ModbusRtu) ProcReadInputRegisters(cid uint8, address, quality uint16, v
 		}
 	}
 
-	logrus.Debugf("ProcReadInputRegisters dataMap = %v", dataMap)
+	util.GetTagLog().Debugf("modbus", "ProcReadInputRegisters dataMap = %v", dataMap)
 	if rtumgr, ok := mapRtuUploadManager.Load(o.devinfo.DevCode); ok {
 		prtumgr := rtumgr.(*RtuUploadManager)
 		if prtumgr != nil {
@@ -367,7 +370,7 @@ func (o *ModbusRtu) ProcResult(err error, req *Request) {
 	} else {
 		//连续采集超过5次都错误,则报告设备离线
 		if req.ErrCnt == 5 {
-			logrus.Errorf("采集设备[DevCode=%s,SlaveID=%d,Tid=%d,Cid=%d]的数据发生错误:%s",
+			util.GetTagLog().Errorf("modbus", "采集设备[DevCode=%s,SlaveID=%d,Tid=%d,Cid=%d]的数据发生错误:%s",
 				req.Rtuinfo.DevCode, req.Rtuinfo.DevID, req.Rtuinfo.TID, req.CID, err.Error())
 
 			//离线状态报告
@@ -375,7 +378,7 @@ func (o *ModbusRtu) ProcResult(err error, req *Request) {
 			if str, err := obj.EnCode(req.Rtuinfo.DevCode, appConfig.GID, GetNextUint64(), err, req.Rtuinfo.TID, nil); err == nil {
 				topic := GetTopic(o.GetDevType(), req.Rtuinfo.DevCode, protocol.TP_MODBUS_DATA)
 				GetMQTTMgr().Publish(topic, str, 0, ToAll)
-				logrus.Debugf("topic:%s,payload:%s", topic, str)
+				util.GetTagLog().Debugf("modbus", "topic:%s,payload:%s", topic, str)
 			}
 		}
 	}
@@ -386,8 +389,8 @@ func (o *ModbusRtu) procRequest(req *Request) {
 	var result []byte
 	defer func() {
 		if err := recover(); err != nil {
-			logrus.Error("procRequest:panic:", err)
-			logrus.Error("stack:", string(debug.Stack()))
+			util.GetTagLog().Error("modbus", "procRequest:panic:", err)
+			util.GetTagLog().Error("modbus", "stack:", string(debug.Stack()))
 		}
 	}()
 	req.TxCnt++
@@ -395,30 +398,30 @@ func (o *ModbusRtu) procRequest(req *Request) {
 	// A bit of access read
 	case modbus.FuncCodeReadCoils:
 		result, err = o.ReadCoils(req.Rtuinfo.DevID, req.Address, req.Quantity)
-		logrus.Debugf("ReadCoils result:= %s", string(result))
+		util.GetTagLog().Debugf("modbus", "ReadCoils result:= %s", string(result))
 		if err == nil {
 			o.ProcReadCoils(req.CID, req.Address, req.Quantity, result)
 		}
 	case modbus.FuncCodeReadDiscreteInputs:
 		result, err = o.ReadDiscreteInputs(req.Rtuinfo.DevID, req.Address, req.Quantity)
-		logrus.Debugf("ReadDiscreteInputs result:= %s", string(result))
+		util.GetTagLog().Debugf("modbus", "ReadDiscreteInputs result:= %s", string(result))
 		if err == nil {
 			o.ProcReadDiscretes(req.CID, req.Address, req.Quantity, result)
 		}
 	// 16-bit access read
 	case modbus.FuncCodeReadHoldingRegisters: //03
 		result, err = o.ReadHoldingRegistersBytes(req.Rtuinfo.DevID, req.Address, req.Quantity)
-		logrus.Debugf("ReadHoldingRegistersBytes result:= %s", string(result))
+		util.GetTagLog().Debugf("modbus", "ReadHoldingRegistersBytes result:= %s", string(result))
 		if err == nil {
 			o.ProcReadHoldingRegisters(req.CID, req.Address, req.Quantity, result)
 		} else {
-			logrus.Errorf("设备上传网关数据为空!可能原因:未配置正确,请检查配置文件;设备线路或设备本身损坏;网关串口损坏.\n")
-			logrus.Errorf("error:%v\n", err)
+			util.GetTagLog().Errorf("modbus", "设备上传网关数据为空!可能原因:未配置正确,请检查配置文件;设备线路或设备本身损坏;网关串口损坏.\n")
+			util.GetTagLog().Errorf("modbus", "error:%v\n", err)
 		}
 
 	case modbus.FuncCodeReadInputRegisters:
 		result, err = o.ReadInputRegistersBytes(req.Rtuinfo.DevID, req.Address, req.Quantity)
-		logrus.Debugf("ReadInputRegistersBytes result:= %s", string(result))
+		util.GetTagLog().Debugf("modbus", "ReadInputRegistersBytes result:= %s", string(result))
 		if err == nil {
 			o.ProcReadInputRegisters(req.CID, req.Address, req.Quantity, result)
 		}
@@ -466,8 +469,8 @@ func (o *ModbusRtu) WriteData(slaveID, funcCode byte, address, quantity uint16,
 	var err error
 	defer func() {
 		if err := recover(); err != nil {
-			logrus.Error("WriteData:panic:", err)
-			logrus.Error("stack:", string(debug.Stack()))
+			util.GetTagLog().Error("modbus", "WriteData:panic:", err)
+			util.GetTagLog().Error("modbus", "stack:", string(debug.Stack()))
 		}
 	}()
 	switch funcCode {
@@ -487,7 +490,7 @@ func (o *ModbusRtu) WriteData(slaveID, funcCode byte, address, quantity uint16,
 	case modbus.FuncCodeWriteMultipleRegisters: //16
 		err = o.WriteMultipleRegistersBytes(slaveID, address, quantity, value)
 	default:
-		logrus.Errorf("不支持的功能码:%d", funcCode)
+		util.GetTagLog().Errorf("modbus", "不支持的功能码:%d", funcCode)
 		err = errors.New("不支持的功能码")
 	}
 	return err
@@ -508,17 +511,17 @@ func (o *ModbusRtu) Send(slaveID byte, request modbus.ProtocolDataUnit) (modbus.
 	if err != nil {
 		return response, err
 	}
-	logrus.Debugf("发送给设备: %s 的数据: [% x]", o.devinfo.DevCode, aduRequest)
+	util.GetTagLog().Debugf("modbus", "发送给设备: %s 的数据: [% x]", o.devinfo.DevCode, aduRequest)
 	aduResponse, err := o.SendRecvData(aduRequest)
-	logrus.Debugf("收到的数据: %s 的数据: [% x]", o.devinfo.DevCode, aduResponse)
+	util.GetTagLog().Debugf("modbus", "收到的数据: %s 的数据: [% x]", o.devinfo.DevCode, aduResponse)
 
 	if err != nil {
-		logrus.Debugf("ReadHoldingRegistersBytes SendRecvData err:= %v", err)
+		util.GetTagLog().Debugf("modbus", "ReadHoldingRegistersBytes SendRecvData err:= %v", err)
 		return response, err
 	}
 	rspSlaveID, pdu, err := modbus.DecodeRTUFrame(aduResponse)
 	if err != nil {
-		logrus.Debugf("ReadHoldingRegistersBytes DecodeRTUFrame err:= %v", err)
+		util.GetTagLog().Debugf("modbus", "ReadHoldingRegistersBytes DecodeRTUFrame err:= %v", err)
 		return response, err
 	}
 	response = modbus.ProtocolDataUnit{FuncCode: pdu[0], Data: pdu[1:]}

+ 3 - 3
edge/ipole/monitor.go

@@ -9,7 +9,7 @@ import (
 	"sync"
 	"time"
 
-	"github.com/sirupsen/logrus"
+	"lc/common/util"
 )
 
 const (
@@ -59,7 +59,7 @@ func CheckProRunning(serverName string) (bool, error) {
 		return false, err
 	}
 
-	logrus.Info("ps len", len(str), str)
+	util.GetTagLog().Info("sys", "ps len", len(str), str)
 	return len(str) > 30, nil
 }
 
@@ -83,6 +83,6 @@ func getFilelist(path string) {
 		return nil
 	})
 	if err != nil {
-		logrus.Infof("filepath.Walk() returned %v\r\n", err)
+		util.GetTagLog().Infof("sys", "filepath.Walk() returned %v\r\n", err)
 	}
 }

+ 2 - 0
edge/ipole/mqtt_init.go

@@ -27,11 +27,13 @@ func GetMQTTMgr() *mqtt.MQTTMgr {
 			CloudUser:     appConfig.Cloud.Mqtt.User,
 			CloudPassword: appConfig.Cloud.Mqtt.Password,
 			CloudTimeout:  appConfig.Cloud.Mqtt.Timeout,
+			CloudOnline:   &MqttOnline{},
 			EdgeServer:    appConfig.Edge.Mqtt.Server,
 			EdgeClientID:  appConfig.GID + "@" + appname + version,
 			EdgeUser:      appConfig.Edge.Mqtt.User,
 			EdgePassword:  appConfig.Edge.Mqtt.Password,
 			EdgeTimeout:   appConfig.Edge.Mqtt.Timeout,
+			EdgeOnline:    &MqttOnline{},
 		})
 	})
 	return _mqttMgr

+ 43 - 14
edge/ipole/mqtthandle.go

@@ -1,8 +1,10 @@
 package main
 
 import (
+	"fmt"
 	"os"
 	"strings"
+	"time"
 
 	"lc/common/mqtt"
 	"lc/common/protocol"
@@ -98,14 +100,20 @@ func HandleTpWModel(m mqtt.Message) {
 	}
 }
 func HandleTpQLog(m mqtt.Message) {
+	util.GetTagLog().Infof("sys", "HandleTpQLog:收到日志查询请求,topic=%s,payload=%s", m.Topic(), m.PayloadString())
 	var obj protocol.Pack_IDObject
 	var ret protocol.Pack_MutilFileObject
-	if err := obj.DeCode(m.PayloadString()); err == nil {
-		//读文件内容
-		ReadMutilFileContent(protocol.TP_GW_LOG, obj.Data.Id, &ret)
-		if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
-			GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_ACK), str, 0, ToCloud)
-		}
+	if err := obj.DeCode(m.PayloadString()); err != nil {
+		util.GetTagLog().Errorf("sys", "HandleTpQLog:DeCode失败,err=%v", err)
+		return
+	}
+	//读文件内容
+	ReadMutilFileContent(protocol.TP_GW_LOG, obj.Data.Id, &ret)
+	if str, err := ret.EnCode(appConfig.GID, obj.Seq); err == nil {
+		GetMQTTMgr().Publish(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_ACK), str, 0, ToCloud)
+		util.GetTagLog().Infof("sys", "HandleTpQLog:日志ACK已发布")
+	} else {
+		util.GetTagLog().Errorf("sys", "HandleTpQLog:EnCode失败,err=%v", err)
 	}
 }
 func HandleTpRLog(m mqtt.Message) {
@@ -135,20 +143,40 @@ type MqttOnline struct {
 }
 
 func (o *MqttOnline) GetOnlineMsg() (string, string) {
-	//发布上线消息
-	var obj protocol.Pack_IDObject
-	str, err := obj.EnCode(appConfig.GID, GetNextUint64(), 0)
-	if err != nil {
-		return "", ""
-	}
-	return GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_ONLINE), str
+	// 手拼 JSON,避开 json-iterator ConfigFastest 的潜在 marshal 问题
+	seq := GetNextUint64()
+	now := protocol.BJNow().Format("2006-01-02 15:04:05")
+	payload := fmt.Sprintf(`{"id":"%s","seq":%d,"gid":"%s","time":"%s","data":{"id":0}}`,
+		appConfig.GID, seq, appConfig.GID, now)
+	topic := GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_ONLINE)
+	util.GetTagLog().Infof("sys", "GetOnlineMsg: topic=%s", topic)
+	return topic, payload
 }
 
 func (o *MqttOnline) GetWillMsg() (string, string) {
-	payload, _ := (&protocol.Pack_IDObject{}).EnCode(appConfig.GID, GetNextUint64(), 0) //遗嘱消息
+	seq := GetNextUint64()
+	now := protocol.BJNow().Format("2006-01-02 15:04:05")
+	payload := fmt.Sprintf(`{"id":"%s","seq":%d,"gid":"%s","time":"%s","data":{"id":0}}`,
+		appConfig.GID, seq, appConfig.GID, now)
 	return GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_WILL), payload
 }
 
+// Heartbeat 定期重发 online 消息,防止 QoS=0 丢包导致云端永久离线
+func Heartbeat(args ...interface{}) interface{} {
+	for {
+		time.Sleep(60 * time.Second)
+		mgr := GetMQTTMgr()
+		if mgr.Cloud != nil && mgr.Cloud.IsConnected() {
+			topic, str := (&MqttOnline{}).GetOnlineMsg()
+			if topic != "" {
+				if err := mgr.Cloud.PublishString(topic, str, 0); err != nil {
+					util.GetTagLog().Errorf("sys", "Heartbeat:发布online失败,topic=%s,err=%v", topic, err)
+				}
+			}
+		}
+	}
+}
+
 // InitCloudMqttSubscribeTopics 初始化网关级别的主题订阅及路由
 func InitCloudMqttSubscribeTopics() {
 	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_APP), mqtt.AtMostOnce, HandleTpQApp, ToCloud)
@@ -161,5 +189,6 @@ func InitCloudMqttSubscribeTopics() {
 	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SET_MODEL), mqtt.AtMostOnce, HandleTpWModel, ToCloud)
 	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG), mqtt.AtMostOnce, HandleTpQLog, ToCloud)
 	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_REMOVE_LOG), mqtt.AtMostOnce, HandleTpRLog, ToCloud)
+	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_LOG_CFG), mqtt.AtMostOnce, HandleTpLogCfg, ToCloud)
 	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_GATEWAY, appConfig.GID, protocol.TP_GW_SYS), mqtt.AtMostOnce, HandleTpQSys, ToCloud)
 }

+ 7 - 7
edge/ipole/radar.go

@@ -2,8 +2,8 @@ package main
 
 //雷达
 import (
-	"github.com/sirupsen/logrus"
 	"lc/common/protocol"
+	"lc/common/util"
 	"lc/edge/ipole/modbus"
 	"math"
 	"net"
@@ -21,7 +21,7 @@ func RadarReceive(args ...interface{}) interface{} {
 		Port: PORT,
 	})
 	if err != nil {
-		logrus.Fatal("Listen failed,", err)
+		util.GetTagLog().Error("radar", "RadarReceive:UDP监听失败,", err)
 		return 0
 	}
 
@@ -29,19 +29,19 @@ func RadarReceive(args ...interface{}) interface{} {
 		var data [1024]byte
 		n, addr, err := UDPConn.ReadFromUDP(data[:]) //同步读
 		if err != nil {
-			logrus.Warningf("Read from udp server:%s failed,err:%s", addr, err)
+			util.GetTagLog().Warnf("radar", "Read from udp server:%s failed,err:%s", addr, err)
 			continue
 		}
-		logrus.Debugf("Read from udp server addr: %s", addr)
+		util.GetTagLog().Debugf("radar", "Read from udp server addr: %s", addr)
 		if b, sa := verify(data[:n-2]); !b {
-			logrus.Warningf("解析来自:%s 的数据出错,===%v ==== err:%s", addr, string(data[:n]), sa)
+			util.GetTagLog().Warnf("radar", "解析来自:%s 的数据出错,===%v ==== err:%s", addr, string(data[:n]), sa)
 			continue
 		} else {
 			l := len(sa)
 			for _, s := range sa[:l-1] {
 				d := strings.Split(s, ",")
 				if len(d) != 5 {
-					logrus.Warningf("解析的数据出错, 长度不为5")
+					util.GetTagLog().Warnf("radar", "解析的数据出错, 长度不为5")
 					continue
 				}
 				go func() {
@@ -54,7 +54,7 @@ func RadarReceive(args ...interface{}) interface{} {
 					data1[3] = math.Abs(data1[3])
 					var obj protocol.Pack_UploadData
 					if str, err := obj.EnCode(addr.String(), appConfig.GID, GetNextUint64(), nil, 0, data1); err == nil {
-						logrus.Debugf("Publish TP_RADAR_DATA")
+						util.GetTagLog().Debugf("radar", "Publish TP_RADAR_DATA")
 						topic := GetTopic(protocol.DT_Radar, appConfig.GID, protocol.TP_RADAR_DATA)
 						GetMQTTMgr().Publish(topic, str, 0, ToCloud) //上传消息 雷达
 					}

+ 4 - 4
edge/ipole/serialmgr.go

@@ -5,9 +5,9 @@ import (
 	"time"
 
 	"github.com/goburrow/serial"
-	"github.com/sirupsen/logrus"
 
 	"lc/common/protocol"
+	"lc/common/util"
 )
 
 var _once sync.Once
@@ -39,10 +39,10 @@ func openSerialPort(sp *protocol.SerialPort) (*serialPort, error) {
 	serial.SetAutoReconnect(SerialDefaultAutoReconnect)
 	err = serial.connect()
 	if err != nil {
-		logrus.Errorf("打开串口[code=%d,address=%s]失败:%s", sp.Code, sp.Address, err.Error())
+		util.GetTagLog().Errorf("sys", "打开串口[code=%d,address=%s]失败:%s", sp.Code, sp.Address, err.Error())
 		return nil, err
 	}
-	logrus.Infof("打开串口[code=%d,address=%s]成功", sp.Code, sp.Address)
+	util.GetTagLog().Infof("sys", "打开串口[code=%d,address=%s]成功", sp.Code, sp.Address)
 	return &serial, nil
 }
 
@@ -78,7 +78,7 @@ func (o *SerialMgr) UpdateSerialPort(sp *protocol.SerialPort) error {
 	if s, ok := o.mapSerial[sp.Code]; ok { //存在则关闭,并重新打开
 		err := s.Close() //关闭
 		if err != nil {
-			logrus.Errorf("关闭串口[code=%d,address=%s]失败:%s", sp.Code, sp.Address, err.Error())
+			util.GetTagLog().Errorf("sys", "关闭串口[code=%d,address=%s]失败:%s", sp.Code, sp.Address, err.Error())
 		}
 		delete(o.mapSerial, sp.Code)
 	}

+ 1 - 2
edge/ipole/uploaddata.go

@@ -2,7 +2,6 @@ package main
 
 import (
 	"fmt"
-	"github.com/sirupsen/logrus"
 	"lc/common/models"
 	"math"
 	"sort"
@@ -46,7 +45,7 @@ func NewRtuUploadManager(rtuinfo *protocol.DevInfo) *RtuUploadManager {
 	if rtuinfo.DevType == 4 && (rtuinfo.TID == 3 || rtuinfo.TID == 5) {
 		o.LiguidDataMgr = &LiguidDataMgr{}
 	}
-	logrus.Infof("NewRtuUploadManager rtuinfo = %v \n", rtuinfo)
+	util.GetTagLog().Infof("sys", "NewRtuUploadManager rtuinfo = %v \n", rtuinfo)
 	o.Timer.AddJobFunc(o.SendData, time.Duration(o.Rtuinfo.SendCloud)*time.Millisecond)
 	return &o
 }

+ 35 - 37
edge/ipole/ym485.go

@@ -7,7 +7,6 @@ import (
 	"time"
 
 	"github.com/go-redis/redis/v7"
-	"github.com/sirupsen/logrus"
 
 	"lc/common/mqtt"
 	"lc/common/protocol"
@@ -51,12 +50,12 @@ func NewYmLampController(info *protocol.DevInfo) Device {
 
 	iot, err := loadModel(info.TID)
 	if err != nil {
-		logrus.Errorf("NewYmLampController:加载模型[tid=%d]文件发生错误:%s", info.TID, err.Error())
+		util.GetTagLog().Errorf("ym485", "NewYmLampController:加载模型[tid=%d]文件发生错误:%s", info.TID, err.Error())
 	} else {
 		if iot.Protocol == YmProtocol {
 			dev.model = iot
 		} else {
-			logrus.Error("NewYmLampController:物模型错误,非YmProtocol协议")
+			util.GetTagLog().Error("ym485", "NewYmLampController:物模型错误,非YmProtocol协议")
 		}
 	}
 	dev.SetTopicHandle()
@@ -73,7 +72,7 @@ func (o *YmLampController) MQTTSubscribe() {
 	GetMQTTMgr().Subscribe(GetTopic(protocol.DT_LAMPCONTROLLER, o.devinfo.DevCode, protocol.TP_YM_SET_ONOFFTIME), mqtt.ExactlyOnce, o.HandleCache, ToCloud) //设置时间
 }
 func (o *YmLampController) HandleCache(m mqtt.Message) {
-	logrus.Infof("Topic:%v,message:%v\n", m.Topic(), m.PayloadString())
+	util.GetTagLog().Infof("ym485", "Topic:%v,message:%v\n", m.Topic(), m.PayloadString())
 	o.downQueue.Put(m)
 }
 func (o *YmLampController) Start() {
@@ -132,20 +131,20 @@ func (o *YmLampController) UpdateModel2(mi *ModelInfo) {
 		return
 	}
 	if mi.Flag == 0 {
-		logrus.Errorf("设备[%s]的物模型[tid=%d]模型文件被删除,下次启动即将生效。", o.devinfo.DevCode, mi.TID)
+		util.GetTagLog().Errorf("ym485", "设备[%s]的物模型[tid=%d]模型文件被删除,下次启动即将生效。", o.devinfo.DevCode, mi.TID)
 		return
 	}
-	logrus.Debugf("YmLampController.UpdateModel2:更新设备[%s]的物模型[%d]", o.devinfo.DevCode, mi.TID)
+	util.GetTagLog().Debugf("ym485", "YmLampController.UpdateModel2:更新设备[%s]的物模型[%d]", o.devinfo.DevCode, mi.TID)
 	iot, err := loadModel(mi.TID)
 	if err != nil {
-		logrus.Errorf("YmLampController.UpdateModel2:加载模型[%d]文件错误:%s", mi.TID, err.Error())
+		util.GetTagLog().Errorf("ym485", "YmLampController.UpdateModel2:加载模型[%d]文件错误:%s", mi.TID, err.Error())
 		return
 	}
 	if iot.Protocol == YmProtocol { //合法的物模型
 		o.model = iot
-		logrus.Infof("YmLampController.UpdateModel2:更新设备[%s]的物模型[%d]成功", o.devinfo.DevCode, mi.TID)
+		util.GetTagLog().Infof("ym485", "YmLampController.UpdateModel2:更新设备[%s]的物模型[%d]成功", o.devinfo.DevCode, mi.TID)
 	} else {
-		logrus.Error("YmLampController.UpdateModel2:物模型错误,TID和文件名tid不一致或协议非ModbusRTU协议")
+		util.GetTagLog().Error("ym485", "YmLampController.UpdateModel2:物模型错误,TID和文件名tid不一致或协议非ModbusRTU协议")
 	}
 }
 
@@ -160,19 +159,19 @@ func (o *YmLampController) GetDevType() string {
 func (o *YmLampController) HandleData() {
 	defer func() {
 		if err := recover(); err != nil {
-			logrus.Error("YmLampController.HandleData:panic:", err)
-			logrus.Error("stack:", string(debug.Stack()))
+			util.GetTagLog().Error("ym485", "YmLampController.HandleData:panic:", err)
+			util.GetTagLog().Error("ym485", "stack:", string(debug.Stack()))
 			time.Sleep(5 * time.Second)
 			go o.HandleData()
 		}
 	}()
 	o.QueryLampControllerAddr() //线上挂了两个灯控,则不能
-	logrus.Infof("lamp state:%d", oldstate)
+	util.GetTagLog().Infof("ym485", "lamp state:%d", oldstate)
 	nextFillTime := time.Time{}
 	for {
 		select {
 		case <-o.ctx.Done():
-			logrus.Errorf("设备[%s]的HandleData退出,原因:%v", o.devinfo.DevCode, o.ctx.Err())
+			util.GetTagLog().Errorf("ym485", "设备[%s]的HandleData退出,原因:%v", o.devinfo.DevCode, o.ctx.Err())
 			return
 		case devinfo_ := <-o.chanDevInfo:
 			o.devinfo = devinfo_
@@ -185,7 +184,7 @@ func (o *YmLampController) HandleData() {
 					if fn, ok := o.mapTopicHandle[mm.Topic()]; ok {
 						fn(mm)
 					} else {
-						logrus.Errorf("YmLampController.Handle:不支持的主题:%s", mm.Topic())
+						util.GetTagLog().Errorf("ym485", "YmLampController.Handle:不支持的主题:%s", mm.Topic())
 					}
 				}
 			} else {
@@ -213,11 +212,11 @@ func (o *YmLampController) SendRecvData(aduRequest []byte, retry int) (aduRespon
 	// time.Sleep(time.Millisecond * 10)
 	serial := GetSerialMgr().GetSerialPort(o.devinfo.Code)
 	if serial == nil {
-		logrus.Errorf("YM串口未找到: Code=%d", o.devinfo.Code) // ← 加这行
+		util.GetTagLog().Errorf("ym485", "YM串口未找到: Code=%d", o.devinfo.Code) // ← 加这行
 		return nil, ErrClosedConnection
 	}
 	if !serial.IsConnected() {
-		logrus.Errorf("YM串口未连接: Code=%d", o.devinfo.Code) // ← 加这行
+		util.GetTagLog().Errorf("ym485", "YM串口未连接: Code=%d", o.devinfo.Code) // ← 加这行
 	}
 	if retry <= 0 {
 		retry = 1
@@ -261,10 +260,10 @@ func (o *YmLampController) QueryLampControllerAddr() {
 	}
 	var ret ym485.QueryAddrACK
 	if err = ret.DeCode(recvbuf); err == nil {
-		logrus.Infof("灯控地址:%s", ret.Addr)
+		util.GetTagLog().Infof("ym485", "灯控地址:%s", ret.Addr)
 	}
 	if ret.Addr != o.devinfo.DevCode {
-		logrus.Errorf("请将DevCode配置为灯控地址!")
+		util.GetTagLog().Errorf("ym485", "请将DevCode配置为灯控地址!")
 	}
 }
 
@@ -281,7 +280,7 @@ func (o *YmLampController) QuerySoftVer() {
 	}
 	var ret ym485.QuerySoftVerAck
 	if err = ret.DeCode(recvbuf); err == nil {
-		logrus.Infof("灯控型号:%s,软件版本:%s", ret.Model, ret.Version)
+		util.GetTagLog().Infof("ym485", "灯控型号:%s,软件版本:%s", ret.Model, ret.Version)
 	}
 }
 
@@ -298,7 +297,7 @@ func (o *YmLampController) QueryBasicsettings() {
 	}
 	var ret ym485.PackingBasic
 	if err = ret.DeCode(recvbuf); err == nil {
-		logrus.Infof("灯控使能状态:%d,是否上电开灯:%d", ret.Enabled, ret.OnOff)
+		util.GetTagLog().Infof("ym485", "灯控使能状态:%d,是否上电开灯:%d", ret.Enabled, ret.OnOff)
 	}
 }
 
@@ -312,7 +311,7 @@ func (o *YmLampController) QueryDeviceState() {
 	}
 	recvbuf, err := o.SendRecvData(buf.Bytes(), 1) //根据地址得到数据
 	if err != nil {
-		logrus.Errorf("单灯状态数据查询错误,DevCode=%s,err=%v", o.devinfo.DevCode, err)
+		util.GetTagLog().Errorf("ym485", "单灯状态数据查询错误,DevCode=%s,err=%v", o.devinfo.DevCode, err)
 		return
 	}
 
@@ -334,7 +333,7 @@ func (o *YmLampController) QueryDeviceState() {
 		//状态改变记录日志
 		if oldstate != o.State[0].State {
 			oldstate = o.State[0].State
-			logrus.Infof("lamp state:%d", oldstate)
+			util.GetTagLog().Infof("ym485", "lamp state:%d", oldstate)
 		}
 	}
 	mapData := make(map[string]*protocol.CHZB_LampData)
@@ -360,7 +359,6 @@ func (o *YmLampController) QueryDeviceState() {
 	var ret1 protocol.Pack_CHZB_UploadData
 	if str, err := ret1.EnCode(o.devinfo.DevCode, appConfig.GID, GetNextUint64(), o.devinfo.TID, mapData); err == nil {
 		topic := GetTopic(protocol.DT_LAMPCONTROLLER, o.devinfo.DevCode, protocol.TP_YM_DATA)
-		logrus.Infof("QueryDeviceState:发布灯控数据,topic=%s,payload=%s", topic, str)
 		GetMQTTMgr().Publish(topic, str, mqtt.AtLeastOnce, ToCloud)
 	}
 }
@@ -378,7 +376,7 @@ func (o *YmLampController) TurnOnOff(flag uint8) error {
 	}
 	recvbuf, err := o.SendRecvData(buf.Bytes(), 3)
 	if err != nil {
-		logrus.Errorf("TurnOnOff:串口通信失败,DevCode=%s,flag=%d,err=%v", o.devinfo.DevCode, flag, err)
+		util.GetTagLog().Errorf("ym485", "TurnOnOff:串口通信失败,DevCode=%s,flag=%d,err=%v", o.devinfo.DevCode, flag, err)
 		return err
 	}
 	var ret ym485.PackingResult
@@ -388,7 +386,7 @@ func (o *YmLampController) TurnOnOff(flag uint8) error {
 			return nil
 		}
 	}
-	logrus.Errorf("TurnOnOff:设备返回失败,DevCode=%s,flag=%d,Result=0x%02X", o.devinfo.DevCode, flag, ret.Result)
+	util.GetTagLog().Errorf("ym485", "TurnOnOff:设备返回失败,DevCode=%s,flag=%d,Result=0x%02X", o.devinfo.DevCode, flag, ret.Result)
 	return errors.New(protocol.FAILED_STR)
 }
 
@@ -427,9 +425,9 @@ func (o *YmLampController) SetBasicsettings(enabled, on uint8) {
 	var ret ym485.PackingResult
 	if err = ret.DeCode(recvbuf); err == nil {
 		if ret.Result == 89 { //字母'Y'
-			logrus.Infoln("设置灯控基本信息执行成功")
+			util.GetTagLog().Info("ym485", "设置灯控基本信息执行成功")
 		} else {
-			logrus.Errorln("设置灯控基本信息执行失败")
+			util.GetTagLog().Error("ym485", "设置灯控基本信息执行失败")
 		}
 	}
 }
@@ -440,7 +438,7 @@ func (o *YmLampController) Switch(Switch, Brightness uint8) error {
 	} else {
 		err := o.SetBrightness(Brightness)
 		if err != nil {
-			logrus.Errorf("调节亮度失败:%s", err.Error())
+			util.GetTagLog().Errorf("ym485", "调节亮度失败:%s", err.Error())
 		}
 		return o.TurnOnOff(Switch)
 	}
@@ -500,7 +498,7 @@ func (o *YmLampController) ConfirmState(t time.Time) {
 func (o *YmLampController) HandleTpYmSetSwitch(m mqtt.Message) {
 	var obj protocol.Pack_CHZB_Switch
 	if err := obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("ym485", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
@@ -512,7 +510,7 @@ func (o *YmLampController) HandleTpYmSetSwitch(m mqtt.Message) {
 	}
 	err := o.Switch(obj.Data.Switch, obj.Data.Brightness)
 	if err != nil {
-		logrus.Errorf("HandleTpYmSetSwitch:开关灯失败,DevCode=%s,Switch=%d,err=%v", o.devinfo.DevCode, obj.Data.Switch, err)
+		util.GetTagLog().Errorf("ym485", "HandleTpYmSetSwitch:开关灯失败,DevCode=%s,Switch=%d,err=%v", o.devinfo.DevCode, obj.Data.Switch, err)
 	} else {
 		o.tswitch = util.MlNow()
 		mapRedisTempLampsOOT := make(map[string]interface{}) //临时开关灯记录,用于排除异常亮灯正常亮灯的情况
@@ -525,7 +523,7 @@ func (o *YmLampController) HandleTpYmSetSwitch(m mqtt.Message) {
 		o.mapTempLampsOOT = &ltr                         //内存
 		mapRedisTempLampsOOT[o.devinfo.DevCode] = ltrstr //redis
 		if err := redisEdgeData.HSet(LampSwitchPrefix+o.devinfo.DevCode, mapRedisTempLampsOOT).Err(); err != nil {
-			logrus.Errorf("手动开关灯时间设置[内容:%v]缓存到redis失败:%s", mapRedisTempLampsOOT, err.Error())
+			util.GetTagLog().Errorf("ym485", "手动开关灯时间设置[内容:%v]缓存到redis失败:%s", mapRedisTempLampsOOT, err.Error())
 		}
 	}
 	var ret protocol.Pack_Ack
@@ -537,14 +535,14 @@ func (o *YmLampController) HandleTpYmSetSwitch(m mqtt.Message) {
 func (o *YmLampController) HandleTpYmSetOnofftime(m mqtt.Message) {
 	var obj protocol.Pack_SetOnOffTime
 	if err := obj.DeCode(m.PayloadString()); err != nil {
-		logrus.Errorf("协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
+		util.GetTagLog().Errorf("ym485", "协议解析错误:%s,协议主题:%s,协议内容:%s", err.Error(), m.Topic(), m.PayloadString())
 		return
 	}
 	if obj.Id != o.devinfo.DevCode {
 		return
 	}
 	if len(obj.Data.OnOffTime) == 0 {
-		logrus.Errorf("Handle_TP_YM_SET_ONOFFTIME:错误,灯控编号[%v],时间段个数:%v", obj.Id, obj.Data.OnOffTime)
+		util.GetTagLog().Errorf("ym485", "Handle_TP_YM_SET_ONOFFTIME:错误,灯控编号[%v],时间段个数:%v", obj.Id, obj.Data.OnOffTime)
 		return
 	}
 	mapRedisOOT := make(map[string]interface{})
@@ -555,7 +553,7 @@ func (o *YmLampController) HandleTpYmSetOnofftime(m mqtt.Message) {
 	//持久缓存到redis,以便于重启后读取进内存中
 	err := redisEdgeData.HSet(LampOotPrefix+o.devinfo.DevCode, mapRedisOOT).Err()
 	if err != nil {
-		logrus.Errorf("灯控时间设置[内容:%v]缓存到redis失败:%s", mapRedisOOT, err.Error())
+		util.GetTagLog().Errorf("ym485", "灯控时间设置[内容:%v]缓存到redis失败:%s", mapRedisOOT, err.Error())
 	}
 	var ret protocol.Pack_Ack
 	if str, err := ret.EnCode(o.devinfo.DevCode, appConfig.GID, obj.Seq, err); err == nil {
@@ -569,7 +567,7 @@ func (o *YmLampController) ReloadOOTFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("YmLampController.ReloadOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("ym485", "YmLampController.ReloadOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	for k, v := range mapdata {
@@ -590,7 +588,7 @@ func (o *YmLampController) ReloadSwitchOOTFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("YmLampController.ReloadSwitchOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("ym485", "YmLampController.ReloadSwitchOOTFromRedis设备[%s]从redis加载时间策略失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	for k, v := range mapdata {
@@ -611,7 +609,7 @@ func (o *YmLampController) ReloadLampAlarmFromRedis() error {
 		if err == redis.Nil {
 			return nil
 		}
-		logrus.Errorf("YmLampController.ReloadLampAlarmFromRedis设备[%s]从redis加载广播恢复截止时间失败:%s", o.devinfo.DevCode, err.Error())
+		util.GetTagLog().Errorf("ym485", "YmLampController.ReloadLampAlarmFromRedis设备[%s]从redis加载广播恢复截止时间失败:%s", o.devinfo.DevCode, err.Error())
 		return err
 	}
 	for k, v := range mapAlarm {