merge release/2.1.0

This commit is contained in:
Zhang Dongliang
2018-08-29 15:52:53 +08:00
72 changed files with 1824 additions and 510 deletions
Generated
+16 -11
View File
@@ -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",
+3 -4
View File
@@ -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)
%:
@:
+2
View File
@@ -0,0 +1,2 @@
DESCRIPTION="Yunion Cloud Region DNS Service"
# SERVICE="yes"
+2 -2
View File
@@ -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")
+14 -2
View File
@@ -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
+1 -2
View File
@@ -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["<action>"], query, data)
+31 -4
View File
@@ -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)
+12 -5
View File
@@ -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)
}
+3 -3
View File
@@ -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)
+8 -1
View File
@@ -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 {
+265 -75
View File
@@ -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
}
+20 -2
View File
@@ -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 {
+1 -1
View File
@@ -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")
+22
View File
@@ -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"
)
+67 -13
View File
@@ -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")
-2
View File
@@ -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
+29 -11
View File
@@ -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)
+181 -18
View File
@@ -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
+1 -1
View File
@@ -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
+14 -5
View File
@@ -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
}
+52
View File
@@ -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
+18
View File
@@ -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")
+1 -1
View File
@@ -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 {
+25
View File
@@ -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)
+1
View File
@@ -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
}
+2 -2
View File
@@ -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()))
}
}
+4 -3
View File
@@ -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)
}
+109
View File
@@ -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{})
}
+11 -5
View File
@@ -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)
}
+1 -5
View File
@@ -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())
+1
View File
@@ -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)
+1 -1
View File
@@ -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 {
+56
View File
@@ -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
}
+80 -87
View File
@@ -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 {
+27 -22
View File
@@ -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 {
+2 -10
View File
@@ -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 {
+3 -6
View File
@@ -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)
+74 -120
View File
@@ -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 {
+30 -26
View File
@@ -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{}) {
+1 -1
View File
@@ -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
}
+1 -6
View File
@@ -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)
+1
View File
@@ -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"
+31 -2
View File
@@ -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
}
+13
View File
@@ -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
}
+6 -2
View File
@@ -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
}
+4
View File
@@ -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
}
+202 -36
View File
@@ -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 {
+31
View File
@@ -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
}
+4
View File
@@ -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)
+16
View File
@@ -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
}
+11
View File
@@ -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
})
}
+57
View File
@@ -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"`
+5 -1
View File
@@ -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
+5
View File
@@ -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())
}
+17
View File
@@ -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)
}
+4
View File
@@ -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
}
+5
View File
@@ -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())
}
+5
View File
@@ -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
}
+1
View File
@@ -76,3 +76,4 @@ func (dc *SDatacenter) GetIStorages() ([]cloudprovider.ICloudStorage, error) {
}
return dc.istorages, nil
}
+4
View File
@@ -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)
}
+4
View File
@@ -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)
}
+24
View File
@@ -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
}
+6
View File
@@ -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 {
+108
View File
@@ -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
}
+1 -6
View File
@@ -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
}
+19 -1
View File
@@ -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
}
+35 -1
View File
@@ -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
}
+7 -1
View File
@@ -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}
+1 -1
View File
@@ -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
+1 -1
View File
@@ -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)
+1 -1
View File
@@ -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
+3
View File
@@ -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