| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015 |
- package controllers
- import (
- "crypto/md5"
- "encoding/base64"
- "fmt"
- "io"
- "os"
- "path/filepath"
- "strconv"
- "strings"
- "time"
- "github.com/astaxie/beego"
- "lc/common/models"
- "lc/common/mqtt"
- "lc/common/protocol"
- )
- type GatewayController struct {
- BaseController
- }
- type ReqGateway struct {
- Code string `json:"code"` //设备编号,禁止修改
- Tenant string `json:"tenant"` //租户ID
- Name string `json:"name"` //设备名称
- Brand int `json:"brand"` //品牌
- Model int `json:"model"` //型号
- State int `json:"state"` //1启用,0禁用
- }
- type Gateway struct {
- Code string `json:"code"` //设备编号,禁止修改
- Name string `json:"name"` //设备名称
- Brand int `json:"brand"` //品牌
- Model int `json:"model"` //型号
- State int `json:"state"` //1启用,0禁用
- }
- type ReqImportGateway struct {
- Tenant string `json:"tenant"` //租户ID
- List []Gateway `json:"list"` //网关列表
- }
- type RespImport struct {
- Code string `json:"code"`
- Error string `json:"error"`
- }
- // CreateGateway @Title 创建网关设备
- // @Description 创建网关设备
- // @Param body controllers.ReqGateway true "数据"
- // @Success 0 {int} BaseResponse.Code "成功"
- // @Failure 1 {int} BaseResponse.Code "失败"
- // @router /v1/create [post]
- func (o *GatewayController) CreateGateway() {
- var obj ReqGateway
- if err := json.Unmarshal(o.Ctx.Input.RequestBody, &obj); err != nil {
- beego.Debug(string(o.Ctx.Input.RequestBody))
- o.Response(Failure, fmt.Sprintf("数据包解析错误:%s", err.Error()), nil)
- return
- }
- oo := models.Gateway{
- ID: obj.Code,
- Name: obj.Name,
- Tenant: obj.Tenant,
- Brand: obj.Brand,
- Model: obj.Model,
- State: obj.State,
- }
- if err := oo.SaveFromWeb(); err != nil {
- o.Response(Failure, fmt.Sprintf("数据插入失败:%s", err.Error()), nil)
- return
- }
- o.Response(Success, "成功", oo.ID)
- }
- // UpdateGateway @Title 更新网关设备
- // @Description 更新网关设备
- // @Param body controllers.ReqGateway true "数据"
- // @Success 0 {int} BaseResponse.Code "成功"
- // @Failure 1 {int} BaseResponse.Code "失败"
- // @router /v1/update [post]
- func (o *GatewayController) UpdateGateway() {
- var obj ReqGateway
- if err := json.Unmarshal(o.Ctx.Input.RequestBody, &obj); err != nil {
- beego.Debug(string(o.Ctx.Input.RequestBody))
- o.Response(Failure, fmt.Sprintf("数据包解析错误:%s", err.Error()), nil)
- return
- }
- oo := models.Gateway{
- ID: obj.Code,
- Name: obj.Name,
- Tenant: obj.Tenant,
- Brand: obj.Brand,
- Model: obj.Model,
- State: obj.State,
- }
- if err := oo.SaveFromWeb(); err != nil {
- o.Response(Failure, fmt.Sprintf("数据更新失败:%s", err.Error()), nil)
- return
- }
- o.Response(Success, "成功", oo.ID)
- }
- // DeleteGateway @Title 删除网关设备
- // @Description 删除网关设备
- // @Param code query string true "设备ID"
- // @Success 0 {int} BaseResponse.Code "成功"
- // @Failure 1 {int} BaseResponse.Code "失败"
- // @router /v1/delete [post]
- func (o *GatewayController) DeleteGateway() {
- code := strings.Trim(o.GetString("code"), " ")
- if len(code) == 0 {
- o.Response(Failure, "code为空", nil)
- return
- }
- c := models.Gateway{
- ID: code,
- }
- if err := c.Delete(); err != nil {
- beego.Error(fmt.Sprintf("删除失败:%s", err.Error()))
- o.Response(Failure, fmt.Sprintf("数据删除失败:%s", err.Error()), nil)
- return
- }
- o.Response(Success, "成功", c.ID)
- }
- // ImportGateway @Title 批量导入网关设备
- // @Description 批量导入网关设备
- // @Param body controllers.ReqImportGateway true "数据"
- // @Success 0 {int} BaseResponse.Code "成功"
- // @Failure 1 {int} BaseResponse.Code "失败"
- // @router /v1/import [post]
- func (o *GatewayController) ImportGateway() {
- var obj ReqImportGateway
- if err := json.Unmarshal(o.Ctx.Input.RequestBody, &obj); err != nil {
- beego.Debug(string(o.Ctx.Input.RequestBody))
- o.Response(Failure, fmt.Sprintf("数据包解析错误:%s", err.Error()), nil)
- return
- }
- var resp []RespImport
- for _, v := range obj.List {
- var aresp RespImport
- aresp.Code = v.Code
- oo := models.Gateway{
- ID: v.Code,
- Name: v.Name,
- Tenant: obj.Tenant,
- Brand: v.Brand,
- Model: v.Model,
- State: v.State,
- }
- err := oo.SaveFromWeb()
- if err != nil {
- aresp.Error = err.Error()
- beego.Error(fmt.Sprintf("zigbee集控器数据导入失败,code=%s,失败原因:%s", v.Code, err.Error()))
- }
- resp = append(resp, aresp)
- }
- 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)
- }
- // DeployUpload @Title 上传部署文件
- // @Description 上传网关部署文件
- // @router /v1/deploy/upload [post]
- func (o *GatewayController) DeployUpload() {
- gid := strings.Trim(o.GetString("gid"), " ")
- tenant := strings.Trim(o.GetString("tenant"), " ")
- file, header, err := o.Ctx.Request.FormFile("file")
- if err != nil {
- o.Response(Failure, "读取上传文件失败: "+err.Error(), nil)
- return
- }
- defer file.Close()
- data, err := io.ReadAll(file)
- if err != nil {
- o.Response(Failure, "读取文件内容失败: "+err.Error(), nil)
- return
- }
- md5Hash := fmt.Sprintf("%x", md5.Sum(data))
- // 保存到临时目录
- tmpDir := filepath.Join(os.TempDir(), "ipole_deploy")
- os.MkdirAll(tmpDir, os.ModePerm)
- tmpFile := filepath.Join(tmpDir, md5Hash)
- if err := os.WriteFile(tmpFile, data, os.ModePerm); err != nil {
- o.Response(Failure, "保存文件失败: "+err.Error(), nil)
- return
- }
- // 创建部署记录(原始SQL绕过GORM零值问题)
- now := time.Now()
- if err := models.G_db.Exec("INSERT INTO t_gateway_deploy (gid, tenant, from_version, to_version, md5, status, phase1_result, phase2_result, error_msg, create_time, update_time) VALUES (?, ?, ?, ?, ?, 0, ?, ?, ?, ?, ?)", gid, tenant, "", "", md5Hash, "", "", "", now, now).Error; err != nil {
- o.Response(Failure, "创建部署记录失败: "+err.Error(), nil)
- return
- }
- o.Response(Success, "文件上传成功", map[string]interface{}{
- "md5": md5Hash,
- "size": len(data),
- "path": tmpFile,
- "fileName": header.Filename,
- })
- }
- // DeployPush @Title 下发部署指令
- // @Description 下发部署指令到网关
- // @router /v1/deploy/push [post]
- func (o *GatewayController) DeployPush() {
- gid := strings.Trim(o.GetString("gid"), " ")
- tenant := strings.Trim(o.GetString("tenant"), " ")
- version := strings.Trim(o.GetString("version"), " ")
- fileName := strings.Trim(o.GetString("fileName"), " ")
- filePath := strings.Trim(o.GetString("filePath"), " ")
- md5Hash := strings.Trim(o.GetString("md5"), " ")
- if gid == "" || tenant == "" || filePath == "" {
- o.Response(Failure, "参数不完整", nil)
- return
- }
- if fileName == "" {
- fileName = "ipole"
- }
- if version == "" {
- version = fmt.Sprintf("v%s", time.Now().Format("20060102-150405"))
- }
- data, err := os.ReadFile(filePath)
- if err != nil {
- o.Response(Failure, "读取部署文件失败: "+err.Error(), nil)
- return
- }
- const chunkSize = 64 * 1024 // 64KB
- totalChunks := (len(data) + chunkSize - 1) / chunkSize
- // 1. 发送部署指令
- var cmd protocol.Pack_DeployCmd
- seq := GetNextUint64()
- cmdStr, err := cmd.EnCodeCmd(gid, gid, seq, version, fileName, md5Hash, totalChunks, chunkSize)
- if err != nil {
- o.Response(Failure, "编码部署指令失败: "+err.Error(), nil)
- return
- }
- topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_DEPLOY)
- GetMqttHandler().PublishString(topic, cmdStr, mqtt.AtLeastOnce)
- // 2. 逐片发送
- for i := 0; i < totalChunks; i++ {
- start := i * chunkSize
- end := start + chunkSize
- if end > len(data) {
- end = len(data)
- }
- chunkData := base64.StdEncoding.EncodeToString(data[start:end])
- seq := GetNextUint64()
- chunkStr, err := cmd.EnCodeChunk(gid, gid, seq, version, i, chunkData)
- if err != nil {
- o.Response(Failure, fmt.Sprintf("编码分片%d失败: %s", i, err.Error()), nil)
- return
- }
- GetMqttHandler().PublishString(topic, chunkStr, mqtt.AtLeastOnce)
- if i%20 == 0 && i > 0 {
- time.Sleep(50 * time.Millisecond)
- }
- }
- models.G_db.Model(&models.GatewayDeploy{}).
- Where("gid = ? AND status = 0", gid).
- Updates(map[string]interface{}{
- "to_version": version,
- "from_version": "",
- "md5": md5Hash,
- "update_time": time.Now(),
- })
- o.Response(Success, "部署指令已下发", map[string]interface{}{
- "totalChunks": totalChunks,
- "chunkSize": chunkSize,
- "fileSize": len(data),
- })
- }
- // DeployStatus @Title 查询部署状态
- // @Description 查询网关部署状态
- // @router /v1/deploy/status [get]
- func (o *GatewayController) DeployStatus() {
- gid := strings.Trim(o.GetString("gid"), " ")
- var deploy models.GatewayDeploy
- err := models.G_db.Where("gid = ?", gid).Order("create_time DESC").First(&deploy).Error
- if err != nil {
- // record not found 返回空状态,方便前端轮询等待
- o.Response(Success, "暂无部署记录", map[string]interface{}{
- "status": -1,
- })
- return
- }
- statusText := map[uint8]string{0: "进行中", 1: "成功", 2: "失败", 3: "已回滚"}
- o.Response(Success, "查询成功", map[string]interface{}{
- "id": deploy.ID,
- "fromVersion": deploy.FromVersion,
- "toVersion": deploy.ToVersion,
- "status": deploy.Status,
- "statusText": statusText[deploy.Status],
- "errorMsg": deploy.ErrorMsg,
- "phase1Result": deploy.Phase1Result,
- "phase2Result": deploy.Phase2Result,
- "createTime": deploy.CreateTime.Format("2006-01-02 15:04:05"),
- })
- }
- // DeployBatchStatus @Title 批量查询部署状态
- // @Description 一次查询多个网关的最近一次部署状态
- // @router /v1/deploy/batch-status [get]
- func (o *GatewayController) DeployBatchStatus() {
- tenant := strings.Trim(o.GetString("tenant"), " ")
- gidsStr := strings.Trim(o.GetString("gids"), " ")
- if tenant == "" {
- o.Response(Failure, "tenant不能为空", nil)
- return
- }
- // 两步查询:先取每个 gid 最新 id,再取完整记录
- query := models.G_db.Table("t_gateway_deploy").
- Where("tenant = ?", tenant)
- if gidsStr != "" {
- gids := strings.Split(gidsStr, ",")
- for i := range gids {
- gids[i] = strings.Trim(gids[i], " ")
- }
- query = query.Where("gid IN (?)", gids)
- }
- var maxIds []int64
- dbResult := query.Select("MAX(id)").Group("gid").Pluck("MAX(id)", &maxIds)
- if dbResult.Error != nil {
- o.Response(Failure, "查询部署记录失败: "+dbResult.Error.Error(), nil)
- return
- }
- var deploys []models.GatewayDeploy
- if len(maxIds) > 0 {
- models.G_db.Where("id IN (?)", maxIds).Find(&deploys)
- }
- result := make(map[string]interface{})
- for _, deploy := range deploys {
- result[deploy.GID] = map[string]interface{}{
- "id": deploy.ID,
- "status": deploy.Status,
- "toVersion": deploy.ToVersion,
- "phase1Result": deploy.Phase1Result,
- "phase2Result": deploy.Phase2Result,
- "errorMsg": deploy.ErrorMsg,
- "createTime": deploy.CreateTime.Format("2006-01-02 15:04:05"),
- }
- }
- // 补充 gids 列表中有但无部署记录的网关(value = null)
- if gidsStr != "" {
- for _, gid := range strings.Split(gidsStr, ",") {
- gid = strings.Trim(gid, " ")
- if gid == "" {
- continue
- }
- if _, ok := result[gid]; !ok {
- result[gid] = nil
- }
- }
- }
- o.Response(Success, "查询成功", result)
- }
|