cgateway.go 33 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015
  1. package controllers
  2. import (
  3. "crypto/md5"
  4. "encoding/base64"
  5. "fmt"
  6. "io"
  7. "os"
  8. "path/filepath"
  9. "strconv"
  10. "strings"
  11. "time"
  12. "github.com/astaxie/beego"
  13. "lc/common/models"
  14. "lc/common/mqtt"
  15. "lc/common/protocol"
  16. )
  17. type GatewayController struct {
  18. BaseController
  19. }
  20. type ReqGateway struct {
  21. Code string `json:"code"` //设备编号,禁止修改
  22. Tenant string `json:"tenant"` //租户ID
  23. Name string `json:"name"` //设备名称
  24. Brand int `json:"brand"` //品牌
  25. Model int `json:"model"` //型号
  26. State int `json:"state"` //1启用,0禁用
  27. }
  28. type Gateway struct {
  29. Code string `json:"code"` //设备编号,禁止修改
  30. Name string `json:"name"` //设备名称
  31. Brand int `json:"brand"` //品牌
  32. Model int `json:"model"` //型号
  33. State int `json:"state"` //1启用,0禁用
  34. }
  35. type ReqImportGateway struct {
  36. Tenant string `json:"tenant"` //租户ID
  37. List []Gateway `json:"list"` //网关列表
  38. }
  39. type RespImport struct {
  40. Code string `json:"code"`
  41. Error string `json:"error"`
  42. }
  43. // CreateGateway @Title 创建网关设备
  44. // @Description 创建网关设备
  45. // @Param body controllers.ReqGateway true "数据"
  46. // @Success 0 {int} BaseResponse.Code "成功"
  47. // @Failure 1 {int} BaseResponse.Code "失败"
  48. // @router /v1/create [post]
  49. func (o *GatewayController) CreateGateway() {
  50. var obj ReqGateway
  51. if err := json.Unmarshal(o.Ctx.Input.RequestBody, &obj); err != nil {
  52. beego.Debug(string(o.Ctx.Input.RequestBody))
  53. o.Response(Failure, fmt.Sprintf("数据包解析错误:%s", err.Error()), nil)
  54. return
  55. }
  56. oo := models.Gateway{
  57. ID: obj.Code,
  58. Name: obj.Name,
  59. Tenant: obj.Tenant,
  60. Brand: obj.Brand,
  61. Model: obj.Model,
  62. State: obj.State,
  63. }
  64. if err := oo.SaveFromWeb(); err != nil {
  65. o.Response(Failure, fmt.Sprintf("数据插入失败:%s", err.Error()), nil)
  66. return
  67. }
  68. o.Response(Success, "成功", oo.ID)
  69. }
  70. // UpdateGateway @Title 更新网关设备
  71. // @Description 更新网关设备
  72. // @Param body controllers.ReqGateway true "数据"
  73. // @Success 0 {int} BaseResponse.Code "成功"
  74. // @Failure 1 {int} BaseResponse.Code "失败"
  75. // @router /v1/update [post]
  76. func (o *GatewayController) UpdateGateway() {
  77. var obj ReqGateway
  78. if err := json.Unmarshal(o.Ctx.Input.RequestBody, &obj); err != nil {
  79. beego.Debug(string(o.Ctx.Input.RequestBody))
  80. o.Response(Failure, fmt.Sprintf("数据包解析错误:%s", err.Error()), nil)
  81. return
  82. }
  83. oo := models.Gateway{
  84. ID: obj.Code,
  85. Name: obj.Name,
  86. Tenant: obj.Tenant,
  87. Brand: obj.Brand,
  88. Model: obj.Model,
  89. State: obj.State,
  90. }
  91. if err := oo.SaveFromWeb(); err != nil {
  92. o.Response(Failure, fmt.Sprintf("数据更新失败:%s", err.Error()), nil)
  93. return
  94. }
  95. o.Response(Success, "成功", oo.ID)
  96. }
  97. // DeleteGateway @Title 删除网关设备
  98. // @Description 删除网关设备
  99. // @Param code query string true "设备ID"
  100. // @Success 0 {int} BaseResponse.Code "成功"
  101. // @Failure 1 {int} BaseResponse.Code "失败"
  102. // @router /v1/delete [post]
  103. func (o *GatewayController) DeleteGateway() {
  104. code := strings.Trim(o.GetString("code"), " ")
  105. if len(code) == 0 {
  106. o.Response(Failure, "code为空", nil)
  107. return
  108. }
  109. c := models.Gateway{
  110. ID: code,
  111. }
  112. if err := c.Delete(); err != nil {
  113. beego.Error(fmt.Sprintf("删除失败:%s", err.Error()))
  114. o.Response(Failure, fmt.Sprintf("数据删除失败:%s", err.Error()), nil)
  115. return
  116. }
  117. o.Response(Success, "成功", c.ID)
  118. }
  119. // ImportGateway @Title 批量导入网关设备
  120. // @Description 批量导入网关设备
  121. // @Param body controllers.ReqImportGateway true "数据"
  122. // @Success 0 {int} BaseResponse.Code "成功"
  123. // @Failure 1 {int} BaseResponse.Code "失败"
  124. // @router /v1/import [post]
  125. func (o *GatewayController) ImportGateway() {
  126. var obj ReqImportGateway
  127. if err := json.Unmarshal(o.Ctx.Input.RequestBody, &obj); err != nil {
  128. beego.Debug(string(o.Ctx.Input.RequestBody))
  129. o.Response(Failure, fmt.Sprintf("数据包解析错误:%s", err.Error()), nil)
  130. return
  131. }
  132. var resp []RespImport
  133. for _, v := range obj.List {
  134. var aresp RespImport
  135. aresp.Code = v.Code
  136. oo := models.Gateway{
  137. ID: v.Code,
  138. Name: v.Name,
  139. Tenant: obj.Tenant,
  140. Brand: v.Brand,
  141. Model: v.Model,
  142. State: v.State,
  143. }
  144. err := oo.SaveFromWeb()
  145. if err != nil {
  146. aresp.Error = err.Error()
  147. beego.Error(fmt.Sprintf("zigbee集控器数据导入失败,code=%s,失败原因:%s", v.Code, err.Error()))
  148. }
  149. resp = append(resp, aresp)
  150. }
  151. o.Response(Success, "成功", resp)
  152. }
  153. // PushSerialConfig @Title 下发串口配置到网关
  154. // @Description 下发 serial.json 配置到指定网关
  155. // @Param body controllers.ReqGatewayConfig true "数据"
  156. // @Success 0 {int} BaseResponse.Code "成功"
  157. // @Failure 1 {int} BaseResponse.Code "失败"
  158. // @router /v1/serial/push [post]
  159. func (o *GatewayController) PushSerialConfig() {
  160. gid := strings.Trim(o.GetString("gid"), " ")
  161. tenant := strings.Trim(o.GetString("tenant"), " ")
  162. if gid == "" || tenant == "" {
  163. o.Response(Failure, "gid和tenant不能为空", nil)
  164. return
  165. }
  166. body := o.Ctx.Input.RequestBody
  167. if len(body) == 0 {
  168. o.Response(Failure, "请求body不能为空,请在表单中导入或添加串口配置后再下发", nil)
  169. return
  170. }
  171. sc := protocol.SerialConfig{Serial: make(map[uint8]*protocol.SerialPort)}
  172. if err := json.Unmarshal(body, &sc); err != nil || len(sc.Serial) == 0 {
  173. o.Response(Failure, "串口配置JSON解析失败或为空", nil)
  174. return
  175. }
  176. content := string(body)
  177. seq := GetNextUint64()
  178. var obj protocol.Pack_SeqFileObject
  179. obj.Data.File = "serial.json"
  180. obj.Data.Content = content
  181. str, err := obj.EnCode(gid, seq)
  182. if err != nil {
  183. o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
  184. return
  185. }
  186. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_SET_SERIAL)
  187. if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
  188. o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
  189. return
  190. }
  191. // 下发成功后,将配置存档到DB
  192. pushedCodes := make(map[int]bool)
  193. for _, v := range sc.Serial {
  194. models.G_db.Save(&models.GatewaySerial{
  195. ID: gid, ComID: int(v.Code),
  196. Interface: v.Interface, Address: v.Address,
  197. BaudRate: v.BaudRate, DataBits: int(v.DataBits),
  198. StopBits: int(v.StopBits), Parity: v.Parity,
  199. Timeout: int(v.Timeout), ProtocolType: int(v.ProtocolType),
  200. })
  201. pushedCodes[int(v.Code)] = true
  202. }
  203. // 标记不在此次推送中的旧串口为删除
  204. var oldSerials []models.GatewaySerial
  205. models.G_db.Where("id = ?", gid).Find(&oldSerials)
  206. for _, s := range oldSerials {
  207. if !pushedCodes[s.ComID] {
  208. models.G_db.Delete(&s)
  209. }
  210. }
  211. // 记录指令
  212. dcr := models.DeviceCmdRecord{ID: seq, GID: gid, DID: gid, Topic: topic, Message: str, State: 0}
  213. if err := models.G_db.Create(&dcr).Error; err != nil {
  214. beego.Error("PushSerialConfig:指令入库失败:", err.Error())
  215. }
  216. o.Response(Success, "串口配置已下发并保存", nil)
  217. }
  218. // PushDevConfig @Title 下发设备配置到网关
  219. // @Description 下发 dev/{code}.json 配置到指定网关
  220. // @Param body controllers.ReqDevConfig true "数据"
  221. // @Success 0 {int} BaseResponse.Code "成功"
  222. // @Failure 1 {int} BaseResponse.Code "失败"
  223. // @router /v1/dev/push [post]
  224. func (o *GatewayController) PushDevConfig() {
  225. gid := strings.Trim(o.GetString("gid"), " ")
  226. tenant := strings.Trim(o.GetString("tenant"), " ")
  227. codeStr := strings.Trim(o.GetString("code"), " ")
  228. if gid == "" || tenant == "" || codeStr == "" {
  229. o.Response(Failure, "gid、tenant、code不能为空", nil)
  230. return
  231. }
  232. code, err := strconv.Atoi(codeStr)
  233. if err != nil {
  234. o.Response(Failure, fmt.Sprintf("code格式错误:%s", err.Error()), nil)
  235. return
  236. }
  237. var mdc protocol.MapDevConfig
  238. // 优先从请求body解析(前端表单直接提交的配置)
  239. body := o.Ctx.Input.RequestBody
  240. if len(body) == 0 {
  241. o.Response(Failure, "请求body不能为空,请在表单中导入或添加设备配置后再下发", nil)
  242. return
  243. }
  244. if err := json.Unmarshal(body, &mdc); err != nil || len(mdc.Rtu) == 0 {
  245. o.Response(Failure, "设备配置JSON解析失败或为空", nil)
  246. return
  247. }
  248. content := string(body)
  249. seq := GetNextUint64()
  250. var obj protocol.Pack_SeqFileObject
  251. obj.Data.File = codeStr + ".json"
  252. obj.Data.Content = content
  253. str, err := obj.EnCode(gid, seq)
  254. if err != nil {
  255. o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
  256. return
  257. }
  258. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_SET_RTU)
  259. if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
  260. o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
  261. return
  262. }
  263. // 下发成功后,将配置存档到DB
  264. pushedDevCodes := make(map[string]bool)
  265. for _, v := range mdc.Rtu {
  266. models.G_db.Save(&models.GatewayDevice{
  267. ID: v.DevCode, Name: v.Name, GID: gid,
  268. ComID: int(v.Code), RtuID: int(v.DevID), TID: int(v.TID),
  269. SendCloud: v.SendCloud, WaitTime: int(v.WaitTime),
  270. ProtocolType: int(v.ProtocolType), DevType: int(v.DevType),
  271. Tenant: tenant, State: 1,
  272. })
  273. pushedDevCodes[v.DevCode] = true
  274. }
  275. // 标记不在此次推送中且同串口下的旧设备为删除
  276. var oldDevices []models.GatewayDevice
  277. models.G_db.Where("g_id = ? AND com_id = ? AND state = 1", gid, code).Find(&oldDevices)
  278. for _, d := range oldDevices {
  279. if !pushedDevCodes[d.ID] {
  280. models.G_db.Model(&d).Update("state", 0)
  281. }
  282. }
  283. // 记录指令
  284. dcr := models.DeviceCmdRecord{ID: seq, GID: gid, DID: gid, Topic: topic, Message: str, State: 0}
  285. if err := models.G_db.Create(&dcr).Error; err != nil {
  286. beego.Error("PushDevConfig:指令入库失败:", err.Error())
  287. }
  288. o.Response(Success, "设备配置已下发并保存", nil)
  289. }
  290. // PushModelConfig @Title 下发物模型配置到网关
  291. // @Description 下发 model/{tid}.json 配置到指定网关
  292. // @Param body controllers.ReqModelConfig true "数据"
  293. // @Success 0 {int} BaseResponse.Code "成功"
  294. // @Failure 1 {int} BaseResponse.Code "失败"
  295. // @router /v1/model/push [post]
  296. func (o *GatewayController) PushModelConfig() {
  297. gid := strings.Trim(o.GetString("gid"), " ")
  298. tenant := strings.Trim(o.GetString("tenant"), " ")
  299. tidStr := strings.Trim(o.GetString("tid"), " ")
  300. if gid == "" || tenant == "" || tidStr == "" {
  301. o.Response(Failure, "gid、tenant、tid不能为空", nil)
  302. return
  303. }
  304. tid, err := strconv.Atoi(tidStr)
  305. if err != nil {
  306. o.Response(Failure, fmt.Sprintf("tid格式错误:%s", err.Error()), nil)
  307. return
  308. }
  309. // 必须从请求body取模型JSON(不下发时以表单数据为准)
  310. body := o.Ctx.Input.RequestBody
  311. if len(body) == 0 {
  312. o.Response(Failure, "请求body不能为空,请在表单中选择JSON文件后再下发", nil)
  313. return
  314. }
  315. fileContent := string(body)
  316. var iot protocol.IotModel
  317. if err := json.Unmarshal(body, &iot); err != nil {
  318. o.Response(Failure, fmt.Sprintf("物模型JSON解析失败:%s", err.Error()), nil)
  319. return
  320. }
  321. device := iot.Device
  322. modelName := iot.Model
  323. protocolName := iot.Protocol
  324. seq := GetNextUint64()
  325. var obj protocol.Pack_SeqFileObject
  326. obj.Data.File = tidStr + ".json"
  327. obj.Data.Content = fileContent
  328. str, err := obj.EnCode(gid, seq)
  329. if err != nil {
  330. o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
  331. return
  332. }
  333. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_SET_MODEL)
  334. if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
  335. o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
  336. return
  337. }
  338. // 记录指令
  339. dcr := models.DeviceCmdRecord{ID: seq, GID: gid, DID: gid, Topic: topic, Message: str, State: 0}
  340. if err := models.G_db.Create(&dcr).Error; err != nil {
  341. beego.Error("PushModelConfig:指令入库失败:", err.Error())
  342. }
  343. // 下发成功后存档到 t_gateway_model
  344. models.G_db.Save(&models.GatewayModel{
  345. GID: gid, TID: uint16(tid),
  346. Device: device, Model: modelName, Protocol: protocolName,
  347. File: fileContent, PushedAt: time.Now(),
  348. })
  349. o.Response(Success, "物模型配置已下发并保存", nil)
  350. }
  351. // SerialConfigList @Title 查询网关串口列表
  352. // @Description 查询指定网关的所有串口配置
  353. // @Param gid query string true "网关ID"
  354. // @Success 0 {int} BaseResponse.Code "成功"
  355. // @router /v1/serial/list [get]
  356. func (o *GatewayController) SerialConfigList() {
  357. gid := strings.Trim(o.GetString("gid"), " ")
  358. if gid == "" {
  359. o.Response(Failure, "gid不能为空", nil)
  360. return
  361. }
  362. var serials []models.GatewaySerial
  363. if err := models.G_db.Where("id = ?", gid).Find(&serials).Error; err != nil {
  364. o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
  365. return
  366. }
  367. o.Response(Success, "成功", serials)
  368. }
  369. // SerialConfigSave @Title 保存串口配置
  370. // @Description 新增或更新单条串口配置
  371. // @Param gid query string true "网关ID"
  372. // @Param comid query int true "串口编号"
  373. // @Success 0 {int} BaseResponse.Code "成功"
  374. // @router /v1/serial/save [post]
  375. func (o *GatewayController) SerialConfigSave() {
  376. gid := strings.Trim(o.GetString("gid"), " ")
  377. comIDStr := strings.Trim(o.GetString("comid"), " ")
  378. if gid == "" || comIDStr == "" {
  379. o.Response(Failure, "gid和comid不能为空", nil)
  380. return
  381. }
  382. comID, err := strconv.Atoi(comIDStr)
  383. if err != nil {
  384. o.Response(Failure, fmt.Sprintf("comid格式错误:%s", err.Error()), nil)
  385. return
  386. }
  387. baudRate, _ := strconv.Atoi(strings.Trim(o.GetString("baudrate"), " "))
  388. dataBits, _ := strconv.Atoi(strings.Trim(o.GetString("databits"), " "))
  389. stopBits, _ := strconv.Atoi(strings.Trim(o.GetString("stopbits"), " "))
  390. timeout, _ := strconv.Atoi(strings.Trim(o.GetString("timeout"), " "))
  391. protocolType, _ := strconv.Atoi(strings.Trim(o.GetString("protocoltype"), " "))
  392. s := models.GatewaySerial{
  393. ID: gid,
  394. ComID: comID,
  395. Interface: strings.Trim(o.GetString("interface"), " "),
  396. Address: strings.Trim(o.GetString("address"), " "),
  397. BaudRate: baudRate,
  398. DataBits: dataBits,
  399. StopBits: stopBits,
  400. Parity: strings.Trim(o.GetString("parity"), " "),
  401. Timeout: timeout,
  402. ProtocolType: protocolType,
  403. }
  404. if err := models.G_db.Save(&s).Error; err != nil {
  405. o.Response(Failure, fmt.Sprintf("保存失败:%s", err.Error()), nil)
  406. return
  407. }
  408. o.Response(Success, "保存成功", nil)
  409. }
  410. // SerialConfigDelete @Title 删除串口配置
  411. // @Description 删除单条串口配置
  412. // @Param gid query string true "网关ID"
  413. // @Param comid query int true "串口编号"
  414. // @router /v1/serial/delete [post]
  415. func (o *GatewayController) SerialConfigDelete() {
  416. gid := strings.Trim(o.GetString("gid"), " ")
  417. comIDStr := strings.Trim(o.GetString("comid"), " ")
  418. if gid == "" || comIDStr == "" {
  419. o.Response(Failure, "gid和comid不能为空", nil)
  420. return
  421. }
  422. comID, err := strconv.Atoi(comIDStr)
  423. if err != nil {
  424. o.Response(Failure, fmt.Sprintf("comid格式错误:%s", err.Error()), nil)
  425. return
  426. }
  427. if err := models.G_db.Where("id = ? AND com_id = ?", gid, comID).Delete(&models.GatewaySerial{}).Error; err != nil {
  428. o.Response(Failure, fmt.Sprintf("删除失败:%s", err.Error()), nil)
  429. return
  430. }
  431. o.Response(Success, "删除成功", nil)
  432. }
  433. // DevConfigList @Title 查询网关设备列表
  434. // @Description 查询指定网关某串口下的设备配置
  435. // @Param gid query string true "网关ID"
  436. // @Param code query int false "串口编号,0查全部"
  437. // @router /v1/dev/list [get]
  438. func (o *GatewayController) DevConfigList() {
  439. gid := strings.Trim(o.GetString("gid"), " ")
  440. if gid == "" {
  441. o.Response(Failure, "gid不能为空", nil)
  442. return
  443. }
  444. codeStr := strings.Trim(o.GetString("code"), " ")
  445. var devices []models.GatewayDevice
  446. db := models.G_db.Where("g_id = ? AND state = 1", gid)
  447. if codeStr != "" {
  448. code, _ := strconv.Atoi(codeStr)
  449. db = db.Where("com_id = ?", code)
  450. }
  451. if err := db.Find(&devices).Error; err != nil {
  452. o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
  453. return
  454. }
  455. o.Response(Success, "成功", devices)
  456. }
  457. // DevConfigSave @Title 保存设备配置
  458. // @Description 新增或更新单条设备配置
  459. // @router /v1/dev/save [post]
  460. func (o *GatewayController) DevConfigSave() {
  461. gid := strings.Trim(o.GetString("gid"), " ")
  462. tenant := strings.Trim(o.GetString("tenant"), " ")
  463. devCode := strings.Trim(o.GetString("devcode"), " ")
  464. if gid == "" || devCode == "" {
  465. o.Response(Failure, "gid和devcode不能为空", nil)
  466. return
  467. }
  468. comID, _ := strconv.Atoi(strings.Trim(o.GetString("comid"), " "))
  469. rtuID, _ := strconv.Atoi(strings.Trim(o.GetString("rtuid"), " "))
  470. tid, _ := strconv.Atoi(strings.Trim(o.GetString("tid"), " "))
  471. protocolType, _ := strconv.Atoi(strings.Trim(o.GetString("protocoltype"), " "))
  472. devType, _ := strconv.Atoi(strings.Trim(o.GetString("devtype"), " "))
  473. sendCloud, _ := strconv.Atoi(strings.Trim(o.GetString("sendcloud"), " "))
  474. waitTime, _ := strconv.Atoi(strings.Trim(o.GetString("waittime"), " "))
  475. d := models.GatewayDevice{
  476. ID: devCode,
  477. Name: strings.Trim(o.GetString("name"), " "),
  478. GID: gid,
  479. ComID: comID,
  480. RtuID: rtuID,
  481. TID: tid,
  482. SendCloud: sendCloud,
  483. WaitTime: waitTime,
  484. ProtocolType: protocolType,
  485. DevType: devType,
  486. Tenant: tenant,
  487. State: 1,
  488. }
  489. if err := models.G_db.Save(&d).Error; err != nil {
  490. o.Response(Failure, fmt.Sprintf("保存失败:%s", err.Error()), nil)
  491. return
  492. }
  493. o.Response(Success, "保存成功", nil)
  494. }
  495. // DevConfigDelete @Title 删除设备配置
  496. // @Description 删除单条设备配置(软删除)
  497. // @Param devcode query string true "设备编码"
  498. // @router /v1/dev/delete [post]
  499. func (o *GatewayController) DevConfigDelete() {
  500. devCode := strings.Trim(o.GetString("devcode"), " ")
  501. if devCode == "" {
  502. o.Response(Failure, "devcode不能为空", nil)
  503. return
  504. }
  505. if err := models.G_db.Model(&models.GatewayDevice{}).Where("id = ?", devCode).Update("state", 0).Error; err != nil {
  506. o.Response(Failure, fmt.Sprintf("删除失败:%s", err.Error()), nil)
  507. return
  508. }
  509. o.Response(Success, "删除成功", nil)
  510. }
  511. // ModelList @Title 查询物模型列表
  512. // @Description 查询网关已下发的物模型(默认查所有,传gid查指定网关)
  513. // @router /v1/model/list [get]
  514. func (o *GatewayController) ModelList() {
  515. gid := strings.Trim(o.GetString("gid"), " ")
  516. type ModelSummary struct {
  517. TID uint16 `json:"tid"`
  518. Device string `json:"device"`
  519. Model string `json:"model"`
  520. Protocol string `json:"protocol"`
  521. File string `json:"file,omitempty"`
  522. }
  523. var list []ModelSummary
  524. if gid != "" {
  525. var gms []models.GatewayModel
  526. if err := models.G_db.Where("g_id = ?", gid).Find(&gms).Error; err != nil {
  527. o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
  528. return
  529. }
  530. for _, v := range gms {
  531. list = append(list, ModelSummary{TID: v.TID, Device: v.Device, Model: v.Model, Protocol: v.Protocol, File: v.File})
  532. }
  533. } else {
  534. var gms []models.GatewayModel
  535. if err := models.G_db.Find(&gms).Error; err != nil {
  536. o.Response(Failure, fmt.Sprintf("查询失败:%s", err.Error()), nil)
  537. return
  538. }
  539. for _, v := range gms {
  540. list = append(list, ModelSummary{TID: v.TID, Device: v.Device, Model: v.Model, Protocol: v.Protocol, File: v.File})
  541. }
  542. }
  543. o.Response(Success, "成功", list)
  544. }
  545. // ModelSave @Title 保存物模型
  546. // @Description 新增或更新物模型到 t_gateway_model
  547. // @router /v1/model/save [post]
  548. func (o *GatewayController) ModelSave() {
  549. gid := strings.Trim(o.GetString("gid"), " ")
  550. tidStr := strings.Trim(o.GetString("tid"), " ")
  551. if gid == "" || tidStr == "" {
  552. o.Response(Failure, "gid和tid不能为空", nil)
  553. return
  554. }
  555. tid, err := strconv.Atoi(tidStr)
  556. if err != nil {
  557. o.Response(Failure, fmt.Sprintf("tid格式错误:%s", err.Error()), nil)
  558. return
  559. }
  560. device := strings.Trim(o.GetString("device"), " ")
  561. modelName := strings.Trim(o.GetString("model"), " ")
  562. protocol := strings.Trim(o.GetString("protocol"), " ")
  563. fileContent := strings.Trim(o.GetString("file"), " ")
  564. if fileContent == "" {
  565. // 尝试从上传的请求body中读取文件内容
  566. fileContent = string(o.Ctx.Input.RequestBody)
  567. }
  568. m := models.GatewayModel{
  569. GID: gid, TID: uint16(tid),
  570. Device: device, Model: modelName, Protocol: protocol,
  571. File: fileContent, PushedAt: time.Now(),
  572. }
  573. if err := models.G_db.Save(&m).Error; err != nil {
  574. o.Response(Failure, fmt.Sprintf("保存失败:%s", err.Error()), nil)
  575. return
  576. }
  577. o.Response(Success, "保存成功", nil)
  578. }
  579. // ModelDelete @Title 删除物模型
  580. // @Description 从 t_gateway_model 删除指定物模型
  581. // @router /v1/model/delete [post]
  582. func (o *GatewayController) ModelDelete() {
  583. gid := strings.Trim(o.GetString("gid"), " ")
  584. tidStr := strings.Trim(o.GetString("tid"), " ")
  585. if gid == "" || tidStr == "" {
  586. o.Response(Failure, "gid和tid不能为空", nil)
  587. return
  588. }
  589. tid, err := strconv.Atoi(tidStr)
  590. if err != nil {
  591. o.Response(Failure, fmt.Sprintf("tid格式错误:%s", err.Error()), nil)
  592. return
  593. }
  594. if err := models.G_db.Where("g_id = ? AND t_id = ?", gid, tid).Delete(&models.GatewayModel{}).Error; err != nil {
  595. o.Response(Failure, fmt.Sprintf("删除失败:%s", err.Error()), nil)
  596. return
  597. }
  598. o.Response(Success, "删除成功", nil)
  599. }
  600. // LogQuery @Title 查询边缘日志
  601. // @Description 触发边缘上传日志文件
  602. // @router /v1/log/query [post]
  603. func (o *GatewayController) LogQuery() {
  604. gid := strings.Trim(o.GetString("gid"), " ")
  605. tenant := strings.Trim(o.GetString("tenant"), " ")
  606. if gid == "" || tenant == "" {
  607. o.Response(Failure, "gid和tenant不能为空", nil)
  608. return
  609. }
  610. seq := GetNextUint64()
  611. var obj protocol.Pack_IDObject
  612. str, err := obj.EnCode(gid, seq, 0)
  613. if err != nil {
  614. o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
  615. return
  616. }
  617. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_LOG)
  618. beego.Info("LogQuery:发布,topic=", topic, ",payload=", str)
  619. if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
  620. beego.Error("LogQuery:MQTT发布失败,topic=", topic, ",err=", err)
  621. o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
  622. return
  623. }
  624. beego.Info("LogQuery:发布成功,topic=", topic)
  625. o.Response(Success, "日志查询指令已发送,请稍后刷新查看", map[string]interface{}{"seq": seq})
  626. }
  627. // LogList @Title 列出已保存的日志文件
  628. // @Description 列出网关已上传的日志文件
  629. // @router /v1/log/list [get]
  630. func (o *GatewayController) LogList() {
  631. gid := strings.Trim(o.GetString("gid"), " ")
  632. tenant := strings.Trim(o.GetString("tenant"), " ")
  633. if gid == "" || tenant == "" {
  634. o.Response(Failure, "gid和tenant不能为空", nil)
  635. return
  636. }
  637. dir := "/opt/ipole_data/gateway_logs/" + tenant + "/" + gid + "/log/"
  638. entries, err := os.ReadDir(dir)
  639. if err != nil {
  640. o.Response(Failure, "该网关暂无日志文件,请先点击「查询日志」", nil)
  641. return
  642. }
  643. type LogFile struct {
  644. Name string `json:"name"`
  645. Size int64 `json:"size"`
  646. Time string `json:"time"`
  647. }
  648. var files []LogFile
  649. for _, e := range entries {
  650. if !e.IsDir() {
  651. info, _ := e.Info()
  652. f := LogFile{Name: e.Name(), Size: info.Size(), Time: info.ModTime().Format("2006-01-02 15:04:05")}
  653. files = append(files, f)
  654. }
  655. }
  656. o.Response(Success, "成功", files)
  657. }
  658. // LogDelete @Title 删除边缘日志
  659. // @Description 远程删除边缘网关的所有日志文件
  660. // @Param gid query string true "网关ID"
  661. // @Param tenant query string true "租户ID"
  662. // @router /v1/log/delete [post]
  663. func (o *GatewayController) LogDelete() {
  664. gid := strings.Trim(o.GetString("gid"), " ")
  665. tenant := strings.Trim(o.GetString("tenant"), " ")
  666. if gid == "" || tenant == "" {
  667. o.Response(Failure, "gid和tenant不能为空", nil)
  668. return
  669. }
  670. seq := GetNextUint64()
  671. var obj protocol.Pack_IDObject
  672. str, err := obj.EnCode(gid, seq, 0)
  673. if err != nil {
  674. o.Response(Failure, fmt.Sprintf("编码失败:%s", err.Error()), nil)
  675. return
  676. }
  677. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_REMOVE_LOG)
  678. beego.Info("LogDelete:发布,topic=", topic)
  679. if err := GetMqttHandler().PublishString(topic, str, mqtt.AtLeastOnce); err != nil {
  680. beego.Error("LogDelete:MQTT发布失败,topic=", topic, ",err=", err)
  681. o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
  682. return
  683. }
  684. beego.Info("LogDelete:发布成功,topic=", topic)
  685. o.Response(Success, "日志删除指令已发送", nil)
  686. }
  687. // LogCfg @Title 设置边缘日志等级
  688. // @Description 云端下发日志等级配置到边缘网关
  689. // @Param gid query string true "网关ID"
  690. // @Param tenant query string true "租户ID"
  691. // @Success 0 {int} BaseResponse.Code "成功"
  692. // @Failure 1 {int} BaseResponse.Code "失败"
  693. // @router /v1/log/cfg [post]
  694. func (o *GatewayController) LogCfg() {
  695. gid := strings.Trim(o.GetString("gid"), " ")
  696. tenant := strings.Trim(o.GetString("tenant"), " ")
  697. if gid == "" || tenant == "" {
  698. o.Response(Failure, "gid和tenant不能为空", nil)
  699. return
  700. }
  701. body := o.Ctx.Input.RequestBody
  702. if len(body) == 0 {
  703. o.Response(Failure, "请求body为空,请提供日志等级配置", nil)
  704. return
  705. }
  706. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_LOG_CFG)
  707. beego.Info("LogCfg:发布,topic=", topic)
  708. if err := GetMqttHandler().PublishString(topic, string(body), mqtt.AtLeastOnce); err != nil {
  709. beego.Error("LogCfg:MQTT发布失败,topic=", topic, ",err=", err)
  710. o.Response(Failure, fmt.Sprintf("MQTT发布失败:%s", err.Error()), nil)
  711. return
  712. }
  713. beego.Info("LogCfg:发布成功,topic=", topic)
  714. o.Response(Success, "日志等级配置已下发", nil)
  715. }
  716. // LogRead @Title 读取日志内容
  717. // @Description 读取指定日志文件内容
  718. // @router /v1/log/read [get]
  719. func (o *GatewayController) LogRead() {
  720. gid := strings.Trim(o.GetString("gid"), " ")
  721. tenant := strings.Trim(o.GetString("tenant"), " ")
  722. fname := strings.Trim(o.GetString("file"), " ")
  723. if gid == "" || tenant == "" || fname == "" {
  724. o.Response(Failure, "gid、tenant、file不能为空", nil)
  725. return
  726. }
  727. // 防止路径遍历
  728. if strings.Contains(fname, "..") || strings.Contains(fname, "/") || strings.Contains(fname, "\\") {
  729. o.Response(Failure, "非法文件名", nil)
  730. return
  731. }
  732. path := "/opt/ipole_data/gateway_logs/" + tenant + "/" + gid + "/log/" + fname
  733. buf, err := os.ReadFile(path)
  734. if err != nil {
  735. o.Response(Failure, fmt.Sprintf("读取失败:%s", err.Error()), nil)
  736. return
  737. }
  738. // 限制返回最近100KB,避免超大文件
  739. content := string(buf)
  740. if len(content) > 100*1024 {
  741. content = content[len(content)-100*1024:]
  742. }
  743. o.Response(Success, "成功", content)
  744. }
  745. // DeployUpload @Title 上传部署文件
  746. // @Description 上传网关部署文件
  747. // @router /v1/deploy/upload [post]
  748. func (o *GatewayController) DeployUpload() {
  749. gid := strings.Trim(o.GetString("gid"), " ")
  750. tenant := strings.Trim(o.GetString("tenant"), " ")
  751. file, header, err := o.Ctx.Request.FormFile("file")
  752. if err != nil {
  753. o.Response(Failure, "读取上传文件失败: "+err.Error(), nil)
  754. return
  755. }
  756. defer file.Close()
  757. data, err := io.ReadAll(file)
  758. if err != nil {
  759. o.Response(Failure, "读取文件内容失败: "+err.Error(), nil)
  760. return
  761. }
  762. md5Hash := fmt.Sprintf("%x", md5.Sum(data))
  763. // 保存到临时目录
  764. tmpDir := filepath.Join(os.TempDir(), "ipole_deploy")
  765. os.MkdirAll(tmpDir, os.ModePerm)
  766. tmpFile := filepath.Join(tmpDir, md5Hash)
  767. if err := os.WriteFile(tmpFile, data, os.ModePerm); err != nil {
  768. o.Response(Failure, "保存文件失败: "+err.Error(), nil)
  769. return
  770. }
  771. // 创建部署记录(原始SQL绕过GORM零值问题)
  772. now := time.Now()
  773. 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 {
  774. o.Response(Failure, "创建部署记录失败: "+err.Error(), nil)
  775. return
  776. }
  777. o.Response(Success, "文件上传成功", map[string]interface{}{
  778. "md5": md5Hash,
  779. "size": len(data),
  780. "path": tmpFile,
  781. "fileName": header.Filename,
  782. })
  783. }
  784. // DeployPush @Title 下发部署指令
  785. // @Description 下发部署指令到网关
  786. // @router /v1/deploy/push [post]
  787. func (o *GatewayController) DeployPush() {
  788. gid := strings.Trim(o.GetString("gid"), " ")
  789. tenant := strings.Trim(o.GetString("tenant"), " ")
  790. version := strings.Trim(o.GetString("version"), " ")
  791. fileName := strings.Trim(o.GetString("fileName"), " ")
  792. filePath := strings.Trim(o.GetString("filePath"), " ")
  793. md5Hash := strings.Trim(o.GetString("md5"), " ")
  794. if gid == "" || tenant == "" || filePath == "" {
  795. o.Response(Failure, "参数不完整", nil)
  796. return
  797. }
  798. if fileName == "" {
  799. fileName = "ipole"
  800. }
  801. if version == "" {
  802. version = fmt.Sprintf("v%s", time.Now().Format("20060102-150405"))
  803. }
  804. data, err := os.ReadFile(filePath)
  805. if err != nil {
  806. o.Response(Failure, "读取部署文件失败: "+err.Error(), nil)
  807. return
  808. }
  809. const chunkSize = 64 * 1024 // 64KB
  810. totalChunks := (len(data) + chunkSize - 1) / chunkSize
  811. // 1. 发送部署指令
  812. var cmd protocol.Pack_DeployCmd
  813. seq := GetNextUint64()
  814. cmdStr, err := cmd.EnCodeCmd(gid, gid, seq, version, fileName, md5Hash, totalChunks, chunkSize)
  815. if err != nil {
  816. o.Response(Failure, "编码部署指令失败: "+err.Error(), nil)
  817. return
  818. }
  819. topic := GetTopic(tenant, protocol.DT_GATEWAY, gid, protocol.TP_GW_DEPLOY)
  820. GetMqttHandler().PublishString(topic, cmdStr, mqtt.AtLeastOnce)
  821. // 2. 逐片发送
  822. for i := 0; i < totalChunks; i++ {
  823. start := i * chunkSize
  824. end := start + chunkSize
  825. if end > len(data) {
  826. end = len(data)
  827. }
  828. chunkData := base64.StdEncoding.EncodeToString(data[start:end])
  829. seq := GetNextUint64()
  830. chunkStr, err := cmd.EnCodeChunk(gid, gid, seq, version, i, chunkData)
  831. if err != nil {
  832. o.Response(Failure, fmt.Sprintf("编码分片%d失败: %s", i, err.Error()), nil)
  833. return
  834. }
  835. GetMqttHandler().PublishString(topic, chunkStr, mqtt.AtLeastOnce)
  836. if i%20 == 0 && i > 0 {
  837. time.Sleep(50 * time.Millisecond)
  838. }
  839. }
  840. models.G_db.Model(&models.GatewayDeploy{}).
  841. Where("gid = ? AND status = 0", gid).
  842. Updates(map[string]interface{}{
  843. "to_version": version,
  844. "from_version": "",
  845. "md5": md5Hash,
  846. "update_time": time.Now(),
  847. })
  848. o.Response(Success, "部署指令已下发", map[string]interface{}{
  849. "totalChunks": totalChunks,
  850. "chunkSize": chunkSize,
  851. "fileSize": len(data),
  852. })
  853. }
  854. // DeployStatus @Title 查询部署状态
  855. // @Description 查询网关部署状态
  856. // @router /v1/deploy/status [get]
  857. func (o *GatewayController) DeployStatus() {
  858. gid := strings.Trim(o.GetString("gid"), " ")
  859. var deploy models.GatewayDeploy
  860. err := models.G_db.Where("gid = ?", gid).Order("create_time DESC").First(&deploy).Error
  861. if err != nil {
  862. // record not found 返回空状态,方便前端轮询等待
  863. o.Response(Success, "暂无部署记录", map[string]interface{}{
  864. "status": -1,
  865. })
  866. return
  867. }
  868. statusText := map[uint8]string{0: "进行中", 1: "成功", 2: "失败", 3: "已回滚"}
  869. o.Response(Success, "查询成功", map[string]interface{}{
  870. "id": deploy.ID,
  871. "fromVersion": deploy.FromVersion,
  872. "toVersion": deploy.ToVersion,
  873. "status": deploy.Status,
  874. "statusText": statusText[deploy.Status],
  875. "errorMsg": deploy.ErrorMsg,
  876. "phase1Result": deploy.Phase1Result,
  877. "phase2Result": deploy.Phase2Result,
  878. "createTime": deploy.CreateTime.Format("2006-01-02 15:04:05"),
  879. })
  880. }
  881. // DeployBatchStatus @Title 批量查询部署状态
  882. // @Description 一次查询多个网关的最近一次部署状态
  883. // @router /v1/deploy/batch-status [get]
  884. func (o *GatewayController) DeployBatchStatus() {
  885. tenant := strings.Trim(o.GetString("tenant"), " ")
  886. gidsStr := strings.Trim(o.GetString("gids"), " ")
  887. if tenant == "" {
  888. o.Response(Failure, "tenant不能为空", nil)
  889. return
  890. }
  891. // 两步查询:先取每个 gid 最新 id,再取完整记录
  892. query := models.G_db.Table("t_gateway_deploy").
  893. Where("tenant = ?", tenant)
  894. if gidsStr != "" {
  895. gids := strings.Split(gidsStr, ",")
  896. for i := range gids {
  897. gids[i] = strings.Trim(gids[i], " ")
  898. }
  899. query = query.Where("gid IN (?)", gids)
  900. }
  901. var maxIds []int64
  902. dbResult := query.Select("MAX(id)").Group("gid").Pluck("MAX(id)", &maxIds)
  903. if dbResult.Error != nil {
  904. o.Response(Failure, "查询部署记录失败: "+dbResult.Error.Error(), nil)
  905. return
  906. }
  907. var deploys []models.GatewayDeploy
  908. if len(maxIds) > 0 {
  909. models.G_db.Where("id IN (?)", maxIds).Find(&deploys)
  910. }
  911. result := make(map[string]interface{})
  912. for _, deploy := range deploys {
  913. result[deploy.GID] = map[string]interface{}{
  914. "id": deploy.ID,
  915. "status": deploy.Status,
  916. "toVersion": deploy.ToVersion,
  917. "phase1Result": deploy.Phase1Result,
  918. "phase2Result": deploy.Phase2Result,
  919. "errorMsg": deploy.ErrorMsg,
  920. "createTime": deploy.CreateTime.Format("2006-01-02 15:04:05"),
  921. }
  922. }
  923. // 补充 gids 列表中有但无部署记录的网关(value = null)
  924. if gidsStr != "" {
  925. for _, gid := range strings.Split(gidsStr, ",") {
  926. gid = strings.Trim(gid, " ")
  927. if gid == "" {
  928. continue
  929. }
  930. if _, ok := result[gid]; !ok {
  931. result[gid] = nil
  932. }
  933. }
  934. }
  935. o.Response(Success, "查询成功", result)
  936. }