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 securejoin "github.com/cyphar/filepath-securejoin"
12 "tangled.org/core/notifier"
13 "tangled.org/core/spindle/config"
14 "tangled.org/core/spindle/db"
15 "tangled.org/core/spindle/models"
16 "tangled.org/core/spindle/secrets"
17)
18
19var (
20 ErrTimedOut = errors.New("timed out")
21 ErrWorkflowFailed = errors.New("workflow failed")
22)
23
24func 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) {
25 l.Info("starting all workflows in parallel", "pipeline", pipelineId)
26
27 // extract secrets
28 var allSecrets []secrets.UnlockedSecret
29 if didSlashRepo, err := securejoin.SecureJoin(pipeline.RepoOwner, pipeline.RepoName); err == nil {
30 if res, err := vault.GetSecretsUnlocked(ctx, secrets.DidSlashRepo(didSlashRepo)); err == nil {
31 allSecrets = res
32 }
33 }
34
35 secretValues := make([]string, len(allSecrets))
36 for i, s := range allSecrets {
37 secretValues[i] = s.Value
38 }
39
40 var wg sync.WaitGroup
41 for eng, wfs := range pipeline.Workflows {
42 workflowTimeout := eng.WorkflowTimeout()
43 l.Info("using workflow timeout", "timeout", workflowTimeout)
44
45 for _, w := range wfs {
46 wg.Add(1)
47 go func() {
48 defer wg.Done()
49
50 wid := models.WorkflowId{
51 PipelineId: pipelineId,
52 Name: w.Name,
53 }
54
55 defer func() {
56 logBucket := cfg.S3.LogBucket
57
58 if logBucket != "" {
59 logFile := filepath.Join(cfg.Server.LogDir, fmt.Sprintf("%s.log", wid.String()))
60 if err := uploadWorkflowLogs(ctx, logFile, "tangled-demo"); err != nil {
61 l.Error("error uploading logs", "err", err)
62 }
63 }
64 }()
65
66 wfLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues)
67 if err != nil {
68 l.Warn("failed to setup step logger; logs will not be persisted", "error", err)
69 wfLogger = models.NullLogger{}
70 } else {
71 l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid)
72 defer wfLogger.Close()
73 }
74
75 err = db.StatusRunning(wid, n)
76 if err != nil {
77 l.Error("failed to set workflow status to running", "wid", wid, "err", err)
78 return
79 }
80
81 err = eng.SetupWorkflow(ctx, wid, &w, wfLogger)
82 if err != nil {
83 // TODO(winter): Should this always set StatusFailed?
84 // In the original, we only do in a subset of cases.
85 l.Error("setting up worklow", "wid", wid, "err", err)
86
87 destroyErr := eng.DestroyWorkflow(ctx, wid)
88 if destroyErr != nil {
89 l.Error("failed to destroy workflow after setup failure", "error", destroyErr)
90 }
91
92 dbErr := db.StatusFailed(wid, err.Error(), -1, n)
93 if dbErr != nil {
94 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr)
95 }
96 return
97 }
98 defer eng.DestroyWorkflow(ctx, wid)
99
100 ctx, cancel := context.WithTimeout(ctx, workflowTimeout)
101 defer cancel()
102
103 for stepIdx, step := range w.Steps {
104 // log start of step
105 if wfLogger != nil {
106 wfLogger.
107 ControlWriter(stepIdx, step, models.StepStatusStart).
108 Write([]byte{0})
109 }
110
111 err = eng.RunStep(ctx, wid, &w, stepIdx, allSecrets, wfLogger)
112
113 // log end of step
114 if wfLogger != nil {
115 wfLogger.
116 ControlWriter(stepIdx, step, models.StepStatusEnd).
117 Write([]byte{0})
118 }
119
120 if err != nil {
121 if errors.Is(err, ErrTimedOut) {
122 dbErr := db.StatusTimeout(wid, n)
123 if dbErr != nil {
124 l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr)
125 }
126 } else {
127 dbErr := db.StatusFailed(wid, err.Error(), -1, n)
128 if dbErr != nil {
129 l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr)
130 }
131 }
132 return
133 }
134 }
135
136 err = db.StatusSuccess(wid, n)
137 if err != nil {
138 l.Error("failed to set workflow status to success", "wid", wid, "err", err)
139 }
140 }()
141 }
142 }
143
144 wg.Wait()
145 l.Info("all workflows completed")
146}
147
148func uploadWorkflowLogs(ctx context.Context, logfile, bucket string) error {
149 s3, err := NewS3(bucket)
150 if err != nil {
151 return fmt.Errorf("error creating s3 client: %w", err)
152 }
153
154 name := filepath.Join(logfile)
155 if err := s3.WriteFile(ctx, name); err != nil {
156 return fmt.Errorf("error saving logs: %w", err)
157 }
158
159 return nil
160}