cluster.jobs.go 7.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308
  1. package main
  2. /*
  3. Cluster job execution on this node.
  4. The jobs package decides what to run and where; this file is the bridge
  5. to the AGI runtime: it resolves the owner, injects the job helper
  6. functions into the VM and captures the script's output.
  7. */
  8. import (
  9. "context"
  10. "encoding/json"
  11. "errors"
  12. "net/http"
  13. "net/http/httptest"
  14. "strings"
  15. "sync"
  16. "time"
  17. "github.com/robertkrimen/otto"
  18. "imuslab.com/arozos/mod/agi"
  19. "imuslab.com/arozos/mod/cluster/jobs"
  20. "imuslab.com/arozos/mod/cluster/metadata"
  21. "imuslab.com/arozos/mod/utils"
  22. )
  23. // clusterJobExecutor runs job scripts through the AGI gateway.
  24. type clusterJobExecutor struct{}
  25. func (e *clusterJobExecutor) UserExists(owner string) bool {
  26. if authAgent == nil {
  27. return false
  28. }
  29. return authAgent.UserExists(owner)
  30. }
  31. // Run executes the wrapped job source as the job owner.
  32. func (e *clusterJobExecutor) Run(ctx context.Context, rec jobs.Record, source string, hooks jobs.ExecHooks) (json.RawMessage, error) {
  33. if AGIGateway == nil {
  34. return nil, errors.New("AGI gateway not ready")
  35. }
  36. u, err := userHandler.GetUserInfoFromUsername(rec.Spec.Owner)
  37. if err != nil {
  38. return nil, err
  39. }
  40. jobCtx := jobs.ContextFor(rec, clusterManager.NodeID())
  41. ctxJSON, _ := json.Marshal(jobCtx)
  42. var mu sync.Mutex
  43. var output json.RawMessage
  44. var stopper *agi.JobStopper
  45. //The AGI runtime needs a request object for its serverless helpers
  46. req := httptest.NewRequest(http.MethodPost, "/system/cluster/jobs/exec", strings.NewReader(""))
  47. rw := httptest.NewRecorder()
  48. inject := func(vm *otto.Otto) {
  49. vm.Set("_job_spec", func(call otto.FunctionCall) otto.Value {
  50. v, _ := vm.ToValue(string(ctxJSON))
  51. return v
  52. })
  53. vm.Set("_job_log", func(call otto.FunctionCall) otto.Value {
  54. line, _ := call.Argument(0).ToString()
  55. if hooks.Log != nil {
  56. hooks.Log(line)
  57. }
  58. return otto.TrueValue()
  59. })
  60. vm.Set("_job_progress", func(call otto.FunctionCall) otto.Value {
  61. p, _ := call.Argument(0).ToFloat()
  62. if hooks.Progress != nil {
  63. hooks.Progress(p)
  64. }
  65. return otto.TrueValue()
  66. })
  67. vm.Set("_job_cancelled", func(call otto.FunctionCall) otto.Value {
  68. cancelled := hooks.Cancelled != nil && hooks.Cancelled()
  69. v, _ := vm.ToValue(cancelled)
  70. return v
  71. })
  72. vm.Set("_job_output", func(call otto.FunctionCall) otto.Value {
  73. raw, _ := call.Argument(0).ToString()
  74. mu.Lock()
  75. if json.Valid([]byte(raw)) {
  76. output = json.RawMessage(raw)
  77. } else {
  78. js, _ := json.Marshal(raw)
  79. output = js
  80. }
  81. mu.Unlock()
  82. return otto.TrueValue()
  83. })
  84. }
  85. done := make(chan error, 1)
  86. go func() {
  87. done <- AGIGateway.ExecuteJobScript(source, rec.Spec.ScriptName, u, rw, req, inject, func(s *agi.JobStopper) {
  88. mu.Lock()
  89. stopper = s
  90. mu.Unlock()
  91. })
  92. }()
  93. select {
  94. case runErr := <-done:
  95. mu.Lock()
  96. defer mu.Unlock()
  97. return output, runErr
  98. case <-ctx.Done():
  99. //Timeout or cancellation: interrupt the VM and wait for it to unwind
  100. mu.Lock()
  101. s := stopper
  102. mu.Unlock()
  103. s.Stop()
  104. select {
  105. case <-done:
  106. case <-time.After(10 * time.Second):
  107. }
  108. mu.Lock()
  109. defer mu.Unlock()
  110. if ctx.Err() == context.DeadlineExceeded {
  111. return output, errors.New("job timed out")
  112. }
  113. return output, errors.New("job cancelled")
  114. }
  115. }
  116. // clusterJobLocality reports how many bytes of the given cluster paths have a
  117. // healthy copy on a node, and their total size.
  118. func clusterJobLocality(paths []string, nodeID string) (int64, int64) {
  119. if clusterMetadata == nil {
  120. return 0, 0
  121. }
  122. var onNode, total int64
  123. for _, p := range paths {
  124. rec, err := clusterMetadata.Stat(p)
  125. if err != nil {
  126. continue
  127. }
  128. if rec.IsDir {
  129. for _, child := range clusterMetadata.ListSubtree(rec.Path) {
  130. if child.IsDir {
  131. continue
  132. }
  133. total += child.Size
  134. if locationOnNode(&child, nodeID) {
  135. onNode += child.Size
  136. }
  137. }
  138. continue
  139. }
  140. total += rec.Size
  141. if locationOnNode(rec, nodeID) {
  142. onNode += rec.Size
  143. }
  144. }
  145. return onNode, total
  146. }
  147. func locationOnNode(rec *metadata.FileRecord, nodeID string) bool {
  148. for _, l := range rec.HealthyLocations() {
  149. if l.NodeID == nodeID {
  150. return true
  151. }
  152. }
  153. return false
  154. }
  155. // clusterJobCaller maps a request to the logged-in user for the jobs API.
  156. func clusterJobCaller(w http.ResponseWriter, r *http.Request) (jobs.Caller, error) {
  157. u, err := userHandler.GetUserInfoFromRequest(w, r)
  158. if err != nil {
  159. return jobs.Caller{}, err
  160. }
  161. return jobs.Caller{Username: u.Username, IsAdmin: u.IsAdmin()}, nil
  162. }
  163. // clusterJobSubmit reads the script from the caller's file system and queues
  164. // the job. POST: name, script (vpath), args (JSON), inputs (JSON array),
  165. // features (comma separated), timeout, priority.
  166. func clusterJobSubmit(w http.ResponseWriter, r *http.Request) {
  167. if clusterJobs == nil {
  168. utils.SendErrorResponse(w, "cluster jobs not available")
  169. return
  170. }
  171. u, err := userHandler.GetUserInfoFromRequest(w, r)
  172. if err != nil {
  173. utils.SendErrorResponse(w, "not logged in")
  174. return
  175. }
  176. scriptVpath, err := utils.PostPara(r, "script")
  177. if err != nil {
  178. utils.SendErrorResponse(w, "script path required")
  179. return
  180. }
  181. source, err := clusterReadScript(u.Username, scriptVpath)
  182. if err != nil {
  183. utils.SendErrorResponse(w, err.Error())
  184. return
  185. }
  186. req := jobs.SubmitRequest{
  187. Name: r.PostFormValue("name"),
  188. ScriptVpath: scriptVpath,
  189. Script: source,
  190. }
  191. if args := r.PostFormValue("args"); strings.TrimSpace(args) != "" {
  192. if !json.Valid([]byte(args)) {
  193. utils.SendErrorResponse(w, "args must be valid JSON")
  194. return
  195. }
  196. req.Args = json.RawMessage(args)
  197. }
  198. if inputs := r.PostFormValue("inputs"); strings.TrimSpace(inputs) != "" {
  199. if err := json.Unmarshal([]byte(inputs), &req.Inputs); err != nil {
  200. req.Inputs = splitList(inputs)
  201. }
  202. }
  203. if features := r.PostFormValue("features"); features != "" {
  204. req.Features = splitList(features)
  205. }
  206. if nodes := r.PostFormValue("nodes"); nodes != "" {
  207. req.Nodes = splitList(nodes)
  208. }
  209. if n, err := utils.PostInt(r, "timeout"); err == nil {
  210. req.TimeoutSec = n
  211. }
  212. if n, err := utils.PostInt(r, "priority"); err == nil {
  213. req.Priority = n
  214. }
  215. if n, err := utils.PostInt(r, "cores"); err == nil {
  216. req.MinCores = n
  217. }
  218. if ds := strings.TrimSpace(r.PostFormValue("dataset")); ds != "" {
  219. req.Dataset = ds
  220. }
  221. if n, err := utils.PostInt(r, "partition"); err == nil {
  222. req.PartitionMax = n
  223. }
  224. rec, err := clusterJobs.Submit(jobs.SpecFrom(req, u.Username))
  225. if err != nil {
  226. utils.SendErrorResponse(w, err.Error())
  227. return
  228. }
  229. js, _ := json.Marshal(rec)
  230. utils.SendJSONResponse(w, string(js))
  231. }
  232. func splitList(s string) []string {
  233. out := []string{}
  234. for _, part := range strings.Split(s, ",") {
  235. if p := strings.TrimSpace(part); p != "" {
  236. out = append(out, p)
  237. }
  238. }
  239. return out
  240. }
  241. // clusterSchedExplain answers /system/cluster/sched/explain?job=<id> with the
  242. // ranking the scheduler would have used for that job.
  243. func clusterSchedExplain(r *http.Request) (interface{}, error) {
  244. if clusterJobs == nil {
  245. return nil, errors.New("cluster jobs not available")
  246. }
  247. id, err := utils.GetPara(r, "job")
  248. if err != nil {
  249. return nil, errors.New("a job id is required")
  250. }
  251. results, err := clusterJobs.Explain(id)
  252. if err != nil {
  253. return nil, err
  254. }
  255. rec, _ := clusterJobs.Get(id)
  256. return map[string]interface{}{
  257. "job": rec.Spec.Name,
  258. "status": rec.State.Status,
  259. "node": rec.State.Node,
  260. "ranking": results,
  261. }, nil
  262. }
  263. // clusterReadScript loads a script from a user's file system.
  264. func clusterReadScript(username string, vpath string) (string, error) {
  265. u, err := userHandler.GetUserInfoFromUsername(username)
  266. if err != nil {
  267. return "", err
  268. }
  269. if !u.CanRead(vpath) {
  270. return "", errors.New("permission denied: " + vpath)
  271. }
  272. fsh, err := u.GetFileSystemHandlerFromVirtualPath(vpath)
  273. if err != nil {
  274. return "", err
  275. }
  276. rpath, err := fsh.FileSystemAbstraction.VirtualPathToRealPath(vpath, u.Username)
  277. if err != nil {
  278. return "", err
  279. }
  280. if !fsh.FileSystemAbstraction.FileExists(rpath) {
  281. return "", errors.New("script not found: " + vpath)
  282. }
  283. content, err := fsh.FileSystemAbstraction.ReadFile(rpath)
  284. if err != nil {
  285. return "", err
  286. }
  287. return string(content), nil
  288. }