Finalize structured diagnostic aggregation

This commit is contained in:
2026-08-27 16:40:17 +00:00
parent 480680b257
commit 4dbbf68051
112 changed files with 656 additions and 893 deletions

View File

@@ -24,7 +24,6 @@ type laneExtractState struct {
reuseEligible bool
values []erasedExtractArtifact
serialized []CheckpointArtifact
warnings []contracts.Warning
diagnostics []contracts.DiagnosticGroup
rejected []contracts.RejectedOutput
incomplete []int
@@ -39,7 +38,6 @@ type laneExtractState struct {
type finalizedExtractResults struct {
accepted []erasedExtractArtifact
serialized []CheckpointArtifact
warnings []contracts.Warning
diagnostics []contracts.DiagnosticGroup
rejected []contracts.RejectedOutput
incomplete []int
@@ -55,8 +53,8 @@ func loadExtract(loader CheckpointLoader, stepID, laneID, moduleKey string, 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)
func recordExtract(recorder CheckpointRecorder, stepID, laneID, moduleKey string, deps []CheckpointFingerprint, outputs []CheckpointArtifact, rejected []contracts.RejectedOutput) error {
return checkpointExtractSucceeded(recorder, stepID, laneID, moduleKey, deps, outputs, rejected)
}
type extractJob struct {
@@ -69,7 +67,6 @@ type extractJobResult struct {
chunkIndex int
value erasedExtractArtifact
serialized CheckpointArtifact
warnings []contracts.Warning
diagnostics []contracts.DiagnosticGroup
rejected *contracts.RejectedOutput
validationIncomplete bool
@@ -356,7 +353,6 @@ func hydrateRequiredLane(input RunInput, loader CheckpointLoader, doc *source.So
return state, err
}
hydrated := resolution.artifacts[0]
local.Warnings = append(local.Warnings, cloneWarnings(checkpoint.Warnings)...)
diagnostics, diagnosticErr := promoteCheckpointDiagnostics(hydrated.Diagnostics, contracts.DiagnosticOrigin{Stage: contracts.DiagnosticOriginStageNormalize, StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Normalize.Module}, nil)
if diagnosticErr != nil {
return state, fmt.Errorf("promote reused accepted diagnostics: %w", diagnosticErr)
@@ -433,7 +429,7 @@ func prepareLaneExtract(input RunInput, loader CheckpointLoader, doc *source.Sou
}
state.diagnostics = append(state.diagnostics, diagnostics...)
}
state.warnings, state.rejected = cloneWarnings(cp.Warnings), cloneRejectedOutputs(cp.Rejected)
state.rejected = cloneRejectedOutputs(cp.Rejected)
}
return state, nil
}
@@ -464,15 +460,14 @@ func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.
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)}
payload := map[string]any{}
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, Diagnostics: contracts.CloneProducerDiagnostics(extracted.Diagnostics)}, nil
return producerAttemptOutput{Value: extractAttemptValue{artifact: artifact, serialized: serializedCandidate, terminal: &terminal}, Candidate: extracted.ModelCandidate, Diagnostics: contracts.CloneProducerDiagnostics(extracted.Diagnostics)}, nil
}, func(validationCtx context.Context, output producerAttemptOutput) (validationReport, error) {
candidate, ok := output.Value.(extractAttemptValue)
if !ok {
@@ -481,7 +476,6 @@ func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.
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 {
@@ -510,7 +504,6 @@ func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.
result.validationSummary = &summary
result.rejected.Validation = cloneValidationSummaryPtr(result.validationSummary)
}
result.warnings = cloneWarnings(terminalResult.Warnings)
return result
}
if err == nil {
@@ -520,7 +513,7 @@ func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.
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)}
payload := map[string]any{"output": debugCheckpointArtifact(candidate.serialized), "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)
@@ -533,7 +526,6 @@ func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.
return result
}
result.value, result.serialized = candidate.artifact, stored
result.warnings = cloneWarnings(terminalResult.Warnings)
result.validationIncomplete = terminalResult.ValidationIncomplete
result.validationSummary = &summary
}
@@ -554,13 +546,11 @@ func finalizeLaneExtract(checkpoints CheckpointRecorder, stepID string, state *l
}
if result.rejected != nil {
state.rejected = append(state.rejected, *result.rejected)
state.warnings = append(state.warnings, result.warnings...)
state.diagnostics = append(state.diagnostics, result.diagnostics...)
continue
}
state.values = append(state.values, result.value)
state.serialized = append(state.serialized, result.serialized)
state.warnings = append(state.warnings, result.warnings...)
state.diagnostics = append(state.diagnostics, result.diagnostics...)
if result.validationIncomplete {
state.incomplete = append(state.incomplete, result.chunkIndex)
@@ -574,7 +564,7 @@ func finalizeLaneExtract(checkpoints CheckpointRecorder, stepID string, state *l
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 {
if err := recordExtract(checkpoints, stepID, lane.ID, lane.Extract.Module, state.deps, state.serialized, state.rejected); err != nil {
return fmt.Errorf("write extract checkpoint for lane %q: %w", lane.ID, err)
}
}
@@ -587,7 +577,6 @@ func (r *Runner) continueLane(ctx context.Context, input RunInput, checkpoints C
results := finalizedExtractResults{
accepted: state.values,
serialized: state.serialized,
warnings: state.warnings,
diagnostics: state.diagnostics,
rejected: state.rejected,
incomplete: state.incomplete,
@@ -595,14 +584,13 @@ func (r *Runner) continueLane(ctx context.Context, input RunInput, checkpoints C
decision: state.decision,
reuseEligible: state.reuseEligible,
}
local.Warnings = append(local.Warnings, cloneWarnings(results.warnings)...)
local.diagnosticGroups = append(local.diagnosticGroups, contracts.CloneDiagnosticCollection(contracts.DiagnosticCollection{Groups: results.diagnostics}).Groups...)
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 {
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), "validation_incomplete_chunks": append([]int(nil), results.incomplete...)}}); err != nil {
return local, &laneRunError{stage: StageExtract, err: err}
}
if len(results.accepted) == 0 {
@@ -667,7 +655,6 @@ func mergeLaneOutput(dst *RunOutput, src RunOutput) error {
}
dst.NormalizeOutputs = append(dst.NormalizeOutputs, cloneSerializedOutputs(src.NormalizeOutputs)...)
dst.Rejected = append(dst.Rejected, cloneRejectedOutputs(src.Rejected)...)
dst.Warnings = append(dst.Warnings, cloneWarnings(src.Warnings)...)
appendDiagnosticGroups(dst, src.diagnosticGroups)
dst.CheckpointEvents = append(dst.CheckpointEvents, src.CheckpointEvents...)
dst.ValidationSummaries = append(dst.ValidationSummaries, cloneValidationSummaries(src.ValidationSummaries)...)