Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
77 changes: 62 additions & 15 deletions internal/run/source.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,9 +60,16 @@ type SourceResult struct {
MergeResult MergeResult
CLIVersion string
// The path to the output OAS spec
OutputPath string
oldSpecPath string
newSpecPath string
OutputPath string
// documentPath is the fully resolved document the source pipeline
// produced, i.e. the path runSourceInner returns to the caller. It is only
// set once the source ran to completion and is what a cached RunSource hit
// hands back. It can differ from OutputPath: with --frozen-workflow-lockfile
// the output location is never written, so OutputPath may point at a
// temp file that does not exist.
documentPath string
oldSpecPath string
newSpecPath string
}

type LintingError struct {
Expand All @@ -80,20 +87,22 @@ func (e *LintingError) Error() string {

func (w *Workflow) RunSource(ctx context.Context, parentStep *workflowTracking.WorkflowStep, sourceID, targetID, targetLanguage string) (string, *SourceResult, error) {
// Fast path: return cached result if this source was already run
if cached, ok := func() (*SourceResult, bool) {
w.sourceMu.Lock()
defer w.sourceMu.Unlock()
if c, ok := w.SourceResults[sourceID]; ok && c.OutputPath != "" {
return c, true
}
return nil, false
}(); ok {
return cached.OutputPath, cached, nil
if path, cached, ok := w.cachedSource(sourceID); ok {
return path, cached, nil
}

// Check if another goroutine is already running this source (diamond dependency case).
// If so, wait for it to finish and return its result.
return w.runSourceOnce(ctx, parentStep, sourceID, targetID, targetLanguage)
}

// runSourceOnce is RunSource without the lock-free fast path. It guarantees
// that a source is run at most once per workflow while its result is usable:
// a caller either joins the run already in flight (diamond dependency case),
// is served the result of a run that completed since it last looked at the
// cache, or registers itself as the single runner.
func (w *Workflow) runSourceOnce(ctx context.Context, parentStep *workflowTracking.WorkflowStep, sourceID, targetID, targetLanguage string) (string, *SourceResult, error) {
w.sourceInflightMu.Lock()
// Check if another goroutine is already running this source.
// If so, wait for it to finish and return its result.
if inflight, ok := w.sourceInflight[sourceID]; ok {
w.sourceInflightMu.Unlock()
<-inflight.done
Expand All @@ -102,6 +111,15 @@ func (w *Workflow) RunSource(ctx context.Context, parentStep *workflowTracking.W
}
return inflight.path, inflight.result, nil
}
// A run may have finished between the caller's cache check and taking the
// lock. Its result is published to SourceResults before its in-flight
// entry is removed (below), so re-checking here, under the same lock that
// guards registration, is enough to never start a second run for a source
// that already completed.
if path, cached, ok := w.cachedSource(sourceID); ok {
w.sourceInflightMu.Unlock()
return path, cached, nil
}
// Register ourselves as the in-flight resolver for this source
inflight := &sourceInflight{done: make(chan struct{})}
w.sourceInflight[sourceID] = inflight
Expand All @@ -115,9 +133,33 @@ func (w *Workflow) RunSource(ctx context.Context, parentStep *workflowTracking.W
inflight.err = err
close(inflight.done)

// The in-flight entry only exists to let concurrent callers share one run.
// Drop it now that the run is over: a successful result is served from
// SourceResults (already published by runSourceInner, which is what makes
// the re-check above sound), a failed one must be re-run by the next caller
// (e.g. the minimum viable spec retry, which mutates the source before
// re-running).
w.sourceInflightMu.Lock()
delete(w.sourceInflight, sourceID)
w.sourceInflightMu.Unlock()

return path, result, err
}

// cachedSource returns the resolved document of a source that already ran to
// completion in this workflow, so that targets sharing a source only run it
// once. The path returned is the document the pipeline produced, not
// SourceResult.OutputPath: the two differ whenever the output location was not
// written (frozen workflow lockfile runs).
func (w *Workflow) cachedSource(sourceID string) (string, *SourceResult, bool) {
w.sourceMu.Lock()
defer w.sourceMu.Unlock()
if c, ok := w.SourceResults[sourceID]; ok && c.documentPath != "" {
return c.documentPath, c, true
}
return "", nil, false
}

func (w *Workflow) runSourceInner(ctx context.Context, parentStep *workflowTracking.WorkflowStep, sourceID, targetID, targetLanguage string) (string, *SourceResult, error) {
source := w.workflow.Sources[sourceID]

Expand Down Expand Up @@ -170,7 +212,10 @@ func (w *Workflow) runSourceInner(ctx context.Context, parentStep *workflowTrack
defer func() {
w.sourceMu.Lock()
w.SourceResults[sourceID] = sourceRes
w.sourceOrder = append(w.sourceOrder, sourceID)
// A source that failed is run again by the next caller; record it once.
if !slices.Contains(w.sourceOrder, sourceID) {
w.sourceOrder = append(w.sourceOrder, sourceID)
}
w.sourceMu.Unlock()
_ = w.OnSourceResult(sourceRes, SourceStepComplete)
}()
Expand Down Expand Up @@ -335,6 +380,8 @@ func (w *Workflow) runSourceInner(ctx context.Context, parentStep *workflowTrack

rootStep.SucceedWorkflow()

sourceRes.documentPath = currentDocument

return currentDocument, sourceRes, nil
}

Expand Down
Loading
Loading