mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
@@ -29,6 +29,7 @@ const (
|
||||
FEISHU_ROBOT = "feishu-robot"
|
||||
DINGTALK_ROBOT = "dingtalk-robot"
|
||||
WORKWX_ROBOT = "workwx-robot"
|
||||
WEBHOOK = "webhook"
|
||||
|
||||
ROBOT = "robot"
|
||||
|
||||
@@ -60,4 +61,7 @@ const (
|
||||
TEMPLATE_TYPE_TITLE = "title"
|
||||
TEMPLATE_TYPE_CONTENT = "content"
|
||||
TEMPLATE_TYPE_REMOTE = "remote"
|
||||
|
||||
CTYPE_ROBOT_YES = "yes"
|
||||
CTYPE_ROBOT_ONLY = "only"
|
||||
)
|
||||
|
||||
@@ -14,6 +14,13 @@
|
||||
|
||||
package notifyclient
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
)
|
||||
|
||||
const (
|
||||
SYSTEM_ERROR = "SYSTEM_ERROR"
|
||||
SYSTEM_WARNING = "SYSTEM_WARNING"
|
||||
@@ -28,3 +35,42 @@ const (
|
||||
|
||||
IMAGE_ACTIVED = "IMAGE_ACTIVED"
|
||||
)
|
||||
|
||||
type SAction string
|
||||
|
||||
var (
|
||||
Event SEvent
|
||||
|
||||
ActionCreate SAction = "create"
|
||||
ActionUpdate SAction = "update"
|
||||
ActionDelete SAction = "delete"
|
||||
ActionRebuildRoot SAction = "rebuild_root"
|
||||
ActionChangeConfig SAction = "change_config"
|
||||
)
|
||||
|
||||
type SEvent struct {
|
||||
resourceType string
|
||||
action SAction
|
||||
}
|
||||
|
||||
func (se SEvent) WithResourceType(manager db.IModelManager) SEvent {
|
||||
se.resourceType = manager.Keyword()
|
||||
return se
|
||||
}
|
||||
|
||||
func (se SEvent) WithAction(a SAction) SEvent {
|
||||
se.action = a
|
||||
return se
|
||||
}
|
||||
|
||||
func (se SEvent) ResourceType() string {
|
||||
return se.resourceType
|
||||
}
|
||||
|
||||
func (se SEvent) Action() string {
|
||||
return string(se.action)
|
||||
}
|
||||
|
||||
func (se SEvent) String() string {
|
||||
return strings.ToUpper(fmt.Sprintf("%s_%s", se.ResourceType(), se.Action()))
|
||||
}
|
||||
|
||||
@@ -29,12 +29,14 @@ import (
|
||||
|
||||
"yunion.io/x/onecloud/pkg/appsrv"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/consts"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
"yunion.io/x/onecloud/pkg/i18n"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/mcclient/auth"
|
||||
"yunion.io/x/onecloud/pkg/mcclient/modulebase"
|
||||
"yunion.io/x/onecloud/pkg/mcclient/modules"
|
||||
npk "yunion.io/x/onecloud/pkg/mcclient/modules/notify"
|
||||
"yunion.io/x/onecloud/pkg/util/stringutils2"
|
||||
)
|
||||
|
||||
var (
|
||||
@@ -97,6 +99,9 @@ func getTemplate(ctx context.Context, topic string, contType string, channel npk
|
||||
}
|
||||
|
||||
func getContent(ctx context.Context, topic string, contType string, channel npk.TNotifyChannel, data jsonutils.JSONObject) (string, error) {
|
||||
if channel == npk.NotifyByWebhook {
|
||||
return "", nil
|
||||
}
|
||||
tmpl, err := getTemplate(ctx, topic, contType, channel)
|
||||
if err != nil {
|
||||
return "", err
|
||||
@@ -110,6 +115,23 @@ func getContent(ctx context.Context, topic string, contType string, channel npk.
|
||||
return buf.String(), nil
|
||||
}
|
||||
|
||||
func NotifyWebhook(ctx context.Context, userCred mcclient.TokenCredential, obj db.IModel, action SAction) error {
|
||||
ret, err := db.FetchCustomizeColumns(obj.GetModelManager(), ctx, userCred, jsonutils.NewDict(), []interface{}{obj}, stringutils2.SSortedStrings{}, false)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(ret) == 0 {
|
||||
return fmt.Errorf("unable to get details for model %q", obj.GetId())
|
||||
}
|
||||
event := Event.WithAction(action).WithResourceType(obj.GetModelManager())
|
||||
msg := jsonutils.NewDict()
|
||||
msg.Set("resource_type", jsonutils.NewString(event.ResourceType()))
|
||||
msg.Set("action", jsonutils.NewString(event.Action()))
|
||||
msg.Set("resource_details", ret[0])
|
||||
RawNotifyWithCtx(ctx, []string{}, false, npk.NotifyByWebhook, npk.NotifyPriorityNormal, event.String(), msg)
|
||||
return nil
|
||||
}
|
||||
|
||||
func NotifyWithCtx(ctx context.Context, recipientId []string, isGroup bool, priority npk.TNotifyPriority, event string, data jsonutils.JSONObject) {
|
||||
notify(ctx, recipientId, isGroup, priority, event, data)
|
||||
}
|
||||
|
||||
@@ -48,6 +48,7 @@ import (
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db/quotas"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/notifyclient"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/policy"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/userdata"
|
||||
"yunion.io/x/onecloud/pkg/cloudprovider"
|
||||
@@ -1599,6 +1600,11 @@ func (self *SGuest) PostUpdate(ctx context.Context, userCred mcclient.TokenCrede
|
||||
log.Errorf("StartRemoteUpdateTask fail: %s", err)
|
||||
}
|
||||
}
|
||||
// notify webhook
|
||||
err := notifyclient.NotifyWebhook(ctx, userCred, self, notifyclient.ActionUpdate)
|
||||
if err != nil {
|
||||
log.Errorf("unable to NotifyWebhook: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func (manager *SGuestManager) checkCreateQuota(
|
||||
|
||||
@@ -272,7 +272,11 @@ func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecCompleteFailed(ctx con
|
||||
func (self *GuestChangeConfigTask) OnGuestChangeCpuMemSpecFinish(ctx context.Context, guest *models.SGuest) {
|
||||
models.HostManager.ClearSchedDescCache(guest.HostId)
|
||||
self.SetStage("OnSyncConfigComplete", nil)
|
||||
err := guest.StartSyncTaskWithoutSyncstatus(ctx, self.UserCred, false, self.GetTaskId())
|
||||
err := notifyclient.NotifyWebhook(ctx, self.UserCred, guest, notifyclient.ActionChangeConfig)
|
||||
if err != nil {
|
||||
log.Errorf("unable to NotifyWebhook: %v", err)
|
||||
}
|
||||
err = guest.StartSyncTaskWithoutSyncstatus(ctx, self.UserCred, false, self.GetTaskId())
|
||||
if err != nil {
|
||||
self.markStageFailed(ctx, guest, jsonutils.NewString(fmt.Sprintf("StartSyncstatus fail %s", err)))
|
||||
return
|
||||
|
||||
@@ -145,6 +145,10 @@ func (self *GuestCreateTask) OnDeployGuestDescComplete(ctx context.Context, obj
|
||||
}
|
||||
|
||||
func (self *GuestCreateTask) notifyServerCreated(ctx context.Context, guest *models.SGuest) {
|
||||
err := notifyclient.NotifyWebhook(ctx, self.UserCred, guest, notifyclient.ActionCreate)
|
||||
if err != nil {
|
||||
log.Errorf("unable to NotifyWebhook: %v", err)
|
||||
}
|
||||
guest.NotifyServerEvent(
|
||||
ctx, self.UserCred, notifyclient.SERVER_CREATED,
|
||||
notify.NotifyPriorityImportant, true, nil, false,
|
||||
|
||||
@@ -345,6 +345,10 @@ func (self *GuestDeleteTask) DeleteGuest(ctx context.Context, guest *models.SGue
|
||||
}
|
||||
|
||||
func (self *GuestDeleteTask) NotifyServerDeleted(ctx context.Context, guest *models.SGuest) {
|
||||
err := notifyclient.NotifyWebhook(ctx, self.UserCred, guest, notifyclient.ActionDelete)
|
||||
if err != nil {
|
||||
log.Errorf("unable to NotifyWebhook: %v", err)
|
||||
}
|
||||
guest.NotifyServerEvent(
|
||||
ctx,
|
||||
self.UserCred,
|
||||
|
||||
@@ -176,6 +176,10 @@ func (self *GuestRebuildRootTask) OnRebuildAllDisksComplete(ctx context.Context,
|
||||
}
|
||||
}
|
||||
db.OpsLog.LogEvent(guest, db.ACT_REBUILD_ROOT, "", self.UserCred)
|
||||
err = notifyclient.NotifyWebhook(ctx, self.UserCred, guest, notifyclient.ActionRebuildRoot)
|
||||
if err != nil {
|
||||
log.Errorf("unable to NotifyWebhook: %v", err)
|
||||
}
|
||||
guest.NotifyServerEvent(
|
||||
ctx,
|
||||
self.UserCred,
|
||||
|
||||
@@ -33,4 +33,5 @@ const (
|
||||
NotifyFeishuRobot = TNotifyChannel("feishu-robot")
|
||||
NotifyByDingTalkRobot = TNotifyChannel("dingtalk-robot")
|
||||
NotifyByWorkwxRobot = TNotifyChannel("workwx-robot")
|
||||
NotifyByWebhook = TNotifyChannel("webhook")
|
||||
)
|
||||
|
||||
+18
-13
@@ -69,7 +69,7 @@ func (cm *SConfigManager) ValidateCreateData(ctx context.Context, userCred mccli
|
||||
if err != nil {
|
||||
return input, err
|
||||
}
|
||||
if !utils.IsInStringArray(input.Type, []string{api.EMAIL, api.MOBILE, api.DINGTALK, api.FEISHU, api.WEBCONSOLE, api.WORKWX, api.FEISHU_ROBOT, api.DINGTALK_ROBOT, api.WORKWX_ROBOT}) {
|
||||
if !utils.IsInStringArray(input.Type, []string{api.EMAIL, api.MOBILE, api.DINGTALK, api.FEISHU, api.WEBCONSOLE, api.WORKWX, api.FEISHU_ROBOT, api.DINGTALK_ROBOT, api.WORKWX_ROBOT, api.WEBHOOK}) {
|
||||
return input, httperrors.NewInputParameterError("unkown type %q", input.Type)
|
||||
}
|
||||
if input.Content == nil {
|
||||
@@ -139,19 +139,15 @@ func (cm *SConfigManager) AllowPerformGetTypes(ctx context.Context, userCred mcc
|
||||
return true
|
||||
}
|
||||
|
||||
func (cm *SConfigManager) PerformGetTypes(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ConfigManagerGetTypesInput) (api.ConfigManagerGetTypesOutput, error) {
|
||||
output := api.ConfigManagerGetTypesOutput{}
|
||||
allContactType, err := cm.allContactType()
|
||||
if err != nil {
|
||||
return output, err
|
||||
}
|
||||
func (cm *SConfigManager) filterContactType(cTypes []string, robot string) []string {
|
||||
var judge func(string) bool
|
||||
switch input.Robot {
|
||||
case "only":
|
||||
ret := make([]string, 0, len(cTypes)/2)
|
||||
switch robot {
|
||||
case api.CTYPE_ROBOT_ONLY:
|
||||
judge = func(ctype string) bool {
|
||||
return strings.Contains(ctype, "robot")
|
||||
}
|
||||
case "yes":
|
||||
case api.CTYPE_ROBOT_YES:
|
||||
judge = func(ctype string) bool {
|
||||
return true
|
||||
}
|
||||
@@ -160,12 +156,21 @@ func (cm *SConfigManager) PerformGetTypes(ctx context.Context, userCred mcclient
|
||||
return !strings.Contains(ctype, "robot")
|
||||
}
|
||||
}
|
||||
for _, ctype := range allContactType {
|
||||
for _, ctype := range cTypes {
|
||||
if judge(ctype) {
|
||||
output.Types = append(output.Types, ctype)
|
||||
ret = append(ret, ctype)
|
||||
}
|
||||
}
|
||||
output.Types = sortContactType(output.Types)
|
||||
return ret
|
||||
}
|
||||
|
||||
func (cm *SConfigManager) PerformGetTypes(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ConfigManagerGetTypesInput) (api.ConfigManagerGetTypesOutput, error) {
|
||||
output := api.ConfigManagerGetTypesOutput{}
|
||||
allContactType, err := cm.allContactType()
|
||||
if err != nil {
|
||||
return output, err
|
||||
}
|
||||
output.Types = sortContactType(cm.filterContactType(allContactType, input.Robot))
|
||||
return output, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -77,12 +77,18 @@ func (nm *SNotificationManager) ValidateCreateData(ctx context.Context, userCred
|
||||
if err != nil {
|
||||
return input, err
|
||||
}
|
||||
if !utils.IsInStringArray(input.ContactType, allContactType) {
|
||||
switch {
|
||||
case input.ContactType == api.WEBHOOK:
|
||||
case utils.IsInStringArray(input.ContactType, ConfigManager.filterContactType(allContactType, "")):
|
||||
//check uids, rids and contacts
|
||||
if len(input.Receivers) == 0 && len(input.Contacts) == 0 {
|
||||
return input, httperrors.NewMissingParameterError("receivers | contacts")
|
||||
}
|
||||
case utils.IsInStringArray(input.ContactType, ConfigManager.filterContactType(allContactType, api.CTYPE_ROBOT_ONLY)):
|
||||
default:
|
||||
return input, httperrors.NewInputParameterError("Unconfigured contact type %q", input.ContactType)
|
||||
}
|
||||
// check uids, rids and contacts
|
||||
if len(input.Receivers) == 0 && len(input.Contacts) == 0 {
|
||||
return input, httperrors.NewMissingParameterError("receivers | contacts")
|
||||
if !utils.IsInStringArray(input.ContactType, allContactType) {
|
||||
}
|
||||
// check receivers
|
||||
if len(input.Receivers) > 0 {
|
||||
@@ -181,6 +187,9 @@ func (n *SNotification) ReceiverNotificationsNotOK() ([]SReceiverNotification, e
|
||||
rnq := ReceiverNotificationManager.Query().Equals("notification_id", n.Id).NotEquals("status", api.RECEIVER_NOTIFICATION_OK)
|
||||
rns := make([]SReceiverNotification, 0, 1)
|
||||
err := db.FetchModelObjects(ReceiverNotificationManager, rnq, &rns)
|
||||
if err == sql.ErrNoRows {
|
||||
return []SReceiverNotification{}, nil
|
||||
}
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -104,10 +104,6 @@ func (self *NotificationSendTask) OnInit(ctx context.Context, obj db.IStandalone
|
||||
contactMap[contact] = &rns[i]
|
||||
}
|
||||
|
||||
if len(contactMap) == 0 {
|
||||
self.taskFailed(ctx, notification, strings.Join(failedRecord, "; "), true)
|
||||
}
|
||||
|
||||
// set status before send
|
||||
now := time.Now()
|
||||
contacts := make([]string, 0, len(contactMap))
|
||||
@@ -138,7 +134,7 @@ func (self *NotificationSendTask) OnInit(ctx context.Context, obj db.IStandalone
|
||||
for _, rn := range contactMap {
|
||||
rn.AfterSend(ctx, true, "")
|
||||
}
|
||||
if len(failedRecord) == len(contacts) {
|
||||
if len(failedRecord) > 0 && len(failedRecord) == len(contacts) {
|
||||
self.taskFailed(ctx, notification, strings.Join(failedRecord, "; "), true)
|
||||
return
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user