Integrate previous-session artifact hydration into prepare stage
This commit is contained in:
@@ -555,6 +555,12 @@ func needsObjectStoreForRun(cfg *config.Config, stages []stage.Stage) bool {
|
|||||||
if cfg.Session.Inputs.AudioS3 != nil && stageRequested("prepare") {
|
if cfg.Session.Inputs.AudioS3 != nil && stageRequested("prepare") {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
if stageRequested("prepare") {
|
||||||
|
requirements := artifacts.CollectPreviousArtifactRequirements(configuredScriptoriumArtifacts(cfg))
|
||||||
|
if len(requirements) > 0 && strings.TrimSpace(cfg.Session.PreviousSessionID) != "" {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
}
|
||||||
if !stageRequested("archive") {
|
if !stageRequested("archive") {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
@@ -569,3 +575,10 @@ func needsObjectStoreForRun(cfg *config.Config, stages []stage.Stage) bool {
|
|||||||
}
|
}
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func configuredScriptoriumArtifacts(cfg *config.Config) map[string]config.ScriptoriumArtifactConfig {
|
||||||
|
if cfg == nil || cfg.Pipeline == nil || cfg.Pipeline.Scriptorium == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return cfg.Pipeline.Scriptorium.Artifacts
|
||||||
|
}
|
||||||
|
|||||||
@@ -199,6 +199,55 @@ func TestExecuteStagesAnalyzeOutputsPersistAsScriptoriumArtifacts(t *testing.T)
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestNeedsObjectStoreForRunPrepareWithPreviousRequirements(t *testing.T) {
|
||||||
|
tests := []struct {
|
||||||
|
name string
|
||||||
|
previousSessionID string
|
||||||
|
want bool
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "previous session configured",
|
||||||
|
previousSessionID: "2026-05-10",
|
||||||
|
want: true,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "previous session missing",
|
||||||
|
previousSessionID: "",
|
||||||
|
want: false,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, tt := range tests {
|
||||||
|
t.Run(tt.name, func(t *testing.T) {
|
||||||
|
cfg := &config.Config{
|
||||||
|
Pipeline: &config.PipelineConfig{
|
||||||
|
Scriptorium: &config.ScriptoriumConfig{
|
||||||
|
Artifacts: map[string]config.ScriptoriumArtifactConfig{
|
||||||
|
"session_recap": {
|
||||||
|
Enabled: true,
|
||||||
|
Inputs: map[string]config.ScriptoriumInputConfig{
|
||||||
|
"previous_recap": {
|
||||||
|
Source: "narratio.previous_session.artifact.session_recap",
|
||||||
|
Required: true,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
Session: &config.SessionConfig{
|
||||||
|
PreviousSessionID: tt.previousSessionID,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
got := needsObjectStoreForRun(cfg, []stage.Stage{countingStage{name: "prepare", runs: new(int)}})
|
||||||
|
if got != tt.want {
|
||||||
|
t.Fatalf("needsObjectStoreForRun() = %v, want %v", got, tt.want)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func TestExecuteStagesArchiveFailsWhenRequiredRecapPromotionMissingForSelectedArtifacts(t *testing.T) {
|
func TestExecuteStagesArchiveFailsWhenRequiredRecapPromotionMissingForSelectedArtifacts(t *testing.T) {
|
||||||
cfg := testConfig(t)
|
cfg := testConfig(t)
|
||||||
cfg.Pipeline.Storage.S3 = &config.StorageS3Config{
|
cfg.Pipeline.Storage.S3 = &config.StorageS3Config{
|
||||||
|
|||||||
@@ -145,6 +145,20 @@ func (prepareStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
previousRequirements := collectPreparePreviousRequirements(env.Config)
|
||||||
|
var previousHydration *previousSessionHydrationResult
|
||||||
|
if len(previousRequirements) > 0 {
|
||||||
|
if err := clearManagedPreviousState(paths); err != nil {
|
||||||
|
return nil, fmt.Errorf("prepare: clear previous-session cache: %w", err)
|
||||||
|
}
|
||||||
|
hydration, err := hydratePreviousSessionArtifacts(ctx, env, paths, previousRequirements)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("prepare: hydrate previous-session artifacts: %w", err)
|
||||||
|
}
|
||||||
|
previousHydration = hydration
|
||||||
|
inputs = append(inputs, hydration.Inputs...)
|
||||||
|
}
|
||||||
|
|
||||||
sort.Slice(inputs, func(i, j int) bool {
|
sort.Slice(inputs, func(i, j int) bool {
|
||||||
if inputs[i].Kind != inputs[j].Kind {
|
if inputs[i].Kind != inputs[j].Kind {
|
||||||
return inputs[i].Kind < inputs[j].Kind
|
return inputs[i].Kind < inputs[j].Kind
|
||||||
@@ -153,13 +167,27 @@ func (prepareStage) Run(ctx context.Context, env *Env, m *manifest.Manifest) (*S
|
|||||||
})
|
})
|
||||||
m.Inputs = inputs
|
m.Inputs = inputs
|
||||||
|
|
||||||
return &StageResult{
|
metadata := map[string]any{
|
||||||
Metadata: map[string]any{
|
|
||||||
"prepared": true,
|
"prepared": true,
|
||||||
"stage": "prepare",
|
"stage": "prepare",
|
||||||
"inputs_count": len(inputs),
|
"inputs_count": len(inputs),
|
||||||
"audio_files_resolved": countAudioInputs(inputs),
|
"audio_files_resolved": countAudioInputs(inputs),
|
||||||
},
|
}
|
||||||
|
if len(previousRequirements) > 0 {
|
||||||
|
metadata["previous_requirements_count"] = len(previousRequirements)
|
||||||
|
if previousHydration != nil {
|
||||||
|
metadata["previous_artifacts_hydrated"] = append([]string(nil), previousHydration.Hydrated...)
|
||||||
|
metadata["previous_artifacts_hydrated_count"] = len(previousHydration.Hydrated)
|
||||||
|
metadata["previous_artifacts_missing_optional"] = append([]string(nil), previousHydration.SkippedMissing...)
|
||||||
|
metadata["previous_artifacts_missing_optional_count"] = len(previousHydration.SkippedMissing)
|
||||||
|
if strings.TrimSpace(previousHydration.PreviousRunID) != "" {
|
||||||
|
metadata["previous_session_run_id"] = strings.TrimSpace(previousHydration.PreviousRunID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return &StageResult{
|
||||||
|
Metadata: metadata,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -354,6 +382,35 @@ func countAudioInputs(inputs []manifest.InputRecord) int {
|
|||||||
return count
|
return count
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func collectPreparePreviousRequirements(cfg *config.Config) []artifacts.PreviousArtifactRequirement {
|
||||||
|
if cfg == nil || cfg.Pipeline == nil || cfg.Pipeline.Scriptorium == nil {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
return artifacts.CollectPreviousArtifactRequirements(cfg.Pipeline.Scriptorium.Artifacts)
|
||||||
|
}
|
||||||
|
|
||||||
|
func clearManagedPreviousState(paths artifacts.SessionPaths) error {
|
||||||
|
previousDir := filepath.Clean(paths.PreviousDir)
|
||||||
|
sessionRoot := filepath.Clean(paths.Root)
|
||||||
|
if strings.TrimSpace(previousDir) == "" || strings.TrimSpace(sessionRoot) == "" {
|
||||||
|
return fmt.Errorf("previous/session root paths are required")
|
||||||
|
}
|
||||||
|
if previousDir == sessionRoot {
|
||||||
|
return fmt.Errorf("refusing to clear session root as previous cache: %q", previousDir)
|
||||||
|
}
|
||||||
|
prefix := sessionRoot + string(filepath.Separator)
|
||||||
|
if !strings.HasPrefix(previousDir, prefix) {
|
||||||
|
return fmt.Errorf("refusing to clear path outside session root: %q", previousDir)
|
||||||
|
}
|
||||||
|
if filepath.Base(previousDir) != config.PathPreviousDirSegment {
|
||||||
|
return fmt.Errorf("refusing to clear non-previous path %q", previousDir)
|
||||||
|
}
|
||||||
|
if err := os.RemoveAll(previousDir); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
return os.MkdirAll(previousDir, 0o755)
|
||||||
|
}
|
||||||
|
|
||||||
func pathsWorkDirForManifest(env *Env, m *manifest.Manifest, sessionID string) string {
|
func pathsWorkDirForManifest(env *Env, m *manifest.Manifest, sessionID string) string {
|
||||||
if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Session == nil {
|
if env == nil || env.Config == nil || env.Config.Pipeline == nil || env.Config.Session == nil {
|
||||||
return ""
|
return ""
|
||||||
|
|||||||
@@ -284,6 +284,188 @@ func TestPrepareStageAudioSourceConflictFails(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestPrepareStageWithoutPreviousRequirementsDoesNotTouchPreviousState(t *testing.T) {
|
||||||
|
env, m := setupPrepareEnv(t)
|
||||||
|
root := filepath.Dir(env.Config.SessionPath)
|
||||||
|
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
|
||||||
|
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
|
||||||
|
|
||||||
|
paths := sessionPathsForEnv(env, m.SessionID)
|
||||||
|
stalePath := filepath.Join(paths.PreviousDir, "stale.txt")
|
||||||
|
writeFile(t, stalePath, "stale")
|
||||||
|
|
||||||
|
_, err := (prepareStage{}).Run(context.Background(), env, m)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("prepare.Run() error = %v", err)
|
||||||
|
}
|
||||||
|
if _, err := os.Stat(stalePath); err != nil {
|
||||||
|
t.Fatalf("expected previous stale file to remain untouched: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPrepareStageOptionalPreviousArtifactWithoutPreviousSessionIDSucceeds(t *testing.T) {
|
||||||
|
env, m := setupPrepareEnv(t)
|
||||||
|
root := filepath.Dir(env.Config.SessionPath)
|
||||||
|
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
|
||||||
|
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
|
||||||
|
configurePreparePreviousArtifactSource(env, false)
|
||||||
|
env.Config.Session.PreviousSessionID = ""
|
||||||
|
|
||||||
|
paths := sessionPathsForEnv(env, m.SessionID)
|
||||||
|
writeFile(t, filepath.Join(paths.PreviousArtifactsDir, "stale.txt"), "stale")
|
||||||
|
|
||||||
|
result, err := (prepareStage{}).Run(context.Background(), env, m)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("prepare.Run() error = %v", err)
|
||||||
|
}
|
||||||
|
if result.Metadata["previous_requirements_count"] != 1 {
|
||||||
|
t.Fatalf("metadata previous_requirements_count = %#v, want 1", result.Metadata["previous_requirements_count"])
|
||||||
|
}
|
||||||
|
if result.Metadata["previous_artifacts_hydrated_count"] != 0 {
|
||||||
|
t.Fatalf("metadata previous_artifacts_hydrated_count = %#v, want 0", result.Metadata["previous_artifacts_hydrated_count"])
|
||||||
|
}
|
||||||
|
if result.Metadata["previous_artifacts_missing_optional_count"] != 1 {
|
||||||
|
t.Fatalf("metadata previous_artifacts_missing_optional_count = %#v, want 1", result.Metadata["previous_artifacts_missing_optional_count"])
|
||||||
|
}
|
||||||
|
if _, err := os.Stat(filepath.Join(paths.PreviousArtifactsDir, "stale.txt")); !os.IsNotExist(err) {
|
||||||
|
t.Fatalf("expected stale previous state to be cleared, stat err = %v", err)
|
||||||
|
}
|
||||||
|
for _, in := range m.Inputs {
|
||||||
|
if in.Kind == preparePreviousInputKindManifest || in.Kind == preparePreviousInputKindArtifact {
|
||||||
|
t.Fatalf("unexpected previous input record: %#v", in)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPrepareStageRequiredPreviousArtifactWithoutPreviousSessionIDFails(t *testing.T) {
|
||||||
|
env, m := setupPrepareEnv(t)
|
||||||
|
root := filepath.Dir(env.Config.SessionPath)
|
||||||
|
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
|
||||||
|
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
|
||||||
|
configurePreparePreviousArtifactSource(env, true)
|
||||||
|
env.Config.Session.PreviousSessionID = ""
|
||||||
|
|
||||||
|
_, err := (prepareStage{}).Run(context.Background(), env, m)
|
||||||
|
if err == nil || !strings.Contains(err.Error(), "previous_session_id is required") {
|
||||||
|
t.Fatalf("error = %v, want previous_session_id required failure", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPrepareStageHydratesRequiredPreviousArtifactAndRecordsInputs(t *testing.T) {
|
||||||
|
env, m := setupPrepareEnv(t)
|
||||||
|
root := filepath.Dir(env.Config.SessionPath)
|
||||||
|
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
|
||||||
|
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
|
||||||
|
env.Config.Session.Campaign = "forsaken"
|
||||||
|
env.Config.Session.PreviousSessionID = "2026-05-10"
|
||||||
|
env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{Bucket: "my-dnd-archive", RootPrefix: "dnd"}
|
||||||
|
configurePreparePreviousArtifactSource(env, true)
|
||||||
|
|
||||||
|
fake := &storage.FakeBackend{}
|
||||||
|
env.ObjectStore = fake
|
||||||
|
seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
|
||||||
|
includeRunPointerObject: true,
|
||||||
|
includeArtifactObject: true,
|
||||||
|
artifactBody: "# prior recap\n",
|
||||||
|
})
|
||||||
|
|
||||||
|
paths := sessionPathsForEnv(env, m.SessionID)
|
||||||
|
result, err := (prepareStage{}).Run(context.Background(), env, m)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("prepare.Run() error = %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if _, err := os.Stat(paths.PreviousManifestPath); err != nil {
|
||||||
|
t.Fatalf("expected previous manifest: %v", err)
|
||||||
|
}
|
||||||
|
recapPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
|
||||||
|
if _, err := os.Stat(recapPath); err != nil {
|
||||||
|
t.Fatalf("expected previous artifact: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var hasPreviousManifest, hasPreviousArtifact bool
|
||||||
|
for _, in := range m.Inputs {
|
||||||
|
if in.Kind == preparePreviousInputKindManifest {
|
||||||
|
hasPreviousManifest = true
|
||||||
|
}
|
||||||
|
if in.Kind == preparePreviousInputKindArtifact {
|
||||||
|
hasPreviousArtifact = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !hasPreviousManifest || !hasPreviousArtifact {
|
||||||
|
t.Fatalf("manifest inputs missing previous provenance: %#v", m.Inputs)
|
||||||
|
}
|
||||||
|
if result.Metadata["previous_artifacts_hydrated_count"] != 1 {
|
||||||
|
t.Fatalf("metadata previous_artifacts_hydrated_count = %#v, want 1", result.Metadata["previous_artifacts_hydrated_count"])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestPrepareStageRerunOverwritesPreviousCache(t *testing.T) {
|
||||||
|
env, m := setupPrepareEnv(t)
|
||||||
|
root := filepath.Dir(env.Config.SessionPath)
|
||||||
|
writeFile(t, filepath.Join(root, "audio", "a.flac"), "a")
|
||||||
|
env.Config.Session.Inputs.AudioFiles = []string{"./audio/a.flac"}
|
||||||
|
env.Config.Session.Campaign = "forsaken"
|
||||||
|
env.Config.Session.PreviousSessionID = "2026-05-10"
|
||||||
|
env.Config.Pipeline.Storage.S3 = &config.StorageS3Config{Bucket: "my-dnd-archive", RootPrefix: "dnd"}
|
||||||
|
configurePreparePreviousArtifactSource(env, true)
|
||||||
|
|
||||||
|
fake := &storage.FakeBackend{}
|
||||||
|
env.ObjectStore = fake
|
||||||
|
seed := seedPreviousCurrentState(t, env, fake, previousStateSeedOptions{
|
||||||
|
includeRunPointerObject: true,
|
||||||
|
includeArtifactObject: true,
|
||||||
|
artifactBody: "# old recap\n",
|
||||||
|
})
|
||||||
|
|
||||||
|
paths := sessionPathsForEnv(env, m.SessionID)
|
||||||
|
if _, err := (prepareStage{}).Run(context.Background(), env, m); err != nil {
|
||||||
|
t.Fatalf("first prepare run error = %v", err)
|
||||||
|
}
|
||||||
|
recapPath := artifacts.SessionPreviousArtifactPath(paths, "artifacts/session_recap.md")
|
||||||
|
firstBytes, err := os.ReadFile(recapPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("read first hydrated artifact: %v", err)
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(string(firstBytes)) != "# old recap" {
|
||||||
|
t.Fatalf("first hydrated content = %q, want %q", strings.TrimSpace(string(firstBytes)), "# old recap")
|
||||||
|
}
|
||||||
|
|
||||||
|
writeFile(t, filepath.Join(paths.PreviousArtifactsDir, "stale.txt"), "stale")
|
||||||
|
fake.SeedObject(storage.FakeObject{Key: seed.ArtifactKey, Data: []byte("# new recap\n")})
|
||||||
|
|
||||||
|
if _, err := (prepareStage{}).Run(context.Background(), env, m); err != nil {
|
||||||
|
t.Fatalf("second prepare run error = %v", err)
|
||||||
|
}
|
||||||
|
secondBytes, err := os.ReadFile(recapPath)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("read second hydrated artifact: %v", err)
|
||||||
|
}
|
||||||
|
if strings.TrimSpace(string(secondBytes)) != "# new recap" {
|
||||||
|
t.Fatalf("second hydrated content = %q, want %q", strings.TrimSpace(string(secondBytes)), "# new recap")
|
||||||
|
}
|
||||||
|
if _, err := os.Stat(filepath.Join(paths.PreviousArtifactsDir, "stale.txt")); !os.IsNotExist(err) {
|
||||||
|
t.Fatalf("expected stale previous cache file to be removed, stat err = %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func configurePreparePreviousArtifactSource(env *Env, required bool) {
|
||||||
|
env.Config.Pipeline.Scriptorium = &config.ScriptoriumConfig{
|
||||||
|
Artifacts: map[string]config.ScriptoriumArtifactConfig{
|
||||||
|
"session_recap": {
|
||||||
|
Enabled: true,
|
||||||
|
OutputPath: "artifacts/session_recap.md",
|
||||||
|
Inputs: map[string]config.ScriptoriumInputConfig{
|
||||||
|
"previous_recap": {
|
||||||
|
Source: "narratio.previous_session.artifact.session_recap",
|
||||||
|
Required: required,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func setupPrepareEnv(t *testing.T) (*Env, *manifest.Manifest) {
|
func setupPrepareEnv(t *testing.T) (*Env, *manifest.Manifest) {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
workspace := t.TempDir()
|
workspace := t.TempDir()
|
||||||
|
|||||||
Reference in New Issue
Block a user