fix(deletion): move retention feature delete to end of job
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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] "
|
||||
|
||||
@@ -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]:
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user