127 lines
3.9 KiB
Go
127 lines
3.9 KiB
Go
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
|
|
}
|