From 5eca3c4c97a8a7b2ee4e2a1254a695a77e81661b Mon Sep 17 00:00:00 2001 From: TangBin Date: Mon, 29 Oct 2018 14:42:29 +0800 Subject: [PATCH] add driver --- pkg/compute/guestdrivers/aws.go | 63 +++++++++++ pkg/compute/hostdrivers/aws.go | 192 ++++++++++++++++++++++++++++++++ 2 files changed, 255 insertions(+) diff --git a/pkg/compute/guestdrivers/aws.go b/pkg/compute/guestdrivers/aws.go index 32dd08a21e..247aad9ab1 100644 --- a/pkg/compute/guestdrivers/aws.go +++ b/pkg/compute/guestdrivers/aws.go @@ -1 +1,64 @@ package guestdrivers + +import ( + "context" + "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/mcclient" +) + +type SAwsGuestDriver struct { + SManagedVirtualizedGuestDriver +} + +func (self *SAwsGuestDriver) GetHypervisor() string { + return models.HYPERVISOR_AWS +} + +func (self *SAwsGuestDriver) ChooseHostStorage(host *models.SHost, backend string) *models.SStorage { + storages := host.GetAttachedStorages("") + for i := 0; i < len(storages); i += 1 { + if storages[i].StorageType == backend { + return &storages[i] + } + } + + for _, stype := range []string{"gp2", "io1", "st1", "sc1", "standard"} { + for i := 0; i < len(storages); i += 1 { + if storages[i].StorageType == stype { + return &storages[i] + } + } + } + return nil +} + +func (self *SAwsGuestDriver) GetDetachDiskStatus() ([]string, error) { + return []string{models.VM_READY, models.VM_RUNNING}, nil +} + +func (self *SAwsGuestDriver) RequestDetachDisk(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { + return guest.StartSyncTask(ctx, task.GetUserCred(), false, task.GetTaskId()) +} + +func (self *SAwsGuestDriver) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) { + return self.SManagedVirtualizedGuestDriver.ValidateCreateData(ctx, userCred, data) +} + +func (self *SAwsGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error { + return nil +} + +func (self *SAwsGuestDriver) OnGuestDeployTaskDataReceived(ctx context.Context, guest *models.SGuest, task taskman.ITask, data jsonutils.JSONObject) error { + return nil +} + +func (self *SAwsGuestDriver) RequestDiskSnapshot(ctx context.Context, guest *models.SGuest, task taskman.ITask, snapshotId, diskId string) error { + return nil +} + +func init() { + driver := SAwsGuestDriver{} + models.RegisterGuestDriver(&driver) +} \ No newline at end of file diff --git a/pkg/compute/hostdrivers/aws.go b/pkg/compute/hostdrivers/aws.go index a58d95a88b..611de61c9c 100644 --- a/pkg/compute/hostdrivers/aws.go +++ b/pkg/compute/hostdrivers/aws.go @@ -1 +1,193 @@ package hostdrivers + +import ( + "context" + "fmt" + "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" + "yunion.io/x/onecloud/pkg/cloudprovider" + "yunion.io/x/onecloud/pkg/httperrors" + + "yunion.io/x/jsonutils" + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" +) + +type SAwsHostDriver struct { + SBaseHostDriver +} + +func (self *SAwsHostDriver) GetHostType() string { + return models.HOST_TYPE_AWS +} + +func (self *SAwsHostDriver) 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 { + return err + } + + osArch, _ := params.GetString("os_arch") + 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(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 + } + }) + return nil +} + +func (self *SAwsHostDriver) RequestPrepareSaveDiskOnHost(ctx context.Context, host *models.SHost, disk *models.SDisk, imageId string, task taskman.ITask) error { + task.ScheduleRun(nil) + return nil +} + +func (self *SAwsHostDriver) 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 { + params := task.GetParams() + osType, _ := params.GetString("properties", "os_type") + + 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), osType, ""); 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 *SAwsHostDriver) 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 + } else { + if size, err := content.Int("size"); err != nil { + return err + } else { + size = size >> 10 + if iDisk, err := iCloudStorage.CreateIDisk(disk.GetName(), int(size), ""); err != nil { + return err + } else { + if _, err := disk.GetModelManager().TableSpec().Update(disk, func() error { + disk.ExternalId = iDisk.GetGlobalId() + + if metaData := iDisk.GetMetadata(); metaData != nil { + meta := make(map[string]string) + if err := metaData.Unmarshal(meta); err != nil { + log.Errorf("Get disk %s Metadata error: %v", disk.Name, err) + } else { + for key, value := range meta { + if err := disk.SetMetadata(ctx, key, value, task.GetUserCred()); err != nil { + log.Errorf("set disk %s mata %s => %s error: %v", disk.Name, key, value, err) + } + } + } + } + + return nil + }); err != nil { + log.Errorf("Update disk externalId err: %v", err) + return err + } + data := jsonutils.NewDict() + data.Add(jsonutils.NewInt(int64(iDisk.GetDiskSizeMB())), "disk_size") + data.Add(jsonutils.NewString(iDisk.GetDiskFormat()), "disk_format") + task.ScheduleRun(data) + } + } + } + return nil +} + +func (self *SAwsHostDriver) RequestDeallocateDiskOnHost(host *models.SHost, storage *models.SStorage, disk *models.SDisk, task taskman.ITask) error { + data := jsonutils.NewDict() + if iCloudStorage, err := storage.GetIStorage(); err != nil { + return err + } else if iDisk, err := iCloudStorage.GetIDisk(disk.GetExternalId()); err != nil { + if err == cloudprovider.ErrNotFound { + task.ScheduleRun(data) + return nil + } + return err + } else if err := iDisk.Delete(); err != nil { + return err + } + task.ScheduleRun(data) + return nil +} + +func (self *SAwsHostDriver) RequestResizeDiskOnHostOnline(host *models.SHost, storage *models.SStorage, disk *models.SDisk, size int64, task taskman.ITask) error { + return self.RequestResizeDiskOnHost(host, storage, disk, size, task) +} + +func (self *SAwsHostDriver) RequestResizeDiskOnHost(host *models.SHost, storage *models.SStorage, disk *models.SDisk, size int64, task taskman.ITask) error { + if iCloudStorage, err := storage.GetIStorage(); err != nil { + return err + } else if iDisk, err := iCloudStorage.GetIDisk(disk.GetExternalId()); err != nil { + return err + } else if err := iDisk.Resize(size >> 10); err != nil { + return err + } else { + task.ScheduleRun(jsonutils.Marshal(map[string]int64{"disk_size": size})) + } + return nil +} + +func init() { + driver := SAwsHostDriver{} + models.RegisterHostDriver(&driver) +}