Monorepo for Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

core / spindle / engine / engine.go
5.2 kB 177 lines
1package engine 2 3import ( 4 "context" 5 "errors" 6 "fmt" 7 "log/slog" 8 "path/filepath" 9 "sync" 10 11 "tangled.org/core/notifier" 12 "tangled.org/core/spindle/config" 13 "tangled.org/core/spindle/db" 14 "tangled.org/core/spindle/models" 15 "tangled.org/core/spindle/secrets" 16) 17 18var ( 19 ErrTimedOut = errors.New("timed out") 20 ErrWorkflowFailed = errors.New("workflow failed") 21) 22 23type workflowFinalizer interface { 24 FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) error 25} 26 27func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { 28 l.Info("starting all workflows in parallel", "pipeline", pipelineId) 29 30 var allSecrets []secrets.UnlockedSecret 31 // never pass secrets to pipelines that run untrusted (e.g. fork) code 32 if pipeline.TrustedSource && pipeline.RepoDid != "" { 33 if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(pipeline.RepoDid.String())); err == nil { 34 allSecrets = res 35 } 36 } else if !pipeline.TrustedSource { 37 l.Info("skipping secrets for untrusted pipeline source", "pipeline", pipelineId) 38 } 39 40 secretValues := make([]string, len(allSecrets)) 41 for i, s := range allSecrets { 42 secretValues[i] = s.Value 43 } 44 45 s3, err := NewS3(cfg.S3.LogBucket) 46 if err != nil { 47 l.Error("error creating s3 client", "err", err) 48 } 49 50 var wg sync.WaitGroup 51 for eng, wfs := range pipeline.Workflows { 52 workflowTimeout := eng.WorkflowTimeout() 53 l.Info("using workflow timeout", "timeout", workflowTimeout) 54 55 for _, w := range wfs { 56 wg.Go(func() { 57 wid := models.WorkflowId{ 58 PipelineId: pipelineId, 59 Name: w.Name, 60 } 61 62 defer func() { 63 if s3 != nil { 64 logFile := filepath.Join(cfg.Server.LogDir, fmt.Sprintf("%s.log", wid.String())) 65 if err := s3.WriteFile(ctx, logFile); err != nil { 66 l.Error("error uploading logs", "err", err) 67 } 68 } 69 }() 70 71 wfLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues) 72 if err != nil { 73 l.Warn("failed to setup step logger; logs will not be persisted", "error", err) 74 wfLogger = models.NullLogger{} 75 } else { 76 l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) 77 defer wfLogger.Close() 78 } 79 80 l.Info("waiting for slot", "wid", wid) 81 slot := WorkflowSlot(NoopSlot{}) 82 if s, ok := eng.(WorkflowSlotter); ok { 83 var err error 84 slot, err = s.AcquireWorkflowSlot(ctx, wid, &w) 85 if err != nil { 86 l.Error("failed to acquire slot", "wid", wid, "err", err) 87 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 88 if dbErr != nil { 89 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 90 } 91 return 92 } 93 } 94 defer slot.Release() 95 96 err = db.StatusRunning(wid, n) 97 if err != nil { 98 l.Error("failed to set workflow status to running", "wid", wid, "err", err) 99 return 100 } 101 102 err = eng.SetupWorkflow(ctx, wid, &w, wfLogger) 103 if err != nil { 104 // TODO(winter): Should this always set StatusFailed? 105 // In the original, we only do in a subset of cases. 106 l.Error("setting up workflow", "wid", wid, "err", err) 107 108 destroyErr := eng.DestroyWorkflow(ctx, wid) 109 if destroyErr != nil { 110 l.Error("failed to destroy workflow after setup failure", "error", destroyErr) 111 } 112 113 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 114 if dbErr != nil { 115 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 116 } 117 return 118 } 119 defer eng.DestroyWorkflow(ctx, wid) 120 121 ctx, cancel := context.WithTimeout(ctx, workflowTimeout) 122 defer cancel() 123 124 for stepIdx, step := range w.Steps { 125 // log start of step 126 if wfLogger != nil { 127 wfLogger. 128 ControlWriter(stepIdx, step, models.StepStatusStart). 129 Write([]byte{0}) 130 } 131 132 err = eng.RunStep(ctx, wid, &w, stepIdx, allSecrets, wfLogger) 133 134 // log end of step 135 if wfLogger != nil { 136 wfLogger. 137 ControlWriter(stepIdx, step, models.StepStatusEnd). 138 Write([]byte{0}) 139 } 140 141 if err != nil { 142 if errors.Is(err, ErrTimedOut) { 143 dbErr := db.StatusTimeout(wid, n) 144 if dbErr != nil { 145 l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) 146 } 147 } else { 148 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 149 if dbErr != nil { 150 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 151 } 152 } 153 return 154 } 155 } 156 157 if finalizer, ok := eng.(workflowFinalizer); ok { 158 if err := finalizer.FinalizeWorkflow(ctx, wid, &w, wfLogger); err != nil { 159 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 160 if dbErr != nil { 161 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 162 } 163 return 164 } 165 } 166 167 err = db.StatusSuccess(wid, n) 168 if err != nil { 169 l.Error("failed to set workflow status to success", "wid", wid, "err", err) 170 } 171 }) 172 } 173 } 174 175 wg.Wait() 176 l.Info("all workflows completed") 177}