This repository has no description
1package mill
2
3import (
4 "context"
5 "errors"
6 "io"
7 "log/slog"
8 "sync"
9 "testing"
10 "time"
11
12 "tangled.org/core/api/tangled"
13 "tangled.org/core/spindle/engine"
14 "tangled.org/core/spindle/models"
15
16 millproto "tangled.org/core/spindle/mill/proto"
17 millv1 "tangled.org/core/spindle/mill/proto/gen"
18)
19
20type scriptedEncoder func(*millproto.Message) error
21
22func (e scriptedEncoder) Encode(msg *millproto.Message) error { return e(msg) }
23
24func testWorkflow(name string) *models.Workflow {
25 return &models.Workflow{
26 Name: name,
27 Environment: map[string]string{},
28 Steps: []models.Step{remoteStep{}},
29 Data: &millWorkflowState{
30 RawWorkflow: tangled.Pipeline_Workflow{Name: name},
31 RawPipeline: tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{}},
32 },
33 }
34}
35
36func testWorkflowWithRunsOn(name string, runsOn []string) *models.Workflow {
37 wf := testWorkflow(name)
38 wf.Data.(*millWorkflowState).RawWorkflow.RunsOn = runsOn
39 return wf
40}
41
42func addCandidateSession(t *testing.T, m *Mill, nodeID string, labels []string, load float64, enc messageEncoder) *millSession {
43 t.Helper()
44 if enc == nil {
45 enc = scriptedEncoder(func(*millproto.Message) error { return nil })
46 }
47 sess := newSession(nodeID, "inc-"+nodeID, labels, enc, slog.New(slog.NewTextHandler(io.Discard, nil)))
48 sess.snapshot = &millv1.NodeSnapshot{
49 Seqno: 1,
50 Engines: map[string]*millv1.EngineAvailability{
51 "dummy": {Available: load < 1.0, Load: map[string]float64{"slots": load}},
52 },
53 }
54 m.mu.Lock()
55 m.sessions[nodeID] = sess
56 m.mu.Unlock()
57 return sess
58}
59
60func assertRankedNodes(t *testing.T, got []*millSession, want []string) {
61 t.Helper()
62 if len(got) != len(want) {
63 t.Fatalf("rankCandidates() returned %d candidates, want %d: got %v want %v", len(got), len(want), sessionIDs(got), want)
64 }
65 for i := range want {
66 if got[i].nodeID != want[i] {
67 t.Fatalf("rankCandidates()[%d] = %q, want %q; full order got %v want %v", i, got[i].nodeID, want[i], sessionIDs(got), want)
68 }
69 }
70}
71
72func sessionIDs(sessions []*millSession) []string {
73 out := make([]string, len(sessions))
74 for i, sess := range sessions {
75 out[i] = sess.nodeID
76 }
77 return out
78}
79
80func sameStringMultiset(a, b []string) bool {
81 if len(a) != len(b) {
82 return false
83 }
84 counts := make(map[string]int, len(a))
85 for _, s := range a {
86 counts[s]++
87 }
88 for _, s := range b {
89 if counts[s] == 0 {
90 return false
91 }
92 counts[s]--
93 }
94 return true
95}
96
97type reserveReply struct {
98 accepted bool
99 rejectClass millv1.RejectClass
100 reason string
101}
102
103func addReplyingCandidateSession(t *testing.T, m *Mill, nodeID string, labels []string, load float64, asked chan<- string, reply reserveReply) *millSession {
104 t.Helper()
105 var sess *millSession
106 sess = addCandidateSession(t, m, nodeID, labels, load, scriptedEncoder(func(msg *millproto.Message) error {
107 rs := msg.GetReserveSeat()
108 if rs == nil {
109 return nil
110 }
111 if asked != nil {
112 asked <- nodeID
113 }
114 sess.deliver(rs.GetLeaseId(), &millproto.Message{ReserveResult: &millv1.ReserveResult{
115 LeaseId: rs.GetLeaseId(),
116 Accepted: reply.accepted,
117 RejectReason: reply.reason,
118 RejectClass: reply.rejectClass,
119 }})
120 return nil
121 }))
122 return sess
123}
124
125func drainAsked(ch <-chan string) []string {
126 var out []string
127 for {
128 select {
129 case nodeID := <-ch:
130 out = append(out, nodeID)
131 default:
132 return out
133 }
134 }
135}
136
137func TestCommitRetriesAfterSessionCloseBeforeCommitted(t *testing.T) {
138 l := slog.New(slog.NewTextHandler(io.Discard, nil))
139 m := New(l, Config{BidTimeout: 25 * time.Millisecond, ReconnectGrace: time.Second})
140 wf := testWorkflow("build")
141 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
142 lease := newLease("lease-1", "node-1", "inc-1", "dummy")
143 lease.wid = wid
144 wf.Data.(*millWorkflowState).Lease = lease
145
146 m.mu.Lock()
147 m.leases[lease.id] = lease
148 m.mu.Unlock()
149
150 var sess1 *millSession
151 firstCommit := make(chan struct{})
152 sess1 = newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error {
153 if msg.GetCommitLease() != nil {
154 close(firstCommit)
155 m.detachSession(sess1)
156 }
157 return nil
158 }), l)
159 m.attachSession(sess1)
160
161 ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
162 defer cancel()
163 done := make(chan error, 1)
164 go func() { done <- m.commitAndWait(ctx, wf, nil) }()
165
166 select {
167 case <-firstCommit:
168 case <-ctx.Done():
169 t.Fatal("first commit was not sent")
170 }
171
172 var sess2 *millSession
173 sess2 = newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error {
174 if msg.GetCommitLease() == nil {
175 return nil
176 }
177 leaseID := msg.GetCommitLease().GetLeaseId()
178 sess2.deliver(leaseID, &millproto.Message{Committed: &millv1.Committed{LeaseId: leaseID}})
179 _ = m.onEventBatch(sess2, &millv1.EventBatch{
180 Epoch: sess2.epoch,
181 Events: []*millv1.Event{
182 {
183 Seqno: 1,
184 LeaseId: leaseID,
185 Payload: &millv1.Event_AttemptResult{
186 AttemptResult: &millv1.AttemptResult{
187 Status: millv1.TerminalStatus_SUCCESS,
188 },
189 },
190 },
191 },
192 })
193 return nil
194 }), l)
195 m.attachSession(sess2)
196 m.sessionReady(sess2)
197
198 select {
199 case err := <-done:
200 if err != nil {
201 t.Fatalf("commitAndWait() error = %v, want success after reconnect", err)
202 }
203 case <-ctx.Done():
204 t.Fatal("commitAndWait() did not finish after reconnect")
205 }
206}
207
208func TestDestroyRunningLeaseDoesNotDropCancelledTerminal(t *testing.T) {
209 l := slog.New(slog.NewTextHandler(io.Discard, nil))
210 m := New(l, Config{})
211 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
212 lease := newLease("lease-1", "node-1", "inc-1", "dummy")
213 lease.wid = wid
214 lease.setState(leaseRunning)
215
216 m.mu.Lock()
217 m.leases[lease.id] = lease
218 m.mu.Unlock()
219
220 m.destroy(wid)
221 if lease.getState() == leaseDone {
222 t.Fatal("destroy sealed the lease before the terminal result")
223 }
224
225 lease.deliverTerminal(&millv1.AttemptResult{
226 Status: millv1.TerminalStatus_CANCELLED,
227 })
228 res := <-lease.terminal
229 if err := terminalError(res.Status); !errors.Is(err, engine.ErrWorkflowCanceled) {
230 t.Fatalf("terminalError() = %v, want ErrCancelled", err)
231 }
232}
233
234func TestPlaceBlocksWhenNoCapacity(t *testing.T) {
235 l := slog.New(slog.NewTextHandler(io.Discard, nil))
236 m := New(l, Config{})
237 wf := testWorkflow("build")
238 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
239
240 // no executors at all: place must block until ctx expires (user sees pending)
241 ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond)
242 defer cancel()
243
244 _, err := m.place(ctx, "dummy", wid, wf)
245 if err != context.DeadlineExceeded {
246 t.Fatalf("place() error = %v, want DeadlineExceeded", err)
247 }
248}
249
250func TestRankCandidatesFiltersRequiredLabelsWithANDSemantics(t *testing.T) {
251 m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{})
252 addCandidateSession(t, m, "linux-high", []string{"linux"}, 0.0, nil)
253 addCandidateSession(t, m, "linux-arm", []string{"linux", "arm64"}, 0.25, nil)
254 addCandidateSession(t, m, "unlabeled", nil, 0.5, nil)
255 addCandidateSession(t, m, "linux-arm-gpu", []string{"linux", "arm64", "gpu"}, 0.75, nil)
256 addCandidateSession(t, m, "linux-arm-full", []string{"linux", "arm64"}, 1.0, nil)
257
258 tests := []struct {
259 name string
260 requiredLabels []string
261 want []string
262 }{
263 {
264 name: "no required labels keeps old capacity ranking",
265 want: []string{"linux-high", "linux-arm", "unlabeled", "linux-arm-gpu"},
266 },
267 {
268 name: "single required label includes every candidate carrying it",
269 requiredLabels: []string{"linux"},
270 want: []string{"linux-high", "linux-arm", "linux-arm-gpu"},
271 },
272 {
273 name: "all required labels must be present",
274 requiredLabels: []string{"linux", "arm64"},
275 want: []string{"linux-arm", "linux-arm-gpu"},
276 },
277 {
278 name: "one missing required label excludes the candidate",
279 requiredLabels: []string{"linux", "arm64", "gpu"},
280 want: []string{"linux-arm-gpu"},
281 },
282 {
283 name: "unknown required label leaves no candidate",
284 requiredLabels: []string{"linux", "arm64", "metal"},
285 },
286 }
287
288 for _, tt := range tests {
289 t.Run(tt.name, func(t *testing.T) {
290 assertRankedNodes(t, m.rankCandidates("dummy", tt.requiredLabels), tt.want)
291 })
292 }
293}
294
295func TestPlaceWithMissingRequiredLabelsStaysPendingWithoutReserve(t *testing.T) {
296 l := slog.New(slog.NewTextHandler(io.Discard, nil))
297 m := New(l, Config{BidTimeout: 10 * time.Millisecond})
298 reserveSent := make(chan struct{}, 1)
299 addCandidateSession(t, m, "linux-only", []string{"linux"}, 0.75, scriptedEncoder(func(msg *millproto.Message) error {
300 if msg.GetReserveSeat() != nil {
301 select {
302 case reserveSent <- struct{}{}:
303 default:
304 }
305 }
306 return nil
307 }))
308 wf := testWorkflowWithRunsOn("build", []string{"linux", "arm64"})
309 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
310
311 ctx, cancel := context.WithTimeout(context.Background(), 120*time.Millisecond)
312 defer cancel()
313 _, err := m.place(ctx, "dummy", wid, wf)
314 if err != context.DeadlineExceeded {
315 t.Fatalf("place() error = %v, want DeadlineExceeded while job remains pending", err)
316 }
317 select {
318 case <-reserveSent:
319 t.Fatal("place() sent ReserveSeat to executor missing a required label")
320 default:
321 }
322}
323
324func TestMaxPendingRejects(t *testing.T) {
325 l := slog.New(slog.NewTextHandler(io.Discard, nil))
326 m := New(l, Config{MaxPending: 1})
327
328 m.mu.Lock()
329 m.pending = 1
330 m.mu.Unlock()
331
332 wf2 := testWorkflow("b")
333 _, err := m.place(context.Background(), "dummy", models.WorkflowId{Name: "b"}, wf2)
334 if err == nil {
335 t.Fatal("place() past maxPending should error")
336 }
337}
338
339func TestCancelledRunningLeaseSurvivesReleaseForReconnectReplay(t *testing.T) {
340 m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
341 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
342 lease := newLease("lease-1", "node-1", "inc-1", "dummy")
343 lease.wid = wid
344 lease.setState(leaseRunning)
345 if err := m.persistLease(lease, leaseRowRunning); err != nil {
346 t.Fatalf("persistLease: %v", err)
347 }
348 m.mu.Lock()
349 m.leases[lease.id] = lease
350 m.mu.Unlock()
351
352 m.destroy(wid)
353 slot := &millSlot{fleet: m, lease: lease}
354 slot.Release()
355
356 m.mu.Lock()
357 _, retained := m.leases[lease.id]
358 m.mu.Unlock()
359 if !retained {
360 t.Fatal("slot release removed a cancellation-requested running lease before its terminal result")
361 }
362 if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 {
363 t.Fatalf("durable leases after slot release = %+v, err = %v; want retained lease", rows, err)
364 }
365
366 var sentMu sync.Mutex
367 var sent []*millproto.Message
368 sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error {
369 sentMu.Lock()
370 sent = append(sent, msg)
371 sentMu.Unlock()
372 return nil
373 }), discardLogger())
374 if _, ok := m.attachSession(sess); !ok {
375 t.Fatal("attachSession rejected reconnect")
376 }
377 m.sessionReady(sess)
378 sentMu.Lock()
379 var replayed bool
380 for _, msg := range sent {
381 if cancel := msg.GetCancelAttempt(); cancel != nil && cancel.GetLeaseId() == lease.id {
382 replayed = true
383 }
384 }
385 sentMu.Unlock()
386 if !replayed {
387 t.Fatal("reconnect did not replay CancelAttempt for retained lease")
388 }
389
390 if err := m.onEventBatch(sess, &millv1.EventBatch{
391 Epoch: sess.epoch,
392 Events: []*millv1.Event{
393 {
394 Seqno: 1,
395 LeaseId: lease.id,
396 Payload: &millv1.Event_AttemptResult{
397 AttemptResult: &millv1.AttemptResult{
398 Status: millv1.TerminalStatus_CANCELLED,
399 },
400 },
401 },
402 },
403 }); err != nil {
404 t.Fatalf("onEventBatch: %v", err)
405 }
406 m.mu.Lock()
407 _, retained = m.leases[lease.id]
408 m.mu.Unlock()
409 if retained {
410 t.Fatal("terminal result did not clean retained cancelled lease")
411 }
412 if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 {
413 t.Fatalf("durable leases after terminal = %+v, err = %v; want none", rows, err)
414 }
415 slot.Release()
416}
417
418func TestSessionRequestCancelledBeforeRegistrationOrSend(t *testing.T) {
419 sent := 0
420 sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(*millproto.Message) error {
421 sent++
422 return nil
423 }), discardLogger())
424 ctx, cancel := context.WithCancel(context.Background())
425 cancel()
426 _, err := sess.request(ctx, "lease-1", &millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: "lease-1"}})
427 if !errors.Is(err, context.Canceled) {
428 t.Fatalf("request error = %v, want context.Canceled", err)
429 }
430 if sent != 0 {
431 t.Fatalf("request sent %d messages for an already-cancelled context, want 0", sent)
432 }
433 sess.mu.Lock()
434 pending := len(sess.pending)
435 sess.mu.Unlock()
436 if pending != 0 {
437 t.Fatalf("request left %d pending waiters, want 0", pending)
438 }
439}
440
441func TestAttachSessionReplacesSilentIncumbentButRejectsActiveDuplicate(t *testing.T) {
442 m := New(discardLogger(), Config{ReconnectGrace: time.Minute})
443 old := newSession("node-1", "inc-old", nil, nopEncoder(), discardLogger())
444 transportClosed := make(chan struct{})
445 old.closeTransport = func() error {
446 close(transportClosed)
447 return nil
448 }
449 if _, ok := m.attachSession(old); !ok {
450 t.Fatal("first attach rejected")
451 }
452 m.mu.Lock()
453 old.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace)
454 m.mu.Unlock()
455
456 replacement := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger())
457 if _, ok := m.attachSession(replacement); !ok {
458 t.Fatal("silent incumbent blocked authenticated replacement")
459 }
460 select {
461 case <-transportClosed:
462 default:
463 t.Fatal("replacing a silent incumbent did not close its transport")
464 }
465 m.mu.Lock()
466 replacement.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace)
467 m.mu.Unlock()
468 if err := replacement.dispatch(m, &millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{Seqno: 1}}); err != nil {
469 t.Fatalf("periodic snapshot dispatch: %v", err)
470 }
471 if _, ok := m.attachSession(newSession("node-1", "inc-dup", nil, nopEncoder(), discardLogger())); ok {
472 t.Fatal("active replacement did not reject a duplicate session")
473 }
474}
475
476func TestPlaceReleasesRemoteReservationWhenInitialPersistenceFails(t *testing.T) {
477 m, bdb := restoreTestMill(t, Config{BidTimeout: time.Second})
478 if err := bdb.Close(); err != nil {
479 t.Fatalf("close db: %v", err)
480 }
481 released := make(chan string, 1)
482 var sess *millSession
483 sess = addCandidateSession(t, m, "node-1", nil, 0, scriptedEncoder(func(msg *millproto.Message) error {
484 switch {
485 case msg.GetReserveSeat() != nil:
486 leaseID := msg.GetReserveSeat().GetLeaseId()
487 sess.deliver(leaseID, &millproto.Message{ReserveResult: &millv1.ReserveResult{
488 LeaseId: leaseID,
489 Accepted: true,
490 }})
491 case msg.GetReleaseLease() != nil:
492 released <- msg.GetReleaseLease().GetLeaseId()
493 }
494 return nil
495 }))
496
497 ctx, cancel := context.WithTimeout(context.Background(), time.Second)
498 defer cancel()
499 slot, err := m.place(
500 ctx,
501 "dummy",
502 models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"},
503 testWorkflow("build"),
504 )
505 if err == nil {
506 t.Fatal("place succeeded after reserved lease persistence failed")
507 }
508 if slot != nil {
509 t.Fatalf("place returned slot %T after persistence failure", slot)
510 }
511 select {
512 case leaseID := <-released:
513 if leaseID == "" {
514 t.Fatal("ReleaseLease had empty lease id")
515 }
516 case <-time.After(time.Second):
517 t.Fatal("persistence failure did not compensate with ReleaseLease")
518 }
519 m.mu.Lock()
520 leases := len(m.leases)
521 m.mu.Unlock()
522 if leases != 0 {
523 t.Fatalf("mill published %d leases after initial persistence failure, want 0", leases)
524 }
525}
526
527func TestGapsAndDuplicates(t *testing.T) {
528 m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
529 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
530 m.attachSession(sess)
531
532 owned := newLease("lease-1", "node-1", "inc-1", "dummy")
533 owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
534 m.mu.Lock()
535 m.leases[owned.id] = owned
536 m.mu.Unlock()
537
538 err := m.onEventBatch(sess, &millv1.EventBatch{
539 Epoch: sess.epoch,
540 Events: []*millv1.Event{
541 {
542 Seqno: 0,
543 LeaseId: owned.id,
544 Payload: &millv1.Event_StatusEvent{
545 StatusEvent: &millv1.StatusEvent{Status: millv1.NonterminalStatus_RUNNING},
546 },
547 },
548 },
549 })
550 if err != nil {
551 t.Fatalf("expected duplicate to be skipped without error, got: %v", err)
552 }
553
554 err = m.onEventBatch(sess, &millv1.EventBatch{
555 Epoch: sess.epoch,
556 Events: []*millv1.Event{
557 {
558 Seqno: 2,
559 LeaseId: owned.id,
560 Payload: &millv1.Event_StatusEvent{
561 StatusEvent: &millv1.StatusEvent{Status: millv1.NonterminalStatus_RUNNING},
562 },
563 },
564 },
565 })
566 if err == nil {
567 t.Fatal("expected error due to seqno gap, got nil")
568 }
569}
570
571func TestAtomicBatchRollback(t *testing.T) {
572 m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
573 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
574 m.attachSession(sess)
575
576 owned := newLease("lease-1", "node-1", "inc-1", "dummy")
577 owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
578 m.mu.Lock()
579 m.leases[owned.id] = owned
580 m.mu.Unlock()
581
582 if _, err := bdb.Exec(`
583 create trigger reject_status_event
584 before insert on events
585 begin
586 select raise(abort, 'forced status event failure');
587 end
588 `); err != nil {
589 t.Fatalf("failed to create fail trigger: %v", err)
590 }
591 defer bdb.Exec("drop trigger reject_status_event")
592
593 err := m.onEventBatch(sess, &millv1.EventBatch{
594 Epoch: sess.epoch,
595 Events: []*millv1.Event{
596 {
597 Seqno: 1,
598 LeaseId: owned.id,
599 Payload: &millv1.Event_StatusEvent{
600 StatusEvent: &millv1.StatusEvent{Status: millv1.NonterminalStatus_RUNNING},
601 },
602 },
603 },
604 })
605 if err == nil {
606 t.Fatal("expected status event insertion to fail due to trigger")
607 }
608
609 m.mu.Lock()
610 seqno := m.nodeSeqno["node-1/inc-1"]
611 m.mu.Unlock()
612 if seqno != 0 {
613 t.Fatalf("expected seqno 0 due to rollback, got %d", seqno)
614 }
615}
616
617func TestTerminalBeforeACK(t *testing.T) {
618 m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
619
620 ackSent := make(chan struct{})
621 sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error {
622 if msg.GetAck() != nil {
623 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
624 st, err := bdb.GetStatus(wid)
625 if err != nil || st.Status != "success" {
626 t.Errorf("expected terminal status success at ACK time, got status: %v, err: %v", st, err)
627 }
628 close(ackSent)
629 }
630 return nil
631 }), discardLogger())
632 m.attachSession(sess)
633
634 owned := newLease("lease-1", "node-1", "inc-1", "dummy")
635 owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
636 m.mu.Lock()
637 m.leases[owned.id] = owned
638 m.mu.Unlock()
639
640 err := m.onEventBatch(sess, &millv1.EventBatch{
641 Epoch: sess.epoch,
642 Events: []*millv1.Event{
643 {
644 Seqno: 1,
645 LeaseId: owned.id,
646 Payload: &millv1.Event_AttemptResult{
647 AttemptResult: &millv1.AttemptResult{Status: millv1.TerminalStatus_SUCCESS},
648 },
649 },
650 },
651 })
652 if err != nil {
653 t.Fatalf("onEventBatch: %v", err)
654 }
655
656 select {
657 case <-ackSent:
658 case <-time.After(2 * time.Second):
659 t.Fatal("ACK was not sent")
660 }
661}
662
663func TestExecutorRestartEmptySnapshot(t *testing.T) {
664 m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
665
666 lease := newLease("lease-1", "node-1", "inc-old", "dummy")
667 lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
668 m.mu.Lock()
669 m.leases[lease.id] = lease
670 m.mu.Unlock()
671 if err := m.persistLease(lease, leaseRowRunning); err != nil {
672 t.Fatalf("persistLease: %v", err)
673 }
674
675 sess := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger())
676 m.attachSession(sess)
677
678 err := m.onSnapshot(sess, &millv1.NodeSnapshot{
679 Seqno: 1,
680 ActiveLeaseIds: nil,
681 })
682 if err != nil {
683 t.Fatalf("onSnapshot: %v", err)
684 }
685
686 m.mu.Lock()
687 _, stillActive := m.leases["lease-1"]
688 m.mu.Unlock()
689 if stillActive {
690 t.Fatal("expected old epoch lease to be reconciled and failed")
691 }
692}
693
694func TestReplacementLostBeforeSnapshotFailsOldEpochLease(t *testing.T) {
695 m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
696 lease := newLease("lease-1", "node-1", "inc-old", "dummy")
697 lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
698 lease.setState(leaseRunning)
699 if err := m.persistLease(lease, leaseRowRunning); err != nil {
700 t.Fatalf("persistLease: %v", err)
701 }
702 m.mu.Lock()
703 m.leases[lease.id] = lease
704 m.mu.Unlock()
705
706 replacement := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger())
707 if _, ok := m.attachSession(replacement); !ok {
708 t.Fatal("attachSession rejected replacement")
709 }
710 replacement.disconnected = true
711 m.failLeasesAfterGrace(replacement)
712
713 m.mu.Lock()
714 _, stillActive := m.leases[lease.id]
715 m.mu.Unlock()
716 if stillActive {
717 t.Fatal("replacement loss stranded old-epoch lease")
718 }
719}
720
721func TestSeqRegression(t *testing.T) {
722 m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
723 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
724 m.attachSession(sess)
725
726 err := m.onSnapshot(sess, &millv1.NodeSnapshot{
727 Seqno: 5,
728 })
729 if err != nil {
730 t.Fatalf("first snapshot: %v", err)
731 }
732
733 err = m.onSnapshot(sess, &millv1.NodeSnapshot{
734 Seqno: 4,
735 })
736 if err == nil {
737 t.Fatal("expected seqno regression to be rejected")
738 }
739}
740
741func TestClaimedSweep(t *testing.T) {
742 m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
743
744 lease := newLease("lease-1", "node-1", "inc-1", "dummy")
745 lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
746 lease.orphaned = true
747 lease.claimed = false
748 m.mu.Lock()
749 m.leases[lease.id] = lease
750 m.mu.Unlock()
751
752 sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
753 m.attachSession(sess)
754 err := m.onSnapshot(sess, &millv1.NodeSnapshot{
755 Seqno: 1,
756 ActiveLeaseIds: []string{"lease-1"},
757 })
758 if err != nil {
759 t.Fatalf("onSnapshot: %v", err)
760 }
761
762 m.detachSession(sess)
763
764 m.sweepUnclaimedOrphans()
765
766 m.mu.Lock()
767 _, stillRunning := m.leases["lease-1"]
768 m.mu.Unlock()
769 if !stillRunning {
770 t.Fatal("claimed lease was incorrectly swept by startup sweep")
771 }
772}
773
774func TestCancelDeadline(t *testing.T) {
775 m, _ := restoreTestMill(t, Config{CancelTimeout: 50 * time.Millisecond})
776
777 sessClosed := make(chan struct{})
778 sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error {
779 return nil
780 }), discardLogger())
781 sess.closeTransport = func() error {
782 close(sessClosed)
783 return nil
784 }
785 m.attachSession(sess)
786
787 lease := newLease("lease-1", "node-1", "inc-1", "dummy")
788 lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
789 lease.setState(leaseRunning)
790 m.mu.Lock()
791 m.leases[lease.id] = lease
792 m.mu.Unlock()
793
794 m.destroy(lease.wid)
795
796 select {
797 case <-sessClosed:
798 case <-time.After(1 * time.Second):
799 t.Fatal("session was not closed after cancel deadline expiration")
800 }
801}
802
803func TestCleanupRetry(t *testing.T) {
804 m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
805
806 lease := newLease("lease-1", "node-1", "inc-1", "dummy")
807 lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
808 if err := m.persistLease(lease, leaseRowRunning); err != nil {
809 t.Fatalf("persist lease: %v", err)
810 }
811 m.mu.Lock()
812 m.leases[lease.id] = lease
813 m.mu.Unlock()
814
815 if _, err := bdb.Exec(`
816 create trigger reject_cleanup_delete
817 before delete on mill_leases
818 begin
819 select raise(abort, 'forced delete failure');
820 end
821 `); err != nil {
822 t.Fatalf("failed to create fail trigger: %v", err)
823 }
824
825 err := m.cleanupLease(lease)
826 if err == nil {
827 t.Fatal("expected cleanupLease to fail")
828 }
829
830 m.mu.Lock()
831 _, stillRunning := m.leases["lease-1"]
832 m.mu.Unlock()
833 if !stillRunning {
834 t.Fatal("lease was removed from memory despite cleanup failure")
835 }
836
837 if _, err := bdb.Exec("drop trigger reject_cleanup_delete"); err != nil {
838 t.Fatalf("drop trigger: %v", err)
839 }
840
841 err = m.cleanupLease(lease)
842 if err != nil {
843 t.Fatalf("expected retry cleanup to succeed, got: %v", err)
844 }
845
846 m.mu.Lock()
847 _, stillRunning = m.leases["lease-1"]
848 m.mu.Unlock()
849 if stillRunning {
850 t.Fatal("lease still in memory after successful cleanup retry")
851 }
852}
853
854func TestBoundedBidding(t *testing.T) {
855 m := New(discardLogger(), Config{TopK: 2, BidTimeout: 10 * time.Millisecond})
856
857 addCandidateSession(t, m, "node-1", nil, 0, nil)
858 addCandidateSession(t, m, "node-2", nil, 0, nil)
859 addCandidateSession(t, m, "node-3", nil, 0, nil)
860 addCandidateSession(t, m, "node-4", nil, 0, nil)
861 addCandidateSession(t, m, "node-5", nil, 0, nil)
862
863 ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
864 defer cancel()
865
866 lease, err := m.bid(ctx, "dummy", models.WorkflowId{}, testWorkflow("build"))
867 if err != nil {
868 t.Fatalf("bid: %v", err)
869 }
870 if lease != nil {
871 t.Fatalf("did not expect a lease, got %+v", lease)
872 }
873}