diff --git a/docs/internal/modules.md b/docs/internal/modules.md index e85c916..fc74f70 100644 --- a/docs/internal/modules.md +++ b/docs/internal/modules.md @@ -57,7 +57,8 @@ rules are defined in the The generic chunker validates the source document, walks units in configured windows, clones each selected unit, and emits deterministic ordered chunk IDs. Overlap changes the next window start but never reorders units. It records the -first and last unit and unit count in chunk metadata. +first and last unit and unit count in chunk metadata, and derives the chunk's +canonical source reference from those unit references. The accepted options and defaults are defined in [Configuration](../config.md#implemented-production-modules). Generic @@ -68,7 +69,7 @@ framework validation canonicalizes the returned unit slices before extraction. The scene chunker prepares a structured Scriptorium request from the full transcript, session, and optional D&D reference inputs. It validates the model's scene boundaries against source-unit IDs and converts them into deterministic -chunks. +chunks with canonical source references spanning each scene's units. Scene validation requires sequential, contiguous, non-overlapping coverage from the first source unit through the last. Each chunk contains JSON scene content diff --git a/docs/internal/overview.md b/docs/internal/overview.md index 5bab987..bbb8583 100644 --- a/docs/internal/overview.md +++ b/docs/internal/overview.md @@ -31,7 +31,7 @@ a sorted set of artifact lanes before the runner constructs any stage module. | `internal/core/artifacts` | Run-manifest and provenance models. | | `internal/core/config` | Defaults, YAML parsing, environment overrides, validation, redaction, and effective pipeline resolution. | | `internal/core/diagnostics` | Scoped run directories, diagnostics writers, atomic writes, and retention decisions. | -| `internal/core/source` | Generic source documents, units, references, lookup, and validation. | +| `internal/core/source` | Generic source documents, units, chunks, canonical references, lookup, validation, and deterministic source digests. | | `internal/core/workspace` | Effective workspace settings, confined paths and writes, checkpoint identity, and checkpoint manifest models. | ## Framework Packages diff --git a/docs/internal/pipeline.md b/docs/internal/pipeline.md index 6623a65..a03ea54 100644 --- a/docs/internal/pipeline.md +++ b/docs/internal/pipeline.md @@ -77,7 +77,8 @@ normalize requests retain access to the original source material. Source validation requires every unit to carry a canonical self-reference to its containing document and its own unit ID. Explicit clone, checkpoint, and debug boundaries retain that reference, and the canonical source digest covers -it deterministically. +it deterministically. Chunks use the same source model and carry one canonical +reference spanning the first selected unit through the last. `pipeline.RunOutput` carries the run manifest, accepted normalized results, rejected results, warnings, checkpoint events, and logical files returned by the @@ -115,13 +116,15 @@ whose results are accepted and used. ## Chunk Canonicalization Before lane execution, generic validation requires unique chunk IDs, matching -source identity, indexes matching returned order, valid ordered boundaries, +source identity, indexes matching returned order, a valid canonical reference, non-empty content and media type, and at least one valid source unit per chunk. -Units may not repeat inside a chunk and must preserve source-document order. +Units may not repeat inside a chunk and must form a contiguous range in +source-document order. The chunk reference must exactly match the source and +the first and last unit references. The runner then rebuilds each chunk's unit slice from the source document by -unit ID. It preserves the module-owned boundaries, content, media type, and -cloned metadata. The framework permits gaps and overlap between separate +unit ID. It preserves the canonical reference, content, media type, and cloned +metadata. The framework permits gaps and overlap between separate chunks; stricter coverage policy belongs to the chunk implementation. ## Validation And Retries diff --git a/docs/operations.md b/docs/operations.md index 0b03564..a59e66f 100644 --- a/docs/operations.md +++ b/docs/operations.md @@ -91,6 +91,14 @@ digests match the current invocation. Changes to input bytes, the resolved pipeline, selected lanes, the runtime LLM profile override, or bound reference content invalidate reuse. +Current checkpoint manifests use workspace schema `notarius.workspace.v2`. +Manifests written with `notarius.workspace.v1` are incompatible because their +chunk provenance has an older shape. On the first explicit resume after an +upgrade, each affected checkpoint is treated as a reuse miss and its workflow +step executes normally. The compatibility check does not migrate or delete the +v1 files; when checkpoint writing is enabled, normal execution refreshes the +affected checkpoint files in the current schema. + Runs do not reuse checkpoints unless explicitly requested. Without reuse, the workflow executes normally and refreshes checkpoint files when checkpointing is enabled. diff --git a/internal/cli/compatibility_test.go b/internal/cli/compatibility_test.go index 2fa2de7..c3fdf47 100644 --- a/internal/cli/compatibility_test.go +++ b/internal/cli/compatibility_test.go @@ -330,8 +330,8 @@ func TestProductionLLMCallersShareScheduledClient(t *testing.T) { {ID: 2, Kind: "segment", Text: "The spell takes effect.", Ref: source.SourceRef{SourceID: "session-alpha", StartUnitID: 2, EndUnitID: 2}}, }, } - chunk := contracts.SourceChunk{ - ID: "session-alpha:chunk:0", SourceID: doc.ID, Index: 0, StartUnitID: 1, EndUnitID: 2, + chunk := source.Chunk{ + ID: "session-alpha:chunk:0", SourceID: doc.ID, Index: 0, Ref: source.SourceRef{SourceID: doc.ID, StartUnitID: 1, EndUnitID: 2}, Content: []byte(`{"scene":"Aria casts Cure Wounds."}`), MediaType: "application/json", Units: append([]source.SourceUnit(nil), doc.Units...), } diff --git a/internal/cli/run_test.go b/internal/cli/run_test.go index 6c900cf..c1d9909 100644 --- a/internal/cli/run_test.go +++ b/internal/cli/run_test.go @@ -3611,8 +3611,8 @@ func (fakeRunChunker) ReferenceSlots() []contracts.ReferenceSlot { func (fakeRunChunker) Chunk(ctx context.Context, req contracts.ChunkRequest) (contracts.ChunkResult, error) { return contracts.ChunkResult{ - Chunks: []contracts.SourceChunk{ - {ID: "chunk-1", SourceID: req.Source.ID, Index: 0, StartUnitID: 1, EndUnitID: 1, Content: []byte(`{"units":[1]}`), MediaType: "application/json", Units: req.Source.Units}, + Chunks: []source.Chunk{ + {ID: "chunk-1", SourceID: req.Source.ID, Index: 0, Ref: source.SourceRef{SourceID: req.Source.ID, StartUnitID: 1, EndUnitID: 1}, Content: []byte(`{"units":[1]}`), MediaType: "application/json", Units: req.Source.Units}, }, }, nil } diff --git a/internal/core/source/digest.go b/internal/core/source/digest.go index cd7df81..80a637a 100644 --- a/internal/core/source/digest.go +++ b/internal/core/source/digest.go @@ -33,3 +33,33 @@ func DigestDocument(doc *SourceDocument) (string, error) { sum := sha256.Sum256(encoded) return "sha256:" + hex.EncodeToString(sum[:]), nil } + +// DigestChunk returns a deterministic digest of a chunk, including its source +// provenance, content, units, and metadata. +func DigestChunk(chunk Chunk) (string, error) { + payload := struct { + ID string `json:"id"` + SourceID string `json:"source_id"` + Index int `json:"index"` + Ref SourceRef `json:"ref"` + Content []byte `json:"content"` + MediaType string `json:"media_type"` + Units []SourceUnit `json:"units"` + Metadata map[string]any `json:"metadata,omitempty"` + }{ + ID: chunk.ID, + SourceID: chunk.SourceID, + Index: chunk.Index, + Ref: chunk.Ref, + Content: chunk.Content, + MediaType: chunk.MediaType, + Units: chunk.Units, + Metadata: chunk.Metadata, + } + encoded, err := json.Marshal(payload) + if err != nil { + return "", fmt.Errorf("encode source chunk for digest: %w", err) + } + sum := sha256.Sum256(encoded) + return "sha256:" + hex.EncodeToString(sum[:]), nil +} diff --git a/internal/core/source/source.go b/internal/core/source/source.go index 9d87125..e39cd87 100644 --- a/internal/core/source/source.go +++ b/internal/core/source/source.go @@ -22,3 +22,14 @@ type SourceRef struct { StartUnitID int `json:"start_unit_id"` EndUnitID int `json:"end_unit_id"` } + +type Chunk struct { + ID string `json:"id"` + SourceID string `json:"source_id"` + Index int `json:"index"` + Ref SourceRef `json:"ref"` + Content []byte `json:"-"` + MediaType string `json:"media_type"` + Units []SourceUnit `json:"units"` + Metadata map[string]any `json:"metadata,omitempty"` +} diff --git a/internal/core/source/source_test.go b/internal/core/source/source_test.go index 0ece58e..8ee0d3f 100644 --- a/internal/core/source/source_test.go +++ b/internal/core/source/source_test.go @@ -222,6 +222,42 @@ func TestDigestDocumentIsDeterministicAndIncludesUnitReference(t *testing.T) { } } +func TestDigestChunkIsDeterministicAndIncludesReference(t *testing.T) { + doc := validDocument() + chunk := Chunk{ + ID: "chunk-1", + SourceID: doc.ID, + Index: 0, + Ref: SourceRef{SourceID: doc.ID, StartUnitID: 1, EndUnitID: 2}, + Content: []byte("chunk content"), + MediaType: "text/plain", + Units: doc.Units, + Metadata: map[string]any{"second": "value", "first": true}, + } + first, err := DigestChunk(chunk) + if err != nil { + t.Fatalf("DigestChunk() error = %v, want nil", err) + } + + chunk.Metadata = map[string]any{"first": true, "second": "value"} + second, err := DigestChunk(chunk) + if err != nil { + t.Fatalf("DigestChunk(reordered metadata) error = %v, want nil", err) + } + if first != second { + t.Fatalf("digests = %q and %q, want deterministic map ordering", first, second) + } + + chunk.Ref.EndUnitID = 1 + changed, err := DigestChunk(chunk) + if err != nil { + t.Fatalf("DigestChunk(changed ref) error = %v, want nil", err) + } + if first == changed { + t.Fatalf("digest = %q after reference change, want different digest", changed) + } +} + func TestValidateRefValid(t *testing.T) { doc := validDocument() ref := SourceRef{ diff --git a/internal/core/workspace/manifest.go b/internal/core/workspace/manifest.go index 84cdc12..d3cbb39 100644 --- a/internal/core/workspace/manifest.go +++ b/internal/core/workspace/manifest.go @@ -2,7 +2,10 @@ package workspace import "time" -const WorkspaceSchemaVersion = "notarius.workspace.v1" +const ( + WorkspaceSchemaVersion = "notarius.workspace.v2" + WorkspaceSchemaVersionV1 = "notarius.workspace.v1" +) type StageName string diff --git a/internal/core/workspace/manifest_test.go b/internal/core/workspace/manifest_test.go index 2990b7f..ab80841 100644 --- a/internal/core/workspace/manifest_test.go +++ b/internal/core/workspace/manifest_test.go @@ -7,6 +7,12 @@ import ( ) func TestStageManifestDefaults(t *testing.T) { + if WorkspaceSchemaVersion != "notarius.workspace.v2" { + t.Fatalf("current schema version = %q, want notarius.workspace.v2", WorkspaceSchemaVersion) + } + if WorkspaceSchemaVersionV1 != "notarius.workspace.v1" { + t.Fatalf("legacy schema version = %q, want notarius.workspace.v1", WorkspaceSchemaVersionV1) + } manifest := NewStageManifest(StageExtract, StatusRunning) if manifest.WorkspaceSchemaVersion != WorkspaceSchemaVersion { diff --git a/internal/framework/checkpoint/loader.go b/internal/framework/checkpoint/loader.go index 026bb49..a4d8a64 100644 --- a/internal/framework/checkpoint/loader.go +++ b/internal/framework/checkpoint/loader.go @@ -78,7 +78,11 @@ func (l *WorkspaceLoader) Chunk(moduleKey string, sourceDigest string) (pipeline if len(chunks) == 0 { return pipeline.ChunkCheckpoint{}, invalidDecision("chunk checkpoint payload has no chunks") } - if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), chunkOutputDigests(chunks)) { + outputDigests, err := chunkOutputDigests(chunks) + if err != nil { + return pipeline.ChunkCheckpoint{}, invalidDecision("chunk checkpoint output cannot be digested: %v", err) + } + if !fingerprintsEqual(coreworkspaceToPipelineFingerprints(manifest.OutputDigests), outputDigests) { return pipeline.ChunkCheckpoint{}, invalidDecision("chunk checkpoint output digests do not match payload") } return pipeline.ChunkCheckpoint{Chunks: chunks, Warnings: cloneWarnings(payload.Warnings)}, reusedDecision() @@ -180,6 +184,9 @@ func (l *WorkspaceLoader) validateManifest(manifest coreworkspace.StageManifest, } func (l *WorkspaceLoader) validateLaneManifest(manifest coreworkspace.StageManifest, stage coreworkspace.StageName, laneID string, moduleKey string, dependencies []pipeline.CheckpointFingerprint, statuses ...coreworkspace.StageStatus) pipeline.CheckpointDecision { + if manifest.WorkspaceSchemaVersion == coreworkspace.WorkspaceSchemaVersionV1 { + return invalidDecision("checkpoint workspace schema version %q is incompatible with %q and must be recomputed", manifest.WorkspaceSchemaVersion, coreworkspace.WorkspaceSchemaVersion) + } if manifest.WorkspaceSchemaVersion != coreworkspace.WorkspaceSchemaVersion { return invalidDecision("checkpoint workspace schema version %q is not supported", manifest.WorkspaceSchemaVersion) } @@ -211,26 +218,25 @@ func (l *WorkspaceLoader) validateLaneManifest(manifest coreworkspace.StageManif return reusedDecision() } -func sourceChunksFromEnvelope(values []chunkEnvelope) ([]contracts.SourceChunk, error) { +func sourceChunksFromEnvelope(values []chunkEnvelope) ([]source.Chunk, error) { if len(values) == 0 { return nil, nil } - out := make([]contracts.SourceChunk, 0, len(values)) + out := make([]source.Chunk, 0, len(values)) for _, value := range values { content, err := contentFromEnvelope(value.Content) if err != nil { return nil, err } - out = append(out, contracts.SourceChunk{ - ID: value.ID, - SourceID: value.SourceID, - Index: value.Index, - StartUnitID: value.StartUnitID, - EndUnitID: value.EndUnitID, - Content: content, - MediaType: value.Content.MediaType, - Units: cloneSourceUnits(value.Units), - Metadata: cloneMetadata(value.Metadata), + out = append(out, source.Chunk{ + ID: value.ID, + SourceID: value.SourceID, + Index: value.Index, + Ref: value.Ref, + Content: content, + MediaType: value.Content.MediaType, + Units: cloneSourceUnits(value.Units), + Metadata: cloneMetadata(value.Metadata), }) } return out, nil diff --git a/internal/framework/checkpoint/recorder.go b/internal/framework/checkpoint/recorder.go index bebcdd9..ae2e661 100644 --- a/internal/framework/checkpoint/recorder.go +++ b/internal/framework/checkpoint/recorder.go @@ -73,7 +73,11 @@ func (r *WorkspaceRecorder) ChunkRunning(moduleKey string, sourceDigest string) return r.writeManifest("chunk/manifest.json", coreworkspace.ChunkManifest{StageManifest: manifest}) } -func (r *WorkspaceRecorder) ChunkSucceeded(moduleKey string, sourceDigest string, chunks []contracts.SourceChunk, warnings []contracts.Warning) error { +func (r *WorkspaceRecorder) ChunkSucceeded(moduleKey string, sourceDigest string, chunks []source.Chunk, warnings []contracts.Warning) error { + outputDigests, err := chunkOutputDigests(chunks) + if err != nil { + return fmt.Errorf("digest chunk checkpoint output: %w", err) + } payload := chunksEnvelope{Chunks: chunkEnvelopes(chunks), Warnings: cloneWarnings(warnings)} if err := r.writePayload("chunk/chunks.json", payload); err != nil { return err @@ -81,7 +85,7 @@ func (r *WorkspaceRecorder) ChunkSucceeded(moduleKey string, sourceDigest string manifest := r.newStageManifest(coreworkspace.StageChunk, coreworkspace.StatusSucceeded) manifest.ModuleKey = moduleKey manifest.DependencyFingerprints = workspaceFingerprints(digestFingerprints("source_document", sourceDigest)) - manifest.OutputDigests = workspaceFingerprints(chunkOutputDigests(chunks)) + manifest.OutputDigests = workspaceFingerprints(outputDigests) manifest.ValidationStatus = validationStatusString(warnings, nil) manifest.CompletedAt = timePtr(r.timestamp()) return r.writeManifest("chunk/manifest.json", coreworkspace.ChunkManifest{ @@ -260,14 +264,13 @@ type chunksEnvelope struct { } type chunkEnvelope struct { - ID string `json:"id"` - SourceID string `json:"source_id"` - Index int `json:"index"` - StartUnitID int `json:"start_unit_id"` - EndUnitID int `json:"end_unit_id"` - Content binaryEnvelope `json:"content"` - Units []source.SourceUnit `json:"units,omitempty"` - Metadata map[string]any `json:"metadata,omitempty"` + ID string `json:"id"` + SourceID string `json:"source_id"` + Index int `json:"index"` + Ref source.SourceRef `json:"ref"` + Content binaryEnvelope `json:"content"` + Units []source.SourceUnit `json:"units,omitempty"` + Metadata map[string]any `json:"metadata,omitempty"` } type extractOutputsEnvelope struct { @@ -320,21 +323,20 @@ type binaryEnvelope struct { Warnings []contracts.Warning `json:"warnings,omitempty"` } -func chunkEnvelopes(chunks []contracts.SourceChunk) []chunkEnvelope { +func chunkEnvelopes(chunks []source.Chunk) []chunkEnvelope { if len(chunks) == 0 { return nil } out := make([]chunkEnvelope, 0, len(chunks)) for _, chunk := range chunks { out = append(out, chunkEnvelope{ - ID: chunk.ID, - SourceID: chunk.SourceID, - Index: chunk.Index, - StartUnitID: chunk.StartUnitID, - EndUnitID: chunk.EndUnitID, - Content: binaryEnvelopeFromContent(chunk.Content, chunk.MediaType, chunk.Metadata, nil), - Units: cloneSourceUnits(chunk.Units), - Metadata: cloneMetadata(chunk.Metadata), + ID: chunk.ID, + SourceID: chunk.SourceID, + Index: chunk.Index, + Ref: chunk.Ref, + Content: binaryEnvelopeFromContent(chunk.Content, chunk.MediaType, chunk.Metadata, nil), + Units: cloneSourceUnits(chunk.Units), + Metadata: cloneMetadata(chunk.Metadata), }) } return out @@ -468,15 +470,19 @@ func extractPayloads(outputs []contracts.ExtractOutput) []contracts.RawPayload { return payloads } -func chunkOutputDigests(chunks []contracts.SourceChunk) []pipeline.CheckpointFingerprint { +func chunkOutputDigests(chunks []source.Chunk) ([]pipeline.CheckpointFingerprint, error) { values := make([]pipeline.CheckpointFingerprint, 0, len(chunks)) for _, chunk := range chunks { + digest, err := source.DigestChunk(chunk) + if err != nil { + return nil, fmt.Errorf("chunk %q: %w", chunk.ID, err) + } values = append(values, pipeline.CheckpointFingerprint{ Name: chunk.ID, - Value: contentDigest(chunk.Content), + Value: digest, }) } - return normalizeFingerprints(values) + return normalizeFingerprints(values), nil } func digestFingerprints(name string, digest string) []pipeline.CheckpointFingerprint { diff --git a/internal/framework/checkpoint/recorder_test.go b/internal/framework/checkpoint/recorder_test.go index bf65ed5..089858a 100644 --- a/internal/framework/checkpoint/recorder_test.go +++ b/internal/framework/checkpoint/recorder_test.go @@ -24,16 +24,15 @@ func TestWorkspaceRecorderWritesSuccessfulCheckpointFiles(t *testing.T) { Digest: "sha256:source", Units: []source.SourceUnit{{ID: 1, Kind: "line", Text: "hello", Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}}}, } - chunks := []contracts.SourceChunk{ + chunks := []source.Chunk{ { - ID: "chunk-1", - SourceID: "source-1", - Index: 0, - StartUnitID: 1, - EndUnitID: 1, - Content: []byte("chunk content"), - MediaType: "text/plain", - Units: doc.Units, + ID: "chunk-1", + SourceID: "source-1", + Index: 0, + Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, + Content: []byte("chunk content"), + MediaType: "text/plain", + Units: doc.Units, }, } @@ -91,16 +90,15 @@ func TestWorkspaceLoaderReusesSuccessfulCheckpointFiles(t *testing.T) { Digest: "sha256:source", Units: []source.SourceUnit{{ID: 1, Kind: "line", Text: "hello", Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}}}, } - chunks := []contracts.SourceChunk{ + chunks := []source.Chunk{ { - ID: "chunk-1", - SourceID: "source-1", - Index: 0, - StartUnitID: 1, - EndUnitID: 1, - Content: []byte("chunk content"), - MediaType: "text/plain", - Units: doc.Units, + ID: "chunk-1", + SourceID: "source-1", + Index: 0, + Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, + Content: []byte("chunk content"), + MediaType: "text/plain", + Units: doc.Units, }, } extractOutput := contracts.ExtractOutput{ @@ -162,6 +160,9 @@ func TestWorkspaceLoaderReusesSuccessfulCheckpointFiles(t *testing.T) { if !decision.Reused || len(chunkCheckpoint.Chunks) != 1 || string(chunkCheckpoint.Chunks[0].Content) != "chunk content" { t.Fatalf("chunk decision = %#v checkpoint=%#v, want reused", decision, chunkCheckpoint) } + if got, want := chunkCheckpoint.Chunks[0].Ref, chunks[0].Ref; got != want { + t.Fatalf("checkpoint chunk ref = %#v, want %#v", got, want) + } extractCheckpoint, decision := loader.Extract("spells", "dnd/spells", extractDeps) if !decision.Reused || len(extractCheckpoint.Outputs) != 1 || string(extractCheckpoint.Outputs[0].Payload.Content) != `{"spell":"cure wounds"}` { t.Fatalf("extract decision = %#v checkpoint=%#v, want reused", decision, extractCheckpoint) @@ -187,7 +188,7 @@ func TestWorkspaceLoaderInvalidatesMissingCorruptAndMismatchedCheckpoints(t *tes t.Run("dependency mismatch", func(t *testing.T) { root := t.TempDir() recorder := newTestRecorder(t, root) - chunks := []contracts.SourceChunk{{ + chunks := []source.Chunk{{ ID: "chunk-1", SourceID: "source-1", Content: []byte("chunk content"), @@ -202,10 +203,41 @@ func TestWorkspaceLoaderInvalidatesMissingCorruptAndMismatchedCheckpoints(t *tes } }) + t.Run("incompatible workspace schema remains untouched", func(t *testing.T) { + root := t.TempDir() + recorder := newTestRecorder(t, root) + doc := &source.SourceDocument{ + ID: "source-1", Kind: "document", Format: "text/plain", Digest: "sha256:source", + Units: []source.SourceUnit{{ID: 1, Kind: "line", Text: "hello", Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}}}, + } + if err := recorder.SourceSucceeded("seriatim", doc); err != nil { + t.Fatalf("SourceSucceeded: %v", err) + } + manifestPath := filepath.Join(root, "source", "manifest.json") + manifest := strings.Replace(string(readFile(t, manifestPath)), coreworkspace.WorkspaceSchemaVersion, coreworkspace.WorkspaceSchemaVersionV1, 1) + if err := os.WriteFile(manifestPath, []byte(manifest), 0o644); err != nil { + t.Fatalf("write legacy manifest: %v", err) + } + beforeManifest := readFile(t, manifestPath) + payloadPath := filepath.Join(root, "source", "source-document.json") + beforePayload := readFile(t, payloadPath) + + loader := &WorkspaceLoader{root: root} + if _, decision := loader.Source("seriatim"); decision.Reused || !strings.Contains(decision.Reason, "incompatible") || !strings.Contains(decision.Reason, coreworkspace.WorkspaceSchemaVersionV1) { + t.Fatalf("decision = %#v, want incompatible legacy schema invalidation", decision) + } + if got := readFile(t, manifestPath); string(got) != string(beforeManifest) { + t.Fatal("legacy manifest changed during reuse decision") + } + if got := readFile(t, payloadPath); string(got) != string(beforePayload) { + t.Fatal("legacy payload changed during reuse decision") + } + }) + t.Run("corrupt payload", func(t *testing.T) { root := t.TempDir() recorder := newTestRecorder(t, root) - chunks := []contracts.SourceChunk{{ + chunks := []source.Chunk{{ ID: "chunk-1", SourceID: "source-1", Content: []byte("chunk content"), diff --git a/internal/framework/contracts/composition_test.go b/internal/framework/contracts/composition_test.go index 652c8a5..d30ec27 100644 --- a/internal/framework/contracts/composition_test.go +++ b/internal/framework/contracts/composition_test.go @@ -141,17 +141,20 @@ func (chunker compositionChunker) Chunk(ctx context.Context, req contracts.Chunk } return contracts.ChunkResult{ - Chunks: []contracts.SourceChunk{ + Chunks: []source.Chunk{ { - ID: req.Source.ID + ":chunk:0", - SourceID: req.Source.ID, - Index: 0, - StartUnitID: req.Source.Units[0].ID, - EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, - Content: []byte(`{"units":[{"id":1,"kind":"unit","text":"First source unit."},{"id":2,"kind":"unit","text":"Second source unit."}]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), req.Source.Units...), - Metadata: map[string]any{"strategy": "whole-document"}, + ID: req.Source.ID + ":chunk:0", + SourceID: req.Source.ID, + Index: 0, + Ref: source.SourceRef{ + SourceID: req.Source.ID, + StartUnitID: req.Source.Units[0].ID, + EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, + }, + Content: []byte(`{"units":[{"id":1,"kind":"unit","text":"First source unit."},{"id":2,"kind":"unit","text":"Second source unit."}]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), req.Source.Units...), + Metadata: map[string]any{"strategy": "whole-document"}, }, }, }, nil diff --git a/internal/framework/contracts/contracts.go b/internal/framework/contracts/contracts.go index 26c5ec8..31c6a0d 100644 --- a/internal/framework/contracts/contracts.go +++ b/internal/framework/contracts/contracts.go @@ -138,18 +138,6 @@ type InputAdapter interface { Parse(ctx context.Context, req ParseRequest) (*source.SourceDocument, error) } -type SourceChunk struct { - ID string `json:"id"` - SourceID string `json:"source_id"` - Index int `json:"index"` - StartUnitID int `json:"start_unit_id"` - EndUnitID int `json:"end_unit_id"` - Content []byte `json:"-"` - MediaType string `json:"media_type"` - Units []source.SourceUnit `json:"units"` - Metadata map[string]any `json:"metadata,omitempty"` -} - type ChunkRequest struct { Source *source.SourceDocument `json:"-"` SourceInput LLMInputMaterial `json:"source_input,omitempty"` @@ -162,8 +150,8 @@ type ChunkRequest struct { } type ChunkResult struct { - Chunks []SourceChunk `json:"chunks"` - Warnings []Warning `json:"warnings,omitempty"` + Chunks []source.Chunk `json:"chunks"` + Warnings []Warning `json:"warnings,omitempty"` } type Chunker interface { @@ -224,7 +212,7 @@ type ReferenceSet struct { type ExtractionRequest struct { Source *source.SourceDocument `json:"-"` - Chunk *SourceChunk `json:"chunk,omitempty"` + Chunk *source.Chunk `json:"chunk,omitempty"` AmbientContext map[string]any `json:"ambient_context,omitempty"` SourceInput LLMInputMaterial `json:"source_input,omitempty"` SessionID string `json:"session_id,omitempty"` @@ -277,8 +265,8 @@ type ValidationRequest struct { Payload RawPayload `json:"payload"` ChunkID string `json:"chunk_id,omitempty"` ChunkIndex int `json:"chunk_index,omitempty"` - Chunk *SourceChunk `json:"chunk,omitempty"` - Chunks []SourceChunk `json:"chunks,omitempty"` + Chunk *source.Chunk `json:"chunk,omitempty"` + Chunks []source.Chunk `json:"chunks,omitempty"` ExtractOutputs []ExtractOutput `json:"extract_outputs,omitempty"` MergeOutput MergeOutput `json:"merge_output,omitempty"` } diff --git a/internal/framework/contracts/contracts_test.go b/internal/framework/contracts/contracts_test.go index 9f81ca8..63d62fe 100644 --- a/internal/framework/contracts/contracts_test.go +++ b/internal/framework/contracts/contracts_test.go @@ -78,22 +78,22 @@ func TestFakeChunkerReturnsSourceChunks(t *testing.T) { chunk := result.Chunks[0] if chunk.ID != "source-1:chunk:0" { - t.Fatalf("SourceChunk.ID = %q, want source-1:chunk:0", chunk.ID) + t.Fatalf("source.Chunk.ID = %q, want source-1:chunk:0", chunk.ID) } if chunk.SourceID != doc.ID { - t.Fatalf("SourceChunk.SourceID = %q, want %q", chunk.SourceID, doc.ID) + t.Fatalf("source.Chunk.SourceID = %q, want %q", chunk.SourceID, doc.ID) } if chunk.Index != 0 { - t.Fatalf("SourceChunk.Index = %d, want 0", chunk.Index) + t.Fatalf("source.Chunk.Index = %d, want 0", chunk.Index) } - if chunk.StartUnitID != 1 || chunk.EndUnitID != 1 { - t.Fatalf("SourceChunk boundaries = %d-%d, want 1-1", chunk.StartUnitID, chunk.EndUnitID) + if chunk.Ref.StartUnitID != 1 || chunk.Ref.EndUnitID != 1 { + t.Fatalf("source.Chunk.Ref = %#v, want source-1:1-1", chunk.Ref) } if chunk.MediaType != "application/json" || string(chunk.Content) != `{"units":[{"id":1,"kind":"section","text":"Source text."}]}` { - t.Fatalf("SourceChunk payload = %q %s, want JSON units", chunk.MediaType, chunk.Content) + t.Fatalf("source.Chunk payload = %q %s, want JSON units", chunk.MediaType, chunk.Content) } if len(chunk.Units) != 1 { - t.Fatalf("len(SourceChunk.Units) = %d, want 1", len(chunk.Units)) + t.Fatalf("len(source.Chunk.Units) = %d, want 1", len(chunk.Units)) } } @@ -130,15 +130,14 @@ func TestFakeExtractorReceivesChunkAndAmbientContext(t *testing.T) { {ID: 2, Kind: "section", Text: "Second source text."}, }, } - chunk := SourceChunk{ - ID: "source-1:chunk:1", - SourceID: doc.ID, - Index: 1, - StartUnitID: 2, - EndUnitID: 2, - Content: []byte(`{"units":[{"id":2,"kind":"section","text":"Second source text."}]}`), - MediaType: "application/json", - Units: []source.SourceUnit{doc.Units[1]}, + chunk := source.Chunk{ + ID: "source-1:chunk:1", + SourceID: doc.ID, + Index: 1, + Ref: source.SourceRef{SourceID: doc.ID, StartUnitID: 2, EndUnitID: 2}, + Content: []byte(`{"units":[{"id":2,"kind":"section","text":"Second source text."}]}`), + MediaType: "application/json", + Units: []source.SourceUnit{doc.Units[1]}, } result, err := extractor.Extract(context.Background(), ExtractionRequest{ @@ -472,16 +471,19 @@ func (chunker fakeChunker) ReferenceSlots() []ReferenceSlot { func (chunker fakeChunker) Chunk(ctx context.Context, req ChunkRequest) (ChunkResult, error) { return ChunkResult{ - Chunks: []SourceChunk{ + Chunks: []source.Chunk{ { - ID: req.Source.ID + ":chunk:0", - SourceID: req.Source.ID, - Index: 0, - StartUnitID: req.Source.Units[0].ID, - EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, - Content: []byte(`{"units":[{"id":1,"kind":"section","text":"Source text."}]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), req.Source.Units...), + ID: req.Source.ID + ":chunk:0", + SourceID: req.Source.ID, + Index: 0, + Ref: source.SourceRef{ + SourceID: req.Source.ID, + StartUnitID: req.Source.Units[0].ID, + EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, + }, + Content: []byte(`{"units":[{"id":1,"kind":"section","text":"Source text."}]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), req.Source.Units...), }, }, }, nil diff --git a/internal/framework/pipeline/checkpoint.go b/internal/framework/pipeline/checkpoint.go index 974d69d..200f248 100644 --- a/internal/framework/pipeline/checkpoint.go +++ b/internal/framework/pipeline/checkpoint.go @@ -21,7 +21,7 @@ type CheckpointRecorder interface { SourceSucceeded(moduleKey string, doc *source.SourceDocument) error SourceFailed(moduleKey string, err error) error ChunkRunning(moduleKey string, sourceDigest string) error - ChunkSucceeded(moduleKey string, sourceDigest string, chunks []contracts.SourceChunk, warnings []contracts.Warning) error + ChunkSucceeded(moduleKey string, sourceDigest string, chunks []source.Chunk, warnings []contracts.Warning) error ChunkRejected(moduleKey string, sourceDigest string, rejected contracts.RejectedOutput) error ChunkFailed(moduleKey string, sourceDigest string, err error) error ExtractRunning(laneID string, moduleKey string, dependencies []CheckpointFingerprint) error @@ -55,7 +55,7 @@ type SourceCheckpoint struct { } type ChunkCheckpoint struct { - Chunks []contracts.SourceChunk + Chunks []source.Chunk Warnings []contracts.Warning } @@ -94,7 +94,7 @@ func (noopCheckpointRecorder) SourceRunning(string) error func (noopCheckpointRecorder) SourceSucceeded(string, *source.SourceDocument) error { return nil } func (noopCheckpointRecorder) SourceFailed(string, error) error { return nil } func (noopCheckpointRecorder) ChunkRunning(string, string) error { return nil } -func (noopCheckpointRecorder) ChunkSucceeded(string, string, []contracts.SourceChunk, []contracts.Warning) error { +func (noopCheckpointRecorder) ChunkSucceeded(string, string, []source.Chunk, []contracts.Warning) error { return nil } func (noopCheckpointRecorder) ChunkRejected(string, string, contracts.RejectedOutput) error { @@ -180,17 +180,21 @@ func digestFingerprints(name string, digest string) []CheckpointFingerprint { return []CheckpointFingerprint{{Name: name, Value: digest}} } -func joinedChunkDigest(chunks []contracts.SourceChunk) string { +func joinedChunkDigest(chunks []source.Chunk) (string, error) { if len(chunks) == 0 { - return "" + return "", nil } values := make([]string, 0, len(chunks)) for _, chunk := range chunks { - values = append(values, chunk.ID+"="+checkpointContentDigest(chunk.Content)) + digest, err := source.DigestChunk(chunk) + if err != nil { + return "", fmt.Errorf("digest chunk %q: %w", chunk.ID, err) + } + values = append(values, chunk.ID+"="+digest) } sort.Strings(values) sum := sha256.Sum256([]byte(strings.Join(values, "\n"))) - return "sha256:" + hex.EncodeToString(sum[:]) + return "sha256:" + hex.EncodeToString(sum[:]), nil } func normalizeCheckpointFingerprints(values []CheckpointFingerprint) []CheckpointFingerprint { diff --git a/internal/framework/pipeline/chunk_validation.go b/internal/framework/pipeline/chunk_validation.go index 8f9a1e5..1baa952 100644 --- a/internal/framework/pipeline/chunk_validation.go +++ b/internal/framework/pipeline/chunk_validation.go @@ -5,10 +5,9 @@ import ( "strings" "gitea.maximumdirect.net/eric/notarius/internal/core/source" - "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" ) -func validateAndCanonicalizeChunkResult(doc *source.SourceDocument, chunks []contracts.SourceChunk) ([]contracts.SourceChunk, error) { +func validateAndCanonicalizeChunkResult(doc *source.SourceDocument, chunks []source.Chunk) ([]source.Chunk, error) { if len(chunks) == 0 { return nil, fmt.Errorf("chunks must not be empty") } @@ -20,7 +19,7 @@ func validateAndCanonicalizeChunkResult(doc *source.SourceDocument, chunks []con sourceUnits[unit.ID] = unit } - canonicalChunks := make([]contracts.SourceChunk, 0, len(chunks)) + canonicalChunks := make([]source.Chunk, 0, len(chunks)) seenChunkIDs := make(map[string]struct{}, len(chunks)) for chunkIndex, chunk := range chunks { if strings.TrimSpace(chunk.ID) == "" { @@ -37,17 +36,6 @@ func validateAndCanonicalizeChunkResult(doc *source.SourceDocument, chunks []con if chunk.Index != chunkIndex { return nil, fmt.Errorf("chunk %q index %d does not match returned order %d", chunk.ID, chunk.Index, chunkIndex) } - startIndex, ok := sourceUnitIndexes[chunk.StartUnitID] - if !ok { - return nil, fmt.Errorf("chunk %q start_unit_id %d was not found in source document %q", chunk.ID, chunk.StartUnitID, doc.ID) - } - endIndex, ok := sourceUnitIndexes[chunk.EndUnitID] - if !ok { - return nil, fmt.Errorf("chunk %q end_unit_id %d was not found in source document %q", chunk.ID, chunk.EndUnitID, doc.ID) - } - if startIndex > endIndex { - return nil, fmt.Errorf("chunk %q start_unit_id %d appears after end_unit_id %d", chunk.ID, chunk.StartUnitID, chunk.EndUnitID) - } if len(chunk.Units) == 0 { return nil, fmt.Errorf("chunk %q units must not be empty", chunk.ID) } @@ -57,6 +45,9 @@ func validateAndCanonicalizeChunkResult(doc *source.SourceDocument, chunks []con if strings.TrimSpace(chunk.MediaType) == "" { return nil, fmt.Errorf("chunk %q media_type must not be empty", chunk.ID) } + if err := source.ValidateRef(doc, chunk.Ref); err != nil { + return nil, fmt.Errorf("chunk %q ref: %w", chunk.ID, err) + } seenUnitIDs := make(map[int]struct{}, len(chunk.Units)) previousSourceIndex := -1 @@ -74,23 +65,33 @@ func validateAndCanonicalizeChunkResult(doc *source.SourceDocument, chunks []con if !ok { return nil, fmt.Errorf("chunk %q source unit %d was not found in source document %q", chunk.ID, unit.ID, doc.ID) } - if sourceIndex <= previousSourceIndex { - return nil, fmt.Errorf("chunk %q source units must appear in source document order", chunk.ID) + if previousSourceIndex >= 0 && sourceIndex != previousSourceIndex+1 { + return nil, fmt.Errorf("chunk %q source units must form a contiguous range in source document order", chunk.ID) + } + if unit.Ref != sourceUnits[unit.ID].Ref { + return nil, fmt.Errorf("chunk %q source unit %d ref does not match source document", chunk.ID, unit.ID) } previousSourceIndex = sourceIndex canonicalUnits = append(canonicalUnits, cloneSourceUnit(sourceUnits[unit.ID])) } + expectedRef := source.SourceRef{ + SourceID: doc.ID, + StartUnitID: canonicalUnits[0].Ref.StartUnitID, + EndUnitID: canonicalUnits[len(canonicalUnits)-1].Ref.EndUnitID, + } + if chunk.Ref != expectedRef { + return nil, fmt.Errorf("chunk %q ref %#v does not match unit span %#v", chunk.ID, chunk.Ref, expectedRef) + } - canonicalChunks = append(canonicalChunks, contracts.SourceChunk{ - ID: chunk.ID, - SourceID: chunk.SourceID, - Index: chunk.Index, - StartUnitID: chunk.StartUnitID, - EndUnitID: chunk.EndUnitID, - Content: append([]byte(nil), chunk.Content...), - MediaType: chunk.MediaType, - Units: canonicalUnits, - Metadata: cloneMetadata(chunk.Metadata), + canonicalChunks = append(canonicalChunks, source.Chunk{ + ID: chunk.ID, + SourceID: chunk.SourceID, + Index: chunk.Index, + Ref: expectedRef, + Content: append([]byte(nil), chunk.Content...), + MediaType: chunk.MediaType, + Units: canonicalUnits, + Metadata: cloneMetadata(chunk.Metadata), }) } diff --git a/internal/framework/pipeline/debug.go b/internal/framework/pipeline/debug.go index e1c8b2e..6c3180e 100644 --- a/internal/framework/pipeline/debug.go +++ b/internal/framework/pipeline/debug.go @@ -104,14 +104,13 @@ type debugSourceDocument struct { } type debugSourceChunk struct { - ID string `json:"id"` - SourceID string `json:"source_id"` - Index int `json:"index"` - StartUnitID int `json:"start_unit_id"` - EndUnitID int `json:"end_unit_id"` - Content debugBinaryEnvelope `json:"content"` - Units []source.SourceUnit `json:"units,omitempty"` - Metadata map[string]any `json:"metadata,omitempty"` + ID string `json:"id"` + SourceID string `json:"source_id"` + Index int `json:"index"` + Ref source.SourceRef `json:"ref"` + Content debugBinaryEnvelope `json:"content"` + Units []source.SourceUnit `json:"units,omitempty"` + Metadata map[string]any `json:"metadata,omitempty"` } type debugExtractOutput struct { @@ -448,20 +447,19 @@ func debugSourceDocumentEnvelope(doc *source.SourceDocument) *debugSourceDocumen } } -func debugSourceChunkEnvelope(chunk contracts.SourceChunk) debugSourceChunk { +func debugSourceChunkEnvelope(chunk source.Chunk) debugSourceChunk { return debugSourceChunk{ - ID: chunk.ID, - SourceID: chunk.SourceID, - Index: chunk.Index, - StartUnitID: chunk.StartUnitID, - EndUnitID: chunk.EndUnitID, - Content: debugContentEnvelope(chunk.Content, chunk.MediaType, chunk.Metadata, nil), - Units: cloneSourceUnits(chunk.Units), - Metadata: redactSensitiveMap(chunk.Metadata), + ID: chunk.ID, + SourceID: chunk.SourceID, + Index: chunk.Index, + Ref: chunk.Ref, + Content: debugContentEnvelope(chunk.Content, chunk.MediaType, chunk.Metadata, nil), + Units: cloneSourceUnits(chunk.Units), + Metadata: redactSensitiveMap(chunk.Metadata), } } -func debugSourceChunkEnvelopes(chunks []contracts.SourceChunk) []debugSourceChunk { +func debugSourceChunkEnvelopes(chunks []source.Chunk) []debugSourceChunk { if len(chunks) == 0 { return nil } diff --git a/internal/framework/pipeline/debug_test.go b/internal/framework/pipeline/debug_test.go index d29c0a5..fca86eb 100644 --- a/internal/framework/pipeline/debug_test.go +++ b/internal/framework/pipeline/debug_test.go @@ -28,3 +28,24 @@ func TestDebugSourceDocumentPreservesUnitReferences(t *testing.T) { t.Fatalf("source document ref = %q after debug mutation, want source-1", got) } } + +func TestDebugSourceChunkPreservesReference(t *testing.T) { + doc := validSourceDocument() + chunk := source.Chunk{ + ID: "chunk-1", SourceID: doc.ID, Index: 0, + Ref: source.SourceRef{SourceID: doc.ID, StartUnitID: 1, EndUnitID: 1}, + Content: []byte("chunk content"), MediaType: "text/plain", Units: doc.Units[:1], + } + envelope := debugSourceChunkEnvelope(chunk) + encoded, err := json.Marshal(envelope) + if err != nil { + t.Fatalf("Marshal(debug source chunk) error = %v, want nil", err) + } + var decoded debugSourceChunk + if err := json.Unmarshal(encoded, &decoded); err != nil { + t.Fatalf("Unmarshal(debug source chunk) error = %v, want nil", err) + } + if decoded.Ref != chunk.Ref { + t.Fatalf("debug chunk ref = %#v, want %#v", decoded.Ref, chunk.Ref) + } +} diff --git a/internal/framework/pipeline/registry_integration_test.go b/internal/framework/pipeline/registry_integration_test.go index 9b47acc..48114f8 100644 --- a/internal/framework/pipeline/registry_integration_test.go +++ b/internal/framework/pipeline/registry_integration_test.go @@ -117,16 +117,19 @@ func (chunker integrationChunker) ReferenceSlots() []contracts.ReferenceSlot { func (chunker integrationChunker) Chunk(ctx context.Context, req contracts.ChunkRequest) (contracts.ChunkResult, error) { return contracts.ChunkResult{ - Chunks: []contracts.SourceChunk{ + Chunks: []source.Chunk{ { - ID: "chunk-0", - SourceID: req.Source.ID, - Index: 0, - StartUnitID: req.Source.Units[0].ID, - EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, - Content: []byte(`{"units":[1]}`), - MediaType: "application/json", - Units: req.Source.Units, + ID: "chunk-0", + SourceID: req.Source.ID, + Index: 0, + Ref: source.SourceRef{ + SourceID: req.Source.ID, + StartUnitID: req.Source.Units[0].ID, + EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, + }, + Content: []byte(`{"units":[1]}`), + MediaType: "application/json", + Units: req.Source.Units, }, }, }, nil diff --git a/internal/framework/pipeline/runner.go b/internal/framework/pipeline/runner.go index 916b0b7..c7b35e0 100644 --- a/internal/framework/pipeline/runner.go +++ b/internal/framework/pipeline/runner.go @@ -173,7 +173,7 @@ func (r *Runner) Run(ctx context.Context, input RunInput) (output RunOutput, err return failOutput(output), fmt.Errorf("build chunker %q: %w", input.Pipeline.Chunk.Module, err) } attachModuleManifestMetadata(&output, "chunker", chunker) - var canonicalChunks []contracts.SourceChunk + var canonicalChunks []source.Chunk var chunkWarnings []contracts.Warning chunkCheckpoint, chunkDecision := checkpointLoader.Chunk(chunker.Key(), doc.Digest) recordCheckpointEvent(&output, checkpointLoader, string(StageChunk), "", chunker.Key(), chunkDecision) @@ -388,7 +388,7 @@ func (r *Runner) Run(ctx context.Context, input RunInput) (output RunOutput, err return output, nil } -func (r *Runner) runLane(ctx context.Context, input RunInput, checkpoints CheckpointRecorder, checkpointLoader CheckpointLoader, doc *source.SourceDocument, sourceInput contracts.LLMInputMaterial, sessionID string, chunks []contracts.SourceChunk, lane ResolvedArtifactLane, output *RunOutput) error { +func (r *Runner) runLane(ctx context.Context, input RunInput, checkpoints CheckpointRecorder, checkpointLoader CheckpointLoader, doc *source.SourceDocument, sourceInput contracts.LLMInputMaterial, sessionID string, chunks []source.Chunk, lane ResolvedArtifactLane, output *RunOutput) error { extractor, err := r.registries.Extractors.Build(lane.Extract.Module) if err != nil { return fmt.Errorf("build extractor %q for lane %q: %w", lane.Extract.Module, lane.ID, err) @@ -406,7 +406,11 @@ func (r *Runner) runLane(ctx context.Context, input RunInput, checkpoints Checkp extractOutputs := make([]contracts.ExtractOutput, 0, len(chunks)) extractWarnings := []contracts.Warning{} extractRejectedStart := len(output.Rejected) - extractDependencies := digestFingerprints("chunks", joinedChunkDigest(chunks)) + chunksDigest, err := joinedChunkDigest(chunks) + if err != nil { + return fmt.Errorf("digest chunks for lane %q: %w", lane.ID, err) + } + extractDependencies := digestFingerprints("chunks", chunksDigest) extractCheckpoint, extractDecision := checkpointLoader.Extract(lane.ID, extractor.Key(), extractDependencies) recordCheckpointEvent(output, checkpointLoader, string(StageExtract), lane.ID, extractor.Key(), extractDecision) extractStarted := time.Now().UTC() @@ -887,8 +891,8 @@ type rawValidationTarget struct { llmClient contracts.StructuredLLMClient chunkID string chunkIndex int - chunk *contracts.SourceChunk - chunks []contracts.SourceChunk + chunk *source.Chunk + chunks []source.Chunk schema contracts.ResponseSchema payload contracts.RawPayload extractOutputs []contracts.ExtractOutput @@ -946,7 +950,7 @@ 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, debug DebugRecorder) ([]contracts.Warning, *contracts.RejectedOutput, error) { +func (r *Runner) validateChunksRaw(ctx context.Context, doc *source.SourceDocument, moduleKey string, chunks []source.Chunk, sourceInput contracts.LLMInputMaterial, sessionID string, references contracts.ReferenceSet, llmClient contracts.StructuredLLMClient, metadata map[string]any, chains []ResolvedValidatorChain, attempt int, debug DebugRecorder) ([]contracts.Warning, *contracts.RejectedOutput, error) { return r.validateRaw(ctx, rawValidationTarget{ stage: StageChunk, moduleKey: moduleKey, @@ -1463,7 +1467,7 @@ func sourceInputMaterial(inputPath string, content []byte) contracts.LLMInputMat ) } -func chunkInputMaterial(sourceInput contracts.LLMInputMaterial, chunk contracts.SourceChunk) contracts.LLMInputMaterial { +func chunkInputMaterial(sourceInput contracts.LLMInputMaterial, chunk source.Chunk) contracts.LLMInputMaterial { return contracts.NewLLMInputMaterial( "source", chunk.MediaType, @@ -1537,7 +1541,7 @@ func cloneResponseSchema(schema contracts.ResponseSchema) contracts.ResponseSche return schema } -func cloneSourceChunkPtr(chunk *contracts.SourceChunk) *contracts.SourceChunk { +func cloneSourceChunkPtr(chunk *source.Chunk) *source.Chunk { if chunk == nil { return nil } @@ -1545,18 +1549,18 @@ func cloneSourceChunkPtr(chunk *contracts.SourceChunk) *contracts.SourceChunk { return &cloned } -func cloneSourceChunk(chunk contracts.SourceChunk) contracts.SourceChunk { +func cloneSourceChunk(chunk source.Chunk) source.Chunk { chunk.Content = append([]byte(nil), chunk.Content...) chunk.Units = cloneSourceUnits(chunk.Units) chunk.Metadata = cloneMetadata(chunk.Metadata) return chunk } -func cloneSourceChunks(chunks []contracts.SourceChunk) []contracts.SourceChunk { +func cloneSourceChunks(chunks []source.Chunk) []source.Chunk { if len(chunks) == 0 { return nil } - out := make([]contracts.SourceChunk, 0, len(chunks)) + out := make([]source.Chunk, 0, len(chunks)) for _, chunk := range chunks { out = append(out, cloneSourceChunk(chunk)) } diff --git a/internal/framework/pipeline/runner_test.go b/internal/framework/pipeline/runner_test.go index b0e6877..481564a 100644 --- a/internal/framework/pipeline/runner_test.go +++ b/internal/framework/pipeline/runner_test.go @@ -249,17 +249,17 @@ func TestRunRejectsChunkerBuildChunkAndEmptyChunkErrors(t *testing.T) { func TestRunRejectsInvalidChunks(t *testing.T) { tests := []struct { name string - chunks []contracts.SourceChunk + chunks []source.Chunk want string }{ { name: "empty chunk id", - chunks: []contracts.SourceChunk{chunkWithUnits("", "source-1", 0, unitWithID("u1"))}, + chunks: []source.Chunk{chunkWithUnits("", "source-1", 0, unitWithID("u1"))}, want: "id must not be empty", }, { name: "duplicate chunk id", - chunks: []contracts.SourceChunk{ + chunks: []source.Chunk{ chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u1")), chunkWithUnits("chunk-0", "source-1", 1, unitWithID("u2")), }, @@ -267,58 +267,95 @@ func TestRunRejectsInvalidChunks(t *testing.T) { }, { name: "wrong source id", - chunks: []contracts.SourceChunk{chunkWithUnits("chunk-0", "other-source", 0, unitWithID("u1"))}, + chunks: []source.Chunk{chunkWithUnits("chunk-0", "other-source", 0, unitWithID("u1"))}, want: "source_id", }, { name: "wrong index", - chunks: []contracts.SourceChunk{chunkWithUnits("chunk-0", "source-1", 1, unitWithID("u1"))}, + chunks: []source.Chunk{chunkWithUnits("chunk-0", "source-1", 1, unitWithID("u1"))}, want: "index", }, + { + name: "missing ref", + chunks: []source.Chunk{func() source.Chunk { + chunk := chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u1")) + chunk.Ref = source.SourceRef{} + return chunk + }()}, + want: "source_id", + }, + { + name: "foreign ref", + chunks: []source.Chunk{func() source.Chunk { + chunk := chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u1")) + chunk.Ref.SourceID = "other-source" + return chunk + }()}, + want: "does not match document id", + }, { name: "unknown start id", - chunks: []contracts.SourceChunk{chunkWithBounds("chunk-0", "source-1", 0, 9, 1, unitWithID("u1"))}, + chunks: []source.Chunk{chunkWithRef("chunk-0", "source-1", 0, 9, 1, unitWithID("u1"))}, want: "start_unit_id", }, { name: "unknown end id", - chunks: []contracts.SourceChunk{chunkWithBounds("chunk-0", "source-1", 0, 1, 9, unitWithID("u1"))}, + chunks: []source.Chunk{chunkWithRef("chunk-0", "source-1", 0, 1, 9, unitWithID("u1"))}, want: "end_unit_id", }, { name: "reversed bounds", - chunks: []contracts.SourceChunk{chunkWithBounds("chunk-0", "source-1", 0, 2, 1, unitWithID("u1"), unitWithID("u2"))}, + chunks: []source.Chunk{chunkWithRef("chunk-0", "source-1", 0, 2, 1, unitWithID("u1"), unitWithID("u2"))}, want: "appears after", }, { name: "empty units", - chunks: []contracts.SourceChunk{{ID: "chunk-0", SourceID: "source-1", Index: 0, StartUnitID: 1, EndUnitID: 1, Content: []byte(`{"units":[]}`), MediaType: "application/json"}}, + chunks: []source.Chunk{{ID: "chunk-0", SourceID: "source-1", Index: 0, Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, Content: []byte(`{"units":[]}`), MediaType: "application/json"}}, want: "units must not be empty", }, { name: "empty content", - chunks: []contracts.SourceChunk{{ID: "chunk-0", SourceID: "source-1", Index: 0, StartUnitID: 1, EndUnitID: 1, MediaType: "application/json", Units: []source.SourceUnit{unitWithID("u1")}}}, + chunks: []source.Chunk{{ID: "chunk-0", SourceID: "source-1", Index: 0, Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, MediaType: "application/json", Units: []source.SourceUnit{unitWithID("u1")}}}, want: "content must not be empty", }, { name: "empty media type", - chunks: []contracts.SourceChunk{{ID: "chunk-0", SourceID: "source-1", Index: 0, StartUnitID: 1, EndUnitID: 1, Content: []byte(`{"units":[1]}`), Units: []source.SourceUnit{unitWithID("u1")}}}, + chunks: []source.Chunk{{ID: "chunk-0", SourceID: "source-1", Index: 0, Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, Content: []byte(`{"units":[1]}`), Units: []source.SourceUnit{unitWithID("u1")}}}, want: "media_type must not be empty", }, { name: "repeated unit inside chunk", - chunks: []contracts.SourceChunk{chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u1"), unitWithID("u1"))}, + chunks: []source.Chunk{chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u1"), unitWithID("u1"))}, want: "repeats source unit", }, { name: "unknown unit", - chunks: []contracts.SourceChunk{chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u9"))}, + chunks: []source.Chunk{chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u9"))}, want: "was not found", }, { name: "units out of source order", - chunks: []contracts.SourceChunk{chunkWithBounds("chunk-0", "source-1", 0, 1, 2, unitWithID("u2"), unitWithID("u1"))}, - want: "source document order", + chunks: []source.Chunk{chunkWithRef("chunk-0", "source-1", 0, 1, 2, unitWithID("u2"), unitWithID("u1"))}, + want: "contiguous range", + }, + { + name: "noncontiguous units", + chunks: []source.Chunk{chunkWithRef("chunk-0", "source-1", 0, 1, 3, unitWithID("u1"), unitWithID("u3"))}, + want: "contiguous range", + }, + { + name: "ref does not match unit span", + chunks: []source.Chunk{chunkWithRef("chunk-0", "source-1", 0, 1, 2, unitWithID("u1"))}, + want: "does not match unit span", + }, + { + name: "unit ref does not match source", + chunks: []source.Chunk{func() source.Chunk { + unit := unitWithID("u1") + unit.Ref.SourceID = "other-source" + return chunkWithRef("chunk-0", "source-1", 0, 1, 1, unit) + }()}, + want: "ref does not match source document", }, } @@ -342,7 +379,7 @@ func TestRunRejectsInvalidChunks(t *testing.T) { func TestRunAllowsPartialCoverageAndOverlappingChunks(t *testing.T) { modules := defaultRunnerModules() - modules.chunker.chunks = []contracts.SourceChunk{ + modules.chunker.chunks = []source.Chunk{ chunkWithUnits("chunk-0", "source-1", 0, unitWithID("u1"), unitWithID("u2")), chunkWithUnits("chunk-1", "source-1", 1, unitWithID("u2")), } @@ -359,20 +396,20 @@ func TestRunAllowsPartialCoverageAndOverlappingChunks(t *testing.T) { func TestRunCanonicalizesChunkUnitsBeforeExtraction(t *testing.T) { modules := defaultRunnerModules() modules.input.doc = sourceDocumentWithUnitMetadata() - modules.chunker.chunks = []contracts.SourceChunk{ + modules.chunker.chunks = []source.Chunk{ { - ID: "chunk-0", - SourceID: "source-1", - Index: 0, - StartUnitID: 1, - EndUnitID: 1, - Content: []byte(`{"units":[{"id":1}]}`), - MediaType: "application/json", + ID: "chunk-0", + SourceID: "source-1", + Index: 0, + Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, + Content: []byte(`{"units":[{"id":1}]}`), + MediaType: "application/json", Units: []source.SourceUnit{ { ID: 1, Kind: "mutated-kind", Text: "mutated text", + Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, Metadata: map[string]any{ "speaker": "chunker-speaker", "note": "chunker note", @@ -423,15 +460,14 @@ func TestRunCanonicalizesChunkUnitsBeforeExtraction(t *testing.T) { func TestRunPreservesChunkMetadataDuringCanonicalization(t *testing.T) { modules := defaultRunnerModules() - modules.chunker.chunks = []contracts.SourceChunk{ + modules.chunker.chunks = []source.Chunk{ { - ID: "chunk-0", - SourceID: "source-1", - Index: 0, - StartUnitID: 1, - EndUnitID: 1, - Content: []byte(`{"units":[{"id":1}]}`), - MediaType: "application/json", + ID: "chunk-0", + SourceID: "source-1", + Index: 0, + Ref: source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 1}, + Content: []byte(`{"units":[{"id":1}]}`), + MediaType: "application/json", Units: []source.SourceUnit{ unitWithID("u1"), }, @@ -975,7 +1011,7 @@ func TestRunPassesPerChunkRawOutputsToMergeAndNormalize(t *testing.T) { func TestRunPassesChunkContentAndMediaTypeToExtractors(t *testing.T) { modules := defaultRunnerModules() - modules.chunker.chunks = []contracts.SourceChunk{ + modules.chunker.chunks = []source.Chunk{ sourceChunkWithContent("chunk-0", 0, []byte(`{"chunk":0}`), "application/vnd.test+json"), } @@ -1027,10 +1063,27 @@ func TestRunDoesNotPassCheckpointPathsToModules(t *testing.T) { } } +func TestCheckpointChunkDigestIncludesCanonicalReference(t *testing.T) { + chunk := sourceChunkWithID("chunk-0", 0) + first, err := joinedChunkDigest([]source.Chunk{chunk}) + if err != nil { + t.Fatalf("joinedChunkDigest() error = %v, want nil", err) + } + + chunk.Ref.EndUnitID = 2 + second, err := joinedChunkDigest([]source.Chunk{chunk}) + if err != nil { + t.Fatalf("joinedChunkDigest(changed ref) error = %v, want nil", err) + } + if first == second { + t.Fatalf("checkpoint chunk digests = %q and %q, want provenance change to alter dependency identity", first, second) + } +} + func TestRunReusesCheckpointedWorkflowOutputs(t *testing.T) { modules := defaultRunnerModules() doc := validSourceDocument() - chunks := []contracts.SourceChunk{sourceChunkWithID("chunk-0", 0)} + chunks := []source.Chunk{sourceChunkWithID("chunk-0", 0)} extractOutput := contracts.ExtractOutput{ LaneID: "alpha", ExtractorKey: "extract-alpha", @@ -1468,7 +1521,7 @@ func TestRunCollectsStageWarnings(t *testing.T) { func TestRunCollectsChunkValidatorWarnings(t *testing.T) { modules := defaultRunnerModules() - modules.chunker.chunks = []contracts.SourceChunk{sourceChunkWithID("chunk-0", 0)} + modules.chunker.chunks = []source.Chunk{sourceChunkWithID("chunk-0", 0)} validator := &runnerChainValidator{ name: "chain-chunk", warnings: []contracts.Warning{{ReasonCode: "chunk-validator-warning", Message: "chunk validator warning"}}, @@ -1904,7 +1957,7 @@ type runnerModules struct { func defaultRunnerModules() *runnerModules { return &runnerModules{ input: &runnerInputAdapter{key: "input", doc: validSourceDocument()}, - chunker: &runnerChunker{key: "chunk", chunks: []contracts.SourceChunk{sourceChunkWithID("chunk-0", 0), sourceChunkWithID("chunk-1", 1)}}, + chunker: &runnerChunker{key: "chunk", chunks: []source.Chunk{sourceChunkWithID("chunk-0", 0), sourceChunkWithID("chunk-1", 1)}}, extractors: map[string]*runnerExtractor{ "extract-alpha": {key: "extract-alpha"}, }, @@ -2013,7 +2066,7 @@ func (adapter *runnerInputAdapter) ManifestMetadata() map[string]any { type runnerChunker struct { key string - chunks []contracts.SourceChunk + chunks []source.Chunk warnings []contracts.Warning err error failureErr error @@ -2505,21 +2558,20 @@ func sourceDocumentWithUnitMetadata() *source.SourceDocument { } } -func sourceChunkWithID(id string, index int) contracts.SourceChunk { +func sourceChunkWithID(id string, index int) source.Chunk { unit := unitWithID("u1") - return contracts.SourceChunk{ - ID: id, - SourceID: "source-1", - Index: index, - StartUnitID: unit.ID, - EndUnitID: unit.ID, - Content: []byte(`{"units":[{"id":1,"kind":"unit","text":"Source unit."}]}`), - MediaType: "application/json", - Units: []source.SourceUnit{unit}, + return source.Chunk{ + ID: id, + SourceID: "source-1", + Index: index, + Ref: unit.Ref, + Content: []byte(`{"units":[{"id":1,"kind":"unit","text":"Source unit."}]}`), + MediaType: "application/json", + Units: []source.SourceUnit{unit}, } } -func sourceChunkWithContent(id string, index int, content []byte, mediaType string) contracts.SourceChunk { +func sourceChunkWithContent(id string, index int, content []byte, mediaType string) source.Chunk { chunk := sourceChunkWithID(id, index) chunk.Content = append([]byte(nil), content...) chunk.MediaType = mediaType @@ -2541,25 +2593,24 @@ func unitWithID(id string) source.SourceUnit { } } -func chunkWithUnits(id string, sourceID string, index int, units ...source.SourceUnit) contracts.SourceChunk { +func chunkWithUnits(id string, sourceID string, index int, units ...source.SourceUnit) source.Chunk { startUnitID, endUnitID := 1, 1 if len(units) > 0 { startUnitID = units[0].ID endUnitID = units[len(units)-1].ID } - return chunkWithBounds(id, sourceID, index, startUnitID, endUnitID, units...) + return chunkWithRef(id, sourceID, index, startUnitID, endUnitID, units...) } -func chunkWithBounds(id string, sourceID string, index int, startUnitID int, endUnitID int, units ...source.SourceUnit) contracts.SourceChunk { - return contracts.SourceChunk{ - ID: id, - SourceID: sourceID, - Index: index, - StartUnitID: startUnitID, - EndUnitID: endUnitID, - Content: []byte(`{"units":[1]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), units...), +func chunkWithRef(id string, sourceID string, index int, startUnitID int, endUnitID int, units ...source.SourceUnit) source.Chunk { + return source.Chunk{ + ID: id, + SourceID: sourceID, + Index: index, + Ref: source.SourceRef{SourceID: sourceID, StartUnitID: startUnitID, EndUnitID: endUnitID}, + Content: []byte(`{"units":[1]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), units...), } } diff --git a/internal/framework/pipeline/walking_skeleton_test.go b/internal/framework/pipeline/walking_skeleton_test.go index ac019de..1cadfdc 100644 --- a/internal/framework/pipeline/walking_skeleton_test.go +++ b/internal/framework/pipeline/walking_skeleton_test.go @@ -231,26 +231,24 @@ func (chunker walkingSkeletonChunker) Chunk(ctx context.Context, req contracts.C return contracts.ChunkResult{}, fmt.Errorf("fixture source must contain at least three units") } return contracts.ChunkResult{ - Chunks: []contracts.SourceChunk{ + Chunks: []source.Chunk{ { - ID: req.Source.ID + ":chunk:0", - SourceID: req.Source.ID, - Index: 0, - StartUnitID: req.Source.Units[0].ID, - EndUnitID: req.Source.Units[1].ID, - Content: []byte(`{"units":[1,2]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), req.Source.Units[:2]...), + ID: req.Source.ID + ":chunk:0", + SourceID: req.Source.ID, + Index: 0, + Ref: source.SourceRef{SourceID: req.Source.ID, StartUnitID: req.Source.Units[0].ID, EndUnitID: req.Source.Units[1].ID}, + Content: []byte(`{"units":[1,2]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), req.Source.Units[:2]...), }, { - ID: req.Source.ID + ":chunk:1", - SourceID: req.Source.ID, - Index: 1, - StartUnitID: req.Source.Units[2].ID, - EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, - Content: []byte(`{"units":[3]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), req.Source.Units[2:]...), + ID: req.Source.ID + ":chunk:1", + SourceID: req.Source.ID, + Index: 1, + Ref: source.SourceRef{SourceID: req.Source.ID, StartUnitID: req.Source.Units[2].ID, EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID}, + Content: []byte(`{"units":[3]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), req.Source.Units[2:]...), }, }, }, nil diff --git a/internal/modules/dnd/chunk/scenes/chunker.go b/internal/modules/dnd/chunk/scenes/chunker.go index e4bbbcd..f10c195 100644 --- a/internal/modules/dnd/chunk/scenes/chunker.go +++ b/internal/modules/dnd/chunk/scenes/chunker.go @@ -135,7 +135,7 @@ func Register(registry *pipeline.ChunkerRegistry) error { }) } -func chunksFromResponse(doc *source.SourceDocument, response chunkResponse) ([]contracts.SourceChunk, error) { +func chunksFromResponse(doc *source.SourceDocument, response chunkResponse) ([]source.Chunk, error) { if response.Scenes == nil { return nil, fmt.Errorf("scenes must be present") } @@ -148,7 +148,7 @@ func chunksFromResponse(doc *source.SourceDocument, response chunkResponse) ([]c unitIndexes[unit.ID] = i } - chunks := make([]contracts.SourceChunk, 0, len(response.Scenes)) + chunks := make([]source.Chunk, 0, len(response.Scenes)) previousEnd := -1 for i, scene := range response.Scenes { normalized, err := normalizeScene(doc, i, scene) @@ -186,15 +186,18 @@ func chunksFromResponse(doc *source.SourceDocument, response chunkResponse) ([]c if err != nil { return nil, err } - chunks = append(chunks, contracts.SourceChunk{ - ID: fmt.Sprintf("scene-%06d", i+1), - SourceID: doc.ID, - Index: i, - StartUnitID: units[0].ID, - EndUnitID: units[len(units)-1].ID, - Content: content, - MediaType: "application/json", - Units: units, + chunks = append(chunks, source.Chunk{ + ID: fmt.Sprintf("scene-%06d", i+1), + SourceID: doc.ID, + Index: i, + Ref: source.SourceRef{ + SourceID: doc.ID, + StartUnitID: units[0].Ref.StartUnitID, + EndUnitID: units[len(units)-1].Ref.EndUnitID, + }, + Content: content, + MediaType: "application/json", + Units: units, Metadata: map[string]any{ "scene_title": normalized.ShortTitle, "primary_mode": normalized.PrimaryMode, diff --git a/internal/modules/dnd/chunk/scenes/chunker_test.go b/internal/modules/dnd/chunk/scenes/chunker_test.go index cbe2049..d4b884c 100644 --- a/internal/modules/dnd/chunk/scenes/chunker_test.go +++ b/internal/modules/dnd/chunk/scenes/chunker_test.go @@ -179,8 +179,8 @@ func TestChunkReturnsSceneChunksFromStructuredOutput(t *testing.T) { if first.SourceID != "session-alpha" || first.Index != 0 { t.Fatalf("first chunk = %#v, want source and index fields", first) } - if first.StartUnitID != 1 || first.EndUnitID != 2 { - t.Fatalf("first boundaries = %d-%d, want 1-2", first.StartUnitID, first.EndUnitID) + if first.Ref != (source.SourceRef{SourceID: "session-alpha", StartUnitID: 1, EndUnitID: 2}) { + t.Fatalf("first ref = %#v, want session-alpha:1-2", first.Ref) } if first.MediaType != "application/json" || len(first.Content) == 0 { t.Fatalf("first payload = media type %q length %d, want JSON content", first.MediaType, len(first.Content)) @@ -574,7 +574,7 @@ func scene(startUnitID int, endUnitID int) sceneResponse { } } -func chunkIDs(chunks []contracts.SourceChunk) []string { +func chunkIDs(chunks []source.Chunk) []string { ids := make([]string, 0, len(chunks)) for _, chunk := range chunks { ids = append(ids, chunk.ID) diff --git a/internal/modules/dnd/extract/spells/extractor_test.go b/internal/modules/dnd/extract/spells/extractor_test.go index 41341bb..7ac6fd0 100644 --- a/internal/modules/dnd/extract/spells/extractor_test.go +++ b/internal/modules/dnd/extract/spells/extractor_test.go @@ -7,6 +7,7 @@ import ( "strings" "testing" + "gitea.maximumdirect.net/eric/notarius/internal/core/source" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" "gitea.maximumdirect.net/eric/notarius/internal/modules/dnd/shared" ) @@ -412,12 +413,12 @@ func spellSourceInput() contracts.LLMInputMaterial { return contracts.NewLLMInputMaterial("source", "application/json", []byte(spellTranscriptJSON), "sha256:transcript", "file:///session-alpha.json") } -func spellChunkInput(chunk *contracts.SourceChunk) contracts.LLMInputMaterial { +func spellChunkInput(chunk *source.Chunk) contracts.LLMInputMaterial { return contracts.NewLLMInputMaterial("source", chunk.MediaType, chunk.Content, "sha256:chunk", "file:///session-alpha.json") } func emptyChunkRequest(req contracts.ExtractionRequest) contracts.ExtractionRequest { - req.Chunk = &contracts.SourceChunk{ + req.Chunk = &source.Chunk{ ID: req.Chunk.ID, SourceID: req.Chunk.SourceID, Index: req.Chunk.Index, diff --git a/internal/modules/dnd/extract/spells/test_helpers_test.go b/internal/modules/dnd/extract/spells/test_helpers_test.go index 6c0f4cc..ff46b18 100644 --- a/internal/modules/dnd/extract/spells/test_helpers_test.go +++ b/internal/modules/dnd/extract/spells/test_helpers_test.go @@ -11,16 +11,19 @@ import ( func promptExtractionRequest() contracts.ExtractionRequest { doc := promptSourceDocument() - chunk := &contracts.SourceChunk{ - ID: "session-alpha:chunk:0", - SourceID: doc.ID, - Index: 0, - StartUnitID: doc.Units[0].ID, - EndUnitID: doc.Units[len(doc.Units)-1].ID, - Content: []byte(`{"units":[1,2]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), doc.Units...), - Metadata: map[string]any{"ignored": "chunk metadata"}, + chunk := &source.Chunk{ + ID: "session-alpha:chunk:0", + SourceID: doc.ID, + Index: 0, + Ref: source.SourceRef{ + SourceID: doc.ID, + StartUnitID: doc.Units[0].ID, + EndUnitID: doc.Units[len(doc.Units)-1].ID, + }, + Content: []byte(`{"units":[1,2]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), doc.Units...), + Metadata: map[string]any{"ignored": "chunk metadata"}, } return contracts.ExtractionRequest{ Source: doc, diff --git a/internal/modules/generic/chunk/units/chunker.go b/internal/modules/generic/chunk/units/chunker.go index 9cebdc6..0db19fb 100644 --- a/internal/modules/generic/chunk/units/chunker.go +++ b/internal/modules/generic/chunk/units/chunker.go @@ -61,7 +61,7 @@ func (c *Chunker) Chunk(ctx context.Context, req contracts.ChunkRequest) (contra } step := opts.maxUnits - opts.overlapUnits - chunks := make([]contracts.SourceChunk, 0, (len(req.Source.Units)+step-1)/step) + chunks := make([]source.Chunk, 0, (len(req.Source.Units)+step-1)/step) for start := 0; start < len(req.Source.Units); start += step { end := start + opts.maxUnits if end > len(req.Source.Units) { @@ -72,15 +72,18 @@ func (c *Chunker) Chunk(ctx context.Context, req contracts.ChunkRequest) (contra if err != nil { return contracts.ChunkResult{}, err } - chunks = append(chunks, contracts.SourceChunk{ - ID: fmt.Sprintf("chunk-%06d", len(chunks)+1), - SourceID: req.Source.ID, - Index: len(chunks), - StartUnitID: units[0].ID, - EndUnitID: units[len(units)-1].ID, - Content: content, - MediaType: "application/json", - Units: units, + chunks = append(chunks, source.Chunk{ + ID: fmt.Sprintf("chunk-%06d", len(chunks)+1), + SourceID: req.Source.ID, + Index: len(chunks), + Ref: source.SourceRef{ + SourceID: req.Source.ID, + StartUnitID: units[0].Ref.StartUnitID, + EndUnitID: units[len(units)-1].Ref.EndUnitID, + }, + Content: content, + MediaType: "application/json", + Units: units, Metadata: map[string]any{ "start_unit_id": units[0].ID, "end_unit_id": units[len(units)-1].ID, diff --git a/internal/modules/generic/chunk/units/chunker_test.go b/internal/modules/generic/chunk/units/chunker_test.go index 959633b..94c1cbc 100644 --- a/internal/modules/generic/chunk/units/chunker_test.go +++ b/internal/modules/generic/chunk/units/chunker_test.go @@ -62,8 +62,8 @@ func TestChunkUsesDefaultsForSingleChunk(t *testing.T) { if got := unitIDs(chunk.Units); !reflect.DeepEqual(got, []int{1, 2, 3}) { t.Fatalf("unit IDs = %#v, want all units", got) } - if chunk.StartUnitID != 1 || chunk.EndUnitID != 3 { - t.Fatalf("chunk boundaries = %d-%d, want 1-3", chunk.StartUnitID, chunk.EndUnitID) + if chunk.Ref != (source.SourceRef{SourceID: "source-1", StartUnitID: 1, EndUnitID: 3}) { + t.Fatalf("chunk ref = %#v, want source-1:1-3", chunk.Ref) } if chunk.MediaType != "application/json" || len(chunk.Content) == 0 { t.Fatalf("chunk payload = media type %q length %d, want JSON content", chunk.MediaType, len(chunk.Content)) @@ -209,7 +209,7 @@ func zeroPad3(value int) string { return fmt.Sprintf("%03d", value) } -func chunkIDs(chunks []contracts.SourceChunk) []string { +func chunkIDs(chunks []source.Chunk) []string { ids := make([]string, 0, len(chunks)) for _, chunk := range chunks { ids = append(ids, chunk.ID) diff --git a/internal/modules/integration/dnd_spells_config_test.go b/internal/modules/integration/dnd_spells_config_test.go index 288ab9c..a481702 100644 --- a/internal/modules/integration/dnd_spells_config_test.go +++ b/internal/modules/integration/dnd_spells_config_test.go @@ -244,16 +244,15 @@ func (dndSpellsChunker) ReferenceSlots() []contracts.ReferenceSlot { func (dndSpellsChunker) Chunk(ctx context.Context, req contracts.ChunkRequest) (contracts.ChunkResult, error) { return contracts.ChunkResult{ - Chunks: []contracts.SourceChunk{ + Chunks: []source.Chunk{ { - ID: req.Source.ID + ":chunk:0", - SourceID: req.Source.ID, - Index: 0, - StartUnitID: req.Source.Units[0].ID, - EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, - Content: []byte(`{"units":[1,2,3]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), req.Source.Units...), + ID: req.Source.ID + ":chunk:0", + SourceID: req.Source.ID, + Index: 0, + Ref: source.SourceRef{SourceID: req.Source.ID, StartUnitID: req.Source.Units[0].ID, EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID}, + Content: []byte(`{"units":[1,2,3]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), req.Source.Units...), }, }, }, nil diff --git a/internal/modules/seriatim/input/transcript/runner_test.go b/internal/modules/seriatim/input/transcript/runner_test.go index 99d8ea3..6c2e81e 100644 --- a/internal/modules/seriatim/input/transcript/runner_test.go +++ b/internal/modules/seriatim/input/transcript/runner_test.go @@ -169,16 +169,15 @@ func (runnerSeriatimChunker) ReferenceSlots() []contracts.ReferenceSlot { func (runnerSeriatimChunker) Chunk(ctx context.Context, req contracts.ChunkRequest) (contracts.ChunkResult, error) { return contracts.ChunkResult{ - Chunks: []contracts.SourceChunk{ + Chunks: []source.Chunk{ { - ID: req.Source.ID + ":chunk:0", - SourceID: req.Source.ID, - Index: 0, - StartUnitID: req.Source.Units[0].ID, - EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID, - Content: []byte(`{"units":[1,2]}`), - MediaType: "application/json", - Units: append([]source.SourceUnit(nil), req.Source.Units...), + ID: req.Source.ID + ":chunk:0", + SourceID: req.Source.ID, + Index: 0, + Ref: source.SourceRef{SourceID: req.Source.ID, StartUnitID: req.Source.Units[0].ID, EndUnitID: req.Source.Units[len(req.Source.Units)-1].ID}, + Content: []byte(`{"units":[1,2]}`), + MediaType: "application/json", + Units: append([]source.SourceUnit(nil), req.Source.Units...), }, }, }, nil