sendMsg.go 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358
  1. package service
  2. import (
  3. "context"
  4. "database/sql"
  5. "encoding/json"
  6. "fmt"
  7. "log"
  8. "strconv"
  9. "strings"
  10. "time"
  11. "app.yhyue.com/moapp/MessageCenter/util"
  12. "go.etcd.io/etcd/client/v3/concurrency"
  13. "app.yhyue.com/moapp/MessageCenter/entity"
  14. "app.yhyue.com/moapp/MessageCenter/rpc/message"
  15. )
  16. // 类型的顺序
  17. const order = "1,4"
  18. const EtcdCount = "Count.%s.%s" //etcd 消息未读数量 Count.用户id.消息类型=数量
  19. func SendMsg(this message.SendMsgRequest) (int64, string) {
  20. //orm := entity.Engine.NewSession()
  21. //defer orm.Close()
  22. //err := orm.Begin()
  23. //fmt.Println(err)
  24. //count := entity.Mysql11.Query("conversation", map[string]interface{}{"receive_id": this.ReceiveUserId, "send_id": this.SendUserId})
  25. r, err := entity.Mysql11.Query("select count(*) as c from conversation where receive_id = ? and send_id = ? ", this.ReceiveUserId, this.SendUserId)
  26. c := 0
  27. log.Println("查询结果", r)
  28. for r.Next() {
  29. err := r.Scan(&c)
  30. if err != nil {
  31. panic(err.Error())
  32. }
  33. }
  34. log.Println("查询数量:", c)
  35. sql3 := `INSERT INTO message(appid,receive_userid,receive_name,send_userid,send_name,title,content,msg_type,link,cite_id,createtime,isRead,isdel)
  36. values ("%s",'%s','%s','%s','%s','%s','%s','%d','%s',0,'%s',0,1);`
  37. 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"))
  38. if c < 1 {
  39. sql1 := `INSERT INTO conversation(appid,` + "`key`" + `,user_id,receive_id,receive_name,send_id,send_name,sort,createtime)
  40. values ('%s','','%s','%s','%s','%s','%s',0,'%s');`
  41. 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"))
  42. ok := entity.Mysql.ExecTx("发送消息事务", func(tx *sql.Tx) bool {
  43. //插入会话表
  44. _, err := entity.Mysql11.Exec(sql1)
  45. sql2 := `INSERT INTO conversation(appid,` + "`key`" + `,user_id,receive_id,receive_name,send_id,send_name,sort,createtime)
  46. values ('%s','','%s','%s','%s','%s','%s',0,'%s');`
  47. 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"))
  48. _, err = entity.Mysql11.Exec(sql2)
  49. //插入消息表
  50. _, err = entity.Mysql11.Exec(sql3)
  51. if err != nil {
  52. return false
  53. }
  54. return true
  55. })
  56. if ok {
  57. return 1, "消息发送成功"
  58. }
  59. }
  60. _, err = entity.Mysql11.Exec(sql3)
  61. if err == nil {
  62. EtcdCountAdd(this.ReceiveUserId, strconv.Itoa(int(this.MsgType)))
  63. return 1, "消息发送成功"
  64. }
  65. return 0, "消息发送失败"
  66. }
  67. func FindUserMsg(this message.FindUserMsgReq) message.FindUserMsgRes {
  68. //orm := entity.Engine
  69. //var messages []*entity.Message
  70. var err error
  71. var count int64
  72. cquery := map[string]interface{}{
  73. "receive_userid": this.UserId,
  74. "isdel": 1,
  75. "appid": this.Appid,
  76. }
  77. if this.MsgType != -1 {
  78. cquery["msg_type"] = this.MsgType
  79. }
  80. if this.Read != -1 {
  81. cquery["isRead"] = this.Read
  82. }
  83. count = entity.Mysql.Count("message", cquery)
  84. //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()
  85. data := message.FindUserMsgRes{}
  86. if count > 0 {
  87. /*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).
  88. OrderBy("createtime desc").
  89. Limit(int(this.PageSize), (int(this.OffSet)-1)*int(this.PageSize)).
  90. Find(&messages)*/
  91. res := entity.Mysql.Find("message", cquery, "", "createtime desc", (int(this.OffSet)-1)*int(this.PageSize), int(this.PageSize))
  92. log.Println("数据:", res)
  93. if res != nil && len(*res) > 0 {
  94. for _, v := range *res {
  95. _id := util.Int64All(v["id"])
  96. id := strconv.FormatInt(_id, 10)
  97. data.Data = append(data.Data, &message.Messages{
  98. Id: id,
  99. Appid: util.ObjToString(v["appId"]),
  100. ReceiveUserId: util.ObjToString(v["receive_userid"]),
  101. ReceiveName: util.ObjToString(v["receive_name"]),
  102. SendUserId: util.ObjToString(v["send_userid"]),
  103. SendName: util.ObjToString(v["send_name"]),
  104. Createtime: util.ObjToString(v["createtime"]),
  105. Title: util.ObjToString(v["title"]),
  106. MsgType: int64(util.IntAll(v["msg_type"])),
  107. Link: util.ObjToString(v["link"]),
  108. CiteId: util.Int64All(v["cite_id"]),
  109. Content: util.ObjToString(v["content"]),
  110. IsRead: util.Int64All(v["isRead"]),
  111. MsgLogId: util.Int64All(v["msg_log_id"]),
  112. })
  113. }
  114. }
  115. }
  116. data.Count = count
  117. if err != nil {
  118. data.Code = 0
  119. data.Message = "查询失败"
  120. } else {
  121. data.Code = 1
  122. data.Message = "查询成功"
  123. }
  124. return data
  125. }
  126. // 指定分类未读消息合计
  127. func ClassCountUnread(msgType int, userId string, appId string) (int64, string, int64) {
  128. //orm := entity.Engine
  129. //count, err := orm.Table("message").Where("msg_type=? and receive_userid=? and isdel=1 and appid=? and isRead=0", msgType, userId, appId).Count()
  130. query := map[string]interface{}{
  131. "msg_type": msgType,
  132. "receive_userid": userId,
  133. "isdel": 1,
  134. "appid": appId,
  135. "isRead": 0,
  136. }
  137. count := entity.Mysql.Count("message", query)
  138. return 1, "查询指定分类未读消息成功", count
  139. }
  140. // etcd数量加1
  141. func EtcdCountAdd(userId, msgType string) {
  142. s1, err := concurrency.NewSession(entity.EtcdCli)
  143. if err != nil {
  144. log.Println(err)
  145. }
  146. defer s1.Close()
  147. keyString := fmt.Sprintf(EtcdCount, userId, msgType)
  148. m1 := concurrency.NewMutex(s1, keyString)
  149. // 会话s1获取锁
  150. if err := m1.Lock(context.TODO()); err != nil {
  151. log.Println("EtcdCountAdd获取锁失败", err)
  152. }
  153. // 操作数量
  154. ctx, cancel := context.WithTimeout(context.Background(), time.Second)
  155. resp, err := s1.Client().Get(ctx, keyString)
  156. cancel()
  157. if err != nil {
  158. fmt.Printf("get from etcd failed, err:%v\n", err)
  159. // 释放锁
  160. if err := m1.Unlock(context.TODO()); err != nil {
  161. log.Println("EtcdCountAdd释放锁失败", err)
  162. }
  163. return
  164. }
  165. var count int
  166. for _, ev := range resp.Kvs {
  167. if ev.Value != nil {
  168. err := json.Unmarshal([]byte(ev.Value), &count)
  169. if err != nil {
  170. log.Println("etcd get err:", err)
  171. } else {
  172. break
  173. }
  174. }
  175. }
  176. ctx, cancel = context.WithTimeout(context.Background(), time.Second)
  177. _, err = s1.Client().Put(ctx, keyString, strconv.Itoa(count+1))
  178. cancel()
  179. if err != nil {
  180. fmt.Printf("put to etcd failed, err:%v\n", err)
  181. // 释放锁
  182. if err := m1.Unlock(context.TODO()); err != nil {
  183. log.Println("EtcdCountAdd释放锁失败2", err)
  184. }
  185. return
  186. }
  187. // 释放锁
  188. if err := m1.Unlock(context.TODO()); err != nil {
  189. log.Println("EtcdCountAdd释放锁失败3", err)
  190. }
  191. }
  192. //单条消息,-1
  193. func EtcdCountMinusOne(userId, msgType string) {
  194. s1, err := concurrency.NewSession(entity.EtcdCli)
  195. if err != nil {
  196. log.Println(err)
  197. }
  198. defer s1.Close()
  199. keyString := fmt.Sprintf(EtcdCount, userId, msgType)
  200. m1 := concurrency.NewMutex(s1, keyString)
  201. // 会话s1获取锁
  202. if err := m1.Lock(context.TODO()); err != nil {
  203. log.Println("EtcdCountMinusOne获取锁失败", err)
  204. }
  205. // 操作数量
  206. ctx, cancel := context.WithTimeout(context.Background(), time.Second)
  207. resp, err := s1.Client().Get(ctx, keyString)
  208. cancel()
  209. if err != nil {
  210. fmt.Printf("get from etcd failed, err:%v\n", err)
  211. // 释放锁
  212. if err := m1.Unlock(context.TODO()); err != nil {
  213. log.Println("EtcdCountMinusOne释放锁失败", err)
  214. }
  215. return
  216. }
  217. var count int
  218. for _, ev := range resp.Kvs {
  219. if ev.Value != nil {
  220. err := json.Unmarshal([]byte(ev.Value), &count)
  221. if err != nil {
  222. log.Println("etcd get err:", err)
  223. } else {
  224. break
  225. }
  226. }
  227. }
  228. ctx, cancel = context.WithTimeout(context.Background(), time.Second)
  229. var fin string
  230. if count > 0 {
  231. fin = strconv.Itoa(count - 1)
  232. } else {
  233. fin = "0"
  234. }
  235. _, err = s1.Client().Put(ctx, keyString, fin)
  236. cancel()
  237. if err != nil {
  238. fmt.Printf("put to etcd failed, err:%v\n", err)
  239. // 释放锁
  240. if err := m1.Unlock(context.TODO()); err != nil {
  241. log.Println("EtcdCountMinusOne释放锁失败2", err)
  242. }
  243. return
  244. }
  245. // 释放锁
  246. if err := m1.Unlock(context.TODO()); err != nil {
  247. log.Println("EtcdCountMinusOne释放锁失败3", err)
  248. }
  249. }
  250. // 消息类别置0
  251. func EtcdSetCountZero(userId, msgType string) {
  252. log.Println(" 消息类别置0", userId, msgType)
  253. s1, err := concurrency.NewSession(entity.EtcdCli)
  254. if err != nil {
  255. log.Println(err)
  256. }
  257. defer s1.Close()
  258. keyString := fmt.Sprintf(EtcdCount, userId, msgType)
  259. m1 := concurrency.NewMutex(s1, keyString)
  260. // 会话s1获取锁
  261. if err := m1.Lock(context.TODO()); err != nil {
  262. log.Println("EtcdSetCountZero获取锁失败", err)
  263. }
  264. // 操作数量
  265. ctx, cancel := context.WithTimeout(context.Background(), time.Second)
  266. _, err = s1.Client().Put(ctx, keyString, "0")
  267. cancel()
  268. if err != nil {
  269. fmt.Printf("put to etcd failed, err:%v\n", err)
  270. // 释放锁
  271. if err := m1.Unlock(context.TODO()); err != nil {
  272. log.Println("EtcdSetCountZero释放锁失败", err)
  273. }
  274. return
  275. }
  276. // 释放锁
  277. if err := m1.Unlock(context.TODO()); err != nil {
  278. log.Println("EtcdSetCountZero释放锁失败", err)
  279. }
  280. }
  281. func MultSave(this message.MultipleSaveMsgReq) (int64, string) {
  282. userIdArr := strings.Split(this.UserIds, ",")
  283. userNameArr := strings.Split(this.UserNames, ",")
  284. log.Println("参数:", len(userIdArr), len(userNameArr))
  285. if len(userIdArr) > 0 {
  286. var errCount int64
  287. for k, v := range userIdArr {
  288. if v == "" {
  289. return errCount, "调用结束"
  290. }
  291. userName := userNameArr[k]
  292. //消息数组
  293. bT := time.Now()
  294. c := entity.Mysql.Count("conversation", map[string]interface{}{"receive_id": v, "send_id": this.SendUserId})
  295. // log.Println("count", c)
  296. //m := entity.Mysql.SelectBySql("select count(id) from conversation where receive_id=? and send_id=?", v, this.SendUserId)
  297. sql3 := `INSERT INTO message(appid,receive_userid,receive_name,send_userid,send_name,title,content,msg_type,link,cite_id,createtime,isRead,isdel,msg_log_id) values ("%s",'%s','%s','%s','%s','%s','%s',%d,'%s',0,'%s',0,1,%d);`
  298. sql3 = fmt.Sprintf(sql3, this.Appid, v, userName, this.SendUserId, this.SendName, this.Title, this.Content, this.MsgType, this.Link, time.Now().Format("2006-01-02 15:04:05"), this.MsgLogId)
  299. if c <= 0 {
  300. sql1 := `INSERT INTO conversation(appid,secret_key,user_id,receive_id,receive_name,send_id,send_name,sort,createtime) values ('%s','','%s','%s','%s','%s','%s',0,'%s');`
  301. sql1 = fmt.Sprintf(sql1, this.Appid, this.SendUserId, v, userName, this.SendUserId, this.SendName, time.Now().Format("2006-01-02 15:04:05"))
  302. ok := entity.Mysql.ExecTx("发送消息事务", func(tx *sql.Tx) bool {
  303. //插入会话表
  304. //开始时间
  305. in1 := entity.Mysql.InsertBySqlByTx(tx, sql1)
  306. sql2 := `INSERT INTO conversation(appid,secret_key,user_id,receive_id,receive_name,send_id,send_name,sort,createtime) values ('%s','','%s','%s','%s','%s','%s',0,'%s');`
  307. sql2 = fmt.Sprintf(sql2, this.Appid, v, this.SendUserId, this.SendName, v, userName, time.Now().Format("2006-01-02 15:04:05"))
  308. in2 := entity.Mysql.InsertBySqlByTx(tx, sql2)
  309. //插入消息表
  310. in3 := entity.Mysql.InsertBySqlByTx(tx, sql3)
  311. log.Println("存储耗时:", time.Since(bT))
  312. log.Println(in1, in2, in3)
  313. return in1 > -1 && in2 > -1 && in3 > -1
  314. })
  315. log.Println("执行事务是否成功:", ok)
  316. if !ok {
  317. errCount++
  318. continue
  319. }
  320. EtcdCountAdd(v, strconv.Itoa(int(this.MsgType)))
  321. } else {
  322. in := entity.Mysql.InsertBySql(sql3)
  323. if in > -1 {
  324. EtcdCountAdd(v, strconv.Itoa(int(this.MsgType)))
  325. } else {
  326. errCount++
  327. }
  328. }
  329. }
  330. return errCount, "发送成功"
  331. }
  332. return 0, "没有要发送的用户"
  333. }