From 964bc84288cc6a7f9596e4d457ed04cb74e1535b Mon Sep 17 00:00:00 2001 From: Jason Ross Date: Sat, 19 Sep 2026 16:10:24 -0500 Subject: [PATCH] fix(deletion): move retention feature delete to end of job --- README.md | 13 +++--- src/audio_scribe/cli.py | 4 +- src/audio_scribe/config.py | 6 +-- src/audio_scribe/jobs/batch.py | 4 +- src/audio_scribe/jobs/runner.py | 21 +++++++-- tests/test_jobs_runner.py | 81 ++++++++++++++++++++++++++++++++- 6 files changed, 109 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index 04683fd..703b95a 100644 --- a/README.md +++ b/README.md @@ -86,14 +86,15 @@ Video and audio are kept by default. | Flag | Effect | |---|---| -| `--no-retain-video` | delete the video once the FLAC is extracted and verified | -| `--no-retain-audio` | delete the FLAC once every output is written | +| `--no-retain-video` | delete the video once the job is complete | +| `--no-retain-audio` | delete the FLAC once the job is complete | | `--no-retain` | both | -Deletion always happens **after** the next stage is committed, so interrupting a -run is still safe and resume still works. What you actually give up is re-run -flexibility: re-transcribing later with a different `--language`, `--model` or -`--task` will need another download. +These are cleanup flags: deletion is the very last step of a job, after every +stage has finished and the outputs are published. A job that fails or is +interrupted keeps its media, so resume still works. What you actually give up is +re-run flexibility: re-transcribing later with a different `--language`, `--model` +or `--task` will need another download. These flags prompt once per run. On a non-TTY they refuse outright rather than proceeding — a piped or cron invocation should never delete gigabytes with nobody diff --git a/src/audio_scribe/cli.py b/src/audio_scribe/cli.py index 3282f66..226af32 100644 --- a/src/audio_scribe/cli.py +++ b/src/audio_scribe/cli.py @@ -153,8 +153,8 @@ def main() -> None: @click.option("--js-runtime", default=None, help="deno (default), node, bun or quickjs.") @click.option("--ov-cache-dir", default=None, type=click.Path(path_type=Path)) @click.option("--no-retain", is_flag=True, help="Delete both video and audio when done.") -@click.option("--no-retain-video", is_flag=True, help="Delete the video after extraction.") -@click.option("--no-retain-audio", is_flag=True, help="Delete the FLAC after transcription.") +@click.option("--no-retain-video", is_flag=True, help="Delete the video when the job is done.") +@click.option("--no-retain-audio", is_flag=True, help="Delete the FLAC when the job is done.") @click.option("-y", "--yes", is_flag=True, help="Confirm destructive flags non-interactively.") @click.option("--force", is_flag=True, help="Re-run every stage.") @click.option( diff --git a/src/audio_scribe/config.py b/src/audio_scribe/config.py index 77ea509..c3e1c1f 100644 --- a/src/audio_scribe/config.py +++ b/src/audio_scribe/config.py @@ -43,10 +43,10 @@ class RunConfig: def retention_summary(self) -> str: parts: list[str] = [] if self.no_retain_video: - parts.append("video after the FLAC is extracted and verified") + parts.append("the video") if self.no_retain_audio: - parts.append("audio after all transcript outputs are written") - return "; ".join(parts) + parts.append("the audio (FLAC)") + return " and ".join(parts) def request(self) -> TranscribeRequest: from audio_scribe.backends.base import ( # pylint: disable=import-outside-toplevel diff --git a/src/audio_scribe/jobs/batch.py b/src/audio_scribe/jobs/batch.py index 4600a11..a83a17d 100644 --- a/src/audio_scribe/jobs/batch.py +++ b/src/audio_scribe/jobs/batch.py @@ -82,8 +82,8 @@ def confirm_retention( out.write( "\nWARNING source media will be deleted after each job completes:\n" f" {config.retention_summary()}\n\n" - " Interrupting a run is safe: media is removed only once the next stage\n" - " is committed, so resume always works.\n\n" + " Interrupting a run is safe: media is removed only after every stage of\n" + " a job has finished, so resume always works.\n\n" " What you lose is re-run flexibility. Re-transcribing later with a\n" " different --language / --model / --task will need another download.\n\n" "Continue? [y/N] " diff --git a/src/audio_scribe/jobs/runner.py b/src/audio_scribe/jobs/runner.py index e093265..1cce319 100644 --- a/src/audio_scribe/jobs/runner.py +++ b/src/audio_scribe/jobs/runner.py @@ -146,6 +146,22 @@ def _execute( # need_outputs is the root condition in plan_stages, so every non-empty plan # ends in OUTPUTS and this always has work to do. _do_transcript(job, config, state, stages=stages, holder=holder) + # Cleanup is last: a failure in any stage above skips it, so the media is + # still there for the resume. + _apply_retention(job, config, state, video) + + +def _apply_retention(job: JobPaths, config: RunConfig, state: JobState, video: Path | None) -> None: + """Delete what the --no-retain-* flags name, now that the job is complete. + + Only a file that is actually on disk is deleted and recorded: stamping + DELETED_BY_POLICY over an artifact that was already gone would make the event + log claim a deletion this run never made. + """ + if config.no_retain_video and video is not None: + _record_deletion(job, state, "video", "--no-retain-video", video) + if config.no_retain_audio and job.audio_file.exists(): + _record_deletion(job, state, "audio", "--no-retain-audio", job.audio_file) def _do_download(job: JobPaths, url: str, config: RunConfig, state: JobState) -> Path: @@ -182,9 +198,6 @@ def _do_audio(job: JobPaths, config: RunConfig, state: JobState, video: Path | N duration_s=ffmpeg.probe_duration(flac), ) store.save(job, state) - # Only now, with the successor verified and committed, is deletion safe. - if config.no_retain_video: - _record_deletion(job, state, "video", "--no-retain-video", video) def _do_transcript( @@ -233,8 +246,6 @@ def _do_transcript( state.outputs = Outputs(status=ArtifactStatus.PRESENT, paths=written) store.save(job, state) outputs_stage.publish(job, written, config.formats, out_dir=config.out_dir, title=state.title) - if config.no_retain_audio: - _record_deletion(job, state, "audio", "--no-retain-audio", job.audio_file) def _rerender(job: JobPaths, config: RunConfig, state: JobState) -> dict[str, str]: diff --git a/tests/test_jobs_runner.py b/tests/test_jobs_runner.py index cdc391a..28015f1 100644 --- a/tests/test_jobs_runner.py +++ b/tests/test_jobs_runner.py @@ -2,7 +2,7 @@ from __future__ import annotations import json from dataclasses import replace -from typing import TYPE_CHECKING, ClassVar +from typing import TYPE_CHECKING, Any, ClassVar import pytest @@ -15,6 +15,8 @@ from audio_scribe.jobs.plan import Stage, Verdicts from audio_scribe.jobs.state import Artifact, ArtifactStatus, JobState from audio_scribe.jobs.verify import Verdict from audio_scribe.stages import download as download_stage +from audio_scribe.stages import outputs as outputs_stage +from audio_scribe.stages import transcribe as transcribe_stage from audio_scribe.transcript import Segment, TranscriptResult if TYPE_CHECKING: @@ -161,13 +163,88 @@ class TestResume: class TestRetention: - def test_video_is_deleted_only_after_the_flac_is_committed( + def test_no_retain_video_removes_the_video_and_keeps_the_flac( self, seeded: paths.JobPaths, config: RunConfig ) -> None: runner.run_job(seeded, "https://y.test/x", replace(config, no_retain_video=True)) assert download_stage.find_video(seeded) is None assert seeded.audio_file.exists() + def test_nothing_is_deleted_until_the_outputs_are_published( + self, seeded: paths.JobPaths, config: RunConfig, monkeypatch: pytest.MonkeyPatch + ) -> None: + # Publishing is the last real stage; both flags must wait for it. + on_disk_at_publish: list[tuple[bool, bool]] = [] + real_publish = outputs_stage.publish + + def spy(*args: Any, **kwargs: Any) -> list[Path]: + on_disk_at_publish.append( + (download_stage.find_video(seeded) is not None, seeded.audio_file.exists()) + ) + return real_publish(*args, **kwargs) + + monkeypatch.setattr(outputs_stage, "publish", spy) + both = replace(config, no_retain_video=True, no_retain_audio=True) + runner.run_job(seeded, "https://y.test/x", both) + assert on_disk_at_publish == [(True, True)] + assert download_stage.find_video(seeded) is None + assert not seeded.audio_file.exists() + + def test_a_failed_transcription_keeps_the_media( + self, seeded: paths.JobPaths, config: RunConfig, monkeypatch: pytest.MonkeyPatch + ) -> None: + def boom(*_a: object, **_k: object) -> None: + raise errors.TranscribeError("the backend fell over") + + monkeypatch.setattr(transcribe_stage, "run", boom) + both = replace(config, no_retain_video=True, no_retain_audio=True) + outcome = runner.run_job(seeded, "https://y.test/x", both) + assert outcome.ok is False + assert download_stage.find_video(seeded) is not None + assert seeded.audio_file.exists() + state = store.load(seeded) + assert state is not None + assert state.video.status is ArtifactStatus.PRESENT + assert state.audio.status is ArtifactStatus.PRESENT + assert store.read_events(seeded) == [] + + def test_the_resume_after_a_failure_does_the_cleanup( + self, seeded: paths.JobPaths, config: RunConfig, monkeypatch: pytest.MonkeyPatch + ) -> None: + real_run = transcribe_stage.run + + def boom(*_a: object, **_k: object) -> None: + raise errors.TranscribeError("the backend fell over") + + both = replace(config, no_retain_video=True, no_retain_audio=True) + monkeypatch.setattr(transcribe_stage, "run", boom) + runner.run_job(seeded, "https://y.test/x", both) + + monkeypatch.setattr(transcribe_stage, "run", real_run) + outcome = runner.run_job(seeded, "https://y.test/x", both) + assert outcome.stages == (Stage.TRANSCRIBE, Stage.OUTPUTS) + assert download_stage.find_video(seeded) is None + assert not seeded.audio_file.exists() + + def test_a_dry_run_deletes_nothing(self, seeded: paths.JobPaths, config: RunConfig) -> None: + runner.run_job(seeded, "https://y.test/x", config) + seeded.transcript_file("srt").unlink() + both = replace(config, no_retain_video=True, no_retain_audio=True) + runner.run_job(seeded, "https://y.test/x", both, dry_run=True) + assert download_stage.find_video(seeded) is not None + assert seeded.audio_file.exists() + + def test_re_rendering_does_not_record_a_second_deletion( + self, seeded: paths.JobPaths, config: RunConfig + ) -> None: + both = replace(config, no_retain_video=True, no_retain_audio=True) + runner.run_job(seeded, "https://y.test/x", both) + seeded.transcript_file("srt").unlink() + outcome = runner.run_job(seeded, "https://y.test/x", both) + assert outcome.stages == (Stage.OUTPUTS,) + deleted = [e["artifact"] for e in store.read_events(seeded) if e["event"] == "deleted"] + assert sorted(deleted) == ["audio", "video"] + def test_the_deletion_is_recorded_as_policy_not_loss( self, seeded: paths.JobPaths, config: RunConfig ) -> None: