client.go 6.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292
  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. mtx sync.RWMutex
  20. curSerialNumber uint32 // 当前业务流水号
  21. asyncTasks map[string]*AsyncTask // key:SessionId_FuncodeRsp_SerialNumber
  22. wsConn *websocket.Conn // 终端WebSocket连接
  23. quoteSubs []request.QuoteGoods // 当前已订阅行情的商品
  24. ch chan interface{} // 接收实时行情订阅channel
  25. unSubscribe chan struct{} // 取消订阅信号
  26. writeChan chan []byte // 推送队列 QuoteServer -> Client
  27. wsCloseChan chan struct{} // 终端WebSocket连接关闭信号
  28. }
  29. // GetSerialNumber 获取可用流水号
  30. func (r *Client) GetSerialNumber() uint32 {
  31. r.mtx.Lock()
  32. defer func() {
  33. r.mtx.Unlock()
  34. }()
  35. r.curSerialNumber += 1
  36. return r.curSerialNumber
  37. }
  38. // SetQuoteSubs 设置商品行情订阅信息
  39. func (r *Client) SetQuoteSubs(req request.QuoteSubscribeReq) {
  40. r.mtx.Lock()
  41. defer r.mtx.Unlock()
  42. r.quoteSubs = req.QuoteGoodses
  43. if r.ch != nil {
  44. // r.unSubscribe <- struct{}{}
  45. global.M2A_Publish.Unsubscribe(publish.Topic_Quote, r.ch)
  46. }
  47. r.ch = global.M2A_Publish.Subscribe(publish.Topic_Quote)
  48. go func() {
  49. for {
  50. select {
  51. case msg, ok := <-r.ch:
  52. if !ok {
  53. return // 管道已关闭,退出循环
  54. }
  55. // 向客户端发送行情信息
  56. if p, ok := msg.(*packet.MiQuotePacket); ok {
  57. DispatchRealQuote(p, r)
  58. }
  59. case <-r.unSubscribe:
  60. // 取消订阅信息
  61. return
  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) WriteWsBuf(buf []byte) (err error) {
  74. r.writeChan <- buf
  75. return
  76. }
  77. // GetAsyncTask 获取目标异步任务
  78. // key:SessionId_FuncodeRsp
  79. func (r *Client) GetAsyncTask(key string) (asyncTask *AsyncTask) {
  80. r.mtx.RLock()
  81. defer func() {
  82. r.mtx.RUnlock()
  83. }()
  84. asyncTask, ok := r.asyncTasks[key]
  85. if ok {
  86. return asyncTask
  87. } else {
  88. return nil
  89. }
  90. }
  91. // SetAsyncTask 设置异步任务
  92. func (r *Client) SetAsyncTask(asyncTask *AsyncTask, key string) {
  93. r.mtx.Lock()
  94. defer func() {
  95. r.mtx.Unlock()
  96. }()
  97. if r.asyncTasks == nil {
  98. r.asyncTasks = make(map[string]*AsyncTask, 0)
  99. }
  100. delete(r.asyncTasks, key)
  101. r.asyncTasks[key] = asyncTask
  102. }
  103. // DeleteAsyncTask 删除异步任务
  104. func (r *Client) DeleteAsyncTask(key string) {
  105. r.mtx.Lock()
  106. defer func() {
  107. r.mtx.Unlock()
  108. }()
  109. delete(r.asyncTasks, key)
  110. }
  111. func (r *Client) GetAllAsyncTask() *map[string]*AsyncTask {
  112. return &r.asyncTasks
  113. }
  114. func (r *Client) SetWebSocket(ws *websocket.Conn) (err error) {
  115. r.mtx.Lock()
  116. defer r.mtx.Unlock()
  117. if r.wsConn != nil {
  118. r.wsConn.Close()
  119. }
  120. r.wsConn = ws
  121. r.writeChan = make(chan []byte, 100)
  122. // 开始读取客户端发送信息
  123. go r.readClientWsMessage()
  124. // 开始推送客户端信息循环
  125. go r.writeClientWsMessage()
  126. r.wsConn.SetCloseHandler(func(code int, text string) error {
  127. close(r.wsCloseChan)
  128. return nil
  129. })
  130. return
  131. }
  132. // readClientWsMessage 处理终端发过来的websocket数据
  133. // 注意: 阻塞式, 直到websocket关闭才退出
  134. func (r *Client) readClientWsMessage() {
  135. for {
  136. // 40秒心跳超时
  137. r.wsConn.SetReadDeadline(time.Now().Add(40 * time.Second))
  138. mt, msg, err := r.wsConn.ReadMessage()
  139. if err != nil {
  140. fmt.Println(err)
  141. break
  142. }
  143. switch mt {
  144. case websocket.PingMessage:
  145. _ = r.wsConn.WriteMessage(mt, msg)
  146. case websocket.CloseMessage:
  147. return
  148. case websocket.BinaryMessage:
  149. if err := r.clientToQuoteAgentMsg(msg); err != nil {
  150. return
  151. }
  152. case websocket.TextMessage:
  153. fmt.Println(string(msg))
  154. _ = r.wsConn.WriteMessage(mt, msg)
  155. }
  156. select {
  157. case <-r.wsCloseChan:
  158. return
  159. }
  160. }
  161. }
  162. // writeClientWsMessage 由于websocket非线程安全,
  163. // 所以由统一协程写入
  164. func (r *Client) writeClientWsMessage() {
  165. // defer r.close()
  166. for {
  167. select {
  168. case buf := <-r.writeChan:
  169. err := r.wsConn.WriteMessage(websocket.BinaryMessage, buf)
  170. if err != nil {
  171. return
  172. }
  173. case <-r.wsCloseChan: // 与终端连接关闭信息
  174. return
  175. }
  176. }
  177. }
  178. // clientToQuoteAgentMsg 处理客户端发上来的包
  179. func (r *Client) clientToQuoteAgentMsg(msg []byte) error {
  180. var p packet.MiQuotePacket
  181. err := p.UnPackHead(msg)
  182. if err != nil {
  183. // logger.Logger().Errorf("[%v c->s] invalid packet: %v", r.ProxyId, err)
  184. return err
  185. }
  186. // 长度15 是心跳包, 目前只允许客户端上行心跳包
  187. if p.Length == 15 {
  188. // logger.Logger().Debugf("%v [%v c->s] %s", r.ClientIP, r.ProxyId, p.HeaderInfo())
  189. // 将心跳发回给客户端
  190. err = r.wsConn.WriteMessage(websocket.BinaryMessage, msg)
  191. if err != nil {
  192. return err
  193. }
  194. }
  195. return nil
  196. }
  197. // MQPacket 与总线交互的数据体
  198. type MQPacket struct {
  199. FunCode uint32 // 功能码
  200. SessionId uint32 // 数据包的sid
  201. Data *[]byte // 业务数据体
  202. }
  203. // AsyncTask 异步任务结构体
  204. type AsyncTask struct {
  205. PacketRsp chan MQPacket // 总线数据处理通道
  206. FuncodeRsp uint32 // 回复功能码
  207. SerialNumber uint32 // 通信流水号
  208. Own *Client
  209. IsEncrypted bool // 是否加密
  210. Rsp chan rsp.MQBodyRsp // 回调
  211. doClose sync.Once // 仅关闭通道一次
  212. }
  213. // Finish 完成
  214. func (r *AsyncTask) Finish() {
  215. r.doClose.Do(
  216. func() {
  217. close(r.Rsp)
  218. key := fmt.Sprintf("%v_%v_%v", r.Own.SessionID, r.FuncodeRsp, r.SerialNumber)
  219. r.Own.DeleteAsyncTask(key)
  220. })
  221. }
  222. type LoginRedis struct {
  223. LoginID string `json:"loginId" redis:"loginId"` // 登陆账号
  224. UserID string `json:"userId" redis:"userId"` // 用户ID
  225. SessionID string `json:"sessionId" redis:"sessionId"` // 终端sid
  226. Token string `json:"token" redis:"token"` // 令牌
  227. Group string `json:"group" redis:"group"` // 终端分组
  228. Addr string `json:"addr" redis:"addr"` // 客户端地址信息 // FIXME: 由于本服务改用短连,所以每次提交请交请求可能会不一样,后期可判断是否在中间件中进行拦截
  229. }
  230. // FromMap Map to Struct
  231. func (r *LoginRedis) FromMap(val map[string]interface{}) error {
  232. return mapstructure.Decode(val, r)
  233. }
  234. // ToMap Struct to Map
  235. func (r *LoginRedis) ToMap() (val map[string]interface{}, err error) {
  236. if marshalContent, err := json.Marshal(r); err != nil {
  237. return nil, err
  238. } else {
  239. d := json.NewDecoder(bytes.NewReader(marshalContent))
  240. d.UseNumber() // 设置将float64转为一个number
  241. if err := d.Decode(&val); err != nil {
  242. fmt.Println(err)
  243. } else {
  244. for k, v := range val {
  245. val[k] = v
  246. }
  247. }
  248. }
  249. return
  250. }