| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677 |
- package controllers
- import (
- "Tech/common"
- "Tech/models"
- "fmt"
- "sync/atomic"
- "github.com/beego/beego/v2/client/orm"
- beego "github.com/beego/beego/v2/server/web"
- "github.com/gorilla/websocket"
- )
- // WebSocketController WebSocket处理器
- type WebSocketController struct {
- beego.Controller
- }
- var wsClientCounter int64
- // WebSocket处理连接
- func (c *WebSocketController) Handle() {
- // 升级HTTP连接到WebSocket
- conn, err := websocket.Upgrade(c.Ctx.ResponseWriter, c.Ctx.Request, nil, 1024, 1024)
- if err != nil {
- c.Data["json"] = map[string]interface{}{
- "code": 500,
- "msg": "WebSocket连接升级失败: " + err.Error(),
- }
- c.ServeJSON()
- return
- }
- // 生成唯一的客户端ID
- clientID := fmt.Sprintf("ws_client_%d", atomic.AddInt64(&wsClientCounter, 1))
- // 注册WebSocket连接
- wsConn := common.GlobalWSManager.Register(clientID, conn)
- defer func() {
- common.GlobalWSManager.Unregister(clientID)
- }()
- // 获取最新的课程运行记录以获取线路号
- o := orm.NewOrm()
- var latestCourse models.CourseRunLog
- err = o.QueryTable(new(models.CourseRunLog)).
- OrderBy("-RunLogID").
- One(&latestCourse)
- lineID := 0
- if err == nil && err != orm.ErrNoRows {
- // 根据运行课程ID查询课程信息以获取线路号
- var courseInfo models.CourseInfo
- err = o.QueryTable(new(models.CourseInfo)).
- Filter("CourseID", latestCourse.RunCourseID).
- One(&courseInfo)
- if err == nil {
- lineID = courseInfo.LineNo
- }
- }
- // 发送连接成功消息
- wsConn.SendChan <- []byte(fmt.Sprintf(`{"event":"connected","data":{"clientID":"%s","message":"WebSocket连接已建立","lineID":%d}}`, clientID, lineID))
- // 立即发送所有缓存的设备数据(首次连接推送全量数据)
- common.SendCurrentDataToClient(clientID)
- // 读取客户端消息
- for {
- _, _, err := conn.ReadMessage()
- if err != nil {
- break
- }
- }
- }
|