Automatic merge from release/2.1.0 -> release/2.2.0

* commit '031ed3a3e8c9704e2545eb61ca3a5c45350b2875':
  支持server-save-image
This commit is contained in:
邱剑
2018-09-04 16:35:26 +08:00
17 changed files with 813 additions and 20 deletions
+13
View File
@@ -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)
}
@@ -201,9 +205,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
+5
View File
@@ -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")
}
@@ -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
}
+49
View File
@@ -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
+23
View File
@@ -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
}
+90
View File
@@ -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(), &quota); 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")
+2
View File
@@ -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
+37
View File
@@ -1561,6 +1561,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)
}
+2
View File
@@ -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
+107
View File
@@ -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"}))
}
}
@@ -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)
}
+27 -20
View File
@@ -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
+107
View File
@@ -267,3 +267,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
}
}
+31
View File
@@ -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 {
+45
View File
@@ -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
})
}
+73
View File
@@ -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
}
+128
View File
@@ -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
}
}
}