From edcd98e5fef151b5231da97a1bd92f18ba3c5657 Mon Sep 17 00:00:00 2001 From: rainzm Date: Tue, 3 Aug 2021 16:54:01 +0800 Subject: [PATCH 1/5] fix(notify): optimization of print log 1. remomve unnecessary log 2. reduce the level of some logs, such as info -> debug --- pkg/notify/models/notification.go | 6 ------ pkg/notify/models/subscriber.go | 2 -- pkg/notify/tasks/notifications_send_task.go | 2 +- 3 files changed, 1 insertion(+), 9 deletions(-) diff --git a/pkg/notify/models/notification.go b/pkg/notify/models/notification.go index a081245027..49d1990380 100644 --- a/pkg/notify/models/notification.go +++ b/pkg/notify/models/notification.go @@ -90,7 +90,6 @@ func (nm *SNotificationManager) ValidateCreateData(ctx context.Context, userCred input.Contacts = []string{""} } } - log.Infof("notify input: %s", jsonutils.Marshal(input)) // check robot if len(input.Robots) > 0 { @@ -161,7 +160,6 @@ func (nm *SNotificationManager) ValidateCreateData(ctx context.Context, userCred if err != nil { return input, errors.Wrapf(err, "unable to generate name for %s", name) } - log.Infof("after validatecreate input: %s", jsonutils.Marshal(input)) return input, nil } @@ -219,7 +217,6 @@ func (nm *SNotificationManager) AllowPerformEventNotify(ctx context.Context, use // TODO: support project and domain func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.NotificationManagerEventNotifyInput) (api.NotificationManagerEventNotifyOutput, error) { - log.Infof("default receiverIds: %s", input.ReceiverIds) var output api.NotificationManagerEventNotifyOutput // check event _, err := parseEvent(input.Event) @@ -250,7 +247,6 @@ func (nm *SNotificationManager) PerformEventNotify(ctx context.Context, userCred if err != nil { return output, errors.Wrap(err, "unable to get receive") } - log.Infof("receiver for topic: %s", receiverIds1) receiverIds = append(receiverIds, receiverIds1...) } @@ -414,7 +410,6 @@ func (nm *SNotificationManager) createWithRobots(ctx context.Context, userCred m func (nm *SNotificationManager) create(ctx context.Context, userCred mcclient.TokenCredential, contactType string, receiverIds, contacts []string, priority, eventId string) error { if len(receiverIds)+len(contacts) == 0 { - log.Infof("%s: no send", contactType) return nil } @@ -441,7 +436,6 @@ func (nm *SNotificationManager) create(ctx context.Context, userCred mcclient.To return errors.Wrap(err, "ReceiverNotificationManager.CreateContact") } } - log.Infof("start NotificationSendTask for %s", contactType) n.SetModelManager(nm, n) task, err := taskman.TaskManager.NewTask(ctx, "NotificationSendTask", n, userCred, nil, "", "") if err != nil { diff --git a/pkg/notify/models/subscriber.go b/pkg/notify/models/subscriber.go index b62d304b89..9a7955e6f0 100644 --- a/pkg/notify/models/subscriber.go +++ b/pkg/notify/models/subscriber.go @@ -105,7 +105,6 @@ func (sm *SSubscriberManager) validateReceivers(ctx context.Context, receivers [ } func (sm *SSubscriberManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, input api.SubscriberCreateInput) (api.SubscriberCreateInput, error) { - log.Infof("before deal: %s", jsonutils.Marshal(input)) var err error // permission check sSystem, sDomain := string(rbacutils.ScopeSystem), string(rbacutils.ScopeDomain) @@ -191,7 +190,6 @@ func (sm *SSubscriberManager) ValidateCreateData(ctx context.Context, userCred m default: return input, httperrors.NewInputParameterError("unkown type %q", input.Type) } - log.Infof("after deal input: %s", jsonutils.Marshal(input)) return input, nil } diff --git a/pkg/notify/tasks/notifications_send_task.go b/pkg/notify/tasks/notifications_send_task.go index 3dd3e962cb..db6987882c 100644 --- a/pkg/notify/tasks/notifications_send_task.go +++ b/pkg/notify/tasks/notifications_send_task.go @@ -216,7 +216,7 @@ type FailedReceiverSpec struct { } func (self *NotificationSendTask) batchSend(ctx context.Context, contactType string, receivers []ReceiverSpec, params rpcapi.SendParams) (fails []FailedReceiverSpec, err error) { - log.Infof("contactType: %s, receivers: %s, params: %s", contactType, receivers, jsonutils.Marshal(params)) + log.Debugf("contactType: %s, receivers: %s, params: %s", contactType, receivers, jsonutils.Marshal(params)) if contactType != apis.ROBOT && contactType != apis.WEBHOOK { return self._batchSend(ctx, contactType, receivers, func(res []*rpcapi.SReceiver) ([]*rpcapi.FailedRecord, error) { return models.NotifyService.BatchSend(ctx, contactType, rpcapi.BatchSendParams{ From c2d7dd6a45147cb728a2c0f30025a127416dc58d Mon Sep 17 00:00:00 2001 From: rainzm Date: Tue, 3 Aug 2021 17:08:57 +0800 Subject: [PATCH 2/5] feat(notify): add ResourceAttributionName for subscriber --- pkg/apis/notify/subscriber.go | 2 ++ pkg/notify/models/subscriber.go | 19 +++++++++++-------- 2 files changed, 13 insertions(+), 8 deletions(-) diff --git a/pkg/apis/notify/subscriber.go b/pkg/apis/notify/subscriber.go index 4632079a85..d55bc4aa33 100644 --- a/pkg/apis/notify/subscriber.go +++ b/pkg/apis/notify/subscriber.go @@ -31,6 +31,8 @@ type SubscriberCreateInput struct { // example: 1e3824756bac4ac084e784ed297ec652 ResourceAttributionId string + ResourceAttributionName string + // description: domain id of resource // example: 1e3824756bac4ac084e784ed297ec652 DomainId string diff --git a/pkg/notify/models/subscriber.go b/pkg/notify/models/subscriber.go index 9a7955e6f0..100ae7d326 100644 --- a/pkg/notify/models/subscriber.go +++ b/pkg/notify/models/subscriber.go @@ -76,14 +76,15 @@ type SSubscriber struct { db.SStandaloneAnonResourceBase db.SEnabledResourceBase - TopicID string `width:"128" charset:"ascii" nullable:"false" index:"true" get:"user" list:"user" create:"required"` - Type string `width:"16" charset:"ascii" nullable:"false" index:"true" get:"user" list:"user" create:"required"` - Identification string `width:"128" charset:"ascii" nullable:"false" index:"true"` - RoleScope string `width:"8" charset:"ascii" nullable:"false" get:"user" list:"user" create:"optional"` - ResourceScope string `width:"8" charset:"ascii" nullable:"false" get:"user" list:"user" create:"required"` - ResourceAttributionId string `width:"128" charset:"ascii" nullable:"false" get:"user" list:"user" create:"optional"` - Scope string `width:"128" charset:"ascii" nullable:"false" create:"required"` - DomainId string `width:"128" charset:"ascii" nullable:"false" create:"optional"` + TopicID string `width:"128" charset:"ascii" nullable:"false" index:"true" get:"user" list:"user" create:"required"` + Type string `width:"16" charset:"ascii" nullable:"false" index:"true" get:"user" list:"user" create:"required"` + Identification string `width:"128" charset:"ascii" nullable:"false" index:"true"` + RoleScope string `width:"8" charset:"ascii" nullable:"false" get:"user" list:"user" create:"optional"` + ResourceScope string `width:"8" charset:"ascii" nullable:"false" get:"user" list:"user" create:"required"` + ResourceAttributionId string `width:"128" charset:"ascii" nullable:"false" get:"user" list:"user" create:"optional"` + ResourceAttributionName string `width:"128" charset:"utf8" list:"user" create:"optional"` + Scope string `width:"128" charset:"ascii" nullable:"false" create:"required"` + DomainId string `width:"128" charset:"ascii" nullable:"false" create:"optional"` } func (sm *SSubscriberManager) validateReceivers(ctx context.Context, receivers []string) ([]string, error) { @@ -149,6 +150,7 @@ func (sm *SSubscriberManager) ValidateCreateData(ctx context.Context, userCred m domainId = tenant.DomainId input.DomainId = domainId input.ResourceAttributionId = tenant.GetId() + input.ResourceAttributionName = tenant.GetName() case api.SUBSCRIBER_SCOPE_DOMAIN: tenant, err := db.TenantCacheManager.FetchDomainByIdOrName(ctx, input.ResourceAttributionId) if err != nil { @@ -157,6 +159,7 @@ func (sm *SSubscriberManager) ValidateCreateData(ctx context.Context, userCred m domainId = tenant.DomainId input.DomainId = domainId input.ResourceAttributionId = tenant.DomainId + input.ResourceAttributionName = tenant.Domain } if input.Scope == sDomain && domainId != userCred.GetDomainId() { return input, httperrors.NewForbiddenError("domain %s admin can't create subscriber for domain %s", userCred.GetDomainId(), domainId) From 3dea565c8add227b0de51cbea15224988aff4953 Mon Sep 17 00:00:00 2001 From: rainzm Date: Tue, 3 Aug 2021 17:10:54 +0800 Subject: [PATCH 3/5] fix(notify): fix for subscriber_receiver --- pkg/notify/models/subscriber_receiver.go | 17 ++++++++++------- pkg/notify/service/handlers.go | 5 +++++ 2 files changed, 15 insertions(+), 7 deletions(-) diff --git a/pkg/notify/models/subscriber_receiver.go b/pkg/notify/models/subscriber_receiver.go index f7927b48fe..ccec29b2fe 100644 --- a/pkg/notify/models/subscriber_receiver.go +++ b/pkg/notify/models/subscriber_receiver.go @@ -46,6 +46,14 @@ type SSubscriberReceiver struct { ReceiverId string `width:"36" charset:"ascii" nullable:"false" index:"true"` } +func (srm *SSubscriberReceiverManager) GetMasterFieldName() string { + return "subscriber_id" +} + +func (srm *SSubscriberReceiverManager) GetSlaveFieldName() string { + return "receiver_id" +} + func (srm *SSubscriberReceiverManager) getBySubscriberId(sId string) ([]SSubscriberReceiver, error) { q := srm.Query().Equals("subscriber_id", sId) srs := make([]SSubscriberReceiver, 0, 2) @@ -65,16 +73,11 @@ func (srm *SSubscriberReceiverManager) create(ctx context.Context, sId, receiver } func (srm *SSubscriberReceiverManager) delete(sId, receiverId string) error { - q := srm.Query().Equals("receiver_id", receiverId) + q := srm.Query().Equals("receiver_id", receiverId).Equals("subscriber_id", sId) var sr SSubscriberReceiver err := q.First(&sr) if err != nil { return err } - srp := &sr - _, err = db.Update(srp, func() error { - srp.MarkDelete() - return nil - }) - return err + return sr.Delete(context.Background(), nil) } diff --git a/pkg/notify/service/handlers.go b/pkg/notify/service/handlers.go index f450ff4cae..0f186d122e 100644 --- a/pkg/notify/service/handlers.go +++ b/pkg/notify/service/handlers.go @@ -73,6 +73,11 @@ func InitHandlers(app *appsrv.Application) { handler := db.NewModelHandler(manager) dispatcher.AddModelDispatcher(API_VERSION, app, handler) } + for _, manager := range []db.IJointModelManager{ + models.SubscriberReceiverManager, + } { + db.RegisterModelManager(manager) + } for _, manager := range []db.IJointModelManager{ models.ReceiverNotificationManager, } { From c89ae949e019197fb15bf3dc13b44e342acd71d0 Mon Sep 17 00:00:00 2001 From: rainzm Date: Tue, 3 Aug 2021 17:12:11 +0800 Subject: [PATCH 4/5] fix(notify): add list option 'enabled' for topic --- pkg/apis/notify/topic.go | 1 + pkg/notify/models/topic.go | 10 +++++++++- 2 files changed, 10 insertions(+), 1 deletion(-) diff --git a/pkg/apis/notify/topic.go b/pkg/apis/notify/topic.go index 024ad1c0b0..b0c327c854 100644 --- a/pkg/apis/notify/topic.go +++ b/pkg/apis/notify/topic.go @@ -18,6 +18,7 @@ import "yunion.io/x/onecloud/pkg/apis" type TopicListInput struct { apis.StandaloneResourceListInput + apis.EnabledResourceBaseListInput } type TopicDetails struct { diff --git a/pkg/notify/models/topic.go b/pkg/notify/models/topic.go index bcfa6da62d..4c5b22edcb 100644 --- a/pkg/notify/models/topic.go +++ b/pkg/notify/models/topic.go @@ -228,7 +228,15 @@ func (sm *STopicManager) InitializeData() error { } func (sm *STopicManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, input notify.TopicListInput) (*sqlchemy.SQuery, error) { - return sm.SStandaloneResourceBaseManager.ListItemFilter(ctx, q, userCred, input.StandaloneResourceListInput) + q, err := sm.SStandaloneResourceBaseManager.ListItemFilter(ctx, q, userCred, input.StandaloneResourceListInput) + if err != nil { + return nil, errors.Wrap(err, "SStandaloneResourceBaseManager.ListItemFilter") + } + q, err = sm.SEnabledResourceBaseManager.ListItemFilter(ctx, q, userCred, input.EnabledResourceBaseListInput) + if err != nil { + return nil, errors.Wrap(err, "SEnabledResourceBaseManager.ListItemFilter") + } + return q, nil } func (sm *STopicManager) FetchCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, objs []interface{}, fields stringutils2.SSortedStrings, isList bool) []notify.TopicDetails { From fc1adcedf5aa4615ecbadb3132ab2259c883c184 Mon Sep 17 00:00:00 2001 From: rainzm Date: Tue, 3 Aug 2021 17:13:20 +0800 Subject: [PATCH 5/5] fix(notify): fix create and list of subscriber --- pkg/notify/models/subscriber.go | 46 ++++++++++++++++++++++++++++----- 1 file changed, 40 insertions(+), 6 deletions(-) diff --git a/pkg/notify/models/subscriber.go b/pkg/notify/models/subscriber.go index 100ae7d326..ea91026e72 100644 --- a/pkg/notify/models/subscriber.go +++ b/pkg/notify/models/subscriber.go @@ -49,6 +49,7 @@ import ( "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/mcclient/modules" + "yunion.io/x/onecloud/pkg/util/logclient" "yunion.io/x/onecloud/pkg/util/rbacutils" "yunion.io/x/onecloud/pkg/util/stringutils2" ) @@ -64,7 +65,7 @@ func init() { "subscribers", ), } - SubscriberManager.SetVirtualObject(ReceiverNotificationManager) + SubscriberManager.SetVirtualObject(SubscriberManager) } type SSubscriberManager struct { @@ -196,6 +197,27 @@ func (sm *SSubscriberManager) ValidateCreateData(ctx context.Context, userCred m return input, nil } +func (s *SSubscriber) PostCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) { + s.SStandaloneAnonResourceBase.PostCreate(ctx, userCred, ownerId, query, data) + var input api.SubscriberCreateInput + _ = data.Unmarshal(&input) + if s.Type == api.SUBSCRIBER_TYPE_RECEIVER { + err := s.SetReceivers(ctx, input.Receivers) + if err != nil { + logclient.AddActionLogWithContext(ctx, s, logclient.ACT_CREATE, err.Error(), userCred, false) + _, err := db.Update(s, func() error { + s.SetEnabled(false) + return nil + }) + if err != nil { + log.Errorf("unable to enable subscriber: %v", err) + } + } + } + logclient.AddActionLogWithContext(ctx, s, logclient.ACT_CREATE, "", userCred, true) + return +} + func (s *SSubscriber) CustomizeCreate(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) error { err := s.SStandaloneAnonResourceBase.CustomizeCreate(ctx, userCred, ownerId, query, data) if err != nil { @@ -205,10 +227,6 @@ func (s *SSubscriber) CustomizeCreate(ctx context.Context, userCred mcclient.Tok _ = data.Unmarshal(&input) switch input.Type { case api.SUBSCRIBER_TYPE_RECEIVER: - err := s.SetReceivers(ctx, input.Receivers) - if err != nil { - return errors.Wrapf(err, "unable to set connect receivers with subscriber %s", s.Id) - } case api.SUBSCRIBER_TYPE_ROBOT: s.Identification = input.Robot case api.SUBSCRIBER_TYPE_ROLE: @@ -339,7 +357,7 @@ func (s *SSubscriber) CustomizeDelete(ctx context.Context, userCred mcclient.Tok } func (s *SSubscriber) receiverIdentifications() ([]api.Identification, error) { - srSubq := SubscriberReceiverManager.Query().Equals("subscription_id", s.Id).SubQuery() + srSubq := SubscriberReceiverManager.Query().Equals("subscriber_id", s.Id).SubQuery() rq := ReceiverManager.Query("id", "name") rq = rq.Join(srSubq, sqlchemy.Equals(srSubq.Field("receiver_id"), rq.Field("id"))) var ret []api.Identification @@ -523,6 +541,8 @@ func (sr *SSubscriber) SetReceivers(ctx context.Context, receiverIds []string) e addReceivers = append(addReceivers, receiverIds[j]) j++ case dbReceivers[i] == receiverIds[j]: + i++ + j++ } } // add @@ -576,3 +596,17 @@ func (s *SSubscriber) PerformDisable(ctx context.Context, userCred mcclient.Toke } return nil, nil } + +func (sm *SSubscriberManager) QueryDistinctExtraField(q *sqlchemy.SQuery, field string) (*sqlchemy.SQuery, error) { + q, err := sm.SStandaloneAnonResourceBaseManager.QueryDistinctExtraField(q, field) + if err == nil { + return q, nil + } + switch field { + case "resource_scope": + return sm.Query("resource_scope").Distinct(), nil + case "type": + return sm.Query("type").Distinct(), nil + } + return q, nil +}