package pipeline_test import ( "context" "encoding/json" "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/generic/chunk/units" "gitea.maximumdirect.net/eric/notarius/internal/modules/generic/merge/appendorder" "gitea.maximumdirect.net/eric/notarius/internal/modules/generic/normalize/noop" jsonoutput "gitea.maximumdirect.net/eric/notarius/internal/modules/generic/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 != units.Key { t.Fatalf("Chunk.Module = %q, want %q", pipeline.Chunk.Module, units.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 := units.Register(chunkers); err != nil { t.Fatalf("register generic chunker: %v", err) } if err := pipeline.RegisterExtractor[defaultArtifact](extractors, pipeline.ModuleSpec{ Key: "extract", Stage: pipeline.StageExtract, ArtifactKind: defaultArtifactKind, Requires: []string{"chunks"}, Provides: []string{"records"}, }, func() (contracts.Extractor[defaultArtifact], error) { return defaultExtractor{}, nil }); err != nil { t.Fatalf("register extractor: %v", err) } if err := appendorder.RegisterTyped(mergers, defaultArtifactKind, func(values []defaultArtifact) (defaultArtifact, error) { if len(values) == 0 { return defaultArtifact{}, nil } return values[0], nil }); err != nil { t.Fatalf("register appendorder merger: %v", err) } if err := noop.RegisterTyped[defaultArtifact](normalizers, defaultArtifactKind); err != nil { t.Fatalf("register noop normalizer: %v", err) } if err := jsonoutput.Register(outputs); err != nil { t.Fatalf("register json output: %v", err) } codecs := pipeline.NewArtifactCodecRegistry() if err := pipeline.RegisterArtifactCodec(codecs, defaultArtifactCodec{}); 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, } } 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) ReferenceSlots() []contracts.ReferenceSlot { return nil } func (defaultExtractor) Extract(ctx context.Context, req contracts.TypedExtractionRequest) (contracts.TypedExtractionResult[defaultArtifact], error) { return contracts.TypedExtractionResult[defaultArtifact]{}, nil } const defaultArtifactKind contracts.ArtifactKind = "test/default" type defaultArtifact struct { Value string `json:"value"` } type defaultArtifactCodec struct{} func (defaultArtifactCodec) Kind() contracts.ArtifactKind { return defaultArtifactKind } func (defaultArtifactCodec) Schema() contracts.ArtifactSchema { return contracts.ArtifactSchema{ID: "urn:notarius:test:default", Name: "default", Version: "1", JSONSchema: []byte(`{"type":"object"}`)} } func (defaultArtifactCodec) MediaType() string { return "application/json" } func (defaultArtifactCodec) EncodeCandidate(value defaultArtifact) ([]byte, error) { return json.Marshal(value) } func (defaultArtifactCodec) Encode(value defaultArtifact) ([]byte, error) { return json.Marshal(value) } func (defaultArtifactCodec) Decode(content []byte) (defaultArtifact, error) { var value defaultArtifact err := json.Unmarshal(content, &value) return value, err }