client.go 9.0 KB

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