diff --git a/Gopkg.lock b/Gopkg.lock index ccd7d5e941..99299205d0 100644 --- a/Gopkg.lock +++ b/Gopkg.lock @@ -666,13 +666,15 @@ version = "v1.2.0" [[projects]] - digest = "1:d867dfa6751c8d7a435821ad3b736310c2ed68945d05b50fb9d23aee0540c8cc" branch = "master" + digest = "1:c1f242537d2898c84fe9bfd8bd94b177cd5ff98a10c397dbbc494f4e2b173be6" name = "github.com/serialx/hashring" packages = ["."] + pruneopts = "UT" revision = "49a4782e9908fe098c907022a1bd7519c79803d6" [[projects]] + digest = "1:d867dfa6751c8d7a435821ad3b736310c2ed68945d05b50fb9d23aee0540c8cc" name = "github.com/sirupsen/logrus" packages = ["."] pruneopts = "UT" @@ -1032,26 +1034,26 @@ [[projects]] branch = "master" - digest = "1:5ede93047a3e04f6ede8fbc011f13ee94c397b727613fd9712cc99ebc47ff906" + digest = "1:d0257638bb52243f9fa293ef07544081259d54edd8e53f084e522b6c22e92637" name = "yunion.io/x/jsonutils" packages = ["."] pruneopts = "UT" - revision = "41e805b221e8fcd9435b706b4dd473df3dd8cff4" + revision = "38477c9cceb895816fe21507d73da30376d358b7" [[projects]] branch = "master" - digest = "1:b707a5d591b21dad9ebdb264281632ed43d34e2b980dce967acb2c7adcc93ad3" + digest = "1:636db180d3fc734536a5437d93bb6fe937feb036c1ed062c25872a47c4283163" name = "yunion.io/x/log" packages = [ ".", "hooks", ] pruneopts = "UT" - revision = "a10b94c5038480920262e0f344656328434955a2" + revision = "5185c49f6d361f61e561fbccba1bedcc7c62748b" [[projects]] branch = "master" - digest = "1:13da776b435ec7873c6f9309191296d920a110145372d2a9e1a58ab9e8b34b0d" + digest = "1:8a2992edca980cb03eb9dc89f62c03c8ab6fe0c5ca4fc1da71eb736aae635f65" name = "yunion.io/x/pkg" packages = [ "gotypes", @@ -1086,29 +1088,30 @@ "utils", ] pruneopts = "UT" - revision = "56b426f0ca15288cc26442f5802c949f137ade60" + revision = "1e97a2736e5b90c0274b7d391b8d97aa9107725c" [[projects]] branch = "master" - digest = "1:3d3e1ecb41c63448df5d66eeafa5f6da720ca6e32543e89faa82b5bc4b9adebc" + digest = "1:e23019db2cda58480528738c86bc7abbff09d10488c2a51dd5d500f8686815f0" name = "yunion.io/x/sqlchemy" packages = ["."] pruneopts = "UT" - revision = "a74ef73e555a1ed2e19b1a47e0567814c289eaf9" + revision = "19a3c90524d3638f8e3580231396dda46088aebc" [[projects]] branch = "master" - digest = "1:6baf7b4ee14e156a9dd6c5c850bb8201af47665e3267e3cdee26edf2a90eaf73" + digest = "1:71e1b62648868a9083f9e4693145a0154a600591ca548aaf64cae32a10e956e5" name = "yunion.io/x/structarg" packages = ["."] pruneopts = "UT" - revision = "d5e5d87357b9bc2164215117763f6b15b2a3e75d" + revision = "4a5eb8e2cdfbf7f6511561a3c710c116b171f405" [solve-meta] analyzer-name = "dep" analyzer-version = 1 input-imports = [ "github.com/aliyun/alibaba-cloud-sdk-go/sdk", + "github.com/aliyun/alibaba-cloud-sdk-go/sdk/errors", "github.com/aliyun/alibaba-cloud-sdk-go/sdk/requests", "github.com/aliyun/aliyun-oss-go-sdk/oss", "github.com/aokoli/goutils", @@ -1147,6 +1150,7 @@ "github.com/mholt/caddy/startupshutdown", "github.com/miekg/dns", "github.com/moul/http2curl", + "github.com/serialx/hashring", "github.com/stretchr/testify/assert", "github.com/vmware/govmomi", "github.com/vmware/govmomi/object", @@ -1158,6 +1162,7 @@ "golang.org/x/crypto/ssh", "gopkg.in/gin-gonic/gin.v1", "k8s.io/api/core/v1", + "k8s.io/apimachinery/pkg/api/errors", "k8s.io/apimachinery/pkg/apis/meta/v1", "k8s.io/apimachinery/pkg/fields", "k8s.io/apimachinery/pkg/labels", diff --git a/Makefile b/Makefile index 923270969f..ee09d3df04 100644 --- a/Makefile +++ b/Makefile @@ -83,10 +83,6 @@ bin_dir: output_dir output_dir: @mkdir -p $(BUILD_DIR) - -dep: - cd $(ROOT_DIR) && dep ensure -v - dep_clean: rm -fr $(GOPATH)/pkg/dep/sources/* @@ -105,5 +101,8 @@ fmt: find . -type f -name "*.go" -not -path "./_output/*" \ -not -path "./vendor/*" | xargs gofmt -s -w +dep: + cd $(ROOT_DIR) && dep ensure -v -update $(shell for p in $$(ls vendor/yunion.io/x/); do echo "yunion.io/x/$$p"; done | xargs) + %: @: diff --git a/build/region-dns/vars b/build/region-dns/vars new file mode 100644 index 0000000000..778d49a5da --- /dev/null +++ b/build/region-dns/vars @@ -0,0 +1,2 @@ +DESCRIPTION="Yunion Cloud Region DNS Service" +# SERVICE="yes" diff --git a/cmd/climc/shell/keypairs.go b/cmd/climc/shell/keypairs.go index a13ae163a7..d56d2122a2 100644 --- a/cmd/climc/shell/keypairs.go +++ b/cmd/climc/shell/keypairs.go @@ -113,9 +113,9 @@ func init() { if len(args.PublicKey) > 0 { content, e := ioutil.ReadFile(args.PublicKey) if e != nil { - params.Add(jsonutils.NewString(args.PublicKey), "public_Key") + params.Add(jsonutils.NewString(args.PublicKey), "public_key") } else { - params.Add(jsonutils.NewString(string(content)), "public_Key") + params.Add(jsonutils.NewString(string(content)), "public_key") } } else { return fmt.Errorf("no public key provided") diff --git a/cmd/climc/shell/servers.go b/cmd/climc/shell/servers.go index 845c02a16c..6e963049b0 100644 --- a/cmd/climc/shell/servers.go +++ b/cmd/climc/shell/servers.go @@ -652,14 +652,26 @@ func init() { }) type ServerRebuildRootOptions struct { - SERVER string `help:"Server to rebuild root"` - Image string `help:"New root Image template ID"` + SERVER string `help:"Server to rebuild root"` + AutoStart bool `help:"Auto start server after it is created"` + Image string `help:"New root Image template ID"` + Keypair string `help:"ssh Keypair used for login"` + Password string `help:"Default user password"` } R(&ServerRebuildRootOptions{}, "server-rebuild-root", "Rebuild VM root image with new template", func(s *mcclient.ClientSession, args *ServerRebuildRootOptions) error { params := jsonutils.NewDict() if len(args.Image) > 0 { params.Add(jsonutils.NewString(args.Image), "image_id") } + if args.AutoStart { + params.Add(jsonutils.JSONTrue, "auto_start") + } + if args.Keypair != "" { + params.Add(jsonutils.NewString(args.Keypair), "keypair") + } + if args.Password != "" { + params.Add(jsonutils.NewString(args.Password), "password") + } srv, err := modules.Servers.PerformAction(s, args.SERVER, "rebuild-root", params) if err != nil { return err diff --git a/pkg/appsrv/dispatcher/dispatcher.go b/pkg/appsrv/dispatcher/dispatcher.go index 1aeff85240..ea95a458f9 100644 --- a/pkg/appsrv/dispatcher/dispatcher.go +++ b/pkg/appsrv/dispatcher/dispatcher.go @@ -209,8 +209,7 @@ func performClassActionHandler(ctx context.Context, w http.ResponseWriter, r *ht if data == nil { data = body.(*jsonutils.JSONDict) } - } - if data == nil { + } else { data = jsonutils.NewDict() } results, err := manager.PerformClassAction(ctx, params[""], query, data) diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index bc1db0d727..6b990ef325 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -252,6 +252,30 @@ func applyListItemsGeneralFilters(manager IModelManager, q *sqlchemy.SQuery, return q, nil } +func applyListItemsGeneralJointFilters(manager IModelManager, q *sqlchemy.SQuery, + userCred mcclient.TokenCredential, jointFilters []string, filterAny bool) (*sqlchemy.SQuery, error) { + for _, f := range jointFilters { + jfc := filterclause.ParseJointFilterClause(f) + if jfc != nil { + jointModelManager := GetModelManager(jfc.GetJointModelName()) + schFields := searchFields(jointModelManager, userCred) + if ok, _ := utils.InStringArray(jfc.GetField(), schFields); ok { + sq := jointModelManager.Query(jfc.RelatedKey) + cond := jfc.GetJointFilter(sq) + if cond != nil { + sq = sq.Filter(cond) + if filterAny { + q = q.Filter(sqlchemy.OR(sqlchemy.In(q.Field("id"), sq))) + } else { + q = q.Filter(sqlchemy.AND(sqlchemy.In(q.Field("id"), sq))) + } + } + } + } + } + return q, nil +} + func listItemQueryFilters(manager IModelManager, ctx context.Context, q *sqlchemy.SQuery, userCred mcclient.TokenCredential, query jsonutils.JSONObject) (*sqlchemy.SQuery, error) { @@ -273,11 +297,15 @@ func listItemQueryFilters(manager IModelManager, ctx context.Context, q *sqlchem return nil, err } } + filterAny, _ := query.Bool("filter_any") filters := jsonutils.GetQueryStringArray(query, "filter") if len(filters) > 0 { - filterAny, _ := query.Bool("filter_any") q, err = applyListItemsGeneralFilters(manager, q, userCred, filters, filterAny) } + jointFilter := jsonutils.GetQueryStringArray(query, "joint_filter") + if len(jointFilter) > 0 { + q, _ = applyListItemsGeneralJointFilters(manager, q, userCred, jointFilter, filterAny) + } return q, nil } @@ -525,7 +553,7 @@ func (dispatcher *DBModelDispatcher) tryGetModelProperty(ctx context.Context, pr } func (dispatcher *DBModelDispatcher) Get(ctx context.Context, idStr string, query jsonutils.JSONObject) (jsonutils.JSONObject, error) { - log.Debugf("Get %s", idStr) + // log.Debugf("Get %s", idStr) userCred := fetchUserCredential(ctx) data, err := dispatcher.tryGetModelProperty(ctx, idStr, query) @@ -542,7 +570,7 @@ func (dispatcher *DBModelDispatcher) Get(ctx context.Context, idStr string, quer } else if err != nil { return nil, err } - log.Debugf("Get found %s", model) + // log.Debugf("Get found %s", model) if !model.AllowGetDetails(ctx, userCred, query) { return nil, httperrors.NewForbiddenError("Not allow to get details") } @@ -987,7 +1015,6 @@ func updateItem(manager IModelManager, item IModel, ctx context.Context, userCre return nil, httperrors.NewGeneralError(err) } item.PreUpdate(ctx, userCred, query, dataDict) - diff, err := manager.TableSpec().Update(item, func() error { filterData := dataDict.CopyIncludes(updateFields(manager, userCred)...) err = filterData.Unmarshal(item) diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index 8e4f45acf6..f221ae6f11 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -289,20 +289,21 @@ func (manager *STaskManager) execTask(taskId string, data jsonutils.JSONObject) } } -func execITask(taskValue reflect.Value, task *STask, data jsonutils.JSONObject, isMulti bool) { +func execITask(taskValue reflect.Value, task *STask, odata jsonutils.JSONObject, isMulti bool) { var err error ctxData := task.GetRequestContext() ctx := ctxData.GetContext() taskFailed := false + data := odata if data != nil { taskStatus, _ := data.GetString("__status__") if len(taskStatus) > 0 && taskStatus != "OK" { taskFailed = true - data, err = data.Get("reason") + data, err = data.Get("__reason__") if err != nil { - data = jsonutils.NewString("Task failed due to unknown remote errors!") + data = jsonutils.NewString(fmt.Sprintf("Task failed due to unknown remote errors! %s", odata)) } } } else { @@ -328,7 +329,13 @@ func execITask(taskValue reflect.Value, task *STask, data jsonutils.JSONObject, if !funcValue.IsValid() || funcValue.IsNil() { msg := fmt.Sprintf("Stage %s not found", stageName) - log.Errorf(msg) + if taskFailed { + // failed handler is optional, ignore the error + log.Warningf(msg) + msg, _ = data.GetString() + } else { + log.Errorf(msg) + } task.SetStageFailed(ctx, msg) task.SaveRequestContext(&ctxData) return @@ -548,7 +555,7 @@ func (self *STask) NotifyParentTaskFailure(ctx context.Context, reason string) { if len(reason) > 100 { reason = reason[:100] + "..." } - body.Add(jsonutils.NewString(fmt.Sprintf("Subtask %s failed: %s", self.TaskName, reason))) + body.Add(jsonutils.NewString(fmt.Sprintf("Subtask %s failed: %s", self.TaskName, reason)), "__reason__") self.NotifyParentTaskComplete(ctx, body, true) } diff --git a/pkg/cloudcommon/db/virtualresource.go b/pkg/cloudcommon/db/virtualresource.go index db56d599f4..8220398d3c 100644 --- a/pkg/cloudcommon/db/virtualresource.go +++ b/pkg/cloudcommon/db/virtualresource.go @@ -161,16 +161,16 @@ func (model *SVirtualResourceBase) AllowPerformMetadata(ctx context.Context, use } func (model *SVirtualResourceBase) GetTenantCache(ctx context.Context) (*STenant, error) { - log.Debugf("Get tenant by Id %s", model.ProjectId) + // log.Debugf("Get tenant by Id %s", model.ProjectId) return TenantCacheManager.FetchTenantById(ctx, model.ProjectId) } func (model *SVirtualResourceBase) getMoreDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, extra *jsonutils.JSONDict) *jsonutils.JSONDict { if userCred.IsSystemAdmin() { - log.Debugf("GetCustomizeColumns") + // log.Debugf("GetCustomizeColumns") tobj, err := model.GetTenantCache(ctx) if err == nil { - log.Debugf("GetTenantFromCache %s", jsonutils.Marshal(tobj)) + // log.Debugf("GetTenantFromCache %s", jsonutils.Marshal(tobj)) extra.Add(jsonutils.NewString(tobj.GetName()), "tenant") } else { log.Errorf("GetTenantCache fail %s", err) diff --git a/pkg/cloudprovider/resources.go b/pkg/cloudprovider/resources.go index f388d66feb..3c333f8946 100644 --- a/pkg/cloudprovider/resources.go +++ b/pkg/cloudprovider/resources.go @@ -18,6 +18,7 @@ type ICloudResource interface { Refresh() error IsEmulated() bool + GetMetadata() *jsonutils.JSONDict } type ICloudRegion interface { @@ -153,7 +154,12 @@ type ICloudVM interface { StopVM(isForce bool) error DeleteVM() error + UpdateVM(name string) error + RebuildRoot(imageId string) error + DeployVM(name string, password string, publicKey string, resetPassword bool, deleteKeypair bool, description string) error + ChangeConfig(instanceId string, ncpu int, vmem int) error GetVNCInfo() (jsonutils.JSONObject, error) + AttachDisk(diskId string) error } type ICloudNic interface { @@ -193,7 +199,8 @@ type ICloudDisk interface { GetCacheMode() string GetMountpoint() string Delete() error - Resize(int64) error + + Resize(newSize int64) error } type ICloudVpc interface { diff --git a/pkg/compute/guestdrivers/aliyun.go b/pkg/compute/guestdrivers/aliyun.go index 9477c541f6..f149f322ec 100644 --- a/pkg/compute/guestdrivers/aliyun.go +++ b/pkg/compute/guestdrivers/aliyun.go @@ -103,7 +103,6 @@ func (self *SAliyunGuestDriver) GetJsonDescAtHost(ctx context.Context, guest *mo imageId := disk.GetTemplateId() scimg := models.StoragecachedimageManager.GetStoragecachedimage(cache.Id, imageId) config.ExternalImageId = scimg.ExternalId - img := scimg.GetCachedimage() config.OsDistribution, _ = img.Info.GetString("properties", "os_distribution") config.OsVersion, _ = img.Info.GetString("properties", "os_version") @@ -118,108 +117,159 @@ func (self *SAliyunGuestDriver) GetJsonDescAtHost(ctx context.Context, guest *mo } type SDiskInfo struct { - Size int - Uuid string + Size int + Uuid string + Metadata map[string]string } func (self *SAliyunGuestDriver) RequestDeployGuestOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, task taskman.ITask) error { config := guest.GetDeployConfigOnHost(ctx, host, task.GetParams()) - onfinish, err := config.GetString("on_finish") + /* onfinish, err := config.GetString("on_finish") if err != nil { return err - } + } */ + action, err := config.GetString("action") if err != nil { return err } - if action != "create" { - return fmt.Errorf("Action %s not supported", action) - } - ihost, err := host.GetIHost() if err != nil { return err } - desc := SAliyunVMCreateConfig{} - err = config.Unmarshal(&desc, "desc") - if err != nil { - return err - } - - taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { - passwd := seclib2.RandomPassword2(12) - - iVM, err := ihost.CreateVM(desc.Name, desc.ExternalImageId, desc.SysDiskSize, desc.Cpu, desc.Memory, desc.ExternalNetworkId, - desc.IpAddr, desc.Description, passwd, desc.StorageType, desc.DataDisks, desc.PublicKey) + if action == "create" { + desc := SAliyunVMCreateConfig{} + err = config.Unmarshal(&desc, "desc") 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 + return err } - if len(guest.SecgrpId) > 0 { - if err := iVM.SyncSecurityGroup(guest.SecgrpId, guest.GetSecgroupName(), guest.GetSecRules()); err != nil { - log.Errorf("SyncSecurityGroup error: %v", err) - return nil, err - } - } + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + passwd := seclib2.RandomPassword2(12) - if onfinish == "none" { - err = iVM.StartVM() + iVM, err := ihost.CreateVM(desc.Name, desc.ExternalImageId, desc.SysDiskSize, desc.Cpu, desc.Memory, desc.ExternalNetworkId, + desc.IpAddr, desc.Description, passwd, desc.StorageType, desc.DataDisks, desc.PublicKey) if err != nil { return nil, err } - } - - encpasswd, err := utils.EncryptAESBase64(guest.Id, passwd) - if err != nil { - log.Errorf("encrypt password failed %s", err) - } - - data := jsonutils.NewDict() - data.Add(jsonutils.NewString(iVM.GetOSType()), "os") - data.Add(jsonutils.NewString("root"), "account") - data.Add(jsonutils.NewString(encpasswd), "key") - - 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() - diskInfo[i] = dinfo + 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 } - data.Add(jsonutils.Marshal(&diskInfo), "disks") + 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 + } + + if len(guest.SecgrpId) > 0 { + if err := iVM.SyncSecurityGroup(guest.SecgrpId, guest.GetSecgroupName(), guest.GetSecRules()); err != nil { + log.Errorf("SyncSecurityGroup error: %v", err) + return nil, err + } + } + + /*if onfinish == "none" { + err = iVM.StartVM() + if err != nil { + return nil, err + } + }*/ + + encpasswd, err := utils.EncryptAESBase64(guest.Id, passwd) + if err != nil { + log.Errorf("encrypt password failed %s", err) + } + + data := jsonutils.NewDict() + data.Add(jsonutils.NewString(iVM.GetOSType()), "os") + data.Add(jsonutils.NewString("root"), "account") + data.Add(jsonutils.NewString(encpasswd), "key") + + 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() + 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, nil + }) + } else if action == "deploy" { + iVM, err := ihost.GetIVMById(guest.GetExternalId()) + if err != nil || iVM == nil { + log.Errorf("cannot find vm %s", err) + return fmt.Errorf("cannot find vm") } - data.Add(jsonutils.NewString(iVM.GetGlobalId()), "uuid") + params := task.GetParams() + log.Debugf("Deploy VM params %s", params.String()) + var name string + if v, e := params.GetString("name"); e != nil { + name = v + } + var description string + if v, e := params.GetString("description"); e != nil { + description = v + } + resetPassword := jsonutils.QueryBoolean(params, "reset_password", false) + deleteKeypair := jsonutils.QueryBoolean(params, "__delete_keypair__", false) + password, _ := params.GetString("password") + if resetPassword && len(password) == 0 { + password = seclib2.RandomPassword2(12) + } + + publicKey := "" + if k, e := config.GetString("public_key"); e != nil { + publicKey = k + } + + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + encpasswd, err := utils.EncryptAESBase64(guest.Id, password) + if err != nil { + log.Errorf("encrypt password failed %s", err) + } + + data := jsonutils.NewDict() + data.Add(jsonutils.NewString("root"), "account") // 用户名 + data.Add(jsonutils.NewString(encpasswd), "key") // 密码 + e := iVM.DeployVM(name, password, publicKey, resetPassword, deleteKeypair, description) + return data, e + }) + } else { + return fmt.Errorf("Action %s not supported", action) + } - return data, nil - }) return nil } @@ -243,6 +293,13 @@ func (self *SAliyunGuestDriver) OnGuestDeployTaskDataReceived(ctx context.Contex disk.DiskSize = diskInfo[i].Size disk.ExternalId = diskInfo[i].Uuid disk.Status = models.DISK_READY + if len(diskInfo[i].Metadata) > 0 { + for key, value := range diskInfo[i].Metadata { + 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 }) if err != nil { @@ -258,6 +315,20 @@ func (self *SAliyunGuestDriver) OnGuestDeployTaskDataReceived(ctx context.Contex if len(uuid) > 0 { guest.SetExternalId(uuid) } + + if metaData, _ := data.Get("metadata"); metaData != nil { + meta := make(map[string]string, 0) + if err := metaData.Unmarshal(meta); err != nil { + log.Errorf("Get guest %s metadata error: %v", guest.Name, err) + } else { + for key, value := range meta { + if err := guest.SetMetadata(ctx, key, value, task.GetUserCred()); err != nil { + log.Errorf("set guest %s mata %s => %s error: %v", guest.Name, key, value, err) + } + } + } + } + guest.SaveDeployInfo(ctx, task.GetUserCred(), data) return nil } @@ -277,3 +348,122 @@ func (self *SAliyunGuestDriver) RequestSyncConfigOnHost(ctx context.Context, gue }) return nil } + +type SAliyunVMChangeConfig struct { + InstanceId string + Cpu int + Memory int +} + +func (self *SAliyunGuestDriver) DoGuestCreateDisksTask(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { + subtask, err := taskman.TaskManager.NewTask(ctx, "AliyunGuestCreateDiskTask", guest, task.GetUserCred(), task.GetParams(), task.GetTaskId(), "", nil) + if err != nil { + return err + } + subtask.ScheduleRun(nil) + return nil +} + +func (self *SAliyunGuestDriver) AllowReconfigGuest() bool { + return true +} + +func (self *SAliyunGuestDriver) RequestChangeVmConfig(ctx context.Context, guest *models.SGuest, task taskman.ITask, vcpuCount, vmemSize int64) error { + config := SAliyunVMChangeConfig{} + config.InstanceId = guest.GetExternalId() + config.Cpu = int(vcpuCount) + config.Memory = int(vmemSize) + ihost, err := guest.GetHost().GetIHost() + if err != nil { + return err + } + + iVM, err := ihost.GetIVMById(config.InstanceId) + if err != nil { + return err + } + + if int(guest.VcpuCount) != config.Cpu || guest.VmemSize != config.Memory { + err = iVM.ChangeConfig(config.InstanceId, config.Cpu, config.Memory) + if err != nil { + return err + } + } + + log.Debugf("VMchangeConfig %s, wait status ready ...", iVM.GetGlobalId()) + err = cloudprovider.WaitStatus(iVM, models.VM_READY, time.Second*5, time.Second*300) + if err != nil { + return err + } + log.Debugf("VMchangeConfig %s, and status is ready", iVM.GetGlobalId()) + return nil +} + +func (self *SAliyunGuestDriver) RequestStartOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) (jsonutils.JSONObject, error) { + ihost, e := host.GetIHost() + if e != nil { + return nil, e + } + + ivm, e := ihost.GetIVMById(guest.GetExternalId()) + if e != nil { + return nil, e + } + + result := jsonutils.NewDict() + if ivm.GetStatus() != models.VM_RUNNING { + if err := ivm.StartVM(); err != nil { + return nil, e + } else { + task.ScheduleRun(result) + } + } else { + result.Add(jsonutils.NewBool(true), "is_running") + } + + return result, e +} + +func (self *SAliyunGuestDriver) RequestRebuildRootDisk(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { + ihost, e := guest.GetHost().GetIHost() + if e != nil { + return e + } + + externalId := guest.GetExternalId() + if len(externalId) <= 0 { + return fmt.Errorf("external id not found") + } + + disks := guest.GetDisks() + if len(disks) <= 0 { + return fmt.Errorf("guest has no disk") + } + + imageId := guest.CategorizeDisks().Root.TemplateId + cacheId := disks[0].GetDisk().GetStorage().GetStoragecache().Id + externalImageId := models.StoragecachedimageManager.GetStoragecachedimage(cacheId, imageId).ExternalId + if len(externalImageId) <= 0 { + return fmt.Errorf("external image (%s) id is not found", imageId) + } + + iVM, err := ihost.GetIVMById(externalId) + if err != nil { + return err + } + + err = iVM.RebuildRoot(externalImageId) + if err != nil { + return err + } + + log.Debugf("VMrebuildRoot %s, wait status ready ...", iVM.GetGlobalId()) + err = cloudprovider.WaitStatus(iVM, models.VM_READY, time.Second*5, time.Second*1800) + if err != nil { + return err + } + log.Debugf("VMrebuildRoot %s, and status is ready", iVM.GetGlobalId()) + + task.ScheduleRun(nil) + return nil +} diff --git a/pkg/compute/hostdrivers/aliyun.go b/pkg/compute/hostdrivers/aliyun.go index 0bcfa18005..4c1d7bf029 100644 --- a/pkg/compute/hostdrivers/aliyun.go +++ b/pkg/compute/hostdrivers/aliyun.go @@ -48,7 +48,7 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host * return nil } -func (self *SAliyunHostDriver) RequestAllocateDiskOnStorage(host *models.SHost, storage *models.SStorage, disk *models.SDisk, task taskman.ITask, content *jsonutils.JSONDict) error { +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 } else { @@ -61,6 +61,20 @@ func (self *SAliyunHostDriver) RequestAllocateDiskOnStorage(host *models.SHost, } 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) @@ -81,9 +95,13 @@ func (self *SAliyunHostDriver) RequestDeallocateDiskOnHost(host *models.SHost, s return err } else if iDisk, err := iCloudStorage.GetIDisk(disk.GetExternalId()); err != nil { return err + } else if err := iDisk.Delete(); err != nil { + return err } else { - return iDisk.Delete() + data := jsonutils.NewDict() + task.ScheduleRun(data) } + return nil } func (self *SAliyunHostDriver) RequestResizeDiskOnHostOnline(host *models.SHost, storage *models.SStorage, disk *models.SDisk, size int64, task taskman.ITask) error { diff --git a/pkg/compute/hostdrivers/kvm.go b/pkg/compute/hostdrivers/kvm.go index 0a5cbca5b2..ed08542a36 100644 --- a/pkg/compute/hostdrivers/kvm.go +++ b/pkg/compute/hostdrivers/kvm.go @@ -74,7 +74,7 @@ func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *mod return nil } -func (self *SKVMHostDriver) RequestAllocateDiskOnStorage(host *models.SHost, storage *models.SStorage, disk *models.SDisk, task taskman.ITask, content *jsonutils.JSONDict) error { +func (self *SKVMHostDriver) RequestAllocateDiskOnStorage(ctx context.Context, host *models.SHost, storage *models.SStorage, disk *models.SDisk, task taskman.ITask, content *jsonutils.JSONDict) error { header := http.Header{} header.Add("X-Task-Id", task.GetTaskId()) header.Add("X-Region-Version", "v2") diff --git a/pkg/compute/models/baremetalstatus.go b/pkg/compute/models/baremetalstatus.go new file mode 100644 index 0000000000..5c4561f14c --- /dev/null +++ b/pkg/compute/models/baremetalstatus.go @@ -0,0 +1,22 @@ +package models + +const ( + BAREMETAL_INIT = "init" + BAREMETAL_PREPARE = "prepare" + BAREMETAL_PREPARE_FAIL = "prepare_fail" + BAREMETAL_READY = "ready" + BAREMETAL_RUNNING = "running" + BAREMETAL_MAINTAINING = "maintaining" + BAREMETAL_START_MAINTAIN = "start_maintain" + BAREMETAL_DELETING = "deleting" + BAREMETAL_DELETE = "delete" + BAREMETAL_DELETE_FAIL = "delete_fail" + BAREMETAL_UNKNOWN = "unknown" + BAREMETAL_SYNCING_STATUS = "syncing_status" + BAREMETAL_SYNC = "sync" + BAREMETAL_SYNC_FAIL = "sync_fail" + BAREMETAL_START_CONVERT = "start_convert" + BAREMETAL_CONVERTING = "converting" + BAREMETAL_START_FAIL = "start_fail" + BAREMETAL_STOP_FAIL = "stop_fail" +) diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index a2e827b7b9..b05fa85b5e 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -297,7 +297,7 @@ func (self *SDisk) StartDiskCreateTask(ctx context.Context, userCred mcclient.To return nil } -func (self *SDisk) StartAllocate(host *SHost, storage *SStorage, taskId string, userCred mcclient.TokenCredential, rebuild bool, snapshot string, task taskman.ITask) error { +func (self *SDisk) StartAllocate(ctx context.Context, host *SHost, storage *SStorage, taskId string, userCred mcclient.TokenCredential, rebuild bool, snapshot string, task taskman.ITask) error { log.Infof("Allocating disk on host %s ...", host.GetName()) templateId := self.GetTemplateId() @@ -326,7 +326,7 @@ func (self *SDisk) StartAllocate(host *SHost, storage *SStorage, taskId string, if rebuild { content.Add(jsonutils.JSONTrue, "rebuild") } - return host.GetHostDriver().RequestAllocateDiskOnStorage(host, storage, self, task, content) + return host.GetHostDriver().RequestAllocateDiskOnStorage(ctx, host, storage, self, task, content) } func (self *SDisk) AllowPerformResize(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { @@ -433,7 +433,7 @@ func (manager *SDiskManager) getDisksByStorage(storage *SStorage) ([]SDisk, erro return disks, nil } -func (manager *SDiskManager) syncCloudDisk(userCred mcclient.TokenCredential, vdisk cloudprovider.ICloudDisk) (*SDisk, error) { +func (manager *SDiskManager) syncCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, vdisk cloudprovider.ICloudDisk) (*SDisk, error) { diskObj, err := manager.FetchByExternalId(vdisk.GetGlobalId()) if err != nil { if err == sql.ErrNoRows { @@ -444,13 +444,13 @@ func (manager *SDiskManager) syncCloudDisk(userCred mcclient.TokenCredential, vd return nil, err } storage := storageObj.(*SStorage) - return manager.newFromCloudDisk(userCred, vdisk, storage) + return manager.newFromCloudDisk(ctx, userCred, vdisk, storage) } else { return nil, err } } else { disk := diskObj.(*SDisk) - err = disk.syncWithCloudDisk(userCred, vdisk) + err = disk.syncWithCloudDisk(ctx, userCred, vdisk) if err != nil { return nil, err } @@ -490,7 +490,7 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To } for i := 0; i < len(commondb); i += 1 { - err = commondb[i].syncWithCloudDisk(userCred, commonext[i]) + err = commondb[i].syncWithCloudDisk(ctx, userCred, commonext[i]) if err != nil { syncResult.UpdateError(err) } else { @@ -501,7 +501,7 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To } for i := 0; i < len(added); i += 1 { - new, err := manager.newFromCloudDisk(userCred, added[i], storage) + new, err := manager.newFromCloudDisk(ctx, userCred, added[i], storage) if err != nil { syncResult.AddError(err) } else { @@ -514,7 +514,7 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To return localDisks, remoteDisks, syncResult } -func (self *SDisk) syncWithCloudDisk(userCred mcclient.TokenCredential, extDisk cloudprovider.ICloudDisk) error { +func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, extDisk cloudprovider.ICloudDisk) error { _, err := self.GetModelManager().TableSpec().Update(self, func() error { extDisk.Refresh() self.Name = extDisk.GetName() @@ -535,11 +535,26 @@ func (self *SDisk) syncWithCloudDisk(userCred mcclient.TokenCredential, extDisk }) if err != nil { log.Errorf("syncWithCloudDisk error %s", err) + return err } - return err + + if metaData := extDisk.GetMetadata(); metaData != nil { + meta := make(map[string]string, 0) + if err := metaData.Unmarshal(meta); err != nil { + log.Errorf("Get VM Metadata error: %v", err) + } else { + for key, value := range meta { + if err := self.SetMetadata(ctx, key, value, userCred); err != nil { + log.Errorf("set disk %s mata %s => %s error: %v", self.Name, key, value, err) + } + } + } + } + + return nil } -func (manager *SDiskManager) newFromCloudDisk(userCred mcclient.TokenCredential, extDisk cloudprovider.ICloudDisk, storage *SStorage) (*SDisk, error) { +func (manager *SDiskManager) newFromCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, extDisk cloudprovider.ICloudDisk, storage *SStorage) (*SDisk, error) { disk := SDisk{} disk.SetModelManager(manager) @@ -562,6 +577,20 @@ func (manager *SDiskManager) newFromCloudDisk(userCred mcclient.TokenCredential, log.Errorf("newFromCloudZone fail %s", err) return nil, err } + + if metaData := extDisk.GetMetadata(); metaData != nil { + meta := make(map[string]string) + if err := metaData.Unmarshal(meta); err != nil { + log.Errorf("Get VM Metadata error: %v", err) + } else { + for key, value := range meta { + if err := disk.SetMetadata(ctx, key, value, userCred); err != nil { + log.Errorf("set disk %s mata %s => %s error: %v", disk.Name, key, value, err) + } + } + } + } + return &disk, nil } @@ -701,6 +730,27 @@ func parseIsoInfo(ctx context.Context, userCred mcclient.TokenCredential, info s return image.Id, nil } +// def get_disk_spec_v2(conf): +// def _get_spec(storages): +// spec = {} +// for adapter, ss in group_by_adapter(storages).items(): +// if len(ss) == 0: +// continue +// spec[adapter] = get_disk_spec(ss) +// return spec + +// spec = {} +// for driver in DISK_DRIVERS: +// storages = [s for s in conf if s['driver'] == driver] +// if len(storages) != 0: +// spec[driver] = _get_spec(storages) +// return spec + +func GetDiskSpecV2(storageInfo jsonutils.JSONObject) *jsonutils.JSONDict { + // ToDo + return nil +} + func (self *SDisk) fetchDiskInfo(diskConfig *SDiskConfig) { if len(diskConfig.ImageId) > 0 { self.TemplateId = diskConfig.ImageId @@ -832,9 +882,8 @@ func (self *SDisk) GetAttachedGuests() []SGuest { q = q.Filter(sqlchemy.Equals(guestdisks.Field("disk_id"), self.Id)) ret := make([]SGuest, 0) - err := q.All(&ret) - if err != nil { - log.Errorf("%s", err) + if err := db.FetchModelObjects(GuestManager, q, &ret); err != nil { + log.Errorf("Fetch Geusts Objects %v", err) return nil } return ret @@ -871,6 +920,11 @@ func (self *SDisk) GetShortDesc() *jsonutils.JSONDict { storage := self.GetStorage() desc.Add(jsonutils.NewString(storage.StorageType), "storage_type") desc.Add(jsonutils.NewString(storage.MediumType), "medium_type") + + if priceKey := self.GetMetadata("price_key", nil); len(priceKey) > 0 { + desc.Add(jsonutils.NewString(priceKey), "price_key") + } + fs := self.GetFsFormat() if len(fs) > 0 { desc.Add(jsonutils.NewString(fs), "fs_format") diff --git a/pkg/compute/models/dnsrecords.go b/pkg/compute/models/dnsrecords.go index a1177a3401..5da570532c 100644 --- a/pkg/compute/models/dnsrecords.go +++ b/pkg/compute/models/dnsrecords.go @@ -6,7 +6,6 @@ import ( "strings" "yunion.io/x/jsonutils" - "yunion.io/x/log" "yunion.io/x/pkg/util/regutils" "yunion.io/x/sqlchemy" @@ -205,7 +204,6 @@ func (man *SDnsRecordManager) QueryDns(projectId, name string) *SDnsRecord { err := q.First(rec) if err != nil { - log.Errorf("QueryDns %q fail: %s", name, err) return nil } return rec diff --git a/pkg/compute/models/guestnetworks.go b/pkg/compute/models/guestnetworks.go index 5ff6910060..882bbb5851 100644 --- a/pkg/compute/models/guestnetworks.go +++ b/pkg/compute/models/guestnetworks.go @@ -19,6 +19,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/compute/options" + "database/sql" ) const ( @@ -263,17 +264,18 @@ func (self *SGuestnetwork) GetJsonDescAtHost(host *SHost) jsonutils.JSONObject { } func (manager *SGuestnetworkManager) GetGuestByAddress(address string) *SGuest { - gn := SGuestnetwork{} - gn.SetModelManager(GuestnetworkManager) - - err := manager.Query().Equals("ip_addr", address).First(&gn) - if err != nil { - log.Errorf("GetGuestByAddress query fail: %s", err) - return nil - } - guest, _ := GuestManager.FetchById(gn.GuestId) - if guest != nil { - return guest.(*SGuest) + networks := manager.TableSpec().Instance() + guests := GuestManager.Query() + q := guests.Join(networks, sqlchemy.AND( + sqlchemy.IsFalse(networks.Field("deleted")), + sqlchemy.Equals(networks.Field("ip_addr"), address), + sqlchemy.Equals(networks.Field("guest_id"), guests.Field("id")), + )) + guest := &SGuest{} + guest.SetModelManager(GuestManager) + err := q.First(guest) + if err == nil { + return guest } return nil } @@ -333,6 +335,22 @@ func (manager *SGuestnetworkManager) DeleteGuestNics(ctx context.Context, guest return nil } +func (manager *SGuestnetworkManager) getGuestNicByIP(ip string) (*SGuestnetwork, error) { + gn := SGuestnetwork{} + q := manager.Query() + q = q.Equals("ip_addr", ip) + err := q.First(&gn) + if err != nil { + if err != sql.ErrNoRows { + log.Errorf("getGuestNicByIP fail %s", err) + return nil, err + } + return nil, nil + } + gn.SetModelManager(manager) + return &gn, nil +} + func (self *SGuestnetwork) LogDetachEvent(userCred mcclient.TokenCredential, guest *SGuest, network *SNetwork) { if network == nil { netTmp, _ := NetworkManager.FetchById(self.NetworkId) diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 8ecba69d47..afd50582dc 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -493,9 +493,11 @@ func (self *SGuest) ValidateUpdateData(ctx context.Context, userCred mcclient.To return nil, err } - // if data.Contains("name") { - // return nil, httperrors.NewInputParameterError("cannot update server name") - // } + if data.Contains("name") { + if name, _ := data.GetString("name"); len(name) < 2 { + return nil, httperrors.NewInputParameterError("name is to short") + } + } /* if self.GetHypervisor() == HYPERVISOR_BAREMETAL { return nil, httperrors.NewInputParameterError("Cannot modify memory for baremetal") } @@ -1247,6 +1249,18 @@ func (self *SGuest) syncWithCloudVM(ctx context.Context, userCred mcclient.Token db.OpsLog.LogEvent(self, db.ACT_UPDATE, diffStr, userCred) } } + if metaData := extVM.GetMetadata(); metaData != nil { + meta := make(map[string]string, 0) + if metaData.Unmarshal(meta); err != nil { + log.Errorf("Get VM Metadata error: %v", err) + } else { + for key, value := range meta { + if err := self.SetMetadata(ctx, key, value, userCred); err != nil { + log.Errorf("set guest %s mata %s => %s error: %v", self.Name, key, value, err) + } + } + } + } return nil } @@ -1277,6 +1291,20 @@ func (manager *SGuestManager) newCloudVM(ctx context.Context, userCred mcclient. if err != nil { log.Errorf("Insert fail %s", err) } + + if metaData := extVM.GetMetadata(); metaData != nil { + meta := make(map[string]string, 0) + if err := metaData.Unmarshal(meta); err != nil { + log.Errorf("Get VM Metadata error: %v", err) + } else { + for key, value := range meta { + if err := guest.SetMetadata(ctx, key, value, userCred); err != nil { + log.Errorf("set guest %s mata %s => %s error: %v", guest.Name, key, value, err) + } + } + } + } + return &guest, nil } @@ -1475,7 +1503,20 @@ func (self *SGuest) SyncVMNics(ctx context.Context, userCred mcclient.TokenCrede if add.net == nil { continue // cannot determine which network it attached to } - err := self.Attach2Network(ctx, userCred, add.net, nil, add.nic.GetIP(), + // check if the IP has been occupied, if yes, release the IP + gn, err := GuestnetworkManager.getGuestNicByIP(add.nic.GetIP()) + if err != nil { + result.AddError(err) + continue + } + if gn != nil { + err = gn.Detach(ctx, userCred) + if err != nil { + result.AddError(err) + continue + } + } + err = self.Attach2Network(ctx, userCred, add.net, nil, add.nic.GetIP(), add.nic.GetMAC(), add.nic.GetDriver(), 0, false, -1, add.reserve, IPAllocationDefault, true) if err != nil { result.AddError(err) @@ -1539,7 +1580,11 @@ func (self *SGuest) PerformDeploy(ctx context.Context, userCred mcclient.TokenCr if !ok { return nil, fmt.Errorf("Parse query body error") } + + // 变更密码/密钥时需要Restart才能生效。更新普通字段不需要Restart + doRestart := false if kwargs.Contains("__delete_keypair__") || kwargs.Contains("keypair") { + doRestart = true var kpId string if !jsonutils.QueryBoolean(kwargs, "__delete_keypair__", false) { keypair, _ := kwargs.GetString("keypair") @@ -1561,8 +1606,9 @@ func (self *SGuest) PerformDeploy(ctx context.Context, userCred mcclient.TokenCr kwargs.Set("reset_password", jsonutils.JSONTrue) } } + if utils.IsInStringArray(self.Status, []string{VM_RUNNING, VM_READY, VM_ADMIN}) { - if self.Status == VM_RUNNING { + if doRestart && self.Status == VM_RUNNING { kwargs.Set("restart", jsonutils.JSONTrue) } err := self.StartGuestDeployTask(ctx, userCred, kwargs, "deploy", "") @@ -1648,7 +1694,7 @@ func (self *SGuest) SyncVMDisks(ctx context.Context, userCred mcclient.TokenCred if len(vdisks[i].GetGlobalId()) == 0 { continue } - disk, err := DiskManager.syncCloudDisk(userCred, vdisks[i]) + disk, err := DiskManager.syncCloudDisk(ctx, userCred, vdisks[i]) if err != nil { result.Error(err) return result @@ -1906,12 +1952,10 @@ func (self *SGuest) createDiskOnHost(ctx context.Context, userCred mcclient.Toke if storage == nil { return nil, fmt.Errorf("No storage to create disk") } - disk, err := self.createDiskOnStorage(ctx, userCred, storage, diskConfig, pendingUsage) if err != nil { return nil, err } - err = self.attach2Disk(disk, userCred, diskConfig.Driver, diskConfig.Cache, diskConfig.Mountpoint) return disk, err } @@ -2325,6 +2369,9 @@ func (self *SGuest) PerformDetachdisk(ctx context.Context, userCred mcclient.Tok disk := iDisk.(*SDisk) if disk != nil { if self.isAttach2Disk(disk) { + if disk.DiskType == DISK_TYPE_SYS { + return nil, httperrors.NewUnsupportOperationError("Cannot detach sys disk") + } detachDiskStatus, err := self.GetDriver().GetDetachDiskStatus() if err != nil { return nil, err @@ -2337,7 +2384,7 @@ func (self *SGuest) PerformDetachdisk(ctx context.Context, userCred mcclient.Tok disk.SetStatus(userCred, DISK_DETACHING, "") } taskData := jsonutils.NewDict() - taskData.Add(jsonutils.NewString(diskId), "disk_id") + taskData.Add(jsonutils.NewString(disk.Id), "disk_id") taskData.Add(jsonutils.NewBool(keepDisk), "keep_disk") self.GetDriver().StartGuestDetachdiskTask(ctx, userCred, self, taskData, "") return nil, nil @@ -2444,16 +2491,27 @@ func (self *SGuest) PerformChangeConfig(ctx context.Context, userCred mcclient.T diskIdx += 1 } - for storageId, needSize := range diskSizes { - iStorage, err := StorageManager.FetchById(storageId) - if err != nil { - return nil, httperrors.NewBadRequestError("Fetch storage error: %s", err) - } - storage := iStorage.(*SStorage) - if storage.GetFreeCapacity() < needSize { - return nil, httperrors.NewInsufficientResourceError("Not enough free space") - } + provider, e := self.GetHost().GetDriver() + if e != nil { + log.Errorf("Get Provider Error: %s", e) + return nil, httperrors.NewInsufficientResourceError("Provider Not Found") } + + if !provider.IsPublicCloud() { + for storageId, needSize := range diskSizes { + iStorage, err := StorageManager.FetchById(storageId) + if err != nil { + return nil, httperrors.NewBadRequestError("Fetch storage error: %s", err) + } + storage := iStorage.(*SStorage) + if storage.GetFreeCapacity() < needSize { + return nil, httperrors.NewInsufficientResourceError("Not enough free space") + } + } + } else { + log.Debugf("Skip storage free capacity validating for public cloud: %s", provider.GetName()) + } + if newDisks.Length() > 0 { confs.Add(newDisks, "create") } @@ -3027,6 +3085,111 @@ func (manager *SGuestManager) FetchGuestById(guestId string) *SGuest { return guest.(*SGuest) } +func (self *SGuest) GetSpec(checkStatus bool) *jsonutils.JSONDict { + if checkStatus { + if !utils.IsInStringArray(self.Status, []string{VM_SCHEDULE_FAILED}) { + return nil + } + } + spec := jsonutils.NewDict() + spec.Set("cpu", jsonutils.NewInt(int64(self.VcpuCount))) + spec.Set("mem", jsonutils.NewInt(int64(self.VmemSize))) + + // get disk spec + guestdisks := self.GetDisks() + diskSpecs := jsonutils.NewArray() + for _, guestdisk := range guestdisks { + disk := guestdisk.GetDisk() + diskSpec := jsonutils.NewDict() + diskSpec.Set("size", jsonutils.NewInt(int64(disk.DiskSize))) + s := disk.GetStorage() + diskSpec.Set("backend", jsonutils.NewString(s.StorageType)) + diskSpec.Set("medium_type", jsonutils.NewString(s.MediumType)) + diskSpecs.Add(diskSpec) + } + spec.Set("disk", diskSpecs) + + // get nic spec + guestnics := self.GetNetworks() + nicSpecs := jsonutils.NewArray() + for _, guestnic := range guestnics { + nicSpec := jsonutils.NewDict() + nicSpec.Set("bandwidth", jsonutils.NewInt(int64(guestnic.getBandwidth()))) + t := "int" + if guestnic.IsExit() { + t = "ext" + } + nicSpec.Set("type", jsonutils.NewString(t)) + nicSpecs.Add(nicSpec) + } + spec.Set("nic", nicSpecs) + + // get isolate device spec + guestgpus := self.GetIsolatedDevices() + gpuSpecs := jsonutils.NewArray() + for _, guestgpu := range guestgpus { + if strings.HasPrefix(guestgpu.DevType, "GPU") { + gs := guestgpu.GetSpec(false) + if gs != nil { + gpuSpecs.Add(gs) + } + } + } + spec.Set("gpu", gpuSpecs) + return spec +} + +func (self *SGuest) GetTemplateId() string { + guestdisks := self.GetDisks() + for _, guestdisk := range guestdisks { + disk := guestdisk.GetDisk() + if disk != nil { + templateId := disk.GetTemplateId() + if len(templateId) > 0 { + return templateId + } + } + } + return "" +} + +func (self *SGuest) GetShortDesc() *jsonutils.JSONDict { + desc := self.SStandaloneResourceBase.GetShortDesc() + desc.Set("mem", jsonutils.NewInt(int64(self.VmemSize))) + desc.Set("cpu", jsonutils.NewInt(int64(self.VcpuCount))) + templateId := self.GetTemplateId() + if len(templateId) > 0 { + desc.Set("cpu", jsonutils.NewString(templateId)) + } + extBw := self.getBandwidth(true) + intBw := self.getBandwidth(false) + if extBw > 0 { + desc.Set("ext_bandwidth", jsonutils.NewInt(int64(extBw))) + } + if intBw > 0 { + desc.Set("int_bandwidth", jsonutils.NewInt(int64(intBw))) + } + + if priceKey := self.GetMetadata("price_key", nil); len(priceKey) > 0 { + desc.Add(jsonutils.NewString(priceKey), "price_key") + } + + desc.Set("hypervisor", jsonutils.NewString(self.GetHypervisor())) + spec := self.GetSpec(false) + if self.GetHypervisor() == HYPERVISOR_BAREMETAL { + host := self.GetHost() + if host != nil { + hostSpec := host.GetSpec(false) + hostSpecIdent := host.GetSpecIdent(hostSpec) + spec.Set("host_spec", jsonutils.NewString(strings.Join(hostSpecIdent, "/"))) + } + } + if spec != nil { + desc.Update(spec) + } + return desc +} + func (self *SGuest) saveOsType(osType string) error { _, err := self.GetModelManager().TableSpec().Update(self, func() error { self.OsType = osType diff --git a/pkg/compute/models/hostdrivers.go b/pkg/compute/models/hostdrivers.go index 3f162d75da..6749e5ac68 100644 --- a/pkg/compute/models/hostdrivers.go +++ b/pkg/compute/models/hostdrivers.go @@ -12,7 +12,7 @@ import ( type IHostDriver interface { GetHostType() string CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, scimg *SStoragecachedimage, task taskman.ITask) error - RequestAllocateDiskOnStorage(host *SHost, storage *SStorage, disk *SDisk, task taskman.ITask, content *jsonutils.JSONDict) 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 RequestResizeDiskOnHost(host *SHost, storage *SStorage, disk *SDisk, size int64, task taskman.ITask) error diff --git a/pkg/compute/models/hostnetworks.go b/pkg/compute/models/hostnetworks.go index abf31b3ed6..5a94918e32 100644 --- a/pkg/compute/models/hostnetworks.go +++ b/pkg/compute/models/hostnetworks.go @@ -92,7 +92,7 @@ func (man *SHostnetworkManager) QueryByAddress(addr string) *sqlchemy.SQuery { func (man *SHostnetworkManager) GetHostNetworkByAddress(addr string) *SHostnetwork { network := SHostnetwork{} err := man.QueryByAddress(addr).First(&network) - if err != nil { + if err == nil { return &network } return nil @@ -107,9 +107,18 @@ func (man *SHostnetworkManager) GetNetworkByAddress(addr string) *SNetwork { } func (man *SHostnetworkManager) GetHostByAddress(addr string) *SHost { - net := man.GetHostNetworkByAddress(addr) - if net == nil { - return nil + networks := man.TableSpec().Instance() + hosts := HostManager.Query() + q := hosts.Join(networks, sqlchemy.AND( + sqlchemy.IsFalse(networks.Field("deleted")), + sqlchemy.Equals(networks.Field("ip_addr"), addr), + sqlchemy.Equals(networks.Field("baremetal_id"), hosts.Field("id")), + )) + host := &SHost{} + host.SetModelManager(HostManager) + err := q.First(host) + if err == nil { + return host } - return net.GetHost() + return nil } diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 6b47c4e370..0874133ba1 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -506,6 +506,58 @@ func (self *SHost) ClearSchedDescCache() error { return HostManager.ClearSchedDescCache(self.Id) } +func (self *SHost) GetSpec(statusCheck bool) *jsonutils.JSONDict { + if statusCheck { + if utils.IsInStringArray(self.Status, []string{BAREMETAL_INIT, BAREMETAL_PREPARE_FAIL, BAREMETAL_PREPARE}) || + self.getBaremetalServer() != nil { + return nil + } + if self.MemSize == 0 || self.CpuCount == 0 { + return nil + } + } + spec := self.GetHardwareSpecification() + spec.Remove("storage_info") + netInfo := jsonutils.NewArray() + nifs := self.GetNetInterfaces() + for _, nif := range nifs { + netDesc := nif.getBaremetalJsonDesc() + nicType, err := netDesc.GetString("nic_type") + if err != nil && nicType != NIC_TYPE_IPMI { + netInfo.Add(netDesc) + } + } + spec.Set("nic_count", jsonutils.NewInt(int64(netInfo.Length()))) + manufacture, err := self.SysInfo.Get("manufacture") + if err != nil { + manufacture = jsonutils.NewString("Unknown") + } + spec.Set("manufacture", manufacture) + model, err := self.SysInfo.Get("model") + if err != nil { + model = jsonutils.NewString("Unknown") + } + spec.Set("model", model) + return nil +} + +func (self *SHost) GetSpecIdent(spec *jsonutils.JSONDict) []string { + // Todo + return []string{} +} + +func (self *SHost) GetHardwareSpecification() *jsonutils.JSONDict { + spec := jsonutils.NewDict() + spec.Set("cpu", jsonutils.NewInt(int64(self.CpuCount))) + spec.Set("mem", jsonutils.NewInt(int64(self.MemSize))) + if self.StorageInfo != nil { + spec.Set("disk", GetDiskSpecV2(self.StorageInfo)) + spec.Set("driver", jsonutils.NewString(self.StorageDriver)) + spec.Set("storage_info", self.StorageInfo) + } + return spec +} + type SStorageCapacity struct { Capacity int Used int diff --git a/pkg/compute/models/isolated_devices.go b/pkg/compute/models/isolated_devices.go index f3182d3f67..d1c7c26c21 100644 --- a/pkg/compute/models/isolated_devices.go +++ b/pkg/compute/models/isolated_devices.go @@ -380,6 +380,24 @@ func (self *SIsolatedDevice) getDesc() *jsonutils.JSONDict { return desc } +func (self *SIsolatedDevice) GetSpec(checkStatus bool) *jsonutils.JSONDict { + if checkStatus { + if len(self.GuestId) > 0 { + return nil + } + host := self.getHost() + if host.Status != BAREMETAL_RUNNING || !host.Enabled { + return nil + } + } + spec := jsonutils.NewDict() + spec.Set("dev_type", jsonutils.NewString(self.DevType)) + spec.Set("model", jsonutils.NewString(self.Model)) + spec.Set("pci_id", jsonutils.NewString(self.VendorDeviceId)) + spec.Set("vendor", jsonutils.NewString(self.getVendor())) + return spec +} + func (self *SIsolatedDevice) GetShortDesc() *jsonutils.JSONDict { desc := self.getDesc() desc.Add(jsonutils.NewString(self.Keyword()), "res_name") diff --git a/pkg/compute/models/secgrouprules.go b/pkg/compute/models/secgrouprules.go index 4347ea48c9..46b032e195 100644 --- a/pkg/compute/models/secgrouprules.go +++ b/pkg/compute/models/secgrouprules.go @@ -328,7 +328,7 @@ func (manager *SSecurityGroupRuleManager) SyncRules(ctx context.Context, userCre cmp := strings.Compare(dbStr, ruleStr) if cmp == 0 { if dbRules[j].Description != rules[i].Description { - if _, err := manager.TableSpec().Update(dbRules[j], func() error { + if _, err := manager.TableSpec().Update(&dbRules[j], func() error { dbRules[j].Description = rules[i].Description return nil }); err != nil { diff --git a/pkg/compute/models/storages.go b/pkg/compute/models/storages.go index 4ac8ea7587..0093f4de8b 100644 --- a/pkg/compute/models/storages.go +++ b/pkg/compute/models/storages.go @@ -12,6 +12,7 @@ import ( "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/compare" + "yunion.io/x/pkg/util/sysutils" "yunion.io/x/sqlchemy" ) @@ -100,6 +101,10 @@ func (self *SStorage) IsLocal() bool { return self.StorageType == STORAGE_LOCAL || self.StorageType == STORAGE_BAREMETAL } +func (manager *SStorageManager) AllowListItems(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool { + return true +} + func (self *SStorage) GetUsedCapacity(isReady tristate.TriState) int { disks := DiskManager.Query().SubQuery() q := disks.Query(sqlchemy.SUM("sum", disks.Field("disk_size"))).Equals("storage_id", self.Id) @@ -654,6 +659,26 @@ func (manager *SStorageManager) ListItemFilter(ctx context.Context, q *sqlchemy. q = q.Filter(sqlchemy.In(q.Field("zone_id"), sq.SubQuery())) } + if jsonutils.QueryBoolean(query, "share", false) { + q = q.Filter(sqlchemy.NotIn(q.Field("storage_type"), sysutils.LOCAL_STORAGE_TYPES)) + } + + if jsonutils.QueryBoolean(query, "local", false) { + q = q.Filter(sqlchemy.In(q.Field("storage_type"), sysutils.LOCAL_STORAGE_TYPES)) + } + + if jsonutils.QueryBoolean(query, "usable", false) { + hostStorageTable := HoststorageManager.Query().SubQuery() + hostTable := HostManager.Query().SubQuery() + sq := hostStorageTable.Query(hostStorageTable.Field("storage_id")).Join(hostTable, + sqlchemy.Equals(hostTable.Field("id"), hostStorageTable.Field("host_id"))). + Filter(sqlchemy.Equals(hostTable.Field("host_status"), HOST_ONLINE)) + + q = q.Filter(sqlchemy.In(q.Field("id"), sq)). + Filter(sqlchemy.In(q.Field("status"), []string{STORAGE_ENABLED, STORAGE_ONLINE})). + Filter(sqlchemy.IsTrue(q.Field("enabled"))) + } + managerStr := jsonutils.GetAnyString(query, []string{"manager", "provider", "manager_id", "provider_id"}) if len(managerStr) > 0 { provider := CloudproviderManager.FetchCloudproviderByIdOrName(managerStr) diff --git a/pkg/compute/models/wires.go b/pkg/compute/models/wires.go index 22fa6d9a6f..d522c53e58 100644 --- a/pkg/compute/models/wires.go +++ b/pkg/compute/models/wires.go @@ -347,6 +347,7 @@ func (self *SWire) getNetworks() ([]SNetwork, error) { func (self *SWire) getGatewayNetworkQuery() *sqlchemy.SQuery { q := self.getNetworkQuery() q = q.IsNotNull("guest_gateway").IsNotEmpty("guest_gateway") + q = q.Equals("status", NETWORK_STATUS_AVAILABLE) return q } diff --git a/pkg/compute/tasks/disk_create_task.go b/pkg/compute/tasks/disk_create_task.go index 1c47d60d5a..be0e2635ed 100644 --- a/pkg/compute/tasks/disk_create_task.go +++ b/pkg/compute/tasks/disk_create_task.go @@ -40,10 +40,10 @@ func (self *DiskCreateTask) OnStorageCacheImageComplete(ctx context.Context, dis } storage := disk.GetStorage() host := storage.GetMasterHost() - db.OpsLog.LogEvent(disk, db.ACT_ALLOCATE, disk.GetShortDesc(), self.GetUserCred()) + db.OpsLog.LogEvent(disk, db.ACT_ALLOCATING, disk.GetShortDesc(), self.GetUserCred()) disk.SetStatus(self.GetUserCred(), models.DISK_STARTALLOC, "") self.SetStage("on_disk_ready", nil) - if err := disk.StartAllocate(host, storage, self.GetTaskId(), self.GetUserCred(), rebuild, snapshot, self); err != nil { + if err := disk.StartAllocate(ctx, host, storage, self.GetTaskId(), self.GetUserCred(), rebuild, snapshot, self); err != nil { self.OnStartAllocateFailed(ctx, disk, jsonutils.NewString(err.Error())) } } diff --git a/pkg/compute/tasks/disk_resize_task.go b/pkg/compute/tasks/disk_resize_task.go index 584c43ab94..dd25da5879 100644 --- a/pkg/compute/tasks/disk_resize_task.go +++ b/pkg/compute/tasks/disk_resize_task.go @@ -32,7 +32,7 @@ func (self *DiskResizeTask) OnInit(ctx context.Context, obj db.IStandaloneModel, } resion := "Cannot find host for disk" if host == nil || host.HostStatus != models.HOST_ONLINE { - disk.SetStatus(self.GetUserCred(), models.DISK_READY, resion) + disk.SetDiskReady(ctx, self.GetUserCred(), resion) self.SetStageFailed(ctx, resion) db.OpsLog.LogEvent(disk, db.ACT_RESIZE_FAIL, resion, self.GetUserCred()) } else { @@ -65,7 +65,7 @@ func (self *DiskResizeTask) OnStartResizeDiskSucc(ctx context.Context, disk *mod } func (self *DiskResizeTask) OnStartResizeDiskFailed(ctx context.Context, disk *models.SDisk, resion error) { - disk.SetStatus(self.GetUserCred(), models.DISK_READY, resion.Error()) + disk.SetDiskReady(ctx, self.GetUserCred(), resion.Error()) self.SetStageFailed(ctx, resion.Error()) db.OpsLog.LogEvent(disk, db.ACT_RESIZE_FAIL, resion.Error(), self.GetUserCred()) } @@ -94,6 +94,7 @@ func (self *DiskResizeTask) OnDiskResizeComplete(ctx context.Context, disk *mode self.OnStartResizeDiskFailed(ctx, disk, err) return } + disk.SetDiskReady(ctx, self.GetUserCred(), "") notes := fmt.Sprintf("%s=>%s", oldStatus, disk.Status) db.OpsLog.LogEvent(disk, db.ACT_UPDATE_STATUS, notes, self.UserCred) self.CleanHostSchedCache(disk) @@ -103,6 +104,6 @@ func (self *DiskResizeTask) OnDiskResizeComplete(ctx context.Context, disk *mode } func (self *DiskResizeTask) OnDiskResizeCompleteFailed(ctx context.Context, disk *models.SDisk, resion error) { - disk.SetStatus(self.UserCred, models.DISK_READY, resion.Error()) + disk.SetDiskReady(ctx, self.GetUserCred(), resion.Error()) db.OpsLog.LogEvent(disk, db.ACT_RESIZE_FAIL, disk.GetShortDesc(), self.UserCred) } diff --git a/pkg/compute/tasks/guest_create_disk_task.go b/pkg/compute/tasks/guest_create_disk_task.go index c7034bd735..bc734944cd 100644 --- a/pkg/compute/tasks/guest_create_disk_task.go +++ b/pkg/compute/tasks/guest_create_disk_task.go @@ -9,6 +9,7 @@ import ( "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/log" ) type GuestCreateDiskTask struct { @@ -118,7 +119,115 @@ func (self *KVMGuestCreateDiskTask) OnConfigSyncComplete(ctx context.Context, ob self.SetStageComplete(ctx, nil) } +type AliyunGuestCreateDiskTask struct { + SGuestBaseTask +} + +func (self *AliyunGuestCreateDiskTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + self.SetStage("on_aliyun_disk_prepared", nil) + self.OnAliyunDiskPrepared(ctx, obj, data) +} + +func (self *AliyunGuestCreateDiskTask) OnAliyunDiskPrepared(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + var diskIndex = 0 + var diskReady = true + for { + diskId, err := self.Params.GetString(fmt.Sprintf("disk.%d.id", diskIndex)) + if !diskReady || err != nil { + break + } + iDisk, err := models.DiskManager.FetchById(diskId) + if err != nil { + self.SetStageFailed(ctx, err.Error()) + return + } + if iDisk == nil { + self.SetStageFailed(ctx, "Disk not found") + return + } + disk := iDisk.(*models.SDisk) + if disk.Status == models.DISK_INIT { + snapInfo, err := self.Params.GetString(fmt.Sprintf("disk.%d.snapshot", diskIndex)) + if err != nil { + snapInfo = "" + } + err = disk.StartDiskCreateTask(ctx, self.UserCred, false, snapInfo, self.GetTaskId()) + if err != nil { + self.SetStageFailed(ctx, err.Error()) + return + } + diskReady = false + break + } + diskIndex += 1 + } + diskIndex = 0 + guest := obj.(*models.SGuest) + + for { + diskId, err := self.Params.GetString(fmt.Sprintf("disk.%d.id", diskIndex)) + if !diskReady || err != nil { + break + } + iDisk, err := models.DiskManager.FetchById(diskId) + if err != nil { + self.SetStageFailed(ctx, err.Error()) + return + } + if iDisk == nil { + self.SetStageFailed(ctx, "Disk not found") + return + } + disk := iDisk.(*models.SDisk) + if disk.Status != models.DISK_READY { + diskReady = false + break + } + + ihost, err := guest.GetHost().GetIHost() + if err != nil { + self.SetStageFailed(ctx, "Host not found") + return + } + + iVM, e := ihost.GetIVMById(guest.GetExternalId()) + if e != nil { + self.SetStageFailed(ctx, "Aliyun VM not found") + return + } + + err = iVM.AttachDisk(disk.GetExternalId()) + if err != nil { + log.Debugf("Attach Disk %s to guest fail: %s", diskId, err) + self.SetStageFailed(ctx, "Attach Disk to guest fail") + return + } + diskIndex += 1 + } + + if diskReady { + if guest.Status == models.VM_RUNNING { + self.SetStage("on_config_sync_complete", nil) + err := guest.StartSyncstatus(ctx, self.UserCred, self.GetTaskId()) + if err != nil { + self.SetStageFailed(ctx, err.Error()) + } + } else { + self.SetStageComplete(ctx, nil) + } + } +} + +func (self *AliyunGuestCreateDiskTask) OnConfigSyncComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *AliyunGuestCreateDiskTask) AttachAliyunDisks(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + func init() { taskman.RegisterTask(GuestCreateDiskTask{}) taskman.RegisterTask(KVMGuestCreateDiskTask{}) + taskman.RegisterTask(AliyunGuestCreateDiskTask{}) } diff --git a/pkg/compute/tasks/guest_delete_task.go b/pkg/compute/tasks/guest_delete_task.go index e663bf990d..548388eaac 100644 --- a/pkg/compute/tasks/guest_delete_task.go +++ b/pkg/compute/tasks/guest_delete_task.go @@ -4,6 +4,9 @@ import ( "context" "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/utils" + "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" @@ -29,15 +32,18 @@ func (self *GuestDeleteTask) OnInit(ctx context.Context, obj db.IStandaloneModel func (self *GuestDeleteTask) OnGuestStopComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) - guestStatus, _ := self.Params.GetString("guest_status") if options.Options.EnablePendingDelete && !guest.PendingDeleted && !jsonutils.QueryBoolean(self.Params, "purge", false) && - !jsonutils.QueryBoolean(self.Params, "override_pending_delete", false) && - !utils.IsInStringArray(guestStatus, []string{models.VM_SCHEDULE_FAILED, models.VM_NETWORK_FAILED, models.VM_DISK_FAILED, + !jsonutils.QueryBoolean(self.Params, "override_pending_delete", false) { + log.Debugf("XXXXXXX Do guest pending delete... XXXXXXX") + guestStatus, _ := self.Params.GetString("guest_status") + if !utils.IsInStringArray(guestStatus, []string{models.VM_SCHEDULE_FAILED, models.VM_NETWORK_FAILED, models.VM_DISK_FAILED, models.VM_CREATE_FAILED, models.VM_DEVICE_FAILED}) { - self.StartPendingDeleteGuest(ctx, guest) - return + self.StartPendingDeleteGuest(ctx, guest) + return + } } + log.Debugf("XXXXXXX Do real delete on guest ... XXXXXXX") self.OnGuestStopCompleteFailed(ctx, guest, data) } diff --git a/pkg/compute/tasks/guest_deploy_task.go b/pkg/compute/tasks/guest_deploy_task.go index 5f5dab563d..2b9c98f190 100644 --- a/pkg/compute/tasks/guest_deploy_task.go +++ b/pkg/compute/tasks/guest_deploy_task.go @@ -55,15 +55,11 @@ func (self *GuestDeployTask) StartDeployGuestOnHost(ctx context.Context, guest * log.Errorf("request_deploy_guest_on_host %s", err) self.OnDeployGuestFail(ctx, guest, err) } else { - self.OnDeployGuestSucc(guest) + guest.SetStatus(self.UserCred, models.VM_DEPLOYING, "") } logclient.AddActionLog(guest, logclient.ACT_VM_DEPLOY, err, self.UserCred) } -func (self *GuestDeployTask) OnDeployGuestSucc(guest *models.SGuest) { - guest.SetStatus(self.UserCred, models.VM_DEPLOYING, "") -} - func (self *GuestDeployTask) OnDeployGuestFail(ctx context.Context, guest *models.SGuest, err error) { guest.SetStatus(self.UserCred, models.VM_DEPLOY_FAILED, err.Error()) self.SetStageFailed(ctx, err.Error()) diff --git a/pkg/compute/tasks/guest_start_task.go b/pkg/compute/tasks/guest_start_task.go index 71147afdac..980af93c9b 100644 --- a/pkg/compute/tasks/guest_start_task.go +++ b/pkg/compute/tasks/guest_start_task.go @@ -74,6 +74,7 @@ func (self *GuestStartTask) OnStartCompleteFailed(ctx context.Context, obj db.IS } func (self *GuestStartTask) onStartGuestFailed(ctx context.Context, guest *models.SGuest, err error) { + guest.SetStatus(self.UserCred, models.VM_START_FAILED, err.Error()) self.SetStageFailed(ctx, err.Error()) self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) logclient.AddActionLog(guest, logclient.ACT_VM_START, err, self.UserCred) diff --git a/pkg/compute/tasks/guest_stop_task.go b/pkg/compute/tasks/guest_stop_task.go index f50838d64c..ebdf37d936 100644 --- a/pkg/compute/tasks/guest_stop_task.go +++ b/pkg/compute/tasks/guest_stop_task.go @@ -52,7 +52,7 @@ func (self *GuestStopTask) OnGuestStopTaskComplete(ctx context.Context, obj db.I if !self.isSubtask() { guest.SetStatus(self.UserCred, models.VM_READY, "") } - db.OpsLog.LogEvent(guest, db.ACT_STOP, nil, self.UserCred) + db.OpsLog.LogEvent(guest, db.ACT_STOP, guest.GetShortDesc(), self.UserCred) models.HostManager.ClearSchedDescCache(guest.HostId) self.SetStageComplete(ctx, nil) if guest.Status == models.VM_READY && guest.DisableDelete.IsFalse() && guest.ShutdownBehavior == models.SHUTDOWN_TERMINATE { diff --git a/pkg/dns/README.md b/pkg/dns/README.md new file mode 100644 index 0000000000..79bfeb66f3 --- /dev/null +++ b/pkg/dns/README.md @@ -0,0 +1,56 @@ +# 行为约定 + +A + +```sh +names='' +names="$names titan" #ok, guest PlainName +names="$names titan.hq.cloud.yunionyun.com" #ok, guest CloudZoneFQDN +names="$names kubenode" #ok, host PlainName +names="$names kubenode.hq.cloud.yunionyun.com" #ok, host CloudZoneFQDN +names="$names whoever-the-ether" #NXDOMAIN, nonexistent PlainName +names="$names whoever-the-ether.titan.hq.cloud.yunionyun.com" #NXDOMAIN, nonexistent CloudZoneFQDN +names="$names mail.google.com" #ok, dnsrecords +names="$names www.douban.com" #ok, pub +names="$names app" #ok, k8s svc PlainName in "default" ns +names="$names app.default" #NXDOMAIN, k8s svc name.namespace +names="$names app.default.hq.cloud.yunionyun.com" #ok, k8s name.ns CloudZoneFQDN +names="$names mon-kafka.system" #NXDOMAIN, k8s svc name.namespace +names="$names mon-kafka.system.hq.cloud.yunionyun.com" #ok, k8s name.ns CloudZoneFQDN + +# TODO source IP +# TODO other TYPEs +for name in $names; do + echo "############### $name" + #dig @192.168.222.171 $name + #dig -p 54 @10.168.222.136 $name + dig -p 54 @192.168.222.171 $name +done +``` + +PTR + +```sh +names='' +names="$names " #ok, guest ip +names="$names " #ok, host ip +names="$names " #ok, ptr records in db +names="$names " #NXDOMAIN, others + +for name in $names; do + echo "############### $name" + dig -p 54 @192.168.222.171 -x $name +done +``` + +# 配置 + + log { + # note that apart from rcode like NXDOMAIN, SERVFAIL, coredns will also + # log NOERROR response when it's NoData as defined by coredns itself + # + # > NoData indicates name found, but not the type: NOERROR in header, SOA in auth. + # + class denial + class error + } diff --git a/pkg/dns/dns.go b/pkg/dns/dns.go index b60ab73c51..248504b9f6 100644 --- a/pkg/dns/dns.go +++ b/pkg/dns/dns.go @@ -18,25 +18,27 @@ import ( "github.com/mholt/caddy" "github.com/miekg/dns" v1 "k8s.io/api/core/v1" + k8serrors "k8s.io/apimachinery/pkg/api/errors" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" "k8s.io/apimachinery/pkg/labels" "k8s.io/client-go/kubernetes" ylog "yunion.io/x/log" - "yunion.io/x/pkg/utils" - "yunion.io/x/sqlchemy" - "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/util/k8s" + "yunion.io/x/pkg/utils" + "yunion.io/x/sqlchemy" ) const ( PluginName string = "yunion" // defaultTTL to apply to all answers - defaultTTL = 10 + defaultTTL = 10 + defaultDbMaxOpenConn = 32 + defaultDbMaxIdleConn = 32 ) var ( @@ -74,15 +76,17 @@ func (r *SRegionDNS) initDB(c *caddy.Controller) error { if err != nil { return err } - dbConn, err := sql.Open(dialect, sqlStr) + sqlDb, err := sql.Open(dialect, sqlStr) if err != nil { return err } - sqlchemy.SetDB(dbConn) + sqlDb.SetMaxOpenConns(defaultDbMaxOpenConn) + sqlDb.SetMaxIdleConns(defaultDbMaxIdleConn) + sqlchemy.SetDB(sqlDb) db.InitAllManagers() c.OnShutdown(func() error { - r.CloseDB() + sqlchemy.CloseDB() return nil }) return nil @@ -103,10 +107,6 @@ func (r *SRegionDNS) initK8s(c *caddy.Controller) { ylog.Infof("Init k8s client success, %d pods in the cluster", len(pods.Items)) } -func (r *SRegionDNS) CloseDB() { - sqlchemy.CloseDB() -} - func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *dns.Msg) (int, error) { var ( records []dns.RR @@ -116,36 +116,26 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d opt := plugin.Options{} state := request.Request{W: w, Req: rmsg, Context: ctx} - zone := plugin.Zones(r.Zones).Matches(state.Name()) - switch state.QType() { case dns.TypeA: - ylog.Debugf("A question: %#v", state) records, err = plugin.A(r, zone, state, nil, opt) case dns.TypeAAAA: - ylog.Debugf("AAAA question: %#v", state) + // TODO fallthrough to next records, err = plugin.AAAA(r, zone, state, nil, opt) case dns.TypeTXT: - ylog.Debugf("TXT question: %#v", state) records, err = plugin.TXT(r, zone, state, opt) case dns.TypeCNAME: - ylog.Debugf("CNAME question: %#v", state) records, err = plugin.CNAME(r, zone, state, opt) case dns.TypePTR: - ylog.Debugf("PTR question: %#v", state) records, err = plugin.PTR(r, zone, state, opt) case dns.TypeMX: - ylog.Debugf("MX question: %#v", state) records, extra, err = plugin.MX(r, zone, state, opt) case dns.TypeSRV: - ylog.Debugf("SRV question: %#v", state) records, extra, err = plugin.SRV(r, zone, state, opt) case dns.TypeSOA: - ylog.Debugf("SOA question: %#v", state) records, err = plugin.SOA(r, zone, state, opt) case dns.TypeNS: - ylog.Debugf("NS question: %#v", state) if state.Name() == zone { records, extra, err = plugin.NS(r, zone, state, opt) break @@ -157,18 +147,19 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d _, err = plugin.A(r, zone, state, nil, opt) } - if r.IsNameError(err) { + if err == errCallNext { if r.Fall.Through(state.Name()) { return plugin.NextOrFailure(r.Name(), r.Next, ctx, w, rmsg) } return plugin.BackendError(r, zone, dns.RcodeNameError, state, nil /* err */, opt) - } - if err != nil { - return plugin.BackendError(r, zone, dns.RcodeServerFailure, state, err, opt) + } else if err == errRefused { + return plugin.BackendError(r, zone, dns.RcodeRefused, state, err, opt) + } else if err == errNotFound { + return plugin.BackendError(r, zone, dns.RcodeNameError, state, err, opt) } if len(records) == 0 { - return plugin.BackendError(r, zone, dns.RcodeSuccess, state, err, opt) + return plugin.BackendError(r, zone, dns.RcodeNameError, state, err, opt) } m := new(dns.Msg) @@ -184,13 +175,14 @@ func (r *SRegionDNS) ServeDNS(ctx context.Context, w dns.ResponseWriter, rmsg *d } var ( - errNoItems = errors.New("no items found") + errRefused = errors.New("refused the query") + errNotFound = errors.New("not found") + errCallNext = errors.New("continue to next") ) // Services implements the ServiceBackend interface func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Options) (services []msg.Service, err error) { switch state.QType() { - case dns.TypeTXT: t, _ := dnsutil.TrimZone(state.Name(), state.Zone) @@ -203,7 +195,6 @@ func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Opti } svc := msg.Service{Text: "0.0.1", TTL: 28800, Key: msg.Path(state.QName(), "coredns")} return []msg.Service{svc}, nil - case dns.TypeNS: ns := r.nsAddr() svc := msg.Service{Host: ns.A.String(), Key: msg.Path(state.QName(), "coredns")} @@ -219,7 +210,6 @@ func (r *SRegionDNS) Services(state request.Request, exact bool, opt plugin.Opti } services, err = r.Records(state, false) - ylog.Debugf("Get records: %#v, error: %v", services, err) return } @@ -230,7 +220,7 @@ func (r *SRegionDNS) Lookup(state request.Request, name string, typ uint16) (*dn // IsNameError implements the ServiceBackend interface func (r *SRegionDNS) IsNameError(err error) bool { - return err == errNoItems + return err == errCallNext } // Records looks up records in region mysql @@ -242,25 +232,22 @@ func (r *SRegionDNS) Records(state request.Request, exact bool) ([]msg.Service, return r.findRecords(req) } -func (r *SRegionDNS) getHostIpWithName(req *recordRequest) []string { +func (r *SRegionDNS) getHostIpWithName(req *recordRequest) string { name := req.QueryName() host, _ := models.HostManager.FetchByName("", name) if host == nil { - return nil + return "" } ip := host.(*models.SHost).AccessIp - if len(ip) == 0 { - return nil - } - return []string{ip} + return ip } func (r *SRegionDNS) getGuestIpWithName(req *recordRequest) []string { ips := []string{} name := req.QueryName() projectId := req.ProjectId() - isExitOnly := req.IsExitOnly() - ips = models.GuestManager.GetIpInProjectWithName(projectId, name, isExitOnly) + wantOnlyExit := false + ips = models.GuestManager.GetIpInProjectWithName(projectId, name, wantOnlyExit) return ips } @@ -268,6 +255,9 @@ func (r *SRegionDNS) getK8sServiceBackends(req *recordRequest) ([]string, error) queryInfo := req.GetK8sQueryInfo() pods, err := r.getK8sServicePods(queryInfo.Namespace, queryInfo.ServiceName) if err != nil { + if k8serrors.IsNotFound(err) { + err = nil + } return nil, err } ips := make([]string, 0) @@ -281,7 +271,11 @@ func (r *SRegionDNS) getK8sServiceBackends(req *recordRequest) ([]string, error) } func (r *SRegionDNS) getK8sServicePods(namespace, name string) ([]v1.Pod, error) { - cli := r.K8sClient + cli, err := k8s.NewClientByFile(r.K8sConfigFile, nil) + if err != nil { + ylog.Errorf("Init kubernetes client error: %v", err) + return nil, err + } svc, err := cli.CoreV1().Services(namespace).Get(name, metav1.GetOptions{}) if err != nil { return nil, err @@ -301,10 +295,9 @@ func (r *SRegionDNS) Name() string { return PluginName } -func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Service, err error) { +func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Service) { ips := models.DnsRecordManager.QueryDnsIps(req.ProjectId(), req.Name(), req.Type()) if len(ips) == 0 { - err = errNoItems return } @@ -317,12 +310,12 @@ func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Servic if req.IsSRV() { parts := strings.SplitN(ip.Addr, ":", 2) if len(parts) != 2 { - err = fmt.Errorf("Invalid SRV records: %q", ip.Addr) + ylog.Errorf("Invalid SRV records: %q", ip.Addr) return } port, e := strconv.Atoi(parts[1]) if e != nil { - err = e + ylog.Errorf("Invalid SRV records: %q", ip.Addr) return } s = msg.Service{Host: parts[0], Port: port, TTL: ttl} @@ -334,17 +327,6 @@ func (r *SRegionDNS) queryLocalDnsRecords(req *recordRequest) (recs []msg.Servic return } -func (r *SRegionDNS) IsCloudNetworkIp(req *recordRequest) bool { - if req.network != nil { - return true - } - return false -} - -func (r *SRegionDNS) IsK8sClientReady() bool { - return r.K8sClient != nil -} - func (r *SRegionDNS) isMyDomain(req *recordRequest) bool { zones := []string{fmt.Sprintf("%s.", r.PrimaryZone)} zone := plugin.Zones(zones).Matches(req.state.Name()) @@ -354,54 +336,65 @@ func (r *SRegionDNS) isMyDomain(req *recordRequest) bool { return false } -func (r *SRegionDNS) findRecords(req *recordRequest) (recs []msg.Service, err error) { +func (r *SRegionDNS) findRecords(req *recordRequest) ([]msg.Service, error) { // 1. try local dns records table - recs, err = r.queryLocalDnsRecords(req) - if len(recs) != 0 { - return + rrs := r.queryLocalDnsRecords(req) + if len(rrs) > 0 { + return rrs, nil } + isPlainName := req.IsPlainName() isMyDomain := r.isMyDomain(req) - - isCloudIp := r.IsCloudNetworkIp(req) - - // 2. not my domain and src ip not in cloud network table - // query from upstream - if !isMyDomain && !isCloudIp { - err = errNoItems - return + if isPlainName { + isCloudIp := req.SrcInCloud() + if isCloudIp { + ips := r.findInternalRecordIps(req) + if len(ips) > 0 { + return ips2DnsRecords(ips), nil + } else { + return nil, errNotFound + } + } else { + return nil, errRefused + } + } else if isMyDomain { + ips := r.findInternalRecordIps(req) + if len(ips) > 0 { + return ips2DnsRecords(ips), nil + } else { + return nil, errNotFound + } + } else { + return nil, errCallNext } - - // 3. internal query - ips, err := r.findInternalRecordIps(req) - return ips2DnsRecords(ips), err } -func (r *SRegionDNS) findInternalRecordIps(req *recordRequest) ([]string, error) { - // 1. try host table - ip := r.getHostIpWithName(req) - if len(ip) != 0 { - return ip, nil +func (r *SRegionDNS) findInternalRecordIps(req *recordRequest) []string { + { + // 1. try host table + ip := r.getHostIpWithName(req) + if len(ip) > 0 { + return []string{ip} + } } - // 2. try guest table - ip = r.getGuestIpWithName(req) - if len(ip) != 0 { - return ip, nil + { + // 2. try guest table + ips := r.getGuestIpWithName(req) + if len(ips) > 0 { + return ips + } } - if !r.IsK8sClientReady() { + if r.K8sClient == nil { ylog.Warningf("K8s client not ready, skip it.") - return nil, errNoItems + return nil } // 3. try k8s service backends ips, err := r.getK8sServiceBackends(req) - if len(ips) != 0 { - return ips, nil - } if err != nil { ylog.Errorf("Get k8s service backends error: %v", err) } - return nil, errNoItems + return ips } func ips2DnsRecords(ips []string) []msg.Service { diff --git a/pkg/dns/parse.go b/pkg/dns/parse.go index 5cf836b646..568a773567 100644 --- a/pkg/dns/parse.go +++ b/pkg/dns/parse.go @@ -7,17 +7,16 @@ import ( "github.com/coredns/coredns/request" "github.com/miekg/dns" - "yunion.io/x/pkg/tristate" - "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/pkg/tristate" ) type recordRequest struct { - state request.Request - domainSegs []string - guest *models.SGuest - host *models.SHost - network *models.SNetwork + state request.Request + domainSegs []string + srcProjectId string + srcInCloud bool + network *models.SNetwork } func parseRequest(state request.Request) (r *recordRequest, err error) { @@ -28,9 +27,19 @@ func parseRequest(state request.Request) (r *recordRequest, err error) { domainSegs: segs, } srcIP := r.SrcIP4() - r.guest = models.GuestnetworkManager.GetGuestByAddress(srcIP) - r.host = models.HostnetworkManager.GetHostByAddress(srcIP) - r.network, _ = models.NetworkManager.GetNetworkOfIP(srcIP, "", tristate.None) + // NOTE the check on networks_tbl is a hack, we should be more specific + // by querying only guest networks, host networks, and others like + // loadbalancer network the to come. + // + // Order matters here, we want to find the srcIP project as accurately + // as possible + if guest := models.GuestnetworkManager.GetGuestByAddress(srcIP); guest != nil { + r.srcProjectId = guest.ProjectId + r.srcInCloud = true + } else if network, _ := models.NetworkManager.GetNetworkOfIP(srcIP, "", tristate.None); network != nil { + r.srcProjectId = network.ProjectId + r.srcInCloud = true + } return } @@ -41,6 +50,11 @@ func (r recordRequest) Name() string { return name } +func (r recordRequest) IsPlainName() bool { + nl := dns.CountLabel(r.Name()) + return nl == 1 +} + func (r recordRequest) QueryName() string { seps := strings.Split(r.Name(), ".") if len(seps) == 0 { @@ -63,20 +77,11 @@ func (r recordRequest) SrcIP4() string { } func (r recordRequest) ProjectId() string { - if r.guest != nil { - return r.guest.ProjectId - } - if r.network != nil { - return r.network.ProjectId - } - return "" + return r.srcProjectId } -func (r recordRequest) IsExitOnly() bool { - if r.guest == nil { - return false - } - return r.guest.IsExitOnly() +func (r recordRequest) SrcInCloud() bool { + return r.srcInCloud } type K8sQueryInfo struct { diff --git a/pkg/dns/reverse.go b/pkg/dns/reverse.go index e6cf7e2276..f6df515fdc 100644 --- a/pkg/dns/reverse.go +++ b/pkg/dns/reverse.go @@ -8,8 +8,6 @@ import ( "github.com/coredns/coredns/plugin/pkg/dnsutil" "github.com/coredns/coredns/request" - "yunion.io/x/log" - "yunion.io/x/onecloud/pkg/compute/models" ) @@ -21,13 +19,7 @@ func (r *SRegionDNS) Reverse(state request.Request, exact bool, opt plugin.Optio return nil, e } records, err := r.getNameForIp(ip, state) - if err != nil { - log.Errorf("Reverse get name for ip: %v", err) - } - if len(records) == 0 { - return records, errNoItems - } - return records, nil + return records, err } func (r *SRegionDNS) getNameForIp(ip string, state request.Request) ([]msg.Service, error) { @@ -53,7 +45,7 @@ func (r *SRegionDNS) getNameForIp(ip string, state request.Request) ([]msg.Servi if guest != nil { return []msg.Service{{Host: r.joinDomain(guest.Name), TTL: defaultTTL}}, nil } - return nil, errNoItems + return nil, errNotFound } func (r *SRegionDNS) joinDomain(name string) string { diff --git a/pkg/dns/setup.go b/pkg/dns/setup.go index 1a6a1ff673..77a9cb1c4b 100644 --- a/pkg/dns/setup.go +++ b/pkg/dns/setup.go @@ -2,7 +2,6 @@ package dns import ( "fmt" - "os" "github.com/coredns/coredns/core/dnsserver" "github.com/coredns/coredns/plugin" @@ -20,18 +19,16 @@ func init() { } func setup(c *caddy.Controller) error { - os.Stderr = os.Stdout - rDNS, err := regionDNSParse(c) if err != nil { return plugin.Error(PluginName, err) } - if rDNS.PrimaryZone == "" { - return fmt.Errorf("dns_domain must provided") + if len(rDNS.PrimaryZone) == 0 { + return fmt.Errorf("dns_domain missing") } if !regutils.MatchDomainName(rDNS.PrimaryZone) { - return fmt.Errorf("dns_domain %q not match domain format", rDNS.PrimaryZone) + return fmt.Errorf("dns_domain %q invalid", rDNS.PrimaryZone) } err = rDNS.initDB(c) diff --git a/pkg/httperrors/errors.go b/pkg/httperrors/errors.go index d57bf8ff79..0daa2cc0b4 100644 --- a/pkg/httperrors/errors.go +++ b/pkg/httperrors/errors.go @@ -6,219 +6,173 @@ import ( "yunion.io/x/onecloud/pkg/util/httputils" ) -func NewJsonClientError(code int, title string, msg string) *httputils.JSONClientError { - err := httputils.JSONClientError{Code: code, Class: title, Details: msg} +func NewJsonClientError(code int, title string, msg string, error httputils.Error) *httputils.JSONClientError { + err := httputils.JSONClientError{Code: code, Class: title, Details: msg, Data: error} return &err } -func NewBadGatewayError(msg string, params ...interface{}) *httputils.JSONClientError { +func errorMessage(msg string, params ...interface{}) (string, httputils.Error) { + fileds := make([]string, len(params)) + for i, v := range params { + fileds[i] = fmt.Sprint(v) + } + + error := httputils.Error{Id: msg, Fields: fileds} if len(params) > 0 { msg = fmt.Sprintf(msg, params...) } - return NewJsonClientError(502, "BadGateway", msg) + + return msg, error +} + +func NewBadGatewayError(msg string, params ...interface{}) *httputils.JSONClientError { + msg, err := errorMessage(msg, params) + return NewJsonClientError(502, "BadGateway", msg, err) } func NewNotImplementedError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(501, "NotImplemented", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(501, "NotImplemented", msg, err) } func NewInternalServerError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(500, "InternalServerError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(500, "InternalServerError", msg, err) } func NewResourceNotReadyError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(500, "ResourceNotReadyError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(500, "ResourceNotReadyError", msg, err) } func NewOutOfResourceError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(500, "NewOutOfResourceError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(500, "NewOutOfResourceError", msg, err) } func NewServerStatusError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(400, "ServerStatusError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(400, "ServerStatusError", msg, err) } func NewPaymentError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(402, "PaymentError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(402, "PaymentError", msg, err) } func NewImageNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "ImageNotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "ImageNotFoundError", msg, err) } func NewResourceNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "ResourceNotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "ResourceNotFoundError", msg, err) } func NewSpecNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "SpecNotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "SpecNotFoundError", msg, err) } func NewActionNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "ActionNotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "ActionNotFoundError", msg, err) } func NewTenantNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "TenantNotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "TenantNotFoundError", msg, err) } func NewUserNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "UserNotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "UserNotFoundError", msg, err) } func NewInvalidStatusError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(400, "InvalidStatusError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(400, "InvalidStatusError", msg, err) } func NewInputParameterError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(400, "InputParameterError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(400, "InputParameterError", msg, err) } func NewInsufficientResourceError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(400, "InsufficientResourceError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(400, "InsufficientResourceError", msg, err) } func NewOutOfQuotaError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(400, "OutOfQuotaError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(400, "OutOfQuotaError", msg, err) } func NewNotSufficientPrivilegeError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(403, "NotSufficientPrivilegeError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(403, "NotSufficientPrivilegeError", msg, err) } func NewUnsupportOperationError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(406, "UnsupportOperationError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(406, "UnsupportOperationError", msg, err) } func NewNotEmptyError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(406, "NotEmptyError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(406, "NotEmptyError", msg, err) } func NewBadRequestError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(400, "BadRequestError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(400, "BadRequestError", msg, err) } func NewUnauthorizedError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(401, "UnauthorizedError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(401, "UnauthorizedError", msg, err) } func NewInvalidCredentialError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(401, "InvalidCredentialError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(401, "InvalidCredentialError", msg, err) } func NewForbiddenError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(403, "ForbiddenError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(403, "ForbiddenError", msg, err) } func NewNotFoundError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(404, "NotFoundError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(404, "NotFoundError", msg, err) } func NewNotAcceptableError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(406, "NotAcceptableError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(406, "NotAcceptableError", msg, err) } func NewDuplicateNameError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(409, "DuplicateNameError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(409, "DuplicateNameError", msg, err) } func NewConflictError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(409, "ConflictError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(409, "ConflictError", msg, err) } func NewResourceBusyError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(409, "ResourceBusyError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(409, "ResourceBusyError", msg, err) } func NewRequireLicenseError(msg string, params ...interface{}) *httputils.JSONClientError { - if len(params) > 0 { - msg = fmt.Sprintf(msg, params...) - } - return NewJsonClientError(402, "RequireLicenseError", msg) + msg, err := errorMessage(msg, params) + return NewJsonClientError(402, "RequireLicenseError", msg, err) } func NewGeneralError(err error) *httputils.JSONClientError { diff --git a/pkg/httperrors/httperrors.go b/pkg/httperrors/httperrors.go index a696ab0394..84dc0d3216 100644 --- a/pkg/httperrors/httperrors.go +++ b/pkg/httperrors/httperrors.go @@ -8,18 +8,22 @@ import ( "yunion.io/x/onecloud/pkg/util/httputils" ) -func HTTPError(w http.ResponseWriter, msg string, statusCode int, class string) { +func HTTPError(w http.ResponseWriter, msg string, statusCode int, class string, error httputils.Error) { w.WriteHeader(statusCode) w.Header().Set("Content-Type", "application/json") body := jsonutils.NewDict() body.Add(jsonutils.NewInt(int64(statusCode)), "code") body.Add(jsonutils.NewString(msg), "details") body.Add(jsonutils.NewString(class), "class") + err := jsonutils.NewDict() + err.Add(jsonutils.NewString(error.Id), "id") + err.Add(jsonutils.NewStringArray(error.Fields), "fields") + body.Add(err, "data") w.Write([]byte(body.String())) } func JsonClientError(w http.ResponseWriter, e *httputils.JSONClientError) { - HTTPError(w, e.Details, e.Code, e.Class) + HTTPError(w, e.Details, e.Code, e.Class, e.Data) } func GeneralServerError(w http.ResponseWriter, e error) { @@ -31,52 +35,52 @@ func GeneralServerError(w http.ResponseWriter, e error) { } } -func BadRequestError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewBadRequestError(msg)) +func BadRequestError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewBadRequestError(msg, params...)) } -func PaymentError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewPaymentError(msg)) +func PaymentError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewPaymentError(msg, params...)) } -func UnauthorizedError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewUnauthorizedError(msg)) +func UnauthorizedError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewUnauthorizedError(msg, params...)) } -func InvalidCredentialError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewInvalidCredentialError(msg)) +func InvalidCredentialError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewInvalidCredentialError(msg, params...)) } -func ForbiddenError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewForbiddenError(msg)) +func ForbiddenError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewForbiddenError(msg, params...)) } -func NotFoundError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewNotFoundError(msg)) +func NotFoundError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewNotFoundError(msg, params...)) } -func NotImplementedError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewNotImplementedError(msg)) +func NotImplementedError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewNotImplementedError(msg, params...)) } -func NotAcceptableError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewNotAcceptableError(msg)) +func NotAcceptableError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewNotAcceptableError(msg, params...)) } -func InvalidInputError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewInputParameterError(msg)) +func InvalidInputError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewInputParameterError(msg, params...)) } -func ConflictError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewConflictError(msg)) +func ConflictError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewConflictError(msg, params...)) } -func InternalServerError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewInternalServerError(msg)) +func InternalServerError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewInternalServerError(msg, params...)) } -func BadGatewayError(w http.ResponseWriter, msg string) { - JsonClientError(w, NewBadGatewayError(msg)) +func BadGatewayError(w http.ResponseWriter, msg string, params ...interface{}) { + JsonClientError(w, NewBadGatewayError(msg, params...)) } func TenantNotFoundError(w http.ResponseWriter, msg string, params ...interface{}) { diff --git a/pkg/mcclient/auth/auth.go b/pkg/mcclient/auth/auth.go index 214be4b0ad..4e0957e66a 100644 --- a/pkg/mcclient/auth/auth.go +++ b/pkg/mcclient/auth/auth.go @@ -103,7 +103,7 @@ func (c *TokenCacheVerify) Verify(cli *mcclient.Client, adminToken, token string if err != nil { return nil, fmt.Errorf("Add %s credential to cache: %#v", cred.GetTokenString(), err) } - log.Infof("Add token: %s", cred) + // log.Debugf("Add token: %s", cred) return cred, nil } diff --git a/pkg/mcclient/modules/mod_domains.go b/pkg/mcclient/modules/mod_domains.go index 84d57ffb34..e7be85cd5e 100644 --- a/pkg/mcclient/modules/mod_domains.go +++ b/pkg/mcclient/modules/mod_domains.go @@ -240,12 +240,7 @@ func (this *DomainManager) DoDomainConfigDelete(s *mcclient.ClientSession, param return ret, httperrors.NewResourceNotFoundError("找不到该认证域") } - driver, err := detail.GetString("driver") - if err != nil { - log.Errorf("got driver from domain detail error: %v", err) - return ret, httperrors.NewInternalServerError("服务器错误,获取认证协议失败,不允许删除") - } - + driver, _ := detail.GetString("driver") if driver != "ldap" { if result, err := UsersV3.List(s, params); err != nil { log.Errorf("user list got error: %v", err) diff --git a/pkg/util/aliyun/aliyun.go b/pkg/util/aliyun/aliyun.go index 26b9efce28..9d002a46d1 100644 --- a/pkg/util/aliyun/aliyun.go +++ b/pkg/util/aliyun/aliyun.go @@ -3,6 +3,7 @@ package aliyun import ( "github.com/aliyun/alibaba-cloud-sdk-go/sdk" "github.com/aliyun/alibaba-cloud-sdk-go/sdk/requests" + "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" diff --git a/pkg/util/aliyun/disk.go b/pkg/util/aliyun/disk.go index aeeca8096c..7349b8c7c0 100644 --- a/pkg/util/aliyun/disk.go +++ b/pkg/util/aliyun/disk.go @@ -54,6 +54,16 @@ type SDisk struct { ZoneId string } +func (self *SDisk) GetMetadata() *jsonutils.JSONDict { + data := jsonutils.NewDict() + + // The pricingInfo key structure is 'RegionId::DiskCategory::DiskType + priceKey := fmt.Sprintf("%s::%s::%s", self.RegionId, self.Category, self.Type) + data.Add(jsonutils.NewString(priceKey), "price_key") + + return data +} + func (self *SRegion) GetDisks(instanceId string, zoneId string, category string, diskIds []string, offset int, limit int) ([]SDisk, int, error) { if limit > 50 || limit <= 0 { limit = 50 @@ -138,6 +148,12 @@ func (self *SDisk) Refresh() error { return jsonutils.Update(self, new) } +func (self *SDisk) ResizeDisk(newSize int64) error { + // newSize 单位为 GB. 范围在20 ~2000. 只能往大调。不能调小 + // https://help.aliyun.com/document_detail/25522.html?spm=a2c4g.11174283.6.897.aHwqkS + return self.storage.zone.region.resizeDisk(self.DiskId, newSize) +} + func (self *SDisk) GetDiskFormat() string { return "vhd" } @@ -185,7 +201,7 @@ func (self *SDisk) GetMountpoint() string { return "" } -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, desc string) (string, error) { params := make(map[string]string) params["ZoneId"] = zoneId params["DiskName"] = name @@ -223,11 +239,24 @@ func (self *SRegion) deleteDisk(diskId string) error { return err } +func (self *SRegion) DeleteDisk(diskId string) error { + params := make(map[string]string) + params["DiskId"] = diskId + + _, err := self.ecsRequest("DeleteDisk", params) + return err +} + func (self *SRegion) resizeDisk(diskId string, size int64) error { params := make(map[string]string) params["DiskId"] = diskId params["NewSize"] = fmt.Sprintf("%d", size) _, err := self.ecsRequest("ResizeDisk", params) - return err + if err != nil { + log.Errorf("ResizeDisk %s to %s GiB fail %s", diskId, size, err) + return err + } + + return nil } diff --git a/pkg/util/aliyun/errors.go b/pkg/util/aliyun/errors.go new file mode 100644 index 0000000000..da618ab69d --- /dev/null +++ b/pkg/util/aliyun/errors.go @@ -0,0 +1,13 @@ +package aliyun + +import ( + aliyunerrors "github.com/aliyun/alibaba-cloud-sdk-go/sdk/errors" +) + +func isError(err error, code string) bool { + aliyunErr, ok := err.(aliyunerrors.Error) + if ! ok { + return false + } + return aliyunErr.ErrorCode() == code +} diff --git a/pkg/util/aliyun/host.go b/pkg/util/aliyun/host.go index 3ee4a40236..24c5a2a1b1 100644 --- a/pkg/util/aliyun/host.go +++ b/pkg/util/aliyun/host.go @@ -19,6 +19,10 @@ type SHost struct { zone *SZone } +func (self *SHost) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SHost) GetIWires() ([]cloudprovider.ICloudWire, error) { return self.zone.GetIWires() } @@ -157,7 +161,7 @@ func (self *SHost) GetManagerId() string { return self.zone.region.client.providerId } -func (self *SHost) getInstanceById(instanceId string) (*SInstance, error) { +func (self *SHost) GetInstanceById(instanceId string) (*SInstance, error) { inst, err := self.zone.region.GetInstance(instanceId) if err != nil { return nil, err @@ -171,7 +175,7 @@ func (self *SHost) CreateVM(name string, imgId string, sysDiskSize int, cpu int, if err != nil { return nil, err } - vm, err := self.getInstanceById(vmId) + vm, err := self.GetInstanceById(vmId) if err != nil { return nil, err } diff --git a/pkg/util/aliyun/image.go b/pkg/util/aliyun/image.go index 9dfb6f0b3f..849e9eb570 100644 --- a/pkg/util/aliyun/image.go +++ b/pkg/util/aliyun/image.go @@ -55,6 +55,10 @@ type SImage struct { Usage string } +func (self *SImage) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SImage) GetId() string { return self.ImageId } diff --git a/pkg/util/aliyun/instance.go b/pkg/util/aliyun/instance.go index 80a2e9c267..7f034089b6 100644 --- a/pkg/util/aliyun/instance.go +++ b/pkg/util/aliyun/instance.go @@ -4,6 +4,8 @@ import ( "fmt" "time" + "yunion.io/x/onecloud/pkg/util/seclib2" + "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/util/osprofile" @@ -110,6 +112,7 @@ type SInstance struct { InternetMaxBandwidthIn int InternetMaxBandwidthOut int IoOptimized bool + KeyPairName string Memory int NetworkInterfaces SNetworkInterfaces OSName string @@ -162,6 +165,20 @@ func (self *SRegion) GetInstances(zoneId string, ids []string, offset int, limit return instances, int(total), nil } +func (self *SInstance) GetMetadata() *jsonutils.JSONDict { + data := jsonutils.NewDict() + + // The pricingInfo key structure is 'RegionId::InstanceType::NetworkType::OSType::IoOptimized' + optimized := "optimized" + if !self.IoOptimized { + optimized = "none" + } + priceKey := fmt.Sprintf("%s::%s::%s::%s::%s", self.RegionId, self.InstanceType, self.InstanceNetworkType, self.OSType, optimized) + data.Add(jsonutils.NewString(priceKey), "price_key") + + return data +} + func (self *SInstance) GetCreateTime() time.Time { return self.CreationTime } @@ -322,6 +339,73 @@ func (self *SInstance) GetHypervisor() string { return models.HYPERVISOR_ALIYUN } +func (self *SInstance) StartVM() error { + err := self.host.zone.region.StartVM(self.InstanceId) + if err != nil { + return err + } + return cloudprovider.WaitStatus(self, models.VM_RUNNING, 5*time.Second, 180*time.Second) // 3minutes +} + +func (self *SInstance) StopVM(isForce bool) error { + err := self.host.zone.region.StopVM(self.InstanceId, isForce) + if err != nil { + return err + } + return cloudprovider.WaitStatus(self, models.VM_READY, 10*time.Second, 300*time.Second) // 5mintues +} + +func (self *SInstance) GetVNCInfo() (jsonutils.JSONObject, error) { + url, err := self.host.zone.region.GetInstanceVNCUrl(self.InstanceId) + if err != nil { + return nil, err + } + passwd := seclib.RandomPassword(6) + err = self.host.zone.region.ModifyInstanceVNCUrlPassword(self.InstanceId, passwd) + if err != nil { + return nil, err + } + ret := jsonutils.NewDict() + ret.Add(jsonutils.NewString(url), "url") + ret.Add(jsonutils.NewString(passwd), "password") + ret.Add(jsonutils.NewString("aliyun"), "protocol") + ret.Add(jsonutils.NewString(self.InstanceId), "instance_id") + return ret, nil +} + +func (self *SInstance) UpdateVM(name string) error { + return self.host.zone.region.UpdateVM(self.InstanceId, name) +} + +func (self *SInstance) DeployVM(name string, password string, publicKey string, resetPassword bool, deleteKeypair bool, description string) error { + var keypairName string + if len(publicKey) > 0 { + key, e := self.host.lookUpAliyunKeypair(publicKey) + if e != nil { + key, e = self.host.importAliyunKeypair(publicKey) + if e != nil { + return e + } + } + + keypairName = key + } + + return self.host.zone.region.DeployVM(self.InstanceId, name, password, keypairName, resetPassword, deleteKeypair, description) +} + +func (self *SInstance) RebuildRoot(imageId string) error { + return self.host.zone.region.ReplaceSystemDisk(self.InstanceId, imageId) +} + +func (self *SInstance) ChangeConfig(instanceId string, ncpu int, vmem int) error { + return self.host.zone.region.ChangeVMConfig(self.ZoneId, self.InstanceId, ncpu, vmem, nil) +} + +func (self *SInstance) AttachDisk(diskId string) error { + return self.host.zone.region.AttachDisk(self.InstanceId, diskId) +} + func (self *SRegion) GetInstance(instanceId string) (*SInstance, error) { instances, _, err := self.GetInstances("", []string{instanceId}, 0, 1) if err != nil { @@ -402,7 +486,10 @@ func (self *SRegion) doStopVM(instanceId string, isForce bool) error { } func (self *SRegion) doDeleteVM(instanceId string) error { - return self.instanceOperation(instanceId, "DeleteInstance", nil) + params := make(map[string]string) + params["TerminateSubscription"] = "false" + params["Force"] = "true" + return self.instanceOperation(instanceId, "DeleteInstance", params) } /*func (self *SRegion) waitInstanceStatus(instanceId string, target string, interval time.Duration, timeout time.Duration) error { @@ -425,8 +512,13 @@ func (self *SInstance) waitStatus(target string, interval time.Duration, timeout }*/ func (self *SRegion) StartVM(instanceId string) error { - status, _ := self.GetInstanceStatus(instanceId) + status, err := self.GetInstanceStatus(instanceId) + if err != nil { + log.Errorf("Fail to get instance status on StartVM: %s", err) + return err + } if status != InstanceStatusStopped { + log.Errorf("StartVM: vm status is %s expect %s", status, InstanceStatusStopped) return cloudprovider.ErrInvalidStatus } return self.doStartVM(instanceId) @@ -437,8 +529,13 @@ func (self *SRegion) StartVM(instanceId string) error { } func (self *SRegion) StopVM(instanceId string, isForce bool) error { - status, _ := self.GetInstanceStatus(instanceId) + status, err := self.GetInstanceStatus(instanceId) + if err != nil { + log.Errorf("Fail to get instance status on StopVM: %s", err) + return err + } if status != InstanceStatusRunning { + log.Errorf("StopVM: vm status is %s expect %s", status, InstanceStatusRunning) return cloudprovider.ErrInvalidStatus } return self.doStopVM(instanceId, isForce) @@ -450,13 +547,13 @@ func (self *SRegion) StopVM(instanceId string, isForce bool) error { func (self *SRegion) DeleteVM(instanceId string) error { status, err := self.GetInstanceStatus(instanceId) - if status == InstanceStatusRunning { - err = self.StopVM(instanceId, true) - if err != nil { - return err - } - } else if status != InstanceStatusStopped { - return cloudprovider.ErrInvalidStatus + if err != nil { + log.Errorf("Fail to get instance status on DeleteVM: %s", err) + return err + } + log.Debugf("Instance status on delete is %s", status) + if status != InstanceStatusStopped { + log.Warningf("DeleteVM: vm status is %s expect %s", status, InstanceStatusStopped) } return self.doDeleteVM(instanceId) // if err != nil { @@ -472,46 +569,115 @@ func (self *SRegion) DeleteVM(instanceId string) error { // } } -func (self *SInstance) StartVM() error { - err := self.host.zone.region.StartVM(self.InstanceId) +func (self *SRegion) DeployVM(instanceId string, name string, password string, keypairName string, resetPassword bool, deleteKeypair bool, description string) error { + instance, err := self.GetInstance(instanceId) if err != nil { return err } - return cloudprovider.WaitStatus(self, models.VM_RUNNING, 5*time.Second, 180*time.Second) // 3minutes -} -func (self *SInstance) StopVM(isForce bool) error { - err := self.host.zone.region.StopVM(self.InstanceId, isForce) - if err != nil { - return err + // 修改密钥时直接返回 + if deleteKeypair { + return self.DetachKeyPair(instanceId, instance.KeyPairName) + } + + if len(keypairName) > 0 { + return self.AttachKeypair(instanceId, keypairName) + } + + params := make(map[string]string) + if resetPassword { + params["Password"] = seclib2.RandomPassword2(12) + } + // 指定密码的情况下,使用指定的密码 + if len(password) > 0 { + params["Password"] = password + } + + if len(name) > 0 && instance.InstanceName != name { + params["InstanceName"] = name + params["HostName"] = name + } + + if len(description) > 0 && instance.Description != description { + params["Description"] = description + } + + if len(params) > 0 { + return self.modifyInstanceAttribute(instanceId, params) + } else { + return nil } - return cloudprovider.WaitStatus(self, models.VM_READY, 10*time.Second, 300*time.Second) // 5mintues } func (self *SInstance) DeleteVM() error { - err := self.host.zone.region.DeleteVM(self.InstanceId) - if err != nil { - return err + for { + err := self.host.zone.region.DeleteVM(self.InstanceId) + if err != nil { + if isError(err, "IncorrectInstanceStatus.Initializing") { + log.Infof("The instance is initializing, try later ...") + time.Sleep(10 * time.Second) + } else { + return err + } + } else { + break + } } return cloudprovider.WaitDeleted(self, 10*time.Second, 300*time.Second) // 5minutes } -func (self *SInstance) GetVNCInfo() (jsonutils.JSONObject, error) { - url, err := self.host.zone.region.GetInstanceVNCUrl(self.InstanceId) - if err != nil { - return nil, err +func (self *SRegion) UpdateVM(instanceId string, hostname string) error { + /* + api: ModifyInstanceAttribute + https://help.aliyun.com/document_detail/25503.html?spm=a2c4g.11186623.4.1.DrgpjW + */ + params := make(map[string]string) + params["HostName"] = hostname + return self.modifyInstanceAttribute(instanceId, params) +} + +func (self *SRegion) modifyInstanceAttribute(instanceId string, params map[string]string) error { + return self.instanceOperation(instanceId, "ModifyInstanceAttribute", params) +} + +func (self *SRegion) ReplaceSystemDisk(instanceId string, image string) error { + params := make(map[string]string) + params["ImageId"] = image + return self.instanceOperation(instanceId, "ReplaceSystemDisk", params) +} + +func (self *SRegion) ChangeVMConfig(zoneId string, instanceId string, ncpu int, vmem int, disks []*SDisk) error { + // todo: support change disk config? + params := make(map[string]string) + instanceTypes, e := self.GetMatchInstanceTypes(ncpu, vmem, 0, zoneId) + if e != nil { + return e } - passwd := seclib.RandomPassword(6) - err = self.host.zone.region.ModifyInstanceVNCUrlPassword(self.InstanceId, passwd) - if err != nil { - return nil, err + + for _, instancetype := range instanceTypes { + params["InstanceType"] = instancetype.InstanceTypeId + params["ClientToken"] = utils.GenRequestId(20) + if err := self.instanceOperation(instanceId, "ModifyInstanceSpec", params); err != nil { + log.Errorf("Failed for %s: %s", instancetype.InstanceTypeId, err) + } else { + return nil + } } - ret := jsonutils.NewDict() - ret.Add(jsonutils.NewString(url), "url") - ret.Add(jsonutils.NewString(passwd), "password") - ret.Add(jsonutils.NewString("aliyun"), "protocol") - ret.Add(jsonutils.NewString(self.InstanceId), "instance_id") - return ret, nil + + return fmt.Errorf("Failed to change vm config, specification not supported") +} + +func (self *SRegion) AttachDisk(instanceId string, diskId string) error { + params := make(map[string]string) + params["InstanceId"] = instanceId + params["DiskId"] = diskId + _, err := self.ecsRequest("AttachDisk", params) + if err != nil { + log.Errorf("AttachDisk %s to %s fail %s", diskId, instanceId, err) + return err + } + + return nil } func (self *SInstance) SyncSecurityGroup(secgroupId string, name string, rules []secrules.SecurityRule) error { diff --git a/pkg/util/aliyun/keypair.go b/pkg/util/aliyun/keypair.go index c8da45b6a0..e3b2a504a8 100644 --- a/pkg/util/aliyun/keypair.go +++ b/pkg/util/aliyun/keypair.go @@ -1,6 +1,7 @@ package aliyun import ( + "encoding/json" "fmt" "yunion.io/x/log" ) @@ -62,3 +63,33 @@ func (self *SRegion) ImportKeypair(name string, pubKey string) (*SKeypair, error } return &keypair, nil } + +func (self *SRegion) AttachKeypair(instanceId string, name string) error { + params := make(map[string]string) + params["RegionId"] = self.RegionId + params["KeyPairName"] = name + instances, _ := json.Marshal(&[...]string{instanceId}) + params["InstanceIds"] = string(instances) + _, err := self.ecsRequest("AttachKeyPair", params) + if err != nil { + log.Errorf("AttachKeyPair fail %s", err) + return err + } + + return nil +} + +func (self *SRegion) DetachKeyPair(instanceId string, name string) error { + params := make(map[string]string) + params["RegionId"] = self.RegionId + params["KeyPairName"] = name + instances, _ := json.Marshal(&[...]string{instanceId}) + params["InstanceIds"] = string(instances) + _, err := self.ecsRequest("DetachKeyPair", params) + if err != nil { + log.Errorf("DetachKeyPair fail %s", err) + return err + } + + return nil +} diff --git a/pkg/util/aliyun/region.go b/pkg/util/aliyun/region.go index 78efb6a653..311f7bf195 100644 --- a/pkg/util/aliyun/region.go +++ b/pkg/util/aliyun/region.go @@ -36,6 +36,10 @@ func (self *SRegion) GetClient() *SAliyunClient { return self.client } +func (self *SRegion) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SRegion) getEcsClient() (*sdk.Client, error) { if self.ecsClient == nil { cli, err := sdk.NewClientWithAccessKey(self.RegionId, self.client.accessKey, self.client.secret) diff --git a/pkg/util/aliyun/securitygroup.go b/pkg/util/aliyun/securitygroup.go index 69515ae79d..3082451737 100644 --- a/pkg/util/aliyun/securitygroup.go +++ b/pkg/util/aliyun/securitygroup.go @@ -77,6 +77,10 @@ func (v PermissionSet) Less(i, j int) bool { return false } +func (self *SSecurityGroup) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SSecurityGroup) GetId() string { return self.SecurityGroupId } @@ -574,3 +578,15 @@ func (self *SRegion) leaveSecurityGroup(secgroupId, instanceId string) error { _, err := self.ecsRequest("LeaveSecurityGroup", params) return err } + +func (self *SRegion) deleteSecurityGroup(secGrpId string) error { + params := make(map[string]string) + params["SecurityGroupId"] = secGrpId + + _, err := self.ecsRequest("DeleteSecurityGroup", params) + if err != nil { + log.Errorf("Delete security group fail %s", err) + return err + } + return nil +} diff --git a/pkg/util/aliyun/shell/disk.go b/pkg/util/aliyun/shell/disk.go index 1196bb5739..19cc4c7693 100644 --- a/pkg/util/aliyun/shell/disk.go +++ b/pkg/util/aliyun/shell/disk.go @@ -21,4 +21,15 @@ func init() { printList(disks, total, args.Offset, args.Limit, []string{}) return nil }) + + type DiskDeleteOptions struct { + Instance string `help:"Instance ID"` + } + shellutils.R(&DiskDeleteOptions{}, "disk-delete", "List disks", func(cli *aliyun.SRegion, args *DiskDeleteOptions) error { + e := cli.DeleteDisk(args.Instance) + if e != nil { + return e + } + return nil + }) } diff --git a/pkg/util/aliyun/shell/instance.go b/pkg/util/aliyun/shell/instance.go index 495e2949cf..4a741d0b66 100644 --- a/pkg/util/aliyun/shell/instance.go +++ b/pkg/util/aliyun/shell/instance.go @@ -82,6 +82,63 @@ func init() { return nil }) + /* + server-change-config 更改系统配置 + server-reset + */ + type InstanceDeployOptions struct { + ID string `help:"instance ID"` + Name string `help:"new instance name"` + Hostname string `help:"new hostname"` + Keypair string `help:"Keypair Name"` + DeleteKeypair bool `help:"Remove SSH keypair"` + Password string `help:"new password"` + ResetPassword bool `help:"Force reset password"` + Description string `help:"new instances description"` + } + + shellutils.R(&InstanceDeployOptions{}, "instance-deploy", "Deploy keypair/password to a stopped virtual server", func(cli *aliyun.SRegion, args *InstanceDeployOptions) error { + err := cli.DeployVM(args.ID, args.Name, args.Password, args.Keypair, args.ResetPassword, args.DeleteKeypair, args.Description) + if err != nil { + return err + } + return nil + }) + + type InstanceRebuildRootOptions struct { + ID string `help:"instance ID"` + Image string `help:"Image ID"` + } + + shellutils.R(&InstanceRebuildRootOptions{}, "instance-rebuild-root", "Reinstall virtual server system image", func(cli *aliyun.SRegion, args *InstanceRebuildRootOptions) error { + err := cli.ReplaceSystemDisk(args.ID, args.Image) + if err != nil { + return err + } + return nil + }) + + type InstanceChangeConfigOptions struct { + ID string `help:"instance ID"` + Ncpu int `help:"number of CPU"` + Vmem int `help:"MiB of memory"` + Disk []int `help:"Data disk sizes int GB"` + } + + shellutils.R(&InstanceChangeConfigOptions{}, "instance-change-config", "Deploy keypair/password to a stopped virtual server", func(cli *aliyun.SRegion, args *InstanceChangeConfigOptions) error { + instance, e := cli.GetInstance(args.ID) + if e != nil { + return e + } + + // todo : add create disks + err := cli.ChangeVMConfig(instance.ZoneId, args.ID, args.Ncpu, args.Vmem, nil) + if err != nil { + return err + } + return nil + }) + type InstanceUpdatePasswordOptions struct { ID string `help:"Instance ID"` PASSWD string `help:"new password"` diff --git a/pkg/util/aliyun/storage.go b/pkg/util/aliyun/storage.go index b54700b1e6..d8dec31739 100644 --- a/pkg/util/aliyun/storage.go +++ b/pkg/util/aliyun/storage.go @@ -15,6 +15,10 @@ type SStorage struct { storageType string } +func (self *SStorage) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SStorage) GetId() string { return fmt.Sprintf("%s-%s-%s", self.zone.region.client.providerId, self.zone.GetId(), self.storageType) } @@ -100,7 +104,7 @@ func (self *SStorage) GetIStoragecache() cloudprovider.ICloudStoragecache { } 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 diff --git a/pkg/util/aliyun/storagecache.go b/pkg/util/aliyun/storagecache.go index 61394948c2..fab2072253 100644 --- a/pkg/util/aliyun/storagecache.go +++ b/pkg/util/aliyun/storagecache.go @@ -5,6 +5,7 @@ import ( "strings" "time" + "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" @@ -19,6 +20,10 @@ type SStoragecache struct { iimages []cloudprovider.ICloudImage } +func (self *SStoragecache) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SStoragecache) GetId() string { return fmt.Sprintf("%s-%s", self.region.client.providerId, self.region.GetId()) } diff --git a/pkg/util/aliyun/vpc.go b/pkg/util/aliyun/vpc.go index 9751039c39..ea645ff2be 100644 --- a/pkg/util/aliyun/vpc.go +++ b/pkg/util/aliyun/vpc.go @@ -45,6 +45,10 @@ type SVpc struct { VpcName string } +func (self *SVpc) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SVpc) GetId() string { return self.VpcId } @@ -184,6 +188,19 @@ func (self *SVpc) GetManagerId() string { } func (self *SVpc) Delete() error { + err := self.fetchSecurityGroups() + if err != nil { + log.Errorf("fetchSecurityGroup for VPC delete fail %s", err) + return err + } + for i := 0; i < len(self.secgroups); i += 1 { + secgroup := self.secgroups[i].(*SSecurityGroup) + err := self.region.deleteSecurityGroup(secgroup.SecurityGroupId) + if err != nil { + log.Errorf("deleteSecurityGroup for VPC delete fail %s", err) + return err + } + } return self.region.DeleteVpc(self.VpcId) } diff --git a/pkg/util/aliyun/vswitch.go b/pkg/util/aliyun/vswitch.go index 9c9e7bbe84..bef2c3d74d 100644 --- a/pkg/util/aliyun/vswitch.go +++ b/pkg/util/aliyun/vswitch.go @@ -34,6 +34,10 @@ type SVSwitch struct { ZoneId string } +func (self *SVSwitch) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SVSwitch) GetId() string { return self.VSwitchId } diff --git a/pkg/util/aliyun/wire.go b/pkg/util/aliyun/wire.go index 903dcc681a..12793359c1 100644 --- a/pkg/util/aliyun/wire.go +++ b/pkg/util/aliyun/wire.go @@ -3,6 +3,7 @@ package aliyun import ( "fmt" + "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" ) @@ -14,6 +15,10 @@ type SWire struct { inetworks []cloudprovider.ICloudNetwork } +func (self *SWire) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SWire) GetId() string { return fmt.Sprintf("%s-%s", self.vpc.GetId(), self.zone.GetId()) } diff --git a/pkg/util/aliyun/zone.go b/pkg/util/aliyun/zone.go index 844b4d6e6f..2241e4d9ce 100644 --- a/pkg/util/aliyun/zone.go +++ b/pkg/util/aliyun/zone.go @@ -3,6 +3,7 @@ package aliyun import ( "fmt" + "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" @@ -111,6 +112,10 @@ type SZone struct { AvailableDedicatedHostTypes SDedicatedHostTypes } +func (self *SZone) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SZone) GetId() string { return self.ZoneId } diff --git a/pkg/util/esxi/datacenter.go b/pkg/util/esxi/datacenter.go index 7da30e1dfb..6c9d1c6d0d 100644 --- a/pkg/util/esxi/datacenter.go +++ b/pkg/util/esxi/datacenter.go @@ -76,3 +76,4 @@ func (dc *SDatacenter) GetIStorages() ([]cloudprovider.ICloudStorage, error) { } return dc.istorages, nil } + diff --git a/pkg/util/esxi/host.go b/pkg/util/esxi/host.go index 4f524948d2..3a9a053c5d 100644 --- a/pkg/util/esxi/host.go +++ b/pkg/util/esxi/host.go @@ -68,6 +68,10 @@ func NewHost(manager *SESXiClient, host *mo.HostSystem, dc *SDatacenter) *SHost return &SHost{SManagedObject: newManagedObject(manager, host, dc)} } +func (self *SHost) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SHost) getHostSystem() *mo.HostSystem { return self.object.(*mo.HostSystem) } diff --git a/pkg/util/esxi/storage.go b/pkg/util/esxi/storage.go index 90846ed5a9..b64c57ce1f 100644 --- a/pkg/util/esxi/storage.go +++ b/pkg/util/esxi/storage.go @@ -18,6 +18,10 @@ func NewDatastore(manager *SESXiClient, ds *mo.Datastore, dc *SDatacenter) *SDat return &SDatastore{SManagedObject: newManagedObject(manager, ds, dc)} } +func (self *SDatastore) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SDatastore) getDatastore() *mo.Datastore { return self.object.(*mo.Datastore) } diff --git a/pkg/util/esxi/virtualmachine.go b/pkg/util/esxi/virtualmachine.go index 2005bbc9bc..34c1a9c682 100644 --- a/pkg/util/esxi/virtualmachine.go +++ b/pkg/util/esxi/virtualmachine.go @@ -27,6 +27,10 @@ func NewVirtualMachine(manager *SESXiClient, vm *mo.VirtualMachine, dc *SDatacen return &SVirtualMachine{SManagedObject: newManagedObject(manager, vm, dc), host: host} } +func (self *SVirtualMachine) GetMetadata() *jsonutils.JSONDict { + return nil +} + func (self *SVirtualMachine) getVirtualMachine() *mo.VirtualMachine { return self.object.(*mo.VirtualMachine) } @@ -65,6 +69,22 @@ func (self *SVirtualMachine) IsEmulated() bool { return false } +func (self *SVirtualMachine) DeployVM(name string, password string, publicKey string, resetPassword bool, deleteKeypair bool, description string) error { + return cloudprovider.ErrNotImplemented +} + +func (self *SVirtualMachine) RebuildRoot(imageId string) error { + return cloudprovider.ErrNotImplemented +} + +func (self *SVirtualMachine) UpdateVM(name string) error { + return cloudprovider.ErrNotImplemented +} + +func (self *SVirtualMachine) AttachDisk(diskId string) error { + return cloudprovider.ErrNotImplemented +} + func (self *SVirtualMachine) getUuid() string { return self.getVirtualMachine().Summary.Config.Uuid } @@ -209,3 +229,7 @@ func (self *SVirtualMachine) acquireVmrcUrl() (jsonutils.JSONObject, error) { ret.Add(jsonutils.NewString(url), "url") return ret, nil } + +func (dc *SVirtualMachine) ChangeConfig(instanceId string, ncpu int, vmem int) error { + return cloudprovider.ErrNotImplemented +} diff --git a/pkg/util/httputils/httputils.go b/pkg/util/httputils/httputils.go index 603f99987a..efd7cfecd5 100644 --- a/pkg/util/httputils/httputils.go +++ b/pkg/util/httputils/httputils.go @@ -32,10 +32,16 @@ var ( cyan = color.New(color.FgHiCyan, color.Bold).PrintlnFunc() ) +type Error struct { + Id string + Fields []string +} + type JSONClientError struct { Code int Class string Details string + Data Error } func (e *JSONClientError) Error() string { diff --git a/vendor/yunion.io/x/log/hooks/file.go b/vendor/yunion.io/x/log/hooks/file.go new file mode 100644 index 0000000000..3a99b0397a --- /dev/null +++ b/vendor/yunion.io/x/log/hooks/file.go @@ -0,0 +1,108 @@ +package hooks + +import ( + "fmt" + "os" + "path/filepath" + + "github.com/sirupsen/logrus" +) + +type LogFileHook struct { + FileDir string + FileName string + fullPath string + file *os.File + written int64 +} + +func (h *LogFileHook) Init() error { + if fi, err := os.Lstat(h.FileDir); err != nil { + if os.IsNotExist(err) { + os.MkdirAll(h.FileDir, 0755) + } else { + return fmt.Errorf("Lstat %s: %s", h.FileDir, err) + } + } else if !fi.Mode().IsDir() { + return fmt.Errorf("%s exists and it's not a directory", h.FileDir) + } + + h.fullPath = filepath.Join(h.FileDir, h.FileName) + file, err := os.OpenFile(h.fullPath, os.O_WRONLY|os.O_APPEND|os.O_CREATE, 0755) + if err != nil { + return fmt.Errorf("OpenFile %s: %s", h.fullPath, err) + } + h.file = file + h.written = 0 + return nil +} + +func (h *LogFileHook) DeInit() { + if h.file != nil { + h.file.Close() + } +} + +func (h *LogFileHook) Levels() []logrus.Level { + return logrus.AllLevels +} + +func (h *LogFileHook) Fire(e *logrus.Entry) error { + if b, err := e.Logger.Formatter.Format(e); err != nil { + return err + } else { + n, err := h.file.Write(b) + h.written += int64(n) + return err + } + return nil +} + +func (h *LogFileHook) Written() int64 { + return h.written +} + +// rotate by size +type LogFileRotateHook struct { + LogFileHook + RotateNum int + RotateSize int64 + filePaths []string +} + +func (h *LogFileRotateHook) Init() error { + if err := h.LogFileHook.Init(); err != nil { + return err + } + h.filePaths = make([]string, h.RotateNum) + for i := 1; i < h.RotateNum; i++ { + fileName := fmt.Sprintf("%s.%d", h.FileName, i) + filePath := filepath.Join(h.FileDir, fileName) + h.filePaths[i] = filePath + } + h.filePaths[0] = filepath.Join(h.FileDir, h.FileName) + return nil +} + +func (h *LogFileRotateHook) rotate() { + for i := h.RotateNum - 1; i > 0; i-- { + filePath0 := h.filePaths[i-1] + if _, err := os.Lstat(filePath0); err != nil { + continue + } + filePath1 := h.filePaths[i] + os.Rename(filePath0, filePath1) + } + h.LogFileHook.DeInit() + h.LogFileHook.Init() +} + +func (h *LogFileRotateHook) Fire(e *logrus.Entry) error { + if err := h.LogFileHook.Fire(e); err != nil { + return err + } + if h.LogFileHook.Written() >= h.RotateSize { + h.rotate() + } + return nil +} diff --git a/vendor/yunion.io/x/log/hooks/stdio.go b/vendor/yunion.io/x/log/hooks/stdio.go index 6fa3dc4381..e5d32000ad 100644 --- a/vendor/yunion.io/x/log/hooks/stdio.go +++ b/vendor/yunion.io/x/log/hooks/stdio.go @@ -12,14 +12,9 @@ type StdioHook struct{} func (hook *StdioHook) Fire(entry *logrus.Entry) error { line, err := entry.String() if err != nil { - fmt.Fprintf(os.Stderr, "Unable to read entry, %v", err) return err } - if entry.Level >= logrus.ErrorLevel { - fmt.Fprintf(os.Stderr, line) - } else { - fmt.Fprintf(os.Stdout, line) - } + fmt.Fprintf(os.Stderr, line) return nil } diff --git a/vendor/yunion.io/x/pkg/gotypes/gotypes.go b/vendor/yunion.io/x/pkg/gotypes/gotypes.go index e589aa7c61..615499fa4a 100644 --- a/vendor/yunion.io/x/pkg/gotypes/gotypes.go +++ b/vendor/yunion.io/x/pkg/gotypes/gotypes.go @@ -120,6 +120,15 @@ func ParseValue(val string, tp reflect.Type) (reflect.Value, error) { } case reflect.String: return reflect.ValueOf(val), nil + case reflect.Ptr: + tpElem := tp.Elem() + rv, err := ParseValue(val, tpElem) + if err != nil { + return reflect.ValueOf(val), fmt.Errorf("Cannot parse %s to %s", val, tp) + } + rvv := reflect.New(tpElem) + rvv.Elem().Set(rv) + return rvv, nil default: if tp == TimeType { tm, e := timeutils.ParseTimeStr(val) @@ -173,7 +182,16 @@ func SetValue(value reflect.Value, valStr string) error { Float32SliceType, Float64SliceType, StringSliceType: reflect.Append(value, reflect.ValueOf(valStr)) default: - return fmt.Errorf("Unsupported type: %v", value.Type()) + if value.Kind() == reflect.Ptr && value.Elem().Kind() != reflect.Slice { + newVal := reflect.New(value.Type().Elem()) + newValElem := newVal.Elem() + if err := SetValue(newValElem, valStr); err != nil { + return err + } + value.Set(newVal) + } else { + return fmt.Errorf("Unsupported type: %v", value.Type()) + } } return nil } diff --git a/vendor/yunion.io/x/pkg/util/filterclause/filterclause.go b/vendor/yunion.io/x/pkg/util/filterclause/filterclause.go index 2a3a1ccd10..a0928f6527 100644 --- a/vendor/yunion.io/x/pkg/util/filterclause/filterclause.go +++ b/vendor/yunion.io/x/pkg/util/filterclause/filterclause.go @@ -16,6 +16,20 @@ type SFilterClause struct { params []string } +type SJointFilterClause struct { + SFilterClause + JointModel string + RelatedKey string +} + +func (jfc *SJointFilterClause) GetJointFilter(q *sqlchemy.SQuery) sqlchemy.ICondition { + return jfc.QueryCondition(q) +} + +func (jfc *SJointFilterClause) GetJointModelName() string { + return jfc.JointModel[:len(jfc.JointModel)-1] +} + func (fc *SFilterClause) QueryCondition(q *sqlchemy.SQuery) sqlchemy.ICondition { field := q.Field(fc.field) if field == nil { @@ -63,11 +77,13 @@ func (fc *SFilterClause) String() string { } var ( - filterClausePattern *regexp.Regexp + filterClausePattern *regexp.Regexp + jointFilterClausePattern *regexp.Regexp ) func init() { filterClausePattern = regexp.MustCompile(`^(\w+)\.(\w+)\((.*)\)`) + jointFilterClausePattern = regexp.MustCompile(`^(\w+)\((\w+)\).(\w+)\.(\w+)\((.*)\)`) } func ParseFilterClause(filter string) *SFilterClause { @@ -79,3 +95,21 @@ func ParseFilterClause(filter string) *SFilterClause { fc := SFilterClause{field: matches[1], funcName: matches[2], params: params} return &fc } + +func ParseJointFilterClause(jointFilter string) *SJointFilterClause { + matches := jointFilterClausePattern.FindStringSubmatch(jointFilter) + if matches == nil { + return nil + } + params := utils.FindWords([]byte(matches[5]), 0) + jfc := SJointFilterClause{ + SFilterClause: SFilterClause{ + field: matches[3], + funcName: matches[4], + params: params, + }, + JointModel: matches[1], + RelatedKey: matches[2], + } + return &jfc +} diff --git a/vendor/yunion.io/x/pkg/util/sysutils/storagetypes.go b/vendor/yunion.io/x/pkg/util/sysutils/storagetypes.go index 476aa53bff..16baa403d2 100644 --- a/vendor/yunion.io/x/pkg/util/sysutils/storagetypes.go +++ b/vendor/yunion.io/x/pkg/util/sysutils/storagetypes.go @@ -9,12 +9,18 @@ const ( STORAGE_NAS = "nas" STORAGE_VSAN = "vsan" + STORAGE_CLOUD = "cloud" STORAGE_CLOUD_SSD = "cloud_ssd" STORAGE_CLOUD_EFFICIENCY = "cloud_efficiency" + + STORAGE_STANDARD = "standard" //Azure hdd storage type + STORAGE_PREMIUM = "premium" //Azure ssd storage type ) var STORAGE_TYPES = []string{STORAGE_LOCAL, STORAGE_BAREMETAL, STORAGE_SHEEPDOG, - STORAGE_RBD, STORAGE_DOCKER, STORAGE_NAS, STORAGE_VSAN, STORAGE_CLOUD_SSD, STORAGE_CLOUD_EFFICIENCY} + STORAGE_RBD, STORAGE_DOCKER, STORAGE_NAS, STORAGE_VSAN, + STORAGE_CLOUD, STORAGE_CLOUD_SSD, STORAGE_CLOUD_EFFICIENCY, + STORAGE_STANDARD, STORAGE_PREMIUM} var LOCAL_STORAGE_TYPES = []string{STORAGE_LOCAL, STORAGE_BAREMETAL} diff --git a/vendor/yunion.io/x/sqlchemy/field_update.go b/vendor/yunion.io/x/sqlchemy/field_update.go index 2087dbb6c0..b6b32203a5 100644 --- a/vendor/yunion.io/x/sqlchemy/field_update.go +++ b/vendor/yunion.io/x/sqlchemy/field_update.go @@ -102,7 +102,7 @@ func (ts *STableSpec) updateFields(dt interface{}, fields map[string]interface{} buf.WriteString(fmt.Sprintf(", `%s` = `%s` + 1", versionField, versionField)) } for _, updatedField := range updatedFields { - buf.WriteString(fmt.Sprintf(", `%s` = NOW()", updatedField)) + buf.WriteString(fmt.Sprintf(", `%s` = UTC_TIMESTAMP()", updatedField)) } buf.WriteString(" WHERE ") first = true diff --git a/vendor/yunion.io/x/sqlchemy/insert.go b/vendor/yunion.io/x/sqlchemy/insert.go index d9f10406c1..c7442588f0 100644 --- a/vendor/yunion.io/x/sqlchemy/insert.go +++ b/vendor/yunion.io/x/sqlchemy/insert.go @@ -36,7 +36,7 @@ func (t *STableSpec) insertSqlPrep(dataFields map[string]interface{}) (string, [ if ok && (dtc.IsCreatedAt || dtc.IsUpdatedAt) { createdAtFields = append(createdAtFields, k) names = append(names, fmt.Sprintf("`%s`", k)) - format = append(format, "NOW()") + format = append(format, "UTC_TIMESTAMP()") } else if ov != nil && !c.IsZero(ov) && !isAutoInc { v := c.ConvertFromValue(ov) values = append(values, v) diff --git a/vendor/yunion.io/x/sqlchemy/update.go b/vendor/yunion.io/x/sqlchemy/update.go index d99ef0b95e..e40e0f6eff 100644 --- a/vendor/yunion.io/x/sqlchemy/update.go +++ b/vendor/yunion.io/x/sqlchemy/update.go @@ -126,7 +126,7 @@ func (us *SUpdateSession) saveUpdate(dt interface{}) (map[string]SUpdateDiff, er buf.WriteString(fmt.Sprintf(", `%s` = `%s` + 1", versionField, versionField)) } for _, updatedField := range updatedFields { - buf.WriteString(fmt.Sprintf(", `%s` = NOW()", updatedField)) + buf.WriteString(fmt.Sprintf(", `%s` = UTC_TIMESTAMP()", updatedField)) } buf.WriteString(" WHERE ") first = true diff --git a/vendor/yunion.io/x/structarg/structarg.go b/vendor/yunion.io/x/structarg/structarg.go index 9eb223f272..7f7f6c6414 100644 --- a/vendor/yunion.io/x/structarg/structarg.go +++ b/vendor/yunion.io/x/structarg/structarg.go @@ -305,6 +305,9 @@ func (this *ArgumentParser) addArgument(f reflect.StructField, v reflect.Value) return fmt.Errorf("positional %s must not have default value", token) } } + if !positional && use_default && required { + return fmt.Errorf("non-positional argument with default value should not have required:true set") + } subcommand, err := strconv.ParseBool(tagMap[TAG_SUBCOMMAND]) if err != nil { subcommand = false