jobs_test.go 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413
  1. package jobs
  2. import (
  3. "context"
  4. "encoding/json"
  5. "errors"
  6. "net/http/httptest"
  7. "path/filepath"
  8. "strings"
  9. "sync"
  10. "testing"
  11. "time"
  12. "imuslab.com/arozos/mod/cluster/capability"
  13. "imuslab.com/arozos/mod/cluster/membership"
  14. "imuslab.com/arozos/mod/cluster/metadata"
  15. )
  16. func init() {
  17. membership.HeartbeatInterval = 500 * time.Millisecond
  18. membership.OnlineWindow = 2 * time.Second
  19. membership.OfflineWindow = 4 * time.Second
  20. }
  21. // fakeExecutor pretends to run scripts: it looks for markers in the source.
  22. type fakeExecutor struct {
  23. mu sync.Mutex
  24. ran []string
  25. users map[string]bool
  26. block chan struct{} // when set, Run waits on it (or on ctx)
  27. failAll bool
  28. }
  29. func newExec(users ...string) *fakeExecutor {
  30. f := &fakeExecutor{users: map[string]bool{}}
  31. for _, u := range users {
  32. f.users[u] = true
  33. }
  34. return f
  35. }
  36. func (f *fakeExecutor) UserExists(owner string) bool {
  37. f.mu.Lock()
  38. defer f.mu.Unlock()
  39. return f.users[owner]
  40. }
  41. func (f *fakeExecutor) Run(ctx context.Context, rec Record, source string, hooks ExecHooks) (json.RawMessage, error) {
  42. f.mu.Lock()
  43. f.ran = append(f.ran, rec.Spec.ID)
  44. block := f.block
  45. fail := f.failAll
  46. f.mu.Unlock()
  47. if hooks.Log != nil {
  48. hooks.Log("started " + rec.Spec.Name)
  49. }
  50. if hooks.Progress != nil {
  51. hooks.Progress(0.5)
  52. }
  53. if block != nil {
  54. select {
  55. case <-block:
  56. case <-ctx.Done():
  57. return nil, ctx.Err()
  58. }
  59. }
  60. if fail {
  61. return nil, context.Canceled
  62. }
  63. if !strings.Contains(source, "function run(") {
  64. return nil, ErrScriptNoRun
  65. }
  66. return json.RawMessage(`{"ok":true,"node":"` + rec.State.Node + `"}`), nil
  67. }
  68. func (f *fakeExecutor) count() int {
  69. f.mu.Lock()
  70. defer f.mu.Unlock()
  71. return len(f.ran)
  72. }
  73. type testNode struct {
  74. id string
  75. m *membership.Manager
  76. meta *metadata.Manager
  77. jobs *Manager
  78. exec Executor
  79. srv *httptest.Server
  80. }
  81. func newTestNode(t *testing.T, id string, exec Executor, caps capability.Manifest) *testNode {
  82. t.Helper()
  83. dir := t.TempDir()
  84. m, err := membership.NewManager(membership.Option{
  85. NodeID: id, DBFile: filepath.Join(dir, "c.db"), KeyFile: filepath.Join(dir, "k"),
  86. Version: "t", DefaultName: "Node " + id,
  87. Capabilities: func() capability.Manifest { return caps },
  88. Health: func() membership.Health { return membership.Health{CPUUsage: 10, RAMUsed: 1, RAMTotal: 4} },
  89. })
  90. if err != nil {
  91. t.Fatalf("NewManager: %v", err)
  92. }
  93. meta, err := metadata.New(metadata.Option{
  94. Membership: m, LeaseDuration: 3 * time.Second, LeaseRenew: time.Second,
  95. LeaseTick: 200 * time.Millisecond, PendingRetry: 300 * time.Millisecond,
  96. })
  97. if err != nil {
  98. t.Fatalf("metadata.New: %v", err)
  99. }
  100. jm, err := New(Option{
  101. Membership: m, Metadata: meta, Executor: exec,
  102. ScheduleInterval: 300 * time.Millisecond, TaskLease: 2 * time.Second, LeaseRenew: 500 * time.Millisecond,
  103. })
  104. if err != nil {
  105. t.Fatalf("jobs.New: %v", err)
  106. }
  107. srv := httptest.NewServer(m.ACNHandler())
  108. cfg := m.Config()
  109. cfg.AdvertiseURL = srv.URL
  110. m.UpdateConfig(cfg)
  111. t.Cleanup(func() { jm.Close(); meta.Close(); m.Close(); srv.Close() })
  112. return &testNode{id: id, m: m, meta: meta, jobs: jm, exec: exec, srv: srv}
  113. }
  114. func manifest(os, arch string, features ...string) capability.Manifest {
  115. f := map[string]bool{}
  116. for _, x := range features {
  117. f[x] = true
  118. }
  119. return capability.Manifest{OS: os, Arch: arch, CPUCores: 4, TotalRAM: 8 << 30, Features: f, DetectedAt: 1}
  120. }
  121. func waitFor(t *testing.T, what string, d time.Duration, cond func() bool) {
  122. t.Helper()
  123. deadline := time.Now().Add(d)
  124. for !cond() {
  125. if time.Now().After(deadline) {
  126. t.Fatalf("timed out waiting for %s", what)
  127. }
  128. time.Sleep(50 * time.Millisecond)
  129. }
  130. }
  131. func cluster2(t *testing.T, execA, execB *fakeExecutor, capsA, capsB capability.Manifest) (*testNode, *testNode) {
  132. t.Helper()
  133. a := newTestNode(t, "node-a", execA, capsA)
  134. b := newTestNode(t, "node-b", execB, capsB)
  135. if _, err := a.m.CreateCluster("Jobs"); err != nil {
  136. t.Fatalf("CreateCluster: %v", err)
  137. }
  138. time.Sleep(1100 * time.Millisecond)
  139. token, _, _ := a.m.NewJoinToken(time.Hour)
  140. if _, err := b.m.JoinCluster(token); err != nil {
  141. t.Fatalf("join: %v", err)
  142. }
  143. waitFor(t, "leader", 10*time.Second, func() bool { return a.meta.IsLeader() && b.meta.Leader() == "node-a" })
  144. return a, b
  145. }
  146. const goodScript = "function run(job) { return {ok: true}; }"
  147. func TestBuildSourceWrapsScript(t *testing.T) {
  148. src := buildSource(goodScript)
  149. for _, want := range []string{"_job_spec()", "JOB.log", "JOB.abortIfCancelled", goodScript, "_job_output", "run(JOB)"} {
  150. if !strings.Contains(src, want) {
  151. t.Errorf("wrapped source missing %q", want)
  152. }
  153. }
  154. }
  155. func TestSubmitValidation(t *testing.T) {
  156. a := newTestNode(t, "solo", newExec("toby"), manifest("linux", "amd64"))
  157. if _, err := a.jobs.Submit(Spec{Owner: "toby", Script: goodScript}); err != metadata.ErrNotInCluster {
  158. t.Errorf("submit outside cluster: %v", err)
  159. }
  160. a.m.CreateCluster("Solo")
  161. waitFor(t, "leader", 5*time.Second, func() bool { return a.meta.IsLeader() })
  162. if _, err := a.jobs.Submit(Spec{Owner: "toby"}); err == nil {
  163. t.Errorf("empty script accepted")
  164. }
  165. if _, err := a.jobs.Submit(Spec{Script: goodScript}); err == nil {
  166. t.Errorf("missing owner accepted")
  167. }
  168. rec, err := a.jobs.Submit(Spec{Owner: "toby", Script: goodScript, ScriptName: "x.agi", Inputs: []string{"cluster:/a/../b"}})
  169. if err != nil {
  170. t.Fatalf("Submit: %v", err)
  171. }
  172. if rec.Spec.ID == "" || rec.Spec.Name != "x.agi" || rec.Spec.TimeoutSec != 3600 || rec.Spec.MaxAttempts != 3 || rec.Spec.Kind != KindRun {
  173. t.Errorf("defaults wrong: %+v", rec.Spec)
  174. }
  175. if rec.Spec.Inputs[0] != "/b" {
  176. t.Errorf("inputs not normalised: %v", rec.Spec.Inputs)
  177. }
  178. }
  179. func TestJobRunsAndReplicates(t *testing.T) {
  180. execA, execB := newExec("toby"), newExec("toby")
  181. a, b := cluster2(t, execA, execB, manifest("linux", "amd64"), manifest("linux", "amd64"))
  182. rec, err := a.jobs.Submit(Spec{Owner: "toby", Name: "hello", Script: goodScript})
  183. if err != nil {
  184. t.Fatalf("Submit: %v", err)
  185. }
  186. waitFor(t, "job to succeed", 15*time.Second, func() bool {
  187. r, ok := a.jobs.Get(rec.Spec.ID)
  188. return ok && r.State.Status == StatusSucceeded
  189. })
  190. final, _ := a.jobs.Get(rec.Spec.ID)
  191. if final.State.Node == "" || final.State.Progress != 1 || len(final.State.Log) == 0 {
  192. t.Errorf("final state incomplete: %+v", final.State)
  193. }
  194. var out map[string]interface{}
  195. if json.Unmarshal(final.State.Output, &out) != nil || out["ok"] != true {
  196. t.Errorf("output wrong: %s", final.State.Output)
  197. }
  198. if execA.count()+execB.count() != 1 {
  199. t.Errorf("job should run exactly once, ran %d times", execA.count()+execB.count())
  200. }
  201. //The record reaches the other node
  202. waitFor(t, "b to see the finished job", 5*time.Second, func() bool {
  203. r, ok := b.jobs.Get(rec.Spec.ID)
  204. return ok && r.State.Status == StatusSucceeded
  205. })
  206. //Listing is owner scoped
  207. if len(b.jobs.List("toby")) != 1 || len(b.jobs.List("someone")) != 0 {
  208. t.Errorf("owner filtering wrong")
  209. }
  210. st := a.jobs.Status("")
  211. if st.Succeeded != 1 || !st.IsLeader || st.LocalSlots < 1 {
  212. t.Errorf("status wrong: %+v", st)
  213. }
  214. }
  215. func TestRequirementsPickTheRightNode(t *testing.T) {
  216. execA, execB := newExec("toby"), newExec("toby")
  217. //Only b has ffmpeg
  218. a, b := cluster2(t, execA, execB, manifest("linux", "amd64"), manifest("linux", "amd64", "ffmpeg"))
  219. rec, err := a.jobs.Submit(Spec{Owner: "toby", Name: "transcode", Script: goodScript,
  220. Requirements: capability.Requirements{Features: []string{"ffmpeg"}}})
  221. if err != nil {
  222. t.Fatalf("Submit: %v", err)
  223. }
  224. waitFor(t, "job to run on b", 15*time.Second, func() bool {
  225. r, ok := a.jobs.Get(rec.Spec.ID)
  226. return ok && r.State.Status == StatusSucceeded
  227. })
  228. final, _ := a.jobs.Get(rec.Spec.ID)
  229. if final.State.Node != "node-b" {
  230. t.Errorf("job should run on the node with ffmpeg, ran on %s", final.State.Node)
  231. }
  232. if execA.count() != 0 || execB.count() != 1 {
  233. t.Errorf("wrong executor ran the job: a=%d b=%d", execA.count(), execB.count())
  234. }
  235. //A requirement nobody satisfies stays queued with an explanation
  236. rec2, _ := a.jobs.Submit(Spec{Owner: "toby", Name: "cuda", Script: goodScript,
  237. Requirements: capability.Requirements{Features: []string{"cuda"}}})
  238. waitFor(t, "queued reason", 10*time.Second, func() bool {
  239. r, ok := a.jobs.Get(rec2.Spec.ID)
  240. return ok && r.State.Status == StatusQueued && strings.Contains(r.State.Reason, "cuda")
  241. })
  242. _ = b
  243. }
  244. func TestCancelRunningJob(t *testing.T) {
  245. execA := newExec("toby")
  246. execA.block = make(chan struct{})
  247. a := newTestNode(t, "solo", execA, manifest("linux", "amd64"))
  248. a.m.CreateCluster("Solo")
  249. waitFor(t, "leader", 5*time.Second, func() bool { return a.meta.IsLeader() })
  250. rec, err := a.jobs.Submit(Spec{Owner: "toby", Name: "slow", Script: goodScript})
  251. if err != nil {
  252. t.Fatalf("Submit: %v", err)
  253. }
  254. waitFor(t, "job running", 10*time.Second, func() bool {
  255. r, ok := a.jobs.Get(rec.Spec.ID)
  256. return ok && r.State.Status == StatusRunning
  257. })
  258. if err := a.jobs.Cancel(rec.Spec.ID, "someone-else", false); err == nil {
  259. t.Errorf("another user cancelled the job")
  260. }
  261. if err := a.jobs.Cancel(rec.Spec.ID, "toby", false); err != nil {
  262. t.Fatalf("Cancel: %v", err)
  263. }
  264. waitFor(t, "cancelled", 10*time.Second, func() bool {
  265. r, ok := a.jobs.Get(rec.Spec.ID)
  266. return ok && r.State.Status == StatusCancelled && r.State.Finished > 0
  267. })
  268. close(execA.block)
  269. }
  270. func TestWaitAndOwnerMissing(t *testing.T) {
  271. execA, execB := newExec(), newExec("toby") //a has no account for toby
  272. a, b := cluster2(t, execA, execB, manifest("linux", "amd64"), manifest("linux", "amd64"))
  273. rec, err := a.jobs.Submit(Spec{Owner: "toby", Name: "wait-me", Script: goodScript})
  274. if err != nil {
  275. t.Fatalf("Submit: %v", err)
  276. }
  277. final, err := a.jobs.Wait(rec.Spec.ID, 20*time.Second)
  278. if err != nil {
  279. t.Fatalf("Wait: %v", err)
  280. }
  281. if final.State.Status != StatusSucceeded || final.State.Node != "node-b" {
  282. t.Errorf("job should have run on the node that knows the user: %+v", final.State)
  283. }
  284. if _, err := a.jobs.Wait("no-such-job", time.Second); err != ErrNotFound {
  285. t.Errorf("waiting for an unknown job: %v", err)
  286. }
  287. _ = b
  288. }
  289. func TestNodesPinAJob(t *testing.T) {
  290. execA, execB := newExec("toby"), newExec("toby")
  291. a, _ := cluster2(t, execA, execB, manifest("linux", "amd64"), manifest("linux", "amd64"))
  292. //One job per node, each pinned to its node, lands exactly there
  293. for _, id := range []string{"node-a", "node-b"} {
  294. rec, err := a.jobs.Submit(SpecFrom(SubmitRequest{Name: "hello " + id, Script: goodScript, Nodes: []string{id}}, "toby"))
  295. if err != nil {
  296. t.Fatalf("Submit(%s): %v", id, err)
  297. }
  298. waitFor(t, "pinned job on "+id, 15*time.Second, func() bool {
  299. r, ok := a.jobs.Get(rec.Spec.ID)
  300. return ok && r.State.Status == StatusSucceeded
  301. })
  302. if r, _ := a.jobs.Get(rec.Spec.ID); r.State.Node != id {
  303. t.Errorf("job pinned to %s ran on %s", id, r.State.Node)
  304. }
  305. }
  306. if execA.count() != 1 || execB.count() != 1 {
  307. t.Errorf("each node should run exactly one job: a=%d b=%d", execA.count(), execB.count())
  308. }
  309. //A job pinned to a node that is not in the cluster waits and says why
  310. rec, _ := a.jobs.Submit(SpecFrom(SubmitRequest{Name: "nowhere", Script: goodScript, Nodes: []string{"node-z"}}, "toby"))
  311. waitFor(t, "queued reason", 10*time.Second, func() bool {
  312. r, ok := a.jobs.Get(rec.Spec.ID)
  313. return ok && r.State.Status == StatusQueued && strings.Contains(r.State.Reason, "limited to other nodes")
  314. })
  315. }
  316. func TestRefusalIsPermanent(t *testing.T) {
  317. tests := []struct {
  318. name string
  319. err error
  320. want bool
  321. }{
  322. {"no error", nil, false},
  323. {"owner missing locally", ErrNoUser, true},
  324. {"cannot execute locally", ErrNoExecutor, true},
  325. {"owner missing on a remote node", errors.New("remote node returned Conflict: {\"error\":\"" + ErrNoUser.Error() + "\"}"), true},
  326. {"record not replicated yet", errors.New("remote node returned Not Found: {\"error\":\"job not found\"}"), false},
  327. {"tunnel still reconnecting", errors.New("target node is not connected through a tunnel on this node"), false},
  328. {"request timed out", context.DeadlineExceeded, false},
  329. }
  330. for _, tc := range tests {
  331. t.Run(tc.name, func(t *testing.T) {
  332. if got := refusalIsPermanent(tc.err); got != tc.want {
  333. t.Errorf("refusalIsPermanent(%v) = %v; want %v", tc.err, got, tc.want)
  334. }
  335. })
  336. }
  337. }
  338. func TestBurstRunsEachJobOnce(t *testing.T) {
  339. execA, execB := newExec("toby"), newExec("toby")
  340. a, _ := cluster2(t, execA, execB, manifest("linux", "amd64"), manifest("linux", "amd64"))
  341. //Every Submit starts its own scheduling pass, so a burst of submissions
  342. //runs many passes at once; each job must still run exactly one time
  343. const total = 30
  344. ids := make(chan string, total)
  345. var wg sync.WaitGroup
  346. for i := 0; i < total; i++ {
  347. wg.Add(1)
  348. go func(i int) {
  349. defer wg.Done()
  350. pin := []string{"node-a", "node-b"}[i%2]
  351. rec, err := a.jobs.Submit(SpecFrom(SubmitRequest{Name: "burst", Script: goodScript, Nodes: []string{pin}}, "toby"))
  352. if err != nil {
  353. t.Errorf("Submit: %v", err)
  354. return
  355. }
  356. ids <- rec.Spec.ID
  357. }(i)
  358. }
  359. wg.Wait()
  360. close(ids)
  361. all := []string{}
  362. for id := range ids {
  363. all = append(all, id)
  364. }
  365. waitFor(t, "every job to finish", 30*time.Second, func() bool {
  366. for _, id := range all {
  367. if r, ok := a.jobs.Get(id); !ok || r.State.Status != StatusSucceeded {
  368. return false
  369. }
  370. }
  371. return true
  372. })
  373. runs := map[string]int{}
  374. for _, f := range []*fakeExecutor{execA, execB} {
  375. f.mu.Lock()
  376. for _, id := range f.ran {
  377. runs[id]++
  378. }
  379. f.mu.Unlock()
  380. }
  381. for _, id := range all {
  382. if runs[id] != 1 {
  383. t.Errorf("job %s ran %d times, want exactly once", id, runs[id])
  384. }
  385. }
  386. }