From fab24f13b5a388d1d6bcd12e062933a363d481f7 Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Mon, 11 Sep 2023 08:57:55 +0800 Subject: [PATCH 1/6] =?UTF-8?q?=E6=9B=B4=E6=96=B0=E5=87=BD=E6=95=B0?= =?UTF-8?q?=E5=90=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- advisor/internal/task/schedule_scheme.go | 181 +++++++++++++++---- common/pkgs/mq/advisor/apis.go | 9 +- common/pkgs/mq/advisor/task/task.go | 78 ++------ common/pkgs/mq/executor/pcm.go | 50 ++--- common/pkgs/mq/executor/storage.go | 6 +- common/pkgs/mq/manager/api.go | 6 +- executor/internal/services/pcm.go | 14 +- executor/internal/services/storage.go | 2 +- executor/internal/task/cache_move_package.go | 8 +- executor/internal/task/pcm_schedule_task.go | 30 +-- executor/internal/task/pcm_upload_img.go | 8 +- go.mod | 1 + go.sum | 2 + 13 files changed, 226 insertions(+), 169 deletions(-) diff --git a/advisor/internal/task/schedule_scheme.go b/advisor/internal/task/schedule_scheme.go index ce5a105..eb717c0 100644 --- a/advisor/internal/task/schedule_scheme.go +++ b/advisor/internal/task/schedule_scheme.go @@ -1,14 +1,22 @@ package task import ( + "fmt" "time" + "gitlink.org.cn/cloudream/common/models" "gitlink.org.cn/cloudream/common/pkgs/logger" "gitlink.org.cn/cloudream/common/pkgs/task" - exectsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/executor/task" + "gitlink.org.cn/cloudream/common/utils/convertto" + "gitlink.org.cn/cloudream/scheduler/common/globals" + "gitlink.org.cn/cloudream/scheduler/common/models/job" + advtsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/advisor/task" + "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/collector" ) type GetScheduleScheme struct { + Job job.NormalJob + preAdjustNodeID int64 } func NewGetScheduleScheme() *GetScheduleScheme { @@ -23,7 +31,9 @@ func (t *GetScheduleScheme) Execute(task *task.Task[TaskContext], ctx TaskContex err := t.do(task.ID(), ctx) if err != nil { //TODO 若任务失败,上报的状态failed字段根据情况修改 - ctx.reporter.Report(task.ID(), exectsk.NewScheduleTaskStatus("failed", err.Error(), 0)) + ctx.reporter.Report(task.ID(), advtsk.NewTaskStatus("failed", err.Error(), true, advtsk.AdjustedScheme{})) + } else { + ctx.reporter.Report(task.ID(), advtsk.NewTaskStatus("failed", err.Error(), false, advtsk.AdjustedScheme{})) } ctx.reporter.ReportNow() @@ -33,41 +43,138 @@ func (t *GetScheduleScheme) Execute(task *task.Task[TaskContext], ctx TaskContex } func (t *GetScheduleScheme) do(taskID string, ctx TaskContext) error { - // pcmCli, err := globals.PCMPool.Acquire() - // if err != nil { - // return fmt.Errorf("new pcm client: %w", err) - // } - // defer pcmCli.Close() + isAvailable, err := t.CheckResourceAvailability() + if err != nil { + return err + } - // resp, err := pcmCli.ScheduleTask(pcm.ScheduleTaskReq{ - - // }) - - // if err != nil { - // return err - // } - - // var prevStatus string - // for { - // tsResp, err := pcmCli.GetTaskStatus(pcm.GetTaskStatusReq{ - // NodeID: t.nodeID, - // PCMJobID: resp.PCMJobID, - // }) - // if err != nil { - // return err - // } - - // if tsResp.Status != prevStatus { - // ctx.reporter.Report(taskID, exectsk.NewScheduleTaskStatus(tsResp.Status, "", resp.PCMJobID)) - // } - - // prevStatus = tsResp.Status - - // // TODO 根据接口result返回情况修改 - // // 根据返回的result判定任务是否完成,若完成 跳出循环,结束任务 - // if tsResp.Status == "Completed" { - // return nil - // } - // } + if isAvailable { + // 确认code、dataset、image是否已经调度到该中心 + } else { + // 重新执行预调度方案,寻找最优节点 + } return nil } + +// 检查预调度节点资源是否足够 +func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { + colCli, err := globals.CollectorMQPool.Acquire() + if err != nil { + return false, fmt.Errorf("new collector client: %w", err) + } + defer colCli.Close() + + neededCPU := t.Job.Info.Resources.CPU + if neededCPU > 0 { + resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ + NodeId: t.preAdjustNodeID, + ResourceType: models.ResourceTypeCPU, + }) + if err != nil { + return false, err + } + + availCPU := resp.Data.(models.CPUResourceData).Available.Value + + if float64(availCPU) < 1.5*neededCPU { + fmt.Printf("Schedule Scheme is wrong: Insufficient cpu") + return false, nil + } + } + + neededNPU := t.Job.Info.Resources.NPU + if neededNPU > 0 { + resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ + NodeId: t.preAdjustNodeID, + ResourceType: models.ResourceTypeNPU, + }) + if err != nil { + return false, err + } + + availNPU := resp.Data.(models.NPUResourceData).Available.Value + + if float64(availNPU) < 1.5*neededNPU { + fmt.Printf("Schedule Scheme is wrong: Insufficient npu") + return false, nil + } + } + + neededGPU := t.Job.Info.Resources.GPU + if neededGPU > 0 { + resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ + NodeId: t.preAdjustNodeID, + ResourceType: models.ResourceTypeGPU, + }) + if err != nil { + return false, err + } + + availGPU := resp.Data.(models.GPUResourceData).Available.Value + + if float64(availGPU) < 1.5*neededGPU { + fmt.Printf("Schedule Scheme is wrong: Insufficient gpu") + return false, nil + } + } + + neededMLU := t.Job.Info.Resources.MLU + if neededMLU > 0 { + resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ + NodeId: t.preAdjustNodeID, + ResourceType: models.ResourceTypeMLU, + }) + if err != nil { + return false, err + } + + availMLU := resp.Data.(models.MLUResourceData).Available.Value + + if float64(availMLU) < 1.5*neededMLU { + fmt.Printf("Schedule Scheme is wrong: Insufficient mlu") + return false, nil + } + } + + neededStorage := t.Job.Info.Resources.Storage + if neededStorage > 0 { + resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ + NodeId: t.preAdjustNodeID, + ResourceType: models.ResourceTypeStorage, + }) + if err != nil { + return false, err + } + + availStorage := resp.Data.(models.StorageResourceData).Available.Value + + bytesStorage := convertto.GBToBytes(availStorage) + + if bytesStorage < int64(1.5*float64(neededStorage)) { + fmt.Printf("Schedule Scheme is wrong: Insufficient storage") + return false, nil + } + } + + neededMemory := t.Job.Info.Resources.Memory + if neededMemory > 0 { + resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ + NodeId: t.preAdjustNodeID, + ResourceType: models.ResourceTypeMemory, + }) + if err != nil { + return false, err + } + + availMemory := resp.Data.(models.MemoryResourceData).Available.Value + + bytesMemory := convertto.GBToBytes(availMemory) + + if bytesMemory < int64(1.5*float64(neededMemory)) { + fmt.Printf("Schedule Scheme is wrong: Insufficient memory") + return false, nil + } + } + return true, nil + +} diff --git a/common/pkgs/mq/advisor/apis.go b/common/pkgs/mq/advisor/apis.go index d67646f..8eefba4 100644 --- a/common/pkgs/mq/advisor/apis.go +++ b/common/pkgs/mq/advisor/apis.go @@ -3,20 +3,19 @@ package advisor import ( "gitlink.org.cn/cloudream/common/models" "gitlink.org.cn/cloudream/common/pkgs/mq" + "gitlink.org.cn/cloudream/scheduler/common/models/job" ) // 获取调度方案 var _ = Register(Service.StartGetScheduleScheme) type StartGetScheduleScheme struct { - // UserID int64 `json:"userID"` - // PackageID int64 `json:"packageID"` + Job job.NormalJob `json:"job"` } -func NewStartGetScheduleScheme() StartGetScheduleScheme { +func NewStartGetScheduleScheme(job job.NormalJob) StartGetScheduleScheme { return StartGetScheduleScheme{ - // UserID: userID, - // PackageID: packageID, + Job: job, } } diff --git a/common/pkgs/mq/advisor/task/task.go b/common/pkgs/mq/advisor/task/task.go index 27e34c1..f96b3f8 100644 --- a/common/pkgs/mq/advisor/task/task.go +++ b/common/pkgs/mq/advisor/task/task.go @@ -1,73 +1,21 @@ package task -type TaskStatus interface{} - -type TaskStatusConst interface { - TaskStatus | ScheduleTaskStatus | UploadImageTaskStatus +type TaskStatus struct { + Status string `json:"status"` + Error string `json:"error"` + IsAdjustment bool `json:"isAdjustment"` + AdjustedScheme AdjustedScheme `json:"adjustedScheme"` } -type ScheduleTaskStatus struct { - Status string `json:"status"` - Error string `json:"error"` - PCMJobID int64 `json:"pcmJobID"` +type AdjustedScheme struct { + NodeID int64 `json:"nodeID"` } -func NewScheduleTaskStatus(status string, err string, pcmJobID int64) ScheduleTaskStatus { - return ScheduleTaskStatus{ - Status: status, - Error: err, - PCMJobID: pcmJobID, - } -} - -type UploadImageTaskStatus struct { - Status string `json:"status"` - Error string `json:"error"` - ImageID int64 `json:"imageID"` -} - -func NewUploadImageTaskStatus(status string, err string, imageID int64) UploadImageTaskStatus { - return UploadImageTaskStatus{ - Status: status, - Error: err, - ImageID: imageID, - } -} - -type CacheMovePackageTaskStatus struct { - Status string `json:"status"` - Error string `json:"error"` -} - -func NewCacheMovePackageTaskStatus(status string, err string) CacheMovePackageTaskStatus { - return CacheMovePackageTaskStatus{ - Status: status, - Error: err, - } -} - -type CreatePackageTaskStatus struct { - Status string `json:"status"` - Error string `json:"error"` - PackageID int64 `json:"packageID"` -} - -func NewCreatePackageTaskStatus(status string, err string, packageID int64) CreatePackageTaskStatus { - return CreatePackageTaskStatus{ - Status: status, - Error: err, - PackageID: packageID, - } -} - -type LoadPackageTaskStatus struct { - Status string `json:"status"` - Error string `json:"error"` -} - -func NewLoadPackageTaskStatus(status string, err string) LoadPackageTaskStatus { - return LoadPackageTaskStatus{ - Status: status, - Error: err, +func NewTaskStatus(status string, err string, isAdjustment bool, adjustedScheme AdjustedScheme) TaskStatus { + return TaskStatus{ + Status: status, + Error: err, + IsAdjustment: isAdjustment, + AdjustedScheme: adjustedScheme, } } diff --git a/common/pkgs/mq/executor/pcm.go b/common/pkgs/mq/executor/pcm.go index b2038ec..8190efb 100644 --- a/common/pkgs/mq/executor/pcm.go +++ b/common/pkgs/mq/executor/pcm.go @@ -18,16 +18,16 @@ type PCMService interface { var _ = Register(PCMService.StartUploadImage) type StartUploadImage struct { - NodeID int64 `json:"nodeID"` + SlwNodeID int64 `json:"slwNodeID"` ImagePath string `json:"imagePath"` } type StartUploadImageResp struct { TaskID string `json:"taskID"` } -func NewStartUploadImage(nodeID int64, imagePath string) StartUploadImage { +func NewStartUploadImage(slwNodeID int64, imagePath string) StartUploadImage { return StartUploadImage{ - NodeID: nodeID, + SlwNodeID: slwNodeID, ImagePath: imagePath, } } @@ -44,12 +44,12 @@ func (c *Client) StartUploadImage(msg StartUploadImage, opts ...mq.RequestOption var _ = Register(PCMService.GetImageList) type GetImageList struct { - NodeID int64 `json:"nodeID"` + SlwNodeID int64 `json:"slwNodeID"` } -func NewGetImageList(nodeID int64) GetImageList { +func NewGetImageList(slwNodeID int64) GetImageList { return GetImageList{ - NodeID: nodeID, + SlwNodeID: slwNodeID, } } @@ -71,14 +71,14 @@ func (c *Client) GetImageList(msg GetImageList, opts ...mq.RequestOption) (*GetI var _ = Register(PCMService.DeleteImage) type DeleteImage struct { - NodeID int64 `json:"nodeID"` - PCMJobID int64 `json:"pcmJobID"` + SlwNodeID int64 `json:"slwNodeID"` + PCMJobID int64 `json:"pcmJobID"` } -func NewDeleteImage(nodeID int64, pcmJobID int64) DeleteImage { +func NewDeleteImage(slwNodeID int64, pcmJobID int64) DeleteImage { return DeleteImage{ - NodeID: nodeID, - PCMJobID: pcmJobID, + SlwNodeID: slwNodeID, + PCMJobID: pcmJobID, } } @@ -100,21 +100,21 @@ func (c *Client) DeleteImage(msg DeleteImage, opts ...mq.RequestOption) (*Delete var _ = Register(PCMService.StartScheduleTask) type StartScheduleTask struct { - NodeID int64 `json:"nodeID"` - Envs []map[string]string `json:"envs"` - ImageID int64 `json:"imageID"` - CMDLine string `json:"cmdLine"` + SlwNodeID int64 `json:"slwNodeID"` + Envs []map[string]string `json:"envs"` + ImageID int64 `json:"imageID"` + CMDLine string `json:"cmdLine"` } type StartScheduleTaskResp struct { TaskID string `json:"taskID"` } -func NewStartScheduleTask(nodeID int64, envs []map[string]string, imageID int64, cmdLine string) StartScheduleTask { +func NewStartScheduleTask(slwNodeID int64, envs []map[string]string, imageID int64, cmdLine string) StartScheduleTask { return StartScheduleTask{ - NodeID: nodeID, - Envs: envs, - ImageID: imageID, - CMDLine: cmdLine, + SlwNodeID: slwNodeID, + Envs: envs, + ImageID: imageID, + CMDLine: cmdLine, } } func NewStartScheduleTaskResp(taskID string) StartScheduleTaskResp { @@ -130,14 +130,14 @@ func (c *Client) StartScheduleTask(msg StartUploadImage, opts ...mq.RequestOptio var _ = Register(PCMService.DeleteTask) type DeleteTask struct { - NodeID int64 `json:"nodeID"` - PCMJobID int64 `json:"pcmJobID"` + SlwNodeID int64 `json:"slwNodeID"` + PCMJobID int64 `json:"pcmJobID"` } -func NewDeleteTask(nodeID int64, pcmJobID int64) DeleteTask { +func NewDeleteTask(slwNodeID int64, pcmJobID int64) DeleteTask { return DeleteTask{ - NodeID: nodeID, - PCMJobID: pcmJobID, + SlwNodeID: slwNodeID, + PCMJobID: pcmJobID, } } diff --git a/common/pkgs/mq/executor/storage.go b/common/pkgs/mq/executor/storage.go index a9b02d7..be2ab4f 100644 --- a/common/pkgs/mq/executor/storage.go +++ b/common/pkgs/mq/executor/storage.go @@ -80,17 +80,17 @@ var _ = Register(StorageService.StartCacheMovePackage) type StartCacheMovePackage struct { UserID int64 `json:"userID"` PackageID int64 `json:"packageID"` - NodeID int64 `json:"nodeID"` + StgNodeID int64 `json:"stgNodeID"` } type StartCacheMovePackageResp struct { TaskID string `json:"taskID"` } -func NewStartCacheMovePackage(userID int64, packageID int64, nodeID int64) StartCacheMovePackage { +func NewStartCacheMovePackage(userID int64, packageID int64, stgNodeID int64) StartCacheMovePackage { return StartCacheMovePackage{ UserID: userID, PackageID: packageID, - NodeID: nodeID, + StgNodeID: stgNodeID, } } func NewStartCacheMovePackageResp(taskID string) StartCacheMovePackageResp { diff --git a/common/pkgs/mq/manager/api.go b/common/pkgs/mq/manager/api.go index 5055c27..6f0014d 100644 --- a/common/pkgs/mq/manager/api.go +++ b/common/pkgs/mq/manager/api.go @@ -61,10 +61,10 @@ func NewReportAdvisorTaskStatus(advisorID string, taskStatus []AdvisorTaskStatus TaskStatus: taskStatus, } } -func NewReportAdvisorTaskStatusResp() ReportExecutorTaskStatusResp { - return ReportExecutorTaskStatusResp{} +func NewReportAdvisorTaskStatusResp() ReportAdvisorTaskStatusResp { + return ReportAdvisorTaskStatusResp{} } -func NewAdvisorTaskStatus[T exectsk.TaskStatusConst](taskID string, status T) AdvisorTaskStatus { +func NewAdvisorTaskStatus(taskID string, status advtsk.TaskStatus) AdvisorTaskStatus { return AdvisorTaskStatus{ TaskID: taskID, Status: status, diff --git a/executor/internal/services/pcm.go b/executor/internal/services/pcm.go index bbf0123..807e4a0 100644 --- a/executor/internal/services/pcm.go +++ b/executor/internal/services/pcm.go @@ -11,7 +11,7 @@ import ( ) func (svc *Service) StartUploadImage(msg *execmq.StartUploadImage) (*execmq.StartUploadImageResp, *mq.CodeMessage) { - tsk := svc.taskManager.StartNew(schtsk.NewPCMUploadImage(msg.NodeID, msg.ImagePath)) + tsk := svc.taskManager.StartNew(schtsk.NewPCMUploadImage(msg.SlwNodeID, msg.ImagePath)) return mq.ReplyOK(execmq.NewStartUploadImageResp(tsk.ID())) } @@ -24,7 +24,7 @@ func (svc *Service) GetImageList(msg *execmq.GetImageList) (*execmq.GetImageList defer pcmCli.Close() resp, err := pcmCli.GetImageList(pcm.GetImageListReq{ - NodeID: msg.NodeID, + SlwNodeID: msg.SlwNodeID, }) if err != nil { logger.Warnf("get image list failed, err: %s", err.Error()) @@ -43,8 +43,8 @@ func (svc *Service) DeleteImage(msg *execmq.DeleteImage) (*execmq.DeleteImageRes defer pcmCli.Close() resp, err := pcmCli.DeleteImage(pcm.DeleteImageReq{ - NodeID: msg.NodeID, - PCMJobID: msg.PCMJobID, + SlwNodeID: msg.SlwNodeID, + PCMJobID: msg.PCMJobID, }) if err != nil { logger.Warnf("delete image failed, err: %s", err.Error()) @@ -54,7 +54,7 @@ func (svc *Service) DeleteImage(msg *execmq.DeleteImage) (*execmq.DeleteImageRes } func (svc *Service) StartScheduleTask(msg *execmq.StartScheduleTask) (*execmq.StartScheduleTaskResp, *mq.CodeMessage) { - tsk := svc.taskManager.StartNew(schtsk.NewPCMScheduleTask(msg.NodeID, msg.Envs, msg.ImageID, msg.CMDLine)) + tsk := svc.taskManager.StartNew(schtsk.NewPCMScheduleTask(msg.SlwNodeID, msg.Envs, msg.ImageID, msg.CMDLine)) return mq.ReplyOK(execmq.NewStartScheduleTaskResp(tsk.ID())) } @@ -67,8 +67,8 @@ func (svc *Service) DeleteTask(msg *execmq.DeleteTask) (*execmq.DeleteTaskResp, defer pcmCli.Close() resp, err := pcmCli.DeleteTask(pcm.DeleteTaskReq{ - NodeID: msg.NodeID, - PCMJobID: msg.PCMJobID, + SlwNodeID: msg.SlwNodeID, + PCMJobID: msg.PCMJobID, }) if err != nil { logger.Warnf("delete task failed, err: %s", err.Error()) diff --git a/executor/internal/services/storage.go b/executor/internal/services/storage.go index c335b4f..da29fcc 100644 --- a/executor/internal/services/storage.go +++ b/executor/internal/services/storage.go @@ -17,7 +17,7 @@ func (svc *Service) StartStorageCreatePackage(msg *execmq.StartStorageCreatePack } func (svc *Service) StartCacheMovePackage(msg *execmq.StartCacheMovePackage) (*execmq.StartCacheMovePackageResp, *mq.CodeMessage) { - tsk := svc.taskManager.StartNew(schtsk.NewCacheMovePackage(msg.UserID, msg.PackageID, msg.NodeID)) + tsk := svc.taskManager.StartNew(schtsk.NewCacheMovePackage(msg.UserID, msg.PackageID, msg.StgNodeID)) // tsk := svc.taskManager.StartNew(task.TaskBody[schtsk.NewCacheMovePackage(msg.UserID, msg.PackageID, msg.NodeID)]) return mq.ReplyOK(execmq.NewStartCacheMovePackageResp(tsk.ID())) } diff --git a/executor/internal/task/cache_move_package.go b/executor/internal/task/cache_move_package.go index c7175e5..05994ab 100644 --- a/executor/internal/task/cache_move_package.go +++ b/executor/internal/task/cache_move_package.go @@ -14,14 +14,14 @@ import ( type CacheMovePackage struct { userID int64 packageID int64 - nodeID int64 + stgNodeID int64 } -func NewCacheMovePackage(userID int64, packageID int64, nodeID int64) *CacheMovePackage { +func NewCacheMovePackage(userID int64, packageID int64, stgNodeID int64) *CacheMovePackage { return &CacheMovePackage{ userID: userID, packageID: packageID, - nodeID: nodeID, + stgNodeID: stgNodeID, } } @@ -54,6 +54,6 @@ func (t *CacheMovePackage) do(ctx TaskContext) error { return stgCli.CacheMovePackage(storage.CacheMovePackageReq{ UserID: t.userID, PackageID: t.packageID, - NodeID: t.packageID, + StgNodeID: t.stgNodeID, }) } diff --git a/executor/internal/task/pcm_schedule_task.go b/executor/internal/task/pcm_schedule_task.go index 919316a..0e28fcc 100644 --- a/executor/internal/task/pcm_schedule_task.go +++ b/executor/internal/task/pcm_schedule_task.go @@ -13,18 +13,18 @@ import ( ) type PCMScheduleTask struct { - nodeID int64 - envs []map[string]string - imageID int64 - cmdLine string + slwNodeID int64 + envs []map[string]string + imageID int64 + cmdLine string } -func NewPCMScheduleTask(nodeID int64, envs []map[string]string, imageID int64, cmdLine string) *PCMScheduleTask { +func NewPCMScheduleTask(slwNodeID int64, envs []map[string]string, imageID int64, cmdLine string) *PCMScheduleTask { return &PCMScheduleTask{ - nodeID: nodeID, - envs: envs, - imageID: imageID, - cmdLine: cmdLine, + slwNodeID: slwNodeID, + envs: envs, + imageID: imageID, + cmdLine: cmdLine, } } @@ -53,10 +53,10 @@ func (t *PCMScheduleTask) do(taskID string, ctx TaskContext) error { defer pcmCli.Close() resp, err := pcmCli.ScheduleTask(pcm.ScheduleTaskReq{ - NodeID: t.nodeID, - Envs: t.envs, - ImageID: t.imageID, - CMDLine: t.cmdLine, + SlwNodeID: t.slwNodeID, + Envs: t.envs, + ImageID: t.imageID, + CMDLine: t.cmdLine, }) if err != nil { @@ -66,8 +66,8 @@ func (t *PCMScheduleTask) do(taskID string, ctx TaskContext) error { var prevStatus string for { tsResp, err := pcmCli.GetTaskStatus(pcm.GetTaskStatusReq{ - NodeID: t.nodeID, - PCMJobID: resp.PCMJobID, + SlwNodeID: t.slwNodeID, + PCMJobID: resp.PCMJobID, }) if err != nil { return err diff --git a/executor/internal/task/pcm_upload_img.go b/executor/internal/task/pcm_upload_img.go index cf52d03..f1cd3e0 100644 --- a/executor/internal/task/pcm_upload_img.go +++ b/executor/internal/task/pcm_upload_img.go @@ -12,13 +12,13 @@ import ( ) type PCMUploadImage struct { - nodeID int64 + slwNodeID int64 imagePath string } -func NewPCMUploadImage(nodeID int64, imagePath string) *PCMUploadImage { +func NewPCMUploadImage(slwNodeID int64, imagePath string) *PCMUploadImage { return &PCMUploadImage{ - nodeID: nodeID, + slwNodeID: slwNodeID, imagePath: imagePath, } } @@ -48,7 +48,7 @@ func (t *PCMUploadImage) do(taskID string, ctx TaskContext) error { defer pcmCli.Close() resp, err := pcmCli.UploadImage(pcm.UploadImageReq{ - NodeID: t.nodeID, + SlwNodeID: t.slwNodeID, ImagePath: t.imagePath, }) if err != nil { diff --git a/go.mod b/go.mod index 3116847..bdc1473 100644 --- a/go.mod +++ b/go.mod @@ -26,6 +26,7 @@ require ( github.com/hashicorp/errwrap v1.1.0 // indirect github.com/hashicorp/go-multierror v1.1.1 // indirect github.com/imdario/mergo v0.3.15 // indirect + github.com/inhies/go-bytesize v0.0.0-20220417184213-4913239db9cf github.com/json-iterator/go v1.1.12 // indirect github.com/klauspost/cpuid/v2 v2.2.4 // indirect github.com/leodido/go-urn v1.2.4 // indirect diff --git a/go.sum b/go.sum index f7a58be..2a1e95d 100644 --- a/go.sum +++ b/go.sum @@ -41,6 +41,8 @@ github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+l github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM= github.com/imdario/mergo v0.3.15 h1:M8XP7IuFNsqUx6VPK2P9OSmsYsI/YFaGil0uD21V3dM= github.com/imdario/mergo v0.3.15/go.mod h1:WBLT9ZmE3lPoWsEzCh9LPo3TiwVN+ZKEjmz+hD27ysY= +github.com/inhies/go-bytesize v0.0.0-20220417184213-4913239db9cf h1:FtEj8sfIcaaBfAKrE1Cwb61YDtYq9JxChK1c7AKce7s= +github.com/inhies/go-bytesize v0.0.0-20220417184213-4913239db9cf/go.mod h1:yrqSXGoD/4EKfF26AOGzscPOgTTJcyAwM2rpixWT+t4= github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnrnM= github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/jtolds/gls v4.20.0+incompatible h1:xdiiI2gbIgH/gLH7ADydsJ1uDOEzR8yvV7C0MuV77Wo= From 2ef2715c73ddb59e90a0a0709bfea149c3c56468 Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Tue, 12 Sep 2023 10:03:07 +0800 Subject: [PATCH 2/6] . --- advisor/internal/services/advisor.go | 6 ++--- advisor/internal/task/schedule_scheme.go | 31 ++++++++++++------------ common/pkgs/mq/advisor/apis.go | 18 +++++++------- common/pkgs/mq/advisor/server.go | 2 +- common/pkgs/mq/advisor/task/task.go | 23 +++++++++++++++--- 5 files changed, 49 insertions(+), 31 deletions(-) diff --git a/advisor/internal/services/advisor.go b/advisor/internal/services/advisor.go index 056769b..5cf9584 100644 --- a/advisor/internal/services/advisor.go +++ b/advisor/internal/services/advisor.go @@ -6,7 +6,7 @@ import ( advmq "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/advisor" ) -func (svc *Service) StartGetScheduleScheme(msg *advmq.StartGetScheduleScheme) (*advmq.StartGetScheduleSchemeResp, *mq.CodeMessage) { - tsk := svc.taskManager.StartNew(schtsk.NewGetScheduleScheme()) - return mq.ReplyOK(advmq.NewStartGetScheduleSchemeResp(tsk.ID())) +func (svc *Service) StartMakeScheduleScheme(msg *advmq.StartMakeScheduleScheme) (*advmq.StartMakeScheduleSchemeResp, *mq.CodeMessage) { + tsk := svc.taskManager.StartNew(schtsk.NewMakeScheduleScheme()) + return mq.ReplyOK(advmq.NewStartMakeScheduleSchemeResp(tsk.ID())) } diff --git a/advisor/internal/task/schedule_scheme.go b/advisor/internal/task/schedule_scheme.go index eb717c0..289cb41 100644 --- a/advisor/internal/task/schedule_scheme.go +++ b/advisor/internal/task/schedule_scheme.go @@ -14,26 +14,27 @@ import ( "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/collector" ) -type GetScheduleScheme struct { +type MakeScheduleScheme struct { Job job.NormalJob preAdjustNodeID int64 } -func NewGetScheduleScheme() *GetScheduleScheme { - return &GetScheduleScheme{} +func NewMakeScheduleScheme() *MakeScheduleScheme { + return &MakeScheduleScheme{} } -func (t *GetScheduleScheme) Execute(task *task.Task[TaskContext], ctx TaskContext, complete CompleteFn) { - log := logger.WithType[GetScheduleScheme]("Task") +func (t *MakeScheduleScheme) Execute(task *task.Task[TaskContext], ctx TaskContext, complete CompleteFn) { + log := logger.WithType[MakeScheduleScheme]("Task") log.Debugf("begin") defer log.Debugf("end") err := t.do(task.ID(), ctx) if err != nil { //TODO 若任务失败,上报的状态failed字段根据情况修改 - ctx.reporter.Report(task.ID(), advtsk.NewTaskStatus("failed", err.Error(), true, advtsk.AdjustedScheme{})) + ctx.reporter.Report(task.ID(), advtsk.NewScheduleSchemeTaskStatus("failed", err.Error(), true, advtsk.AdjustedScheme{})) } else { - ctx.reporter.Report(task.ID(), advtsk.NewTaskStatus("failed", err.Error(), false, advtsk.AdjustedScheme{})) + ///////// 修改 + ctx.reporter.Report(task.ID(), advtsk.NewScheduleSchemeTaskStatus("failed", "", false, advtsk.AdjustedScheme{})) } ctx.reporter.ReportNow() @@ -42,7 +43,7 @@ func (t *GetScheduleScheme) Execute(task *task.Task[TaskContext], ctx TaskContex }) } -func (t *GetScheduleScheme) do(taskID string, ctx TaskContext) error { +func (t *MakeScheduleScheme) do(taskID string, ctx TaskContext) error { isAvailable, err := t.CheckResourceAvailability() if err != nil { return err @@ -57,7 +58,7 @@ func (t *GetScheduleScheme) do(taskID string, ctx TaskContext) error { } // 检查预调度节点资源是否足够 -func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { +func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { colCli, err := globals.CollectorMQPool.Acquire() if err != nil { return false, fmt.Errorf("new collector client: %w", err) @@ -67,7 +68,7 @@ func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { neededCPU := t.Job.Info.Resources.CPU if neededCPU > 0 { resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - NodeId: t.preAdjustNodeID, + SlwNodeID: t.preAdjustNodeID, ResourceType: models.ResourceTypeCPU, }) if err != nil { @@ -85,7 +86,7 @@ func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { neededNPU := t.Job.Info.Resources.NPU if neededNPU > 0 { resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - NodeId: t.preAdjustNodeID, + SlwNodeID: t.preAdjustNodeID, ResourceType: models.ResourceTypeNPU, }) if err != nil { @@ -103,7 +104,7 @@ func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { neededGPU := t.Job.Info.Resources.GPU if neededGPU > 0 { resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - NodeId: t.preAdjustNodeID, + SlwNodeID: t.preAdjustNodeID, ResourceType: models.ResourceTypeGPU, }) if err != nil { @@ -121,7 +122,7 @@ func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { neededMLU := t.Job.Info.Resources.MLU if neededMLU > 0 { resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - NodeId: t.preAdjustNodeID, + SlwNodeID: t.preAdjustNodeID, ResourceType: models.ResourceTypeMLU, }) if err != nil { @@ -139,7 +140,7 @@ func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { neededStorage := t.Job.Info.Resources.Storage if neededStorage > 0 { resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - NodeId: t.preAdjustNodeID, + SlwNodeID: t.preAdjustNodeID, ResourceType: models.ResourceTypeStorage, }) if err != nil { @@ -159,7 +160,7 @@ func (t *GetScheduleScheme) CheckResourceAvailability() (bool, error) { neededMemory := t.Job.Info.Resources.Memory if neededMemory > 0 { resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - NodeId: t.preAdjustNodeID, + SlwNodeID: t.preAdjustNodeID, ResourceType: models.ResourceTypeMemory, }) if err != nil { diff --git a/common/pkgs/mq/advisor/apis.go b/common/pkgs/mq/advisor/apis.go index 8eefba4..a62d6e7 100644 --- a/common/pkgs/mq/advisor/apis.go +++ b/common/pkgs/mq/advisor/apis.go @@ -7,30 +7,30 @@ import ( ) // 获取调度方案 -var _ = Register(Service.StartGetScheduleScheme) +var _ = Register(Service.StartMakeScheduleScheme) -type StartGetScheduleScheme struct { +type StartMakeScheduleScheme struct { Job job.NormalJob `json:"job"` } -func NewStartGetScheduleScheme(job job.NormalJob) StartGetScheduleScheme { - return StartGetScheduleScheme{ +func NewStartGetScheduleScheme(job job.NormalJob) StartMakeScheduleScheme { + return StartMakeScheduleScheme{ Job: job, } } -type StartGetScheduleSchemeResp struct { +type StartMakeScheduleSchemeResp struct { TaskID string `json:"taskID"` } -func NewStartGetScheduleSchemeResp(taskID string) StartGetScheduleSchemeResp { - return StartGetScheduleSchemeResp{ +func NewStartMakeScheduleSchemeResp(taskID string) StartMakeScheduleSchemeResp { + return StartMakeScheduleSchemeResp{ TaskID: taskID, } } -func (c *Client) StartGetScheduleScheme(msg StartGetScheduleScheme, opts ...mq.RequestOption) (*StartGetScheduleSchemeResp, error) { - return mq.Request[StartGetScheduleSchemeResp](c.rabbitCli, msg, opts...) +func (c *Client) StartMakeScheduleScheme(msg StartMakeScheduleScheme, opts ...mq.RequestOption) (*StartMakeScheduleSchemeResp, error) { + return mq.Request[StartMakeScheduleSchemeResp](c.rabbitCli, msg, opts...) } func init() { diff --git a/common/pkgs/mq/advisor/server.go b/common/pkgs/mq/advisor/server.go index aa05d35..8213542 100644 --- a/common/pkgs/mq/advisor/server.go +++ b/common/pkgs/mq/advisor/server.go @@ -10,7 +10,7 @@ const ( ) type Service interface { - StartGetScheduleScheme(msg *StartGetScheduleScheme) (*StartGetScheduleSchemeResp, *mq.CodeMessage) + StartMakeScheduleScheme(msg *StartMakeScheduleScheme) (*StartMakeScheduleSchemeResp, *mq.CodeMessage) } type Server struct { diff --git a/common/pkgs/mq/advisor/task/task.go b/common/pkgs/mq/advisor/task/task.go index f96b3f8..08fec7a 100644 --- a/common/pkgs/mq/advisor/task/task.go +++ b/common/pkgs/mq/advisor/task/task.go @@ -1,6 +1,23 @@ package task -type TaskStatus struct { +import ( + "gitlink.org.cn/cloudream/common/pkgs/types" + myreflect "gitlink.org.cn/cloudream/common/utils/reflect" +) + +type TaskStatus interface{} + +// 增加了新类型后需要在这里也同步添加 +type TaskStatusConst interface { + TaskStatus | ScheduleSchemeTaskStatus +} + +// 增加了新类型后需要在这里也同步添加 +var TaskStatusTypeUnion = types.NewTypeUnion[TaskStatus]( + myreflect.TypeOf[ScheduleSchemeTaskStatus](), +) + +type ScheduleSchemeTaskStatus struct { Status string `json:"status"` Error string `json:"error"` IsAdjustment bool `json:"isAdjustment"` @@ -11,8 +28,8 @@ type AdjustedScheme struct { NodeID int64 `json:"nodeID"` } -func NewTaskStatus(status string, err string, isAdjustment bool, adjustedScheme AdjustedScheme) TaskStatus { - return TaskStatus{ +func NewScheduleSchemeTaskStatus(status string, err string, isAdjustment bool, adjustedScheme AdjustedScheme) TaskStatus { + return ScheduleSchemeTaskStatus{ Status: status, Error: err, IsAdjustment: isAdjustment, From ee55a173a3d72a2147638730decc450b5c3399bc Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Tue, 12 Sep 2023 11:23:27 +0800 Subject: [PATCH 3/6] =?UTF-8?q?=E6=9B=B4=E6=96=B0advisor=20mq?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- advisor/internal/task/schedule_scheme.go | 48 +++++++++++------------ common/pkgs/mq/advisor/apis.go | 31 ++++++++------- common/pkgs/mq/advisor/server.go | 4 +- common/pkgs/mq/executor/pcm.go | 50 ++++++++++++------------ common/pkgs/mq/executor/storage.go | 6 +-- 5 files changed, 70 insertions(+), 69 deletions(-) diff --git a/advisor/internal/task/schedule_scheme.go b/advisor/internal/task/schedule_scheme.go index 289cb41..27f3b14 100644 --- a/advisor/internal/task/schedule_scheme.go +++ b/advisor/internal/task/schedule_scheme.go @@ -67,10 +67,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { neededCPU := t.Job.Info.Resources.CPU if neededCPU > 0 { - resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - SlwNodeID: t.preAdjustNodeID, - ResourceType: models.ResourceTypeCPU, - }) + resp, err := colCli.GetOneResourceData(collector.NewGetOneResourceData( + t.preAdjustNodeID, + models.ResourceTypeCPU, + )) if err != nil { return false, err } @@ -85,10 +85,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { neededNPU := t.Job.Info.Resources.NPU if neededNPU > 0 { - resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - SlwNodeID: t.preAdjustNodeID, - ResourceType: models.ResourceTypeNPU, - }) + resp, err := colCli.GetOneResourceData(collector.NewGetOneResourceData( + t.preAdjustNodeID, + models.ResourceTypeNPU, + )) if err != nil { return false, err } @@ -103,10 +103,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { neededGPU := t.Job.Info.Resources.GPU if neededGPU > 0 { - resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - SlwNodeID: t.preAdjustNodeID, - ResourceType: models.ResourceTypeGPU, - }) + resp, err := colCli.GetOneResourceData(collector.NewGetOneResourceData( + t.preAdjustNodeID, + models.ResourceTypeGPU, + )) if err != nil { return false, err } @@ -121,10 +121,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { neededMLU := t.Job.Info.Resources.MLU if neededMLU > 0 { - resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - SlwNodeID: t.preAdjustNodeID, - ResourceType: models.ResourceTypeMLU, - }) + resp, err := colCli.GetOneResourceData(collector.NewGetOneResourceData( + t.preAdjustNodeID, + models.ResourceTypeMLU, + )) if err != nil { return false, err } @@ -139,10 +139,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { neededStorage := t.Job.Info.Resources.Storage if neededStorage > 0 { - resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - SlwNodeID: t.preAdjustNodeID, - ResourceType: models.ResourceTypeStorage, - }) + resp, err := colCli.GetOneResourceData(collector.NewGetOneResourceData( + t.preAdjustNodeID, + models.ResourceTypeStorage, + )) if err != nil { return false, err } @@ -159,10 +159,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { neededMemory := t.Job.Info.Resources.Memory if neededMemory > 0 { - resp, err := colCli.GetOneResourceData(collector.GetOneResourceData{ - SlwNodeID: t.preAdjustNodeID, - ResourceType: models.ResourceTypeMemory, - }) + resp, err := colCli.GetOneResourceData(collector.NewGetOneResourceData( + t.preAdjustNodeID, + models.ResourceTypeMemory, + )) if err != nil { return false, err } diff --git a/common/pkgs/mq/advisor/apis.go b/common/pkgs/mq/advisor/apis.go index a62d6e7..5c8d3ed 100644 --- a/common/pkgs/mq/advisor/apis.go +++ b/common/pkgs/mq/advisor/apis.go @@ -1,38 +1,39 @@ package advisor import ( - "gitlink.org.cn/cloudream/common/models" "gitlink.org.cn/cloudream/common/pkgs/mq" "gitlink.org.cn/cloudream/scheduler/common/models/job" ) +type ApiService interface { + StartMakeScheduleScheme(msg *StartMakeScheduleScheme) (*StartMakeScheduleSchemeResp, *mq.CodeMessage) +} + // 获取调度方案 var _ = Register(Service.StartMakeScheduleScheme) type StartMakeScheduleScheme struct { + mq.MessageBodyBase Job job.NormalJob `json:"job"` } -func NewStartGetScheduleScheme(job job.NormalJob) StartMakeScheduleScheme { - return StartMakeScheduleScheme{ +type StartMakeScheduleSchemeResp struct { + mq.MessageBodyBase + TaskID string `json:"taskID"` +} + +func NewStartGetScheduleScheme(job job.NormalJob) *StartMakeScheduleScheme { + return &StartMakeScheduleScheme{ Job: job, } } -type StartMakeScheduleSchemeResp struct { - TaskID string `json:"taskID"` -} - -func NewStartMakeScheduleSchemeResp(taskID string) StartMakeScheduleSchemeResp { - return StartMakeScheduleSchemeResp{ +func NewStartMakeScheduleSchemeResp(taskID string) *StartMakeScheduleSchemeResp { + return &StartMakeScheduleSchemeResp{ TaskID: taskID, } } -func (c *Client) StartMakeScheduleScheme(msg StartMakeScheduleScheme, opts ...mq.RequestOption) (*StartMakeScheduleSchemeResp, error) { - return mq.Request[StartMakeScheduleSchemeResp](c.rabbitCli, msg, opts...) -} - -func init() { - mq.RegisterUnionType(models.ResourceDataTypeUnion) +func (c *Client) StartMakeScheduleScheme(msg *StartMakeScheduleScheme, opts ...mq.RequestOption) (*StartMakeScheduleSchemeResp, error) { + return mq.Request(Service.StartMakeScheduleScheme, c.rabbitCli, msg, opts...) } diff --git a/common/pkgs/mq/advisor/server.go b/common/pkgs/mq/advisor/server.go index 1824570..2c59b41 100644 --- a/common/pkgs/mq/advisor/server.go +++ b/common/pkgs/mq/advisor/server.go @@ -6,11 +6,11 @@ import ( ) const ( - ServerQueueName = "Scheduler-Collector" + ServerQueueName = "Scheduler-Advisor" ) type Service interface { - StartMakeScheduleScheme(msg *StartMakeScheduleScheme) (*StartMakeScheduleSchemeResp, *mq.CodeMessage) + ApiService } type Server struct { diff --git a/common/pkgs/mq/executor/pcm.go b/common/pkgs/mq/executor/pcm.go index 4f5ddd9..b384c07 100644 --- a/common/pkgs/mq/executor/pcm.go +++ b/common/pkgs/mq/executor/pcm.go @@ -19,7 +19,7 @@ var _ = Register(Service.StartUploadImage) type StartUploadImage struct { mq.MessageBodyBase - NodeID int64 `json:"nodeID"` + SlwNodeID int64 `json:"slwNodeID"` ImagePath string `json:"imagePath"` } type StartUploadImageResp struct { @@ -27,9 +27,9 @@ type StartUploadImageResp struct { TaskID string `json:"taskID"` } -func NewStartUploadImage(nodeID int64, imagePath string) *StartUploadImage { +func NewStartUploadImage(slwNodeID int64, imagePath string) *StartUploadImage { return &StartUploadImage{ - NodeID: nodeID, + SlwNodeID: slwNodeID, ImagePath: imagePath, } } @@ -47,16 +47,16 @@ var _ = Register(Service.GetImageList) type GetImageList struct { mq.MessageBodyBase - NodeID int64 `json:"nodeID"` + SlwNodeID int64 `json:"slwNodeID"` } type GetImageListResp struct { mq.MessageBodyBase ImageIDs []int64 `json:"imageIDs"` } -func NewGetImageList(nodeID int64) *GetImageList { +func NewGetImageList(slwNodeID int64) *GetImageList { return &GetImageList{ - NodeID: nodeID, + SlwNodeID: slwNodeID, } } func NewGetImageListResp(imageIDs []int64) *GetImageListResp { @@ -73,18 +73,18 @@ var _ = Register(Service.DeleteImage) type DeleteImage struct { mq.MessageBodyBase - NodeID int64 `json:"nodeID"` - PCMJobID int64 `json:"pcmJobID"` + SlwNodeID int64 `json:"slwNodeID"` + PCMJobID int64 `json:"pcmJobID"` } type DeleteImageResp struct { mq.MessageBodyBase Result string `json:"result"` } -func NewDeleteImage(nodeID int64, pcmJobID int64) *DeleteImage { +func NewDeleteImage(slwNodeID int64, pcmJobID int64) *DeleteImage { return &DeleteImage{ - NodeID: nodeID, - PCMJobID: pcmJobID, + SlwNodeID: slwNodeID, + PCMJobID: pcmJobID, } } func NewDeleteImageResp(result string) *DeleteImageResp { @@ -101,22 +101,22 @@ var _ = Register(Service.StartScheduleTask) type StartScheduleTask struct { mq.MessageBodyBase - NodeID int64 `json:"nodeID"` - Envs []map[string]string `json:"envs"` - ImageID int64 `json:"imageID"` - CMDLine string `json:"cmdLine"` + SlwNodeID int64 `json:"slwNodeID"` + Envs []map[string]string `json:"envs"` + ImageID int64 `json:"imageID"` + CMDLine string `json:"cmdLine"` } type StartScheduleTaskResp struct { mq.MessageBodyBase TaskID string `json:"taskID"` } -func NewStartScheduleTask(nodeID int64, envs []map[string]string, imageID int64, cmdLine string) *StartScheduleTask { +func NewStartScheduleTask(slwNodeID int64, envs []map[string]string, imageID int64, cmdLine string) *StartScheduleTask { return &StartScheduleTask{ - NodeID: nodeID, - Envs: envs, - ImageID: imageID, - CMDLine: cmdLine, + SlwNodeID: slwNodeID, + Envs: envs, + ImageID: imageID, + CMDLine: cmdLine, } } func NewStartScheduleTaskResp(taskID string) *StartScheduleTaskResp { @@ -133,18 +133,18 @@ var _ = Register(Service.DeleteTask) type DeleteTask struct { mq.MessageBodyBase - NodeID int64 `json:"nodeID"` - PCMJobID int64 `json:"pcmJobID"` + SlwNodeID int64 `json:"slwNodeID"` + PCMJobID int64 `json:"pcmJobID"` } type DeleteTaskResp struct { mq.MessageBodyBase Result string `json:"result"` } -func NewDeleteTask(nodeID int64, pcmJobID int64) *DeleteTask { +func NewDeleteTask(slwNodeID int64, pcmJobID int64) *DeleteTask { return &DeleteTask{ - NodeID: nodeID, - PCMJobID: pcmJobID, + SlwNodeID: slwNodeID, + PCMJobID: pcmJobID, } } func NewDeleteTaskResp(result string) *DeleteTaskResp { diff --git a/common/pkgs/mq/executor/storage.go b/common/pkgs/mq/executor/storage.go index 02257f6..2f90701 100644 --- a/common/pkgs/mq/executor/storage.go +++ b/common/pkgs/mq/executor/storage.go @@ -85,18 +85,18 @@ type StartCacheMovePackage struct { mq.MessageBodyBase UserID int64 `json:"userID"` PackageID int64 `json:"packageID"` - NodeID int64 `json:"nodeID"` + StgNodeID int64 `json:"stgNodeID"` } type StartCacheMovePackageResp struct { mq.MessageBodyBase TaskID string `json:"taskID"` } -func NewStartCacheMovePackage(userID int64, packageID int64, nodeID int64) *StartCacheMovePackage { +func NewStartCacheMovePackage(userID int64, packageID int64, stgNodeID int64) *StartCacheMovePackage { return &StartCacheMovePackage{ UserID: userID, PackageID: packageID, - NodeID: nodeID, + StgNodeID: stgNodeID, } } func NewStartCacheMovePackageResp(taskID string) *StartCacheMovePackageResp { From 6c6ab4bc1829fc6e5bb43e97e2e1fd57a371662e Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Tue, 12 Sep 2023 15:34:43 +0800 Subject: [PATCH 4/6] =?UTF-8?q?=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- advisor/internal/task/schedule_scheme.go | 50 ++++++++++++-------- executor/internal/task/cache_move_package.go | 2 +- 2 files changed, 30 insertions(+), 22 deletions(-) diff --git a/advisor/internal/task/schedule_scheme.go b/advisor/internal/task/schedule_scheme.go index 27f3b14..52f1805 100644 --- a/advisor/internal/task/schedule_scheme.go +++ b/advisor/internal/task/schedule_scheme.go @@ -2,16 +2,18 @@ package task import ( "fmt" + "strconv" "time" "gitlink.org.cn/cloudream/common/models" "gitlink.org.cn/cloudream/common/pkgs/logger" "gitlink.org.cn/cloudream/common/pkgs/task" - "gitlink.org.cn/cloudream/common/utils/convertto" "gitlink.org.cn/cloudream/scheduler/common/globals" "gitlink.org.cn/cloudream/scheduler/common/models/job" advtsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/advisor/task" "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/collector" + + "github.com/inhies/go-bytesize" ) type MakeScheduleScheme struct { @@ -75,10 +77,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { return false, err } - availCPU := resp.Data.(models.CPUResourceData).Available.Value + availCPU := resp.Data.(models.CPUResourceData).Available - if float64(availCPU) < 1.5*neededCPU { - fmt.Printf("Schedule Scheme is wrong: Insufficient cpu") + if float64(availCPU.Value) < 1.5*neededCPU { + logger.Warnf("Schedule Scheme is wrong: Insufficient CPU resources. Available CPU: %d%s", availCPU.Value, availCPU.Unit) return false, nil } } @@ -93,10 +95,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { return false, err } - availNPU := resp.Data.(models.NPUResourceData).Available.Value + availNPU := resp.Data.(models.NPUResourceData).Available - if float64(availNPU) < 1.5*neededNPU { - fmt.Printf("Schedule Scheme is wrong: Insufficient npu") + if float64(availNPU.Value) < 1.5*neededNPU { + logger.Warnf("Schedule Scheme is wrong: Insufficient NPU resources. Available NPU: %d%s", availNPU.Value, availNPU.Unit) return false, nil } } @@ -111,10 +113,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { return false, err } - availGPU := resp.Data.(models.GPUResourceData).Available.Value + availGPU := resp.Data.(models.GPUResourceData).Available - if float64(availGPU) < 1.5*neededGPU { - fmt.Printf("Schedule Scheme is wrong: Insufficient gpu") + if float64(availGPU.Value) < 1.5*neededGPU { + logger.Warnf("Schedule Scheme is wrong: Insufficient GPU resources. Available GPU: %d%s", availGPU.Value, availGPU.Unit) return false, nil } } @@ -129,10 +131,10 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { return false, err } - availMLU := resp.Data.(models.MLUResourceData).Available.Value + availMLU := resp.Data.(models.MLUResourceData).Available - if float64(availMLU) < 1.5*neededMLU { - fmt.Printf("Schedule Scheme is wrong: Insufficient mlu") + if float64(availMLU.Value) < 1.5*neededMLU { + logger.Warnf("Schedule Scheme is wrong: Insufficient MLU resources. Available MLU: %d%s", availMLU.Value, availMLU.Unit) return false, nil } } @@ -147,12 +149,15 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { return false, err } - availStorage := resp.Data.(models.StorageResourceData).Available.Value + availStorage := resp.Data.(models.StorageResourceData).Available - bytesStorage := convertto.GBToBytes(availStorage) + bytesStorage, err := bytesize.Parse(strconv.FormatFloat(availStorage.Value, 'f', -1, 64) + availStorage.Unit) + if err != nil { + return false, err + } - if bytesStorage < int64(1.5*float64(neededStorage)) { - fmt.Printf("Schedule Scheme is wrong: Insufficient storage") + if int64(bytesStorage) < int64(1.5*float64(neededStorage)) { + logger.Warnf("Schedule Scheme is wrong: Insufficient storage resources. Available storage: %f%s", availStorage.Value, availStorage.Unit) return false, nil } } @@ -167,12 +172,15 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { return false, err } - availMemory := resp.Data.(models.MemoryResourceData).Available.Value + availMemory := resp.Data.(models.MemoryResourceData).Available - bytesMemory := convertto.GBToBytes(availMemory) + bytesMemory, err := bytesize.Parse(strconv.FormatFloat(availMemory.Value, 'f', -1, 64) + availMemory.Unit) + if err != nil { + return false, err + } - if bytesMemory < int64(1.5*float64(neededMemory)) { - fmt.Printf("Schedule Scheme is wrong: Insufficient memory") + if int64(bytesMemory) < int64(1.5*float64(neededMemory)) { + logger.Warnf("Schedule Scheme is wrong: Insufficient memory resources. Available memory: %f%s", availMemory.Value, availMemory.Unit) return false, nil } } diff --git a/executor/internal/task/cache_move_package.go b/executor/internal/task/cache_move_package.go index 05994ab..dbaacc2 100644 --- a/executor/internal/task/cache_move_package.go +++ b/executor/internal/task/cache_move_package.go @@ -54,6 +54,6 @@ func (t *CacheMovePackage) do(ctx TaskContext) error { return stgCli.CacheMovePackage(storage.CacheMovePackageReq{ UserID: t.userID, PackageID: t.packageID, - StgNodeID: t.stgNodeID, + NodeID: t.stgNodeID, }) } From c8026090878bfa50c10797d682e81f4dcc2a29e6 Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Tue, 12 Sep 2023 17:20:17 +0800 Subject: [PATCH 5/6] =?UTF-8?q?=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- advisor/internal/task/schedule_scheme.go | 25 ++++++++++++++---------- common/pkgs/mq/advisor/task/task.go | 16 ++++++++------- 2 files changed, 24 insertions(+), 17 deletions(-) diff --git a/advisor/internal/task/schedule_scheme.go b/advisor/internal/task/schedule_scheme.go index 52f1805..342cf1a 100644 --- a/advisor/internal/task/schedule_scheme.go +++ b/advisor/internal/task/schedule_scheme.go @@ -2,7 +2,6 @@ package task import ( "fmt" - "strconv" "time" "gitlink.org.cn/cloudream/common/models" @@ -80,7 +79,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availCPU := resp.Data.(models.CPUResourceData).Available if float64(availCPU.Value) < 1.5*neededCPU { - logger.Warnf("Schedule Scheme is wrong: Insufficient CPU resources. Available CPU: %d%s", availCPU.Value, availCPU.Unit) + logger.WithField("jobID", t.Job.JobID). + Infof("insufficient CPU resources: wanted CPU:%f, available CPU: %d%s", 1.5*neededCPU, availCPU.Value, availCPU.Unit) return false, nil } } @@ -98,7 +98,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availNPU := resp.Data.(models.NPUResourceData).Available if float64(availNPU.Value) < 1.5*neededNPU { - logger.Warnf("Schedule Scheme is wrong: Insufficient NPU resources. Available NPU: %d%s", availNPU.Value, availNPU.Unit) + logger.WithField("jobID", t.Job.JobID). + Infof("insufficient NPU resources: wanted NPU:%f, available NPU: %d%s", 1.5*neededNPU, availNPU.Value, availNPU.Unit) return false, nil } } @@ -116,7 +117,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availGPU := resp.Data.(models.GPUResourceData).Available if float64(availGPU.Value) < 1.5*neededGPU { - logger.Warnf("Schedule Scheme is wrong: Insufficient GPU resources. Available GPU: %d%s", availGPU.Value, availGPU.Unit) + logger.WithField("jobID", t.Job.JobID). + Infof("insufficient GPU resources: wanted GPU:%f, available GPU: %d%s", 1.5*neededGPU, availGPU.Value, availGPU.Unit) return false, nil } } @@ -134,7 +136,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availMLU := resp.Data.(models.MLUResourceData).Available if float64(availMLU.Value) < 1.5*neededMLU { - logger.Warnf("Schedule Scheme is wrong: Insufficient MLU resources. Available MLU: %d%s", availMLU.Value, availMLU.Unit) + logger.WithField("jobID", t.Job.JobID). + Infof("insufficient MLU resources: wanted MLU:%f, available MLU: %d%s", 1.5*neededMLU, availMLU.Value, availMLU.Unit) return false, nil } } @@ -151,13 +154,14 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availStorage := resp.Data.(models.StorageResourceData).Available - bytesStorage, err := bytesize.Parse(strconv.FormatFloat(availStorage.Value, 'f', -1, 64) + availStorage.Unit) + bytesStorage, err := bytesize.Parse(fmt.Sprintf("%f%s", availStorage.Value, availStorage.Unit)) if err != nil { return false, err } if int64(bytesStorage) < int64(1.5*float64(neededStorage)) { - logger.Warnf("Schedule Scheme is wrong: Insufficient storage resources. Available storage: %f%s", availStorage.Value, availStorage.Unit) + logger.WithField("jobID", t.Job.JobID). + Infof("insufficient storage resources: wanted storage:%s, available storage: %f%s", bytesize.New(1.5*float64(neededStorage)), availStorage.Value, availStorage.Unit) return false, nil } } @@ -174,16 +178,17 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availMemory := resp.Data.(models.MemoryResourceData).Available - bytesMemory, err := bytesize.Parse(strconv.FormatFloat(availMemory.Value, 'f', -1, 64) + availMemory.Unit) + bytesMemory, err := bytesize.Parse(fmt.Sprintf("%f%s", availMemory.Value, availMemory.Unit)) if err != nil { return false, err } if int64(bytesMemory) < int64(1.5*float64(neededMemory)) { - logger.Warnf("Schedule Scheme is wrong: Insufficient memory resources. Available memory: %f%s", availMemory.Value, availMemory.Unit) + logger.WithField("jobID", t.Job.JobID). + Infof("insufficient memory resources: wanted memory:%s, available memory: %f%s", bytesize.New(1.5*float64(neededMemory)), availMemory.Value, availMemory.Unit) return false, nil } + } return true, nil - } diff --git a/common/pkgs/mq/advisor/task/task.go b/common/pkgs/mq/advisor/task/task.go index 08fec7a..55651dd 100644 --- a/common/pkgs/mq/advisor/task/task.go +++ b/common/pkgs/mq/advisor/task/task.go @@ -5,19 +5,21 @@ import ( myreflect "gitlink.org.cn/cloudream/common/utils/reflect" ) -type TaskStatus interface{} - -// 增加了新类型后需要在这里也同步添加 -type TaskStatusConst interface { - TaskStatus | ScheduleSchemeTaskStatus +type TaskStatus interface { + Noop() } +type TaskStatusBase struct{} + +func (b *TaskStatusBase) Noop() {} + // 增加了新类型后需要在这里也同步添加 var TaskStatusTypeUnion = types.NewTypeUnion[TaskStatus]( myreflect.TypeOf[ScheduleSchemeTaskStatus](), ) type ScheduleSchemeTaskStatus struct { + TaskStatusBase Status string `json:"status"` Error string `json:"error"` IsAdjustment bool `json:"isAdjustment"` @@ -28,8 +30,8 @@ type AdjustedScheme struct { NodeID int64 `json:"nodeID"` } -func NewScheduleSchemeTaskStatus(status string, err string, isAdjustment bool, adjustedScheme AdjustedScheme) TaskStatus { - return ScheduleSchemeTaskStatus{ +func NewScheduleSchemeTaskStatus(status string, err string, isAdjustment bool, adjustedScheme AdjustedScheme) *ScheduleSchemeTaskStatus { + return &ScheduleSchemeTaskStatus{ Status: status, Error: err, IsAdjustment: isAdjustment, From 418e2798b8695ee27cd02b2f7cabad5766fb588e Mon Sep 17 00:00:00 2001 From: songjc <969378911@qq.com> Date: Tue, 12 Sep 2023 17:26:09 +0800 Subject: [PATCH 6/6] =?UTF-8?q?=E6=9B=B4=E6=96=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- advisor/internal/task/schedule_scheme.go | 24 ++++++++++++------------ 1 file changed, 12 insertions(+), 12 deletions(-) diff --git a/advisor/internal/task/schedule_scheme.go b/advisor/internal/task/schedule_scheme.go index 342cf1a..dd2a992 100644 --- a/advisor/internal/task/schedule_scheme.go +++ b/advisor/internal/task/schedule_scheme.go @@ -79,8 +79,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availCPU := resp.Data.(models.CPUResourceData).Available if float64(availCPU.Value) < 1.5*neededCPU { - logger.WithField("jobID", t.Job.JobID). - Infof("insufficient CPU resources: wanted CPU:%f, available CPU: %d%s", 1.5*neededCPU, availCPU.Value, availCPU.Unit) + logger.WithField("JobID", t.Job.JobID). + Infof("insufficient CPU resources, want: %f, available: %d%s", 1.5*neededCPU, availCPU.Value, availCPU.Unit) return false, nil } } @@ -98,8 +98,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availNPU := resp.Data.(models.NPUResourceData).Available if float64(availNPU.Value) < 1.5*neededNPU { - logger.WithField("jobID", t.Job.JobID). - Infof("insufficient NPU resources: wanted NPU:%f, available NPU: %d%s", 1.5*neededNPU, availNPU.Value, availNPU.Unit) + logger.WithField("JobID", t.Job.JobID). + Infof("insufficient NPU resources, want: %f, available: %d%s", 1.5*neededNPU, availNPU.Value, availNPU.Unit) return false, nil } } @@ -117,8 +117,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availGPU := resp.Data.(models.GPUResourceData).Available if float64(availGPU.Value) < 1.5*neededGPU { - logger.WithField("jobID", t.Job.JobID). - Infof("insufficient GPU resources: wanted GPU:%f, available GPU: %d%s", 1.5*neededGPU, availGPU.Value, availGPU.Unit) + logger.WithField("JobID", t.Job.JobID). + Infof("insufficient GPU resources, want: %f, available: %d%s", 1.5*neededGPU, availGPU.Value, availGPU.Unit) return false, nil } } @@ -136,8 +136,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { availMLU := resp.Data.(models.MLUResourceData).Available if float64(availMLU.Value) < 1.5*neededMLU { - logger.WithField("jobID", t.Job.JobID). - Infof("insufficient MLU resources: wanted MLU:%f, available MLU: %d%s", 1.5*neededMLU, availMLU.Value, availMLU.Unit) + logger.WithField("JobID", t.Job.JobID). + Infof("insufficient MLU resources, want: %f, available: %d%s", 1.5*neededMLU, availMLU.Value, availMLU.Unit) return false, nil } } @@ -160,8 +160,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { } if int64(bytesStorage) < int64(1.5*float64(neededStorage)) { - logger.WithField("jobID", t.Job.JobID). - Infof("insufficient storage resources: wanted storage:%s, available storage: %f%s", bytesize.New(1.5*float64(neededStorage)), availStorage.Value, availStorage.Unit) + logger.WithField("JobID", t.Job.JobID). + Infof("insufficient storage resources, want: %s, available: %f%s", bytesize.New(1.5*float64(neededStorage)), availStorage.Value, availStorage.Unit) return false, nil } } @@ -184,8 +184,8 @@ func (t *MakeScheduleScheme) CheckResourceAvailability() (bool, error) { } if int64(bytesMemory) < int64(1.5*float64(neededMemory)) { - logger.WithField("jobID", t.Job.JobID). - Infof("insufficient memory resources: wanted memory:%s, available memory: %f%s", bytesize.New(1.5*float64(neededMemory)), availMemory.Value, availMemory.Unit) + logger.WithField("JobID", t.Job.JobID). + Infof("insufficient memory resources, want: %s, available: %f%s", bytesize.New(1.5*float64(neededMemory)), availMemory.Value, availMemory.Unit) return false, nil }