Integrate chunk validation retries
This commit is contained in:
@@ -23,6 +23,14 @@ type chunkPlanExecution struct {
|
||||
summary artifacts.ChunkPlanSummary
|
||||
}
|
||||
|
||||
type generatedChunkPlanCandidate struct {
|
||||
plan source.ChunkPlan
|
||||
chunks []source.Chunk
|
||||
record ChunkPlanRecord
|
||||
producerWarnings []contracts.Warning
|
||||
terminal *attemptTerminalRecorder
|
||||
}
|
||||
|
||||
func effectiveChunkCacheMode(mode ChunkCacheMode) ChunkCacheMode {
|
||||
if mode == "" {
|
||||
return ChunkCacheBypass
|
||||
@@ -58,17 +66,31 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
|
||||
case ChunkPlanHit:
|
||||
plan, chunks, validationErr := validateAndMaterializeChunkPlan(doc, record.Plan)
|
||||
if validationErr == nil {
|
||||
if err := result.setCandidate(record, "reused"); err != nil {
|
||||
return result, fmt.Errorf("clone reused chunk plan record: %w", err)
|
||||
report, err := r.validateChunkReport(ctx, doc, chunker.Key(), chunks, sourceInput, sessionID, input.pipeline.ChunkReferences.ReferenceSet, input.Metadata, input.Prepared.chunkValidators, 1, input.Debug)
|
||||
if err != nil {
|
||||
result.setValidation(report.Warnings(), nil, err)
|
||||
return result, err
|
||||
}
|
||||
validationWarnings, rejection, err := r.validateChunks(ctx, doc, chunker.Key(), chunks, sourceInput, sessionID, input.pipeline.ChunkReferences.ReferenceSet, input.Metadata, input.Prepared.chunkValidators, 1, input.Debug)
|
||||
result.plan = &plan
|
||||
result.chunks = chunks
|
||||
result.warnings = append(cloneWarnings(record.Warnings), validationWarnings...)
|
||||
result.rejection = rejection
|
||||
result.accepted = rejection == nil && err == nil
|
||||
result.setValidation(validationWarnings, rejection, err)
|
||||
return result, err
|
||||
if report.FirstRejection() == nil {
|
||||
if incomplete := firstIncompleteValidation(report); incomplete != nil && input.pipeline.ChunkValidationPolicy.ValidatorFailure == ValidatorFailureFailRun {
|
||||
return result, validatorFailureError(*incomplete)
|
||||
}
|
||||
if err := result.setCandidate(record, "reused"); err != nil {
|
||||
return result, fmt.Errorf("clone reused chunk plan record: %w", err)
|
||||
}
|
||||
result.plan = &plan
|
||||
result.chunks = chunks
|
||||
result.warnings = append(cloneWarnings(record.Warnings), report.Warnings()...)
|
||||
result.accepted = true
|
||||
result.setValidation(report.Warnings(), nil, nil)
|
||||
if firstIncompleteValidation(report) != nil {
|
||||
result.summary.ValidationStatus = "incomplete"
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
// A cache hit is not model material. Its rejection is discarded and
|
||||
// generation begins with the ordinary initial request below.
|
||||
result.setValidation(report.Warnings(), chunkRejection(report, 1, chunker.Key()), nil)
|
||||
}
|
||||
result.lookup = ChunkPlanDecision{Status: ChunkPlanInvalid, Reason: chunkPlanLookupReason(ChunkPlanInvalid)}
|
||||
result.summary.LookupStatus = "invalid"
|
||||
@@ -83,37 +105,41 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
|
||||
result.summary.PublicationStatus = "not_published"
|
||||
}
|
||||
|
||||
var producerWarnings []contracts.Warning
|
||||
retryResult, err := runSimpleRetry(ctx, input.pipeline.Chunk.Retries, func(attempt int) (retryAttemptResult, error) {
|
||||
terminal, err := runProducerAttempts(ctx, producerAttemptConfig{
|
||||
Retries: input.pipeline.Chunk.Retries,
|
||||
Policy: input.pipeline.ChunkValidationPolicy,
|
||||
AllowStructuralRetry: input.pipeline.ChunkExecutionClass == contracts.ExecutionClassLLMBacked,
|
||||
}, func(attemptCtx context.Context, request producerAttemptRequest) (producerAttemptOutput, error) {
|
||||
attempt := request.Number
|
||||
attemptStarted := time.Now().UTC()
|
||||
attemptPath := path.Join("chunk", fmt.Sprintf("attempt-%02d", attempt))
|
||||
attemptCtx, llmScope := withDebugLLMScope(ctx, attemptPath)
|
||||
terminal := newAttemptTerminalRecorder(input.Debug, attemptPath, "chunk", llmScope, debugTimedEnvelope{Stage: string(StageChunk), ModuleKey: chunker.Key(), Attempt: attempt, StartedAt: attemptStarted})
|
||||
attemptCtx, llmScope := withDebugLLMScope(attemptCtx, attemptPath)
|
||||
attemptTerminal := newAttemptTerminalRecorder(input.Debug, attemptPath, "chunk", llmScope, debugTimedEnvelope{Stage: string(StageChunk), ModuleKey: chunker.Key(), Attempt: attempt, StartedAt: attemptStarted})
|
||||
requestMetadata, metadataErr := cloneMetadata(input.Metadata)
|
||||
if metadataErr != nil {
|
||||
return retryAttemptResult{}, terminal.record(nil, fmt.Errorf("clone chunk request metadata: %w", metadataErr))
|
||||
return producerAttemptOutput{}, attemptTerminal.record(nil, fmt.Errorf("clone chunk request metadata: %w", metadataErr))
|
||||
}
|
||||
chunkResult, callErr := chunker.Plan(attemptCtx, contracts.ChunkRequest{
|
||||
Source: doc, SourceInput: sourceInput.Clone(), SessionID: sessionID,
|
||||
References: CloneReferenceSet(input.pipeline.ChunkReferences.ReferenceSet),
|
||||
LLMProfile: input.pipeline.Chunk.LLMProfile, StructuredOutputRepairAttempts: cloneStructuredOutputRepairAttempts(input.pipeline.Chunk.StructuredOutputRepairAttempts), Metadata: requestMetadata,
|
||||
LLMProfile: input.pipeline.Chunk.LLMProfile, StructuredOutputRepairAttempts: cloneStructuredOutputRepairAttempts(input.pipeline.Chunk.StructuredOutputRepairAttempts), Correction: request.Correction, Metadata: requestMetadata,
|
||||
})
|
||||
if callErr != nil {
|
||||
return retryAttemptResult{}, terminal.record(nil, fmt.Errorf("chunk source with chunker %q: %w", chunker.Key(), callErr))
|
||||
return producerAttemptOutput{}, attemptTerminal.record(nil, fmt.Errorf("chunk source with chunker %q: %w", chunker.Key(), callErr))
|
||||
}
|
||||
plan, chunks, validationErr := validateAndMaterializeChunkPlan(doc, chunkResult.Plan)
|
||||
if validationErr != nil {
|
||||
attemptErr := fmt.Errorf("validate chunk plan from chunker %q: %w", chunker.Key(), validationErr)
|
||||
payload := map[string]any{"plan": debugChunkPlanEnvelope(chunkResult.Plan), "warnings": debugWarningEnvelopes(chunkResult.Warnings)}
|
||||
return retryAttemptResult{}, terminal.record(payload, attemptErr)
|
||||
return producerAttemptOutput{}, attemptTerminal.record(payload, fmt.Errorf("%w: %v", contracts.ErrInvalidStructuredOutput, attemptErr))
|
||||
}
|
||||
planDigest, digestErr := source.DigestChunkPlan(plan)
|
||||
if digestErr != nil {
|
||||
return retryAttemptResult{}, terminal.record(nil, fmt.Errorf("digest generated chunk plan: %w", digestErr))
|
||||
return producerAttemptOutput{}, attemptTerminal.record(nil, fmt.Errorf("digest generated chunk plan: %w", digestErr))
|
||||
}
|
||||
producerMetadata, _, metadataErr := moduleManifestMetadata(chunker)
|
||||
if metadataErr != nil {
|
||||
return retryAttemptResult{}, terminal.record(nil, fmt.Errorf("clone chunker manifest metadata: %w", metadataErr))
|
||||
return producerAttemptOutput{}, attemptTerminal.record(nil, fmt.Errorf("clone chunker manifest metadata: %w", metadataErr))
|
||||
}
|
||||
profile := ""
|
||||
if input.pipeline.ChunkExecutionClass == contracts.ExecutionClassLLMBacked {
|
||||
@@ -129,42 +155,27 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
|
||||
},
|
||||
Warnings: cloneWarnings(chunkResult.Warnings), CreatedAt: time.Now().UTC(),
|
||||
}
|
||||
action := "generated"
|
||||
if mode == ChunkCacheRefresh {
|
||||
action = "refreshed"
|
||||
return producerAttemptOutput{Value: generatedChunkPlanCandidate{plan: plan, chunks: chunks, record: candidate, producerWarnings: cloneWarnings(chunkResult.Warnings), terminal: &attemptTerminal}, Candidate: chunkResult.ModelCandidate, Warnings: cloneWarnings(chunkResult.Warnings)}, nil
|
||||
}, func(validationCtx context.Context, output producerAttemptOutput) (validationReport, error) {
|
||||
candidate, ok := output.Value.(generatedChunkPlanCandidate)
|
||||
if !ok {
|
||||
return validationReport{}, fmt.Errorf("chunk attempt candidate has incompatible type")
|
||||
}
|
||||
if mode == ChunkCacheBypass {
|
||||
action = "bypassed"
|
||||
}
|
||||
if candidateErr := result.setCandidate(candidate, action); candidateErr != nil {
|
||||
return retryAttemptResult{}, terminal.record(nil, fmt.Errorf("clone generated chunk plan record: %w", candidateErr))
|
||||
}
|
||||
validationWarnings, rejected, validationErr := r.validateChunks(attemptCtx, doc, chunker.Key(), chunks, sourceInput, sessionID, input.pipeline.ChunkReferences.ReferenceSet, input.Metadata, input.Prepared.chunkValidators, attempt, input.Debug)
|
||||
attemptWarnings := append(cloneWarnings(chunkResult.Warnings), validationWarnings...)
|
||||
report, validationErr := r.validateChunkReport(validationCtx, doc, chunker.Key(), candidate.chunks, sourceInput, sessionID, input.pipeline.ChunkReferences.ReferenceSet, input.Metadata, input.Prepared.chunkValidators, candidate.terminal.envelope.Attempt, input.Debug)
|
||||
attemptWarnings := append(cloneWarnings(output.Warnings), report.Warnings()...)
|
||||
payload := map[string]any{
|
||||
"plan": debugChunkPlanEnvelope(plan), "materialized_chunks": debugSourceChunkEnvelopes(chunks),
|
||||
"warnings": debugWarningEnvelopes(attemptWarnings), "rejection": debugRejectedOutputPtr(rejected),
|
||||
"plan": debugChunkPlanEnvelope(candidate.plan), "materialized_chunks": debugSourceChunkEnvelopes(candidate.chunks),
|
||||
"warnings": debugWarningEnvelopes(attemptWarnings), "rejection": debugRejectedOutputPtr(chunkRejection(report, candidate.terminal.envelope.Attempt, chunker.Key())),
|
||||
}
|
||||
if validationErr != nil {
|
||||
result.setValidation(validationWarnings, rejected, validationErr)
|
||||
return retryAttemptResult{}, terminal.record(payload, validationErr)
|
||||
return report, candidate.terminal.record(payload, validationErr)
|
||||
}
|
||||
if rejected != nil {
|
||||
result.setValidation(validationWarnings, rejected, nil)
|
||||
if debugErr := terminal.record(payload, nil); debugErr != nil {
|
||||
return retryAttemptResult{}, debugErr
|
||||
if report.FirstRejection() == nil && input.pipeline.ChunkValidationPolicy.ValidatorFailure == ValidatorFailureFailRun {
|
||||
if failure := firstIncompleteValidation(report); failure != nil {
|
||||
return report, candidate.terminal.record(payload, validatorFailureError(*failure))
|
||||
}
|
||||
return retryAttemptResult{rejection: rejected, warnings: attemptWarnings}, nil
|
||||
}
|
||||
result.chunks = chunks
|
||||
result.plan = &plan
|
||||
result.warnings = attemptWarnings
|
||||
producerWarnings = cloneWarnings(chunkResult.Warnings)
|
||||
result.setValidation(validationWarnings, nil, nil)
|
||||
if debugErr := terminal.record(payload, nil); debugErr != nil {
|
||||
return retryAttemptResult{}, debugErr
|
||||
}
|
||||
return retryAttemptResult{accepted: true}, nil
|
||||
return report, candidate.terminal.record(payload, nil)
|
||||
})
|
||||
if err != nil {
|
||||
if result.summary.ValidationStatus == "not_run" {
|
||||
@@ -172,19 +183,46 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
|
||||
}
|
||||
return result, err
|
||||
}
|
||||
result.accepted = retryResult.accepted
|
||||
result.rejection = retryResult.rejection
|
||||
if !retryResult.accepted {
|
||||
result.warnings = cloneWarnings(retryResult.warnings)
|
||||
if terminal.Action == producerTerminalRejected {
|
||||
result.rejection = terminal.Rejection
|
||||
if result.rejection != nil {
|
||||
result.rejection.Stage = string(StageChunk)
|
||||
result.rejection.ModuleKey = chunker.Key()
|
||||
}
|
||||
result.warnings = cloneWarnings(terminal.Warnings)
|
||||
result.setValidation(terminal.Validation.Warnings(), result.rejection, nil)
|
||||
return result, nil
|
||||
}
|
||||
candidate, ok := terminal.Value.(generatedChunkPlanCandidate)
|
||||
if !ok {
|
||||
return result, fmt.Errorf("chunk attempt terminal has incompatible value")
|
||||
}
|
||||
if err := result.setCandidate(candidate.record, "generated"); err != nil {
|
||||
return result, fmt.Errorf("clone generated chunk plan record: %w", err)
|
||||
}
|
||||
if mode == ChunkCacheRefresh {
|
||||
result.action = "refreshed"
|
||||
result.summary.Action = "refreshed"
|
||||
}
|
||||
if mode == ChunkCacheBypass {
|
||||
result.action = "bypassed"
|
||||
result.summary.Action = "bypassed"
|
||||
}
|
||||
result.accepted = true
|
||||
result.plan = &candidate.plan
|
||||
result.chunks = candidate.chunks
|
||||
result.warnings = cloneWarnings(terminal.Warnings)
|
||||
result.setValidation(terminal.Validation.Warnings(), nil, nil)
|
||||
if terminal.ValidationIncomplete {
|
||||
result.summary.ValidationStatus = "incomplete"
|
||||
}
|
||||
|
||||
if mode == ChunkCacheAuto || mode == ChunkCacheRefresh {
|
||||
if (mode == ChunkCacheAuto || mode == ChunkCacheRefresh) && !terminal.ValidationIncomplete {
|
||||
record, cloneErr := cloneChunkPlanRecord(*result.record)
|
||||
if cloneErr != nil {
|
||||
return result, fmt.Errorf("clone chunk plan record for publication: %w", cloneErr)
|
||||
}
|
||||
record.Warnings = cloneWarnings(producerWarnings)
|
||||
record.Warnings = cloneWarnings(candidate.producerWarnings)
|
||||
if err := input.ChunkPlans.Save(record); err != nil {
|
||||
return result, fmt.Errorf("save chunk plan: %w", err)
|
||||
}
|
||||
@@ -193,6 +231,14 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func chunkRejection(report validationReport, attempt int, moduleKey string) *contracts.RejectedOutput {
|
||||
rejection := report.FirstRejection()
|
||||
if rejection == nil {
|
||||
return nil
|
||||
}
|
||||
return &contracts.RejectedOutput{Stage: string(StageChunk), ModuleKey: moduleKey, ValidatorName: rejection.validatorName, ReasonCode: rejection.reasonCode, Message: rejection.message, AttemptCount: attempt, DiagnosticArtifactPath: rejection.diagnosticPath}
|
||||
}
|
||||
|
||||
func chunkPlanLookupStatus(status ChunkPlanStatus) string {
|
||||
switch status {
|
||||
case ChunkPlanHit:
|
||||
|
||||
Reference in New Issue
Block a user