subscribe.go 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195
  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. }
  60. }()
  61. switch send.Module {
  62. case "hangup":
  63. this.SendHangupSubscribe(send.DeviceId)
  64. case "autofight":
  65. this.SendAutoFightSubscribe(send.DeviceId)
  66. case "guildgame":
  67. this.SendGuildGameSubscribe(send.OpenIds)
  68. case "alliance":
  69. this.SendAllianceSubscribe(send.OpenIds)
  70. case "warheavens":
  71. this.SendWarHeavensSubscribe(send.OpenIds)
  72. default:
  73. logger.Info("[push] unknown type=%s, openId=%s",
  74. send.Module, send.DeviceId)
  75. }
  76. }
  77. func (this *SubscribeSender) SendHangupSubscribe(deviceId string) {
  78. subscribeOne := new(subscribe.Subscribe)
  79. err := mgm.Coll(&subscribe.Subscribe{}).First(bson.M{"device_id": deviceId, "modules.hangup.enabled": true}, subscribeOne)
  80. if err != nil {
  81. return
  82. }
  83. msg := make(map[string]map[string]string)
  84. msg["thing1"] = make(map[string]string)
  85. msg["thing1"]["value"] = "挂机奖励时长已满"
  86. msg["thing3"] = make(map[string]string)
  87. msg["thing3"]["value"] = "您的挂机奖励时长已满,请打开游戏领取"
  88. this.sendWechatSubscribe(subscribeOne.OpenId, "hangup", msg)
  89. }
  90. func (this *SubscribeSender) SendAutoFightSubscribe(deviceId string) {
  91. subscribeOne := new(subscribe.Subscribe)
  92. err := mgm.Coll(&subscribe.Subscribe{}).First(bson.M{"device_id": deviceId, "modules.autofight.enabled": true}, subscribeOne)
  93. if err != nil {
  94. return
  95. }
  96. msg := make(map[string]map[string]string)
  97. msg["thing2"] = make(map[string]string)
  98. msg["thing2"]["value"] = "离线闯关结束,请收取您的离线闯关奖励。"
  99. msg["thing1"] = make(map[string]string)
  100. msg["thing1"]["value"] = "离线闯关提醒"
  101. this.sendWechatSubscribe(subscribeOne.OpenId, "autofight", msg)
  102. }
  103. func (this *SubscribeSender) SendGuildGameSubscribe(openIds []string) {
  104. msg := make(map[string]map[string]string)
  105. msg["thing2"] = make(map[string]string)
  106. msg["thing2"]["value"] = "门派攻防战即将开始,为门派荣誉而战吧。"
  107. msg["thing4"] = make(map[string]string)
  108. msg["thing4"]["value"] = "门派攻防战"
  109. for _, openId := range openIds {
  110. this.sendWechatSubscribe(openId, "guildgame", msg)
  111. }
  112. }
  113. func (this *SubscribeSender) SendAllianceSubscribe(openIds []string) {
  114. msg := make(map[string]map[string]string)
  115. msg["thing2"] = make(map[string]string)
  116. msg["thing2"]["value"] = "三界战场即将开始,争夺人间道统。"
  117. msg["thing4"] = make(map[string]string)
  118. msg["thing4"]["value"] = "三界争峰"
  119. for _, openId := range openIds {
  120. this.sendWechatSubscribe(openId, "alliance", msg)
  121. }
  122. }
  123. func (this *SubscribeSender) SendWarHeavensSubscribe(openIds []string) {
  124. msg := make(map[string]map[string]string)
  125. msg["thing2"] = make(map[string]string)
  126. msg["thing2"]["value"] = "独战诸仙即将开始!"
  127. msg["thing4"] = make(map[string]string)
  128. msg["thing4"]["value"] = "决战诸仙"
  129. for _, openId := range openIds {
  130. this.sendWechatSubscribe(openId, "warheavens", msg)
  131. }
  132. }
  133. func (this *SubscribeSender) sendWechatSubscribe(openId string, module string, msg map[string]map[string]string) {
  134. templateID := subscribeTemplateID(module)
  135. if templateID == "" {
  136. logger.Info("[push] subscribe template id not configured, channel=%s, module=%s", os.Getenv("CHANNEL"), module)
  137. return
  138. }
  139. this.wxHelper.SendWechatSubscribe(openId, templateID, msg)
  140. }
  141. func subscribeTemplateID(module string) string {
  142. templates, err := loadSubscribeTemplateConfig()
  143. if err != nil {
  144. logger.Info("[push] load subscribe template config failed: %v", err)
  145. return ""
  146. }
  147. channel := strings.ToLower(strings.TrimSpace(os.Getenv("CHANNEL")))
  148. moduleKey := strings.ToLower(strings.TrimSpace(module))
  149. if channel != "" {
  150. if templateID := strings.TrimSpace(templates[channel][moduleKey]); templateID != "" {
  151. return templateID
  152. }
  153. }
  154. return strings.TrimSpace(templates["default"][moduleKey])
  155. }
  156. func loadSubscribeTemplateConfig() (map[string]map[string]string, error) {
  157. subscribeTemplateOnce.Do(func() {
  158. data, err := os.ReadFile(subscribeTemplateConfigFile)
  159. if err != nil {
  160. subscribeTemplateLoadErr = err
  161. return
  162. }
  163. subscribeTemplateLoadErr = yaml.Unmarshal(data, &subscribeTemplateMap)
  164. })
  165. return subscribeTemplateMap, subscribeTemplateLoadErr
  166. }