mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
fix(region): download cached image from source host before migrating
This commit is contained in:
@@ -28,7 +28,6 @@ import (
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/pkg/utils"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/apis"
|
||||
api "yunion.io/x/onecloud/pkg/apis/compute"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/cmdline"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
@@ -39,7 +38,6 @@ import (
|
||||
"yunion.io/x/onecloud/pkg/compute/options"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/mcclient/auth"
|
||||
"yunion.io/x/onecloud/pkg/util/httputils"
|
||||
"yunion.io/x/onecloud/pkg/util/k8s/tokens"
|
||||
)
|
||||
@@ -164,6 +162,15 @@ func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *mod
|
||||
format, _ := params.GetString("format")
|
||||
isForce := jsonutils.QueryBoolean(params, "is_force", false)
|
||||
|
||||
srcHostId, _ := params.GetString("source_host_id")
|
||||
var srcHost *models.SHost
|
||||
if srcHostId != "" {
|
||||
srcHost = models.HostManager.FetchHostById(srcHostId)
|
||||
if srcHost == nil {
|
||||
return errors.Errorf("Source host %s not found", srcHostId)
|
||||
}
|
||||
}
|
||||
|
||||
type contentStruct struct {
|
||||
ImageId string
|
||||
Format string
|
||||
@@ -176,37 +183,31 @@ func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *mod
|
||||
content.ImageId = imageId
|
||||
content.Format = format
|
||||
|
||||
if options.Options.ImageCacheFromHost {
|
||||
obj, err := models.CachedimageManager.FetchById(imageId)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Fetch cached image by image_id %s", imageId)
|
||||
}
|
||||
cacheImage := obj.(*models.SCachedimage)
|
||||
srcHostCacheImage, err := cacheImage.ChooseSourceStoragecacheInRange(api.HOST_TYPE_HYPERVISOR, []string{host.Id}, []interface{}{host.GetZone()})
|
||||
obj, err := models.CachedimageManager.FetchById(imageId)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Fetch cached image by image_id %s", imageId)
|
||||
}
|
||||
cacheImage := obj.(*models.SCachedimage)
|
||||
rangeObjs := []interface{}{host.GetZone()}
|
||||
if srcHost != nil {
|
||||
rangeObjs = append(rangeObjs, srcHost)
|
||||
}
|
||||
srcHostCacheImage, err := cacheImage.ChooseSourceStoragecacheInRange(api.HOST_TYPE_HYPERVISOR, []string{host.Id}, rangeObjs)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Choose source storagecache")
|
||||
}
|
||||
if srcHostCacheImage != nil {
|
||||
err = srcHostCacheImage.AddDownloadRefcount()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if srcHostCacheImage != nil {
|
||||
err = srcHostCacheImage.AddDownloadRefcount()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
srcHost, err := srcHostCacheImage.GetHost()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
content.SrcUrl = fmt.Sprintf("%s/download/images/%s", srcHost.ManagerUri, imageId)
|
||||
}
|
||||
} else {
|
||||
// from glance service
|
||||
glanceURL, err := auth.GetServiceURL(apis.SERVICE_TYPE_IMAGE, "", host.GetZone().GetName(), "")
|
||||
|
||||
srcHost, err := srcHostCacheImage.GetHost()
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Get %s service url", apis.SERVICE_TYPE_IMAGE)
|
||||
}
|
||||
content.SrcUrl = fmt.Sprintf("%s/images/%s", glanceURL, imageId)
|
||||
if content.Format != "" {
|
||||
content.SrcUrl = fmt.Sprintf("%s?format=%s", content.SrcUrl, content.Format)
|
||||
return errors.Wrapf(err, "Get storage cached image %s host", srcHostCacheImage.GetId())
|
||||
}
|
||||
content.SrcUrl = fmt.Sprintf("%s/download/images/%s", srcHost.ManagerUri, imageId)
|
||||
|
||||
}
|
||||
|
||||
url := fmt.Sprintf("%s/disks/image_cache", host.ManagerUri)
|
||||
@@ -222,7 +223,7 @@ func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *mod
|
||||
|
||||
_, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false)
|
||||
if err != nil {
|
||||
return err
|
||||
return errors.Wrapf(err, "POST %s", url)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -170,6 +170,11 @@ func (self *SCachedimage) GetHypervisor() string {
|
||||
return osType
|
||||
}
|
||||
|
||||
func (self *SCachedimage) GetChecksum() string {
|
||||
checksum, _ := self.Info.GetString("checksum")
|
||||
return checksum
|
||||
}
|
||||
|
||||
func (self *SCachedimage) getStoragecacheQuery() *sqlchemy.SQuery {
|
||||
q := StoragecachedimageManager.Query().Equals("cachedimage_id", self.Id)
|
||||
return q
|
||||
@@ -450,6 +455,8 @@ func (self *SCachedimage) ChooseSourceStoragecacheInRange(hostType string, exclu
|
||||
|
||||
for _, rangeObj := range rangeObjs {
|
||||
switch v := rangeObj.(type) {
|
||||
case *SHost:
|
||||
q = q.Filter(sqlchemy.Equals(host.Field("id"), v.Id))
|
||||
case *SZone:
|
||||
q = q.Filter(sqlchemy.Equals(host.Field("zone_id"), v.Id))
|
||||
case *SCloudprovider:
|
||||
|
||||
@@ -385,6 +385,10 @@ func (self *SStoragecache) getCachedImageSize() int64 {
|
||||
}
|
||||
|
||||
func (self *SStoragecache) StartImageCacheTask(ctx context.Context, userCred mcclient.TokenCredential, imageId string, format string, isForce bool, parentTaskId string) error {
|
||||
return self.StartImageCacheTaskFromHost(ctx, userCred, imageId, format, isForce, "", parentTaskId)
|
||||
}
|
||||
|
||||
func (self *SStoragecache) StartImageCacheTaskFromHost(ctx context.Context, userCred mcclient.TokenCredential, imageId string, format string, isForce bool, srcHostId string, parentTaskId string) error {
|
||||
StoragecachedimageManager.Register(ctx, userCred, self.Id, imageId, "")
|
||||
data := jsonutils.NewDict()
|
||||
data.Add(jsonutils.NewString(imageId), "image_id")
|
||||
@@ -408,6 +412,10 @@ func (self *SStoragecache) StartImageCacheTask(ctx context.Context, userCred mcc
|
||||
if isForce {
|
||||
data.Add(jsonutils.JSONTrue, "is_force")
|
||||
}
|
||||
|
||||
if srcHostId != "" {
|
||||
data.Add(jsonutils.NewString(srcHostId), "source_host_id")
|
||||
}
|
||||
task, err := taskman.TaskManager.NewTask(ctx, "StorageCacheImageTask", self, userCred, data, parentTaskId, "", nil)
|
||||
if err != nil {
|
||||
log.Errorf("create StorageCacheImageTask fail %s", err)
|
||||
|
||||
@@ -50,7 +50,6 @@ type ComputeOptions struct {
|
||||
LoadbalancerPendingDeleteCheckInterval int `default:"3600" help:"Interval between checks of pending deleted loadbalancer objects, defaults to 1h"`
|
||||
|
||||
ImageCacheStoragePolicy string `default:"least_used" choices:"best_fit|least_used" help:"Policy to choose storage for image cache, best_fit or least_used"`
|
||||
ImageCacheFromHost bool `default:"false" help:"Download cached image from host"`
|
||||
MetricsRetentionDays int32 `default:"30" help:"Retention days for monitoring metrics in influxdb"`
|
||||
|
||||
DefaultBandwidth int `default:"1000" help:"Default bandwidth"`
|
||||
|
||||
@@ -138,8 +138,8 @@ func (self *GuestMigrateTask) SaveScheduleResult(ctx context.Context, obj ISched
|
||||
if len(disk.TemplateId) > 0 && isLocalStorage {
|
||||
targetStorageCache := targetHost.GetLocalStoragecache()
|
||||
if targetStorageCache != nil {
|
||||
err := targetStorageCache.StartImageCacheTask(
|
||||
ctx, self.UserCred, disk.TemplateId, disk.DiskFormat, false, self.GetTaskId())
|
||||
err := targetStorageCache.StartImageCacheTaskFromHost(
|
||||
ctx, self.UserCred, disk.TemplateId, disk.DiskFormat, false, guest.HostId, self.GetTaskId())
|
||||
if err != nil {
|
||||
self.TaskFailed(ctx, guest, jsonutils.NewString(err.Error()))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user