Merge pull request #9868 from rainzm/notify/alert

Send notification when abnormal login occurs
This commit is contained in:
Zexi Li
2021-01-14 09:38:18 +08:00
committed by GitHub
18 changed files with 379 additions and 61 deletions
@@ -0,0 +1,5 @@
{{- if .admin -}}
域{{ .domain }}下的账号{{ .user }}由于异常登录已被锁定,请您核实用户使用情况。如果需要为用户解锁,请到用户列表启用该用户。
{{- else -}}
您的账号{{ .user }}由于异常登录已被锁定,请联系管理员解锁账号。
{{- end -}}
@@ -0,0 +1,5 @@
{{- if .admin -}}
The account {{ .user }} in domain {{ .domain }} has been locked due to abnormal login. Please verify the situation and, if you need to unlock this account, enable it in the user list.
{{- else -}}
The account {{ .user }} has been locked due to abnormal login. Please contact the administrator to unlock your account.
{{- end -}}
@@ -0,0 +1 @@
安全告警
@@ -0,0 +1 @@
Security Alerts
+2
View File
@@ -58,6 +58,8 @@ const (
NOTIFICATION_STATUS_OK = "ok"
NOTIFICATION_STATUS_PART_OK = "part_ok"
NOTIFICATION_TAG_ALERT = "alert"
TEMPLATE_TYPE_TITLE = "title"
TEMPLATE_TYPE_CONTENT = "content"
TEMPLATE_TYPE_REMOTE = "remote"
+7
View File
@@ -46,6 +46,12 @@ type NotificationCreateInput struct {
// description: message content or jsonobject
// required: ture
Message string `json:"message"`
// description: notification tag
// required: false
// example: alert
Tag string `json:"tag"`
Metadata map[string]interface{} `json:"metadata"`
IgnoreNonexistentReceiver bool `json:"ignore_nonexistent_receiver"`
}
type ReceiveDetail struct {
@@ -73,4 +79,5 @@ type NotificationListInput struct {
ContactType string
ReceiverId string
Tag string
}
+2
View File
@@ -34,6 +34,8 @@ const (
SERVER_PANICKED = "SERVER_PANICKED"
IMAGE_ACTIVED = "IMAGE_ACTIVED"
USER_LOGIN_EXCEPTION = "USER_LOGIN_EXCEPTION"
)
type SAction string
+100 -30
View File
@@ -53,9 +53,43 @@ var (
notifyAdminUsers []string
notifyAdminGroups []string
notifyclientI18nTable = i18n.Table{}
notifyclientI18nTable = i18n.Table{}
AdminSessionGenerator SAdminSessionGenerator = getAdminSesion
UserLangFetcher SUserLangFetcher = getUserLang
)
type SAdminSessionGenerator func(ctx context.Context, region string, apiVersion string) (*mcclient.ClientSession, error)
type SUserLangFetcher func(uids []string) (map[string]string, error)
func getAdminSesion(ctx context.Context, region string, apiVersion string) (*mcclient.ClientSession, error) {
return auth.GetAdminSession(ctx, region, apiVersion), nil
}
func getUserLang(uids []string) (map[string]string, error) {
s, err := AdminSessionGenerator(context.Background(), consts.GetRegion(), "")
if err != nil {
return nil, err
}
uidLang := make(map[string]string)
if len(uids) > 0 {
params := jsonutils.NewDict()
params.Set("filter", jsonutils.NewString(fmt.Sprintf("id.in(%s)", strings.Join(uids, ","))))
params.Set("details", jsonutils.JSONFalse)
params.Set("scope", jsonutils.NewString("system"))
params.Set("system", jsonutils.JSONTrue)
ret, err := modules.UsersV3.List(s, params)
if err != nil {
return nil, err
}
for i := range ret.Data {
id, _ := ret.Data[i].GetString("id")
langStr, _ := ret.Data[i].GetString("lang")
uidLang[id] = langStr
}
}
return uidLang, nil
}
const (
SUFFIX = "suffix"
)
@@ -149,6 +183,37 @@ func Notify(recipientId []string, isGroup bool, priority npk.TNotifyPriority, ev
notify(context.Background(), recipientId, isGroup, priority, event, data)
}
func NotifyWithTag(ctx context.Context, params SNotifyParams) {
p := sNotifyParams{
recipientId: params.RecipientId,
isGroup: params.IsGroup,
event: params.Event,
data: params.Data,
priority: params.Priority,
tag: params.Tag,
metadata: params.Metadata,
ignoreNonexistentReceiver: params.IgnoreNonexistentReceiver,
}
notifyWithChannel(ctx, p,
npk.NotifyByEmail,
npk.NotifyByDingTalk,
npk.NotifyByFeishu,
npk.NotifyByWorkwx,
npk.NotifyByWebConsole,
)
}
type SNotifyParams struct {
RecipientId []string
IsGroup bool
Priority npk.TNotifyPriority
Event string
Data jsonutils.JSONObject
Tag string
Metadata map[string]interface{}
IgnoreNonexistentReceiver bool
}
func NotifyWithContact(ctx context.Context, contacts []string, channel npk.TNotifyChannel, priority npk.TNotifyPriority, event string, data jsonutils.JSONObject) {
p := sNotifyParams{
contacts: contacts,
@@ -217,7 +282,6 @@ type sTarget struct {
func lang(ctx context.Context, contactType npk.TNotifyChannel, reIds []string, contacts []string) (map[language.Tag]*sTarget, error) {
contextLang := i18n.Lang(ctx)
s := auth.GetAdminSession(context.Background(), consts.GetRegion(), "")
langMap := make(map[language.Tag]*sTarget)
insertReid := func(lang language.Tag, id string) {
t := langMap[lang]
@@ -241,22 +305,9 @@ func lang(ctx context.Context, contactType npk.TNotifyChannel, reIds []string, c
uids = append(uids, contacts...)
}
uidLang := make(map[string]string)
if len(uids) > 0 {
params := jsonutils.NewDict()
params.Set("filter", jsonutils.NewString(fmt.Sprintf("id.in(%s)", strings.Join(uids, ","))))
params.Set("details", jsonutils.JSONFalse)
params.Set("scope", jsonutils.NewString("system"))
params.Set("system", jsonutils.JSONTrue)
ret, err := modules.UsersV3.List(s, params)
if err != nil {
return nil, err
}
for i := range ret.Data {
id, _ := ret.Data[i].GetString("id")
langStr, _ := ret.Data[i].GetString("lang")
uidLang[id] = langStr
}
uidLang, err := UserLangFetcher(uids)
if err != nil {
return nil, errors.Wrap(err, "unable to feth UserLang")
}
insert := func(id string, insertFunc func(language.Tag, string)) {
langStr := uidLang[id]
@@ -289,7 +340,10 @@ func lang(ctx context.Context, contactType npk.TNotifyChannel, reIds []string, c
func genMsgViaLang(ctx context.Context, p sNotifyParams) ([]npk.SNotifyMessage, error) {
reIds := make([]string, 0)
s := auth.GetAdminSession(context.Background(), consts.GetRegion(), "")
s, err := AdminSessionGenerator(context.Background(), consts.GetRegion(), "")
if err != nil {
return nil, err
}
if p.isGroup {
// fetch uid
uidSet := sets.NewString()
@@ -333,6 +387,9 @@ func genMsgViaLang(ctx context.Context, p sNotifyParams) ([]npk.SNotifyMessage,
body, _ = p.data.GetString()
}
msg.Msg = body
msg.Tag = p.tag
msg.Metadata = p.metadata
msg.IgnoreNonexistentReceiver = p.ignoreNonexistentReceiver
msgs = append(msgs, msg)
}
return msgs, nil
@@ -346,8 +403,12 @@ func intelliNotify(ctx context.Context, p sNotifyParams) {
}
for i := range msgs {
msg := msgs[i]
log.Infof("msg: %s", jsonutils.Marshal(msg))
notifyClientWorkerMan.Run(func() {
s := auth.GetAdminSession(context.Background(), consts.GetRegion(), "")
s, err := AdminSessionGenerator(context.Background(), consts.GetRegion(), "")
if err != nil {
log.Errorf("fail to get session: %v", err)
}
for {
err := npk.Notifications.Send(s, msg)
if err == nil {
@@ -388,14 +449,17 @@ func intelliNotify(ctx context.Context, p sNotifyParams) {
}
type sNotifyParams struct {
recipientId []string
isGroup bool
contacts []string
channel npk.TNotifyChannel
priority npk.TNotifyPriority
event string
data jsonutils.JSONObject
createReceiver bool
recipientId []string
isGroup bool
contacts []string
channel npk.TNotifyChannel
priority npk.TNotifyPriority
event string
data jsonutils.JSONObject
createReceiver bool
tag string
metadata map[string]interface{}
ignoreNonexistentReceiver bool
}
func rawNotify(ctx context.Context, p sNotifyParams) {
@@ -516,7 +580,10 @@ func NotifyRobotWithCtx(ctx context.Context, recipientId []string, isGroup bool,
}
func notifyRobot(ctx context.Context, robot string, recipientId []string, isGroup bool, priority npk.TNotifyPriority, event string, data jsonutils.JSONObject) error {
s := auth.GetAdminSession(ctx, consts.GetRegion(), "")
s, err := AdminSessionGenerator(ctx, consts.GetRegion(), "")
if err != nil {
return err
}
params := jsonutils.NewDict()
params.Set("robot", jsonutils.NewString(robot))
result, err := modules.NotifyConfig.PerformClassAction(s, "get-types", params)
@@ -656,7 +723,10 @@ func getIdentityId(s *mcclient.ClientSession, idName string, manager modulebase.
}
func FetchNotifyAdminRecipients(ctx context.Context, region string, users []string, groups []string) {
s := auth.GetAdminSession(ctx, region, "v1")
s, err := AdminSessionGenerator(ctx, region, "v1")
if err != nil {
log.Errorf("unable to get admin session: %v", err)
}
notifyAdminUsers = make([]string, 0)
for _, u := range users {
+79
View File
@@ -16,15 +16,23 @@ package sql
import (
"context"
"time"
"github.com/pkg/errors"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/util/sets"
api "yunion.io/x/onecloud/pkg/apis/identity"
noapi "yunion.io/x/onecloud/pkg/apis/notify"
"yunion.io/x/onecloud/pkg/cloudcommon/notifyclient"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/keystone/driver"
"yunion.io/x/onecloud/pkg/keystone/models"
o "yunion.io/x/onecloud/pkg/keystone/options"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/mcclient/modules/notify"
)
type SSQLDriver struct {
@@ -64,6 +72,7 @@ func (sql *SSQLDriver) Authenticate(ctx context.Context, ident mcclient.SAuthent
localUser.SaveFailedAuth()
if o.Options.PasswordErrorLockCount > 0 && localUser.FailedAuthCount > o.Options.PasswordErrorLockCount {
models.UserManager.LockUser(usrExt.Id, "too many failed auth attempts")
sql.alertNotify(ctx, usrExt, time.Now())
return nil, errors.Wrap(httperrors.ErrTooManyAttempts, "user locked")
}
return nil, errors.Wrap(err, "usrExt.VerifyPassword")
@@ -72,6 +81,76 @@ func (sql *SSQLDriver) Authenticate(ctx context.Context, ident mcclient.SAuthent
return usrExt, nil
}
func (sql *SSQLDriver) alertNotify(ctx context.Context, uext *api.SUserExtended, triggerTime time.Time) {
// get all users
daUserIds, err := getDomainAdminUserIds(uext.DomainName)
if err != nil {
log.Errorf("unable to get user with role domainadmin in domain %s: %v", uext.DomainName, err)
}
aUserIds, err := getAdminUserIds()
if err != nil {
log.Errorf("unable to get user with role admin: %v", err)
}
userSet := sets.NewString(daUserIds...)
userSet.Insert(aUserIds...)
userSet.Insert(uext.Id)
data := jsonutils.NewDict()
data.Set("user", jsonutils.NewString(uext.Name))
data.Set("domain", jsonutils.NewString(uext.DomainName))
metadata := map[string]interface{}{
"trigger_time": triggerTime,
}
// user
p := notifyclient.SNotifyParams{
RecipientId: userSet.UnsortedList(),
Priority: notify.NotifyPriorityCritical,
Event: notifyclient.USER_LOGIN_EXCEPTION,
Data: data,
Tag: noapi.NOTIFICATION_TAG_ALERT,
Metadata: metadata,
IgnoreNonexistentReceiver: true,
}
notifyclient.NotifyWithTag(ctx, p)
}
func fetchRoleId(name string) (string, error) {
id := struct {
Id string
}{}
q := models.RoleManager.Query().Equals("name", name)
err := q.First(&id)
if err != nil {
return "", err
}
return id.Id, nil
}
func getAdminUserIds() ([]string, error) {
return getUserIdsWithRole("admin", "")
}
func getUserIdsWithRole(roleName string, domainId string) ([]string, error) {
roleId, err := fetchRoleId(roleName)
if err != nil {
return nil, errors.Wrapf(err, "unable to fetch roleid of %s", roleName)
}
ras, _, err := models.AssignmentManager.FetchAll("", "", roleId, "", "", domainId, []string{}, []string{}, []string{}, []string{}, []string{}, []string{}, false, true, false, false, false, 0, 0)
if err != nil {
return nil, err
}
userIds := make([]string, 0, len(ras))
for i := range ras {
userIds = append(userIds, ras[i].User.Id)
}
log.Infof("%s User: %v", roleName, userIds)
return userIds, nil
}
func getDomainAdminUserIds(domainId string) ([]string, error) {
return getUserIdsWithRole("domainadmin", domainId)
}
func (sql *SSQLDriver) Sync(ctx context.Context) error {
return nil
}
+12 -5
View File
@@ -458,6 +458,7 @@ func roleAssignmentHandler(ctx context.Context, w http.ResponseWriter, r *http.R
input.Role.Id,
input.Scope.Domain.Id,
input.Scope.Project.Id,
"",
input.Users,
input.Groups,
input.Roles,
@@ -481,7 +482,7 @@ func roleAssignmentHandler(ctx context.Context, w http.ResponseWriter, r *http.R
}
func (manager *SAssignmentManager) queryAll(
userId, groupId, roleId, domainId, projectId string,
userId, groupId, roleId, domainId, projectId string, projectDomainId string,
users, groups, roles, domains, projects, projectDomains []string,
) *sqlchemy.SQuery {
assigments := manager.Query().SubQuery()
@@ -563,6 +564,12 @@ func (manager *SAssignmentManager) queryAll(
))
q = q.In("project_id", subq.SubQuery()).In("type", []string{api.AssignmentUserProject, api.AssignmentGroupProject})
}
if len(projectDomainId) > 0 {
subq := ProjectManager.Query("id")
domainQ := DomainManager.Query("id", "name").Equals("id", projectDomainId).SubQuery()
subq = subq.Join(domainQ, sqlchemy.Equals(subq.Field("domain_id"), domainQ.Field("id")))
q = q.In("project_id", subq.SubQuery()).In("type", []string{api.AssignmentUserProject, api.AssignmentGroupProject})
}
if len(projectDomains) > 0 {
subq := ProjectManager.Query("id")
domainQ := DomainManager.Query("id", "name").SubQuery()
@@ -638,17 +645,17 @@ func (assign *sAssignmentInternal) getRoleAssignment(domains, projects, groups,
}
func (manager *SAssignmentManager) FetchAll(
userId, groupId, roleId, domainId, projectId string,
userId, groupId, roleId, domainId, projectId string, projectDomainId string,
userStrs, groupStrs, roleStrs, domainStrs, projectStrs, projectDomainStrs []string,
includeNames, effective, includeSub, includeSystem, includePolicies bool,
limit, offset int) ([]api.SRoleAssignment, int64, error) {
var q *sqlchemy.SQuery
if effective {
usrq := manager.queryAll(userId, "", roleId, domainId, projectId, userStrs, nil, roleStrs, domainStrs, projectStrs, projectDomainStrs).In("type", []string{api.AssignmentUserProject, api.AssignmentUserDomain})
usrq := manager.queryAll(userId, "", roleId, domainId, projectId, projectDomainId, userStrs, nil, roleStrs, domainStrs, projectStrs, projectDomainStrs).In("type", []string{api.AssignmentUserProject, api.AssignmentUserDomain})
memberships := UsergroupManager.Query("user_id", "group_id").SubQuery()
grpproj := manager.queryAll("", groupId, roleId, domainId, projectId, nil, groupStrs, roleStrs, domainStrs, projectStrs, projectDomainStrs).In("type", []string{api.AssignmentGroupProject, api.AssignmentGroupDomain}).SubQuery()
grpproj := manager.queryAll("", groupId, roleId, domainId, projectId, projectDomainId, nil, groupStrs, roleStrs, domainStrs, projectStrs, projectDomainStrs).In("type", []string{api.AssignmentGroupProject, api.AssignmentGroupDomain}).SubQuery()
q2 := grpproj.Query(
grpproj.Field("type"),
memberships.Field("user_id"),
@@ -672,7 +679,7 @@ func (manager *SAssignmentManager) FetchAll(
q = sqlchemy.Union(usrq, q2).Query().Distinct()
} else {
q = manager.queryAll(userId, groupId, roleId, domainId, projectId, userStrs, groupStrs, roleStrs, domainStrs, projectStrs, projectDomainStrs).Distinct()
q = manager.queryAll(userId, groupId, roleId, domainId, projectId, projectDomainId, userStrs, groupStrs, roleStrs, domainStrs, projectStrs, projectDomainStrs).Distinct()
}
if !includeSystem {
+7 -3
View File
@@ -30,7 +30,7 @@ var (
defaultClient *mcclient.Client = nil
)
func getDefaultClient() *mcclient.Client {
func GetDefaultClient() *mcclient.Client {
if defaultClient == nil {
defaultClient = mcclient.NewClient("", 300, options.Options.DebugClient, true, "", "")
refreshDefaultClientServiceCatalog()
@@ -53,7 +53,11 @@ func GetDefaultAdminCred() mcclient.TokenCredential {
return defaultAdminCred
}
func getDefaultAdminCred() mcclient.TokenCredential {
func GetDefaultAdminSSimpleToken() *mcclient.SSimpleToken {
return getDefaultAdminCred()
}
func getDefaultAdminCred() *mcclient.SSimpleToken {
token := mcclient.SSimpleToken{}
usr, _ := UserManager.FetchUserExtended("", api.SystemAdminUser, api.DEFAULT_DOMAIN_ID, "")
token.UserId = usr.Id
@@ -72,5 +76,5 @@ func getDefaultAdminCred() mcclient.TokenCredential {
}
func GetDefaultClientSession(ctx context.Context, token mcclient.TokenCredential, region, apiVersion string) *mcclient.ClientSession {
return getDefaultClient().NewSession(ctx, region, "", "", token, apiVersion)
return GetDefaultClient().NewSession(ctx, region, "", "", token, apiVersion)
}
+24
View File
@@ -1167,3 +1167,27 @@ func (user *SUser) PerformUnlinkIdp(
}
return nil, nil
}
func GetUserLangForKeyStone(uids []string) (map[string]string, error) {
simpleUsers := make([]struct {
Id string
Lang string
}, 0, len(uids))
q := UserManager.Query()
if len(uids) == 0 {
return nil, nil
} else if len(uids) == 1 {
q = q.Equals("id", uids[0])
} else {
q = q.In("id", uids)
}
err := q.All(&simpleUsers)
if err != nil {
return nil, err
}
ret := make(map[string]string, len(simpleUsers))
for i := range simpleUsers {
ret[simpleUsers[i].Id] = simpleUsers[i].Lang
}
return ret, nil
}
+4
View File
@@ -25,6 +25,7 @@ import (
app_common "yunion.io/x/onecloud/pkg/cloudcommon/app"
"yunion.io/x/onecloud/pkg/cloudcommon/cronman"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/notifyclient"
common_options "yunion.io/x/onecloud/pkg/cloudcommon/options"
"yunion.io/x/onecloud/pkg/cloudcommon/policy"
"yunion.io/x/onecloud/pkg/keystone/cronjobs"
@@ -34,6 +35,7 @@ import (
"yunion.io/x/onecloud/pkg/keystone/saml"
_ "yunion.io/x/onecloud/pkg/keystone/tasks"
"yunion.io/x/onecloud/pkg/keystone/tokens"
"yunion.io/x/onecloud/pkg/keystone/util"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/util/logclient"
)
@@ -50,6 +52,8 @@ func StartService() {
policy.DefaultPolicyFetcher = localPolicyFetcher
logclient.DefaultSessionGenerator = models.GetDefaultClientSession
cronman.DefaultAdminSessionGenerator = models.GetDefaultAdminCred
notifyclient.AdminSessionGenerator = util.GetDefaulAdminSession
notifyclient.UserLangFetcher = models.GetUserLangForKeyStone
models.InitSyncWorkers()
+1 -1
View File
@@ -293,7 +293,7 @@ func (t *SAuthToken) getTokenV3(
token.Token.Projects[i].Domain.Id = extProjs[i].DomainId
token.Token.Projects[i].Domain.Name = extProjs[i].DomainName
}*/
assigns, _, err := models.AssignmentManager.FetchAll(user.Id, "", "", "", "",
assigns, _, err := models.AssignmentManager.FetchAll(user.Id, "", "", "", "", "",
nil, nil, nil, nil, nil, nil,
true, true, true, true, true, 0, 0)
if err != nil {
+15
View File
@@ -0,0 +1,15 @@
// 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 util // import "yunion.io/x/onecloud/pkg/keystone/util"
+62
View File
@@ -0,0 +1,62 @@
// 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 util
import (
"context"
"time"
"yunion.io/x/pkg/utils"
api "yunion.io/x/onecloud/pkg/apis/identity"
"yunion.io/x/onecloud/pkg/keystone/models"
"yunion.io/x/onecloud/pkg/keystone/tokens"
"yunion.io/x/onecloud/pkg/mcclient"
)
var (
authToken *tokens.SAuthToken
simpleToken *mcclient.SSimpleToken
)
func getDefaultAdminCredWithToken() (mcclient.TokenCredential, error) {
if simpleToken == nil {
simpleToken = models.GetDefaultAdminSSimpleToken()
}
var err error
if now := time.Now(); authToken == nil || authToken.ExpiresAt.Sub(now) < time.Duration(3600) {
authTokenTmp := &tokens.SAuthToken{
UserId: simpleToken.GetUserId(),
Method: api.AUTH_METHOD_TOKEN,
ProjectId: simpleToken.GetProjectId(),
ExpiresAt: now.Add(24 * time.Hour),
AuditIds: []string{utils.GenRequestId(16)},
}
simpleToken.Token, err = authTokenTmp.EncodeFernetToken()
if err != nil {
return nil, err
}
authToken = authTokenTmp
}
return simpleToken, nil
}
func GetDefaulAdminSession(ctx context.Context, region, apiVersion string) (*mcclient.ClientSession, error) {
cred, err := getDefaultAdminCredWithToken()
if err != nil {
return nil, err
}
return models.GetDefaultClientSession(ctx, cred, region, apiVersion), nil
}
+30 -21
View File
@@ -29,24 +29,30 @@ var (
)
type SNotifyMessage struct {
Uid []string `json:"uid,omitempty"`
Gid []string `json:"gid,omitempty"`
ContactType TNotifyChannel `json:"contact_type,omitempty"`
Contacts []string `json:"contracts"`
Topic string `json:"topic,omitempty"`
Priority TNotifyPriority `json:"priority,omitempty"`
Msg string `json:"msg,omitempty"`
Remark string `json:"remark,omitempty"`
Broadcast bool `json:"broadcast,omitempty"`
Uid []string `json:"uid,omitempty"`
Gid []string `json:"gid,omitempty"`
ContactType TNotifyChannel `json:"contact_type,omitempty"`
Contacts []string `json:"contracts"`
Topic string `json:"topic,omitempty"`
Priority TNotifyPriority `json:"priority,omitempty"`
Msg string `json:"msg,omitempty"`
Remark string `json:"remark,omitempty"`
Broadcast bool `json:"broadcast,omitempty"`
Tag string `json:"tag"`
Metadata map[string]interface{} `json:"metadata"`
IgnoreNonexistentReceiver bool `json:"ignore_nonexistent_receiver"`
}
type SNotifyV2Message struct {
Receivers []string `json:"receivers"`
Contacts []string `json:"contacts"`
ContactType string `json:"contact_type"`
Topic string `json:"topic"`
Priority string `json:"priority"`
Message string `json:"message"`
Receivers []string `json:"receivers"`
Contacts []string `json:"contacts"`
ContactType string `json:"contact_type"`
Topic string `json:"topic"`
Priority string `json:"priority"`
Message string `json:"message"`
Tag string `json:"tag"`
Metadata map[string]interface{} `json:"metadata"`
IgnoreNonexistentReceiver bool `json:"ignore_nonexistent_receiver"`
}
type NotificationManager struct {
@@ -75,12 +81,15 @@ func (manager *NotificationManager) Send(s *mcclient.ClientSession, msg SNotifyM
receiverIds = append(receiverIds, msg.Uid...)
v2msg := SNotifyV2Message{
Receivers: receiverIds,
Contacts: msg.Contacts,
ContactType: string(msg.ContactType),
Topic: msg.Topic,
Priority: string(msg.Priority),
Message: msg.Msg,
Receivers: receiverIds,
Contacts: msg.Contacts,
ContactType: string(msg.ContactType),
Topic: msg.Topic,
Priority: string(msg.Priority),
Message: msg.Msg,
Tag: msg.Tag,
Metadata: msg.Metadata,
IgnoreNonexistentReceiver: msg.IgnoreNonexistentReceiver,
}
params := jsonutils.Marshal(&v2msg)
+22 -1
View File
@@ -70,6 +70,7 @@ type SNotification struct {
Message string `create:"required"`
ReceivedAt time.Time `nullable:"true" list:"user" get:"user"`
SendTimes int
Tag string `width:"16" nullable:"true" index:"true" create:"optional"`
}
const (
@@ -77,6 +78,9 @@ const (
)
func (nm *SNotificationManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.NotificationCreateInput) (api.NotificationCreateInput, error) {
if len(input.Tag) == 0 && utils.IsInStringArray(input.Tag, []string{api.NOTIFICATION_TAG_ALERT}) {
return input, httperrors.NewInputParameterError("invalid tag")
}
// compatible
if len(input.Receivers) != 0 && input.ContactType == api.WEBCONSOLE {
input.Contacts = input.Receivers
@@ -115,9 +119,14 @@ func (nm *SNotificationManager) ValidateCreateData(ctx context.Context, userCred
if idSet.Has(re) || nameSet.Has(re) {
continue
}
return input, httperrors.NewInputParameterError("no such receiver whose uid is %q", re)
if !input.IgnoreNonexistentReceiver {
return input, httperrors.NewInputParameterError("no such receiver whose uid is %q", re)
}
}
input.Receivers = idSet.UnsortedList()
if len(input.Receivers)+len(input.Contacts) == 0 {
return input, httperrors.NewInputParameterError("no valid receiver or contact")
}
}
nowStr := time.Now().Format("2006-01-02 15:04:05")
if len(input.Priority) == 0 {
@@ -160,6 +169,15 @@ func (n *SNotification) CustomizeCreate(ctx context.Context, userCred mcclient.T
}
func (n *SNotification) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) {
if data.Contains("metadata") {
metadata := make(map[string]interface{})
err := data.Unmarshal(&metadata, "metadata")
if err != nil {
log.Errorf("unable to unmarshal to metadata: %v", err)
} else {
n.SetAllMetadata(ctx, metadata, userCred)
}
}
n.SetStatus(userCred, api.NOTIFICATION_STATUS_RECEIVED, "")
task, err := taskman.TaskManager.NewTask(ctx, "NotificationSendTask", n, userCred, nil, "", "")
if err != nil {
@@ -442,6 +460,9 @@ func (nm *SNotificationManager) ListItemFilter(ctx context.Context, q *sqlchemy.
subq := ReceiverNotificationManager.Query("notification_id").Equals("receiver_id", input.ReceiverId).SubQuery()
q = q.Join(subq, sqlchemy.Equals(q.Field("id"), subq.Field("notification_id")))
}
if len(input.Tag) > 0 {
q = q.Equals("tag", input.Tag)
}
return q, nil
}