mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge branch 'master' into hotfix/qx-storagecache-docs
This commit is contained in:
@@ -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),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
|
||||
@@ -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),
|
||||
|
||||
Reference in New Issue
Block a user