Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450package engine
import ( "context" "errors" "fmt" "io" "log/slog" "os" "path/filepath" "slices" "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/models" "tangled.org/core/spindle/observability" "tangled.org/core/spindle/secrets")
const testRunToken = "eyJhbGciOiJFUzI1NiJ9.eyJzdWIiOiJkaWQ6cGxjOmJvbHRsZXNzIn0.c2lnbmF0dXJl"
type staticVault struct { secrets []secrets.UnlockedSecret}
func (v staticVault) AddSecret(context.Context, secrets.UnlockedSecret) error { return nil }func (v staticVault) RemoveSecret(context.Context, secrets.Secret[any]) error { return nil }func (v staticVault) RemoveAllSecrets(context.Context, secrets.RepoIdentifier) error { return nil}func (v staticVault) GetSecretsLocked(context.Context, secrets.RepoIdentifier) ([]secrets.LockedSecret, error) { return nil, nil}func (v staticVault) GetSecretsUnlocked(context.Context, secrets.RepoIdentifier) ([]secrets.UnlockedSecret, error) { return v.secrets, nil}
type echoStep struct{}
func (echoStep) Name() string { return "echo" }func (echoStep) Command() string { return "env" }func (echoStep) Kind() models.StepKind { return models.StepKindUser }
// echoEngine prints every credential it is handed, the way `set -x` prints the// environment of a real step. With a gate it holds each workflow at its slot so// a test can observe what happens while a run waits for capacity.type slotWait struct { name string grant chan struct{}}
type echoEngine struct { gate chan slotWait timeout time.Duration
mu sync.Mutex secrets map[string][]secrets.UnlockedSecret ran map[string]chan struct{} ranClosed map[string]bool order []string}
func newEchoEngine() *echoEngine { return &echoEngine{ secrets: map[string][]secrets.UnlockedSecret{}, ran: map[string]chan struct{}{}, ranClosed: map[string]bool{}, }}
func (e *echoEngine) InitWorkflow(twf tangled.Pipeline_Workflow, _ tangled.Pipeline) (*models.Workflow, error) { e.mu.Lock() defer e.mu.Unlock() if e.ran[twf.Name] == nil { e.ran[twf.Name] = make(chan struct{}) } return &models.Workflow{Name: twf.Name, Steps: []models.Step{echoStep{}}}, nil}
func (e *echoEngine) AcquireWorkflowSlot(ctx context.Context, _ models.WorkflowId, wf *models.Workflow, _ AcquireMode) (WorkflowSlot, error) { e.note("slot:" + wf.Name) if e.gate == nil { return NoopSlot{}, nil } grant := make(chan struct{}) select { case e.gate <- slotWait{name: wf.Name, grant: grant}: case <-ctx.Done(): return nil, ctx.Err() } select { case <-grant: return NoopSlot{}, nil case <-ctx.Done(): return nil, ctx.Err() }}
func (e *echoEngine) SetupWorkflow(context.Context, models.WorkflowId, *models.Workflow, models.WorkflowLogger) error { return nil}
func (e *echoEngine) WorkflowTimeout() time.Duration { if e.timeout == 0 { return time.Minute } return e.timeout}
func (e *echoEngine) DestroyWorkflow(context.Context, models.WorkflowId) error { return nil }
func (e *echoEngine) RunStep(_ context.Context, wid models.WorkflowId, w *models.Workflow, idx int, unlocked []secrets.UnlockedSecret, wfLogger models.WorkflowLogger) error { e.mu.Lock() e.secrets[wid.Name] = unlocked e.order = append(e.order, "steps:"+wid.Name) if !e.ranClosed[wid.Name] { e.ranClosed[wid.Name] = true close(e.ran[wid.Name]) } e.mu.Unlock()
for _, s := range unlocked { fmt.Fprintf(wfLogger.DataWriter(idx, "stdout"), "TANGLED_ID_TOKEN=%s\n", s.Value) } if _, ok := w.Environment[models.IDTokenEnvVar]; ok { return errors.New("the id token was merged into the workflow env") } return nil}
func (e *echoEngine) note(event string) { e.mu.Lock() defer e.mu.Unlock() e.order = append(e.order, event)}
func (e *echoEngine) events() []string { e.mu.Lock() defer e.mu.Unlock() return append([]string(nil), e.order...)}
func (e *echoEngine) ranWorkflow(name string) <-chan struct{} { e.mu.Lock() defer e.mu.Unlock() if e.ran[name] == nil { e.ran[name] = make(chan struct{}) } return e.ran[name]}
func (e *echoEngine) runSecrets(name string) []secrets.UnlockedSecret { e.mu.Lock() defer e.mu.Unlock() return e.secrets[name]}
func secretValue(all []secrets.UnlockedSecret, key string) string { for _, s := range all { if s.Key == key { return s.Value } } return ""}
// startEchoPipeline runs StartWorkflows to completion. It reports errors instead// of failing a test so it can be driven from another goroutine.func startEchoPipeline(eng *echoEngine, pipeline *models.Pipeline, names ...string) (logPath string, cleanup func(), err error) { records := make([]tangled.Pipeline_Workflow, len(names)) for i, name := range names { records[i] = tangled.Pipeline_Workflow{Name: name} } return startEchoPipelineRecords(eng, pipeline, records...)}
func startEchoPipelineRecords(eng *echoEngine, pipeline *models.Pipeline, records ...tangled.Pipeline_Workflow) (logPath string, cleanup func(), err error) { ctx := observability.WithMetrics(context.Background(), observability.NewMetrics())
logDir, err := os.MkdirTemp("", "spindle-logs") if err != nil { return "", func() {}, err } repoDir, err := os.MkdirTemp("", "spindle-repos") if err != nil { os.RemoveAll(logDir) return "", func() {}, err } dbDir, err := os.MkdirTemp("", "spindle-db") if err != nil { os.RemoveAll(logDir) os.RemoveAll(repoDir) return "", func() {}, err } cleanup = func() { os.RemoveAll(logDir) os.RemoveAll(repoDir) os.RemoveAll(dbDir) }
d, err := db.Make(ctx, filepath.Join(dbDir, "spindle.db")) if err != nil { cleanup() return "", func() {}, err }
cfg := &config.Config{} cfg.Server.LogDir = logDir cfg.Server.RepoDir = repoDir
n := notifier.New() pipelineId := models.PipelineId("3mrkp6iz6os2o") workflows := map[models.Engine][]models.Workflow{} for _, record := range records { wf, err := eng.InitWorkflow(record, tangled.Pipeline{}) if err != nil { cleanup() return "", func() {}, err } if record.Audience != nil { wf.Audience = *record.Audience } if err := d.StatusPending(models.WorkflowId{PipelineId: pipelineId, Name: record.Name}, &n); err != nil { cleanup() return "", func() {}, err } workflows[eng] = append(workflows[eng], *wf) } pipeline.Workflows = workflows
StartWorkflows( slog.New(slog.NewTextHandler(io.Discard, nil)), staticVault{secrets: []secrets.UnlockedSecret{{Key: "VAULT_SECRET", Value: "vault-value"}}}, cfg, nil, nil, d, &n, nil, nil, ctx, pipeline, pipelineId, )
if err := d.Close(); err != nil { cleanup() return "", func() {}, err } return models.LogFilePath(logDir, models.WorkflowId{PipelineId: pipelineId, Name: records[0].Name}), cleanup, nil}
func runEchoPipeline(t *testing.T, eng *echoEngine, pipeline *models.Pipeline, names ...string) string { t.Helper() logPath, cleanup, err := startEchoPipeline(eng, pipeline, names...) if err != nil { cleanup() t.Fatal(err) } t.Cleanup(cleanup) return logPath}
func runEchoPipelineRecords(t *testing.T, eng *echoEngine, pipeline *models.Pipeline, records ...tangled.Pipeline_Workflow) string { t.Helper() logPath, cleanup, err := startEchoPipelineRecords(eng, pipeline, records...) if err != nil { cleanup() t.Fatal(err) } t.Cleanup(cleanup) return logPath}
func audienceRecord(name, audience string) tangled.Pipeline_Workflow { return tangled.Pipeline_Workflow{Name: name, Audience: &audience}}
func TestWorkflowLogMasksTheRunIDToken(t *testing.T) { eng := newEchoEngine() logPath := runEchoPipelineRecords(t, eng, &models.Pipeline{ RepoDid: "did:plc:boltless", TrustedSource: true, MintIDToken: func(time.Duration, string) (string, error) { return testRunToken, nil }, }, audienceRecord("build", "vault"))
if value := secretValue(eng.runSecrets("build"), models.IDTokenEnvVar); value != testRunToken { t.Fatalf("%s handed to the step = %q, want the run's id token", models.IDTokenEnvVar, value) }
logged, err := os.ReadFile(logPath) if err != nil { t.Fatal(err) } if strings.Contains(string(logged), testRunToken) { t.Fatal("the persisted log kept the id token in the clear") } if !strings.Contains(string(logged), "***") { t.Fatalf("the persisted log has no masked entry: %s", logged) }}
func TestForkRunsGetNoIDTokenEvenWhenOneWasMinted(t *testing.T) { eng := newEchoEngine() logPath := runEchoPipelineRecords(t, eng, &models.Pipeline{ RepoDid: "did:plc:boltless", TrustedSource: false, MintIDToken: func(time.Duration, string) (string, error) { return testRunToken, nil }, }, audienceRecord("build", "vault"))
if value := secretValue(eng.runSecrets("build"), models.IDTokenEnvVar); value != "" { t.Fatalf("%s handed to a fork run = %q, want none", models.IDTokenEnvVar, value) } if value := secretValue(eng.runSecrets("build"), "VAULT_SECRET"); value != "" { t.Fatalf("VAULT_SECRET handed to a fork run = %q, want none", value) }
logged, err := os.ReadFile(logPath) if err != nil { t.Fatal(err) } if strings.Contains(string(logged), testRunToken) { t.Fatalf("a fork run leaked the id token: %s", logged) }}
func TestWorkflowProceedsWhenTheIssuerWithholdsTheToken(t *testing.T) { eng := newEchoEngine() // what an identity the issuer refuses to sign looks like by the time a // workflow runs: a minter that yields no token and no error runEchoPipelineRecords(t, eng, &models.Pipeline{ RepoDid: "did:plc:boltless", TrustedSource: true, MintIDToken: func(time.Duration, string) (string, error) { return "", nil }, }, audienceRecord("build", "vault"))
if value := secretValue(eng.runSecrets("build"), models.IDTokenEnvVar); value != "" { t.Fatalf("%s = %q, want none", models.IDTokenEnvVar, value) } if value := secretValue(eng.runSecrets("build"), "VAULT_SECRET"); value != "vault-value" { t.Fatalf("VAULT_SECRET = %q, want the run to keep its other credentials", value) }}
func TestIDTokenIsMintedPerWorkflowOnceItHoldsItsSlot(t *testing.T) { eng := newEchoEngine() eng.gate = make(chan slotWait) eng.timeout = 3 * time.Hour
var mu sync.Mutex var budgets []time.Duration var audiences []string minter := func(budget time.Duration, audience string) (string, error) { mu.Lock() defer mu.Unlock() budgets = append(budgets, budget) audiences = append(audiences, audience) return fmt.Sprintf("%s.%d", testRunToken, len(budgets)), nil }
type result struct { cleanup func() err error } done := make(chan result, 1) go func() { _, cleanup, err := startEchoPipelineRecords(eng, &models.Pipeline{ RepoDid: "did:plc:boltless", TrustedSource: true, MintIDToken: minter, }, audienceRecord("first", "aud-first"), audienceRecord("second", "aud-second")) done <- result{cleanup: cleanup, err: err} }()
mintedCount := func() int { mu.Lock() defer mu.Unlock() return len(budgets) }
// each workflow is released only after the test has seen that no token was // minted for it while it waited waiting := make([]slotWait, 0, 2) for range 2 { select { case w := <-eng.gate: if got := mintedCount(); got != len(waiting) { t.Fatalf("minted %d tokens while %d workflows held a slot, before releasing %q: a token must not spend its lifetime waiting for capacity", got, len(waiting), w.name) } waiting = append(waiting, w) close(w.grant) case <-time.After(20 * time.Second): t.Fatal("a workflow never waited for its slot") } }
tokens := map[string]string{} for _, name := range []string{"first", "second"} { select { case <-eng.ranWorkflow(name): case <-time.After(20 * time.Second): t.Fatalf("%s never ran its steps", name) } if token := secretValue(eng.runSecrets(name), models.IDTokenEnvVar); token != "" { tokens[name] = token } }
res := <-done if res.cleanup != nil { defer res.cleanup() } if res.err != nil { t.Fatal(res.err) }
if got := mintedCount(); got != 2 { t.Fatalf("minted %d tokens for 2 workflows, want one each", got) } for _, budget := range budgets { if budget != 3*time.Hour { t.Fatalf("mint budget = %s, want the workflow time budget 3h", budget) } } if len(tokens) != 2 { t.Fatalf("workflows carried tokens %v, want one freshly minted token each", tokens) } if tokens["first"] == tokens["second"] { t.Fatalf("both workflows carry %q, want a fresh token per workflow", tokens["first"]) } if !slices.Contains(audiences, "aud-first") || !slices.Contains(audiences, "aud-second") { t.Fatalf("minted for audiences %v, want each workflow's own", audiences) }}
func TestWorkflowsWithoutAnAudienceGetNoIDToken(t *testing.T) { eng := newEchoEngine() runEchoPipeline(t, eng, &models.Pipeline{ RepoDid: "did:plc:boltless", TrustedSource: true, MintIDToken: func(time.Duration, string) (string, error) { t.Error("the minter ran for a workflow that declared no audience") return testRunToken, nil }, }, "build")
if value := secretValue(eng.runSecrets("build"), models.IDTokenEnvVar); value != "" { t.Fatalf("%s = %q, want none: the audience declaration is the trigger", models.IDTokenEnvVar, value) } if value := secretValue(eng.runSecrets("build"), "VAULT_SECRET"); value != "vault-value" { t.Fatalf("VAULT_SECRET = %q, want the run to keep its other credentials", value) }}