Something went wrong. Try again.
Monorepo for Tangled
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182package mill
import ( "context" "log/slog" "net/http" "net/http/httptest" "path/filepath" "strings" "sync" "testing" "time"
"tangled.org/core/api/tangled" "tangled.org/core/notifier" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" "tangled.org/core/spindle/mill/executor" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" "tangled.org/core/spindle/secrets")
const executorRunToken = "eyJhbGciOiJFUzI1NiJ9.eyJzdWIiOiJkaWQ6cGxjOnRlc3RyZXBvIn0.c2lnbmF0dXJl"
type noSecretsVault struct{}
func (noSecretsVault) AddSecret(context.Context, secrets.UnlockedSecret) error { return nil }func (noSecretsVault) RemoveSecret(context.Context, secrets.Secret[any]) error { return nil }func (noSecretsVault) RemoveAllSecrets(context.Context, secrets.RepoIdentifier) error { return nil}func (noSecretsVault) GetSecretsLocked(context.Context, secrets.RepoIdentifier) ([]secrets.LockedSecret, error) { return nil, nil}func (noSecretsVault) GetSecretsUnlocked(context.Context, secrets.RepoIdentifier) ([]secrets.UnlockedSecret, error) { return nil, nil}
// recordingEngine stands where a real engine would on the executor: it captures// the credentials its steps execute with.type recordingEngine struct { mu sync.Mutex secrets []secrets.UnlockedSecret ran chan struct{}}
type recordingStep struct{}
func (recordingStep) Name() string { return "run" }func (recordingStep) Command() string { return "true" }func (recordingStep) Kind() models.StepKind { return models.StepKindUser }
func (e *recordingEngine) InitWorkflow(twf tangled.Pipeline_Workflow, _ tangled.Pipeline) (*models.Workflow, error) { return &models.Workflow{Name: twf.Name, Steps: []models.Step{recordingStep{}}}, nil}
func (e *recordingEngine) AcquireWorkflowSlot(context.Context, models.WorkflowId, *models.Workflow, engine.AcquireMode) (engine.WorkflowSlot, error) { return engine.NoopSlot{}, nil}
func (e *recordingEngine) SetupWorkflow(context.Context, models.WorkflowId, *models.Workflow, models.WorkflowLogger) error { return nil}
func (e *recordingEngine) WorkflowTimeout() time.Duration { return time.Minute }
func (e *recordingEngine) DestroyWorkflow(context.Context, models.WorkflowId) error { return nil }
func (e *recordingEngine) RunStep(_ context.Context, _ models.WorkflowId, _ *models.Workflow, _ int, unlocked []secrets.UnlockedSecret, _ models.WorkflowLogger) error { e.mu.Lock() e.secrets = unlocked e.mu.Unlock() select { case e.ran <- struct{}{}: default: } return nil}
func (e *recordingEngine) runSecrets() []secrets.UnlockedSecret { e.mu.Lock() defer e.mu.Unlock() return e.secrets}
func TestExecutorHostedStepsCarryTheRunIDToken(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() var logs strings.Builder l := slog.New(slog.NewTextHandler(&logs, nil))
millDir := t.TempDir() bdb, err := db.Make(ctx, filepath.Join(millDir, "mill.db")) if err != nil { t.Fatalf("mill db: %v", err) } bn := notifier.New() mill := New(l, Config{LogDir: millDir, ReconnectGrace: time.Minute, BidTimeout: 2 * time.Second}) mill.Attach(bdb, &bn, testQuotaManager(t, bdb)) mill.RegisterMetrics(observability.NewMetrics()) registerTestExecutor(t, bdb, "exec-1", HashToken("test-token"), nil)
srv := httptest.NewServer(http.HandlerFunc(mill.HandleExecutorConn)) defer srv.Close() wsURL := "ws" + strings.TrimPrefix(srv.URL, "http")
execDir := t.TempDir() edb, err := db.Make(ctx, filepath.Join(execDir, "exec.db")) if err != nil { t.Fatalf("exec db: %v", err) } en := notifier.New() cfg := &config.Config{} cfg.Server.LogDir = execDir cfg.Server.Hostname = "exec-1" cfg.ArtifactStores.Disk.Dir = filepath.Join(execDir, "artifacts") cfg.Mill.ArtifactStore = "disk" cfg.Mill.URL = wsURL cfg.Mill.Seats = 2 cfg.Mill.SharedSecret = "test-token"
runner := &recordingEngine{ran: make(chan struct{}, 1)} exec, err := executor.New(cfg, map[string]models.Engine{"dummy": runner}, edb, &en, l, nil, nil) if err != nil { t.Fatalf("executor.New: %v", err) } exec.RegisterMetrics(observability.NewMetrics()) go exec.Connect(observability.WithMetrics(ctx, observability.NewMetrics()))
// the mill mints nothing itself here: it is handed the run's token exactly // as runJob hands it over, and has to ship it to the executor be := NewEngine("dummy", mill) pipelineId := models.PipelineId("3mrkp6iz6os2o") wid := models.WorkflowId{PipelineId: pipelineId, Name: "build"} audience := "//iam.googleapis.com/projects/tangled/locations/global/workloadIdentityPools/ci/providers/spindle" twf := tangled.Pipeline_Workflow{Name: "build", Audience: &audience, Raw: "steps:\n - name: run\n command: true\n"} wf, err := be.InitWorkflow(twf, testPipeline()) if err != nil { t.Fatalf("InitWorkflow: %v", err) } wf.Audience = audience if err := bdb.StatusPending(wid, &bn); err != nil { t.Fatal(err) }
runCtx, runCancel := context.WithTimeout(observability.WithMetrics(ctx, observability.NewMetrics()), 20*time.Second) defer runCancel() go engine.StartWorkflows( l, noSecretsVault{}, &config.Config{}, testQuotaManager(t, bdb), nil, bdb, &bn, nil, nil, runCtx, &models.Pipeline{ RepoDid: "did:plc:testrepo", TrustedSource: true, MintIDToken: func(time.Duration, string) (string, error) { return executorRunToken, nil }, Workflows: map[models.Engine][]models.Workflow{be: {*wf}}, }, pipelineId, )
select { case <-runner.ran: case <-time.After(15 * time.Second): evs, _ := bdb.GetEvents(0, 100) t.Fatalf("the executor never ran the workflow; mill events: %+v; logs:\n%s", evs, logs.String()) }
var got string for _, s := range runner.runSecrets() { if s.Key == models.IDTokenEnvVar { got = s.Value } } if got != executorRunToken { t.Fatalf("%s on the executor = %q, want the run's id token: an executor-hosted run must receive it", models.IDTokenEnvVar, got) } if !waitForStatus(t, bdb, wid, "success") { t.Fatal("workflow never reached success") }}