This repository has no description
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}