Files
notarius/internal/framework/pipeline/runner_concurrent.go

667 lines
29 KiB
Go

package pipeline
import (
"context"
"errors"
"fmt"
"path"
"sort"
"strings"
"sync"
"time"
"gitea.maximumdirect.net/eric/notarius/internal/core/artifacts"
"gitea.maximumdirect.net/eric/notarius/internal/core/fileio"
"gitea.maximumdirect.net/eric/notarius/internal/core/source"
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
)
type laneExtractState struct {
index int
prepared preparedLaneExecutor
deps []CheckpointFingerprint
decision CheckpointDecision
reuseEligible bool
values []erasedExtractArtifact
serialized []CheckpointArtifact
warnings []contracts.Warning
rejected []contracts.RejectedOutput
incomplete []int
validationSummaries []artifacts.ValidationSummary
results map[int]extractJobResult
remaining int
failed bool
terminal bool
output RunOutput
}
type finalizedExtractResults struct {
accepted []erasedExtractArtifact
serialized []CheckpointArtifact
warnings []contracts.Warning
rejected []contracts.RejectedOutput
incomplete []int
validationSummaries []artifacts.ValidationSummary
decision CheckpointDecision
reuseEligible bool
}
func loadExtract(loader CheckpointLoader, stepID, laneID, moduleKey string, deps []CheckpointFingerprint) (ExtractCheckpoint, CheckpointDecision) {
if stepAware, ok := loader.(StepCheckpointLoader); ok {
return stepAware.ExtractForStep(stepID, laneID, moduleKey, deps)
}
return loader.Extract(laneID, moduleKey, deps)
}
func recordExtract(recorder CheckpointRecorder, stepID, laneID, moduleKey string, deps []CheckpointFingerprint, outputs []CheckpointArtifact, rejected []contracts.RejectedOutput, warnings []contracts.Warning) error {
return checkpointExtractSucceeded(recorder, stepID, laneID, moduleKey, deps, outputs, rejected, warnings)
}
type extractJob struct {
lane *laneExtractState
chunk source.Chunk
}
type extractJobResult struct {
laneIndex int
chunkIndex int
value erasedExtractArtifact
serialized CheckpointArtifact
warnings []contracts.Warning
rejected *contracts.RejectedOutput
validationIncomplete bool
validationSummary *artifacts.ValidationSummary
err error
}
type extractAttemptValue struct {
artifact erasedExtractArtifact
serialized CheckpointArtifact
terminal *attemptTerminalRecorder
}
type laneCompletion struct {
index int
output RunOutput
err error
}
type orderedRunError struct {
stage int
lane int
chunk int
err error
}
func (r *Runner) runLanes(parent context.Context, input RunInput, step PreparedPipelineStep, checkpoints CheckpointRecorder, loader CheckpointLoader, doc *source.SourceDocument, sourceInput contracts.LLMInputMaterial, sessionID string, chunks []source.Chunk) (RunOutput, error) {
output := RunOutput{Manifest: manifestFromPipeline(input)}
states, err := initializeLaneStates(input, step, checkpoints, loader, doc, chunks, &output)
if err != nil {
if mergeErr := mergeTerminalLaneStates(&output, states); mergeErr != nil {
return output, errors.Join(err, mergeErr)
}
return output, err
}
completedOutputs, runErrors := r.runLaneEngine(parent, laneEngineConfig{
input: input, checkpoints: checkpoints, loader: loader, doc: doc,
sourceInput: sourceInput, sessionID: sessionID, chunks: chunks, states: states,
})
if err := mergeCompletedLanes(&output, completedOutputs); err != nil {
return output, err
}
if err := selectRunError(parent, runErrors); err != nil {
return output, err
}
return output, nil
}
func initializeLaneStates(input RunInput, step PreparedPipelineStep, checkpoints CheckpointRecorder, loader CheckpointLoader, doc *source.SourceDocument, chunks []source.Chunk, output *RunOutput) ([]*laneExtractState, error) {
states := make([]*laneExtractState, len(step.lanes))
for i, prepared := range step.lanes {
if prepared.typed == nil {
return states, fmt.Errorf("typed lane %q executor is not prepared", prepared.resolved.ID)
}
if err := setTypedLaneManifestMetadata(output, prepared.resolved.ID, prepared.typed.extractor, prepared.typed.merger, prepared.typed.normalizer); err != nil {
return states, err
}
if input.CheckpointPolicy.requiresReusable(input.stepID, prepared.resolved.ID) && !input.CheckpointPolicy.forced(input.stepID, prepared.resolved.ID) {
if !laneReferencesReuseEligible(input, prepared.resolved) {
return states, fmt.Errorf("required reusable checkpoint unavailable for step %q lane %q: generated input lineage is validation-incomplete", input.stepID, prepared.resolved.ID)
}
state, err := hydrateRequiredLane(input, loader, doc, i, prepared)
states[i] = state
if err != nil {
return states, err
}
continue
}
state, err := prepareLaneExtract(input, loader, doc, chunks, i, prepared, output)
if err != nil {
return states, err
}
if !state.decision.Reused && state.reuseEligible {
if err := checkpointExtractRunning(checkpoints, input.stepID, prepared.resolved.ID, prepared.resolved.Extract.Module, state.deps); err != nil {
return states, fmt.Errorf("write extract checkpoint for lane %q: %w", prepared.resolved.ID, err)
}
} else if err := finalizeLaneExtract(checkpoints, input.stepID, state); err != nil {
return states, err
}
states[i] = state
}
return states, nil
}
type laneEngineConfig struct {
input RunInput
checkpoints CheckpointRecorder
loader CheckpointLoader
doc *source.SourceDocument
sourceInput contracts.LLMInputMaterial
sessionID string
chunks []source.Chunk
states []*laneExtractState
}
type laneEngine struct {
runner *Runner
laneEngineConfig
ctx context.Context
cancel context.CancelFunc
workerCount int
jobs chan extractJob
results chan extractJobResult
completions chan laneCompletion
continuations chan *laneExtractState
continuationWorkers sync.WaitGroup
pending []*laneExtractState
completedOutputs []RunOutput
runErrors []orderedRunError
launched int
completed int
}
func (r *Runner) runLaneEngine(parent context.Context, config laneEngineConfig) ([]RunOutput, []orderedRunError) {
engine := newLaneEngine(r, parent, config)
return engine.run()
}
func newLaneEngine(r *Runner, parent context.Context, config laneEngineConfig) *laneEngine {
workerCount := config.input.ExtractWorkers
if workerCount < 1 {
workerCount = 1
}
ctx, cancel := context.WithCancel(parent)
return &laneEngine{
runner: r, laneEngineConfig: config, ctx: ctx, cancel: cancel, workerCount: workerCount,
jobs: make(chan extractJob, workerCount), results: make(chan extractJobResult, workerCount),
completions: make(chan laneCompletion, len(config.states)), continuations: make(chan *laneExtractState, workerCount),
completedOutputs: make([]RunOutput, len(config.states)),
}
}
func (e *laneEngine) run() ([]RunOutput, []orderedRunError) {
defer e.cancel()
e.initializeCollection()
e.startExtractWorkers()
e.startContinuationWorkers()
e.collect()
close(e.continuations)
e.continuationWorkers.Wait()
return e.completedOutputs, e.runErrors
}
func (e *laneEngine) initializeCollection() {
for _, state := range e.states {
if state.terminal {
e.completedOutputs[state.index] = state.output
} else if state.decision.Reused {
e.pending = append(e.pending, state)
}
}
}
func (e *laneEngine) startExtractWorkers() {
var workers sync.WaitGroup
for i := 0; i < e.workerCount; i++ {
workers.Add(1)
go func() {
defer workers.Done()
for job := range e.jobs {
if e.ctx.Err() != nil {
continue
}
e.results <- e.runner.runExtractJob(e.ctx, e.input, e.doc, e.sourceInput, e.sessionID, job)
}
}()
}
go e.dispatchExtractJobs()
go func() { workers.Wait(); close(e.results) }()
}
func (e *laneEngine) dispatchExtractJobs() {
defer close(e.jobs)
for chunkIndex := range e.chunks {
for laneIndex := range e.states {
state := e.states[laneIndex]
if state.terminal || state.decision.Reused {
continue
}
select {
case e.jobs <- extractJob{lane: state, chunk: e.chunks[chunkIndex]}:
case <-e.ctx.Done():
return
}
}
}
}
func (e *laneEngine) startContinuationWorkers() {
for i := 0; i < e.workerCount; i++ {
e.continuationWorkers.Add(1)
go func() {
defer e.continuationWorkers.Done()
for state := range e.continuations {
if err := e.ctx.Err(); err != nil {
e.completions <- laneCompletion{index: state.index, err: err}
continue
}
laneOutput, err := e.runner.continueLane(e.ctx, e.input, e.checkpoints, e.loader, e.doc, e.sourceInput, e.sessionID, e.chunks, state)
e.completions <- laneCompletion{index: state.index, output: laneOutput, err: err}
}
}()
}
}
// collect is the sole owner of pending, launch, and completion accounting. It
// must drain closed extract results and one completion for every launched
// continuation, even after cancellation.
func (e *laneEngine) collect() {
resultChannel := (<-chan extractJobResult)(e.results)
for resultChannel != nil || len(e.pending) > 0 || e.completed < e.launched {
var continuationChannel chan<- *laneExtractState
var nextContinuation *laneExtractState
if len(e.pending) > 0 && e.ctx.Err() == nil {
continuationChannel = e.continuations
nextContinuation = e.pending[0]
} else if e.ctx.Err() != nil {
e.pending = nil
}
select {
case continuationChannel <- nextContinuation:
e.pending = e.pending[1:]
e.launched++
case result, ok := <-resultChannel:
if !ok {
resultChannel = nil
continue
}
e.handleExtractResult(result)
case completion := <-e.completions:
e.handleCompletion(completion)
}
}
}
func (e *laneEngine) handleExtractResult(result extractJobResult) {
state := e.states[result.laneIndex]
state.remaining--
if result.err != nil {
state.failed = true
e.runErrors = append(e.runErrors, orderedRunError{stage: 0, lane: state.index, chunk: result.chunkIndex, err: result.err})
if state.reuseEligible {
_ = checkpointExtractFailed(e.checkpoints, e.input.stepID, state.prepared.resolved.ID, state.prepared.resolved.Extract.Module, state.deps, result.err)
}
e.cancel()
} else {
state.results[result.chunkIndex] = result
}
if state.remaining != 0 || state.failed || e.ctx.Err() != nil {
return
}
if err := finalizeLaneExtract(e.checkpoints, e.input.stepID, state); err != nil {
state.failed = true
e.runErrors = append(e.runErrors, orderedRunError{stage: 0, lane: state.index, chunk: len(e.chunks), err: err})
e.cancel()
return
}
e.pending = append(e.pending, state)
}
func (e *laneEngine) handleCompletion(completion laneCompletion) {
e.completed++
e.completedOutputs[completion.index] = completion.output
if completion.err != nil {
e.runErrors = append(e.runErrors, classifyLaneError(completion.index, len(e.chunks), completion.err))
e.cancel()
}
}
func hydrateRequiredLane(input RunInput, loader CheckpointLoader, doc *source.SourceDocument, index int, prepared preparedLaneExecutor) (*laneExtractState, error) {
lane, typed := prepared.resolved, prepared.typed
local := RunOutput{Manifest: manifestFromPipeline(input)}
checkpoint, decision := loader.AcceptedNormalize(input.stepID, lane.ID, lane.Normalize.Module)
if decision.Reused {
decision = checkpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused)
if checkpoint.Output.LaneID != lane.ID || checkpoint.Output.ModuleKey != lane.Normalize.Module || checkpoint.Output.SourceID != doc.ID {
decision = checkpointDecision(CheckpointDecisionExecuted, CheckpointReasonArtifactPayloadInvalid)
}
}
resolution, err := resolveCheckpointDecision(&local, loader, input.CheckpointPolicy, StageNormalize, input.stepID, lane.ID, lane.Normalize.Module, decision, typed.codec, []CheckpointArtifact{checkpoint.Output})
decision = resolution.decision
state := &laneExtractState{index: index, prepared: prepared, decision: decision, reuseEligible: true, terminal: true, output: local}
if err != nil {
return state, err
}
hydrated := resolution.artifacts[0]
local.Warnings = append(local.Warnings, cloneWarnings(checkpoint.Warnings)...)
local.NormalizeOutputs = append(local.NormalizeOutputs, contracts.SerializedOutput{
StepID: input.stepID,
LaneID: lane.ID,
NormalizerKey: lane.Normalize.Module,
SourceID: doc.ID,
Artifact: contracts.CloneSerializedArtifact(hydrated.Artifact),
})
local.normalizeReuseEligibility = map[generatedOutputKey]bool{generatedOutputKeyFor(input.stepID, lane.ID): true}
state.output = local
return state, nil
}
func mergeTerminalLaneStates(output *RunOutput, states []*laneExtractState) error {
for _, state := range states {
if state == nil || !state.terminal {
continue
}
if err := mergeLaneOutput(output, state.output); err != nil {
return err
}
}
return nil
}
func mergeCompletedLanes(output *RunOutput, completedOutputs []RunOutput) error {
for i := range completedOutputs {
if err := mergeLaneOutput(output, completedOutputs[i]); err != nil {
return err
}
}
return nil
}
func prepareLaneExtract(input RunInput, loader CheckpointLoader, doc *source.SourceDocument, chunks []source.Chunk, index int, prepared preparedLaneExecutor, output *RunOutput) (*laneExtractState, error) {
lane, typed := prepared.resolved, prepared.typed
digest, err := joinedChunkDigest(chunks)
if err != nil {
return nil, fmt.Errorf("digest chunks for lane %q: %w", lane.ID, err)
}
extractReferences := operationReferenceSet(input, lane.ExtractReferences)
deps := append(digestFingerprints("chunks", digest), generatedReferenceDependencies(extractReferences)...)
state := &laneExtractState{index: index, prepared: prepared, deps: normalizeCheckpointFingerprints(deps), reuseEligible: referenceTargetReuseEligible(input, lane.ExtractReferences), remaining: len(chunks), results: make(map[int]extractJobResult, len(chunks))}
if !state.reuseEligible {
state.decision = checkpointDecision(CheckpointDecisionExecuted, CheckpointReasonValidationIncompleteLineage)
recordCheckpointEvent(output, loader, string(StageExtract), input.stepID, lane.ID, lane.Extract.Module, state.decision)
return state, nil
}
cp, decision := loadExtract(loader, input.stepID, lane.ID, lane.Extract.Module, state.deps)
resolution, err := resolveCheckpointDecision(output, loader, input.CheckpointPolicy, StageExtract, input.stepID, lane.ID, lane.Extract.Module, decision, typed.codec, cp.Outputs)
if err != nil {
return nil, err
}
decision = resolution.decision
state.decision = decision
if decision.Reused {
state.remaining = 0
for i, stored := range resolution.artifacts {
value := resolution.values[i]
artifact := erasedExtractArtifact{LaneID: lane.ID, ExtractorKey: lane.Extract.Module, SourceID: doc.ID, ChunkID: stored.ChunkID, ChunkIndex: stored.ChunkIndex, ChunkRef: stored.ChunkRef, Value: value}
if stored.ChunkIndex >= 0 && stored.ChunkIndex < len(chunks) && artifact.ChunkRef == (source.SourceRef{}) {
artifact.ChunkRef = chunks[stored.ChunkIndex].Ref
}
state.values = append(state.values, artifact)
state.serialized = append(state.serialized, cloneCheckpointArtifact(stored))
}
state.warnings, state.rejected = cloneWarnings(cp.Warnings), cloneRejectedOutputs(cp.Rejected)
}
return state, nil
}
func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.SourceDocument, sourceInput contracts.LLMInputMaterial, sessionID string, job extractJob) extractJobResult {
state := job.lane
chunk, cloneErr := cloneSourceChunk(job.chunk)
lane, typed := state.prepared.resolved, state.prepared.typed
result := extractJobResult{laneIndex: state.index, chunkIndex: job.chunk.Index}
if cloneErr != nil {
result.err = fmt.Errorf("clone chunk %q for extraction: %w", job.chunk.ID, cloneErr)
return result
}
terminalResult, err := runProducerAttempts(ctx, producerAttemptConfig{Retries: lane.Extract.Retries, Policy: lane.ExtractValidationPolicy, AllowStructuralRetry: lane.ExtractExecutionClass == contracts.ExecutionClassLLMBacked}, func(attemptCtx context.Context, request producerAttemptRequest) (producerAttemptOutput, error) {
attempt := request.Number
started := time.Now().UTC()
attemptPath := path.Join("extract", fileio.EncodePathComponent(lane.ID), fmt.Sprintf("chunk-%06d", chunk.Index+1), fmt.Sprintf("attempt-%02d", attempt))
attemptCtx, llmScope := withDebugLLMScope(attemptCtx, attemptPath)
terminal := newAttemptTerminalRecorder(input.Debug, attemptPath, "extract", llmScope, debugTimedEnvelope{Stage: string(StageExtract), StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Extract.Module, Attempt: attempt, AttemptKind: string(request.Kind), StartedAt: started})
requestMetadata, metadataErr := cloneMetadata(input.Metadata)
if metadataErr != nil {
return producerAttemptOutput{}, terminal.record(nil, fmt.Errorf("clone extract request metadata: %w", metadataErr))
}
extractReferences := operationReferenceSet(input, lane.ExtractReferences)
extracted, callErr := typed.extract(attemptCtx, typed.extractor, contracts.TypedExtractionRequest{Source: doc, Chunk: &chunk, SourceInput: chunkInputMaterial(sourceInput, chunk), SessionID: sessionID, References: CloneReferenceSet(extractReferences), LLMProfile: lane.Extract.LLMProfile, StructuredOutputRepairAttempts: cloneStructuredOutputRepairAttempts(lane.Extract.StructuredOutputRepairAttempts), Correction: request.Correction, Metadata: requestMetadata})
if callErr != nil {
attemptErr := fmt.Errorf("extract lane %q chunk %q with extractor %q: %w", lane.ID, chunk.ID, lane.Extract.Module, callErr)
return producerAttemptOutput{}, terminal.record(nil, attemptErr)
}
artifact := erasedExtractArtifact{LaneID: lane.ID, ExtractorKey: lane.Extract.Module, SourceID: doc.ID, ChunkID: chunk.ID, ChunkIndex: chunk.Index, ChunkRef: chunk.Ref, Value: extracted.Value}
attemptWarnings := cloneWarnings(extracted.Warnings)
serializedCandidate, encodeErr := serializeCandidateArtifact(typed.codec, artifact.LaneID, artifact.ExtractorKey, artifact.SourceID, artifact.Value)
if encodeErr != nil {
attemptErr := fmt.Errorf("serialize extract candidate for lane %q chunk %q: %w", lane.ID, chunk.ID, encodeErr)
payload := map[string]any{"warnings": debugWarningEnvelopes(attemptWarnings)}
return producerAttemptOutput{}, terminal.record(payload, attemptErr)
}
serializedCandidate.ChunkID, serializedCandidate.ChunkIndex, serializedCandidate.ChunkRef = artifact.ChunkID, artifact.ChunkIndex, artifact.ChunkRef
return producerAttemptOutput{Value: extractAttemptValue{artifact: artifact, serialized: serializedCandidate, terminal: &terminal}, Candidate: extracted.ModelCandidate, Warnings: attemptWarnings}, nil
}, func(validationCtx context.Context, output producerAttemptOutput) (validationReport, error) {
candidate, ok := output.Value.(extractAttemptValue)
if !ok {
return validationReport{}, fmt.Errorf("extract attempt has incompatible value")
}
report, validationErr := r.validateTypedReport(validationCtx, typed.codec, typedValidationTarget{stage: StageExtract, stepID: input.stepID, laneID: lane.ID, moduleKey: lane.Extract.Module, source: doc, sourceID: doc.ID, sourceInput: chunkInputMaterial(sourceInput, chunk), sessionID: sessionID, references: operationReferenceSet(input, lane.ExtractReferences), metadata: input.Metadata, chunk: &chunk, ref: chunk.Ref, value: candidate.artifact.Value, candidate: &candidate.serialized}, state.prepared.extractValidators, candidate.terminal.envelope.Attempt, input.Debug)
payload := map[string]any{
"output": debugCheckpointArtifact(candidate.serialized),
"warnings": debugWarningEnvelopes(append(cloneWarnings(output.Warnings), report.Warnings()...)),
"rejection": debugRejectedOutputPtr(typedRejection(report, typedValidationTarget{stage: StageExtract, stepID: input.stepID, laneID: lane.ID, moduleKey: lane.Extract.Module, chunk: &chunk}, candidate.terminal.envelope.Attempt)),
}
if validationErr != nil {
return report, candidate.terminal.record(payload, validationErr)
}
return report, candidate.terminal.record(payload, nil)
})
result.err = err
summary := validationSummary(terminalResult, StageExtract, input.stepID, lane.ID, lane.Extract.Module, chunk.ID, chunk.Index)
if debugErr := writeProducerTerminalDebug(input.Debug, path.Join("extract", fileio.EncodePathComponent(lane.ID), fmt.Sprintf("chunk-%06d", chunk.Index+1), "terminal.json"), terminalResult, lane.ExtractValidationPolicy, summary); debugErr != nil {
result.err = errors.Join(result.err, debugErr)
return result
}
if err == nil && terminalResult.Action == producerTerminalRejected {
result.rejected = terminalResult.Rejection
if result.rejected != nil {
result.rejected.Stage, result.rejected.StepID, result.rejected.LaneID, result.rejected.ModuleKey, result.rejected.ChunkID, result.rejected.ChunkIndex = string(StageExtract), input.stepID, lane.ID, lane.Extract.Module, chunk.ID, chunk.Index
result.validationSummary = &summary
result.rejected.Validation = cloneValidationSummaryPtr(result.validationSummary)
}
result.warnings = cloneWarnings(terminalResult.Warnings)
return result
}
if err == nil {
candidate, ok := terminalResult.Value.(extractAttemptValue)
if !ok {
result.err = fmt.Errorf("extract attempt terminal has incompatible value")
return result
}
stored, encodeErr := checkpointArtifact(typed.codec, candidate.artifact.LaneID, candidate.artifact.ExtractorKey, candidate.artifact.SourceID, candidate.artifact.Value)
payload := map[string]any{"output": debugCheckpointArtifact(candidate.serialized), "warnings": debugWarningEnvelopes(terminalResult.Warnings), "rejection": debugRejectedOutputPtr(nil)}
if encodeErr != nil {
attemptErr := fmt.Errorf("serialize accepted extract output for lane %q chunk %q: %w", lane.ID, chunk.ID, encodeErr)
result.err = candidate.terminal.record(payload, attemptErr)
return result
}
stored.ChunkID, stored.ChunkIndex, stored.ChunkRef = candidate.artifact.ChunkID, candidate.artifact.ChunkIndex, candidate.artifact.ChunkRef
if debugErr := candidate.terminal.record(payload, nil); debugErr != nil {
result.err = debugErr
return result
}
result.value, result.serialized = candidate.artifact, stored
result.warnings = cloneWarnings(terminalResult.Warnings)
result.validationIncomplete = terminalResult.ValidationIncomplete
result.validationSummary = &summary
}
return result
}
func finalizeLaneExtract(checkpoints CheckpointRecorder, stepID string, state *laneExtractState) error {
lane := state.prepared.resolved
indexes := make([]int, 0, len(state.results))
for index := range state.results {
indexes = append(indexes, index)
}
sort.Ints(indexes)
for _, index := range indexes {
result := state.results[index]
if result.validationSummary != nil {
state.validationSummaries = append(state.validationSummaries, artifacts.CloneValidationSummary(*result.validationSummary))
}
if result.rejected != nil {
state.rejected = append(state.rejected, *result.rejected)
state.warnings = append(state.warnings, result.warnings...)
continue
}
state.values = append(state.values, result.value)
state.serialized = append(state.serialized, result.serialized)
state.warnings = append(state.warnings, result.warnings...)
if result.validationIncomplete {
state.incomplete = append(state.incomplete, result.chunkIndex)
}
}
sort.SliceStable(state.values, func(i, j int) bool { return state.values[i].ChunkIndex < state.values[j].ChunkIndex })
sort.SliceStable(state.serialized, func(i, j int) bool { return state.serialized[i].ChunkIndex < state.serialized[j].ChunkIndex })
sort.SliceStable(state.rejected, func(i, j int) bool { return state.rejected[i].ChunkIndex < state.rejected[j].ChunkIndex })
sort.Ints(state.incomplete)
if len(state.incomplete) > 0 {
state.reuseEligible = false
}
if !state.decision.Reused && state.reuseEligible {
if err := recordExtract(checkpoints, stepID, lane.ID, lane.Extract.Module, state.deps, state.serialized, state.rejected, state.warnings); err != nil {
return fmt.Errorf("write extract checkpoint for lane %q: %w", lane.ID, err)
}
}
return nil
}
func (r *Runner) continueLane(ctx context.Context, input RunInput, checkpoints CheckpointRecorder, loader CheckpointLoader, doc *source.SourceDocument, sourceInput contracts.LLMInputMaterial, sessionID string, chunks []source.Chunk, state *laneExtractState) (RunOutput, error) {
lane := state.prepared.resolved
local := RunOutput{Manifest: manifestFromPipeline(input)}
results := finalizedExtractResults{
accepted: state.values,
serialized: state.serialized,
warnings: state.warnings,
rejected: state.rejected,
incomplete: state.incomplete,
validationSummaries: state.validationSummaries,
decision: state.decision,
reuseEligible: state.reuseEligible,
}
local.Warnings = append(local.Warnings, cloneWarnings(results.warnings)...)
local.Rejected = append(local.Rejected, cloneRejectedOutputs(results.rejected)...)
local.ValidationSummaries = append(local.ValidationSummaries, cloneValidationSummaries(results.validationSummaries)...)
if err := writeDebugTimed(input.Debug, path.Join("extract", fileio.EncodePathComponent(lane.ID), "input.json"), debugTimedEnvelope{Stage: string(StageExtract), StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Extract.Module, StartedAt: time.Now().UTC(), Payload: map[string]any{"reused": results.decision.Reused, "decision": results.decision, "source": debugSourceDocumentEnvelope(doc), "chunks": debugSourceChunkEnvelopes(chunks), "options": redactSensitiveMap(lane.Extract.Options), "metadata": redactSensitiveMap(input.Metadata)}}); err != nil {
return local, &laneRunError{stage: StageExtract, err: err}
}
if err := writeDebugTimed(input.Debug, path.Join("extract", fileio.EncodePathComponent(lane.ID), "output.json"), debugTimedEnvelope{Stage: string(StageExtract), StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Extract.Module, StartedAt: time.Now().UTC(), Payload: map[string]any{"reused": results.decision.Reused, "outputs": debugCheckpointArtifacts(results.serialized), "rejected": debugRejectedOutputEnvelopes(results.rejected), "warnings": debugWarningEnvelopes(results.warnings), "validation_incomplete_chunks": append([]int(nil), results.incomplete...)}}); err != nil {
return local, &laneRunError{stage: StageExtract, err: err}
}
if len(results.accepted) == 0 {
return local, nil
}
err := r.continueTypedLane(ctx, input, checkpoints, loader, doc, sourceInput, sessionID, state.prepared, results, &local)
return local, err
}
func classifyLaneError(lane, sentinel int, err error) orderedRunError {
stage := 1
var laneErr *laneRunError
if errors.As(err, &laneErr) {
switch laneErr.stage {
case StageExtract:
stage = 0
case StageNormalize:
stage = 2
}
} else if strings.Contains(strings.ToLower(err.Error()), "normalize") {
stage = 2
}
return orderedRunError{stage: stage, lane: lane, chunk: sentinel, err: err}
}
func selectRunError(parent context.Context, values []orderedRunError) error {
if err := parent.Err(); err != nil {
return err
}
if len(values) == 0 {
return nil
}
hasReal := false
for _, value := range values {
if !errors.Is(value.err, context.Canceled) && !errors.Is(value.err, context.DeadlineExceeded) {
hasReal = true
break
}
}
filtered := values[:0]
for _, value := range values {
if hasReal && (errors.Is(value.err, context.Canceled) || errors.Is(value.err, context.DeadlineExceeded)) {
continue
}
filtered = append(filtered, value)
}
sort.SliceStable(filtered, func(i, j int) bool {
if filtered[i].stage != filtered[j].stage {
return filtered[i].stage < filtered[j].stage
}
if filtered[i].lane != filtered[j].lane {
return filtered[i].lane < filtered[j].lane
}
return filtered[i].chunk < filtered[j].chunk
})
return filtered[0].err
}
func mergeLaneOutput(dst *RunOutput, src RunOutput) error {
if dst == nil {
return nil
}
dst.NormalizeOutputs = append(dst.NormalizeOutputs, cloneSerializedOutputs(src.NormalizeOutputs)...)
dst.Rejected = append(dst.Rejected, cloneRejectedOutputs(src.Rejected)...)
dst.Warnings = append(dst.Warnings, cloneWarnings(src.Warnings)...)
dst.CheckpointEvents = append(dst.CheckpointEvents, src.CheckpointEvents...)
dst.ValidationSummaries = append(dst.ValidationSummaries, cloneValidationSummaries(src.ValidationSummaries)...)
if len(src.normalizeReuseEligibility) > 0 {
if dst.normalizeReuseEligibility == nil {
dst.normalizeReuseEligibility = make(map[generatedOutputKey]bool, len(src.normalizeReuseEligibility))
}
for key, eligible := range src.normalizeReuseEligibility {
dst.normalizeReuseEligibility[key] = eligible
}
}
for i := range dst.Manifest.ArtifactLanes {
for j := range src.Manifest.ArtifactLanes {
if dst.Manifest.ArtifactLanes[i].ID == src.Manifest.ArtifactLanes[j].ID && dst.Manifest.ArtifactLanes[i].StepID == src.Manifest.ArtifactLanes[j].StepID && src.Manifest.ArtifactLanes[j].Metadata != nil {
metadata, err := cloneMetadata(src.Manifest.ArtifactLanes[j].Metadata)
if err != nil {
return fmt.Errorf("clone lane %q manifest metadata: %w", dst.Manifest.ArtifactLanes[i].ID, err)
}
dst.Manifest.ArtifactLanes[i].Metadata = metadata
}
}
}
return nil
}