forked from
tangled.org/core
Monorepo for Tangled
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}