mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
temp
This commit is contained in:
@@ -0,0 +1,19 @@
|
||||
package sshkeys
|
||||
|
||||
import "yunion.io/x/jsonutils"
|
||||
|
||||
type SSHKeys struct {
|
||||
PublicKey string
|
||||
DeletePublicKey string
|
||||
AdminPublicKey string
|
||||
ProjectPublicKey string
|
||||
}
|
||||
|
||||
func GetKeys(data jsonutils.JSONObject) *SSHKeys {
|
||||
var ret = new(SSHKeys)
|
||||
ret.PublicKey, _ = data.GetString("public_key")
|
||||
ret.DeletePublicKey, _ = data.GetString("delete_public_key")
|
||||
ret.AdminPublicKey, _ = data.GetString("admin_public_key")
|
||||
ret.ProjectPublicKey, _ = data.GetString("project_public_key")
|
||||
return ret
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package workmanager
|
||||
import (
|
||||
"context"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
@@ -22,11 +23,11 @@ func (w *SWorkManager) done() {
|
||||
atomic.AddInt32(&w.curCount, -1)
|
||||
}
|
||||
|
||||
func (w *SWorkManager) DelayTask(task DelayTaskFunc, ctx context.Context, params jsonutils.JSONObject) {
|
||||
func (w *SWorkManager) DelayTask(task func()) {
|
||||
w.add()
|
||||
go func() {
|
||||
defer w.done()
|
||||
task(ctx, params)
|
||||
task()
|
||||
}()
|
||||
}
|
||||
|
||||
@@ -34,6 +35,7 @@ func (w *SWorkManager) Stop() {
|
||||
log.Infof("WorkManager To stop, wait for workers ...")
|
||||
for w.curCount > 0 {
|
||||
log.Warningf("Busy workers count %d, waiting stopped", w.curCount)
|
||||
time.Sleep(1)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
package guestfs
|
||||
|
||||
import (
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/sshkeys"
|
||||
)
|
||||
|
||||
type SDeployInfo struct {
|
||||
publicKey *sshkeys.SSHKeys
|
||||
deploys jsonutils.JSONObject
|
||||
password string
|
||||
isInit bool
|
||||
}
|
||||
@@ -14,10 +14,13 @@ import (
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/httpclients"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/sshkeys"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/workmanager"
|
||||
"yunion.io/x/onecloud/pkg/hostman"
|
||||
"yunion.io/x/onecloud/pkg/hostman/guestfs"
|
||||
"yunion.io/x/onecloud/pkg/httperrors"
|
||||
"yunion.io/x/pkg/util/regutils"
|
||||
"yunion.io/x/pkg/util/seclib"
|
||||
)
|
||||
|
||||
type SGuestManager struct {
|
||||
@@ -195,7 +198,36 @@ func (m *SGuestManager) Monitor(sid, cmd string, callback func(string)) error {
|
||||
}
|
||||
|
||||
func (m *SGuestManager) DoDeploy(ctx context.Context, sid string, body jsonutils.JSONObject, isInit bool) {
|
||||
// TODO
|
||||
guest, ok := m.Servers[sid]
|
||||
if ok {
|
||||
desc, _ := body.Get("desc")
|
||||
if desc != nil {
|
||||
guest.SaveDesc(desc)
|
||||
}
|
||||
if jsonutils.QueryBoolean(body, "k8s_pod", false) {
|
||||
TaskComplete(ctx, nil)
|
||||
return
|
||||
}
|
||||
// TODO
|
||||
publicKey := sshkeys.GetKeys(body)
|
||||
deploys, _ := body.Get("deploys")
|
||||
password, _ := body.GetString("password")
|
||||
resetPassword := jsonutils.QueryBoolean(body, "reset_password", false)
|
||||
if resetPassword && len(password) == 0 {
|
||||
password = seclib.RandomPassword(12)
|
||||
}
|
||||
|
||||
guestInfo, err := guest.DeployFs(&guestfs.SDeployInfo{
|
||||
publicKey, deploys, password, isInit})
|
||||
if err != nil {
|
||||
log.Errorf("Deploy guest fs error: %s", err)
|
||||
TaskFailed(ctx, err.Error())
|
||||
} else {
|
||||
TaskComplete(ctx, guestInfo)
|
||||
}
|
||||
} else {
|
||||
TaskFailed(ctx, fmt.Sprinft("Guest %s not found", sid))
|
||||
}
|
||||
}
|
||||
|
||||
// delay cpuset balance
|
||||
|
||||
@@ -19,6 +19,8 @@ import (
|
||||
"yunion.io/x/onecloud/pkg/appctx"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/httpclients"
|
||||
"yunion.io/x/onecloud/pkg/hostman"
|
||||
"yunion.io/x/onecloud/pkg/hostman/guestfs"
|
||||
"yunion.io/x/onecloud/pkg/hostman/hostinfo"
|
||||
"yunion.io/x/onecloud/pkg/hostman/monitor"
|
||||
"yunion.io/x/onecloud/pkg/hostman/options"
|
||||
)
|
||||
@@ -194,7 +196,7 @@ func (s *SKVMGuestInstance) onAsyncScriptStart(ctx context.Context, isStarted bo
|
||||
} else {
|
||||
log.Infof("Async start server %s failed: %s!!!", s.GetName(), err)
|
||||
if ctx != nil {
|
||||
s.TaskFailed(ctx, fmt.Sprintf("Async start server failed: %s", err))
|
||||
TaskFailed(ctx, fmt.Sprintf("Async start server failed: %s", err))
|
||||
}
|
||||
s.SyncStatus()
|
||||
}
|
||||
@@ -284,7 +286,7 @@ func (s *SKVMGuestInstance) onGetQemuVersion(ctx context.Context, version string
|
||||
migratePort, _ := s.Desc.Get("live_migrate_dest_port")
|
||||
body := jsonutils.NewDict(
|
||||
jsonutils.JSONPair{"live_migrate_dest_port", migratePort})
|
||||
s.TaskComplete(ctx, body)
|
||||
TaskComplete(ctx, body)
|
||||
} else if jsonutils.QueryBoolean(s.Desc, "is_slave", false) {
|
||||
// TODO
|
||||
} else if jsonutils.QueryBoolean(s.Desc, "is_master", false) && ctx == nil {
|
||||
@@ -341,7 +343,7 @@ func (s *SKVMGuestInstance) SyncStatus() {
|
||||
httpclients.GetDefaultComputeClient().UpdateServerStatus(s.GetId(), status)
|
||||
}
|
||||
|
||||
func (s *SKVMGuestInstance) TaskFailed(ctx context.Context, reason string) error {
|
||||
func TaskFailed(ctx context.Context, reason string) error {
|
||||
if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil {
|
||||
httpclients.GetDefaultComputeClient().TaskFail(ctx, taskId.(string), reason)
|
||||
return nil
|
||||
@@ -351,7 +353,7 @@ func (s *SKVMGuestInstance) TaskFailed(ctx context.Context, reason string) error
|
||||
}
|
||||
}
|
||||
|
||||
func (s *SKVMGuestInstance) TaskComplete(ctx context.Context, data jsonutils.JSONObject) error {
|
||||
func TaskComplete(ctx context.Context, data jsonutils.JSONObject) error {
|
||||
if taskId := ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID); taskId != nil {
|
||||
httpclients.GetDefaultComputeClient().TaskComplete(ctx, taskId.(string), data, 0)
|
||||
return nil
|
||||
@@ -379,3 +381,16 @@ func (s *SKVMGuestInstance) StartGuest(ctx context.Context, params jsonutils.JSO
|
||||
s.asyncScriptStart(ctx, params)
|
||||
})
|
||||
}
|
||||
|
||||
func (s *SKVMGuestInstance) DeployFs(deployInfo *guestfs.SDeployInfo) (jsonutils.JSONObject, error) {
|
||||
disks, _ := s.Desc.GetArray("disks")
|
||||
if len(disks) > 0 {
|
||||
storageId, _ := disks[0].GetString("storage_id")
|
||||
diskId, _ := disks[0].GetString("disk_id")
|
||||
|
||||
disk := hostinfo.GetStorageManager().GetStorageDisk(storageId, diskId)
|
||||
return disk.DeployGuestFs(s.Desc, deployInfo)
|
||||
} else {
|
||||
return nil, fmt.Errorf("Guest dosen't have disk ??")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11,12 +11,15 @@ import (
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/qemutils"
|
||||
"yunion.io/x/onecloud/pkg/hostman"
|
||||
"yunion.io/x/onecloud/pkg/hostman/options"
|
||||
"yunion.io/x/onecloud/pkg/hostman/storageman"
|
||||
)
|
||||
|
||||
type SHostInfo struct {
|
||||
isRegistered bool
|
||||
|
||||
Cpu *SCPUInfo
|
||||
|
||||
storageManager *storageman.SStorageManager
|
||||
}
|
||||
|
||||
func (h *SHostInfo) Start() error {
|
||||
@@ -290,3 +293,7 @@ func Init() error {
|
||||
func Instance() *SHostInfo {
|
||||
return hostInfo
|
||||
}
|
||||
|
||||
func GetStorageManager() *storageman.SStorageManager {
|
||||
return hostInfo.storageManager
|
||||
}
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
package storageman
|
||||
|
||||
import (
|
||||
"path"
|
||||
"sync"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/onecloud/pkg/hostman/guestfs"
|
||||
)
|
||||
|
||||
type IStorage interface {
|
||||
StorageType() string
|
||||
GetPath() string
|
||||
|
||||
// Find owner disks first, if not found, call create disk
|
||||
GetDiskById(diskId string) IDisk
|
||||
CreateDisk(diskId string) IDisk
|
||||
}
|
||||
|
||||
type IDisk interface {
|
||||
GetId() string
|
||||
Probe() bool
|
||||
DeployGuestFs(guestDesc *jsonutils.JSONDict,
|
||||
deployInfo *guestfs.SDeployInfo) (jsonutils.JSONObject, error)
|
||||
}
|
||||
|
||||
type SBaseStorage struct {
|
||||
Manager *SStorageManager
|
||||
StorageId string
|
||||
Path string
|
||||
StorageName string
|
||||
StorageConf *jsonutils.JSONDict
|
||||
StoragecacheId string
|
||||
|
||||
Disks []IDisk
|
||||
DiskLock *sync.Mutex
|
||||
}
|
||||
|
||||
func (s *SBaseStorage) GetPath() string {
|
||||
return s.Path
|
||||
}
|
||||
|
||||
func NewBaseStorage(manager *SStorageManager, path string) *SBaseStorage {
|
||||
var ret = new(SBaseStorage)
|
||||
ret.Disks = make([]IDisk, 0)
|
||||
ret.DiskLock = new(sync.Mutex)
|
||||
ret.Manager = manager
|
||||
ret.Path = path
|
||||
return ret
|
||||
}
|
||||
|
||||
type SBaseDisk struct {
|
||||
Id string
|
||||
Storage IStorage
|
||||
}
|
||||
|
||||
func (d *SBaseDisk) getPath() string {
|
||||
return path.Join(d.Storage.GetPath(), d.Id)
|
||||
}
|
||||
|
||||
func (d *SBaseDisk) DeployGuestFs(
|
||||
guestDesc *jsonutils.JSONDict,
|
||||
deployInfo *guestfs.SDeployInfo) (jsonutils.JSONObject, error) {
|
||||
// TODO
|
||||
var kvmDisk = NewKVMGuestDisk(d.getPath())
|
||||
defer kvmDisk.Disconnect()
|
||||
kvmDisk.Connect()
|
||||
}
|
||||
|
||||
func NewBaseDisk(storage IStorage, id string) *SBaseDisk {
|
||||
var ret = new(SBaseDisk)
|
||||
ret.Storage = storage
|
||||
ret.Id = id
|
||||
return ret
|
||||
}
|
||||
|
||||
type SStorageManager struct {
|
||||
storages map[string]IStorage
|
||||
}
|
||||
|
||||
func NewStorageManager() *SStorageManager {
|
||||
var ret = new(SStorageManager)
|
||||
// TODO
|
||||
ret.storages = make(map[string]IStorage, 0)
|
||||
|
||||
return ret
|
||||
}
|
||||
|
||||
func (m *SStorageManager) GetStorageDisk(storageId, diskId string) IDisk {
|
||||
if storage, ok := m.storages[storageId]; ok {
|
||||
return storage.GetDiskById(diskId)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
package storageman
|
||||
|
||||
import (
|
||||
"os"
|
||||
)
|
||||
|
||||
type SLocalStorage struct {
|
||||
*SBaseStorage
|
||||
}
|
||||
|
||||
func NewLocalStorage(manager *SStorageManager, path string) *SLocalStorage {
|
||||
var ret = new(SLocalStorage)
|
||||
ret.SBaseStorage = NewBaseStorage(manager, path)
|
||||
ret.StartSnapshotRecycle()
|
||||
return ret
|
||||
}
|
||||
|
||||
func (s *SLocalStorage) StorageType() string {
|
||||
return "local"
|
||||
}
|
||||
|
||||
func (s *SLocalStorage) GetDiskById(diskId string) IDisk {
|
||||
for i := 0; i < len(s.Disks); i++ {
|
||||
if s.Disks[i].GetId() == diskId {
|
||||
return s.Disks[i]
|
||||
}
|
||||
}
|
||||
return s.CreateDisk(diskId)
|
||||
}
|
||||
|
||||
func (s *SLocalStorage) CreateDisk(diskId string) IDisk {
|
||||
s.DiskLock.Lock()
|
||||
defer s.DiskLock.Unlock()
|
||||
var disk = NewLocalDisk(s, diskId)
|
||||
if disk.Probe() {
|
||||
s.Disks = append(s.Disks, disk)
|
||||
return disk
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *SLocalStorage) StartSnapshotRecycle() {
|
||||
//TODO
|
||||
}
|
||||
|
||||
type SLocalDisk struct {
|
||||
*SBaseDisk
|
||||
}
|
||||
|
||||
func NewLocalDisk(storage IStorage, id string) *SLocalDisk {
|
||||
var ret = new(SLocalDisk)
|
||||
ret.SBaseDisk = NewBaseDisk(storage, id)
|
||||
return ret
|
||||
}
|
||||
|
||||
func (d *SLocalDisk) GetId() string {
|
||||
return d.Id
|
||||
}
|
||||
|
||||
func (d *SLocalDisk) Probe() bool {
|
||||
if _, err := os.Stat(d.getPath()); !os.IsNotExist(err) {
|
||||
return true
|
||||
}
|
||||
// TODO alter ??
|
||||
return false
|
||||
}
|
||||
Reference in New Issue
Block a user