From 2cfd7ad087e7fd2592a73f6ab94b2204dc151406 Mon Sep 17 00:00:00 2001 From: Jason Ross Date: Sun, 13 Sep 2026 14:24:24 -0500 Subject: [PATCH] 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) --- pyproject.toml | 11 +- src/ccn_transcribe/jobs/plan.py | 63 ++++++++++ src/ccn_transcribe/jobs/state.py | 166 +++++++++++++++++++++++++ src/ccn_transcribe/jobs/store.py | 118 ++++++++++++++++++ src/ccn_transcribe/jobs/verify.py | 113 +++++++++++++++++ src/ccn_transcribe/paths.py | 115 +++++++++++++++++ tests/test_jobs_plan.py | 130 +++++++++++++++++++ tests/test_jobs_state.py | 137 ++++++++++++++++++++ tests/test_jobs_store.py | 144 +++++++++++++++++++++ tests/test_jobs_verify.py | 199 ++++++++++++++++++++++++++++++ tests/test_paths.py | 124 +++++++++++++++++++ 11 files changed, 1318 insertions(+), 2 deletions(-) create mode 100644 src/ccn_transcribe/jobs/plan.py create mode 100644 src/ccn_transcribe/jobs/state.py create mode 100644 src/ccn_transcribe/jobs/store.py create mode 100644 src/ccn_transcribe/jobs/verify.py create mode 100644 src/ccn_transcribe/paths.py create mode 100644 tests/test_jobs_plan.py create mode 100644 tests/test_jobs_state.py create mode 100644 tests/test_jobs_store.py create mode 100644 tests/test_jobs_verify.py create mode 100644 tests/test_paths.py diff --git a/pyproject.toml b/pyproject.toml index 1799f82..b5a51ea 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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"] diff --git a/src/ccn_transcribe/jobs/plan.py b/src/ccn_transcribe/jobs/plan.py new file mode 100644 index 0000000..3d2168b --- /dev/null +++ b/src/ccn_transcribe/jobs/plan.py @@ -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]) diff --git a/src/ccn_transcribe/jobs/state.py b/src/ccn_transcribe/jobs/state.py new file mode 100644 index 0000000..93b270a --- /dev/null +++ b/src/ccn_transcribe/jobs/state.py @@ -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", {})), + ) diff --git a/src/ccn_transcribe/jobs/store.py b/src/ccn_transcribe/jobs/store.py new file mode 100644 index 0000000..8f6e44b --- /dev/null +++ b/src/ccn_transcribe/jobs/store.py @@ -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() diff --git a/src/ccn_transcribe/jobs/verify.py b/src/ccn_transcribe/jobs/verify.py new file mode 100644 index 0000000..a15e7e9 --- /dev/null +++ b/src/ccn_transcribe/jobs/verify.py @@ -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 diff --git a/src/ccn_transcribe/paths.py b/src/ccn_transcribe/paths.py new file mode 100644 index 0000000..3f8dc76 --- /dev/null +++ b/src/ccn_transcribe/paths.py @@ -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) diff --git a/tests/test_jobs_plan.py b/tests/test_jobs_plan.py new file mode 100644 index 0000000..61a2cf2 --- /dev/null +++ b/tests/test_jobs_plan.py @@ -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) diff --git a/tests/test_jobs_state.py b/tests/test_jobs_state.py new file mode 100644 index 0000000..880660a --- /dev/null +++ b/tests/test_jobs_state.py @@ -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 diff --git a/tests/test_jobs_store.py b/tests/test_jobs_store.py new file mode 100644 index 0000000..ad1a0be --- /dev/null +++ b/tests/test_jobs_store.py @@ -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"] diff --git a/tests/test_jobs_verify.py b/tests/test_jobs_verify.py new file mode 100644 index 0000000..bae7659 --- /dev/null +++ b/tests/test_jobs_verify.py @@ -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 diff --git a/tests/test_paths.py b/tests/test_paths.py new file mode 100644 index 0000000..2c4874b --- /dev/null +++ b/tests/test_paths.py @@ -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"