diff --git a/advisor/internal/reporter/reporter.go b/advisor/internal/reporter/reporter.go index bcd6ca0..1d71fe7 100644 --- a/advisor/internal/reporter/reporter.go +++ b/advisor/internal/reporter/reporter.go @@ -15,7 +15,7 @@ import ( type Reporter struct { advisorID schmod.AdvisorID reportInterval time.Duration - taskStatus map[string]advtsk.AdvTaskStatus + taskStatus map[string]advtsk.TaskStatus taskStatusLock sync.Mutex reportNow chan bool } @@ -24,12 +24,12 @@ func NewReporter(advisorID schmod.AdvisorID, reportInterval time.Duration) *Repo return &Reporter{ advisorID: advisorID, reportInterval: reportInterval, - taskStatus: make(map[string]advtsk.AdvTaskStatus), + taskStatus: make(map[string]advtsk.TaskStatus), reportNow: make(chan bool), } } -func (r *Reporter) Report(taskID string, taskStatus advtsk.AdvTaskStatus) { +func (r *Reporter) Report(taskID string, taskStatus advtsk.TaskStatus) { r.taskStatusLock.Lock() defer r.taskStatusLock.Unlock() @@ -65,7 +65,7 @@ func (r *Reporter) Serve() error { for taskID, status := range r.taskStatus { taskStatus = append(taskStatus, mgrmq.NewAdvisorTaskStatus(taskID, status)) } - r.taskStatus = make(map[string]advtsk.AdvTaskStatus) + r.taskStatus = make(map[string]advtsk.TaskStatus) r.taskStatusLock.Unlock() _, err := magCli.ReportAdvisorTaskStatus(mgrmq.NewReportAdvisorTaskStatus(r.advisorID, taskStatus)) diff --git a/advisor/internal/task/task.go b/advisor/internal/task/task.go index 00e99df..03feb0a 100644 --- a/advisor/internal/task/task.go +++ b/advisor/internal/task/task.go @@ -39,7 +39,7 @@ func NewManager(reporter *reporter.Reporter, scheduleSvc *scheduler.Service) Man } } -func (m *Manager) StartByInfo(info advtsk.AdvTaskInfo) (*Task, error) { +func (m *Manager) StartByInfo(info advtsk.TaskInfo) (*Task, error) { infoType := myreflect.TypeOfValue(info) ctor, ok := taskFromInfoCtors[infoType] @@ -50,10 +50,10 @@ func (m *Manager) StartByInfo(info advtsk.AdvTaskInfo) (*Task, error) { return m.StartNew(ctor(info)), nil } -var taskFromInfoCtors map[reflect.Type]func(advtsk.AdvTaskInfo) TaskBody = make(map[reflect.Type]func(advtsk.AdvTaskInfo) task.TaskBody[TaskContext]) +var taskFromInfoCtors map[reflect.Type]func(advtsk.TaskInfo) TaskBody = make(map[reflect.Type]func(advtsk.TaskInfo) task.TaskBody[TaskContext]) -func Register[TInfo advtsk.AdvTaskInfo, TTaskBody TaskBody](ctor func(info TInfo) TTaskBody) { - taskFromInfoCtors[myreflect.TypeOf[TInfo]()] = func(info advtsk.AdvTaskInfo) TaskBody { +func Register[TInfo advtsk.TaskInfo, TTaskBody TaskBody](ctor func(info TInfo) TTaskBody) { + taskFromInfoCtors[myreflect.TypeOf[TInfo]()] = func(info advtsk.TaskInfo) TaskBody { return ctor(info.(TInfo)) } } diff --git a/client/internal/http/jobset.go b/client/internal/http/jobset.go index 0e4c104..65f5c88 100644 --- a/client/internal/http/jobset.go +++ b/client/internal/http/jobset.go @@ -8,6 +8,7 @@ import ( "gitlink.org.cn/cloudream/common/consts/errorcode" "gitlink.org.cn/cloudream/common/pkgs/logger" schsdk "gitlink.org.cn/cloudream/common/sdks/scheduler" + "gitlink.org.cn/cloudream/common/utils/serder" ) type JobSetService struct { @@ -35,14 +36,14 @@ func (s *JobSetService) Submit(ctx *gin.Context) { return } - jobSetInfo, err := schsdk.JobSetInfoFromJSON(bodyData) + jobSetInfo, err := serder.JSONToObjectEx[schsdk.JobSetInfo](bodyData) if err != nil { log.Warnf("parsing request body: %s", err.Error()) ctx.JSON(http.StatusOK, Failed(errorcode.OperationFailed, "parse request body failed")) return } - jobsetID, uploadScheme, err := s.svc.JobSetSvc().Submit(*jobSetInfo) + jobsetID, uploadScheme, err := s.svc.JobSetSvc().Submit(jobSetInfo) if err != nil { log.Warnf("submitting jobset: %s", err.Error()) ctx.JSON(http.StatusOK, Failed(errorcode.OperationFailed, "submit jobset failed")) diff --git a/common/models/job/job.go b/common/models/job/job.go index 46286e0..12bb888 100644 --- a/common/models/job/job.go +++ b/common/models/job/job.go @@ -2,9 +2,9 @@ package jobmod import ( "github.com/samber/lo" - "gitlink.org.cn/cloudream/common/pkgs/mq" "gitlink.org.cn/cloudream/common/pkgs/types" schsdk "gitlink.org.cn/cloudream/common/sdks/scheduler" + "gitlink.org.cn/cloudream/common/utils/serder" ) type FileScheduleAction string @@ -75,7 +75,7 @@ var JobTypeUnion = types.NewTypeUnion[Job]( (*NormalJob)(nil), (*ResourceJob)(nil), ) -var _ = mq.RegisterUnionType(JobTypeUnion) +var _ = serder.UseTypeUnionExternallyTagged(&JobTypeUnion) // TODO var _ = serder.RegisterNewTaggedTypeUnion(JobTypeUnion, "Type", "type") diff --git a/common/models/job/state.go b/common/models/job/state.go index 14d1761..91f9365 100644 --- a/common/models/job/state.go +++ b/common/models/job/state.go @@ -1,8 +1,8 @@ package jobmod import ( - "gitlink.org.cn/cloudream/common/pkgs/mq" "gitlink.org.cn/cloudream/common/pkgs/types" + "gitlink.org.cn/cloudream/common/utils/serder" ) type JobState interface { @@ -20,7 +20,7 @@ var JobStateTypeUnion = types.NewTypeUnion[JobState]( (*StateFailed)(nil), (*StateSuccess)(nil), ) -var _ = mq.RegisterUnionType(JobStateTypeUnion) +var _ = serder.UseTypeUnionExternallyTagged(&JobStateTypeUnion) // TODO var _ = serder.RegisterNewTaggedTypeUnion(JobStateTypeUnion, "Type", "type") diff --git a/common/pkgs/mq/advisor/task.go b/common/pkgs/mq/advisor/task.go index 54bea85..e04ea93 100644 --- a/common/pkgs/mq/advisor/task.go +++ b/common/pkgs/mq/advisor/task.go @@ -15,7 +15,7 @@ var _ = Register(Service.StartTask) type StartTask struct { mq.MessageBodyBase - Info advtsk.AdvTaskInfo `json:"info"` + Info advtsk.TaskInfo `json:"info"` } type StartTaskResp struct { mq.MessageBodyBase @@ -23,7 +23,7 @@ type StartTaskResp struct { TaskID string `json:"taskID"` } -func NewStartTask(info advtsk.AdvTaskInfo) *StartTask { +func NewStartTask(info advtsk.TaskInfo) *StartTask { return &StartTask{ Info: info, } diff --git a/common/pkgs/mq/advisor/task/make_adjust_scheme.go b/common/pkgs/mq/advisor/task/make_adjust_scheme.go index aecea9d..1d073d4 100644 --- a/common/pkgs/mq/advisor/task/make_adjust_scheme.go +++ b/common/pkgs/mq/advisor/task/make_adjust_scheme.go @@ -4,8 +4,6 @@ import ( jobmod "gitlink.org.cn/cloudream/scheduler/common/models/job" ) -var _ = Register[*MakeAdjustScheme, *MakeAdjustSchemeStatus]() - type MakeAdjustScheme struct { TaskInfoBase Job jobmod.NormalJob `json:"job"` @@ -29,3 +27,7 @@ func NewMakeAdjustSchemeStatus(err string, scheme jobmod.JobScheduleScheme) *Mak Scheme: scheme, } } + +func init() { + Register[*MakeAdjustScheme, *MakeAdjustSchemeStatus]() +} diff --git a/common/pkgs/mq/advisor/task/task.go b/common/pkgs/mq/advisor/task/task.go index 6e28f92..0e1b289 100644 --- a/common/pkgs/mq/advisor/task/task.go +++ b/common/pkgs/mq/advisor/task/task.go @@ -1,48 +1,40 @@ package task import ( - "gitlink.org.cn/cloudream/common/pkgs/mq" "gitlink.org.cn/cloudream/common/pkgs/types" myreflect "gitlink.org.cn/cloudream/common/utils/reflect" + "gitlink.org.cn/cloudream/common/utils/serder" ) // 任务。 -// 由于json-iter库的缺陷,这个类型名必须加一点前缀,否则会和executor中的重名导致代码异常 -type AdvTaskInfo interface { +type TaskInfo interface { Noop() } // 增加了新类型后需要在这里也同步添加 -var TaskInfoTypeUnion = types.NewTypeUnion[AdvTaskInfo]() +var TaskInfoTypeUnion = serder.UseTypeUnionExternallyTagged(types.Ref(types.NewTypeUnion[TaskInfo]())) type TaskInfoBase struct{} func (s *TaskInfoBase) Noop() {} // 任务上报的状态 -// 由于json-iter库的缺陷,这个类型名必须加一点前缀,否则会和executor中的重名导致代码异常 -type AdvTaskStatus interface { +type TaskStatus interface { Noop() } // 增加了新类型后需要在这里也同步添加 -var TaskStatusTypeUnion = types.NewTypeUnion[AdvTaskStatus]() +var TaskStatusTypeUnion = serder.UseTypeUnionExternallyTagged(types.Ref(types.NewTypeUnion[TaskStatus]())) type TaskStatusBase struct{} func (s *TaskStatusBase) Noop() {} -// 注:此函数必须以var _ = Register[xxx, xxx]()的形式被调用,这样才能保证init中RegisterUnionType时 -// TypeUnion不是空的。(因为包级变量初始化比init函数调用先进行) -func Register[TTaskInfo AdvTaskInfo, TTaskStatus AdvTaskStatus]() any { +// 只能在init函数中调用,因为包级变量初始化会比init函数调用先进行 +func Register[TTaskInfo TaskInfo, TTaskStatus TaskStatus]() any { TaskInfoTypeUnion.Add(myreflect.TypeOf[TTaskInfo]()) TaskStatusTypeUnion.Add(myreflect.TypeOf[TTaskStatus]()) return nil } - -func init() { - mq.RegisterUnionType(TaskInfoTypeUnion) - mq.RegisterUnionType(TaskStatusTypeUnion) -} diff --git a/common/pkgs/mq/executor/task.go b/common/pkgs/mq/executor/task.go index 6ad0d72..4d5da43 100644 --- a/common/pkgs/mq/executor/task.go +++ b/common/pkgs/mq/executor/task.go @@ -16,7 +16,7 @@ var _ = Register(Service.StartTask) type StartTask struct { mq.MessageBodyBase - Info exectsk.ExeTaskInfo `json:"info"` + Info exectsk.TaskInfo `json:"info"` } type StartTaskResp struct { mq.MessageBodyBase @@ -24,7 +24,7 @@ type StartTaskResp struct { TaskID string `json:"taskID"` } -func NewStartTask(info exectsk.ExeTaskInfo) *StartTask { +func NewStartTask(info exectsk.TaskInfo) *StartTask { return &StartTask{ Info: info, } diff --git a/common/pkgs/mq/executor/task/cache_move_package.go b/common/pkgs/mq/executor/task/cache_move_package.go index 7c2ef9c..79b4117 100644 --- a/common/pkgs/mq/executor/task/cache_move_package.go +++ b/common/pkgs/mq/executor/task/cache_move_package.go @@ -2,8 +2,6 @@ package task import cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" -var _ = Register[*CacheMovePackage, *CacheMovePackageStatus]() - type CacheMovePackage struct { TaskInfoBase UserID int64 `json:"userID"` @@ -29,3 +27,7 @@ func NewCacheMovePackageStatus(err string, cacheInfos []cdssdk.ObjectCacheInfo) CacheInfos: cacheInfos, } } + +func init() { + Register[*CacheMovePackage, *CacheMovePackageStatus]() +} diff --git a/common/pkgs/mq/executor/task/storage_create_package.go b/common/pkgs/mq/executor/task/storage_create_package.go index f294527..f88ad24 100644 --- a/common/pkgs/mq/executor/task/storage_create_package.go +++ b/common/pkgs/mq/executor/task/storage_create_package.go @@ -2,8 +2,6 @@ package task import cdssdk "gitlink.org.cn/cloudream/common/sdks/storage" -var _ = Register[*StorageCreatePackage, *StorageCreatePackageStatus]() - type StorageCreatePackage struct { TaskInfoBase UserID int64 `json:"userID"` @@ -37,3 +35,7 @@ func NewStorageCreatePackageStatus(status string, err string, packageID int64) * PackageID: packageID, } } + +func init() { + Register[*StorageCreatePackage, *StorageCreatePackageStatus]() +} diff --git a/common/pkgs/mq/executor/task/storage_load_package.go b/common/pkgs/mq/executor/task/storage_load_package.go index 900ee4a..9056895 100644 --- a/common/pkgs/mq/executor/task/storage_load_package.go +++ b/common/pkgs/mq/executor/task/storage_load_package.go @@ -1,7 +1,5 @@ package task -var _ = Register[*StorageLoadPackage, *StorageLoadPackageStatus]() - type StorageLoadPackage struct { TaskInfoBase UserID int64 `json:"userID"` @@ -27,3 +25,7 @@ func NewStorageLoadPackageStatus(err string, fullPath string) *StorageLoadPackag FullPath: fullPath, } } + +func init() { + Register[*StorageLoadPackage, *StorageLoadPackageStatus]() +} diff --git a/common/pkgs/mq/executor/task/submit_task.go b/common/pkgs/mq/executor/task/submit_task.go index 73786b2..f08a4fc 100644 --- a/common/pkgs/mq/executor/task/submit_task.go +++ b/common/pkgs/mq/executor/task/submit_task.go @@ -5,8 +5,6 @@ import ( schsdk "gitlink.org.cn/cloudream/common/sdks/scheduler" ) -var _ = Register[*SubmitTask, *SubmitTaskStatus]() - type SubmitTask struct { TaskInfoBase PCMParticipantID pcmsdk.ParticipantID `json:"pcmParticipantID"` @@ -37,3 +35,7 @@ func NewSubmitTaskStatus(status pcmsdk.TaskStatus, err string) *SubmitTaskStatus Error: err, } } + +func init() { + Register[*SubmitTask, *SubmitTaskStatus]() +} diff --git a/common/pkgs/mq/executor/task/task.go b/common/pkgs/mq/executor/task/task.go index 00e6f35..73045b5 100644 --- a/common/pkgs/mq/executor/task/task.go +++ b/common/pkgs/mq/executor/task/task.go @@ -1,48 +1,40 @@ package task import ( - "gitlink.org.cn/cloudream/common/pkgs/mq" "gitlink.org.cn/cloudream/common/pkgs/types" myreflect "gitlink.org.cn/cloudream/common/utils/reflect" + "gitlink.org.cn/cloudream/common/utils/serder" ) // 任务 -// 由于json-iter库的缺陷,这个类型名必须加一点前缀,否则会和advisor中的重名导致代码异常 -type ExeTaskInfo interface { +type TaskInfo interface { Noop() } // 增加了新类型后需要在这里也同步添加 -var TaskInfoTypeUnion = types.NewTypeUnion[ExeTaskInfo]() +var TaskInfoTypeUnion = serder.UseTypeUnionExternallyTagged(types.Ref(types.NewTypeUnion[TaskInfo]())) type TaskInfoBase struct{} func (s *TaskInfoBase) Noop() {} // 任务上报的状态 -// 由于json-iter库的缺陷,这个类型名必须加一点前缀,否则会和advisor中的重名导致代码异常 -type ExeTaskStatus interface { +type TaskStatus interface { Noop() } // 增加了新类型后需要在这里也同步添加 -var TaskStatusTypeUnion = types.NewTypeUnion[ExeTaskStatus]() +var TaskStatusTypeUnion = serder.UseTypeUnionExternallyTagged(types.Ref(types.NewTypeUnion[TaskStatus]())) type TaskStatusBase struct{} func (s *TaskStatusBase) Noop() {} -// 注:此函数必须以var _ = Register[xxx, xxx]()的形式被调用,这样才能保证init中RegisterUnionType时 -// TypeUnion不是空的。(因为包级变量初始化比init函数调用先进行) -func Register[TTaskInfo ExeTaskInfo, TTaskStatus ExeTaskStatus]() any { +// 只能在init函数中调用,因为包级变量初始化会比init函数调用先进行 +func Register[TTaskInfo TaskInfo, TTaskStatus TaskStatus]() any { TaskInfoTypeUnion.Add(myreflect.TypeOf[TTaskInfo]()) TaskStatusTypeUnion.Add(myreflect.TypeOf[TTaskStatus]()) return nil } - -func init() { - mq.RegisterUnionType(TaskInfoTypeUnion) - mq.RegisterUnionType(TaskStatusTypeUnion) -} diff --git a/common/pkgs/mq/executor/task/upload_image.go b/common/pkgs/mq/executor/task/upload_image.go index af5d889..1231e1e 100644 --- a/common/pkgs/mq/executor/task/upload_image.go +++ b/common/pkgs/mq/executor/task/upload_image.go @@ -4,8 +4,6 @@ import ( pcmsdk "gitlink.org.cn/cloudream/common/sdks/pcm" ) -var _ = Register[*UploadImage, *UploadImageStatus]() - type UploadImage struct { TaskInfoBase PCMParticipantID pcmsdk.ParticipantID `json:"pcmParticipantID"` @@ -33,3 +31,7 @@ func NewUploadImageStatus(status string, err string, pcmImageID pcmsdk.ImageID, Name: name, } } + +func init() { + Register[*UploadImage, *UploadImageStatus]() +} diff --git a/common/pkgs/mq/manager/advisor.go b/common/pkgs/mq/manager/advisor.go index 17843fb..b9860dd 100644 --- a/common/pkgs/mq/manager/advisor.go +++ b/common/pkgs/mq/manager/advisor.go @@ -25,7 +25,7 @@ type ReportAdvisorTaskStatusResp struct { } type AdvisorTaskStatus struct { TaskID string - Status advtsk.AdvTaskStatus + Status advtsk.TaskStatus } func NewReportAdvisorTaskStatus(advisorID schmod.AdvisorID, taskStatus []AdvisorTaskStatus) *ReportAdvisorTaskStatus { @@ -37,7 +37,7 @@ func NewReportAdvisorTaskStatus(advisorID schmod.AdvisorID, taskStatus []Advisor func NewReportAdvisorTaskStatusResp() *ReportAdvisorTaskStatusResp { return &ReportAdvisorTaskStatusResp{} } -func NewAdvisorTaskStatus(taskID string, status exectsk.ExeTaskStatus) AdvisorTaskStatus { +func NewAdvisorTaskStatus(taskID string, status exectsk.TaskStatus) AdvisorTaskStatus { return AdvisorTaskStatus{ TaskID: taskID, Status: status, diff --git a/common/pkgs/mq/manager/executor.go b/common/pkgs/mq/manager/executor.go index a8052d8..39f7d7b 100644 --- a/common/pkgs/mq/manager/executor.go +++ b/common/pkgs/mq/manager/executor.go @@ -24,7 +24,7 @@ type ReportExecutorTaskStatusResp struct { } type ExecutorTaskStatus struct { TaskID string - Status exectsk.ExeTaskStatus + Status exectsk.TaskStatus } func NewReportExecutorTaskStatus(executorID schmod.ExecutorID, taskStatus []ExecutorTaskStatus) *ReportExecutorTaskStatus { @@ -36,7 +36,7 @@ func NewReportExecutorTaskStatus(executorID schmod.ExecutorID, taskStatus []Exec func NewReportExecutorTaskStatusResp() *ReportExecutorTaskStatusResp { return &ReportExecutorTaskStatusResp{} } -func NewExecutorTaskStatus(taskID string, status exectsk.ExeTaskStatus) ExecutorTaskStatus { +func NewExecutorTaskStatus(taskID string, status exectsk.TaskStatus) ExecutorTaskStatus { return ExecutorTaskStatus{ TaskID: taskID, Status: status, diff --git a/executor/internal/reporter/reporter.go b/executor/internal/reporter/reporter.go index 12bf6b1..e0da653 100644 --- a/executor/internal/reporter/reporter.go +++ b/executor/internal/reporter/reporter.go @@ -15,7 +15,7 @@ import ( type Reporter struct { executorID schmod.ExecutorID reportInterval time.Duration - taskStatus map[string]exectsk.ExeTaskStatus + taskStatus map[string]exectsk.TaskStatus taskStatusLock sync.Mutex reportNow chan bool } @@ -24,12 +24,12 @@ func NewReporter(executorID schmod.ExecutorID, reportInterval time.Duration) Rep return Reporter{ executorID: executorID, reportInterval: reportInterval, - taskStatus: make(map[string]exectsk.ExeTaskStatus), + taskStatus: make(map[string]exectsk.TaskStatus), reportNow: make(chan bool), } } -func (r *Reporter) Report(taskID string, taskStatus exectsk.ExeTaskStatus) { +func (r *Reporter) Report(taskID string, taskStatus exectsk.TaskStatus) { r.taskStatusLock.Lock() defer r.taskStatusLock.Unlock() @@ -65,7 +65,7 @@ func (r *Reporter) Serve() error { for taskID, status := range r.taskStatus { taskStatus = append(taskStatus, mgrmq.NewExecutorTaskStatus(taskID, status)) } - r.taskStatus = make(map[string]exectsk.ExeTaskStatus) + r.taskStatus = make(map[string]exectsk.TaskStatus) r.taskStatusLock.Unlock() _, err := magCli.ReportExecutorTaskStatus(mgrmq.NewReportExecutorTaskStatus(r.executorID, taskStatus)) diff --git a/executor/internal/task/task.go b/executor/internal/task/task.go index 324f03e..63bbf29 100644 --- a/executor/internal/task/task.go +++ b/executor/internal/task/task.go @@ -36,7 +36,7 @@ func NewManager(reporter *reporter.Reporter) Manager { } } -func (m *Manager) StartByInfo(info exectsk.ExeTaskInfo) (*Task, error) { +func (m *Manager) StartByInfo(info exectsk.TaskInfo) (*Task, error) { infoType := myreflect.TypeOfValue(info) ctor, ok := taskFromInfoCtors[infoType] @@ -47,10 +47,10 @@ func (m *Manager) StartByInfo(info exectsk.ExeTaskInfo) (*Task, error) { return m.StartNew(ctor(info)), nil } -var taskFromInfoCtors map[reflect.Type]func(exectsk.ExeTaskInfo) TaskBody = make(map[reflect.Type]func(exectsk.ExeTaskInfo) task.TaskBody[TaskContext]) +var taskFromInfoCtors map[reflect.Type]func(exectsk.TaskInfo) TaskBody = make(map[reflect.Type]func(exectsk.TaskInfo) task.TaskBody[TaskContext]) -func Register[TInfo exectsk.ExeTaskInfo, TTaskBody TaskBody](ctor func(info TInfo) TTaskBody) { - taskFromInfoCtors[myreflect.TypeOf[TInfo]()] = func(info exectsk.ExeTaskInfo) TaskBody { +func Register[TInfo exectsk.TaskInfo, TTaskBody TaskBody](ctor func(info TInfo) TTaskBody) { + taskFromInfoCtors[myreflect.TypeOf[TInfo]()] = func(info exectsk.TaskInfo) TaskBody { return ctor(info.(TInfo)) } } diff --git a/manager/internal/advisormgr/advisormgr.go b/manager/internal/advisormgr/advisormgr.go index 3ea541e..0ba7231 100644 --- a/manager/internal/advisormgr/advisormgr.go +++ b/manager/internal/advisormgr/advisormgr.go @@ -25,7 +25,7 @@ type AdvisorInfo struct { lastReportTime time.Time } -type OnTaskUpdatedCallbackFn func(jobID schsdk.JobID, fullTaskID string, taskStatus advtsk.AdvTaskStatus) +type OnTaskUpdatedCallbackFn func(jobID schsdk.JobID, fullTaskID string, taskStatus advtsk.TaskStatus) type OnTimeoutCallbackFn func(jobID schsdk.JobID, fullTaskID string) type Manager struct { @@ -86,7 +86,7 @@ func (m *Manager) Report(advID schmod.AdvisorID, taskStatus []mgrmq.AdvisorTaskS } // 启动一个Task,并将其关联到指定的Job。返回一个在各Executor之间唯一的TaskID -func (m *Manager) StartTask(jobID schsdk.JobID, info advtsk.AdvTaskInfo) (string, error) { +func (m *Manager) StartTask(jobID schsdk.JobID, info advtsk.TaskInfo) (string, error) { m.lock.Lock() defer m.lock.Unlock() diff --git a/manager/internal/executormgr/executormgr.go b/manager/internal/executormgr/executormgr.go index e9cefde..33a5edb 100644 --- a/manager/internal/executormgr/executormgr.go +++ b/manager/internal/executormgr/executormgr.go @@ -26,7 +26,7 @@ type ExecutorInfo struct { lastReportTime time.Time } -type OnTaskUpdatedCallbackFn func(jobID schsdk.JobID, fullTaskID string, taskStatus exetsk.ExeTaskStatus) +type OnTaskUpdatedCallbackFn func(jobID schsdk.JobID, fullTaskID string, taskStatus exetsk.TaskStatus) type OnTimeoutCallbackFn func(jobID schsdk.JobID, fullTaskID string) type Manager struct { @@ -87,7 +87,7 @@ func (m *Manager) Report(execID schmod.ExecutorID, taskStatus []mgrmq.ExecutorTa } // 启动一个Task,并将其关联到指定的Job。返回一个在各Executor之间唯一的TaskID -func (m *Manager) StartTask(jobID schsdk.JobID, info exetsk.ExeTaskInfo) (string, error) { +func (m *Manager) StartTask(jobID schsdk.JobID, info exetsk.TaskInfo) (string, error) { m.lock.Lock() defer m.lock.Unlock() diff --git a/manager/internal/jobmgr/event/advisor_task_updated.go b/manager/internal/jobmgr/event/advisor_task_updated.go index 7a0bbff..1e3931c 100644 --- a/manager/internal/jobmgr/event/advisor_task_updated.go +++ b/manager/internal/jobmgr/event/advisor_task_updated.go @@ -5,17 +5,17 @@ import advtsk "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq/advisor/task" // advisor上报任务进度 type AdvisorTaskUpdated struct { FullTaskID string - TaskStatus advtsk.AdvTaskStatus + TaskStatus advtsk.TaskStatus } -func NewAdvisorTaskUpdated(fullTaskID string, taskStatus advtsk.AdvTaskStatus) *AdvisorTaskUpdated { +func NewAdvisorTaskUpdated(fullTaskID string, taskStatus advtsk.TaskStatus) *AdvisorTaskUpdated { return &AdvisorTaskUpdated{ FullTaskID: fullTaskID, TaskStatus: taskStatus, } } -func AssertAdvisorTaskStatus[T advtsk.AdvTaskStatus](evt Event, fullTaskID string) (T, error) { +func AssertAdvisorTaskStatus[T advtsk.TaskStatus](evt Event, fullTaskID string) (T, error) { var ret T if evt == nil { return ret, ErrUnconcernedTask diff --git a/manager/internal/jobmgr/event/executor_task_updated.go b/manager/internal/jobmgr/event/executor_task_updated.go index 4ae50a3..fe299d6 100644 --- a/manager/internal/jobmgr/event/executor_task_updated.go +++ b/manager/internal/jobmgr/event/executor_task_updated.go @@ -7,17 +7,17 @@ import ( // executor上报任务进度 type ExecutorTaskUpdated struct { FullTaskID string - TaskStatus exectsk.ExeTaskStatus + TaskStatus exectsk.TaskStatus } -func NewExecutorTaskUpdated(fullTaskID string, taskStatus exectsk.ExeTaskStatus) *ExecutorTaskUpdated { +func NewExecutorTaskUpdated(fullTaskID string, taskStatus exectsk.TaskStatus) *ExecutorTaskUpdated { return &ExecutorTaskUpdated{ FullTaskID: fullTaskID, TaskStatus: taskStatus, } } -func AssertExecutorTaskStatus[T exectsk.ExeTaskStatus](evt Event, fullTaskID string) (T, error) { +func AssertExecutorTaskStatus[T exectsk.TaskStatus](evt Event, fullTaskID string) (T, error) { var ret T if evt == nil { return ret, ErrUnconcernedTask diff --git a/manager/internal/jobmgr/jobmgr.go b/manager/internal/jobmgr/jobmgr.go index ae855e6..c6cfa40 100644 --- a/manager/internal/jobmgr/jobmgr.go +++ b/manager/internal/jobmgr/jobmgr.go @@ -181,7 +181,7 @@ func (m *Manager) LocalFileUploaded(jobSetID schsdk.JobSetID, localPath string, return nil } -func (m *Manager) executorTaskUpdated(jobID schsdk.JobID, fullTaskID string, taskStatus exectsk.ExeTaskStatus) { +func (m *Manager) executorTaskUpdated(jobID schsdk.JobID, fullTaskID string, taskStatus exectsk.TaskStatus) { m.pubLock.Lock() defer m.pubLock.Unlock() @@ -205,7 +205,7 @@ func (m *Manager) executorTaskTimeout(jobID schsdk.JobID, fullTaskID string) { job.Handler.OnEvent(event.ToJob(jobID), event.NewExecutorTaskTimeout(fullTaskID)) } -func (m *Manager) advisorTaskUpdated(jobID schsdk.JobID, fullTaskID string, taskStatus advtsk.AdvTaskStatus) { +func (m *Manager) advisorTaskUpdated(jobID schsdk.JobID, fullTaskID string, taskStatus advtsk.TaskStatus) { m.pubLock.Lock() defer m.pubLock.Unlock()