forked from JointCloud/pcm-coordinator
106 lines
3.2 KiB
Go
106 lines
3.2 KiB
Go
package schedule
|
|
|
|
import (
|
|
"context"
|
|
"gitlink.org.cn/JointCloud/pcm-coordinator/internal/scheduler/schedulers"
|
|
"gitlink.org.cn/JointCloud/pcm-coordinator/internal/scheduler/schedulers/option"
|
|
"gitlink.org.cn/JointCloud/pcm-coordinator/internal/scheduler/service/executor"
|
|
"gitlink.org.cn/JointCloud/pcm-coordinator/internal/svc"
|
|
"gitlink.org.cn/JointCloud/pcm-coordinator/internal/types"
|
|
"gitlink.org.cn/JointCloud/pcm-coordinator/pkg/constants"
|
|
"strconv"
|
|
"strings"
|
|
|
|
"github.com/zeromicro/go-zero/core/logx"
|
|
)
|
|
|
|
type ScheduleSubmitLogic struct {
|
|
logx.Logger
|
|
ctx context.Context
|
|
svcCtx *svc.ServiceContext
|
|
}
|
|
|
|
func NewScheduleSubmitLogic(ctx context.Context, svcCtx *svc.ServiceContext) *ScheduleSubmitLogic {
|
|
return &ScheduleSubmitLogic{
|
|
Logger: logx.WithContext(ctx),
|
|
ctx: ctx,
|
|
svcCtx: svcCtx,
|
|
}
|
|
}
|
|
|
|
func (l *ScheduleSubmitLogic) ScheduleSubmit(req *types.ScheduleReq) (resp *types.ScheduleResp, err error) {
|
|
resp = &types.ScheduleResp{}
|
|
opt := &option.AiOption{
|
|
AdapterId: req.AiOption.AdapterId,
|
|
ClusterIds: req.AiOption.AiClusterIds,
|
|
TaskName: req.AiOption.TaskName,
|
|
ResourceType: req.AiOption.ResourceType,
|
|
Replica: req.AiOption.Replica,
|
|
ComputeCard: req.AiOption.ComputeCard,
|
|
Tops: req.AiOption.Tops,
|
|
TaskType: req.AiOption.TaskType,
|
|
DatasetsName: req.AiOption.Datasets,
|
|
AlgorithmName: req.AiOption.Algorithm,
|
|
StrategyName: req.AiOption.Strategy,
|
|
ClusterToStaticWeight: req.AiOption.StaticWeightMap,
|
|
Params: req.AiOption.Params,
|
|
Envs: req.AiOption.Envs,
|
|
Cmd: req.AiOption.Cmd,
|
|
}
|
|
aiSchdl, err := schedulers.NewAiScheduler(l.ctx, "", l.svcCtx.Scheduler, opt)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
results, err := l.svcCtx.Scheduler.AssignAndSchedule(aiSchdl, executor.SUBMIT_MODE_JOINT_CLOUD, nil)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
switch opt.GetOptionType() {
|
|
case option.AI:
|
|
rs := (results).([]*schedulers.AiResult)
|
|
var synergystatus int64
|
|
if len(rs) > 1 {
|
|
synergystatus = 1
|
|
}
|
|
|
|
taskId, err := l.svcCtx.Scheduler.CreateTask(req.AiOption.TaskName, "", 0, synergystatus, req.AiOption.Strategy, "", req.Token, "", &l.svcCtx.Config)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
adapterName, err := l.svcCtx.Scheduler.AiStorages.GetAdapterNameById(rs[0].AdapterId)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
for _, r := range rs {
|
|
scheResult := &types.ScheduleResult{}
|
|
scheResult.ClusterId = r.ClusterId
|
|
scheResult.TaskId = strconv.FormatInt(taskId, 10)
|
|
scheResult.JobId = r.JobId
|
|
scheResult.Strategy = r.Strategy
|
|
scheResult.Card = strings.ToUpper(r.Card)
|
|
scheResult.Replica = r.Replica
|
|
scheResult.Msg = r.Msg
|
|
|
|
opt.ComputeCard = strings.ToUpper(r.Card)
|
|
|
|
clusterName, _ := l.svcCtx.Scheduler.AiStorages.GetClusterNameById(r.ClusterId)
|
|
|
|
err := l.svcCtx.Scheduler.AiStorages.SaveAiTask(taskId, opt, adapterName, r.ClusterId, clusterName, r.JobId, constants.Saved, r.Msg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
l.svcCtx.Scheduler.AiStorages.AddNoticeInfo(r.AdapterId, adapterName, r.ClusterId, clusterName, r.TaskName, "create", "任务创建中")
|
|
|
|
resp.Results = append(resp.Results, scheResult)
|
|
}
|
|
|
|
}
|
|
|
|
return resp, nil
|
|
}
|