Persist diagnostics in reusable pipeline state

This commit is contained in:
2026-08-27 15:40:45 +00:00
parent 5175cb0722
commit 1f1967c8d2
18 changed files with 215 additions and 54 deletions

View File

@@ -279,6 +279,15 @@ type CheckpointArtifact struct {
ChunkRef source.SourceRef
Artifact contracts.SerializedArtifact
SchemaDigest string
Diagnostics []CheckpointDiagnostic
}
// CheckpointDiagnostic stores a producer-local diagnostic alongside a
// reusable artifact. The runner supplies the current run's origin when it
// promotes this value into a diagnostic group.
type CheckpointDiagnostic struct {
Diagnostic contracts.ProducerDiagnostic `json:"diagnostic"`
ValidatorKey string `json:"validator_key,omitempty"`
}
type ExtractCheckpoint struct {

View File

@@ -8,7 +8,7 @@ import (
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
)
const ChunkPlanSchemaVersion = "notarius.chunk-plan.v2"
const ChunkPlanSchemaVersion = "notarius.chunk-plan.v3"
type ChunkPlanProducer struct {
InputModule string `json:"input_module"`
@@ -19,13 +19,14 @@ type ChunkPlanProducer struct {
}
type ChunkPlanRecord struct {
SchemaVersion string `json:"schema_version"`
SourceDigest string `json:"source_digest"`
PlanDigest string `json:"plan_digest"`
Plan source.ChunkPlan `json:"plan"`
Producer ChunkPlanProducer `json:"producer"`
Warnings []contracts.Warning `json:"warnings,omitempty"`
CreatedAt time.Time `json:"created_at"`
SchemaVersion string `json:"schema_version"`
SourceDigest string `json:"source_digest"`
PlanDigest string `json:"plan_digest"`
Plan source.ChunkPlan `json:"plan"`
Producer ChunkPlanProducer `json:"producer"`
Warnings []contracts.Warning `json:"warnings,omitempty"`
Diagnostics []contracts.ProducerDiagnostic `json:"diagnostics,omitempty"`
CreatedAt time.Time `json:"created_at"`
}
type ChunkPlanStore interface {

View File

@@ -9,33 +9,53 @@ import (
)
func terminalDiagnosticGroups(terminal producerAttemptTerminal, origin contracts.DiagnosticOrigin, chunk *source.Chunk) ([]contracts.DiagnosticGroup, error) {
groups, err := promoteProducerDiagnostics(terminal.Diagnostics, origin, chunk)
if err != nil {
return nil, err
return promoteCheckpointDiagnostics(terminalCheckpointDiagnostics(terminal), origin, chunk)
}
func terminalCheckpointDiagnostics(terminal producerAttemptTerminal) []CheckpointDiagnostic {
diagnostics := make([]CheckpointDiagnostic, 0, len(terminal.Diagnostics))
for _, diagnostic := range terminal.Diagnostics {
diagnostics = append(diagnostics, CheckpointDiagnostic{Diagnostic: diagnostic})
}
for _, record := range terminal.Validation.Diagnostics() {
validatorOrigin := origin
validatorOrigin.ValidatorKey = record.validatorName
promoted, err := promoteProducerDiagnostics([]contracts.ProducerDiagnostic{record.diagnostic}, validatorOrigin, chunk)
if err != nil {
return nil, fmt.Errorf("validator %q diagnostic: %w", record.validatorName, err)
}
groups = append(groups, promoted...)
diagnostics = append(diagnostics, CheckpointDiagnostic{Diagnostic: record.diagnostic, ValidatorKey: record.validatorName})
}
if terminal.Action == producerTerminalIncompleteAccepted {
for _, record := range incompleteValidationDiagnostics(terminal.Validation) {
validatorOrigin := origin
validatorOrigin.ValidatorKey = record.validatorName
promoted, err := promoteProducerDiagnostics([]contracts.ProducerDiagnostic{record.diagnostic}, validatorOrigin, chunk)
if err != nil {
return nil, fmt.Errorf("validator %q incomplete diagnostic: %w", record.validatorName, err)
}
groups = append(groups, promoted...)
diagnostics = append(diagnostics, CheckpointDiagnostic{Diagnostic: record.diagnostic, ValidatorKey: record.validatorName})
}
}
return cloneCheckpointDiagnostics(diagnostics)
}
func promoteCheckpointDiagnostics(diagnostics []CheckpointDiagnostic, origin contracts.DiagnosticOrigin, chunk *source.Chunk) ([]contracts.DiagnosticGroup, error) {
groups := make([]contracts.DiagnosticGroup, 0, len(diagnostics))
for _, record := range diagnostics {
validatorOrigin := origin
validatorOrigin.ValidatorKey = record.ValidatorKey
promoted, err := promoteProducerDiagnostics([]contracts.ProducerDiagnostic{record.Diagnostic}, validatorOrigin, chunk)
if err != nil {
return nil, fmt.Errorf("checkpoint diagnostic: %w", err)
}
groups = append(groups, promoted...)
}
return groups, nil
}
func cloneCheckpointDiagnostics(diagnostics []CheckpointDiagnostic) []CheckpointDiagnostic {
if len(diagnostics) == 0 {
return nil
}
cloned := make([]CheckpointDiagnostic, len(diagnostics))
for index, diagnostic := range diagnostics {
cloned[index] = CheckpointDiagnostic{
Diagnostic: contracts.CloneProducerDiagnostics([]contracts.ProducerDiagnostic{diagnostic.Diagnostic})[0],
ValidatorKey: diagnostic.ValidatorKey,
}
}
return cloned
}
func promoteProducerDiagnostics(diagnostics []contracts.ProducerDiagnostic, origin contracts.DiagnosticOrigin, chunk *source.Chunk) ([]contracts.DiagnosticGroup, error) {
if err := contracts.ValidateProducerDiagnostics(diagnostics); err != nil {
return nil, err

View File

@@ -85,7 +85,12 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
result.warnings = append(cloneWarnings(record.Warnings), report.Warnings()...)
result.accepted = true
result.setValidation(report.Warnings(), nil, nil)
cachedTerminal := producerAttemptTerminal{Action: producerTerminalAccepted, Validation: report}
cachedTerminal := producerAttemptTerminal{Action: producerTerminalAccepted, Diagnostics: contracts.CloneProducerDiagnostics(record.Diagnostics), Validation: report}
diagnostics, diagnosticErr := terminalDiagnosticGroups(cachedTerminal, contracts.DiagnosticOrigin{Stage: contracts.DiagnosticOriginStageChunk, ModuleKey: chunker.Key()}, nil)
if diagnosticErr != nil {
return result, fmt.Errorf("promote reused chunk diagnostics: %w", diagnosticErr)
}
result.diagnostics = diagnostics
summary := validationSummary(cachedTerminal, StageChunk, "", "", chunker.Key(), "", 0)
result.validation = &summary
return result, nil
@@ -242,6 +247,7 @@ func (r *Runner) runChunkPlan(ctx context.Context, input RunInput, doc *source.S
return result, fmt.Errorf("clone chunk plan record for publication: %w", cloneErr)
}
record.Warnings = cloneWarnings(candidate.producerWarnings)
record.Diagnostics = contracts.CloneProducerDiagnostics(candidate.producerDiagnostics)
if err := input.ChunkPlans.Save(record); err != nil {
return result, fmt.Errorf("save chunk plan: %w", err)
}
@@ -319,6 +325,7 @@ func cloneChunkPlanRecord(record ChunkPlanRecord) (ChunkPlanRecord, error) {
}
record.Producer.Metadata = metadata
record.Warnings = cloneWarnings(record.Warnings)
record.Diagnostics = contracts.CloneProducerDiagnostics(record.Diagnostics)
return record, nil
}

View File

@@ -233,6 +233,37 @@ func TestRunnerChunkPlanHitUsesStoredProducerProvenance(t *testing.T) {
}
}
func TestRunnerReusesChunkPlanDiagnosticsWithCurrentValidatorDiagnostics(t *testing.T) {
prepared, plan := preparedTerminalDebugPipeline(t)
record := chunkPlanRecord(t, prepared, plan)
record.Diagnostics = []contracts.ProducerDiagnostic{{
Disposition: contracts.DiagnosticDispositionWarning,
Category: contracts.DiagnosticCategoryConfiguration,
ReasonCode: "stored_chunk_diagnostic",
OccurrenceCount: 1,
Samples: []contracts.DiagnosticSample{{Scope: "reference", Message: "Stored configuration signal."}},
}}
validator := &countingChunkValidator{result: contracts.ValidationResult{Approved: true, Diagnostics: []contracts.ProducerDiagnostic{{
Disposition: contracts.DiagnosticDispositionAdvisory,
Category: contracts.DiagnosticCategoryDataQuality,
ReasonCode: "current_validator_diagnostic",
OccurrenceCount: 1,
Samples: []contracts.DiagnosticSample{{Scope: "chunk", Message: "Current validator finding."}},
}}}}
prepared.chunkValidators.validators = []preparedValidator{{resolved: ResolvedValidator{Binding: Binding(validator.Name()), Target: ValidatorTargetChunk}, chunk: validator}}
output, err := New().Run(context.Background(), RunInput{Prepared: prepared, RawInput: []byte("input"), ChunkCacheMode: ChunkCacheAuto, ChunkPlans: &recordingChunkPlanStore{record: record, decision: ChunkPlanDecision{Status: ChunkPlanHit}}})
if err != nil {
t.Fatal(err)
}
if validator.calls != 1 || len(output.Diagnostics.Groups) != 2 {
t.Fatalf("validator calls = %d diagnostics = %#v", validator.calls, output.Diagnostics)
}
if output.Diagnostics.Groups[0].ReasonCode != "stored_chunk_diagnostic" || output.Diagnostics.Groups[1].ReasonCode != "current_validator_diagnostic" || output.Diagnostics.Groups[1].Origin.ValidatorKey != validator.Name() {
t.Fatalf("diagnostic groups = %#v", output.Diagnostics.Groups)
}
}
func TestRunnerProvidesAcceptedChunkMapToOutput(t *testing.T) {
for _, test := range []struct {
name string

View File

@@ -357,6 +357,11 @@ func hydrateRequiredLane(input RunInput, loader CheckpointLoader, doc *source.So
}
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)
}
appendDiagnosticGroups(&local, diagnostics)
local.NormalizeOutputs = append(local.NormalizeOutputs, contracts.SerializedOutput{
StepID: input.stepID,
LaneID: lane.ID,
@@ -421,6 +426,12 @@ func prepareLaneExtract(input RunInput, loader CheckpointLoader, doc *source.Sou
}
state.values = append(state.values, artifact)
state.serialized = append(state.serialized, cloneCheckpointArtifact(stored))
chunk := source.Chunk{ID: stored.ChunkID, Index: stored.ChunkIndex}
diagnostics, diagnosticErr := promoteCheckpointDiagnostics(stored.Diagnostics, contracts.DiagnosticOrigin{Stage: contracts.DiagnosticOriginStageExtract, StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Extract.Module}, &chunk)
if diagnosticErr != nil {
return nil, fmt.Errorf("promote reused extract diagnostics: %w", diagnosticErr)
}
state.diagnostics = append(state.diagnostics, diagnostics...)
}
state.warnings, state.rejected = cloneWarnings(cp.Warnings), cloneRejectedOutputs(cp.Rejected)
}
@@ -516,6 +527,7 @@ func (r *Runner) runExtractJob(ctx context.Context, input RunInput, doc *source.
return result
}
stored.ChunkID, stored.ChunkIndex, stored.ChunkRef = candidate.artifact.ChunkID, candidate.artifact.ChunkIndex, candidate.artifact.ChunkRef
stored.Diagnostics = terminalCheckpointDiagnostics(terminalResult)
if debugErr := candidate.terminal.record(payload, nil); debugErr != nil {
result.err = debugErr
return result

View File

@@ -158,6 +158,13 @@ func TestRunnerContinuesFromFreshAndReusedExtractResults(t *testing.T) {
return erasedTypedResult{
Value: typedValueForLane(0, request.Chunk.Index),
Warnings: []contracts.Warning{{Scope: "extract", ReasonCode: "observed", Message: "accepted extract"}},
Diagnostics: []contracts.ProducerDiagnostic{{
Disposition: contracts.DiagnosticDispositionObservation,
Category: contracts.DiagnosticCategoryNormalization,
ReasonCode: "accepted_extract_normalized",
OccurrenceCount: 1,
Samples: []contracts.DiagnosticSample{{Scope: "extract", Message: "accepted extract"}},
}},
}, nil
})
@@ -209,6 +216,9 @@ func TestRunnerContinuesFromFreshAndReusedExtractResults(t *testing.T) {
if !reflect.DeepEqual(reused.Warnings, fresh.Warnings) {
t.Fatalf("reused warnings = %#v, want fresh warnings %#v", reused.Warnings, fresh.Warnings)
}
if !reflect.DeepEqual(reused.Diagnostics, fresh.Diagnostics) {
t.Fatalf("reused diagnostics = %#v, want fresh diagnostics %#v", reused.Diagnostics, fresh.Diagnostics)
}
}
func TestRunnerPromotesOnlyAcceptedExtractRetryWarnings(t *testing.T) {

View File

@@ -14,11 +14,12 @@ import (
)
type terminalChunker struct {
key string
plan source.ChunkPlan
warnings []contracts.Warning
err error
calls *int
key string
plan source.ChunkPlan
warnings []contracts.Warning
diagnostics []contracts.ProducerDiagnostic
err error
calls *int
}
func (c terminalChunker) Key() string { return c.key }
@@ -29,7 +30,7 @@ func (c terminalChunker) Plan(context.Context, contracts.ChunkRequest) (contract
if c.calls != nil {
(*c.calls)++
}
return contracts.ChunkPlanResult{Plan: source.CloneChunkPlan(c.plan), Warnings: cloneWarnings(c.warnings)}, c.err
return contracts.ChunkPlanResult{Plan: source.CloneChunkPlan(c.plan), Warnings: cloneWarnings(c.warnings), Diagnostics: contracts.CloneProducerDiagnostics(c.diagnostics)}, c.err
}
type terminalChunkValidator struct {

View File

@@ -37,6 +37,7 @@ func recordNormalize(recorder CheckpointRecorder, stepID, laneID, moduleKey stri
}
func cloneCheckpointArtifact(output CheckpointArtifact) CheckpointArtifact {
output.Artifact = contracts.CloneSerializedArtifact(output.Artifact)
output.Diagnostics = cloneCheckpointDiagnostics(output.Diagnostics)
return output
}
@@ -270,6 +271,11 @@ func (r *Runner) runMergeStage(ctx context.Context, input RunInput, checkpoints
serializedMerge = mergeResolution.artifacts[0]
mergeWarnings = cloneWarnings(mergeCP.Warnings)
output.Warnings = append(output.Warnings, mergeWarnings...)
diagnostics, diagnosticErr := promoteCheckpointDiagnostics(serializedMerge.Diagnostics, contracts.DiagnosticOrigin{Stage: contracts.DiagnosticOriginStageMerge, StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Merge.Module}, nil)
if diagnosticErr != nil {
return stageResult, fmt.Errorf("promote reused merge diagnostics: %w", diagnosticErr)
}
appendDiagnosticGroups(output, diagnostics)
} else {
if stageResult.reuseEligible {
if err := checkpointMergeRunning(checkpoints, input.stepID, lane.ID, lane.Merge.Module, mergeDeps); err != nil {
@@ -358,6 +364,7 @@ func (r *Runner) runMergeStage(ctx context.Context, input RunInput, checkpoints
}
return stageResult, candidate.terminal.record(payload, attemptErr)
}
stored.Diagnostics = terminalCheckpointDiagnostics(terminalResult)
if debugErr := candidate.terminal.record(payload, nil); debugErr != nil {
if stageResult.reuseEligible {
_ = checkpointMergeFailed(checkpoints, input.stepID, lane.ID, lane.Merge.Module, mergeDeps, debugErr)
@@ -420,6 +427,11 @@ func (r *Runner) runNormalizeStage(ctx context.Context, input RunInput, checkpoi
serializedNormalize = normalizeResolution.artifacts[0]
normalizeWarnings = cloneWarnings(normalizeCP.Warnings)
output.Warnings = append(output.Warnings, normalizeWarnings...)
diagnostics, diagnosticErr := promoteCheckpointDiagnostics(serializedNormalize.Diagnostics, contracts.DiagnosticOrigin{Stage: contracts.DiagnosticOriginStageNormalize, StepID: input.stepID, LaneID: lane.ID, ModuleKey: lane.Normalize.Module}, nil)
if diagnosticErr != nil {
return stageResult, fmt.Errorf("promote reused normalize diagnostics: %w", diagnosticErr)
}
appendDiagnosticGroups(output, diagnostics)
} else {
if stageResult.reuseEligible {
if err := checkpointNormalizeRunning(checkpoints, input.stepID, lane.ID, lane.Normalize.Module, normalizeDeps); err != nil {
@@ -528,6 +540,7 @@ func (r *Runner) runNormalizeStage(ctx context.Context, input RunInput, checkpoi
}
return stageResult, candidate.terminal.record(payload, attemptErr)
}
stored.Diagnostics = terminalCheckpointDiagnostics(terminalResult)
if debugErr := candidate.terminal.record(payload, nil); debugErr != nil {
if stageResult.reuseEligible {
_ = checkpointNormalizeFailed(checkpoints, input.stepID, lane.ID, lane.Normalize.Module, normalizeDeps, debugErr)