package 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") } }