diff --git a/pkg/compute/guestdrivers/aws.go b/pkg/compute/guestdrivers/aws.go index 83a69c5294..920c28b495 100644 --- a/pkg/compute/guestdrivers/aws.go +++ b/pkg/compute/guestdrivers/aws.go @@ -3,10 +3,12 @@ package guestdrivers import ( "context" "fmt" + "time" "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/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/mcclient" ) @@ -49,6 +51,43 @@ func (self *SAwsGuestDriver) ValidateCreateData(ctx context.Context, userCred mc return self.SManagedVirtualizedGuestDriver.ValidateCreateData(ctx, userCred, data) } +func fetchAwsIVMinfo(desc SManagedVMCreateConfig, iVM cloudprovider.ICloudVM, guestId string) *jsonutils.JSONDict { + data := jsonutils.NewDict() + data.Add(jsonutils.NewString(iVM.GetOSType()), "os") + if len(desc.OsDistribution) > 0 { + data.Add(jsonutils.NewString(desc.OsDistribution), "distro") + } + if len(desc.OsVersion) > 0 { + data.Add(jsonutils.NewString(desc.OsVersion), "version") + } + + idisks, err := iVM.GetIDisks() + + if err != nil { + log.Errorf("GetiDisks error %s", err) + } else { + diskInfo := make([]SDiskInfo, len(idisks)) + for i := 0; i < len(idisks); i += 1 { + dinfo := SDiskInfo{} + dinfo.Uuid = idisks[i].GetGlobalId() + dinfo.Size = idisks[i].GetDiskSizeMB() + dinfo.DiskType = idisks[i].GetDiskType() + if metaData := idisks[i].GetMetadata(); metaData != nil { + dinfo.Metadata = make(map[string]string, 0) + if err := metaData.Unmarshal(dinfo.Metadata); err != nil { + log.Errorf("Get disk %s metadata info error: %v", idisks[i].GetName(), err) + } + } + diskInfo[i] = dinfo + } + data.Add(jsonutils.Marshal(&diskInfo), "disks") + } + + data.Add(jsonutils.NewString(iVM.GetGlobalId()), "uuid") + data.Add(iVM.GetMetadata(), "metadata") + return data +} + func (self *SAwsGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error { config := guest.GetDeployConfigOnHost(ctx, host, task.GetParams()) log.Debugf("RequestDeployGuestOnHost: %s", config) @@ -58,16 +97,141 @@ func (self *SAwsGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest return err } + ihost, err := host.GetIHost() + if err != nil { + return err + } + + desc := SManagedVMCreateConfig{} + err = config.Unmarshal(&desc, "desc") + if err != nil { + return err + } + publicKey, _ := config.GetString("public_key") + passwd, _ := config.GetString("password") + switch action { case "create": - return nil + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error){ + nets := guest.GetNetworks() + net := nets[0].GetNetwork() + vpc := net.GetVpc() + + ivpc, err := vpc.GetIVpc() + if err != nil { + log.Errorf("getIVPC fail %s", err) + return nil, err + } + + secgrpId, err := ivpc.SyncSecurityGroup(desc.SecGroupId, desc.SecGroupName, desc.SecRules) + if err != nil { + log.Errorf("SyncSecurityGroup fail %s", err) + return nil, err + } + + iVM, err := ihost.CreateVM(desc.Name, desc.ExternalImageId, desc.SysDiskSize, desc.Cpu, desc.Memory, desc.ExternalNetworkId, + desc.IpAddr, desc.Description, "", desc.StorageType, desc.DataDisks, publicKey, secgrpId) + if err != nil { + return nil, err + } + log.Debugf("VMcreated %s, wait status ready ...", iVM.GetGlobalId()) + err = cloudprovider.WaitStatus(iVM, models.VM_READY, time.Second*5, time.Second*1800) + if err != nil { + return nil, err + } + log.Debugf("VMcreated %s, and status is ready", iVM.GetGlobalId()) + + iVM, err = ihost.GetIVMById(iVM.GetGlobalId()) + if err != nil { + log.Errorf("cannot find vm %s", err) + return nil, err + } + + data := fetchAwsIVMinfo(desc, iVM, guest.Id) + return data, nil + }) case "deploy": - return nil + iVM, err := ihost.GetIVMById(guest.GetExternalId()) + if err != nil || iVM == nil { + log.Errorf("cannot find vm %s", err) + return fmt.Errorf("cannot find vm") + } + + params := task.GetParams() + log.Debugf("Deploy VM params %s", params.String()) + + name, _ := params.GetString("name") + description, _ := params.GetString("description") + publicKey, _ := config.GetString("public_key") + deleteKeypair := jsonutils.QueryBoolean(params, "__delete_keypair__", false) + + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error){ + err := iVM.DeployVM(name, passwd, publicKey, deleteKeypair, description) + if err != nil { + return nil, err + } + + data := fetchIVMinfo(desc, iVM, guest.Id, passwd) + return data, nil + }) case "rebuild": - return nil + iVM, err := ihost.GetIVMById(guest.GetExternalId()) + if err != nil || iVM == nil { + log.Errorf("cannot find vm %s", err) + return fmt.Errorf("cannot find vm") + } + + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error){ + diskId, err := iVM.RebuildRoot(desc.ExternalImageId, passwd, publicKey, desc.SysDiskSize) + if err != nil { + return nil, err + } + + log.Debugf("VMrebuildRoot %s new diskID %s, wait status ready ...", iVM.GetGlobalId(), diskId) + + err = cloudprovider.WaitStatus(iVM, models.VM_READY, time.Second*5, time.Second*1800) + if err != nil { + return nil, err + } + log.Debugf("VMrebuildRoot %s, and status is ready", iVM.GetGlobalId()) + + maxWaitSecs := 300 + waited := 0 + + for { + // hack, wait disk number consistent + idisks, err := iVM.GetIDisks() + if err != nil { + log.Errorf("fail to find VM idisks %s", err) + return nil, err + } + if len(idisks) < len(desc.DataDisks)+1 { + if waited > maxWaitSecs { + log.Errorf("inconsistent disk number, wait timeout, must be something wrong on remote") + return nil, cloudprovider.ErrTimeout + } + log.Debugf("inconsistent disk number???? %d != %d", len(idisks), len(desc.DataDisks)+1) + time.Sleep(time.Second * 5) + waited += 5 + } else { + if idisks[0].GetGlobalId() != diskId { + log.Errorf("system disk id inconsistent %s != %s", idisks[0].GetGlobalId(), diskId) + return nil, fmt.Errorf("inconsistent sys disk id after rebuild root") + } + + break + } + } + + data := fetchIVMinfo(desc, iVM, guest.Id, passwd) + + return data, nil + }) default: - return nil + log.Errorf("RequestDeployGuestOnHost: Action %s not supported", action) + return fmt.Errorf("Action %s not supported", action) } + return nil } diff --git a/pkg/util/aws/disk.go b/pkg/util/aws/disk.go index ac1e81a703..aa50d851d2 100644 --- a/pkg/util/aws/disk.go +++ b/pkg/util/aws/disk.go @@ -367,7 +367,7 @@ func (self *SRegion) resetDisk(diskId, snapshotId string) error { return cloudprovider.ErrNotImplemented } -func (self *SRegion) CreateDisk(zoneId string, category string, name string, sizeGb int, desc string) (string, error) { +func (self *SRegion) CreateDisk(zoneId string, category string, name string, sizeGb int, snapshotId string, desc string) (string, error) { tagspec := TagSpec{ResourceType: "volume"} tagspec.SetNameTag(name) tagspec.SetDescTag(desc) @@ -377,6 +377,9 @@ func (self *SRegion) CreateDisk(zoneId string, category string, name string, siz params.SetAvailabilityZone(zoneId) params.SetVolumeType(category) params.SetSize(int64(sizeGb)) + if len(snapshotId) > 0 { + params.SetSnapshotId(snapshotId) + } params.SetTagSpecifications([]*ec2.TagSpecification{ec2Tags}) ret, err := self.ec2Client.CreateVolume(params) diff --git a/pkg/util/aws/image.go b/pkg/util/aws/image.go index 911c3b8102..bd5ae080f9 100644 --- a/pkg/util/aws/image.go +++ b/pkg/util/aws/image.go @@ -32,6 +32,12 @@ type ImageImportTask struct { TaskId string } +type RootDevice struct { + SnapshotId string + Size int // GB + Category string // VolumeType +} + type SImage struct { storageCache *SStoragecache @@ -49,6 +55,7 @@ type SImage struct { Size int Status ImageStatusType Usage string + RootDevice RootDevice } func (self *SImage) GetId() string { @@ -241,6 +248,15 @@ func (self *SRegion) GetImages(status ImageStatusType, owner ImageOwnerType, ima log.Debugf(err.Error()) } + var rootDevice RootDevice + for _, block := range image.BlockDeviceMappings { + if len(*image.RootDeviceName) > 0 && *block.DeviceName == *image.RootDeviceName { + rootDevice.SnapshotId = *block.Ebs.SnapshotId + rootDevice.Category = *block.Ebs.VolumeType + rootDevice.Size = int(*block.Ebs.VolumeSize) + } + } + images = append(images, SImage{ storageCache: self.getStoragecache(), Architecture: *image.Architecture, @@ -253,6 +269,7 @@ func (self *SRegion) GetImages(status ImageStatusType, owner ImageOwnerType, ima Status: ImageStatusType(*image.State), CreationTime: *image.CreationDate, Size: size, + RootDevice: rootDevice, // Usage: "", // OSName: *image.Platform, }) diff --git a/pkg/util/aws/instance.go b/pkg/util/aws/instance.go index 1ceef457f2..dad2bd9e9f 100644 --- a/pkg/util/aws/instance.go +++ b/pkg/util/aws/instance.go @@ -340,7 +340,16 @@ func (self *SInstance) UpdateVM(name string) error { } func (self *SInstance) RebuildRoot(imageId string, passwd string, publicKey string, sysSizeGB int) (string, error) { - panic("implement me") + if len(publicKey) > 0 || len(passwd) > 0{ + return "", fmt.Errorf("aws rebuild root not support specific publickey/password") + } + + diskId, err := self.host.zone.region.ReplaceSystemDisk(self.InstanceId, imageId, sysSizeGB) + if err != nil { + return "", err + } + + return diskId, nil } func (self *SInstance) DeployVM(name string, password string, publicKey string, deleteKeypair bool, description string) error { @@ -692,8 +701,53 @@ func (self *SRegion) UpdateVM(instanceId string, hostname string) error { return fmt.Errorf("aws not support change hostname.") } -func (self *SRegion) ReplaceSystemDisk(instanceId string, imageId string, passwd string, keypairName string, sysDiskSizeGB int) (string, error) { - return "", cloudprovider.ErrNotSupported +func (self *SRegion) ReplaceSystemDisk(instanceId string, imageId string, sysDiskSizeGB int) (string, error) { + instance, err := self.GetInstance(instanceId) + if err != nil { + return "", err + } + + disks, _, err := self.GetDisks(instanceId, instance.ZoneId, "", nil, 0, 0) + if err != nil { + return "", err + } + + var rootDisk *SDisk + for _,disk := range disks { + if disk.Type == models.DISK_TYPE_SYS { + rootDisk = &disk + } + } + + if rootDisk == nil { + return "", fmt.Errorf("can not find root disk of instance %s", instanceId) + } + + image,err := self.GetImage(imageId) + if err != nil { + return "", err + } + + diskId, err := self.CreateDisk(instance.ZoneId, rootDisk.Category, rootDisk.GetName(), sysDiskSizeGB, image.RootDevice.SnapshotId, "") + if err != nil { + return "", err + } + + self.ec2Client.WaitUntilVolumeAvailable(&ec2.DescribeVolumesInput{VolumeIds: []*string{&diskId}}) + // todo: 检查instance状态 + err = instance.DetachDisk(rootDisk.DiskId) + if err != nil { + return "", err + } + + self.ec2Client.WaitUntilInstanceStopped(&ec2.DescribeInstancesInput{InstanceIds: []*string{&instanceId}}) + err = instance.AttachDisk(diskId) + if err != nil { + return "", err + } + + self.ec2Client.WaitUntilInstanceStopped(&ec2.DescribeInstancesInput{InstanceIds: []*string{&instanceId}}) + return diskId, nil } func (self *SRegion) ChangeVMConfig(zoneId string, instanceId string, ncpu int, vmem int, disks []*SDisk) error { diff --git a/pkg/util/aws/shell/instance.go b/pkg/util/aws/shell/instance.go index 8076c49360..ba5430eda4 100644 --- a/pkg/util/aws/shell/instance.go +++ b/pkg/util/aws/shell/instance.go @@ -130,13 +130,11 @@ func init() { type InstanceRebuildRootOptions struct { ID string `help:"instance ID"` Image string `help:"Image ID"` - Password string `help:"pasword"` - Keypair string `help:"keypair name"` Size int `help:"system disk size in GB"` } shellutils.R(&InstanceRebuildRootOptions{}, "instance-rebuild-root", "Reinstall virtual server system image", func(cli *aws.SRegion, args *InstanceRebuildRootOptions) error { - diskID, err := cli.ReplaceSystemDisk(args.ID, args.Image, args.Password, args.Keypair, args.Size) + diskID, err := cli.ReplaceSystemDisk(args.ID, args.Image, args.Size) if err != nil { return err } diff --git a/pkg/util/aws/storage.go b/pkg/util/aws/storage.go index 0a191e89b3..400a18b67d 100644 --- a/pkg/util/aws/storage.go +++ b/pkg/util/aws/storage.go @@ -100,7 +100,7 @@ func (self *SStorage) GetManagerId() string { } func (self *SStorage) CreateIDisk(name string, sizeGb int, desc string) (cloudprovider.ICloudDisk, error) { - diskId, err := self.zone.region.CreateDisk(self.zone.ZoneId, self.storageType, name, sizeGb, desc) + diskId, err := self.zone.region.CreateDisk(self.zone.ZoneId, self.storageType, name, sizeGb, "", desc) if err != nil { log.Errorf("createDisk fail %s", err) return nil, err