client.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477
  1. package client
  2. import (
  3. "bytes"
  4. "encoding/json"
  5. "fmt"
  6. "mtp20access/global"
  7. rsp "mtp20access/model/mq/response"
  8. "mtp20access/model/quote/request"
  9. "mtp20access/packet"
  10. "mtp20access/publish"
  11. "sync"
  12. "time"
  13. "github.com/gorilla/websocket"
  14. "github.com/mitchellh/mapstructure"
  15. "go.uber.org/zap"
  16. )
  17. var Clients map[int]*Client // key:SessionID
  18. type Client struct {
  19. LoginRedis
  20. // FIXME: - 这里是否需要定义多个锁?
  21. mtx sync.RWMutex
  22. curSerialNumber uint32 // 当前业务流水号
  23. asyncTasks map[string]*AsyncTask // key:SessionId_FuncodeRsp_SerialNumber
  24. wsQuoteConn *websocket.Conn // 终端行情WebSocket连接
  25. quoteWriteChan chan []byte // 推送队列 Access -> Client
  26. quoteChan chan interface{} // 接收实时行情订阅channel QuoteServer -> Access
  27. quoteSubs []request.QuoteGoods // 当前已订阅行情的商品
  28. qmtx sync.RWMutex
  29. // wsQuoteCloseChan chan struct{} // 终端WebSocket连接关闭信号
  30. wsTradeConn *websocket.Conn // 终端交易WebSocket连接
  31. tradeWriteChan chan []byte // 推送队列 Access -> Client
  32. tradeChan chan interface{} // 接收交易通知channel RabbitMQ -> Access
  33. tmtx sync.RWMutex
  34. // wsTradeCloseChan chan struct{} // 终端WebSocket连接关闭信号
  35. }
  36. // GetSerialNumber 获取可用流水号
  37. func (r *Client) GetSerialNumber() uint32 {
  38. r.mtx.Lock()
  39. defer func() {
  40. r.mtx.Unlock()
  41. }()
  42. r.curSerialNumber += 1
  43. return r.curSerialNumber
  44. }
  45. // SetQuoteSubs 设置商品行情订阅信息
  46. func (r *Client) SetQuoteSubs(req request.QuoteSubscribeReq) {
  47. r.qmtx.Lock()
  48. defer r.qmtx.Unlock()
  49. r.quoteSubs = req.QuoteGoodses
  50. // 从订阅中心订阅行情通知
  51. if r.quoteChan != nil {
  52. global.M2A_Publish.Unsubscribe(publish.Topic_Quote, r.quoteChan)
  53. }
  54. r.quoteChan = global.M2A_Publish.Subscribe(publish.Topic_Quote)
  55. go func() {
  56. for {
  57. msg, ok := <-r.quoteChan
  58. if !ok {
  59. return // 管道已关闭,退出循环
  60. }
  61. // 向客户端发送行情信息
  62. if p, ok := msg.(*packet.MiQuotePacket); ok {
  63. DispatchRealQuote(p, r)
  64. }
  65. }
  66. }()
  67. }
  68. // SetQuoteSubs 设置交易通知订阅信息
  69. func (r *Client) SetTradeNft() {
  70. r.tmtx.Lock()
  71. defer r.tmtx.Unlock()
  72. // 从订阅中心订阅交易通知
  73. if r.tradeChan != nil {
  74. global.M2A_Publish.Unsubscribe(publish.Topic_Trading, r.tradeChan)
  75. }
  76. r.tradeChan = global.M2A_Publish.Subscribe(publish.Topic_Trading)
  77. go func() {
  78. for {
  79. msg, ok := <-r.tradeChan
  80. if !ok {
  81. return // 管道已关闭,退出循环
  82. }
  83. // 向客户端发送交易通知信息
  84. if m, ok := msg.(map[string][]byte); ok {
  85. // 判断是否自己的通知
  86. if p, ok := m[r.SessionID]; ok {
  87. // 发送给客户端
  88. if r.tradeWriteChan != nil {
  89. r.tradeWriteChan <- p
  90. }
  91. }
  92. }
  93. }
  94. }()
  95. }
  96. // GetQuoteSubs 获取已订阅行情商品信息列表
  97. func (r *Client) GetQuoteSubs() (quoteSubs []request.QuoteGoods) {
  98. r.qmtx.Lock()
  99. defer r.qmtx.Unlock()
  100. return r.quoteSubs
  101. }
  102. // // WriteWsBuf 向客户端发送实时行情
  103. // func (r *Client) WriteQuoteWsBuf(buf []byte) (err error) {
  104. // if r.quoteWriteChan != nil {
  105. // r.quoteWriteChan <- buf
  106. // }
  107. // return
  108. // }
  109. // func (r *Client) WriteTradeWsBuf(buf []byte) (err error) {
  110. // if r.tradeWriteChan != nil {
  111. // r.tradeWriteChan <- buf
  112. // } else {
  113. // global.M2A_LOG.Info("WriteTradeWsBuf::tradeWriteChan == nil", zap.Any("LoginID", r.LoginID), zap.Any("SessionID", r.SessionID))
  114. // }
  115. // if r.wsTradeConn != nil {
  116. // // r.tradeWriteChan <- buf
  117. // } else {
  118. // global.M2A_LOG.Info("WriteTradeWsBuf::wsTradeConn == nil", zap.Any("LoginID", r.LoginID), zap.Any("SessionID", r.SessionID))
  119. // global.M2A_LOG.Info("当前所有Client----")
  120. // for _, item := range Clients {
  121. // global.M2A_LOG.Info("client",
  122. // zap.Any("LoginID", item.LoginID),
  123. // zap.Any("Group", item.Group), zap.Any("SessionID", item.SessionID),
  124. // zap.Any("wsTradeConn", item.wsTradeConn), zap.Any("tradeWriteChan", item.tradeWriteChan),
  125. // zap.Any("wsQuoteConn", item.wsQuoteConn), zap.Any("quoteWriteChan", item.quoteWriteChan))
  126. // }
  127. // global.M2A_LOG.Info("----------------")
  128. // }
  129. // return
  130. // }
  131. // GetAsyncTask 获取目标异步任务
  132. // key:SessionId_FuncodeRsp
  133. func (r *Client) GetAsyncTask(key string) (asyncTask *AsyncTask) {
  134. r.mtx.RLock()
  135. defer func() {
  136. r.mtx.RUnlock()
  137. }()
  138. asyncTask, ok := r.asyncTasks[key]
  139. if ok {
  140. return asyncTask
  141. } else {
  142. return nil
  143. }
  144. }
  145. // SetAsyncTask 设置异步任务
  146. func (r *Client) SetAsyncTask(asyncTask *AsyncTask, key string) {
  147. r.mtx.Lock()
  148. defer func() {
  149. r.mtx.Unlock()
  150. }()
  151. if r.asyncTasks == nil {
  152. r.asyncTasks = make(map[string]*AsyncTask, 0)
  153. }
  154. delete(r.asyncTasks, key)
  155. r.asyncTasks[key] = asyncTask
  156. }
  157. // DeleteAsyncTask 删除异步任务
  158. func (r *Client) DeleteAsyncTask(key string) {
  159. r.mtx.Lock()
  160. defer func() {
  161. r.mtx.Unlock()
  162. }()
  163. delete(r.asyncTasks, key)
  164. }
  165. func (r *Client) GetAllAsyncTask() *map[string]*AsyncTask {
  166. return &r.asyncTasks
  167. }
  168. // **************** Quote WebSocket ****************
  169. func (r *Client) SetQuoteWebSocket(ws *websocket.Conn) (err error) {
  170. r.qmtx.Lock()
  171. defer r.qmtx.Unlock()
  172. // if r.wsQuoteConn != nil {
  173. // r.wsQuoteConn.Close()
  174. // }
  175. r.wsQuoteConn = ws
  176. r.quoteWriteChan = make(chan []byte, 100)
  177. // 开始读取客户端发送信息
  178. go r.readClientWsQuoteMessage()
  179. // 开始推送客户端信息循环
  180. go r.writeClientWsQuoteMessage()
  181. // r.wsQuoteConn.SetCloseHandler(func(code int, text string) error {
  182. // close(r.wsCloseChan)
  183. // return nil
  184. // })
  185. return
  186. }
  187. // readClientWsQuoteMessage 处理终端发过来的websocket数据
  188. // 注意: 阻塞式, 直到websocket关闭才退出
  189. func (r *Client) readClientWsQuoteMessage() {
  190. for {
  191. fmt.Println("readClientWsQuoteMessage start")
  192. // 40秒心跳超时
  193. r.wsQuoteConn.SetReadDeadline(time.Now().Add(40 * time.Second))
  194. mt, msg, err := r.wsQuoteConn.ReadMessage()
  195. if err != nil {
  196. fmt.Printf("readClientWsQuoteMessage: %v\n", err)
  197. return
  198. }
  199. switch mt {
  200. case websocket.PingMessage:
  201. _ = r.wsQuoteConn.WriteMessage(mt, msg)
  202. case websocket.CloseMessage:
  203. return
  204. case websocket.BinaryMessage:
  205. if err := r.clientToQuoteAgentMsg(msg); err != nil {
  206. return
  207. }
  208. case websocket.TextMessage:
  209. fmt.Println(string(msg))
  210. _ = r.wsQuoteConn.WriteMessage(mt, msg)
  211. }
  212. // FIXME: - 这里要判断是否有问题
  213. // select {
  214. // case <-r.wsCloseChan:
  215. // return
  216. // }
  217. }
  218. }
  219. // writeClientWsQuoteMessage 由于websocket非线程安全,
  220. // 所以由统一协程写入
  221. func (r *Client) writeClientWsQuoteMessage() {
  222. // defer r.close()
  223. // for {
  224. // select {
  225. // case buf := <-r.quoteWriteChan:
  226. // err := r.wsQuoteConn.WriteMessage(websocket.BinaryMessage, buf)
  227. // if err != nil {
  228. // return
  229. // }
  230. // case <-r.wsQuoteCloseChan: // 与终端连接关闭信息
  231. // return
  232. // }
  233. // }
  234. for {
  235. buf, ok := <-r.quoteWriteChan
  236. if !ok {
  237. // 通道已关闭
  238. return
  239. }
  240. err := r.wsQuoteConn.WriteMessage(websocket.BinaryMessage, buf)
  241. if err != nil {
  242. return
  243. }
  244. }
  245. }
  246. // clientToQuoteAgentMsg 处理客户端发上来的包
  247. func (r *Client) clientToQuoteAgentMsg(msg []byte) error {
  248. var p packet.MiQuotePacket
  249. err := p.UnPackHead(msg)
  250. if err != nil {
  251. // logger.Logger().Errorf("[%v c->s] invalid packet: %v", r.ProxyId, err)
  252. return err
  253. }
  254. // 长度15 是心跳包, 目前只允许客户端上行心跳包
  255. if p.Length == 15 {
  256. // logger.Logger().Debugf("%v [%v c->s] %s", r.ClientIP, r.ProxyId, p.HeaderInfo())
  257. // 将心跳发回给客户端
  258. err = r.wsQuoteConn.WriteMessage(websocket.BinaryMessage, msg)
  259. if err != nil {
  260. return err
  261. }
  262. }
  263. return nil
  264. }
  265. // **************** Trade WebSocket ****************
  266. func (r *Client) SetTradeWebSocket(ws *websocket.Conn) (err error) {
  267. r.tmtx.Lock()
  268. defer r.tmtx.Unlock()
  269. // if r.wsTradeConn != nil {
  270. // r.wsTradeConn.Close()
  271. // }
  272. r.wsTradeConn = ws
  273. r.tradeWriteChan = make(chan []byte, 100)
  274. r.wsTradeConn.SetCloseHandler(func(code int, text string) error {
  275. global.M2A_LOG.Info("交易长连发生断开", zap.Any("LoginID", r.LoginID), zap.Any("SessionID", r.SessionID), zap.Any("code", code), zap.Any("text", text))
  276. return nil
  277. })
  278. // 开始读取客户端发送信息
  279. go r.readClientWsTradeMessage()
  280. // 开始推送客户端信息循环
  281. go r.writeClientWsTradeMessage()
  282. // 设置交易通知订阅信息
  283. r.SetTradeNft()
  284. global.M2A_LOG.Info("交易长连客户端接入", zap.Any("LoginID", r.LoginID), zap.Any("SessionID", r.SessionID))
  285. return
  286. }
  287. // readClientWsQuoteMessage 处理终端发过来的websocket数据
  288. // 注意: 阻塞式, 直到websocket关闭才退出
  289. func (r *Client) readClientWsTradeMessage() {
  290. for {
  291. fmt.Println("readClientWsTradeMessage start")
  292. // 40秒心跳超时
  293. r.wsTradeConn.SetReadDeadline(time.Now().Add(40 * time.Second))
  294. mt, msg, err := r.wsTradeConn.ReadMessage()
  295. if err != nil {
  296. global.M2A_LOG.Info("readClientWsTradeMessage::ReadMessage发生错误", zap.Any("err", err))
  297. return
  298. }
  299. switch mt {
  300. case websocket.PingMessage:
  301. if err = r.wsTradeConn.WriteMessage(mt, msg); err != nil {
  302. global.M2A_LOG.Info("readClientWsTradeMessage::PingMessage发生错误", zap.Any("err", err))
  303. }
  304. case websocket.CloseMessage:
  305. return
  306. case websocket.BinaryMessage:
  307. if err := r.clientToTradeAgentMsg(msg); err != nil {
  308. global.M2A_LOG.Info("readClientWsTradeMessage::BinaryMessage发生错误", zap.Any("err", err))
  309. return
  310. }
  311. case websocket.TextMessage:
  312. if err = r.wsTradeConn.WriteMessage(mt, msg); err != nil {
  313. global.M2A_LOG.Info("readClientWsTradeMessage::TextMessage发生错误", zap.Any("err", err))
  314. return
  315. }
  316. }
  317. }
  318. }
  319. // writeClientWsQuoteMessage 由于websocket非线程安全,
  320. // 所以由统一协程写入
  321. func (r *Client) writeClientWsTradeMessage() {
  322. for {
  323. buf, ok := <-r.tradeWriteChan
  324. if !ok {
  325. // 通道已关闭
  326. global.M2A_LOG.Info("writeClientWsTradeMessage 通道已关闭", zap.Any("LoginID", r.LoginID), zap.Any("SessionID", r.SessionID))
  327. return
  328. }
  329. err := r.wsTradeConn.WriteMessage(websocket.BinaryMessage, buf)
  330. if err != nil {
  331. global.M2A_LOG.Info("writeClientWsTradeMessage 向客户端发送信息时发生错误", zap.Any("LoginID", r.LoginID), zap.Any("SessionID", r.SessionID), zap.Any("err", err))
  332. return
  333. }
  334. }
  335. }
  336. // clientToTradeAgentMsg 处理客户端发上来的包
  337. func (r *Client) clientToTradeAgentMsg(msg []byte) error {
  338. var p packet.MiPacket
  339. err := p.UnPackHead(msg)
  340. if err != nil {
  341. // logger.Logger().Errorf("[%v c->s] invalid packet: %v", r.ProxyId, err)
  342. return err
  343. }
  344. // 长度24 是心跳包, 目前只允许客户端上行心跳包
  345. if p.Length == 24 {
  346. // logger.Logger().Debugf("%v [%v c->s] %s", r.ClientIP, r.ProxyId, p.HeaderInfo())
  347. // 将心跳发回给客户端
  348. err = r.wsTradeConn.WriteMessage(websocket.BinaryMessage, msg)
  349. if err != nil {
  350. return err
  351. }
  352. }
  353. return nil
  354. }
  355. // **************** RabbitMQ ****************
  356. // MQPacket 与总线交互的数据体
  357. type MQPacket struct {
  358. FunCode uint32 // 功能码
  359. SessionId uint32 // 数据包的sid
  360. Data *[]byte // 业务数据体
  361. }
  362. // AsyncTask 异步任务结构体
  363. type AsyncTask struct {
  364. PacketRsp chan MQPacket // 总线数据处理通道
  365. FuncodeRsp uint32 // 回复功能码
  366. SerialNumber uint32 // 通信流水号
  367. Own *Client
  368. IsEncrypted bool // 是否加密
  369. Rsp chan rsp.MQBodyRsp // 回调
  370. doClose sync.Once // 仅关闭通道一次
  371. }
  372. // Finish 完成
  373. func (r *AsyncTask) Finish() {
  374. r.doClose.Do(
  375. func() {
  376. close(r.Rsp)
  377. key := fmt.Sprintf("%v_%v_%v", r.Own.SessionID, r.FuncodeRsp, r.SerialNumber)
  378. r.Own.DeleteAsyncTask(key)
  379. })
  380. }
  381. type LoginRedis struct {
  382. LoginID string `json:"loginId" redis:"loginId"` // 登陆账号
  383. UserID string `json:"userId" redis:"userId"` // 用户ID
  384. SessionID string `json:"sessionId" redis:"sessionId"` // 终端sid
  385. Token string `json:"token" redis:"token"` // 令牌
  386. Group string `json:"group" redis:"group"` // 终端分组
  387. Addr string `json:"addr" redis:"addr"` // 客户端地址信息 // FIXME: 由于本服务改用短连,所以每次提交请交请求可能会不一样,后期可判断是否在中间件中进行拦截
  388. OldToken string `json:"-" redis:"-"` // 旧接入Token
  389. }
  390. // FromMap Map to Struct
  391. func (r *LoginRedis) FromMap(val map[string]interface{}) error {
  392. return mapstructure.Decode(val, r)
  393. }
  394. // ToMap Struct to Map
  395. func (r *LoginRedis) ToMap() (val map[string]interface{}, err error) {
  396. if marshalContent, err := json.Marshal(r); err != nil {
  397. return nil, err
  398. } else {
  399. d := json.NewDecoder(bytes.NewReader(marshalContent))
  400. d.UseNumber() // 设置将float64转为一个number
  401. if err := d.Decode(&val); err != nil {
  402. fmt.Println(err)
  403. } else {
  404. for k, v := range val {
  405. val[k] = v
  406. }
  407. }
  408. }
  409. return
  410. }