diff --git a/pkg/cloudcommon/sshkeys/sshkeys.go b/pkg/cloudcommon/sshkeys/sshkeys.go new file mode 100644 index 0000000000..8bd4c4b697 --- /dev/null +++ b/pkg/cloudcommon/sshkeys/sshkeys.go @@ -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 +} diff --git a/pkg/cloudcommon/workmanager/manager.go b/pkg/cloudcommon/workmanager/manager.go index f41449ba38..f0ffd344ed 100644 --- a/pkg/cloudcommon/workmanager/manager.go +++ b/pkg/cloudcommon/workmanager/manager.go @@ -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) } } diff --git a/pkg/hostman/guestfs/core.go b/pkg/hostman/guestfs/core.go new file mode 100644 index 0000000000..1a037d2c13 --- /dev/null +++ b/pkg/hostman/guestfs/core.go @@ -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 +} diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 2b8d02c7e6..e80fe56efd 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -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 diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index 6715c6f3e9..9a6b54026e 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -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 ??") + } +} diff --git a/pkg/hostman/hostinfo/hostinfo.go b/pkg/hostman/hostinfo/hostinfo.go index 869e4ef901..0da82ee4ac 100644 --- a/pkg/hostman/hostinfo/hostinfo.go +++ b/pkg/hostman/hostinfo/hostinfo.go @@ -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 +} diff --git a/pkg/hostman/storageman/core.go b/pkg/hostman/storageman/core.go new file mode 100644 index 0000000000..b1dd36bb55 --- /dev/null +++ b/pkg/hostman/storageman/core.go @@ -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 +} diff --git a/pkg/hostman/storageman/localstorage.go b/pkg/hostman/storageman/localstorage.go new file mode 100644 index 0000000000..0e41d91a28 --- /dev/null +++ b/pkg/hostman/storageman/localstorage.go @@ -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 +}