client.go 3.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135
  1. package global
  2. import (
  3. "bytes"
  4. "encoding/json"
  5. "fmt"
  6. rsp "mtp20access/model/mq/response"
  7. "sync"
  8. "github.com/mitchellh/mapstructure"
  9. )
  10. type LoginRedis struct {
  11. LoginID string `json:"loginId" redis:"loginId"` // 登陆账号
  12. UserID string `json:"userId" redis:"userId"` // 用户ID
  13. SessionID string `json:"sessionId" redis:"sessionId"` // 终端sid
  14. Token string `json:"token" redis:"token"` // 令牌
  15. Group string `json:"group" redis:"group"` // 终端分组
  16. Addr string `json:"addr" redis:"addr"` // 客户端地址信息 // FIXME: 由于本服务改用短连,所以每次提交请交请求可能会不一样,后期可判断是否在中间件中进行拦截
  17. }
  18. // FromMap Map to Struct
  19. func (r *LoginRedis) FromMap(val map[string]interface{}) error {
  20. return mapstructure.Decode(val, r)
  21. }
  22. // ToMap Struct to Map
  23. func (r *LoginRedis) ToMap() (val map[string]interface{}, err error) {
  24. if marshalContent, err := json.Marshal(r); err != nil {
  25. return nil, err
  26. } else {
  27. d := json.NewDecoder(bytes.NewReader(marshalContent))
  28. d.UseNumber() // 设置将float64转为一个number
  29. if err := d.Decode(&val); err != nil {
  30. fmt.Println(err)
  31. } else {
  32. for k, v := range val {
  33. val[k] = v
  34. }
  35. }
  36. }
  37. return
  38. }
  39. type Client struct {
  40. LoginRedis
  41. mtx sync.RWMutex
  42. curSerialNumber uint32 // 当前业务流水号
  43. asyncTasks map[string]*AsyncTask // key:SessionId_FuncodeRsp_SerialNumber
  44. }
  45. // GetSerialNumber 获取可用流水号
  46. func (r *Client) GetSerialNumber() uint32 {
  47. r.mtx.Lock()
  48. defer func() {
  49. r.mtx.Unlock()
  50. }()
  51. r.curSerialNumber += 1
  52. return r.curSerialNumber
  53. }
  54. // GetAsyncTask 获取目标异步任务
  55. // key:SessionId_FuncodeRsp
  56. func (r *Client) GetAsyncTask(key string) (asyncTask *AsyncTask) {
  57. r.mtx.RLock()
  58. defer func() {
  59. r.mtx.RUnlock()
  60. }()
  61. asyncTask, ok := r.asyncTasks[key]
  62. if ok {
  63. return asyncTask
  64. } else {
  65. return nil
  66. }
  67. }
  68. // SetAsyncTask 设置异步任务
  69. func (r *Client) SetAsyncTask(asyncTask *AsyncTask, key string) {
  70. r.mtx.Lock()
  71. defer func() {
  72. r.mtx.Unlock()
  73. }()
  74. if r.asyncTasks == nil {
  75. r.asyncTasks = make(map[string]*AsyncTask, 0)
  76. }
  77. delete(r.asyncTasks, key)
  78. r.asyncTasks[key] = asyncTask
  79. }
  80. // DeleteAsyncTask 删除异步任务
  81. func (r *Client) DeleteAsyncTask(key string) {
  82. r.mtx.Lock()
  83. defer func() {
  84. r.mtx.Unlock()
  85. }()
  86. delete(r.asyncTasks, key)
  87. }
  88. func (r *Client) GetAllAsyncTask() *map[string]*AsyncTask {
  89. return &r.asyncTasks
  90. }
  91. // MQPacket 与总线交互的数据体
  92. type MQPacket struct {
  93. FunCode uint32 // 功能码
  94. SessionId uint32 // 数据包的sid
  95. Data *[]byte // 业务数据体
  96. }
  97. // AsyncTask 异步任务结构体
  98. type AsyncTask struct {
  99. PacketRsp chan MQPacket // 总线数据处理通道
  100. FuncodeRsp uint32 // 回复功能码
  101. SerialNumber uint32 // 通信流水号
  102. Own *Client
  103. IsEncrypted bool // 是否加密
  104. Rsp chan rsp.MQBodyRsp // 回调
  105. doClose sync.Once // 仅关闭通道一次
  106. }
  107. // Finish 完成
  108. func (r *AsyncTask) Finish() {
  109. r.doClose.Do(
  110. func() {
  111. close(r.Rsp)
  112. key := fmt.Sprintf("%v_%v_%v", r.Own.SessionID, r.FuncodeRsp, r.SerialNumber)
  113. r.Own.DeleteAsyncTask(key)
  114. })
  115. }