diff --git a/pkg/image/models/image_subs.go b/pkg/image/models/image_subs.go index 0d44c8cd1b..1a91f1a57b 100644 --- a/pkg/image/models/image_subs.go +++ b/pkg/image/models/image_subs.go @@ -389,3 +389,18 @@ func (self *SImageSubformat) SetStatusSeeding(seeding bool) { torrent.SetTorrentSeeding(filePath, seeding) } } + +func (subimg *SImageSubformat) verifyStatusSelf(ctx context.Context) error { + if len(subimg.Location) == 0 { + return nil + } + filePath := subimg.Location + _, rc, err := GetImage(ctx, filePath) + if err != nil { + subimg.SetStatus(api.IMAGE_STATUS_UNKNOWN) + return errors.Wrap(err, "GetImage") + } + defer rc.Close() + subimg.SetStatus(api.IMAGE_STATUS_ACTIVE) + return nil +} diff --git a/pkg/image/models/images.go b/pkg/image/models/images.go index 8ec70a9f2b..0e6f1db125 100644 --- a/pkg/image/models/images.go +++ b/pkg/image/models/images.go @@ -2328,3 +2328,69 @@ func (img *SImage) markDataImage(userCred mcclient.TokenCredential) error { db.OpsLog.LogEvent(img, db.ACT_UPDATE, diff, userCred) return nil } + +func (manager *SImageManager) FetchImages(filter func(q *sqlchemy.SQuery) *sqlchemy.SQuery) ([]SImage, error) { + q := manager.Query() + if filter != nil { + q = filter(q) + } + images := make([]SImage, 0) + err := db.FetchModelObjects(manager, q, &images) + if err != nil { + return nil, errors.Wrap(err, "db.FetchModelObjects") + } + return images, nil +} + +func (manager *SImageManager) VerifyActiveImageStatus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { + images, err := manager.FetchImages(func(q *sqlchemy.SQuery) *sqlchemy.SQuery { + return q.In("status", []string{api.IMAGE_STATUS_ACTIVE, api.IMAGE_STATUS_UNKNOWN}) + }) + if err != nil { + log.Errorf("FetchImages failed: %s", err) + return + } + for i := range images { + img := &images[i] + err := img.verifyStatus(ctx, userCred) + if err != nil { + log.Errorf("VerifyStatus %s(%s) failed: %s", img.Id, img.Name, err) + } + } +} + +func (img *SImage) verifyStatus(ctx context.Context, userCred mcclient.TokenCredential) error { + errs := make([]error, 0) + subImages := ImageSubformatManager.GetAllSubImages(img.Id) + for i := range subImages { + err := subImages[i].verifyStatusSelf(ctx) + if err != nil { + errs = append(errs, err) + } + } + err := img.verifyStatusSelf(ctx, userCred) + if err != nil { + errs = append(errs, err) + } + if len(errs) > 0 { + return errors.NewAggregate(errs) + } + return nil +} + +func (img *SImage) verifyStatusSelf(ctx context.Context, userCred mcclient.TokenCredential) error { + if len(img.Location) == 0 { + return nil + } + filePath := img.Location + _, rc, err := GetImage(ctx, filePath) + if err != nil { + img.SetStatus(ctx, userCred, api.IMAGE_STATUS_UNKNOWN, errors.Wrap(err, "verifyStatusSelf").Error()) + return errors.Wrap(err, "GetImage") + } + defer rc.Close() + if img.Status != api.IMAGE_STATUS_ACTIVE { + img.SetStatus(ctx, userCred, api.IMAGE_STATUS_ACTIVE, "verifyStatusSelf") + } + return nil +} diff --git a/pkg/image/options/options.go b/pkg/image/options/options.go index 52f6a26616..3f291d69f7 100644 --- a/pkg/image/options/options.go +++ b/pkg/image/options/options.go @@ -55,6 +55,8 @@ type SImageOptions struct { S3CheckImageStatus bool `help:"Enable s3 check image status"` ImageStreamWorkerCount int `help:"Image stream worker count" default:"10"` + + VerifyImageStatusIntervalMinutes int `help:"verify image status periodically, default 15 minutes" default:"15"` } var ( diff --git a/pkg/image/service/service.go b/pkg/image/service/service.go index 8ed467375a..25d005a0bb 100644 --- a/pkg/image/service/service.go +++ b/pkg/image/service/service.go @@ -158,6 +158,8 @@ func StartService() { cron.AddJobEveryFewHour("AutoPurgeSplitable", 4, 30, 0, db.AutoPurgeSplitable, false) + cron.AddJobAtIntervals("MarkDataImage", time.Duration(options.Options.VerifyImageStatusIntervalMinutes)*time.Minute, models.ImageManager.VerifyActiveImageStatus) + cron.Start() }