subscribe.go 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199
  1. package sender
  2. import (
  3. "dsbqj-admin/model/mongo/subscribe"
  4. "dsbqj-admin/pkg/helper/wechat"
  5. "dsbqj-admin/pkg/logger"
  6. "github.com/kamva/mgm/v3"
  7. "go.mongodb.org/mongo-driver/bson"
  8. yaml "gopkg.in/yaml.v2"
  9. "log"
  10. "os"
  11. "strings"
  12. "sync"
  13. )
  14. const subscribeTemplateConfigFile = "config/subscribe_template.yaml"
  15. var (
  16. subscribeTemplateOnce sync.Once
  17. subscribeTemplateMap map[string]map[string]string
  18. subscribeTemplateLoadErr error
  19. )
  20. type SubscribeSend struct {
  21. DeviceId string
  22. OpenIds []string
  23. Module string
  24. }
  25. type SubscribeSender struct {
  26. workerNum int
  27. queue chan *SubscribeSend
  28. wxHelper *wechat.WechatHelper
  29. }
  30. func NewSubscribeSender(workerNum int) *SubscribeSender {
  31. return &SubscribeSender{
  32. workerNum: workerNum,
  33. queue: make(chan *SubscribeSend, 1000),
  34. wxHelper: wechat.NewWechatHelper(),
  35. }
  36. }
  37. func (this *SubscribeSender) Send(send *SubscribeSend) {
  38. select {
  39. case this.queue <- send:
  40. default:
  41. logger.Info("subscribe sender queue full, drop task=%+v", send)
  42. }
  43. }
  44. func (this *SubscribeSender) Start() {
  45. for i := 0; i < this.workerNum; i++ {
  46. go this.worker(i)
  47. }
  48. }
  49. func (this *SubscribeSender) worker(id int) {
  50. for task := range this.queue {
  51. logger.Info("Subscribe worker(%d) do safe send module %s, openid_len %d device id %s", id, task.Module, len(task.OpenIds), task.DeviceId)
  52. this.safeSend(task)
  53. }
  54. }
  55. func (this *SubscribeSender) safeSend(send *SubscribeSend) {
  56. defer func() {
  57. if r := recover(); r != nil {
  58. log.Printf("send subscribe panic: %v", r)
  59. logger.Info("[订阅推送] 发送订阅消息异常, module=%s, device_id=%s, err=%v", send.Module, send.DeviceId, r)
  60. }
  61. }()
  62. switch send.Module {
  63. case "hangup":
  64. this.SendHangupSubscribe(send.DeviceId)
  65. case "autofight":
  66. this.SendAutoFightSubscribe(send.DeviceId)
  67. case "guildgame":
  68. this.SendGuildGameSubscribe(send.OpenIds)
  69. case "alliance":
  70. this.SendAllianceSubscribe(send.OpenIds)
  71. case "warheavens":
  72. this.SendWarHeavensSubscribe(send.OpenIds)
  73. default:
  74. logger.Info("[push] unknown type=%s, openId=%s",
  75. send.Module, send.DeviceId)
  76. }
  77. }
  78. func (this *SubscribeSender) SendHangupSubscribe(deviceId string) {
  79. subscribeOne := new(subscribe.Subscribe)
  80. err := mgm.Coll(&subscribe.Subscribe{}).First(bson.M{"device_id": deviceId, "modules.hangup.enabled": true}, subscribeOne)
  81. if err != nil {
  82. logger.Info("[订阅推送] 查询订阅记录失败, module=%s, device_id=%s, err=%v", "hangup", deviceId, err)
  83. return
  84. }
  85. msg := make(map[string]map[string]string)
  86. msg["thing1"] = make(map[string]string)
  87. msg["thing1"]["value"] = "挂机奖励时长已满"
  88. msg["thing3"] = make(map[string]string)
  89. msg["thing3"]["value"] = "您的挂机奖励时长已满,请打开游戏领取"
  90. this.sendWechatSubscribe(subscribeOne.OpenId, "hangup", msg)
  91. }
  92. func (this *SubscribeSender) SendAutoFightSubscribe(deviceId string) {
  93. subscribeOne := new(subscribe.Subscribe)
  94. err := mgm.Coll(&subscribe.Subscribe{}).First(bson.M{"device_id": deviceId, "modules.autofight.enabled": true}, subscribeOne)
  95. if err != nil {
  96. logger.Info("[订阅推送] 查询订阅记录失败, module=%s, device_id=%s, err=%v", "autofight", deviceId, err)
  97. return
  98. }
  99. msg := make(map[string]map[string]string)
  100. msg["thing2"] = make(map[string]string)
  101. msg["thing2"]["value"] = "离线闯关结束,请收取您的离线闯关奖励。"
  102. msg["thing1"] = make(map[string]string)
  103. msg["thing1"]["value"] = "离线闯关提醒"
  104. this.sendWechatSubscribe(subscribeOne.OpenId, "autofight", msg)
  105. }
  106. func (this *SubscribeSender) SendGuildGameSubscribe(openIds []string) {
  107. msg := make(map[string]map[string]string)
  108. msg["thing2"] = make(map[string]string)
  109. msg["thing2"]["value"] = "门派攻防战即将开始,为门派荣誉而战吧。"
  110. msg["thing4"] = make(map[string]string)
  111. msg["thing4"]["value"] = "门派攻防战"
  112. for _, openId := range openIds {
  113. this.sendWechatSubscribe(openId, "guildgame", msg)
  114. }
  115. }
  116. func (this *SubscribeSender) SendAllianceSubscribe(openIds []string) {
  117. msg := make(map[string]map[string]string)
  118. msg["thing2"] = make(map[string]string)
  119. msg["thing2"]["value"] = "三界战场即将开始,争夺人间道统。"
  120. msg["thing4"] = make(map[string]string)
  121. msg["thing4"]["value"] = "三界争峰"
  122. for _, openId := range openIds {
  123. this.sendWechatSubscribe(openId, "alliance", msg)
  124. }
  125. }
  126. func (this *SubscribeSender) SendWarHeavensSubscribe(openIds []string) {
  127. msg := make(map[string]map[string]string)
  128. msg["thing2"] = make(map[string]string)
  129. msg["thing2"]["value"] = "独战诸仙即将开始!"
  130. msg["thing4"] = make(map[string]string)
  131. msg["thing4"]["value"] = "决战诸仙"
  132. for _, openId := range openIds {
  133. this.sendWechatSubscribe(openId, "warheavens", msg)
  134. }
  135. }
  136. func (this *SubscribeSender) sendWechatSubscribe(openId string, module string, msg map[string]map[string]string) {
  137. templateID := subscribeTemplateID(module)
  138. if templateID == "" {
  139. logger.Info("[push] subscribe template id not configured, channel=%s, module=%s", os.Getenv("CHANNEL"), module)
  140. return
  141. }
  142. this.wxHelper.SendWechatSubscribe(openId, templateID, msg)
  143. }
  144. func subscribeTemplateID(module string) string {
  145. templates, err := loadSubscribeTemplateConfig()
  146. if err != nil {
  147. logger.Info("[push] load subscribe template config failed: %v", err)
  148. return ""
  149. }
  150. channel := strings.ToLower(strings.TrimSpace(os.Getenv("CHANNEL")))
  151. moduleKey := strings.ToLower(strings.TrimSpace(module))
  152. if channel != "" {
  153. if templateID := strings.TrimSpace(templates[channel][moduleKey]); templateID != "" {
  154. return templateID
  155. }
  156. }
  157. return strings.TrimSpace(templates["default"][moduleKey])
  158. }
  159. func loadSubscribeTemplateConfig() (map[string]map[string]string, error) {
  160. subscribeTemplateOnce.Do(func() {
  161. data, err := os.ReadFile(subscribeTemplateConfigFile)
  162. if err != nil {
  163. subscribeTemplateLoadErr = err
  164. return
  165. }
  166. subscribeTemplateLoadErr = yaml.Unmarshal(data, &subscribeTemplateMap)
  167. })
  168. return subscribeTemplateMap, subscribeTemplateLoadErr
  169. }