405 lines
12 KiB
Go
405 lines
12 KiB
Go
package pipeline
|
|
|
|
import (
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"fmt"
|
|
"sort"
|
|
"strings"
|
|
)
|
|
|
|
const (
|
|
DefaultChunkModule = "generic"
|
|
DefaultMergeModule = "appendorder"
|
|
DefaultNormalizeModule = "noop"
|
|
DefaultOutputModule = "json"
|
|
DefaultLLMProfile = "default"
|
|
)
|
|
|
|
type ModuleBinding struct {
|
|
Module string `json:"module"`
|
|
LLMProfile string `json:"llm_profile,omitempty"`
|
|
Options map[string]any `json:"options,omitempty"`
|
|
}
|
|
|
|
type ArtifactLaneProfile struct {
|
|
Extract ModuleBinding `json:"extract"`
|
|
Merge ModuleBinding `json:"merge,omitempty"`
|
|
Normalize ModuleBinding `json:"normalize,omitempty"`
|
|
Validators []ModuleBinding `json:"validators,omitempty"`
|
|
}
|
|
|
|
type PipelineProfile struct {
|
|
ID string `json:"id"`
|
|
Input ModuleBinding `json:"input"`
|
|
Chunk ModuleBinding `json:"chunk,omitempty"`
|
|
Artifacts map[string]ArtifactLaneProfile `json:"artifacts"`
|
|
Output ModuleBinding `json:"output,omitempty"`
|
|
}
|
|
|
|
type ResolveOptions struct {
|
|
Only []string
|
|
}
|
|
|
|
type ResolvedArtifactLane struct {
|
|
ID string
|
|
Extract ModuleBinding
|
|
Merge ModuleBinding
|
|
Normalize ModuleBinding
|
|
Validators []ModuleBinding
|
|
}
|
|
|
|
type ResolvedPipeline struct {
|
|
ID string
|
|
Digest string
|
|
Input ModuleBinding
|
|
Chunk ModuleBinding
|
|
ArtifactLanes []ResolvedArtifactLane
|
|
Output ModuleBinding
|
|
}
|
|
|
|
type ModuleCatalog struct {
|
|
Inputs *InputAdapterRegistry
|
|
Chunkers *ChunkerRegistry
|
|
Extractors *ExtractorRegistry
|
|
Mergers *MergerRegistry
|
|
Normalizers *NormalizerRegistry
|
|
Validators *ValidatorRegistry
|
|
Outputs *OutputEncoderRegistry
|
|
}
|
|
|
|
func Binding(module string) ModuleBinding {
|
|
return ModuleBinding{Module: strings.TrimSpace(module)}
|
|
}
|
|
|
|
func ResolvePipeline(profile PipelineProfile, options ResolveOptions, catalog ModuleCatalog) (ResolvedPipeline, error) {
|
|
pipelineID := strings.TrimSpace(profile.ID)
|
|
if pipelineID == "" {
|
|
return ResolvedPipeline{}, fmt.Errorf("pipeline id must not be empty")
|
|
}
|
|
|
|
if len(profile.Artifacts) == 0 {
|
|
return ResolvedPipeline{}, fmt.Errorf("pipeline %q must declare at least one artifact lane", pipelineID)
|
|
}
|
|
|
|
input := resolveBinding(profile.Input, "")
|
|
if input.Module == "" {
|
|
return ResolvedPipeline{}, fmt.Errorf("pipeline %q input module must not be empty", pipelineID)
|
|
}
|
|
inputModuleSpec, err := inputSpec(catalog, input.Module)
|
|
if err != nil {
|
|
return ResolvedPipeline{}, moduleLookupError(pipelineID, "", StageInput, input.Module, err)
|
|
}
|
|
|
|
capabilities := newCapabilitySet()
|
|
if missing, ok := capabilities.missing(inputModuleSpec.Requires); ok {
|
|
return ResolvedPipeline{}, capabilityError(pipelineID, "", StageInput, input.Module, missing)
|
|
}
|
|
capabilities.add(inputModuleSpec.Provides...)
|
|
|
|
chunk := resolveBinding(profile.Chunk, DefaultChunkModule)
|
|
chunkSpec, err := chunkerSpec(catalog, chunk.Module)
|
|
if err != nil {
|
|
return ResolvedPipeline{}, moduleLookupError(pipelineID, "", StageChunk, chunk.Module, err)
|
|
}
|
|
if missing, ok := capabilities.missing(chunkSpec.Requires); ok {
|
|
return ResolvedPipeline{}, capabilityError(pipelineID, "", StageChunk, chunk.Module, missing)
|
|
}
|
|
capabilities.add(chunkSpec.Provides...)
|
|
|
|
lanesByID, selectedLaneIDs, err := selectedArtifactLanes(pipelineID, profile.Artifacts, options)
|
|
if err != nil {
|
|
return ResolvedPipeline{}, err
|
|
}
|
|
if len(selectedLaneIDs) == 0 {
|
|
return ResolvedPipeline{}, fmt.Errorf("pipeline %q must select at least one artifact lane", pipelineID)
|
|
}
|
|
|
|
resolved := ResolvedPipeline{
|
|
ID: pipelineID,
|
|
Input: input,
|
|
Chunk: chunk,
|
|
Output: resolveBinding(profile.Output, DefaultOutputModule),
|
|
}
|
|
outputCapabilities := capabilities.clone()
|
|
|
|
for _, laneID := range selectedLaneIDs {
|
|
laneProfile := lanesByID[laneID]
|
|
lane, laneCapabilities, err := resolveArtifactLane(pipelineID, laneID, laneProfile, capabilities, catalog)
|
|
if err != nil {
|
|
return ResolvedPipeline{}, err
|
|
}
|
|
resolved.ArtifactLanes = append(resolved.ArtifactLanes, lane)
|
|
outputCapabilities.addSet(laneCapabilities)
|
|
}
|
|
|
|
outputSpec, err := outputSpec(catalog, resolved.Output.Module)
|
|
if err != nil {
|
|
return ResolvedPipeline{}, moduleLookupError(pipelineID, "", StageOutput, resolved.Output.Module, err)
|
|
}
|
|
if missing, ok := outputCapabilities.missing(outputSpec.Requires); ok {
|
|
return ResolvedPipeline{}, capabilityError(pipelineID, "", StageOutput, resolved.Output.Module, missing)
|
|
}
|
|
|
|
digest, err := resolvedPipelineDigest(resolved)
|
|
if err != nil {
|
|
return ResolvedPipeline{}, fmt.Errorf("pipeline %q digest: %w", pipelineID, err)
|
|
}
|
|
resolved.Digest = digest
|
|
return resolved, nil
|
|
}
|
|
|
|
func resolveArtifactLane(pipelineID, laneID string, profile ArtifactLaneProfile, inherited capabilitySet, catalog ModuleCatalog) (ResolvedArtifactLane, capabilitySet, error) {
|
|
lane := ResolvedArtifactLane{
|
|
ID: laneID,
|
|
Extract: resolveBinding(profile.Extract, ""),
|
|
Merge: resolveBinding(profile.Merge, DefaultMergeModule),
|
|
Normalize: resolveBinding(profile.Normalize, DefaultNormalizeModule),
|
|
Validators: resolveBindings(profile.Validators, ""),
|
|
}
|
|
if lane.Extract.Module == "" {
|
|
return ResolvedArtifactLane{}, nil, fmt.Errorf("pipeline %q lane %q extract module must not be empty", pipelineID, laneID)
|
|
}
|
|
|
|
capabilities := inherited.clone()
|
|
|
|
extractSpec, err := extractorSpec(catalog, lane.Extract.Module)
|
|
if err != nil {
|
|
return ResolvedArtifactLane{}, nil, moduleLookupError(pipelineID, laneID, StageExtract, lane.Extract.Module, err)
|
|
}
|
|
if missing, ok := capabilities.missing(extractSpec.Requires); ok {
|
|
return ResolvedArtifactLane{}, nil, capabilityError(pipelineID, laneID, StageExtract, lane.Extract.Module, missing)
|
|
}
|
|
capabilities.add(extractSpec.Provides...)
|
|
|
|
mergeSpec, err := mergerSpec(catalog, lane.Merge.Module)
|
|
if err != nil {
|
|
return ResolvedArtifactLane{}, nil, moduleLookupError(pipelineID, laneID, StageMerge, lane.Merge.Module, err)
|
|
}
|
|
if missing, ok := capabilities.missing(mergeSpec.Requires); ok {
|
|
return ResolvedArtifactLane{}, nil, capabilityError(pipelineID, laneID, StageMerge, lane.Merge.Module, missing)
|
|
}
|
|
capabilities.add(mergeSpec.Provides...)
|
|
|
|
normalizeSpec, err := normalizerSpec(catalog, lane.Normalize.Module)
|
|
if err != nil {
|
|
return ResolvedArtifactLane{}, nil, moduleLookupError(pipelineID, laneID, StageNormalize, lane.Normalize.Module, err)
|
|
}
|
|
if missing, ok := capabilities.missing(normalizeSpec.Requires); ok {
|
|
return ResolvedArtifactLane{}, nil, capabilityError(pipelineID, laneID, StageNormalize, lane.Normalize.Module, missing)
|
|
}
|
|
capabilities.add(normalizeSpec.Provides...)
|
|
|
|
for _, validator := range lane.Validators {
|
|
validatorSpec, err := validatorSpec(catalog, validator.Module)
|
|
if err != nil {
|
|
return ResolvedArtifactLane{}, nil, moduleLookupError(pipelineID, laneID, StageValidate, validator.Module, err)
|
|
}
|
|
if missing, ok := capabilities.missing(validatorSpec.Requires); ok {
|
|
return ResolvedArtifactLane{}, nil, capabilityError(pipelineID, laneID, StageValidate, validator.Module, missing)
|
|
}
|
|
capabilities.add(validatorSpec.Provides...)
|
|
}
|
|
|
|
return lane, capabilities, nil
|
|
}
|
|
|
|
func resolveBinding(binding ModuleBinding, defaultModule string) ModuleBinding {
|
|
module := strings.TrimSpace(binding.Module)
|
|
if module == "" {
|
|
module = defaultModule
|
|
}
|
|
llmProfile := strings.TrimSpace(binding.LLMProfile)
|
|
if llmProfile == "" {
|
|
llmProfile = DefaultLLMProfile
|
|
}
|
|
return ModuleBinding{
|
|
Module: module,
|
|
LLMProfile: llmProfile,
|
|
Options: cloneOptions(binding.Options),
|
|
}
|
|
}
|
|
|
|
func resolveBindings(bindings []ModuleBinding, defaultModule string) []ModuleBinding {
|
|
if len(bindings) == 0 {
|
|
return nil
|
|
}
|
|
|
|
resolved := make([]ModuleBinding, 0, len(bindings))
|
|
for _, binding := range bindings {
|
|
resolvedBinding := resolveBinding(binding, defaultModule)
|
|
resolved = append(resolved, resolvedBinding)
|
|
}
|
|
return resolved
|
|
}
|
|
|
|
func cloneOptions(options map[string]any) map[string]any {
|
|
if len(options) == 0 {
|
|
return nil
|
|
}
|
|
|
|
copied := make(map[string]any, len(options))
|
|
for key, value := range options {
|
|
copied[key] = value
|
|
}
|
|
return copied
|
|
}
|
|
|
|
func selectedArtifactLanes(pipelineID string, artifacts map[string]ArtifactLaneProfile, options ResolveOptions) (map[string]ArtifactLaneProfile, []string, error) {
|
|
lanesByID := make(map[string]ArtifactLaneProfile, len(artifacts))
|
|
for rawLaneID, lane := range artifacts {
|
|
laneID := strings.TrimSpace(rawLaneID)
|
|
if laneID == "" {
|
|
return nil, nil, fmt.Errorf("pipeline %q artifact lane id must not be empty", pipelineID)
|
|
}
|
|
if _, ok := lanesByID[laneID]; ok {
|
|
return nil, nil, fmt.Errorf("pipeline %q artifact lane %q is duplicated after trimming", pipelineID, laneID)
|
|
}
|
|
lanesByID[laneID] = lane
|
|
}
|
|
|
|
if len(options.Only) == 0 {
|
|
keys := make([]string, 0, len(lanesByID))
|
|
for laneID := range lanesByID {
|
|
keys = append(keys, laneID)
|
|
}
|
|
sort.Strings(keys)
|
|
return lanesByID, keys, nil
|
|
}
|
|
|
|
selected := make(map[string]struct{}, len(options.Only))
|
|
for _, rawLaneID := range options.Only {
|
|
laneID := strings.TrimSpace(rawLaneID)
|
|
if laneID == "" {
|
|
return nil, nil, fmt.Errorf("pipeline %q selected artifact lane id must not be empty", pipelineID)
|
|
}
|
|
if _, ok := lanesByID[laneID]; !ok {
|
|
return nil, nil, fmt.Errorf("pipeline %q selected artifact lane %q is not declared", pipelineID, laneID)
|
|
}
|
|
selected[laneID] = struct{}{}
|
|
}
|
|
|
|
keys := make([]string, 0, len(selected))
|
|
for laneID := range selected {
|
|
keys = append(keys, laneID)
|
|
}
|
|
sort.Strings(keys)
|
|
return lanesByID, keys, nil
|
|
}
|
|
|
|
func resolvedPipelineDigest(resolved ResolvedPipeline) (string, error) {
|
|
withoutDigest := struct {
|
|
ID string
|
|
Input ModuleBinding
|
|
Chunk ModuleBinding
|
|
ArtifactLanes []ResolvedArtifactLane
|
|
Output ModuleBinding
|
|
}{
|
|
ID: resolved.ID,
|
|
Input: resolved.Input,
|
|
Chunk: resolved.Chunk,
|
|
ArtifactLanes: resolved.ArtifactLanes,
|
|
Output: resolved.Output,
|
|
}
|
|
encoded, err := json.Marshal(withoutDigest)
|
|
if err != nil {
|
|
return "", err
|
|
}
|
|
|
|
sum := sha256.Sum256(encoded)
|
|
return "sha256:" + hex.EncodeToString(sum[:]), nil
|
|
}
|
|
|
|
func inputSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Inputs, key)
|
|
}
|
|
|
|
func chunkerSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Chunkers, key)
|
|
}
|
|
|
|
func extractorSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Extractors, key)
|
|
}
|
|
|
|
func mergerSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Mergers, key)
|
|
}
|
|
|
|
func normalizerSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Normalizers, key)
|
|
}
|
|
|
|
func validatorSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Validators, key)
|
|
}
|
|
|
|
func outputSpec(catalog ModuleCatalog, key string) (ModuleSpec, error) {
|
|
return registrySpec(catalog.Outputs, key)
|
|
}
|
|
|
|
type specRegistry interface {
|
|
Spec(key string) (ModuleSpec, bool)
|
|
}
|
|
|
|
func registrySpec(registry specRegistry, key string) (ModuleSpec, error) {
|
|
if registry == nil {
|
|
return ModuleSpec{}, fmt.Errorf("module %q is not registered", key)
|
|
}
|
|
spec, ok := registry.Spec(key)
|
|
if !ok {
|
|
return ModuleSpec{}, fmt.Errorf("module %q is not registered", key)
|
|
}
|
|
return spec, nil
|
|
}
|
|
|
|
func moduleLookupError(pipelineID, laneID string, stage ModuleStage, module string, err error) error {
|
|
if laneID != "" {
|
|
return fmt.Errorf("pipeline %q lane %q %s module %q: %w", pipelineID, laneID, stage, module, err)
|
|
}
|
|
return fmt.Errorf("pipeline %q %s module %q: %w", pipelineID, stage, module, err)
|
|
}
|
|
|
|
func capabilityError(pipelineID, laneID string, stage ModuleStage, module, capability string) error {
|
|
if laneID != "" {
|
|
return fmt.Errorf("pipeline %q lane %q %s module %q requires missing capability %q", pipelineID, laneID, stage, module, capability)
|
|
}
|
|
return fmt.Errorf("pipeline %q %s module %q requires missing capability %q", pipelineID, stage, module, capability)
|
|
}
|
|
|
|
type capabilitySet map[string]struct{}
|
|
|
|
func newCapabilitySet() capabilitySet {
|
|
return make(capabilitySet)
|
|
}
|
|
|
|
func (set capabilitySet) clone() capabilitySet {
|
|
copied := make(capabilitySet, len(set))
|
|
for capability := range set {
|
|
copied[capability] = struct{}{}
|
|
}
|
|
return copied
|
|
}
|
|
|
|
func (set capabilitySet) add(values ...string) {
|
|
for _, value := range values {
|
|
set[value] = struct{}{}
|
|
}
|
|
}
|
|
|
|
func (set capabilitySet) addSet(other capabilitySet) {
|
|
for value := range other {
|
|
set[value] = struct{}{}
|
|
}
|
|
}
|
|
|
|
func (set capabilitySet) missing(required []string) (string, bool) {
|
|
for _, capability := range required {
|
|
if _, ok := set[capability]; !ok {
|
|
return capability, true
|
|
}
|
|
}
|
|
return "", false
|
|
}
|