Make command line references pipeline scoped

This commit is contained in:
2026-08-29 12:34:31 +00:00
parent f208dbe954
commit deebc89255
9 changed files with 824 additions and 238 deletions

View File

@@ -4,6 +4,7 @@ import (
"bytes"
"os"
"path/filepath"
"slices"
"strings"
"testing"
@@ -14,20 +15,16 @@ import (
func TestReferenceSelectorsParseAndApplyAllDocumentedForms(t *testing.T) {
tests := []struct {
name string
selector string
only []string
wantStage pipeline.ModuleStage
wantLane string
wantSlot string
name string
selector string
want []string
}{
{name: "flat", selector: "alpha-slot", wantStage: pipeline.StageExtract, wantLane: "alpha", wantSlot: "alpha-slot"},
{name: "chunk", selector: "chunk.chunk-slot", wantStage: pipeline.StageChunk, wantSlot: "chunk-slot"},
{name: "merge", selector: "merge.alpha-merge", only: []string{"alpha"}, wantStage: pipeline.StageMerge, wantLane: "alpha", wantSlot: "alpha-merge"},
{name: "lane", selector: "alpha.alpha-slot", wantStage: pipeline.StageExtract, wantLane: "alpha", wantSlot: "alpha-slot"},
{name: "lane extract", selector: "alpha.extract.alpha-slot", wantStage: pipeline.StageExtract, wantLane: "alpha", wantSlot: "alpha-slot"},
{name: "lane merge", selector: "alpha.merge.alpha-merge", wantStage: pipeline.StageMerge, wantLane: "alpha", wantSlot: "alpha-merge"},
{name: "lane normalize", selector: "alpha.normalize.alpha-normalize", wantStage: pipeline.StageNormalize, wantLane: "alpha", wantSlot: "alpha-normalize"},
{name: "pipeline", selector: "shared", want: []string{"alpha.extract.shared", "alpha.merge.shared", "alpha.normalize.shared", "beta.extract.shared", "beta.merge.shared", "beta.normalize.shared"}},
{name: "chunk", selector: "chunk.chunk-slot", want: []string{"chunk.chunk-slot"}},
{name: "lane", selector: "alpha.shared", want: []string{"alpha.extract.shared", "alpha.merge.shared", "alpha.normalize.shared"}},
{name: "lane extract", selector: "alpha.extract.alpha-slot", want: []string{"alpha.extract.alpha-slot"}},
{name: "lane merge", selector: "alpha.merge.alpha-merge", want: []string{"alpha.merge.alpha-merge"}},
{name: "lane normalize", selector: "alpha.normalize.alpha-normalize", want: []string{"alpha.normalize.alpha-normalize"}},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
@@ -37,70 +34,135 @@ func TestReferenceSelectorsParseAndApplyAllDocumentedForms(t *testing.T) {
if err != nil {
t.Fatal(err)
}
overrides, _, err := resolveCLIReferenceRequests(cfg, "demo", tt.only, catalog, []cliReferenceRequest{{Selector: selector, Source: "reference.txt"}}, nil)
overrides, _, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog, []cliReferenceRequest{{Selector: selector, Source: "reference.txt"}}, nil)
if err != nil {
t.Fatalf("resolve selector: %v", err)
}
if len(overrides) != 1 {
t.Fatalf("overrides = %#v, want one binding", overrides)
if got := referenceContractBindingLabels(overrides); !slices.Equal(got, tt.want) {
t.Fatalf("binding targets = %#v, want %#v", got, tt.want)
}
got := overrides[0]
if got.Stage != tt.wantStage || got.LaneID != tt.wantLane || got.SlotName != tt.wantSlot || got.BindingSource != contracts.ReferenceBindingSourceCLI {
t.Fatalf("binding = %#v, want %s/%s/%s from CLI", got, tt.wantStage, tt.wantLane, tt.wantSlot)
}
})
}
}
func TestReferenceSelectorsRejectAmbiguityWithSpecificSuggestions(t *testing.T) {
cfg := referenceContractConfig()
catalog := referenceContractCatalog(t, true, true)
for _, tt := range []struct {
name string
selector string
want []string
}{
{name: "flat shared slot", selector: "shared", want: []string{"alpha.extract.shared", "beta.extract.shared"}},
{name: "lane shared slot", selector: "alpha.shared", want: []string{"alpha.extract.shared", "alpha.merge.shared", "alpha.normalize.shared"}},
{name: "all mergers", selector: "merge.shared", want: []string{"alpha.merge.shared", "beta.merge.shared"}},
} {
t.Run(tt.name, func(t *testing.T) {
selector, err := parseReferenceSelector(tt.selector, "--reference")
if err != nil {
t.Fatal(err)
}
_, _, err = resolveCLIReferenceRequests(cfg, "demo", nil, catalog, []cliReferenceRequest{{Selector: selector, Source: "reference.txt"}}, nil)
if err == nil {
t.Fatal("resolve selector succeeded, want ambiguity error")
}
for _, fragment := range tt.want {
if !strings.Contains(err.Error(), fragment) {
t.Fatalf("error = %q, want suggestion %q", err, fragment)
for _, binding := range overrides {
if binding.Source != "reference.txt" || binding.BindingSource != contracts.ReferenceBindingSourceCLI {
t.Fatalf("binding = %#v, want CLI source", binding)
}
}
})
}
}
func TestReferenceSelectorsRespectSelectedLanesBeforeMaterialization(t *testing.T) {
func TestReferenceSelectorSpecificityAndFinalOccurrenceChooseConcreteBindings(t *testing.T) {
cfg := referenceContractConfig()
catalog := referenceContractCatalog(t, true, true)
requests := []cliReferenceRequest{
{Selector: mustParseReferenceSelector(t, "shared", "--reference"), Source: "pipeline-first.txt"},
{Selector: mustParseReferenceSelector(t, "shared", "--reference"), Source: "pipeline-final.txt"},
{Selector: mustParseReferenceSelector(t, "alpha.shared", "--reference"), Source: "lane.txt"},
{Selector: mustParseReferenceSelector(t, "alpha.extract.shared", "--reference"), Source: "binding.txt"},
}
overrides, unbinds, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog, requests, nil)
if err != nil {
t.Fatal(err)
}
if len(unbinds) != 0 {
t.Fatalf("unbinds = %#v, want none", unbinds)
}
want := map[string]string{
"alpha.extract.shared": "binding.txt",
"alpha.merge.shared": "lane.txt",
"alpha.normalize.shared": "lane.txt",
"beta.extract.shared": "pipeline-final.txt",
"beta.merge.shared": "pipeline-final.txt",
"beta.normalize.shared": "pipeline-final.txt",
}
for _, binding := range overrides {
label := referenceContractBindingLabel(binding)
if binding.Source != want[label] {
t.Fatalf("binding %s source = %q, want %q", label, binding.Source, want[label])
}
delete(want, label)
}
if len(want) != 0 {
t.Fatalf("missing bindings: %#v", want)
}
}
func TestCompleteDNDSharedCLIReferencesExpandAcrossCompatibleTargets(t *testing.T) {
cfg := loadMaintainedExample(t, repositoryPath("examples", "dnd-complete.config.yml"))
catalog := catalogFromRegistries(productionTestComponents(t).registries)
sources := map[string]string{
"party": "/references/party.txt",
"players": "/references/players.txt",
"glossary": "/references/glossary.txt",
"spell_catalog": "/references/spells.json",
}
requests := make([]cliReferenceRequest, 0, len(sources))
for _, slot := range []string{"party", "players", "glossary", "spell_catalog"} {
requests = append(requests, cliReferenceRequest{
Selector: mustParseReferenceSelector(t, slot, "--reference"),
Source: sources[slot],
})
}
overrides, unbinds, err := resolveCLIReferenceRequests(cfg, "dnd-session", nil, catalog, requests, nil)
if err != nil {
t.Fatalf("expand complete D&D references: %v", err)
}
if len(unbinds) != 0 {
t.Fatalf("unbinds = %#v, want none", unbinds)
}
actual := make(map[string]pipeline.ReferenceBinding, len(overrides))
for _, binding := range overrides {
actual[referenceContractBindingLabel(binding)] = binding
}
targets, err := selectedReferenceTargets(cfg, "dnd-session", nil, catalog)
if err != nil {
t.Fatal(err)
}
matched := make(map[string]int, len(sources))
for _, target := range targets {
for slot, sourcePath := range sources {
if _, ok := target.slots[slot]; !ok {
continue
}
matched[slot]++
label := targetLabel(target) + "." + slot
binding, ok := actual[label]
if !ok || binding.Source != sourcePath || binding.BindingSource != contracts.ReferenceBindingSourceCLI {
t.Fatalf("binding %q = %#v, want CLI source %q", label, binding, sourcePath)
}
}
}
for slot := range sources {
if matched[slot] < 2 {
t.Fatalf("reference %q matched %d target(s), want a shared D&D reference", slot, matched[slot])
}
}
if _, err := cfg.Resolve(config.ResolveInput{PipelineID: "dnd-session", Catalog: catalog, ReferenceOverrides: overrides}); err != nil {
t.Fatalf("resolve complete D&D CLI references: %v", err)
}
}
func TestReferenceSelectorsRejectInvalidOrUnselectedScopesBeforeMaterialization(t *testing.T) {
cfg := referenceContractConfig()
catalog := referenceContractCatalog(t, true, true)
for _, tt := range []struct {
name string
selector string
only []string
want string
}{
{name: "unselected lane", selector: "beta.extract.beta-slot", want: `reference lane "beta" is not selected`},
{name: "pipeline slot", selector: "missing", want: `reference slot "missing" is not declared by any selected target`},
{name: "lane slot", selector: "alpha.missing", want: `reference slot "missing" is not declared by selected lane "alpha"`},
{name: "binding slot", selector: "alpha.extract.missing", want: `reference slot "missing" is not declared`},
{name: "former merge shorthand", selector: "merge.shared", want: `reference lane "merge" is not selected`},
{name: "unselected lane", selector: "beta.extract.beta-slot", only: []string{"alpha"}, want: `reference lane "beta" is not selected`},
{name: "unknown lane", selector: "missing.extract.beta-slot", want: `reference lane "missing" is not selected`},
} {
t.Run(tt.name, func(t *testing.T) {
selector, err := parseReferenceSelector(tt.selector, "--reference")
if err != nil {
t.Fatal(err)
}
_, _, err = resolveCLIReferenceRequests(cfg, "demo", []string{"alpha"}, catalog, []cliReferenceRequest{{Selector: selector, Source: filepath.Join(t.TempDir(), "missing.txt")}}, nil)
selector := mustParseReferenceSelector(t, tt.selector, "--reference")
_, _, err := resolveCLIReferenceRequests(cfg, "demo", tt.only, catalog, []cliReferenceRequest{{Selector: selector, Source: filepath.Join(t.TempDir(), "missing.txt")}}, nil)
if err == nil || !strings.Contains(err.Error(), tt.want) || strings.Contains(err.Error(), "missing.txt") {
t.Fatalf("error = %v, want selection failure before file access", err)
t.Fatalf("error = %v, want selection failure containing %q before file access", err, tt.want)
}
})
}
@@ -131,40 +193,50 @@ func TestReferenceSyntaxErrorsReturnTwo(t *testing.T) {
}
}
func TestReferenceOverridesUseFinalExactTargetBinding(t *testing.T) {
func TestReferenceBindAndUnbindSpecificity(t *testing.T) {
cfg := referenceContractConfig()
catalog := referenceContractCatalog(t, true, true)
alphaShared, err := parseReferenceSelector("alpha.extract.shared", "--reference")
if err != nil {
t.Fatal(err)
}
betaShared, err := parseReferenceSelector("beta.extract.shared", "--reference")
if err != nil {
t.Fatal(err)
}
overrides, unbinds, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog, []cliReferenceRequest{
{Selector: alphaShared, Source: "alpha-first.txt"},
{Selector: alphaShared, Source: "alpha-final.txt"},
{Selector: betaShared, Source: "beta-only.txt"},
}, nil)
if err != nil {
t.Fatal(err)
}
if len(unbinds) != 0 {
t.Fatalf("unbinds = %#v, want none", unbinds)
}
effective, err := cfg.Resolve(config.ResolveInput{PipelineID: "demo", Catalog: catalog, ReferenceOverrides: overrides})
if err != nil {
t.Fatalf("resolve pipeline: %v", err)
}
alpha := referenceContractLane(t, effective.ResolvedPipeline, "alpha")
beta := referenceContractLane(t, effective.ResolvedPipeline, "beta")
if source := referenceContractBindingSource(alpha.ExtractReferences.Bindings, "shared"); source != "alpha-final.txt" {
t.Fatalf("alpha shared source = %q, want final exact-target override", source)
}
if source := referenceContractBindingSource(beta.ExtractReferences.Bindings, "shared"); source != "beta-only.txt" {
t.Fatalf("beta shared source = %q, want target-specific override", source)
}
t.Run("specific unbind carves out broad binding", func(t *testing.T) {
overrides, unbinds, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog,
[]cliReferenceRequest{{Selector: mustParseReferenceSelector(t, "shared", "--reference"), Source: "shared.txt"}},
[]cliReferenceUnbindRequest{{Selector: mustParseReferenceSelector(t, "alpha.extract.shared", "--without-reference")}},
)
if err != nil {
t.Fatal(err)
}
if got := referenceContractBindingLabels(overrides); slices.Contains(got, "alpha.extract.shared") || len(got) != 5 {
t.Fatalf("overrides = %#v, want all shared targets except alpha extract", got)
}
if got := referenceContractUnbindLabels(unbinds); !slices.Equal(got, []string{"alpha.extract.shared"}) {
t.Fatalf("unbinds = %#v, want alpha extract", got)
}
})
t.Run("specific binding restores broad unbind", func(t *testing.T) {
overrides, unbinds, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog,
[]cliReferenceRequest{{Selector: mustParseReferenceSelector(t, "alpha.extract.shared", "--reference"), Source: "alpha.txt"}},
[]cliReferenceUnbindRequest{{Selector: mustParseReferenceSelector(t, "shared", "--without-reference")}},
)
if err != nil {
t.Fatal(err)
}
if got := referenceContractBindingLabels(overrides); !slices.Equal(got, []string{"alpha.extract.shared"}) {
t.Fatalf("overrides = %#v, want alpha extract", got)
}
if got := referenceContractUnbindLabels(unbinds); slices.Contains(got, "alpha.extract.shared") || len(got) != 5 {
t.Fatalf("unbinds = %#v, want all shared targets except alpha extract", got)
}
})
t.Run("same specificity conflicts", func(t *testing.T) {
_, _, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog,
[]cliReferenceRequest{{Selector: mustParseReferenceSelector(t, "alpha.shared", "--reference"), Source: "alpha.txt"}},
[]cliReferenceUnbindRequest{{Selector: mustParseReferenceSelector(t, "alpha.shared", "--without-reference")}},
)
if err == nil || !strings.Contains(err.Error(), "same specificity") {
t.Fatalf("error = %v, want same-specificity conflict", err)
}
})
}
func TestReferenceUnbindsRemoveOptionalAndProtectRequiredSlots(t *testing.T) {
@@ -257,6 +329,52 @@ func TestReferenceMaterializationSeparatesCLIAndConfigPathOrigins(t *testing.T)
}
}
func TestPipelineScopedCLIReferenceProtectsGeneratedHandoff(t *testing.T) {
cfg := referenceContractConfig()
profile := cfg.Pipelines["demo"]
alpha := profile.Artifacts["alpha"]
beta := profile.Artifacts["beta"]
alpha.Extract.References["shared"] = pipeline.GeneratedReference("produce", "beta")
profile.Artifacts = nil
profile.Steps = []pipeline.PipelineStepProfile{
{ID: "produce", Artifacts: map[string]pipeline.ArtifactLaneProfile{"beta": beta}},
{ID: "consume", Artifacts: map[string]pipeline.ArtifactLaneProfile{"alpha": alpha}},
}
cfg.Pipelines["demo"] = profile
catalog := referenceContractCatalog(t, true, true)
t.Run("binding conflicts before file access", func(t *testing.T) {
overrides, _, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog, []cliReferenceRequest{{
Selector: mustParseReferenceSelector(t, "shared", "--reference"),
Source: filepath.Join(t.TempDir(), "never-read.json"),
}}, nil)
if err != nil {
t.Fatalf("expand CLI reference: %v", err)
}
_, err = cfg.Resolve(config.ResolveInput{PipelineID: "demo", Catalog: catalog, ReferenceOverrides: overrides})
if err == nil || !strings.Contains(err.Error(), "conflicting generated and external bindings") || strings.Contains(err.Error(), "never-read.json") {
t.Fatalf("resolve error = %v, want generated/external conflict before file access", err)
}
})
t.Run("unbind leaves generated source intact", func(t *testing.T) {
_, unbinds, err := resolveCLIReferenceRequests(cfg, "demo", nil, catalog, nil, []cliReferenceUnbindRequest{{
Selector: mustParseReferenceSelector(t, "shared", "--without-reference"),
}})
if err != nil {
t.Fatalf("expand CLI unbind: %v", err)
}
effective, err := cfg.Resolve(config.ResolveInput{PipelineID: "demo", Catalog: catalog, ReferenceUnbinds: unbinds})
if err != nil {
t.Fatalf("resolve generated reference with CLI unbind: %v", err)
}
binding := referenceContractFindBinding(referenceContractLane(t, effective.ResolvedPipeline, "alpha").ExtractReferences.Bindings, "shared")
if binding == nil || binding.Artifact == nil || binding.Artifact.Step != "produce" || binding.Artifact.Lane != "beta" {
t.Fatalf("generated binding = %#v, want preserved produce/beta handoff", binding)
}
})
}
func TestReferenceTargetLookupUsesArtifactVariantsAndReportsMissingContext(t *testing.T) {
cfg := referenceContractConfig()
full := referenceContractCatalog(t, true, true)
@@ -368,7 +486,7 @@ func referenceContractCatalog(t *testing.T, includeBetaMerger, includeBetaNormal
register(registries.Chunkers.RegisterWithSpec(pipeline.ModuleSpec{Key: "reference/chunk", Stage: pipeline.StageChunk, ExecutionClass: contracts.ExecutionClassDeterministic, Requires: []string{"source"}, Provides: []string{"chunks"}, ReferenceSlots: []contracts.ReferenceSlot{{Name: "chunk-slot"}, {Name: "required-chunk", Required: true}}}, func() (contracts.Chunker, error) { return stateTestChunker{}, nil }))
register(pipeline.RegisterArtifactCodec(registries.ArtifactCodecs, referenceContractCodecA{}))
register(pipeline.RegisterArtifactCodec(registries.ArtifactCodecs, referenceContractCodecB{}))
register(pipeline.RegisterExtractor(registries.Extractors, pipeline.ModuleSpec{Key: "reference/extract-alpha", Stage: pipeline.StageExtract, ExecutionClass: contracts.ExecutionClassDeterministic, Requires: []string{"chunks"}, Provides: []string{"artifact"}, ArtifactKind: referenceContractKindAlpha, ReferenceSlots: []contracts.ReferenceSlot{{Name: "shared"}, {Name: "alpha-slot"}, {Name: "required-extract", Required: true}}}, func() (contracts.Extractor[stateTestArtifact], error) { return stateTestExtractor{}, nil }))
register(pipeline.RegisterExtractor(registries.Extractors, pipeline.ModuleSpec{Key: "reference/extract-alpha", Stage: pipeline.StageExtract, ExecutionClass: contracts.ExecutionClassDeterministic, Requires: []string{"chunks"}, Provides: []string{"artifact"}, ArtifactKind: referenceContractKindAlpha, ReferenceSlots: []contracts.ReferenceSlot{{Name: "shared", AcceptedArtifactKinds: []contracts.ArtifactKind{referenceContractKindBeta}, AcceptedMediaTypes: []string{"application/json"}}, {Name: "alpha-slot"}, {Name: "required-extract", Required: true}}}, func() (contracts.Extractor[stateTestArtifact], error) { return stateTestExtractor{}, nil }))
register(pipeline.RegisterExtractor(registries.Extractors, pipeline.ModuleSpec{Key: "reference/extract-beta", Stage: pipeline.StageExtract, ExecutionClass: contracts.ExecutionClassDeterministic, Requires: []string{"chunks"}, Provides: []string{"artifact"}, ArtifactKind: referenceContractKindBeta, ReferenceSlots: []contracts.ReferenceSlot{{Name: "shared"}, {Name: "beta-slot"}, {Name: "required-extract", Required: true}}}, func() (contracts.Extractor[stateTestArtifact], error) { return stateTestExtractor{}, nil }))
register(pipeline.RegisterMerger(registries.Mergers, pipeline.ModuleSpec{Key: "reference/shared-merge", Stage: pipeline.StageMerge, ExecutionClass: contracts.ExecutionClassDeterministic, Requires: []string{"artifact"}, Provides: []string{"merged"}, ArtifactKind: referenceContractKindAlpha, ReferenceSlots: []contracts.ReferenceSlot{{Name: "shared"}, {Name: "alpha-merge"}, {Name: "required-merge", Required: true}}}, func() (contracts.Merger[stateTestArtifact], error) { return stateTestMerger{}, nil }))
if includeBetaMerger {
@@ -444,6 +562,42 @@ func referenceContractBindingSource(bindings []pipeline.ReferenceBinding, slot s
return ""
}
func mustParseReferenceSelector(t *testing.T, value, flagName string) cliReferenceSelector {
t.Helper()
selector, err := parseReferenceSelector(value, flagName)
if err != nil {
t.Fatal(err)
}
return selector
}
func referenceContractBindingLabels(bindings []pipeline.ReferenceBinding) []string {
labels := make([]string, 0, len(bindings))
for _, binding := range bindings {
labels = append(labels, referenceContractBindingLabel(binding))
}
return labels
}
func referenceContractBindingLabel(binding pipeline.ReferenceBinding) string {
if binding.Stage == pipeline.StageChunk {
return "chunk." + binding.SlotName
}
return binding.LaneID + "." + string(binding.Stage) + "." + binding.SlotName
}
func referenceContractUnbindLabels(unbinds []pipeline.ReferenceUnbind) []string {
labels := make([]string, 0, len(unbinds))
for _, unbind := range unbinds {
labels = append(labels, referenceContractBindingLabel(pipeline.ReferenceBinding{
Stage: unbind.Stage,
LaneID: unbind.LaneID,
SlotName: unbind.SlotName,
}))
}
return labels
}
func referenceContractFindBinding(bindings []pipeline.ReferenceBinding, slot string) *pipeline.ReferenceBinding {
for i := range bindings {
if bindings[i].SlotName == slot {

View File

@@ -177,7 +177,7 @@ func runPipelineCommand(args []string, stdout, stderr io.Writer, opts Options) i
fs.Var(&llmProfile, "llm-profile", "LLM profile override")
fs.Var(&reasoningEffort, "reasoning-effort", "reasoning effort override")
fs.Var(&chunkCache, "chunk_cache", "chunk plan cache mode: auto, bypass, or refresh")
fs.Var(&referenceFlags, "reference", "reference binding, as slot=path, chunk.slot=path, merge.slot=path, lane.slot=path, lane.extract.slot=path, lane.merge.slot=path, or lane.normalize.slot=path")
fs.Var(&referenceFlags, "reference", "reference binding, as slot=path, chunk.slot=path, lane.slot=path, lane.extract.slot=path, lane.merge.slot=path, or lane.normalize.slot=path")
fs.Var(&withoutReferenceFlags, "without-reference", "unbind a reference, using the same selector forms as --reference")
fs.Var(&recomputeStep, "recompute-step", "recompute one ordered pipeline step and dependent lanes")
if err := validateRunFlagValues(args); err != nil {
@@ -1287,11 +1287,21 @@ type cliReferenceUnbindRequest struct {
}
type cliReferenceSelector struct {
Scope cliReferenceSelectorScope
LaneID string
Stage pipeline.ModuleStage
SlotName string
}
type cliReferenceSelectorScope uint8
const (
cliReferenceScopePipeline cliReferenceSelectorScope = iota
cliReferenceScopeLane
cliReferenceScopeChunk
cliReferenceScopeBinding
)
func parseReferenceFlags(values []string) ([]cliReferenceRequest, error) {
if len(values) == 0 {
return nil, nil
@@ -1300,7 +1310,7 @@ func parseReferenceFlags(values []string) ([]cliReferenceRequest, error) {
for _, raw := range values {
name, source, ok := strings.Cut(raw, "=")
if !ok {
return nil, fmt.Errorf("--reference must use slot=path or lane.slot=path")
return nil, fmt.Errorf("--reference must use slot=path, lane.slot=path, or lane.stage.slot=path")
}
if strings.TrimSpace(source) == "" {
return nil, fmt.Errorf("--reference path must not be empty; use --without-reference to unbind")
@@ -1350,17 +1360,14 @@ func parseReferenceSelector(raw string, flagName string) (cliReferenceSelector,
}
switch len(parts) {
case 1:
return cliReferenceSelector{SlotName: strings.TrimSpace(parts[0])}, nil
return cliReferenceSelector{Scope: cliReferenceScopePipeline, SlotName: strings.TrimSpace(parts[0])}, nil
case 2:
first := strings.TrimSpace(parts[0])
slotName := strings.TrimSpace(parts[1])
if first == string(pipeline.StageChunk) {
return cliReferenceSelector{Stage: pipeline.StageChunk, SlotName: slotName}, nil
return cliReferenceSelector{Scope: cliReferenceScopeChunk, Stage: pipeline.StageChunk, SlotName: slotName}, nil
}
if first == string(pipeline.StageMerge) {
return cliReferenceSelector{Stage: pipeline.StageMerge, SlotName: slotName}, nil
}
return cliReferenceSelector{LaneID: first, SlotName: slotName}, nil
return cliReferenceSelector{Scope: cliReferenceScopeLane, LaneID: first, SlotName: slotName}, nil
case 3:
laneID := strings.TrimSpace(parts[0])
stage := pipeline.ModuleStage(strings.TrimSpace(parts[1]))
@@ -1368,9 +1375,9 @@ func parseReferenceSelector(raw string, flagName string) (cliReferenceSelector,
if stage != pipeline.StageExtract && stage != pipeline.StageMerge && stage != pipeline.StageNormalize {
return cliReferenceSelector{}, fmt.Errorf("%s lane-qualified selector must use lane.extract.slot, lane.merge.slot, or lane.normalize.slot", flagName)
}
return cliReferenceSelector{LaneID: laneID, Stage: stage, SlotName: slotName}, nil
return cliReferenceSelector{Scope: cliReferenceScopeBinding, LaneID: laneID, Stage: stage, SlotName: slotName}, nil
default:
return cliReferenceSelector{}, fmt.Errorf("%s must use slot, chunk.slot, merge.slot, lane.slot, lane.extract.slot, lane.merge.slot, or lane.normalize.slot", flagName)
return cliReferenceSelector{}, fmt.Errorf("%s must use slot, chunk.slot, lane.slot, lane.extract.slot, lane.merge.slot, or lane.normalize.slot", flagName)
}
}
@@ -1391,37 +1398,153 @@ func resolveCLIReferenceRequests(
return nil, nil, err
}
overrides := make([]pipeline.ReferenceBinding, 0, len(referenceRequests))
// Broad CLI selectors are only presentation syntax. Collapse them into one
// highest-specificity action per concrete framework target before pipeline
// resolution so the generic reference contract stays stage-and-lane exact.
actions := make(map[cliReferenceTargetKey]resolvedCLIReferenceAction)
for _, request := range referenceRequests {
target, err := resolveCLIReferenceTarget(targets, request.Selector)
matches, err := resolveCLIReferenceTargets(targets, request.Selector)
if err != nil {
return nil, nil, err
}
overrides = append(overrides, pipeline.ReferenceBinding{
Stage: target.stage,
LaneID: target.laneID,
SlotName: request.Selector.SlotName,
Source: request.Source,
BindingSource: contracts.ReferenceBindingSourceCLI,
})
for _, target := range matches {
candidate := resolvedCLIReferenceAction{
kind: cliReferenceActionBind,
selector: request.Selector,
target: target,
slotName: request.Selector.SlotName,
source: request.Source,
specificity: request.Selector.specificity(),
}
if err := mergeCLIReferenceAction(actions, candidate); err != nil {
return nil, nil, err
}
}
}
unbinds := make([]pipeline.ReferenceUnbind, 0, len(unbindRequests))
for _, request := range unbindRequests {
target, err := resolveCLIReferenceTarget(targets, request.Selector)
matches, err := resolveCLIReferenceTargets(targets, request.Selector)
if err != nil {
return nil, nil, err
}
unbinds = append(unbinds, pipeline.ReferenceUnbind{
Stage: target.stage,
LaneID: target.laneID,
SlotName: request.Selector.SlotName,
})
for _, target := range matches {
candidate := resolvedCLIReferenceAction{
kind: cliReferenceActionUnbind,
selector: request.Selector,
target: target,
slotName: request.Selector.SlotName,
specificity: request.Selector.specificity(),
}
if err := mergeCLIReferenceAction(actions, candidate); err != nil {
return nil, nil, err
}
}
}
resolved := make([]resolvedCLIReferenceAction, 0, len(actions))
for _, action := range actions {
resolved = append(resolved, action)
}
sort.Slice(resolved, func(i, j int) bool {
left, right := resolved[i], resolved[j]
if left.target.laneID != right.target.laneID {
return left.target.laneID < right.target.laneID
}
if left.target.stage != right.target.stage {
return referenceStageOrder(left.target.stage) < referenceStageOrder(right.target.stage)
}
return left.slotName < right.slotName
})
overrides := make([]pipeline.ReferenceBinding, 0, len(resolved))
unbinds := make([]pipeline.ReferenceUnbind, 0, len(resolved))
for _, action := range resolved {
switch action.kind {
case cliReferenceActionBind:
overrides = append(overrides, pipeline.ReferenceBinding{
Stage: action.target.stage,
LaneID: action.target.laneID,
SlotName: action.slotName,
Source: action.source,
BindingSource: contracts.ReferenceBindingSourceCLI,
})
case cliReferenceActionUnbind:
unbinds = append(unbinds, pipeline.ReferenceUnbind{
Stage: action.target.stage,
LaneID: action.target.laneID,
SlotName: action.slotName,
})
}
}
return overrides, unbinds, nil
}
type cliReferenceActionKind uint8
const (
cliReferenceActionBind cliReferenceActionKind = iota
cliReferenceActionUnbind
)
type cliReferenceTargetKey struct {
stage pipeline.ModuleStage
laneID string
slotName string
}
type resolvedCLIReferenceAction struct {
kind cliReferenceActionKind
selector cliReferenceSelector
target selectedReferenceTarget
slotName string
source string
specificity int
}
func mergeCLIReferenceAction(actions map[cliReferenceTargetKey]resolvedCLIReferenceAction, candidate resolvedCLIReferenceAction) error {
key := cliReferenceTargetKey{stage: candidate.target.stage, laneID: candidate.target.laneID, slotName: candidate.slotName}
current, ok := actions[key]
if !ok || candidate.specificity > current.specificity {
actions[key] = candidate
return nil
}
if candidate.specificity < current.specificity {
return nil
}
if candidate.kind != current.kind {
return fmt.Errorf("reference target %q slot %q is both bound by %q and unbound by %q at the same specificity", targetLabel(candidate.target), candidate.slotName, current.selector.String(), candidate.selector.String())
}
actions[key] = candidate
return nil
}
func (selector cliReferenceSelector) specificity() int {
switch selector.Scope {
case cliReferenceScopePipeline:
return 0
case cliReferenceScopeLane:
return 1
case cliReferenceScopeChunk, cliReferenceScopeBinding:
return 2
default:
return -1
}
}
func (selector cliReferenceSelector) String() string {
switch selector.Scope {
case cliReferenceScopePipeline:
return selector.SlotName
case cliReferenceScopeLane:
return selector.LaneID + "." + selector.SlotName
case cliReferenceScopeChunk:
return "chunk." + selector.SlotName
case cliReferenceScopeBinding:
return selector.LaneID + "." + string(selector.Stage) + "." + selector.SlotName
default:
return selector.SlotName
}
}
type selectedReferenceTarget struct {
laneID string
stage pipeline.ModuleStage
@@ -1638,114 +1761,77 @@ func referenceSlotSet(slots []contracts.ReferenceSlot) map[string]struct{} {
return slotSet
}
func resolveCLIReferenceTarget(targets []selectedReferenceTarget, selector cliReferenceSelector) (selectedReferenceTarget, error) {
func resolveCLIReferenceTargets(targets []selectedReferenceTarget, selector cliReferenceSelector) ([]selectedReferenceTarget, error) {
slotName := strings.TrimSpace(selector.SlotName)
if slotName == "" {
return selectedReferenceTarget{}, fmt.Errorf("reference slot must not be empty")
return nil, fmt.Errorf("reference slot must not be empty")
}
if selector.Stage == pipeline.StageChunk {
switch selector.Scope {
case cliReferenceScopePipeline:
matches := make([]selectedReferenceTarget, 0, len(targets))
for _, target := range targets {
if _, ok := target.slots[slotName]; ok {
matches = append(matches, target)
}
}
if len(matches) == 0 {
return nil, fmt.Errorf("reference slot %q is not declared by any selected target", slotName)
}
return matches, nil
case cliReferenceScopeChunk:
for _, target := range targets {
if target.stage != pipeline.StageChunk {
continue
}
if _, ok := target.slots[slotName]; !ok {
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is not declared by chunk module %q", slotName, target.module)
return nil, fmt.Errorf("reference slot %q is not declared by chunk module %q", slotName, target.module)
}
return target, nil
}
return selectedReferenceTarget{}, fmt.Errorf("reference chunk target is not selected")
}
if selector.Stage == pipeline.StageExtract || selector.Stage == pipeline.StageMerge || selector.Stage == pipeline.StageNormalize {
if selector.LaneID == "" && selector.Stage == pipeline.StageMerge {
return resolveCLIReferenceStageTarget(targets, selector.Stage, slotName)
return []selectedReferenceTarget{target}, nil
}
return nil, fmt.Errorf("reference chunk target is not selected")
case cliReferenceScopeLane:
laneSelected := false
matches := make([]selectedReferenceTarget, 0, 3)
for _, target := range targets {
if target.laneID == selector.LaneID && target.stage == selector.Stage {
if _, ok := target.slots[slotName]; !ok {
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is not declared by selected %s target %q", slotName, selector.Stage, targetLabel(target))
}
return target, nil
if target.laneID != selector.LaneID {
continue
}
laneSelected = true
if _, ok := target.slots[slotName]; ok {
matches = append(matches, target)
}
}
return selectedReferenceTarget{}, fmt.Errorf("reference lane %q is not selected", selector.LaneID)
}
if strings.TrimSpace(selector.LaneID) != "" {
return resolveCLIReferenceLaneTarget(targets, strings.TrimSpace(selector.LaneID), slotName)
}
return resolveCLIReferenceFlatTarget(targets, slotName)
}
func resolveCLIReferenceStageTarget(targets []selectedReferenceTarget, stage pipeline.ModuleStage, slotName string) (selectedReferenceTarget, error) {
matches := make([]selectedReferenceTarget, 0, 2)
for _, target := range targets {
if target.stage != stage {
continue
if !laneSelected {
return nil, fmt.Errorf("reference lane %q is not selected", selector.LaneID)
}
if _, ok := target.slots[slotName]; ok {
matches = append(matches, target)
if len(matches) == 0 {
return nil, fmt.Errorf("reference slot %q is not declared by selected lane %q", slotName, selector.LaneID)
}
}
switch len(matches) {
case 0:
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is not declared by any selected %s target", slotName, stage)
case 1:
return matches[0], nil
return matches, nil
case cliReferenceScopeBinding:
laneSelected := false
for _, target := range targets {
if target.laneID != selector.LaneID {
continue
}
laneSelected = true
if target.stage != selector.Stage {
continue
}
if _, ok := target.slots[slotName]; !ok {
return nil, fmt.Errorf("reference slot %q is not declared by selected %s target %q", slotName, selector.Stage, targetLabel(target))
}
return []selectedReferenceTarget{target}, nil
}
if !laneSelected {
return nil, fmt.Errorf("reference lane %q is not selected", selector.LaneID)
}
return nil, fmt.Errorf("reference %s target is not selected for lane %q", selector.Stage, selector.LaneID)
default:
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is declared by multiple selected %s targets (%s); use a more specific selector such as %s", slotName, stage, targetList(matches), selectorSuggestions(matches, slotName))
return nil, fmt.Errorf("reference selector has unknown scope")
}
}
func resolveCLIReferenceLaneTarget(targets []selectedReferenceTarget, laneID string, slotName string) (selectedReferenceTarget, error) {
laneSelected := false
matches := make([]selectedReferenceTarget, 0, 2)
for _, target := range targets {
if target.laneID != laneID {
continue
}
laneSelected = true
if _, ok := target.slots[slotName]; ok {
matches = append(matches, target)
}
}
if !laneSelected {
return selectedReferenceTarget{}, fmt.Errorf("reference lane %q is not selected", laneID)
}
switch len(matches) {
case 0:
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is not declared by selected lane %q", slotName, laneID)
case 1:
return matches[0], nil
default:
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is declared by multiple selected targets in lane %q (%s); use a more specific selector such as %s", slotName, laneID, targetList(matches), selectorSuggestions(matches, slotName))
}
}
func resolveCLIReferenceFlatTarget(targets []selectedReferenceTarget, slotName string) (selectedReferenceTarget, error) {
matches := make([]selectedReferenceTarget, 0, 2)
for _, target := range targets {
if _, ok := target.slots[slotName]; ok {
matches = append(matches, target)
}
}
switch len(matches) {
case 0:
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is not declared by any selected reference target", slotName)
case 1:
return matches[0], nil
default:
return selectedReferenceTarget{}, fmt.Errorf("reference slot %q is declared by multiple selected targets (%s); use a more specific selector such as %s", slotName, targetList(matches), selectorSuggestions(matches, slotName))
}
}
func targetList(targets []selectedReferenceTarget) string {
labels := make([]string, 0, len(targets))
for _, target := range targets {
labels = append(labels, targetLabel(target))
}
sort.Strings(labels)
return strings.Join(labels, ", ")
}
func targetLabel(target selectedReferenceTarget) string {
if target.stage == pipeline.StageChunk {
return "chunk"
@@ -1753,17 +1839,19 @@ func targetLabel(target selectedReferenceTarget) string {
return target.laneID + "." + string(target.stage)
}
func selectorSuggestions(targets []selectedReferenceTarget, slotName string) string {
suggestions := make([]string, 0, len(targets))
for _, target := range targets {
if target.stage == pipeline.StageChunk {
suggestions = append(suggestions, "chunk."+slotName)
continue
}
suggestions = append(suggestions, target.laneID+"."+string(target.stage)+"."+slotName)
func referenceStageOrder(stage pipeline.ModuleStage) int {
switch stage {
case pipeline.StageChunk:
return 0
case pipeline.StageExtract:
return 1
case pipeline.StageMerge:
return 2
case pipeline.StageNormalize:
return 3
default:
return 4
}
sort.Strings(suggestions)
return strings.Join(suggestions, " or ")
}
func sortedPipelineIDs(cfg config.Config) []string {