This repository has no description
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}