client.go 11 KB

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