You can not select more than 25 topics Topics must start with a chinese character,a letter or number, can include dashes ('-') and can be up to 35 characters long.

grampus.go 28 kB

3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
2 years ago
3 years ago
2 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847
  1. package repo
  2. import (
  3. "code.gitea.io/gitea/modules/auth"
  4. "code.gitea.io/gitea/modules/git"
  5. "code.gitea.io/gitea/modules/grampus"
  6. "code.gitea.io/gitea/modules/modelarts"
  7. "code.gitea.io/gitea/modules/timeutil"
  8. "code.gitea.io/gitea/modules/util"
  9. "encoding/json"
  10. "errors"
  11. "github.com/unknwon/com"
  12. "io/ioutil"
  13. "net/http"
  14. "os"
  15. "path"
  16. "strconv"
  17. "strings"
  18. "time"
  19. "code.gitea.io/gitea/models"
  20. "code.gitea.io/gitea/modules/base"
  21. "code.gitea.io/gitea/modules/cloudbrain"
  22. "code.gitea.io/gitea/modules/context"
  23. "code.gitea.io/gitea/modules/log"
  24. "code.gitea.io/gitea/modules/setting"
  25. )
  26. const (
  27. tplGrampusTrainJobShow base.TplName = "repo/grampus/trainjob/show"
  28. //GPU
  29. tplGrampusTrainJobGPUNew base.TplName = "repo/grampus/trainjob/gpu/new"
  30. //NPU
  31. tplGrampusTrainJobNPUNew base.TplName = "repo/grampus/trainjob/npu/new"
  32. )
  33. func GrampusTrainJobGPUNew(ctx *context.Context) {
  34. ctx.Data["datasetType"] = models.TypeCloudBrainOne
  35. err := grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  36. if err != nil {
  37. ctx.ServerError("get new train-job info failed", err)
  38. return
  39. }
  40. waitCount := cloudbrain.GetWaitingCloudbrainCount(models.TypeC2Net, models.GPUResource, models.JobTypeTrain)
  41. ctx.Data["WaitCount"] = waitCount
  42. ctx.HTML(http.StatusOK, tplGrampusTrainJobGPUNew)
  43. }
  44. func GrampusTrainJobNPUNew(ctx *context.Context) {
  45. ctx.Data["datasetType"] = models.TypeCloudBrainTwo
  46. err := grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  47. if err != nil {
  48. ctx.ServerError("get new train-job info failed", err)
  49. return
  50. }
  51. waitCount := cloudbrain.GetWaitingCloudbrainCount(models.TypeC2Net, models.NPUResource, models.JobTypeTrain)
  52. ctx.Data["WaitCount"] = waitCount
  53. ctx.HTML(200, tplGrampusTrainJobNPUNew)
  54. }
  55. func grampusTrainJobNewDataPrepare(ctx *context.Context, processType string) error {
  56. ctx.Data["PageIsCloudBrain"] = true
  57. t := time.Now()
  58. var displayJobName = cutString(ctx.User.Name, 5) + t.Format("2006010215") + strconv.Itoa(int(t.Unix()))[5:]
  59. ctx.Data["display_job_name"] = displayJobName
  60. //get valid images
  61. images, err := grampus.GetImages(processType)
  62. if err != nil {
  63. log.Error("GetImages failed:", err.Error())
  64. } else {
  65. ctx.Data["images"] = images.Infos
  66. }
  67. grampus.InitSpecialPool()
  68. ctx.Data["GPUEnabled"] = true
  69. ctx.Data["NPUEnabled"] = true
  70. includeCenters := make(map[string]struct{})
  71. excludeCenters := make(map[string]struct{})
  72. if grampus.SpecialPools != nil {
  73. for _, pool := range grampus.SpecialPools.Pools {
  74. if pool.IsExclusive {
  75. if !IsUserInOrgPool(ctx.User.ID, pool) {
  76. ctx.Data[pool.Type+"Enabled"] = false
  77. }
  78. } else {
  79. if strings.Contains(strings.ToLower(processType), strings.ToLower(pool.Type)) {
  80. if IsUserInOrgPool(ctx.User.ID, pool) {
  81. for _, center := range pool.Pool {
  82. includeCenters[center.Queue] = struct{}{}
  83. }
  84. } else {
  85. for _, center := range pool.Pool {
  86. excludeCenters[center.Queue] = struct{}{}
  87. }
  88. }
  89. }
  90. }
  91. }
  92. }
  93. //get valid resource specs
  94. specs, err := grampus.GetResourceSpecs(processType)
  95. grampusSpecs := getFilterSpecBySpecialPool(specs, includeCenters, excludeCenters)
  96. if err != nil {
  97. log.Error("GetResourceSpecs failed:", err.Error())
  98. } else {
  99. ctx.Data["flavor_infos"] = grampusSpecs
  100. }
  101. //get branches
  102. branches, _, err := ctx.Repo.GitRepo.GetBranches(0, 0)
  103. if err != nil {
  104. log.Error("GetBranches error:", err.Error())
  105. } else {
  106. ctx.Data["branches"] = branches
  107. }
  108. ctx.Data["branchName"] = ctx.Repo.BranchName
  109. if processType == grampus.ProcessorTypeGPU {
  110. ctx.Data["datasetType"] = models.TypeCloudBrainOne
  111. } else if processType == grampus.ProcessorTypeNPU {
  112. ctx.Data["datasetType"] = models.TypeCloudBrainTwo
  113. }
  114. return nil
  115. }
  116. func getFilterSpecBySpecialPool(specs *models.GetGrampusResourceSpecsResult, includeCenters map[string]struct{}, excludeCenters map[string]struct{}) []models.GrampusSpec {
  117. if len(includeCenters) == 0 && len(excludeCenters) == 0 {
  118. return specs.Infos
  119. }
  120. var grampusSpecs []models.GrampusSpec
  121. for _, info := range specs.Infos {
  122. if isInIncludeCenters(info, includeCenters) || (len(excludeCenters) != 0 && isNotAllInExcludeCenters(info, excludeCenters)) {
  123. grampusSpecs = append(grampusSpecs, info)
  124. }
  125. }
  126. return grampusSpecs
  127. }
  128. func isInIncludeCenters(grampusSpec models.GrampusSpec, centers map[string]struct{}) bool {
  129. for _, center := range grampusSpec.Centers {
  130. if _, ok := centers[center.ID]; ok {
  131. return true
  132. }
  133. }
  134. return false
  135. }
  136. func isNotAllInExcludeCenters(grampusSpec models.GrampusSpec, centers map[string]struct{}) bool {
  137. for _, center := range grampusSpec.Centers {
  138. if _, ok := centers[center.ID]; !ok {
  139. return true
  140. }
  141. }
  142. return false
  143. }
  144. func IsUserInOrgPool(userId int64, pool *models.SpecialPool) bool {
  145. org, _ := models.GetOrgByName(pool.Org)
  146. if org != nil {
  147. isOrgMember, _ := models.IsOrganizationMember(org.ID, userId)
  148. return isOrgMember
  149. }
  150. return false
  151. }
  152. func grampusParamCheckCreateTrainJob(form auth.CreateGrampusTrainJobForm) error {
  153. if !strings.HasSuffix(strings.TrimSpace(form.BootFile), ".py") {
  154. log.Error("the boot file(%s) must be a python file", form.BootFile)
  155. return errors.New("启动文件必须是python文件")
  156. }
  157. if form.BranchName == "" {
  158. log.Error("the branch must not be null!", form.BranchName)
  159. return errors.New("代码分支不能为空!")
  160. }
  161. return nil
  162. }
  163. func GrampusTrainJobGpuCreate(ctx *context.Context, form auth.CreateGrampusTrainJobForm) {
  164. displayJobName := form.DisplayJobName
  165. jobName := util.ConvertDisplayJobNameToJobName(displayJobName)
  166. uuid := form.Attachment
  167. description := form.Description
  168. bootFile := strings.TrimSpace(form.BootFile)
  169. params := form.Params
  170. repo := ctx.Repo.Repository
  171. codeLocalPath := setting.JobPath + jobName + cloudbrain.CodeMountPath + "/"
  172. codeMinioPath := setting.CBCodePathPrefix + jobName + cloudbrain.CodeMountPath + "/"
  173. dataMinioPath := setting.Attachment.Minio.BasePath + path.Join(uuid[0:1], uuid[1:2]) + "/" + uuid
  174. branchName := form.BranchName
  175. flavorName := form.FlavorName
  176. image := strings.TrimSpace(form.Image)
  177. if !jobNamePattern.MatchString(displayJobName) {
  178. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  179. ctx.RenderWithErr(ctx.Tr("repo.cloudbrain_jobname_err"), tplGrampusTrainJobGPUNew, &form)
  180. return
  181. }
  182. errStr := checkSpecialPool(ctx, "GPU")
  183. if errStr != "" {
  184. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  185. ctx.RenderWithErr(errStr, tplGrampusTrainJobGPUNew, &form)
  186. return
  187. }
  188. //check count limit
  189. count, err := models.GetGrampusCountByUserID(ctx.User.ID, string(models.JobTypeTrain), models.GPUResource)
  190. if err != nil {
  191. log.Error("GetGrampusCountByUserID failed:%v", err, ctx.Data["MsgID"])
  192. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  193. ctx.RenderWithErr("system error", tplGrampusTrainJobGPUNew, &form)
  194. return
  195. } else {
  196. if count >= 1 {
  197. log.Error("the user already has running or waiting task", ctx.Data["MsgID"])
  198. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  199. ctx.RenderWithErr("you have already a running or waiting task, can not create more", tplGrampusTrainJobGPUNew, &form)
  200. return
  201. }
  202. }
  203. //check param
  204. if err := grampusParamCheckCreateTrainJob(form); err != nil {
  205. log.Error("paramCheckCreateTrainJob failed:(%v)", err, ctx.Data["MsgID"])
  206. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  207. ctx.RenderWithErr(err.Error(), tplGrampusTrainJobGPUNew, &form)
  208. return
  209. }
  210. //check whether the task name in the project is duplicated
  211. tasks, err := models.GetCloudbrainsByDisplayJobName(repo.ID, string(models.JobTypeTrain), displayJobName)
  212. if err == nil {
  213. if len(tasks) != 0 {
  214. log.Error("the job name did already exist", ctx.Data["MsgID"])
  215. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  216. ctx.RenderWithErr("the job name did already exist", tplGrampusTrainJobGPUNew, &form)
  217. return
  218. }
  219. } else {
  220. if !models.IsErrJobNotExist(err) {
  221. log.Error("system error, %v", err, ctx.Data["MsgID"])
  222. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  223. ctx.RenderWithErr("system error", tplGrampusTrainJobGPUNew, &form)
  224. return
  225. }
  226. }
  227. //check dataset
  228. attachment, err := models.GetAttachmentByUUID(uuid)
  229. if err != nil {
  230. log.Error("GetAttachmentByUUID failed:", err.Error(), ctx.Data["MsgID"])
  231. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  232. ctx.RenderWithErr("dataset is not exist", tplGrampusTrainJobGPUNew, &form)
  233. return
  234. }
  235. //prepare code and out path
  236. _, err = ioutil.ReadDir(codeLocalPath)
  237. if err == nil {
  238. os.RemoveAll(codeLocalPath)
  239. }
  240. if err := downloadZipCode(ctx, codeLocalPath, branchName); err != nil {
  241. log.Error("downloadZipCode failed, server timed out: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
  242. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  243. ctx.RenderWithErr("Create task failed, internal error", tplGrampusTrainJobGPUNew, &form)
  244. return
  245. }
  246. //todo: upload code (send to file_server todo this work?)
  247. //upload code
  248. if err := uploadCodeToMinio(codeLocalPath+"/", jobName, cloudbrain.CodeMountPath+"/"); err != nil {
  249. log.Error("Failed to uploadCodeToMinio: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
  250. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  251. ctx.RenderWithErr("Create task failed, internal error", tplGrampusTrainJobGPUNew, &form)
  252. return
  253. }
  254. modelPath := setting.JobPath + jobName + cloudbrain.ModelMountPath + "/"
  255. if err := mkModelPath(modelPath); err != nil {
  256. log.Error("Failed to mkModelPath: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
  257. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  258. ctx.RenderWithErr("Create task failed, internal error", tplGrampusTrainJobGPUNew, &form)
  259. return
  260. }
  261. //init model readme
  262. if err := uploadCodeToMinio(modelPath, jobName, cloudbrain.ModelMountPath+"/"); err != nil {
  263. log.Error("Failed to uploadCodeToMinio: %s (%v)", repo.FullName(), err, ctx.Data["MsgID"])
  264. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  265. ctx.RenderWithErr("Create task failed, internal error", tplGrampusTrainJobGPUNew, &form)
  266. return
  267. }
  268. //prepare command
  269. command, err := generateCommand(repo.Name, grampus.ProcessorTypeGPU, codeMinioPath+cloudbrain.DefaultBranchName+".zip", dataMinioPath, bootFile, params, setting.CBCodePathPrefix+jobName+cloudbrain.ModelMountPath+"/", attachment.Name)
  270. if err != nil {
  271. log.Error("Failed to generateCommand: %s (%v)", displayJobName, err, ctx.Data["MsgID"])
  272. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  273. ctx.RenderWithErr("Create task failed, internal error", tplGrampusTrainJobGPUNew, &form)
  274. return
  275. }
  276. commitID, _ := ctx.Repo.GitRepo.GetBranchCommitID(branchName)
  277. req := &grampus.GenerateTrainJobReq{
  278. JobName: jobName,
  279. DisplayJobName: displayJobName,
  280. ComputeResource: models.GPUResource,
  281. ProcessType: grampus.ProcessorTypeGPU,
  282. Command: command,
  283. ResourceSpecId: form.FlavorID,
  284. ImageUrl: image,
  285. Description: description,
  286. BootFile: bootFile,
  287. Uuid: uuid,
  288. CommitID: commitID,
  289. BranchName: branchName,
  290. Params: form.Params,
  291. FlavorName: flavorName,
  292. EngineName: image,
  293. DatasetName: attachment.Name,
  294. IsLatestVersion: modelarts.IsLatestVersion,
  295. VersionCount: modelarts.VersionCount,
  296. WorkServerNumber: 1,
  297. }
  298. err = grampus.GenerateTrainJob(ctx, req)
  299. if err != nil {
  300. log.Error("GenerateTrainJob failed:%v", err.Error(), ctx.Data["MsgID"])
  301. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeGPU)
  302. ctx.RenderWithErr(err.Error(), tplGrampusTrainJobGPUNew, &form)
  303. return
  304. }
  305. ctx.Redirect(setting.AppSubURL + ctx.Repo.RepoLink + "/modelarts/train-job")
  306. }
  307. func checkSpecialPool(ctx *context.Context, resourceType string) string {
  308. grampus.InitSpecialPool()
  309. if grampus.SpecialPools != nil {
  310. for _, pool := range grampus.SpecialPools.Pools {
  311. if pool.IsExclusive && pool.Type == resourceType {
  312. org, _ := models.GetOrgByName(pool.Org)
  313. if org != nil {
  314. isOrgMember, _ := models.IsOrganizationMember(org.ID, ctx.User.ID)
  315. if !isOrgMember {
  316. return ctx.Tr("repo.grampus.no_operate_right")
  317. }
  318. }
  319. }
  320. }
  321. }
  322. return ""
  323. }
  324. func GrampusTrainJobNpuCreate(ctx *context.Context, form auth.CreateGrampusTrainJobForm) {
  325. displayJobName := form.DisplayJobName
  326. jobName := util.ConvertDisplayJobNameToJobName(displayJobName)
  327. uuid := form.Attachment
  328. description := form.Description
  329. bootFile := strings.TrimSpace(form.BootFile)
  330. params := form.Params
  331. repo := ctx.Repo.Repository
  332. codeLocalPath := setting.JobPath + jobName + modelarts.CodePath
  333. codeObsPath := grampus.JobPath + jobName + modelarts.CodePath
  334. dataObsPath := setting.BasePath + path.Join(uuid[0:1], uuid[1:2]) + "/" + uuid + "/"
  335. branchName := form.BranchName
  336. isLatestVersion := modelarts.IsLatestVersion
  337. flavorName := form.FlavorName
  338. versionCount := modelarts.VersionCount
  339. engineName := form.EngineName
  340. if !jobNamePattern.MatchString(displayJobName) {
  341. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  342. ctx.RenderWithErr(ctx.Tr("repo.cloudbrain_jobname_err"), tplGrampusTrainJobNPUNew, &form)
  343. return
  344. }
  345. errStr := checkSpecialPool(ctx, "NPU")
  346. if errStr != "" {
  347. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  348. ctx.RenderWithErr(errStr, tplGrampusTrainJobGPUNew, &form)
  349. return
  350. }
  351. //check count limit
  352. count, err := models.GetGrampusCountByUserID(ctx.User.ID, string(models.JobTypeTrain), models.NPUResource)
  353. if err != nil {
  354. log.Error("GetGrampusCountByUserID failed:%v", err, ctx.Data["MsgID"])
  355. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  356. ctx.RenderWithErr("system error", tplGrampusTrainJobNPUNew, &form)
  357. return
  358. } else {
  359. if count >= 1 {
  360. log.Error("the user already has running or waiting task", ctx.Data["MsgID"])
  361. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  362. ctx.RenderWithErr("you have already a running or waiting task, can not create more", tplGrampusTrainJobNPUNew, &form)
  363. return
  364. }
  365. }
  366. //check param
  367. if err := grampusParamCheckCreateTrainJob(form); err != nil {
  368. log.Error("paramCheckCreateTrainJob failed:(%v)", err)
  369. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  370. ctx.RenderWithErr(err.Error(), tplGrampusTrainJobNPUNew, &form)
  371. return
  372. }
  373. //check whether the task name in the project is duplicated
  374. tasks, err := models.GetCloudbrainsByDisplayJobName(repo.ID, string(models.JobTypeTrain), displayJobName)
  375. if err == nil {
  376. if len(tasks) != 0 {
  377. log.Error("the job name did already exist", ctx.Data["MsgID"])
  378. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  379. ctx.RenderWithErr("the job name did already exist", tplGrampusTrainJobNPUNew, &form)
  380. return
  381. }
  382. } else {
  383. if !models.IsErrJobNotExist(err) {
  384. log.Error("system error, %v", err, ctx.Data["MsgID"])
  385. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  386. ctx.RenderWithErr("system error", tplGrampusTrainJobNPUNew, &form)
  387. return
  388. }
  389. }
  390. //check dataset
  391. attachment, err := models.GetAttachmentByUUID(uuid)
  392. if err != nil {
  393. log.Error("GetAttachmentByUUID failed:", err.Error(), ctx.Data["MsgID"])
  394. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  395. ctx.RenderWithErr("dataset is not exist", tplGrampusTrainJobNPUNew, &form)
  396. return
  397. }
  398. //prepare code and out path
  399. _, err = ioutil.ReadDir(codeLocalPath)
  400. if err == nil {
  401. os.RemoveAll(codeLocalPath)
  402. }
  403. if err := downloadZipCode(ctx, codeLocalPath, branchName); err != nil {
  404. log.Error("downloadZipCode failed, server timed out: %s (%v)", repo.FullName(), err)
  405. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  406. ctx.RenderWithErr("Create task failed, server timed out", tplGrampusTrainJobNPUNew, &form)
  407. return
  408. }
  409. //todo: upload code (send to file_server todo this work?)
  410. if err := obsMkdir(setting.CodePathPrefix + jobName + modelarts.OutputPath); err != nil {
  411. log.Error("Failed to obsMkdir_output: %s (%v)", repo.FullName(), err)
  412. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  413. ctx.RenderWithErr("Failed to obsMkdir_output", tplGrampusTrainJobNPUNew, &form)
  414. return
  415. }
  416. if err := uploadCodeToObs(codeLocalPath, jobName, ""); err != nil {
  417. log.Error("Failed to uploadCodeToObs: %s (%v)", repo.FullName(), err)
  418. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  419. ctx.RenderWithErr("Failed to uploadCodeToObs", tplGrampusTrainJobNPUNew, &form)
  420. return
  421. }
  422. //prepare command
  423. command, err := generateCommand(repo.Name, grampus.ProcessorTypeNPU, codeObsPath+cloudbrain.DefaultBranchName+".zip", dataObsPath+"'"+attachment.Name+"'", bootFile, params, setting.CodePathPrefix+jobName+modelarts.OutputPath, attachment.Name)
  424. if err != nil {
  425. log.Error("Failed to generateCommand: %s (%v)", displayJobName, err, ctx.Data["MsgID"])
  426. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  427. ctx.RenderWithErr("Create task failed, internal error", tplGrampusTrainJobNPUNew, &form)
  428. return
  429. }
  430. commitID, _ := ctx.Repo.GitRepo.GetBranchCommitID(branchName)
  431. req := &grampus.GenerateTrainJobReq{
  432. JobName: jobName,
  433. DisplayJobName: displayJobName,
  434. ComputeResource: models.NPUResource,
  435. ProcessType: grampus.ProcessorTypeNPU,
  436. Command: command,
  437. ResourceSpecId: form.FlavorID,
  438. ImageId: form.ImageID,
  439. DataUrl: dataObsPath,
  440. Description: description,
  441. CodeObsPath: codeObsPath,
  442. BootFileUrl: codeObsPath + bootFile,
  443. BootFile: bootFile,
  444. WorkServerNumber: form.WorkServerNumber,
  445. Uuid: uuid,
  446. CommitID: commitID,
  447. IsLatestVersion: isLatestVersion,
  448. BranchName: branchName,
  449. Params: form.Params,
  450. FlavorName: flavorName,
  451. EngineName: engineName,
  452. VersionCount: versionCount,
  453. TotalVersionCount: modelarts.TotalVersionCount,
  454. DatasetName: attachment.Name,
  455. }
  456. err = grampus.GenerateTrainJob(ctx, req)
  457. if err != nil {
  458. log.Error("GenerateTrainJob failed:%v", err.Error())
  459. grampusTrainJobNewDataPrepare(ctx, grampus.ProcessorTypeNPU)
  460. ctx.RenderWithErr(err.Error(), tplGrampusTrainJobNPUNew, &form)
  461. return
  462. }
  463. ctx.Redirect(setting.AppSubURL + ctx.Repo.RepoLink + "/modelarts/train-job")
  464. }
  465. func GrampusStopJob(ctx *context.Context) {
  466. var ID = ctx.Params(":jobid")
  467. var resultCode = "0"
  468. var errorMsg = ""
  469. var status = ""
  470. task := ctx.Cloudbrain
  471. for {
  472. if task.Status == string(models.GrampusStatusStopped) || task.Status == string(models.GrampusStatusFailed) || task.Status == string(models.GrampusStatusSucceeded) {
  473. log.Error("the job(%s) has been stopped", task.JobName, ctx.Data["msgID"])
  474. resultCode = "-1"
  475. errorMsg = "system error"
  476. break
  477. }
  478. res, err := grampus.StopJob(task.JobID)
  479. if err != nil {
  480. log.Error("StopJob(%s) failed:%v", task.JobName, err, ctx.Data["msgID"])
  481. resultCode = strconv.Itoa(res.ErrorCode)
  482. errorMsg = res.ErrorMsg
  483. break
  484. }
  485. task.Status = string(models.GrampusStatusStopped)
  486. if task.EndTime == 0 {
  487. task.EndTime = timeutil.TimeStampNow()
  488. }
  489. task.ComputeAndSetDuration()
  490. err = models.UpdateJob(task)
  491. if err != nil {
  492. log.Error("UpdateJob(%s) failed:%v", task.JobName, err, ctx.Data["msgID"])
  493. resultCode = "-1"
  494. errorMsg = "system error"
  495. break
  496. }
  497. status = task.Status
  498. break
  499. }
  500. ctx.JSON(200, map[string]interface{}{
  501. "result_code": resultCode,
  502. "error_msg": errorMsg,
  503. "status": status,
  504. "id": ID,
  505. "StatusOK": 0,
  506. })
  507. }
  508. func GrampusTrainJobDel(ctx *context.Context) {
  509. var listType = ctx.Query("listType")
  510. if err := deleteGrampusJob(ctx); err != nil {
  511. log.Error("deleteGrampusJob failed: %v", err, ctx.Data["msgID"])
  512. ctx.ServerError(err.Error(), err)
  513. return
  514. }
  515. var isAdminPage = ctx.Query("isadminpage")
  516. var isHomePage = ctx.Query("ishomepage")
  517. if ctx.IsUserSiteAdmin() && isAdminPage == "true" {
  518. ctx.Redirect(setting.AppSubURL + "/admin" + "/cloudbrains")
  519. } else if isHomePage == "true" {
  520. ctx.Redirect(setting.AppSubURL + "/cloudbrains")
  521. } else {
  522. ctx.Redirect(setting.AppSubURL + ctx.Repo.RepoLink + "/modelarts/train-job?listType=" + listType)
  523. }
  524. }
  525. func deleteGrampusJob(ctx *context.Context) error {
  526. task := ctx.Cloudbrain
  527. if task.Status != string(models.GrampusStatusStopped) && task.Status != string(models.GrampusStatusSucceeded) && task.Status != string(models.GrampusStatusFailed) {
  528. log.Error("the job(%s) has not been stopped", task.JobName, ctx.Data["msgID"])
  529. return errors.New("the job has not been stopped")
  530. }
  531. err := models.DeleteJob(task)
  532. if err != nil {
  533. log.Error("DeleteJob failed: %v", err, ctx.Data["msgID"])
  534. return err
  535. }
  536. storageType := models.TypeCloudBrainOne
  537. if task.ComputeResource == models.NPUResource {
  538. storageType = models.TypeCloudBrainTwo
  539. }
  540. deleteJobStorage(task.JobName, storageType)
  541. return nil
  542. }
  543. func GrampusTrainJobShow(ctx *context.Context) {
  544. ctx.Data["PageIsCloudBrain"] = true
  545. var task *models.Cloudbrain
  546. task, err := models.GetCloudbrainByJobIDWithDeleted(ctx.Params(":jobid"))
  547. if err != nil {
  548. log.Error("GetCloudbrainByJobID failed:" + err.Error())
  549. ctx.NotFound(ctx.Req.URL.RequestURI(), nil)
  550. return
  551. }
  552. if task.DeletedAt.IsZero() { //normal record
  553. result, err := grampus.GetJob(task.JobID)
  554. if err != nil {
  555. log.Error("GetJob failed:" + err.Error())
  556. ctx.NotFound(ctx.Req.URL.RequestURI(), nil)
  557. return
  558. }
  559. if result != nil {
  560. if len(result.JobInfo.Tasks[0].CenterID) == 1 && len(result.JobInfo.Tasks[0].CenterName) == 1 {
  561. task.AiCenter = result.JobInfo.Tasks[0].CenterID[0] + "+" + result.JobInfo.Tasks[0].CenterName[0]
  562. }
  563. task.Status = grampus.TransTrainJobStatus(result.JobInfo.Status)
  564. if task.Status != result.JobInfo.Status || result.JobInfo.Status == models.GrampusStatusRunning {
  565. task.Duration = result.JobInfo.RunSec
  566. task.TrainJobDuration = models.ConvertDurationToStr(task.Duration)
  567. if task.StartTime == 0 && result.JobInfo.StartedAt > 0 {
  568. task.StartTime = timeutil.TimeStamp(result.JobInfo.StartedAt)
  569. }
  570. if task.EndTime == 0 && models.IsTrainJobTerminal(task.Status) && task.StartTime > 0 {
  571. task.EndTime = task.StartTime.Add(task.Duration)
  572. }
  573. task.CorrectCreateUnix()
  574. err = models.UpdateJob(task)
  575. if err != nil {
  576. log.Error("UpdateJob failed:" + err.Error())
  577. }
  578. }
  579. }
  580. }
  581. if len(task.Parameters) > 0 {
  582. var parameters models.Parameters
  583. err := json.Unmarshal([]byte(task.Parameters), &parameters)
  584. if err != nil {
  585. log.Error("Failed to Unmarshal Parameters: %s (%v)", task.Parameters, err)
  586. ctx.ServerError("system error", err)
  587. return
  588. }
  589. if len(parameters.Parameter) > 0 {
  590. paramTemp := ""
  591. for _, Parameter := range parameters.Parameter {
  592. param := Parameter.Label + " = " + Parameter.Value + "; "
  593. paramTemp = paramTemp + param
  594. }
  595. task.Parameters = paramTemp[:len(paramTemp)-2]
  596. } else {
  597. task.Parameters = ""
  598. }
  599. }
  600. taskList := make([]*models.Cloudbrain, 0)
  601. taskList = append(taskList, task)
  602. ctx.Data["version_list_task"] = taskList
  603. ctx.Data["canDownload"] = cloudbrain.CanModifyJob(ctx, task)
  604. ctx.Data["displayJobName"] = task.DisplayJobName
  605. aiCenterInfo := strings.Split(task.AiCenter, "+")
  606. if len(aiCenterInfo) == 2 {
  607. ctx.Data["ai_center"] = aiCenterInfo[1]
  608. }
  609. ctx.HTML(http.StatusOK, tplGrampusTrainJobShow)
  610. }
  611. func GrampusGetLog(ctx *context.Context) {
  612. jobID := ctx.Params(":jobid")
  613. job, err := models.GetCloudbrainByJobID(jobID)
  614. if err != nil {
  615. log.Error("GetCloudbrainByJobID failed: %v", err, ctx.Data["MsgID"])
  616. ctx.ServerError(err.Error(), err)
  617. return
  618. }
  619. content, err := grampus.GetTrainJobLog(job.JobID)
  620. if err != nil {
  621. log.Error("GetTrainJobLog failed: %v", err, ctx.Data["MsgID"])
  622. ctx.ServerError(err.Error(), err)
  623. return
  624. }
  625. ctx.JSON(http.StatusOK, map[string]interface{}{
  626. "JobName": job.JobName,
  627. "Content": content,
  628. })
  629. return
  630. }
  631. func generateCommand(repoName, processorType, codeRemotePath, dataRemotePath, bootFile, paramSrc, outputRemotePath, datasetName string) (string, error) {
  632. var command string
  633. workDir := grampus.NpuWorkDir
  634. if processorType == grampus.ProcessorTypeGPU {
  635. workDir = grampus.GpuWorkDir
  636. }
  637. command += "pwd;cd " + workDir + grampus.CommandPrepareScript
  638. //download code & dataset
  639. if processorType == grampus.ProcessorTypeNPU {
  640. commandDownload := "./downloader_for_obs " + setting.Bucket + " " + codeRemotePath + " " + grampus.CodeArchiveName + " " + dataRemotePath + " '" + datasetName + "';"
  641. command += commandDownload
  642. } else if processorType == grampus.ProcessorTypeGPU {
  643. commandDownload := "./downloader_for_minio " + setting.Grampus.Env + " " + codeRemotePath + " " + grampus.CodeArchiveName + " " + dataRemotePath + " '" + datasetName + "';"
  644. command += commandDownload
  645. }
  646. //check download result
  647. commandCheckRes := "bash -c \"[[ $? -eq 0 ]] && exit 0 || exit -1;\";"
  648. command += commandCheckRes
  649. //unzip code & dataset
  650. toolUnzip := "unzip -q '"
  651. if strings.HasSuffix(datasetName, ".tar.gz") {
  652. toolUnzip = "tar -zxvf '"
  653. }
  654. commandUnzip := "cd " + workDir + "code;unzip -q master.zip;echo \"start to unzip dataset\";cd " + workDir + "dataset;" + toolUnzip + datasetName + "';"
  655. command += commandUnzip
  656. //check unzip result
  657. commandCheckRes = "bash -c \"[[ $? -eq 0 ]] && exit 0 || exit -1;\";"
  658. command += commandCheckRes
  659. command += "echo \"unzip finished;start to exec code;\";"
  660. //exec code
  661. var parameters models.Parameters
  662. var paramCode string
  663. param := make([]models.Parameter, 0)
  664. if len(paramSrc) != 0 {
  665. err := json.Unmarshal([]byte(paramSrc), &parameters)
  666. if err != nil {
  667. log.Error("Failed to Unmarshal params: %s (%v)", paramSrc, err)
  668. return command, err
  669. }
  670. for _, parameter := range parameters.Parameter {
  671. param = append(param, models.Parameter{
  672. Label: parameter.Label,
  673. Value: parameter.Value,
  674. })
  675. paramCode += " --" + parameter.Label + "=" + parameter.Value
  676. }
  677. }
  678. commandCode := "cd " + workDir + "code/" + strings.ToLower(repoName) + ";python " + bootFile + paramCode + ";"
  679. command += commandCode
  680. //get exec result
  681. commandGetRes := "result=$?;"
  682. command += commandGetRes
  683. //upload models
  684. if processorType == grampus.ProcessorTypeNPU {
  685. commandUpload := "cd " + workDir + "script_for_grampus/;./uploader_for_obs " + setting.Bucket + " " + outputRemotePath + " " + workDir + "output/;"
  686. command += commandUpload
  687. } else if processorType == grampus.ProcessorTypeGPU {
  688. commandUpload := "cd " + workDir + "script_for_grampus/;./uploader_for_minio " + setting.Grampus.Env + " " + outputRemotePath + " " + workDir + "output/;"
  689. command += commandUpload
  690. }
  691. //check exec result
  692. commandCheckRes = "bash -c \"[[ $result -eq 0 ]] && exit 0 || exit -1\""
  693. command += commandCheckRes
  694. return command, nil
  695. }
  696. func downloadZipCode(ctx *context.Context, codePath, branchName string) error {
  697. archiveType := git.ZIP
  698. archivePath := codePath
  699. if !com.IsDir(archivePath) {
  700. if err := os.MkdirAll(archivePath, os.ModePerm); err != nil {
  701. log.Error("MkdirAll failed:" + err.Error())
  702. return err
  703. }
  704. }
  705. // Get corresponding commit.
  706. var (
  707. commit *git.Commit
  708. err error
  709. )
  710. gitRepo := ctx.Repo.GitRepo
  711. if err != nil {
  712. log.Error("OpenRepository failed:" + err.Error())
  713. return err
  714. }
  715. if gitRepo.IsBranchExist(branchName) {
  716. commit, err = gitRepo.GetBranchCommit(branchName)
  717. if err != nil {
  718. log.Error("GetBranchCommit failed:" + err.Error())
  719. return err
  720. }
  721. }
  722. archivePath = path.Join(archivePath, grampus.CodeArchiveName)
  723. if !com.IsFile(archivePath) {
  724. if err := commit.CreateArchive(archivePath, git.CreateArchiveOpts{
  725. Format: archiveType,
  726. Prefix: setting.Repository.PrefixArchiveFiles,
  727. }); err != nil {
  728. log.Error("CreateArchive failed:" + err.Error())
  729. return err
  730. }
  731. }
  732. return nil
  733. }