This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / mill / mill_test.go
25 kB 873 lines
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}