diff --git a/cmd/climc/shell/cloudproviders.go b/cmd/climc/shell/cloudproviders.go index b8f0019741..f8eb5bf518 100644 --- a/cmd/climc/shell/cloudproviders.go +++ b/cmd/climc/shell/cloudproviders.go @@ -33,10 +33,10 @@ func init() { type CloudproviderCreateOptions struct { NAME string `help:"Name of cloud provider"` - ACCOUNT string `help:"Account to access the cloud provider"` - SECRET string `help:"Secret to access the cloud provider, clientId/clientScret/subscriptionId for Azure"` + ACCOUNT string `help:"Account to access the cloud provider, tenantId/subscriptionId for Azure"` + SECRET string `help:"Secret to access the cloud provider, clientId/clientScret for Azure"` PROVIDER string `help:"Driver for cloud provider" choices:"VMware|Aliyun|Azure"` - AccessURL string `helo:"hello" metavar:"Azure choices: "` + AccessURL string `helo:"hello" metavar:"Azure choices: "` Desc string `help:"Description"` Enabled bool `help:"Enabled the provider automatically"` } diff --git a/cmd/climc/shell/k8s/deployment.go b/cmd/climc/shell/k8s/deployment.go index ee001943b8..21e8338e7d 100644 --- a/cmd/climc/shell/k8s/deployment.go +++ b/cmd/climc/shell/k8s/deployment.go @@ -2,6 +2,7 @@ package k8s import ( "fmt" + "io/ioutil" "strconv" "strings" @@ -105,6 +106,31 @@ func initDeployment() { printObjectYAML(ret) return nil }) + + type createFromFileOpt struct { + resourceGetOptions + FILE string `help:"K8s resource YAML or JSON file"` + } + R(&createFromFileOpt{}, "k8s-create", "Create resource by file", func(s *mcclient.ClientSession, args *createFromFileOpt) error { + params := args.ClusterParams() + params.Add(jsonutils.NewString(args.NAME), "name") + + content, err := ioutil.ReadFile(args.FILE) + if err != nil { + return err + } + namespace := args.Namespace + if namespace != "" { + params.Add(jsonutils.NewString(namespace), "namespace") + } + params.Add(jsonutils.NewString(string(content)), "content") + ret, err := k8s.DeployFromFile.Create(s, params) + if err != nil { + return err + } + printObjectYAML(ret) + return nil + }) } type portMapping struct { diff --git a/cmd/climc/shell/k8s/raw.go b/cmd/climc/shell/k8s/raw.go index ac4d6daf52..cebcd968a8 100644 --- a/cmd/climc/shell/k8s/raw.go +++ b/cmd/climc/shell/k8s/raw.go @@ -57,7 +57,10 @@ func initRaw() { if err != nil { return err } - body := jsonutils.Marshal(string(content)) + body, err := jsonutils.Parse(content) + if err != nil { + return err + } err = k8s.RawResource.Put(s, args.KIND, args.Namespace, args.NAME, body, args.Cluster) if err != nil { return err diff --git a/cmd/climc/shell/usages.go b/cmd/climc/shell/usages.go index 6a909cf124..41aa649f85 100644 --- a/cmd/climc/shell/usages.go +++ b/cmd/climc/shell/usages.go @@ -80,4 +80,16 @@ func init() { printObject(result) return nil }) + + R(&ResourceUsageOptions{}, "cloud-region-usage", "Show general usage of a cloud region", func(s *mcclient.ClientSession, args *ResourceUsageOptions) error { + params := fetchHostTypeOptions(&args.GeneralUsageOptions) + params.Add(jsonutils.NewString("cloudregions"), "range_type") + params.Add(jsonutils.NewString(args.ID), "range_id") + result, err := modules.Usages.GetGeneralUsage(s, params) + if err != nil { + return err + } + printObject(result) + return nil + }) } diff --git a/pkg/appsrv/radix.go b/pkg/appsrv/radix.go index 380bee48c1..67ddd20a95 100644 --- a/pkg/appsrv/radix.go +++ b/pkg/appsrv/radix.go @@ -3,24 +3,26 @@ package appsrv import ( "fmt" "strings" + ) type RadixNode struct { data interface{} + fullPath []string next []*RadixNode parent *RadixNode matchNext *RadixNode - matchTable []string + // matchTable []string segment string } func NewRadix() *RadixNode { return &RadixNode{data: nil, - next: make([]*RadixNode, 0), - matchNext: nil, - matchTable: nil, - parent: nil, - segment: ""} + fullPath: nil, + next: make([]*RadixNode, 0), + matchNext: nil, + parent: nil, + segment: ""} } func (r *RadixNode) String() string { @@ -50,11 +52,21 @@ func isMatchSegment(seg string) bool { } func (r *RadixNode) Add(segments []string, data interface{}) error { + return r.add(segments, segments, data) +} + +func (r *RadixNode) add(path []string, segments []string, data interface{}) error { + // log.Debugf("add %#v %#v", path, segments) + if len(segments) == 0 { if r.data != nil { return fmt.Errorf("Duplicate data for node %s", r.String()) } else { r.data = data + r.fullPath = make([]string, len(path)) + for i := 0; i < len(path); i += 1 { + r.fullPath[i] = path[i] + } return nil } } @@ -65,18 +77,18 @@ func (r *RadixNode) Add(segments []string, data interface{}) error { return fmt.Errorf("%s has been registered, %s conflict with %s", r.matchNext.String(), r.matchNext.segment, segments[0]) } */ nextNode = r.matchNext - nextNode.matchTable = append(nextNode.matchTable, segments[0]) + // nextNode.matchTable = append(nextNode.matchTable, segments[0]) } else { nextNode = NewRadix() nextNode.segment = "<*>" nextNode.parent = r - nextNode.matchTable = []string{segments[0]} + // nextNode.matchTable = []string{segments[0]} r.matchNext = nextNode } } else { - for _, node := range r.next { - if node.segment == segments[0] { - nextNode = node + for i := 0; i < len(r.next); i += 1 { + if r.next[i].segment == segments[0] { + nextNode = r.next[i] break } } @@ -87,42 +99,46 @@ func (r *RadixNode) Add(segments []string, data interface{}) error { r.next = append(r.next, nextNode) } } - return nextNode.Add(segments[1:], data) + return nextNode.add(path, segments[1:], data) } func (r *RadixNode) Match(segments []string, params map[string]string) interface{} { + data, allPaths := r.match(segments) + // log.Debugf("%#v", allPaths) + for i := 0; i < len(segments); i += 1 { + for j := 0; j < len(allPaths); j += 1 { + if i < len(allPaths[j]) && isMatchSegment(allPaths[j][i]) { + params[allPaths[j][i]] = segments[i] + } + } + } + return data +} + +func (r *RadixNode) match(segments []string) (interface{}, [][]string) { if len(segments) == 0 { - return r.data + return r.data, r.getAllFullPaths() } else { - var ret interface{} = nil + var retData interface{} = nil + var retPath [][]string = nil exactMatch := false for _, node := range r.next { if node.segment == segments[0] { - ret = node.Match(segments[1:], params) - if ret != nil { - // log.Debugf("Match %s ret %#v", node.segment, ret) + retData, retPath = node.match(segments[1:]) + if retData != nil { exactMatch = true - } else { - // log.Debugf("No match %s ret %#v", node.segment, ret) } break } } - if ret != nil { - return ret + if retData != nil { + return retData, retPath } else { if !exactMatch && r.matchNext != nil { - ret = r.matchNext.Match(segments[1:], params) - if ret != nil { - for _, segname := range r.matchNext.matchTable { - if _, ok := params[segname]; !ok { - params[segname] = segments[0] - } - } - } - return ret + retData, retPath = r.matchNext.match(segments[1:]) + return retData, retPath } else { - return r.data + return r.data, r.getAllFullPaths() } } } @@ -139,3 +155,20 @@ func (r *RadixNode) Walk(f func(path string, data interface{})) { r.matchNext.Walk(f) } } + +func (r *RadixNode) getAllFullPaths() [][]string { + if r.fullPath != nil { + return [][]string{r.fullPath} + }else { + ret := make([][]string, 0) + for _, node := range r.next { + fp := node.getAllFullPaths() + ret = append(ret, fp...) + } + if r.matchNext != nil { + fp := r.matchNext.getAllFullPaths() + ret = append(ret, fp...) + } + return ret + } +} diff --git a/pkg/appsrv/radix_test.go b/pkg/appsrv/radix_test.go index fe978ceec3..ef87a46398 100644 --- a/pkg/appsrv/radix_test.go +++ b/pkg/appsrv/radix_test.go @@ -61,10 +61,16 @@ func TestRadixNode(t *testing.T) { func TestParams(t *testing.T) { r := NewRadix() - r.Add([]string{"POST", "clouds", ""}, "classAction") - r.Add([]string{"POST", "clouds", "", ""}, "objectAction") + r.Add([]string{"POST", "clouds", ""}, "classAction") + r.Add([]string{"POST", "clouds", "", "sync"}, "objectSyncAction") + r.Add([]string{"POST", "clouds", "", ""}, "objectAction") params := make(map[string]string) - ret := r.Match([]string{"POST", "clouds", "id", "sync"}, params) + ret := r.Match([]string{"POST", "clouds", "myid", "sync"}, params) + t.Logf("match: %s", ret) + t.Logf("params: %s", params) + + params = make(map[string]string) + ret = r.Match([]string{"POST", "clouds", "myid", "start"}, params) t.Logf("match: %s", ret) t.Logf("params: %s", params) diff --git a/pkg/cloudcommon/database.go b/pkg/cloudcommon/database.go index fd2886a742..9cc03cb4f6 100644 --- a/pkg/cloudcommon/database.go +++ b/pkg/cloudcommon/database.go @@ -20,8 +20,8 @@ func InitDB(options *DBOptions) { } sqlchemy.SetDB(dbConn) - // lm := lockman.NewInMemoryLockManager() - lm := lockman.NewNoopLockManager() + lm := lockman.NewInMemoryLockManager() + // lm := lockman.NewNoopLockManager() lockman.Init(lm) } diff --git a/pkg/cloudcommon/db/db_dispatcher.go b/pkg/cloudcommon/db/db_dispatcher.go index a7b4d99f9a..0198aaa062 100644 --- a/pkg/cloudcommon/db/db_dispatcher.go +++ b/pkg/cloudcommon/db/db_dispatcher.go @@ -1105,7 +1105,9 @@ func (dispatcher *DBModelDispatcher) Delete(ctx context.Context, idstr string, q return nil, httperrors.NewGeneralError(err) } log.Debugf("Delete %s", model.GetShortDesc()) + lockman.LockObject(ctx, model) defer lockman.ReleaseObject(ctx, model) + return deleteItem(dispatcher.modelManager, model, ctx, userCred, query, data) } diff --git a/pkg/cloudcommon/db/fetch.go b/pkg/cloudcommon/db/fetch.go index 05f8f49f13..0297e57216 100644 --- a/pkg/cloudcommon/db/fetch.go +++ b/pkg/cloudcommon/db/fetch.go @@ -31,10 +31,10 @@ func fetchById(manager IModelManager, idStr string) (IModel, error) { } } -func fetchByName(manager IModelManager, ownerProjId string, idStr string) (IModel, error) { +func fetchByName(manager IModelManager, owner string, idStr string) (IModel, error) { q := manager.Query() q = manager.FilterByName(q, idStr) - q = manager.FilterByOwner(q, ownerProjId) + q = manager.FilterByOwner(q, owner) count := q.Count() if count == 1 { obj, err := NewModelObject(manager) @@ -101,7 +101,7 @@ func fetchItemByName(manager IModelManager, ctx context.Context, userCred mcclie } } q = manager.FilterByName(q, idStr) - q = manager.FilterByOwner(q, userCred.GetProjectId()) + q = manager.FilterByOwner(q, manager.GetOwnerId(userCred)) count := q.Count() if count == 1 { item, err := NewModelObject(manager) diff --git a/pkg/cloudcommon/db/interface.go b/pkg/cloudcommon/db/interface.go index 71097f4614..bcc4c00779 100644 --- a/pkg/cloudcommon/db/interface.go +++ b/pkg/cloudcommon/db/interface.go @@ -37,7 +37,7 @@ type IModelManager interface { FilterById(q *sqlchemy.SQuery, idStr string) *sqlchemy.SQuery FilterByNotId(q *sqlchemy.SQuery, idStr string) *sqlchemy.SQuery FilterByName(q *sqlchemy.SQuery, name string) *sqlchemy.SQuery - FilterByOwner(q *sqlchemy.SQuery, ownerProjId string) *sqlchemy.SQuery + FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery GetOwnerId(userCred mcclient.TokenCredential) string diff --git a/pkg/cloudcommon/db/lockman/example/main.go b/pkg/cloudcommon/db/lockman/example/main.go new file mode 100644 index 0000000000..6623ac1e5d --- /dev/null +++ b/pkg/cloudcommon/db/lockman/example/main.go @@ -0,0 +1,60 @@ +package main + +import ( + "time" + "sync" + "context" + + "yunion.io/x/log" + "yunion.io/x/pkg/util/stringutils" + + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" + "math/rand" +) + + +type FakeObject struct { + Id string +} + +func (o *FakeObject) GetId() string { + return o.Id +} + +func (o *FakeObject) Keyword() string { + return "fake" +} + +func run(ctx context.Context, obj lockman.ILockedObject, id int, sleep time.Duration) { + log.Infof("ready to run at %d [%p]", id, ctx) + lockman.LockObject(ctx, obj) + defer lockman.ReleaseObject(ctx, obj) + log.Infof("Acquire obj at %d [%p]", id, ctx) + time.Sleep(sleep) + log.Infof("Release obj at %d [%p]", id, ctx) +} + +func main() { + lockman.Init(lockman.NewInMemoryLockManager()) + objId := stringutils.UUID4() + cycle := 10 + + var wg sync.WaitGroup + + log.Infof("Start") + + for id := 0; id <= 3; id += 1 { + wg.Add(1) + go func(localId int) { + log.Infof("Start %d", localId) + ctx := context.WithValue(context.Background(), "ID", localId) + for i := 0; i < cycle; i += 1 { + obj := &FakeObject{Id: objId} + run(ctx, obj, localId, time.Duration(rand.Intn(1000))*time.Millisecond) + } + wg.Done() + }(id) + } + + wg.Wait() +} \ No newline at end of file diff --git a/pkg/cloudcommon/db/lockman/fifo.go b/pkg/cloudcommon/db/lockman/fifo.go new file mode 100644 index 0000000000..1aba01fb5f --- /dev/null +++ b/pkg/cloudcommon/db/lockman/fifo.go @@ -0,0 +1,46 @@ +package lockman + +import ( + "container/list" + "yunion.io/x/log" +) + +type FIFO struct { + fifo *list.List +} + +func NewFIFO() *FIFO { + return &FIFO{fifo: list.New()} +} + +func (f *FIFO) Push(ele interface{}) { + f.fifo.PushBack(ele) +} + +func (f *FIFO) Pop(ele interface{}) interface{} { + e := f.fifo.Front() + for e != nil && e.Value != ele { + e = e.Next() + } + if e != nil { + v := f.fifo.Remove(e) + if v != ele { + log.Fatalf("remove element not identical!!") + } + } + return nil +} + +func (f *FIFO) Len() int { + return f.fifo.Len() +} + +type ElementInspectFunc func(ele interface{}) + +func (f *FIFO) Enum(eif ElementInspectFunc) { + e := f.fifo.Front() + for e != nil { + eif(e.Value) + e = e.Next() + } +} \ No newline at end of file diff --git a/pkg/cloudcommon/db/lockman/fifo_test.go b/pkg/cloudcommon/db/lockman/fifo_test.go new file mode 100644 index 0000000000..742f2723e2 --- /dev/null +++ b/pkg/cloudcommon/db/lockman/fifo_test.go @@ -0,0 +1,19 @@ +package lockman + +import "testing" + +func TestFIFO_Pop(t *testing.T) { + fifo := NewFIFO() + data := []int{1, 2, 3} + for i := 0; i < len(data); i += 1 { + fifo.Push(data[i]) + } + t.Logf("FIFO size: %d", fifo.Len()) + fifo.Enum(func(ele interface{}) { + t.Logf("FIFO ele %#v", ele) + }) + for i := 0; i < len(data); i += 1 { + fifo.Pop(data[i]) + } + t.Logf("FIFO size: %d", fifo.Len()) +} diff --git a/pkg/cloudcommon/db/lockman/inmemory.go b/pkg/cloudcommon/db/lockman/inmemory.go index c9cf2befc9..89ebe455f9 100644 --- a/pkg/cloudcommon/db/lockman/inmemory.go +++ b/pkg/cloudcommon/db/lockman/inmemory.go @@ -2,55 +2,93 @@ package lockman import ( "context" - "runtime/debug" "sync" "yunion.io/x/log" - "yunion.io/x/pkg/util/fifoutils" + "runtime/debug" ) -type SInMemoryLockOwner struct { +const ( + debug_log = false +) + +/*type SInMemoryLockOwner struct { owner context.Context - ready chan bool -} +}*/ type SInMemoryLockRecord struct { - lock *sync.Mutex - holder context.Context - counter int - queue *fifoutils.FIFO + lock *sync.Mutex + cond *sync.Cond + holder context.Context + depth int + waiter *FIFO } func newInMemoryLockRecord(ctx context.Context) *SInMemoryLockRecord { - rec := SInMemoryLockRecord{lock: &sync.Mutex{}, queue: fifoutils.NewFIFO(), holder: ctx, counter: 0} + lock := &sync.Mutex{} + cond := sync.NewCond(lock) + rec := SInMemoryLockRecord{lock: lock, cond: cond, holder: ctx, depth: 0, waiter: NewFIFO()} return &rec } -func (rec *SInMemoryLockRecord) lockContext(ctx context.Context) *SInMemoryLockOwner { +func (rec *SInMemoryLockRecord) lockContext(ctx context.Context) { rec.lock.Lock() defer rec.lock.Unlock() + if rec.holder == nil { + rec.holder = ctx + rec.depth = 1 + return + } + + if debug_log { + log.Debugf("rec.hold=[%p] ctx=[%p] %v", rec.holder, ctx, rec.holder==ctx) + } + if rec.holder == ctx { - rec.counter += 1 - log.Infof("lockContext: same ctx, counter: %d [%p]", rec.counter, rec.holder) - if rec.counter > 32 { - // MUST BE BUG + rec.depth += 1 + if debug_log { + log.Infof("lockContext: same ctx, depth: %d [%p]", rec.depth, rec.holder) + } + if rec.depth > 32 { + // XXX MUST BE BUG ??? debug.PrintStack() panic("Too many recursive locks!!!") } - return nil + return } - // check - for i := 0; i < rec.queue.Len(); i += 1 { - ele := rec.queue.ElementAt(i).(*SInMemoryLockOwner) - if ele.owner == ctx { - log.Fatalf("try to lock from a wait context????") - } - } - owner := SInMemoryLockOwner{owner: ctx, ready: make(chan bool)} - rec.queue.Push(&owner) - return &owner + // check + rec.waiter.Enum(func(ele interface{}) { + electx := ele.(context.Context) + if electx == ctx { + log.Fatalf("try to lock from a waiter context????") + } + }) + + rec.waiter.Push(ctx) + + if debug_log { + log.Debugf("waiter size %d after push", rec.waiter.Len()) + log.Debugf("Start to wait ... [%p]", ctx) + } + + for rec.holder != nil { + rec.cond.Wait() + } + + if debug_log { + log.Debugf("End of wait ... [%p]", ctx) + } + + rec.waiter.Pop(ctx) + + if debug_log { + log.Debugf("waiter size %d after pop", rec.waiter.Len()) + } + + rec.holder = ctx + rec.depth = 1 } func (rec *SInMemoryLockRecord) unlockContext(ctx context.Context) (needClean bool) { @@ -61,31 +99,27 @@ func (rec *SInMemoryLockRecord) unlockContext(ctx context.Context) (needClean bo log.Fatalf("try to unlock a wait context???") } - rec.counter -= 1 + if debug_log { + log.Debugf("unlockContext depth %d [%p]", rec.depth, ctx) + } - if rec.counter <= 0 { - if rec.queue.Len() == 0 { + rec.depth -= 1 + + if rec.depth <= 0 { + if debug_log { + log.Debugf("depth 0, to release lock for context [%p]", ctx) + } + + rec.holder = nil + if rec.waiter.Len() == 0 { return true } - newHolder := rec.queue.Pop().(*SInMemoryLockOwner) - rec.holder = newHolder.owner - rec.counter = 1 - newHolder.notify() + rec.cond.Signal() } return false } -func (owner *SInMemoryLockOwner) wait() { - // log.Infof("wait for notify %p", owner.owner) - <-owner.ready -} - -func (owner *SInMemoryLockOwner) notify() { - // log.Infof("notify %p", owner.owner) - owner.ready <- true -} - type SInMemoryLockManager struct { tableLock *sync.Mutex lockTable map[string]*SInMemoryLockRecord @@ -117,10 +151,7 @@ func (lockman *SInMemoryLockManager) getRecord(ctx context.Context, key string, func (lockman *SInMemoryLockManager) LockKey(ctx context.Context, key string) { record := lockman.getRecordWithLock(ctx, key) - owner := record.lockContext(ctx) - if owner != nil { - owner.wait() - } + record.lockContext(ctx) } func (lockman *SInMemoryLockManager) UnlockKey(ctx context.Context, key string) { @@ -137,4 +168,4 @@ func (lockman *SInMemoryLockManager) UnlockKey(ctx context.Context, key string) if needClean { delete(lockman.lockTable, key) } -} +} \ No newline at end of file diff --git a/pkg/cloudcommon/db/lockman/inmemory_test.go b/pkg/cloudcommon/db/lockman/inmemory_test.go index 27ee308556..ad7b89866a 100644 --- a/pkg/cloudcommon/db/lockman/inmemory_test.go +++ b/pkg/cloudcommon/db/lockman/inmemory_test.go @@ -24,10 +24,10 @@ func (o *FakeObject) Keyword() string { func run(t *testing.T, ctx context.Context, obj ILockedObject, id int, sleep time.Duration) { t.Logf("ready to run at %d [%p]", id, ctx) LockObject(ctx, obj) - defer LockObject(ctx, obj) - t.Logf("Acquire obj at %s [%p]", id, ctx) + defer ReleaseObject(ctx, obj) + t.Logf("Acquire obj at %d [%p]", id, ctx) time.Sleep(sleep) - t.Logf("Release obj at %s [%p]", id, ctx) + t.Logf("Release obj at %d [%p]", id, ctx) } func TestInMemoryLockManager(t *testing.T) { diff --git a/pkg/cloudcommon/db/modelbase.go b/pkg/cloudcommon/db/modelbase.go index 6b5f48f0aa..4e17ff01e9 100644 --- a/pkg/cloudcommon/db/modelbase.go +++ b/pkg/cloudcommon/db/modelbase.go @@ -104,7 +104,7 @@ func (manager *SModelBaseManager) FilterByName(q *sqlchemy.SQuery, name string) return q } -func (manager *SModelBaseManager) FilterByOwner(q *sqlchemy.SQuery, ownerProjId string) *sqlchemy.SQuery { +func (manager *SModelBaseManager) FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery { return q } diff --git a/pkg/cloudcommon/db/namevalidator.go b/pkg/cloudcommon/db/namevalidator.go index aef4032450..71e9703e9b 100644 --- a/pkg/cloudcommon/db/namevalidator.go +++ b/pkg/cloudcommon/db/namevalidator.go @@ -7,11 +7,11 @@ import ( "yunion.io/x/pkg/util/stringutils" ) -func isNameUnique(manager IModelManager, ownerProjId string, name string) bool { +func isNameUnique(manager IModelManager, owner string, name string) bool { q := manager.Query() q = manager.FilterByName(q, name) if !globalVirtualResourceNamespace { - q = manager.FilterByOwner(q, ownerProjId) + q = manager.FilterByOwner(q, owner) } return q.Count() == 0 } diff --git a/pkg/cloudcommon/db/opslog.go b/pkg/cloudcommon/db/opslog.go index 457908e27f..97c5f3b3ce 100644 --- a/pkg/cloudcommon/db/opslog.go +++ b/pkg/cloudcommon/db/opslog.go @@ -349,9 +349,9 @@ func (self *SOpsLogManager) FilterByName(q *sqlchemy.SQuery, name string) *sqlch return q } -func (self *SOpsLogManager) FilterByOwner(q *sqlchemy.SQuery, ownerProjId string) *sqlchemy.SQuery { - if len(ownerProjId) > 0 { - return q.Equals("owner_project_id", ownerProjId) +func (self *SOpsLogManager) FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery { + if len(owner) > 0 { + return q.Equals("owner_project_id", owner) } else { return q } diff --git a/pkg/cloudcommon/db/sharablevirtual.go b/pkg/cloudcommon/db/sharablevirtual.go index 33e8170198..62c38c442d 100644 --- a/pkg/cloudcommon/db/sharablevirtual.go +++ b/pkg/cloudcommon/db/sharablevirtual.go @@ -22,8 +22,8 @@ func NewSharableVirtualResourceBaseManager(dt interface{}, tableName string, key return SSharableVirtualResourceBaseManager{SVirtualResourceBaseManager: NewVirtualResourceBaseManager(dt, tableName, keyword, keywordPlural)} } -func (manager *SSharableVirtualResourceBaseManager) FilterByOwner(q *sqlchemy.SQuery, ownerProjId string) *sqlchemy.SQuery { - q = q.Filter(sqlchemy.OR(sqlchemy.Equals(q.Field("tenant_id"), ownerProjId), sqlchemy.IsTrue(q.Field("is_public")))) +func (manager *SSharableVirtualResourceBaseManager) FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery { + q = q.Filter(sqlchemy.OR(sqlchemy.Equals(q.Field("tenant_id"), owner), sqlchemy.IsTrue(q.Field("is_public")))) q = q.Filter(sqlchemy.OR(sqlchemy.IsNull(q.Field("pending_deleted")), sqlchemy.IsFalse(q.Field("pending_deleted")))) q = q.Filter(sqlchemy.OR(sqlchemy.IsNull(q.Field("is_system")), sqlchemy.IsFalse(q.Field("is_system")))) return q diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index f221ae6f11..26df5fb08f 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -87,7 +87,7 @@ func (manager *STaskManager) FilterByName(q *sqlchemy.SQuery, name string) *sqlc return q } -func (manager *STaskManager) FilterByOwner(q *sqlchemy.SQuery, ownerProjId string) *sqlchemy.SQuery { +func (manager *STaskManager) FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery { return q } diff --git a/pkg/cloudcommon/db/virtualresource.go b/pkg/cloudcommon/db/virtualresource.go index 8220398d3c..7270321b4b 100644 --- a/pkg/cloudcommon/db/virtualresource.go +++ b/pkg/cloudcommon/db/virtualresource.go @@ -49,8 +49,8 @@ func (model *SVirtualResourceBase) GetOwnerProjectId() string { return model.ProjectId } -func (manager *SVirtualResourceBaseManager) FilterByOwner(q *sqlchemy.SQuery, ownerProjId string) *sqlchemy.SQuery { - q = q.Equals("tenant_id", ownerProjId) +func (manager *SVirtualResourceBaseManager) FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery { + q = q.Equals("tenant_id", owner) q = q.Filter(sqlchemy.OR(sqlchemy.IsNull(q.Field("pending_deleted")), sqlchemy.IsFalse(q.Field("pending_deleted")))) q = q.Filter(sqlchemy.OR(sqlchemy.IsNull(q.Field("is_system")), sqlchemy.IsFalse(q.Field("is_system")))) return q @@ -282,8 +282,10 @@ func (model *SVirtualResourceBase) VirtualModelManager() IVirtualModelManager { func (model *SVirtualResourceBase) CancelPendingDelete(ctx context.Context, userCred mcclient.TokenCredential) error { ownerProjId := model.GetOwnerProjectId() + lockman.LockClass(ctx, model.GetModelManager(), ownerProjId) defer lockman.ReleaseClass(ctx, model.GetModelManager(), ownerProjId) + _, err := model.GetModelManager().TableSpec().Update(model, func() error { model.Name = GenerateName(model.GetModelManager(), ownerProjId, model.Name) model.PendingDeleted = false diff --git a/pkg/cloudcommon/validators/errors.go b/pkg/cloudcommon/validators/errors.go index 57c7271ece..9045fc93ed 100644 --- a/pkg/cloudcommon/validators/errors.go +++ b/pkg/cloudcommon/validators/errors.go @@ -1,9 +1,14 @@ package validators import ( + "database/sql" "fmt" + + "yunion.io/x/onecloud/pkg/httperrors" ) +var returnHttpError = true + type ErrType uintptr const ( @@ -21,7 +26,7 @@ const ( var errTypeToString = map[ErrType]string{ ERR_SUCCESS: "No error", ERR_GENERAL: "General error", - ERR_MISSING_KEY: "Missing_key error", + ERR_MISSING_KEY: "Missing key error", ERR_INVALID_TYPE: "Invalid type error", ERR_INVALID_CHOICE: "Invalid choice error", ERR_NOT_IN_RANGE: "Not in range error", @@ -85,16 +90,30 @@ func newModelManagerError(modelKeyword string) error { } func newModelNotFoundError(modelKeyword, idOrName string, err error) error { - msg := fmt.Sprintf("cannot find %q with id/name %q: %s", - modelKeyword, idOrName, err) + msg := fmt.Sprintf("cannot find %q with id/name %q", + modelKeyword, idOrName) + if err != sql.ErrNoRows { + msg += ": " + err.Error() + } return newError(ERR_MODEL_NOT_FOUND, msg) } func newError(typ ErrType, msg string) error { - return &ValidateError{ + err := &ValidateError{ ErrType: typ, Msg: msg, } + if returnHttpError { + switch typ { + case ERR_SUCCESS: + return nil + case ERR_GENERAL, ERR_MODEL_MANAGER: + return httperrors.NewInternalServerError(msg) + default: + return httperrors.NewInputParameterError(msg) + } + } + return err } func IsModelNotFoundError(err error) bool { diff --git a/pkg/cloudcommon/validators/validators_test.go b/pkg/cloudcommon/validators/validators_test.go index 29f110a2bc..508e4ea2e4 100644 --- a/pkg/cloudcommon/validators/validators_test.go +++ b/pkg/cloudcommon/validators/validators_test.go @@ -47,6 +47,8 @@ type C struct { } func testS(t *testing.T, v IValidator, c *C) { + returnHttpError = false + j, _ := jsonutils.ParseString(c.In) jd := j.(*jsonutils.JSONDict) err := v.Validate(jd) @@ -143,6 +145,91 @@ func TestStringChoicesValidator(t *testing.T) { } } +func TestStringMultiChoicesValidator(t *testing.T) { + type MultiChoicesC struct { + *C + KeepDup bool + } + choices := NewChoices("choice0", "choice1") + cases := []*MultiChoicesC{ + { + C: &C{ + Name: "missing non-optional", + In: `{}`, + Out: `{}`, + Optional: false, + Err: ERR_MISSING_KEY, + ValueWant: "", + }, + }, + { + C: &C{ + Name: "missing optional", + In: `{}`, + Out: `{}`, + Optional: true, + ValueWant: "", + }, + }, + { + C: &C{ + Name: "missing with default", + In: `{}`, + Out: `{s: "choice0,choice1"}`, + Default: "choice0,choice1", + ValueWant: "choice0,choice1", + }, + }, + { + C: &C{ + Name: "good choices", + In: `{"s": "choice0,choice1"}`, + Out: `{"s": "choice0,choice1"}`, + ValueWant: "choice0,choice1", + }, + }, + { + C: &C{ + Name: "keep dup", + In: `{"s": "choice0,choice0,choice1,choice0"}`, + Out: `{"s": "choice0,choice0,choice1,choice0"}`, + ValueWant: "choice0,choice0,choice1,choice0", + }, + KeepDup: true, + }, + { + C: &C{ + Name: "strip dup", + In: `{"s": "choice0,choice0,choice1,choice0"}`, + Out: `{"s": "choice0,choice1"}`, + ValueWant: "choice0,choice1", + }, + }, + { + C: &C{ + Name: "invalid choice", + In: `{"s": "choice0,choicex"}`, + Out: `{"s": "choice0,choicex"}`, + Err: ERR_INVALID_CHOICE, + ValueWant: "", + }, + }, + } + for _, c := range cases { + t.Run(c.Name, func(t *testing.T) { + v := NewStringMultiChoicesValidator("s", choices).Sep(",").KeepDup(c.KeepDup) + if c.Default != nil { + s := c.Default.(string) + v.Default(s) + } + if c.Optional { + v.Optional(true) + } + testS(t, v, c.C) + }) + } +} + func TestBoolValidator(t *testing.T) { cases := []*C{ { diff --git a/pkg/compute/guestdrivers/azure.go b/pkg/compute/guestdrivers/azure.go index d94159232d..02f594819c 100644 --- a/pkg/compute/guestdrivers/azure.go +++ b/pkg/compute/guestdrivers/azure.go @@ -21,6 +21,10 @@ type SAzureGuestDriver struct { SManagedVirtualizedGuestDriver } +const ( + DEFAULT_USER = "yunion" +) + func init() { driver := SAzureGuestDriver{} models.RegisterGuestDriver(&driver) @@ -157,7 +161,7 @@ func (self *SAzureGuestDriver) RequestDeployGuestOnHost(ctx context.Context, gue log.Errorf("encrypt password failed %s", err) } else { data.Add(jsonutils.NewString(iVM.GetOSType()), "os") - data.Add(jsonutils.NewString("root"), "account") + data.Add(jsonutils.NewString(DEFAULT_USER), "account") data.Add(jsonutils.NewString(encpasswd), "key") if len(desc.OsDistribution) > 0 { @@ -218,8 +222,8 @@ func (self *SAzureGuestDriver) RequestDeployGuestOnHost(ctx context.Context, gue } data := jsonutils.NewDict() - data.Add(jsonutils.NewString("root"), "account") // 用户名 - data.Add(jsonutils.NewString(encpasswd), "key") // 密码 + data.Add(jsonutils.NewString(DEFAULT_USER), "account") // 用户名 + data.Add(jsonutils.NewString(encpasswd), "key") // 密码 e := iVM.DeployVM(name, password, publicKey, resetPassword, deleteKeypair, description) return data, e }) diff --git a/pkg/compute/hostdrivers/aliyun.go b/pkg/compute/hostdrivers/aliyun.go index 421e3d1b87..2a7ad29a72 100644 --- a/pkg/compute/hostdrivers/aliyun.go +++ b/pkg/compute/hostdrivers/aliyun.go @@ -10,6 +10,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" ) type SAliyunHostDriver struct { @@ -25,7 +26,7 @@ func (self *SAliyunHostDriver) GetHostType() string { return models.HOST_TYPE_ALIYUN } -func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error { +func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error { params := task.GetParams() imageId, err := params.GetString("image_id") if err != nil { @@ -36,9 +37,16 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host * osType, _ := params.GetString("os_type") osDist, _ := params.GetString("os_distribution") + isForce := jsonutils.QueryBoolean(params, "is_force", false) userCred := task.GetUserCred() taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + + lockman.LockRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId)) + defer lockman.ReleaseRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId)) + + scimg := models.StoragecachedimageManager.Register(ctx, task.GetUserCred(), storageCache.Id, imageId) + iStorageCache, err := storageCache.GetIStorageCache() if err != nil { return nil, err @@ -49,6 +57,8 @@ func (self *SAliyunHostDriver) CheckAndSetCacheImage(ctx context.Context, host * if err != nil { return nil, err } else { + scimg.SetExternalId(extImgId) + ret := jsonutils.NewDict() ret.Add(jsonutils.NewString(extImgId), "image_id") return ret, nil diff --git a/pkg/compute/hostdrivers/azure.go b/pkg/compute/hostdrivers/azure.go index a0f6b949c5..4990830042 100644 --- a/pkg/compute/hostdrivers/azure.go +++ b/pkg/compute/hostdrivers/azure.go @@ -2,9 +2,11 @@ package hostdrivers import ( "context" + "fmt" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/httperrors" @@ -22,7 +24,7 @@ func (self *SAzureHostDriver) GetHostType() string { return models.HOST_TYPE_AZURE } -func (self *SAzureHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error { +func (self *SAzureHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error { params := task.GetParams() imageId, err := params.GetString("image_id") if err != nil { @@ -36,6 +38,10 @@ func (self *SAzureHostDriver) CheckAndSetCacheImage(ctx context.Context, host *m isForce := jsonutils.QueryBoolean(params, "is_force", false) userCred := task.GetUserCred() taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + lockman.LockRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId)) + defer lockman.ReleaseRawObject(ctx, "cachedimages", fmt.Sprintf("%s-%s", storageCache.Id, imageId)) + + scimg := models.StoragecachedimageManager.Register(ctx, task.GetUserCred(), storageCache.Id, imageId) iStorageCache, err := storageCache.GetIStorageCache() if err != nil { return nil, err diff --git a/pkg/compute/hostdrivers/hostdrivers.go b/pkg/compute/hostdrivers/hostdrivers.go deleted file mode 100644 index a58d95a88b..0000000000 --- a/pkg/compute/hostdrivers/hostdrivers.go +++ /dev/null @@ -1 +0,0 @@ -package hostdrivers diff --git a/pkg/compute/hostdrivers/kvm.go b/pkg/compute/hostdrivers/kvm.go index db6d253a3c..b317f44a49 100644 --- a/pkg/compute/hostdrivers/kvm.go +++ b/pkg/compute/hostdrivers/kvm.go @@ -25,7 +25,7 @@ func (self *SKVMHostDriver) GetHostType() string { return models.HOST_TYPE_HYPERVISOR } -func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, scimg *models.SStoragecachedimage, task taskman.ITask) error { +func (self *SKVMHostDriver) CheckAndSetCacheImage(ctx context.Context, host *models.SHost, storageCache *models.SStoragecache, task taskman.ITask) error { params := task.GetParams() imageId, err := params.GetString("image_id") if err != nil { diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 33c857efbe..26283436b1 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -108,11 +108,16 @@ func (manager *SDiskManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQu if !ok { return nil, fmt.Errorf("Invalid querystring formst: %v", query) } - if jsonutils.QueryBoolean(query, "unused", false) { + if query.Contains("unused") { guestdisks := GuestdiskManager.Query().SubQuery() sq := guestdisks.Query(guestdisks.Field("disk_id")) - q = q.Filter(sqlchemy.NotIn(q.Field("id"), sq)) + if jsonutils.QueryBoolean(query, "unused", false) { + q = q.Filter(sqlchemy.NotIn(q.Field("id"), sq)) + } else { + q = q.Filter(sqlchemy.In(q.Field("id"), sq)) + } } + storages := StorageManager.Query().SubQuery() if jsonutils.QueryBoolean(query, "share", false) { sq := storages.Query(storages.Field("id")).Filter(sqlchemy.NotIn(storages.Field("storage_type"), STORAGE_LOCAL_TYPES)) @@ -521,21 +526,9 @@ func (self *SDisk) GetCloudprovider() *SCloudprovider { } func (self *SDisk) GetPathAtHost(host *SHost) string { - storage := self.GetStorage() - if storage.StorageType == STORAGE_RBD { - pool, _ := storage.StorageConf.GetString("pool") - monHost, _ := storage.StorageConf.GetString("mon_host") - key, _ := storage.StorageConf.GetString("key") - for _, keyword := range []string{"@", ":", "="} { - monHost = strings.Replace(monHost, keyword, fmt.Sprintf("\\%s", keyword), -1) - key = strings.Replace(key, keyword, fmt.Sprintf("\\%s", keyword), -1) - } - return fmt.Sprintf("rbd:%s/%s:mon_host=%s:key=%s", pool, self.Id, monHost, key) - } else if storage.StorageType == STORAGE_LOCAL || storage.StorageType == STORAGE_NAS { - hostStorage := host.GetHoststorageOfId(self.StorageId) - if hostStorage != nil { - return path.Join(hostStorage.MountPoint, self.Id) - } + hostStorage := host.GetHoststorageOfId(self.StorageId) + if hostStorage != nil { + return path.Join(hostStorage.MountPoint, self.Id) } return "" } diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index 38f415cb14..02d10139eb 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -19,6 +19,7 @@ import ( "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" ) const ( @@ -485,6 +486,9 @@ func (self *SElasticip) PerformAssociate(ctx context.Context, userCred mcclient. server := vmObj.(*SGuest) + lockman.LockObject(ctx, server) + defer lockman.ReleaseObject(ctx, server) + if server.PendingDeleted { return nil, httperrors.NewInvalidStatusError("cannot associate pending delete server") } @@ -741,11 +745,30 @@ func (u EipUsage) Total() int { return u.PublicIPCount + u.EIPCount } -func (manager *SElasticipManager) TotalCount(projectId string) EipUsage { +func (manager *SElasticipManager) usageQ(q *sqlchemy.SQuery, rangeObj db.IStandaloneModel, hostTypes []string) *sqlchemy.SQuery { + if rangeObj == nil { + return q + } + zones := ZoneManager.Query().SubQuery() + hosts := HostManager.Query().SubQuery() + sq := zones.Query(zones.Field("cloudregion_id")). + Join(hosts, sqlchemy.AND( + sqlchemy.IsFalse(hosts.Field("deleted")), + sqlchemy.IsTrue(hosts.Field("enabled")), + sqlchemy.Equals(hosts.Field("zone_id"), zones.Field("id")))) + sq = AttachUsageQuery(sq, hosts, hosts.Field("id"), hostTypes, rangeObj) + q = q.Filter(sqlchemy.In(q.Field("cloudregion_id"), sq.Distinct())) + return q +} + +func (manager *SElasticipManager) TotalCount(projectId string, rangeObj db.IStandaloneModel, hostTypes []string) EipUsage { usage := EipUsage{} q1 := manager.Query().Equals("mode", EIP_MODE_INSTANCE_PUBLICIP) + q1 = manager.usageQ(q1, rangeObj, hostTypes) q2 := manager.Query().Equals("mode", EIP_MODE_STANDALONE_EIP) + q2 = manager.usageQ(q2, rangeObj, hostTypes) q3 := manager.Query().Equals("mode", EIP_MODE_STANDALONE_EIP).IsNotEmpty("associate_id") + q3 = manager.usageQ(q3, rangeObj, hostTypes) if len(projectId) > 0 { q1 = q1.Equals("tenant_id", projectId) q2 = q2.Equals("tenant_id", projectId) diff --git a/pkg/compute/models/guestdisks.go b/pkg/compute/models/guestdisks.go index 1f69a83d0d..26e0c74c0c 100644 --- a/pkg/compute/models/guestdisks.go +++ b/pkg/compute/models/guestdisks.go @@ -136,7 +136,8 @@ func (self *SGuestdisk) GetJsonDescAtHost(host *SHost) jsonutils.JSONObject { desc.Add(jsonutils.NewString(storagecacheimg.Path), "image_path") } } - if host.HostType == HOST_TYPE_HYPERVISOR && disk.IsLocal() { + storage := disk.GetStorage() + if host.HostType == HOST_TYPE_HYPERVISOR && disk.IsLocal() || (storage != nil && storage.StorageType == STORAGE_RBD) { desc.Add(jsonutils.NewString(disk.StorageId), "storage_id") localpath := disk.GetPathAtHost(host) if len(localpath) == 0 { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 33451ac727..d930e8ce5b 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -319,6 +319,20 @@ func (manager *SGuestManager) ListItemFilter(ctx context.Context, q *sqlchemy.SQ q = q.In("host_id", sq) } + withEip, _ := queryDict.GetString("with_eip") + withoutEip, _ := queryDict.GetString("without_eip") + if len(withEip) > 0 || len(withoutEip) > 0 { + eips := ElasticipManager.Query().SubQuery() + sq := eips.Query(eips.Field("associate_id")).Equals("associate_type", EIP_ASSOCIATE_TYPE_SERVER) + sq = sq.IsNotNull("associate_id").IsNotEmpty("associate_id") + + if utils.ToBool(withEip) { + q = q.In("id", sq) + } else if utils.ToBool(withoutEip) { + q = q.NotIn("id", sq) + } + } + gpu, _ := queryDict.GetString("gpu") if len(gpu) != 0 { isodev := IsolatedDeviceManager.Query().SubQuery() @@ -1629,8 +1643,10 @@ func (self *SGuest) PerformSaveImage(ctx context.Context, userCred mcclient.Toke properties.Add(jsonutils.NewString(self.OsType), "os_type") kwargs.Add(properties, "properties") kwargs.Add(jsonutils.NewBool(restart), "restart") + lockman.LockObject(ctx, disks.Root) defer lockman.ReleaseObject(ctx, disks.Root) + if imageId, err := disks.Root.PrepareSaveImage(ctx, userCred, kwargs); err != nil { return nil, err } else { @@ -2023,7 +2039,7 @@ func (self *SGuest) CreateDisksOnHost(ctx context.Context, userCred mcclient.Tok func (self *SGuest) createDiskOnStorage(ctx context.Context, userCred mcclient.TokenCredential, storage *SStorage, diskConfig *SDiskConfig, pendingUsage quotas.IQuota) (*SDisk, error) { lockman.LockObject(ctx, storage) - defer lockman.LockObject(ctx, storage) + defer lockman.ReleaseObject(ctx, storage) lockman.LockClass(ctx, QuotaManager, self.ProjectId) defer lockman.ReleaseClass(ctx, QuotaManager, self.ProjectId) diff --git a/pkg/compute/models/hostdrivers.go b/pkg/compute/models/hostdrivers.go index 3bec75d3b4..7a46ba820d 100644 --- a/pkg/compute/models/hostdrivers.go +++ b/pkg/compute/models/hostdrivers.go @@ -11,7 +11,7 @@ import ( type IHostDriver interface { GetHostType() string - CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, scimg *SStoragecachedimage, task taskman.ITask) error + CheckAndSetCacheImage(ctx context.Context, host *SHost, storagecache *SStoragecache, task taskman.ITask) error RequestPrepareSaveDiskOnHost(ctx context.Context, host *SHost, disk *SDisk, imageId string, task taskman.ITask) error RequestSaveUploadImageOnHost(ctx context.Context, host *SHost, disk *SDisk, imageId string, task taskman.ITask, data jsonutils.JSONObject) error RequestAllocateDiskOnStorage(ctx context.Context, host *SHost, storage *SStorage, disk *SDisk, task taskman.ITask, content *jsonutils.JSONDict) error diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 8d995a54ce..16d0e43494 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -583,7 +583,8 @@ func (self *SHost) GetAttachedStorageCapacity() SStorageCapacity { func _getLeastUsedStorage(storages []SStorage, backends []string) *SStorage { var best *SStorage var bestCap int - for _, s := range storages { + for i := 0; i < len(storages); i++ { + s := storages[i] if len(backends) > 0 { in, _ := utils.InStringArray(s.StorageType, backends) if !in { diff --git a/pkg/compute/models/isolated_devices.go b/pkg/compute/models/isolated_devices.go index d1c7c26c21..d88332cf5c 100644 --- a/pkg/compute/models/isolated_devices.go +++ b/pkg/compute/models/isolated_devices.go @@ -346,6 +346,7 @@ func (manager *SIsolatedDeviceManager) totalCountQ( sqlchemy.IsFalse(hosts.Field("deleted")), sqlchemy.IsTrue(hosts.Field("enabled")), )) + q = q.Filter(sqlchemy.Equals(hosts.Field("id"), q.Field("host_id"))) if len(devType) != 0 { q.In("dev_type", devType) } diff --git a/pkg/compute/models/keypairs.go b/pkg/compute/models/keypairs.go index 4563ff639f..84cb4f5b28 100644 --- a/pkg/compute/models/keypairs.go +++ b/pkg/compute/models/keypairs.go @@ -116,8 +116,8 @@ func totalKeypairCount(userId string) int { return q.Count() } -func (manager *SKeypairManager) FilterByOwner(q *sqlchemy.SQuery, ownerId string) *sqlchemy.SQuery { - return q.Equals("owner_id", ownerId) +func (manager *SKeypairManager) FilterByOwner(q *sqlchemy.SQuery, owner string) *sqlchemy.SQuery { + return q.Equals("owner_id", owner) } func (self *SKeypair) GetOwnerProjectId() string { diff --git a/pkg/compute/models/quotas.go b/pkg/compute/models/quotas.go index a648831646..f18d91ce89 100644 --- a/pkg/compute/models/quotas.go +++ b/pkg/compute/models/quotas.go @@ -71,7 +71,7 @@ func (self *SQuota) FetchUsage(projectId string) error { diskSize := totalDiskSize(projectId, tristate.None, tristate.None, false) net := totalGuestNicCount(projectId, nil, false) guest := totalGuestResourceCount(projectId, nil, nil, "", false, false, "") - eipUsage := ElasticipManager.TotalCount(projectId) + eipUsage := ElasticipManager.TotalCount(projectId, nil, nil) // XXX // keypair belongs to user // keypair := totalKeypairCount(projectId) diff --git a/pkg/compute/models/secgroups.go b/pkg/compute/models/secgroups.go index a572574bb9..9f2efa7bd4 100644 --- a/pkg/compute/models/secgroups.go +++ b/pkg/compute/models/secgroups.go @@ -270,8 +270,10 @@ func (manager *SSecurityGroupManager) DelaySync(ctx context.Context, userCred mc log.Errorf("DelaySync secgroup failed") } else { needSync := false + lockman.LockObject(ctx, secgrp) defer lockman.ReleaseObject(ctx, secgrp) + if secgrp.IsDirty { if _, err := secgrp.GetModelManager().TableSpec().Update(secgrp, func() error { secgrp.IsDirty = false diff --git a/pkg/compute/models/usage.go b/pkg/compute/models/usage.go index 6483fe6e7a..e45bce93de 100644 --- a/pkg/compute/models/usage.go +++ b/pkg/compute/models/usage.go @@ -42,6 +42,10 @@ func AttachUsageQuery( sqlchemy.Equals(hosts.Field("id"), aggHosts.Field("host_id")), sqlchemy.IsFalse(aggHosts.Field("deleted")))). Filter(sqlchemy.Equals(aggHosts.Field("schedtag_id"), rangeObjId)) + case "cloudregion": + zones := ZoneManager.Query().SubQuery() + q = q.Join(zones, sqlchemy.Equals(hosts.Field("zone_id"), zones.Field("id"))) + q = q.Filter(sqlchemy.Equals(zones.Field("cloudregion_id"), rangeObjId)) } return q } diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index b3a6c10dba..65758a2313 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -73,6 +73,7 @@ func (self *GuestBatchCreateTask) onSchedulerRequestFail(ctx context.Context, gu func (self *GuestBatchCreateTask) onScheduleFail(ctx context.Context, guest *models.SGuest, msg string) { lockman.LockObject(ctx, guest) defer lockman.ReleaseObject(ctx, guest) + reason := "No matching resources" if len(msg) > 0 { reason = fmt.Sprintf("%s: %s", reason, msg) diff --git a/pkg/compute/tasks/guest_change_config_task.go b/pkg/compute/tasks/guest_change_config_task.go index 78963011d8..84b112e43c 100644 --- a/pkg/compute/tasks/guest_change_config_task.go +++ b/pkg/compute/tasks/guest_change_config_task.go @@ -181,8 +181,10 @@ func (self *GuestChangeConfigTask) OnCreateDisksComplete(ctx context.Context, ob if addMem > 0 { cancelUsage.Memory = addMem } + lockman.LockClass(ctx, guest.GetModelManager(), guest.ProjectId) defer lockman.ReleaseClass(ctx, guest.GetModelManager(), guest.ProjectId) + err = models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, guest.ProjectId, &pendingUsage, &cancelUsage) if err != nil { self.SetStageFailed(ctx, err.Error()) diff --git a/pkg/compute/tasks/storage_cache_image_task.go b/pkg/compute/tasks/storage_cache_image_task.go index fc24e13ea0..bd3d476fab 100644 --- a/pkg/compute/tasks/storage_cache_image_task.go +++ b/pkg/compute/tasks/storage_cache_image_task.go @@ -34,7 +34,7 @@ func (self *StorageCacheImageTask) OnInit(ctx context.Context, obj db.IStandalon self.SetStage("on_image_cache_complete", nil) host, _ := storageCache.GetHost() - err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, scimg, self) + err := host.GetHostDriver().CheckAndSetCacheImage(ctx, host, storageCache, self) if err != nil { errData := taskman.Error2TaskData(err) self.OnImageCacheCompleteFailed(ctx, storageCache, errData) @@ -46,8 +46,8 @@ func (self *StorageCacheImageTask) OnImageCacheComplete(ctx context.Context, obj storageCache := obj.(*models.SStoragecache) imageId, _ := self.Params.GetString("image_id") scimg := models.StoragecachedimageManager.Register(ctx, self.UserCred, storageCache.Id, imageId) - extImgId, _ := data.GetString("image_id") - self.OnCacheSucc(ctx, storageCache, imageId, scimg, extImgId) + // extImgId, _ := data.GetString("image_id") + self.OnCacheSucc(ctx, storageCache, imageId, scimg) } func (self *StorageCacheImageTask) OnImageCacheCompleteFailed(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { @@ -70,11 +70,8 @@ func (self *StorageCacheImageTask) OnCacheFailed(ctx context.Context, cache *mod self.SetStageFailed(ctx, err.Error()) } -func (self *StorageCacheImageTask) OnCacheSucc(ctx context.Context, cache *models.SStoragecache, imageId string, scimg *models.SStoragecachedimage, extImgId string) { +func (self *StorageCacheImageTask) OnCacheSucc(ctx context.Context, cache *models.SStoragecache, imageId string, scimg *models.SStoragecachedimage) { scimg.SetStatus(self.UserCred, models.CACHED_IMAGE_STATUS_READY, "cached") - if len(cache.ExternalId) > 0 && len(extImgId) > 0 && scimg.ExternalId != extImgId { - scimg.SetExternalId(extImgId) - } models.CachedimageManager.ImageAddRefCount(imageId) db.OpsLog.LogEvent(cache, db.ACT_CACHED_IMAGE, imageId, self.UserCred) self.SetStageComplete(ctx, nil) diff --git a/pkg/compute/usages/handler.go b/pkg/compute/usages/handler.go index 5f41c23aa7..b6212002df 100644 --- a/pkg/compute/usages/handler.go +++ b/pkg/compute/usages/handler.go @@ -98,12 +98,13 @@ func addHandler(prefix, rangeObjKey string, hf appsrv.FilterHandler, app *appsrv func AddUsageHandler(prefix string, app *appsrv.Application) { prefix = fmt.Sprintf("%s/usages", prefix) for key, f := range map[string]appsrv.FilterHandler{ - "": rangeObjHandler(nil, ReportGeneralUsage), - "zone": rangeObjHandler(models.ZoneManager, ReportZoneUsage), - "wire": rangeObjHandler(models.WireManager, ReportWireUsage), - "schedtag": rangeObjHandler(models.SchedtagManager, ReportSchedtagUsage), - "host": rangeObjHandler(models.HostManager, ReportHostUsage), - "vcenter": rangeObjHandler(models.VCenterManager, ReportVCenterUsage), + "": rangeObjHandler(nil, ReportGeneralUsage), + "zone": rangeObjHandler(models.ZoneManager, ReportZoneUsage), + "wire": rangeObjHandler(models.WireManager, ReportWireUsage), + "schedtag": rangeObjHandler(models.SchedtagManager, ReportSchedtagUsage), + "host": rangeObjHandler(models.HostManager, ReportHostUsage), + "vcenter": rangeObjHandler(models.VCenterManager, ReportVCenterUsage), + "cloudregion": rangeObjHandler(models.CloudregionManager, ReportCloudRegionUsage), } { addHandler(prefix, key, auth.Authenticate(f), app) } @@ -150,8 +151,11 @@ func ReportSchedtagUsage(userCred mcclient.TokenCredential, schedtag db.IStandal } func ReportZoneUsage(userCred mcclient.TokenCredential, zone db.IStandaloneModel, hostTypes []string) (Usage, error) { - //return - return nil, nil + return ReportGeneralUsage(userCred, zone, hostTypes) +} + +func ReportCloudRegionUsage(userCred mcclient.TokenCredential, cloudRegion db.IStandaloneModel, hostTypes []string) (Usage, error) { + return ReportGeneralUsage(userCred, cloudRegion, hostTypes) } //func ReportGuestUsage() @@ -212,7 +216,7 @@ func ReportGeneralUsage(userCred mcclient.TokenCredential, rangeObj db.IStandalo IsolatedDeviceUsage(rangeObj, hostTypes), WireUsage(rangeObj, hostTypes), NetworkUsage(userCred, rangeObj), - EipUsage(), + EipUsage(rangeObj, hostTypes), ) return @@ -338,12 +342,12 @@ func IsolatedDeviceUsage(rangeObj db.IStandaloneModel, hostType []string) Usage return count } -func EipUsage() Usage { - eipUsage := models.ElasticipManager.TotalCount("") +func EipUsage(rangeObj db.IStandaloneModel, hostTypes []string) Usage { + eipUsage := models.ElasticipManager.TotalCount("", rangeObj, hostTypes) count := make(map[string]interface{}) count["eip.all"] = eipUsage.Total() count["eip.public_ip"] = eipUsage.PublicIPCount count["eip.floating_ip"] = eipUsage.EIPCount count["eip.floating_ip.used"] = eipUsage.EIPUsedCount return count -} \ No newline at end of file +} diff --git a/pkg/mcclient/modules/k8s/deployment.go b/pkg/mcclient/modules/k8s/deployment.go index e62df70244..c72fe53250 100644 --- a/pkg/mcclient/modules/k8s/deployment.go +++ b/pkg/mcclient/modules/k8s/deployment.go @@ -4,15 +4,28 @@ import ( "yunion.io/x/onecloud/pkg/mcclient/modules" ) -var Deployments *DeploymentManager +var ( + Deployments *DeploymentManager + DeployFromFile *DeployFromFileManager +) type DeploymentManager struct { *NamespaceResourceManager } +type DeployFromFileManager struct { + *NamespaceResourceManager +} + func init() { Deployments = &DeploymentManager{ NewNamespaceResourceManager("deployment", "deployments", NewNamespaceCols(), NewColumns())} + + DeployFromFile = &DeployFromFileManager{ + NewNamespaceResourceManager("deployfromfile", "deployfromfiles", + NewNamespaceCols(), NewColumns())} + modules.Register(Deployments) + modules.Register(DeployFromFile) } diff --git a/pkg/mcclient/modules/k8s/raw.go b/pkg/mcclient/modules/k8s/raw.go index d291ab9488..b77db278a8 100644 --- a/pkg/mcclient/modules/k8s/raw.go +++ b/pkg/mcclient/modules/k8s/raw.go @@ -68,10 +68,8 @@ func (m *RawResourceManager) Get(s *mcclient.ClientSession, kind string, namespa } func (m *RawResourceManager) Put(s *mcclient.ClientSession, kind string, namespace string, name string, body jsonutils.JSONObject, cluster string) error { - newBody := jsonutils.NewDict() - newBody.Add(body, "raw") ctx := newRawResource(kind, namespace, name, cluster) - _, err := m.request(s, "PUT", ctx.path(), newBody) + _, err := m.request(s, "PUT", ctx.path(), body) return err } diff --git a/pkg/mcclient/modules/mod_webconsole.go b/pkg/mcclient/modules/mod_webconsole.go index 52a145384a..0c1b79f985 100644 --- a/pkg/mcclient/modules/mod_webconsole.go +++ b/pkg/mcclient/modules/mod_webconsole.go @@ -13,6 +13,8 @@ var ( func init() { WebConsole = WebConsoleManager{NewWebConsoleManager()} + + register(&WebConsole) } type WebConsoleManager struct { diff --git a/pkg/mcclient/options/base.go b/pkg/mcclient/options/base.go index 1062e83346..eafefcceec 100644 --- a/pkg/mcclient/options/base.go +++ b/pkg/mcclient/options/base.go @@ -3,6 +3,7 @@ package options import ( "fmt" "reflect" + "time" "yunion.io/x/jsonutils" "yunion.io/x/pkg/gotypes" @@ -88,6 +89,11 @@ func optionsStructRvToParams(rv reflect.Value) (*jsonutils.JSONDict, error) { if ft.Anonymous { continue } + if f.Type() == gotypes.TimeType { + t := f.Interface().(time.Time) + p.Set(name, jsonutils.NewTimeString(t)) + continue + } // TODO msg := fmt.Sprintf("do not know what to do with non-anonymous struct field: %s", ft.Name) panic(msg) @@ -142,23 +148,24 @@ func ListStructToParams(v interface{}) (*jsonutils.JSONDict, error) { } type BaseListOptions struct { - Limit *int `default:"20" help:"Page limit"` - Offset *int `default:"0" help:"Page offset"` - OrderBy []string `help:"Name of the field to be ordered by"` - Order string `help:"List order" choices:"desc|asc"` - Details *bool `help:"Show more details" default:"false"` - Search string `help:"Filter results by a simple keyword search"` - Meta *bool `help:"Piggyback metadata information"` - Filter []string `help:"Filters"` - JointFilter []string `help:"Filters with joint table col; joint_tbl.related_key(origin_key).filter_col.filter_cond(filters)"` - FilterAny *bool `help:"If true, match if any of the filters matches; otherwise, match if all of the filters match"` - Admin *bool `help:"Is an admin call?"` - Tenant string `help:"Tenant ID or Name"` - User string `help:"User ID or Name"` - System *bool `help:"Show system resource"` - PendingDelete *bool `help:"Show pending deleted resource"` - Field []string `help:"Show only specified fields"` - ShowEmulated *bool `help:"Show all resources including the emulated resources"` + Limit *int `default:"20" help:"Page limit"` + Offset *int `default:"0" help:"Page offset"` + OrderBy []string `help:"Name of the field to be ordered by"` + Order string `help:"List order" choices:"desc|asc"` + Details *bool `help:"Show more details" default:"false"` + Search string `help:"Filter results by a simple keyword search"` + Meta *bool `help:"Piggyback metadata information"` + Filter []string `help:"Filters"` + JointFilter []string `help:"Filters with joint table col; joint_tbl.related_key(origin_key).filter_col.filter_cond(filters)"` + FilterAny *bool `help:"If true, match if any of the filters matches; otherwise, match if all of the filters match"` + Admin *bool `help:"Is an admin call?"` + Tenant string `help:"Tenant ID or Name"` + User string `help:"User ID or Name"` + System *bool `help:"Show system resource"` + PendingDelete *bool `help:"Show only pending deleted resource"` + PendingDeleteAll *bool `help:"Show all resources including pending deleted" json:"-"` + Field []string `help:"Show only specified fields"` + ShowEmulated *bool `help:"Show all resources including the emulated resources"` } func (opts *BaseListOptions) Params() (*jsonutils.JSONDict, error) { @@ -169,10 +176,16 @@ func (opts *BaseListOptions) Params() (*jsonutils.JSONDict, error) { if len(opts.Filter) == 0 { params.Remove("filter_any") } + if BoolV(opts.PendingDeleteAll) { + params.Set("pending_delete", jsonutils.NewString("all")) + } if opts.Admin == nil { - requiresSystem := len(opts.Tenant) > 0 || BoolV(opts.System) || BoolV(opts.PendingDelete) + requiresSystem := len(opts.Tenant) > 0 || + BoolV(opts.System) || + BoolV(opts.PendingDelete) || + BoolV(opts.PendingDeleteAll) if requiresSystem { - params.Set("admin", jsonutils.NewBool(true)) + params.Set("admin", jsonutils.JSONTrue) } } return params, nil diff --git a/pkg/mcclient/options/servers.go b/pkg/mcclient/options/servers.go index a5f9f70a11..879424c738 100644 --- a/pkg/mcclient/options/servers.go +++ b/pkg/mcclient/options/servers.go @@ -21,6 +21,8 @@ type ServerListOptions struct { Hypervisor string `help:"Show server of hypervisor" choices:"kvm|esxi|container|baremetal|aliyun|azure"` Manager string `help:"Show servers imported from manager"` Region string `help:"Show servers in cloudregion"` + WithEip *bool `help:"Show Servers with EIP"` + WithoutEip *bool `help:"Show Servers without EIP"` BaseListOptions } @@ -210,6 +212,7 @@ type ServerDeployOptions struct { Deploy []string `help:"Specify deploy files in virtual server file system" json:"-"` ResetPassword *bool `help:"Force reset password"` Password string `help:"Default user password"` + AutoStart *bool `help:"Auto start server after deployed"` } func (opts *ServerDeployOptions) Params() (*jsonutils.JSONDict, error) { @@ -252,11 +255,12 @@ type ServerMonitorOptions struct { } type ServerSaveImageOptions struct { - ID string `help:"ID or name of server" json:"-"` - IMAGE string `help:"Image name" json:"name"` - Public *bool `help:"Make the image public available" json:"is_public"` - Format string `help:"image format" choices:"vmdk|qcow2"` - Notes string `help:"Notes about the image"` + ID string `help:"ID or name of server" json:"-"` + IMAGE string `help:"Image name" json:"name"` + Public *bool `help:"Make the image public available" json:"is_public"` + Format string `help:"image format" choices:"vmdk|qcow2"` + Notes string `help:"Notes about the image"` + AutoStart *bool `help:"Auto start server after image saved"` } type ServerRebuildRootOptions struct { diff --git a/pkg/util/aliyun/disk.go b/pkg/util/aliyun/disk.go index 70a55f4350..111e36bfd2 100644 --- a/pkg/util/aliyun/disk.go +++ b/pkg/util/aliyun/disk.go @@ -122,6 +122,9 @@ func (self *SDisk) Resize(size int64) error { } func (self *SDisk) GetName() string { + if len(self.DiskName) > 0 { + return self.DiskName + } return self.DiskId } diff --git a/pkg/util/aliyun/storagecache.go b/pkg/util/aliyun/storagecache.go index 666b46574f..d945c8d054 100644 --- a/pkg/util/aliyun/storagecache.go +++ b/pkg/util/aliyun/storagecache.go @@ -5,7 +5,6 @@ import ( "os" "strings" "time" - "github.com/aliyun/aliyun-oss-go-sdk/oss" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -87,12 +86,14 @@ func (self *SStoragecache) GetIImages() ([]cloudprovider.ICloudImage, error) { } func (self *SStoragecache) UploadImage(userCred mcclient.TokenCredential, imageId string, osArch, osType, osDist string, extId string, isForce bool) (string, error) { + if len(extId) > 0 { status, _ := self.region.GetImageStatus(extId) if status == ImageStatusAvailable && !isForce { return extId, nil } } + return self.uploadImage(userCred, imageId, osArch, osType, osDist, isForce) } diff --git a/pkg/util/azure/azure.go b/pkg/util/azure/azure.go index 85527e143e..de4a4d0f25 100644 --- a/pkg/util/azure/azure.go +++ b/pkg/util/azure/azure.go @@ -2,6 +2,8 @@ package azure import ( "context" + "fmt" + "regexp" "strings" "github.com/Azure/azure-sdk-for-go/services/preview/subscription/mgmt/2018-03-01-preview/subscription" @@ -37,7 +39,7 @@ const ( EIP_RESOURCE = "eip" ) -var DefaultResourceGroups = map[string]string{ +var defaultResourceGroups = map[string]string{ DISK_RESOURCE: "YunionDiskResource", INSTANCE_RESOURCE: "YunionInstanceResource", VPC_RESOURCE: "YunionVpcResource", @@ -88,6 +90,18 @@ func NewAzureClient(providerId string, providerName string, accessKey string, se } } +func pareResourceGroupWithName(s string, resourceType string) (string, string, string) { + valid := regexp.MustCompile("resourceGroups/(.+)/providers/.+/(.+)$") + if resourceGroups := valid.FindStringSubmatch(s); len(resourceGroups) == 3 { + return s, resourceGroups[1], resourceGroups[2] + } + if len(s) == 0 { + log.Errorf("pareResourceGroupWithName[%s] error", resourceType) + } + globalId := fmt.Sprintf("resourceGroups/%s/providers/%s/%s", defaultResourceGroups[resourceType], resourceType, s) + return globalId, defaultResourceGroups[resourceType], s +} + func (self *SAzureClient) isResourceGroupExist(resourceGroup string) (bool, error) { groupClient := resources.NewGroupsClientWithBaseURI(self.baseUrl, self.subscriptionId) groupClient.Authorizer = self.authorizer @@ -113,7 +127,7 @@ func (self *SAzureClient) createResourceGroup(resourceGruop string) error { } func (self *SAzureClient) fetchAzueResourceGroup() error { - for _, value := range DefaultResourceGroups { + for _, value := range defaultResourceGroups { if exist, err := self.isResourceGroupExist(value); err != nil { log.Errorf("Check ResourceGroup error: %v", err) } else if !exist { diff --git a/pkg/util/azure/disk.go b/pkg/util/azure/disk.go index 0b199d2d58..ebee8edfcc 100644 --- a/pkg/util/azure/disk.go +++ b/pkg/util/azure/disk.go @@ -2,7 +2,6 @@ package azure import ( "context" - "fmt" "time" "yunion.io/x/jsonutils" @@ -66,23 +65,23 @@ type SDisk struct { Tags map[string]string } -func (self *SRegion) CreateDisk(storageType string, name string, sizeGb int32, desc string) error { +func (self *SRegion) CreateDisk(storageType string, name string, sizeGb int32, desc string) (string, error) { return self.createDisk(storageType, name, sizeGb, desc) } -func (self *SRegion) createDisk(storageType string, name string, sizeGb int32, desc string) error { +func (self *SRegion) createDisk(storageType string, name string, sizeGb int32, desc string) (string, error) { computeClient := compute.NewDisksClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) computeClient.Authorizer = self.client.authorizer sku := compute.DiskSku{Name: compute.StorageAccountTypes(storageType)} properties := compute.DiskProperties{DiskSizeGB: &sizeGb, CreationData: &compute.CreationData{CreateOption: "Empty"}} disk := compute.Disk{Name: &name, Location: &self.Name, DiskProperties: &properties, Sku: &sku} - resourceGroup, diskName := PareResourceGroupWithName(name, DISK_RESOURCE) + diskId, resourceGroup, diskName := pareResourceGroupWithName(name, DISK_RESOURCE) if result, err := computeClient.CreateOrUpdate(context.Background(), resourceGroup, diskName, disk); err != nil { - return err + return "", err } else if err := result.WaitForCompletion(context.Background(), computeClient.Client); err != nil { - return err + return "", err } else { - return nil + return diskId, nil } } @@ -93,7 +92,7 @@ func (self *SRegion) DeleteDisk(diskId string) error { func (self *SRegion) deleteDisk(diskId string) error { diskClient := compute.NewDisksClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) diskClient.Authorizer = self.client.authorizer - resourceGroup, name := PareResourceGroupWithName(diskId, DISK_RESOURCE) + _, resourceGroup, name := pareResourceGroupWithName(diskId, DISK_RESOURCE) if result, err := diskClient.Delete(context.Background(), resourceGroup, name); err != nil { return err } else if err := result.WaitForCompletion(context.Background(), diskClient.Client); err != nil { @@ -109,7 +108,7 @@ func (self *SRegion) ResizeDisk(diskId string, sizeGb int32) error { func (self *SRegion) resizeDisk(diskId string, sizeGb int32) error { diskClient := compute.NewDisksClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) diskClient.Authorizer = self.client.authorizer - resourceGroup, diskName := PareResourceGroupWithName(diskId, DISK_RESOURCE) + _, resourceGroup, diskName := pareResourceGroupWithName(diskId, DISK_RESOURCE) params := compute.DiskUpdate{ DiskUpdateProperties: &compute.DiskUpdateProperties{ DiskSizeGB: &sizeGb, @@ -123,10 +122,11 @@ func (self *SRegion) resizeDisk(diskId string, sizeGb int32) error { return nil } -func (self *SRegion) GetDisk(resourceGroup string, diskName string) (*SDisk, error) { +func (self *SRegion) GetDisk(diskId string) (*SDisk, error) { disk := SDisk{} computeClient := compute.NewDisksClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) computeClient.Authorizer = self.client.authorizer + _, resourceGroup, diskName := pareResourceGroupWithName(diskId, DISK_RESOURCE) if _disk, err := computeClient.Get(context.Background(), resourceGroup, diskName); err != nil { return nil, err } else if err := jsonutils.Update(&disk, _disk); err != nil { @@ -174,13 +174,8 @@ func (self *SDisk) GetId() string { return self.ID } -func (self *SRegion) getDisk(resourceGroup string, diskName string) (*SDisk, error) { - return self.GetDisk(resourceGroup, diskName) -} - func (self *SDisk) Refresh() error { - resourceGropu, diskName := PareResourceGroupWithName(self.ID, DISK_RESOURCE) - if disk, err := self.storage.zone.region.GetDisk(resourceGropu, diskName); err != nil { + if disk, err := self.storage.zone.region.GetDisk(self.ID); err != nil { return cloudprovider.ErrNotFound } else { return jsonutils.Update(self, disk) @@ -203,8 +198,8 @@ func (self *SDisk) GetName() string { } func (self *SDisk) GetGlobalId() string { - resourceGroup, _ := PareResourceGroupWithName(self.ID, DISK_RESOURCE) - return fmt.Sprintf("resourceGroups/%s/providers/disk/%s", resourceGroup, self.Name) + globalId, _, _ := pareResourceGroupWithName(self.ID, DISK_RESOURCE) + return globalId } func (self *SDisk) IsEmulated() bool { diff --git a/pkg/util/azure/eip.go b/pkg/util/azure/eip.go index ad1904f169..d69f196fa2 100644 --- a/pkg/util/azure/eip.go +++ b/pkg/util/azure/eip.go @@ -2,7 +2,6 @@ package azure import ( "context" - "fmt" "github.com/Azure/azure-sdk-for-go/services/network/mgmt/2018-06-01/network" "yunion.io/x/jsonutils" @@ -50,7 +49,7 @@ type SEipAddress struct { func (region *SRegion) AllocateEIP(eipName string) (*SEipAddress, error) { eip := SEipAddress{region: region} - resourceGroup, eipName := PareResourceGroupWithName(eipName, EIP_RESOURCE) + _, resourceGroup, eipName := pareResourceGroupWithName(eipName, EIP_RESOURCE) networkClient := network.NewPublicIPAddressesClientWithBaseURI(region.client.baseUrl, region.SubscriptionID) networkClient.Authorizer = region.client.authorizer params := network.PublicIPAddress{ @@ -106,7 +105,7 @@ func (region *SRegion) GetEips() ([]SEipAddress, error) { func (region *SRegion) GetEip(eipId string) (*SEipAddress, error) { eip := SEipAddress{region: region} - resourceGroup, eipName := PareResourceGroupWithName(eipId, EIP_RESOURCE) + _, resourceGroup, eipName := pareResourceGroupWithName(eipId, EIP_RESOURCE) if len(eipName) == 0 { return nil, cloudprovider.ErrNotFound } @@ -131,13 +130,11 @@ func (self *SEipAddress) Associate(instanceId string) error { } func (region *SRegion) AssociateEip(eipId string, instanceId string) error { - resourceGroup, instanceName := PareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) - if instance, err := region.GetInstance(resourceGroup, instanceName); err != nil { + if instance, err := region.GetInstance(instanceId); err != nil { return err } else { nicId := instance.Properties.NetworkProfile.NetworkInterfaces[0].ID - resourceGroup, nicName := PareResourceGroupWithName(nicId, NIC_RESOURCE) - if nic, err := region.getNetworkInterface(resourceGroup, nicName); err != nil { + if nic, err := region.getNetworkInterface(nicId); err != nil { return err } else { oldIPConf := nic.Properties.IPConfigurations[0] @@ -161,6 +158,7 @@ func (region *SRegion) AssociateEip(eipId string, instanceId string) error { IPConfigurations: &InterfaceIPConfiguration, }, } + _, resourceGroup, nicName := pareResourceGroupWithName(nic.ID, NIC_RESOURCE) if result, err := interfaceClinet.CreateOrUpdate(context.Background(), resourceGroup, nicName, params); err != nil { return err } else if err := result.WaitForCompletion(context.Background(), interfaceClinet.Client); err != nil { @@ -173,7 +171,7 @@ func (region *SRegion) AssociateEip(eipId string, instanceId string) error { func (region *SRegion) GetIEipById(eipId string) (cloudprovider.ICloudEIP, error) { eip := SEipAddress{region: region} - resourceGroup, eipName := PareResourceGroupWithName(eipId, EIP_RESOURCE) + _, resourceGroup, eipName := pareResourceGroupWithName(eipId, EIP_RESOURCE) if len(eipName) == 0 { return nil, cloudprovider.ErrNotFound } @@ -199,7 +197,7 @@ func (self *SEipAddress) Delete() error { } func (region *SRegion) DeallocateEIP(eipId string) error { - resourceGroup, eipName := PareResourceGroupWithName(eipId, EIP_RESOURCE) + _, resourceGroup, eipName := pareResourceGroupWithName(eipId, EIP_RESOURCE) networkClient := network.NewPublicIPAddressesClientWithBaseURI(region.client.baseUrl, region.SubscriptionID) networkClient.Authorizer = region.client.authorizer if result, err := networkClient.Delete(context.Background(), resourceGroup, eipName); err != nil { @@ -221,8 +219,7 @@ func (region *SRegion) DissociateEip(eipId string) error { log.Debugf("eip %s not associate any instance", eip.Name) return nil } else { - resourceGroup, nicName := PareResourceGroupWithName(eip.Properties.IPConfiguration.ID, NIC_RESOURCE) - if nic, err := region.getNetworkInterface(resourceGroup, nicName); err != nil { + if nic, err := region.getNetworkInterface(eip.Properties.IPConfiguration.ID); err != nil { return err } else { oldIPConf := nic.Properties.IPConfigurations[0] @@ -245,6 +242,7 @@ func (region *SRegion) DissociateEip(eipId string) error { IPConfigurations: &InterfaceIPConfiguration, }, } + _, resourceGroup, nicName := pareResourceGroupWithName(nic.ID, NIC_RESOURCE) if result, err := interfaceClinet.CreateOrUpdate(context.Background(), resourceGroup, nicName, params); err != nil { return err } else if err := result.WaitForCompletion(context.Background(), interfaceClinet.Client); err != nil { @@ -268,8 +266,8 @@ func (self *SEipAddress) GetBandwidth() int { } func (self *SEipAddress) GetGlobalId() string { - resourceGroup, eipName := PareResourceGroupWithName(self.ID, EIP_RESOURCE) - return fmt.Sprintf("resourceGroups/%s/providers/eip/%s", resourceGroup, eipName) + globalId, _, _ := pareResourceGroupWithName(self.ID, EIP_RESOURCE) + return globalId } func (self *SEipAddress) GetId() string { diff --git a/pkg/util/azure/host.go b/pkg/util/azure/host.go index a2af38e019..4390b516d3 100644 --- a/pkg/util/azure/host.go +++ b/pkg/util/azure/host.go @@ -3,6 +3,7 @@ package azure import ( "context" "fmt" + "strings" "time" "github.com/Azure/azure-sdk-for-go/services/compute/mgmt/2018-06-01/compute" @@ -49,8 +50,7 @@ func (self *SHost) CreateVM(name string, imgId string, sysDiskSize int, cpu int, if err != nil { return nil, err } - resourceGroup, instanceName := PareResourceGroupWithName(vmId, INSTANCE_RESOURCE) - if vm, err := self.zone.region.GetInstance(resourceGroup, instanceName); err != nil { + if vm, err := self.zone.region.GetInstance(vmId); err != nil { return nil, err } else { return vm, err @@ -93,12 +93,12 @@ func (self *SHost) _createVM(name string, imgId string, sysDiskSize int, cpu int for i := 0; i < len(diskSizes); i++ { diskName := fmt.Sprintf("vdisk_%s_%d", name, time.Now().UnixNano()) size := int32(diskSizes[i] >> 10) - index := int32(i) + lun := int32(i) DataDisks = append(DataDisks, compute.DataDisk{ Name: &diskName, DiskSizeGB: &size, CreateOption: compute.Empty, - Lun: &index, + Lun: &lun, }) } @@ -150,12 +150,18 @@ func (self *SHost) _createVM(name string, imgId string, sysDiskSize int, cpu int for _, profile := range self.zone.region.getHardwareProfile(cpu, memMB) { params.HardwareProfile.VMSize = compute.VirtualMachineSizeTypes(profile) log.Debugf("Try HardwareProfile : %s", profile) - resourceGroup, instanceName := PareResourceGroupWithName(name, INSTANCE_RESOURCE) + instanceId, resourceGroup, instanceName := pareResourceGroupWithName(name, INSTANCE_RESOURCE) result, err := computeClient.CreateOrUpdate(context.Background(), resourceGroup, instanceName, params) if err != nil { log.Errorf("Failed for %s: %s", profile, err) } else if err := result.WaitForCompletion(context.Background(), computeClient.Client); err != nil { - return "", err + if strings.Index(err.Error(), "OSProvisioningTimedOut") == -1 { + return "", err + } else if instance, err := self.zone.region.GetInstance(instanceId); err != nil { + return "", err + } else { + return instance.ID, nil + } } else if vm, err := result.Result(computeClient); err != nil { return "", err } else { @@ -218,9 +224,8 @@ func (self *SHost) GetIStorages() ([]cloudprovider.ICloudStorage, error) { return self.zone.GetIStorages() } -func (self *SHost) GetIVMById(gid string) (cloudprovider.ICloudVM, error) { - resourceGroup, instanceName := PareResourceGroupWithName(gid, INSTANCE_RESOURCE) - if instance, err := self.zone.region.GetInstance(resourceGroup, instanceName); err != nil { +func (self *SHost) GetIVMById(instanceId string) (cloudprovider.ICloudVM, error) { + if instance, err := self.zone.region.GetInstance(instanceId); err != nil { return nil, err } else { instance.host = self diff --git a/pkg/util/azure/image.go b/pkg/util/azure/image.go index df5474f9cb..ee49d984b6 100644 --- a/pkg/util/azure/image.go +++ b/pkg/util/azure/image.go @@ -2,7 +2,6 @@ package azure import ( "context" - "fmt" "github.com/Azure/azure-sdk-for-go/services/compute/mgmt/2018-06-01/compute" "yunion.io/x/jsonutils" @@ -81,8 +80,8 @@ func (self *SImage) IsEmulated() bool { } func (self *SImage) GetGlobalId() string { - resourceGroup, imageName := PareResourceGroupWithName(self.ID, IMAGE_RESOURCE) - return fmt.Sprintf("resourceGroups/%s/providers/image/%s", resourceGroup, imageName) + globalId, _, _ := pareResourceGroupWithName(self.ID, IMAGE_RESOURCE) + return globalId } func (self *SImage) GetStatus() string { @@ -124,7 +123,7 @@ func (self *SRegion) GetImage(imageId string) (*SImage, error) { image := SImage{} imageClient := compute.NewImagesClientWithBaseURI(self.client.baseUrl, self.SubscriptionID) imageClient.Authorizer = self.client.authorizer - resourceGroup, imageName := PareResourceGroupWithName(imageId, IMAGE_RESOURCE) + _, resourceGroup, imageName := pareResourceGroupWithName(imageId, IMAGE_RESOURCE) if result, err := imageClient.Get(context.Background(), resourceGroup, imageName, ""); err != nil { if result.Response.StatusCode == 404 { return nil, cloudprovider.ErrNotFound @@ -161,7 +160,7 @@ func (self *SRegion) CreateImageByBlob(imageName, osType, blobURI string, diskSi StorageProfile: &storageProfile, }, } - resourceGroup, imageName := PareResourceGroupWithName(imageName, IMAGE_RESOURCE) + _, resourceGroup, imageName := pareResourceGroupWithName(imageName, IMAGE_RESOURCE) if result, err := imageClient.CreateOrUpdate(context.Background(), resourceGroup, imageName, params); err != nil { log.Errorf("Create image from blob error: %v", err) return nil, err @@ -190,7 +189,7 @@ func (self *SRegion) GetImages() ([]SImage, error) { func (self *SRegion) DeleteImage(imageId string) error { imageClient := compute.NewImagesClientWithBaseURI(self.client.baseUrl, self.SubscriptionID) imageClient.Authorizer = self.client.authorizer - resourceGroup, imageName := PareResourceGroupWithName(imageId, IMAGE_RESOURCE) + _, resourceGroup, imageName := pareResourceGroupWithName(imageId, IMAGE_RESOURCE) if result, err := imageClient.Delete(context.Background(), resourceGroup, imageName); err != nil { return err } else if err := result.WaitForCompletion(context.Background(), imageClient.Client); err != nil { diff --git a/pkg/util/azure/instance.go b/pkg/util/azure/instance.go index 9f666cf1ef..952be055bb 100644 --- a/pkg/util/azure/instance.go +++ b/pkg/util/azure/instance.go @@ -3,7 +3,6 @@ package azure import ( "context" "fmt" - "regexp" "strings" "time" @@ -11,6 +10,7 @@ import ( "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/util/seclib2" "yunion.io/x/pkg/util/osprofile" "yunion.io/x/pkg/util/secrules" @@ -165,23 +165,15 @@ type SInstance struct { Tags map[string]string } -func PareResourceGroupWithName(s string, resourceType string) (string, string) { - valid := regexp.MustCompile("resourceGroups/(.+)/providers/.+/(.+)$") - if resourceGroups := valid.FindStringSubmatch(s); len(resourceGroups) == 3 { - return resourceGroups[1], resourceGroups[2] - } - log.Errorf("PareResourceGroupWithName[%s] error", s) - return DefaultResourceGroups[resourceType], s -} - -func (self *SRegion) GetInstance(resourceGroup string, VMName string) (*SInstance, error) { +func (self *SRegion) GetInstance(instanceId string) (*SInstance, error) { instance := SInstance{} computeClient := compute.NewVirtualMachinesClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) computeClient.Authorizer = self.client.authorizer - if len(VMName) == 0 { + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + if len(instanceName) == 0 { return nil, cloudprovider.ErrNotFound } - if _instance, err := computeClient.Get(context.Background(), resourceGroup, VMName, "instanceView"); err != nil { + if _instance, err := computeClient.Get(context.Background(), resourceGroup, instanceName, "instanceView"); err != nil { if _instance.Response.StatusCode == 404 { return nil, cloudprovider.ErrNotFound } @@ -214,7 +206,7 @@ func (self *SRegion) GetInstances() ([]SInstance, error) { } func (self *SRegion) doDeleteVM(instanceId string) error { - resourceGroup, instanceName := PareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) computeClient := compute.NewVirtualMachinesClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) computeClient.Authorizer = self.client.authorizer if resulte, err := computeClient.Delete(context.Background(), resourceGroup, instanceName); err != nil { @@ -251,17 +243,16 @@ func (self *SInstance) IsEmulated() bool { func (self *SInstance) getDisks() ([]SDisk, error) { disks := make([]SDisk, 0) - resourceGroup, diskName := PareResourceGroupWithName(self.Properties.StorageProfile.OsDisk.ManagedDisk.ID, DISK_RESOURCE) - if osdisk, err := self.getDiskWithStore(resourceGroup, diskName); err != nil { - log.Errorf("Failed to find instance %s os disk: %s", self.Name, diskName) + diskId := self.Properties.StorageProfile.OsDisk.ManagedDisk.ID + if osdisk, err := self.getDiskWithStore(diskId); err != nil { + log.Errorf("Failed to find instance %s os disk: %s", self.Name, diskId) } else { disks = append(disks, *osdisk) } for _, _disk := range self.Properties.StorageProfile.DataDisks { - resourceGroup, diskName := PareResourceGroupWithName(_disk.ManagedDisk.ID, DISK_RESOURCE) - if disk, err := self.getDiskWithStore(resourceGroup, diskName); err != nil { - log.Errorf("Failed to find instance %s data disk: %s", self.Name, diskName) + if disk, err := self.getDiskWithStore(_disk.ManagedDisk.ID); err != nil { + log.Errorf("Failed to find instance %s data disk: %s", self.Name, _disk.ManagedDisk.ID) return nil, err } else { disks = append(disks, *disk) @@ -273,9 +264,8 @@ func (self *SInstance) getDisks() ([]SDisk, error) { func (self *SInstance) getNics() ([]SInstanceNic, error) { nics := make([]SInstanceNic, 0) for _, _nic := range self.Properties.NetworkProfile.NetworkInterfaces { - resourceGroup, nicName := PareResourceGroupWithName(_nic.ID, NIC_RESOURCE) - if nic, err := self.host.zone.region.getNetworkInterface(resourceGroup, nicName); err != nil { - log.Errorf("Failed to find instance %s nic: %s", self.Name, nicName) + if nic, err := self.host.zone.region.getNetworkInterface(_nic.ID); err != nil { + log.Errorf("Failed to find instance %s nic: %s", self.Name, _nic.ID) return nil, err } else { nic.instance = self @@ -286,8 +276,7 @@ func (self *SInstance) getNics() ([]SInstanceNic, error) { } func (self *SInstance) Refresh() error { - resourceGroup, instanceName := PareResourceGroupWithName(self.ID, INSTANCE_RESOURCE) - if instance, err := self.host.zone.region.GetInstance(resourceGroup, instanceName); err != nil { + if instance, err := self.host.zone.region.GetInstance(self.ID); err != nil { return err } else if err := jsonutils.Update(self, instance); err != nil { return err @@ -296,9 +285,6 @@ func (self *SInstance) Refresh() error { } func (self *SInstance) GetStatus() string { - if len(self.Properties.InstanceView.Statuses) == 0 { - self.Refresh() - } for _, statuses := range self.Properties.InstanceView.Statuses { if code := strings.Split(statuses.Code, "/"); len(code) == 2 { if code[0] == "PowerState" { @@ -324,27 +310,226 @@ func (self *SInstance) GetIHost() cloudprovider.ICloudHost { } func (self *SInstance) AttachDisk(diskId string) error { + return self.host.zone.region.AttachDisk(self.ID, diskId) +} + +func (region *SRegion) UpdateInstance(instanceId string, params compute.VirtualMachineUpdate) error { + computeClient := compute.NewVirtualMachinesClientWithBaseURI(region.client.baseUrl, region.client.subscriptionId) + computeClient.Authorizer = region.client.authorizer + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + if result, err := computeClient.Update(context.Background(), resourceGroup, instanceName, params); err != nil { + return err + } else if err := result.WaitForCompletion(context.Background(), computeClient.Client); err != nil { + return err + } return nil } +func (region *SRegion) AttachDisk(instanceId, diskId string) error { + if instance, err := region.GetInstance(instanceId); err != nil { + return err + } else if disk, err := region.GetDisk(diskId); err != nil { + return err + } else { + dataDisks := []compute.DataDisk{} + maxLun := int32(0) + for i := 0; i < len(instance.Properties.StorageProfile.DataDisks); i++ { + _disk := instance.Properties.StorageProfile.DataDisks[i] + if disk.ID == _disk.ManagedDisk.ID { + return nil + } else { + if maxLun < _disk.Lun { + maxLun = _disk.Lun + } + dataDisks = append(dataDisks, compute.DataDisk{ + Lun: &_disk.Lun, + CreateOption: compute.DiskCreateOptionTypesAttach, + ManagedDisk: &compute.ManagedDiskParameters{ + ID: &_disk.ManagedDisk.ID, + }, + }) + } + } + maxLun++ + dataDisks = append(dataDisks, compute.DataDisk{ + Lun: &maxLun, + CreateOption: compute.DiskCreateOptionTypesAttach, + ManagedDisk: &compute.ManagedDiskParameters{ + ID: &disk.ID, + }, + }) + params := compute.VirtualMachineUpdate{ + VirtualMachineProperties: &compute.VirtualMachineProperties{ + StorageProfile: &compute.StorageProfile{ + DataDisks: &dataDisks, + }, + }, + } + return region.UpdateInstance(instanceId, params) + } +} + func (self *SInstance) DetachDisk(diskId string) error { - return nil + return self.host.zone.region.DetachDisk(self.ID, diskId) +} + +func (region *SRegion) DetachDisk(instanceId, diskId string) error { + if instance, err := region.GetInstance(instanceId); err != nil { + return err + } else if disk, err := region.GetDisk(diskId); err != nil { + return err + } else { + dataDisks := []compute.DataDisk{} + for i := 0; i < len(instance.Properties.StorageProfile.DataDisks); i++ { + if instance.Properties.StorageProfile.DataDisks[i].ManagedDisk.ID == disk.ID { + continue + } + dataDisks = append(dataDisks, compute.DataDisk{ + Lun: &instance.Properties.StorageProfile.DataDisks[i].Lun, + ManagedDisk: &compute.ManagedDiskParameters{ + ID: &instance.Properties.StorageProfile.DataDisks[i].ManagedDisk.ID, + }, + }) + } + params := compute.VirtualMachineUpdate{ + VirtualMachineProperties: &compute.VirtualMachineProperties{ + StorageProfile: &compute.StorageProfile{ + DataDisks: &dataDisks, + }, + }, + } + return region.UpdateInstance(instanceId, params) + } } func (self *SInstance) ChangeConfig(instanceId string, ncpu int, vmem int) error { - return nil + return self.host.zone.region.ChangeVMConfig(instanceId, ncpu, vmem) +} + +func (region *SRegion) ChangeVMConfig(instanceId string, ncpu int, vmem int) error { + for _, vmSize := range region.getHardwareProfile(ncpu, vmem) { + params := compute.VirtualMachineUpdate{ + VirtualMachineProperties: &compute.VirtualMachineProperties{ + HardwareProfile: &compute.HardwareProfile{ + VMSize: compute.VirtualMachineSizeTypes(vmSize), + }, + }, + } + log.Debugf("Try HardwareProfile : %s", vmSize) + if err := region.UpdateInstance(instanceId, params); err == nil { + return nil + } + + } + return fmt.Errorf("Failed to change vm config, specification not supported") } func (self *SInstance) DeployVM(name string, password string, publicKey string, resetPassword bool, deleteKeypair bool, description string) error { + return self.host.zone.region.DeployVM(self.ID, name, password, publicKey, resetPassword, deleteKeypair, description) +} + +type VirtualMachineExtensionProperties struct { + Publisher string + Type string + TypeHandlerVersion string +} + +type SVirtualMachineExtension struct { + Location string + Properties VirtualMachineExtensionProperties +} + +func (region *SRegion) resetLoginInfo(instanceId string, setting map[string]string) error { + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + extensionClient := compute.NewVirtualMachineExtensionsClientWithBaseURI(region.client.baseUrl, region.SubscriptionID) + extensionClient.Authorizer = region.client.authorizer + extension := SVirtualMachineExtension{} + if result, err := extensionClient.Get(context.Background(), resourceGroup, instanceName, "enablevmaccess", ""); err != nil { + return err + } else if err := jsonutils.Update(&extension, result); err != nil { + return err + } + + params := compute.VirtualMachineExtension{ + Location: ®ion.Name, + VirtualMachineExtensionProperties: &compute.VirtualMachineExtensionProperties{ + Publisher: &extension.Properties.Publisher, + Type: &extension.Properties.Type, + TypeHandlerVersion: &extension.Properties.TypeHandlerVersion, + ProtectedSettings: setting, + }, + } + + if result, err := extensionClient.CreateOrUpdate(context.Background(), resourceGroup, instanceName, "enablevmaccess", params); err != nil { + return err + } else if err := result.WaitForCompletion(context.Background(), extensionClient.Client); err != nil { + return err + } return nil } +func (region *SRegion) resetPublicKey(instanceId string, username, publicKey string) error { + setting := map[string]string{ + "username": username, + "ssh_key": publicKey, + } + return region.resetLoginInfo(instanceId, setting) +} + +func (region *SRegion) resetPassword(instanceId, username, password string) error { + setting := map[string]string{ + "username": username, + "password": password, + } + return region.resetLoginInfo(instanceId, setting) +} + +func (region *SRegion) DeployVM(instanceId, name, password, publicKey string, resetPassword bool, deleteKeypair bool, description string) error { + if instance, err := region.GetInstance(instanceId); err != nil { + return err + } else { + if deleteKeypair { + return nil + } + if len(publicKey) > 0 { + return region.resetPublicKey(instanceId, instance.Properties.OsProfile.AdminUsername, publicKey) + } else if resetPassword { + if len(password) == 0 { + password = seclib2.RandomPassword2(12) + } + return region.resetPassword(instanceId, instance.Properties.OsProfile.AdminUsername, password) + } + return nil + } +} + func (self *SInstance) RebuildRoot(imageId string) error { + return self.host.zone.region.RebuildRoot(self.ID) +} + +func (region *SRegion) RebuildRoot(instanceId string) error { + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + computeClient := compute.NewVirtualMachinesClientWithBaseURI(region.client.baseUrl, region.client.subscriptionId) + computeClient.Authorizer = region.client.authorizer + if result, err := computeClient.Redeploy(context.Background(), resourceGroup, instanceName); err != nil { + return err + } else if err := result.WaitForCompletion(context.Background(), computeClient.Client); err != nil { + if strings.Index(err.Error(), "OSProvisioningTimedOut") > 0 { + if instance, err := region.GetInstance(instanceId); err != nil { + return err + } else if status := instance.GetStatus(); status == models.VM_RUNNING { + region.StopVM(instanceId, true) + } + return nil + } else { + return err + } + } return nil } func (self *SInstance) UpdateVM(name string) error { - return nil + return cloudprovider.ErrNotSupported } func (self *SInstance) GetId() string { @@ -356,13 +541,12 @@ func (self *SInstance) GetName() string { } func (self *SInstance) GetGlobalId() string { - resourceGroup, instanceName := PareResourceGroupWithName(self.ID, INSTANCE_RESOURCE) - return fmt.Sprintf("resourceGroups/%s/providers/server/%s", resourceGroup, instanceName) + globalId, _, _ := pareResourceGroupWithName(self.ID, INSTANCE_RESOURCE) + return globalId } func (self *SRegion) GetInstanceStatus(instanceId string) (string, error) { - resourceGroup, instanceName := PareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) - instance, err := self.GetInstance(resourceGroup, instanceName) + instance, err := self.GetInstance(self.ID) if err != nil { return "", err } @@ -377,14 +561,12 @@ func (self *SRegion) DeleteVM(instanceId string) error { return err } } else if status != models.VM_READY { - log.Debugf("instance %s status: %s", instanceId, status) return cloudprovider.ErrInvalidStatus } return self.doDeleteVM(instanceId) } func (self *SInstance) DeleteVM() error { - log.Debugf("delete: %s %s", self.ID, self.Name) if err := self.host.zone.region.DeleteVM(self.ID); err != nil { return err } @@ -409,8 +591,8 @@ func (self *SInstance) DeleteVM() error { return nil } -func (self *SInstance) getDiskWithStore(resourceGroup string, diskName string) (*SDisk, error) { - if disk, err := self.host.zone.region.GetDisk(resourceGroup, diskName); err != nil { +func (self *SInstance) getDiskWithStore(diskId string) (*SDisk, error) { + if disk, err := self.host.zone.region.GetDisk(diskId); err != nil { return nil, err } else if store, err := self.host.zone.getStorageByType(string(disk.Sku.Name)); err != nil { log.Errorf("fail to find storage for disk(%s) : %v", disk.Name, err) @@ -448,18 +630,14 @@ func (self *SInstance) GetOSType() string { func (self *SInstance) GetINics() ([]cloudprovider.ICloudNic, error) { nics := make([]cloudprovider.ICloudNic, 0) - for _, _nic := range self.Properties.NetworkProfile.NetworkInterfaces { - resourceGroup, nicName := PareResourceGroupWithName(_nic.ID, NIC_RESOURCE) - if nic, err := self.host.zone.region.getNetworkInterface(resourceGroup, nicName); err != nil { - return nics, err - } else { - nic.instance = self - nics = append(nics, nic) + if _nics, err := self.getNics(); err != nil { + return nil, err + } else { + for i := 0; i < len(_nics); i++ { + _nics[i].instance = self + nics = append(nics, &_nics[i]) } } - for _, nic := range nics { - log.Debugf("find nic %s for instance %s", nic.GetIP(), self.Name) - } return nics, nil } @@ -497,22 +675,12 @@ func (self *SInstance) fetchVMSize() error { } func (self *SInstance) GetVcpuCount() int8 { - if self.vmSize == nil { - if err := self.fetchVMSize(); err != nil { - log.Errorf("fail to fetch vmSize: %v", err) - return 0 - } - } + self.fetchVMSize() return int8(self.vmSize.NumberOfCores) } func (self *SInstance) GetVmemSizeMB() int { - if self.vmSize == nil { - if err := self.fetchVMSize(); err != nil { - log.Errorf("fail to fetch vmSize: %v", err) - return 0 - } - } + self.fetchVMSize() return int(self.vmSize.MemoryInMB) } @@ -520,18 +688,13 @@ func (self *SInstance) GetCreateTime() time.Time { return self.CreationTime } -func (self *SInstance) GetEIP() cloudprovider.ICloudEIP { - return nil - //return &self.EipAddress -} - func (self *SInstance) GetVNCInfo() (jsonutils.JSONObject, error) { ret := jsonutils.NewDict() return ret, nil } func (self *SRegion) StartVM(instanceId string) error { - resourceGroup, instanceName := PareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) computeClient := compute.NewVirtualMachinesClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) computeClient.Authorizer = self.client.authorizer if result, err := computeClient.Start(context.Background(), resourceGroup, instanceName); err != nil { @@ -561,7 +724,7 @@ func (self *SRegion) StopVM(instanceId string, isForce bool) error { } func (self *SRegion) doStopVM(instanceId string, isForce bool) error { - resourceGroup, instanceName := PareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) + _, resourceGroup, instanceName := pareResourceGroupWithName(instanceId, INSTANCE_RESOURCE) computeClient := compute.NewVirtualMachinesClientWithBaseURI(self.client.baseUrl, self.client.subscriptionId) computeClient.Authorizer = self.client.authorizer if result, err := computeClient.PowerOff(context.Background(), resourceGroup, instanceName); err != nil { diff --git a/pkg/util/azure/instancenic.go b/pkg/util/azure/instancenic.go index ff61ed57be..226df9930c 100644 --- a/pkg/util/azure/instancenic.go +++ b/pkg/util/azure/instancenic.go @@ -52,7 +52,7 @@ func (self *SInstanceNic) GetIP() string { } func (region *SRegion) DeleteNetworkInterface(interfaceId string) error { - resourceGroup, nicName := PareResourceGroupWithName(interfaceId, NIC_RESOURCE) + _, resourceGroup, nicName := pareResourceGroupWithName(interfaceId, NIC_RESOURCE) networkClient := network.NewInterfacesClientWithBaseURI(region.client.baseUrl, region.SubscriptionID) networkClient.Authorizer = region.client.authorizer if result, err := networkClient.Delete(context.Background(), resourceGroup, nicName); err != nil { diff --git a/pkg/util/azure/network.go b/pkg/util/azure/network.go index 4abe20b91b..4eb311f473 100644 --- a/pkg/util/azure/network.go +++ b/pkg/util/azure/network.go @@ -2,7 +2,6 @@ package azure import ( "context" - "fmt" "strings" "github.com/Azure/azure-sdk-for-go/services/network/mgmt/2018-06-01/network" @@ -36,8 +35,8 @@ func (self *SNetwork) GetName() string { } func (self *SNetwork) GetGlobalId() string { - resourceGroup, networkName := PareResourceGroupWithName(self.ID, NETWORK_RESOURCE) - return fmt.Sprintf("resourceGroups/%s/providers/network/%s", resourceGroup, networkName) + globalId, _, _ := pareResourceGroupWithName(self.ID, NETWORK_RESOURCE) + return globalId } func (self *SNetwork) IsEmulated() bool { @@ -71,7 +70,7 @@ func (self *SNetwork) Delete() error { region := self.wire.vpc.region networkClient := network.NewVirtualNetworksClientWithBaseURI(region.client.baseUrl, region.SubscriptionID) networkClient.Authorizer = region.client.authorizer - resourceGroup, vpcName := PareResourceGroupWithName(vpc.ID, VPC_RESOURCE) + _, resourceGroup, vpcName := pareResourceGroupWithName(vpc.ID, VPC_RESOURCE) if result, err := networkClient.CreateOrUpdate(context.Background(), resourceGroup, vpcName, params); err != nil { return err } else if err := result.WaitForCompletion(context.Background(), networkClient.Client); err != nil { diff --git a/pkg/util/azure/networkinterface.go b/pkg/util/azure/networkinterface.go index e5337be5a1..cf28b9fb44 100644 --- a/pkg/util/azure/networkinterface.go +++ b/pkg/util/azure/networkinterface.go @@ -8,10 +8,11 @@ import ( "yunion.io/x/jsonutils" ) -func (self *SRegion) getNetworkInterface(resourceGroup string, nicName string) (*SInstanceNic, error) { +func (self *SRegion) getNetworkInterface(interfaceId string) (*SInstanceNic, error) { nic := SInstanceNic{} networkClient := network.NewInterfacesClientWithBaseURI(self.client.baseUrl, self.SubscriptionID) networkClient.Authorizer = self.client.authorizer + _, resourceGroup, nicName := pareResourceGroupWithName(interfaceId, NIC_RESOURCE) if _nic, err := networkClient.Get(context.Background(), resourceGroup, nicName, ""); err != nil { return nil, err } else if err := jsonutils.Update(&nic, _nic); err != nil { @@ -35,7 +36,7 @@ func (self *SRegion) GetNetworkInterfaces() ([]SInstanceNic, error) { func (self *SRegion) isNetworkInstanceNameAvaliable(nicName string) bool { networkClinet := network.NewInterfacesClientWithBaseURI(self.client.baseUrl, self.SubscriptionID) networkClinet.Authorizer = self.client.authorizer - resourceGroup, nicName := PareResourceGroupWithName(nicName, NIC_RESOURCE) + _, resourceGroup, nicName := pareResourceGroupWithName(nicName, NIC_RESOURCE) if result, err := networkClinet.Get(context.Background(), resourceGroup, nicName, ""); err != nil || result.Response.StatusCode == 404 { return true } @@ -81,7 +82,7 @@ func (self *SRegion) CreateNetworkInterface(nicName string, ipAddr string, subne }, } //log.Debugf("create params: %", jsonutils.Marshal(params).PrettyString()) - resourceGroup, nicName := PareResourceGroupWithName(nicName, NIC_RESOURCE) + _, resourceGroup, nicName := pareResourceGroupWithName(nicName, NIC_RESOURCE) if result, err := networkClinet.CreateOrUpdate(context.Background(), resourceGroup, nicName, params); err != nil { return nil, err } else if err := result.WaitForCompletion(context.Background(), networkClinet.Client); err != nil { diff --git a/pkg/util/azure/region.go b/pkg/util/azure/region.go index e63e3bf2d4..2da24cfc46 100644 --- a/pkg/util/azure/region.go +++ b/pkg/util/azure/region.go @@ -138,8 +138,7 @@ func (self *SRegion) CreateIVpc(name string, desc string, cidr string) (cloudpro addressSpace := network.AddressSpace{AddressPrefixes: &addressPrefixes} properties := network.VirtualNetworkPropertiesFormat{AddressSpace: &addressSpace} parameters := network.VirtualNetwork{Name: &name, Location: &self.Name, VirtualNetworkPropertiesFormat: &properties} - resourceGroup, vpcName := PareResourceGroupWithName(name, VPC_RESOURCE) - vpcId := fmt.Sprintf("resourceGroups/%s/providers/vpc/%s", resourceGroup, vpcName) + vpcId, resourceGroup, vpcName := pareResourceGroupWithName(name, VPC_RESOURCE) if result, err := vpcClient.CreateOrUpdate(context.Background(), resourceGroup, vpcName, parameters); err != nil { return nil, err } else if err := result.WaitForCompletion(context.Background(), vpcClient.Client); err != nil { @@ -197,10 +196,10 @@ func (self *SRegion) GetIVpcById(id string) (cloudprovider.ICloudVpc, error) { if ivpcs, err := self.GetIVpcs(); err != nil { return nil, err } else { - resourceGroup, vpcName := PareResourceGroupWithName(id, VPC_RESOURCE) + _, resourceGroup, vpcName := pareResourceGroupWithName(id, VPC_RESOURCE) for i := 0; i < len(ivpcs); i++ { vpcId := ivpcs[i].GetId() - _resourceGroup, _vpcName := PareResourceGroupWithName(vpcId, VPC_RESOURCE) + _, _resourceGroup, _vpcName := pareResourceGroupWithName(vpcId, VPC_RESOURCE) if _resourceGroup == resourceGroup && _vpcName == vpcName { return ivpcs[i], nil } diff --git a/pkg/util/azure/securitygroup.go b/pkg/util/azure/securitygroup.go index 0b8deb200f..c84756cddf 100644 --- a/pkg/util/azure/securitygroup.go +++ b/pkg/util/azure/securitygroup.go @@ -169,7 +169,7 @@ func (self *SSecurityGroup) IsEmulated() bool { } func (self *SSecurityGroup) Refresh() error { - resourceGroup, secgrpName := PareResourceGroupWithName(self.ID, SECGRP_RESOURCE) + _, resourceGroup, secgrpName := pareResourceGroupWithName(self.ID, SECGRP_RESOURCE) networkClient := network.NewSecurityGroupsClientWithBaseURI(self.vpc.region.client.baseUrl, self.vpc.region.SubscriptionID) networkClient.Authorizer = self.vpc.region.client.authorizer if secgrp, err := networkClient.Get(context.Background(), resourceGroup, secgrpName, ""); err != nil { diff --git a/pkg/util/azure/shell/disk.go b/pkg/util/azure/shell/disk.go index 6c802301bf..2a08538479 100644 --- a/pkg/util/azure/shell/disk.go +++ b/pkg/util/azure/shell/disk.go @@ -7,9 +7,6 @@ import ( func init() { type DiskListOptions struct { - // Instance string `help:"Instance ID"` - // Zone string `help:"Zone ID"` - // Category string `help:"Disk category"` Offset int `help:"List offset"` Limit int `help:"List limit"` } @@ -31,10 +28,9 @@ func init() { } shellutils.R(&DiskCreateOptions{}, "disk-create", "Create disk", func(cli *azure.SRegion, args *DiskCreateOptions) error { - resourceGroup, diskName := azure.PareResourceGroupWithName(args.NAME, azure.DISK_RESOURCE) - if err := cli.CreateDisk(args.StorageType, args.NAME, args.SizeGb, args.Desc); err != nil { + if diskId, err := cli.CreateDisk(args.StorageType, args.NAME, args.SizeGb, args.Desc); err != nil { return err - } else if disk, err := cli.GetDisk(resourceGroup, diskName); err != nil { + } else if disk, err := cli.GetDisk(diskId); err != nil { return err } else { printObject(disk) @@ -42,11 +38,20 @@ func init() { return nil }) - type DiskDeleteOptions struct { + type DiskOptions struct { ID string `help:"Disk ID"` } - shellutils.R(&DiskDeleteOptions{}, "disk-delete", "Delete disks", func(cli *azure.SRegion, args *DiskDeleteOptions) error { + shellutils.R(&DiskOptions{}, "disk-show", "Show disk", func(cli *azure.SRegion, args *DiskOptions) error { + if disk, err := cli.GetDisk(args.ID); err != nil { + return err + } else { + printObject(disk) + return nil + } + }) + + shellutils.R(&DiskOptions{}, "disk-delete", "Delete disks", func(cli *azure.SRegion, args *DiskOptions) error { return cli.DeleteDisk(args.ID) }) diff --git a/pkg/util/azure/shell/eip.go b/pkg/util/azure/shell/eip.go index 8119d58e49..8c943453f1 100644 --- a/pkg/util/azure/shell/eip.go +++ b/pkg/util/azure/shell/eip.go @@ -58,7 +58,12 @@ func init() { err := cli.AssociateEip(args.ID, args.INSTANCE) return err }) - shellutils.R(&EipAssociateOptions{}, "eip-dissociate", "Dissociate an EIP", func(cli *azure.SRegion, args *EipAssociateOptions) error { + + type EipDissociateOptions struct { + ID string `help:"EIP allocation ID"` + } + + shellutils.R(&EipDissociateOptions{}, "eip-dissociate", "Dissociate an EIP", func(cli *azure.SRegion, args *EipDissociateOptions) error { err := cli.DissociateEip(args.ID) return err }) diff --git a/pkg/util/azure/shell/instance.go b/pkg/util/azure/shell/instance.go index 6e68a1eeaf..39260bdeb1 100644 --- a/pkg/util/azure/shell/instance.go +++ b/pkg/util/azure/shell/instance.go @@ -43,8 +43,7 @@ func init() { ID string `help:"Instance ID"` } shellutils.R(&InstanceShowOptions{}, "instance-show", "Show intance detail", func(cli *azure.SRegion, args *InstanceShowOptions) error { - resourceGroup, instanceName := azure.PareResourceGroupWithName(args.ID, azure.INSTANCE_RESOURCE) - if instance, err := cli.GetInstance(resourceGroup, instanceName); err != nil { + if instance, err := cli.GetInstance(args.ID); err != nil { return err } else { printObject(instance) @@ -52,4 +51,42 @@ func init() { } }) + type InstanceRebuildOptions struct { + ID string `help:"Instance ID"` + } + shellutils.R(&InstanceRebuildOptions{}, "instance-rebuild", "Rebuild intance root", func(cli *azure.SRegion, args *InstanceRebuildOptions) error { + return cli.RebuildRoot(args.ID) + }) + + type InstanceDiskOptions struct { + ID string `help:"Instance ID"` + DISK string `help:"Disk ID"` + } + shellutils.R(&InstanceDiskOptions{}, "instance-attach-disk", "Attach a disk to intance", func(cli *azure.SRegion, args *InstanceDiskOptions) error { + return cli.AttachDisk(args.ID, args.DISK) + }) + + shellutils.R(&InstanceDiskOptions{}, "instance-detach-disk", "Attach a disk to intance", func(cli *azure.SRegion, args *InstanceDiskOptions) error { + return cli.DetachDisk(args.ID, args.DISK) + }) + + type InstanceConfigOptions struct { + ID string `help:"Instance ID"` + NCPU int `help:"Number of cpu core"` + MEMERY int `helo:"Instance memery in mb"` + } + + shellutils.R(&InstanceConfigOptions{}, "instance-change-conf", "Attach a disk to intance", func(cli *azure.SRegion, args *InstanceConfigOptions) error { + return cli.ChangeVMConfig(args.ID, args.NCPU, args.MEMERY) + }) + + type InstanceDeployOptions struct { + ID string `help:"Instance ID"` + Password string `help:"Password for instance"` + PublicKey string `helo:"Deploy ssh_key for instance"` + } + + shellutils.R(&InstanceDeployOptions{}, "instance-reset-password", "Reset intance password", func(cli *azure.SRegion, args *InstanceDeployOptions) error { + return cli.DeployVM(args.ID, "", args.Password, args.PublicKey, true, false, "") + }) } diff --git a/pkg/util/azure/storage.go b/pkg/util/azure/storage.go index 4fcb024582..491544134e 100644 --- a/pkg/util/azure/storage.go +++ b/pkg/util/azure/storage.go @@ -49,10 +49,9 @@ func (self *SStorage) GetCapacityMB() int { } func (self *SStorage) CreateIDisk(name string, sizeGb int, desc string) (cloudprovider.ICloudDisk, error) { - resourceGroup, diskName := PareResourceGroupWithName(name, DISK_RESOURCE) - if err := self.zone.region.createDisk(self.storageType, diskName, int32(sizeGb), desc); err != nil { + if diskId, err := self.zone.region.createDisk(self.storageType, name, int32(sizeGb), desc); err != nil { return nil, err - } else if disk, err := self.zone.region.GetDisk(resourceGroup, diskName); err != nil { + } else if disk, err := self.zone.region.GetDisk(diskId); err != nil { return nil, err } else { disk.storage = self @@ -60,9 +59,8 @@ func (self *SStorage) CreateIDisk(name string, sizeGb int, desc string) (cloudpr } } -func (self *SStorage) GetIDisk(idStr string) (cloudprovider.ICloudDisk, error) { - resourceGroup, diskName := PareResourceGroupWithName(idStr, DISK_RESOURCE) - if disk, err := self.zone.region.GetDisk(resourceGroup, diskName); err != nil { +func (self *SStorage) GetIDisk(diskId string) (cloudprovider.ICloudDisk, error) { + if disk, err := self.zone.region.GetDisk(diskId); err != nil { return nil, err } else { disk.storage = self diff --git a/pkg/util/azure/storagecache.go b/pkg/util/azure/storagecache.go index 553321b25b..e33202887b 100644 --- a/pkg/util/azure/storagecache.go +++ b/pkg/util/azure/storagecache.go @@ -107,7 +107,7 @@ func (self *SRegion) CreateStorageAccount(resourceGroup, storageAccount string) sku := storageaccount.Sku{Name: storageaccount.SkuName("Standard_GRS")} params := storageaccount.AccountCreateParameters{Sku: &sku, Location: &self.Name, Kind: storageaccount.Kind("Storage")} if len(resourceGroup) == 0 { - resourceGroup = DefaultResourceGroups[STORAGE_RESOURCE] + resourceGroup = defaultResourceGroups[STORAGE_RESOURCE] } if len(storageAccount) == 0 { storageAccount = fmt.Sprintf("%s%s", self.Name, DefaultStorageAccount) @@ -160,7 +160,7 @@ func (self *SRegion) getStorageAccountKey(resourceGroup, storageAccount string) func (self *SRegion) CheckBlobContainer(resourceGroup, storageAccount, blobName string) error { if len(resourceGroup) == 0 { - resourceGroup = DefaultResourceGroups[STORAGE_RESOURCE] + resourceGroup = defaultResourceGroups[STORAGE_RESOURCE] } if len(storageAccount) == 0 { storageAccount = fmt.Sprintf("%s%s", self.Name, DefaultStorageAccount) @@ -251,7 +251,7 @@ func (self *SRegion) getContainerFiles(storageAccount, accessKey, containerName func (self *SRegion) ListContainerFiles(resourceGroup, storageAccount, blobName string) ([]Blob, error) { if len(resourceGroup) == 0 { - resourceGroup = DefaultResourceGroups[STORAGE_RESOURCE] + resourceGroup = defaultResourceGroups[STORAGE_RESOURCE] } if len(storageAccount) == 0 { storageAccount = fmt.Sprintf("%s%s", self.Name, DefaultStorageAccount) @@ -354,7 +354,7 @@ func (self *SRegion) uploadContainerFileByPath(storageAccount, accessKey, contai func (self *SRegion) UploadContainerFiles(resourceGroup, storageAccount, containerName, filePath string) (string, error) { if len(resourceGroup) == 0 { - resourceGroup = DefaultResourceGroups[STORAGE_RESOURCE] + resourceGroup = defaultResourceGroups[STORAGE_RESOURCE] } if len(storageAccount) == 0 { storageAccount = fmt.Sprintf("%s%s", self.Name, DefaultStorageAccount) @@ -395,12 +395,12 @@ func (self *SStoragecache) uploadImage(userCred mcclient.TokenCredential, imageI storageAccount := fmt.Sprintf("%s%s", self.region.Name, DefaultStorageAccount) - if err := self.region.CheckBlobContainer(DefaultResourceGroups[STORAGE_RESOURCE], storageAccount, DefaultBlobContainer); err != nil { + if err := self.region.CheckBlobContainer(defaultResourceGroups[STORAGE_RESOURCE], storageAccount, DefaultBlobContainer); err != nil { return "", err } size, _ := meta.Int("size") - accessKey, err := self.region.getStorageAccountKey(DefaultResourceGroups[STORAGE_RESOURCE], storageAccount) + accessKey, err := self.region.getStorageAccountKey(defaultResourceGroups[STORAGE_RESOURCE], storageAccount) if err != nil { return "", err } diff --git a/pkg/util/azure/vpc.go b/pkg/util/azure/vpc.go index 7af3ac2570..7e50cf305b 100644 --- a/pkg/util/azure/vpc.go +++ b/pkg/util/azure/vpc.go @@ -2,7 +2,6 @@ package azure import ( "context" - "fmt" "strings" "yunion.io/x/jsonutils" @@ -64,8 +63,8 @@ func (self *SVpc) GetName() string { } func (self *SVpc) GetGlobalId() string { - resourceGroup, vpcName := PareResourceGroupWithName(self.ID, VPC_RESOURCE) - return fmt.Sprintf("resourceGroups/%s/providers/vpc/%s", resourceGroup, vpcName) + globalId, _, _ := pareResourceGroupWithName(self.ID, VPC_RESOURCE) + return globalId } func (self *SVpc) IsEmulated() bool { @@ -83,7 +82,7 @@ func (self *SVpc) GetCidrBlock() string { func (self *SVpc) Delete() error { vpcClient := network.NewVirtualNetworksClientWithBaseURI(self.region.client.baseUrl, self.region.client.subscriptionId) vpcClient.Authorizer = self.region.client.authorizer - resourceGroup, vpcName := PareResourceGroupWithName(self.ID, VPC_RESOURCE) + _, resourceGroup, vpcName := pareResourceGroupWithName(self.ID, VPC_RESOURCE) if result, err := vpcClient.Delete(context.Background(), resourceGroup, vpcName); err != nil { return err } else if err := result.WaitForCompletion(context.Background(), vpcClient.Client); err != nil { @@ -196,7 +195,7 @@ func (self *SVpc) GetStatus() string { } func (self *SVpc) Refresh() error { - resourceGroup, vpcName := PareResourceGroupWithName(self.ID, VPC_RESOURCE) + _, resourceGroup, vpcName := pareResourceGroupWithName(self.ID, VPC_RESOURCE) vpcClient := network.NewVirtualNetworksClientWithBaseURI(self.region.client.baseUrl, self.region.SubscriptionID) vpcClient.Authorizer = self.region.client.authorizer if result, err := vpcClient.Get(context.Background(), resourceGroup, vpcName, ""); err != nil { diff --git a/pkg/util/azure/wire.go b/pkg/util/azure/wire.go index 34afd04920..b26d95bc1a 100644 --- a/pkg/util/azure/wire.go +++ b/pkg/util/azure/wire.go @@ -78,13 +78,13 @@ func (self *SRegion) createNetwork(vpc *SVpc, subnetName string, cidr string, de networkClient := network.NewVirtualNetworksClientWithBaseURI(self.client.baseUrl, self.SubscriptionID) networkClient.Authorizer = self.client.authorizer - resourceGroup, vpcName := PareResourceGroupWithName(vpc.ID, VPC_RESOURCE) + networkId, resourceGroup, vpcName := pareResourceGroupWithName(vpc.ID, VPC_RESOURCE) if result, err := networkClient.CreateOrUpdate(context.Background(), resourceGroup, vpcName, params); err != nil { return "", err } else if err := result.WaitForCompletion(context.Background(), networkClient.Client); err != nil { return "", err } - return fmt.Sprintf("/subscriptions/%s/resourceGroups/%s/providers/Microsoft.Network/virtualNetworks/%s/subnets/%s", self.SubscriptionID, resourceGroup, vpc.Name, subnetName), nil + return networkId, nil } func (self *SWire) CreateINetwork(name string, cidr string, desc string) (cloudprovider.ICloudNetwork, error) { @@ -141,11 +141,11 @@ func (self *SWire) getNetworkById(networkId string) *SNetwork { log.Errorf("getNetworkById error: %v", err) return nil } else { - resourceGroup, networkName := PareResourceGroupWithName(networkId, NETWORK_RESOURCE) + _, resourceGroup, networkName := pareResourceGroupWithName(networkId, NETWORK_RESOURCE) log.Debugf("search for networks %d", len(networks)) for i := 0; i < len(networks); i++ { network := networks[i].(*SNetwork) - _resourceGroup, _networkName := PareResourceGroupWithName(network.ID, NETWORK_RESOURCE) + _, _resourceGroup, _networkName := pareResourceGroupWithName(network.ID, NETWORK_RESOURCE) if resourceGroup == _resourceGroup && networkName == _networkName { return network } diff --git a/pkg/webconsole/command/ipmi_command.go b/pkg/webconsole/command/ipmi_command.go index 94e656b754..da250ba9f1 100644 --- a/pkg/webconsole/command/ipmi_command.go +++ b/pkg/webconsole/command/ipmi_command.go @@ -3,6 +3,8 @@ package command import ( "fmt" "os/exec" + + o "yunion.io/x/onecloud/pkg/webconsole/options" ) type IpmiInfo struct { @@ -27,7 +29,7 @@ func NewIpmitoolSolCommand(info *IpmiInfo) (*IpmitoolSol, error) { if info.Password == "" { return nil, fmt.Errorf("Empty password") } - name := "ipmitool" + name := o.Options.IpmitoolPath cmd := NewBaseCommand(name, "-I", "lanplus") cmd.AppendArgs("-H", info.IpAddr) cmd.AppendArgs("-U", info.Username) diff --git a/pkg/webconsole/command/kube_command.go b/pkg/webconsole/command/kube_command.go index 34a9087a6d..c90b9bae8e 100644 --- a/pkg/webconsole/command/kube_command.go +++ b/pkg/webconsole/command/kube_command.go @@ -6,6 +6,8 @@ import ( "os/exec" "yunion.io/x/log" + + o "yunion.io/x/onecloud/pkg/webconsole/options" ) type Kubectl struct { @@ -14,7 +16,7 @@ type Kubectl struct { } func NewKubectlCommand(kubeconfig, namespace string) *Kubectl { - name := "kubectl" + name := o.Options.KubectlPath if len(namespace) == 0 { namespace = "default" } @@ -92,7 +94,7 @@ func NewPodBashCommand(kubeconfig, namespace, pod, container string) ICommand { TTY(). Pod(pod). Container(container). - Command("bash", "-i", "-l") + Command("sh") } type KubectlLog struct { @@ -120,8 +122,18 @@ func (c *KubectlLog) Pod(name string) *KubectlLog { return c } +func (c *KubectlLog) Container(name string) *KubectlLog { + if name == "" { + return c + } + // -c, --container='': Print the logs of this container + c.AppendArgs("-c", name) + return c +} + func NewPodLogCommand(kubeconfig, namespace, pod, container string) ICommand { return NewKubectlCommand(kubeconfig, namespace).Logs(). Follow(). - Pod(pod) + Pod(pod). + Container(container) } diff --git a/pkg/webconsole/options/options.go b/pkg/webconsole/options/options.go index 7561ce79eb..45e9235f00 100644 --- a/pkg/webconsole/options/options.go +++ b/pkg/webconsole/options/options.go @@ -11,5 +11,7 @@ var ( type WebConsoleOptions struct { cloudcommon.Options - ApiServer string `help:"API server url to handle websocket connection, usually with public access" default:"http://webconsole.yunion.io"` + ApiServer string `help:"API server url to handle websocket connection, usually with public access" default:"http://webconsole.yunion.io"` + KubectlPath string `help:"kubectl binary path used to connect k8s cluster" default:"/usr/bin/kubectl"` + IpmitoolPath string `help:"ipmitool binary path used to connect baremetal sol" default:"/usr/bin/ipmitool"` } diff --git a/pkg/webconsole/service/service.go b/pkg/webconsole/service/service.go index 29dc081d04..3836f618bb 100644 --- a/pkg/webconsole/service/service.go +++ b/pkg/webconsole/service/service.go @@ -17,6 +17,12 @@ import ( "yunion.io/x/onecloud/pkg/webconsole/server" ) +func ensureBinExists(binPath string) { + if _, err := os.Stat(binPath); os.IsNotExist(err) { + log.Fatalf("Binary %s not exists", binPath) + } +} + func StartService() { cloudcommon.ParseOptions(&o.Options, &o.Options.Options, os.Args, "webconsole.conf") @@ -28,6 +34,10 @@ func StartService() { log.Fatalf("invalid --api-server %s", o.Options.ApiServer) } + for _, binPath := range []string{o.Options.KubectlPath, o.Options.IpmitoolPath} { + ensureBinExists(binPath) + } + cloudcommon.InitAuth(&o.Options.Options, func() { log.Infof("Auth complete") })