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 }