Harden pipeline state and plan release upgrades

This commit is contained in:
2026-08-30 18:51:20 +00:00
parent 3da97ca50c
commit c812fe3655
18 changed files with 1381 additions and 1237 deletions

View File

@@ -437,29 +437,9 @@ func materializeS3AudioInputs(ctx context.Context, env *Env, m *manifest.Manifes
return s3AudioMaterializationStats{}, fmt.Errorf("run id is required for s3 audio input")
}
sessionPrefix := artifacts.S3SessionPrefix(env.Config.Pipeline.Storage.S3.RootPrefix, campaign, sessionID)
audioPrefix := artifacts.S3AudioPrefix(sessionPrefix, env.Config.Session.Inputs.AudioS3.Prefix)
objects, err := env.ObjectStore.List(ctx, audioPrefix)
audioObjects, err := listS3AudioObjects(ctx, env, campaign, sessionID)
if err != nil {
return s3AudioMaterializationStats{}, fmt.Errorf("list s3 audio objects under %q: %w", audioPrefix, err)
}
audioObjects := make([]storage.ObjectInfo, 0, len(objects))
for _, obj := range objects {
key := strings.TrimSpace(obj.Key)
if key == "" || strings.HasSuffix(key, "/") {
continue
}
if !isFlac(key) {
continue
}
audioObjects = append(audioObjects, obj)
}
sort.Slice(audioObjects, func(i, j int) bool {
return audioObjects[i].Key < audioObjects[j].Key
})
if len(audioObjects) == 0 {
return s3AudioMaterializationStats{}, fmt.Errorf("no .flac files found under s3 audio prefix %q", audioPrefix)
return s3AudioMaterializationStats{}, err
}
spoolAudioDir := strings.TrimSpace(m.LocalSpoolDir)
@@ -475,16 +455,10 @@ func materializeS3AudioInputs(ctx context.Context, env *Env, m *manifest.Manifes
return s3AudioMaterializationStats{}, fmt.Errorf("create work audio directory %q: %w", workAudioDir, err)
}
seenBase := map[string]string{}
stats := s3AudioMaterializationStats{}
cacheEnabled := env.Config.Pipeline.Cache.S3Audio == nil || *env.Config.Pipeline.Cache.S3Audio
for _, obj := range audioObjects {
base := path.Base(obj.Key)
if prev, exists := seenBase[base]; exists && prev != obj.Key {
return s3AudioMaterializationStats{}, fmt.Errorf("duplicate s3 audio basename %q from %q and %q", base, prev, obj.Key)
}
seenBase[base] = obj.Key
spoolPath := filepath.Join(spoolAudioDir, base)
workPath := filepath.Join(workAudioDir, base)
result, err := audio.MaterializeS3Audio(ctx, audio.S3MaterializeRequest{
@@ -525,6 +499,42 @@ func materializeS3AudioInputs(ctx context.Context, env *Env, m *manifest.Manifes
return stats, nil
}
func listS3AudioObjects(ctx context.Context, env *Env, campaign, sessionID string) ([]storage.ObjectInfo, error) {
if env == nil || env.ObjectStore == nil || env.Config == nil || env.Config.Pipeline == nil ||
env.Config.Session == nil || env.Config.Pipeline.Storage.S3 == nil || env.Config.Session.Inputs.AudioS3 == nil {
return nil, fmt.Errorf("s3 audio input requires object store and resolved storage configuration")
}
sessionPrefix := artifacts.S3SessionPrefix(env.Config.Pipeline.Storage.S3.RootPrefix, campaign, sessionID)
audioPrefix := artifacts.S3AudioPrefix(sessionPrefix, env.Config.Session.Inputs.AudioS3.Prefix)
objects, err := env.ObjectStore.List(ctx, audioPrefix)
if err != nil {
return nil, fmt.Errorf("list s3 audio objects under %q: %w", audioPrefix, err)
}
audioObjects := make([]storage.ObjectInfo, 0, len(objects))
seenBase := map[string]string{}
for _, obj := range objects {
key := strings.TrimSpace(obj.Key)
if key == "" || strings.HasSuffix(key, "/") || !isFlac(key) {
continue
}
base := path.Base(key)
if previous, exists := seenBase[base]; exists && previous != key {
return nil, fmt.Errorf("duplicate s3 audio basename %q from %q and %q", base, previous, key)
}
seenBase[base] = key
obj.Key = key
audioObjects = append(audioObjects, obj)
}
sort.Slice(audioObjects, func(i, j int) bool {
return audioObjects[i].Key < audioObjects[j].Key
})
if len(audioObjects) == 0 {
return nil, fmt.Errorf("no .flac files found under s3 audio prefix %q", audioPrefix)
}
return audioObjects, nil
}
func countAudioInputs(inputs []manifest.InputRecord) int {
count := 0
for _, in := range inputs {