Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202package executor
import ( "context" "encoding/json" "errors" "io" "log/slog" "testing"
"tangled.org/core/api/tangled" "tangled.org/core/spindle/engine" millv1 "tangled.org/core/spindle/mill/proto/gen" "tangled.org/core/spindle/models")
type placementEngine struct { *fakeEngine validationErr error validated bool}
func (e *placementEngine) ValidateWorkflowPlacement(*models.Workflow) error { e.validated = true return e.validationErr}
func TestHandleReserveRejectsMissingTriggerMetadata(t *testing.T) { enc := newCaptureEncoder() e := &Executor{ l: slog.New(slog.NewTextHandler(io.Discard, nil)), enc: enc, seats: 1, engines: map[string]models.Engine{"microvm": &fakeEngine{}}, active: make(map[string]*reservation), } twf, err := json.Marshal(tangled.Pipeline_Workflow{Name: "build"}) if err != nil { t.Fatal(err) } tpl, err := json.Marshal(tangled.Pipeline{}) if err != nil { t.Fatal(err) }
// engines deref TriggerMetadata unconditionally. this must be a reject, // not a panic that takes the whole executor down e.handleReserve(context.Background(), &millv1.ReserveSeat{ LeaseId: "lease-1", TargetEngine: "microvm", RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), Knot: "k", Rkey: "r", })
result := (<-enc.messages).GetReserveResult() if result == nil { t.Fatal("handleReserve() did not send ReserveResult") } if result.GetAccepted() { t.Fatal("handleReserve() accepted a pipeline without trigger metadata") } if result.GetRejectClass() != millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE { t.Fatalf("reject class = %v, want incompatible", result.GetRejectClass()) } if len(e.active) != 0 { t.Fatalf("active reservations = %d, want 0", len(e.active)) }}
func TestHandleReserveValidatesPlacementBeforeAcquiringSlot(t *testing.T) { enc := newCaptureEncoder() validationErr := errors.New("image architecture is not native") eng := &placementEngine{fakeEngine: &fakeEngine{}, validationErr: validationErr} e := &Executor{ l: slog.New(slog.NewTextHandler(io.Discard, nil)), enc: enc, seats: 1, engines: map[string]models.Engine{"microvm": eng}, active: make(map[string]*reservation), } twf, err := json.Marshal(tangled.Pipeline_Workflow{Name: "build"}) if err != nil { t.Fatal(err) } tpl, err := json.Marshal(tangled.Pipeline{TriggerMetadata: &tangled.Pipeline_TriggerMetadata{}}) if err != nil { t.Fatal(err) }
e.handleReserve(context.Background(), &millv1.ReserveSeat{ LeaseId: "lease-1", TargetEngine: "microvm", RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), Knot: "k", Rkey: "r", RepoDid: "did:web:example.com", })
result := (<-enc.messages).GetReserveResult() if result == nil { t.Fatal("handleReserve() did not send ReserveResult") } if result.GetAccepted() { t.Fatal("handleReserve() accepted placement validation failure") } if result.GetRejectClass() != millv1.RejectClass_REJECT_CLASS_INCOMPATIBLE { t.Fatalf("reject class = %v, want incompatible", result.GetRejectClass()) } if !eng.validated { t.Fatal("placement validator was not called") } if eng.acquireCalled { t.Fatal("slot acquisition ran after placement validation failed") } if len(e.active) != 0 { t.Fatalf("active reservations = %d, want 0", len(e.active)) }}
func TestHandleReservePreservesTypedInitFailure(t *testing.T) { enc := newCaptureEncoder() eng := &fakeEngine{ initErr: engine.ClassifiedFailure( engine.FailureClassUser, engine.FailureReasonWorkflowInvalid, errors.New("invalid workflow"), ), } e := &Executor{ l: slog.New(slog.NewTextHandler(io.Discard, nil)), enc: enc, seats: 1, engines: map[string]models.Engine{"microvm": eng}, active: make(map[string]*reservation), }
e.handleReserve(context.Background(), testReserveSeat(t, "lease-1", "microvm")) result := (<-enc.messages).GetReserveResult() if result == nil || result.GetAccepted() { t.Fatalf("reserve result = %+v, want rejection", result) } if result.GetFailureClass() != string(engine.FailureClassUser) || result.GetFailureReason() != string(engine.FailureReasonWorkflowInvalid) { t.Fatalf( "failure attribution = %q/%q, want user/workflow_invalid", result.GetFailureClass(), result.GetFailureReason(), ) }}
func TestHandleReserveKeepsTriggerMetadataForExecution(t *testing.T) { e := testExecutor(t) e.enc = newCaptureEncoder() e.seats = 1 e.engines = map[string]models.Engine{"microvm": &fakeEngine{}}
repoDID := "did:web:example.com" metadata := &tangled.Pipeline_TriggerMetadata{ Kind: "push", Push: &tangled.Pipeline_PushTriggerData{ NewSha: "0123456789abcdef", Ref: "refs/heads/main", }, Repo: &tangled.Pipeline_TriggerRepo{RepoDid: &repoDID}, } twf, err := json.Marshal(tangled.Pipeline_Workflow{Name: "build"}) if err != nil { t.Fatal(err) } tpl, err := json.Marshal(tangled.Pipeline{TriggerMetadata: metadata}) if err != nil { t.Fatal(err) }
e.handleReserve(context.Background(), &millv1.ReserveSeat{ LeaseId: "lease-1", TargetEngine: "microvm", RawWorkflowJson: string(twf), RawPipelineJson: string(tpl), Knot: "k", Rkey: "r", RepoDid: repoDID, })
result := (<-e.enc.(*captureEncoder).messages).GetReserveResult() if result == nil || !result.GetAccepted() { t.Fatalf("ReserveResult = %+v, want accepted", result) } res := e.active["lease-1"] if res == nil || res.pipeline.TriggerMetadata == nil || res.pipeline.TriggerMetadata.Push == nil { t.Fatal("reservation dropped trigger metadata") } if got := res.pipeline.TriggerMetadata.Push.NewSha; got != metadata.Push.NewSha { t.Fatalf("trigger commit = %q, want %q", got, metadata.Push.NewSha) } res.ttlTimer.Stop()}