package pipeline_test import ( "context" "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" "gitea.maximumdirect.net/eric/notarius/internal/modules/chunk/generic" "gitea.maximumdirect.net/eric/notarius/internal/modules/merge/appendorder" "gitea.maximumdirect.net/eric/notarius/internal/modules/normalize/noop" jsonoutput "gitea.maximumdirect.net/eric/notarius/internal/modules/output/json" ) func TestPipelineConfigResolvesWithProductionDefaultsRegistered(t *testing.T) { cfg := config.Default() cfg.Pipelines = map[string]pipeline.PipelineProfile{ "defaults": { Input: pipeline.Binding("input"), Artifacts: map[string]pipeline.ArtifactLaneProfile{ "events": {Extract: pipeline.Binding("extract")}, }, }, } resolved, err := cfg.Resolve(config.ResolveInput{ PipelineID: "defaults", Catalog: defaultModuleCatalog(t), }) if err != nil { t.Fatalf("Resolve() error = %v, want nil", err) } pipeline := resolved.ResolvedPipeline if pipeline.Chunk.Module != generic.Key { t.Fatalf("Chunk.Module = %q, want %q", pipeline.Chunk.Module, generic.Key) } if pipeline.Output.Module != jsonoutput.Key { t.Fatalf("Output.Module = %q, want %q", pipeline.Output.Module, jsonoutput.Key) } lane := pipeline.ArtifactLanes[0] if lane.Merge.Module != appendorder.Key { t.Fatalf("Merge.Module = %q, want %q", lane.Merge.Module, appendorder.Key) } if lane.Normalize.Module != noop.Key { t.Fatalf("Normalize.Module = %q, want %q", lane.Normalize.Module, noop.Key) } } func defaultModuleCatalog(t *testing.T) pipeline.ModuleCatalog { t.Helper() inputs := pipeline.NewInputAdapterRegistry() chunkers := pipeline.NewChunkerRegistry() extractors := pipeline.NewExtractorRegistry() mergers := pipeline.NewMergerRegistry() normalizers := pipeline.NewNormalizerRegistry() outputs := pipeline.NewOutputEncoderRegistry() if err := inputs.RegisterWithSpec(pipeline.ModuleSpec{ Key: "input", Stage: pipeline.StageInput, Provides: []string{"source"}, }, func() (contracts.InputAdapter, error) { return defaultInput{}, nil }); err != nil { t.Fatalf("register input: %v", err) } if err := generic.Register(chunkers); err != nil { t.Fatalf("register generic chunker: %v", err) } if err := extractors.RegisterWithSpec(pipeline.ModuleSpec{ Key: "extract", Stage: pipeline.StageExtract, Requires: []string{"chunks"}, Provides: []string{"records"}, }, func() (contracts.Extractor, error) { return defaultExtractor{}, nil }); err != nil { t.Fatalf("register extractor: %v", err) } if err := appendorder.Register(mergers); err != nil { t.Fatalf("register appendorder merger: %v", err) } if err := noop.Register(normalizers); err != nil { t.Fatalf("register noop normalizer: %v", err) } if err := jsonoutput.Register(outputs); err != nil { t.Fatalf("register json output: %v", err) } return pipeline.ModuleCatalog{ Inputs: inputs, Chunkers: chunkers, Extractors: extractors, Mergers: mergers, Normalizers: normalizers, Outputs: outputs, } } type defaultInput struct{} func (defaultInput) Key() string { return "input" } func (defaultInput) Parse(ctx context.Context, req contracts.ParseRequest) (*source.SourceDocument, error) { return nil, nil } type defaultExtractor struct{} func (defaultExtractor) Key() string { return "extract" } func (defaultExtractor) ArtifactType() string { return "record" } func (defaultExtractor) SchemaVersion() string { return "v1" } func (defaultExtractor) ReferenceSlots() []contracts.ReferenceSlot { return nil } func (defaultExtractor) Validators() []contracts.Validator { return nil } func (defaultExtractor) Extract(ctx context.Context, req contracts.ExtractionRequest) (contracts.ExtractionResult, error) { return contracts.ExtractionResult{}, nil }