diff --git a/pkg/cloudcommon/workmanager/manager.go b/pkg/cloudcommon/workmanager/manager.go index d4d10e396d..7f90d65c7c 100644 --- a/pkg/cloudcommon/workmanager/manager.go +++ b/pkg/cloudcommon/workmanager/manager.go @@ -24,6 +24,7 @@ import ( "yunion.io/x/log" "yunion.io/x/onecloud/pkg/appctx" + "yunion.io/x/onecloud/pkg/appsrv" ) type DelayTaskFunc func(context.Context, interface{}) (jsonutils.JSONObject, error) @@ -35,6 +36,7 @@ type SWorkManager struct { onFailed OnTaskFailed onCompleted OnTaskCompleted + worker *appsrv.SWorkerManager } func (w *SWorkManager) add() { @@ -45,16 +47,26 @@ func (w *SWorkManager) done() { atomic.AddInt32(&w.curCount, -1) } +func (w *SWorkManager) DelayTask(ctx context.Context, task DelayTaskFunc, params interface{}) { + w.delayTask(ctx, task, params, w.worker) +} + +func (w *SWorkManager) DelayTaskWithWorker( + ctx context.Context, task DelayTaskFunc, params interface{}, worker *appsrv.SWorkerManager, +) { + w.delayTask(ctx, task, params, worker) +} + // If delay task is not panic and task func return err is nil // task complete will be called, otherwise called task failed // Params is interface for receive any type, task func should do type assertion -func (w *SWorkManager) DelayTask(ctx context.Context, task DelayTaskFunc, params interface{}) { +func (w *SWorkManager) delayTask(ctx context.Context, task DelayTaskFunc, params interface{}, worker *appsrv.SWorkerManager) { if ctx == nil || ctx.Value(appctx.APP_CONTEXT_KEY_TASK_ID) == nil { - w.DelayTaskWithoutReqctx(ctx, task, params) + w.delayTaskWithoutReqctx(ctx, task, params, worker) return } else { w.add() - go func() { + worker.Run(func() { defer w.done() defer func() { if r := recover(); r != nil { @@ -85,14 +97,20 @@ func (w *SWorkManager) DelayTask(ctx context.Context, task DelayTaskFunc, params log.Infof("DelayTask complete: %v", res) w.onCompleted(ctx, res) } - }() + }, nil, nil) } } -// response task by self, did not callback func (w *SWorkManager) DelayTaskWithoutReqctx(ctx context.Context, task DelayTaskFunc, params interface{}) { + w.delayTaskWithoutReqctx(ctx, task, params, w.worker) +} + +// response task by self, did not callback +func (w *SWorkManager) delayTaskWithoutReqctx( + ctx context.Context, task DelayTaskFunc, params interface{}, worker *appsrv.SWorkerManager, +) { w.add() - go func() { + w.worker.Run(func() { defer w.done() defer func() { if r := recover(); r != nil { @@ -109,7 +127,7 @@ func (w *SWorkManager) DelayTaskWithoutReqctx(ctx context.Context, task DelayTas log.Errorln("DelayTaskWithoutReqctx error: ", err) w.onFailed(ctx, err.Error()) } - }() + }, nil, nil) } func (w *SWorkManager) Stop() { @@ -120,9 +138,14 @@ func (w *SWorkManager) Stop() { } } -func NewWorkManger(onFailed OnTaskFailed, onCompleted OnTaskCompleted) *SWorkManager { +func NewWorkManger(onFailed OnTaskFailed, onCompleted OnTaskCompleted, workerCount int) *SWorkManager { + if workerCount <= 0 { + workerCount = 1 + } return &SWorkManager{ onFailed: onFailed, onCompleted: onCompleted, + worker: appsrv.NewWorkerManager( + "RequestWorker", workerCount, appsrv.DEFAULT_BACKLOG, false), } } diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index 9feba8b57a..90ed77b4ed 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -174,11 +174,14 @@ func deleteGuest(ctx context.Context, w http.ResponseWriter, r *http.Request) { } func guestCreate(ctx context.Context, sid string, body jsonutils.JSONObject) (interface{}, error) { - err := guestman.GetGuestManager().PrepareCreate(sid) - if err != nil { - return nil, err + if guestman.GetGuestManager().IsGuestExist(sid) { + return nil, httperrors.NewBadRequestError("Guest %s is exist", sid) } - hostutils.DelayTask(ctx, guestman.GetGuestManager().GuestDeploy, &guestman.SGuestDeploy{sid, body, true}) + hostutils.DelayTaskWithWorker(ctx, + guestman.GetGuestManager().GuestCreate, + &guestman.SGuestDeploy{sid, body, true}, + guestman.NbdWorker, + ) return nil, nil } @@ -187,7 +190,11 @@ func guestDeploy(ctx context.Context, sid string, body jsonutils.JSONObject) (in if err != nil { return nil, err } - hostutils.DelayTask(ctx, guestman.GetGuestManager().GuestDeploy, &guestman.SGuestDeploy{sid, body, false}) + hostutils.DelayTaskWithWorker(ctx, + guestman.GetGuestManager().GuestDeploy, + &guestman.SGuestDeploy{sid, body, false}, + guestman.NbdWorker, + ) return nil, nil } diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 77b96fd09f..50db11bed0 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -29,6 +29,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/util/regutils" + "yunion.io/x/pkg/util/seclib" compute "yunion.io/x/onecloud/pkg/apis/compute" appsrv "yunion.io/x/onecloud/pkg/appsrv" @@ -41,12 +42,12 @@ import ( cgrouputils "yunion.io/x/onecloud/pkg/util/cgrouputils" fileutils2 "yunion.io/x/onecloud/pkg/util/fileutils2" netutils2 "yunion.io/x/onecloud/pkg/util/netutils2" - seclib2 "yunion.io/x/onecloud/pkg/util/seclib2" timeutils2 "yunion.io/x/onecloud/pkg/util/timeutils2" ) var ( LAST_USED_PORT = 0 + NbdWorker = appsrv.NewWorkerManager("nbd_worker", 1, appsrv.DEFAULT_BACKLOG, false) ) const ( @@ -334,6 +335,61 @@ func (m *SGuestManager) Monitor(sid, cmd string, callback func(string)) error { } } +func (m *SGuestManager) GuestCreate(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { + deployParams, ok := params.(*SGuestDeploy) + if !ok { + return nil, hostutils.ParamsError + } + + var guest *SKVMGuestInstance + err := func() error { + m.ServersLock.Lock() + defer m.ServersLock.Unlock() + if _, ok := m.GetServer(deployParams.Sid); ok { + return httperrors.NewBadRequestError("Guest %s exists", deployParams.Sid) + } + guest = NewKVMGuestInstance(deployParams.Sid, m) + desc, _ := deployParams.Body.Get("desc") + if desc != nil { + err := guest.SaveDesc(desc) + if err != nil { + return errors.Wrap(err, "save desc") + } + } + m.SaveServer(deployParams.Sid, guest) + return guest.PrepareDir() + }() + if err != nil { + return nil, errors.Wrap(err, "prepare guest") + } + return m.startDeploy(ctx, deployParams, guest) +} + +func (m *SGuestManager) startDeploy( + ctx context.Context, deployParams *SGuestDeploy, guest *SKVMGuestInstance) (jsonutils.JSONObject, error) { + + if jsonutils.QueryBoolean(deployParams.Body, "k8s_pod", false) { + return nil, nil + } + publicKey := deployapi.GetKeys(deployParams.Body) + deploys, _ := deployParams.Body.GetArray("deploys") + password, _ := deployParams.Body.GetString("password") + resetPassword := jsonutils.QueryBoolean(deployParams.Body, "reset_password", false) + if resetPassword && len(password) == 0 { + password = seclib.RandomPassword(12) + } + enableCloudInit := jsonutils.QueryBoolean(deployParams.Body, "enable_cloud_init", false) + + guestInfo, err := guest.DeployFs(deployapi.NewDeployInfo( + publicKey, deployapi.JsonDeploysToStructs(deploys), password, deployParams.IsInit, false, + options.HostOptions.LinuxDefaultRootUser, options.HostOptions.WindowsDefaultAdminUser, enableCloudInit)) + if err != nil { + return nil, errors.Wrap(err, "Deploy guest fs") + } else { + return guestInfo, nil + } +} + // Delay process func (m *SGuestManager) GuestDeploy(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { deployParams, ok := params.(*SGuestDeploy) @@ -347,26 +403,7 @@ func (m *SGuestManager) GuestDeploy(ctx context.Context, params interface{}) (js if desc != nil { guest.SaveDesc(desc) } - if jsonutils.QueryBoolean(deployParams.Body, "k8s_pod", false) { - return nil, nil - } - publicKey := deployapi.GetKeys(deployParams.Body) - deploys, _ := deployParams.Body.GetArray("deploys") - password, _ := deployParams.Body.GetString("password") - resetPassword := jsonutils.QueryBoolean(deployParams.Body, "reset_password", false) - if resetPassword && len(password) == 0 { - password = seclib2.RandomPassword2(12) - } - enableCloudInit := jsonutils.QueryBoolean(deployParams.Body, "enable_cloud_init", false) - - guestInfo, err := guest.DeployFs(deployapi.NewDeployInfo( - publicKey, deployapi.JsonDeploysToStructs(deploys), password, deployParams.IsInit, false, - options.HostOptions.LinuxDefaultRootUser, options.HostOptions.WindowsDefaultAdminUser, enableCloudInit)) - if err != nil { - return nil, errors.Wrap(err, "Deploy guest fs") - } else { - return guestInfo, nil - } + return m.startDeploy(ctx, deployParams, guest) } else { return nil, fmt.Errorf("Guest %s not found", deployParams.Sid) } diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index c006fd3f37..88e6ba8b16 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -66,6 +66,8 @@ func (host *SHostService) OnExitService() { func (host *SHostService) RunService() { app := app_common.InitApp(&options.HostOptions.BaseOptions, false) + + hostutils.Init() hostInstance := hostinfo.Instance() if err := hostInstance.Init(); err != nil { log.Fatalf(err.Error()) diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index 1b2c571173..1c202f2534 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -179,7 +179,14 @@ func DelayTaskWithoutReqctx(ctx context.Context, task workmanager.DelayTaskFunc, wm.DelayTaskWithoutReqctx(ctx, task, params) } -func init() { - wm = workmanager.NewWorkManger(TaskFailed, TaskComplete) - k8sWm = workmanager.NewWorkManger(K8sTaskFailed, K8sTaskComplete) +func DelayTaskWithWorker( + ctx context.Context, task workmanager.DelayTaskFunc, + params interface{}, worker *appsrv.SWorkerManager, +) { + wm.DelayTaskWithWorker(ctx, task, params, worker) +} + +func Init() { + wm = workmanager.NewWorkManger(TaskFailed, TaskComplete, options.HostOptions.DefaultRequestWorkerCount) + k8sWm = workmanager.NewWorkManger(K8sTaskFailed, K8sTaskComplete, options.HostOptions.DefaultRequestWorkerCount) } diff --git a/pkg/hostman/options/options.go b/pkg/hostman/options/options.go index 684c128d63..55f36bf667 100644 --- a/pkg/hostman/options/options.go +++ b/pkg/hostman/options/options.go @@ -105,7 +105,8 @@ type SHostOptions struct { MaxReservedMemory int `default:"10240" help:"host reserved memory"` - DeployServerSocketPath string `help:"Deploy server listen socket path" default:"/var/run/deploy.sock"` + DeployServerSocketPath string `help:"Deploy server listen socket path" default:"/var/run/deploy.sock"` + DefaultRequestWorkerCount int `default:"8" help:"default request worker count"` } var HostOptions SHostOptions diff --git a/pkg/hostman/storageman/imagecachemanager_rbd.go b/pkg/hostman/storageman/imagecachemanager_rbd.go index 521957724f..17aaf23fcc 100644 --- a/pkg/hostman/storageman/imagecachemanager_rbd.go +++ b/pkg/hostman/storageman/imagecachemanager_rbd.go @@ -112,7 +112,7 @@ func (c *SRbdImageCacheManager) PrefetchImageCache(ctx context.Context, data int if err != nil { return nil, err } - format, _ := body.GetString("format") + format := "qcow2" srcUrl, _ := body.GetString("src_url") zone, _ := body.GetString("zone") diff --git a/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go b/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go index 40e7c2d8ad..d33cfd5716 100644 --- a/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go +++ b/pkg/scheduler/algorithm/predicates/disk_schedtag_predicate.go @@ -105,7 +105,7 @@ func (p *DiskSchedtagPredicate) IsResourceFitInput(u *core.Unit, c core.Candidat if c.Getter().ResourceType() == computeapi.HostResourceTypePrepaidRecycle { return nil } - if !(len(d.Backend) == 0 || d.Backend == computeapi.STORAGE_LOCAL) { + if len(d.Backend) != 0 { if storage.StorageType != d.Backend { return &FailReason{ fmt.Sprintf("Storage %s backend %s != %s", storage.Name, storage.StorageType, d.Backend),