Merge branch 'release/2.2.0' of ssh://git.yunion.io/~quxuan/onecloud into feature/qx-azure

This commit is contained in:
屈轩
2018-09-13 16:22:53 +08:00
72 changed files with 1043 additions and 356 deletions
+3 -3
View File
@@ -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: <https://management.chinacloudapi.cn、https://management.azure.com、https://management.usgovcloudapi.net、https://management.microsoftazure.de>"`
AccessURL string `helo:"hello" metavar:"Azure choices: <AzureGermanCloud、AzureChinaCloud、AzureUSGovernmentCloud、AzurePublicCloud>"`
Desc string `help:"Description"`
Enabled bool `help:"Enabled the provider automatically"`
}
+26
View File
@@ -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 {
+4 -1
View File
@@ -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
+12
View File
@@ -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
})
}
+64 -31
View File
@@ -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
}
}
+9 -3
View File
@@ -61,10 +61,16 @@ func TestRadixNode(t *testing.T) {
func TestParams(t *testing.T) {
r := NewRadix()
r.Add([]string{"POST", "clouds", "<action>"}, "classAction")
r.Add([]string{"POST", "clouds", "<resid>", "<action>"}, "objectAction")
r.Add([]string{"POST", "clouds", "<cls_action>"}, "classAction")
r.Add([]string{"POST", "clouds", "<resid>", "sync"}, "objectSyncAction")
r.Add([]string{"POST", "clouds", "<resid>", "<obj_action>"}, "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)
+2 -2
View File
@@ -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)
}
+2
View File
@@ -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)
}
+3 -3
View File
@@ -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)
+1 -1
View File
@@ -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
@@ -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()
}
+46
View File
@@ -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()
}
}
+19
View File
@@ -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())
}
+79 -48
View File
@@ -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)
}
}
}
+3 -3
View File
@@ -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) {
+1 -1
View File
@@ -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
}
+2 -2
View File
@@ -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
}
+3 -3
View File
@@ -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
}
+2 -2
View File
@@ -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
+1 -1
View File
@@ -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
}
+4 -2
View File
@@ -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
+23 -4
View File
@@ -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 {
@@ -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{
{
+7 -3
View File
@@ -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
})
+11 -1
View File
@@ -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
+7 -1
View File
@@ -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
-1
View File
@@ -1 +0,0 @@
package hostdrivers
+1 -1
View File
@@ -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 {
+10 -17
View File
@@ -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 ""
}
+24 -1
View File
@@ -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)
+2 -1
View File
@@ -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 {
+17 -1
View File
@@ -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)
+1 -1
View File
@@ -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
+2 -1
View File
@@ -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 {
+1
View File
@@ -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)
}
+2 -2
View File
@@ -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 {
+1 -1
View File
@@ -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)
+2
View File
@@ -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
+4
View File
@@ -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
}
@@ -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)
@@ -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())
@@ -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)
+16 -12
View File
@@ -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
}
}
+14 -1
View File
@@ -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)
}
+1 -3
View File
@@ -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
}
+2
View File
@@ -13,6 +13,8 @@ var (
func init() {
WebConsole = WebConsoleManager{NewWebConsoleManager()}
register(&WebConsole)
}
type WebConsoleManager struct {
+32 -19
View File
@@ -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
+9 -5
View File
@@ -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 {
+3
View File
@@ -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
}
+2 -1
View File
@@ -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)
}
+16 -2
View File
@@ -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 {
+13 -18
View File
@@ -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 {
+11 -13
View File
@@ -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 {
+14 -9
View File
@@ -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
+5 -6
View File
@@ -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 {
+231 -68
View File
@@ -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: &region.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 {
+1 -1
View File
@@ -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 {
+3 -4
View File
@@ -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 {
+4 -3
View File
@@ -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 {
+3 -4
View File
@@ -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
}
+1 -1
View File
@@ -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 {
+13 -8
View File
@@ -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)
})
+6 -1
View File
@@ -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
})
+39 -2
View File
@@ -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, "")
})
}
+4 -6
View File
@@ -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
+6 -6
View File
@@ -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
}
+4 -5
View File
@@ -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 {
+4 -4
View File
@@ -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
}
+3 -1
View File
@@ -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)
+15 -3
View File
@@ -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)
}
+3 -1
View File
@@ -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"`
}
+10
View File
@@ -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")
})