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.4 kB 160 lines
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}