This repository has no description
0

Configure Feed

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

core / spindle / mill / executor / reserved_test.go
15 kB 546 lines
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}