354 lines
14 KiB
Go
354 lines
14 KiB
Go
package pipeline
|
|
|
|
import (
|
|
"fmt"
|
|
"reflect"
|
|
"strings"
|
|
|
|
"gitea.maximumdirect.net/eric/notarius/internal/framework/contracts"
|
|
)
|
|
|
|
// PreparedPipeline owns the constructed, run-local implementation set for one
|
|
// resolved pipeline. Its implementation values are private so execution cannot
|
|
// replace or reconfigure them after preparation.
|
|
type PreparedPipeline struct {
|
|
Input ModuleBinding
|
|
Chunk ModuleBinding
|
|
ArtifactLanes []PreparedArtifactLane
|
|
Output ModuleBinding
|
|
|
|
resolved ResolvedPipeline
|
|
dependencies ModuleDependencies
|
|
input contracts.InputAdapter
|
|
chunker contracts.Chunker
|
|
chunkValidators preparedValidatorChain
|
|
lanes []preparedLaneExecutor
|
|
output contracts.OutputEncoder
|
|
checkpointFingerprints []CheckpointFingerprint
|
|
}
|
|
|
|
type PreparedArtifactLane struct {
|
|
Resolved ResolvedArtifactLane
|
|
}
|
|
|
|
type preparedLaneExecutor struct {
|
|
resolved ResolvedArtifactLane
|
|
typed *preparedTypedLane
|
|
extractValidators preparedValidatorChain
|
|
mergeValidators preparedValidatorChain
|
|
normalizeValidators preparedValidatorChain
|
|
}
|
|
|
|
type preparedTypedLane struct {
|
|
extractor any
|
|
merger any
|
|
normalizer any
|
|
extract typedExtractOperation
|
|
merge typedMergeOperation
|
|
normalize typedNormalizeOperation
|
|
codec artifactCodecEntry
|
|
}
|
|
|
|
type preparedValidatorChain struct {
|
|
resolved ResolvedValidatorChain
|
|
validators []preparedValidator
|
|
}
|
|
|
|
type preparedValidator struct {
|
|
resolved ResolvedValidator
|
|
typed any
|
|
typedValidate typedValidateOperation
|
|
chunk contracts.ChunkValidator
|
|
serialized contracts.SerializedValidator
|
|
}
|
|
|
|
// Prepare validates all configured options and constructs every selected
|
|
// module and validator before any operation method can run.
|
|
func Prepare(resolved ResolvedPipeline, registries Registries, deps ModuleDependencies) (*PreparedPipeline, error) {
|
|
if err := validateResolvedPipeline(resolved); err != nil {
|
|
return nil, err
|
|
}
|
|
if err := validateRegistrySet(resolved, registries); err != nil {
|
|
return nil, err
|
|
}
|
|
stable := cloneResolvedPipeline(resolved)
|
|
prepared := &PreparedPipeline{
|
|
Input: cloneModuleBinding(stable.Input),
|
|
Chunk: cloneModuleBinding(stable.Chunk),
|
|
Output: cloneModuleBinding(stable.Output),
|
|
resolved: stable,
|
|
dependencies: deps,
|
|
}
|
|
request := func(binding ModuleBinding, references contracts.ReferenceSet) BuildRequest {
|
|
return BuildRequest{Dependencies: deps, Options: cloneOptions(binding.Options), References: references}
|
|
}
|
|
|
|
input, err := registries.Inputs.BuildWithRequest(stable.Input.Module, request(stable.Input, contracts.ReferenceSet{}))
|
|
if err != nil {
|
|
return nil, constructionError(stable.ID, "", StageInput, stable.Input.Module, "", err)
|
|
}
|
|
prepared.input = input
|
|
|
|
chunker, err := registries.Chunkers.BuildWithRequest(stable.Chunk.Module, request(stable.Chunk, stable.ChunkReferences.ReferenceSet))
|
|
if err != nil {
|
|
return nil, constructionError(stable.ID, "", StageChunk, stable.Chunk.Module, "", err)
|
|
}
|
|
prepared.chunker = chunker
|
|
prepared.chunkValidators, err = prepareValidatorChain(stable, registries, deps, StageChunk, "", stable.Chunk.Module, stable.ChunkReferences.ReferenceSet)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
prepared.ArtifactLanes = make([]PreparedArtifactLane, 0, len(stable.ArtifactLanes))
|
|
prepared.lanes = make([]preparedLaneExecutor, 0, len(stable.ArtifactLanes))
|
|
for _, lane := range stable.ArtifactLanes {
|
|
executor, err := prepareLane(stable, lane, registries, deps)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
prepared.ArtifactLanes = append(prepared.ArtifactLanes, PreparedArtifactLane{Resolved: cloneResolvedArtifactLane(lane)})
|
|
prepared.lanes = append(prepared.lanes, executor)
|
|
}
|
|
|
|
output, err := registries.Outputs.BuildWithRequest(stable.Output.Module, request(stable.Output, contracts.ReferenceSet{}))
|
|
if err != nil {
|
|
return nil, constructionError(stable.ID, "", StageOutput, stable.Output.Module, "", err)
|
|
}
|
|
prepared.output = output
|
|
prepared.checkpointFingerprints, err = collectPreparedCheckpointFingerprints(prepared)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
return prepared, nil
|
|
}
|
|
|
|
func prepareLane(pipeline ResolvedPipeline, lane ResolvedArtifactLane, registries Registries, deps ModuleDependencies) (preparedLaneExecutor, error) {
|
|
executor := preparedLaneExecutor{resolved: cloneResolvedArtifactLane(lane)}
|
|
request := func(binding ModuleBinding, references contracts.ReferenceSet) BuildRequest {
|
|
return BuildRequest{Dependencies: deps, Options: cloneOptions(binding.Options), References: references}
|
|
}
|
|
extractEntry, ok := registries.Extractors.typedEntry(lane.Extract.Module)
|
|
if !ok {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageExtract, lane.Extract.Module, "", fmt.Errorf("typed construction entry is not registered"))
|
|
}
|
|
module, err := buildErasedModule(extractEntry.builder, request(lane.Extract, lane.ExtractReferences.ReferenceSet), lane.Extract.Module, "extractor")
|
|
if err != nil {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageExtract, lane.Extract.Module, "", err)
|
|
}
|
|
codec, _, codecErr := registries.ArtifactCodecs.entry(lane.ArtifactKind)
|
|
if codecErr != nil {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageExtract, lane.Extract.Module, "", codecErr)
|
|
}
|
|
executor.typed = &preparedTypedLane{extractor: module, extract: extractEntry.extract, codec: codec}
|
|
|
|
executor.extractValidators, err = prepareValidatorChain(pipeline, registries, deps, StageExtract, lane.ID, lane.Extract.Module, lane.ExtractReferences.ReferenceSet)
|
|
if err != nil {
|
|
return preparedLaneExecutor{}, err
|
|
}
|
|
|
|
mergeEntry, ok := registries.Mergers.typedEntry(lane.Merge.Module, lane.ArtifactKind)
|
|
if !ok {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageMerge, lane.Merge.Module, "", fmt.Errorf("typed construction entry is not registered"))
|
|
}
|
|
module, err = buildErasedModule(mergeEntry.builder, request(lane.Merge, lane.MergeReferences.ReferenceSet), lane.Merge.Module, "merger")
|
|
if err != nil {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageMerge, lane.Merge.Module, "", err)
|
|
}
|
|
executor.typed.merger = module
|
|
executor.typed.merge = mergeEntry.merge
|
|
executor.mergeValidators, err = prepareValidatorChain(pipeline, registries, deps, StageMerge, lane.ID, lane.Merge.Module, lane.MergeReferences.ReferenceSet)
|
|
if err != nil {
|
|
return preparedLaneExecutor{}, err
|
|
}
|
|
|
|
normalizeEntry, ok := registries.Normalizers.typedEntry(lane.Normalize.Module, lane.ArtifactKind)
|
|
if !ok {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageNormalize, lane.Normalize.Module, "", fmt.Errorf("typed construction entry is not registered"))
|
|
}
|
|
module, err = buildErasedModule(normalizeEntry.builder, request(lane.Normalize, lane.NormalizeReferences.ReferenceSet), lane.Normalize.Module, "normalizer")
|
|
if err != nil {
|
|
return preparedLaneExecutor{}, constructionError(pipeline.ID, lane.ID, StageNormalize, lane.Normalize.Module, "", err)
|
|
}
|
|
executor.typed.normalizer = module
|
|
executor.typed.normalize = normalizeEntry.normalize
|
|
executor.normalizeValidators, err = prepareValidatorChain(pipeline, registries, deps, StageNormalize, lane.ID, lane.Normalize.Module, lane.NormalizeReferences.ReferenceSet)
|
|
if err != nil {
|
|
return preparedLaneExecutor{}, err
|
|
}
|
|
return executor, nil
|
|
}
|
|
|
|
func prepareValidatorChain(pipeline ResolvedPipeline, registries Registries, deps ModuleDependencies, stage ModuleStage, laneID, moduleKey string, references contracts.ReferenceSet) (preparedValidatorChain, error) {
|
|
resolved := resolvedValidatorChain(stage, laneID, moduleKey, pipeline.ValidatorChains)
|
|
prepared := preparedValidatorChain{resolved: resolved}
|
|
for _, validator := range resolved.Validators {
|
|
request := BuildRequest{Dependencies: deps, Options: cloneOptions(validator.Binding.Options), References: references}
|
|
built, err := buildPreparedValidator(registries.Validators, validator, request)
|
|
if err != nil {
|
|
return preparedValidatorChain{}, constructionError(pipeline.ID, laneID, stage, moduleKey, validator.Binding.Module, err)
|
|
}
|
|
prepared.validators = append(prepared.validators, built)
|
|
}
|
|
return prepared, nil
|
|
}
|
|
|
|
func buildPreparedValidator(registry *ValidatorRegistry, resolved ResolvedValidator, request BuildRequest) (preparedValidator, error) {
|
|
prepared := preparedValidator{resolved: resolved}
|
|
key := resolved.Binding.Module
|
|
var implementation any
|
|
var err error
|
|
switch resolved.Target {
|
|
case ValidatorTargetTyped:
|
|
entry, ok := registry.typedEntry(key, resolved.ArtifactKind)
|
|
if !ok {
|
|
return preparedValidator{}, fmt.Errorf("typed construction entry is not registered")
|
|
}
|
|
implementation, err = entry.builder(cloneBuildRequest(request))
|
|
prepared.typed = implementation
|
|
prepared.typedValidate = entry.validate
|
|
case ValidatorTargetChunk:
|
|
entry, ok := registry.chunkEntry(key)
|
|
if !ok {
|
|
return preparedValidator{}, fmt.Errorf("chunk construction entry is not registered")
|
|
}
|
|
prepared.chunk, err = entry.builder(cloneBuildRequest(request))
|
|
implementation = prepared.chunk
|
|
case ValidatorTargetSerialized:
|
|
entry, ok := registry.serializedEntry(key)
|
|
if !ok {
|
|
return preparedValidator{}, fmt.Errorf("serialized construction entry is not registered")
|
|
}
|
|
prepared.serialized, err = entry.builder(cloneBuildRequest(request))
|
|
implementation = prepared.serialized
|
|
default:
|
|
return preparedValidator{}, fmt.Errorf("validator construction target %q is not supported", resolved.Target)
|
|
}
|
|
if err != nil {
|
|
return preparedValidator{}, err
|
|
}
|
|
if isNilImplementation(implementation) {
|
|
return preparedValidator{}, fmt.Errorf("validator %q builder returned nil", key)
|
|
}
|
|
identity, ok := implementation.(interface {
|
|
Name() string
|
|
ExecutionClass() contracts.ExecutionClass
|
|
})
|
|
if !ok {
|
|
return preparedValidator{}, fmt.Errorf("validator %q builder returned incompatible implementation %T", key, implementation)
|
|
}
|
|
if identity.Name() != key {
|
|
return preparedValidator{}, fmt.Errorf("validator %q returned name %q", key, identity.Name())
|
|
}
|
|
if identity.ExecutionClass() != resolved.ExecutionClass {
|
|
return preparedValidator{}, fmt.Errorf("validator %q returned execution class %q, want %q", key, identity.ExecutionClass(), resolved.ExecutionClass)
|
|
}
|
|
return prepared, nil
|
|
}
|
|
|
|
func buildErasedModule(builder func(BuildRequest) (any, error), request BuildRequest, key, kind string) (any, error) {
|
|
implementation, err := builder(cloneBuildRequest(request))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
if isNilImplementation(implementation) {
|
|
return nil, fmt.Errorf("%s %q builder returned nil", kind, key)
|
|
}
|
|
identity, ok := implementation.(interface{ Key() string })
|
|
if !ok {
|
|
return nil, fmt.Errorf("%s %q builder returned incompatible implementation %T", kind, key, implementation)
|
|
}
|
|
if identity.Key() != key {
|
|
return nil, fmt.Errorf("%s %q returned key %q", kind, key, identity.Key())
|
|
}
|
|
return implementation, nil
|
|
}
|
|
|
|
func isNilImplementation(value any) bool {
|
|
if value == nil {
|
|
return true
|
|
}
|
|
reflected := reflect.ValueOf(value)
|
|
switch reflected.Kind() {
|
|
case reflect.Chan, reflect.Func, reflect.Interface, reflect.Map, reflect.Pointer, reflect.Slice:
|
|
return reflected.IsNil()
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
func constructionError(pipelineID, laneID string, stage ModuleStage, moduleKey, validatorKey string, cause error) error {
|
|
scope := fmt.Sprintf("pipeline %q %s module %q", pipelineID, stage, moduleKey)
|
|
if laneID != "" {
|
|
scope = fmt.Sprintf("pipeline %q lane %q %s module %q", pipelineID, laneID, stage, moduleKey)
|
|
}
|
|
if strings.TrimSpace(validatorKey) != "" {
|
|
scope += fmt.Sprintf(" validator %q", validatorKey)
|
|
}
|
|
return fmt.Errorf("prepare %s: %w", scope, cause)
|
|
}
|
|
|
|
func (registries Registries) catalog() ModuleCatalog {
|
|
return ModuleCatalog{
|
|
Inputs: registries.Inputs, Chunkers: registries.Chunkers, ArtifactCodecs: registries.ArtifactCodecs,
|
|
Extractors: registries.Extractors, Mergers: registries.Mergers, Normalizers: registries.Normalizers,
|
|
Validators: registries.Validators, ValidatorChains: registries.ValidatorChains, Outputs: registries.Outputs,
|
|
}
|
|
}
|
|
|
|
func validateRegistrySet(resolved ResolvedPipeline, registries Registries) error {
|
|
if registries.Inputs == nil {
|
|
return fmt.Errorf("input registry must not be nil")
|
|
}
|
|
if registries.Chunkers == nil {
|
|
return fmt.Errorf("chunker registry must not be nil")
|
|
}
|
|
if registries.Extractors == nil {
|
|
return fmt.Errorf("extractor registry must not be nil")
|
|
}
|
|
if registries.Mergers == nil {
|
|
return fmt.Errorf("merger registry must not be nil")
|
|
}
|
|
if registries.Normalizers == nil {
|
|
return fmt.Errorf("normalizer registry must not be nil")
|
|
}
|
|
if registries.Validators == nil {
|
|
for _, chain := range resolved.ValidatorChains {
|
|
if len(chain.Validators) > 0 {
|
|
return fmt.Errorf("validator registry must not be nil")
|
|
}
|
|
}
|
|
}
|
|
if registries.Outputs == nil {
|
|
return fmt.Errorf("output encoder registry must not be nil")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func cloneResolvedPipeline(in ResolvedPipeline) ResolvedPipeline {
|
|
out := in
|
|
out.Input = cloneModuleBinding(in.Input)
|
|
out.Chunk = cloneModuleBinding(in.Chunk)
|
|
out.Output = cloneModuleBinding(in.Output)
|
|
out.ChunkReferences = CloneReferenceTarget(in.ChunkReferences)
|
|
out.ValidatorChains = cloneResolvedValidatorChains(in.ValidatorChains)
|
|
if len(in.ArtifactLanes) > 0 {
|
|
out.ArtifactLanes = make([]ResolvedArtifactLane, len(in.ArtifactLanes))
|
|
for i, lane := range in.ArtifactLanes {
|
|
out.ArtifactLanes[i] = cloneResolvedArtifactLane(lane)
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
func cloneResolvedArtifactLane(in ResolvedArtifactLane) ResolvedArtifactLane {
|
|
out := in
|
|
out.Extract = cloneModuleBinding(in.Extract)
|
|
out.Merge = cloneModuleBinding(in.Merge)
|
|
out.Normalize = cloneModuleBinding(in.Normalize)
|
|
out.Validators = cloneModuleBindings(in.Validators)
|
|
out.ExtractReferences = CloneReferenceTarget(in.ExtractReferences)
|
|
out.MergeReferences = CloneReferenceTarget(in.MergeReferences)
|
|
out.NormalizeReferences = CloneReferenceTarget(in.NormalizeReferences)
|
|
return out
|
|
}
|