This repository has no description
0

Configure Feed

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

core / spindle / engine / engine.go
4.9 kB 166 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 23func 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) { 24 l.Info("starting all workflows in parallel", "pipeline", pipelineId) 25 26 var allSecrets []secrets.UnlockedSecret 27 // never pass secrets to pipelines that run untrusted (e.g. fork) code 28 if pipeline.TrustedSource && pipeline.RepoDid != "" { 29 if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(pipeline.RepoDid.String())); err == nil { 30 allSecrets = res 31 } 32 } else if !pipeline.TrustedSource { 33 l.Info("skipping secrets for untrusted pipeline source", "pipeline", pipelineId) 34 } 35 36 secretValues := make([]string, len(allSecrets)) 37 for i, s := range allSecrets { 38 secretValues[i] = s.Value 39 } 40 41 s3, err := NewS3(cfg.S3.LogBucket) 42 if err != nil { 43 l.Error("error creating s3 client", "err", err) 44 } 45 46 var wg sync.WaitGroup 47 for eng, wfs := range pipeline.Workflows { 48 workflowTimeout := eng.WorkflowTimeout() 49 l.Info("using workflow timeout", "timeout", workflowTimeout) 50 51 for _, w := range wfs { 52 wg.Go(func() { 53 wid := models.WorkflowId{ 54 PipelineId: pipelineId, 55 Name: w.Name, 56 } 57 58 defer func() { 59 if s3 != nil { 60 logFile := filepath.Join(cfg.Server.LogDir, fmt.Sprintf("%s.log", wid.String())) 61 if err := s3.WriteFile(ctx, logFile); err != nil { 62 l.Error("error uploading logs", "err", err) 63 } 64 } 65 }() 66 67 wfLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues) 68 if err != nil { 69 l.Warn("failed to setup step logger; logs will not be persisted", "error", err) 70 wfLogger = models.NullLogger{} 71 } else { 72 l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) 73 defer wfLogger.Close() 74 } 75 76 l.Info("waiting for slot", "wid", wid) 77 slot := WorkflowSlot(NoopSlot{}) 78 if s, ok := eng.(WorkflowSlotter); ok { 79 var err error 80 slot, err = s.AcquireWorkflowSlot(ctx, wid, &w) 81 if err != nil { 82 l.Error("failed to acquire slot", "wid", wid, "err", err) 83 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 84 if dbErr != nil { 85 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 86 } 87 return 88 } 89 } 90 defer slot.Release() 91 92 err = db.StatusRunning(wid, n) 93 if err != nil { 94 l.Error("failed to set workflow status to running", "wid", wid, "err", err) 95 return 96 } 97 98 err = eng.SetupWorkflow(ctx, wid, &w, wfLogger) 99 if err != nil { 100 // TODO(winter): Should this always set StatusFailed? 101 // In the original, we only do in a subset of cases. 102 l.Error("setting up workflow", "wid", wid, "err", err) 103 104 destroyErr := eng.DestroyWorkflow(ctx, wid) 105 if destroyErr != nil { 106 l.Error("failed to destroy workflow after setup failure", "error", destroyErr) 107 } 108 109 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 110 if dbErr != nil { 111 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 112 } 113 return 114 } 115 // don't put this after the workflowTimeout deadline assignment 116 // below. engines that implement "ssh-after-fail" rely on the 117 // unbounded ctx for retaining the workflow after it fails. 118 defer eng.DestroyWorkflow(ctx, wid) 119 120 ctx, cancel := context.WithTimeout(ctx, workflowTimeout) 121 defer cancel() 122 123 for stepIdx, step := range w.Steps { 124 // log start of step 125 if wfLogger != nil { 126 wfLogger. 127 ControlWriter(stepIdx, step, models.StepStatusStart). 128 Write([]byte{0}) 129 } 130 131 err = eng.RunStep(ctx, wid, &w, stepIdx, allSecrets, wfLogger) 132 133 // log end of step 134 if wfLogger != nil { 135 wfLogger. 136 ControlWriter(stepIdx, step, models.StepStatusEnd). 137 Write([]byte{0}) 138 } 139 140 if err != nil { 141 if errors.Is(err, ErrTimedOut) { 142 dbErr := db.StatusTimeout(wid, n) 143 if dbErr != nil { 144 l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) 145 } 146 } else { 147 dbErr := db.StatusFailed(wid, err.Error(), -1, n) 148 if dbErr != nil { 149 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 150 } 151 } 152 return 153 } 154 } 155 156 err = db.StatusSuccess(wid, n) 157 if err != nil { 158 l.Error("failed to set workflow status to success", "wid", wid, "err", err) 159 } 160 }) 161 } 162 } 163 164 wg.Wait() 165 l.Info("all workflows completed") 166}