Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756package engine
import ( "context" "crypto/sha256" "encoding/hex" "errors" "fmt" "io" "log/slog" "os" "sync" "time"
"go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/codes" "go.opentelemetry.io/otel/trace" "strings" "tangled.org/core/notifier" "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" "tangled.org/core/spindle/observability" "tangled.org/core/spindle/quota" "tangled.org/core/spindle/secrets" "tangled.org/core/spindle/storage")
var ( ErrTimedOut = errors.New("timed out") ErrWorkflowFailed = errors.New("workflow failed") ErrWorkflowCanceled = errors.New("workflow canceled"))
type FailureClass string
const ( FailureClassNone FailureClass = "none" FailureClassUser FailureClass = "user" FailureClassInfrastructure FailureClass = "infrastructure" FailureClassPolicy FailureClass = "policy")
type FailureReason string
const ( FailureReasonSuccess FailureReason = "success" FailureReasonCommandFailed FailureReason = "command_failed" FailureReasonResourceExhausted FailureReason = "resource_exhausted" FailureReasonConfigurationFailed FailureReason = "configuration_failed" FailureReasonWorkflowInvalid FailureReason = "workflow_invalid" FailureReasonSetupFailed FailureReason = "setup_failed" FailureReasonRuntimeFailed FailureReason = "runtime_failed" FailureReasonCapacityUnavailable FailureReason = "capacity_unavailable" FailureReasonQuotaDenied FailureReason = "quota_denied" FailureReasonExecutorLost FailureReason = "executor_lost" FailureReasonTimeout FailureReason = "timeout" FailureReasonCancelled FailureReason = "cancelled")
type WorkflowFailure struct { Class FailureClass Reason FailureReason Err error}
func (e *WorkflowFailure) Error() string { if e.Err == nil { return string(e.Reason) } return e.Err.Error()}
func (e *WorkflowFailure) Unwrap() error { return e.Err}
func ClassifiedFailure(class FailureClass, reason FailureReason, err error) error { if err == nil { return nil } var classified *WorkflowFailure if errors.As(err, &classified) { return err } return &WorkflowFailure{Class: class, Reason: reason, Err: err}}
func FailureAttribution(result string, err error) (string, string) { switch result { case "success": return string(FailureClassNone), string(FailureReasonSuccess) case "timeout": return string(FailureClassUser), string(FailureReasonTimeout) case "cancelled": return string(FailureClassUser), string(FailureReasonCancelled) }
var failure *WorkflowFailure if errors.As(err, &failure) { return string(failure.Class), string(failure.Reason) } if errors.Is(err, ErrNoWorkflowSlots) { return string(FailureClassInfrastructure), string(FailureReasonCapacityUnavailable) } return string(FailureClassInfrastructure), string(FailureReasonRuntimeFailed)}
var ( activeMu sync.Mutex activeCancels = make(map[models.WorkflowId]context.CancelCauseFunc))
func CancelWorkflow(wid models.WorkflowId) { activeMu.Lock() cancel, ok := activeCancels[wid] activeMu.Unlock() if ok { cancel(ErrWorkflowCanceled) }}
// user cancel, timeout is DeadlineExceededfunc isCanceled(wfCtx context.Context) bool { return errors.Is(context.Cause(wfCtx), ErrWorkflowCanceled)}
// for when recording early wf cancellationsfunc writeWfError(db *db.DB, n *notifier.Notifier, l *slog.Logger, wfCtx context.Context, wid models.WorkflowId, phase string, err error) { l = l.With("wid", wid, "phase", phase) switch { case isCanceled(wfCtx): l.InfoContext(wfCtx, "workflow canceled") if dbErr := db.StatusCancelled(wid, "User canceled the workflow", -1, n); dbErr != nil { l.ErrorContext(wfCtx, "failed to set workflow status to cancelled", "err", dbErr) } case errors.Is(err, ErrTimedOut) || errors.Is(wfCtx.Err(), context.DeadlineExceeded): l.InfoContext(wfCtx, "workflow timed out") if dbErr := db.StatusTimeout(wid, n); dbErr != nil { l.ErrorContext(wfCtx, "failed to set workflow status to timeout", "err", dbErr) } default: l.ErrorContext(wfCtx, "workflow failed", "err", err) if dbErr := db.StatusFailed(wid, err.Error(), -1, n); dbErr != nil { l.ErrorContext(wfCtx, "failed to set workflow status to failed", "err", dbErr) } }}
// the mill streams the executor's real log file into place itself// a local logger here would only write competing linestype workflowLoggerProvider interface { WorkflowLogger(wid models.WorkflowId) models.WorkflowLogger}
// for engines that manage status updates outside StartWorkflowstype RemoteStatusEngine interface { AuthorsRemoteStatus()}
type WorkflowQuotaReporter interface { QuotaResources(wf *models.Workflow) quota.Resources}
type WorkflowQuotaStoreBinder interface { BindWorkflowQuotaStore(wf *models.Workflow, store quota.ReservationStore) error}type metricEngineNamer interface { MetricEngineName() string}
func engineName(eng models.Engine) string { if eng == nil { return "unknown" } if named, ok := eng.(metricEngineNamer); ok { return boundEngineName(named.MetricEngineName()) } name := strings.TrimPrefix(fmt.Sprintf("%T", eng), "*") if idx := strings.Index(name, "."); idx != -1 { name = name[:idx] } return boundEngineName(name)}
func boundEngineName(name string) string { switch name { case "dummy", "microvm", "nixery": return name default: return "unknown" }}
func workflowResult(ctx context.Context, err error) string { if isCanceled(ctx) { return "cancelled" } if errors.Is(err, ErrTimedOut) || errors.Is(ctx.Err(), context.DeadlineExceeded) { return "timeout" } if err != nil { return "failure" } return "success"}
func reportWorkflowStatusError(l *slog.Logger, database *db.DB, n *notifier.Notifier, wid models.WorkflowId, err error) { if errors.Is(err, ErrTimedOut) { dbErr := database.StatusTimeout(wid, n) if dbErr != nil { l.Error("failed to set workflow status to timeout", "wid", wid, "err", dbErr) } } else if errors.Is(err, ErrWorkflowCanceled) { dbErr := database.StatusCancelled(wid, err.Error(), -1, n) if dbErr != nil { l.Error("failed to set workflow status to cancelled", "wid", wid, "err", dbErr) } } else { dbErr := database.StatusFailed(wid, err.Error(), -1, n) if dbErr != nil { l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) } }}
func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, qm *quota.Manager, stores *artifactstore.Stores, db *db.DB, n *notifier.Notifier, cacheStore storage.Storage, cacheController CacheController, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { l.Info("starting all workflows in parallel", "pipeline", pipelineId)
isTrustedRepo := pipeline.TrustedSource && pipeline.RepoDid != "" var allSecrets []secrets.UnlockedSecret // never pass secrets to pipelines that run untrusted (e.g. fork) code if isTrustedRepo { if res, err := vault.GetSecretsUnlocked(ctx, secrets.RepoIdentifier(pipeline.RepoDid.String())); err == nil { allSecrets = res } } else if !pipeline.TrustedSource { l.Info("skipping secrets for untrusted pipeline source", "pipeline", pipelineId) } if cacheController == nil && cacheStore != nil && db != nil { cacheController = NewLocalCacheController( db, cacheStore, cfg.Server.RepoDir, cfg.Cache.MaxBytesPerOwner, cfg.Cache.MaxEntriesPerOwner, l, ) } cacheEnabled := cacheStore != nil && cacheController != nil && isTrustedRepo if cacheStore != nil && !pipeline.TrustedSource { l.Info("skipping caches for untrusted pipeline source", "pipeline", pipelineId) }
secretValues := make([]string, len(allSecrets)) for i, s := range allSecrets { secretValues[i] = s.Value }
// wid.String() is lossy so two different names can map to the same key // eg. "foo bar" and "foo-bar"... wfCounts := make(map[string]int) for _, wfs := range pipeline.Workflows { for _, w := range wfs { wid := models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, } wfCounts[wid.String()]++ } } var wg sync.WaitGroup for eng, wfs := range pipeline.Workflows { workflowTimeout := eng.WorkflowTimeout() cacheRunner, cachesSupported := eng.(CacheRunner) l.Info("using workflow timeout", "timeout", workflowTimeout)
for _, w := range wfs { w := w wid := models.WorkflowId{ PipelineId: pipelineId, Name: w.Name, }
if wfCounts[wid.String()] > 1 { l.Warn("skipping workflow due to name collision", "wid", wid, "key", wid.String()) dbErr := db.StatusFailed(wid, fmt.Sprintf("colliding workflow name: %s; rename to something else", wid.String()), -1, n) if dbErr != nil { l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) } observability.GetMetrics(ctx).RecordWorkflowTerminal( engineName(eng), "failure", string(FailureClassUser), string(FailureReasonWorkflowInvalid), ) continue }
wg.Go(func() { repoDID := w.RepoDID if repoDID == "" && pipeline != nil { repoDID = pipeline.RepoDid.String() } wl := l.With( "owner_did", w.OwnerDID, "repo_did", repoDID, "pipeline_id", pipelineId.AtUri().String(), "workflow_id", wid.String(), )
if st, err := db.GetStatus(wid); err == nil && models.StatusKind(st.Status).IsFinish() { wl.Info("skipping finished workflow", "wid", wid, "status", st.Status) return }
requestedResources := quota.Resources{} _, hasQuotaReporter := eng.(WorkflowQuotaReporter) if reporter, ok := eng.(WorkflowQuotaReporter); ok { for resource, amount := range reporter.QuotaResources(&w) { if amount != 0 { requestedResources[resource] = amount } } } reqWorkflows := requestedResources[quota.ResourceWorkflows] reqMemoryMiB := requestedResources[quota.ResourceMemoryMiB] reqVCPUs := requestedResources[quota.ResourceVCPUs] reqDiskMiB := requestedResources[quota.ResourceDiskMiB] reqCacheBytes := requestedResources[quota.ResourceCacheStorageBytes]
var err error var wfCtx context.Context = ctx engName := engineName(eng)
var span trace.Span wfCtx, span = observability.Tracer().Start(wfCtx, "workflow.run", trace.WithAttributes( attribute.String(observability.WorkflowEngineKey, engName), )) if span.IsRecording() { attrs := []attribute.KeyValue{ attribute.String(observability.PipelineIDKey, pipelineId.AtUri().String()), attribute.String(observability.WorkflowIDKey, wid.String()), } if w.OwnerDID != "" { attrs = append(attrs, attribute.String(observability.OwnerDIDKey, w.OwnerDID)) } if repoDID != "" { attrs = append(attrs, attribute.String(observability.RepoDIDKey, repoDID)) } attrs = append(attrs, observability.RequestedResourceAttrs(reqWorkflows, reqVCPUs, reqMemoryMiB, reqDiskMiB, reqCacheBytes)...) span.SetAttributes(attrs...) } defer span.End()
_, remoteStatus := eng.(RemoteStatusEngine) metrics := observability.GetMetrics(ctx) if !remoteStatus { metrics.RecordWorkflowStart(engName) } startTime := time.Now() result := "success" remoteExecutionStarted := false
defer func() { if err != nil { span.SetStatus(codes.Error, "workflow failed") } else if isCanceled(wfCtx) { span.SetStatus(codes.Error, "workflow cancelled") } else { span.SetStatus(codes.Ok, "success") }
duration := time.Since(startTime) resStr := result if err != nil { resStr = workflowResult(wfCtx, err) } else if isCanceled(wfCtx) { resStr = "cancelled" }
var cpuUsec, memoryCurrent, memoryPeak, swapCurrent, swapPeak, pidsCurrent, ioReadBytes, ioWriteBytes, ioReadOps, ioWriteOps, volumeAllocated uint64 var cgroupAvailable, volumeAvailable bool if reporter, ok := eng.(WorkflowResourceUsageReporter); ok { if usage, ok := reporter.WorkflowResourceUsage(&w); ok { cpuUsec = usage.CPUUsec memoryCurrent = usage.MemoryCurrentBytes memoryPeak = usage.MemoryPeakBytes swapCurrent = usage.SwapCurrentBytes swapPeak = usage.SwapPeakBytes pidsCurrent = usage.PIDsCurrent ioReadBytes = usage.IOReadBytes ioWriteBytes = usage.IOWriteBytes ioReadOps = usage.IOReadOps ioWriteOps = usage.IOWriteOps volumeAllocated = usage.VolumeAllocatedBytes cgroupAvailable = usage.CgroupAvailable volumeAvailable = usage.VolumeAvailable } }
wl.InfoContext(wfCtx, "workflow finished", "result", resStr, "duration_seconds", duration.Seconds(), "requested_workflows", reqWorkflows, "requested_vcpus", reqVCPUs, "requested_memory_mib", reqMemoryMiB, "requested_disk_mib", reqDiskMiB, "requested_cache_bytes", reqCacheBytes, "actual_cpu_usec", cpuUsec, "actual_memory_current_bytes", memoryCurrent, "actual_memory_peak_bytes", memoryPeak, "actual_swap_current_bytes", swapCurrent, "actual_swap_peak_bytes", swapPeak, "actual_pids_current", pidsCurrent, "actual_io_read_bytes", ioReadBytes, "actual_io_write_bytes", ioWriteBytes, "actual_io_read_ops", ioReadOps, "actual_io_write_ops", ioWriteOps, "actual_volume_allocated_bytes", volumeAllocated, "actual_cgroup_available", cgroupAvailable, "actual_volume_available", volumeAvailable, )
if span.IsRecording() { span.SetAttributes( observability.ActualResourceAttrs( cpuUsec, memoryCurrent, memoryPeak, swapCurrent, swapPeak, pidsCurrent, ioReadBytes, ioWriteBytes, ioReadOps, ioWriteOps, volumeAllocated, cgroupAvailable, volumeAvailable, )..., ) }
failureClass, failureReason := FailureAttribution(resStr, err) if remoteStatus { if !remoteExecutionStarted { metrics.RecordWorkflowTerminal(engName, resStr, failureClass, failureReason) } return }
if observability.NotifyWorkflowTerminal(wfCtx, engName, resStr, failureClass, failureReason) { metrics.RecordWorkflowExecutionEnd(wfCtx, engName, resStr, duration) } else { metrics.RecordWorkflowEnd(wfCtx, engName, resStr, failureClass, failureReason, duration) } }()
var wfLogger models.WorkflowLogger closeLog := func() {} if p, ok := eng.(workflowLoggerProvider); ok { wfLogger = p.WorkflowLogger(wid) } else if fileLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues); err != nil { wl.WarnContext(wfCtx, "failed to setup step logger; logs will not be persisted", "error", err) wfLogger = models.NullLogger{} } else { wl.InfoContext(wfCtx, "setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) wfLogger = fileLogger var closeOnce sync.Once closeLog = func() { closeOnce.Do(func() { if err := fileLogger.Close(); err != nil { wl.ErrorContext(wfCtx, "failed to close workflow log", "wid", wid, "err", err) } }) } defer archiveWorkflowLog(wl, stores, db, cfg.Server.LogDir, wid, repoDID) defer closeLog() }
timeoutCtx, timeoutCancel := context.WithTimeout(wfCtx, workflowTimeout) defer timeoutCancel()
var userCancel context.CancelCauseFunc wfCtx, userCancel = context.WithCancelCause(timeoutCtx) defer userCancel(nil)
// allow wf context to be cancelled properly by manual cancel activeMu.Lock() activeCancels[wid] = userCancel activeMu.Unlock() defer func() { activeMu.Lock() delete(activeCancels, wid) activeMu.Unlock() }()
slot := WorkflowSlot(NoopSlot{}) var publishTerminalStatus func() var quotaLease quota.Lease destroyWorkflow := false slotAcquired := false setTerminalError := func(phase string, workflowErr error) { publishTerminalStatus = func() { writeWfError(db, n, wl, wfCtx, wid, phase, workflowErr) } } defer func() { closeLog() if publishTerminalStatus != nil { publishTerminalStatus() } if destroyWorkflow { if err := eng.DestroyWorkflow(ctx, wid); err != nil { wl.ErrorContext(wfCtx, "failed to destroy workflow", "wid", wid, "err", err) } } if slotAcquired { slot.Release() } if quotaLease != nil { quotaLease.Release() } }()
if qm != nil && hasQuotaReporter { resID := quota.WorkflowReservationID(w.RunID, w.OwnerDID, w.RepoDID, wid.Knot, wid.Rkey, wid.Name) req := quota.ReserveRequest{ ID: resID, Kind: quota.KindWorkflow, Key: resID, Identity: quota.Identity{ OwnerDID: w.OwnerDID, RepoDID: w.RepoDID, }, Resources: requestedResources, } quotaLease, err = qm.Acquire(wfCtx, req) if err != nil { if errors.Is(err, quota.ErrDenied) { err = ClassifiedFailure(FailureClassPolicy, FailureReasonQuotaDenied, err) } setTerminalError("acquiring quota", err) return } }
wl.InfoContext(wfCtx, "waiting for slot", "wid", wid)
if s, ok := eng.(WorkflowSlotter); ok { slot, err = s.AcquireWorkflowSlot(wfCtx, wid, &w, Wait) if err != nil { setTerminalError("waiting for slot", err) return } slotAcquired = true remoteExecutionStarted = true }
if !remoteStatus { err = db.StatusRunning(wid, n) if err != nil { wl.ErrorContext(wfCtx, "failed to set workflow status to running", "wid", wid, "err", err) return } if delay, ok, delayErr := db.WorkflowStartupDelay(wfCtx, wid); delayErr != nil { wl.WarnContext(wfCtx, "failed to measure workflow startup delay", "err", delayErr) } else if ok { metrics.RecordWorkflowStartupDelay(wfCtx, engName, "local", delay) } }
wl.InfoContext(wfCtx, "workflow started", "engine", engName, "requested_workflows", reqWorkflows, "requested_vcpus", reqVCPUs, "requested_memory_mib", reqMemoryMiB, "requested_disk_mib", reqDiskMiB, "requested_cache_bytes", reqCacheBytes, )
err = eng.SetupWorkflow(wfCtx, wid, &w, wfLogger) if err != nil { err = ClassifiedFailure(FailureClassInfrastructure, FailureReasonSetupFailed, err) destroyWorkflow = !isCanceled(wfCtx) if !remoteStatus { setTerminalError("setting up workflow", err) } return } destroyWorkflow = true
var bindings []models.CacheBinding var trackedStore *trackedCacheStore if cacheEnabled && len(w.Caches) > 0 { bindings, err = cacheController.Plan(wfCtx, pipeline, &w) if err != nil { l.Warn("cache planning failed", "wid", wid, "err", err) } else { w.CacheBindings = bindings trackedStore = newTrackedCacheStore(cacheStore, cacheController, l, bindings) if cachesSupported { wfLogger.ControlWriter(CacheRestoreStepIdx, CacheRestoreStep, models.StepStatusStart).Write([]byte{0}) // caches are an optimization, never a reason to fail the workflow if err := cacheRunner.RestoreCache(wfCtx, wid, &w, trackedStore, bindings, wfLogger); err != nil { l.Warn("cache restore failed", "wid", wid, "err", err) } wfLogger.ControlWriter(CacheRestoreStepIdx, CacheRestoreStep, models.StepStatusEnd).Write([]byte{0}) } else if !remoteStatus { l.Warn("engine does not support caches, skipping restore", "wid", wid) } } }
cleanupCaches := func() { if trackedStore != nil { trackedStore.cleanup(context.WithoutCancel(wfCtx)) trackedStore = nil } } // dont save on timeouts, their context is already dead saveCaches := func(failed bool) { if trackedStore == nil || !cachesSupported { return } toSave := make([]models.CacheBinding, 0, len(bindings)) for _, binding := range bindings { if binding.SaveKey != "" && binding.SaveOn(failed) { toSave = append(toSave, binding) } } if len(toSave) == 0 { return } wfLogger.ControlWriter(CacheSaveStepIdx, CacheSaveStep, models.StepStatusStart).Write([]byte{0}) if err := cacheRunner.SaveCache(wfCtx, wid, &w, trackedStore, toSave, wfLogger); err != nil { l.Warn("cache save failed", "wid", wid, "err", err) } wfLogger.ControlWriter(CacheSaveStepIdx, CacheSaveStep, models.StepStatusEnd).Write([]byte{0}) } for stepIdx, step := range w.Steps { if wfLogger != nil { wfLogger. ControlWriter(stepIdx, step, models.StepStatusStart). Write([]byte{0}) }
if !remoteStatus { metrics.RecordStepStart(engName) } stepStart := time.Now()
var stepSpan trace.Span var stepCtx context.Context stepCtx, stepSpan = observability.Tracer().Start(wfCtx, "step.run") if stepSpan.IsRecording() { attrs := []attribute.KeyValue{ attribute.String(observability.WorkflowEngineKey, engName), attribute.String(observability.PipelineIDKey, pipelineId.AtUri().String()), attribute.String(observability.WorkflowIDKey, wid.String()), attribute.Int(observability.StepIndexKey, stepIdx), attribute.String(observability.StepNameKey, step.Name()), } if w.OwnerDID != "" { attrs = append(attrs, attribute.String(observability.OwnerDIDKey, w.OwnerDID)) } if repoDID != "" { attrs = append(attrs, attribute.String(observability.RepoDIDKey, repoDID)) } stepSpan.SetAttributes(attrs...) }
err = eng.RunStep(stepCtx, wid, &w, stepIdx, allSecrets, wfLogger)
if wfLogger != nil { wfLogger. ControlWriter(stepIdx, step, models.StepStatusEnd). Write([]byte{0}) }
stepResult := "success" if err != nil { stepResult = workflowResult(wfCtx, err) stepSpan.SetStatus(codes.Error, "step failed") } else { stepSpan.SetStatus(codes.Ok, "success") } stepSpan.End()
if !remoteStatus { metrics.RecordStepEnd(stepCtx, engName, stepResult, time.Since(stepStart)) }
if err != nil { if !errors.Is(err, ErrTimedOut) && !errors.Is(wfCtx.Err(), context.DeadlineExceeded) && !isCanceled(wfCtx) { saveCaches(true) } cleanupCaches() if !remoteStatus { setTerminalError("running step", err) } return } }
saveCaches(false) cleanupCaches()
if isCanceled(wfCtx) { if !remoteStatus { setTerminalError("before success", nil) } return }
if !remoteStatus { publishTerminalStatus = func() { if err := db.StatusSuccess(wid, n); err != nil { wl.ErrorContext(wfCtx, "failed to set workflow status to success", "wid", wid, "err", err) } } } }) } }
wg.Wait() l.Info("all workflows completed")}
func archiveWorkflowLog(l *slog.Logger, stores *artifactstore.Stores, database *db.DB, logDir string, wid models.WorkflowId, repoDID string) { if stores == nil { return } logPath := models.LogFilePath(logDir, wid) file, err := os.Open(logPath) if err != nil { l.Error("open workflow log for archival", "wid", wid, "err", err) return } hash := sha256.New() if _, err := io.Copy(hash, file); err != nil { _ = file.Close() l.Error("hash workflow log", "wid", wid, "err", err) return } _ = file.Close()
ref := wid.String() + ".log" uploadCtx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) defer cancel() errs := stores.PutFile(uploadCtx, ref, logPath) for _, err := range errs { l.Error("archive workflow log", "wid", wid, "err", err) } if len(errs) == len(stores.Names()) { return } digest := "sha256:" + hex.EncodeToString(hash.Sum(nil)) if err := database.SaveArtifactRef(wid.String(), repoDID, wid, ref, digest); err != nil { l.Error("save workflow log artifact", "wid", wid, "err", err) }}