fix(verify): catch a truncated container, and checkpoint aria2 often enough to resume
Three defects found by truncating a real download and killing one mid-flight. Matroska and WebM write the duration into the header, so ffprobe reports the full 212.8s of a video.mkv cut off after a kilobyte -- exit 0, plausible answer. verify_video therefore returned OK, the corrupt file was never re-downloaded, ffmpeg extracted the 0.02s of audio it could find, and the run reported "1 ok" with an empty transcript. The size recorded at download time is the only evidence the bytes are still there, so it is now checked whenever it is known rather than only as a fallback when the duration is unreadable. The audio stage now takes the expected duration and rejects an extraction that does not match it. ffmpeg exits 0 on a truncated container, so without this a damaged source yields a confident transcript of near-silence, which is a worse outcome than a failed job. It compares against the video's own probed duration rather than state.duration_s, which can come from playlist metadata. aria2 saves its control file every 60s by default. Since that file is what a resume reads, a kill -9 inside the first minute preserved a control file recording zero completed pieces: measured 0/13 on the sample, so the "resume" re-downloaded the lot while reporting a partial. At --auto-save-interval=20 the same kill preserves 2/13 pieces and the resumed download is byte-identical to a clean one. The duration tolerance moves to media/ffmpeg.py, which both callers already import, instead of being restated per call site. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -170,7 +170,11 @@ def _do_audio(job: JobPaths, config: RunConfig, state: JobState, video: Path | N
|
||||
hint="Re-run without --force-stage audio, or let the download stage run.",
|
||||
)
|
||||
state.attempts["audio"] = state.attempts.get("audio", 0) + 1
|
||||
flac = audio_stage.extract(job, video, profile=config.audio_profile)
|
||||
# The video's own probed duration, not state.duration_s: the latter can come
|
||||
# from playlist metadata, and a wrong expectation would reject a good extract.
|
||||
flac = audio_stage.extract(
|
||||
job, video, profile=config.audio_profile, expected_s=state.video.duration_s
|
||||
)
|
||||
state.audio = Artifact(
|
||||
status=ArtifactStatus.PRESENT,
|
||||
path=str(flac.relative_to(job.root)),
|
||||
|
||||
@@ -23,9 +23,8 @@ if TYPE_CHECKING:
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# A FLAC this far off the video's duration means extraction was cut short.
|
||||
_DURATION_TOLERANCE = 0.005
|
||||
_MIN_TOLERANCE_S = 1.0
|
||||
# Slack for filesystem reporting; the recorded size is the exact byte count.
|
||||
_SIZE_SLACK = 4096
|
||||
|
||||
|
||||
class Verdict(StrEnum):
|
||||
@@ -75,12 +74,18 @@ def verify_video(job: JobPaths, artifact: Artifact, declared_size: int | None) -
|
||||
if path is None or not path.exists() or path.stat().st_size == 0:
|
||||
return Verdict.PARTIAL if _resumable_partial_exists(job) else Verdict.MISSING
|
||||
|
||||
whole = declared_size is None or abs(path.stat().st_size - declared_size) < _SIZE_SLACK
|
||||
if not whole:
|
||||
# Matroska and WebM write the duration into the header, so a file cut off
|
||||
# after a kilobyte still probes as its full length. The size recorded at
|
||||
# download time is the only thing that catches a truncation.
|
||||
return Verdict.CORRUPT
|
||||
|
||||
duration = _probe_duration(path)
|
||||
if duration is None:
|
||||
# Some mkv/webm muxes declare no duration; fall back to the expected size.
|
||||
if declared_size and abs(path.stat().st_size - declared_size) < 4096:
|
||||
return Verdict.OK
|
||||
return Verdict.CORRUPT
|
||||
# Some mkv/webm muxes declare no duration; the recorded size is then the
|
||||
# only evidence, and it already matched.
|
||||
return Verdict.OK if declared_size else Verdict.CORRUPT
|
||||
return Verdict.OK if duration > 0 else Verdict.CORRUPT
|
||||
|
||||
|
||||
@@ -95,10 +100,8 @@ def verify_audio(job: JobPaths, artifact: Artifact, video_duration: float | None
|
||||
duration = _probe_duration(path)
|
||||
if duration is None or duration <= 0:
|
||||
return Verdict.CORRUPT
|
||||
if video_duration:
|
||||
tolerance = max(_MIN_TOLERANCE_S, video_duration * _DURATION_TOLERANCE)
|
||||
if abs(duration - video_duration) > tolerance:
|
||||
return Verdict.CORRUPT
|
||||
if video_duration and not ffmpeg.durations_match(duration, video_duration):
|
||||
return Verdict.CORRUPT
|
||||
return Verdict.OK
|
||||
|
||||
|
||||
|
||||
@@ -96,6 +96,17 @@ def probe_duration(path: Path) -> float | None:
|
||||
return None
|
||||
|
||||
|
||||
# A derived stream this far off its source means the extraction was cut short.
|
||||
# The absolute floor covers short clips, where container rounding alone exceeds 0.5%.
|
||||
DURATION_TOLERANCE = 0.005
|
||||
MIN_TOLERANCE_S = 1.0
|
||||
|
||||
|
||||
def durations_match(actual: float, expected: float) -> bool:
|
||||
"""Whether a derived stream is as long as the source it came from."""
|
||||
return abs(actual - expected) <= max(MIN_TOLERANCE_S, expected * DURATION_TOLERANCE)
|
||||
|
||||
|
||||
def has_audio_stream(path: Path) -> bool:
|
||||
try:
|
||||
return probe_field(path, "stream=index", "a:0") is not None
|
||||
|
||||
@@ -72,6 +72,7 @@ def aria2c_args(
|
||||
timeout: int = 30,
|
||||
connect_timeout: int = 15,
|
||||
retry_wait: int = 3,
|
||||
auto_save: int = 20,
|
||||
log_file: Path | None = None,
|
||||
extra: Sequence[str] = (),
|
||||
) -> list[str]:
|
||||
@@ -88,6 +89,10 @@ def aria2c_args(
|
||||
# Sixteen writers at out-of-order offsets thrash a 16M cache.
|
||||
f"--disk-cache={disk_cache}",
|
||||
f"--file-allocation={file_allocation}",
|
||||
# Default 60s. The control file is what a resume reads, so anything not
|
||||
# checkpointed when the process is killed is re-downloaded; a kill in the
|
||||
# first minute otherwise preserves a file recording zero completed pieces.
|
||||
f"--auto-save-interval={auto_save}",
|
||||
]
|
||||
if connections != DEFAULT_CONNECTIONS:
|
||||
args += [f"-x{connections}", f"-s{connections}"]
|
||||
|
||||
@@ -20,7 +20,14 @@ if TYPE_CHECKING:
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def extract(job: JobPaths, video: Path, *, profile: str = "source", compression: int = 5) -> Path:
|
||||
def extract(
|
||||
job: JobPaths,
|
||||
video: Path,
|
||||
*,
|
||||
profile: str = "source",
|
||||
compression: int = 5,
|
||||
expected_s: float | None = None,
|
||||
) -> Path:
|
||||
job.ensure()
|
||||
staging = job.tmp_dir / "audio.flac"
|
||||
staging.unlink(missing_ok=True)
|
||||
@@ -33,6 +40,17 @@ def extract(job: JobPaths, video: Path, *, profile: str = "source", compression:
|
||||
"extracted FLAC has no readable duration",
|
||||
hint="The source audio track may be damaged.",
|
||||
)
|
||||
if expected_s and not ffmpeg.durations_match(duration, expected_s):
|
||||
# ffmpeg exits 0 on a truncated container, extracting the few frames it
|
||||
# found. Without this the pipeline transcribes near-silence and calls it
|
||||
# a success, which is worse than failing.
|
||||
staging.unlink(missing_ok=True)
|
||||
raise errors.DecodeError(
|
||||
f"extracted {duration:.1f}s of audio, which does not match the "
|
||||
f"{expected_s:.1f}s source",
|
||||
hint="The source media is probably truncated; re-run with --force-stage download.",
|
||||
)
|
||||
|
||||
staging.replace(job.audio_file)
|
||||
log.info("extracted %.1fs of audio to %s", duration, job.audio_file.name)
|
||||
return job.audio_file
|
||||
|
||||
@@ -64,3 +64,31 @@ def not_media(tmp_path: Path) -> Path:
|
||||
path = tmp_path / "notes.txt"
|
||||
path.write_text("this is definitely not a media file")
|
||||
return path
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def sine_video(tmp_path_factory: pytest.TempPathFactory) -> Path:
|
||||
"""A real mkv with audio -- truncating it leaves a readable header duration."""
|
||||
path = tmp_path_factory.mktemp("media") / "clip.mkv"
|
||||
_run(
|
||||
[
|
||||
"ffmpeg",
|
||||
"-nostdin",
|
||||
"-loglevel",
|
||||
"error",
|
||||
"-y",
|
||||
"-f",
|
||||
"lavfi",
|
||||
"-i",
|
||||
"color=c=black:s=64x64:d=1",
|
||||
"-f",
|
||||
"lavfi",
|
||||
"-i",
|
||||
"sine=frequency=440:duration=1:sample_rate=44100",
|
||||
"-r",
|
||||
"5",
|
||||
"-shortest",
|
||||
str(path),
|
||||
]
|
||||
)
|
||||
return path
|
||||
|
||||
@@ -112,6 +112,30 @@ class TestVerifyVideo:
|
||||
monkeypatch.setattr(ffmpeg, "probe_duration", no_duration)
|
||||
assert verify.verify_video(job, _present(dst, job.root), 999_999) is verify.Verdict.CORRUPT
|
||||
|
||||
def test_truncation_is_caught_by_size_even_though_it_probes(
|
||||
self, job: paths.JobPaths, sine_video: Path
|
||||
) -> None:
|
||||
# Matroska writes the duration in its header, so ffprobe reports the full
|
||||
# length of a file that was cut off after a kilobyte. The recorded size is
|
||||
# the only evidence that the bytes are gone.
|
||||
dst = job.media_dir / "video.mkv"
|
||||
whole = sine_video.read_bytes()
|
||||
dst.write_bytes(whole[:1024])
|
||||
assert ffmpeg.probe_duration(dst) == pytest.approx(1.0, abs=0.2)
|
||||
assert (
|
||||
verify.verify_video(job, _present(dst, job.root), len(whole)) is verify.Verdict.CORRUPT
|
||||
)
|
||||
|
||||
def test_a_whole_file_matching_its_recorded_size_is_ok(
|
||||
self, job: paths.JobPaths, sine_video: Path
|
||||
) -> None:
|
||||
dst = job.media_dir / "video.mkv"
|
||||
dst.write_bytes(sine_video.read_bytes())
|
||||
assert (
|
||||
verify.verify_video(job, _present(dst, job.root), dst.stat().st_size)
|
||||
is verify.Verdict.OK
|
||||
)
|
||||
|
||||
|
||||
class TestVerifyAudio:
|
||||
def _flac(self, job: paths.JobPaths, src: Path) -> job_state.Artifact:
|
||||
|
||||
@@ -170,3 +170,21 @@ class TestPreflight:
|
||||
monkeypatch.setattr(subprocess, "run", boom)
|
||||
with pytest.raises(errors.PreflightError):
|
||||
ffmpeg.has_flac_encoder()
|
||||
|
||||
|
||||
class TestDurationsMatch:
|
||||
def test_an_exact_match(self) -> None:
|
||||
assert ffmpeg.durations_match(212.861, 212.861) is True
|
||||
|
||||
def test_within_the_relative_tolerance(self) -> None:
|
||||
assert ffmpeg.durations_match(3600.0, 3610.0) is True
|
||||
|
||||
def test_beyond_the_relative_tolerance(self) -> None:
|
||||
assert ffmpeg.durations_match(3600.0, 3700.0) is False
|
||||
|
||||
def test_short_clips_get_an_absolute_floor(self) -> None:
|
||||
# 0.5% of two seconds is 10ms; rounding in container metadata exceeds that.
|
||||
assert ffmpeg.durations_match(2.0, 2.9) is True
|
||||
|
||||
def test_a_near_empty_extraction_never_matches(self) -> None:
|
||||
assert ffmpeg.durations_match(0.02, 212.861) is False
|
||||
|
||||
@@ -25,6 +25,15 @@ class TestAria2cArgs:
|
||||
def test_stall_abort_can_be_disabled_for_genuinely_slow_links(self) -> None:
|
||||
assert "--lowest-speed-limit=0" in opts.aria2c_args(lowest_speed="0")
|
||||
|
||||
def test_checkpoints_progress_often_enough_to_resume(self) -> None:
|
||||
# aria2 saves the control file every 60s by default, so a hard kill in the
|
||||
# first minute leaves a control file recording zero completed pieces and
|
||||
# the "resume" silently re-downloads the lot. Measured: 0/13 pieces.
|
||||
assert "--auto-save-interval=20" in opts.aria2c_args()
|
||||
|
||||
def test_the_checkpoint_interval_is_tunable(self) -> None:
|
||||
assert "--auto-save-interval=5" in opts.aria2c_args(auto_save=5)
|
||||
|
||||
def test_raises_the_disk_cache(self) -> None:
|
||||
assert "--disk-cache=64M" in opts.aria2c_args()
|
||||
|
||||
|
||||
@@ -78,6 +78,24 @@ class TestAudioStage:
|
||||
(job.tmp_dir / "audio.flac").write_bytes(b"junk from a killed run")
|
||||
assert audio_stage.extract(job, sine_wav).exists()
|
||||
|
||||
def test_rejects_an_extraction_shorter_than_its_source(
|
||||
self, job: paths.JobPaths, sine_wav: Path
|
||||
) -> None:
|
||||
# A truncated video yields a fraction of a second of real audio while the
|
||||
# container still claims the full length; without this the pipeline
|
||||
# transcribes near-silence and reports success.
|
||||
with pytest.raises(errors.DecodeError, match="does not match"):
|
||||
audio_stage.extract(job, sine_wav, expected_s=212.861)
|
||||
assert not job.audio_file.exists()
|
||||
|
||||
def test_accepts_an_extraction_matching_its_source(
|
||||
self, job: paths.JobPaths, sine_wav: Path
|
||||
) -> None:
|
||||
assert audio_stage.extract(job, sine_wav, expected_s=1.0).exists()
|
||||
|
||||
def test_no_expectation_means_no_comparison(self, job: paths.JobPaths, sine_wav: Path) -> None:
|
||||
assert audio_stage.extract(job, sine_wav, expected_s=None).exists()
|
||||
|
||||
|
||||
class FakeBackend:
|
||||
name: ClassVar[str] = "fake"
|
||||
|
||||
Reference in New Issue
Block a user