package pipeline import ( "context" "encoding/json" "reflect" "strings" "testing" "gitea.maximumdirect.net/eric/notarius/internal/framework/contracts" ) type acceptedCheckpointLoader struct { CheckpointLoader accepted map[string]NormalizeCheckpoint acceptedDecision map[string]CheckpointDecision extractDeps map[string][]CheckpointFingerprint } func newAcceptedCheckpointLoader() *acceptedCheckpointLoader { return &acceptedCheckpointLoader{ CheckpointLoader: NoopCheckpointLoader(), accepted: make(map[string]NormalizeCheckpoint), acceptedDecision: make(map[string]CheckpointDecision), extractDeps: make(map[string][]CheckpointFingerprint), } } func (l *acceptedCheckpointLoader) Enabled() bool { return true } func (l *acceptedCheckpointLoader) AcceptedNormalize(stepID, laneID, _ string) (NormalizeCheckpoint, CheckpointDecision) { key := CheckpointLaneKey(stepID, laneID) return l.accepted[key], l.acceptedDecision[key] } func (l *acceptedCheckpointLoader) Extract(laneID string, _ string, dependencies []CheckpointFingerprint) (ExtractCheckpoint, CheckpointDecision) { l.extractDeps[laneID] = append([]CheckpointFingerprint(nil), dependencies...) return ExtractCheckpoint{}, NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonMissing) } func (l *acceptedCheckpointLoader) Merge(_ string, _ string, _ []CheckpointFingerprint) (MergeCheckpoint, CheckpointDecision) { return MergeCheckpoint{}, NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonMissing) } func (l *acceptedCheckpointLoader) Normalize(_ string, _ string, _ []CheckpointFingerprint) (NormalizeCheckpoint, CheckpointDecision) { return NormalizeCheckpoint{}, NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonMissing) } func TestRunnerHydratesRequiredNormalizedArtifact(t *testing.T) { value := codecNotes{Items: []string{"canonical producer value"}} input, _, _ := handoffFixture(t, value) prepared := input.Prepared producer := &prepared.Steps[0].lanes[0] consumer := &prepared.Steps[1].lanes[0] doc := prepared.input.(*typedTestInput).doc stored, err := checkpointArtifact(producer.typed.codec, producer.resolved.ID, producer.resolved.Normalize.Module, doc.ID, value) if err != nil { t.Fatal(err) } operationCalls := 0 producer.typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { operationCalls++ return erasedTypedResult{Value: value}, nil } producer.typed.merge = func(context.Context, any, contracts.TypedMergeRequest[any]) (erasedTypedResult, error) { operationCalls++ return erasedTypedResult{Value: value}, nil } producer.typed.normalize = func(context.Context, any, contracts.TypedNormalizeRequest[any]) (erasedTypedResult, error) { operationCalls++ return erasedTypedResult{Value: value}, nil } validatorCalls := 0 validator := preparedValidator{typedValidate: func(context.Context, any, typedValidationTarget) (contracts.ValidationResult, error) { validatorCalls++ return contracts.ValidationResult{Approved: true}, nil }} producer.extractValidators.validators = []preparedValidator{validator} producer.mergeValidators.validators = []preparedValidator{validator} producer.normalizeValidators.validators = []preparedValidator{validator} var received contracts.ReferenceSet consumer.typed.extract = func(_ context.Context, _ any, request contracts.TypedExtractionRequest) (erasedTypedResult, error) { received = CloneReferenceSet(request.References) return erasedTypedResult{Value: codecScore{Value: 3}}, nil } loader := newAcceptedCheckpointLoader() producerKey := CheckpointLaneKey(producer.resolved.StepID, producer.resolved.ID) consumerKey := CheckpointLaneKey(consumer.resolved.StepID, consumer.resolved.ID) loader.accepted[producerKey] = NormalizeCheckpoint{Output: stored, Warnings: []contracts.Warning{{Scope: "normalize", ReasonCode: "stored-warning", Message: "stored normalize warning"}}} loader.acceptedDecision[producerKey] = NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused) policy := CheckpointExecutionPolicy{ RequireReusableLanes: map[string]struct{}{producerKey: {}}, ForcedLanes: map[string]struct{}{consumerKey: {}}, } output, err := New().Run(context.Background(), RunInput{Prepared: prepared, RawInput: []byte("input"), Checkpoint: loader, CheckpointPolicy: policy}) if err != nil { t.Fatalf("Run() error = %v", err) } if operationCalls != 0 || validatorCalls != 0 { t.Fatalf("hydrated producer calls = operations %d validators %d, want zero", operationCalls, validatorCalls) } item := received.Slots["producer-output"].Items[0] if string(item.Content) != string(stored.Artifact.Content) || item.Producer.StepID != producer.resolved.StepID || item.Producer.LaneID != producer.resolved.ID { t.Fatalf("consumer generated reference = %#v, want exact hydrated producer bytes and identity", item) } if len(output.Warnings) != 1 || output.Warnings[0].ReasonCode != "stored-warning" { t.Fatalf("hydrated warnings = %#v, want normalize checkpoint warnings only", output.Warnings) } assertAcceptedNormalizeEvent(t, output.CheckpointEvents, producer.resolved.StepID, producer.resolved.ID, CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused) for _, event := range output.CheckpointEvents { if event.StepID == producer.resolved.StepID && event.LaneID == producer.resolved.ID && event.Stage != string(StageNormalize) { t.Fatalf("hydrated producer synthesized checkpoint event: %#v", event) } } freshInput, _, _ := handoffFixture(t, value) freshPrepared := freshInput.Prepared freshPrepared.Steps[0].lanes[0].typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { return erasedTypedResult{Value: value}, nil } freshPrepared.Steps[0].lanes[0].typed.normalize = func(context.Context, any, contracts.TypedNormalizeRequest[any]) (erasedTypedResult, error) { return erasedTypedResult{Value: value}, nil } freshPrepared.Steps[1].lanes[0].typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { return erasedTypedResult{Value: codecScore{Value: 3}}, nil } freshLoader := &handoffDependencyLoader{CheckpointLoader: NoopCheckpointLoader()} freshOutput, err := New().Run(context.Background(), RunInput{Prepared: freshPrepared, RawInput: []byte("input"), Checkpoint: freshLoader}) if err != nil { t.Fatalf("fresh Run() error = %v", err) } hydratedDependencies := loader.extractDeps[consumer.resolved.ID] if generatedFingerprintCount(hydratedDependencies) != 1 { t.Fatalf("hydrated consumer dependencies = %#v, want generated producer fingerprint", hydratedDependencies) } var matchedFreshDependencies bool for _, dependencies := range freshLoader.extract { if reflect.DeepEqual(dependencies, hydratedDependencies) { matchedFreshDependencies = true break } } if !matchedFreshDependencies { t.Fatalf("consumer dependencies differ: fresh %#v hydrated %#v", freshLoader.extract, hydratedDependencies) } if !reflect.DeepEqual(freshOutput.Manifest.References, output.Manifest.References) { t.Fatalf("generated provenance differs: fresh %#v hydrated %#v", freshOutput.Manifest.References, output.Manifest.References) } } func TestRunnerRejectsInvalidRequiredNormalizedArtifactBeforeConsumer(t *testing.T) { const contentSentinel = "sensitive-campaign-payload-74291" tests := []struct { name string decision CheckpointDecision mutate func(*CheckpointArtifact) wantCode CheckpointReasonCode }{ {"missing", NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonMissing), nil, CheckpointReasonMissing}, {"rejected status", NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonStatusNotReusable), nil, CheckpointReasonStatusNotReusable}, {"corrupt payload", NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused), func(v *CheckpointArtifact) { v.Artifact.Content = []byte(`{"items":[`) }, CheckpointReasonArtifactPayloadInvalid}, {"non canonical", NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused), func(v *CheckpointArtifact) { v.Artifact.Content = []byte(`{"items": ["stored"]}`) }, CheckpointReasonArtifactNotCanonical}, {"wrong codec identity", NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused), func(v *CheckpointArtifact) { v.Artifact.Kind = "test/score" }, CheckpointReasonArtifactCodecIncompatible}, {"wrong content digest", NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonArtifactDigestMismatch), nil, CheckpointReasonArtifactDigestMismatch}, } for _, test := range tests { t.Run(test.name, func(t *testing.T) { input, _, _ := handoffFixture(t, codecNotes{Items: []string{contentSentinel}}) prepared := input.Prepared producer := &prepared.Steps[0].lanes[0] consumer := &prepared.Steps[1].lanes[0] doc := prepared.input.(*typedTestInput).doc stored, err := checkpointArtifact(producer.typed.codec, producer.resolved.ID, producer.resolved.Normalize.Module, doc.ID, codecNotes{Items: []string{contentSentinel}}) if err != nil { t.Fatal(err) } if test.mutate != nil { test.mutate(&stored) } consumerCalls := 0 consumer.typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { consumerCalls++ return erasedTypedResult{Value: codecScore{Value: 1}}, nil } loader := newAcceptedCheckpointLoader() producerKey := CheckpointLaneKey(producer.resolved.StepID, producer.resolved.ID) loader.accepted[producerKey] = NormalizeCheckpoint{Output: stored} loader.acceptedDecision[producerKey] = test.decision policy := CheckpointExecutionPolicy{RequireReusableLanes: map[string]struct{}{producerKey: {}}} output, runErr := New().Run(context.Background(), RunInput{Prepared: prepared, RawInput: []byte("input"), Checkpoint: loader, CheckpointPolicy: policy}) if runErr == nil || !strings.Contains(runErr.Error(), string(test.wantCode)) || consumerCalls != 0 { t.Fatalf("Run() error = %v consumer calls = %d, want %q before consumer", runErr, consumerCalls, test.wantCode) } assertAcceptedNormalizeEvent(t, output.CheckpointEvents, producer.resolved.StepID, producer.resolved.ID, CheckpointDecisionExecuted, test.wantCode) encoded, err := json.Marshal(struct { Manifest any Events any Error string }{output.Manifest, output.CheckpointEvents, runErr.Error()}) if err != nil { t.Fatal(err) } if strings.Contains(string(encoded), contentSentinel) || strings.Contains(string(encoded), string(stored.Artifact.Content)) { t.Fatalf("failed hydration diagnostics leaked artifact content: %s", encoded) } }) } } func TestRunnerRetainsEarlierHydratedProducerWhenLaterRequiredProducerFails(t *testing.T) { prepared := preparedPipelineWithSharedProducerStep(t) first := &prepared.Steps[0].lanes[0] second := &prepared.Steps[0].lanes[1] consumer := &prepared.Steps[1].lanes[0] doc := prepared.input.(*typedTestInput).doc stored, err := checkpointArtifact(first.typed.codec, first.resolved.ID, first.resolved.Normalize.Module, doc.ID, codecNotes{Items: []string{"retained producer"}}) if err != nil { t.Fatal(err) } consumerCalls := 0 consumer.typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { consumerCalls++ return erasedTypedResult{Value: codecScore{Value: 1}}, nil } loader := newAcceptedCheckpointLoader() firstKey := CheckpointLaneKey(first.resolved.StepID, first.resolved.ID) secondKey := CheckpointLaneKey(second.resolved.StepID, second.resolved.ID) loader.accepted[firstKey] = NormalizeCheckpoint{Output: stored, Warnings: []contracts.Warning{{Scope: "normalize", ReasonCode: "retained-warning", Message: "retained warning"}}} loader.acceptedDecision[firstKey] = NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused) loader.acceptedDecision[secondKey] = NewCheckpointDecision(CheckpointDecisionExecuted, CheckpointReasonMissing) policy := CheckpointExecutionPolicy{RequireReusableLanes: map[string]struct{}{firstKey: {}, secondKey: {}}} output, runErr := New().Run(context.Background(), RunInput{Prepared: prepared, RawInput: []byte("input"), Checkpoint: loader, CheckpointPolicy: policy}) if runErr == nil || !strings.Contains(runErr.Error(), string(CheckpointReasonMissing)) { t.Fatalf("Run() error = %v, want missing required producer", runErr) } if consumerCalls != 0 { t.Fatalf("consumer calls = %d, want zero", consumerCalls) } if len(output.NormalizeOutputs) != 1 || output.NormalizeOutputs[0].LaneID != first.resolved.ID || string(output.NormalizeOutputs[0].Artifact.Content) != string(stored.Artifact.Content) { t.Fatalf("retained normalize outputs = %#v, want first producer", output.NormalizeOutputs) } if len(output.Warnings) != 1 || output.Warnings[0].ReasonCode != "retained-warning" { t.Fatalf("retained warnings = %#v", output.Warnings) } type decisionExpectation struct { step, lane string category CheckpointDecisionCategory reason CheckpointReasonCode } want := []decisionExpectation{ {first.resolved.StepID, first.resolved.ID, CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused}, {second.resolved.StepID, second.resolved.ID, CheckpointDecisionExecuted, CheckpointReasonMissing}, } var got []decisionExpectation for _, event := range output.CheckpointEvents { if event.Stage == string(StageNormalize) { got = append(got, decisionExpectation{event.StepID, event.LaneID, event.Category, event.ReasonCode}) } } if !reflect.DeepEqual(got, want) { t.Fatalf("normalize decisions = %#v, want %#v", got, want) } var manifestGot []decisionExpectation for _, decision := range output.Manifest.CheckpointDecisions { if decision.Stage == string(StageNormalize) { manifestGot = append(manifestGot, decisionExpectation{ decision.StepID, decision.LaneID, CheckpointDecisionCategory(decision.Category), CheckpointReasonCode(decision.ReasonCode), }) } } if !reflect.DeepEqual(manifestGot, want) { t.Fatalf("manifest normalize decisions = %#v, want %#v", manifestGot, want) } } func preparedPipelineWithSharedProducerStep(t *testing.T) *PreparedPipeline { t.Helper() catalog := typedResolutionCatalog(t, completeTypedCatalogOptions()) base := typedResolutionProfile() profile := base profile.Artifacts = nil profile.Steps = []PipelineStepProfile{ {ID: "producers", Artifacts: map[string]ArtifactLaneProfile{"first": base.Artifacts["notes"], "second": base.Artifacts["notes"]}}, {ID: "consumer", Artifacts: map[string]ArtifactLaneProfile{"score": base.Artifacts["score"]}}, } resolved, err := ResolvePipeline(profile, ResolveOptions{}, catalog) if err != nil { t.Fatalf("ResolvePipeline() error = %v", err) } prepared, err := Prepare(resolved, registriesFromModuleCatalog(catalog), ModuleDependencies{}) if err != nil { t.Fatalf("Prepare() error = %v", err) } doc := typedTestDocumentWithUnits(1) prepared.input.(*typedTestInput).doc = doc prepared.chunker = &typedTestChunker{key: "typed/chunk", plan: typedTestPlan(doc)} return prepared } func TestForcedRequiredLaneExecutesInsteadOfHydrating(t *testing.T) { prepared := preparedOrderedPipeline(t, 1, orderedLaneSpec{id: "unrelated", profile: "score"}, orderedLaneSpec{id: "producer", profile: "notes"}, orderedLaneSpec{id: "consumer", profile: "score"}, ) unrelated := &prepared.Steps[0].lanes[0] producer := &prepared.Steps[1].lanes[0] consumer := &prepared.Steps[2].lanes[0] installGeneratedReferenceTarget(&consumer.resolved.ExtractReferences, StageExtract, consumer, "step-2", "producer") doc := prepared.input.(*typedTestInput).doc unrelatedArtifact, err := checkpointArtifact(unrelated.typed.codec, unrelated.resolved.ID, unrelated.resolved.Normalize.Module, doc.ID, codecScore{Value: 9}) if err != nil { t.Fatal(err) } producerArtifact, err := checkpointArtifact(producer.typed.codec, producer.resolved.ID, producer.resolved.Normalize.Module, doc.ID, codecNotes{Items: []string{"stale"}}) if err != nil { t.Fatal(err) } unrelatedCalls, producerCalls := 0, 0 unrelated.typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { unrelatedCalls++ return erasedTypedResult{Value: codecScore{Value: 9}}, nil } producer.typed.extract = func(context.Context, any, contracts.TypedExtractionRequest) (erasedTypedResult, error) { producerCalls++ return erasedTypedResult{Value: codecNotes{Items: []string{"fresh"}}}, nil } producer.typed.normalize = func(context.Context, any, contracts.TypedNormalizeRequest[any]) (erasedTypedResult, error) { return erasedTypedResult{Value: codecNotes{Items: []string{"fresh"}}}, nil } var consumerReferences contracts.ReferenceSet consumer.typed.extract = func(_ context.Context, _ any, request contracts.TypedExtractionRequest) (erasedTypedResult, error) { consumerReferences = CloneReferenceSet(request.References) return erasedTypedResult{Value: codecScore{Value: 1}}, nil } loader := newAcceptedCheckpointLoader() unrelatedKey := CheckpointLaneKey(unrelated.resolved.StepID, unrelated.resolved.ID) producerKey := CheckpointLaneKey(producer.resolved.StepID, producer.resolved.ID) consumerKey := CheckpointLaneKey(consumer.resolved.StepID, consumer.resolved.ID) loader.accepted[unrelatedKey] = NormalizeCheckpoint{Output: unrelatedArtifact} loader.acceptedDecision[unrelatedKey] = NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused) loader.accepted[producerKey] = NormalizeCheckpoint{Output: producerArtifact} loader.acceptedDecision[producerKey] = NewCheckpointDecision(CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused) policy := CheckpointExecutionPolicy{ ForcedLanes: map[string]struct{}{producerKey: {}, consumerKey: {}}, RequireReusableLanes: map[string]struct{}{unrelatedKey: {}, producerKey: {}}, } output, err := New().Run(context.Background(), RunInput{Prepared: prepared, RawInput: []byte("input"), Checkpoint: loader, CheckpointPolicy: policy}) if err != nil { t.Fatalf("Run() error = %v", err) } if producerCalls == 0 || unrelatedCalls != 0 { t.Fatalf("calls producer=%d unrelated=%d, want forced producer execution and unrelated hydration", producerCalls, unrelatedCalls) } if got := string(consumerReferences.Slots["producer-output"].Items[0].Content); !strings.Contains(got, "fresh") || strings.Contains(got, "stale") { t.Fatalf("forced producer reference = %q, want freshly executed output", got) } assertAcceptedNormalizeEvent(t, output.CheckpointEvents, unrelated.resolved.StepID, unrelated.resolved.ID, CheckpointDecisionReused, CheckpointReasonAcceptedArtifactReused) assertAcceptedNormalizeEvent(t, output.CheckpointEvents, producer.resolved.StepID, producer.resolved.ID, CheckpointDecisionForcedRecompute, CheckpointReasonRecomputeStep) } func installGeneratedReferenceTarget(target *ResolvedReferenceTarget, stage ModuleStage, consumer *preparedLaneExecutor, producerStep, producerLane string) { *target = ResolvedReferenceTarget{ Stage: stage, StepID: consumer.resolved.StepID, LaneID: consumer.resolved.ID, Module: consumer.resolved.Extract.Module, Bindings: []ReferenceBinding{{ Stage: stage, LaneID: consumer.resolved.ID, SlotName: "producer-output", Artifact: &ArtifactReference{Step: producerStep, Lane: producerLane}, }}, ReferenceSet: contracts.ReferenceSet{Slots: map[string]contracts.ResolvedReferenceSlot{ "producer-output": {Slot: contracts.ReferenceSlot{Name: "producer-output", AcceptedArtifactKinds: []contracts.ArtifactKind{"test/notes"}, AcceptedMediaTypes: []string{"application/json"}}}, }}, } } func assertAcceptedNormalizeEvent(t *testing.T, events []CheckpointEvent, stepID, laneID string, category CheckpointDecisionCategory, code CheckpointReasonCode) { t.Helper() for _, event := range events { if event.Stage == string(StageNormalize) && event.StepID == stepID && event.LaneID == laneID { if event.Category != category || event.ReasonCode != code { t.Fatalf("accepted normalize event = %#v, want %q/%q", event, category, code) } return } } t.Fatalf("accepted normalize event missing from %#v", events) }