From 67fc4d1ba9b8fbfb9e59292ee64d09319056e929 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=B1=88=E8=BD=A9?= Date: Wed, 29 Aug 2018 14:24:05 +0800 Subject: [PATCH] =?UTF-8?q?=E6=94=AF=E6=8C=81server-save-image?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/cloudprovider/resources.go | 13 +++ pkg/compute/guestdrivers/baremetals.go | 5 + pkg/compute/guestdrivers/virtualization.go | 9 ++ pkg/compute/hostdrivers/aliyun.go | 49 ++++++++ pkg/compute/hostdrivers/kvm.go | 23 ++++ pkg/compute/models/disks.go | 90 +++++++++++++++ pkg/compute/models/guestdrivers.go | 2 + pkg/compute/models/guests.go | 37 ++++++ pkg/compute/models/hostdrivers.go | 2 + pkg/compute/tasks/disk_save_task.go | 107 +++++++++++++++++ pkg/compute/tasks/guest_save_image_task.go | 65 +++++++++++ pkg/mcclient/modules/mod_images.go | 47 ++++---- pkg/util/aliyun/disk.go | 107 +++++++++++++++++ pkg/util/aliyun/image.go | 31 +++++ pkg/util/aliyun/shell/snapshot.go | 45 ++++++++ pkg/util/aliyun/snapshot.go | 73 ++++++++++++ pkg/util/aliyun/storagecache.go | 128 +++++++++++++++++++++ 17 files changed, 813 insertions(+), 20 deletions(-) create mode 100644 pkg/compute/tasks/disk_save_task.go create mode 100644 pkg/compute/tasks/guest_save_image_task.go create mode 100644 pkg/util/aliyun/shell/snapshot.go create mode 100644 pkg/util/aliyun/snapshot.go diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index 3c333f8946..54877f315a 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -56,6 +56,7 @@ type ICloudZone interface { type ICloudImage interface { ICloudResource + Delete() error GetIStoragecache() ICloudStoragecache } @@ -66,6 +67,9 @@ type ICloudStoragecache interface { GetManagerId() string + CreateIImage(snapshotId, imageName, imageDesc string) (ICloudImage, error) + + DownloadImage(userCred mcclient.TokenCredential, imageId string, extId string) (jsonutils.JSONObject, error) UploadImage(userCred mcclient.TokenCredential, imageId string, extId string, isForce bool) (string, error) } @@ -200,9 +204,18 @@ type ICloudDisk interface { GetMountpoint() string Delete() error + CreateISnapshot(name string, desc string) (ICloudSnapshot, error) + GetISnapshot(idStr string) (ICloudSnapshot, error) + GetISnapshots() ([]ICloudSnapshot, error) + Resize(newSize int64) error } +type ICloudSnapshot interface { + ICloudResource + Delete() error +} + type ICloudVpc interface { ICloudResource diff --git a/pkg/compute/guestdrivers/baremetals.go b/pkg/compute/guestdrivers/baremetals.go index 2ec935e955..d12ebc0716 100644 --- a/pkg/compute/guestdrivers/baremetals.go +++ b/pkg/compute/guestdrivers/baremetals.go @@ -9,6 +9,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" ) @@ -172,3 +173,7 @@ func (self *SBaremetalGuestDriver) StartGuestDetachdiskTask(ctx context.Context, func (self *SBaremetalGuestDriver) StartSuspendTask(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error { return fmt.Errorf("Cannot suspend a baremetal serer") } + +func (self *SBaremetalGuestDriver) StartGuestSaveImage(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error { + return httperrors.NewUnsupportOperationError("Cannot save image for baremtal") +} diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index 7efdd6d475..a9ef20f3f1 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -193,3 +193,12 @@ func (self *SVirtualizedGuestDriver) StartSuspendTask(ctx context.Context, userC task.ScheduleRun(nil) return nil } + +func (self *SVirtualizedGuestDriver) StartGuestSaveImage(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, params *jsonutils.JSONDict, parentTaskId string) error { + if task, err := taskman.TaskManager.NewTask(ctx, "GuestSaveImageTask", guest, userCred, params, parentTaskId, "", nil); err != nil { + return err + } else { + task.ScheduleRun(nil) + } + return nil +} diff --git a/pkg/compute/hostdrivers/aliyun.go b/pkg/compute/hostdrivers/aliyun.go index 4c1d7bf029..e0977b298b 100644 --- a/pkg/compute/hostdrivers/aliyun.go +++ b/pkg/compute/hostdrivers/aliyun.go @@ -2,11 +2,13 @@ package hostdrivers import ( "context" + "fmt" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/httperrors" ) type SAliyunHostDriver struct { @@ -48,6 +50,53 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host * return nil } +func (self *SAliyunHostDriver) RequestPrepareSaveDiskOnHost(ctx context.Context, host *models.SHost, disk *models.SDisk, imageId string, task taskman.ITask) error { + task.ScheduleRun(nil) + return nil +} + +func (self *SAliyunHostDriver) RequestSaveUploadImageOnHost(ctx context.Context, host *models.SHost, disk *models.SDisk, imageId string, task taskman.ITask, data jsonutils.JSONObject) error { + if iDisk, err := disk.GetIDisk(); err != nil { + return err + } else if iStorage, err := disk.GetIStorage(); err != nil { + return err + } else if iStoragecache := iStorage.GetIStoragecache(); iStoragecache == nil { + return httperrors.NewResourceNotFoundError("fail to find iStoragecache for storage: %s", iStorage.GetName()) + } else { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + if snapshot, err := iDisk.CreateISnapshot(fmt.Sprintf("Snapshot-%s", imageId), "PrepareSaveImage"); err != nil { + return nil, err + } else { + scimg := models.StoragecachedimageManager.Register(ctx, task.GetUserCred(), iStoragecache.GetId(), imageId) + if scimg.Status != models.CACHED_IMAGE_STATUS_READY { + scimg.SetStatus(task.GetUserCred(), models.CACHED_IMAGE_STATUS_CACHING, "request_prepare_save_disk_on_host") + } + if iImage, err := iStoragecache.CreateIImage(snapshot.GetId(), fmt.Sprintf("Image-%s", imageId), ""); err != nil { + log.Errorf("fail to create iImage: %v", err) + scimg.SetStatus(task.GetUserCred(), models.CACHED_IMAGE_STATUS_CACHE_FAILED, err.Error()) + return nil, err + } else { + scimg.SetExternalId(iImage.GetId()) + if result, err := iStoragecache.DownloadImage(task.GetUserCred(), imageId, iImage.GetId()); err != nil { + scimg.SetStatus(task.GetUserCred(), models.CACHED_IMAGE_STATUS_CACHE_FAILED, err.Error()) + return nil, err + } else { + if err := iImage.Delete(); err != nil { + log.Errorf("Delete iImage %s failed: %v", iImage.GetId(), err) + } + if err := snapshot.Delete(); err != nil { + log.Errorf("Delete snapshot %s failed: %v", snapshot.GetId(), err) + } + scimg.SetStatus(task.GetUserCred(), models.CACHED_IMAGE_STATUS_READY, "") + return result, nil + } + } + } + }) + } + return nil +} + func (self *SAliyunHostDriver) RequestAllocateDiskOnStorage(ctx context.Context, host *models.SHost, storage *models.SStorage, disk *models.SDisk, task taskman.ITask, content *jsonutils.JSONDict) error { if iCloudStorage, err := storage.GetIStorage(); err != nil { return err diff --git a/pkg/compute/hostdrivers/kvm.go b/pkg/compute/hostdrivers/kvm.go index ed08542a36..9bcd07cce3 100644 --- a/pkg/compute/hostdrivers/kvm.go +++ b/pkg/compute/hostdrivers/kvm.go @@ -124,3 +124,26 @@ func (self *SKVMHostDriver) RequestResizeDiskOnHostOnline(host *models.SHost, st } return nil } + +func (self *SKVMHostDriver) RequestPrepareSaveDiskOnHost(ctx context.Context, host *models.SHost, disk *models.SDisk, imageId string, task taskman.ITask) error { + body := jsonutils.NewDict() + body.Add(jsonutils.Marshal(map[string]string{"image_id": imageId}), "disk") + url := fmt.Sprintf("/disks/%s/save-prepare/%s", disk.StorageId, disk.Id) + header := http.Header{"X-Task-Id": []string{task.GetTaskId()}, "X-Region-Version": []string{"v2"}} + _, err := host.Request(task.GetUserCred(), "POST", url, header, body) + return err +} + +func (self *SKVMHostDriver) RequestSaveUploadImageOnHost(ctx context.Context, host *models.SHost, disk *models.SDisk, imageId string, task taskman.ITask, data jsonutils.JSONObject) error { + body := jsonutils.NewDict() + backup, _ := data.GetString("backup") + content := map[string]string{"image_path": backup, "image_id": imageId, "storagecached_id": disk.GetStorage().StoragecacheId} + if data.Contains("format") { + content["format"], _ = data.GetString("format") + } + body.Add(jsonutils.Marshal(content), "disk") + url := fmt.Sprintf("/disks/%s/upload", disk.StorageId) + header := http.Header{"X-Task-Id": []string{task.GetTaskId()}, "X-Region-Version": []string{"v2"}} + _, err := host.Request(task.GetUserCred(), "POST", url, header, body) + return err +} diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index b05fa85b5e..badf077755 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -27,6 +27,8 @@ 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/mcclient/modules" ) const ( @@ -358,6 +360,94 @@ func (self *SDisk) PerformResize(ctx context.Context, userCred mcclient.TokenCre } } +func (self *SDisk) GetIStorage() (cloudprovider.ICloudStorage, error) { + if storage := self.GetStorage(); storage == nil { + return nil, httperrors.NewResourceNotFoundError("fail to find storage for disk %s", self.GetName()) + } else if provider, err := storage.GetDriver(); err != nil { + return nil, err + } else { + return provider.GetIStorageById(storage.GetExternalId()) + } +} + +func (self *SDisk) GetIDisk() (cloudprovider.ICloudDisk, error) { + if iStorage, err := self.GetIStorage(); err != nil { + log.Errorf("fail to find iStorage: %v", err) + return nil, err + } else { + return iStorage.GetIDisk(self.GetExternalId()) + } +} + +func (self *SDisk) GetZone() *SZone { + if storage := self.GetStorage(); storage != nil { + return storage.getZone() + } + return nil +} + +func (self *SDisk) PrepareSaveImage(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (string, error) { + if zone := self.GetZone(); zone == nil { + return "", httperrors.NewResourceNotFoundError("No zone for this disk") + } + data.Add(jsonutils.NewString(self.DiskFormat), "disk_format") + name, _ := data.GetString("name") + s := auth.GetAdminSession(options.Options.Region, "") + if imageList, err := modules.Images.List(s, jsonutils.Marshal(map[string]string{"name": name, "admin": "true"})); err != nil { + return "", err + } else if imageList.Total > 0 { + return "", httperrors.NewConflictError("Duplicate image name %s", name) + } + quota := SQuota{Image: 1} + if _, err := QuotaManager.CheckQuota(ctx, userCred, userCred.GetProjectId(), "a); err != nil { + return "", err + } + data.Add(jsonutils.NewInt(int64(self.DiskSize)), "virtual_size") + if result, err := modules.Images.Create(s, data); err != nil { + return "", err + } else if imageId, err := result.GetString("id"); err != nil { + return "", err + } else { + return imageId, nil + } +} + +func (self *SDisk) AllowPerformSave(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return self.IsOwner(userCred) +} + +func (self *SDisk) PerformSave(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if self.Status != DISK_READY { + return nil, httperrors.NewResourceNotReadyError("Save disk when disk is READY") + + } + if self.GetRuningGuestCount() > 0 { + return nil, httperrors.NewResourceNotReadyError("Save disk when not being USED") + } + + if name, err := data.GetString("name"); err != nil || len(name) == 0 { + return nil, httperrors.NewInputParameterError("Image name is required") + } + kwargs := data.(*jsonutils.JSONDict) + if imageId, err := self.PrepareSaveImage(ctx, userCred, kwargs); err != nil { + return nil, err + } else { + kwargs.Add(jsonutils.NewString(imageId), "image_id") + return nil, self.StartDiskSaveTask(ctx, userCred, kwargs, "") + } +} + +func (self *SDisk) StartDiskSaveTask(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, parentTaskId string) error { + self.SetStatus(userCred, DISK_START_SAVE, "") + if task, err := taskman.TaskManager.NewTask(ctx, "DiskSaveTask", self, userCred, data, parentTaskId, "", nil); err != nil { + log.Errorf("Start DiskSaveTask failed:%v", err) + return err + } else { + task.ScheduleRun(nil) + } + return nil +} + func (self *SDisk) ValidateDeleteCondition(ctx context.Context) error { if self.GetGuestDiskCount() > 0 { return httperrors.NewNotEmptyError("Virtual disk used by virtual servers") diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index 48406df24f..4d841ddae7 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -61,6 +61,8 @@ type IGuestDriver interface { StartDeleteGuestTask(ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, params *jsonutils.JSONDict, parentTaskId string) error + StartGuestSaveImage(ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, params *jsonutils.JSONDict, parentTaskId string) error + RequestStopGuestForDelete(ctx context.Context, guest *SGuest, task taskman.ITask) error RequestDetachDisksFromGuestForDelete(ctx context.Context, guest *SGuest, task taskman.ITask) error diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index afd50582dc..1423bfa7ba 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1560,6 +1560,43 @@ func (self *SGuest) attach2Disk(disk *SDisk, userCred mcclient.TokenCredential, return err } +func (self *SGuest) AllowPerformSaveImage(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + log.Infof("permission: %s", self.IsOwner(userCred)) + return self.IsOwner(userCred) +} + +func (self *SGuest) PerformSaveImage(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if !utils.IsInStringArray(self.Status, []string{VM_READY}) { + return nil, httperrors.NewInputParameterError("Cannot save image in status %s", self.Status) + } else if !data.Contains("name") { + return nil, httperrors.NewInputParameterError("Image name is required") + } else if disks := self.CategorizeDisks(); disks.Root == nil { + return nil, httperrors.NewInputParameterError("No root image") + } else { + kwargs := data.(*jsonutils.JSONDict) + restart := self.Status == VM_RUNNING + properties := jsonutils.NewDict() + if notes, err := data.GetString("notes"); err != nil && len(notes) > 0 { + properties.Add(jsonutils.NewString(notes), "notes") + } + properties.Add(jsonutils.NewString(self.OsType), "os_type") + kwargs.Add(properties, "properties") + kwargs.Add(jsonutils.NewBool(restart), "restart") + lockman.LockObject(ctx, disks.Root) + defer lockman.ReleaseObject(ctx, disks.Root) + if imageId, err := disks.Root.PrepareSaveImage(ctx, userCred, kwargs); err != nil { + return nil, err + } else { + kwargs.Add(jsonutils.NewString(imageId), "image_id") + } + return nil, self.StartGuestSaveImage(ctx, userCred, kwargs, "") + } +} + +func (self *SGuest) StartGuestSaveImage(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict, parentTaskId string) error { + return self.GetDriver().StartGuestSaveImage(ctx, userCred, self, data, parentTaskId) +} + func (self *SGuest) AllowPerformSync(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { return self.IsOwner(userCred) } diff --git a/pkg/compute/models/hostdrivers.go b/pkg/compute/models/hostdrivers.go index 6749e5ac68..1d125b0cba 100644 --- a/pkg/compute/models/hostdrivers.go +++ b/pkg/compute/models/hostdrivers.go @@ -12,6 +12,8 @@ import ( type IHostDriver interface { GetHostType() string CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, scimg *SStoragecachedimage, 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 RequestDeallocateDiskOnHost(host *SHost, storage *SStorage, disk *SDisk, task taskman.ITask) error RequestResizeDiskOnHostOnline(host *SHost, storage *SStorage, disk *SDisk, size int64, task taskman.ITask) error diff --git a/pkg/compute/tasks/disk_save_task.go b/pkg/compute/tasks/disk_save_task.go new file mode 100644 index 0000000000..8990ae8ff2 --- /dev/null +++ b/pkg/compute/tasks/disk_save_task.go @@ -0,0 +1,107 @@ +package tasks + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/compute/options" + "yunion.io/x/onecloud/pkg/mcclient/auth" + mc "yunion.io/x/onecloud/pkg/mcclient/modules" +) + +type DiskSaveTask struct { + SDiskBaseTask +} + +func init() { + taskman.RegisterTask(DiskSaveTask{}) +} + +func (self *DiskSaveTask) GetMasterHost(disk *models.SDisk) *models.SHost { + if guests := disk.GetGuests(); len(guests) == 1 { + if host := guests[0].GetHost(); host == nil { + if storage := disk.GetStorage(); storage != nil { + return storage.GetMasterHost() + } + } else { + return host + } + } + return nil +} + +func (self *DiskSaveTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + disk := obj.(*models.SDisk) + if host := self.GetMasterHost(disk); host == nil { + resion := "Cannot find host for disk" + disk.SetDiskReady(ctx, self.GetUserCred(), resion) + self.TaskFailed(ctx, resion) + db.OpsLog.LogEvent(disk, db.ACT_SAVE_FAIL, resion, self.GetUserCred()) + } else { + disk.SetStatus(self.GetUserCred(), models.DISK_START_SAVE, "") + for _, guest := range disk.GetGuests() { + guest.SetStatus(self.GetUserCred(), models.VM_SAVE_DISK, "") + } + self.StartBackupDisk(ctx, disk, host) + } +} + +func (self *DiskSaveTask) StartBackupDisk(ctx context.Context, disk *models.SDisk, host *models.SHost) { + self.SetStage("on_disk_backup_complete", nil) + disk.SetStatus(self.GetUserCred(), models.DISK_SAVING, "") + imageId, _ := self.GetParams().GetString("image_id") + if err := host.GetHostDriver().RequestPrepareSaveDiskOnHost(ctx, host, disk, imageId, self); err != nil { + log.Errorf("Backup failed: %v", err) + disk.SetDiskReady(ctx, self.GetUserCred(), err.Error()) + self.TaskFailed(ctx, err.Error()) + db.OpsLog.LogEvent(disk, db.ACT_SAVE_FAIL, err.Error(), self.GetUserCred()) + } +} + +func (self *DiskSaveTask) OnDiskBackupCompleteFailed(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { + disk.SetDiskReady(ctx, self.GetUserCred(), data.String()) + db.OpsLog.LogEvent(disk, db.ACT_SAVE_FAIL, data.String(), self.GetUserCred()) +} + +func (self *DiskSaveTask) OnDiskBackupComplete(ctx context.Context, disk *models.SDisk, data *jsonutils.JSONDict) { + disk.SetDiskReady(ctx, self.GetUserCred(), "") + db.OpsLog.LogEvent(disk, db.ACT_SAVE, disk.GetShortDesc(), self.GetUserCred()) + self.SetStageComplete(ctx, nil) + imageId, _ := self.GetParams().GetString("image_id") + if host := self.GetMasterHost(disk); host == nil { + log.Errorf("Saved disk Host mast not be nil") + self.TaskFailed(ctx, "Saved disk Host mast not be nil") + } else { + if self.Params.Contains("format") { + format, _ := self.Params.Get("format") + data.Add(format, "format") + } + if err := self.UploadDisk(ctx, host, disk, imageId, data); err != nil { + log.Errorf("UploadDisk failed: %v", err) + self.TaskFailed(ctx, err.Error()) + } + self.RefreshImageCache(ctx, imageId) + } +} + +func (self *DiskSaveTask) RefreshImageCache(ctx context.Context, imageId string) { + models.CachedimageManager.GetImageById(ctx, self.GetUserCred(), imageId, true) +} + +func (self *DiskSaveTask) UploadDisk(ctx context.Context, host *models.SHost, disk *models.SDisk, imageId string, data *jsonutils.JSONDict) error { + return host.GetHostDriver().RequestSaveUploadImageOnHost(ctx, host, disk, imageId, self, jsonutils.Marshal(data)) +} + +func (self *DiskSaveTask) TaskFailed(ctx context.Context, resion string) { + self.SetStageFailed(ctx, resion) + if imageId, err := self.GetParams().GetString("image_id"); err != nil && len(imageId) > 0 { + log.Errorf("save disk task failed, set image %s killed", imageId) + s := auth.GetAdminSession(options.Options.Region, "") + mc.Images.Update(s, imageId, jsonutils.Marshal(map[string]string{"status": "killed"})) + } +} diff --git a/pkg/compute/tasks/guest_save_image_task.go b/pkg/compute/tasks/guest_save_image_task.go new file mode 100644 index 0000000000..c6c13bfc2b --- /dev/null +++ b/pkg/compute/tasks/guest_save_image_task.go @@ -0,0 +1,65 @@ +package tasks + +import ( + "context" + "fmt" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" +) + +type GuestSaveImageTask struct { + SGuestBaseTask +} + +func init() { + taskman.RegisterTask(GuestSaveImageTask{}) +} + +func (self *GuestSaveImageTask) OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + log.Infof("Saving server image: %s", guest.Name) + if restart, _ := self.GetParams().Bool("restart"); restart { + self.SetStage("on_stop_server_complete", nil) + guest.StartGuestStopTask(ctx, self.GetUserCred(), false, self.GetTaskId()) + } else { + self.OnStopServerComplete(ctx, guest, nil) + } +} + +func (self *GuestSaveImageTask) OnStopServerComplete(ctx context.Context, guest *models.SGuest, body jsonutils.JSONObject) { + if guest.Status != models.VM_READY { + resion := fmt.Sprintf("Server %s not in ready status", guest.Name) + log.Errorf(resion) + self.SetStageFailed(ctx, resion) + } else { + self.SetStage("on_save_root_image_complete", nil) + guest.SetStatus(self.GetUserCred(), models.VM_START_SAVE_DISK, "") + disks := guest.CategorizeDisks() + if err := disks.Root.StartDiskSaveTask(ctx, self.GetUserCred(), self.GetParams(), self.GetTaskId()); err != nil { + self.SetStageFailed(ctx, err.Error()) + } + } +} + +func (self *GuestSaveImageTask) OnSaveRootImageComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + if restart, _ := self.GetParams().Bool("restart"); restart { + self.SetStage("on_start_server_complete", nil) + guest.StartGueststartTask(ctx, self.GetUserCred(), nil, self.GetTaskId()) + } else { + self.SetStageComplete(ctx, nil) + } +} + +func (self *GuestSaveImageTask) OnSaveRootImageCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + log.Errorf("Guest save root image failed: %s", data.PrettyString()) + guest.SetStatus(self.GetUserCred(), models.VM_SAVE_DISK_FAILED, data.PrettyString()) + self.SetStageFailed(ctx, data.PrettyString()) +} + +func (self *GuestSaveImageTask) OnStartServerCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} diff --git a/pkg/mcclient/modules/mod_images.go b/pkg/mcclient/modules/mod_images.go index c7423c2436..ae06aff2ab 100644 --- a/pkg/mcclient/modules/mod_images.go +++ b/pkg/mcclient/modules/mod_images.go @@ -428,24 +428,32 @@ func (this *ImageManager) _create(s *mcclient.ClientSession, params jsonutils.JS if !exists { return nil, fmt.Errorf("Unsupported image format %s", format) } - osType, err := params.GetString("properties", "os_type") - if err != nil { - return nil, fmt.Errorf("Can't get os_type from params: %s", params.String()) - } - exists, _ = utils.InStringArray(osType, []string{"Windows", "Linux", "Freebsd", "Android", "macOS", "VMWare"}) - if !exists { - return nil, fmt.Errorf("OS type must be specified") - } - name, _ := params.GetString("name") - if len(name) == 0 { - return nil, fmt.Errorf("Missing name") - } - dupName, e := this.IsNameDuplicate(s, name) - if dupName { - return nil, fmt.Errorf("Duplicate name %s", name) - } - if e != nil { - return nil, fmt.Errorf("Check name duplicate error %s", e) + imageId, _ := params.GetString("image_id") + path := fmt.Sprintf("/%s", this.URLPath()) + method := "POST" + if len(imageId) == 0 { + osType, err := params.GetString("properties", "os_type") + if err != nil { + return nil, fmt.Errorf("Can't get os_type from params: %s", params.String()) + } + exists, _ = utils.InStringArray(osType, []string{"Windows", "Linux", "Freebsd", "Android", "macOS", "VMWare"}) + if !exists { + return nil, fmt.Errorf("OS type must be specified") + } + name, _ := params.GetString("name") + if len(name) == 0 { + return nil, fmt.Errorf("Missing name") + } + dupName, e := this.IsNameDuplicate(s, name) + if dupName { + return nil, fmt.Errorf("Duplicate name %s", name) + } + if e != nil { + return nil, fmt.Errorf("Check name duplicate error %s", e) + } + } else { + path = fmt.Sprintf("/%s/%s", this.URLPath(), imageId) + method = "PUT" } headers, e := setImageMeta(params) if e != nil { @@ -467,8 +475,7 @@ func (this *ImageManager) _create(s *mcclient.ClientSession, params jsonutils.JS headers.Add("Content-Length", fmt.Sprintf("%d", size)) } } - path := fmt.Sprintf("/%s", this.URLPath()) - resp, err := this.rawRequest(s, "POST", path, headers, body) + resp, err := this.rawRequest(s, method, path, headers, body) _, json, err := s.ParseJSONResponse(resp, err) if err != nil { return nil, err diff --git a/pkg/util/aliyun/disk.go b/pkg/util/aliyun/disk.go index 7349b8c7c0..7975916e80 100644 --- a/pkg/util/aliyun/disk.go +++ b/pkg/util/aliyun/disk.go @@ -260,3 +260,110 @@ func (self *SRegion) resizeDisk(diskId string, size int64) error { return nil } + +func (self *SDisk) CreateISnapshot(name, desc string) (cloudprovider.ICloudSnapshot, error) { + if snapshotId, err := self.storage.zone.region.CreateSnapshot(self.DiskId, name, desc); err != nil { + log.Errorf("createSnapshot fail %s", err) + return nil, err + } else if snapshot, err := self.getSnapshot(snapshotId); err != nil { + return nil, err + } else { + snapshot.disk = self + if err := cloudprovider.WaitStatus(snapshot, string(SnapshotStatusAccoplished), 15*time.Second, 3600*time.Second); err != nil { + return nil, err + } + return snapshot, nil + } +} + +func (self *SRegion) CreateSnapshot(diskId, name, desc string) (string, error) { + params := make(map[string]string) + params["RegionId"] = self.RegionId + params["DiskId"] = diskId + params["SnapshotName"] = name + params["Description"] = desc + + if body, err := self.ecsRequest("CreateSnapshot", params); err != nil { + log.Errorf("CreateSnapshot fail %s", err) + return "", err + } else { + return body.GetString("SnapshotId") + } +} + +func (self *SDisk) GetISnapshot(snapshotId string) (cloudprovider.ICloudSnapshot, error) { + if snapshot, err := self.getSnapshot(snapshotId); err != nil { + return nil, err + } else { + snapshot.disk = self + return snapshot, nil + } +} + +func (self *SDisk) getSnapshot(snapshotId string) (*SSnapshot, error) { + if snapshots, total, err := self.storage.zone.region.GetSnapshots("", "", "", []string{snapshotId}, 0, 1); err != nil { + return nil, err + } else if total != 1 { + return nil, cloudprovider.ErrNotFound + } else { + return &snapshots[0], nil + } +} + +func (self *SDisk) GetISnapshots() ([]cloudprovider.ICloudSnapshot, error) { + snapshots := make([]SSnapshot, 0) + for { + if parts, total, err := self.storage.zone.region.GetSnapshots("", self.DiskId, "", []string{}, 0, 20); err != nil { + log.Errorf("GetDisks fail %s", err) + return nil, err + } else { + snapshots = append(snapshots, parts...) + if len(snapshots) >= total { + break + } + } + } + isnapshots := make([]cloudprovider.ICloudSnapshot, len(snapshots)) + for i := 0; i < len(snapshots); i++ { + snapshots[i].disk = self + isnapshots[i] = &snapshots[i] + } + return isnapshots, nil +} + +func (self *SRegion) GetSnapshots(instanceId string, diskId string, snapshotName string, snapshotIds []string, offset int, limit int) ([]SSnapshot, int, error) { + if limit > 50 || limit <= 0 { + limit = 50 + } + params := make(map[string]string) + params["RegionId"] = self.RegionId + params["PageSize"] = fmt.Sprintf("%d", limit) + params["PageNumber"] = fmt.Sprintf("%d", (offset/limit)+1) + + if len(instanceId) > 0 { + params["InstanceId"] = instanceId + } + if len(diskId) > 0 { + params["diskId"] = diskId + } + if len(snapshotName) > 0 { + params["SnapshotName"] = snapshotName + } + if snapshotIds != nil && len(snapshotIds) > 0 { + params["SnapshotIds"] = jsonutils.Marshal(snapshotIds).String() + } + + if body, err := self.ecsRequest("DescribeSnapshots", params); err != nil { + log.Errorf("GetSnapshots fail %s", err) + return nil, 0, err + } else { + snapshots := make([]SSnapshot, 0) + if err := body.Unmarshal(&snapshots, "Snapshots", "Snapshot"); err != nil { + log.Errorf("Unmarshal snapshot details fail %s", err) + return nil, 0, err + } + total, _ := body.Int("TotalCount") + return snapshots, int(total), nil + } + +} diff --git a/pkg/util/aliyun/image.go b/pkg/util/aliyun/image.go index 849e9eb570..dfdddbcbfd 100644 --- a/pkg/util/aliyun/image.go +++ b/pkg/util/aliyun/image.go @@ -5,6 +5,7 @@ import ( "strings" "time" + "github.com/aliyun/aliyun-oss-go-sdk/oss" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -71,6 +72,10 @@ func (self *SImage) IsEmulated() bool { return false } +func (self *SImage) Delete() error { + return self.storageCache.region.DeleteImage(self.ImageId) +} + func (self *SImage) GetGlobalId() string { return fmt.Sprintf("%s-%s") } @@ -102,6 +107,32 @@ func (self *SImage) Refresh() error { return jsonutils.Update(self, new) } +type ImageExportTask struct { + ImageId string + RegionId string + // RequestId string + TaskId string +} + +func (self *SRegion) ExportImage(imageId string, bucket *oss.Bucket) (*ImageExportTask, error) { + params := make(map[string]string) + params["RegionId"] = self.RegionId + params["ImageId"] = imageId + params["OssBucket"] = bucket.BucketName + params["OssPrefix"] = fmt.Sprintf("%sexport", strings.Replace(imageId, "-", "", -1)) + + if body, err := self.ecsRequest("ExportImage", params); err != nil { + return nil, err + } else { + result := ImageExportTask{} + if err := body.Unmarshal(&result); err != nil { + log.Errorf("unmarshal result error %s", err) + return nil, err + } + return &result, nil + } +} + // {"ImageId":"m-j6c1qlpa7oebbg1n2k60","RegionId":"cn-hongkong","RequestId":"F8B2F6A1-F6AA-4C92-A54C-C4A309CF811F","TaskId":"t-j6c1qlpa7oebbg1rcl9t"} type ImageImportTask struct { diff --git a/pkg/util/aliyun/shell/snapshot.go b/pkg/util/aliyun/shell/snapshot.go new file mode 100644 index 0000000000..5bfe219871 --- /dev/null +++ b/pkg/util/aliyun/shell/snapshot.go @@ -0,0 +1,45 @@ +package shell + +import ( + "yunion.io/x/onecloud/pkg/util/aliyun" + "yunion.io/x/onecloud/pkg/util/shellutils" +) + +func init() { + type SnapshotListOptions struct { + DiskId string `help:"Disk ID"` + InstanceId string `help:"Instance ID"` + SnapshotIds []string `helo:"Snapshot ids"` + Name string `help:"Snapshot Name"` + Limit int `help:"page size"` + Offset int `help:"page offset"` + } + shellutils.R(&SnapshotListOptions{}, "snapshot-list", "List snapshot", func(cli *aliyun.SRegion, args *SnapshotListOptions) error { + if snapshots, total, err := cli.GetSnapshots(args.InstanceId, args.DiskId, args.Name, args.SnapshotIds, args.Offset, args.Limit); err != nil { + return err + } else { + printList(snapshots, total, args.Offset, args.Limit, []string{}) + return nil + } + }) + + type SnapshotDeleteOptions struct { + ID string `help:"Snapshot ID"` + } + + shellutils.R(&SnapshotDeleteOptions{}, "snapshot-delete", "Delete snapshot", func(cli *aliyun.SRegion, args *SnapshotDeleteOptions) error { + return cli.DeleteSnapshot(args.ID) + }) + + type SnapshotCreateOptions struct { + DiskId string `help:"Disk ID"` + Name string `help:"Snapeshot Name"` + Desc string `help:"Snapshot Desc"` + } + + shellutils.R(&SnapshotCreateOptions{}, "snapshot-create", "Create snapshot", func(cli *aliyun.SRegion, args *SnapshotCreateOptions) error { + _, err := cli.CreateSnapshot(args.DiskId, args.Name, args.Desc) + return err + }) + +} diff --git a/pkg/util/aliyun/snapshot.go b/pkg/util/aliyun/snapshot.go new file mode 100644 index 0000000000..bfab3ae94d --- /dev/null +++ b/pkg/util/aliyun/snapshot.go @@ -0,0 +1,73 @@ +package aliyun + +import ( + "fmt" + + "yunion.io/x/jsonutils" +) + +type SnapshotStatusType string + +const ( + SnapshotStatusAccoplished SnapshotStatusType = "accomplished" + SnapshotStatusProgress SnapshotStatusType = "progressing" +) + +type SSnapshot struct { + disk *SDisk + Progress string + SnapshotId string + SnapshotName string + SourceDiskId string + SourceDiskSize int32 + SourceDiskType string + Status SnapshotStatusType + Usage string +} + +func (self *SSnapshot) GetId() string { + return self.SnapshotId +} + +func (self *SSnapshot) GetName() string { + return self.SnapshotName +} + +func (self *SSnapshot) GetStatus() string { + return string(self.Status) +} + +func (self *SSnapshot) Refresh() error { + if snapshot, err := self.disk.getSnapshot(self.SnapshotId); err != nil { + return err + } else if err := jsonutils.Update(self, snapshot); err != nil { + return err + } + return nil +} + +func (self *SSnapshot) GetGlobalId() string { + return fmt.Sprintf("%s", self.SnapshotId) +} + +func (self *SSnapshot) IsEmulated() bool { + return false +} + +func (self *SRegion) DeleteSnapshot(snapshotId string) error { + params := make(map[string]string) + params["SnapshotId"] = snapshotId + _, err := self.ecsRequest("DeleteSnapshot", params) + return err +} + +func (self *SSnapshot) Delete() error { + if self.disk == nil { + return fmt.Errorf("not init disk for snapshot %s", self.SnapshotId) + } + return self.disk.storage.zone.region.DeleteSnapshot(self.SnapshotId) +} + +func (self *SSnapshot) GetMetadata() *jsonutils.JSONDict { + return nil +} diff --git a/pkg/util/aliyun/storagecache.go b/pkg/util/aliyun/storagecache.go index fab2072253..cab1cbd7ab 100644 --- a/pkg/util/aliyun/storagecache.go +++ b/pkg/util/aliyun/storagecache.go @@ -2,13 +2,17 @@ package aliyun import ( "fmt" + "os" "strings" "time" + "github.com/aliyun/aliyun-oss-go-sdk/oss" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" + compute "yunion.io/x/onecloud/pkg/compute/models" "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/mcclient/modules" @@ -173,3 +177,127 @@ func (self *SStoragecache) uploadImage(userCred mcclient.TokenCredential, imageI return task.ImageId, nil } + +func (self *SStoragecache) CreateIImage(snapshoutId, imageName, imageDesc string) (cloudprovider.ICloudImage, error) { + if imageId, err := self.region.createIImage(snapshoutId, imageName, imageDesc); err != nil { + return nil, err + } else if image, err := self.region.GetImage(imageId); err != nil { + return nil, err + } else { + image.storageCache = self + iimage := make([]cloudprovider.ICloudImage, 1) + iimage[0] = image + if err := cloudprovider.WaitStatus(iimage[0], compute.IMAGE_STATUS_ACTIVE, 15*time.Second, 3600*time.Second); err != nil { + return nil, err + } + return iimage[0], nil + } +} + +func (self *SRegion) CheckBucket(bucketName string) (*oss.Bucket, error) { + return self.checkBucket(bucketName) +} + +func (self *SRegion) checkBucket(bucketName string) (*oss.Bucket, error) { + oss, err := self.GetOssClient() + if err != nil { + log.Errorf("GetOssClient err %s", err) + return nil, err + } + if exist, err := oss.IsBucketExist(bucketName); err != nil { + log.Errorf("IsBucketExist err %s", err) + return nil, err + } else if !exist { + log.Debugf("Bucket %s not exists, to create ...", bucketName) + if err := oss.CreateBucket(bucketName); err != nil { + log.Errorf("Create bucket error %s", err) + return nil, err + } + } + log.Debugf("Bucket %s exists", bucketName) + if bucket, err := oss.Bucket(bucketName); err != nil { + log.Errorf("Bucket error %s %s", bucketName, err) + return nil, err + } else { + return bucket, nil + } +} + +func (self *SRegion) createIImage(snapshoutId, imageName, imageDesc string) (string, error) { + params := make(map[string]string) + params["RegionId"] = self.RegionId + params["OssBucket"] = strings.ToLower(fmt.Sprintf("imgcache-%s", self.GetId())) + params["SnapshotId"] = snapshoutId + params["ImageName"] = imageName + params["Description"] = imageDesc + + if _, err := self.checkBucket(params["OssBucket"]); err != nil { + return "", err + } + + if body, err := self.ecsRequest("CreateImage", params); err != nil { + log.Errorf("CreateImage fail %s", err) + return "", err + } else { + log.Infof("%s", body) + return body.GetString("ImageId") + } +} + +func (self *SStoragecache) DownloadImage(userCred mcclient.TokenCredential, imageId string, extId string) (jsonutils.JSONObject, error) { + return self.downloadImage(userCred, imageId, extId) +} + +// 定义进度条监听器。 +type OssProgressListener struct { +} + +// 定义进度变更事件处理函数。 +func (listener *OssProgressListener) ProgressChanged(event *oss.ProgressEvent) { + switch event.EventType { + case oss.TransferStartedEvent: + log.Debugf("Transfer Started, ConsumedBytes: %d, TotalBytes %d.\n", + event.ConsumedBytes, event.TotalBytes) + case oss.TransferDataEvent: + log.Debugf("\rTransfer Data, ConsumedBytes: %d, TotalBytes %d, %d%%.", + event.ConsumedBytes, event.TotalBytes, event.ConsumedBytes*100/event.TotalBytes) + case oss.TransferCompletedEvent: + log.Debugf("\nTransfer Completed, ConsumedBytes: %d, TotalBytes %d.\n", + event.ConsumedBytes, event.TotalBytes) + case oss.TransferFailedEvent: + log.Debugf("\nTransfer Failed, ConsumedBytes: %d, TotalBytes %d.\n", + event.ConsumedBytes, event.TotalBytes) + default: + } +} + +func (self *SStoragecache) downloadImage(userCred mcclient.TokenCredential, imageId string, extId string) (jsonutils.JSONObject, error) { + tmpImageFile := fmt.Sprintf("/tmp/%s", extId) + bucketName := strings.ToLower(fmt.Sprintf("imgcache-%s", self.region.GetId())) + if bucket, err := self.region.checkBucket(bucketName); err != nil { + return nil, err + } else if _, err := self.region.GetImage(extId); err != nil { + return nil, err + } else if task, err := self.region.ExportImage(extId, bucket); err != nil { + return nil, err + } else if err := self.region.waitTaskStatus(ExportImageTask, task.TaskId, "Finished", 15*time.Second, 3600*time.Second); err != nil { + return nil, err + } else if imageList, err := bucket.ListObjects(oss.Prefix(fmt.Sprintf("%sexport", strings.Replace(extId, "-", "", -1)))); err != nil { + return nil, err + } else if len(imageList.Objects) != 1 { + return nil, httperrors.NewResourceNotFoundError("exported image not find") + } else if err := bucket.DownloadFile(imageList.Objects[0].Key, tmpImageFile, 12*1024*1024, oss.Routines(3), oss.Progress(&OssProgressListener{})); err != nil { + return nil, err + } else { + s := auth.GetAdminSession(options.Options.Region, "") + params := jsonutils.Marshal(map[string]string{"image_id": imageId, "disk-format": "raw"}) + if file, err := os.Open(tmpImageFile); err != nil { + return nil, err + } else if result, err := modules.Images.Upload(s, params, file, imageList.Objects[0].Size); err != nil { + return nil, err + } else { + os.Remove(tmpImageFile) + return result, nil + } + } +}