Something went wrong. Try again.
Monorepo for Tangled
Something went wrong. Try again.
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252package engine
import ( "context" "errors" "fmt" "log/slog" "path/filepath" "sync"
"tangled.org/core/notifier" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/storage")
var ( ErrTimedOut = errors.New("timed out") ErrWorkflowFailed = errors.New("workflow failed"))
type workflowFinalizer interface { FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) error}
func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, cacheStore storage.Storage, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { l.Info("starting all workflows in parallel", "pipeline", pipelineId)
isTrustedRepo := pipeline.TrustedSource && pipeline.RepoDid != "" var allSecrets []secrets.UnlockedSecret // never pass secrets to pipelines that run untrusted (e.g. fork) code if isTrustedRepo { if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(pipeline.RepoDid.String())); err == nil { allSecrets = res } } else if !pipeline.TrustedSource { l.Info("skipping secrets for untrusted pipeline source", "pipeline", pipelineId) } // untrusted runs cant read or write shared caches cacheOwnerDID := "" cacheEnabled := cacheStore != nil && isTrustedRepo if cacheEnabled { repo, err := db.GetRepoByDid(pipeline.RepoDid) if err != nil { l.Warn("cache owner lookup failed; caching disabled", "repo", pipeline.RepoDid, "err", err) cacheEnabled = false } else { cacheOwnerDID = repo.Owner.String() } } else if cacheStore != nil && !pipeline.TrustedSource { l.Info("skipping caches for untrusted pipeline source", "pipeline", pipelineId) }
// hash the checked commit, not the working tree cacheRepoPath, cacheRev := "", "" if tm := pipeline.TriggerMetadata; tm != nil { if rev, err := models.ExtractCommitSHA(*tm); err == nil { did := pipeline.RepoDid.String() if tm.SourceRepo != nil && *tm.SourceRepo != "" { did = *tm.SourceRepo } cacheRepoPath, cacheRev = filepath.Join(cfg.Server.RepoDir, did), rev } else { l.Warn("cannot resolve pipeline commit; cache hashing disabled", "err", err) } }
secretValues := make([]string, len(allSecrets)) for i, s := range allSecrets { secretValues[i] = s.Value }
s3, err := NewS3(cfg.S3.LogBucket) if err != nil { l.Error("error creating s3 client", "err", err) }
var wg sync.WaitGroup for eng, wfs := range pipeline.Workflows { workflowTimeout := eng.WorkflowTimeout() cacheRunner, cachesSupported := eng.(CacheRunner) l.Info("using workflow timeout", "timeout", workflowTimeout)
for _, w := range wfs { wg.Go(func() { wid := models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }
defer func() { if s3 != nil { logFile := filepath.Join(cfg.Server.LogDir, fmt.Sprintf("%s.log", wid.String())) if err := s3.WriteFile(ctx, logFile); err != nil { l.Error("error uploading logs", "err", err) } } }()
wfLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues) if err != nil { l.Warn("failed to setup step logger; logs will not be persisted", "error", err) wfLogger = models.NullLogger{} } else { l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) defer wfLogger.Close() }
l.Info("waiting for slot", "wid", wid) slot := WorkflowSlot(NoopSlot{}) if s, ok := eng.(WorkflowSlotter); ok { var err error slot, err = s.AcquireWorkflowSlot(ctx, wid, &w) if err != nil { l.Error("failed to acquire slot", "wid", wid, "err", err) dbErr := db.StatusFailed(wid, err.Error(), -1, n) if dbErr != nil { l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) } return } } defer slot.Release()
err = db.StatusRunning(wid, n) if err != nil { l.Error("failed to set workflow status to running", "wid", wid, "err", err) return }
err = eng.SetupWorkflow(ctx, wid, &w, wfLogger) if err != nil { // TODO(winter): Should this always set StatusFailed? // In the original, we only do in a subset of cases. l.Error("setting up workflow", "wid", wid, "err", err)
destroyErr := eng.DestroyWorkflow(ctx, wid) if destroyErr != nil { l.Error("failed to destroy workflow after setup failure", "error", destroyErr) }
dbErr := db.StatusFailed(wid, err.Error(), -1, n) if dbErr != nil { l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) } return } defer eng.DestroyWorkflow(ctx, wid)
ctx, cancel := context.WithTimeout(ctx, workflowTimeout) defer cancel()
var resolvedCaches []ResolvedCache if cacheEnabled && len(w.Caches) > 0 { if cachesSupported { resolvedCaches = ResolveCaches(ctx, l, db, pipeline.RepoDid.String(), w.Engine, cacheRepoPath, cacheRev, w.Caches) wfLogger.ControlWriter(CacheRestoreStepIdx, CacheRestoreStep, models.StepStatusStart).Write([]byte{0}) // caches are an optimization, never a reason to fail the workflow restoreStore := cacheStoreForRestore(cacheStore, db, l, resolvedCaches) if err := cacheRunner.RestoreCache(ctx, wid, &w, restoreStore, resolvedCaches, wfLogger); err != nil { l.Warn("cache restore failed", "wid", wid, "err", err) } wfLogger.ControlWriter(CacheRestoreStepIdx, CacheRestoreStep, models.StepStatusEnd).Write([]byte{0}) } else { l.Warn("engine does not support caches, skipping restore", "wid", wid) } }
// dont save on timeouts, their context is already dead saveCaches := func(failed bool) { toSave := resolvedCaches[:0] for _, rc := range resolvedCaches { if rc.saveOn(failed) { toSave = append(toSave, rc) } } if len(toSave) == 0 { return } wfLogger.ControlWriter(CacheSaveStepIdx, CacheSaveStep, models.StepStatusStart).Write([]byte{0}) saveStore, err := prepareCacheSaves(ctx, cacheStore, db, l, cacheOwnerDID, pipeline.RepoDid.String(), w.Engine, toSave) if err != nil { l.Warn("cache metadata setup failed", "wid", wid, "err", err) } else { if err := cacheRunner.SaveCache(ctx, wid, &w, saveStore, toSave, wfLogger); err != nil { l.Warn("cache save failed", "wid", wid, "err", err) } saveStore.cleanup(context.WithoutCancel(ctx)) } wfLogger.ControlWriter(CacheSaveStepIdx, CacheSaveStep, models.StepStatusEnd).Write([]byte{0}) }
for stepIdx, step := range w.Steps { // log start of step if wfLogger != nil { wfLogger. ControlWriter(stepIdx, step, models.StepStatusStart). Write([]byte{0}) }
err = eng.RunStep(ctx, wid, &w, stepIdx, allSecrets, wfLogger)
// log end of step if wfLogger != nil { wfLogger. ControlWriter(stepIdx, step, models.StepStatusEnd). Write([]byte{0}) }
if err != nil { if errors.Is(err, ErrTimedOut) { dbErr := db.StatusTimeout(wid, n) if dbErr != nil { l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) } } else { saveCaches(true) dbErr := db.StatusFailed(wid, err.Error(), -1, n) if dbErr != nil { l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) } } return } }
saveCaches(false)
if finalizer, ok := eng.(workflowFinalizer); ok { if err := finalizer.FinalizeWorkflow(ctx, wid, &w, wfLogger); err != nil { dbErr := db.StatusFailed(wid, err.Error(), -1, n) if dbErr != nil { l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) } return } }
err = db.StatusSuccess(wid, n) if err != nil { l.Error("failed to set workflow status to success", "wid", wid, "err", err) } }) } }
wg.Wait() l.Info("all workflows completed")}