Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553package executor
import ( "context" "encoding/json" "errors" "fmt" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/gorilla/websocket" "google.golang.org/protobuf/proto" "io" "log/slog" "net/http" "net/http/httptest" "os" "path/filepath" "strings" "testing" "time"
"tangled.org/core/api/tangled" "tangled.org/core/notifier" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" millproto "tangled.org/core/spindle/mill/proto" millv1 "tangled.org/core/spindle/mill/proto/gen" "tangled.org/core/spindle/models" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/storage")
type captureEncoder struct { messages chan *millproto.Message}
func newCaptureEncoder() *captureEncoder { return &captureEncoder{messages: make(chan *millproto.Message, 4)}}
func (e *captureEncoder) Encode(msg *millproto.Message) error { e.messages <- msg return nil}
type blockingEncoder struct { started chan struct{} release chan struct{}}
type artifactWriterFunc func(context.Context, string, io.Reader) error
func (f artifactWriterFunc) Put(ctx context.Context, ref string, r io.Reader) error { return f(ctx, ref, r)}
func (e *blockingEncoder) Encode(*millproto.Message) error { select { case e.started <- struct{}{}: default: } <-e.release return nil}
type fakeSlot struct { released int release func()}
func (s *fakeSlot) Release() { if s.release != nil { s.release() } s.released++}
type fakeEngine struct { setupCalled bool runCalled bool destroyCalled bool acquireCalled bool acquireStarted chan struct{} releaseAcquire chan struct{} slot engine.WorkflowSlot secrets chan []secrets.UnlockedSecret done chan struct{} initErr error}
func (e *fakeEngine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) { if e.initErr != nil { return nil, e.initErr } return &models.Workflow{Name: twf.Name}, nil}func (e *fakeEngine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, l models.WorkflowLogger) error { e.setupCalled = true return nil}func (e *fakeEngine) WorkflowTimeout() time.Duration { return 7 * time.Minute }func (e *fakeEngine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) error { e.destroyCalled = true return nil}func (e *fakeEngine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, s []secrets.UnlockedSecret, l models.WorkflowLogger) error { e.runCalled = true if e.secrets != nil { e.secrets <- s } if e.done != nil { close(e.done) } return nil}func (e *fakeEngine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, mode engine.AcquireMode) (engine.WorkflowSlot, error) { e.acquireCalled = true if e.acquireStarted != nil { close(e.acquireStarted) } if e.releaseAcquire != nil { select { case <-e.releaseAcquire: case <-ctx.Done(): return nil, ctx.Err() } } if e.slot != nil { return e.slot, nil } return &fakeSlot{}, nil}
type resourceReportingEngine struct { *fakeEngine resources quota.Resources usage engine.WorkflowResourceUsage}
func (e *resourceReportingEngine) QuotaResources(*models.Workflow) quota.Resources { return e.resources}
func (e *resourceReportingEngine) WorkflowResourceUsage(*models.Workflow) (engine.WorkflowResourceUsage, bool) { return e.usage, true}
type fakeCacheEngine struct { *fakeEngine restored bool saved bool}
func (e *fakeCacheEngine) RestoreCache(context.Context, models.WorkflowId, *models.Workflow, storage.Storage, []models.CacheBinding, models.WorkflowLogger) error { e.restored = true return nil}
func (e *fakeCacheEngine) SaveCache(context.Context, models.WorkflowId, *models.Workflow, storage.Storage, []models.CacheBinding, models.WorkflowLogger) error { e.saved = true return nil}
type fakeStep struct{}
func (fakeStep) Name() string { return "test" }func (fakeStep) Command() string { return "true" }func (fakeStep) Kind() models.StepKind { return models.StepKindUser }
func testReserveSeat(t *testing.T, leaseID, engineName string) *millv1.ReserveSeat { t.Helper() twf, err := json.Marshal(tangled.Pipeline_Workflow{Name: "build"}) if err != nil { t.Fatal(err) } repoDID := "did:web:example.com" tpl, err := json.Marshal(tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &repoDID}, }}) if err != nil { t.Fatal(err) } return &millv1.ReserveSeat{ LeaseId: leaseID, TargetEngine: engineName, RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), Knot: "knot.example", Rkey: "rkey", RepoDid: repoDID, }}
func TestNewFailsWhenOutboxCannotInitialize(t *testing.T) { d := testDB(t) if err := d.Close(); err != nil { t.Fatal(err) } n := notifier.New() cfg := &config.Config{} if _, err := New(cfg, nil, d, &n, slog.New(slog.NewTextHandler(io.Discard, nil)), nil, nil); err == nil { t.Fatal("New succeeded with an unavailable outbox database") }}
func TestReservedEngineHandsBackHeldSlotOnce(t *testing.T) { inner := &fakeEngine{} slot := &fakeSlot{} re := newReservedEngine(inner, slot)
got, err := re.(engine.WorkflowSlotter).AcquireWorkflowSlot(context.Background(), models.WorkflowId{}, nil, engine.Wait) if err != nil { t.Fatal(err) } if got != engine.WorkflowSlot(slot) { t.Fatal("AcquireWorkflowSlot() did not return the held slot") } if inner.acquireCalled { t.Fatal("wrapper must not call the inner engine's AcquireWorkflowSlot") }
if _, err := re.(engine.WorkflowSlotter).AcquireWorkflowSlot(context.Background(), models.WorkflowId{}, nil, engine.Wait); err == nil { t.Fatal("second AcquireWorkflowSlot() should error") }}
func TestReservedEngineForwardsResourceReporting(t *testing.T) { inner := &resourceReportingEngine{ fakeEngine: &fakeEngine{}, resources: quota.Resources{ quota.ResourceWorkflows: 1, quota.ResourceVCPUs: 4, }, usage: engine.WorkflowResourceUsage{ CPUUsec: 1234, MemoryPeakBytes: 4096, CgroupAvailable: true, VolumeAllocatedBytes: 8192, VolumeAvailable: true, }, } re := newReservedEngine(inner, &fakeSlot{}) wf := &models.Workflow{}
resources := re.(engine.WorkflowQuotaReporter).QuotaResources(wf) if resources[quota.ResourceWorkflows] != 1 || resources[quota.ResourceVCPUs] != 4 { t.Fatalf("QuotaResources() = %v, want workflow and vCPU resources", resources) }
usage, ok := re.(engine.WorkflowResourceUsageReporter).WorkflowResourceUsage(wf) if !ok { t.Fatal("WorkflowResourceUsage() unavailable") } if usage != inner.usage { t.Fatalf("WorkflowResourceUsage() = %+v, want %+v", usage, inner.usage) }}
func TestReservedEngineForwardsCacheRunner(t *testing.T) { inner := &fakeCacheEngine{fakeEngine: &fakeEngine{}} re := newReservedEngine(inner, &fakeSlot{}) runner, ok := re.(engine.CacheRunner) if !ok { t.Fatal("reserved engine dropped CacheRunner") }
ctx := context.Background() if err := runner.RestoreCache(ctx, models.WorkflowId{}, nil, nil, nil, nil); err != nil { t.Fatal(err) } if err := runner.SaveCache(ctx, models.WorkflowId{}, nil, nil, nil, nil); err != nil { t.Fatal(err) } if !inner.restored || !inner.saved { t.Fatalf("cache calls were not forwarded: restored=%t saved=%t", inner.restored, inner.saved) }}
func TestHandleCommitIsIdempotent(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc e.active["lease-1"] = &reservation{leaseID: "lease-1", committed: true}
e.handleCommit(context.Background(), &millv1.CommitLease{LeaseId: "lease-1"}) msg := <-enc.messages if got := msg.GetCommitted().GetLeaseId(); got != "lease-1" { t.Fatalf("Committed lease = %q, want lease-1", got) }}
func TestHandleCommitRejectsMissingReservation(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc
e.handleCommit(context.Background(), &millv1.CommitLease{LeaseId: "expired"}) result := (<-enc.messages).GetReserveResult() if result == nil { t.Fatal("missing reservation commit did not receive a ReserveResult") } if result.GetLeaseId() != "expired" || result.GetAccepted() { t.Fatalf("ReserveResult = %+v, want correlated rejection", result) }}
func TestHandleCancelFinalizesExpiredReservation(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc
e.handleCancel("lease-expired") var ack *millv1.CancelAck for ack == nil { select { case msg := <-enc.messages: ack = msg.GetCancelAck() case <-time.After(time.Second): t.Fatal("cancel acknowledgement timed out") } } if ack.GetLeaseId() != "lease-expired" { t.Fatalf("CancelAck = %+v, want lease-expired", ack) } rows, err := e.db.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("cancel terminal outbox rows = %d, want 1", len(rows)) } var entry millv1.Event if err := proto.Unmarshal(rows[0].Payload, &entry); err != nil { t.Fatal(err) } if got := entry.GetAttemptResult().GetStatus(); got != millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED { t.Fatalf("cancel terminal = %v, want CANCELLED", got) }}
func TestTerminalOutboxEntrySuppressesLaterEvents(t *testing.T) { e := testExecutor(t) if err := e.appendTerminal("lease-1", string(models.StatusKindSuccess), nil); err != nil { t.Fatal(err) }
e.handleCancel("lease-1") if err := e.appendStatus("lease-1", &tangled.PipelineStatus{Status: string(models.StatusKindRunning)}); err != nil { t.Fatal(err) }
rows, err := e.db.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("outbox rows after terminal = %d, want 1", len(rows)) } var entry millv1.Event if err := proto.Unmarshal(rows[0].Payload, &entry); err != nil { t.Fatal(err) } if got := entry.GetAttemptResult().GetStatus(); got != millv1.TerminalStatus_TERMINAL_STATUS_SUCCESS { t.Fatalf("terminal status = %v, want SUCCESS", got) }}
func TestTerminalGuardSurvivesExecutorRestart(t *testing.T) { d := testDB(t) first := &Executor{ db: d, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, } if err := first.initOutbox(); err != nil { t.Fatal(err) } if err := first.appendTerminal("lease-1", string(models.StatusKindSuccess), nil); err != nil { t.Fatal(err) }
restarted := &Executor{ db: d, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, } if err := restarted.initOutbox(); err != nil { t.Fatal(err) } restarted.handleCancel("lease-1")
rows, err := d.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("outbox rows after restart and duplicate cancel = %d, want 1", len(rows)) }}
func TestPendingArtifactRecoveryCompletesWithoutDeadlock(t *testing.T) { d := testDB(t) logDir := t.TempDir() wid := models.WorkflowId{ PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "3abc"}, Name: "build", } logPath := models.LogFilePath(logDir, wid) if err := os.MkdirAll(filepath.Dir(logPath), 0755); err != nil { t.Fatal(err) } if err := os.WriteFile(logPath, []byte("finished"), 0600); err != nil { t.Fatal(err) } if err := d.SavePendingArtifact( "lease-1", wid, string(models.StatusKindSuccess), "", 0, "logs/lease-1.log", "sha256:test", string(engine.FailureClassNone), string(engine.FailureReasonSuccess), true, ); err != nil { t.Fatal(err) }
uploaded := make(chan string, 1) e := &Executor{ db: d, cfg: &config.Config{Server: config.Server{LogDir: logDir}}, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), writer: artifactWriterFunc(func(_ context.Context, ref string, r io.Reader) error { if _, err := io.ReadAll(r); err != nil { return err } uploaded <- ref return nil }), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatalf("initOutbox: %v", err) } done := make(chan error, 1) go func() { done <- e.recoverPendingArtifacts(context.Background()) }() select { case err := <-done: if err != nil { t.Fatalf("recoverPendingArtifacts: %v", err) } case <-time.After(time.Second): t.Fatal("pending artifact recovery deadlocked") } select { case ref := <-uploaded: if ref != "logs/lease-1.log" { t.Fatalf("uploaded ref = %q", ref) } default: t.Fatal("pending artifact was not uploaded") } pending, err := d.ListPendingArtifacts() if err != nil { t.Fatal(err) } if len(pending) != 0 { t.Fatalf("pending artifacts after recovery = %d, want 0", len(pending)) } rows, err := d.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("outbox rows after recovery = %d, want 1", len(rows)) }}
func TestPendingArtifactRecoveryKeepsRowWhenLogIsMissing(t *testing.T) { d := testDB(t) wid := models.WorkflowId{ PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "3abc"}, Name: "build", } if err := d.SavePendingArtifact( "lease-1", wid, string(models.StatusKindSuccess), "", 0, "logs/lease-1.log", "sha256:test", string(engine.FailureClassNone), string(engine.FailureReasonSuccess), true, ); err != nil { t.Fatal(err) } e := &Executor{ db: d, cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}}, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), writer: artifactWriterFunc(func(context.Context, string, io.Reader) error { t.Fatal("writer called without a log file") return nil }), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatalf("initOutbox: %v", err) } if err := e.recoverPendingArtifacts(context.Background()); err != nil { t.Fatalf("recoverPendingArtifacts: %v", err) } pending, err := d.ListPendingArtifacts() if err != nil { t.Fatal(err) } if len(pending) != 1 { t.Fatalf("pending artifacts with missing log = %d, want 1", len(pending)) } rows, err := d.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 0 { t.Fatalf("outbox rows with missing log = %d, want 0", len(rows)) }}
func TestPendingArtifactRecoveryUsesAggregateDeadline(t *testing.T) { d := testDB(t) logDir := t.TempDir() wids := []models.WorkflowId{ {PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "3abc"}, Name: "build"}, {PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "3def"}, Name: "test"}, } for i, wid := range wids { logPath := models.LogFilePath(logDir, wid) if err := os.MkdirAll(filepath.Dir(logPath), 0755); err != nil { t.Fatal(err) } if err := os.WriteFile(logPath, []byte("finished"), 0600); err != nil { t.Fatal(err) } leaseID := fmt.Sprintf("lease-%d", i+1) if err := d.SavePendingArtifact( leaseID, wid, string(models.StatusKindSuccess), "", 0, "logs/"+leaseID+".log", "sha256:test", string(engine.FailureClassNone), string(engine.FailureReasonSuccess), true, ); err != nil { t.Fatal(err) } }
var deadlines []time.Time e := &Executor{ db: d, cfg: &config.Config{Server: config.Server{LogDir: logDir}}, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), writer: artifactWriterFunc(func(ctx context.Context, _ string, _ io.Reader) error { deadline, ok := ctx.Deadline() if !ok { t.Fatal("artifact recovery context has no deadline") } deadlines = append(deadlines, deadline) return errors.New("unavailable") }), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatalf("initOutbox: %v", err) } started := time.Now() if err := e.recoverPendingArtifacts(context.Background()); err != nil { t.Fatalf("recoverPendingArtifacts: %v", err) } finished := time.Now() if len(deadlines) != len(wids) { t.Fatalf("upload attempts = %d, want %d", len(deadlines), len(wids)) } if !deadlines[0].Equal(deadlines[1]) { t.Fatalf("artifact recovery deadlines differ: %v, %v", deadlines[0], deadlines[1]) } wantTimeout := 2 * time.Minute if deadlines[0].Before(started.Add(wantTimeout)) || deadlines[0].After(finished.Add(wantTimeout)) { t.Fatalf("artifact recovery deadline = %v, want %v after startup", deadlines[0], wantTimeout) }}
func TestHandleCommitPreservesPreauthorizedSecrets(t *testing.T) { d := testDB(t) n := notifier.New() enc := newCaptureEncoder() e := &Executor{ cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}}, db: d, n: &n, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, enc: enc, } if err := e.initOutbox(); err != nil { t.Fatal(err) } e.lifecycleCtx = context.Background()
inner := &fakeEngine{secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{})} slot := &fakeSlot{} repoDid, err := syntax.ParseDID("did:web:example.com") if err != nil { t.Fatal(err) } repoDidString := repoDid.String() res := &reservation{ leaseID: "lease-1", wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, realEngine: inner, slot: slot, wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}}, pipeline: &models.Pipeline{ RepoDid: repoDid, TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &repoDidString}, }, TrustedSource: true, }, } e.active[res.leaseID] = res
e.handleCommit(context.Background(), &millv1.CommitLease{ LeaseId: res.leaseID, Secrets: []*millv1.Secret{{Key: "TOKEN", Value: "secret-value"}}, }) if got := (<-enc.messages).GetCommitted().GetLeaseId(); got != res.leaseID { t.Fatalf("Committed lease = %q, want %q", got, res.leaseID) } if got := e.maskSecrets(res, "value=secret-value"); got != "value=***" { t.Fatalf("masked log = %q", got) } select { case got := <-inner.secrets: if len(got) != 1 || got[0].Key != "TOKEN" || got[0].Value != "secret-value" { t.Fatalf("RunStep secrets = %+v", got) } case <-time.After(2 * time.Second): t.Fatal("RunStep did not receive CommitLease secrets") } select { case <-inner.done: case <-time.After(2 * time.Second): t.Fatal("workflow did not finish") } if res.stopTail != nil { res.stopTail() } e.jobsWG.Wait()}
func TestConnectHandshakesWhilePendingArtifactRecoveryIsBlocked(t *testing.T) { handshake := make(chan error, 1) srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { conn, err := websocket.Upgrade(w, r, nil, 1024, 1024) if err != nil { handshake <- err return } defer conn.Close() stream := millproto.NewWSStream(conn) dec := millproto.NewDecoder(stream) enc := millproto.NewEncoder(stream) helloMsg, err := dec.Decode() if err != nil { handshake <- err return } hello := helloMsg.GetHello() if hello == nil { handshake <- errors.New("first executor message was not Hello") return } if err := enc.Encode(&millproto.Message{Resume: &millv1.Resume{Epoch: hello.GetEpoch()}}); err != nil { handshake <- err return } snapshotMsg, err := dec.Decode() if err != nil { handshake <- err return } if snapshotMsg.GetNodeSnapshot() == nil { handshake <- errors.New("executor did not send NodeSnapshot after Resume") return } handshake <- nil for { if _, err := dec.Decode(); err != nil { return } } })) t.Cleanup(srv.Close)
e := testSessionExecutor(t, "ws"+strings.TrimPrefix(srv.URL, "http")) n := notifier.New() e.n = &n logDir := t.TempDir() e.cfg.Server.LogDir = logDir wid := models.WorkflowId{ PipelineId: models.PipelineId{Knot: "knot.example", Rkey: "3abc"}, Name: "build", } logPath := models.LogFilePath(logDir, wid) if err := os.MkdirAll(filepath.Dir(logPath), 0755); err != nil { t.Fatal(err) } if err := os.WriteFile(logPath, []byte("finished"), 0600); err != nil { t.Fatal(err) } if err := e.db.SavePendingArtifact( "lease-1", wid, string(models.StatusKindSuccess), "", 0, "logs/lease-1.log", "sha256:test", string(engine.FailureClassNone), string(engine.FailureReasonSuccess), true, ); err != nil { t.Fatal(err) } recoveryStarted := make(chan struct{}) e.writer = artifactWriterFunc(func(ctx context.Context, _ string, _ io.Reader) error { close(recoveryStarted) <-ctx.Done() return ctx.Err() })
ctx, cancel := context.WithCancel(context.Background()) done := make(chan struct{}) go func() { e.Connect(ctx) close(done) }() select { case <-recoveryStarted: case <-time.After(time.Second): cancel() t.Fatal("pending artifact recovery did not start") } select { case err := <-handshake: if err != nil { cancel() t.Fatalf("mill handshake: %v", err) } case <-time.After(time.Second): cancel() t.Fatal("blocked artifact recovery delayed the mill handshake") } cancel() select { case <-done: case <-time.After(2 * time.Second): t.Fatal("Connect did not stop after cancellation") }}
func TestHandleCommitRejectsCachesAndSecretsForForkSource(t *testing.T) { d := testDB(t) n := notifier.New() enc := newCaptureEncoder() cache, err := storage.NewDisk(t.TempDir()) if err != nil { t.Fatal(err) } e := &Executor{ db: d, n: &n, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}}, cache: cache, enc: enc, } if err := e.initOutbox(); err != nil { t.Fatal(err) } e.lifecycleCtx = context.Background()
inner := &fakeCacheEngine{fakeEngine: &fakeEngine{ secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{}), }} targetRepoDid, err := syntax.ParseDID("did:plc:target") if err != nil { t.Fatal(err) } sourceRepoDid := "did:plc:fork" cacheObjectID := "11111111-1111-4111-8111-111111111111" res := &reservation{ leaseID: "lease-1", wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, realEngine: inner, slot: &fakeSlot{}, wf: &models.Workflow{ Name: "build", Steps: []models.Step{fakeStep{}}, Caches: []models.CacheEntry{{Key: "deps", Paths: []string{"deps"}}}, }, pipeline: &models.Pipeline{ RepoDid: targetRepoDid, TriggerMetadata: &tangled.Pipeline_TriggerMetadata{ SourceRepo: &sourceRepoDid, Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &sourceRepoDid}, }, }, } e.active[res.leaseID] = res
e.handleCommit(context.Background(), &millv1.CommitLease{ LeaseId: res.leaseID, Secrets: []*millv1.Secret{{Key: "TOKEN", Value: "secret-value"}}, CacheBindings: []*millv1.CacheBinding{{ EntryIndex: 0, RestoreId: cacheObjectID, RestoreKey: "objects/" + targetRepoDid.String() + "/" + cacheObjectID, SaveId: cacheObjectID, SaveKey: "objects/" + targetRepoDid.String() + "/" + cacheObjectID, }}, }) if got := (<-enc.messages).GetCommitted().GetLeaseId(); got != res.leaseID { t.Fatalf("Committed lease = %q, want %q", got, res.leaseID) } e.jobsWG.Wait()
if inner.restored || inner.saved { t.Fatalf("fork cache access = restored %t, saved %t", inner.restored, inner.saved) } if got := <-inner.secrets; len(got) != 0 { t.Fatalf("fork secrets = %+v, want none", got) }}
func TestRunSessionCancellationClosesStalledWebsocket(t *testing.T) { connected := make(chan struct{}) release := make(chan struct{}) srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { conn, err := websocket.Upgrade(w, r, nil, 1024, 1024) if err != nil { return } defer conn.Close() close(connected) <-release })) t.Cleanup(func() { close(release) srv.Close() })
e := testSessionExecutor(t, "ws"+strings.TrimPrefix(srv.URL, "http")) ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { done <- e.runSession(ctx) }() <-connected cancel()
select { case <-done: case <-time.After(2 * time.Second): t.Fatal("runSession did not return after context cancellation") }}
func testSessionExecutor(t *testing.T, url string) *Executor { d := testDB(t) e := &Executor{ millURL: url, seats: 1, engines: make(map[string]models.Engine), db: d, cfg: &config.Config{Server: config.Server{Dev: true}}, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatal(err) } return e}
func testExecutor(t *testing.T) *Executor { d := testDB(t) e := &Executor{ db: d, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatal(err) } return e}
func TestReserveHonorsExecutorSeatCapacity(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc e.seats = 1 e.engines = map[string]models.Engine{"dummy": &fakeEngine{}}
e.handleReserve(context.Background(), testReserveSeat(t, "lease-1", "dummy")) first := (<-enc.messages).GetReserveResult() if first == nil || !first.GetAccepted() { t.Fatalf("first reserve result = %+v, want accepted", first) } if snapshot := (<-enc.messages).GetNodeSnapshot(); snapshot == nil { t.Fatal("accepted reservation did not publish a snapshot") }
e.handleReserve(context.Background(), testReserveSeat(t, "lease-2", "dummy")) second := (<-enc.messages).GetReserveResult() if second == nil || second.GetAccepted() { t.Fatalf("second reserve result = %+v, want transient rejection", second) } if second.GetRejectClass() != millv1.RejectClass_REJECT_CLASS_TRANSIENT { t.Fatalf("second reject class = %v, want transient", second.GetRejectClass()) } if second.GetRejectReason() != "no executor seats available" { t.Fatalf("second reject reason = %q", second.GetRejectReason()) }
cleanup, ok := e.takeUncommittedReservation("lease-1", true) if !ok { t.Fatal("first reservation missing during cleanup") } cleanup()}
func testDB(t *testing.T) *db.DB { d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db")) if err != nil { t.Fatal(err) } return d}
func TestFinishJobReportsCancelledReservationAsCancelled(t *testing.T) { d := testDB(t) res := &reservation{ leaseID: "lease-1", wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, cancelled: true, } e := &Executor{ db: d, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: map[string]*reservation{res.leaseID: res}, maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatal(err) }
e.finishJob(res, &tangled.PipelineStatus{ Pipeline: string(res.wid.PipelineId.AtUri()), Workflow: res.wid.Name, Status: string(models.StatusKindFailed), })
rows, err := d.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("outbox rows = %d, want 1", len(rows)) }
var entry millv1.Event if err := proto.Unmarshal(rows[0].Payload, &entry); err != nil { t.Fatal(err) } got := entry.GetAttemptResult().GetStatus() if got != millv1.TerminalStatus_TERMINAL_STATUS_CANCELLED { t.Fatalf("terminal status = %v, want CANCELLED", got) }}
func TestFinishJobWaitsForEngineCleanup(t *testing.T) { d := testDB(t) res := &reservation{ leaseID: "lease-1", wid: models.WorkflowId{ PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build", }, runDone: make(chan struct{}), } e := &Executor{ db: d, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: map[string]*reservation{res.leaseID: res}, maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatal(err) }
done := make(chan error, 1) go func() { done <- e.finishJob(res, &tangled.PipelineStatus{ Pipeline: string(res.wid.PipelineId.AtUri()), Workflow: res.wid.Name, Status: string(models.StatusKindSuccess), }) }()
select { case err := <-done: t.Fatalf("finishJob returned before engine cleanup: %v", err) case <-time.After(50 * time.Millisecond): } rows, err := d.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 0 { t.Fatalf("terminal rows before engine cleanup = %d, want 0", len(rows)) }
close(res.runDone) select { case err := <-done: if err != nil { t.Fatal(err) } case <-time.After(time.Second): t.Fatal("finishJob did not resume after engine cleanup") }}
func TestReplayRejectsMalformedOutboxRow(t *testing.T) { d := testDB(t) e := &Executor{ db: d, enc: newCaptureEncoder(), l: slog.New(slog.NewTextHandler(io.Discard, nil)), } if err := e.initOutbox(); err != nil { t.Fatal(err) } if _, err := d.AppendOutboxRow([]byte("not protobuf"), true); err != nil { t.Fatal(err) } if err := e.replay(0); err == nil { t.Fatal("replay accepted a malformed row and would leave a permanent seqno gap") }}
func TestReplayDropsRowsAlreadyAcknowledgedByResume(t *testing.T) { e := testExecutor(t) if err := e.appendTerminal("lease-1", string(models.StatusKindSuccess), nil); err != nil { t.Fatal(err) } rows, err := e.db.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("outbox rows = %d, want 1", len(rows)) }
e.enc = newCaptureEncoder() if err := e.replay(rows[0].Seqno); err != nil { t.Fatalf("replay: %v", err) } rows, err = e.db.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 0 { t.Fatalf("outbox rows after resume = %d, want 0", len(rows)) } if len(e.terminalSeqnos) != 0 { t.Fatalf("terminal guards after resume = %d, want 0", len(e.terminalSeqnos)) } select { case <-e.outboxIdle(): default: t.Fatal("resumed acknowledgement left the outbox busy") }}
func TestSocketCancellationIndependence(t *testing.T) { d := testDB(t) n := notifier.New() e := &Executor{ db: d, n: &n, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}}, } if err := e.initOutbox(); err != nil { t.Fatal(err) }
lifecycleCtx, cancelLifecycle := context.WithCancel(context.Background()) defer cancelLifecycle() e.lifecycleCtx = lifecycleCtx
inner := &fakeEngine{secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{})} slot := &fakeSlot{} repoDid, err := syntax.ParseDID("did:web:example.com") if err != nil { t.Fatal(err) } res := &reservation{ leaseID: "lease-1", wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, realEngine: inner, slot: slot, wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}}, pipeline: &models.Pipeline{RepoDid: repoDid}, } e.active[res.leaseID] = res
sessionCtx, cancelSession := context.WithCancel(lifecycleCtx)
e.handleCommit(sessionCtx, &millv1.CommitLease{ LeaseId: "lease-1", })
cancelSession()
// session disconnect must not cancel the running job select { case <-inner.done: case <-time.After(2 * time.Second): t.Fatal("workflow did not complete even though websocket session was cancelled") }
e.jobsWG.Wait()}
func TestMonotonicSnapshots(t *testing.T) { enc := newCaptureEncoder() e := &Executor{ enc: enc, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, }
e.pushSnapshot() msg1 := <-enc.messages seq1 := msg1.GetNodeSnapshot().GetSeqno() if seq1 != 1 { t.Fatalf("first seq = %d, want 1", seq1) }
e.pushSnapshot() msg2 := <-enc.messages seq2 := msg2.GetNodeSnapshot().GetSeqno() if seq2 != 2 { t.Fatalf("second seq = %d, want 2", seq2) }}
func TestDrainWaitsForActiveReservation(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc e.engines = map[string]models.Engine{"dummy": &fakeEngine{}}
res := &reservation{leaseID: "lease-1", committed: true} e.active[res.leaseID] = res done := make(chan error, 1) go func() { done <- e.Drain(context.Background()) }()
snapshot := (<-enc.messages).GetNodeSnapshot() if snapshot == nil || snapshot.GetEngines()["dummy"].GetAvailable() { t.Fatalf("drain snapshot = %+v, want unavailable engine", snapshot) } if got := snapshot.GetActiveLeaseIds(); len(got) != 1 || got[0] != res.leaseID { t.Fatalf("active leases = %v, want [%s]", got, res.leaseID) } select { case err := <-done: t.Fatalf("Drain returned before reservation cleanup: %v", err) default: }
e.mu.Lock() cleanup := e.removeReservationLocked(res, false) e.mu.Unlock() cleanup() select { case err := <-done: if err != nil { t.Fatalf("Drain: %v", err) } case <-time.After(2 * time.Second): t.Fatal("Drain did not return after reservation cleanup") }}
func TestDrainWaitsForReservationCleanup(t *testing.T) { releaseStarted := make(chan struct{}) releaseSlot := make(chan struct{}) slot := &fakeSlot{release: func() { close(releaseStarted) <-releaseSlot }} e := testExecutor(t) e.enc = newCaptureEncoder() res := &reservation{leaseID: "lease-1", slot: slot} e.active[res.leaseID] = res
cleanup, ok := e.takeUncommittedReservation(res.leaseID, true) if !ok { t.Fatal("failed to take reservation") } cleanupDone := make(chan struct{}) go func() { cleanup() close(cleanupDone) }() <-releaseStarted
done := make(chan error, 1) go func() { done <- e.Drain(context.Background()) }() <-e.enc.(*captureEncoder).messages select { case err := <-done: t.Fatalf("Drain returned before slot release completed: %v", err) default: }
close(releaseSlot) <-cleanupDone select { case err := <-done: if err != nil { t.Fatalf("Drain: %v", err) } case <-time.After(2 * time.Second): t.Fatal("Drain did not return after slot release") } if slot.released != 1 { t.Fatalf("slot releases = %d, want 1", slot.released) }}
func TestDrainWaitsForWorkflowCleanup(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc e.jobsWG.Add(1)
done := make(chan error, 1) go func() { done <- e.Drain(context.Background()) }() <-enc.messages select { case err := <-done: t.Fatalf("Drain returned before workflow cleanup: %v", err) default: }
e.jobsWG.Done() select { case err := <-done: if err != nil { t.Fatalf("Drain: %v", err) } case <-time.After(2 * time.Second): t.Fatal("Drain did not return after workflow cleanup") }}
func TestDrainWaitsForTerminalAcknowledgement(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc if err := e.appendTerminal("lease-1", string(models.StatusKindSuccess), nil); err != nil { t.Fatal(err) } batch := (<-enc.messages).GetEventBatch() if batch == nil || len(batch.GetEvents()) != 1 { t.Fatalf("terminal batch = %+v, want one event", batch) }
done := make(chan error, 1) go func() { done <- e.Drain(context.Background()) }() <-enc.messages select { case err := <-done: t.Fatalf("Drain returned before terminal acknowledgement: %v", err) default: }
e.handleAck(&millv1.Ack{Epoch: e.epoch, UpToSeqno: batch.GetEvents()[0].GetSeqno()}) select { case err := <-done: if err != nil { t.Fatalf("Drain: %v", err) } case <-time.After(2 * time.Second): t.Fatal("Drain did not return after terminal acknowledgement") }}
func TestDrainReturnsWithDurableOutboxAfterReconnectGrace(t *testing.T) { e := testExecutor(t) e.cfg = &config.Config{} e.cfg.Mill.ReconnectGrace = 10 * time.Millisecond if err := e.appendTerminal("lease-1", string(models.StatusKindSuccess), nil); err != nil { t.Fatal(err) }
if err := e.Drain(context.Background()); err != nil { t.Fatalf("Drain with disconnected Mill: %v", err) }
rows, err := e.db.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 1 { t.Fatalf("durable outbox rows = %d, want 1", len(rows)) }}
func TestDrainAllowsReconnectToAcknowledgeOutbox(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.cfg = &config.Config{} e.cfg.Mill.ReconnectGrace = 2 * time.Second e.enc = enc if err := e.appendTerminal("lease-1", string(models.StatusKindSuccess), nil); err != nil { t.Fatal(err) } batch := (<-enc.messages).GetEventBatch() if batch == nil || len(batch.GetEvents()) != 1 { t.Fatalf("terminal batch = %+v, want one event", batch) }
done := make(chan error, 1) go func() { done <- e.Drain(context.Background()) }() <-enc.messages e.connMu.Lock() e.enc = nil e.connMu.Unlock() select { case err := <-done: t.Fatalf("Drain returned before reconnect grace elapsed: %v", err) case <-time.After(50 * time.Millisecond): }
e.connMu.Lock() e.enc = enc e.connMu.Unlock() e.handleAck(&millv1.Ack{Epoch: e.epoch, UpToSeqno: batch.GetEvents()[0].GetSeqno()}) select { case err := <-done: if err != nil { t.Fatalf("Drain after reconnect acknowledgement: %v", err) } case <-time.After(2 * time.Second): t.Fatal("Drain did not accept acknowledgement during reconnect grace") }
rows, err := e.db.ListOutboxRows() if err != nil { t.Fatal(err) } if len(rows) != 0 { t.Fatalf("durable outbox rows = %d, want 0", len(rows)) }}
func TestReserveCannotRacePastDrain(t *testing.T) { enc := newCaptureEncoder() e := testExecutor(t) e.enc = enc slot := &fakeSlot{} inner := &fakeEngine{ acquireStarted: make(chan struct{}), releaseAcquire: make(chan struct{}), slot: slot, } e.engines = map[string]models.Engine{"dummy": inner}
reserve := testReserveSeat(t, "lease-1", "dummy") reserveDone := make(chan struct{}) go func() { e.handleReserve(context.Background(), reserve) close(reserveDone) }() <-inner.acquireStarted
if err := e.Drain(context.Background()); err != nil { t.Fatalf("Drain: %v", err) } snapshot := (<-enc.messages).GetNodeSnapshot() if snapshot == nil || snapshot.GetEngines()["dummy"].GetAvailable() { t.Fatalf("drain snapshot = %+v, want unavailable engine", snapshot) } close(inner.releaseAcquire) <-reserveDone
rejected := (<-enc.messages).GetReserveResult() if rejected == nil || rejected.GetAccepted() || rejected.GetRejectReason() != "draining" { t.Fatalf("reserve result = %+v, want draining rejection", rejected) } if slot.released != 1 { t.Fatalf("slot releases = %d, want 1", slot.released) } if len(e.active) != 0 { t.Fatalf("active reservations = %d, want 0", len(e.active)) }}
func TestDrainCancellationIncludesUnavailableSnapshot(t *testing.T) { enc := &blockingEncoder{ started: make(chan struct{}, 1), release: make(chan struct{}), } e := testExecutor(t) e.enc = enc ctx, cancel := context.WithCancel(context.Background()) done := make(chan error, 1) go func() { done <- e.Drain(ctx) }()
<-enc.started cancel() select { case err := <-done: if !errors.Is(err, context.Canceled) { t.Fatalf("Drain error = %v, want context canceled", err) } case <-time.After(2 * time.Second): t.Fatal("Drain ignored cancellation while sending the unavailable snapshot") } close(enc.release)}
func TestTimerRace(t *testing.T) { d := testDB(t) e := &Executor{ db: d, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, } if err := e.initOutbox(); err != nil { t.Fatal(err) }
twf, _ := json.Marshal(tangled.Pipeline_Workflow{Name: "build"}) tpl, _ := json.Marshal(tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{}})
inner := &fakeEngine{} e.engines = map[string]models.Engine{"microvm": inner}
e.handleReserve(context.Background(), &millv1.ReserveSeat{ LeaseId: "lease-1", TargetEngine: "microvm", RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), TtlSeconds: 1, RepoDid: "did:web:example.com", })
e.mu.Lock() res := e.active["lease-1"] e.mu.Unlock()
if res == nil { t.Fatal("reservation was not added") }
deadline := time.Now().Add(5 * time.Second) for { e.mu.Lock() activeLen := len(e.active) e.mu.Unlock() if activeLen == 0 { break } if time.Now().After(deadline) { t.Fatal("reservation was leaked and never expired") } time.Sleep(10 * time.Millisecond) }}
func TestStructuredShutdown(t *testing.T) { d := testDB(t) n := notifier.New() e := &Executor{ db: d, n: &n, l: slog.New(slog.NewTextHandler(io.Discard, nil)), active: make(map[string]*reservation), maxOutboxBytes: 10 * 1024 * 1024, cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}}, } if err := e.initOutbox(); err != nil { t.Fatal(err) }
ctx, cancel := context.WithCancel(context.Background()) e.lifecycleCtx = ctx
inner := &fakeEngine{secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{})} slot := &fakeSlot{} repoDid, err := syntax.ParseDID("did:web:example.com") if err != nil { t.Fatal(err) } res := &reservation{ leaseID: "lease-1", wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"}, realEngine: inner, slot: slot, wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}}, pipeline: &models.Pipeline{RepoDid: repoDid}, } e.active[res.leaseID] = res
e.handleCommit(ctx, &millv1.CommitLease{ LeaseId: "lease-1", })
cancel()
e.jobsWG.Wait()
select { case <-inner.done: default: t.Fatal("shutdown returned but running job did not finish") }}