Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402package mill
import ( "context" "encoding/json" "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/engines/dummy" "tangled.org/core/spindle/mill/executor" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability")
func TestEndToEndDummyJob(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() l := slog.New(slog.NewTextHandler(io.Discard, nil))
millDir := t.TempDir() bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) if err != nil { t.Fatalf("mill db: %v", err) } bn := notifier.New() mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) mill.Attach(bdb, &bn, testQuotaManager(t, bdb)) millMetrics := observability.NewMetrics() mill.RegisterMetrics(millMetrics) registerTestExecutor(t, bdb, "exec-1", HashToken("test-token"), nil)
srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) defer srv.Close() wsURL := "ws" + strings.TrimPrefix(srv.URL, "http")
execDir := t.TempDir() edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) if err != nil { t.Fatalf("exec db: %v", err) } en := notifier.New() cfg := &config.Config{} cfg.Server.LogDir = execDir cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") cfg.Server.Hostname = "exec-1" cfg.Mill.URL = wsURL cfg.Mill.Seats = 2 cfg.Mill.SharedSecret = "test-token"
dummyEng := dummy.New(l) dummyEng.StepDelay = 50 * time.Millisecond engines := map[string]models.Engine{"dummy": dummyEng} exec, err := executor.New(cfg, engines, edb, &en, l, nil, nil) if err != nil { t.Fatalf("executor.New: %v", err) } execMetrics := observability.NewMetrics() exec.RegisterMetrics(execMetrics) go exec.Connect(observability.WithMetrics(ctx, execMetrics))
be := NewEngine("dummy", mill) twf := tangled.Pipeline_Workflow{ Name: "build", Raw: "steps:\n - name: hello\n command: echo hi\n", } wf, err := be.InitWorkflow(twf, testPipeline()) if err != nil { t.Fatalf("InitWorkflow: %v", err) } wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "rkey1"}, Name: "build"} if err := bdb.StatusPending(wid, &bn); err != nil { t.Fatal(err) }
placeCtx, placeCancel := context.WithTimeout(ctx, 10*time.Second) defer placeCancel()
slot, err := mill.place(placeCtx, "dummy", wid, wf) if err != nil { t.Fatalf("place: %v", err) } defer slot.Release()
logPath := models.LogFilePath(millDir, wid) ch := bn.Subscribe() defer bn.Unsubscribe(ch)
sawLogContent := make(chan bool, 1) go func() { ticker := time.NewTicker(20 * time.Millisecond) defer ticker.Stop() timeout := time.After(10 * time.Second) for { select { case <-ch: case <-ticker.C: case <-timeout: sawLogContent <- false return } data, err := os.ReadFile(logPath) if err == nil && strings.Contains(string(data), "echo hi") { sawLogContent <- true return } } }()
if err := mill.commitAndWait(placeCtx, wf, nil); err != nil { t.Fatalf("commitAndWait: %v, want success", err) }
if !<-sawLogContent { t.Fatal("mill live log file never received expected content while job was running") }
if !waitForFileRemoval(t, logPath) { t.Fatalf("mill live log file was not removed after terminal artifact was recorded") }
archived, err := bdb.GetFinishedLog(wid) if err != nil { t.Fatalf("GetFinishedLog: %v", err) } if archived.Workflow != wid.Name || archived.Ref == "" { t.Fatalf("archived log = %+v, want artifact for %s", archived, wid) }
if !waitForStatus(t, bdb, wid, "running") { events, _ := bdb.GetEvents(0, 1000) t.Logf("mill events after completion: %+v", events) t.Fatal("mill never saw streamed running status") }
millFamilies, err := millMetrics.Registry().Gather() if err != nil { t.Fatal(err) } var millTotal, millClassified, millStartupSamples float64 for _, family := range millFamilies { switch family.GetName() { case "spindle_workflows_total": for _, metric := range family.GetMetric() { millTotal += metric.GetCounter().GetValue() } case "spindle_workflow_terminations_total": for _, metric := range family.GetMetric() { millClassified += metric.GetCounter().GetValue() } case "spindle_workflow_startup_delay_seconds": for _, metric := range family.GetMetric() { millStartupSamples += float64(metric.GetHistogram().GetSampleCount()) } } } if millTotal != 1 || millClassified != 1 || millStartupSamples != 1 { t.Fatalf( "mill metrics = total %v classified %v startup samples %v, want exactly one of each", millTotal, millClassified, millStartupSamples, ) }
executorFamilies, err := execMetrics.Registry().Gather() if err != nil { t.Fatal(err) } var executorTotal, executorDurations float64 for _, family := range executorFamilies { switch family.GetName() { case "spindle_workflows_total": for _, metric := range family.GetMetric() { executorTotal += metric.GetCounter().GetValue() } case "spindle_workflow_duration_seconds": for _, metric := range family.GetMetric() { executorDurations += float64(metric.GetHistogram().GetSampleCount()) } } } if executorTotal != 0 || executorDurations != 1 { t.Fatalf("executor metrics = terminal %v duration samples %v, want 0 and 1", executorTotal, executorDurations) }}
func TestExecutorConfiguredLabelsAreStoredOnSession(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() l := slog.New(slog.NewTextHandler(io.Discard, nil))
millDir := t.TempDir() bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) if err != nil { t.Fatalf("mill db: %v", err) } bn := notifier.New() mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) mill.Attach(bdb, &bn, testQuotaManager(t, bdb)) registerTestExecutor(t, bdb, "exec-labels", HashToken("test-token"), []string{"linux", "arm64", "gpu"})
srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) defer srv.Close() wsURL := "ws" + strings.TrimPrefix(srv.URL, "http")
execDir := t.TempDir() edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) if err != nil { t.Fatalf("exec db: %v", err) } en := notifier.New() cfg := &config.Config{} cfg.Server.Dev = true cfg.Server.LogDir = execDir cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") cfg.Server.Hostname = "exec-labels" cfg.Mill.URL = wsURL cfg.Mill.Seats = 2 cfg.Mill.SharedSecret = "test-token" cfg.Mill.Labels = []string{"linux", "arm64", "gpu"}
engines := map[string]models.Engine{"dummy": dummy.New(l)} exec, err := executor.New(cfg, engines, edb, &en, l, nil, nil) if err != nil { t.Fatalf("executor.New: %v", err) } go exec.Connect(ctx)
if !waitForSessionLabels(t, mill, "exec-labels", []string{"linux", "arm64", "gpu"}) { t.Fatal("mill session never stored executor labels from hello") }}
func TestEndToEndDummyJobUsesRequiredLabelsAcrossExecutors(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() l := slog.New(slog.NewTextHandler(io.Discard, nil))
millDir := t.TempDir() bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) if err != nil { t.Fatalf("mill db: %v", err) } bn := notifier.New() mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) mill.Attach(bdb, &bn, testQuotaManager(t, bdb)) registerTestExecutor(t, bdb, "exec-x86", HashToken("token-x86"), []string{"linux/amd64", "kvm"}) registerTestExecutor(t, bdb, "exec-arm", HashToken("token-arm"), []string{"linux/arm64", "kvm"})
srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) defer srv.Close() wsURL := "ws" + strings.TrimPrefix(srv.URL, "http")
startExecutor := func(name, token string, labels []string) { t.Helper() execDir := t.TempDir() edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) if err != nil { t.Fatalf("%s exec db: %v", name, err) } en := notifier.New() cfg := &config.Config{} cfg.Server.Dev = true cfg.Server.LogDir = execDir cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") cfg.Server.Hostname = name cfg.Mill.URL = wsURL cfg.Mill.Seats = 1 cfg.Mill.SharedSecret = token cfg.Mill.Labels = labels
engines := map[string]models.Engine{"dummy": dummy.New(l)} exec, err := executor.New(cfg, engines, edb, &en, l, nil, nil) if err != nil { t.Fatalf("executor.New: %v", err) } go exec.Connect(ctx) } startExecutor("exec-x86", "token-x86", []string{"linux/amd64", "kvm"}) startExecutor("exec-arm", "token-arm", []string{"linux/arm64", "kvm"})
if !waitForSessionLabels(t, mill, "exec-x86", []string{"linux/amd64", "kvm"}) { t.Fatal("x86 executor did not connect with labels") } if !waitForSessionLabels(t, mill, "exec-arm", []string{"linux/arm64", "kvm"}) { t.Fatal("arm executor did not connect with labels") }
be := NewEngine("dummy", mill) twf := tangled.Pipeline_Workflow{ Name: "build-arm", RunsOn: []string{"linux/arm64"}, Raw: "steps:\n - name: hello\n command: echo hi\n", } wf, err := be.InitWorkflow(twf, testPipeline()) if err != nil { t.Fatalf("InitWorkflow: %v", err) } wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "rkey-arm"}, Name: "build-arm"}
placeCtx, placeCancel := context.WithTimeout(ctx, 10*time.Second) defer placeCancel()
slot, err := mill.place(placeCtx, "dummy", wid, wf) if err != nil { t.Fatalf("place: %v", err) } defer slot.Release()
lease := wf.Data.(*millWorkflowState).Lease if lease == nil { t.Fatal("place did not attach lease to workflow state") } if lease.nodeID != "exec-arm" { t.Fatalf("placed on %q, want exec-arm", lease.nodeID) }
if err := mill.commitAndWait(placeCtx, wf, nil); err != nil { t.Fatalf("commitAndWait: %v, want success", err) }}
func waitForSessionLabels(t *testing.T, m *Mill, nodeID string, want []string) bool { t.Helper() deadline := time.Now().Add(3 * time.Second) for time.Now().Before(deadline) { m.mu.Lock() sess := m.sessions[nodeID] var got []string if sess != nil { got = append([]string(nil), sess.labels...) } m.mu.Unlock() if sameStringMultiset(got, want) { return true } time.Sleep(10 * time.Millisecond) } return false}
func waitForStatus(t *testing.T, d *db.DB, wid models.WorkflowId, want string) bool { t.Helper() deadline := time.Now().Add(3 * time.Second) aturi := string(wid.PipelineId.AtUri()) for time.Now().Before(deadline) { evs, err := d.GetEvents(0, 1000) if err != nil { t.Fatalf("GetEvents: %v", err) } for _, ev := range evs { if ev.Nsid != tangled.PipelineStatusNSID { continue } var st tangled.PipelineStatus if err := json.Unmarshal(ev.EventJson, &st); err != nil { continue } if st.Pipeline == aturi && st.Workflow == wid.Name && st.Status == want { return true } } time.Sleep(50 * time.Millisecond) } return false}
func waitForLogFileContent(t *testing.T, path, want string) bool { t.Helper() deadline := time.Now().Add(5 * time.Second) for time.Now().Before(deadline) { data, err := os.ReadFile(path) if err == nil && strings.Contains(string(data), want) { return true } time.Sleep(20 * time.Millisecond) } return false}
func waitForFileRemoval(t *testing.T, path string) bool { t.Helper() deadline := time.Now().Add(5 * time.Second) for time.Now().Before(deadline) { _, err := os.Stat(path) if os.IsNotExist(err) { return true } time.Sleep(20 * time.Millisecond) } return false}