Persist retryable post-publish cleanup obligations

This commit is contained in:
2026-08-10 20:29:00 +00:00
parent 0cf2cbfeb3
commit eac7e155a5
12 changed files with 438 additions and 121 deletions

View File

@@ -19,14 +19,32 @@ import (
type publishSuccessStage struct {
metadata map[string]any
targets *cleanupSeed
}
func (publishSuccessStage) Name() string { return "publish" }
func (publishSuccessStage) Declares() stage.IODecl { return stage.IODecl{} }
func (s publishSuccessStage) Run(_ context.Context, _ *stage.Env, _ *manifest.Manifest) (*stage.StageResult, error) {
func (s publishSuccessStage) Run(_ context.Context, _ *stage.Env, m *manifest.Manifest) (*stage.StageResult, error) {
if s.targets != nil {
s.targets.runWorkDir = m.LocalWorkDir
s.targets.spoolAudioDir = m.LocalSpoolDir
if err := os.MkdirAll(filepath.Join(m.LocalWorkDir, "logs"), 0o755); err != nil {
return nil, err
}
if err := os.WriteFile(filepath.Join(m.LocalWorkDir, "logs", "stage.log"), []byte("log\n"), 0o644); err != nil {
return nil, err
}
if err := os.MkdirAll(m.LocalSpoolDir, 0o755); err != nil {
return nil, err
}
if err := os.WriteFile(filepath.Join(m.LocalSpoolDir, "speaker.flac"), []byte("flac\n"), 0o644); err != nil {
return nil, err
}
}
md := map[string]any{
"stage": "publish",
"uploaded": true,
"published_run_id": m.RunID,
"remote_commit_key": "dnd/campaigns/sample-campaign/sessions/2026-05-03/runs/20260519T010203Z-a1b2c3d4/commit.json",
"current_commit_pointer_key": "dnd/campaigns/sample-campaign/sessions/2026-05-03/current/commit-pointer.json",
}
@@ -49,7 +67,7 @@ func TestPostPublishCleanupDisabledKeepsLocalDirs(t *testing.T) {
cfg.Pipeline.Spool.DeleteAudioAfterPublish = false
cfg.Pipeline.Workspace.CleanupAfterPublish = false
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -63,7 +81,7 @@ func TestPostPublishCleanupSpoolOnly(t *testing.T) {
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = false
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -77,7 +95,7 @@ func TestPostPublishCleanupWorkdirOnly(t *testing.T) {
cfg.Pipeline.Spool.DeleteAudioAfterPublish = false
cfg.Pipeline.Workspace.CleanupAfterPublish = true
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -93,7 +111,7 @@ func TestPostPublishCleanupBothPolicies(t *testing.T) {
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = true
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -103,6 +121,115 @@ func TestPostPublishCleanupBothPolicies(t *testing.T) {
assertExists(t, seed.previousCachePath)
}
func TestPostPublishCleanupRetriesWhenInitialObligationSaveFails(t *testing.T) {
cfg, seed := cleanupFixtureConfig(t)
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = false
failed := false
store := &cleanupFailingManifestStore{
delegate: &manifest.LocalStore{},
fail: func(m *manifest.Manifest) error {
if !failed && m.PostPublishCleanup != nil {
failed = true
return errors.New("injected obligation save failure")
}
return nil
},
}
_, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{
Env: &Env{ManifestStore: store, ObjectStore: &storage.FakeBackend{}},
})
if err == nil || !strings.Contains(err.Error(), "post-publish cleanup incomplete") {
t.Fatalf("executeStages() error = %v, want incomplete cleanup", err)
}
assertExists(t, seed.spoolAudioDir)
assertCleanupPending(t, cfg)
store.fail = nil
if _, err := executeStages(context.Background(), cfg, nil, RunOptions{Env: &Env{ManifestStore: store, ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("retry executeStages() error = %v", err)
}
assertMissing(t, seed.spoolAudioDir)
assertCleanupComplete(t, cfg)
}
func TestPostPublishCleanupRetriesFailedDeletionWithoutTouchingOtherRuns(t *testing.T) {
cfg, seed := cleanupFixtureConfig(t)
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = true
originalRemove := removeRunScopedDirFn
removeRunScopedDirFn = func(root, target, policy string) error {
if policy == "pipeline.workspace.cleanup_after_publish" {
return errors.New("injected deletion failure")
}
return originalRemove(root, target, policy)
}
t.Cleanup(func() { removeRunScopedDirFn = originalRemove })
_, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}})
if err == nil || !strings.Contains(err.Error(), "post-publish cleanup incomplete") {
t.Fatalf("executeStages() error = %v, want incomplete cleanup", err)
}
assertMissing(t, seed.spoolAudioDir)
assertExists(t, seed.runWorkDir)
assertCleanupPending(t, cfg)
assertExists(t, seed.otherRunDir)
assertExists(t, seed.previousCachePath)
removeRunScopedDirFn = originalRemove
if _, err := executeStages(context.Background(), cfg, nil, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("retry executeStages() error = %v", err)
}
assertMissing(t, seed.runWorkDir)
assertExists(t, seed.otherRunDir)
assertExists(t, seed.previousCachePath)
assertCleanupComplete(t, cfg)
}
func TestPostPublishCleanupRetriesWhenCompletionEvidenceSaveFails(t *testing.T) {
cfg, seed := cleanupFixtureConfig(t)
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = false
failed := false
store := &cleanupFailingManifestStore{
delegate: &manifest.LocalStore{},
fail: func(m *manifest.Manifest) error {
if m.PostPublishCleanup == nil {
return nil
}
for _, target := range m.PostPublishCleanup.Targets {
if !failed && target.Completed {
failed = true
return errors.New("injected completion evidence failure")
}
}
return nil
},
}
_, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{
Env: &Env{ManifestStore: store, ObjectStore: &storage.FakeBackend{}},
})
if err == nil || !strings.Contains(err.Error(), "post-publish cleanup incomplete") {
t.Fatalf("executeStages() error = %v, want incomplete cleanup", err)
}
assertMissing(t, seed.spoolAudioDir)
assertCleanupPending(t, cfg)
store.fail = nil
if _, err := executeStages(context.Background(), cfg, nil, RunOptions{Env: &Env{ManifestStore: store, ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("retry executeStages() error = %v", err)
}
assertCleanupComplete(t, cfg)
if _, err := executeStages(context.Background(), cfg, nil, RunOptions{Env: &Env{ManifestStore: store, ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("idempotent retry executeStages() error = %v", err)
}
assertMissing(t, seed.spoolAudioDir)
}
func TestPostPublishCleanupNotRunWhenPublishFails(t *testing.T) {
cfg, seed := cleanupFixtureConfig(t)
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
@@ -122,7 +249,7 @@ func TestPostPublishCleanupNotRunWhenPublishSkipped(t *testing.T) {
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = true
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{metadata: map[string]any{"skipped": true}}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{metadata: map[string]any{"skipped": true}, targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -135,7 +262,7 @@ func TestPostPublishCleanupNotRunWhenCommitPointerMissing(t *testing.T) {
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = true
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{metadata: map[string]any{"current_commit_pointer_key": ""}}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{metadata: map[string]any{"current_commit_pointer_key": ""}, targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -149,7 +276,7 @@ func TestPostPublishCleanupNotRunWhenPublishUploadDisabled(t *testing.T) {
cfg.Pipeline.Workspace.CleanupAfterPublish = true
cfg.Pipeline.Publish.UploadRun = boolPtr(false)
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
if _, err := executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{targets: &seed}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}}); err != nil {
t.Fatalf("executeStages() error = %v", err)
}
@@ -182,12 +309,28 @@ func TestPostPublishCleanupFailsOnUnsafePath(t *testing.T) {
if err != nil {
t.Fatalf("Load() error = %v", err)
}
m.LocalSpoolDir = filepath.Join(filepath.Dir(cfg.Pipeline.Spool.Root), "outside-spool")
m.MarkStageSucceeded("publish", time.Now().UTC(), nil)
m.Stages["publish"].Metadata = map[string]any{
"uploaded": true,
"published_run_id": m.RunID,
"remote_commit_key": "dnd/campaigns/sample-campaign/sessions/2026-05-03/runs/20260516T010203Z-1a2b3c4d/commit.json",
"current_commit_pointer_key": "dnd/campaigns/sample-campaign/sessions/2026-05-03/current/commit-pointer.json",
}
m.PostPublishCleanup = &manifest.PostPublishCleanup{
CommittedRunID: m.RunID,
RemoteCommitKey: m.Stages["publish"].Metadata["remote_commit_key"].(string),
CurrentCommitPointerKey: m.Stages["publish"].Metadata["current_commit_pointer_key"].(string),
Targets: []manifest.CleanupTarget{{
Policy: "pipeline.spool.delete_audio_after_publish",
Root: cfg.Pipeline.Spool.Root,
Path: filepath.Join(filepath.Dir(cfg.Pipeline.Spool.Root), "outside-spool"),
}},
}
if err := store.Save(context.Background(), manifestPath, m); err != nil {
t.Fatalf("Save() error = %v", err)
}
_, err = executeStages(context.Background(), cfg, []stage.Stage{publishSuccessStage{}}, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}})
_, err = executeStages(context.Background(), cfg, nil, RunOptions{Env: &Env{ObjectStore: &storage.FakeBackend{}}})
if err == nil || !strings.Contains(err.Error(), "refusing to delete path outside root") {
t.Fatalf("executeStages() error = %v, want safe-path failure", err)
}
@@ -220,14 +363,15 @@ func TestPostPublishCleanupNotRunWhenCommittedManifestUploadFails(t *testing.T)
cfg, seed, _ := publishStageCleanupFixture(t)
cfg.Pipeline.Spool.DeleteAudioAfterPublish = true
cfg.Pipeline.Workspace.CleanupAfterPublish = true
failKey := artifacts.S3RunSessionManifestKey(seed.sessionPrefix, seed.runID)
publishStageImpl, err := stage.Select("publish")
if err != nil {
t.Fatalf("Select(publish) error = %v", err)
}
_, err = executeStages(context.Background(), cfg, []stage.Stage{publishStageImpl}, RunOptions{
Env: &Env{ObjectStore: &failKeyStore{delegate: &storage.FakeBackend{}, failKey: failKey}},
Env: &Env{ObjectStore: &failKeyStore{delegate: &storage.FakeBackend{}, fail: func(key string) bool {
return strings.HasSuffix(key, "/session-manifest.json")
}}},
})
if err == nil || !strings.Contains(err.Error(), "immutable object") {
t.Fatalf("executeStages() error = %v, want committed-manifest failure", err)
@@ -268,6 +412,28 @@ type cleanupSeed struct {
sessionPrefix string
}
type cleanupFailingManifestStore struct {
delegate manifest.Store
fail func(*manifest.Manifest) error
}
func (s *cleanupFailingManifestStore) Create(ctx context.Context, sessionID string) (*manifest.Manifest, error) {
return s.delegate.Create(ctx, sessionID)
}
func (s *cleanupFailingManifestStore) Load(ctx context.Context, path string) (*manifest.Manifest, error) {
return s.delegate.Load(ctx, path)
}
func (s *cleanupFailingManifestStore) Save(ctx context.Context, path string, m *manifest.Manifest) error {
if s.fail != nil {
if err := s.fail(m); err != nil {
return err
}
}
return s.delegate.Save(ctx, path, m)
}
func cleanupFixtureConfig(t *testing.T) (*config.Config, cleanupSeed) {
t.Helper()
@@ -388,6 +554,11 @@ func writePublishFixtureRunFiles(t *testing.T, runWorkDir, sessionRoot string) {
type failKeyStore struct {
delegate *storage.FakeBackend
failKey string
fail func(string) bool
}
func (s *failKeyStore) fails(key string) bool {
return strings.TrimSpace(key) == strings.TrimSpace(s.failKey) || (s.fail != nil && s.fail(key))
}
func (s *failKeyStore) List(ctx context.Context, prefix string) ([]storage.ObjectInfo, error) {
@@ -403,21 +574,21 @@ func (s *failKeyStore) Download(ctx context.Context, key, localPath string) erro
}
func (s *failKeyStore) Upload(ctx context.Context, localPath, key string, opts storage.UploadOptions) (storage.ObjectInfo, error) {
if strings.TrimSpace(key) == strings.TrimSpace(s.failKey) {
if s.fails(key) {
return storage.ObjectInfo{}, errors.New("forced upload failure")
}
return s.delegate.Upload(ctx, localPath, key, opts)
}
func (s *failKeyStore) UploadReader(ctx context.Context, source io.Reader, key string, opts storage.UploadOptions) (storage.ObjectInfo, error) {
if strings.TrimSpace(key) == strings.TrimSpace(s.failKey) {
if s.fails(key) {
return storage.ObjectInfo{}, errors.New("forced upload failure")
}
return s.delegate.UploadReader(ctx, source, key, opts)
}
func (s *failKeyStore) UploadConditional(ctx context.Context, source io.Reader, key string, opts storage.UploadOptions, condition storage.WriteCondition) (storage.ObjectInfo, error) {
if strings.TrimSpace(key) == strings.TrimSpace(s.failKey) {
if s.fails(key) {
return storage.ObjectInfo{}, errors.New("forced upload failure")
}
return s.delegate.UploadConditional(ctx, source, key, opts, condition)
@@ -440,3 +611,36 @@ func assertMissing(t *testing.T, path string) {
t.Fatalf("expected path to be removed %q, stat err=%v", path, err)
}
}
func assertCleanupPending(t *testing.T, cfg *config.Config) {
t.Helper()
m, err := (&manifest.LocalStore{}).Load(context.Background(), manifestPathFor(cfg))
if err != nil {
t.Fatalf("Load() error = %v", err)
}
if m.PostPublishCleanup == nil {
t.Fatal("expected a persisted cleanup obligation")
}
for _, target := range m.PostPublishCleanup.Targets {
if !target.Completed {
return
}
}
t.Fatal("expected at least one cleanup target to remain incomplete")
}
func assertCleanupComplete(t *testing.T, cfg *config.Config) {
t.Helper()
m, err := (&manifest.LocalStore{}).Load(context.Background(), manifestPathFor(cfg))
if err != nil {
t.Fatalf("Load() error = %v", err)
}
if m.PostPublishCleanup == nil {
t.Fatal("expected a persisted cleanup obligation")
}
for _, target := range m.PostPublishCleanup.Targets {
if !target.Completed {
t.Fatalf("cleanup target remains incomplete: %#v", target)
}
}
}