fix(notify): addNotifyGroup

This commit is contained in:
马鸿飞
2023-08-08 11:36:37 +08:00
parent aa2292f819
commit aaacb52164
12 changed files with 270 additions and 48 deletions
+2 -1
View File
@@ -255,8 +255,9 @@ func (cm *SConfigManager) availableContactTypes(domainId string) ([]string, erro
return sets.NewString(ret...).UnsortedList(), nil
}
func (cm *SConfigManager) allContactType() ([]string, error) {
func (cm *SConfigManager) allContactType(domainid string) ([]string, error) {
q := cm.Query("type")
q = q.Equals("domain_id", domainid)
allTypes := make([]struct {
Type string
}, 0, 3)
+1 -1
View File
@@ -179,7 +179,7 @@ func (eq *SEmailQueue) doSend(ctx context.Context) {
eq.setStatus(ctx, api.EmailSending, nil)
driver := GetDriver(api.EMAIL)
err = driver.Send(ctx, api.SendParams{
EmailMsg: msg,
EmailMsg: *msg,
})
if err != nil {
eq.setStatus(ctx, api.EmailFail, err)
+1
View File
@@ -35,6 +35,7 @@ func InitDB() error {
TopicManager,
RobotManager,
SubscriberManager,
NotificationGroupManager,
} {
err := manager.InitializeData()
if err != nil {
+14 -2
View File
@@ -199,7 +199,7 @@ func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred
var output api.NotificationManagerEventNotifyOutput
// contact type
contactTypes := input.ContactTypes
cts, err := ConfigManager.allContactType()
cts, err := ConfigManager.allContactType(userCred.GetProjectDomainId())
if err != nil {
return output, errors.Wrap(err, "unable to fetch allContactType")
}
@@ -707,6 +707,12 @@ func (n *SNotification) GetTemplate(ctx context.Context, topicId, lang string, n
return out, errors.Wrapf(err, "get topic by id")
}
topic := topicModel.(*STopic)
groupKeys := []string{}
if topic.GroupKeys != nil {
groupKeys = *topic.GroupKeys
}
out.GroupTimes = uint(topic.GroupTimes)
rtStr, aStr, resultStr := event.ResourceType(), string(event.Action()), string(event.Result())
msgObj, err := jsonutils.ParseString(no.Message)
if err != nil {
@@ -727,7 +733,13 @@ func (n *SNotification) GetTemplate(ctx context.Context, topicId, lang string, n
Message: webhookMsg.String(),
}, nil
}
for _, key := range groupKeys {
keyValue, _ := msg.GetString(key)
if len(keyValue) > 0 {
out.GroupKey += keyValue
}
}
out.GroupTimes = uint(topic.GroupTimes)
if lang == "" {
lang = getLangSuffix(ctx)
}
+136
View File
@@ -0,0 +1,136 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package models // import "yunion.io/x/onecloud/pkg/notify/models"
import (
"context"
"fmt"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
apis "yunion.io/x/onecloud/pkg/apis/notify"
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
)
type SNotificationGroupManager struct {
db.SLogBaseManager
}
var NotificationGroupManager *SNotificationGroupManager
func init() {
NotificationGroupManager = &SNotificationGroupManager{
SLogBaseManager: db.NewLogBaseManager(
SNotificationGroup{},
"notification_group_tbl",
"notification_group",
"notification_groups",
"created_at",
consts.OpsLogWithClickhouse,
),
}
NotificationGroupManager.SetVirtualObject(NotificationGroupManager)
}
// 站内信
type SNotificationGroup struct {
db.SLogBase
GroupKey string `width:"128" nullable:"false" create:"required" list:"user" get:"user"`
Title string
// swagger:ignore
Message string
ReceiverId string `width:"128" nullable:"false" create:"required" list:"user" get:"user"`
Body jsonutils.JSONObject
Header jsonutils.JSONObject
MsgKey string
ContactType string `width:"32" nullable:"false" create:"required" list:"user" get:"user"`
Contact string `width:"128" nullable:"false" create:"required" list:"user" get:"user"`
CreatedAt time.Time
DomainId string `width:"128" nullable:"false" create:"required" list:"user" get:"user"`
}
func (ng *SNotificationGroupManager) TaskCreate(ctx context.Context, contactType string, args apis.SendParams) error {
if contactType == apis.WEBCONSOLE {
return nil
}
insertNotificationGroup := SNotificationGroup{
ContactType: contactType,
Body: args.Body,
Header: args.Header,
MsgKey: args.MsgKey,
ReceiverId: args.ReceiverId,
Title: args.Title,
Message: args.Message,
GroupKey: args.GroupKey,
Contact: args.Receivers.Contact,
CreatedAt: time.Now(),
DomainId: args.DomainId,
}
if contactType == apis.EMAIL {
insertNotificationGroup.Title = args.EmailMsg.Subject
insertNotificationGroup.Message = args.EmailMsg.Body
insertNotificationGroup.Contact = args.EmailMsg.To[0]
}
return NotificationGroupManager.TableSpec().Insert(ctx, &insertNotificationGroup)
}
func (ng *SNotificationGroupManager) TaskSend(ctx context.Context, input apis.SNotificationGroupSearchInput) (*apis.SendParams, error) {
q := ng.Query()
q = q.Between("created_at", input.StartTime, input.EndTime)
q = q.Equals("group_key", input.GroupKey)
q = q.Equals("receiver_id", input.ReceiverId)
q = q.Equals("contact_type", input.ContactType)
ngs := []SNotificationGroup{}
err := db.FetchModelObjects(ng, q, &ngs)
if err != nil {
return nil, errors.Wrap(err, "fetch notification groups")
}
if len(ngs) <= 1 {
return nil, errors.Wrapf(errors.ErrNotFound, "notification groups just found :%d", len(ngs))
}
sendParams := &apis.SendParams{
Body: ngs[0].Body,
Header: ngs[0].Header,
MsgKey: ngs[0].MsgKey,
Title: ngs[0].Title,
ReceiverId: ngs[0].ReceiverId,
Receivers: apis.SNotifyReceiver{
Contact: ngs[0].Contact,
},
DomainId: ngs[0].DomainId,
}
msg := ""
joinStr := " \n"
sendParams.Message = msg
if input.ContactType == apis.EMAIL {
joinStr = " <br>"
}
for _, ng := range ngs {
msg += fmt.Sprintf("%s %s", ng.Message, joinStr)
}
sendParams.Message = msg
if input.ContactType == apis.EMAIL {
sendParams.EmailMsg = apis.SEmailMessage{
Subject: sendParams.Title,
Body: msg,
To: []string{ngs[0].Contact},
}
}
return sendParams, nil
}
+30 -8
View File
@@ -73,14 +73,16 @@ func init() {
type STopic struct {
db.SEnabledStatusStandaloneResourceBase
Type string `width:"20" nullable:"false" create:"required" update:"user" list:"user"`
Resources uint64 `nullable:"false"`
Actions uint32 `nullable:"false"`
Results tristate.TriState `default:"true"`
TitleCn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
TitleEn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
ContentCn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
ContentEn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
Type string `width:"20" nullable:"false" create:"required" update:"user" list:"user"`
Resources uint64 `nullable:"false"`
Actions uint32 `nullable:"false"`
Results tristate.TriState `default:"true"`
TitleCn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
TitleEn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
ContentCn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
ContentEn string `length:"medium" nullable:"true" charset:"utf8" list:"user" update:"user" create:"optional"`
GroupKeys *api.STopicGroupKeys `nullable:"true"`
GroupTimes uint32 `nullable:"true"`
WebconsoleDisable tristate.TriState
}
@@ -311,6 +313,9 @@ func (sm *STopicManager) InitializeData() error {
t.ContentEn = api.COMMON_CONTENT_EN
t.TitleCn = api.COMMON_TITLE_CN
t.TitleEn = api.COMMON_TITLE_EN
groupKeys := []string{"action_display"}
t.GroupKeys = (*api.STopicGroupKeys)(&groupKeys)
t.GroupTimes = 60
case DefaultSystemExceptionEvent:
t.addResources(
notify.TOPIC_RESOURCE_HOST,
@@ -383,6 +388,9 @@ func (sm *STopicManager) InitializeData() error {
t.ContentEn = api.SYNC_ACCOUNT_STATUS_CONTENT_EN
t.TitleCn = api.SYNC_ACCOUNT_STATUS_TITLE_CN
t.TitleEn = api.SYNC_ACCOUNT_STATUS_TITLE_EN
groupKeys := []string{"name"}
t.GroupKeys = (*api.STopicGroupKeys)(&groupKeys)
t.GroupTimes = 60
case DefaultNetOutOfSync:
t.addResources(
notify.TOPIC_RESOURCE_NET,
@@ -396,6 +404,9 @@ func (sm *STopicManager) InitializeData() error {
t.ContentEn = api.NET_OUT_OF_SYNC_CONTENT_EN
t.TitleCn = api.NET_OUT_OF_SYNC_TITLE_CN
t.TitleEn = api.NET_OUT_OF_SYNC_TITLE_EN
groupKeys := []string{"service_name"}
t.GroupKeys = (*api.STopicGroupKeys)(&groupKeys)
t.GroupTimes = 60
case DefaultMysqlOutOfSync:
t.addResources(
notify.TOPIC_RESOURCE_DBINSTANCE,
@@ -409,6 +420,9 @@ func (sm *STopicManager) InitializeData() error {
t.ContentEn = api.MYSQL_OUT_OF_SYNC_CONTENT_EN
t.TitleCn = api.MYSQL_OUT_OF_SYNC_TITLE_CN
t.TitleEn = api.MYSQL_OUT_OF_SYNC_TITLE_EN
groupKeys := []string{"ip"}
t.GroupKeys = (*api.STopicGroupKeys)(&groupKeys)
t.GroupTimes = 60
case DefaultServiceAbnormal:
t.addResources(
notify.TOPIC_RESOURCE_SERVICE,
@@ -422,6 +436,9 @@ func (sm *STopicManager) InitializeData() error {
t.ContentEn = api.SERVICE_ABNORMAL_CONTENT_EN
t.TitleCn = api.SERVICE_ABNORMAL_TITLE_CN
t.TitleEn = api.SERVICE_ABNORMAL_TITLE_EN
groupKeys := []string{"service_name"}
t.GroupKeys = (*api.STopicGroupKeys)(&groupKeys)
t.GroupTimes = 60
case DefaultServerPanicked:
t.addResources(
notify.TOPIC_RESOURCE_SERVER,
@@ -435,6 +452,9 @@ func (sm *STopicManager) InitializeData() error {
t.ContentEn = api.SERVER_PANICKED_CONTENT_EN
t.TitleCn = api.SERVER_PANICKED_TITLE_CN
t.TitleEn = api.SERVER_PANICKED_TITLE_EN
groupKeys := []string{"name"}
t.GroupKeys = (*api.STopicGroupKeys)(&groupKeys)
t.GroupTimes = 60
case DefaultPasswordExpire:
t.addResources(
notify.TOPIC_RESOURCE_USER,
@@ -497,6 +517,8 @@ func (sm *STopicManager) InitializeData() error {
topic.Type = t.Type
topic.Results = t.Results
topic.WebconsoleDisable = t.WebconsoleDisable
topic.GroupKeys = t.GroupKeys
topic.GroupTimes = t.GroupTimes
if len(topic.ContentCn) == 0 {
if len(t.ContentCn) == 0 {
t.ContentCn = api.COMMON_CONTENT_CN
+1 -1
View File
@@ -152,7 +152,7 @@ func (emailSender *SEmailSender) Send(ctx context.Context, args api.SendParams)
gmsg.SetHeader("From", models.ConfigMap[api.EMAIL].Content.SenderAddress)
gmsg.SetHeader("To", args.Receivers.Contact)
gmsg.SetHeader("Subject", args.Title)
gmsg.SetBody("text/html", args.Message)
gmsg.SetBody("text/html", args.EmailMsg.Body)
dialer.StartTLSPolicy = gomail.MandatoryStartTLS
if err := dialer.DialAndSend(gmsg); err != nil {
return errors.Wrap(err, "send email")
+1
View File
@@ -66,6 +66,7 @@ func InitHandlers(app *appsrv.Application) {
models.RobotManager,
models.SubscriberManager,
models.EmailQueueManager,
models.NotificationGroupManager,
} {
db.RegisterModelManager(manager)
handler := db.NewModelHandler(manager)
+68 -33
View File
@@ -19,6 +19,7 @@ import (
"database/sql"
"fmt"
"strings"
"sync"
"time"
"yunion.io/x/jsonutils"
@@ -63,6 +64,12 @@ type ReceiverSpec struct {
rNotificaion *models.SReceiverNotification
}
var notificationSendMap sync.Map
var notificationGroupLock sync.Mutex
func init() {
notificationGroupLock = sync.Mutex{}
}
func (self *NotificationSendTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
notification := obj.(*models.SNotification)
if notification.Status == apis.NOTIFICATION_STATUS_OK {
@@ -102,7 +109,7 @@ func (self *NotificationSendTask) OnInit(ctx context.Context, obj db.IStandalone
}
// check receiver enabled
if !receiver.IsEnabled() {
sendFail(&rns[i], fmt.Sprintf("disabled receiver"))
sendFail(&rns[i], "disabled receiver")
continue
}
// check contact enabled
@@ -134,35 +141,6 @@ func (self *NotificationSendTask) OnInit(ctx context.Context, obj db.IStandalone
sendFail(&rns[i], fmt.Sprintf("unverified contactType %q", notification.ContactType))
continue
}
// contact, err := receiver.GetContact(notification.ContactType)
// if err != nil {
// logclient.AddSimpleActionLog(notification, logclient.ACT_SEND_NOTIFICATION, errors.Wrapf(err, "GetContact(%s)", notification.ContactType), self.GetUserCred(), false)
// continue
// }
// notifyRecv := receiver.GetNotifyReceiver()
// notifyRecv.Lang, _ = receiver.GetTemplateLang(ctx)
// notifyRecv.Contact = contact
// cv := ReceiverSpec{
// receiver: notifyRecv,
// rNotificaion: &rns[i],
// }
// switch notifyRecv.Lang {
// case "":
// receivers = append(receivers, cv)
// case apis.TEMPLATE_LANG_EN:
// receiversEn = append(receiversEn, cv)
// case apis.TEMPLATE_LANG_CN:
// receiversCn = append(receiversCn, cv)
// }
// lang, err := receiver.GetTemplateLang(ctx)
// if err != nil {
// reason := fmt.Sprintf("fail to GetTemplateLang: %s", err.Error())
// sendFail(&rns[i], reason)
// continue
// }
lang, err := receiver.GetTemplateLang(ctx)
if err != nil {
reason := fmt.Sprintf("fail to GetTemplateLang: %s", err.Error())
@@ -219,8 +197,11 @@ func (self *NotificationSendTask) OnInit(ctx context.Context, obj db.IStandalone
switch lang {
case apis.TEMPLATE_LANG_CN:
p.Message += "\n来自 " + options.Options.ApiServer
tz, _ := time.LoadLocation(options.Options.TimeZone)
p.Message += "\n发生于 " + time.Now().In(tz).Format("2006-01-02 15:04:05")
case apis.TEMPLATE_LANG_EN:
p.Message += "\nfrom " + options.Options.ApiServer
p.Message += "\nat " + time.Now().In(time.UTC).Format("2006-01-02 15:04:05")
}
}
@@ -271,6 +252,9 @@ type FailedReceiverSpec struct {
}
func (notificationSendTask *NotificationSendTask) batchSend(ctx context.Context, notification *models.SNotification, receivers []ReceiverSpec, params apis.SendParams) (fails []FailedReceiverSpec, err error) {
if notification.ContactType == apis.WEBCONSOLE {
return
}
for i := range receivers {
if receivers[i].receiver.IsRobot() {
robot := receivers[i].receiver.(*models.SRobot)
@@ -288,7 +272,7 @@ func (notificationSendTask *NotificationSendTask) batchSend(ctx context.Context,
params.Receivers.Contact, _ = receiver.GetContact(notification.ContactType)
driver := models.GetDriver(notification.ContactType)
if notification.ContactType == apis.EMAIL {
params.EmailMsg = &apis.SEmailMessage{
params.EmailMsg = apis.SEmailMessage{
To: []string{receiver.Email},
Subject: params.Title,
Body: params.Message,
@@ -299,7 +283,23 @@ func (notificationSendTask *NotificationSendTask) batchSend(ctx context.Context,
mobile := strings.Join(mobileArr, "")
params.Receivers.Contact = mobile
}
err = driver.Send(ctx, params)
params.ReceiverId = receiver.Id
if len(params.GroupKey) > 0 && params.GroupTimes > 0 {
notificationGroupLock.Lock()
if _, ok := notificationSendMap.Load(params.GroupKey + receiver.Id + notification.ContactType); ok {
err = models.NotificationGroupManager.TaskCreate(ctx, notification.ContactType, params)
} else {
err = models.NotificationGroupManager.TaskCreate(ctx, notification.ContactType, params)
if err != nil {
fails = append(fails, FailedReceiverSpec{ReceiverSpec: receivers[i], Reason: err.Error()})
}
err = driver.Send(ctx, params)
createTimeTicker(ctx, driver, params, receiver.Id, notification.ContactType)
}
notificationGroupLock.Unlock()
} else {
err = driver.Send(ctx, params)
}
if err != nil {
fails = append(fails, FailedReceiverSpec{ReceiverSpec: receivers[i], Reason: err.Error()})
}
@@ -313,6 +313,41 @@ func (notificationSendTask *NotificationSendTask) batchSend(ctx context.Context,
}
}
}
return fails, nil
}
func createTimeTicker(ctx context.Context, driver models.ISenderDriver, params apis.SendParams, receiverId, contactType string) {
// 创建一个计时器,每秒触发一次
// params.GroupTimes = 5
timer := time.NewTicker(time.Duration(params.GroupTimes) * time.Minute)
notificationSendMap.Store(params.GroupKey+receiverId+contactType, apis.SNotificationGroupSearchInput{
GroupKey: params.GroupKey,
ReceiverId: receiverId,
ContactType: contactType,
StartTime: time.Now(),
EndTime: time.Now().Add(time.Duration(params.GroupTimes) * time.Minute),
})
// 启动一个goroutine来处理计时器触发的事件
go func() {
for {
// 等待计时器触发的事件
<-timer.C
// 处理计时器触发的事件
arrValue, ok := notificationSendMap.Load(params.GroupKey + receiverId + contactType)
if !ok {
return
}
input := arrValue.(apis.SNotificationGroupSearchInput)
// 组装聚合后的消息
sendParams, err := models.NotificationGroupManager.TaskSend(ctx, input)
if err != nil {
log.Errorln("TaskSend err:", err)
return
}
driverT := models.GetDriver(contactType)
driverT.Send(ctx, *sendParams)
notificationSendMap.Delete(params.GroupKey + receiverId + contactType)
}
}()
}
+1 -1
View File
@@ -114,7 +114,7 @@ func (self *VerificationSendTask) OnInit(ctx context.Context, obj db.IStandalone
if len(emailMsg.Body) == 0 {
emailMsg.Body = param.Message
}
param.EmailMsg = emailMsg
param.EmailMsg = *emailMsg
driver := models.GetDriver(contactType)
err = driver.Send(ctx, param)
// err = models.NotifyService.Send(ctx, contactType, param)