Consolidate configuration and resolver tests

This commit is contained in:
2026-07-18 23:38:59 +00:00
parent bbc83ab042
commit 0cca3b1f5d
10 changed files with 264 additions and 831 deletions

View File

@@ -18,109 +18,102 @@ import (
"gitea.maximumdirect.net/eric/notarius/internal/modules/seriatim/input/transcript"
)
func TestPipelineConfigLoadsAndResolvesWithDNDSpellsExtractor(t *testing.T) {
data, err := os.ReadFile("testdata/pipeline.yml")
if err != nil {
t.Fatalf("ReadFile(pipeline.yml) error = %v, want nil", err)
}
fileCfg, err := config.ParseFileConfigYAML(data)
if err != nil {
t.Fatalf("ParseFileConfigYAML() error = %v, want nil", err)
func TestDNDSpellCapabilityFailures(t *testing.T) {
tests := []struct {
name string
mutate func(pipeline.ModuleSpec, pipeline.ModuleSpec) (pipeline.ModuleSpec, pipeline.ModuleSpec)
wantModule string
wantCap string
}{
{
name: "spell extractor requires transcript source",
mutate: func(input, extractor pipeline.ModuleSpec) (pipeline.ModuleSpec, pipeline.ModuleSpec) {
input.Provides = withoutCapability(input.Provides, "source.transcript")
return input, extractor
},
wantModule: spells.Key,
wantCap: "source.transcript",
},
{
name: "append-order merger requires spell casts",
mutate: func(input, extractor pipeline.ModuleSpec) (pipeline.ModuleSpec, pipeline.ModuleSpec) {
extractor.Provides = withoutCapability(extractor.Provides, "dnd.spell_casts")
return input, extractor
},
wantModule: pipeline.DefaultMergeModule,
wantCap: "dnd.spell_casts",
},
}
cfg := config.Default()
if err := cfg.ApplyFileConfig(fileCfg); err != nil {
t.Fatalf("ApplyFileConfig() error = %v, want nil", err)
}
resolved, err := cfg.Resolve(config.ResolveInput{
PipelineID: "dnd-spells-fixture",
Catalog: dndSpellsTestCatalog(t, dndSpellsCatalogSpecs{}),
})
if err != nil {
t.Fatalf("Resolve() error = %v, want nil", err)
}
if len(resolved.ResolvedPipeline.ArtifactLanes) != 1 {
t.Fatalf("len(ArtifactLanes) = %d, want 1", len(resolved.ResolvedPipeline.ArtifactLanes))
}
lane := resolved.ResolvedPipeline.ArtifactLanes[0]
if lane.ID != "spells" {
t.Fatalf("lane ID = %q, want spells", lane.ID)
}
if lane.Extract.Module != spells.Key {
t.Fatalf("extract module = %q, want %q", lane.Extract.Module, spells.Key)
}
if resolved.ResolvedPipeline.Digest == "" {
t.Fatal("resolved digest is empty")
}
again, err := cfg.Resolve(config.ResolveInput{
PipelineID: "dnd-spells-fixture",
Catalog: dndSpellsTestCatalog(t, dndSpellsCatalogSpecs{}),
})
if err != nil {
t.Fatalf("second Resolve() error = %v, want nil", err)
}
if resolved.ResolvedPipeline.Digest != again.ResolvedPipeline.Digest {
t.Fatalf("resolved digest = %q, second digest = %q; want stable digest", resolved.ResolvedPipeline.Digest, again.ResolvedPipeline.Digest)
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
inputSpec, extractorSpec := tt.mutate(transcript.ModuleSpec(), spells.ModuleSpec())
_, err := pipeline.ResolvePipeline(dndCapabilityProfile(), pipeline.ResolveOptions{}, dndCapabilityCatalog(t, inputSpec, extractorSpec))
if err == nil || !strings.Contains(err.Error(), "missing capability") || !strings.Contains(err.Error(), tt.wantCap) || !strings.Contains(err.Error(), tt.wantModule) {
t.Fatalf("ResolvePipeline() error = %v, want %s missing %s capability", err, tt.wantModule, tt.wantCap)
}
})
}
}
func TestPipelineConfigRejectsMissingTranscriptCapabilityForDNDSpells(t *testing.T) {
inputSpec := transcript.ModuleSpec()
inputSpec.Provides = withoutCapability(inputSpec.Provides, "source.transcript")
chunkSpec := dndSpellsChunkerSpec()
chunkSpec.Requires = nil
_, err := loadDNDSpellsPipelineConfig(t).Resolve(config.ResolveInput{
PipelineID: "dnd-spells-fixture",
Catalog: dndSpellsTestCatalog(t, dndSpellsCatalogSpecs{
input: inputSpec,
chunk: chunkSpec,
}),
})
if err == nil {
t.Fatal("Resolve() error = nil, want missing capability error")
}
if !strings.Contains(err.Error(), "missing capability") ||
!strings.Contains(err.Error(), "source.transcript") ||
!strings.Contains(err.Error(), spells.Key) {
t.Fatalf("Resolve() error = %q, want dnd/spells missing source.transcript capability", err.Error())
func dndCapabilityProfile() pipeline.PipelineProfile {
return pipeline.PipelineProfile{
ID: "dnd-capability",
Input: pipeline.Binding(transcript.Key),
Chunk: pipeline.Binding("fake/chunk"),
Artifacts: map[string]pipeline.ArtifactLaneProfile{
"spells": {Extract: pipeline.Binding(spells.Key)},
},
}
}
func TestPipelineConfigRejectsMissingSpellCastsCapabilityForAppendOrder(t *testing.T) {
extractorSpec := spells.ModuleSpec()
extractorSpec.Provides = withoutCapability(extractorSpec.Provides, "dnd.spell_casts")
func dndCapabilityCatalog(t *testing.T, inputSpec, extractorSpec pipeline.ModuleSpec) pipeline.ModuleCatalog {
t.Helper()
inputs := pipeline.NewInputAdapterRegistry()
if err := inputs.RegisterWithSpec(inputSpec, func() (contracts.InputAdapter, error) { return transcript.New(), nil }); err != nil {
t.Fatalf("register capability input: %v", err)
}
_, err := loadDNDSpellsPipelineConfig(t).Resolve(config.ResolveInput{
PipelineID: "dnd-spells-fixture",
Catalog: dndSpellsTestCatalog(t, dndSpellsCatalogSpecs{
extractor: extractorSpec,
}),
})
if err == nil {
t.Fatal("Resolve() error = nil, want missing capability error")
chunkers := pipeline.NewChunkerRegistry()
if err := chunkers.RegisterWithSpec(pipeline.ModuleSpec{Key: "fake/chunk", Stage: pipeline.StageChunk, Provides: []string{"chunks"}}, func() (contracts.Chunker, error) { return dndSpellsChunker{}, nil }); err != nil {
t.Fatalf("register capability chunker: %v", err)
}
if !strings.Contains(err.Error(), "missing capability") ||
!strings.Contains(err.Error(), "dnd.spell_casts") ||
!strings.Contains(err.Error(), pipeline.DefaultMergeModule) {
t.Fatalf("Resolve() error = %q, want appendorder missing dnd.spell_casts capability", err.Error())
}
}
func TestPipelineConfigRejectsUnknownLaneSelection(t *testing.T) {
_, err := loadDNDSpellsPipelineConfig(t).Resolve(config.ResolveInput{
PipelineID: "dnd-spells-fixture",
Only: []string{"missing"},
Catalog: dndSpellsTestCatalog(t, dndSpellsCatalogSpecs{}),
})
if err == nil {
t.Fatal("Resolve() error = nil, want unknown lane error")
extractors := pipeline.NewExtractorRegistry()
extractorSpec.ArtifactKind = dnd.SpellListKind
if err := pipeline.RegisterExtractor[dnd.SpellList](extractors, extractorSpec, func() (contracts.Extractor[dnd.SpellList], error) {
return configExtractor{key: extractorSpec.Key}, nil
}); err != nil {
t.Fatalf("register capability extractor: %v", err)
}
if !strings.Contains(err.Error(), "selected artifact lane") || !strings.Contains(err.Error(), "missing") {
t.Fatalf("Resolve() error = %q, want unknown lane context", err.Error())
codecs := pipeline.NewArtifactCodecRegistry()
if err := pipeline.RegisterArtifactCodec(codecs, spellcodec.New()); err != nil {
t.Fatalf("register capability codec: %v", err)
}
mergers := pipeline.NewMergerRegistry()
if err := pipeline.RegisterMerger[dnd.SpellList](mergers, pipeline.ModuleSpec{
Key: pipeline.DefaultMergeModule, Stage: pipeline.StageMerge, ArtifactKind: dnd.SpellListKind, Requires: []string{"dnd.spell_casts"},
}, func() (contracts.Merger[dnd.SpellList], error) { return appendorder.NewTyped(appendSpellLists) }); err != nil {
t.Fatalf("register capability merger: %v", err)
}
normalizers := pipeline.NewNormalizerRegistry()
if err := pipeline.RegisterNormalizer[dnd.SpellList](normalizers, pipeline.ModuleSpec{Key: pipeline.DefaultNormalizeModule, Stage: pipeline.StageNormalize, ArtifactKind: dnd.SpellListKind}, func() (contracts.Normalizer[dnd.SpellList], error) {
return noop.NewTyped[dnd.SpellList](), nil
}); err != nil {
t.Fatalf("register capability normalizer: %v", err)
}
outputs := pipeline.NewOutputEncoderRegistry()
if err := outputs.RegisterWithSpec(pipeline.ModuleSpec{Key: pipeline.DefaultOutputModule, Stage: pipeline.StageOutput}, func() (contracts.OutputEncoder, error) { return dndSpellsOutput{}, nil }); err != nil {
t.Fatalf("register capability output: %v", err)
}
return pipeline.ModuleCatalog{
Inputs: inputs, Chunkers: chunkers, ArtifactCodecs: codecs, Extractors: extractors,
Mergers: mergers, Normalizers: normalizers, ValidatorChains: pipeline.NewValidatorChainRegistry(), Outputs: outputs,
}
}

View File

@@ -1,306 +0,0 @@
package transcript
import (
"context"
"encoding/json"
"os"
"reflect"
"strings"
"testing"
"gitea.maximumdirect.net/eric/notarius/internal/core/config"
"gitea.maximumdirect.net/eric/notarius/internal/core/source"
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
"gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline"
)
func TestPipelineConfigLoadsAndResolvesWithSeriatimInput(t *testing.T) {
cfg := loadPipelineConfig(t)
resolved, err := cfg.Resolve(config.ResolveInput{
PipelineID: "seriatim-fixture",
Catalog: seriatimTestCatalog(t, ModuleSpec()),
})
if err != nil {
t.Fatalf("Resolve() error = %v, want nil", err)
}
if resolved.ResolvedPipeline.Input.Module != Key {
t.Fatalf("resolved input module = %q, want %q", resolved.ResolvedPipeline.Input.Module, Key)
}
if resolved.ResolvedPipeline.Digest == "" {
t.Fatal("resolved digest is empty")
}
again, err := cfg.Resolve(config.ResolveInput{
PipelineID: "seriatim-fixture",
Catalog: seriatimTestCatalog(t, ModuleSpec()),
})
if err != nil {
t.Fatalf("second Resolve() error = %v, want nil", err)
}
if resolved.ResolvedPipeline.Digest != again.ResolvedPipeline.Digest {
t.Fatalf("resolved digest = %q, second digest = %q; want stable digest", resolved.ResolvedPipeline.Digest, again.ResolvedPipeline.Digest)
}
}
func TestPipelineConfigRejectsMissingSeriatimCapability(t *testing.T) {
spec := ModuleSpec()
spec.Provides = withoutCapability(spec.Provides, "transcript.timestamps")
cfg := loadPipelineConfig(t)
_, err := cfg.Resolve(config.ResolveInput{
PipelineID: "seriatim-fixture",
Catalog: seriatimTestCatalog(t, spec),
})
if err == nil {
t.Fatal("Resolve() error = nil, want missing capability error")
}
if !strings.Contains(err.Error(), "missing capability") || !strings.Contains(err.Error(), "transcript.timestamps") {
t.Fatalf("Resolve() error = %q, want missing transcript.timestamps capability", err.Error())
}
}
func TestPipelineConfigRejectsUnknownLaneSelection(t *testing.T) {
cfg := loadPipelineConfig(t)
_, err := cfg.Resolve(config.ResolveInput{
PipelineID: "seriatim-fixture",
Only: []string{"missing"},
Catalog: seriatimTestCatalog(t, ModuleSpec()),
})
if err == nil {
t.Fatal("Resolve() error = nil, want unknown lane error")
}
if !strings.Contains(err.Error(), "selected artifact lane") || !strings.Contains(err.Error(), "missing") {
t.Fatalf("Resolve() error = %q, want unknown lane context", err.Error())
}
}
func loadPipelineConfig(t *testing.T) config.Config {
t.Helper()
data, err := os.ReadFile("testdata/pipeline.yml")
if err != nil {
t.Fatalf("ReadFile(pipeline.yml) error = %v, want nil", err)
}
fileCfg, err := config.ParseFileConfigYAML(data)
if err != nil {
t.Fatalf("ParseFileConfigYAML() error = %v, want nil", err)
}
cfg := config.Default()
if err := cfg.ApplyFileConfig(fileCfg); err != nil {
t.Fatalf("ApplyFileConfig() error = %v, want nil", err)
}
profile, ok := cfg.Pipelines["seriatim-fixture"]
if !ok {
t.Fatal("pipeline seriatim-fixture was not loaded")
}
if profile.Input.Module != Key {
t.Fatalf("loaded input module = %q, want %q", profile.Input.Module, Key)
}
return cfg
}
func seriatimTestCatalog(t *testing.T, inputSpec pipeline.ModuleSpec) pipeline.ModuleCatalog {
t.Helper()
inputs := pipeline.NewInputAdapterRegistry()
chunkers := pipeline.NewChunkerRegistry()
extractors := pipeline.NewExtractorRegistry()
mergers := pipeline.NewMergerRegistry()
normalizers := pipeline.NewNormalizerRegistry()
outputs := pipeline.NewOutputEncoderRegistry()
if reflect.DeepEqual(inputSpec, ModuleSpec()) {
if err := Register(inputs); err != nil {
t.Fatalf("register seriatim input: %v", err)
}
} else if err := inputs.RegisterWithSpec(inputSpec, func() (contracts.InputAdapter, error) {
return New(), nil
}); err != nil {
t.Fatalf("register seriatim input override: %v", err)
}
mustRegisterChunker(t, chunkers, pipeline.ModuleSpec{
Key: "fake/chunk",
Stage: pipeline.StageChunk,
Requires: []string{"source.transcript"},
Provides: []string{"chunks"},
})
mustRegisterExtractor(t, extractors, pipeline.ModuleSpec{
Key: "fake/extract",
Stage: pipeline.StageExtract,
ArtifactKind: seriatimArtifactKind,
Requires: []string{"chunks", "transcript.speaker", "transcript.timestamps"},
Provides: []string{"fake.artifacts"},
})
mustRegisterMerger(t, mergers, pipeline.ModuleSpec{
Key: pipeline.DefaultMergeModule,
Stage: pipeline.StageMerge,
ArtifactKind: seriatimArtifactKind,
Requires: []string{"fake.artifacts"},
})
mustRegisterNormalizer(t, normalizers, pipeline.ModuleSpec{
Key: pipeline.DefaultNormalizeModule,
Stage: pipeline.StageNormalize,
ArtifactKind: seriatimArtifactKind,
})
mustRegisterOutput(t, outputs, pipeline.ModuleSpec{
Key: pipeline.DefaultOutputModule,
Stage: pipeline.StageOutput,
})
codecs := pipeline.NewArtifactCodecRegistry()
if err := pipeline.RegisterArtifactCodec(codecs, seriatimArtifactCodec{}); err != nil {
t.Fatalf("register artifact codec: %v", err)
}
return pipeline.ModuleCatalog{
Inputs: inputs,
Chunkers: chunkers,
ArtifactCodecs: codecs,
Extractors: extractors,
Mergers: mergers,
Normalizers: normalizers,
ValidatorChains: pipeline.NewValidatorChainRegistry(),
Outputs: outputs,
}
}
func mustRegisterChunker(t *testing.T, registry *pipeline.ChunkerRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := registry.RegisterWithSpec(spec, func() (contracts.Chunker, error) {
return fakeChunker{}, nil
}); err != nil {
t.Fatalf("register chunker: %v", err)
}
}
func mustRegisterExtractor(t *testing.T, registry *pipeline.ExtractorRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := pipeline.RegisterExtractor[seriatimArtifact](registry, spec, func() (contracts.Extractor[seriatimArtifact], error) {
return fakeExtractor{}, nil
}); err != nil {
t.Fatalf("register extractor: %v", err)
}
}
func mustRegisterMerger(t *testing.T, registry *pipeline.MergerRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := pipeline.RegisterMerger[seriatimArtifact](registry, spec, func() (contracts.Merger[seriatimArtifact], error) {
return fakeMerger{}, nil
}); err != nil {
t.Fatalf("register merger: %v", err)
}
}
func mustRegisterNormalizer(t *testing.T, registry *pipeline.NormalizerRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := pipeline.RegisterNormalizer[seriatimArtifact](registry, spec, func() (contracts.Normalizer[seriatimArtifact], error) {
return fakeNormalizer{}, nil
}); err != nil {
t.Fatalf("register normalizer: %v", err)
}
}
func mustRegisterOutput(t *testing.T, registry *pipeline.OutputEncoderRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := registry.RegisterWithSpec(spec, func() (contracts.OutputEncoder, error) {
return fakeOutput{}, nil
}); err != nil {
t.Fatalf("register output: %v", err)
}
}
type fakeChunker struct{}
func (fakeChunker) Key() string { return "fake/chunk" }
func (fakeChunker) ReferenceSlots() []contracts.ReferenceSlot { return nil }
func (fakeChunker) Plan(ctx context.Context, req contracts.ChunkRequest) (contracts.ChunkPlanResult, error) {
return contracts.ChunkPlanResult{}, nil
}
type fakeExtractor struct{}
func (fakeExtractor) Key() string { return "fake/extract" }
func (fakeExtractor) ReferenceSlots() []contracts.ReferenceSlot { return nil }
func (fakeExtractor) Extract(ctx context.Context, req contracts.TypedExtractionRequest) (contracts.TypedExtractionResult[seriatimArtifact], error) {
return contracts.TypedExtractionResult[seriatimArtifact]{}, nil
}
type fakeMerger struct{}
func (fakeMerger) Key() string { return pipeline.DefaultMergeModule }
func (fakeMerger) Merge(ctx context.Context, req contracts.TypedMergeRequest[seriatimArtifact]) (contracts.TypedMergeResult[seriatimArtifact], error) {
if len(req.ExtractOutputs) == 0 {
return contracts.TypedMergeResult[seriatimArtifact]{}, nil
}
return contracts.TypedMergeResult[seriatimArtifact]{Value: req.ExtractOutputs[0].Value}, nil
}
type fakeNormalizer struct{}
func (fakeNormalizer) Key() string { return pipeline.DefaultNormalizeModule }
func (fakeNormalizer) ReferenceSlots() []contracts.ReferenceSlot { return nil }
func (fakeNormalizer) Normalize(ctx context.Context, req contracts.TypedNormalizeRequest[seriatimArtifact]) (contracts.TypedNormalizeResult[seriatimArtifact], error) {
return contracts.TypedNormalizeResult[seriatimArtifact]{Value: req.MergeOutput.Value}, nil
}
type fakeOutput struct{}
func (fakeOutput) Key() string { return pipeline.DefaultOutputModule }
func (fakeOutput) Encode(ctx context.Context, req contracts.OutputRequest) (contracts.OutputResult, error) {
return contracts.OutputResult{}, nil
}
func withoutCapability(capabilities []string, capability string) []string {
filtered := make([]string, 0, len(capabilities))
for _, candidate := range capabilities {
if candidate != capability {
filtered = append(filtered, candidate)
}
}
return filtered
}
var (
_ contracts.Chunker = fakeChunker{}
_ contracts.Extractor[seriatimArtifact] = fakeExtractor{}
_ contracts.OutputEncoder = fakeOutput{}
)
const seriatimArtifactKind contracts.ArtifactKind = "test/seriatim-event"
type seriatimArtifact struct {
Value string `json:"value"`
SourceRefs []source.SourceRef `json:"source_refs"`
}
type seriatimArtifactCodec struct{}
func (seriatimArtifactCodec) Kind() contracts.ArtifactKind { return seriatimArtifactKind }
func (seriatimArtifactCodec) Schema() contracts.ArtifactSchema {
return contracts.ArtifactSchema{ID: "fake.event", Name: "fake_event", Version: "v1", JSONSchema: []byte(`{"type":"object"}`)}
}
func (seriatimArtifactCodec) MediaType() string { return "application/json" }
func (seriatimArtifactCodec) EncodeCandidate(value seriatimArtifact) ([]byte, error) {
return json.Marshal(value)
}
func (seriatimArtifactCodec) Encode(value seriatimArtifact) ([]byte, error) {
return json.Marshal(value)
}
func (seriatimArtifactCodec) Decode(content []byte) (seriatimArtifact, error) {
var value seriatimArtifact
err := json.Unmarshal(content, &value)
return value, err
}

View File

@@ -0,0 +1,166 @@
package transcript
import (
"context"
"encoding/json"
"reflect"
"testing"
"gitea.maximumdirect.net/eric/notarius/internal/core/config"
"gitea.maximumdirect.net/eric/notarius/internal/core/source"
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
"gitea.maximumdirect.net/eric/notarius/internal/framework/pipeline"
)
func loadPipelineConfig(t *testing.T) config.Config {
t.Helper()
cfg := config.Default()
cfg.Pipelines["seriatim-fixture"] = pipeline.PipelineProfile{
ID: "seriatim-fixture",
Input: pipeline.Binding(Key),
Chunk: pipeline.Binding("fake/chunk"),
Artifacts: map[string]pipeline.ArtifactLaneProfile{
"events": {Extract: pipeline.Binding("fake/extract")},
},
}
return cfg
}
func seriatimTestCatalog(t *testing.T, inputSpec pipeline.ModuleSpec) pipeline.ModuleCatalog {
t.Helper()
inputs := pipeline.NewInputAdapterRegistry()
chunkers := pipeline.NewChunkerRegistry()
extractors := pipeline.NewExtractorRegistry()
mergers := pipeline.NewMergerRegistry()
normalizers := pipeline.NewNormalizerRegistry()
outputs := pipeline.NewOutputEncoderRegistry()
if reflect.DeepEqual(inputSpec, ModuleSpec()) {
if err := Register(inputs); err != nil {
t.Fatalf("register seriatim input: %v", err)
}
} else if err := inputs.RegisterWithSpec(inputSpec, func() (contracts.InputAdapter, error) {
return New(), nil
}); err != nil {
t.Fatalf("register seriatim input override: %v", err)
}
mustRegisterChunker(t, chunkers, pipeline.ModuleSpec{Key: "fake/chunk", Stage: pipeline.StageChunk, Requires: []string{"source.transcript"}, Provides: []string{"chunks"}})
mustRegisterExtractor(t, extractors, pipeline.ModuleSpec{
Key: "fake/extract", Stage: pipeline.StageExtract, ArtifactKind: seriatimArtifactKind,
Requires: []string{"chunks", "transcript.speaker", "transcript.timestamps"}, Provides: []string{"fake.artifacts"},
})
mustRegisterMerger(t, mergers, pipeline.ModuleSpec{Key: pipeline.DefaultMergeModule, Stage: pipeline.StageMerge, ArtifactKind: seriatimArtifactKind, Requires: []string{"fake.artifacts"}})
mustRegisterNormalizer(t, normalizers, pipeline.ModuleSpec{Key: pipeline.DefaultNormalizeModule, Stage: pipeline.StageNormalize, ArtifactKind: seriatimArtifactKind})
mustRegisterOutput(t, outputs, pipeline.ModuleSpec{Key: pipeline.DefaultOutputModule, Stage: pipeline.StageOutput})
codecs := pipeline.NewArtifactCodecRegistry()
if err := pipeline.RegisterArtifactCodec(codecs, seriatimArtifactCodec{}); err != nil {
t.Fatalf("register artifact codec: %v", err)
}
return pipeline.ModuleCatalog{Inputs: inputs, Chunkers: chunkers, ArtifactCodecs: codecs, Extractors: extractors, Mergers: mergers, Normalizers: normalizers, ValidatorChains: pipeline.NewValidatorChainRegistry(), Outputs: outputs}
}
func mustRegisterChunker(t *testing.T, registry *pipeline.ChunkerRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := registry.RegisterWithSpec(spec, func() (contracts.Chunker, error) { return fakeChunker{}, nil }); err != nil {
t.Fatalf("register chunker: %v", err)
}
}
func mustRegisterExtractor(t *testing.T, registry *pipeline.ExtractorRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := pipeline.RegisterExtractor[seriatimArtifact](registry, spec, func() (contracts.Extractor[seriatimArtifact], error) { return fakeExtractor{}, nil }); err != nil {
t.Fatalf("register extractor: %v", err)
}
}
func mustRegisterMerger(t *testing.T, registry *pipeline.MergerRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := pipeline.RegisterMerger[seriatimArtifact](registry, spec, func() (contracts.Merger[seriatimArtifact], error) { return fakeMerger{}, nil }); err != nil {
t.Fatalf("register merger: %v", err)
}
}
func mustRegisterNormalizer(t *testing.T, registry *pipeline.NormalizerRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := pipeline.RegisterNormalizer[seriatimArtifact](registry, spec, func() (contracts.Normalizer[seriatimArtifact], error) { return fakeNormalizer{}, nil }); err != nil {
t.Fatalf("register normalizer: %v", err)
}
}
func mustRegisterOutput(t *testing.T, registry *pipeline.OutputEncoderRegistry, spec pipeline.ModuleSpec) {
t.Helper()
if err := registry.RegisterWithSpec(spec, func() (contracts.OutputEncoder, error) { return fakeOutput{}, nil }); err != nil {
t.Fatalf("register output: %v", err)
}
}
type fakeChunker struct{}
func (fakeChunker) Key() string { return "fake/chunk" }
func (fakeChunker) ReferenceSlots() []contracts.ReferenceSlot { return nil }
func (fakeChunker) Plan(context.Context, contracts.ChunkRequest) (contracts.ChunkPlanResult, error) {
return contracts.ChunkPlanResult{}, nil
}
type fakeExtractor struct{}
func (fakeExtractor) Key() string { return "fake/extract" }
func (fakeExtractor) ReferenceSlots() []contracts.ReferenceSlot { return nil }
func (fakeExtractor) Extract(context.Context, contracts.TypedExtractionRequest) (contracts.TypedExtractionResult[seriatimArtifact], error) {
return contracts.TypedExtractionResult[seriatimArtifact]{}, nil
}
type fakeMerger struct{}
func (fakeMerger) Key() string { return pipeline.DefaultMergeModule }
func (fakeMerger) Merge(_ context.Context, req contracts.TypedMergeRequest[seriatimArtifact]) (contracts.TypedMergeResult[seriatimArtifact], error) {
if len(req.ExtractOutputs) == 0 {
return contracts.TypedMergeResult[seriatimArtifact]{}, nil
}
return contracts.TypedMergeResult[seriatimArtifact]{Value: req.ExtractOutputs[0].Value}, nil
}
type fakeNormalizer struct{}
func (fakeNormalizer) Key() string { return pipeline.DefaultNormalizeModule }
func (fakeNormalizer) ReferenceSlots() []contracts.ReferenceSlot { return nil }
func (fakeNormalizer) Normalize(_ context.Context, req contracts.TypedNormalizeRequest[seriatimArtifact]) (contracts.TypedNormalizeResult[seriatimArtifact], error) {
return contracts.TypedNormalizeResult[seriatimArtifact]{Value: req.MergeOutput.Value}, nil
}
type fakeOutput struct{}
func (fakeOutput) Key() string { return pipeline.DefaultOutputModule }
func (fakeOutput) Encode(context.Context, contracts.OutputRequest) (contracts.OutputResult, error) {
return contracts.OutputResult{}, nil
}
const seriatimArtifactKind contracts.ArtifactKind = "test/seriatim-event"
type seriatimArtifact struct {
Value string `json:"value"`
SourceRefs []source.SourceRef `json:"source_refs"`
}
type seriatimArtifactCodec struct{}
func (seriatimArtifactCodec) Kind() contracts.ArtifactKind { return seriatimArtifactKind }
func (seriatimArtifactCodec) Schema() contracts.ArtifactSchema {
return contracts.ArtifactSchema{ID: "fake.event", Name: "fake_event", Version: "v1", JSONSchema: []byte(`{"type":"object"}`)}
}
func (seriatimArtifactCodec) MediaType() string { return "application/json" }
func (seriatimArtifactCodec) EncodeCandidate(value seriatimArtifact) ([]byte, error) {
return json.Marshal(value)
}
func (seriatimArtifactCodec) Encode(value seriatimArtifact) ([]byte, error) {
return json.Marshal(value)
}
func (seriatimArtifactCodec) Decode(content []byte) (seriatimArtifact, error) {
var value seriatimArtifact
err := json.Unmarshal(content, &value)
return value, err
}

View File

@@ -1,19 +0,0 @@
version: 3
output:
directory: ./notarius-output
cache:
chunk_plans:
mode: bypass
checkpoints: {}
debug:
directory: ./notarius-debug
pipelines:
seriatim-fixture:
input: seriatim
chunk: fake/chunk
artifacts:
events:
extract: fake/extract
merge: appendorder
normalize: noop
output: json