Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657package mill
import ( "bytes" "context" "errors" "io" "log/slog" "strings" "sync" "testing" "time"
"tangled.org/core/api/tangled" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/models" "tangled.org/core/spindle/storage"
millproto "tangled.org/core/spindle/mill/proto" millv1 "tangled.org/core/spindle/mill/proto/gen")
type scriptedEncoder func(*millproto.Message) error
func (e scriptedEncoder) Encode(msg *millproto.Message) error { return e(msg) }
func waitForCacheDeletion(t *testing.T, cache storage.Storage, index *db.DB, key string) { t.Helper() deadline := time.Now().Add(2 * time.Second) for time.Now().Before(deadline) { reader, err := cache.Get(context.Background(), key) if errors.Is(err, storage.ErrNotExist) { if index != nil { pending, queryErr := index.PendingCacheObjectDeletions(context.Background(), 10) if queryErr != nil { t.Fatal(queryErr) } if len(pending) == 0 { return } } else { return } } else if err == nil { _ = reader.Close() } time.Sleep(time.Millisecond) } t.Fatalf("cache object %q was not deleted", key)}
func testWorkflow(name string) *models.Workflow { repoDid := "did:web:example.com" return &models.Workflow{ Name: name, Environment: map[string]string{}, Steps: []models.Step{remoteStep{}}, RepoDID: repoDid, Data: &millWorkflowState{ RawWorkflow: tangled.Pipeline_Workflow{Name: name}, RawPipeline: testPipeline(), }, }}func TestMillEngineRequiresSharedCacheStoreID(t *testing.T) { twf := tangled.Pipeline_Workflow{ Name: "build", Raw: "cache:\n - key: deps\n paths: [deps]\n", } for _, tc := range []struct { name string storeID string want int }{ {name: "unidentified store", want: 0}, {name: "shared store", storeID: "fleet-cache", want: 1}, } { t.Run(tc.name, func(t *testing.T) { m := New(discardLogger(), Config{CacheStoreID: tc.storeID}) wf, err := NewEngine("microvm", m).InitWorkflow(twf, tangled.Pipeline{}) if err != nil { t.Fatal(err) } if got := len(wf.Caches); got != tc.want { t.Fatalf("cache entries = %d, want %d", got, tc.want) } }) }}
func TestEngineInitWorkflowCarriesRepositoryIdentity(t *testing.T) { m := &Mill{l: slog.New(slog.NewTextHandler(io.Discard, nil))} wf, err := NewEngine("dummy", m).InitWorkflow(tangled.Pipeline_Workflow{Name: "build"}, testPipeline()) if err != nil { t.Fatal(err) } if wf.OwnerDID != "did:plc:testowner" || wf.RepoDID != "did:plc:testrepo" { t.Fatalf("workflow identity = %q/%q, want source owner and repository", wf.OwnerDID, wf.RepoDID) }}
func testWorkflowWithRunsOn(name string, runsOn []string) *models.Workflow { wf := testWorkflow(name) wf.Data.(*millWorkflowState).RawWorkflow.RunsOn = runsOn return wf}
func addCandidateSession(t *testing.T, m *Mill, nodeID string, labels []string, load float64, enc messageEncoder) *millSession { t.Helper() if enc == nil { enc = scriptedEncoder(func(*millproto.Message) error { return nil }) } sess := newSession(nodeID, "inc-"+nodeID, labels, enc, slog.New(slog.NewTextHandler(io.Discard, nil))) sess.snapshot = &millv1.NodeSnapshot{ Seqno: 1, Engines: map[string]*millv1.EngineAvailability{ "dummy": {Available: load < 1.0, Load: map[string]float64{"slots": load}}, }, } m.mu.Lock() m.sessions[nodeID] = sess m.mu.Unlock() return sess}
func assertRankedNodes(t *testing.T, got []*millSession, want []string) { t.Helper() if len(got) != len(want) { t.Fatalf("rankCandidates() returned %d candidates, want %d: got %v want %v", len(got), len(want), sessionIDs(got), want) } for i := range want { if got[i].nodeID != want[i] { t.Fatalf("rankCandidates()[%d] = %q, want %q; full order got %v want %v", i, got[i].nodeID, want[i], sessionIDs(got), want) } }}
func sessionIDs(sessions []*millSession) []string { out := make([]string, len(sessions)) for i, sess := range sessions { out[i] = sess.nodeID } return out}
func sameStringMultiset(a, b []string) bool { if len(a) != len(b) { return false } counts := make(map[string]int, len(a)) for _, s := range a { counts[s]++ } for _, s := range b { if counts[s] == 0 { return false } counts[s]-- } return true}
type reserveReply struct { accepted bool rejectClass millv1.RejectClass reason string failureClass string failureReason string}
func addReplyingCandidateSession(t *testing.T, m *Mill, nodeID string, labels []string, load float64, asked chan<- string, reply reserveReply) *millSession { t.Helper() var sess *millSession sess = addCandidateSession(t, m, nodeID, labels, load, scriptedEncoder(func(msg *millproto.Message) error { rs := msg.GetReserveSeat() if rs == nil { return nil } if asked != nil { asked <- nodeID } sess.deliver(rs.GetLeaseId(), &millproto.Message{ReserveResult: &millv1.ReserveResult{ LeaseId: rs.GetLeaseId(), Accepted: reply.accepted, RejectReason: reply.reason, RejectClass: reply.rejectClass, FailureClass: reply.failureClass, FailureReason: reply.failureReason, }}) return nil })) return sess}
func drainAsked(ch <-chan string) []string { var out []string for { select { case nodeID := <-ch: out = append(out, nodeID) default: return out } }}
func TestCommitRetriesAfterSessionCloseBeforeCommitted(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{BidTimeout: 25 * time.Millisecond, ReconnectGrace: time.Second}) wf := testWorkflow("build") restoreID := "11111111-1111-4111-8111-111111111111" saveID := "22222222-2222-4222-8222-222222222222" wf.CacheBindings = []models.CacheBinding{{ EntryIndex: 0, Hash: "abc123", RestoreID: restoreID, RestoreKey: "objects/did:web:example.com/" + restoreID, RestoreName: "deps-abc123", SaveID: saveID, SaveKey: "objects/did:web:example.com/" + saveID, }} wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = wid wf.Data.(*millWorkflowState).Lease = lease
m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
var sess1 *millSession firstCommit := make(chan struct{}) sess1 = newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error { if msg.GetCommitLease() != nil { close(firstCommit) m.detachSession(sess1) } return nil }), l) m.attachSession(sess1)
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second) defer cancel() done := make(chan error, 1) go func() { done <- m.commitAndWait(ctx, wf, nil) }()
select { case <-firstCommit: case <-ctx.Done(): t.Fatal("first commit was not sent") }
var sess2 *millSession sess2 = newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error { if msg.GetCommitLease() == nil { return nil } bindings := msg.GetCommitLease().GetCacheBindings() if len(bindings) != 1 || bindings[0].GetRestoreId() != restoreID || bindings[0].GetSaveId() != saveID { t.Errorf("CommitLease cache bindings = %+v", bindings) } leaseID := msg.GetCommitLease().GetLeaseId() sess2.deliver(leaseID, &millproto.Message{Committed: &millv1.Committed{LeaseId: leaseID}}) _ = m.onEventBatch(sess2, &millv1.EventBatch{ Epoch: sess2.epoch, Events: []*millv1.Event{ { Seqno: 1, LeaseId: leaseID, Payload: &millv1.Event_AttemptResult{ AttemptResult: &millv1.AttemptResult{ Status: millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS, }, }, }, }, }) return nil }), l) m.attachSession(sess2) m.sessionReady(sess2)
select { case err := <-done: if err != nil { t.Fatalf("commitAndWait() error = %v, want success after reconnect", err) } case <-ctx.Done(): t.Fatal("commitAndWait() did not finish after reconnect") }}
func TestDestroyRunningLeaseDoesNotDropCancelledTerminal(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{}) wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = wid lease.setState(leaseRunning)
m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(wid) if lease.getState() == leaseDone { t.Fatal("destroy sealed the lease before the terminal result") }
lease.deliverTerminal(&millv1.AttemptResult{ Status: millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED, }) res := <-lease.terminal if err := terminalError(res.Status); !errors.Is(err, engine.ErrWorkflowCanceled) { t.Fatalf("terminalError() = %v, want ErrCancelled", err) }}
func TestPlaceBlocksWhenNoCapacity(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{}) wf := testWorkflow("build") wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
// with no executors at all, place must block until ctx expires and the user sees pending ctx, cancel := context.WithTimeout(context.Background(), 150*time.Millisecond) defer cancel()
_, err := m.place(ctx, "dummy", wid, wf) if err != context.DeadlineExceeded { t.Fatalf("place() error = %v, want DeadlineExceeded", err) }}
func TestRankCandidatesFiltersRequiredLabelsWithANDSemantics(t *testing.T) { m := New(slog.New(slog.NewTextHandler(io.Discard, nil)), Config{}) addCandidateSession(t, m, "linux-high", []string{"linux"}, 0.0, nil) addCandidateSession(t, m, "linux-arm", []string{"linux", "arm64"}, 0.25, nil) addCandidateSession(t, m, "unlabeled", nil, 0.5, nil) addCandidateSession(t, m, "linux-arm-gpu", []string{"linux", "arm64", "gpu"}, 0.75, nil) addCandidateSession(t, m, "linux-arm-full", []string{"linux", "arm64"}, 1.0, nil)
tests := []struct { name string requiredLabels []string want []string }{ { name: "no required labels keeps old capacity ranking", want: []string{"linux-high", "linux-arm", "unlabeled", "linux-arm-gpu"}, }, { name: "single required label includes every candidate carrying it", requiredLabels: []string{"linux"}, want: []string{"linux-high", "linux-arm", "linux-arm-gpu"}, }, { name: "all required labels must be present", requiredLabels: []string{"linux", "arm64"}, want: []string{"linux-arm", "linux-arm-gpu"}, }, { name: "one missing required label excludes the candidate", requiredLabels: []string{"linux", "arm64", "gpu"}, want: []string{"linux-arm-gpu"}, }, { name: "unknown required label leaves no candidate", requiredLabels: []string{"linux", "arm64", "metal"}, }, }
for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { assertRankedNodes(t, m.rankCandidates("dummy", tt.requiredLabels, false), tt.want) }) }}
func TestRankCandidatesRequiresSharedCacheStore(t *testing.T) { m := New(discardLogger(), Config{CacheStoreID: "disk:/cache"}) matching := addCandidateSession(t, m, "matching", []string{"linux"}, 0.5, nil) matching.cacheStoreID = "disk:/cache" other := addCandidateSession(t, m, "other", []string{"linux"}, 0.1, nil) other.cacheStoreID = "s3:bucket/cache"
assertRankedNodes(t, m.rankCandidates("dummy", []string{"linux"}, true), []string{"matching"}) assertRankedNodes(t, m.rankCandidates("dummy", []string{"linux"}, false), []string{"other", "matching"})}
func TestPlaceWithMissingRequiredLabelsStaysPendingWithoutReserve(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{BidTimeout: 10 * time.Millisecond}) reserveSent := make(chan struct{}, 1) addCandidateSession(t, m, "linux-only", []string{"linux"}, 0.75, scriptedEncoder(func(msg *millproto.Message) error { if msg.GetReserveSeat() != nil { select { case reserveSent <- struct{}{}: default: } } return nil })) wf := testWorkflowWithRunsOn("build", []string{"linux", "arm64"}) wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}
ctx, cancel := context.WithTimeout(context.Background(), 120*time.Millisecond) defer cancel() _, err := m.place(ctx, "dummy", wid, wf) if err != context.DeadlineExceeded { t.Fatalf("place() error = %v, want DeadlineExceeded while job remains pending", err) } select { case <-reserveSent: t.Fatal("place() sent ReserveSeat to executor missing a required label") default: }}
func TestMaxPendingRejects(t *testing.T) { l := slog.New(slog.NewTextHandler(io.Discard, nil)) m := New(l, Config{MaxPending: 1})
m.mu.Lock() m.pending = 1 m.mu.Unlock()
wf2 := testWorkflow("b") _, err := m.place(context.Background(), "dummy", models.WorkflowId{Name: "b"}, wf2) if err == nil { t.Fatal("place() past maxPending should error") }}
func TestCancelledRunningLeaseSurvivesReleaseForReconnectReplay(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = wid lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(wid) slot := &millSlot{fleet: m, lease: lease} slot.Release()
m.mu.Lock() _, retained := m.leases[lease.id] m.mu.Unlock() if !retained { t.Fatal("slot release removed a cancellation-requested running lease before its terminal result") } if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 1 { t.Fatalf("durable leases after slot release = %+v, err = %v; want retained lease", rows, err) }
var sentMu sync.Mutex var sent []*millproto.Message sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error { sentMu.Lock() sent = append(sent, msg) sentMu.Unlock() return nil }), discardLogger()) if _, ok := m.attachSession(sess); !ok { t.Fatal("attachSession rejected reconnect") } m.sessionReady(sess) sentMu.Lock() var replayed bool for _, msg := range sent { if cancel := msg.GetCancelAttempt(); cancel != nil && cancel.GetLeaseId() == lease.id { replayed = true } } sentMu.Unlock() if !replayed { t.Fatal("reconnect did not replay CancelAttempt for retained lease") }
if err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{ { Seqno: 1, LeaseId: lease.id, Payload: &millv1.Event_AttemptResult{ AttemptResult: &millv1.AttemptResult{ Status: millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED, }, }, }, }, }); err != nil { t.Fatalf("onEventBatch: %v", err) } m.mu.Lock() _, retained = m.leases[lease.id] m.mu.Unlock() if retained { t.Fatal("terminal result did not clean retained cancelled lease") } if rows, err := bdb.ListMillLeases(); err != nil || len(rows) != 0 { t.Fatalf("durable leases after terminal = %+v, err = %v; want none", rows, err) } slot.Release()}
func TestSessionRequestCancelledBeforeRegistrationOrSend(t *testing.T) { sent := 0 sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(*millproto.Message) error { sent++ return nil }), discardLogger()) ctx, cancel := context.WithCancel(context.Background()) cancel() _, err := sess.request(ctx, "lease-1", &millproto.Message{ReleaseLease: &millv1.ReleaseLease{LeaseId: "lease-1"}}) if !errors.Is(err, context.Canceled) { t.Fatalf("request error = %v, want context.Canceled", err) } if sent != 0 { t.Fatalf("request sent %d messages for an already-cancelled context, want 0", sent) } sess.mu.Lock() pending := len(sess.pending) sess.mu.Unlock() if pending != 0 { t.Fatalf("request left %d pending waiters, want 0", pending) }}
func TestAttachSessionReplacesSilentIncumbentButRejectsActiveDuplicate(t *testing.T) { m := New(discardLogger(), Config{ReconnectGrace: time.Minute}) old := newSession("node-1", "inc-old", nil, nopEncoder(), discardLogger()) transportClosed := make(chan struct{}) old.closeTransport = func() error { close(transportClosed) return nil } if _, ok := m.attachSession(old); !ok { t.Fatal("first attach rejected") } m.mu.Lock() old.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace) m.mu.Unlock()
replacement := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(replacement); !ok { t.Fatal("silent incumbent blocked authenticated replacement") } select { case <-transportClosed: default: t.Fatal("replacing a silent incumbent did not close its transport") } m.mu.Lock() replacement.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace) m.mu.Unlock() if err := replacement.dispatch(m, &millproto.Message{NodeSnapshot: &millv1.NodeSnapshot{Seqno: 1}}); err != nil { t.Fatalf("periodic snapshot dispatch: %v", err) } if _, ok := m.attachSession(newSession("node-1", "inc-dup", nil, nopEncoder(), discardLogger())); ok { t.Fatal("active replacement did not reject a duplicate session") }}
func TestSilentSessionReplacementWithoutSnapshotFailsLeaseAtDeadline(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) old := newSession("node-1", "inc-old", nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(old); !ok { t.Fatal("first attach rejected") } lease := newLease("lease-1", old.nodeID, old.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease old.lastSeen = time.Now().Add(-2 * m.cfg.ReconnectGrace) m.mu.Unlock()
replacement := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger()) attachStarted := time.Now() if _, ok := m.attachSession(replacement); !ok { t.Fatal("silent incumbent blocked replacement") } attachFinished := time.Now() if !replacement.recovering || replacement.graceTimer == nil || replacement.graceDeadline.IsZero() { t.Fatal("replacement without snapshot has no recovery deadline") } if replacement.graceDeadline.Before(attachStarted.Add(m.cfg.ReconnectGrace)) || replacement.graceDeadline.After(attachFinished.Add(m.cfg.ReconnectGrace)) { t.Fatalf("replacement deadline = %v, want one reconnect grace after attach", replacement.graceDeadline) } m.failLeasesAfterGrace(replacement)
m.mu.Lock() _, stillActive := m.leases[lease.id] m.mu.Unlock() if stillActive { t.Fatal("replacement without snapshot stranded incumbent lease") }}
func TestPlaceReleasesRemoteReservationWhenInitialPersistenceFails(t *testing.T) { m, bdb := restoreTestMill(t, Config{BidTimeout: time.Second}) if err := bdb.Close(); err != nil { t.Fatalf("close db: %v", err) } released := make(chan string, 1) var sess *millSession sess = addCandidateSession(t, m, "node-1", nil, 0, scriptedEncoder(func(msg *millproto.Message) error { switch { case msg.GetReserveSeat() != nil: leaseID := msg.GetReserveSeat().GetLeaseId() sess.deliver(leaseID, &millproto.Message{ReserveResult: &millv1.ReserveResult{ LeaseId: leaseID, Accepted: true, }}) case msg.GetReleaseLease() != nil: released <- msg.GetReleaseLease().GetLeaseId() } return nil }))
ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() slot, err := m.place( ctx, "dummy", models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, testWorkflow("build"), ) if err == nil { t.Fatal("place succeeded after reserved lease persistence failed") } if slot != nil { t.Fatalf("place returned slot %T after persistence failure", slot) } select { case leaseID := <-released: if leaseID == "" { t.Fatal("ReleaseLease had empty lease id") } case <-time.After(time.Second): t.Fatal("persistence failure did not compensate with ReleaseLease") } m.mu.Lock() leases := len(m.leases) m.mu.Unlock() if leases != 0 { t.Fatalf("mill published %d leases after initial persistence failure, want 0", leases) }}
func TestGapsAndDuplicates(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess)
owned := newLease("lease-1", "node-1", "inc-1", "dummy") owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock()
err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{ { Seqno: 0, LeaseId: owned.id, Payload: &millv1.Event_StatusEvent{ StatusEvent: &millv1.StatusEvent{Status: millv1.NonterminalStatus_NONTERMINAL_STATUS_RUNNING}, }, }, }, }) if err != nil { t.Fatalf("expected duplicate to be skipped without error, got: %v", err) }
err = m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{ { Seqno: 2, LeaseId: owned.id, Payload: &millv1.Event_StatusEvent{ StatusEvent: &millv1.StatusEvent{Status: millv1.NonterminalStatus_NONTERMINAL_STATUS_RUNNING}, }, }, }, }) if err == nil { t.Fatal("expected error due to seqno gap, got nil") }}
func TestAtomicBatchRollback(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess)
owned := newLease("lease-1", "node-1", "inc-1", "dummy") owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock()
if _, err := bdb.Exec(` create trigger reject_status_event before insert on events begin select raise(abort, 'forced status event failure'); end `); err != nil { t.Fatalf("failed to create fail trigger: %v", err) } defer bdb.Exec("drop trigger reject_status_event")
err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{ { Seqno: 1, LeaseId: owned.id, Payload: &millv1.Event_StatusEvent{ StatusEvent: &millv1.StatusEvent{Status: millv1.NonterminalStatus_NONTERMINAL_STATUS_RUNNING}, }, }, }, }) if err == nil { t.Fatal("expected status event insertion to fail due to trigger") }
m.mu.Lock() seqno := m.nodeSeqno["node-1/inc-1"] m.mu.Unlock() if seqno != 0 { t.Fatalf("expected seqno 0 due to rollback, got %d", seqno) }}
func TestCacheUpdatesCommitWithStreamCursor(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) cache, err := storage.NewDisk(t.TempDir()) if err != nil { t.Fatal(err) } m.AttachCache(cache)
at := time.Date(2026, 8, 11, 20, 0, 0, 0, time.UTC) old := db.CacheEntry{ ID: "old", StorageKey: "objects/old", OwnerDID: "did:plc:owner", RepoDID: "did:plc:repo", Engine: "dummy", CacheKey: "deps", CacheHash: "hash", SizeBytes: 3, State: "ready", CreatedAt: at.Add(-time.Hour), LastUsedAt: at.Add(-time.Hour), } pending := old pending.ID = "pending" pending.StorageKey = "objects/pending" pending.SizeBytes = 0 pending.State = "pending" pending.CreatedAt = at pending.LastUsedAt = at for _, entry := range []db.CacheEntry{old, pending} { if err := bdb.InsertCacheEntry(context.Background(), entry); err != nil { t.Fatal(err) } if err := cache.Put(context.Background(), entry.StorageKey, bytes.NewReader([]byte(entry.ID))); err != nil { t.Fatal(err) } }
sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess) owned := newLease("lease-1", "node-1", "inc-1", "dummy") owned.wid = models.WorkflowId{ PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build", } if err := m.persistLease(owned, leaseRowRunning); err != nil { t.Fatal(err) } if err := bdb.SaveMillCacheCapabilities(owned.id, []db.MillCacheCapability{ {Action: "restore", CacheID: old.ID, StorageKey: old.StorageKey}, {Action: "save", CacheID: pending.ID, StorageKey: pending.StorageKey}, }); err != nil { t.Fatal(err) } m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock()
batch := &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{ { Seqno: 1, LeaseId: owned.id, Payload: &millv1.Event_CacheUpdate{CacheUpdate: &millv1.CacheUpdate{ Action: millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_USED, Id: old.ID, }}, }, { Seqno: 2, LeaseId: owned.id, Payload: &millv1.Event_CacheUpdate{CacheUpdate: &millv1.CacheUpdate{ Action: millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_STORED, Id: pending.ID, SizeBytes: int64(len(pending.ID)), }}, }, }, } if err := m.onEventBatch(sess, batch); err != nil { t.Fatalf("onEventBatch: %v", err) } if err := m.onEventBatch(sess, batch); err != nil { t.Fatalf("replay onEventBatch: %v", err) }
ready, err := bdb.FindCacheEntry(context.Background(), old.RepoDID, old.Engine, old.CacheKey, old.CacheHash) if err != nil { t.Fatal(err) } if ready.ID != pending.ID || ready.SizeBytes != int64(len(pending.ID)) { t.Fatalf("ready cache = %+v", ready) }
waitForCacheDeletion(t, cache, nil, old.StorageKey) if cursor, err := bdb.GetExecutorCursor(sess.nodeID, sess.epoch); err != nil || cursor != 2 { t.Fatalf("cursor = %d, %v; want 2", cursor, err) }}func TestCacheUpdateCannotDeleteUnplannedObject(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) cache, err := storage.NewDisk(t.TempDir()) if err != nil { t.Fatal(err) } m.AttachCache(cache) victim := "objects/did:web:victim.example/11111111-1111-4111-8111-111111111111" if err := cache.Put(context.Background(), victim, strings.NewReader("victim")); err != nil { t.Fatal(err) }
sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess) lease := newLease("lease-1", sess.nodeID, sess.epoch, "dummy") lease.wid = models.WorkflowId{ PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build", } if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatal(err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
err = m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{{ Seqno: 1, LeaseId: lease.id, Payload: &millv1.Event_CacheUpdate{CacheUpdate: &millv1.CacheUpdate{ Action: millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_STORED, Id: "unknown", }}, }}, }) if err == nil { t.Fatal("unplanned cache update succeeded") } reader, err := cache.Get(context.Background(), victim) if err != nil { t.Fatalf("victim object was deleted: %v", err) } _ = reader.Close() if cursor, err := bdb.GetExecutorCursor(sess.nodeID, sess.epoch); err != nil || cursor != 0 { t.Fatalf("cursor = %d, %v; want 0", cursor, err) }}
func TestMissingPendingRowQueuesAndDeletesPlannedObject(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) cache, err := storage.NewDisk(t.TempDir()) if err != nil { t.Fatal(err) } m.AttachCache(cache) id := "11111111-1111-4111-8111-111111111111" key := "objects/did:web:example.com/" + id if err := cache.Put(context.Background(), key, strings.NewReader("orphan")); err != nil { t.Fatal(err) }
sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess) lease := newLease("lease-1", sess.nodeID, sess.epoch, "dummy") lease.wid = models.WorkflowId{ PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build", } if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatal(err) } if err := bdb.SaveMillCacheCapabilities(lease.id, []db.MillCacheCapability{{ Action: "save", CacheID: id, StorageKey: key, }}); err != nil { t.Fatal(err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
if err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{{ Seqno: 1, LeaseId: lease.id, Payload: &millv1.Event_CacheUpdate{CacheUpdate: &millv1.CacheUpdate{ Action: millv1.CacheUpdateAction_CACHE_UPDATE_ACTION_STORED, Id: id, }}, }}, }); err != nil { t.Fatal(err) } waitForCacheDeletion(t, cache, bdb, key)}
func TestTerminalBeforeACK(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
ackSent := make(chan struct{}) sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(msg *millproto.Message) error { if msg.GetAck() != nil { wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} st, err := bdb.GetStatus(wid) if err != nil || st.Status != "success" { t.Errorf("expected terminal status success at ACK time, got status: %v, err: %v", st, err) } close(ackSent) } return nil }), discardLogger()) m.attachSession(sess)
owned := newLease("lease-1", "node-1", "inc-1", "dummy") owned.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} m.mu.Lock() m.leases[owned.id] = owned m.mu.Unlock()
err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{ { Seqno: 1, LeaseId: owned.id, Payload: &millv1.Event_AttemptResult{ AttemptResult: &millv1.AttemptResult{Status: millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS}, }, }, }, }) if err != nil { t.Fatalf("onEventBatch: %v", err) }
select { case <-ackSent: case <-time.After(2 * time.Second): t.Fatal("ACK was not sent") }}
func TestTerminalWithIncompleteIdentityAdvancesCursor(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) sess := newSession("node-1", "inc-1", nil, scriptedEncoder(func(*millproto.Message) error { return nil }), discardLogger()) m.attachSession(sess)
lease := newLease("lease-1", sess.nodeID, sess.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Rkey: "r"}, Name: "build"} m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{{ Seqno: 1, LeaseId: lease.id, Payload: &millv1.Event_AttemptResult{AttemptResult: &millv1.AttemptResult{ Status: millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS, LogArtifact: &millv1.LogArtifact{ Ref: "logs/" + lease.id + ".log", Hash: "sha256:test", }, }}, }}, }) if err != nil { t.Fatalf("onEventBatch: %v", err) }
m.mu.Lock() seqno := m.nodeSeqno[sess.nodeID+"/"+sess.epoch] m.mu.Unlock() if seqno != 1 { t.Fatalf("cursor = %d, want 1", seqno) } var artifacts int if err := bdb.QueryRow(`select count(*) from mill_artifacts`).Scan(&artifacts); err != nil { t.Fatal(err) } if artifacts != 0 { t.Fatalf("artifact rows = %d, want 0", artifacts) }}
func TestExecutorRestartEmptySnapshot(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
lease := newLease("lease-1", "node-1", "inc-old", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock() if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) }
sess := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger()) m.attachSession(sess)
err := m.onSnapshot(sess, &millv1.NodeSnapshot{ Seqno: 1, ActiveLeaseIds: nil, }) if err != nil { t.Fatalf("onSnapshot: %v", err) }
m.mu.Lock() _, stillActive := m.leases["lease-1"] m.mu.Unlock() if stillActive { t.Fatal("expected old epoch lease to be reconciled and failed") }}
func TestReplacementLostBeforeSnapshotFailsOldEpochLease(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) lease := newLease("lease-1", "node-1", "inc-old", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
replacement := newSession("node-1", "inc-new", nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(replacement); !ok { t.Fatal("attachSession rejected replacement") } m.detachSession(replacement) m.failLeasesAfterGrace(replacement)
m.mu.Lock() _, stillActive := m.leases[lease.id] m.mu.Unlock() if stillActive { t.Fatal("replacement loss stranded old-epoch lease") }}
func TestSeqRegression(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess)
err := m.onSnapshot(sess, &millv1.NodeSnapshot{ Seqno: 5, }) if err != nil { t.Fatalf("first snapshot: %v", err) }
err = m.onSnapshot(sess, &millv1.NodeSnapshot{ Seqno: 4, }) if err == nil { t.Fatal("expected seqno regression to be rejected") }}
func TestClaimedSweep(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.orphaned = true lease.claimed = false m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess) err := m.onSnapshot(sess, &millv1.NodeSnapshot{ Seqno: 1, ActiveLeaseIds: []string{"lease-1"}, }) if err != nil { t.Fatalf("onSnapshot: %v", err) }
m.detachSession(sess)
m.sweepUnclaimedOrphans()
m.mu.Lock() _, stillRunning := m.leases["lease-1"] m.mu.Unlock() if !stillRunning { t.Fatal("claimed lease was incorrectly swept by startup sweep") }}
func TestUnacknowledgedCancelClosesSession(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: 50 * time.Millisecond}) registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil)
sessClosed := make(chan struct{}) sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) sess.closeTransport = func() error { close(sessClosed) return nil } m.attachSession(sess)
lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(lease.wid)
select { case <-sessClosed: case <-time.After(2 * time.Second): t.Fatal("session was not closed after the cancel went unacknowledged") } if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || !ok { t.Fatalf("executor identity rejected after session timeout: ok = %v, err = %v", ok, err) }}
func TestPendingCancelBlocksPlacementUntilAcknowledged(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) sess := addCandidateSession(t, m, "node-1", nil, 0, nil) lease := newLease("lease-1", "node-1", sess.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
assertRankedNodes(t, m.rankCandidates("dummy", nil, false), []string{"node-1"}) m.destroy(lease.wid) assertRankedNodes(t, m.rankCandidates("dummy", nil, false), nil)
m.onCancelAck(sess, &millv1.CancelAck{LeaseId: lease.id}) assertRankedNodes(t, m.rankCandidates("dummy", nil, false), []string{"node-1"})}
func TestTerminalSettlesPendingCancelPlacementGate(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) sess := addCandidateSession(t, m, "node-1", nil, 0, nil) lease := newLease("lease-1", sess.nodeID, sess.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(lease.wid) assertRankedNodes(t, m.rankCandidates("dummy", nil, false), nil) if err := m.onEventBatch(sess, &millv1.EventBatch{ Epoch: sess.epoch, Events: []*millv1.Event{{ Seqno: 1, LeaseId: lease.id, Payload: &millv1.Event_AttemptResult{AttemptResult: &millv1.AttemptResult{ Status: millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED, }}, }}, }); err != nil { t.Fatalf("onEventBatch: %v", err) } assertRankedNodes(t, m.rankCandidates("dummy", nil, false), []string{"node-1"})}
func TestCancelSentAfterTerminalDoesNotLeakPlacementGate(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) sess := addCandidateSession(t, m, "node-1", nil, 0, nil) lease := newLease("lease-1", sess.nodeID, sess.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
lease.requestCancel("done") if err := m.finishLiveLease(lease, string(models.StatusKindCancelled), "done"); err != nil { t.Fatalf("finishLiveLease: %v", err) } m.sendCancel(sess, lease, "done")
m.mu.Lock() pending := len(sess.pendingCancels) m.mu.Unlock() if pending != 0 { t.Fatalf("pending cancels after terminal = %d, want 0", pending) } assertRankedNodes(t, m.rankCandidates("dummy", nil, false), []string{"node-1"})}
func TestValidSnapshotRecoversAfterProtocolError(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) first := addCandidateSession(t, m, "node-1", nil, 0, nil) lease := newLease("lease-1", first.nodeID, first.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.noteSessionError(first, protoErrf("bad stream")) m.detachSession(first) replacement := newSession(first.nodeID, first.epoch, nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(replacement); !ok { t.Fatal("reconnect was rejected") } m.sessionReady(replacement) if err := m.onSnapshot(replacement, &millv1.NodeSnapshot{ Seqno: 1, ActiveLeaseIds: []string{lease.id}, Engines: map[string]*millv1.EngineAvailability{ "dummy": {Available: true}, }, }); err != nil { t.Fatalf("onSnapshot: %v", err) }
m.mu.Lock() recovering := replacement.recovering m.mu.Unlock() if recovering { t.Fatal("valid reconnect snapshot did not end recovery") } assertRankedNodes(t, m.rankCandidates("dummy", nil, false), []string{"node-1"})}
func TestReconnectRetriesKeepOriginalFailureDeadline(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) first := addCandidateSession(t, m, "node-1", nil, 0, nil) lease := newLease("lease-1", first.nodeID, first.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.detachSession(first) deadline := first.graceDeadline second := newSession(first.nodeID, first.epoch, nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(second); !ok { t.Fatal("first reconnect was rejected") } m.detachSession(second) third := newSession(first.nodeID, first.epoch, nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(third); !ok { t.Fatal("second reconnect was rejected") } if !third.graceDeadline.Equal(deadline) { t.Fatalf("reconnect moved failure deadline from %v to %v", deadline, third.graceDeadline) }
m.failLeasesAfterGrace(third) m.mu.Lock() _, tracked := m.leases[lease.id] m.mu.Unlock() if tracked { t.Fatal("reconnect retries kept the lease past its original failure deadline") }}
func TestUnacknowledgedCancelReconnectsReachBoundedTerminal(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) first := addCandidateSession(t, m, "node-1", nil, 0, nil) lease := newLease("lease-1", first.nodeID, first.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(lease.wid) m.checkCancelAck(first, lease) m.detachSession(first) deadline := first.graceDeadline
second := newSession(first.nodeID, first.epoch, nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(second); !ok { t.Fatal("first reconnect was rejected") } m.sessionReady(second) if err := m.onSnapshot(second, &millv1.NodeSnapshot{ Seqno: 1, ActiveLeaseIds: []string{lease.id}, Engines: map[string]*millv1.EngineAvailability{ "dummy": {Available: true}, }, }); err != nil { t.Fatalf("onSnapshot: %v", err) } if !second.recovering { t.Fatal("snapshot cleared recovery while cancel was still pending") } m.checkCancelAck(second, lease) m.detachSession(second)
third := newSession(first.nodeID, first.epoch, nil, nopEncoder(), discardLogger()) if _, ok := m.attachSession(third); !ok { t.Fatal("second reconnect was rejected") } m.sessionReady(third) if !third.graceDeadline.Equal(deadline) { t.Fatalf("cancel reconnect moved failure deadline from %v to %v", deadline, third.graceDeadline) } m.failLeasesAfterGrace(third)
select { case result := <-lease.terminal: if result.GetStatus() != millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED { t.Fatalf("terminal status = %v, want CANCELLED", result.GetStatus()) } case <-time.After(time.Second): t.Fatal("unacknowledged cancel did not reach a bounded terminal") }}
func TestAckedCancelFinishesOnlyItsLeaseAfterTeardownDeadline(t *testing.T) { m, bdb := restoreTestMill(t, Config{ ReconnectGrace: time.Minute, CancelAckTimeout: 30 * time.Millisecond, CancelTeardownTimeout: 150 * time.Millisecond, }) registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil)
sessClosed := make(chan struct{}) sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) sess.closeTransport = func() error { close(sessClosed) return nil } m.attachSession(sess)
lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(lease.wid) m.onCancelAck(sess, &millv1.CancelAck{LeaseId: lease.id})
time.Sleep(80 * time.Millisecond) select { case <-sessClosed: t.Fatal("session was closed while an acked cancel was still tearing down") default: }
select { case res := <-lease.terminal: if res.GetStatus() != millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED { t.Fatalf("terminal status = %v, want CANCELLED", res.GetStatus()) } case <-time.After(2 * time.Second): t.Fatal("teardown deadline passed without the mill finishing the lease") }
select { case <-sessClosed: t.Fatal("teardown deadline closed the session, failing the node's other jobs") default: } if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || !ok { t.Fatalf("executor identity rejected for a slow teardown: ok = %v, err = %v", ok, err) }}
func TestLateCancelAckTimerKeepsReconnectGrace(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil)
sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(sess) if err := m.onSnapshot(sess, &millv1.NodeSnapshot{Seqno: 1}); err != nil { t.Fatalf("onSnapshot: %v", err) }
lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persistLease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
m.destroy(lease.wid) m.onCancelAck(sess, &millv1.CancelAck{LeaseId: lease.id}) m.detachSession(sess) m.checkCancelAck(sess, lease)
m.mu.Lock() recovering := sess.recovering timerArmed := sess.graceTimer != nil m.mu.Unlock() if !recovering || !timerArmed { t.Fatal("late cancel timer disarmed reconnect grace") } m.failLeasesAfterGrace(sess) select { case result := <-lease.terminal: if result.GetStatus() != millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED { t.Fatalf("terminal status = %v, want CANCELLED", result.GetStatus()) } case <-time.After(time.Second): t.Fatal("reconnect grace did not finish the disconnected lease") } if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || !ok { t.Fatalf("executor identity rejected for a cancel it could not answer: ok = %v, err = %v", ok, err) }}
func TestCancelAckDeadlineDoesNotCloseReplacementSession(t *testing.T) { m, _ := restoreTestMill(t, Config{ReconnectGrace: time.Minute, CancelAckTimeout: time.Minute}) oldSession := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) m.attachSession(oldSession) lease := newLease("lease-1", "node-1", oldSession.epoch, "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} lease.setState(leaseRunning) m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock() m.destroy(lease.wid) m.detachSession(oldSession)
replacementClosed := make(chan struct{}) replacement := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) replacement.closeTransport = func() error { close(replacementClosed) return nil } if _, ok := m.attachSession(replacement); !ok { t.Fatal("replacement session was rejected") } replacement.snapshot = &millv1.NodeSnapshot{ Seqno: 1, Engines: map[string]*millv1.EngineAvailability{ "dummy": {Available: true}, }, } m.sessionReady(replacement) assertRankedNodes(t, m.rankCandidates("dummy", nil, false), nil)
m.checkCancelAck(oldSession, lease) select { case <-replacementClosed: t.Fatal("cancel timer from old session closed its replacement") default: }}
func TestCleanupRetry(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute})
lease := newLease("lease-1", "node-1", "inc-1", "dummy") lease.wid = models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"} if err := m.persistLease(lease, leaseRowRunning); err != nil { t.Fatalf("persist lease: %v", err) } m.mu.Lock() m.leases[lease.id] = lease m.mu.Unlock()
if _, err := bdb.Exec(` create trigger reject_cleanup_delete before delete on mill_leases begin select raise(abort, 'forced delete failure'); end `); err != nil { t.Fatalf("failed to create fail trigger: %v", err) }
err := m.cleanupLease(lease) if err == nil { t.Fatal("expected cleanupLease to fail") }
m.mu.Lock() _, stillRunning := m.leases["lease-1"] m.mu.Unlock() if !stillRunning { t.Fatal("lease was removed from memory despite cleanup failure") }
if _, err := bdb.Exec("drop trigger reject_cleanup_delete"); err != nil { t.Fatalf("drop trigger: %v", err) }
err = m.cleanupLease(lease) if err != nil { t.Fatalf("expected retry cleanup to succeed, got: %v", err) }
m.mu.Lock() _, stillRunning = m.leases["lease-1"] m.mu.Unlock() if stillRunning { t.Fatal("lease still in memory after successful cleanup retry") }}
func TestBidPreservesConsistentTypedIncompatibility(t *testing.T) { m := New(discardLogger(), Config{TopK: 1, BidTimeout: time.Second}) addReplyingCandidateSession(t, m, "node-1", nil, 0, nil, reserveReply{ rejectClass: millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE, reason: "invalid workflow", failureClass: string(engine.FailureClassUser), failureReason: string(engine.FailureReasonWorkflowInvalid), })
wid := models.WorkflowId{ PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "pipeline"}, Name: "build", } _, err := m.bid(context.Background(), "dummy", wid, testWorkflow("build")) if err == nil { t.Fatal("bid succeeded despite incompatible executor") } var failure *engine.WorkflowFailure if !errors.As(err, &failure) { t.Fatalf("bid error %T did not preserve typed attribution", err) } if failure.Class != engine.FailureClassUser || failure.Reason != engine.FailureReasonWorkflowInvalid { t.Fatalf("failure attribution = %s/%s", failure.Class, failure.Reason) }}
func TestBoundedBidding(t *testing.T) { m := New(discardLogger(), Config{TopK: 2, BidTimeout: 10 * time.Millisecond})
addCandidateSession(t, m, "node-1", nil, 0, nil) addCandidateSession(t, m, "node-2", nil, 0, nil) addCandidateSession(t, m, "node-3", nil, 0, nil) addCandidateSession(t, m, "node-4", nil, 0, nil) addCandidateSession(t, m, "node-5", nil, 0, nil)
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) defer cancel()
lease, err := m.bid(ctx, "dummy", models.WorkflowId{}, testWorkflow("build")) if err != nil { t.Fatalf("bid: %v", err) } if lease != nil { t.Fatalf("did not expect a lease, got %+v", lease) }}
func TestProtocolErrorsDoNotRejectExecutorIdentity(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) registerTestExecutor(t, bdb, "node-1", HashToken("tok-1"), nil) sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger())
for range 3 { m.noteSessionError(sess, protoErrf("gap in stream seqnos. Expected 1, got 9")) } if _, _, ok, err := bdb.ResolveExecutorToken(HashToken("tok-1")); err != nil || !ok { t.Fatalf("executor identity rejected after protocol errors: ok = %v, err = %v", ok, err) }}
func TestAttachSessionRereadsDatabaseCursor(t *testing.T) { m, bdb := restoreTestMill(t, Config{ReconnectGrace: time.Minute}) if err := bdb.ApplyEventBatch(nil, func(tx *db.EventBatchTx) error { return tx.AdvanceCursor("node-1", "inc-1", 3) }); err != nil { t.Fatalf("AdvanceCursor: %v", err) }
sess := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) if resume, ok := m.attachSession(sess); !ok || resume != 3 { t.Fatalf("attachSession resume = %d, %v; want 3 from the db cursor", resume, ok) }
// an operator skip-forward must take effect on the next reconnect without // a mill restart if _, err := bdb.SetExecutorCursors("node-1", 28); err != nil { t.Fatalf("SetExecutorCursors: %v", err) } m.mu.Lock() sess.disconnected = true m.mu.Unlock() replacement := newSession("node-1", "inc-1", nil, nopEncoder(), discardLogger()) if resume, ok := m.attachSession(replacement); !ok || resume != 28 { t.Fatalf("attachSession resume after reset = %d, %v; want 28", resume, ok) }}