diff --git a/pkg/cloudevent/models/cloudevents.go b/pkg/cloudevent/models/cloudevents.go index 6a3842306b..453f305d7f 100644 --- a/pkg/cloudevent/models/cloudevents.go +++ b/pkg/cloudevent/models/cloudevents.go @@ -16,34 +16,27 @@ package models import ( "context" - "database/sql" - "fmt" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" - "yunion.io/x/pkg/errors" "yunion.io/x/sqlchemy" api "yunion.io/x/onecloud/pkg/apis/cloudevent" "yunion.io/x/onecloud/pkg/cloudcommon/db" - "yunion.io/x/onecloud/pkg/cloudevent/options" "yunion.io/x/onecloud/pkg/cloudprovider" - "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" - "yunion.io/x/onecloud/pkg/mcclient/auth" - "yunion.io/x/onecloud/pkg/mcclient/modulebase" ) type SCloudeventManager struct { - db.SVirtualResourceBaseManager + db.SModelBaseManager } var CloudeventManager *SCloudeventManager -var mods map[string]modulebase.Manager func init() { CloudeventManager = &SCloudeventManager{ - SVirtualResourceBaseManager: db.NewVirtualResourceBaseManager( + SModelBaseManager: db.NewModelBaseManager( SCloudevent{}, "cloudevents_tbl", "cloudevent", @@ -54,7 +47,7 @@ func init() { } type SCloudevent struct { - db.SVirtualResourceBase + db.SModelBase Service string `width:"64" charset:"utf8" nullable:"true" list:"user"` ResourceType string `width:"64" charset:"utf8" nullable:"true" list:"user"` @@ -63,8 +56,11 @@ type SCloudevent struct { Request jsonutils.JSONObject `charset:"utf8" nullable:"true" list:"user"` Account string `width:"64" charset:"utf8" nullable:"true" list:"user"` Success bool `nullable:"false" list:"user"` + CreatedAt time.Time `nullable:"false" created_at:"true" index:"true" get:"user" list:"user"` CloudproviderId string `width:"64" charset:"utf8" nullable:"true" list:"user"` + Manager string `width:"128" charset:"utf8" nullable:"false" index:"true" list:"user"` + Provider string `width:"64" charset:"ascii" nullable:"false" list:"user"` } func (self *SCloudeventManager) AllowCreateItem(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { @@ -80,116 +76,20 @@ func (self *SCloudevent) AllowUpdateItem(ctx context.Context, userCred mcclient. } func (manager *SCloudeventManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, input *api.CloudeventListInput) (*sqlchemy.SQuery, error) { - q, err := manager.SVirtualResourceBaseManager.ListItemFilter(ctx, q, userCred, input.JSON(input)) + q, err := manager.SModelBaseManager.ListItemFilter(ctx, q, userCred, input.JSON(input)) if err != nil { return nil, err } - if len(input.Cloudprovider) > 0 { - providerObj, err := CloudproviderManager.FetchByIdOrName(userCred, input.Cloudprovider) - if err != nil { - if err == sql.ErrNoRows { - return nil, httperrors.NewResourceNotFoundError2(CloudproviderManager.Keyword(), input.Cloudprovider) - } else { - return nil, httperrors.NewGeneralError(err) - } - } - q = q.Equals("cloudprovider_id", providerObj.GetId()) - } if len(input.Providers) > 0 { - sq := CloudproviderManager.Query().SubQuery() - q = q.Join(sq, sqlchemy.Equals(q.Field("cloudprovider_id"), sq.Field("id"))). - Filter(sqlchemy.In(sq.Field("provider"), input.Providers)) + q = q.In("provider", input.Providers) } - //过滤已删除的cloudprovider日志 - sq := CloudproviderManager.Query("id").SubQuery() - q = q.In("cloudprovider_id", sq) + return q, nil } -func (self *SCloudevent) GetCloudprovider() (*SCloudprovider, error) { - cloudprovider, err := CloudproviderManager.FetchById(self.CloudproviderId) - if err != nil { - return nil, err - } - return cloudprovider.(*SCloudprovider), nil -} - func (self *SCloudevent) GetCustomizeColumns(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) *jsonutils.JSONDict { - extra := self.SStatusStandaloneResourceBase.GetCustomizeColumns(ctx, userCred, query) - extra, _ = self.getMoreDetails(ctx, userCred, query, extra) - return extra -} - -func (self *SCloudevent) getMoreDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, extra *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { - cloudprovider, err := self.GetCloudprovider() - if err != nil { - return nil, err - } - info := jsonutils.Marshal(map[string]string{ - "provider": cloudprovider.Provider, - "manager": cloudprovider.Name, - }) - extra.Update(info) - return extra, nil -} - -func (self *SCloudevent) GetExtraDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (*jsonutils.JSONDict, error) { - extra, err := self.SVirtualResourceBase.GetExtraDetails(ctx, userCred, query) - if err != nil { - return nil, err - } - return self.getMoreDetails(ctx, userCred, query, extra) -} - -func (self *SCloudeventManager) fetchMods(ctx context.Context, userCred mcclient.TokenCredential) { - if len(mods) > 0 { - return - } - s := auth.GetAdminSession(ctx, options.Options.Region, "v2") - mods = map[string]modulebase.Manager{} - rms, _ := modulebase.GetRegisterdModules() - for _, _mods := range rms { - for _, _mod := range _mods { - if _, ok := mods[_mod]; !ok { - mod, err := modulebase.GetModule(s, _mod) - if err != nil { - log.Errorf("failed to get mod %s error: %v", _mod, err) - continue - } - mods[mod.GetKeyword()] = mod - } - } - } - return -} - -func (manager *SCloudeventManager) setEventInfo(session *mcclient.ClientSession, mod modulebase.Manager, event *SCloudevent) error { - params := jsonutils.NewDict() - params.Add(jsonutils.NewString(fmt.Sprintf("external_id.equals(%s)", event.Name)), "filter") - result, err := mod.List(session, params) - if err != nil { - return errors.Wrapf(err, "mod.List for %s by externalId: %s", mod.KeyString(), event.Name) - } - if len(result.Data) != 1 { - return errors.Wrapf(err, "found %d %s by externalId: %s", len(result.Data), mod.KeyString(), event.Name) - } - - data := struct { - Name string - TenantId string - }{} - err = result.Data[0].Unmarshal(&data) - if err != nil { - return errors.Wrapf(err, "result.Data[0].Unmarshal %s", result.Data[0]) - } - if len(data.Name) > 0 { - event.Name = event.Name - } - if len(data.TenantId) > 0 { - event.ProjectId = data.TenantId - } - return nil + return self.SModelBase.GetCustomizeColumns(ctx, userCred, query) } func (manager *SCloudeventManager) SyncCloudevent(ctx context.Context, userCred mcclient.TokenCredential, cloudprovider *SCloudprovider, iEvents []cloudprovider.ICloudEvent) int { @@ -203,16 +103,12 @@ func (manager *SCloudeventManager) SyncCloudevent(ctx context.Context, userCred RequestId: iEvent.GetRequestId(), Request: iEvent.GetRequest(), Success: iEvent.IsSuccess(), + Manager: cloudprovider.Name, + Provider: cloudprovider.Provider, CloudproviderId: cloudprovider.Id, } - event.Name = iEvent.GetName() - event.Status = "ready" - event.ProjectId = userCred.GetProjectId() event.CreatedAt = iEvent.GetCreatedAt() - event.ProjectId = userCred.GetProjectId() - event.DomainId = userCred.GetDomainId() - event.SetModelManager(manager, event) err := manager.TableSpec().Insert(event) if err != nil { @@ -223,3 +119,11 @@ func (manager *SCloudeventManager) SyncCloudevent(ctx context.Context, userCred } return count } + +func (manager *SCloudeventManager) GetPagingConfig() *db.SPagingConfig { + return &db.SPagingConfig{ + Order: sqlchemy.SQL_ORDER_DESC, + MarkerField: "created_at", + DefaultLimit: 20, + } +} diff --git a/pkg/cloudevent/models/cloudproviders.go b/pkg/cloudevent/models/cloudproviders.go index d5e7f61a0d..a95012f846 100644 --- a/pkg/cloudevent/models/cloudproviders.go +++ b/pkg/cloudevent/models/cloudproviders.go @@ -56,10 +56,11 @@ func init() { type SCloudprovider struct { db.SEnabledStatusStandaloneResourceBase - HealthStatus string `width:"16" charset:"ascii" default:"normal" nullable:"false" list:"domain"` - SyncStatus string - LastSync time.Time - LastSyncEndAt time.Time + HealthStatus string `width:"16" charset:"ascii" default:"normal" nullable:"false" list:"domain"` + SyncStatus string + LastSync time.Time + LastSyncEndAt time.Time + LastSyncTimeAt time.Time AccessUrl string `width:"64" charset:"ascii" nullable:"true" list:"domain" update:"domain"` Account string `width:"128" charset:"ascii" nullable:"false" list:"domain"` @@ -182,6 +183,18 @@ func (self *SCloudprovider) MarkEndSync(userCred mcclient.TokenCredential) error return nil } +func (self *SCloudprovider) SetLastSyncTimeAt(userCred mcclient.TokenCredential, last time.Time) error { + _, err := db.Update(self, func() error { + self.LastSyncTimeAt = last + return nil + }) + if err != nil { + log.Errorf("Failed to SetLastSyncTimeAt error: %v", err) + return err + } + return nil +} + func (manager *SCloudproviderManager) newFromRegionProvider(ctx context.Context, userCred mcclient.TokenCredential, cloudprovider SCloudprovider) error { cloudprovider.SyncStatus = api.CLOUD_PROVIDER_SYNC_STATUS_IDLE return manager.TableSpec().Insert(&cloudprovider) @@ -226,19 +239,23 @@ func (self *SCloudprovider) GetNextTimeRange() (time.Time, time.Time, error) { if err != nil { return start, end, errors.Wrap(err, "q.CountWithError") } - if count == 0 { + if !self.LastSyncTimeAt.IsZero() { + start = self.LastSyncTimeAt + } else if count == 0 { start = time.Now().AddDate(0, 0, -1*factory.GetMaxCloudEventKeepDays()) } else { - provider := SCloudprovider{} - err = q.First(&provider) + event := &SCloudevent{} + err = q.First(event) if err != nil { return start, end, errors.Wrap(err, "q.First") } - start = provider.CreatedAt - if start.Before(time.Now().AddDate(0, 0, factory.GetMaxCloudEventKeepDays()*-1)) { - start = time.Now().AddDate(0, 0, factory.GetMaxCloudEventKeepDays()*-1) - } + start = event.CreatedAt } + // 避免cloudevent过长时间未运行,再次运行时记录的最后一条时间距离现在间隔太长 + if start.Before(time.Now().AddDate(0, 0, factory.GetMaxCloudEventKeepDays()*-1)) { + start = time.Now().AddDate(0, 0, factory.GetMaxCloudEventKeepDays()*-1) + } + if options.Options.OneSyncForHours > factory.GetMaxCloudEventSyncDays()*24 { end = start.Add(time.Duration(factory.GetMaxCloudEventSyncDays()*24) * time.Hour) } else { diff --git a/pkg/cloudevent/tasks/cloudevent_sync_task.go b/pkg/cloudevent/tasks/cloudevent_sync_task.go index 45fcb758e0..84398e3f3a 100644 --- a/pkg/cloudevent/tasks/cloudevent_sync_task.go +++ b/pkg/cloudevent/tasks/cloudevent_sync_task.go @@ -97,6 +97,7 @@ func (self *CloudeventSyncTask) OnInit(ctx context.Context, obj db.IStandaloneMo log.Infof("Sync %d events for %s(%s) from %s(%d) hours", _count, provider.Name, provider.Id, start.Format("2006-01-02T15:04:05Z"), duration/time.Hour) count += _count + provider.SetLastSyncTimeAt(self.UserCred, end) if time.Now().Sub(end) < duration { break } diff --git a/pkg/mcclient/modules/mod_cloudevents.go b/pkg/mcclient/modules/mod_cloudevents.go index 4d03d8a6ff..63ecd33180 100644 --- a/pkg/mcclient/modules/mod_cloudevents.go +++ b/pkg/mcclient/modules/mod_cloudevents.go @@ -22,8 +22,8 @@ var ( func init() { Cloudevents = NewCloudeventManager("cloudevent", "cloudevents", - []string{"ID", "Name", "Status", "Service", "Success", - "Resource_Type", "Action", "Cloudprovider_Id", "Cloudprovider"}, + []string{"Action", "Service", "Success", + "Resource_Type", "Cloudprovider_Id", "Manager", "Provider"}, []string{}) register(&Cloudevents)