实现准备调整阶段的执行流程

This commit is contained in:
Sydonian 2023-09-15 17:21:36 +08:00
parent 945be88b23
commit 28312e34b5
13 changed files with 494 additions and 221 deletions

View File

@ -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
}
// 任务

View File

@ -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
}

View File

@ -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
}

View File

@ -0,0 +1,12 @@
package event
// 终止处理
type Abort struct {
JobID string
}
func NewAbort(jobID string) *Abort {
return &Abort{
JobID: jobID,
}
}

View File

@ -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
}

View File

@ -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,
}
}

View File

@ -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
}

View File

@ -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,
}
}

View File

@ -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,
}
}

View File

@ -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进行处理。需要加锁

View File

@ -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)
}
}
})
}

View File

@ -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
}

View File

@ -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