diff --git a/pkg/apis/notify/const.go b/pkg/apis/notify/const.go index 3c3f3c0a8a..d625dbb70e 100644 --- a/pkg/apis/notify/const.go +++ b/pkg/apis/notify/const.go @@ -106,6 +106,16 @@ const ( TOPIC_RESOURCE_ELASTICCACHE = "elasticcache" TOPIC_RESOURCE_SCHEDULEDTASK = "scheduledtask" TOPIC_RESOURCE_BAREMETAL = "baremetal" + TOPIC_RESOURCE_VPC = "vpc" + TOPIC_RESOURCE_DNSZONE = "dns_zone" + TOPIC_RESOURCE_NATGATEWAY = "natgateway" + TOPIC_RESOURCE_WEBAPP = "webapp" + TOPIC_RESOURCE_CDNDOMAIN = "cdn_domain" + TOPIC_RESOURCE_FILESYSTEM = "file_system" + TOPIC_RESOURCE_WAF = "waf_instance" + TOPIC_RESOURCE_KAFKA = "kafka" + TOPIC_RESOURCE_ELASTICSEARCH = "elastic_search" + TOPIC_RESOURCE_MONGODB = "mongodb" SUBSCRIBER_TYPE_ROLE = "role" SUBSCRIBER_TYPE_ROBOT = "robot" diff --git a/pkg/apis/notify/event.go b/pkg/apis/notify/event.go index fce8c12ab8..dcac06cf74 100644 --- a/pkg/apis/notify/event.go +++ b/pkg/apis/notify/event.go @@ -40,6 +40,10 @@ var ( ActionCreateBackupServer SAction = "add_backup_server" ActionDelBackupServer SAction = "delete_backup_server" + ActionSyncCreate SAction = "sync_create" + ActionSyncUpdate SAction = "sync_update" + ActionSyncDelete SAction = "sync_delete" + ResultFailed SResult = "failed" ResultSucceed SResult = "succeed" ) diff --git a/pkg/cloudcommon/notifyclient/events.go b/pkg/cloudcommon/notifyclient/events.go index 4834760213..d843d36ee5 100644 --- a/pkg/cloudcommon/notifyclient/events.go +++ b/pkg/cloudcommon/notifyclient/events.go @@ -58,6 +58,10 @@ var ( ActionSyncStatus = api.ActionSyncStatus ActionPendingDelete = api.ActionPendingDelete + + ActionSyncCreate = api.ActionSyncCreate + ActionSyncUpdate = api.ActionSyncUpdate + ActionSyncDelete = api.ActionSyncDelete ) type SEvent struct { diff --git a/pkg/cloudcommon/notifyclient/notify.go b/pkg/cloudcommon/notifyclient/notify.go index 47b5d35df6..fc724e1584 100644 --- a/pkg/cloudcommon/notifyclient/notify.go +++ b/pkg/cloudcommon/notifyclient/notify.go @@ -249,16 +249,21 @@ func (t *eventTask) Run() { } func EventNotify(ctx context.Context, userCred mcclient.TokenCredential, ep SEventNotifyParam) { - ret, err := db.FetchCustomizeColumns(ep.Obj.GetModelManager(), ctx, userCred, jsonutils.NewDict(), []interface{}{ep.Obj}, stringutils2.SSortedStrings{}, false) - if err != nil { - log.Errorf("unable to FetchCustomizeColumns: %v", err) - return + var objDetails *jsonutils.JSONDict + if ep.Action == ActionDelete || ep.Action == ActionSyncDelete { + objDetails = jsonutils.Marshal(ep.Obj).(*jsonutils.JSONDict) + } else { + ret, err := db.FetchCustomizeColumns(ep.Obj.GetModelManager(), ctx, userCred, jsonutils.NewDict(), []interface{}{ep.Obj}, stringutils2.SSortedStrings{}, false) + if err != nil { + log.Errorf("unable to FetchCustomizeColumns: %v", err) + return + } + if len(ret) == 0 { + log.Errorf("unable to FetchCustomizeColumns: details of model %q is empty", ep.Obj.GetId()) + return + } + objDetails = ret[0] } - if len(ret) == 0 { - log.Errorf("unable to FetchCustomizeColumns: details of model %q is empty", ep.Obj.GetId()) - return - } - objDetails := ret[0] if ep.ObjDetailsDecorator != nil { ep.ObjDetailsDecorator(ctx, objDetails) } diff --git a/pkg/compute/models/app.go b/pkg/compute/models/app.go index a1537d69cf..a8c755eb49 100644 --- a/pkg/compute/models/app.go +++ b/pkg/compute/models/app.go @@ -29,6 +29,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -286,6 +287,10 @@ func (self *SCloudregion) newFromCloudApp(ctx context.Context, userCred mcclient SyncCloudProject(userCred, &app, provider.GetOwnerId(), ext, provider.Id) db.OpsLog.LogEvent(&app, db.ACT_CREATE, app.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &app, + Action: notifyclient.ActionSyncCreate, + }) return &app, nil } @@ -319,11 +324,19 @@ func (a *SApp) purge(ctx context.Context, userCred mcclient.TokenCredential) err } func (a *SApp) syncRemoveCloudApp(ctx context.Context, userCred mcclient.TokenCredential) error { - return a.purge(ctx, userCred) + err := a.purge(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: a, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (a *SApp) SyncWithCloudApp(ctx context.Context, userCred mcclient.TokenCredential, provider *SCloudprovider, ext cloudprovider.ICloudApp) error { - _, err := db.UpdateWithLock(ctx, a, func() error { + diff, err := db.UpdateWithLock(ctx, a, func() error { a.ExternalId = ext.GetGlobalId() a.Status = ext.GetStatus() a.Type = ext.GetType() @@ -342,6 +355,12 @@ func (a *SApp) SyncWithCloudApp(ctx context.Context, userCred mcclient.TokenCred if result.IsError() { return result.AllError() } + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: a, + Action: notifyclient.ActionSyncUpdate, + }) + } return nil } diff --git a/pkg/compute/models/buckets.go b/pkg/compute/models/buckets.go index 723dde8a12..51df47bba5 100644 --- a/pkg/compute/models/buckets.go +++ b/pkg/compute/models/buckets.go @@ -39,6 +39,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/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" "yunion.io/x/onecloud/pkg/httperrors" @@ -232,6 +233,10 @@ func (manager *SBucketManager) newFromCloudBucket( } SyncCloudProject(userCred, &bucket, provider.GetOwnerId(), extBucket, provider.Id) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &bucket, + Action: notifyclient.ActionSyncCreate, + }) bucket.SyncShareState(ctx, userCred, provider.getAccountShareInfo()) syncVirtualResourceMetadata(ctx, userCred, &bucket, extBucket) @@ -308,6 +313,12 @@ func (bucket *SBucket) syncWithCloudBucket( syncVirtualResourceMetadata(ctx, userCred, bucket, extBucket) db.OpsLog.LogSyncUpdate(bucket, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: bucket, + Action: notifyclient.ActionSyncUpdate, + }) + } if !oStats.Equals(extBucket.GetStats()) { db.OpsLog.LogEvent(bucket, api.BUCKET_OPS_STATS_CHANGE, bucket.GetShortDesc(ctx), userCred) @@ -332,6 +343,10 @@ func (bucket *SBucket) syncRemoveCloudBucket( if err != nil { return errors.Wrap(err, "RealDelete") } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: bucket, + Action: notifyclient.ActionSyncDelete, + }) return nil } diff --git a/pkg/compute/models/cdn_domains.go b/pkg/compute/models/cdn_domains.go index cfe53acf85..1c8054b72f 100644 --- a/pkg/compute/models/cdn_domains.go +++ b/pkg/compute/models/cdn_domains.go @@ -29,6 +29,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -171,7 +172,15 @@ func (self *SCDNDomain) syncRemoveCloudCDNDomain(ctx context.Context, userCred m if err != nil { return errors.Wrapf(err, "ValidateDeleteCondition") } - return self.RealDelete(ctx, userCred) + err = self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SCDNDomain) GetICloudCDNDomain() (cloudprovider.ICloudCDNDomain, error) { @@ -200,6 +209,12 @@ func (self *SCDNDomain) SyncWithCloudCDNDomain(ctx context.Context, userCred mcc return err } db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } syncMetadata(ctx, userCred, self, ext) if provider := self.GetCloudprovider(); provider != nil { @@ -234,6 +249,10 @@ func (self *SCloudprovider) newFromCloudCDNDomain(ctx context.Context, userCred domain.SyncShareState(ctx, userCred, self.getAccountShareInfo()) db.OpsLog.LogEvent(&domain, db.ACT_CREATE, domain.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &domain, + Action: notifyclient.ActionSyncCreate, + }) return &domain, nil } diff --git a/pkg/compute/models/dbinstances.go b/pkg/compute/models/dbinstances.go index e6cacd8e98..34107a214f 100644 --- a/pkg/compute/models/dbinstances.go +++ b/pkg/compute/models/dbinstances.go @@ -38,6 +38,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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -1492,7 +1493,15 @@ func (manager *SDBInstanceManager) SyncDBInstances(ctx context.Context, userCred } func (self *SDBInstance) syncRemoveCloudDBInstance(ctx context.Context, userCred mcclient.TokenCredential) error { - return self.Purge(ctx, userCred) + err := self.Purge(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SDBInstance) ValidateDeleteCondition(ctx context.Context, info jsonutils.JSONObject) error { @@ -1704,6 +1713,12 @@ func (self *SDBInstance) SyncWithCloudDBInstance(ctx context.Context, userCred m } syncVirtualResourceMetadata(ctx, userCred, self, extInstance) db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } return nil } @@ -1793,6 +1808,11 @@ func (manager *SDBInstanceManager) newFromCloudDBInstance(ctx context.Context, u db.OpsLog.LogEvent(&instance, db.ACT_CREATE, instance.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &instance, + Action: notifyclient.ActionSyncCreate, + }) + return &instance, nil } diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index eef8d5544f..e2173d5e48 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -1463,7 +1463,15 @@ func (self *SDisk) syncRemoveCloudDisk(ctx context.Context, userCred mcclient.To if err != nil { return err } - return self.RealDelete(ctx, userCred) + err = self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, provider cloudprovider.ICloudProvider, extDisk cloudprovider.ICloudDisk, index int, syncOwnerId mcclient.IIdentityProvider, managerId string) error { @@ -1536,6 +1544,13 @@ func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.Toke } db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } + syncVirtualResourceMetadata(ctx, userCred, self, extDisk) if len(guests) == 0 { @@ -1617,6 +1632,11 @@ func (manager *SDiskManager) newFromCloudDisk(ctx context.Context, userCred mccl db.OpsLog.LogEvent(&disk, db.ACT_CREATE, disk.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &disk, + Action: notifyclient.ActionSyncCreate, + }) + return &disk, nil } diff --git a/pkg/compute/models/dns_zonecaches.go b/pkg/compute/models/dns_zonecaches.go index 99709f1ba0..f9a44e9186 100644 --- a/pkg/compute/models/dns_zonecaches.go +++ b/pkg/compute/models/dns_zonecaches.go @@ -31,6 +31,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -223,6 +224,10 @@ func (self *SDnsZoneCache) syncRemove(ctx context.Context, userCred mcclient.Tok if err != nil { return errors.Wrapf(err, "dnsZone.RealDelete for %s(%s)", dnsZone.Name, dnsZone.Id) } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) return self.RealDelete(ctx, userCred) } @@ -288,6 +293,10 @@ func (self *SDnsZoneCache) SyncWithCloudDnsZone(ctx context.Context, userCred mc dnsZone.AddVpc(ctx, add.(string)) } } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) return nil } diff --git a/pkg/compute/models/dns_zones.go b/pkg/compute/models/dns_zones.go index e27b504e48..10a533a214 100644 --- a/pkg/compute/models/dns_zones.go +++ b/pkg/compute/models/dns_zones.go @@ -457,6 +457,10 @@ func (manager *SDnsZoneManager) newFromCloudDnsZone(ctx context.Context, userCre dnsZone.AddVpc(ctx, vpcId) } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: dnsZone, + Action: notifyclient.ActionSyncCreate, + }) SyncCloudDomain(userCred, dnsZone, account.GetOwnerId()) dnsZone.SyncShareState(ctx, userCred, account.getAccountShareInfo()) diff --git a/pkg/compute/models/elastic_search.go b/pkg/compute/models/elastic_search.go index 619904d098..59f42d27f1 100644 --- a/pkg/compute/models/elastic_search.go +++ b/pkg/compute/models/elastic_search.go @@ -31,6 +31,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -392,7 +393,15 @@ func (self *SElasticSearch) GetIElasticSearch() (cloudprovider.ICloudElasticSear } func (self *SElasticSearch) syncRemoveCloudElasticSearch(ctx context.Context, userCred mcclient.TokenCredential) error { - return self.RealDelete(ctx, userCred) + err := self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } // 同步资源属性 @@ -471,6 +480,12 @@ func (self *SElasticSearch) SyncWithCloudElasticSearch(ctx context.Context, user if err != nil { return errors.Wrapf(err, "db.Update") } + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } syncVirtualResourceMetadata(ctx, userCred, self, ext) if provider := self.GetCloudprovider(); provider != nil { @@ -570,6 +585,10 @@ func (self *SCloudregion) newFromCloudElasticSearch(ctx context.Context, userCre return nil, errors.Wrapf(err, "newFromCloudElasticSearch.Insert") } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &es, + Action: notifyclient.ActionSyncCreate, + }) // 同步标签 syncVirtualResourceMetadata(ctx, userCred, &es, ext) // 同步项目归属 diff --git a/pkg/compute/models/elasticcache_instances.go b/pkg/compute/models/elasticcache_instances.go index 7273007356..0e31d20ec1 100644 --- a/pkg/compute/models/elasticcache_instances.go +++ b/pkg/compute/models/elasticcache_instances.go @@ -37,6 +37,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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -576,7 +577,15 @@ func (self *SElasticcache) syncRemoveCloudElasticcache(ctx context.Context, user self.SetStatus(userCred, api.ELASTIC_CACHE_STATUS_ERROR, "sync to delete") return errors.Wrap(err, "ValidateDeleteCondition") } - return self.SVirtualResourceBase.Delete(ctx, userCred) + err = self.SVirtualResourceBase.Delete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SElasticcache) SyncWithCloudElasticcache(ctx context.Context, userCred mcclient.TokenCredential, provider *SCloudprovider, extInstance cloudprovider.ICloudElasticcache) error { @@ -620,6 +629,12 @@ func (self *SElasticcache) SyncWithCloudElasticcache(ctx context.Context, userCr } syncVirtualResourceMetadata(ctx, userCred, self, extInstance) db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } return nil } @@ -739,6 +754,11 @@ func (manager *SElasticcacheManager) newFromCloudElasticcache(ctx context.Contex SyncCloudProject(userCred, &instance, ownerId, extInstance, provider.Id) db.OpsLog.LogEvent(&instance, db.ACT_CREATE, instance.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &instance, + Action: notifyclient.ActionSyncCreate, + }) + return &instance, nil } diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index 606011aedc..f6bb60b5db 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -34,6 +34,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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -388,7 +389,15 @@ func (self *SElasticip) syncRemoveCloudEip(ctx context.Context, userCred mcclien lockman.LockObject(ctx, self) defer lockman.ReleaseObject(ctx, self) - return self.RealDelete(ctx, userCred) + err := self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SElasticip) SyncInstanceWithCloudEip(ctx context.Context, userCred mcclient.TokenCredential, ext cloudprovider.ICloudEIP) error { @@ -485,6 +494,13 @@ func (self *SElasticip) SyncWithCloudEip(ctx context.Context, userCred mcclient. } db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } + err = self.SyncInstanceWithCloudEip(ctx, userCred, ext) if err != nil { return errors.Wrap(err, "fail to sync associated instance of EIP") @@ -563,6 +579,10 @@ func (manager *SElasticipManager) newFromCloudEip(ctx context.Context, userCred } db.OpsLog.LogEvent(&eip, db.ACT_CREATE, eip.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &eip, + Action: notifyclient.ActionSyncCreate, + }) return &eip, nil } diff --git a/pkg/compute/models/filesystem.go b/pkg/compute/models/filesystem.go index 3a6eb08688..c690c216c7 100644 --- a/pkg/compute/models/filesystem.go +++ b/pkg/compute/models/filesystem.go @@ -31,6 +31,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudcommon/validators" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" @@ -381,7 +382,15 @@ func (self *SFileSystem) syncRemove(ctx context.Context, userCred mcclient.Token return self.SetStatus(userCred, api.NAS_STATUS_UNKNOWN, "sync to delete") } - return self.RealDelete(ctx, userCred) + err = self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SFileSystem) CustomizeDelete(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) error { @@ -436,7 +445,7 @@ func (self *SFileSystem) SyncAllWithCloudFileSystem(ctx context.Context, userCre } func (self *SFileSystem) SyncWithCloudFileSystem(ctx context.Context, userCred mcclient.TokenCredential, fs cloudprovider.ICloudFileSystem) error { - _, err := db.Update(self, func() error { + diff, err := db.Update(self, func() error { self.Status = fs.GetStatus() self.StorageType = fs.GetStorageType() self.Protocol = fs.GetProtocol() @@ -456,6 +465,12 @@ func (self *SFileSystem) SyncWithCloudFileSystem(ctx context.Context, userCred m if err != nil { return errors.Wrapf(err, "db.Update") } + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } syncMetadata(ctx, userCred, self, fs) return nil } @@ -490,7 +505,7 @@ func (self *SCloudregion) newFromCloudFileSystem(ctx context.Context, userCred m if zoneId := fs.GetZoneId(); len(zoneId) > 0 { nas.ZoneId, _ = self.getZoneIdBySuffix(zoneId) } - return func() (*SFileSystem, error) { + fileSystem, err := func() (*SFileSystem, error) { lockman.LockRawObject(ctx, FileSystemManager.Keyword(), "name") defer lockman.ReleaseRawObject(ctx, FileSystemManager.Keyword(), "name") @@ -502,6 +517,14 @@ func (self *SCloudregion) newFromCloudFileSystem(ctx context.Context, userCred m return &nas, FileSystemManager.TableSpec().Insert(ctx, &nas) }() + if err != nil { + return nil, err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &nas, + Action: notifyclient.ActionSyncCreate, + }) + return fileSystem, nil } func (self *SFileSystem) AllowPerformSyncstatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 7cbca100cb..a2c03078cb 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -46,6 +46,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" @@ -2408,7 +2409,14 @@ func (self *SGuest) syncRemoveCloudVM(ctx context.Context, userCred mcclient.Tok if options.SyncPurgeRemovedResources.Contains(self.Keyword()) { log.Debugf("purge removed resource %s", self.Name) - return self.purge(ctx, userCred) + err := self.purge(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) } if !lostNamePattern.MatchString(self.Name) { @@ -2527,6 +2535,13 @@ func (self *SGuest) syncWithCloudVM(ctx context.Context, userCred mcclient.Token db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } + syncVirtualResourceMetadata(ctx, userCred, self, extVM) SyncCloudProject(userCred, self, syncOwnerId, extVM, host.ManagerId) @@ -2630,6 +2645,11 @@ func (manager *SGuestManager) newCloudVM(ctx context.Context, userCred mcclient. db.OpsLog.LogEvent(&guest, db.ACT_CREATE, guest.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &guest, + Action: notifyclient.ActionSyncCreate, + }) + if guest.Status == api.VM_RUNNING { db.OpsLog.LogEvent(&guest, db.ACT_START, guest.GetShortDesc(ctx), userCred) } diff --git a/pkg/compute/models/kafka.go b/pkg/compute/models/kafka.go index 58b34d55c1..4d07c8a403 100644 --- a/pkg/compute/models/kafka.go +++ b/pkg/compute/models/kafka.go @@ -31,6 +31,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -389,7 +390,15 @@ func (self *SKafka) GetIKafka() (cloudprovider.ICloudKafka, error) { } func (self *SKafka) syncRemoveCloudKafka(ctx context.Context, userCred mcclient.TokenCredential) error { - return self.RealDelete(ctx, userCred) + err := self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } // 同步资源属性 @@ -469,6 +478,12 @@ func (self *SKafka) SyncWithCloudKafka(ctx context.Context, userCred mcclient.To if err != nil { return errors.Wrapf(err, "db.Update") } + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } syncVirtualResourceMetadata(ctx, userCred, self, ext) if provider := self.GetCloudprovider(); provider != nil { @@ -568,6 +583,10 @@ func (self *SCloudregion) newFromCloudKafka(ctx context.Context, userCred mcclie if err != nil { return nil, errors.Wrapf(err, "newFromCloudKafka.Insert") } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &kafka, + Action: notifyclient.ActionSyncCreate, + }) // 同步标签 syncVirtualResourceMetadata(ctx, userCred, &kafka, ext) diff --git a/pkg/compute/models/loadbalancercachedcertificates.go b/pkg/compute/models/loadbalancercachedcertificates.go index 9a2b744f7e..f5928c69bc 100644 --- a/pkg/compute/models/loadbalancercachedcertificates.go +++ b/pkg/compute/models/loadbalancercachedcertificates.go @@ -30,6 +30,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudcommon/validators" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" @@ -344,6 +345,10 @@ func (man *SCachedLoadbalancerCertificateManager) newFromCloudLoadbalancerCertif SyncCloudProject(userCred, &lbcert, projectId, extCertificate, lbcert.ManagerId) db.OpsLog.LogEvent(&lbcert, db.ACT_CREATE, lbcert.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &lbcert, + Action: notifyclient.ActionSyncCreate, + }) return &lbcert, nil } @@ -375,6 +380,12 @@ func (lbcert *SCachedLoadbalancerCertificate) SyncWithCloudLoadbalancerCertifica } } db.OpsLog.LogSyncUpdate(lbcert, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: lbcert, + Action: notifyclient.ActionSyncUpdate, + }) + } SyncCloudProject(userCred, lbcert, projectId, extCertificate, lbcert.ManagerId) @@ -390,6 +401,10 @@ func (lbcert *SCachedLoadbalancerCertificate) syncRemoveCloudLoadbalancerCertifi err = lbcert.SetStatus(userCred, api.LB_STATUS_UNKNOWN, "sync to delete") } else { err = lbcert.DoPendingDelete(ctx, userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: lbcert, + Action: notifyclient.ActionSyncDelete, + }) } return err } diff --git a/pkg/compute/models/loadbalancers.go b/pkg/compute/models/loadbalancers.go index 49aea50929..146bd6b89a 100644 --- a/pkg/compute/models/loadbalancers.go +++ b/pkg/compute/models/loadbalancers.go @@ -33,6 +33,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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" @@ -1022,6 +1023,11 @@ func (man *SLoadbalancerManager) newFromCloudLoadbalancer(ctx context.Context, u db.OpsLog.LogEvent(&lb, db.ACT_CREATE, lb.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &lb, + Action: notifyclient.ActionSyncCreate, + }) + lb.syncLoadbalancerNetwork(ctx, userCred, lbNetworkIds) return &lb, nil } @@ -1035,6 +1041,10 @@ func (lb *SLoadbalancer) syncRemoveCloudLoadbalancer(ctx context.Context, userCr return lb.SetStatus(userCred, api.LB_STATUS_UNKNOWN, "sync to delete") } else { lb.LBPendingDelete(ctx, userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: lb, + Action: notifyclient.ActionSyncDelete, + }) return nil } } @@ -1217,6 +1227,13 @@ func (lb *SLoadbalancer) SyncWithCloudLoadbalancer(ctx context.Context, userCred db.OpsLog.LogSyncUpdate(lb, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: lb, + Action: notifyclient.ActionSyncUpdate, + }) + } + networkIds := getExtLbNetworkIds(extLb, lb.ManagerId) SyncCloudProject(userCred, lb, syncOwnerId, extLb, provider.Id) lb.syncLoadbalancerNetwork(ctx, userCred, networkIds) diff --git a/pkg/compute/models/mongodb.go b/pkg/compute/models/mongodb.go index c99988d5db..acd537bcd4 100644 --- a/pkg/compute/models/mongodb.go +++ b/pkg/compute/models/mongodb.go @@ -32,6 +32,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/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -457,7 +458,15 @@ func (self *SCloudregion) SyncMongoDBs(ctx context.Context, userCred mcclient.To } func (self *SMongoDB) syncRemoveCloudMongoDB(ctx context.Context, userCred mcclient.TokenCredential) error { - return self.RealDelete(ctx, userCred) + err := self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SMongoDB) ValidateDeleteCondition(ctx context.Context, info jsonutils.JSONObject) error { @@ -529,6 +538,12 @@ func (self *SMongoDB) SyncWithCloudMongoDB(ctx context.Context, userCred mcclien if err != nil { return errors.Wrapf(err, "db.Update") } + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } syncVirtualResourceMetadata(ctx, userCred, self, ext) if provider := self.GetCloudprovider(); provider != nil { SyncCloudProject(userCred, self, provider.GetOwnerId(), ext, provider.Id) @@ -617,6 +632,10 @@ func (self *SCloudregion) newFromCloudMongoDB(ctx context.Context, userCred mccl if err != nil { return nil, errors.Wrapf(err, "newFromCloudMongoDB.Insert") } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &ins, + Action: notifyclient.ActionSyncCreate, + }) syncVirtualResourceMetadata(ctx, userCred, &ins, ext) SyncCloudProject(userCred, &ins, provider.GetOwnerId(), ext, provider.Id) diff --git a/pkg/compute/models/natgateways.go b/pkg/compute/models/natgateways.go index e6053c8421..3da0fffee0 100644 --- a/pkg/compute/models/natgateways.go +++ b/pkg/compute/models/natgateways.go @@ -32,6 +32,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudcommon/validators" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" @@ -454,7 +455,15 @@ func (self *SNatGateway) syncRemoveCloudNatGateway(ctx context.Context, userCred if err != nil { // cannot delete return self.SetStatus(userCred, api.NAT_STATUS_UNKNOWN, "sync to delete") } - return self.purge(ctx, userCred) + err = self.purge(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SNatGateway) ValidateDeleteCondition(ctx context.Context, info jsonutils.JSONObject) error { @@ -507,6 +516,12 @@ func (self *SNatGateway) SyncWithCloudNatGateway(ctx context.Context, userCred m SyncCloudDomain(userCred, self, provider.GetOwnerId()) db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } return nil } @@ -566,6 +581,10 @@ func (manager *SNatGatewayManager) newFromCloudNatGateway(ctx context.Context, u syncMetadata(ctx, userCred, &nat, extNat) db.OpsLog.LogEvent(&nat, db.ACT_CREATE, nat.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &nat, + Action: notifyclient.ActionSyncCreate, + }) return &nat, nil } diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index 3681e92e99..2127e59642 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -53,6 +53,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -655,7 +656,12 @@ func (self *SNetwork) syncRemoveCloudNetwork(ctx context.Context, userCred mccli err = self.SetStatus(userCred, api.NETWORK_STATUS_UNKNOWN, "Sync to remove") } else { err = self.RealDelete(ctx, userCred) - + if err == nil { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + } } return err } @@ -680,6 +686,12 @@ func (self *SNetwork) SyncWithCloudNetwork(ctx context.Context, userCred mcclien return err } db.OpsLog.LogSyncUpdate(self, diff, userCred) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } syncVirtualResourceMetadata(ctx, userCred, self, extNet) SyncCloudProject(userCred, self, syncOwnerId, extNet, vpc.ManagerId) @@ -751,6 +763,10 @@ func (manager *SNetworkManager) newFromCloudNetwork(ctx context.Context, userCre } db.OpsLog.LogEvent(&net, db.ACT_CREATE, net.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &net, + Action: notifyclient.ActionSyncCreate, + }) return &net, nil } diff --git a/pkg/compute/models/secgroupcache.go b/pkg/compute/models/secgroupcache.go index 58d951ba94..53e84147fb 100644 --- a/pkg/compute/models/secgroupcache.go +++ b/pkg/compute/models/secgroupcache.go @@ -31,6 +31,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" @@ -478,6 +479,10 @@ func (manager *SSecurityGroupCacheManager) SyncSecurityGroupCaches(ctx context.C syncResult.DeleteError(err) } else { syncResult.Delete() + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &removed[i], + Action: notifyclient.ActionSyncDelete, + }) } } diff --git a/pkg/compute/models/secgroups.go b/pkg/compute/models/secgroups.go index e76d457bda..af26fc555f 100644 --- a/pkg/compute/models/secgroups.go +++ b/pkg/compute/models/secgroups.go @@ -35,6 +35,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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -629,6 +630,10 @@ func (self *SSecurityGroup) PostCreate(ctx context.Context, userCred mcclient.To SecurityGroupRuleManager.TableSpec().Insert(ctx, rule) } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionCreate, + }) } func (manager *SSecurityGroupManager) FetchSecgroupById(secId string) (*SSecurityGroup, error) { @@ -1140,6 +1145,10 @@ func (manager *SSecurityGroupManager) newFromCloudSecgroup(ctx context.Context, rules, _ := secgroup.SyncSecurityGroupRules(ctx, userCred, dest) db.OpsLog.LogEvent(&secgroup, db.ACT_CREATE, secgroup.GetShortDesc(ctx), userCred) + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: &secgroup, + Action: notifyclient.ActionSyncCreate, + }) return &secgroup, rules, nil } diff --git a/pkg/compute/models/vpcs.go b/pkg/compute/models/vpcs.go index 1cbe89e3d5..67e45da80a 100644 --- a/pkg/compute/models/vpcs.go +++ b/pkg/compute/models/vpcs.go @@ -36,6 +36,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/validators" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" @@ -520,6 +521,12 @@ func (self *SVpc) syncRemoveCloudVpc(ctx context.Context, userCred mcclient.Toke } } else { err = self.RealDelete(ctx, userCred) + if err == nil { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionDelete, + }) + } } return err } diff --git a/pkg/compute/models/waf_instances.go b/pkg/compute/models/waf_instances.go index 4dd3b7654b..407be3e6a5 100644 --- a/pkg/compute/models/waf_instances.go +++ b/pkg/compute/models/waf_instances.go @@ -27,6 +27,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudcommon/validators" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" @@ -382,7 +383,15 @@ func (self *SWafInstance) RealDelete(ctx context.Context, userCred mcclient.Toke } func (self *SWafInstance) syncRemove(ctx context.Context, userCred mcclient.TokenCredential) error { - return self.RealDelete(ctx, userCred) + err := self.RealDelete(ctx, userCred) + if err != nil { + return err + } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncDelete, + }) + return nil } func (self *SWafInstance) GetIRegion() (cloudprovider.ICloudRegion, error) { @@ -409,13 +418,19 @@ func (self *SWafInstance) GetICloudWafInstance() (cloudprovider.ICloudWafInstanc } func (self *SWafInstance) SyncWithCloudWafInstance(ctx context.Context, userCred mcclient.TokenCredential, ext cloudprovider.ICloudWafInstance) error { - _, err := db.Update(self, func() error { + diff, err := db.Update(self, func() error { self.ExternalId = ext.GetGlobalId() self.SetEnabled(ext.GetEnabled()) self.DefaultAction = ext.GetDefaultAction() self.Status = ext.GetStatus() return nil }) + if len(diff) > 0 { + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionSyncUpdate, + }) + } return err } @@ -444,6 +459,10 @@ func (self *SCloudregion) newFromCloudWafInstance(ctx context.Context, userCred if err != nil { return nil, err } + notifyclient.EventNotify(ctx, userCred, notifyclient.SEventNotifyParam{ + Obj: waf, + Action: notifyclient.ActionSyncCreate, + }) return waf, nil } diff --git a/pkg/compute/tasks/cdn_domain_delete_task.go b/pkg/compute/tasks/cdn_domain_delete_task.go index 067f3263a6..b05e79b7fd 100644 --- a/pkg/compute/tasks/cdn_domain_delete_task.go +++ b/pkg/compute/tasks/cdn_domain_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -68,5 +69,10 @@ func (self *CDNDomainDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneM return } + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: domain, + Action: notifyclient.ActionDelete, + }) + self.taskComplete(ctx, domain) } diff --git a/pkg/compute/tasks/disk_create_task.go b/pkg/compute/tasks/disk_create_task.go index 31a7bfb2c0..e2e5dd0cf7 100644 --- a/pkg/compute/tasks/disk_create_task.go +++ b/pkg/compute/tasks/disk_create_task.go @@ -26,6 +26,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" ) @@ -106,6 +107,10 @@ func (self *DiskCreateTask) OnDiskReady(ctx context.Context, disk *models.SDisk, disk.SetStatus(self.UserCred, api.DISK_READY, "") self.CleanHostSchedCache(disk) db.OpsLog.LogEvent(disk, db.ACT_ALLOCATE, disk.GetShortDesc(ctx), self.UserCred) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: disk, + Action: notifyclient.ActionCreate, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/disk_delete_task.go b/pkg/compute/tasks/disk_delete_task.go index 465dc62448..b80a53f1ff 100644 --- a/pkg/compute/tasks/disk_delete_task.go +++ b/pkg/compute/tasks/disk_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/compute/options" "yunion.io/x/onecloud/pkg/util/logclient" @@ -176,6 +177,10 @@ func (self *DiskDeleteTask) OnGuestDiskDeleteComplete(ctx context.Context, obj d disk := obj.(*models.SDisk) self.CleanHostSchedCache(disk) db.OpsLog.LogEvent(disk, db.ACT_DELOCATE, disk.GetShortDesc(ctx), self.UserCred) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: disk, + Action: notifyclient.ActionDelete, + }) if len(disk.SnapshotId) > 0 && disk.GetMetadata("merge_snapshot", nil) == "true" { models.SnapshotManager.AddRefCount(disk.SnapshotId, -1) } diff --git a/pkg/compute/tasks/dnszone_create_task.go b/pkg/compute/tasks/dnszone_create_task.go index 0def7d906e..9796728cab 100644 --- a/pkg/compute/tasks/dnszone_create_task.go +++ b/pkg/compute/tasks/dnszone_create_task.go @@ -106,7 +106,10 @@ func (self *DnsZoneCreateTask) OnInit(ctx context.Context, obj db.IStandaloneMod func (self *DnsZoneCreateTask) OnSyncRecordSetComplete(ctx context.Context, dnsZone *models.SDnsZone, data jsonutils.JSONObject) { dnsZone.SetStatus(self.GetUserCred(), api.DNS_ZONE_STATUS_AVAILABLE, "") - notifyclient.NotifyWebhook(ctx, self.UserCred, dnsZone, notifyclient.ActionCreate) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: dnsZone, + Action: notifyclient.ActionCreate, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/dnszone_delete_task.go b/pkg/compute/tasks/dnszone_delete_task.go index ae1c42f5e5..54222511d6 100644 --- a/pkg/compute/tasks/dnszone_delete_task.go +++ b/pkg/compute/tasks/dnszone_delete_task.go @@ -75,6 +75,9 @@ func (self *DnsZoneDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneMod dnsZone.RealDelete(ctx, self.GetUserCred()) logclient.AddActionLogWithContext(ctx, dnsZone, logclient.ACT_DELETE, nil, self.UserCred, true) - notifyclient.NotifyWebhook(ctx, self.UserCred, dnsZone, notifyclient.ActionDelete) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: dnsZone, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/elastic_search_delete_task.go b/pkg/compute/tasks/elastic_search_delete_task.go index 35438601a1..4b183f68be 100644 --- a/pkg/compute/tasks/elastic_search_delete_task.go +++ b/pkg/compute/tasks/elastic_search_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -68,5 +69,9 @@ func (self *ElasticSearchDeleteTask) OnInit(ctx context.Context, obj db.IStandal func (self *ElasticSearchDeleteTask) taskComplete(ctx context.Context, es *models.SElasticSearch) { es.RealDelete(ctx, self.GetUserCred()) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: es, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/filesystem_create_task.go b/pkg/compute/tasks/filesystem_create_task.go index e382ec697c..a0a2c16af8 100644 --- a/pkg/compute/tasks/filesystem_create_task.go +++ b/pkg/compute/tasks/filesystem_create_task.go @@ -40,6 +40,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -110,6 +111,10 @@ func (self *FileSystemCreateTask) OnInit(ctx context.Context, obj db.IStandalone func (self *FileSystemCreateTask) OnSyncstatusComplete(ctx context.Context, fs *models.SFileSystem, data jsonutils.JSONObject) { self.SetStageComplete(ctx, nil) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionCreate, + }) } func (self *FileSystemCreateTask) OnSyncstatusCompleteFailed(ctx context.Context, fs *models.SFileSystem, data jsonutils.JSONObject) { diff --git a/pkg/compute/tasks/filesystem_delete_task.go b/pkg/compute/tasks/filesystem_delete_task.go index 3a8f66cc48..90cda90eb2 100644 --- a/pkg/compute/tasks/filesystem_delete_task.go +++ b/pkg/compute/tasks/filesystem_delete_task.go @@ -38,6 +38,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -102,5 +103,9 @@ func (self *FileSystemDeleteTask) OnInit(ctx context.Context, obj db.IStandalone func (self *FileSystemDeleteTask) taskComplete(ctx context.Context, fs *models.SFileSystem) { fs.RealDelete(ctx, self.GetUserCred()) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/kafka_delete_task.go b/pkg/compute/tasks/kafka_delete_task.go index e5f1b1e9b3..64f5bf18e1 100644 --- a/pkg/compute/tasks/kafka_delete_task.go +++ b/pkg/compute/tasks/kafka_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -67,5 +68,9 @@ func (self *KafkaDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel func (self *KafkaDeleteTask) taskComplete(ctx context.Context, kafka *models.SKafka) { kafka.RealDelete(ctx, self.GetUserCred()) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: self, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/mongodb_delete_task.go b/pkg/compute/tasks/mongodb_delete_task.go index abc7feae2f..5fa4ea6b2d 100644 --- a/pkg/compute/tasks/mongodb_delete_task.go +++ b/pkg/compute/tasks/mongodb_delete_task.go @@ -38,6 +38,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -80,5 +81,9 @@ func (self *MongoDBDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneMod func (self *MongoDBDeleteTask) taskComplete(ctx context.Context, mongodb *models.SMongoDB) { mongodb.RealDelete(ctx, self.GetUserCred()) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: mongodb, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/nat_create_task.go b/pkg/compute/tasks/nat_create_task.go index 05079324e2..a02eebc224 100644 --- a/pkg/compute/tasks/nat_create_task.go +++ b/pkg/compute/tasks/nat_create_task.go @@ -186,6 +186,10 @@ func (self *NatGatewayCreateTask) OnDeployEipComplete(ctx context.Context, nat * } func (self *NatGatewayCreateTask) OnSyncstatusComplete(ctx context.Context, nat *models.SNatGateway, data jsonutils.JSONObject) { + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: nat, + Action: notifyclient.ActionCreate, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/nat_delete_task.go b/pkg/compute/tasks/nat_delete_task.go index 6c7124e4f4..f715b90a40 100644 --- a/pkg/compute/tasks/nat_delete_task.go +++ b/pkg/compute/tasks/nat_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -132,5 +133,9 @@ func (self *NatGatewayDeleteTask) doDeleteNatGateway(ctx context.Context, nat *m func (self *NatGatewayDeleteTask) taskComplete(ctx context.Context, nat *models.SNatGateway) { nat.RealDelete(ctx, self.GetUserCred()) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: nat, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/vpc_create_task.go b/pkg/compute/tasks/vpc_create_task.go index 10f4d043ff..e7738c23e0 100644 --- a/pkg/compute/tasks/vpc_create_task.go +++ b/pkg/compute/tasks/vpc_create_task.go @@ -23,6 +23,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" ) @@ -61,6 +62,10 @@ func (self *VpcCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, func (self *VpcCreateTask) OnCreateVpcComplete(ctx context.Context, vpc *models.SVpc, data jsonutils.JSONObject) { logclient.AddActionLogWithStartable(self, vpc, logclient.ACT_ALLOCATE, nil, self.UserCred, true) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: vpc, + Action: notifyclient.ActionCreate, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/vpc_delete_task.go b/pkg/compute/tasks/vpc_delete_task.go index 6eeed07287..2ade6554df 100644 --- a/pkg/compute/tasks/vpc_delete_task.go +++ b/pkg/compute/tasks/vpc_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" ) @@ -69,6 +70,10 @@ func (self *VpcDeleteTask) OnDeleteVpcComplete(ctx context.Context, vpc *models. } db.OpsLog.LogEvent(vpc, db.ACT_DELOCATING, vpc.GetShortDesc(ctx), self.UserCred) logclient.AddActionLogWithStartable(self, vpc, logclient.ACT_DELETE, nil, self.UserCred, true) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: vpc, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/waf_delete_task.go b/pkg/compute/tasks/waf_delete_task.go index fb95b2a2c7..a7513471cb 100644 --- a/pkg/compute/tasks/waf_delete_task.go +++ b/pkg/compute/tasks/waf_delete_task.go @@ -24,6 +24,7 @@ import ( api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/logclient" @@ -69,5 +70,9 @@ func (self *WafDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel, func (self *WafDeleteTask) taskComplete(ctx context.Context, waf *models.SWafInstance) { waf.RealDelete(ctx, self.GetUserCred()) + notifyclient.EventNotify(ctx, self.UserCred, notifyclient.SEventNotifyParam{ + Obj: waf, + Action: notifyclient.ActionDelete, + }) self.SetStageComplete(ctx, nil) } diff --git a/pkg/notify/models/event_template.go b/pkg/notify/models/event_template.go index 2f1d475374..15b22d7e12 100644 --- a/pkg/notify/models/event_template.go +++ b/pkg/notify/models/event_template.go @@ -434,6 +434,56 @@ func init() { "snapshot policy", "快照策略", }, + sI18nElme{ + api.TOPIC_RESOURCE_VPC, + "VPC", + "VPC", + }, + sI18nElme{ + api.TOPIC_RESOURCE_DNSZONE, + "DNS zone", + "DNS zone", + }, + sI18nElme{ + api.TOPIC_RESOURCE_NATGATEWAY, + "nat gateway", + "nat网关", + }, + sI18nElme{ + api.TOPIC_RESOURCE_WEBAPP, + "webapp", + "应用程序服务", + }, + sI18nElme{ + api.TOPIC_RESOURCE_CDNDOMAIN, + "CDN domain", + "CDN domain", + }, + sI18nElme{ + api.TOPIC_RESOURCE_FILESYSTEM, + "file system", + "文件系统", + }, + sI18nElme{ + api.TOPIC_RESOURCE_WAF, + "WAF", + "WAF", + }, + sI18nElme{ + api.TOPIC_RESOURCE_KAFKA, + "Kafka", + "Kafka", + }, + sI18nElme{ + api.TOPIC_RESOURCE_ELASTICSEARCH, + "Elasticsearch", + "Elasticsearch", + }, + sI18nElme{ + api.TOPIC_RESOURCE_MONGODB, + "MongoDB", + "MongoDB", + }, sI18nElme{ string(api.ActionCreate), "created", diff --git a/pkg/notify/models/topic.go b/pkg/notify/models/topic.go index 8797c8626b..c7ad16aaf8 100644 --- a/pkg/notify/models/topic.go +++ b/pkg/notify/models/topic.go @@ -88,6 +88,7 @@ const ( DefaultScalingPolicyExecute = "scaling policy execute" DefaultSnapshotPolicyExecute = "snapshot policy execute" DefaultResourceOperationFailed = "resource operation failed" + DefaultResourceSync = "resource sync" ) func (sm *STopicManager) InitializeData() error { @@ -101,6 +102,7 @@ func (sm *STopicManager) InitializeData() error { DefaultScalingPolicyExecute, DefaultSnapshotPolicyExecute, DefaultResourceOperationFailed, + DefaultResourceSync, ) q := sm.Query() topics := make([]STopic, 0, initSNames.Len()) @@ -108,12 +110,17 @@ func (sm *STopicManager) InitializeData() error { if err != nil { return errors.Wrap(err, "unable to FetchModelObjects") } + nameTopicMap := make(map[string]*STopic, len(topics)) for i := range topics { t := &topics[i] initSNames.Delete(t.Name) + nameTopicMap[t.Name] = t + } + for _, name := range initSNames.UnsortedList() { + nameTopicMap[name] = nil } ctx := context.Background() - for _, name := range initSNames.UnsortedList() { + for name, topic := range nameTopicMap { t := new(STopic) t.Name = name t.Enabled = tristate.True @@ -137,6 +144,14 @@ func (sm *STopicManager) InitializeData() error { notify.TOPIC_RESOURCE_ELASTICCACHE, notify.TOPIC_RESOURCE_BAREMETAL, notify.TOPIC_RESOURCE_SECGROUP, + notify.TOPIC_RESOURCE_FILESYSTEM, + notify.TOPIC_RESOURCE_NATGATEWAY, + notify.TOPIC_RESOURCE_VPC, + notify.TOPIC_RESOURCE_CDNDOMAIN, + notify.TOPIC_RESOURCE_WAF, + notify.TOPIC_RESOURCE_KAFKA, + notify.TOPIC_RESOURCE_ELASTICSEARCH, + notify.TOPIC_RESOURCE_MONGODB, ) t.addAction( notify.ActionCreate, @@ -217,10 +232,50 @@ func (sm *STopicManager) InitializeData() error { notify.ActionMigrate, ) t.Type = notify.TOPIC_TYPE_RESOURCE + case DefaultResourceSync: + t.addResources( + notify.TOPIC_RESOURCE_SERVER, + notify.TOPIC_RESOURCE_DISK, + notify.TOPIC_RESOURCE_DBINSTANCE, + notify.TOPIC_RESOURCE_ELASTICCACHE, + notify.TOPIC_RESOURCE_LOADBALANCER, + notify.TOPIC_RESOURCE_EIP, + notify.TOPIC_RESOURCE_VPC, + notify.TOPIC_RESOURCE_NETWORK, + notify.TOPIC_RESOURCE_LOADBALANCERCERTIFICATE, + notify.TOPIC_RESOURCE_DNSZONE, + notify.TOPIC_RESOURCE_NATGATEWAY, + notify.TOPIC_RESOURCE_BUCKET, + notify.TOPIC_RESOURCE_FILESYSTEM, + notify.TOPIC_RESOURCE_WEBAPP, + notify.TOPIC_RESOURCE_CDNDOMAIN, + notify.TOPIC_RESOURCE_WAF, + notify.TOPIC_RESOURCE_KAFKA, + notify.TOPIC_RESOURCE_ELASTICSEARCH, + notify.TOPIC_RESOURCE_MONGODB, + ) + t.addAction( + notify.ActionSyncCreate, + notify.ActionSyncUpdate, + notify.ActionSyncDelete, + ) + t.Type = notify.TOPIC_TYPE_RESOURCE } - err := sm.TableSpec().Insert(ctx, t) - if err != nil { - return errors.Wrapf(err, "unable to insert %s", name) + if topic == nil { + err := sm.TableSpec().Insert(ctx, t) + if err != nil { + return errors.Wrapf(err, "unable to insert %s", name) + } + } else { + _, err := db.Update(topic, func() error { + topic.Resources = t.Resources + topic.Actions = t.Actions + topic.Type = t.Type + return nil + }) + if err != nil { + return errors.Wrapf(err, "unable to update topic %s", topic.Name) + } } } return nil @@ -416,6 +471,16 @@ func init() { notify.TOPIC_RESOURCE_ELASTICCACHE: 16, notify.TOPIC_RESOURCE_SCHEDULEDTASK: 17, notify.TOPIC_RESOURCE_BAREMETAL: 18, + notify.TOPIC_RESOURCE_VPC: 19, + notify.TOPIC_RESOURCE_DNSZONE: 20, + notify.TOPIC_RESOURCE_NATGATEWAY: 21, + notify.TOPIC_RESOURCE_WEBAPP: 22, + notify.TOPIC_RESOURCE_CDNDOMAIN: 23, + notify.TOPIC_RESOURCE_FILESYSTEM: 24, + notify.TOPIC_RESOURCE_WAF: 25, + notify.TOPIC_RESOURCE_KAFKA: 26, + notify.TOPIC_RESOURCE_ELASTICSEARCH: 27, + notify.TOPIC_RESOURCE_MONGODB: 28, }, ) converter.registerAction( @@ -435,6 +500,9 @@ func init() { notify.ActionMigrate: 12, notify.ActionCreateBackupServer: 13, notify.ActionDelBackupServer: 14, + notify.ActionSyncCreate: 15, + notify.ActionSyncUpdate: 16, + notify.ActionSyncDelete: 17, }, ) }