Cleanup and complete the validator refactor

This commit is contained in:
2026-07-07 19:36:48 -05:00
parent fc8e03f98c
commit c5f2b14ff4
6 changed files with 58 additions and 89 deletions

View File

@@ -127,12 +127,12 @@ func (r *Runner) Run(ctx context.Context, input RunInput) (output RunOutput, err
if err != nil {
return false, nil, fmt.Errorf("validate chunks from chunker %q: %w", chunker.Key(), err)
}
rejection, err := r.validateChunksRaw(ctx, doc, chunker.Key(), chunks, sourceInput, sessionID, input.Pipeline.ChunkReferences.ReferenceSet, input.LLMClient, input.Metadata, input.Pipeline.ValidatorChains, attempt)
validationWarnings, rejection, err := r.validateChunksRaw(ctx, doc, chunker.Key(), chunks, sourceInput, sessionID, input.Pipeline.ChunkReferences.ReferenceSet, input.LLMClient, input.Metadata, input.Pipeline.ValidatorChains, attempt)
if err != nil || rejection != nil {
return false, rejection, err
}
canonicalChunks = chunks
chunkWarnings = cloneWarnings(chunkResult.Warnings)
chunkWarnings = append(cloneWarnings(chunkResult.Warnings), validationWarnings...)
return true, nil, nil
})
if err != nil {
@@ -455,39 +455,21 @@ func runWithRetry(ctx context.Context, retries int, run func(attempt int) (bool,
return false, lastRejection, nil
}
func (r *Runner) validateChunksRaw(ctx context.Context, doc *source.SourceDocument, moduleKey string, chunks []contracts.SourceChunk, sourceInput contracts.LLMInputMaterial, sessionID string, references contracts.ReferenceSet, llmClient contracts.StructuredLLMClient, metadata map[string]any, chains []ResolvedValidatorChain, attempt int) (*contracts.RejectedOutput, error) {
for index := range chunks {
chunk := chunks[index]
_, rejection, err := r.validateRaw(ctx, rawValidationTarget{
stage: StageChunk,
moduleKey: moduleKey,
source: doc,
sourceID: doc.ID,
sourceInput: sourceInput.Clone(),
sessionID: sessionID,
references: references,
llmClient: llmClient,
chunkID: chunk.ID,
chunkIndex: chunk.Index,
chunk: &chunk,
chunks: chunks,
payload: contracts.RawPayload{
Content: append([]byte(nil), chunk.Content...),
MediaType: chunk.MediaType,
Metadata: cloneMetadata(chunk.Metadata),
},
metadata: metadata,
chains: chains,
attempt: attempt,
})
if err != nil {
return nil, err
}
if rejection != nil {
return rejection, nil
}
}
return nil, nil
func (r *Runner) validateChunksRaw(ctx context.Context, doc *source.SourceDocument, moduleKey string, chunks []contracts.SourceChunk, sourceInput contracts.LLMInputMaterial, sessionID string, references contracts.ReferenceSet, llmClient contracts.StructuredLLMClient, metadata map[string]any, chains []ResolvedValidatorChain, attempt int) ([]contracts.Warning, *contracts.RejectedOutput, error) {
return r.validateRaw(ctx, rawValidationTarget{
stage: StageChunk,
moduleKey: moduleKey,
source: doc,
sourceID: doc.ID,
sourceInput: sourceInput.Clone(),
sessionID: sessionID,
references: references,
llmClient: llmClient,
chunks: chunks,
metadata: metadata,
chains: chains,
attempt: attempt,
})
}
func (r *Runner) validateRaw(ctx context.Context, target rawValidationTarget) ([]contracts.Warning, *contracts.RejectedOutput, error) {

View File

@@ -801,8 +801,8 @@ func TestRunPassesValidationRequestContextToValidators(t *testing.T) {
t.Fatalf("Run() error = %v, want nil", err)
}
if len(chunkValidator.requests) != 2 {
t.Fatalf("chunk validator requests = %d, want one per chunk", len(chunkValidator.requests))
if len(chunkValidator.requests) != 1 {
t.Fatalf("chunk validator requests = %d, want one collection request", len(chunkValidator.requests))
}
chunkReq := chunkValidator.requests[0]
if chunkReq.Stage != string(StageChunk) || chunkReq.ModuleKey != "chunk" || chunkReq.SourceID != "source-1" || chunkReq.SessionID != "session-123" {
@@ -811,8 +811,12 @@ func TestRunPassesValidationRequestContextToValidators(t *testing.T) {
if chunkReq.LLMClient == nil || string(chunkReq.SourceInput.Content) != string(rawInput) {
t.Fatalf("chunk validation source/client = %#v, want full source input and LLM client", chunkReq.SourceInput)
}
if chunkReq.Chunk == nil || chunkReq.Chunk.ID != "chunk-0" || len(chunkReq.Chunks) != 2 || string(chunkReq.Payload.Content) != string(chunkReq.Chunk.Content) {
t.Fatalf("chunk validation chunk fields = %#v chunks=%#v payload=%s, want chunk payload and all chunks", chunkReq.Chunk, chunkReq.Chunks, chunkReq.Payload.Content)
if chunkReq.Chunk != nil || chunkReq.ChunkID != "" || len(chunkReq.Chunks) != 2 || chunkReq.Chunks[0].ID != "chunk-0" || chunkReq.Chunks[1].ID != "chunk-1" {
t.Fatalf("chunk validation chunk fields = chunk=%#v chunk_id=%q chunks=%#v, want whole chunk collection", chunkReq.Chunk, chunkReq.ChunkID, chunkReq.Chunks)
}
chunkReq.Chunks[0].Content[0] = 'X'
if got := string(modules.chunker.chunks[0].Content); got == string(chunkReq.Chunks[0].Content) {
t.Fatalf("chunk validation chunks alias module output content = %q", got)
}
if item := chunkReq.References.Slots["scene_guide"].Items[0]; string(item.Content) != "chunk reference text" {
t.Fatalf("chunk validation references = %#v, want chunk references", chunkReq.References)
@@ -860,7 +864,6 @@ func TestRunPassesValidationRequestContextToValidators(t *testing.T) {
t.Fatalf("normalize validation references = %#v, want normalize references", normalizeReq.References)
}
chunkReq.Payload.Content[0] = 'X'
chunkReq.Chunks[0].Content[0] = 'Y'
if got := string(modules.chunker.chunks[0].Content); got != `{"units":[{"id":1,"kind":"unit","text":"Source unit."}]}` {
t.Fatalf("validator request mutated original chunk content: %q", got)
@@ -1211,6 +1214,27 @@ func TestRunCollectsStageWarnings(t *testing.T) {
}
}
func TestRunCollectsChunkValidatorWarnings(t *testing.T) {
modules := defaultRunnerModules()
modules.chunker.chunks = []contracts.SourceChunk{sourceChunkWithID("chunk-0", 0)}
validator := &runnerChainValidator{
name: "chain-chunk",
warnings: []contracts.Warning{{ReasonCode: "chunk-validator-warning", Message: "chunk validator warning"}},
}
modules.validators[validator.name] = validator
pipeline := resolvedPipeline()
setResolvedValidatorChain(t, &pipeline, StageChunk, "", "chunk", resolvedValidatorForTest(validator))
output, err := New(newRunnerRegistries(t, modules)).Run(context.Background(), RunInput{Pipeline: pipeline})
if err != nil {
t.Fatalf("Run() error = %v, want nil", err)
}
if got := warningReasons(output.Warnings); !reflect.DeepEqual(got, []string{"chunk-validator-warning"}) {
t.Fatalf("warning reasons = %#v, want chunk validator warning", got)
}
}
func TestRunOutputEncoderReceivesManifestAndRawOutputs(t *testing.T) {
modules := defaultRunnerModules()