sendMsg.go 8.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277
  1. package service
  2. import (
  3. "context"
  4. "database/sql"
  5. "encoding/json"
  6. "fmt"
  7. "go.etcd.io/etcd/client/v3/concurrency"
  8. "log"
  9. "strconv"
  10. "time"
  11. "app.yhyue.com/moapp/MessageCenter/entity"
  12. "app.yhyue.com/moapp/MessageCenter/rpc/message"
  13. )
  14. // 类型的顺序
  15. const order = "1,4"
  16. const EtcdCount = "Count.%s.%s" //etcd 消息未读数量 Count.用户id.消息类型=数量
  17. func SendMsg(this message.SendMsgRequest) (int64, string) {
  18. //orm := entity.Engine.NewSession()
  19. //defer orm.Close()
  20. //err := orm.Begin()
  21. //fmt.Println(err)
  22. count := entity.Mysql.Count("conversation", map[string]interface{}{"receive_id": this.ReceiveUserId, "send_id": this.SendUserId})
  23. sql3 := `INSERT INTO message(appid,receive_userid,receive_name,send_userid,send_name,title,content,msg_type,link,cite_id,createtime,isRead,isdel)
  24. values ("%s",'%s','%s','%s','%s','%s','%s','%d','%s',0,'%s',0,1);`
  25. sql3 = fmt.Sprintf(sql3, this.Appid, this.ReceiveUserId, this.ReceiveName, this.SendUserId, this.SendName, this.Title, this.Content, this.MsgType, this.Link, time.Now().Format("2006-01-02 15:04:05"))
  26. if count < 1 {
  27. sql1 := `INSERT INTO conversation(appid,` + "`key`" + `,user_id,receive_id,receive_name,send_id,send_name,sort,createtime)
  28. values ('%s','','%s','%s','%s','%s','%s',0,'%s');`
  29. sql1 = fmt.Sprintf(sql1, this.Appid, this.SendUserId, this.ReceiveUserId, this.ReceiveName, this.SendUserId, this.SendName, time.Now().Format("2006-01-02 15:04:05"))
  30. ok := entity.Mysql.ExecTx("发送消息事务", func(tx *sql.Tx) bool {
  31. //插入会话表
  32. _, err := entity.Mysql.DB.Exec(sql1)
  33. sql2 := `INSERT INTO conversation(appid,` + "`key`" + `,user_id,receive_id,receive_name,send_id,send_name,sort,createtime)
  34. values ('%s','','%s','%s','%s','%s','%s',0,'%s');`
  35. sql2 = fmt.Sprintf(sql2, this.Appid, this.ReceiveUserId, this.SendUserId, this.SendName, this.ReceiveUserId, this.ReceiveName, time.Now().Format("2006-01-02 15:04:05"))
  36. _, err = entity.Mysql.DB.Exec(sql2)
  37. //插入消息表
  38. _, err = entity.Mysql.DB.Exec(sql3)
  39. if err != nil {
  40. return false
  41. }
  42. return true
  43. })
  44. if ok {
  45. return 1, "消息发送成功"
  46. }
  47. }
  48. _, err := entity.Mysql.DB.Exec(sql3)
  49. if err == nil {
  50. go EtcdCountAdd(this.ReceiveUserId, strconv.Itoa(int(this.MsgType)))
  51. return 1, "消息发送成功"
  52. }
  53. return 0, "消息发送失败"
  54. }
  55. func FindUserMsg(this message.FindUserMsgReq) message.FindUserMsgRes {
  56. orm := entity.Engine
  57. var messages []*entity.Message
  58. var err error
  59. var count int64
  60. q := ""
  61. if this.MsgType != -1 {
  62. q += fmt.Sprintf(" and msg_type = %d", this.MsgType)
  63. }
  64. if this.Read != -1 {
  65. q += fmt.Sprintf(" and isRead = %d", this.Read)
  66. }
  67. count, err = orm.Table("message").Where("((receive_userid = ? and send_userid = ?) or (receive_userid = ? and send_userid = ?)) and isdel = ? and appid = ?"+q, this.UserId, this.ReceiveUserId, this.ReceiveUserId, this.UserId, 1, this.Appid).Count()
  68. data := message.FindUserMsgRes{}
  69. if count > 0 {
  70. err = orm.Table("message").Select("*").Where("((receive_userid = ? and send_userid = ?) or (receive_userid = ? and send_userid = ?)) and isdel = ? and appid = ?"+q, this.UserId, this.ReceiveUserId, this.ReceiveUserId, this.UserId, 1, this.Appid).
  71. OrderBy("createtime desc").
  72. Limit(int(this.PageSize), (int(this.OffSet)-1)*int(this.PageSize)).
  73. Find(&messages)
  74. //log.Println("数据:", messages)
  75. for _, v := range messages {
  76. data.Data = append(data.Data, &message.Messages{
  77. Id: v.Id,
  78. Appid: v.AppId,
  79. ReceiveUserId: v.ReceiveUserid,
  80. ReceiveName: v.ReceiveName,
  81. SendUserId: v.SendUserid,
  82. SendName: v.SendName,
  83. Createtime: v.CreateTime.Format("2006-01-02 15:04:05"),
  84. Title: v.Title,
  85. MsgType: int64(v.MsgType),
  86. Link: v.Link,
  87. CiteId: int64(v.CiteId),
  88. Content: v.Content,
  89. IsRead: int64(v.IsRead),
  90. })
  91. }
  92. }
  93. data.Count = count
  94. if err != nil {
  95. data.Code = 0
  96. data.Message = "查询失败"
  97. } else {
  98. data.Code = 1
  99. data.Message = "查询成功"
  100. }
  101. return data
  102. }
  103. // 指定分类未读消息合计
  104. func ClassCountUnread(msgType int, userId string, appId string) (int64, string, int64) {
  105. //orm := entity.Engine
  106. //count, err := orm.Table("message").Where("msg_type=? and receive_userid=? and isdel=1 and appid=? and isRead=0", msgType, userId, appId).Count()
  107. query := map[string]interface{}{
  108. "msg_type": msgType,
  109. "receive_userid": userId,
  110. "isdel": 1,
  111. "appid": appId,
  112. "isRead": 0,
  113. }
  114. count := entity.Mysql.Count("message", query)
  115. return 1, "查询指定分类未读消息成功", count
  116. }
  117. // etcd数量加1
  118. func EtcdCountAdd(userId, msgType string) {
  119. s1, err := concurrency.NewSession(entity.EtcdCli)
  120. if err != nil {
  121. log.Fatal(err)
  122. }
  123. defer s1.Close()
  124. keyString := fmt.Sprintf(EtcdCount, userId, msgType)
  125. m1 := concurrency.NewMutex(s1, keyString)
  126. // 会话s1获取锁
  127. if err := m1.Lock(context.TODO()); err != nil {
  128. log.Fatal("获取锁失败", err)
  129. }
  130. // 操作数量
  131. ctx, cancel := context.WithTimeout(context.Background(), time.Second)
  132. resp, err := s1.Client().Get(ctx, keyString)
  133. cancel()
  134. if err != nil {
  135. fmt.Printf("get from etcd failed, err:%v\n", err)
  136. // 释放锁
  137. if err := m1.Unlock(context.TODO()); err != nil {
  138. log.Fatal("释放锁失败", err)
  139. }
  140. return
  141. }
  142. var count int
  143. for _, ev := range resp.Kvs {
  144. if ev.Value != nil {
  145. err := json.Unmarshal([]byte(ev.Value), &count)
  146. if err != nil {
  147. log.Println("etcd get err:", err)
  148. } else {
  149. break
  150. }
  151. }
  152. }
  153. ctx, cancel = context.WithTimeout(context.Background(), time.Second)
  154. _, err = s1.Client().Put(ctx, keyString, strconv.Itoa(count+1))
  155. cancel()
  156. if err != nil {
  157. fmt.Printf("put to etcd failed, err:%v\n", err)
  158. // 释放锁
  159. if err := m1.Unlock(context.TODO()); err != nil {
  160. log.Fatal("释放锁失败", err)
  161. }
  162. return
  163. }
  164. // 释放锁
  165. if err := m1.Unlock(context.TODO()); err != nil {
  166. log.Fatal("释放锁失败", err)
  167. }
  168. }
  169. //单条消息,-1
  170. func EtcdCountMinusOne(userId, msgType string) {
  171. s1, err := concurrency.NewSession(entity.EtcdCli)
  172. if err != nil {
  173. log.Fatal(err)
  174. }
  175. defer s1.Close()
  176. keyString := fmt.Sprintf(EtcdCount, userId, msgType)
  177. m1 := concurrency.NewMutex(s1, keyString)
  178. // 会话s1获取锁
  179. if err := m1.Lock(context.TODO()); err != nil {
  180. log.Fatal("获取锁失败", err)
  181. }
  182. // 操作数量
  183. ctx, cancel := context.WithTimeout(context.Background(), time.Second)
  184. resp, err := s1.Client().Get(ctx, keyString)
  185. cancel()
  186. if err != nil {
  187. fmt.Printf("get from etcd failed, err:%v\n", err)
  188. // 释放锁
  189. if err := m1.Unlock(context.TODO()); err != nil {
  190. log.Fatal("释放锁失败", err)
  191. }
  192. return
  193. }
  194. var count int
  195. for _, ev := range resp.Kvs {
  196. if ev.Value != nil {
  197. err := json.Unmarshal([]byte(ev.Value), &count)
  198. if err != nil {
  199. log.Println("etcd get err:", err)
  200. } else {
  201. break
  202. }
  203. }
  204. }
  205. ctx, cancel = context.WithTimeout(context.Background(), time.Second)
  206. var fin string
  207. if count > 0 {
  208. fin = strconv.Itoa(count - 1)
  209. } else {
  210. fin = "0"
  211. }
  212. _, err = s1.Client().Put(ctx, keyString, fin)
  213. cancel()
  214. if err != nil {
  215. fmt.Printf("put to etcd failed, err:%v\n", err)
  216. // 释放锁
  217. if err := m1.Unlock(context.TODO()); err != nil {
  218. log.Fatal("释放锁失败", err)
  219. }
  220. return
  221. }
  222. // 释放锁
  223. if err := m1.Unlock(context.TODO()); err != nil {
  224. log.Fatal("释放锁失败", err)
  225. }
  226. }
  227. // 消息类别置0
  228. func EtcdSetCountZero(userId, msgType string) {
  229. log.Println(" 消息类别置0", userId, msgType)
  230. s1, err := concurrency.NewSession(entity.EtcdCli)
  231. if err != nil {
  232. log.Fatal(err)
  233. }
  234. defer s1.Close()
  235. keyString := fmt.Sprintf(EtcdCount, userId, msgType)
  236. m1 := concurrency.NewMutex(s1, keyString)
  237. // 会话s1获取锁
  238. if err := m1.Lock(context.TODO()); err != nil {
  239. log.Fatal("获取锁失败", err)
  240. }
  241. // 操作数量
  242. ctx, cancel := context.WithTimeout(context.Background(), time.Second)
  243. _, err = s1.Client().Put(ctx, keyString, "0")
  244. cancel()
  245. if err != nil {
  246. fmt.Printf("put to etcd failed, err:%v\n", err)
  247. // 释放锁
  248. if err := m1.Unlock(context.TODO()); err != nil {
  249. log.Fatal("释放锁失败", err)
  250. }
  251. return
  252. }
  253. // 释放锁
  254. if err := m1.Unlock(context.TODO()); err != nil {
  255. log.Fatal("释放锁失败", err)
  256. }
  257. }