From 28312e34b5e34da099df40f3d4ee2fd0b781e8a4 Mon Sep 17 00:00:00 2001 From: Sydonian <794346190@qq.com> Date: Fri, 15 Sep 2023 17:21:36 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=9E=E7=8E=B0=E5=87=86=E5=A4=87=E8=B0=83?= =?UTF-8?q?=E6=95=B4=E9=98=B6=E6=AE=B5=E7=9A=84=E6=89=A7=E8=A1=8C=E6=B5=81?= =?UTF-8?q?=E7=A8=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- common/models/job/job.go | 9 +- common/models/job/state.go | 12 +- manager/internal/jobmgr/command.go | 95 -------- manager/internal/jobmgr/event/abort.go | 12 ++ .../event/advisor_report_task_status.go | 35 +++ manager/internal/jobmgr/event/event.go | 49 +++++ .../event/executor_report_task_status.go | 37 ++++ .../internal/jobmgr/event/job_completed.go | 16 ++ .../jobmgr/event/local_file_uploaded.go | 18 ++ manager/internal/jobmgr/jobmgr.go | 28 +-- .../internal/jobmgr/prescheduling_handler.go | 193 +++++++++-------- .../jobmgr/ready_to_adjust_handler.go | 202 ++++++++++++++++++ manager/internal/jobmgr/state_handler.go | 9 +- 13 files changed, 494 insertions(+), 221 deletions(-) delete mode 100644 manager/internal/jobmgr/command.go create mode 100644 manager/internal/jobmgr/event/abort.go create mode 100644 manager/internal/jobmgr/event/advisor_report_task_status.go create mode 100644 manager/internal/jobmgr/event/event.go create mode 100644 manager/internal/jobmgr/event/executor_report_task_status.go create mode 100644 manager/internal/jobmgr/event/job_completed.go create mode 100644 manager/internal/jobmgr/event/local_file_uploaded.go create mode 100644 manager/internal/jobmgr/ready_to_adjust_handler.go diff --git a/common/models/job/job.go b/common/models/job/job.go index 05c3b30..642bfee 100644 --- a/common/models/job/job.go +++ b/common/models/job/job.go @@ -17,8 +17,7 @@ const ( ) type FileScheduleScheme struct { - Action FileScheduleAction `json:"action"` - TargetStorageID int64 `json:"targetStorageID"` + Action FileScheduleAction `json:"action"` } // 任务调度方案 @@ -53,13 +52,13 @@ func NewJobSet(jobSetID string, jobRefs []JobSetJobRef, preScheduleScheme JobSet } } -func (j *JobSet) GetJobIDByLocalJobID(localJobID string) string { +func (j *JobSet) FindRefByLocalJobID(localJobID string) *JobSetJobRef { ref, ok := lo.Find(j.JobRefs, func(item JobSetJobRef) bool { return item.LocalJobID == localJobID }) if !ok { - return "" + return nil } - return ref.JobID + return &ref } // 任务 diff --git a/common/models/job/state.go b/common/models/job/state.go index 6598566..736f505 100644 --- a/common/models/job/state.go +++ b/common/models/job/state.go @@ -21,7 +21,7 @@ var JobStateTypeUnion = types.NewTypeUnion[JobState]( myreflect.TypeOf[StateExecuting](), myreflect.TypeOf[StateResourcing](), myreflect.TypeOf[StateFailed](), - myreflect.TypeOf[StateFinished](), + myreflect.TypeOf[StateSuccess](), ) type FileSchedulingStep string @@ -67,6 +67,10 @@ type StateMakingAdjustScheme struct { JobStateBase } +func NewStateMakingAdjustScheme() *StateMakingAdjustScheme { + return &StateMakingAdjustScheme{} +} + type StateAdjusting struct { JobStateBase } @@ -75,6 +79,10 @@ type StateReadyToExecute struct { JobStateBase } +func NewStateReadyToExecute() *StateReadyToExecute { + return &StateReadyToExecute{} +} + type StateExecuting struct { JobStateBase } @@ -96,6 +104,6 @@ func NewStateFailed(err string, lastState JobState) *StateFailed { } } -type StateFinished struct { +type StateSuccess struct { JobStateBase } diff --git a/manager/internal/jobmgr/command.go b/manager/internal/jobmgr/command.go deleted file mode 100644 index f239028..0000000 --- a/manager/internal/jobmgr/command.go +++ /dev/null @@ -1,95 +0,0 @@ -package jobmgr - -import ( - advtsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/advisor/task" - exectsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/executor/task" -) - -type Command interface{} - -// 终止处理 -type AbortCmd struct { -} - -func NewAbortCmd() *AbortCmd { - return &AbortCmd{} -} - -// 本地文件上传结束 -type LocalFileUploadedCmd struct { - LocalPath string - Error string - PackageID int64 -} - -func NewLocalFileUploadedCmd(localPath string, err string, packageID int64) *LocalFileUploadedCmd { - return &LocalFileUploadedCmd{ - LocalPath: localPath, - Error: err, - PackageID: packageID, - } -} - -// executor上报任务进度 -type ExecutorReportTaskStatusCmd struct { - FullTaskID string - TaskStatus exectsk.TaskStatus -} - -func NewExecutorReportTaskStatusCmd(fullTaskID string, taskStatus exectsk.TaskStatus) *ExecutorReportTaskStatusCmd { - return &ExecutorReportTaskStatusCmd{ - FullTaskID: fullTaskID, - TaskStatus: taskStatus, - } -} - -// advisor上报任务进度 -type AdvisorReportTaskStatusCmd struct { - FullTaskID string - TaskStatus advtsk.TaskStatus -} - -func NewAdvisorReportTaskStatusCmd(fullTaskID string, taskStatus advtsk.TaskStatus) *AdvisorReportTaskStatusCmd { - return &AdvisorReportTaskStatusCmd{ - FullTaskID: fullTaskID, - TaskStatus: taskStatus, - } -} - -func AssertExecutorTaskStatus[T exectsk.TaskStatus](cmd Command, fullTaskID string) (T, bool) { - var ret T - if cmd == nil { - return ret, false - } - - execTaskCmd, ok := cmd.(*ExecutorReportTaskStatusCmd) - if !ok { - return ret, false - } - - if execTaskCmd.FullTaskID != fullTaskID { - return ret, false - } - - status, ok := execTaskCmd.TaskStatus.(T) - return status, ok -} - -func AssertAdvisorTaskStatus[T advtsk.TaskStatus](cmd Command, fullTaskID string) (T, bool) { - var ret T - if cmd == nil { - return ret, false - } - - execTaskCmd, ok := cmd.(*AdvisorReportTaskStatusCmd) - if !ok { - return ret, false - } - - if execTaskCmd.FullTaskID != fullTaskID { - return ret, false - } - - status, ok := execTaskCmd.TaskStatus.(T) - return status, ok -} diff --git a/manager/internal/jobmgr/event/abort.go b/manager/internal/jobmgr/event/abort.go new file mode 100644 index 0000000..55c9bad --- /dev/null +++ b/manager/internal/jobmgr/event/abort.go @@ -0,0 +1,12 @@ +package event + +// 终止处理 +type Abort struct { + JobID string +} + +func NewAbort(jobID string) *Abort { + return &Abort{ + JobID: jobID, + } +} diff --git a/manager/internal/jobmgr/event/advisor_report_task_status.go b/manager/internal/jobmgr/event/advisor_report_task_status.go new file mode 100644 index 0000000..04c9fa4 --- /dev/null +++ b/manager/internal/jobmgr/event/advisor_report_task_status.go @@ -0,0 +1,35 @@ +package event + +import advtsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/advisor/task" + +// advisor上报任务进度 +type AdvisorReportTaskStatus struct { + FullTaskID string + TaskStatus advtsk.TaskStatus +} + +func NewAdvisorReportTaskStatusCmd(fullTaskID string, taskStatus advtsk.TaskStatus) *AdvisorReportTaskStatus { + return &AdvisorReportTaskStatus{ + FullTaskID: fullTaskID, + TaskStatus: taskStatus, + } +} + +func AssertAdvisorTaskStatus[T advtsk.TaskStatus](cmd Event, fullTaskID string) (T, bool) { + var ret T + if cmd == nil { + return ret, false + } + + execTaskCmd, ok := cmd.(*AdvisorReportTaskStatus) + if !ok { + return ret, false + } + + if execTaskCmd.FullTaskID != fullTaskID { + return ret, false + } + + status, ok := execTaskCmd.TaskStatus.(T) + return status, ok +} diff --git a/manager/internal/jobmgr/event/event.go b/manager/internal/jobmgr/event/event.go new file mode 100644 index 0000000..10f904e --- /dev/null +++ b/manager/internal/jobmgr/event/event.go @@ -0,0 +1,49 @@ +package event + +type Event interface{} + +type BroadcastType string + +const ( + BroadcastAll BroadcastType = "All" + BroadcastJobSet BroadcastType = "JobSet" + BroadcastJob BroadcastType = "Job" +) + +type Broadcast struct { + Type BroadcastType + JobSetID string + JobID string +} + +func (b *Broadcast) ToAll() bool { + return b.Type == BroadcastAll +} + +func (b *Broadcast) ToJobSet() bool { + return b.Type == BroadcastJobSet +} + +func (b *Broadcast) ToJob() bool { + return b.Type == BroadcastJob +} + +func ToAll() Broadcast { + return Broadcast{ + Type: BroadcastAll, + } +} + +func ToJobSet(jobSetID string) Broadcast { + return Broadcast{ + Type: BroadcastJobSet, + JobSetID: jobSetID, + } +} + +func ToJob(jobID string) Broadcast { + return Broadcast{ + Type: BroadcastJob, + JobID: jobID, + } +} diff --git a/manager/internal/jobmgr/event/executor_report_task_status.go b/manager/internal/jobmgr/event/executor_report_task_status.go new file mode 100644 index 0000000..d732178 --- /dev/null +++ b/manager/internal/jobmgr/event/executor_report_task_status.go @@ -0,0 +1,37 @@ +package event + +import ( + exectsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/executor/task" +) + +// executor上报任务进度 +type ExecutorReportTaskStatus struct { + FullTaskID string + TaskStatus exectsk.TaskStatus +} + +func NewExecutorReportTaskStatus(fullTaskID string, taskStatus exectsk.TaskStatus) *ExecutorReportTaskStatus { + return &ExecutorReportTaskStatus{ + FullTaskID: fullTaskID, + TaskStatus: taskStatus, + } +} + +func AssertExecutorTaskStatus[T exectsk.TaskStatus](cmd Event, fullTaskID string) (T, bool) { + var ret T + if cmd == nil { + return ret, false + } + + execTaskCmd, ok := cmd.(*ExecutorReportTaskStatus) + if !ok { + return ret, false + } + + if execTaskCmd.FullTaskID != fullTaskID { + return ret, false + } + + status, ok := execTaskCmd.TaskStatus.(T) + return status, ok +} diff --git a/manager/internal/jobmgr/event/job_completed.go b/manager/internal/jobmgr/event/job_completed.go new file mode 100644 index 0000000..aa94584 --- /dev/null +++ b/manager/internal/jobmgr/event/job_completed.go @@ -0,0 +1,16 @@ +package event + +import ( + jobmod "gitlink.org.cn/cloudream/scheduler/common/models/job" +) + +// 任务结束,包括成功或者失败 +type JobCompleted struct { + Job jobmod.Job +} + +func NewJobCompleted(job jobmod.Job) *JobCompleted { + return &JobCompleted{ + Job: job, + } +} diff --git a/manager/internal/jobmgr/event/local_file_uploaded.go b/manager/internal/jobmgr/event/local_file_uploaded.go new file mode 100644 index 0000000..f4f8273 --- /dev/null +++ b/manager/internal/jobmgr/event/local_file_uploaded.go @@ -0,0 +1,18 @@ +package event + +// 本地文件上传结束 +type LocalFileUploaded struct { + JobSetID string + LocalPath string + Error string + PackageID int64 +} + +func NewLocalFileUploaded(jobSetID string, localPath string, err string, packageID int64) *LocalFileUploaded { + return &LocalFileUploaded{ + JobSetID: jobSetID, + LocalPath: localPath, + Error: err, + PackageID: packageID, + } +} diff --git a/manager/internal/jobmgr/jobmgr.go b/manager/internal/jobmgr/jobmgr.go index b40d401..85b2f88 100644 --- a/manager/internal/jobmgr/jobmgr.go +++ b/manager/internal/jobmgr/jobmgr.go @@ -12,6 +12,7 @@ import ( "gitlink.org.cn/cloudream/scheduler/manager/internal/advisormgr" "gitlink.org.cn/cloudream/scheduler/manager/internal/executormgr" "gitlink.org.cn/cloudream/scheduler/manager/internal/imagemgr" + "gitlink.org.cn/cloudream/scheduler/manager/internal/jobmgr/event" ) type mgrJob struct { @@ -135,27 +136,8 @@ func (m *Manager) LocalFileUploaded(jobSetID string, localPath string, err strin m.pubLock.Lock() defer m.pubLock.Unlock() - jobSet, ok := m.jobSets[jobSetID] - if !ok { - return fmt.Errorf("job set not found") - } - - for _, ref := range jobSet.JobRefs { - job, ok := m.jobs[ref.JobID] - if !ok { - continue - } - - norJob, ok := job.Job.(*jobmod.NormalJob) - if !ok { - continue - } - - // 无需仔细判断 - _, ok = norJob.State.(*jobmod.StatePreScheduling) - if ok { - job.Handler.SendCommand(job.Job.GetJobID(), NewLocalFileUploadedCmd(localPath, err, packageID)) - } + for _, h := range m.handlers { + h.OnEvent(event.ToJobSet(jobSetID), event.NewLocalFileUploaded(jobSetID, localPath, err, packageID)) } return nil @@ -170,7 +152,7 @@ func (m *Manager) ExecutorReportTaskStatus(jobID string, fullTaskID string, task return } - job.Handler.SendCommand(jobID, NewExecutorReportTaskStatusCmd(fullTaskID, taskStatus)) + job.Handler.OnEvent(event.ToJob(jobID), event.NewExecutorReportTaskStatus(fullTaskID, taskStatus)) } func (m *Manager) AdvisorReportTaskStatus(jobID string, fullTaskID string, taskStatus advtsk.TaskStatus) { @@ -182,7 +164,7 @@ func (m *Manager) AdvisorReportTaskStatus(jobID string, fullTaskID string, taskS return } - job.Handler.SendCommand(jobID, NewAdvisorReportTaskStatusCmd(fullTaskID, taskStatus)) + job.Handler.OnEvent(event.ToJob(jobID), event.NewAdvisorReportTaskStatusCmd(fullTaskID, taskStatus)) } // 根据job状态选择handler进行处理。需要加锁 diff --git a/manager/internal/jobmgr/prescheduling_handler.go b/manager/internal/jobmgr/prescheduling_handler.go index c336abc..84c779d 100644 --- a/manager/internal/jobmgr/prescheduling_handler.go +++ b/manager/internal/jobmgr/prescheduling_handler.go @@ -11,14 +11,16 @@ import ( jobmod "gitlink.org.cn/cloudream/scheduler/common/models/job" colmq "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/collector" exectsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/executor/task" + "gitlink.org.cn/cloudream/scheduler/manager/internal/jobmgr/event" ) var ErrPreScheduleFailed = fmt.Errorf("pre schedule failed") type preSchedulingJob struct { - Job *jobmod.NormalJob - State *jobmod.StatePreScheduling - SlwNodeInfo *models.SlwNode + job *jobmod.NormalJob + state *jobmod.StatePreScheduling + slwNodeInfo *models.SlwNode + mgr *Manager } type PreSchedulingHandler struct { @@ -39,89 +41,93 @@ func NewPreSchedulingHandler(mgr *Manager) *PreSchedulingHandler { func (h *PreSchedulingHandler) Handle(job jobmod.Job) { h.cmdChan.Send(func() { - allCompleted, err := func() (bool, error) { - norJob, ok := job.(*jobmod.NormalJob) - if !ok { - return true, fmt.Errorf("unknow job: %v", reflect.TypeOf(job)) - } - - preSchState, ok := norJob.GetState().(*jobmod.StatePreScheduling) - if !ok { - return true, fmt.Errorf("unknow state: %v", reflect.TypeOf(norJob.GetState())) - } - - colCli, err := globals.CollectorMQPool.Acquire() - if err != nil { - return true, fmt.Errorf("new collector client: %w", err) - } - defer colCli.Close() - - getNodeResp, err := colCli.GetSlwNodeInfo(colmq.NewGetSlwNodeInfo(preSchState.Scheme.TargetSlwNodeID)) - if err != nil { - return true, fmt.Errorf("getting slw node info: %w", err) - } - - preJob := &preSchedulingJob{ - Job: norJob, - State: preSchState, - SlwNodeInfo: &getNodeResp.SlwNode, - } - h.jobs[job.GetJobID()] = preJob - return h.onCommand(nil, preJob) - }() - - if allCompleted { - if err != nil { - job.SetState(jobmod.NewStateFailed(err.Error(), job.GetState())) - } else { - job.SetState(jobmod.NewStateReadyToAdjust()) - } - - h.mgr.pubLock.Lock() - h.mgr.handleState(job) - h.mgr.pubLock.Unlock() + norJob, ok := job.(*jobmod.NormalJob) + if !ok { + h.changeJobState(job, jobmod.NewStateFailed(fmt.Sprintf("unknow job: %v", reflect.TypeOf(job)), job.GetState())) + return } + + preSchState, ok := norJob.GetState().(*jobmod.StatePreScheduling) + if !ok { + h.changeJobState(job, jobmod.NewStateFailed(fmt.Sprintf("unknow state: %v", reflect.TypeOf(job.GetState())), job.GetState())) + return + } + + colCli, err := globals.CollectorMQPool.Acquire() + if err != nil { + h.changeJobState(job, jobmod.NewStateFailed(fmt.Sprintf("new collector client: %s", err), job.GetState())) + return + } + defer colCli.Close() + + getNodeResp, err := colCli.GetSlwNodeInfo(colmq.NewGetSlwNodeInfo(preSchState.Scheme.TargetSlwNodeID)) + if err != nil { + h.changeJobState(job, jobmod.NewStateFailed(fmt.Sprintf("getting slw node info: %s", err.Error()), job.GetState())) + return + } + + preJob := &preSchedulingJob{ + job: norJob, + state: preSchState, + slwNodeInfo: &getNodeResp.SlwNode, + } + h.jobs[job.GetJobID()] = preJob + + h.onJobEvent(nil, preJob) }) } -func (h *PreSchedulingHandler) onCommand(cmd Command, job *preSchedulingJob) (bool, error) { - err := h.doPackageScheduling(nil, job.Job.JobID, - job.Job.Info.Files.Dataset, &job.Job.Files.Dataset, - &job.State.Scheme.Dataset, &job.State.Dataset, - job.SlwNodeInfo, +func (h *PreSchedulingHandler) onJobEvent(evt event.Event, job *preSchedulingJob) { + err := h.doPackageScheduling(evt, job, + job.job.Info.Files.Dataset, &job.job.Files.Dataset, + &job.state.Scheme.Dataset, &job.state.Dataset, ) if err != nil { - job.State.Dataset.Error = err.Error() - return true, ErrPreScheduleFailed + job.state.Dataset.Error = err.Error() + h.changeJobState(job.job, jobmod.NewStateFailed(err.Error(), job.state)) + return } - err = h.doPackageScheduling(nil, job.Job.JobID, - job.Job.Info.Files.Code, &job.Job.Files.Code, - &job.State.Scheme.Code, &job.State.Code, - job.SlwNodeInfo, + err = h.doPackageScheduling(evt, job, + job.job.Info.Files.Code, &job.job.Files.Code, + &job.state.Scheme.Code, &job.state.Code, ) if err != nil { - job.State.Code.Error = err.Error() - return true, ErrPreScheduleFailed + job.state.Code.Error = err.Error() + h.changeJobState(job.job, jobmod.NewStateFailed(err.Error(), job.state)) + return } - err = h.doImageScheduling(nil, job.Job.JobID, - job.Job.Info.Files.Image, &job.Job.Files.Image, - &job.State.Scheme.Image, &job.State.Image, - job.SlwNodeInfo, + err = h.doImageScheduling(evt, job, + job.job.Info.Files.Image, &job.job.Files.Image, + &job.state.Scheme.Image, &job.state.Image, ) if err != nil { - job.State.Image.Error = err.Error() - return true, ErrPreScheduleFailed + job.state.Image.Error = err.Error() + h.changeJobState(job.job, jobmod.NewStateFailed(err.Error(), job.state)) + return } // 如果三种文件都调度完成,则可以进入下个阶段了 - return job.State.Dataset.Step == jobmod.StepCompleted && - job.State.Code.Step == jobmod.StepCompleted && - job.State.Image.Step == jobmod.StepCompleted, nil + if job.state.Dataset.Step == jobmod.StepCompleted && + job.state.Code.Step == jobmod.StepCompleted && + job.state.Image.Step == jobmod.StepCompleted { + + h.changeJobState(job.job, jobmod.NewStateReadyToAdjust()) + } } -func (h *PreSchedulingHandler) doPackageScheduling(cmd Command, jobID string, fileInfo models.JobFileInfo, file *jobmod.PackageJobFile, scheme *jobmod.FileScheduleScheme, state *jobmod.FileSchedulingState, nodeInfo *models.SlwNode) error { +func (h *PreSchedulingHandler) changeJobState(job jobmod.Job, state jobmod.JobState) { + job.SetState(state) + + delete(h.jobs, job.GetJobID()) + + h.mgr.pubLock.Lock() + h.mgr.handleState(job) + h.mgr.pubLock.Unlock() +} + +func (h *PreSchedulingHandler) doPackageScheduling(cmd event.Event, job *preSchedulingJob, fileInfo models.JobFileInfo, file *jobmod.PackageJobFile, scheme *jobmod.FileScheduleScheme, state *jobmod.FileSchedulingState) error { if state.Step == jobmod.StepBegin { if scheme.Action == jobmod.ActionNo { state.Step = jobmod.StepCompleted @@ -146,7 +152,7 @@ func (h *PreSchedulingHandler) doPackageScheduling(cmd Command, jobID string, fi return nil } - localFileCmd, ok := cmd.(*LocalFileUploadedCmd) + localFileCmd, ok := cmd.(*event.LocalFileUploaded) if !ok { return nil } @@ -158,7 +164,7 @@ func (h *PreSchedulingHandler) doPackageScheduling(cmd Command, jobID string, fi file.PackageID = localFileCmd.PackageID if scheme.Action == jobmod.ActionMove { - fullTaskID, err := h.mgr.execMgr.StartTask(jobID, exectsk.NewCacheMovePackage(0, file.PackageID, nodeInfo.StgNodeID)) + fullTaskID, err := h.mgr.execMgr.StartTask(job.job.JobID, exectsk.NewCacheMovePackage(0, file.PackageID, job.slwNodeInfo.StgNodeID)) if err != nil { return err } @@ -170,7 +176,7 @@ func (h *PreSchedulingHandler) doPackageScheduling(cmd Command, jobID string, fi } if scheme.Action == jobmod.ActionLoad { - fullTaskID, err := h.mgr.execMgr.StartTask(jobID, exectsk.NewStorageLoadPackage(0, file.PackageID, nodeInfo.StorageID)) + fullTaskID, err := h.mgr.execMgr.StartTask(job.job.JobID, exectsk.NewStorageLoadPackage(0, file.PackageID, job.slwNodeInfo.StorageID)) if err != nil { return err } @@ -185,7 +191,7 @@ func (h *PreSchedulingHandler) doPackageScheduling(cmd Command, jobID string, fi if state.Step == jobmod.StepMoving { if scheme.Action == jobmod.ActionLoad { - fullTaskID, err := h.mgr.execMgr.StartTask(jobID, exectsk.NewStorageLoadPackage(0, file.PackageID, nodeInfo.StorageID)) + fullTaskID, err := h.mgr.execMgr.StartTask(job.job.JobID, exectsk.NewStorageLoadPackage(0, file.PackageID, job.slwNodeInfo.StorageID)) if err != nil { return err } @@ -206,7 +212,7 @@ func (h *PreSchedulingHandler) doPackageScheduling(cmd Command, jobID string, fi return nil } -func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, fileInfo models.JobFileInfo, file *jobmod.ImageJobFile, scheme *jobmod.FileScheduleScheme, state *jobmod.FileSchedulingState, nodeInfo *models.SlwNode) error { +func (h *PreSchedulingHandler) doImageScheduling(evt event.Event, job *preSchedulingJob, fileInfo models.JobFileInfo, file *jobmod.ImageJobFile, scheme *jobmod.FileScheduleScheme, state *jobmod.FileSchedulingState) error { if state.Step == jobmod.StepBegin { if scheme.Action == jobmod.ActionNo { state.Step = jobmod.StepCompleted @@ -235,11 +241,11 @@ func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, file } if state.Step == jobmod.StepUploading { - if cmd == nil { + if evt == nil { return nil } - localFileCmd, ok := cmd.(*LocalFileUploadedCmd) + localFileCmd, ok := evt.(*event.LocalFileUploaded) if !ok { return nil } @@ -252,7 +258,7 @@ func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, file // 要导入镜像,则需要先将镜像移动到指点节点的缓存中 if scheme.Action == jobmod.ActionImportImage { - fullTaskID, err := h.mgr.execMgr.StartTask(jobID, exectsk.NewCacheMovePackage(0, file.PackageID, nodeInfo.StgNodeID)) + fullTaskID, err := h.mgr.execMgr.StartTask(job.job.JobID, exectsk.NewCacheMovePackage(0, file.PackageID, job.slwNodeInfo.StgNodeID)) if err != nil { return err } @@ -267,7 +273,7 @@ func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, file } if state.Step == jobmod.StepMoving { - cacheMoveRet, ok := AssertExecutorTaskStatus[*exectsk.CacheMovePackageStatus](cmd, state.FullTaskID) + cacheMoveRet, ok := event.AssertExecutorTaskStatus[*exectsk.CacheMovePackageStatus](evt, state.FullTaskID) if !ok { return nil } @@ -276,7 +282,7 @@ func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, file return fmt.Errorf("there must be only 1 object in the package that will be imported") } - fullTaskID, err := h.mgr.execMgr.StartTask(jobID, exectsk.NewUploadImage(nodeInfo.ID, stgsdk.MakeIPFSFilePath(cacheMoveRet.CacheInfos[0].FileHash))) + fullTaskID, err := h.mgr.execMgr.StartTask(job.job.JobID, exectsk.NewUploadImage(job.slwNodeInfo.ID, stgsdk.MakeIPFSFilePath(cacheMoveRet.CacheInfos[0].FileHash))) if err != nil { return err } @@ -287,7 +293,7 @@ func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, file } if state.Step == jobmod.StepImageImporting { - uploadImageRet, ok := AssertExecutorTaskStatus[*exectsk.UploadImageStatus](cmd, state.FullTaskID) + uploadImageRet, ok := event.AssertExecutorTaskStatus[*exectsk.UploadImageStatus](evt, state.FullTaskID) if !ok { return nil } @@ -309,24 +315,25 @@ func (h *PreSchedulingHandler) doImageScheduling(cmd Command, jobID string, file return nil } -func (h *PreSchedulingHandler) SendCommand(jobID string, cmd Command) { - job, ok := h.jobs[jobID] - if !ok { - return - } - +func (h *PreSchedulingHandler) OnEvent(broadcast event.Broadcast, evt event.Event) { h.cmdChan.Send(func() { - allCompleted, err := h.onCommand(cmd, job) - if allCompleted { - if err != nil { - job.Job.SetState(jobmod.NewStateFailed(err.Error(), job.State)) - } else { - job.Job.SetState(jobmod.NewStateReadyToAdjust()) + if broadcast.ToAll() { + for _, job := range h.jobs { + h.onJobEvent(evt, job) } - h.mgr.pubLock.Lock() - h.mgr.handleState(job.Job) - h.mgr.pubLock.Unlock() + } else if broadcast.ToJobSet() { + for _, job := range h.jobs { + if job.job.JobSetID != broadcast.JobSetID { + continue + } + + h.onJobEvent(evt, job) + } + } else if broadcast.ToJob() { + if job, ok := h.jobs[broadcast.JobID]; ok { + h.onJobEvent(evt, job) + } } }) } diff --git a/manager/internal/jobmgr/ready_to_adjust_handler.go b/manager/internal/jobmgr/ready_to_adjust_handler.go new file mode 100644 index 0000000..cbfb386 --- /dev/null +++ b/manager/internal/jobmgr/ready_to_adjust_handler.go @@ -0,0 +1,202 @@ +package jobmgr + +import ( + "fmt" + "reflect" + + "gitlink.org.cn/cloudream/common/models" + "gitlink.org.cn/cloudream/common/pkgs/actor" + jobmod "gitlink.org.cn/cloudream/scheduler/common/models/job" + "gitlink.org.cn/cloudream/scheduler/manager/internal/jobmgr/event" +) + +type readyToAdjustJob struct { + job jobmod.Job + state *jobmod.StateReadyToAdjust +} + +type ReadyToAdjustHandler struct { + mgr *Manager + + jobs map[string]*readyToAdjustJob + + cmdChan actor.CommandChannel +} + +func NewReadyToAdjustHandler(mgr *Manager) *ReadyToAdjustHandler { + return &ReadyToAdjustHandler{ + mgr: mgr, + jobs: make(map[string]*readyToAdjustJob), + cmdChan: *actor.NewCommandChannel(), + } +} + +func (h *ReadyToAdjustHandler) Handle(job jobmod.Job) { + h.cmdChan.Send(func() { + state, ok := job.GetState().(*jobmod.StateReadyToAdjust) + if !ok { + h.changeJobState(job, jobmod.NewStateFailed(fmt.Sprintf("unknow state: %v", reflect.TypeOf(job.GetState())), job.GetState())) + return + } + + rjob := &readyToAdjustJob{ + job: job, + state: state, + } + h.jobs[job.GetJobID()] = rjob + + h.onJobEvent(nil, rjob) + }) +} + +func (h *ReadyToAdjustHandler) onJobEvent(evt event.Event, job *readyToAdjustJob) { + if norJob, ok := job.job.(*jobmod.NormalJob); ok { + h.onNormalJobEvent(evt, job, norJob) + } else if resJob, ok := job.job.(*jobmod.ResourceJob); ok { + h.onResourceJobEvent(evt, job, resJob) + } +} + +func (h *ReadyToAdjustHandler) onNormalJobEvent(evt event.Event, job *readyToAdjustJob, norJob *jobmod.NormalJob) { + h.mgr.pubLock.Lock() + jobSet, ok := h.mgr.jobSets[job.job.GetJobSetID()] + h.mgr.pubLock.Unlock() + if !ok { + h.changeJobState(job.job, jobmod.NewStateFailed("job set not found", job.state)) + return + } + + needWait := true + + // 无论发生什么事件,都检查一下前置任务的状态 + if resFile, ok := norJob.Info.Files.Dataset.(*models.ResourceJobFileInfo); ok { + ref := jobSet.FindRefByLocalJobID(resFile.ResourceLocalJobID) + if ref == nil { + h.changeJobState(job.job, jobmod.NewStateFailed(fmt.Sprintf("job %s not found in job set", job.job.GetJobID()), job.state)) + return + } + + h.mgr.pubLock.Lock() + waitJob := h.mgr.jobs[ref.JobID] + h.mgr.pubLock.Unlock() + if waitJob == nil { + h.changeJobState(job.job, jobmod.NewStateFailed(fmt.Sprintf("job %s not found in job set", job.job.GetJobID()), job.state)) + return + } + + if _, ok = waitJob.Job.GetState().(*jobmod.StateSuccess); ok { + waitResJob, ok := waitJob.Job.(*jobmod.ResourceJob) + if !ok { + h.changeJobState(job.job, jobmod.NewStateFailed( + fmt.Sprintf("job %s is not a resource job(%v)", job.job.GetJobID(), reflect.TypeOf(waitJob)), + job.state, + )) + return + } + + norJob.Files.Dataset.PackageID = waitResJob.ResourcePackageID + needWait = needWait || false + } else if _, ok = waitJob.Job.GetState().(*jobmod.StateFailed); ok { + h.changeJobState(job.job, jobmod.NewStateFailed( + fmt.Sprintf("job %s is failed", job.job.GetJobID(), reflect.TypeOf(waitJob)), + job.state, + )) + return + } + } + + if !needWait { + h.changeJobState(job.job, jobmod.NewStateMakingAdjustScheme()) + } +} + +func (h *ReadyToAdjustHandler) onResourceJobEvent(evt event.Event, job *readyToAdjustJob, resJob *jobmod.ResourceJob) { + h.mgr.pubLock.Lock() + jobSet, ok := h.mgr.jobSets[job.job.GetJobSetID()] + h.mgr.pubLock.Unlock() + if !ok { + h.changeJobState(job.job, jobmod.NewStateFailed("job set not found", job.state)) + return + } + + needWait := true + + ref := jobSet.FindRefByLocalJobID(resJob.Info.TargetLocalJobID) + if ref == nil { + h.changeJobState(job.job, jobmod.NewStateFailed(fmt.Sprintf("job %s not found in job set", job.job.GetJobID()), job.state)) + return + } + + h.mgr.pubLock.Lock() + waitJob := h.mgr.jobs[ref.JobID] + h.mgr.pubLock.Unlock() + if waitJob == nil { + h.changeJobState(job.job, jobmod.NewStateFailed(fmt.Sprintf("job %s not found in job set", job.job.GetJobID()), job.state)) + return + } + + // 无论发生什么事件,都检查一下前置任务的状态 + if _, ok = waitJob.Job.GetState().(*jobmod.StateSuccess); ok { + needWait = needWait || false + + } else if _, ok = waitJob.Job.GetState().(*jobmod.StateFailed); ok { + h.changeJobState(job.job, jobmod.NewStateFailed( + fmt.Sprintf("job %s is failed", job.job.GetJobID(), reflect.TypeOf(waitJob)), + job.state, + )) + return + } + + if !needWait { + h.changeJobState(job.job, jobmod.NewStateReadyToExecute()) + } +} + +func (h *ReadyToAdjustHandler) changeJobState(job jobmod.Job, state jobmod.JobState) { + job.SetState(state) + + delete(h.jobs, job.GetJobID()) + + h.mgr.pubLock.Lock() + h.mgr.handleState(job) + h.mgr.pubLock.Unlock() +} + +func (h *ReadyToAdjustHandler) OnEvent(broadcast event.Broadcast, evt event.Event) { + h.cmdChan.Send(func() { + if broadcast.ToAll() { + for _, job := range h.jobs { + h.onJobEvent(evt, job) + } + + } else if broadcast.ToJobSet() { + for _, job := range h.jobs { + if job.job.GetJobSetID() != broadcast.JobSetID { + continue + } + + h.onJobEvent(evt, job) + } + } else if broadcast.ToJob() { + if job, ok := h.jobs[broadcast.JobID]; ok { + h.onJobEvent(evt, job) + } + } + }) +} + +func (h *ReadyToAdjustHandler) Serve() { + cmdChan := h.cmdChan.BeginChanReceive() + defer h.cmdChan.CloseChanReceive() + + for { + select { + case cmd := <-cmdChan: + cmd() + } + } +} + +func (h *ReadyToAdjustHandler) Stop() { + // TODO 支持STOP +} diff --git a/manager/internal/jobmgr/state_handler.go b/manager/internal/jobmgr/state_handler.go index ff1be37..9209fa7 100644 --- a/manager/internal/jobmgr/state_handler.go +++ b/manager/internal/jobmgr/state_handler.go @@ -1,12 +1,15 @@ package jobmgr -import jobmod "gitlink.org.cn/cloudream/scheduler/common/models/job" +import ( + jobmod "gitlink.org.cn/cloudream/scheduler/common/models/job" + "gitlink.org.cn/cloudream/scheduler/manager/internal/jobmgr/event" +) type StateHandler interface { // 处理Job。在此期间全局锁已锁定 Handle(job jobmod.Job) - // 向指定的Job发送一个命令。在此期间全局锁已锁定 - SendCommand(jobID string, cmd Command) + // 外部发生了一个事件 + OnEvent(broadcast event.Broadcast, evt event.Event) // 运行Handler Serve() // 停止此Handler