Merge pull request #11791 from rainzm/subscriber/fix

fix for subscriber
This commit is contained in:
Zexi Li
2021-08-03 20:03:29 +08:00
committed by GitHub
8 changed files with 79 additions and 31 deletions
+2
View File
@@ -31,6 +31,8 @@ type SubscriberCreateInput struct {
// example: 1e3824756bac4ac084e784ed297ec652
ResourceAttributionId string
ResourceAttributionName string
// description: domain id of resource
// example: 1e3824756bac4ac084e784ed297ec652
DomainId string
+1
View File
@@ -18,6 +18,7 @@ import "yunion.io/x/onecloud/pkg/apis"
type TopicListInput struct {
apis.StandaloneResourceListInput
apis.EnabledResourceBaseListInput
}
type TopicDetails struct {
-6
View File
@@ -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 {
+51 -16
View File
@@ -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 {
@@ -76,14 +77,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) {
@@ -105,7 +107,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)
@@ -150,6 +151,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 {
@@ -158,6 +160,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)
@@ -191,10 +194,30 @@ 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
}
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 {
@@ -204,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:
@@ -338,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
@@ -522,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
@@ -575,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
}
+10 -7
View File
@@ -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)
}
+9 -1
View File
@@ -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 {
+5
View File
@@ -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,
} {
+1 -1
View File
@@ -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{