This repository has no description
0

Configure Feed

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

core / spindle / mill / integration_test.go
10 kB 345 lines
1package mill 2 3import ( 4 "context" 5 "encoding/json" 6 "io" 7 "log/slog" 8 "net/http" 9 "net/http/httptest" 10 "os" 11 "path/filepath" 12 "strings" 13 "testing" 14 "time" 15 16 "tangled.org/core/api/tangled" 17 "tangled.org/core/notifier" 18 "tangled.org/core/spindle/config" 19 "tangled.org/core/spindle/db" 20 "tangled.org/core/spindle/engines/dummy" 21 "tangled.org/core/spindle/mill/executor" 22 "tangled.org/core/spindle/models" 23) 24 25func TestEndToEndDummyJob(t *testing.T) { 26 ctx, cancel := context.WithCancel(context.Background()) 27 defer cancel() 28 l := slog.New(slog.NewTextHandler(io.Discard, nil)) 29 30 millDir := t.TempDir() 31 bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) 32 if err != nil { 33 t.Fatalf("mill db: %v", err) 34 } 35 bn := notifier.New() 36 mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) 37 mill.Attach(bdb, &bn) 38 if err := bdb.AddExecutorToken("exec-1", HashToken("test-token"), nil, nil); err != nil { 39 t.Fatalf("register executor token: %v", err) 40 } 41 42 srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) 43 defer srv.Close() 44 wsURL := "ws" + strings.TrimPrefix(srv.URL, "http") 45 46 execDir := t.TempDir() 47 edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) 48 if err != nil { 49 t.Fatalf("exec db: %v", err) 50 } 51 en := notifier.New() 52 cfg := &config.Config{} 53 cfg.Server.Dev = true 54 cfg.Server.LogDir = execDir 55 cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") 56 cfg.Server.Hostname = "exec-1" 57 cfg.Mill.URL = wsURL 58 cfg.Mill.Seats = 2 59 cfg.Mill.SharedSecret = "test-token" 60 61 dummyEng := dummy.New(l) 62 dummyEng.StepDelay = 50 * time.Millisecond 63 engines := map[string]models.Engine{"dummy": dummyEng} 64 exec, err := executor.New(cfg, engines, edb, &en, l) 65 if err != nil { 66 t.Fatalf("executor.New: %v", err) 67 } 68 go exec.Connect(ctx) 69 70 be := NewEngine("dummy", mill) 71 twf := tangled.Pipeline_Workflow{ 72 Name: "build", 73 Raw: "steps:\n - name: hello\n command: echo hi\n", 74 } 75 wf, err := be.InitWorkflow(twf, tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{}}) 76 if err != nil { 77 t.Fatalf("InitWorkflow: %v", err) 78 } 79 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "rkey1"}, Name: "build"} 80 81 placeCtx, placeCancel := context.WithTimeout(ctx, 10*time.Second) 82 defer placeCancel() 83 84 slot, err := mill.place(placeCtx, "dummy", wid, wf) 85 if err != nil { 86 t.Fatalf("place: %v", err) 87 } 88 defer slot.Release() 89 90 logPath := models.LogFilePath(millDir, wid) 91 ch := bn.Subscribe() 92 defer bn.Unsubscribe(ch) 93 94 sawLogContent := make(chan bool, 1) 95 go func() { 96 ticker := time.NewTicker(20 * time.Millisecond) 97 defer ticker.Stop() 98 timeout := time.After(10 * time.Second) 99 for { 100 select { 101 case <-ch: 102 case <-ticker.C: 103 case <-timeout: 104 sawLogContent <- false 105 return 106 } 107 data, err := os.ReadFile(logPath) 108 if err == nil && strings.Contains(string(data), "echo hi") { 109 sawLogContent <- true 110 return 111 } 112 } 113 }() 114 115 if err := mill.commitAndWait(placeCtx, wf, nil); err != nil { 116 t.Fatalf("commitAndWait: %v, want success", err) 117 } 118 119 if !<-sawLogContent { 120 t.Fatal("mill live log file never received expected content while job was running") 121 } 122 123 if !waitForFileRemoval(t, logPath) { 124 t.Fatalf("mill live log file was not removed after terminal artifact was recorded") 125 } 126 127 if !waitForStatus(t, bdb, wid, "running") { 128 events, _ := bdb.GetEvents(0, 1000) 129 t.Logf("mill events after completion: %+v", events) 130 t.Fatal("mill never saw streamed running status") 131 } 132} 133 134func TestExecutorConfiguredLabelsAreStoredOnSession(t *testing.T) { 135 ctx, cancel := context.WithCancel(context.Background()) 136 defer cancel() 137 l := slog.New(slog.NewTextHandler(io.Discard, nil)) 138 139 millDir := t.TempDir() 140 bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) 141 if err != nil { 142 t.Fatalf("mill db: %v", err) 143 } 144 bn := notifier.New() 145 mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) 146 mill.Attach(bdb, &bn) 147 if err := bdb.AddExecutorToken("exec-labels", HashToken("test-token"), nil, []string{"linux", "arm64", "gpu"}); err != nil { 148 t.Fatalf("register executor token: %v", err) 149 } 150 151 srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) 152 defer srv.Close() 153 wsURL := "ws" + strings.TrimPrefix(srv.URL, "http") 154 155 execDir := t.TempDir() 156 edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) 157 if err != nil { 158 t.Fatalf("exec db: %v", err) 159 } 160 en := notifier.New() 161 cfg := &config.Config{} 162 cfg.Server.Dev = true 163 cfg.Server.LogDir = execDir 164 cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") 165 cfg.Server.Hostname = "exec-labels" 166 cfg.Mill.URL = wsURL 167 cfg.Mill.Seats = 2 168 cfg.Mill.SharedSecret = "test-token" 169 cfg.Mill.Labels = []string{"linux", "arm64", "gpu"} 170 171 engines := map[string]models.Engine{"dummy": dummy.New(l)} 172 exec, err := executor.New(cfg, engines, edb, &en, l) 173 if err != nil { 174 t.Fatalf("executor.New: %v", err) 175 } 176 go exec.Connect(ctx) 177 178 if !waitForSessionLabels(t, mill, "exec-labels", []string{"linux", "arm64", "gpu"}) { 179 t.Fatal("mill session never stored executor labels from hello") 180 } 181} 182 183func TestEndToEndDummyJobUsesRequiredLabelsAcrossExecutors(t *testing.T) { 184 ctx, cancel := context.WithCancel(context.Background()) 185 defer cancel() 186 l := slog.New(slog.NewTextHandler(io.Discard, nil)) 187 188 millDir := t.TempDir() 189 bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) 190 if err != nil { 191 t.Fatalf("mill db: %v", err) 192 } 193 bn := notifier.New() 194 mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) 195 mill.Attach(bdb, &bn) 196 if err := bdb.AddExecutorToken("exec-x86", HashToken("token-x86"), nil, []string{"linux/amd64", "kvm"}); err != nil { 197 t.Fatalf("register x86 executor token: %v", err) 198 } 199 if err := bdb.AddExecutorToken("exec-arm", HashToken("token-arm"), nil, []string{"linux/arm64", "kvm"}); err != nil { 200 t.Fatalf("register arm executor token: %v", err) 201 } 202 203 srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) 204 defer srv.Close() 205 wsURL := "ws" + strings.TrimPrefix(srv.URL, "http") 206 207 startExecutor := func(name, token string, labels []string) { 208 t.Helper() 209 execDir := t.TempDir() 210 edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) 211 if err != nil { 212 t.Fatalf("%s exec db: %v", name, err) 213 } 214 en := notifier.New() 215 cfg := &config.Config{} 216 cfg.Server.Dev = true 217 cfg.Server.LogDir = execDir 218 cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") 219 cfg.Server.Hostname = name 220 cfg.Mill.URL = wsURL 221 cfg.Mill.Seats = 1 222 cfg.Mill.SharedSecret = token 223 cfg.Mill.Labels = labels 224 225 engines := map[string]models.Engine{"dummy": dummy.New(l)} 226 exec, err := executor.New(cfg, engines, edb, &en, l) 227 if err != nil { 228 t.Fatalf("executor.New: %v", err) 229 } 230 go exec.Connect(ctx) 231 } 232 startExecutor("exec-x86", "token-x86", []string{"linux/amd64", "kvm"}) 233 startExecutor("exec-arm", "token-arm", []string{"linux/arm64", "kvm"}) 234 235 if !waitForSessionLabels(t, mill, "exec-x86", []string{"linux/amd64", "kvm"}) { 236 t.Fatal("x86 executor did not connect with labels") 237 } 238 if !waitForSessionLabels(t, mill, "exec-arm", []string{"linux/arm64", "kvm"}) { 239 t.Fatal("arm executor did not connect with labels") 240 } 241 242 be := NewEngine("dummy", mill) 243 twf := tangled.Pipeline_Workflow{ 244 Name: "build-arm", 245 RunsOn: []string{"linux/arm64"}, 246 Raw: "steps:\n - name: hello\n command: echo hi\n", 247 } 248 wf, err := be.InitWorkflow(twf, tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{}}) 249 if err != nil { 250 t.Fatalf("InitWorkflow: %v", err) 251 } 252 wid := models.WorkflowId{PipelineId: models.PipelineId{Knot: "knot.test", Rkey: "rkey-arm"}, Name: "build-arm"} 253 254 placeCtx, placeCancel := context.WithTimeout(ctx, 10*time.Second) 255 defer placeCancel() 256 257 slot, err := mill.place(placeCtx, "dummy", wid, wf) 258 if err != nil { 259 t.Fatalf("place: %v", err) 260 } 261 defer slot.Release() 262 263 lease := wf.Data.(*millWorkflowState).Lease 264 if lease == nil { 265 t.Fatal("place did not attach lease to workflow state") 266 } 267 if lease.nodeID != "exec-arm" { 268 t.Fatalf("placed on %q, want exec-arm", lease.nodeID) 269 } 270 271 if err := mill.commitAndWait(placeCtx, wf, nil); err != nil { 272 t.Fatalf("commitAndWait: %v, want success", err) 273 } 274} 275 276func waitForSessionLabels(t *testing.T, m *Mill, nodeID string, want []string) bool { 277 t.Helper() 278 deadline := time.Now().Add(3 * time.Second) 279 for time.Now().Before(deadline) { 280 m.mu.Lock() 281 sess := m.sessions[nodeID] 282 var got []string 283 if sess != nil { 284 got = append([]string(nil), sess.labels...) 285 } 286 m.mu.Unlock() 287 if sameStringMultiset(got, want) { 288 return true 289 } 290 time.Sleep(10 * time.Millisecond) 291 } 292 return false 293} 294 295func waitForStatus(t *testing.T, d *db.DB, wid models.WorkflowId, want string) bool { 296 t.Helper() 297 deadline := time.Now().Add(3 * time.Second) 298 aturi := string(wid.PipelineId.AtUri()) 299 for time.Now().Before(deadline) { 300 evs, err := d.GetEvents(0, 1000) 301 if err != nil { 302 t.Fatalf("GetEvents: %v", err) 303 } 304 for _, ev := range evs { 305 if ev.Nsid != tangled.PipelineStatusNSID { 306 continue 307 } 308 var st tangled.PipelineStatus 309 if err := json.Unmarshal(ev.EventJson, &st); err != nil { 310 continue 311 } 312 if st.Pipeline == aturi && st.Workflow == wid.Name && st.Status == want { 313 return true 314 } 315 } 316 time.Sleep(50 * time.Millisecond) 317 } 318 return false 319} 320 321func waitForLogFileContent(t *testing.T, path, want string) bool { 322 t.Helper() 323 deadline := time.Now().Add(5 * time.Second) 324 for time.Now().Before(deadline) { 325 data, err := os.ReadFile(path) 326 if err == nil && strings.Contains(string(data), want) { 327 return true 328 } 329 time.Sleep(20 * time.Millisecond) 330 } 331 return false 332} 333 334func waitForFileRemoval(t *testing.T, path string) bool { 335 t.Helper() 336 deadline := time.Now().Add(5 * time.Second) 337 for time.Now().Before(deadline) { 338 _, err := os.Stat(path) 339 if os.IsNotExist(err) { 340 return true 341 } 342 time.Sleep(20 * time.Millisecond) 343 } 344 return false 345}