Something went wrong. Try again.
Monorepo for Tangled
Something went wrong. Try again.
Go
at sl/md-editor
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153package 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")
var ( ErrTimedOut = errors.New("timed out") ErrWorkflowFailed = errors.New("workflow failed"))
func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, workflowSem chan struct{}, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { l.Info("starting all workflows in parallel", "pipeline", pipelineId)
// extract secrets var allSecrets []secrets.UnlockedSecret if pipeline.RepoDid != "" { if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(pipeline.RepoDid.String())); err == nil { allSecrets = res } }
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() l.Info("using workflow timeout", "timeout", workflowTimeout)
for _, w := range wfs { wg.Add(1) go func() { defer wg.Done()
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() }
err = db.StatusRunning(wid, n) if err != nil { l.Error("failed to set workflow status to running", "wid", wid, "err", err) return }
// acquire semaphore slot before starting the container workflowSem <- struct{}{} defer func() { <-workflowSem }()
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()
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 { 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")}