Invalidate downstream results when stages are replaced
This commit is contained in:
@@ -3,6 +3,7 @@ package app
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
@@ -19,12 +20,17 @@ import (
|
||||
)
|
||||
|
||||
type materializingNotariusRunner struct {
|
||||
cfg *config.NotariusConfig
|
||||
requests []notarius.RunRequest
|
||||
cfg *config.NotariusConfig
|
||||
requests []notarius.RunRequest
|
||||
failuresRemaining int
|
||||
}
|
||||
|
||||
func (r *materializingNotariusRunner) Run(_ context.Context, req notarius.RunRequest) (notarius.RunResult, error) {
|
||||
r.requests = append(r.requests, req)
|
||||
if r.failuresRemaining > 0 {
|
||||
r.failuresRemaining--
|
||||
return notarius.RunResult{}, errors.New("notarius execution failed")
|
||||
}
|
||||
externalRunID := fmt.Sprintf("notarius-run-%d", len(r.requests))
|
||||
bundle := filepath.Join(req.OutputRoot, externalRunID)
|
||||
lanesDir := filepath.Join(bundle, "lanes")
|
||||
@@ -94,9 +100,121 @@ func TestExtractLifecycleDisabledThenEnabled(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractLifecycleChangedOutcomeRerunsSucceededDownstream(t *testing.T) {
|
||||
cfg, env, runner := extractionLifecycleFixture(t, false)
|
||||
analyzeRuns := 0
|
||||
plan := extractionLifecyclePlan(t, &analyzeRuns)
|
||||
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
|
||||
t.Fatalf("disabled executeStages() error = %v", err)
|
||||
}
|
||||
if analyzeRuns != 1 || len(runner.requests) != 0 {
|
||||
t.Fatalf("disabled run analyze=%d Notarius=%d, want 1 and 0", analyzeRuns, len(runner.requests))
|
||||
}
|
||||
|
||||
cfg.Pipeline.Notarius.Enabled = true
|
||||
summary, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
|
||||
if err != nil {
|
||||
t.Fatalf("enabled executeStages() error = %v", err)
|
||||
}
|
||||
if analyzeRuns != 2 || len(runner.requests) != 1 {
|
||||
t.Fatalf("enabled run analyze=%d Notarius=%d, want 2 and 1", analyzeRuns, len(runner.requests))
|
||||
}
|
||||
if len(summary.Executed) != 2 || len(summary.Skipped) != 0 {
|
||||
t.Fatalf("enabled summary = %#v, want extract and analyze executed", summary)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractLifecycleFailureInvalidatesAndOrdinaryRetryRerunsDownstream(t *testing.T) {
|
||||
cfg, env, runner := extractionLifecycleFixture(t, true)
|
||||
runner.failuresRemaining = 1
|
||||
markLifecycleStageSucceeded(t, cfg, "analyze")
|
||||
analyzeRuns := 0
|
||||
plan := extractionLifecyclePlan(t, &analyzeRuns)
|
||||
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err == nil || !strings.Contains(err.Error(), "notarius execution failed") {
|
||||
t.Fatalf("failed executeStages() error = %v", err)
|
||||
}
|
||||
failed := loadLifecycleManifest(t, cfg)
|
||||
if failed.Stages["extract"].Status != manifest.StatusFailed || failed.Stages["analyze"].Status != manifest.StatusStale {
|
||||
t.Fatalf("failed lifecycle extract=%#v analyze=%#v", failed.Stages["extract"], failed.Stages["analyze"])
|
||||
}
|
||||
if failed.Stages["analyze"].Error == nil || failed.Stages["analyze"].Error.Message != staleReasonFailure {
|
||||
t.Fatalf("analyze stale reason = %#v, want %q", failed.Stages["analyze"].Error, staleReasonFailure)
|
||||
}
|
||||
|
||||
summary, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
|
||||
if err != nil {
|
||||
t.Fatalf("retry executeStages() error = %v", err)
|
||||
}
|
||||
if analyzeRuns != 1 || len(runner.requests) != 2 || len(summary.Executed) != 2 {
|
||||
t.Fatalf("retry analyze=%d Notarius=%d summary=%#v", analyzeRuns, len(runner.requests), summary)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractLifecycleForcedSelfSkipInvalidatesDownstream(t *testing.T) {
|
||||
cfg, env, _ := extractionLifecycleFixture(t, false)
|
||||
markLifecycleStageSucceeded(t, cfg, "analyze")
|
||||
plan, _ := BuildSingleStagePlan("extract")
|
||||
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env, Force: true}); err != nil {
|
||||
t.Fatalf("executeStages() error = %v", err)
|
||||
}
|
||||
loaded := loadLifecycleManifest(t, cfg)
|
||||
if loaded.Stages["extract"].Status != manifest.StatusSkipped || loaded.Stages["analyze"].Status != manifest.StatusStale {
|
||||
t.Fatalf("forced self-skip extract=%#v analyze=%#v", loaded.Stages["extract"], loaded.Stages["analyze"])
|
||||
}
|
||||
if loaded.Stages["analyze"].Error == nil || loaded.Stages["analyze"].Error.Message != staleReasonForcedReplacement {
|
||||
t.Fatalf("analyze stale reason = %#v, want %q", loaded.Stages["analyze"].Error, staleReasonForcedReplacement)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractLifecycleForcedFailureInvalidatesDownstream(t *testing.T) {
|
||||
cfg, env, runner := extractionLifecycleFixture(t, true)
|
||||
runner.failuresRemaining = 1
|
||||
markLifecycleStageSucceeded(t, cfg, "analyze")
|
||||
plan, _ := BuildSingleStagePlan("extract")
|
||||
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env, Force: true}); err == nil {
|
||||
t.Fatal("executeStages() error = nil, want forced extraction failure")
|
||||
}
|
||||
loaded := loadLifecycleManifest(t, cfg)
|
||||
if loaded.Stages["extract"].Status != manifest.StatusFailed || loaded.Stages["analyze"].Status != manifest.StatusStale {
|
||||
t.Fatalf("forced failure extract=%#v analyze=%#v", loaded.Stages["extract"], loaded.Stages["analyze"])
|
||||
}
|
||||
if loaded.Stages["analyze"].Error == nil || loaded.Stages["analyze"].Error.Message != staleReasonForcedReplacement {
|
||||
t.Fatalf("analyze stale reason = %#v, want %q", loaded.Stages["analyze"].Error, staleReasonForcedReplacement)
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractLifecycleRepeatedSelfSkipPreservesSucceededDownstream(t *testing.T) {
|
||||
cfg, env, runner := extractionLifecycleFixture(t, false)
|
||||
analyzeRuns := 0
|
||||
plan := extractionLifecyclePlan(t, &analyzeRuns)
|
||||
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
|
||||
t.Fatalf("first executeStages() error = %v", err)
|
||||
}
|
||||
second, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
|
||||
if err != nil {
|
||||
t.Fatalf("second executeStages() error = %v", err)
|
||||
}
|
||||
if analyzeRuns != 1 || len(runner.requests) != 0 {
|
||||
t.Fatalf("repeated disabled run analyze=%d Notarius=%d, want 1 and 0", analyzeRuns, len(runner.requests))
|
||||
}
|
||||
if len(second.Executed) != 1 || len(second.Skipped) != 2 {
|
||||
t.Fatalf("second summary = %#v, want executed self-skip and skipped analyze", second)
|
||||
}
|
||||
loaded := loadLifecycleManifest(t, cfg)
|
||||
if loaded.Stages["analyze"].Status != manifest.StatusSucceeded {
|
||||
t.Fatalf("analyze = %#v, want succeeded", loaded.Stages["analyze"])
|
||||
}
|
||||
}
|
||||
|
||||
func TestExtractLifecycleSkipsCurrentResumableResult(t *testing.T) {
|
||||
cfg, env, runner := extractionLifecycleFixture(t, true)
|
||||
plan, _ := BuildSingleStagePlan("extract")
|
||||
analyzeRuns := 0
|
||||
plan := extractionLifecyclePlan(t, &analyzeRuns)
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err != nil {
|
||||
t.Fatalf("first executeStages() error = %v", err)
|
||||
}
|
||||
@@ -104,8 +222,8 @@ func TestExtractLifecycleSkipsCurrentResumableResult(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatalf("resume executeStages() error = %v", err)
|
||||
}
|
||||
if len(resumed.Executed) != 0 || len(resumed.Skipped) != 1 || resumed.Skipped[0] != "extract" || len(runner.requests) != 1 {
|
||||
t.Fatalf("resume summary = %#v requests=%d", resumed, len(runner.requests))
|
||||
if len(resumed.Executed) != 0 || len(resumed.Skipped) != 2 || len(runner.requests) != 1 || analyzeRuns != 1 {
|
||||
t.Fatalf("resume summary = %#v requests=%d analyze=%d", resumed, len(runner.requests), analyzeRuns)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -162,7 +280,8 @@ func TestExtractLifecycleResumesAndRerunsObsoleteResults(t *testing.T) {
|
||||
|
||||
func TestExtractLifecycleUnsafeResumeErrorPreservesSuccess(t *testing.T) {
|
||||
cfg, env, runner := extractionLifecycleFixture(t, true)
|
||||
plan, _ := BuildSingleStagePlan("extract")
|
||||
analyzeRuns := 0
|
||||
plan := extractionLifecyclePlan(t, &analyzeRuns)
|
||||
first, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env})
|
||||
if err != nil {
|
||||
t.Fatalf("first executeStages() error = %v", err)
|
||||
@@ -176,7 +295,10 @@ func TestExtractLifecycleUnsafeResumeErrorPreservesSuccess(t *testing.T) {
|
||||
if err := store.Save(context.Background(), first.ManifestPath, persisted); err != nil {
|
||||
t.Fatalf("Save() error = %v", err)
|
||||
}
|
||||
before, _ := json.Marshal(persisted.Stages["extract"])
|
||||
before, _ := json.Marshal(map[string]*manifest.StageRecord{
|
||||
"extract": persisted.Stages["extract"],
|
||||
"analyze": persisted.Stages["analyze"],
|
||||
})
|
||||
|
||||
if _, err := executeStages(context.Background(), cfg, plan, RunOptions{Env: env}); err == nil || !strings.Contains(err.Error(), "unsafe") {
|
||||
t.Fatalf("executeStages() error = %v, want unsafe resume failure", err)
|
||||
@@ -185,9 +307,39 @@ func TestExtractLifecycleUnsafeResumeErrorPreservesSuccess(t *testing.T) {
|
||||
if err != nil {
|
||||
t.Fatalf("Load(after) error = %v", err)
|
||||
}
|
||||
after, _ := json.Marshal(afterManifest.Stages["extract"])
|
||||
if string(before) != string(after) || len(runner.requests) != 1 {
|
||||
t.Fatalf("successful extract record changed: before=%s after=%s requests=%d", before, after, len(runner.requests))
|
||||
after, _ := json.Marshal(map[string]*manifest.StageRecord{
|
||||
"extract": afterManifest.Stages["extract"],
|
||||
"analyze": afterManifest.Stages["analyze"],
|
||||
})
|
||||
if string(before) != string(after) || len(runner.requests) != 1 || analyzeRuns != 1 {
|
||||
t.Fatalf("successful records changed: before=%s after=%s requests=%d analyze=%d", before, after, len(runner.requests), analyzeRuns)
|
||||
}
|
||||
}
|
||||
|
||||
func extractionLifecyclePlan(t *testing.T, analyzeRuns *int) []stage.Stage {
|
||||
t.Helper()
|
||||
plan, err := BuildSingleStagePlan("extract")
|
||||
if err != nil {
|
||||
t.Fatalf("BuildSingleStagePlan(extract) error = %v", err)
|
||||
}
|
||||
return append(plan, countingStage{name: "analyze", runs: analyzeRuns})
|
||||
}
|
||||
|
||||
func loadLifecycleManifest(t *testing.T, cfg *config.Config) *manifest.Manifest {
|
||||
t.Helper()
|
||||
loaded, err := (&manifest.LocalStore{}).Load(context.Background(), manifestPathFor(cfg))
|
||||
if err != nil {
|
||||
t.Fatalf("Load() error = %v", err)
|
||||
}
|
||||
return loaded
|
||||
}
|
||||
|
||||
func markLifecycleStageSucceeded(t *testing.T, cfg *config.Config, name string) {
|
||||
t.Helper()
|
||||
loaded := loadLifecycleManifest(t, cfg)
|
||||
loaded.MarkStageSucceeded(name, time.Now().UTC(), nil)
|
||||
if err := (&manifest.LocalStore{}).Save(context.Background(), manifestPathFor(cfg), loaded); err != nil {
|
||||
t.Fatalf("Save() error = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user