|
- package repo
-
- import (
- "bufio"
- "encoding/json"
- "errors"
- "io"
- "io/ioutil"
- "net/http"
- "os"
- "strconv"
- "strings"
- "time"
-
- "code.gitea.io/gitea/models"
- "code.gitea.io/gitea/modules/aisafety"
- "code.gitea.io/gitea/modules/cloudbrain"
- "code.gitea.io/gitea/modules/context"
- "code.gitea.io/gitea/modules/git"
- "code.gitea.io/gitea/modules/grampus"
- "code.gitea.io/gitea/modules/log"
- "code.gitea.io/gitea/modules/modelarts"
- "code.gitea.io/gitea/modules/setting"
- "code.gitea.io/gitea/modules/storage"
- "code.gitea.io/gitea/modules/timeutil"
- "code.gitea.io/gitea/modules/util"
- "code.gitea.io/gitea/services/cloudbrain/resource"
- "code.gitea.io/gitea/services/reward/point/account"
- uuid "github.com/satori/go.uuid"
- )
-
- const (
- tplModelSafetyTestCreateGrampusGpu = "repo/modelsafety/newgrampusgpu"
- tplModelSafetyTestCreateGrampusNpu = "repo/modelsafety/newgrampusnpu"
- tplModelSafetyTestCreateGpu = "repo/modelsafety/newgpu"
- tplModelSafetyTestCreateNpu = "repo/modelsafety/newnpu"
- tplModelSafetyTestShow = "repo/modelsafety/show"
- )
-
- func CloudBrainAiSafetyCreateTest(ctx *context.Context) {
- log.Info("start to create CloudBrainAiSafetyCreate")
- uuid := uuid.NewV4()
- id := uuid.String()
- seriaNoParas := ctx.Query("serialNo")
- fileName := ctx.Query("fileName")
-
- //if jobType == string(models.JobTypeBenchmark) {
- req := aisafety.TaskReq{
- UnionId: id,
- EvalName: "test1",
- EvalContent: "test1",
- TLPath: "test1",
- Indicators: []string{"ACC", "ASS"},
- CDName: "CIFAR10_1000_FGSM",
- BDName: "CIFAR10_1000基础数据集",
- }
- aisafety.GetAlgorithmList()
- if seriaNoParas != "" {
- aisafety.GetTaskStatus(seriaNoParas)
- } else {
- jsonStr, err := getJsonContent("http://192.168.207.34:8065/Test_zap1234/openi_aisafety/raw/branch/master/result/" + fileName)
- serialNo, err := aisafety.CreateSafetyTask(req, jsonStr)
- if err == nil {
- log.Info("serialNo=" + serialNo)
- time.Sleep(time.Duration(2) * time.Second)
- aisafety.GetTaskStatus(serialNo)
- } else {
- log.Info("CreateSafetyTask error," + err.Error())
- }
- }
- }
-
- func GetAiSafetyTaskByJob(job *models.Cloudbrain) {
- if job == nil {
- log.Error("GetCloudbrainByJobID failed")
- return
- }
- syncAiSafetyTaskStatus(job)
- }
-
- func GetAiSafetyTaskTmpl(ctx *context.Context) {
- ctx.Data["id"] = ctx.Params(":jobid")
- ctx.HTML(200, tplModelSafetyTestShow)
- }
-
- func GetAiSafetyTask(ctx *context.Context) {
- var ID = ctx.Params(":jobid")
- job, err := models.GetCloudbrainByIDWithDeleted(ID)
- if err != nil {
- log.Error("GetCloudbrainByJobID failed:" + err.Error())
- return
- }
- syncAiSafetyTaskStatus(job)
- job, err = models.GetCloudbrainByIDWithDeleted(ID)
- job.BenchmarkType = "CV"
- job.BenchmarkTypeName = "Classification"
- ctx.JSON(200, job)
- }
-
- func StopAiSafetyTask(ctx *context.Context) {
- var ID = ctx.Params(":jobid")
- task, err := models.GetCloudbrainByIDWithDeleted(ID)
- result := make(map[string]interface{})
- result["code"] = -1
- if err != nil {
- log.Error("GetCloudbrainByJobID failed:" + err.Error())
- result["msg"] = "No such task."
- ctx.JSON(200, result)
- return
- }
- if isTaskNotFinished(task.Status) {
- if task.Type == models.TypeCloudBrainTwo {
- //queryTaskStatusFromCloudbrainTwo(job)
- } else if task.Type == models.TypeCloudBrainOne {
- if task.Status == string(models.JobStopped) || task.Status == string(models.JobFailed) || task.Status == string(models.JobSucceeded) {
- log.Error("the job(%s) has been stopped", task.JobName, ctx.Data["msgID"])
- result["msg"] = "cloudbrain.Already_stopped"
- ctx.JSON(200, result)
- return
- }
- err := cloudbrain.StopJob(task.JobID)
- if err != nil {
- log.Error("StopJob(%s) failed:%v", task.JobName, err, ctx.Data["msgID"])
- result["msg"] = "cloudbrain.Stopped_failed"
- ctx.JSON(200, result)
- return
- }
- task.Status = string(models.JobStopped)
- if task.EndTime == 0 {
- task.EndTime = timeutil.TimeStampNow()
- }
- task.ComputeAndSetDuration()
- err = models.UpdateJob(task)
- if err != nil {
- log.Error("UpdateJob(%s) failed:%v", task.JobName, err, ctx.Data["msgID"])
- result["msg"] = "cloudbrain.Stopped_success_update_status_fail"
- ctx.JSON(200, result)
- return
- }
- }
- } else {
- if task.Status == string(models.ModelSafetyTesting) {
- //修改为Failed
- task.Status = string(models.JobStopped)
- if task.EndTime == 0 {
- task.EndTime = timeutil.TimeStampNow()
- }
- task.ComputeAndSetDuration()
- err = models.UpdateJob(task)
- if err != nil {
- log.Error("UpdateJob(%s) failed:%v", task.JobName, err, ctx.Data["msgID"])
- result["msg"] = "cloudbrain.Stopped_success_update_status_fail"
- ctx.JSON(200, result)
- return
- }
- } else {
- log.Info("The job is finished. status=" + task.Status)
- }
- }
-
- }
-
- func DelAiSafetyTask(ctx *context.Context) {
- var ID = ctx.Params(":jobid")
- task, err := models.GetCloudbrainByIDWithDeleted(ID)
- result := make(map[string]interface{})
- result["code"] = 1
- if err != nil {
- log.Error("GetCloudbrainByJobID failed:" + err.Error())
- result["msg"] = "No such task."
- ctx.JSON(200, result)
- return
- }
- if task.Status != string(models.JobStopped) && task.Status != string(models.JobFailed) && task.Status != string(models.JobSucceeded) {
- log.Error("the job(%s) has not been stopped", task.JobName, ctx.Data["msgID"])
- result["msg"] = "the job(" + task.JobName + ") has not been stopped"
- ctx.JSON(200, result)
- return
- }
- if task.Type == models.TypeCloudBrainOne {
- DeleteCloudbrainJobStorage(task.JobName, models.TypeCloudBrainOne)
- }
- err = models.DeleteJob(task)
- if err != nil {
- result["msg"] = err.Error()
- ctx.JSON(200, result)
- return
- }
- result["code"] = 0
- result["msg"] = "Succeed"
- ctx.JSON(200, result)
- }
-
- func syncAiSafetyTaskStatus(job *models.Cloudbrain) {
- if isTaskNotFinished(job.Status) {
- if job.Type == models.TypeCloudBrainTwo {
- queryTaskStatusFromCloudbrainTwo(job)
- } else if job.Type == models.TypeCloudBrainOne {
- queryTaskStatusFromCloudbrain(job)
- } else if job.Type == models.TypeC2Net {
- queryTaskStatusFromGrampus(job)
- }
- } else {
- if job.Status == string(models.ModelSafetyTesting) {
- queryTaskStatusFromModelSafetyTestServer(job)
- } else {
- log.Info("The job is finished. status=" + job.Status)
- }
- }
- }
-
- func TimerHandleModelSafetyTestTask() {
- tasks, err := models.GetModelSafetyTestTask()
- if err == nil {
- if tasks != nil && len(tasks) > 0 {
- for _, job := range tasks {
- syncAiSafetyTaskStatus(job)
- }
- } else {
- log.Info("query running model safety test task 0.")
- }
- } else {
- log.Info("query running model safety test task err." + err.Error())
- }
- }
-
- func queryTaskStatusFromGrampus(task *models.Cloudbrain) {
- if task.DeletedAt.IsZero() { //normal record
- result, err := grampus.GetJob(task.JobID)
- if err != nil {
- log.Error("GetJob failed:" + err.Error())
- return
- }
-
- if result != nil {
- if len(result.JobInfo.Tasks[0].CenterID) == 1 && len(result.JobInfo.Tasks[0].CenterName) == 1 {
- task.AiCenter = result.JobInfo.Tasks[0].CenterID[0] + "+" + result.JobInfo.Tasks[0].CenterName[0]
- }
- task.Status = grampus.TransTrainJobStatus(result.JobInfo.Status)
- if task.Status != models.GrampusStatusSucceeded {
- if task.Status != result.JobInfo.Status || result.JobInfo.Status == models.GrampusStatusRunning {
- task.Duration = result.JobInfo.RunSec
- if task.Duration < 0 {
- task.Duration = 0
- }
- task.TrainJobDuration = models.ConvertDurationToStr(task.Duration)
-
- if task.StartTime == 0 && result.JobInfo.StartedAt > 0 {
- task.StartTime = timeutil.TimeStamp(result.JobInfo.StartedAt)
- }
- if task.EndTime == 0 && models.IsTrainJobTerminal(task.Status) && task.StartTime > 0 {
- task.EndTime = task.StartTime.Add(task.Duration)
- }
- task.CorrectCreateUnix()
- err = models.UpdateJob(task)
- if err != nil {
- log.Error("UpdateJob failed:" + err.Error())
- }
- }
- } else {
- task.Status = string(models.ModelSafetyTesting)
- err = models.UpdateJob(task)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- //send msg to beihang
- sendGPUInferenceResultToTest(task)
- }
- }
- }
-
- }
-
- func queryTaskStatusFromCloudbrainTwo(job *models.Cloudbrain) {
- log.Info("The task not finished,name=" + job.DisplayJobName)
- result, err := modelarts.GetTrainJob(job.JobID, strconv.FormatInt(job.VersionID, 10))
- if err != nil {
- log.Info("query train job error." + err.Error())
- return
- }
-
- job.Status = modelarts.TransTrainJobStatus(result.IntStatus)
- job.Duration = result.Duration
- job.TrainJobDuration = result.TrainJobDuration
- if job.Status != string(models.ModelArtsTrainJobCompleted) {
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- } else {
- job.Status = string(models.ModelSafetyTesting)
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- //send msg to beihang
- sendNPUInferenceResultToTest(job)
- }
-
- }
-
- func sendNPUInferenceResultToTest(job *models.Cloudbrain) {
- datasetname := job.DatasetName
- datasetnames := strings.Split(datasetname, ";")
- indicator := job.LabelName
-
- req := aisafety.TaskReq{
- UnionId: job.JobID,
- EvalName: job.DisplayJobName,
- EvalContent: job.Description,
- TLPath: "test",
- Indicators: strings.Split(indicator, ";"),
- CDName: datasetnames[1],
- BDName: datasetnames[0],
- }
- jsonContent := ""
- VersionOutputPath := modelarts.GetOutputPathByCount(modelarts.TotalVersionCount)
- resultPath := modelarts.JobPath + job.JobName + modelarts.ResultPath + VersionOutputPath + "/result.json"
- body, err := storage.ObsDownloadAFile(setting.Bucket, resultPath)
- if err != nil {
- log.Info("ObsDownloadAFile error." + err.Error() + " resultPath=" + resultPath)
- } else {
- defer body.Close()
- var data []byte
- p := make([]byte, 4096)
- var readErr error
- var readCount int
- for {
- readCount, readErr = body.Read(p)
- if readCount > 0 {
- data = append(data, p[:readCount]...)
- }
- if readErr != nil || readCount == 0 {
- break
- }
- }
- jsonContent = string(data)
- }
-
- if jsonContent != "" {
- serialNo, err := aisafety.CreateSafetyTask(req, jsonContent)
- if err == nil {
- //update serial no to db
- job.PreVersionName = serialNo
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- }
- } else {
- log.Info("The json is null. so set it failed.")
- //update task failed.
- job.Status = string(models.ModelArtsTrainJobFailed)
- err := models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- }
- }
-
- func queryTaskStatusFromCloudbrain(job *models.Cloudbrain) {
-
- log.Info("The task not finished,name=" + job.DisplayJobName)
- jobResult, err := cloudbrain.GetJob(job.JobID)
-
- result, err := models.ConvertToJobResultPayload(jobResult.Payload)
- if err != nil {
- log.Error("ConvertToJobResultPayload failed:", err)
- return
- }
- job.Status = result.JobStatus.State
- if result.JobStatus.State != string(models.JobWaiting) && result.JobStatus.State != string(models.JobFailed) {
- taskRoles := result.TaskRoles
- taskRes, _ := models.ConvertToTaskPod(taskRoles[cloudbrain.SubTaskName].(map[string]interface{}))
- job.Status = taskRes.TaskStatuses[0].State
- }
-
- if result.JobStatus.State != string(models.JobSucceeded) {
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- } else {
- //
- job.Status = string(models.ModelSafetyTesting)
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- //send msg to beihang
- sendGPUInferenceResultToTest(job)
- }
- }
-
- func queryTaskStatusFromModelSafetyTestServer(job *models.Cloudbrain) {
- result, err := aisafety.GetTaskStatus(job.PreVersionName)
- if err == nil {
- if result.Code == "0" {
- if result.Data.Status == 1 {
- log.Info("The task is running....")
- } else {
- if result.Data.Code == 0 {
- job.ResultJson = result.Data.StandardJson
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- }
- }
- } else {
- log.Info("The task is failed.")
- job.Status = string(models.JobFailed)
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- }
- } else {
- log.Info("The task not found.....")
- }
- }
-
- func sendGPUInferenceResultToTest(job *models.Cloudbrain) {
- datasetname := job.DatasetName
- datasetnames := strings.Split(datasetname, ";")
- indicator := job.LabelName
-
- req := aisafety.TaskReq{
- UnionId: job.JobID,
- EvalName: job.DisplayJobName,
- EvalContent: job.Description,
- TLPath: "test",
- Indicators: strings.Split(indicator, ";"),
- CDName: datasetnames[1],
- BDName: datasetnames[0],
- }
-
- resultDir := "/model"
- prefix := "/" + setting.CBCodePathPrefix + job.JobName + resultDir
- files, err := storage.GetOneLevelAllObjectUnderDirMinio(setting.Attachment.Minio.Bucket, prefix, "")
- if err != nil {
- log.Error("query cloudbrain one model failed: %v", err)
- return
- }
- jsonContent := ""
- for _, file := range files {
- if strings.HasSuffix(file.FileName, "result.json") {
- path := storage.GetMinioPath(job.JobName+resultDir+"/", file.FileName)
- log.Info("path=" + path)
- reader, err := os.Open(path)
- defer reader.Close()
- if err == nil {
- r := bufio.NewReader(reader)
- for {
- line, error := r.ReadString('\n')
- if error == io.EOF {
- log.Info("read file completed.")
- break
- }
- if error != nil {
- log.Info("read file error." + error.Error())
- break
- }
- jsonContent += line
- }
- }
- break
- }
- }
- if jsonContent != "" {
- serialNo, err := aisafety.CreateSafetyTask(req, jsonContent)
- if err == nil {
- //update serial no to db
- job.PreVersionName = serialNo
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- }
- } else {
- log.Info("The json is null. so set it failed.")
- //update task failed.
- job.Status = string(models.JobFailed)
- err = models.UpdateJob(job)
- if err != nil {
- log.Error("UpdateJob failed:", err)
- }
- }
- }
-
- func isTaskNotFinished(status string) bool {
- if status == string(models.ModelArtsTrainJobRunning) || status == string(models.ModelArtsTrainJobWaiting) {
- return true
- }
- if status == string(models.JobWaiting) || status == string(models.JobRunning) {
- return true
- }
-
- if status == string(models.ModelArtsTrainJobUnknown) || status == string(models.ModelArtsTrainJobInit) {
- return true
- }
- if status == string(models.ModelArtsTrainJobImageCreating) || status == string(models.ModelArtsTrainJobSubmitTrying) {
- return true
- }
- return false
- }
-
- func AiSafetyCreateForGetGPU(ctx *context.Context) {
- t := time.Now()
- ctx.Data["PageIsCloudBrain"] = true
- ctx.Data["IsCreate"] = true
- ctx.Data["datasetType"] = models.TypeCloudBrainOne
- ctx.Data["BaseDataSetName"] = setting.ModelSafetyTest.BaseDataSetName
- ctx.Data["BaseDataSetUUID"] = setting.ModelSafetyTest.BaseDataSetUUID
- ctx.Data["CombatDataSetName"] = setting.ModelSafetyTest.CombatDataSetName
- ctx.Data["CombatDataSetUUID"] = setting.ModelSafetyTest.CombatDataSetUUID
- var displayJobName = jobNamePrefixValid(cutString(ctx.User.Name, 5)) + t.Format("2006010215") + strconv.Itoa(int(t.Unix()))[5:]
- ctx.Data["display_job_name"] = displayJobName
- prepareCloudbrainOneSpecs(ctx)
- queuesDetail, _ := cloudbrain.GetQueuesDetail()
- if queuesDetail != nil {
- ctx.Data["QueuesDetail"] = queuesDetail
- }
- ctx.HTML(200, tplModelSafetyTestCreateGpu)
- }
- func AiSafetyCreateForGetGrampusGPU(ctx *context.Context) {
- ctx.Data["PageIsCloudBrain"] = true
- ctx.Data["IsCreate"] = true
- ctx.Data["datasetType"] = models.TypeCloudBrainOne
- ctx.Data["BaseDataSetName"] = setting.ModelSafetyTest.BaseDataSetName
- ctx.Data["BaseDataSetUUID"] = setting.ModelSafetyTest.BaseDataSetUUID
- ctx.Data["CombatDataSetName"] = setting.ModelSafetyTest.CombatDataSetName
- ctx.Data["CombatDataSetUUID"] = setting.ModelSafetyTest.CombatDataSetUUID
- err := GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- if err != nil {
- ctx.ServerError("get new train-job info failed", err)
- return
- }
- ctx.HTML(200, tplModelSafetyTestCreateGrampusGpu)
- }
-
- func AiSafetyCreateForGetGrampusNPU(ctx *context.Context) {
- ctx.Data["PageIsCloudBrain"] = true
- ctx.Data["IsCreate"] = true
-
- ctx.Data["datasetType"] = models.TypeCloudBrainTwo
- ctx.Data["BaseDataSetName"] = setting.ModelSafetyTest.BaseDataSetName
- ctx.Data["BaseDataSetUUID"] = setting.ModelSafetyTest.BaseDataSetUUID
- ctx.Data["CombatDataSetName"] = setting.ModelSafetyTest.CombatDataSetName
- ctx.Data["CombatDataSetUUID"] = setting.ModelSafetyTest.CombatDataSetUUID
-
- err := GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
- if err != nil {
- ctx.ServerError("get new train-job info failed", err)
- return
- }
-
- ctx.HTML(200, tplModelSafetyTestCreateGrampusNpu)
- }
-
- func AiSafetyCreateForGetNPU(ctx *context.Context) {
- t := time.Now()
- ctx.Data["PageIsCloudBrain"] = true
- ctx.Data["IsCreate"] = true
- var displayJobName = jobNamePrefixValid(cutString(ctx.User.Name, 5)) + t.Format("2006010215") + strconv.Itoa(int(t.Unix()))[5:]
- ctx.Data["display_job_name"] = displayJobName
- ctx.Data["datasetType"] = models.TypeCloudBrainTwo
- ctx.Data["BaseDataSetName"] = setting.ModelSafetyTest.BaseDataSetName
- ctx.Data["BaseDataSetUUID"] = setting.ModelSafetyTest.BaseDataSetUUID
- ctx.Data["CombatDataSetName"] = setting.ModelSafetyTest.CombatDataSetName
- ctx.Data["CombatDataSetUUID"] = setting.ModelSafetyTest.CombatDataSetUUID
-
- var resourcePools modelarts.ResourcePool
- if err := json.Unmarshal([]byte(setting.ResourcePools), &resourcePools); err != nil {
- ctx.ServerError("json.Unmarshal failed:", err)
- }
- ctx.Data["resource_pools"] = resourcePools.Info
-
- var engines modelarts.Engine
- if err := json.Unmarshal([]byte(setting.Engines), &engines); err != nil {
- ctx.ServerError("json.Unmarshal failed:", err)
- }
- ctx.Data["engines"] = engines.Info
-
- var versionInfos modelarts.VersionInfo
- if err := json.Unmarshal([]byte(setting.EngineVersions), &versionInfos); err != nil {
- ctx.ServerError("json.Unmarshal failed:", err)
- }
- ctx.Data["engine_versions"] = versionInfos.Version
-
- prepareCloudbrainTwoInferenceSpecs(ctx)
- waitCount := cloudbrain.GetWaitingCloudbrainCount(models.TypeCloudBrainTwo, "")
- ctx.Data["WaitCount"] = waitCount
- ctx.HTML(200, tplModelSafetyTestCreateNpu)
- }
-
- func AiSafetyCreateForPost(ctx *context.Context) {
- ctx.Data["PageIsCloudBrain"] = true
- displayJobName := ctx.Query("display_job_name")
- jobName := util.ConvertDisplayJobNameToJobName(displayJobName)
-
- taskType := ctx.QueryInt("type")
- description := ctx.Query("description")
- ctx.Data["description"] = description
-
- repo := ctx.Repo.Repository
-
- tpname := tplCloudBrainModelSafetyNewNpu
- if taskType == models.TypeCloudBrainOne {
- tpname = tplCloudBrainModelSafetyNewGpu
- }
-
- tasks, err := models.GetCloudbrainsByDisplayJobName(repo.ID, string(models.JobTypeModelSafety), displayJobName)
- if err == nil {
- if len(tasks) != 0 {
- log.Error("the job name did already exist", ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("the job name did already exist", tpname, nil)
- return
- }
- } else {
- if !models.IsErrJobNotExist(err) {
- log.Error("system error, %v", err, ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("system error", tpname, nil)
- return
- }
- }
-
- if !jobNamePattern.MatchString(jobName) {
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("repo.cloudbrain_jobname_err"), tpname, nil)
- return
- }
-
- count, err := models.GetModelSafetyCountByUserID(ctx.User.ID)
- if err != nil {
- log.Error("GetCloudbrainCountByUserID failed:%v", err, ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("system error", tpname, nil)
- return
- } else {
- if count >= 1 {
- log.Error("the user already has running or waiting task", ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("repo.cloudbrain.morethanonejob"), tpname, nil)
- return
- }
- }
- BootFile := ctx.Query("boot_file")
- bootFileExist, err := ctx.Repo.FileExists(BootFile, cloudbrain.DefaultBranchName)
- if err != nil || !bootFileExist {
- log.Error("Get bootfile error:", err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("repo.cloudbrain_bootfile_err"), tpname, nil)
- return
- }
-
- if taskType == models.TypeCloudBrainTwo {
- createForNPU(ctx, jobName)
- } else if taskType == models.TypeCloudBrainOne {
- createForGPU(ctx, jobName)
- } else if taskType == models.TypeC2Net {
- ComputeResource := ctx.Query("compute_resource")
- if ComputeResource == models.NPUResource {
- createForGrampusNPU(ctx, jobName)
- } else if ComputeResource == models.GPUResource {
- createForGrampusGPU(ctx, jobName)
- }
- }
- ctx.Redirect(setting.AppSubURL + ctx.Repo.RepoLink + "/cloudbrain/benchmark")
- }
-
- func createForGrampusGPU(ctx *context.Context, jobName string) {
- BootFile := ctx.Query("boot_file")
- displayJobName := ctx.Query("display_job_name")
- description := ctx.Query("description")
- image := strings.TrimSpace(ctx.Query("image"))
- srcDataset := ctx.Query("src_dataset") //uuid
- combatDataset := ctx.Query("combat_dataset") //uuid
- evaluationIndex := ctx.Query("evaluationIndex")
- Params := ctx.Query("run_para_list")
- specId := ctx.QueryInt64("spec_id")
- TrainUrl := ctx.Query("train_url")
- CkptName := ctx.Query("ckpt_name")
- ModelName := ctx.Query("ModelName")
- ModelVersion := ctx.Query("ModelVersion")
- repo := ctx.Repo.Repository
- codeLocalPath := setting.JobPath + jobName + cloudbrain.CodeMountPath + "/"
- codeMinioPath := setting.CBCodePathPrefix + jobName + cloudbrain.CodeMountPath + "/"
- //check specification
- spec, err := resource.GetAndCheckSpec(ctx.User.ID, specId, models.FindSpecsOptions{
- JobType: models.JobTypeTrain,
- ComputeResource: models.GPU,
- Cluster: models.C2NetCluster,
- })
- if err != nil || spec == nil {
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr("Resource specification not available", tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- if !account.IsPointBalanceEnough(ctx.User.ID, spec.UnitPrice) {
- log.Error("point balance is not enough,userId=%d specId=%d", ctx.User.ID, spec.ID)
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(ctx.Tr("points.insufficient_points_balance"), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- //check dataset
- uuid := srcDataset + ";" + combatDataset
- datasetInfos, datasetNames, err := models.GetDatasetInfo(uuid, models.GPU)
- if err != nil {
- log.Error("GetDatasetInfo failed: %v", err, ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.error.dataset_select"), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- //prepare code and out path
- _, err = ioutil.ReadDir(codeLocalPath)
- if err == nil {
- os.RemoveAll(codeLocalPath)
- }
-
- if err := downloadZipCode(ctx, codeLocalPath, cloudbrain.DefaultBranchName); err != nil {
- log.Error("downloadZipCode failed, server timed out: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.load_code_failed"), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- //todo: upload code (send to file_server todo this work?)
- //upload code
- if err := uploadCodeToMinio(codeLocalPath+"/", jobName, cloudbrain.CodeMountPath+"/"); err != nil {
- log.Error("Failed to uploadCodeToMinio: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.load_code_failed"), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- modelPath := setting.JobPath + jobName + cloudbrain.ModelMountPath + "/"
- if err := mkModelPath(modelPath); err != nil {
- log.Error("Failed to mkModelPath: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.load_code_failed"), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- //init model readme
- if err := uploadCodeToMinio(modelPath, jobName, cloudbrain.ModelMountPath+"/"); err != nil {
- log.Error("Failed to uploadCodeToMinio: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.load_code_failed"), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- var datasetRemotePath, allFileName string
- for _, datasetInfo := range datasetInfos {
- if datasetRemotePath == "" {
- datasetRemotePath = datasetInfo.DataLocalPath
- allFileName = datasetInfo.FullName
- } else {
- datasetRemotePath = datasetRemotePath + ";" + datasetInfo.DataLocalPath
- allFileName = allFileName + ";" + datasetInfo.FullName
- }
-
- }
-
- //prepare command
- preTrainModelPath := getPreTrainModelPath(TrainUrl, CkptName)
-
- command, err := generateCommand(repo.Name, grampus.ProcessorTypeGPU, codeMinioPath+cloudbrain.DefaultBranchName+".zip", datasetRemotePath, BootFile, Params, setting.CBCodePathPrefix+jobName+cloudbrain.ModelMountPath+"/", allFileName, preTrainModelPath, CkptName)
- if err != nil {
- log.Error("Failed to generateCommand: %s (%v)", displayJobName, err, ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr("Create task failed, internal error", tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
-
- commitID, _ := ctx.Repo.GitRepo.GetBranchCommitID(cloudbrain.DefaultBranchName)
-
- req := &grampus.GenerateTrainJobReq{
- JobName: jobName,
- DisplayJobName: displayJobName,
- ComputeResource: models.GPUResource,
- ProcessType: grampus.ProcessorTypeGPU,
- Command: command,
- ImageUrl: image,
- Description: description,
- BootFile: BootFile,
- Uuid: uuid,
- CommitID: commitID,
- BranchName: cloudbrain.DefaultBranchName,
- Params: Params,
- EngineName: image,
- DatasetNames: datasetNames,
- DatasetInfos: datasetInfos,
-
- IsLatestVersion: modelarts.IsLatestVersion,
- VersionCount: modelarts.VersionCountOne,
- WorkServerNumber: 1,
- Spec: spec,
- ModelName: ModelName,
- LabelName: evaluationIndex,
- CkptName: CkptName,
- ModelVersion: ModelVersion,
- PreTrainModelUrl: TrainUrl,
- }
- err = grampus.GenerateTrainJob(ctx, req)
- if err != nil {
- log.Error("GenerateTrainJob failed:%v", err.Error(), ctx.Data["MsgID"])
- GrampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
- ctx.RenderWithErr(err.Error(), tplCloudBrainModelSafetyNewGrampusGpu, nil)
- return
- }
- }
-
- func createForGrampusNPU(ctx *context.Context, jobName string) {
-
- }
-
- func createForNPU(ctx *context.Context, jobName string) {
- VersionOutputPath := modelarts.GetOutputPathByCount(modelarts.TotalVersionCount)
- BootFile := ctx.Query("boot_file")
- displayJobName := ctx.Query("display_job_name")
- description := ctx.Query("description")
-
- srcDataset := ctx.Query("src_dataset") //uuid
- combatDataset := ctx.Query("combat_dataset") //uuid
- evaluationIndex := ctx.Query("evaluationIndex")
- Params := ctx.Query("run_para_list")
- specId := ctx.QueryInt64("spec_id")
-
- engineID := ctx.QueryInt("EngineID")
- poolID := ctx.Query("PoolID")
- repo := ctx.Repo.Repository
-
- trainUrl := ctx.Query("train_url")
- modelName := ctx.Query("ModelName")
- modelVersion := ctx.Query("ModelVersion")
- ckptName := ctx.Query("ckpt_name")
- ckptUrl := "/" + trainUrl + ckptName
- log.Info("ckpt url:" + ckptUrl)
-
- FlavorName := ctx.Query("FlavorName")
- EngineName := ctx.Query("EngineName")
-
- isLatestVersion := modelarts.IsLatestVersion
- VersionCount := modelarts.VersionCountOne
-
- codeLocalPath := setting.JobPath + jobName + modelarts.CodePath
- codeObsPath := "/" + setting.Bucket + modelarts.JobPath + jobName + modelarts.CodePath
- resultObsPath := "/" + setting.Bucket + modelarts.JobPath + jobName + modelarts.ResultPath + VersionOutputPath + "/"
- logObsPath := "/" + setting.Bucket + modelarts.JobPath + jobName + modelarts.LogPath + VersionOutputPath + "/"
- log.Info("ckpt url:" + ckptUrl)
- spec, err := resource.GetAndCheckSpec(ctx.User.ID, specId, models.FindSpecsOptions{
- JobType: models.JobTypeInference,
- ComputeResource: models.NPU,
- Cluster: models.OpenICluster,
- AiCenterCode: models.AICenterOfCloudBrainTwo})
- if err != nil || spec == nil {
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("Resource specification not available", tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
- if !account.IsPointBalanceEnough(ctx.User.ID, spec.UnitPrice) {
- log.Error("point balance is not enough,userId=%d specId=%d ", ctx.User.ID, spec.ID)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("points.insufficient_points_balance"), tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
-
- //todo: del the codeLocalPath
- _, err = ioutil.ReadDir(codeLocalPath)
- if err == nil {
- os.RemoveAll(codeLocalPath)
- }
-
- gitRepo, _ := git.OpenRepository(repo.RepoPath())
- commitID, _ := gitRepo.GetBranchCommitID(cloudbrain.DefaultBranchName)
-
- if err := downloadCode(repo, codeLocalPath, cloudbrain.DefaultBranchName); err != nil {
- log.Error("Create task failed, server timed out: %s (%v)", repo.FullName(), err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.load_code_failed"), tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
-
- //todo: upload code (send to file_server todo this work?)
- if err := obsMkdir(setting.CodePathPrefix + jobName + modelarts.ResultPath + VersionOutputPath + "/"); err != nil {
- log.Error("Failed to obsMkdir_result: %s (%v)", repo.FullName(), err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("Failed to obsMkdir_result", tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
-
- if err := obsMkdir(setting.CodePathPrefix + jobName + modelarts.LogPath + VersionOutputPath + "/"); err != nil {
- log.Error("Failed to obsMkdir_log: %s (%v)", repo.FullName(), err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("Failed to obsMkdir_log", tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
-
- if err := uploadCodeToObs(codeLocalPath, jobName, ""); err != nil {
- log.Error("Failed to uploadCodeToObs: %s (%v)", repo.FullName(), err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.load_code_failed"), tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
-
- var parameters models.Parameters
- param := make([]models.Parameter, 0)
- param = append(param, models.Parameter{
- Label: modelarts.ResultUrl,
- Value: "s3:/" + resultObsPath,
- }, models.Parameter{
- Label: modelarts.CkptUrl,
- Value: "s3:/" + ckptUrl,
- })
- uuid := srcDataset + ";" + combatDataset
- datasUrlList, dataUrl, datasetNames, isMultiDataset, err := getDatasUrlListByUUIDS(uuid)
- if err != nil {
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(err.Error(), tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
- dataPath := dataUrl
- jsondatas, err := json.Marshal(datasUrlList)
- if err != nil {
- log.Error("Failed to Marshal: %v", err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("json error:"+err.Error(), tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
- if isMultiDataset {
- param = append(param, models.Parameter{
- Label: modelarts.MultiDataUrl,
- Value: string(jsondatas),
- })
- }
-
- existDeviceTarget := false
- if len(Params) != 0 {
- err := json.Unmarshal([]byte(Params), ¶meters)
- if err != nil {
- log.Error("Failed to Unmarshal params: %s (%v)", Params, err)
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("运行参数错误", tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
-
- for _, parameter := range parameters.Parameter {
- if parameter.Label == modelarts.DeviceTarget {
- existDeviceTarget = true
- }
- if parameter.Label != modelarts.TrainUrl && parameter.Label != modelarts.DataUrl {
- param = append(param, models.Parameter{
- Label: parameter.Label,
- Value: parameter.Value,
- })
- }
- }
- }
- if !existDeviceTarget {
- param = append(param, models.Parameter{
- Label: modelarts.DeviceTarget,
- Value: modelarts.Ascend,
- })
- }
-
- req := &modelarts.GenerateInferenceJobReq{
- JobName: jobName,
- DisplayJobName: displayJobName,
- DataUrl: dataPath,
- Description: description,
- CodeObsPath: codeObsPath,
- BootFileUrl: codeObsPath + BootFile,
- BootFile: BootFile,
- TrainUrl: trainUrl,
- WorkServerNumber: 1,
- EngineID: int64(engineID),
- LogUrl: logObsPath,
- PoolID: poolID,
- Uuid: uuid,
- Parameters: param, //modelarts train parameters
- CommitID: commitID,
- BranchName: cloudbrain.DefaultBranchName,
- Params: Params,
- FlavorName: FlavorName,
- EngineName: EngineName,
- LabelName: evaluationIndex,
- IsLatestVersion: isLatestVersion,
- VersionCount: VersionCount,
- TotalVersionCount: modelarts.TotalVersionCount,
- ModelName: modelName,
- ModelVersion: modelVersion,
- CkptName: ckptName,
- ResultUrl: resultObsPath,
- Spec: spec,
- DatasetName: datasetNames,
- JobType: string(models.JobTypeModelSafety),
- }
-
- err = modelarts.GenerateInferenceJob(ctx, req)
- if err != nil {
- log.Error("GenerateTrainJob failed:%v", err.Error())
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(err.Error(), tplCloudBrainModelSafetyNewNpu, nil)
- return
- }
- }
-
- func createForGPU(ctx *context.Context, jobName string) {
- BootFile := ctx.Query("boot_file")
- displayJobName := ctx.Query("display_job_name")
- description := ctx.Query("description")
- image := strings.TrimSpace(ctx.Query("image"))
- srcDataset := ctx.Query("src_dataset") //uuid
- combatDataset := ctx.Query("combat_dataset") //uuid
- evaluationIndex := ctx.Query("evaluationIndex")
- Params := ctx.Query("run_para_list")
- specId := ctx.QueryInt64("spec_id")
- TrainUrl := ctx.Query("train_url")
- CkptName := ctx.Query("ckpt_name")
- ckptUrl := setting.Attachment.Minio.RealPath + TrainUrl + CkptName
- log.Info("ckpt url:" + ckptUrl)
- spec, err := resource.GetAndCheckSpec(ctx.User.ID, specId, models.FindSpecsOptions{
- JobType: models.JobTypeBenchmark,
- ComputeResource: models.GPU,
- Cluster: models.OpenICluster,
- AiCenterCode: models.AICenterOfCloudBrainOne})
- if err != nil || spec == nil {
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("Resource specification not available", tplCloudBrainModelSafetyNewGpu, nil)
- return
- }
-
- repo := ctx.Repo.Repository
- codePath := setting.JobPath + jobName + cloudbrain.CodeMountPath
- os.RemoveAll(codePath)
-
- if err := downloadCode(repo, codePath, cloudbrain.DefaultBranchName); err != nil {
- log.Error("downloadCode failed, %v", err, ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("system error", tplCloudBrainModelSafetyNewGpu, nil)
- return
- }
-
- err = uploadCodeToMinio(codePath+"/", jobName, cloudbrain.CodeMountPath+"/")
- if err != nil {
- log.Error("uploadCodeToMinio failed, %v", err, ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr("system error", tplCloudBrainModelSafetyNewGpu, nil)
- return
- }
-
- uuid := srcDataset + ";" + combatDataset
- datasetInfos, datasetNames, err := models.GetDatasetInfo(uuid)
- log.Info("uuid=" + uuid)
- if err != nil {
- log.Error("GetDatasetInfo failed: %v", err, ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(ctx.Tr("cloudbrain.error.dataset_select"), tplCloudBrainBenchmarkNew, nil)
- return
- }
- command, err := getGpuModelSafetyCommand(BootFile, Params, CkptName, displayJobName)
- if err != nil {
- log.Error("Get Command failed: %v", err, ctx.Data["MsgID"])
- modelSafetyNewDataPrepare(ctx)
- //ctx.RenderWithErr(ctx.Tr("cloudbrain.error.dataset_select"), tplCloudBrainBenchmarkNew, nil) TODO
- return
- }
- log.Info("Command=" + command)
-
- req := cloudbrain.GenerateCloudBrainTaskReq{
- Ctx: ctx,
- DisplayJobName: displayJobName,
- JobName: jobName,
- Image: image,
- Command: command,
- Uuids: uuid,
- DatasetNames: datasetNames,
- DatasetInfos: datasetInfos,
- CodePath: storage.GetMinioPath(jobName, cloudbrain.CodeMountPath+"/"),
- ModelPath: setting.Attachment.Minio.RealPath + TrainUrl,
- BenchmarkPath: storage.GetMinioPath(jobName, cloudbrain.BenchMarkMountPath+"/"),
- Snn4ImageNetPath: storage.GetMinioPath(jobName, cloudbrain.Snn4imagenetMountPath+"/"),
- BrainScorePath: storage.GetMinioPath(jobName, cloudbrain.BrainScoreMountPath+"/"),
- JobType: string(models.JobTypeModelSafety),
- Description: description,
- BranchName: cloudbrain.DefaultBranchName,
- BootFile: BootFile,
- Params: Params,
- CommitID: "",
- ResultPath: storage.GetMinioPath(jobName, cloudbrain.ResultPath+"/"),
- Spec: spec,
- LabelName: evaluationIndex,
- }
-
- err = cloudbrain.GenerateTask(req)
- if err != nil {
- modelSafetyNewDataPrepare(ctx)
- ctx.RenderWithErr(err.Error(), tplCloudBrainBenchmarkNew, nil)
- return
- }
- //ctx.Redirect(setting.AppSubURL + ctx.Repo.RepoLink + "/cloudbrain/modelsafety_test")
- }
-
- func getGpuModelSafetyCommand(BootFile string, params string, CkptName string, DisplayJobName string) (string, error) {
- var command string
- bootFile := strings.TrimSpace(BootFile)
-
- if !strings.HasSuffix(bootFile, ".py") {
- log.Error("bootFile(%s) format error", bootFile)
- return command, errors.New("bootFile format error")
- }
-
- var parameters models.Parameters
- var param string
- if len(params) != 0 {
- err := json.Unmarshal([]byte(params), ¶meters)
- if err != nil {
- log.Error("Failed to Unmarshal params: %s (%v)", params, err)
- return command, err
- }
-
- for _, parameter := range parameters.Parameter {
- param += " --" + parameter.Label + "=" + parameter.Value
- }
- }
-
- param += " --modelname" + "=" + CkptName
-
- command += "python /code/" + bootFile + param + " > " + cloudbrain.ResultPath + "/" + DisplayJobName + "-" + cloudbrain.LogFile
-
- return command, nil
- }
-
- func modelSafetyNewDataPrepare(ctx *context.Context) error {
- ctx.Data["PageIsCloudBrain"] = true
-
- ctx.Data["boot_file"] = ctx.Query("boot_file")
- ctx.Data["display_job_name"] = ctx.Query("display_job_name")
- ctx.Data["description"] = ctx.Query("description")
- ctx.Data["image"] = strings.TrimSpace(ctx.Query("image"))
- ctx.Data["src_dataset"] = ctx.Query("src_dataset") //uuid
- ctx.Data["combat_dataset"] = ctx.Query("combat_dataset") //uuid
- ctx.Data["evaluationIndex"] = ctx.Query("evaluationIndex")
- ctx.Data["run_para_list"] = ctx.Query("run_para_list")
- ctx.Data["spec_id"] = ctx.QueryInt64("spec_id")
- ctx.Data["train_url"] = ctx.Query("train_url")
- ctx.Data["ckpt_name"] = ctx.Query("ckpt_name")
-
- prepareCloudbrainOneSpecs(ctx)
-
- return nil
- }
-
- func getJsonContent(url string) (string, error) {
-
- resp, err := http.Get(url)
- if err != nil || resp.StatusCode != 200 {
- log.Info("Get organizations url error=" + err.Error())
- return "", err
- }
- bytes, err := ioutil.ReadAll(resp.Body)
- resp.Body.Close()
- if err != nil {
- log.Info("Get organizations url error=" + err.Error())
- return "", err
- }
- str := string(bytes)
- //log.Info("json str =" + str)
-
- return str, nil
- }
|