add image download/upload method

This commit is contained in:
TangBin
2018-10-23 16:26:53 +08:00
parent aa10e8c968
commit 7db7bcb2ca
4 changed files with 199 additions and 19 deletions
+26 -2
View File
@@ -1,12 +1,12 @@
package aws
import (
"fmt"
"github.com/aws/aws-sdk-go/service/ec2"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
"fmt"
"github.com/aws/aws-sdk-go/service/ec2"
)
type ImageStatusType string
@@ -137,6 +137,30 @@ func (self *SRegion) ImportImage(name string, osArch string, osType string, osDi
return &ImageImportTask{ImageId: *ret.ImageId, RegionId: self.RegionId, TaskId: *ret.ImportTaskId}, nil
}
type ImageExportTask struct {
ImageId string
RegionId string
TaskId string
}
func (self *SRegion) ExportImage(instanceId string,imageId string) (*ImageExportTask, error) {
params := &ec2.CreateInstanceExportTaskInput{}
params.SetInstanceId(instanceId)
params.SetDescription(fmt.Sprintf("image %s export from aws", imageId))
params.SetTargetEnvironment("vmware")
spec := &ec2.ExportToS3TaskSpecification{}
spec.SetContainerFormat("ova")
spec.SetDiskImageFormat("RAW")
spec.SetS3Bucket("imgcache-onecloud")
params.SetExportToS3Task(spec)
ret, err := self.ec2Client.CreateInstanceExportTask(params)
if err != nil {
return nil, err
}
return &ImageExportTask{ImageId: imageId, RegionId: self.RegionId, TaskId: *ret.ExportTask.ExportTaskId}, nil
}
func (self *SRegion) GetImage(imageId string) (*SImage, error) {
images, _, err := self.GetImages("", ImageOwnerSelf, []string{imageId}, "", 0, 1)
if err != nil {
+18
View File
@@ -444,6 +444,24 @@ func (self *SRegion) GetInstance(instanceId string) (*SInstance, error) {
return &instances[0], nil
}
func (self *SRegion) GetInstanceIdByImageId(imageId string) (string, error) {
params := &ec2.DescribeInstancesInput{}
filters := []*ec2.Filter{}
filters = AppendSingleValueFilter(filters, "image-id", imageId)
params.SetFilters(filters)
ret, err := self.ec2Client.DescribeInstances(params)
if err != nil {
return "", err
}
for _, item := range ret.Reservations {
for _, instance := range item.Instances {
return *instance.InstanceId, nil
}
}
return "", fmt.Errorf("instance launch with image %s not found", imageId)
}
func (self *SRegion) CreateInstance(name string, imageId string, instanceType string, SubnetId string, securityGroupId string,
zoneId string, desc string, passwd string, disks []SDisk, ipAddr string,
keypair string) (string, error) {
+43 -5
View File
@@ -1,19 +1,23 @@
package aws
import (
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/jsonutils"
"fmt"
"yunion.io/x/onecloud/pkg/compute/models"
sdk "github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/service/ec2"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/aws/aws-sdk-go/aws/credentials"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/aws/aws-sdk-go/service/ec2"
"github.com/aws/aws-sdk-go/service/iam"
"github.com/aws/aws-sdk-go/service/s3"
"yunion.io/x/jsonutils"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
)
type SRegion struct {
client *SAwsClient
ec2Client *ec2.EC2
iamClient *iam.IAM
s3Client *s3.S3
izones []cloudprovider.ICloudZone
ivpcs []cloudprovider.ICloudVpc
@@ -48,6 +52,40 @@ func (self *SRegion) getEc2Client() (*ec2.EC2, error) {
return self.ec2Client, nil
}
func (self *SRegion) getIamClient() (*iam.IAM, error) {
if self.iamClient == nil {
s, err := session.NewSession(&sdk.Config{
Region: sdk.String(self.RegionId),
Credentials: credentials.NewStaticCredentials(self.client.accessKey, self.client.secret, ""),
})
if err != nil {
return nil, err
}
self.iamClient = iam.New(s)
}
return self.iamClient, nil
}
func (self *SRegion) getS3Client() (*s3.S3, error) {
if self.s3Client == nil {
s, err := session.NewSession(&sdk.Config{
Region: sdk.String(self.RegionId),
Credentials: credentials.NewStaticCredentials(self.client.accessKey, self.client.secret, ""),
})
if err != nil {
return nil, err
}
self.s3Client = s3.New(s)
}
return self.s3Client, nil
}
/////////////////////////////////////////////////////////////////////////////
func (self *SRegion) fetchZones() error {
// todo: 这里将过滤出指定region下全部的zones。是否只过滤出可用的zone即可? The state of the Availability Zone (available | information | impaired | unavailable)
+112 -12
View File
@@ -1,15 +1,20 @@
package aws
import (
"bytes"
"fmt"
"github.com/aws/aws-sdk-go/service/ec2"
"github.com/aws/aws-sdk-go/service/iam"
"github.com/aws/aws-sdk-go/service/s3"
"io/ioutil"
"time"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/options"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules"
)
type SStoragecache struct {
@@ -122,18 +127,114 @@ func (self *SStoragecache) fetchImages() error {
func (self *SStoragecache) uploadImage(userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, isForce bool) (string, error) {
// todo: implement me
bucketName := "imgcache-onecloud"
err := self.region.initVmimport()
if err != nil {
return "", err
}
// first upload image to oss
s := auth.GetAdminSession(options.Options.Region, "")
meta, reader, err := modules.Images.Download(s, imageId)
if err != nil {
return "", err
}
log.Infof("meta data %s", meta)
s3Client, err := self.region.getS3Client()
if err != nil {
return "", nil
}
// 内存?
f, err := ioutil.ReadAll(reader)
params := &s3.PutObjectInput{}
params.SetBucket(bucketName)
params.SetKey(imageId)
params.SetBody(bytes.NewReader(f))
_, err = s3Client.PutObject(params)
if err != nil {
return "", nil
}
imageBaseName := imageId
if imageBaseName[0] >= '0' && imageBaseName[0] <= '9' {
imageBaseName = fmt.Sprintf("img%s", imageId)
}
imageName := imageBaseName
nameIdx := 1
// check image name, avoid name conflict
for {
_, err = self.region.GetImageByName(imageName)
if err != nil {
if err == cloudprovider.ErrNotFound {
break
} else {
return "", err
}
}
imageName = fmt.Sprintf("%s-%d", imageBaseName, nameIdx)
nameIdx += 1
}
task, err := self.region.ImportImage(imageName, osArch, osType, osDist, bucketName, imageId)
if err != nil {
log.Errorf("ImportImage error %s %s %s", imageId, bucketName, err)
return "", err
}
// todo:// 等待镜像导入完成
err = self.region.ec2Client.WaitUntilImageExists(&ec2.DescribeImagesInput{ImageIds:[]*string{&task.ImageId}})
return task.ImageId, err
return "", nil
}
func (self *SStoragecache) downloadImage(userCred mcclient.TokenCredential, imageId string, extId string) (jsonutils.JSONObject, error) {
// todo: implement me
return nil, nil
// aws 导出镜像限制比较多。https://docs.aws.amazon.com/zh_cn/vm-import/latest/userguide/vmexport.html
bucketName := "imgcache-onecloud"
if err := self.region.checkBucket(bucketName); err != nil {
return nil, err
}
instanceId, err := self.region.GetInstanceIdByImageId(extId)
if err != nil {
return nil, err
}
task, err := self.region.ExportImage(instanceId, imageId)
if err != nil {
return nil, err
}
taskParams := &ec2.DescribeExportTasksInput{}
taskParams.SetExportTaskIds([]*string{&task.TaskId})
if err := self.region.ec2Client.WaitUntilExportTaskCompleted(taskParams); err != nil {
return nil, err
}
s3Client, err := self.region.getS3Client()
if err != nil {
return nil, err
}
i := &s3.GetObjectInput{}
i.SetBucket(bucketName)
i.SetKey(fmt.Sprintf("%s.%s", task.TaskId, "ova"))
ret, err := s3Client.GetObject(i)
if err != nil {
return nil, err
}
s := auth.GetAdminSession(options.Options.Region, "")
params := jsonutils.Marshal(map[string]string{"image_id": imageId, "disk-format": "raw"})
if result, err := modules.Images.Upload(s, params, ret.Body, IntVal(ret.ContentLength)); err != nil {
return nil, err
} else {
return result, nil
}
}
func (self *SRegion) CheckBucket(bucketName string) error {
@@ -155,11 +256,11 @@ func (self *SRegion) checkBucket(bucketName string) error {
}
func (self *SRegion) IsBucketExist(bucketName string) (bool, error) {
session, err := self.client.getDefaultSession()
s3Client, err := self.getS3Client()
if err != nil {
return false, err
}
s3Client := s3.New(session)
params := &s3.ListBucketsInput{}
ret, err := s3Client.ListBuckets(params)
if err != nil {
@@ -177,11 +278,11 @@ func (self *SRegion) IsBucketExist(bucketName string) (bool, error) {
func (self *SRegion) initVmimportRole() error {
/*需要api access token 具备iam Full access权限*/
session, err := self.client.getDefaultSession()
iamClient, err := self.getIamClient()
if err != nil {
return err
}
iamClient := iam.New(session)
// search role vmimport
rolename := "vmimport"
ret, _ := iamClient.GetRole(&iam.GetRoleInput{RoleName: &rolename})
@@ -217,11 +318,11 @@ func (self *SRegion) initVmimportRole() error {
func (self *SRegion) initVmimportRolePolicy() error {
/*需要api access token 具备iam Full access权限*/
session, err := self.client.getDefaultSession()
iamClient, err := self.getIamClient()
if err != nil {
return err
}
iamClient := iam.New(session)
roleName := "vmimport"
policyName := "vmimport"
ret, err := iamClient.GetRolePolicy(&iam.GetRolePolicyInput{RoleName: &roleName, PolicyName: &policyName})
@@ -277,12 +378,11 @@ func (self *SRegion) initVmimportBucket() error {
return nil
}
// todo: 这里提取出get s3 session.
session, err := self.client.getDefaultSession()
s3Client, err := self.getS3Client()
if err != nil {
return err
}
s3Client := s3.New(session)
_, err = s3Client.CreateBucket(&s3.CreateBucketInput{Bucket: &bucketName})
return err
}