feat(jobs): add job identity, state, verification and resume planning

Artifacts on disk are the source of truth. state.json deliberately has no
top-level "stage" field -- persisting one is how "marked done but the file
is gone" bugs happen -- so the resume point is computed from what verifies.

The distinction that makes --no-retain safe is deleted_by_policy vs missing.
verify_* short-circuits on a policy deletion before touching the filesystem,
because probing a deliberately absent file would raise and degrade the whole
feature into "re-download everything".

A .part without its .aria2 control file is treated as unresumable: aria2
writes segments out of order, so such a file is sparse with holes rather
than a valid prefix, and resuming from its length yields a corrupt video.

Planning walks stages backwards. A policy deletion satisfies a stage that is
not re-running, but not one that is -- so --force-stage transcribe correctly
walks back to re-download. Saved segments let a deleted subtitle file be
re-rendered without re-transcribing a long recording.

state.json is written tmp -> fsync -> replace -> fsync(dir), with a test that
a failed replace leaves the previous record intact.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-09-13 14:24:24 -05:00
co-authored by Claude Opus 5
parent 15a28495b8
commit 2cfd7ad087
11 changed files with 1318 additions and 2 deletions
+9 -2
View File
@@ -63,10 +63,17 @@ strict = true
files = ["src", "tests"]
[tool.pylint.design]
max-attributes = 10
max-attributes = 15
[tool.pylint.main]
disable = ["C0114", "C0115", "C0116", "C0301", "W0611", "W0612", "R0801"]
# W0621: pytest fixtures shadow their names by design.
# C1803: asserting == () is more precise than truthiness about what a call returns.
# R0903: dataclasses and pytest grouping classes legitimately have few methods.
disable = [
"C0114", "C0115", "C0116", "C0301",
"W0611", "W0612", "W0621",
"R0801", "R0903", "C1803",
]
[tool.pytest.ini_options]
testpaths = ["tests"]
+63
View File
@@ -0,0 +1,63 @@
"""Resume planning.
Walk the stages backwards and enter at the first one whose output does not verify.
A stage that is *not* re-running is satisfied by an artifact deleted on purpose;
a stage that *is* re-running is not, and we walk further back to re-materialize it.
"""
from __future__ import annotations
from dataclasses import dataclass
from enum import StrEnum
from ccn_transcribe.jobs.verify import Verdict
class Stage(StrEnum):
DOWNLOAD = "download"
AUDIO = "audio"
TRANSCRIBE = "transcribe"
OUTPUTS = "outputs"
@dataclass(frozen=True, slots=True)
class Verdicts:
video: Verdict
audio: Verdict
outputs: Verdict
def plan_stages(
verdicts: Verdicts,
*,
params_changed: bool = False,
segments_available: bool = False,
force: frozenset[Stage] = frozenset(),
) -> tuple[Stage, ...]:
"""The stages to run, in pipeline order."""
need_outputs = (
verdicts.outputs is not Verdict.OK
or params_changed
or bool(force & {Stage.OUTPUTS, Stage.TRANSCRIBE, Stage.AUDIO, Stage.DOWNLOAD})
)
# Re-rendering from saved segments avoids re-transcribing a long recording
# just because one subtitle file was deleted.
need_transcribe = need_outputs and (
params_changed
or not segments_available
or bool(force & {Stage.TRANSCRIBE, Stage.AUDIO, Stage.DOWNLOAD})
)
need_audio = (
Stage.AUDIO in force
or Stage.DOWNLOAD in force
or (need_transcribe and verdicts.audio is not Verdict.OK)
)
need_download = Stage.DOWNLOAD in force or (need_audio and verdicts.video is not Verdict.OK)
wanted = {
Stage.DOWNLOAD: need_download,
Stage.AUDIO: need_audio,
Stage.TRANSCRIBE: need_transcribe,
Stage.OUTPUTS: need_outputs,
}
return tuple(stage for stage in Stage if wanted[stage])
+166
View File
@@ -0,0 +1,166 @@
"""The persisted job record.
Artifacts on disk are the source of truth; this record holds only what the
filesystem cannot answer -- above all *why* an artifact is absent, which is what
keeps a policy deletion from being mistaken for data loss.
There is deliberately no top-level "stage" field. Persisting one is how
"marked done but the file is gone" bugs happen; the resume point is computed.
"""
from __future__ import annotations
from dataclasses import asdict, dataclass, field
from datetime import UTC, datetime
from enum import StrEnum
from typing import Any
from ccn_transcribe import errors
SCHEMA_VERSION = 1
def _now() -> str:
return datetime.now(UTC).isoformat(timespec="seconds")
class ArtifactStatus(StrEnum):
PRESENT = "present"
MISSING = "missing"
DELETED_BY_POLICY = "deleted_by_policy"
@dataclass
class Artifact:
status: ArtifactStatus = ArtifactStatus.MISSING
path: str | None = None
size: int | None = None
duration_s: float | None = None
reason: str | None = None
deleted_at: str | None = None
@dataclass
class Outputs:
status: ArtifactStatus = ArtifactStatus.MISSING
paths: dict[str, str] = field(default_factory=dict[str, str])
@dataclass
class TranscriptionRecord:
model_id: str
model_revision: str | None
backend: str
device: str
language: str | None
task: str
num_beams: int
segments: int | None = None
audio_s: float | None = None
wall_s: float | None = None
rtf: float | None = None
started_at: str | None = None
finished_at: str | None = None
# Changing any of these invalidates an existing transcript. Backend and device are
# excluded on purpose: falling back from GPU to CPU must not force a re-run.
_RERUN_FIELDS = ("model_id", "language", "task", "num_beams")
@dataclass
class JobState:
job_id: str
source_url: str
schema_version: int = SCHEMA_VERSION
title: str | None = None
duration_s: float | None = None
created_at: str = field(default_factory=_now)
updated_at: str = field(default_factory=_now)
video: Artifact = field(default_factory=Artifact)
audio: Artifact = field(default_factory=Artifact)
outputs: Outputs = field(default_factory=Outputs)
transcription: TranscriptionRecord | None = None
attempts: dict[str, int] = field(default_factory=dict[str, int])
last_error: dict[str, str] | None = None
@classmethod
def new(cls, *, job_id: str, source_url: str) -> JobState:
return cls(job_id=job_id, source_url=source_url)
def params_changed(self, wanted: TranscriptionRecord) -> bool:
if self.transcription is None:
return True
return any(
getattr(self.transcription, name) != getattr(wanted, name) for name in _RERUN_FIELDS
)
def touch(self) -> None:
self.updated_at = _now()
def to_dict(self) -> dict[str, Any]:
return {
"schema_version": self.schema_version,
"job_id": self.job_id,
"source_url": self.source_url,
"title": self.title,
"duration_s": self.duration_s,
"created_at": self.created_at,
"updated_at": self.updated_at,
"artifacts": {
"video": asdict(self.video),
"audio": asdict(self.audio),
"outputs": asdict(self.outputs),
},
"transcription": asdict(self.transcription) if self.transcription else None,
"attempts": dict(self.attempts),
"last_error": self.last_error,
}
@classmethod
def from_dict(cls, doc: dict[str, Any]) -> JobState:
version = int(doc.get("schema_version", SCHEMA_VERSION))
if version > SCHEMA_VERSION:
raise errors.SchemaTooNewError(
f"state.json is schema v{version}, but this build understands v{SCHEMA_VERSION}",
hint="Upgrade ccn-transcribe, or point --workdir somewhere else.",
)
artifacts = doc.get("artifacts", {})
record = doc.get("transcription")
return cls(
job_id=doc["job_id"],
source_url=doc["source_url"],
schema_version=SCHEMA_VERSION,
title=doc.get("title"),
duration_s=doc.get("duration_s"),
created_at=doc.get("created_at", _now()),
updated_at=doc.get("updated_at", _now()),
video=_artifact(artifacts.get("video")),
audio=_artifact(artifacts.get("audio")),
outputs=_outputs(artifacts.get("outputs")),
transcription=TranscriptionRecord(**record) if record else None,
attempts=dict(doc.get("attempts", {})),
last_error=doc.get("last_error"),
)
def _artifact(raw: dict[str, Any] | None) -> Artifact:
if not raw:
return Artifact()
return Artifact(
status=ArtifactStatus(raw.get("status", ArtifactStatus.MISSING)),
path=raw.get("path"),
size=raw.get("size"),
duration_s=raw.get("duration_s"),
reason=raw.get("reason"),
deleted_at=raw.get("deleted_at"),
)
def _outputs(raw: dict[str, Any] | None) -> Outputs:
if not raw:
return Outputs()
return Outputs(
status=ArtifactStatus(raw.get("status", ArtifactStatus.MISSING)),
paths=dict(raw.get("paths", {})),
)
+118
View File
@@ -0,0 +1,118 @@
"""Durable job state: atomic writes, an append-only event log, and a job lock."""
from __future__ import annotations
import fcntl
import json
import logging
import os
from contextlib import contextmanager
from datetime import UTC, datetime
from typing import TYPE_CHECKING, Any
from ccn_transcribe import errors
from ccn_transcribe.jobs.state import JobState
if TYPE_CHECKING:
from collections.abc import Generator
from ccn_transcribe.paths import JobPaths
log = logging.getLogger(__name__)
def _fsync_dir(path: os.PathLike[str]) -> None:
fd = os.open(path, os.O_RDONLY)
try:
os.fsync(fd)
finally:
os.close(fd)
def save(job: JobPaths, state: JobState) -> None:
"""Write state.json atomically.
tmp -> fsync -> replace -> fsync(dir) so an interrupted write can never leave a
truncated record where a good one used to be.
"""
state.touch()
job.root.mkdir(parents=True, exist_ok=True)
tmp = job.state_file.with_suffix(".json.tmp")
with tmp.open("w", encoding="utf-8") as fh:
json.dump(state.to_dict(), fh, indent=2, ensure_ascii=False)
fh.flush()
os.fsync(fh.fileno())
try:
tmp.replace(job.state_file)
except OSError:
tmp.unlink(missing_ok=True)
raise
_fsync_dir(job.root)
def load(job: JobPaths) -> JobState | None:
"""Read state.json, or None when there is nothing usable to read."""
if not job.state_file.exists():
return None
try:
doc = json.loads(job.state_file.read_text(encoding="utf-8"))
except (OSError, ValueError):
# Artifacts are the source of truth, so an unreadable record is survivable.
log.warning("state.json for %s is unreadable; rebuilding from artifacts", job.root.name)
return None
return JobState.from_dict(doc)
def append_event(job: JobPaths, event: dict[str, Any]) -> None:
"""Append one durable event.
This is the record of *why* an artifact is absent. Without it a damaged
state.json would downgrade a policy deletion to "missing" and trigger a
pointless re-download.
"""
job.root.mkdir(parents=True, exist_ok=True)
payload = {"at": datetime.now(UTC).isoformat(timespec="seconds"), **event}
with job.events_file.open("a", encoding="utf-8") as fh:
fh.write(json.dumps(payload, ensure_ascii=False) + "\n")
fh.flush()
os.fsync(fh.fileno())
def read_events(job: JobPaths) -> list[dict[str, Any]]:
if not job.events_file.exists():
return []
events: list[dict[str, Any]] = []
for line in job.events_file.read_text(encoding="utf-8").splitlines():
if not line.strip():
continue
try:
events.append(json.loads(line))
except ValueError:
# One torn line must not cost us the rest of the history.
log.warning("skipping an unreadable line in %s", job.events_file.name)
return events
@contextmanager
def job_lock(job: JobPaths) -> Generator[None]:
"""Hold an exclusive lock for the job's lifetime.
flock rather than a PID file: the kernel drops it when the process dies, so
there is no stale-lock case and no heuristic to get wrong.
"""
job.root.mkdir(parents=True, exist_ok=True)
handle = job.lock_file.open("w", encoding="utf-8")
try:
try:
fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
except OSError as exc:
raise errors.LockedError(
f"job {job.root.name} is already running in another process",
job_id=job.root.name,
hint="Wait for it to finish, or use a different --workdir.",
) from exc
handle.write(str(os.getpid()))
handle.flush()
yield
finally:
handle.close()
+113
View File
@@ -0,0 +1,113 @@
"""Artifact verification.
A stage counts as complete only if its artifact is actually on disk and readable.
Trusting a recorded marker is how a truncated download or a deleted transcript
gets treated as done.
"""
from __future__ import annotations
import logging
from enum import StrEnum
from typing import TYPE_CHECKING
from ccn_transcribe import errors
from ccn_transcribe.jobs.state import Artifact, ArtifactStatus, Outputs
from ccn_transcribe.media import ffmpeg
if TYPE_CHECKING:
from collections.abc import Iterator, Sequence
from pathlib import Path
from ccn_transcribe.paths import JobPaths
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
class Verdict(StrEnum):
OK = "ok"
SATISFIED_ABSENT = "satisfied_absent"
PARTIAL = "partial"
CORRUPT = "corrupt"
MISSING = "missing"
def _resolve(job: JobPaths, artifact: Artifact) -> Path | None:
if artifact.path is None:
return None
return job.root / artifact.path
def _part_files(job: JobPaths) -> Iterator[Path]:
for directory in (job.media_dir, job.tmp_dir):
if directory.is_dir():
yield from directory.glob("*.part")
def _resumable_partial_exists(job: JobPaths) -> bool:
"""A .part is only resumable alongside its aria2 control file.
aria2 writes segments out of order, so a .part without the control file is a
sparse file with holes rather than a valid prefix; resuming from its length
would silently produce a corrupt video.
"""
return any(part.with_suffix(part.suffix + ".aria2").exists() for part in _part_files(job))
def _probe_duration(path: Path) -> float | None:
try:
return ffmpeg.probe_duration(path)
except errors.ProbeError:
return None
def verify_video(job: JobPaths, artifact: Artifact, declared_size: int | None) -> Verdict:
if artifact.status is ArtifactStatus.DELETED_BY_POLICY:
# Short-circuit before touching the filesystem: the obligation is
# discharged, and probing a deliberately absent file would only error.
return Verdict.SATISFIED_ABSENT
path = _resolve(job, artifact)
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
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
return Verdict.OK if duration > 0 else Verdict.CORRUPT
def verify_audio(job: JobPaths, artifact: Artifact, video_duration: float | None) -> Verdict:
if artifact.status is ArtifactStatus.DELETED_BY_POLICY:
return Verdict.SATISFIED_ABSENT
path = _resolve(job, artifact)
if path is None or not path.exists() or path.stat().st_size == 0:
return Verdict.MISSING
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
return Verdict.OK
def verify_outputs(job: JobPaths, outputs: Outputs, wanted: Sequence[str]) -> Verdict:
for name in wanted:
relative = outputs.paths.get(name)
if relative is None:
return Verdict.MISSING
path = job.root / relative
if not path.exists() or path.stat().st_size == 0:
return Verdict.MISSING
return Verdict.OK
+115
View File
@@ -0,0 +1,115 @@
"""Job identity and on-disk layout.
Job ids come from yt-dlp's extractor key plus video id rather than the URL, since
one video has many URL forms. ``index.json`` maps normalized URLs to job ids so
resuming never needs a network round trip.
"""
from __future__ import annotations
import hashlib
import re
from dataclasses import dataclass
from typing import TYPE_CHECKING
from urllib.parse import urlsplit, urlunsplit
if TYPE_CHECKING:
from pathlib import Path
_UNSAFE = re.compile(r"[^A-Za-z0-9_.-]")
_NOT_ALNUM = re.compile(r"[^a-z0-9]+")
_MAX_ID = 96
_TRUNCATE_TO = 64
_DIGEST_LEN = 10
def normalize_url(url: str) -> str:
"""Collapse URL spellings that name the same resource."""
parts = urlsplit(url.strip())
return urlunsplit((parts.scheme.lower(), parts.netloc.lower(), parts.path, parts.query, ""))
def _digest(raw: str) -> str:
return hashlib.sha256(raw.encode()).hexdigest()[:_DIGEST_LEN]
def job_id(extractor_key: str, video_id: str) -> str:
"""A filesystem-safe, collision-resistant directory name for one video."""
extractor = _NOT_ALNUM.sub("", extractor_key.lower())[:24] or "generic"
safe = _UNSAFE.sub("_", video_id)
# Sanitizing is lossy, so two distinct ids can collapse to one name; a digest
# keeps them in separate job directories.
if safe != video_id or len(safe) > _MAX_ID:
safe = f"{safe[:_TRUNCATE_TO]}-{_digest(video_id)}"
return f"{extractor}-{safe}"
def url_job_id(url: str) -> str:
"""Fallback identity for sources yt-dlp cannot name (direct file URLs)."""
return f"url-{hashlib.sha256(normalize_url(url).encode()).hexdigest()[:16]}"
@dataclass(frozen=True, slots=True)
class JobPaths:
root: Path
@property
def state_file(self) -> Path:
return self.root / "state.json"
@property
def lock_file(self) -> Path:
return self.root / ".lock"
@property
def events_file(self) -> Path:
return self.root / "events.jsonl"
@property
def info_file(self) -> Path:
return self.root / "info.json"
@property
def media_dir(self) -> Path:
return self.root / "media"
@property
def out_dir(self) -> Path:
return self.root / "out"
@property
def logs_dir(self) -> Path:
return self.root / "logs"
@property
def tmp_dir(self) -> Path:
# Inside the job dir so os.replace stays within one filesystem.
return self.root / "tmp"
@property
def audio_file(self) -> Path:
return self.media_dir / "audio.flac"
def transcript_file(self, extension: str) -> Path:
return self.out_dir / f"transcript.{extension}"
def ensure(self) -> None:
for directory in (self.media_dir, self.out_dir, self.logs_dir, self.tmp_dir):
directory.mkdir(parents=True, exist_ok=True)
@dataclass(frozen=True, slots=True)
class Workspace:
root: Path
@property
def jobs_dir(self) -> Path:
return self.root / "jobs"
@property
def index_file(self) -> Path:
return self.root / "index.json"
def job(self, identifier: str) -> JobPaths:
return JobPaths(self.jobs_dir / identifier)
+130
View File
@@ -0,0 +1,130 @@
from __future__ import annotations
import pytest
from ccn_transcribe.jobs.plan import Stage, Verdicts, plan_stages
from ccn_transcribe.jobs.verify import Verdict
OK = Verdict.OK
GONE = Verdict.SATISFIED_ABSENT
MISS = Verdict.MISSING
PART = Verdict.PARTIAL
BAD = Verdict.CORRUPT
def plan(video: Verdict, audio: Verdict, outputs: Verdict, **kw: object) -> tuple[Stage, ...]:
return plan_stages(Verdicts(video=video, audio=audio, outputs=outputs), **kw) # type: ignore[arg-type]
class TestNothingToDo:
def test_everything_present_is_a_no_op(self) -> None:
assert plan(OK, OK, OK) == ()
def test_media_deleted_by_policy_is_still_a_no_op(self) -> None:
# The whole point of --no-retain: done is done, do not re-fetch.
assert plan(GONE, GONE, OK) == ()
class TestEntryPoint:
def test_missing_outputs_with_good_audio_skips_download(self) -> None:
assert plan(OK, OK, MISS) == (Stage.TRANSCRIBE, Stage.OUTPUTS)
def test_policy_deleted_video_does_not_trigger_a_download(self) -> None:
assert plan(GONE, OK, MISS) == (Stage.TRANSCRIBE, Stage.OUTPUTS)
def test_missing_audio_with_good_video_extracts_only(self) -> None:
assert plan(OK, MISS, MISS) == (Stage.AUDIO, Stage.TRANSCRIBE, Stage.OUTPUTS)
def test_nothing_present_runs_everything(self) -> None:
assert plan(MISS, MISS, MISS) == (
Stage.DOWNLOAD,
Stage.AUDIO,
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
def test_a_partial_download_still_enters_at_download(self) -> None:
assert plan(PART, MISS, MISS)[0] is Stage.DOWNLOAD
def test_a_corrupt_video_is_re_downloaded(self) -> None:
assert plan(BAD, MISS, MISS)[0] is Stage.DOWNLOAD
def test_a_corrupt_audio_is_re_extracted(self) -> None:
assert plan(OK, BAD, MISS) == (Stage.AUDIO, Stage.TRANSCRIBE, Stage.OUTPUTS)
class TestSegmentReuse:
def test_outputs_alone_are_re_rendered_when_segments_survive(self) -> None:
# Deleting transcript.srt should not cost a re-transcription of a
# two-hour lecture when transcript.json still holds the segments.
assert plan(OK, OK, MISS, segments_available=True) == (Stage.OUTPUTS,)
def test_policy_deleted_audio_is_fine_when_segments_survive(self) -> None:
assert plan(GONE, GONE, MISS, segments_available=True) == (Stage.OUTPUTS,)
def test_changed_params_ignore_saved_segments(self) -> None:
assert Stage.TRANSCRIBE in plan(OK, OK, MISS, segments_available=True, params_changed=True)
class TestParamsChanged:
def test_a_new_language_forces_a_re_transcription(self) -> None:
assert plan(OK, OK, OK, params_changed=True) == (Stage.TRANSCRIBE, Stage.OUTPUTS)
def test_it_re_extracts_audio_that_was_policy_deleted(self) -> None:
assert plan(OK, GONE, OK, params_changed=True) == (
Stage.AUDIO,
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
def test_it_re_downloads_when_both_were_policy_deleted(self) -> None:
# A stage that is re-running is NOT satisfied by a policy deletion.
assert plan(GONE, GONE, OK, params_changed=True) == (
Stage.DOWNLOAD,
Stage.AUDIO,
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
class TestForce:
def test_forcing_transcribe_reruns_it(self) -> None:
assert plan(OK, OK, OK, force=frozenset({Stage.TRANSCRIBE})) == (
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
def test_forcing_transcribe_walks_back_past_policy_deletions(self) -> None:
assert plan(GONE, GONE, OK, force=frozenset({Stage.TRANSCRIBE})) == (
Stage.DOWNLOAD,
Stage.AUDIO,
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
def test_forcing_download_reruns_the_whole_pipeline(self) -> None:
assert plan(OK, OK, OK, force=frozenset({Stage.DOWNLOAD})) == (
Stage.DOWNLOAD,
Stage.AUDIO,
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
def test_forcing_audio_keeps_the_existing_download(self) -> None:
assert plan(OK, OK, OK, force=frozenset({Stage.AUDIO})) == (
Stage.AUDIO,
Stage.TRANSCRIBE,
Stage.OUTPUTS,
)
def test_forcing_outputs_only_re_renders(self) -> None:
assert plan(OK, OK, OK, force=frozenset({Stage.OUTPUTS}), segments_available=True) == (
Stage.OUTPUTS,
)
class TestStageOrder:
@pytest.mark.parametrize("segments", [True, False])
def test_stages_always_come_out_in_pipeline_order(self, segments: bool) -> None:
got = plan(MISS, MISS, MISS, segments_available=segments)
assert list(got) == sorted(got, key=list(Stage).index)
+137
View File
@@ -0,0 +1,137 @@
from __future__ import annotations
import pytest
from ccn_transcribe import errors
from ccn_transcribe.jobs import state as job_state
def make() -> job_state.JobState:
return job_state.JobState.new(job_id="youtube-abc", source_url="https://y.test/watch?v=abc")
class TestNew:
def test_starts_with_everything_missing(self) -> None:
job = make()
assert job.video.status is job_state.ArtifactStatus.MISSING
assert job.audio.status is job_state.ArtifactStatus.MISSING
assert job.outputs.status is job_state.ArtifactStatus.MISSING
def test_records_identity_and_timestamps(self) -> None:
job = make()
assert job.job_id == "youtube-abc"
assert job.source_url == "https://y.test/watch?v=abc"
assert job.created_at and job.updated_at
def test_carries_the_current_schema_version(self) -> None:
assert make().schema_version == job_state.SCHEMA_VERSION
class TestRoundTrip:
def test_survives_serialization(self) -> None:
job = make()
job.video = job_state.Artifact(
job_state.ArtifactStatus.PRESENT, path="media/v.mkv", duration_s=12.5
)
job.title = "A title"
assert job_state.JobState.from_dict(job.to_dict()) == job
def test_enums_serialize_as_plain_strings(self) -> None:
doc = make().to_dict()
assert doc["artifacts"]["video"]["status"] == "missing"
def test_outputs_paths_round_trip(self) -> None:
job = make()
job.outputs = job_state.Outputs(
job_state.ArtifactStatus.PRESENT, paths={"srt": "out/transcript.srt"}
)
assert job_state.JobState.from_dict(job.to_dict()).outputs.paths == {
"srt": "out/transcript.srt"
}
def test_transcription_record_round_trips(self) -> None:
job = make()
job.transcription = job_state.TranscriptionRecord(
model_id="m",
model_revision="r",
backend="openvino",
device="GPU",
language="en",
task="transcribe",
num_beams=1,
)
restored = job_state.JobState.from_dict(job.to_dict())
assert restored.transcription is not None
assert restored.transcription.device == "GPU"
def test_unknown_keys_are_ignored_not_fatal(self) -> None:
doc = make().to_dict()
doc["something_from_the_future"] = 1
assert job_state.JobState.from_dict(doc).job_id == "youtube-abc"
class TestSchemaVersioning:
def test_a_newer_schema_is_refused(self) -> None:
doc = make().to_dict()
doc["schema_version"] = job_state.SCHEMA_VERSION + 1
# Guessing at a future layout risks destroying the user's state dir.
with pytest.raises(errors.SchemaTooNewError):
job_state.JobState.from_dict(doc)
def test_the_current_schema_is_accepted(self) -> None:
assert job_state.JobState.from_dict(make().to_dict()) is not None
class TestParamsChanged:
def _record(self, **kw: object) -> job_state.TranscriptionRecord:
base: dict[str, object] = {
"model_id": "m",
"model_revision": "r",
"backend": "openvino",
"device": "GPU",
"language": "en",
"task": "transcribe",
"num_beams": 1,
}
base.update(kw)
return job_state.TranscriptionRecord(**base) # type: ignore[arg-type]
def test_no_prior_run_counts_as_changed(self) -> None:
assert make().params_changed(self._record()) is True
def test_identical_params_are_unchanged(self) -> None:
job = make()
job.transcription = self._record()
assert job.params_changed(self._record()) is False
@pytest.mark.parametrize(
("field", "value"),
[("model_id", "other"), ("language", "de"), ("task", "translate"), ("num_beams", 5)],
)
def test_a_different_request_forces_a_rerun(self, field: str, value: object) -> None:
job = make()
job.transcription = self._record()
assert job.params_changed(self._record(**{field: value})) is True
def test_device_and_backend_do_not_force_a_rerun(self) -> None:
# Falling back from GPU to CPU must not invalidate a good transcript.
job = make()
job.transcription = self._record()
assert job.params_changed(self._record(device="CPU", backend="x")) is False
class TestDefaultsFromSparseDocuments:
def test_a_document_with_no_artifacts_key_gets_defaults(self) -> None:
doc = make().to_dict()
del doc["artifacts"]
restored = job_state.JobState.from_dict(doc)
assert restored.video.status is job_state.ArtifactStatus.MISSING
assert not restored.outputs.paths
def test_an_empty_artifact_entry_gets_defaults(self) -> None:
doc = make().to_dict()
doc["artifacts"] = {"video": None, "audio": {}, "outputs": None}
restored = job_state.JobState.from_dict(doc)
assert restored.video.status is job_state.ArtifactStatus.MISSING
assert restored.audio.status is job_state.ArtifactStatus.MISSING
assert restored.outputs.status is job_state.ArtifactStatus.MISSING
+144
View File
@@ -0,0 +1,144 @@
from __future__ import annotations
import json
import os
from typing import TYPE_CHECKING
import pytest
from ccn_transcribe import errors, paths
from ccn_transcribe.jobs import state as job_state
from ccn_transcribe.jobs import store
if TYPE_CHECKING:
from pathlib import Path
@pytest.fixture
def job(tmp_path: Path) -> paths.JobPaths:
j = paths.Workspace(tmp_path).job("youtube-abc")
j.ensure()
return j
def _state() -> job_state.JobState:
return job_state.JobState.new(job_id="youtube-abc", source_url="https://y.test/x")
class TestSaveLoad:
def test_round_trip(self, job: paths.JobPaths) -> None:
want = _state()
want.title = "Hello"
store.save(job, want)
assert store.load(job).title == "Hello" # type: ignore[union-attr]
def test_load_returns_none_when_absent(self, job: paths.JobPaths) -> None:
assert store.load(job) is None
def test_save_leaves_no_temp_file(self, job: paths.JobPaths) -> None:
store.save(job, _state())
assert list(job.root.glob("*.tmp")) == []
def test_save_writes_readable_json(self, job: paths.JobPaths) -> None:
store.save(job, _state())
assert json.loads(job.state_file.read_text())["job_id"] == "youtube-abc"
def test_save_bumps_updated_at(self, job: paths.JobPaths) -> None:
st = _state()
st.updated_at = "1999-01-01T00:00:00+00:00"
store.save(job, st)
assert store.load(job).updated_at != "1999-01-01T00:00:00+00:00" # type: ignore[union-attr]
def test_a_failed_save_leaves_the_previous_state_intact(
self, job: paths.JobPaths, monkeypatch: pytest.MonkeyPatch
) -> None:
# The durability claim: a crash mid-save must never truncate good state.
first = _state()
first.title = "original"
store.save(job, first)
def boom(*_a: object, **_k: object) -> None:
raise OSError("disk gave up")
monkeypatch.setattr(os, "replace", boom)
second = _state()
second.title = "replacement"
with pytest.raises(OSError, match="disk gave up"):
store.save(job, second)
assert store.load(job).title == "original" # type: ignore[union-attr]
def test_corrupt_state_is_treated_as_absent(
self, job: paths.JobPaths, caplog: pytest.LogCaptureFixture
) -> None:
# Artifacts are the source of truth, so an unreadable record is recoverable.
job.state_file.write_text("{not json at all")
assert store.load(job) is None
assert "state.json" in caplog.text
def test_a_newer_schema_still_raises(self, job: paths.JobPaths) -> None:
doc = _state().to_dict()
doc["schema_version"] = job_state.SCHEMA_VERSION + 1
job.state_file.write_text(json.dumps(doc))
with pytest.raises(errors.SchemaTooNewError):
store.load(job)
class TestEvents:
def test_append_and_read(self, job: paths.JobPaths) -> None:
store.append_event(job, {"event": "deleted", "artifact": "video"})
[got] = store.read_events(job)
assert got["event"] == "deleted"
assert got["artifact"] == "video"
def test_events_accumulate(self, job: paths.JobPaths) -> None:
store.append_event(job, {"event": "a"})
store.append_event(job, {"event": "b"})
assert [e["event"] for e in store.read_events(job)] == ["a", "b"]
def test_reading_with_no_file_is_empty(self, job: paths.JobPaths) -> None:
assert store.read_events(job) == []
def test_each_event_gets_a_timestamp(self, job: paths.JobPaths) -> None:
store.append_event(job, {"event": "a"})
assert "at" in store.read_events(job)[0]
def test_a_damaged_line_does_not_lose_the_rest(self, job: paths.JobPaths) -> None:
store.append_event(job, {"event": "good"})
with job.events_file.open("a") as fh:
fh.write("{truncated\n")
store.append_event(job, {"event": "also-good"})
assert [e["event"] for e in store.read_events(job)] == ["good", "also-good"]
class TestLock:
def test_lock_is_acquired_and_released(self, job: paths.JobPaths) -> None:
with store.job_lock(job):
assert job.lock_file.exists()
with store.job_lock(job):
pass
def test_a_second_holder_is_refused(self, job: paths.JobPaths) -> None:
with store.job_lock(job), pytest.raises(errors.LockedError), store.job_lock(job):
pass
def test_the_error_names_the_job(self, job: paths.JobPaths) -> None:
with store.job_lock(job):
try:
with store.job_lock(job):
pass
except errors.LockedError as exc:
assert "youtube-abc" in str(exc)
def test_lock_released_even_if_the_body_raises(self, job: paths.JobPaths) -> None:
with pytest.raises(ValueError, match="boom"), store.job_lock(job):
raise ValueError("boom")
with store.job_lock(job):
pass
def test_blank_lines_in_the_event_log_are_skipped(job: paths.JobPaths) -> None:
store.append_event(job, {"event": "a"})
with job.events_file.open("a") as fh:
fh.write("\n \n")
assert [e["event"] for e in store.read_events(job)] == ["a"]
+199
View File
@@ -0,0 +1,199 @@
from __future__ import annotations
from typing import TYPE_CHECKING
import pytest
from ccn_transcribe import paths
from ccn_transcribe.jobs import state as job_state
from ccn_transcribe.jobs import verify
from ccn_transcribe.media import ffmpeg
if TYPE_CHECKING:
from pathlib import Path
@pytest.fixture
def job(tmp_path: Path) -> paths.JobPaths:
j = paths.Workspace(tmp_path).job("youtube-abc")
j.ensure()
return j
def _present(path: Path, root: Path) -> job_state.Artifact:
return job_state.Artifact(job_state.ArtifactStatus.PRESENT, path=str(path.relative_to(root)))
class TestVerifyVideo:
def test_ok_for_a_real_file(self, job: paths.JobPaths, sine_wav: Path) -> None:
dst = job.media_dir / "video.mkv"
dst.write_bytes(sine_wav.read_bytes())
assert verify.verify_video(job, _present(dst, job.root), None) is verify.Verdict.OK
def test_deleted_by_policy_is_satisfied(self, job: paths.JobPaths) -> None:
art = job_state.Artifact(job_state.ArtifactStatus.DELETED_BY_POLICY, path="media/video.mkv")
assert verify.verify_video(job, art, None) is verify.Verdict.SATISFIED_ABSENT
def test_deleted_by_policy_never_touches_the_filesystem(
self, job: paths.JobPaths, monkeypatch: pytest.MonkeyPatch
) -> None:
# probe_duration raises on a missing path, so a policy-deleted video must
# short-circuit before any probe or the whole feature degrades to
# "re-download everything".
def boom(_p: Path) -> float:
raise AssertionError("must not probe a policy-deleted artifact")
monkeypatch.setattr(ffmpeg, "probe_duration", boom)
art = job_state.Artifact(job_state.ArtifactStatus.DELETED_BY_POLICY, path="media/video.mkv")
assert verify.verify_video(job, art, None) is verify.Verdict.SATISFIED_ABSENT
def test_missing_file_is_missing_not_an_error(self, job: paths.JobPaths) -> None:
art = job_state.Artifact(job_state.ArtifactStatus.PRESENT, path="media/video.mkv")
assert verify.verify_video(job, art, None) is verify.Verdict.MISSING
def test_no_recorded_path_is_missing(self, job: paths.JobPaths) -> None:
assert (
verify.verify_video(job, job_state.Artifact(job_state.ArtifactStatus.MISSING), None)
is verify.Verdict.MISSING
)
def test_zero_byte_file_is_missing(self, job: paths.JobPaths) -> None:
dst = job.media_dir / "video.mkv"
dst.touch()
assert verify.verify_video(job, _present(dst, job.root), None) is verify.Verdict.MISSING
def test_unreadable_file_is_corrupt(self, job: paths.JobPaths) -> None:
dst = job.media_dir / "video.mkv"
dst.write_bytes(b"not a video at all")
assert verify.verify_video(job, _present(dst, job.root), None) is verify.Verdict.CORRUPT
def test_part_with_control_file_is_partial(self, job: paths.JobPaths) -> None:
(job.media_dir / "video.mkv.part").write_bytes(b"x" * 10)
(job.media_dir / "video.mkv.part.aria2").write_bytes(b"ctl")
art = job_state.Artifact(job_state.ArtifactStatus.MISSING)
assert verify.verify_video(job, art, None) is verify.Verdict.PARTIAL
def test_part_without_control_file_is_missing(self, job: paths.JobPaths) -> None:
# aria2 writes segments out of order, so a .part with no control file is a
# sparse file with holes, not a valid prefix -- resuming it corrupts output.
(job.media_dir / "video.mkv.part").write_bytes(b"x" * 10)
art = job_state.Artifact(job_state.ArtifactStatus.MISSING)
assert verify.verify_video(job, art, None) is verify.Verdict.MISSING
def test_part_in_tmp_is_also_found(self, job: paths.JobPaths) -> None:
(job.tmp_dir / "video.mkv.part").write_bytes(b"x" * 10)
(job.tmp_dir / "video.mkv.part.aria2").write_bytes(b"ctl")
assert (
verify.verify_video(job, job_state.Artifact(job_state.ArtifactStatus.MISSING), None)
is verify.Verdict.PARTIAL
)
def test_falls_back_to_declared_size_when_duration_is_unknown(
self, job: paths.JobPaths, monkeypatch: pytest.MonkeyPatch
) -> None:
dst = job.media_dir / "video.mkv"
dst.write_bytes(b"x" * 5000)
def no_duration(_p: Path) -> None:
return None
monkeypatch.setattr(ffmpeg, "probe_duration", no_duration)
assert verify.verify_video(job, _present(dst, job.root), 5000) is verify.Verdict.OK
def test_size_mismatch_without_duration_is_corrupt(
self, job: paths.JobPaths, monkeypatch: pytest.MonkeyPatch
) -> None:
dst = job.media_dir / "video.mkv"
dst.write_bytes(b"x" * 10)
def no_duration(_p: Path) -> None:
return None
monkeypatch.setattr(ffmpeg, "probe_duration", no_duration)
assert verify.verify_video(job, _present(dst, job.root), 999_999) is verify.Verdict.CORRUPT
class TestVerifyAudio:
def _flac(self, job: paths.JobPaths, src: Path) -> job_state.Artifact:
ffmpeg.extract_flac(src, job.audio_file)
return _present(job.audio_file, job.root)
def test_ok_for_a_real_flac(self, job: paths.JobPaths, sine_wav: Path) -> None:
assert verify.verify_audio(job, self._flac(job, sine_wav), None) is verify.Verdict.OK
def test_deleted_by_policy_is_satisfied(self, job: paths.JobPaths) -> None:
art = job_state.Artifact(
job_state.ArtifactStatus.DELETED_BY_POLICY, path="media/audio.flac"
)
assert verify.verify_audio(job, art, None) is verify.Verdict.SATISFIED_ABSENT
def test_missing_is_missing(self, job: paths.JobPaths) -> None:
art = job_state.Artifact(job_state.ArtifactStatus.PRESENT, path="media/audio.flac")
assert verify.verify_audio(job, art, None) is verify.Verdict.MISSING
def test_zero_byte_is_missing(self, job: paths.JobPaths) -> None:
job.audio_file.touch()
art = _present(job.audio_file, job.root)
assert verify.verify_audio(job, art, None) is verify.Verdict.MISSING
def test_matches_the_video_duration(self, job: paths.JobPaths, sine_wav: Path) -> None:
assert verify.verify_audio(job, self._flac(job, sine_wav), 1.0) is verify.Verdict.OK
def test_truncated_extraction_is_corrupt(self, job: paths.JobPaths, sine_wav: Path) -> None:
# A one-second FLAC against a one-hour video means extraction was cut short.
assert verify.verify_audio(job, self._flac(job, sine_wav), 3600.0) is verify.Verdict.CORRUPT
def test_garbage_is_corrupt(self, job: paths.JobPaths) -> None:
job.audio_file.write_bytes(b"nonsense")
assert (
verify.verify_audio(job, _present(job.audio_file, job.root), None)
is verify.Verdict.CORRUPT
)
class TestVerifyOutputs:
def _write(self, job: paths.JobPaths, *names: str) -> job_state.Outputs:
out: dict[str, str] = {}
for n in names:
p = job.transcript_file(n)
p.write_text("content")
out[n] = str(p.relative_to(job.root))
return job_state.Outputs(job_state.ArtifactStatus.PRESENT, paths=out)
def test_ok_when_every_wanted_format_exists(self, job: paths.JobPaths) -> None:
got = self._write(job, "txt", "srt")
assert verify.verify_outputs(job, got, ("txt", "srt")) is verify.Verdict.OK
def test_missing_when_a_wanted_format_is_absent(self, job: paths.JobPaths) -> None:
got = self._write(job, "txt")
assert verify.verify_outputs(job, got, ("txt", "srt")) is verify.Verdict.MISSING
def test_missing_when_a_recorded_file_was_deleted(self, job: paths.JobPaths) -> None:
got = self._write(job, "txt")
job.transcript_file("txt").unlink()
assert verify.verify_outputs(job, got, ("txt",)) is verify.Verdict.MISSING
def test_empty_file_is_missing(self, job: paths.JobPaths) -> None:
got = self._write(job, "txt")
job.transcript_file("txt").write_text("")
assert verify.verify_outputs(job, got, ("txt",)) is verify.Verdict.MISSING
def test_nothing_recorded_is_missing(self, job: paths.JobPaths) -> None:
assert (
verify.verify_outputs(
job, job_state.Outputs(job_state.ArtifactStatus.MISSING), ("txt",)
)
is verify.Verdict.MISSING
)
def test_extra_formats_on_disk_do_not_matter(self, job: paths.JobPaths) -> None:
got = self._write(job, "txt", "srt", "vtt")
assert verify.verify_outputs(job, got, ("txt",)) is verify.Verdict.OK
def test_partial_scan_tolerates_missing_directories(tmp_path: Path) -> None:
# A job dir that was never fully created must not raise during verification.
bare = paths.Workspace(tmp_path).job("youtube-bare")
bare.root.mkdir(parents=True)
art = job_state.Artifact(job_state.ArtifactStatus.MISSING)
assert verify.verify_video(bare, art, None) is verify.Verdict.MISSING
+124
View File
@@ -0,0 +1,124 @@
from __future__ import annotations
from typing import TYPE_CHECKING
import pytest
from ccn_transcribe import paths
if TYPE_CHECKING:
from pathlib import Path
class TestJobId:
def test_combines_extractor_and_video_id(self) -> None:
assert paths.job_id("Youtube", "dQw4w9WgXcQ") == "youtube-dQw4w9WgXcQ"
def test_extractor_key_is_lowercased(self) -> None:
# Extractor-key casing has changed across yt-dlp releases; folding it
# keeps a version bump from orphaning existing job directories.
assert paths.job_id("YouTube", "x") == paths.job_id("youtube", "x")
def test_extractor_punctuation_is_stripped(self) -> None:
assert paths.job_id("Some:Site!", "x").startswith("somesite-")
def test_empty_extractor_falls_back_to_generic(self) -> None:
assert paths.job_id("", "x").startswith("generic-")
def test_unsafe_characters_are_replaced(self) -> None:
assert "/" not in paths.job_id("youtube", "a/b")
def test_ids_that_sanitize_alike_stay_distinct(self) -> None:
# "a/b" and "a:b" both sanitize to "a_b"; without a digest suffix they
# would share one job directory and corrupt each other's state.
assert paths.job_id("youtube", "a/b") != paths.job_id("youtube", "a:b")
def test_clean_ids_get_no_digest_suffix(self) -> None:
assert paths.job_id("youtube", "abc-123_x.y") == "youtube-abc-123_x.y"
def test_very_long_ids_are_truncated_but_stay_distinct(self) -> None:
long_a = "z" * 200 + "a"
long_b = "z" * 200 + "b"
assert len(paths.job_id("youtube", long_a)) < 120
assert paths.job_id("youtube", long_a) != paths.job_id("youtube", long_b)
def test_is_deterministic(self) -> None:
assert paths.job_id("youtube", "a/b") == paths.job_id("youtube", "a/b")
@pytest.mark.parametrize("hostile", ["..", "../..", "../../etc/passwd", "/etc/passwd"])
def test_cannot_escape_the_jobs_directory(self, hostile: str, tmp_path: Path) -> None:
jid = paths.job_id("youtube", hostile)
resolved = (tmp_path / "jobs" / jid).resolve()
assert resolved.parent == (tmp_path / "jobs").resolve()
class TestUrlJobId:
def test_is_stable(self) -> None:
url = "https://example.com/a.mp4"
assert paths.url_job_id(url) == paths.url_job_id(url)
def test_differs_between_urls(self) -> None:
assert paths.url_job_id("https://a.test/x") != paths.url_job_id("https://a.test/y")
def test_is_filesystem_safe(self) -> None:
jid = paths.url_job_id("https://a.test/x?q=1&r=2#frag")
assert jid.startswith("url-")
assert all(c.isalnum() or c in "-_." for c in jid)
class TestNormalizeUrl:
def test_strips_surrounding_whitespace(self) -> None:
assert paths.normalize_url(" https://a.test/x ") == "https://a.test/x"
def test_drops_the_fragment(self) -> None:
assert paths.normalize_url("https://a.test/x#t=30") == "https://a.test/x"
def test_keeps_the_query(self) -> None:
# YouTube identifies the video in the query string.
assert paths.normalize_url("https://y.test/watch?v=abc") == "https://y.test/watch?v=abc"
def test_lowercases_scheme_and_host_only(self) -> None:
assert paths.normalize_url("HTTPS://Example.COM/Path") == "https://example.com/Path"
def test_equivalent_urls_map_to_one_job(self) -> None:
a = paths.normalize_url("https://a.test/x#one")
b = paths.normalize_url(" https://a.test/x#two ")
assert paths.url_job_id(a) == paths.url_job_id(b)
class TestWorkspace:
def test_job_directory_lives_under_jobs(self, tmp_path: Path) -> None:
ws = paths.Workspace(tmp_path)
assert ws.job("youtube-x").root == tmp_path / "jobs" / "youtube-x"
def test_index_file_location(self, tmp_path: Path) -> None:
assert paths.Workspace(tmp_path).index_file == tmp_path / "index.json"
def test_ensure_creates_the_tree(self, tmp_path: Path) -> None:
job = paths.Workspace(tmp_path).job("youtube-x")
job.ensure()
for d in (job.media_dir, job.out_dir, job.logs_dir, job.tmp_dir):
assert d.is_dir()
def test_ensure_is_idempotent(self, tmp_path: Path) -> None:
job = paths.Workspace(tmp_path).job("youtube-x")
job.ensure()
job.ensure()
assert job.media_dir.is_dir()
def test_artifact_locations(self, tmp_path: Path) -> None:
job = paths.Workspace(tmp_path).job("youtube-x")
assert job.state_file.name == "state.json"
assert job.lock_file.name == ".lock"
assert job.events_file.name == "events.jsonl"
assert job.info_file.name == "info.json"
assert job.audio_file == job.media_dir / "audio.flac"
def test_tmp_is_inside_the_job_dir(self, tmp_path: Path) -> None:
# os.replace is only atomic within one filesystem.
job = paths.Workspace(tmp_path).job("youtube-x")
assert job.tmp_dir.parent == job.root
def test_transcript_paths_are_named_per_format(self, tmp_path: Path) -> None:
job = paths.Workspace(tmp_path).job("youtube-x")
assert job.transcript_file("srt") == job.out_dir / "transcript.srt"