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