diff --git a/pkg/cloudcommon/db/lockman/example/example b/pkg/cloudcommon/db/lockman/example/example new file mode 100755 index 0000000000..b515a7e452 Binary files /dev/null and b/pkg/cloudcommon/db/lockman/example/example differ diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index 52b7936e55..e9d9322122 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -6,7 +6,6 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/pkg/util/secrules" - "context" ) type ICloudResource interface { @@ -77,7 +76,7 @@ type ICloudStoragecache interface { DownloadImage(userCred mcclient.TokenCredential, imageId string, extId string) (jsonutils.JSONObject, error) - UploadImage(ctx context.Context, userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, extId string, isForce bool) (string, error) + UploadImage(userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, extId string, isForce bool) (string, error) } type ICloudStorage interface { diff --git a/pkg/compute/hostdrivers/aliyun.go b/pkg/compute/hostdrivers/aliyun.go index 5901e8fe9f..15842e7a9b 100644 --- a/pkg/compute/hostdrivers/aliyun.go +++ b/pkg/compute/hostdrivers/aliyun.go @@ -10,6 +10,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" ) type SAliyunHostDriver struct { @@ -24,7 +25,7 @@ func (self *SAliyunHostDriver) GetHostType() string { return models.HOST_TYPE_ALIYUN } -func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error { +func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error { params := task.GetParams() imageId, err := params.GetString("image_id") if err != nil { @@ -35,19 +36,28 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host * osType, _ := params.GetString("os_type") osDist, _ := params.GetString("os_distribution") + isForce := jsonutils.QueryBoolean(params, "is_force", false) userCred := task.GetUserCred() taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + + lockman.LockRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId)) + defer lockman.ReleaseRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId)) + + scimg := models.StoragecachedimageManager.Register(ctx, task.GetUserCred(), storageCache.Id, imageId) + iStorageCache, err := storageCache.GetIStorageCache() if err != nil { return nil, err } - extImgId, err := iStorageCache.UploadImage(ctx, userCred, imageId, osArch, osType, osDist, scimg.ExternalId, isForce) + extImgId, err := iStorageCache.UploadImage(userCred, imageId, osArch, osType, osDist, scimg.ExternalId, isForce) if err != nil { return nil, err } else { + scimg.SetExternalId(extImgId) + ret := jsonutils.NewDict() ret.Add(jsonutils.NewString(extImgId), "image_id") return ret, nil diff --git a/pkg/compute/hostdrivers/hostdrivers.go b/pkg/compute/hostdrivers/hostdrivers.go deleted file mode 100644 index a58d95a88b..0000000000 --- a/pkg/compute/hostdrivers/hostdrivers.go +++ /dev/null @@ -1 +0,0 @@ -package hostdrivers diff --git a/pkg/compute/hostdrivers/kvm.go b/pkg/compute/hostdrivers/kvm.go index 9bcd07cce3..5f776295b0 100644 --- a/pkg/compute/hostdrivers/kvm.go +++ b/pkg/compute/hostdrivers/kvm.go @@ -24,7 +24,7 @@ func (self *SKVMHostDriver) GetHostType() string { return models.HOST_TYPE_HYPERVISOR } -func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error { +func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error { params := task.GetParams() imageId, err := params.GetString("image_id") if err != nil { diff --git a/pkg/compute/models/hostdrivers.go b/pkg/compute/models/hostdrivers.go index 1d125b0cba..bc1f0674d2 100644 --- a/pkg/compute/models/hostdrivers.go +++ b/pkg/compute/models/hostdrivers.go @@ -11,7 +11,7 @@ import ( type IHostDriver interface { GetHostType() string - CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, scimg *SStoragecachedimage, task taskman.ITask) error + CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, task taskman.ITask) error RequestPrepareSaveDiskOnHost(ctx context.Context, host *SHost, disk *SDisk, imageId string, task taskman.ITask) error RequestSaveUploadImageOnHost(ctx context.Context, host *SHost, disk *SDisk, imageId string, task taskman.ITask, data jsonutils.JSONObject) error RequestAllocateDiskOnStorage(ctx context.Context, host *SHost, storage *SStorage, disk *SDisk, task taskman.ITask, content *jsonutils.JSONDict) error diff --git a/pkg/compute/tasks/storage_cache_image_task.go b/pkg/compute/tasks/storage_cache_image_task.go index fc24e13ea0..bd3d476fab 100644 --- a/pkg/compute/tasks/storage_cache_image_task.go +++ b/pkg/compute/tasks/storage_cache_image_task.go @@ -34,7 +34,7 @@ func (self *StorageCacheImageTask) OnInit(ctx context.Context, obj db.IStandalon self.SetStage("on_image_cache_complete", nil) host, _ := storageCache.GetHost() - err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, scimg, self) + err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, self) if err != nil { errData := taskman.Error2TaskData(err) self.OnImageCacheCompleteFailed(ctx, storageCache, errData) @@ -46,8 +46,8 @@ func (self *StorageCacheImageTask) OnImageCacheComplete(ctx context.Context, obj storageCache := obj.(*models.SStoragecache) imageId, _ := self.Params.GetString("image_id") scimg := models.StoragecachedimageManager.Register(ctx, self.UserCred, storageCache.Id, imageId) - extImgId, _ := data.GetString("image_id") - self.OnCacheSucc(ctx, storageCache, imageId, scimg, extImgId) + // extImgId, _ := data.GetString("image_id") + self.OnCacheSucc(ctx, storageCache, imageId, scimg) } func (self *StorageCacheImageTask) OnImageCacheCompleteFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { @@ -70,11 +70,8 @@ func (self *StorageCacheImageTask) OnCacheFailed(ctx context.Context, cache *mod self.SetStageFailed(ctx, err.Error()) } -func (self *StorageCacheImageTask) OnCacheSucc(ctx context.Context, cache *models.SStoragecache, imageId string, scimg *models.SStoragecachedimage, extImgId string) { +func (self *StorageCacheImageTask) OnCacheSucc(ctx context.Context, cache *models.SStoragecache, imageId string, scimg *models.SStoragecachedimage) { scimg.SetStatus(self.UserCred, models.CACHED_IMAGE_STATUS_READY, "cached") - if len(cache.ExternalId) > 0 && len(extImgId) > 0 && scimg.ExternalId != extImgId { - scimg.SetExternalId(extImgId) - } models.CachedimageManager.ImageAddRefCount(imageId) db.OpsLog.LogEvent(cache, db.ACT_CACHED_IMAGE, imageId, self.UserCred) self.SetStageComplete(ctx, nil) diff --git a/pkg/util/aliyun/storagecache.go b/pkg/util/aliyun/storagecache.go index fef99d1c21..d945c8d054 100644 --- a/pkg/util/aliyun/storagecache.go +++ b/pkg/util/aliyun/storagecache.go @@ -5,8 +5,6 @@ import ( "os" "strings" "time" - "context" - "github.com/aliyun/aliyun-oss-go-sdk/oss" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -17,7 +15,6 @@ 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/cloudcommon/db/lockman" ) type SStoragecache struct { @@ -88,9 +85,7 @@ func (self *SStoragecache) GetIImages() ([]cloudprovider.ICloudImage, error) { return self.iimages, nil } -func (self *SStoragecache) UploadImage(ctx context.Context, userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, extId string, isForce bool) (string, error) { - lockman.LockRawObject(ctx, "image", imageId) - defer lockman.ReleaseRawObject(ctx, "image", imageId) +func (self *SStoragecache) UploadImage(userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, extId string, isForce bool) (string, error) { if len(extId) > 0 { status, _ := self.region.GetImageStatus(extId) @@ -98,6 +93,7 @@ func (self *SStoragecache) UploadImage(ctx context.Context, userCred mcclient.To return extId, nil } } + return self.uploadImage(userCred, imageId, osArch, osType, osDist, isForce) }