This repository has no description
1package executor
2
3import (
4 "context"
5 "encoding/json"
6 "github.com/bluesky-social/indigo/atproto/syntax"
7 "github.com/gorilla/websocket"
8 "google.golang.org/protobuf/proto"
9 "io"
10 "log/slog"
11 "net/http"
12 "net/http/httptest"
13 "path/filepath"
14 "strings"
15 "testing"
16 "time"
17
18 "tangled.org/core/api/tangled"
19 "tangled.org/core/notifier"
20 "tangled.org/core/spindle/config"
21 "tangled.org/core/spindle/db"
22 "tangled.org/core/spindle/engine"
23 millproto "tangled.org/core/spindle/mill/proto"
24 millv1 "tangled.org/core/spindle/mill/proto/gen"
25 "tangled.org/core/spindle/models"
26 "tangled.org/core/spindle/secrets"
27)
28
29type captureEncoder struct {
30 messages chan *millproto.Message
31}
32
33func newCaptureEncoder() *captureEncoder {
34 return &captureEncoder{messages: make(chan *millproto.Message, 4)}
35}
36
37func (e *captureEncoder) Encode(msg *millproto.Message) error {
38 e.messages <- msg
39 return nil
40}
41
42type fakeSlot struct{ released int }
43
44func (s *fakeSlot) Release() { s.released++ }
45
46type fakeEngine struct {
47 setupCalled bool
48 runCalled bool
49 destroyCalled bool
50 acquireCalled bool
51 secrets chan []secrets.UnlockedSecret
52 done chan struct{}
53}
54
55func (e *fakeEngine) InitWorkflow(twf tangled.Pipeline_Workflow, tpl tangled.Pipeline) (*models.Workflow, error) {
56 return &models.Workflow{Name: twf.Name}, nil
57}
58func (e *fakeEngine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, l models.WorkflowLogger) error {
59 e.setupCalled = true
60 return nil
61}
62func (e *fakeEngine) WorkflowTimeout() time.Duration { return 7 * time.Minute }
63func (e *fakeEngine) DestroyWorkflow(ctx context.Context, wid models.WorkflowId) error {
64 e.destroyCalled = true
65 return nil
66}
67func (e *fakeEngine) RunStep(ctx context.Context, wid models.WorkflowId, w *models.Workflow, idx int, s []secrets.UnlockedSecret, l models.WorkflowLogger) error {
68 e.runCalled = true
69 if e.secrets != nil {
70 e.secrets <- s
71 }
72 if e.done != nil {
73 close(e.done)
74 }
75 return nil
76}
77func (e *fakeEngine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, mode engine.AcquireMode) (engine.WorkflowSlot, error) {
78 e.acquireCalled = true
79 return &fakeSlot{}, nil
80}
81
82type fakeStep struct{}
83
84func (fakeStep) Name() string { return "test" }
85func (fakeStep) Command() string { return "true" }
86func (fakeStep) Kind() models.StepKind { return models.StepKindUser }
87
88func TestNewFailsWhenOutboxCannotInitialize(t *testing.T) {
89 d := testDB(t)
90 if err := d.Close(); err != nil {
91 t.Fatal(err)
92 }
93 n := notifier.New()
94 cfg := &config.Config{}
95 if _, err := New(cfg, nil, d, &n, slog.New(slog.NewTextHandler(io.Discard, nil))); err == nil {
96 t.Fatal("New succeeded with an unavailable outbox database")
97 }
98}
99
100func TestReservedEngineHandsBackHeldSlotOnce(t *testing.T) {
101 inner := &fakeEngine{}
102 slot := &fakeSlot{}
103 re := newReservedEngine(inner, slot)
104
105 got, err := re.(engine.WorkflowSlotter).AcquireWorkflowSlot(context.Background(), models.WorkflowId{}, nil, engine.Wait)
106 if err != nil {
107 t.Fatal(err)
108 }
109 if got != engine.WorkflowSlot(slot) {
110 t.Fatal("AcquireWorkflowSlot() did not return the held slot")
111 }
112 if inner.acquireCalled {
113 t.Fatal("wrapper must not call the inner engine's AcquireWorkflowSlot")
114 }
115
116 if _, err := re.(engine.WorkflowSlotter).AcquireWorkflowSlot(context.Background(), models.WorkflowId{}, nil, engine.Wait); err == nil {
117 t.Fatal("second AcquireWorkflowSlot() should error")
118 }
119}
120
121func TestHandleCommitIsIdempotent(t *testing.T) {
122 enc := newCaptureEncoder()
123 e := testExecutor(t)
124 e.enc = enc
125 e.active["lease-1"] = &reservation{leaseID: "lease-1", committed: true}
126
127 e.handleCommit(context.Background(), &millv1.CommitLease{LeaseId: "lease-1"})
128 msg := <-enc.messages
129 if got := msg.GetCommitted().GetLeaseId(); got != "lease-1" {
130 t.Fatalf("Committed lease = %q, want lease-1", got)
131 }
132}
133
134func TestHandleCommitRejectsMissingReservation(t *testing.T) {
135 enc := newCaptureEncoder()
136 e := testExecutor(t)
137 e.enc = enc
138
139 e.handleCommit(context.Background(), &millv1.CommitLease{LeaseId: "expired"})
140 result := (<-enc.messages).GetReserveResult()
141 if result == nil {
142 t.Fatal("missing reservation commit did not receive a ReserveResult")
143 }
144 if result.GetLeaseId() != "expired" || result.GetAccepted() {
145 t.Fatalf("ReserveResult = %+v, want correlated rejection", result)
146 }
147}
148
149func TestHandleCancelFinalizesExpiredReservation(t *testing.T) {
150 enc := newCaptureEncoder()
151 e := testExecutor(t)
152 e.enc = enc
153
154 e.handleCancel("lease-expired")
155 var ack *millv1.CancelAck
156 for ack == nil {
157 select {
158 case msg := <-enc.messages:
159 ack = msg.GetCancelAck()
160 case <-time.After(time.Second):
161 t.Fatal("cancel acknowledgement timed out")
162 }
163 }
164 if ack.GetLeaseId() != "lease-expired" {
165 t.Fatalf("CancelAck = %+v, want lease-expired", ack)
166 }
167 rows, err := e.db.ListOutboxRows()
168 if err != nil {
169 t.Fatal(err)
170 }
171 if len(rows) != 1 {
172 t.Fatalf("cancel terminal outbox rows = %d, want 1", len(rows))
173 }
174 var entry millv1.Event
175 if err := proto.Unmarshal(rows[0].Payload, &entry); err != nil {
176 t.Fatal(err)
177 }
178 if got := entry.GetAttemptResult().GetStatus(); got != millv1.TerminalStatus_CANCELLED {
179 t.Fatalf("cancel terminal = %v, want CANCELLED", got)
180 }
181}
182
183func TestHandleCommitPreservesPreauthorizedSecrets(t *testing.T) {
184 d := testDB(t)
185 n := notifier.New()
186 enc := newCaptureEncoder()
187 e := &Executor{
188 cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}},
189 db: d,
190 n: &n,
191 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
192 active: make(map[string]*reservation),
193 maxOutboxBytes: 10 * 1024 * 1024,
194 enc: enc,
195 }
196 if err := e.initOutbox(); err != nil {
197 t.Fatal(err)
198 }
199 e.lifecycleCtx = context.Background()
200
201 inner := &fakeEngine{secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{})}
202 slot := &fakeSlot{}
203 repoDid, err := syntax.ParseDID("did:web:example.com")
204 if err != nil {
205 t.Fatal(err)
206 }
207 res := &reservation{
208 leaseID: "lease-1",
209 wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"},
210 realEngine: inner,
211 slot: slot,
212 wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}},
213 repoDid: repoDid,
214 }
215 e.active[res.leaseID] = res
216
217 e.handleCommit(context.Background(), &millv1.CommitLease{
218 LeaseId: res.leaseID,
219 Secrets: []*millv1.Secret{{Key: "TOKEN", Value: "secret-value"}},
220 })
221 if got := (<-enc.messages).GetCommitted().GetLeaseId(); got != res.leaseID {
222 t.Fatalf("Committed lease = %q, want %q", got, res.leaseID)
223 }
224 select {
225 case got := <-inner.secrets:
226 if len(got) != 1 || got[0].Key != "TOKEN" || got[0].Value != "secret-value" {
227 t.Fatalf("RunStep secrets = %+v", got)
228 }
229 case <-time.After(2 * time.Second):
230 t.Fatal("RunStep did not receive CommitLease secrets")
231 }
232 select {
233 case <-inner.done:
234 case <-time.After(2 * time.Second):
235 t.Fatal("workflow did not finish")
236 }
237 if res.stopTail != nil {
238 res.stopTail()
239 }
240 e.jobsWG.Wait()
241}
242
243func TestRunSessionCancellationClosesStalledWebsocket(t *testing.T) {
244 connected := make(chan struct{})
245 release := make(chan struct{})
246 srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
247 conn, err := websocket.Upgrade(w, r, nil, 1024, 1024)
248 if err != nil {
249 return
250 }
251 defer conn.Close()
252 close(connected)
253 <-release
254 }))
255 t.Cleanup(func() {
256 close(release)
257 srv.Close()
258 })
259
260 e := testSessionExecutor(t, "ws"+strings.TrimPrefix(srv.URL, "http"))
261 ctx, cancel := context.WithCancel(context.Background())
262 done := make(chan error, 1)
263 go func() { done <- e.runSession(ctx) }()
264 <-connected
265 cancel()
266
267 select {
268 case <-done:
269 case <-time.After(2 * time.Second):
270 t.Fatal("runSession did not return after context cancellation")
271 }
272}
273
274func testSessionExecutor(t *testing.T, url string) *Executor {
275 d := testDB(t)
276 e := &Executor{
277 millURL: url,
278 seats: 1,
279 engines: make(map[string]models.Engine),
280 db: d,
281 cfg: &config.Config{Server: config.Server{Dev: true}},
282 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
283 active: make(map[string]*reservation),
284 maxOutboxBytes: 10 * 1024 * 1024,
285 }
286 if err := e.initOutbox(); err != nil {
287 t.Fatal(err)
288 }
289 return e
290}
291
292func testExecutor(t *testing.T) *Executor {
293 d := testDB(t)
294 e := &Executor{
295 db: d,
296 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
297 active: make(map[string]*reservation),
298 maxOutboxBytes: 10 * 1024 * 1024,
299 }
300 if err := e.initOutbox(); err != nil {
301 t.Fatal(err)
302 }
303 return e
304}
305
306func testDB(t *testing.T) *db.DB {
307 d, err := db.Make(context.Background(), filepath.Join(t.TempDir(), "spindle.db"))
308 if err != nil {
309 t.Fatal(err)
310 }
311 return d
312}
313
314func TestFinishJobReportsCancelledReservationAsCancelled(t *testing.T) {
315 d := testDB(t)
316 res := &reservation{
317 leaseID: "lease-1",
318 wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"},
319 cancelled: true,
320 }
321 e := &Executor{
322 db: d,
323 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
324 active: map[string]*reservation{res.leaseID: res},
325 maxOutboxBytes: 10 * 1024 * 1024,
326 }
327 if err := e.initOutbox(); err != nil {
328 t.Fatal(err)
329 }
330
331 e.finishJob(res, &tangled.PipelineStatus{
332 Pipeline: string(res.wid.PipelineId.AtUri()),
333 Workflow: res.wid.Name,
334 Status: string(models.StatusKindFailed),
335 })
336
337 rows, err := d.ListOutboxRows()
338 if err != nil {
339 t.Fatal(err)
340 }
341 if len(rows) != 1 {
342 t.Fatalf("outbox rows = %d, want 1", len(rows))
343 }
344
345 var entry millv1.Event
346 if err := proto.Unmarshal(rows[0].Payload, &entry); err != nil {
347 t.Fatal(err)
348 }
349 got := entry.GetAttemptResult().GetStatus()
350 if got != millv1.TerminalStatus_CANCELLED {
351 t.Fatalf("terminal status = %v, want CANCELLED", got)
352 }
353}
354
355func TestReplayRejectsMalformedOutboxRow(t *testing.T) {
356 d := testDB(t)
357 e := &Executor{
358 db: d,
359 enc: newCaptureEncoder(),
360 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
361 }
362 if err := e.initOutbox(); err != nil {
363 t.Fatal(err)
364 }
365 if _, err := d.AppendOutboxRow([]byte("not protobuf"), true); err != nil {
366 t.Fatal(err)
367 }
368 if err := e.replay(0); err == nil {
369 t.Fatal("replay accepted a malformed row and would leave a permanent seqno gap")
370 }
371}
372
373func TestSocketCancellationIndependence(t *testing.T) {
374 d := testDB(t)
375 n := notifier.New()
376 e := &Executor{
377 db: d,
378 n: &n,
379 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
380 active: make(map[string]*reservation),
381 maxOutboxBytes: 10 * 1024 * 1024,
382 cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}},
383 }
384 if err := e.initOutbox(); err != nil {
385 t.Fatal(err)
386 }
387
388 lifecycleCtx, cancelLifecycle := context.WithCancel(context.Background())
389 defer cancelLifecycle()
390 e.lifecycleCtx = lifecycleCtx
391
392 inner := &fakeEngine{secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{})}
393 slot := &fakeSlot{}
394 repoDid, err := syntax.ParseDID("did:web:example.com")
395 if err != nil {
396 t.Fatal(err)
397 }
398 res := &reservation{
399 leaseID: "lease-1",
400 wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"},
401 realEngine: inner,
402 slot: slot,
403 wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}},
404 repoDid: repoDid,
405 }
406 e.active[res.leaseID] = res
407
408 sessionCtx, cancelSession := context.WithCancel(lifecycleCtx)
409
410 e.handleCommit(sessionCtx, &millv1.CommitLease{
411 LeaseId: "lease-1",
412 })
413
414 cancelSession()
415
416 // session disconnect must not cancel the running job
417 select {
418 case <-inner.done:
419 case <-time.After(2 * time.Second):
420 t.Fatal("workflow did not complete even though websocket session was cancelled")
421 }
422
423 e.jobsWG.Wait()
424}
425
426func TestMonotonicSnapshots(t *testing.T) {
427 enc := newCaptureEncoder()
428 e := &Executor{
429 enc: enc,
430 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
431 active: make(map[string]*reservation),
432 maxOutboxBytes: 10 * 1024 * 1024,
433 }
434
435 e.pushSnapshot()
436 msg1 := <-enc.messages
437 seq1 := msg1.GetNodeSnapshot().GetSeqno()
438 if seq1 != 1 {
439 t.Fatalf("first seq = %d, want 1", seq1)
440 }
441
442 e.pushSnapshot()
443 msg2 := <-enc.messages
444 seq2 := msg2.GetNodeSnapshot().GetSeqno()
445 if seq2 != 2 {
446 t.Fatalf("second seq = %d, want 2", seq2)
447 }
448}
449
450func TestTimerRace(t *testing.T) {
451 d := testDB(t)
452 e := &Executor{
453 db: d,
454 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
455 active: make(map[string]*reservation),
456 maxOutboxBytes: 10 * 1024 * 1024,
457 }
458 if err := e.initOutbox(); err != nil {
459 t.Fatal(err)
460 }
461
462 twf, _ := json.Marshal(tangled.Pipeline_Workflow{Name: "build"})
463 tpl, _ := json.Marshal(tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{}})
464
465 inner := &fakeEngine{}
466 e.engines = map[string]models.Engine{"microvm": inner}
467
468 e.handleReserve(context.Background(), &millv1.ReserveSeat{
469 LeaseId: "lease-1",
470 TargetEngine: "microvm",
471 RawWorkflowJson: string(twf),
472 RawPipelineJson: string(tpl),
473 TtlSeconds: 1,
474 })
475
476 e.mu.Lock()
477 res := e.active["lease-1"]
478 e.mu.Unlock()
479
480 if res == nil {
481 t.Fatal("reservation was not added")
482 }
483
484 deadline := time.Now().Add(5 * time.Second)
485 for {
486 e.mu.Lock()
487 activeLen := len(e.active)
488 e.mu.Unlock()
489 if activeLen == 0 {
490 break
491 }
492 if time.Now().After(deadline) {
493 t.Fatal("reservation was leaked and never expired")
494 }
495 time.Sleep(10 * time.Millisecond)
496 }
497}
498
499func TestStructuredShutdown(t *testing.T) {
500 d := testDB(t)
501 n := notifier.New()
502 e := &Executor{
503 db: d,
504 n: &n,
505 l: slog.New(slog.NewTextHandler(io.Discard, nil)),
506 active: make(map[string]*reservation),
507 maxOutboxBytes: 10 * 1024 * 1024,
508 cfg: &config.Config{Server: config.Server{LogDir: t.TempDir()}},
509 }
510 if err := e.initOutbox(); err != nil {
511 t.Fatal(err)
512 }
513
514 ctx, cancel := context.WithCancel(context.Background())
515 e.lifecycleCtx = ctx
516
517 inner := &fakeEngine{secrets: make(chan []secrets.UnlockedSecret, 1), done: make(chan struct{})}
518 slot := &fakeSlot{}
519 repoDid, err := syntax.ParseDID("did:web:example.com")
520 if err != nil {
521 t.Fatal(err)
522 }
523 res := &reservation{
524 leaseID: "lease-1",
525 wid: models.WorkflowId{PipelineId: models.PipelineId{Knot: "k", Rkey: "r"}, Name: "build"},
526 realEngine: inner,
527 slot: slot,
528 wf: &models.Workflow{Name: "build", Steps: []models.Step{fakeStep{}}},
529 repoDid: repoDid,
530 }
531 e.active[res.leaseID] = res
532
533 e.handleCommit(ctx, &millv1.CommitLease{
534 LeaseId: "lease-1",
535 })
536
537 cancel()
538
539 e.jobsWG.Wait()
540
541 select {
542 case <-inner.done:
543 default:
544 t.Fatal("shutdown returned but running job did not finish")
545 }
546}