client.go 10.0 KB

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