87 lines
3.8 KiB
Go
87 lines
3.8 KiB
Go
package pipeline
|
|
|
|
import (
|
|
"fmt"
|
|
"sort"
|
|
"strings"
|
|
)
|
|
|
|
type checkpointFingerprintComponent struct {
|
|
scope string
|
|
module any
|
|
}
|
|
|
|
// CheckpointFingerprints returns a defensive copy of the stable semantic
|
|
// identities contributed by the pipeline's prepared modules and validators.
|
|
func (p *PreparedPipeline) CheckpointFingerprints() []CheckpointFingerprint {
|
|
if p == nil {
|
|
return nil
|
|
}
|
|
return append([]CheckpointFingerprint(nil), p.checkpointFingerprints...)
|
|
}
|
|
|
|
func collectPreparedCheckpointFingerprints(prepared *PreparedPipeline) ([]CheckpointFingerprint, error) {
|
|
components := []checkpointFingerprintComponent{
|
|
{scope: "input:" + prepared.resolved.Input.Module, module: prepared.input},
|
|
{scope: "chunk:" + prepared.resolved.Chunk.Module, module: prepared.chunker},
|
|
}
|
|
components = appendValidatorFingerprintComponents(components, "chunk:"+prepared.resolved.Chunk.Module, prepared.chunkValidators)
|
|
for _, lane := range prepared.lanes {
|
|
laneScope := func(stage ModuleStage, moduleKey string) string {
|
|
return string(stage) + ":" + lane.resolved.ID + ":" + moduleKey
|
|
}
|
|
components = append(components, checkpointFingerprintComponent{scope: laneScope(StageExtract, lane.resolved.Extract.Module), module: lane.typed.extractor})
|
|
components = appendValidatorFingerprintComponents(components, laneScope(StageExtract, lane.resolved.Extract.Module), lane.extractValidators)
|
|
components = append(components, checkpointFingerprintComponent{scope: laneScope(StageMerge, lane.resolved.Merge.Module), module: lane.typed.merger})
|
|
components = appendValidatorFingerprintComponents(components, laneScope(StageMerge, lane.resolved.Merge.Module), lane.mergeValidators)
|
|
components = append(components, checkpointFingerprintComponent{scope: laneScope(StageNormalize, lane.resolved.Normalize.Module), module: lane.typed.normalizer})
|
|
components = appendValidatorFingerprintComponents(components, laneScope(StageNormalize, lane.resolved.Normalize.Module), lane.normalizeValidators)
|
|
}
|
|
components = append(components, checkpointFingerprintComponent{scope: "output:" + prepared.resolved.Output.Module, module: prepared.output})
|
|
|
|
var values []CheckpointFingerprint
|
|
seen := make(map[string]struct{})
|
|
for _, item := range components {
|
|
provider, ok := item.module.(CheckpointFingerprintProvider)
|
|
if !ok {
|
|
continue
|
|
}
|
|
for _, fingerprint := range provider.CheckpointFingerprints() {
|
|
name := strings.TrimSpace(fingerprint.Name)
|
|
value := strings.TrimSpace(fingerprint.Value)
|
|
if name == "" || value == "" {
|
|
return nil, fmt.Errorf("pipeline %q checkpoint fingerprint for %s must have a non-empty name and value", prepared.resolved.ID, item.scope)
|
|
}
|
|
scopedName := item.scope + ":" + name
|
|
if _, exists := seen[scopedName]; exists {
|
|
return nil, fmt.Errorf("pipeline %q checkpoint fingerprint %q is duplicated", prepared.resolved.ID, scopedName)
|
|
}
|
|
seen[scopedName] = struct{}{}
|
|
values = append(values, CheckpointFingerprint{Name: scopedName, Value: value})
|
|
}
|
|
}
|
|
sort.Slice(values, func(i, j int) bool { return values[i].Name < values[j].Name })
|
|
return values, nil
|
|
}
|
|
|
|
func appendValidatorFingerprintComponents(components []checkpointFingerprintComponent, ownerScope string, chain preparedValidatorChain) []checkpointFingerprintComponent {
|
|
for index, validator := range chain.validators {
|
|
scope := fmt.Sprintf("%s:validator:%d:%s", ownerScope, index, validator.resolved.Binding.Module)
|
|
components = append(components, checkpointFingerprintComponent{scope: scope, module: preparedValidatorImplementation(validator)})
|
|
}
|
|
return components
|
|
}
|
|
|
|
func preparedValidatorImplementation(validator preparedValidator) any {
|
|
switch validator.resolved.Target {
|
|
case ValidatorTargetTyped:
|
|
return validator.typed
|
|
case ValidatorTargetChunk:
|
|
return validator.chunk
|
|
case ValidatorTargetSerialized:
|
|
return validator.serialized
|
|
default:
|
|
return nil
|
|
}
|
|
}
|