Recently Written · git

subplz-web

git clone https://github.com/equwal/subplz-web

Log | Files | Refs


backend/runner.py (24893 bytes)

1 """Runs one job: stage inputs, drive the aligner, collect the artifacts.
2 
3 Nothing here knows which alignment backend is in use. The command to run, how
4 to read its progress and where it writes the subtitles all come from
5 `aligner.Aligner`, so swapping subplz out does not touch this file.
6 
7 The backend is always a subprocess. That boundary is also what lets the public
8 deployment move alignment onto separate GPU workers without code changes.
9 """
10 
11 from __future__ import annotations
12 
13 import json
14 import logging
15 import os
16 import re
17 import shutil
18 import subprocess
19 import time
20 from dataclasses import dataclass
21 from datetime import timezone
22 from pathlib import Path
23 
24 from . import billing, languages, render
25 from .aligner import AlignRequest, aligner
26 from .db import Artifact, Job, JobStatus, SessionLocal, utcnow
27 from .settings import settings
28 from .storage import storage
29 
30 log = logging.getLogger(__name__)
31 
32 # Audio containers subplz/ffmpeg handle. Checked at upload time.
33 AUDIO_SUFFIXES = {
34     ".m4b", ".m4a", ".mp3", ".opus", ".ogg", ".oga", ".flac", ".wav",
35     ".aac", ".wma", ".mka", ".mkv", ".mp4", ".webm", ".avi", ".mov",
36 }
37 # Formats subplz reads directly, plus the ones convert.py turns into epub on
38 # the way in (fb2/mobi/azw3).
39 TEXT_SUFFIXES = {
40     ".epub", ".txt", ".srt", ".vtt", ".ass",
41     ".fb2", ".fb2.zip", ".mobi", ".azw", ".azw3", ".prc",
42 }
43 
44 
45 def text_suffix(name: str) -> str:
46     """Suffix used when staging a text file, keeping the .fb2.zip double."""
47     lowered = name.lower()
48     if lowered.endswith(".fb2.zip"):
49         return ".fb2.zip"
50     return Path(lowered).suffix
51 
52 # Fixed stem for staged inputs: keeps non-ASCII filenames out of the subprocess
53 # command line entirely, and makes the output path deterministic.
54 STAGE_STEM = "source"
55 
56 
57 class JobFailed(RuntimeError):
58     pass
59 
60 
61 @dataclass
62 class Paths:
63     root: Path
64     inp: Path
65     out: Path
66 
67     @classmethod
68     def for_job(cls, job_id: str) -> "Paths":
69         root = settings.data_dir / "work" / job_id
70         return cls(root=root, inp=root / "input", out=root / "out")
71 
72     def create(self) -> None:
73         self.inp.mkdir(parents=True, exist_ok=True)
74         self.out.mkdir(parents=True, exist_ok=True)
75 
76 
77 def input_prefix(job_id: str) -> str:
78     """Storage prefix mirroring a job's staged `input/` tree.
79 
80     Only used when the queue is external: the worker that runs the job is
81     probably not the machine that received the upload.
82     """
83     return f"{job_id}/input"
84 
85 
86 def ensure_inputs(job_id: str) -> None:
87     """Make sure the staged inputs are on this machine's disk."""
88     paths = Paths.for_job(job_id)
89     paths.create()
90     if any(p.is_file() for p in paths.inp.rglob("*")):
91         return  # this process received the upload
92 
93     prefix = input_prefix(job_id)
94     keys = storage.list_prefix(prefix)
95     if not keys:
96         raise JobFailed(
97             "the uploaded files for this job are no longer available. "
98             "Upload them again."
99         )
100     for key in keys:
101         storage.fetch_to(key, paths.inp / key[len(prefix) + 1 :])
102 
103 
104 def staged_audio_path(job_id: str, original_name: str) -> Path:
105     return Paths.for_job(job_id).inp / f"{STAGE_STEM}{Path(original_name).suffix.lower()}"
106 
107 
108 def staged_text_path(job_id: str, original_name: str) -> Path:
109     return Paths.for_job(job_id).inp / f"{STAGE_STEM}{text_suffix(original_name)}"
110 
111 
112 def staged_cover_path(job_id: str, original_name: str) -> Path:
113     """Where a user-supplied cover image is staged.
114 
115     Kept out of the `input/` root alongside audio and text so the aligner's
116     directory scan cannot mistake a picture for content.
117     """
118     inp = Paths.for_job(job_id).inp
119     return inp / "cover" / f"cover{Path(original_name).suffix.lower()}"
120 
121 
122 def find_cover(job_id: str) -> Path | None:
123     folder = Paths.for_job(job_id).inp / "cover"
124     if not folder.is_dir():
125         return None
126     return next((p for p in sorted(folder.iterdir()) if p.is_file()), None)
127 
128 
129 def staged_part_path(job_id: str, index: int, original_name: str) -> Path:
130     """Where chapter file `index` (1-based) of a multi-part audiobook is staged.
131 
132     Parts get plain numeric names so the ffmpeg concat list never needs quoting
133     and playback order is unambiguous.
134     """
135     inp = Paths.for_job(job_id).inp
136     return inp / "parts" / f"{index:04d}{Path(original_name).suffix.lower()}"
137 
138 
139 # Tolerate damaged input rather than aborting the whole book. ffmpeg's default
140 # max_error_rate (0.667) kills a run over one bad chapter.
141 _FFMPEG_TOLERANT = [
142     "-max_error_rate", "1.0",
143     "-err_detect", "ignore_err",
144     "-fflags", "+discardcorrupt",
145 ]
146 
147 
148 def merge_parts(job_id: str, parts: list[Path]) -> Path:
149     """Join a per-chapter audiobook into one file, one chapter per part.
150 
151     Keeping chapter marks matters twice over: subplz processes an m4b chapter by
152     chapter, and the runner reads chapter completions to drive progress.
153     """
154     paths = Paths.for_job(job_id)
155     dest = paths.inp / f"{STAGE_STEM}.m4b"
156     if dest.exists():
157         return dest
158 
159     durations: list[float] = []
160     for p in parts:
161         d = probe_duration(p)
162         if d is None:
163             raise JobFailed(f"could not read the duration of part {p.name}")
164         durations.append(d)
165 
166     scratch = paths.root / "merge"
167     scratch.mkdir(parents=True, exist_ok=True)
168 
169     listing = scratch / "parts.txt"
170     # Absolute paths: ffmpeg resolves a relative entry against the directory of
171     # the list file, not the working directory.
172     listing.write_text(
173         "".join(f"file '{p.resolve().as_posix()}'\n" for p in parts),
174         encoding="utf-8",
175     )
176 
177     # Chapter marks at the part boundaries, in milliseconds.
178     meta = [";FFMETADATA1"]
179     start_ms = 0
180     for i, seconds in enumerate(durations, start=1):
181         end_ms = start_ms + int(round(seconds * 1000))
182         meta += [
183             "[CHAPTER]",
184             "TIMEBASE=1/1000",
185             f"START={start_ms}",
186             f"END={end_ms}",
187             f"title=Part {i:02d}",
188         ]
189         start_ms = end_ms
190     metadata = scratch / "chapters.txt"
191     metadata.write_text("\n".join(meta) + "\n", encoding="utf-8")
192 
193     cmd = [
194         "ffmpeg", "-hide_banner", "-v", "error", "-y",
195         *_FFMPEG_TOLERANT,
196         "-f", "concat", "-safe", "0", "-i", str(listing),
197         "-i", str(metadata),
198         "-map", "0:a:0", "-map_metadata", "1", "-map_chapters", "1",
199         # Drop cover art and any stray tracks the parts carry, so the only
200         # media stream is the audio. (MP4 still writes a bin_data chapter
201         # track; that one is how the container stores chapters and must stay.)
202         "-vn", "-sn",
203         # Keep listenable quality: the aligner downsamples to 16 kHz mono itself
204         # when it reads the file, but the rendered video uses this same audio,
205         # and 16 kHz mono would sound awful on YouTube.
206         "-c:a", "aac", "-b:a", "128k",
207         str(dest),
208     ]
209     proc = subprocess.run(cmd, capture_output=True, timeout=settings.job_timeout_seconds)
210     if proc.returncode != 0 or not dest.exists():
211         err = proc.stderr.decode("utf-8", errors="replace").strip().splitlines()
212         tail = " | ".join(err[-3:]) if err else "no stderr"
213         raise JobFailed(f"could not join the {len(parts)} audio parts: {tail}")
214 
215     shutil.rmtree(scratch, ignore_errors=True)
216     return dest
217 
218 
219 def normalize_for_alignment(job_id: str, source: Path) -> Path:
220     """Produce the 16 kHz mono copy the aligner gets fed.
221 
222     Two jobs at once:
223 
224     * Format. The acoustic model consumes 16 kHz mono no matter what, and
225       handing subplz a 44.1 kHz stereo file has crashed ctranslate2 here
226       (integer divide by zero, part-way through a chapter). Doing the
227       conversion ourselves, once, keeps the aligner on the input that works.
228     * Repair. The error-tolerant flags drop corrupt frames rather than letting
229       ffmpeg abort the run, which is what a damaged chapter would otherwise do.
230 
231     The video keeps using `source`, which stays at listenable quality.
232     """
233     dest = Paths.for_job(job_id).inp / f"{STAGE_STEM}.align.m4a"
234     if dest.exists():
235         return dest
236 
237     cmd = [
238         "ffmpeg", "-hide_banner", "-v", "error", "-y",
239         *_FFMPEG_TOLERANT,
240         "-i", str(source),
241         "-map", "0:a:0", "-map_chapters", "0",
242         "-vn", "-sn",
243         "-c:a", "aac", "-b:a", "64k", "-ar", "16000", "-ac", "1",
244         str(dest),
245     ]
246     proc = subprocess.run(
247         cmd, capture_output=True, timeout=settings.job_timeout_seconds
248     )
249     if proc.returncode != 0 or not dest.exists():
250         err = proc.stderr.decode("utf-8", errors="replace").strip().splitlines()
251         tail = " | ".join(err[-3:]) if err else "no stderr"
252         raise JobFailed(f"could not prepare the audio for alignment: {tail}")
253     return dest
254 
255 
256 def probe_chapters(path: Path) -> int:
257     """Chapter count, or 1 for a flat file.
258 
259     subplz processes chaptered m4b files one chapter at a time, so this is what
260     turns a stream of per-chapter progress bars into an overall percentage.
261     """
262     try:
263         out = subprocess.run(
264             ["ffprobe", "-v", "error", "-show_chapters", "-print_format", "json",
265              str(path)],
266             capture_output=True, text=True, encoding="utf-8", errors="replace",
267             timeout=120, check=True,
268         )
269         return max(1, len(json.loads(out.stdout or "{}").get("chapters", [])))
270     except (subprocess.SubprocessError, ValueError, OSError):
271         return 1
272 
273 
274 def probe_duration(path: Path) -> float | None:
275     """Audio length in seconds, via ffprobe. Used for progress and metadata."""
276     try:
277         out = subprocess.run(
278             [
279                 "ffprobe", "-v", "error",
280                 "-show_entries", "format=duration",
281                 "-of", "default=noprint_wrappers=1:nokey=1",
282                 str(path),
283             ],
284             capture_output=True, text=True, timeout=120, check=True,
285         )
286         return float(out.stdout.strip())
287     except (subprocess.SubprocessError, ValueError, FileNotFoundError, OSError):
288         return None
289 
290 
291 def _staged(inp: Path, suffixes: set[str]) -> Path:
292     """The one staged input with a suffix in `suffixes`."""
293     for p in sorted(inp.iterdir()):
294         if p.is_file() and p.suffix.lower() in suffixes:
295             return p
296     raise JobFailed(f"no staged input found with a supported extension in {inp.name}")
297 
298 
299 def build_request(job: Job, paths: Paths, audio: Path, chapters: int) -> AlignRequest:
300     """Describe the job in backend-neutral terms."""
301     languages.require(job.language)  # reject an unsupported code before we run
302     return AlignRequest(
303         audio=audio,
304         text=_staged(paths.inp, TEXT_SUFFIXES),
305         out_dir=paths.out,
306         language=job.language,
307         model=job.model,
308         device=settings.device,
309         threads=settings.threads or max(1, (os.cpu_count() or 4) - 1),
310         chapters=chapters,
311     )
312 
313 
314 def _set(job_id: str, **fields) -> None:
315     with SessionLocal() as s:
316         job = s.get(Job, job_id)
317         if job is None:
318             return
319         for k, v in fields.items():
320             setattr(job, k, v)
321         s.commit()
322 
323 
324 def _is_canceled(job_id: str) -> bool:
325     with SessionLocal() as s:
326         job = s.get(Job, job_id)
327         return job is not None and job.status == JobStatus.canceled
328 
329 
330 def run_job(job_id: str) -> None:
331     """Execute one job end to end. Never raises; failures land on the Job row."""
332     with SessionLocal() as s:
333         job = s.get(Job, job_id)
334         if job is None:
335             return
336         if job.status == JobStatus.canceled:
337             return
338         audio_name, text_name = job.audio_filename, job.text_filename
339         parts_count = job.audio_parts or 1
340 
341     paths = Paths.for_job(job_id)
342     log_path = paths.root / "subplz.log"
343     _set(job_id, status=JobStatus.running, started_at=utcnow(),
344          stage="Preparing", progress=0.01)
345 
346     try:
347         # No-op when this process received the upload; pulls from shared storage
348         # when the job was enqueued on another machine.
349         ensure_inputs(job_id)
350 
351         if parts_count > 1:
352             _set(job_id, stage=f"Joining {parts_count} audio parts", progress=0.02)
353             parts = sorted(p for p in (paths.inp / "parts").iterdir() if p.is_file())
354             if len(parts) != parts_count:
355                 raise JobFailed(
356                     f"expected {parts_count} audio parts but found {len(parts)}"
357                 )
358             audio_in = merge_parts(job_id, parts)
359         else:
360             audio_in = staged_audio_path(job_id, audio_name)
361 
362         if not audio_in.exists():
363             raise JobFailed(f"staged audio missing: {audio_in.name}")
364 
365         duration = probe_duration(audio_in)
366         chapters = probe_chapters(audio_in)
367 
368         # The aligner gets a normalised 16 kHz mono copy; the video keeps the
369         # full-quality one.
370         _set(job_id, stage="Preparing audio", progress=0.03)
371         align_audio = normalize_for_alignment(job_id, audio_in)
372 
373         _set(job_id, audio_duration_seconds=duration,
374              stage=f"Starting {aligner.name}", progress=0.04)
375 
376         with SessionLocal() as s:
377             job = s.get(Job, job_id)
378             request = build_request(job, paths, align_audio, chapters)
379 
380         returncode = _stream_aligner(job_id, request, log_path)
381 
382         if _is_canceled(job_id):
383             _cleanup_inputs(paths)
384             return
385 
386         if returncode != 0:
387             raise JobFailed(
388                 f"{aligner.name} exited with code {returncode}. "
389                 f"See the run log for details."
390             )
391 
392         _set(job_id, stage="Collecting output", progress=0.92)
393         _collect_artifacts(job_id, paths, log_path, duration, request, audio_in)
394 
395         _set(job_id, status=JobStatus.succeeded, stage="Done", progress=1.0,
396              finished_at=utcnow())
397         # Only on success: a failed job keeps its inputs so it can be retried
398         # without re-uploading hundreds of megabytes.
399         _cleanup_inputs(paths)
400 
401     except Exception as exc:  # noqa: BLE001 - surfaced to the user on the Job row
402         if _is_canceled(job_id):
403             _cleanup_inputs(paths)
404             return
405         # Keep the log even on failure: it is the only way to debug an alignment.
406         try:
407             _store_log(job_id, log_path)
408         except Exception:
409             pass
410         _set(job_id, status=JobStatus.failed, error=str(exc)[:4000],
411              stage="Failed", finished_at=utcnow())
412         # Our failure, not theirs: a credit spent on this run goes back.
413         try:
414             billing.refund_job(job_id)
415         except Exception:  # noqa: BLE001
416             log.exception("job %s: could not refund after failure", job_id)
417 
418 
419 def _stream_aligner(job_id: str, request: AlignRequest, log_path: Path) -> int:
420     """Run the aligner, mirroring its output to a log file and to job progress.
421 
422     The backend decides what its output means; this only moves bytes and
423     watches for cancellation and the time limit.
424     """
425     cmd = aligner.build_command(request)
426     reader = aligner.progress_reader(request)
427 
428     log_path.parent.mkdir(parents=True, exist_ok=True)
429     deadline = time.monotonic() + settings.job_timeout_seconds
430 
431     with log_path.open("w", encoding="utf-8", errors="replace") as log:
432         log.write("$ " + " ".join(cmd) + "\n\n")
433         log.flush()
434 
435         proc = subprocess.Popen(
436             cmd,
437             stdout=subprocess.PIPE,
438             stderr=subprocess.STDOUT,
439             text=True,
440             encoding="utf-8",
441             errors="replace",
442             bufsize=1,
443             env=aligner.environment(),
444             cwd=str(log_path.parent),
445         )
446 
447         try:
448             assert proc.stdout is not None
449             for raw in proc.stdout:
450                 log.write(raw)
451                 log.flush()
452 
453                 if time.monotonic() > deadline:
454                     proc.kill()
455                     raise JobFailed(
456                         f"job exceeded the {settings.job_timeout_seconds}s time limit"
457                     )
458 
459                 if _is_canceled(job_id):
460                     proc.kill()
461                     return proc.wait()
462 
463                 line = raw.strip()
464                 if not line:
465                     continue
466 
467                 update = reader.feed(line)
468                 if update is None:
469                     continue
470                 fields: dict = {"stage": update.stage}
471                 if update.fraction is not None:
472                     fields["progress"] = update.fraction
473                 _set(job_id, **fields)
474         finally:
475             if proc.poll() is None:
476                 proc.kill()
477             proc.wait()
478 
479         return proc.returncode
480 
481 
482 def _collect_artifacts(job_id: str, paths: Paths, log_path: Path,
483                        duration: float | None, request: AlignRequest,
484                        video_audio: Path) -> None:
485     with SessionLocal() as s:
486         job = s.get(Job, job_id)
487         if job is None:
488             raise JobFailed("job vanished while collecting output")
489         language, model = job.language, job.model
490         audio_name, text_name = job.audio_filename, job.text_filename
491         splitter = job.splitter
492         parts_count = job.audio_parts or 1
493 
494     produced = aligner.locate_output(request)
495     if produced is None or not produced.exists():
496         # The backend may exit 0 and still have failed, so ask it why before
497         # falling back to guessing at the content.
498         reason = aligner.failure_reason(request)
499         if reason:
500             raise JobFailed(f"{aligner.name} failed: {reason}")
501         raise JobFailed(
502             f"{aligner.name} finished but produced no subtitle file. The audio "
503             "and text may not match, or the language may be wrong for this book."
504         )
505 
506     # Give the download a filename the user will recognise. A per-chapter
507     # audiobook has no single audio name worth using, so the book names it.
508     stem = Path(text_name if parts_count > 1 else audio_name).stem
509     download_name = f"{stem}.{language}{aligner.output_suffix}"
510 
511     cues, first_cue, last_cue = _summarize_srt(produced)
512 
513     srt_key = f"{job_id}/{download_name}"
514     size = storage.put_file(srt_key, produced)
515 
516     # Three deliverables, because they serve three different jobs:
517     #   .srt  - the timing file on its own (HoshiReader whispersync)
518     #   .mp4  - clean video, no subtitle track, for YouTube (captions are
519     #           uploaded separately there)
520     #   .mkv  - the same video with the subtitles embedded, for local players
521     # The mkv is a stream copy of the mp4, so the second file is nearly free.
522     video_name = video_key = None
523     embed_name = embed_key = None
524     video_size = embed_size = 0
525 
526     if settings.render_video:
527         try:
528             _set(job_id, stage="Rendering video", progress=0.94)
529             scratch = paths.root / "video"
530             # A cover the user supplied wins over the one inside the epub.
531             cover = find_cover(job_id) or render.extract_cover(request.text, scratch)
532             canvas = render.build_canvas(cover, scratch / "canvas.png")
533 
534             video_name = f"{stem}.{language}.mp4"
535             out_video = scratch / video_name
536             render.render_video(
537                 audio=video_audio, canvas=canvas,
538                 dest=out_video, duration=duration,
539             )
540             video_key = f"{job_id}/{video_name}"
541             video_size = storage.put_file(video_key, out_video)
542 
543             try:
544                 _set(job_id, stage="Embedding subtitles", progress=0.97)
545                 embed_name = f"{stem}.{language}.mkv"
546                 out_embed = scratch / embed_name
547                 render.mux_subtitles(out_video, produced, out_embed)
548                 embed_key = f"{job_id}/{embed_name}"
549                 embed_size = storage.put_file(embed_key, out_embed)
550             except Exception as exc:  # noqa: BLE001
551                 embed_name = embed_key = None
552                 log.warning("job %s: subtitle embed failed: %s", job_id, exc)
553 
554         except Exception as exc:  # noqa: BLE001 - never fail a job over the video
555             # The subtitles are the product, so a failed render is not fatal -
556             # but it must not be silent either, or it looks like it never ran.
557             video_name = video_key = None
558             log.warning("job %s: video render failed: %s", job_id, exc,
559                         exc_info=True)
560 
561     metadata = {
562         "job_id": job_id,
563         "created_at": utcnow().astimezone(timezone.utc).isoformat(),
564         "source": {
565             "audio_filename": audio_name,
566             "audio_parts": parts_count,
567             "text_filename": text_name,
568             "audio_duration_seconds": duration,
569         },
570         "alignment": {
571             "backend": aligner.name,
572             "mode": "forced alignment against supplied text (nothing transcribed from scratch)",
573             "model": model,
574             "device": settings.device,
575             "language": language,
576             "language_name": (languages.get(language).name if languages.get(language) else language),
577             "sentence_splitter": splitter,
578         },
579         "output": {
580             "filename": download_name,
581             "format": "srt",
582             "cue_count": cues,
583             "first_cue_start": first_cue,
584             "last_cue_end": last_cue,
585             "size_bytes": size,
586         },
587         "video": (
588             {
589                 "filename": video_name,
590                 "container": "mp4",
591                 "subtitles": "none - upload the .srt to YouTube separately",
592                 "for": "youtube",
593                 "size_bytes": video_size,
594             }
595             if video_name
596             else None
597         ),
598         "video_embedded": (
599             {
600                 "filename": embed_name,
601                 "container": "mkv",
602                 "subtitles": "embedded SRT track, enabled by default",
603                 "for": "local playback (MPV, VLC, Jellyfin)",
604                 "size_bytes": embed_size,
605             }
606             if embed_name
607             else None
608         ),
609     }
610     meta_path = paths.root / "metadata.json"
611     meta_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2),
612                          encoding="utf-8")
613     meta_key = f"{job_id}/metadata.json"
614     meta_size = storage.put_file(meta_key, meta_path)
615 
616     rows = [
617         Artifact(job_id=job_id, kind="srt", filename=download_name,
618                  storage_key=srt_key, size_bytes=size),
619         Artifact(job_id=job_id, kind="metadata", filename="metadata.json",
620                  storage_key=meta_key, size_bytes=meta_size),
621     ]
622     if video_name and video_key:
623         rows.append(
624             Artifact(job_id=job_id, kind="video", filename=video_name,
625                      storage_key=video_key, size_bytes=video_size)
626         )
627     if embed_name and embed_key:
628         rows.append(
629             Artifact(job_id=job_id, kind="video_embedded", filename=embed_name,
630                      storage_key=embed_key, size_bytes=embed_size)
631         )
632     with SessionLocal() as s:
633         for r in rows:
634             s.add(r)
635         s.commit()
636 
637     _store_log(job_id, log_path)
638 
639 
640 def _store_log(job_id: str, log_path: Path) -> None:
641     if not log_path.exists():
642         return
643     with SessionLocal() as s:
644         existing = (
645             s.query(Artifact)
646             .filter(Artifact.job_id == job_id, Artifact.kind == "log")
647             .first()
648         )
649         if existing is not None:
650             return
651     key = f"{job_id}/subplz.log"
652     size = storage.put_file(key, log_path)
653     with SessionLocal() as s:
654         s.add(Artifact(job_id=job_id, kind="log", filename="subplz.log",
655                        storage_key=key, size_bytes=size))
656         s.commit()
657 
658 
659 _TIME = re.compile(r"(\d{2}):(\d{2}):(\d{2}),(\d{3})")
660 
661 
662 def _summarize_srt(path: Path) -> tuple[int, float | None, float | None]:
663     """Cue count and span - cheap sanity signal that the alignment covered the book."""
664     def to_seconds(m: re.Match) -> float:
665         h, mi, s, ms = (int(g) for g in m.groups())
666         return h * 3600 + mi * 60 + s + ms / 1000
667 
668     text = path.read_text(encoding="utf-8", errors="replace")
669     stamps = list(_TIME.finditer(text))
670     count = text.count(" --> ")
671     if not stamps:
672         return count, None, None
673     return count, to_seconds(stamps[0]), to_seconds(stamps[-1])
674 
675 
676 def _cleanup_inputs(paths: Paths) -> None:
677     """Drop the uploaded media once the run is over; keep artifacts and the log."""
678     shutil.rmtree(paths.inp, ignore_errors=True)
679     shutil.rmtree(paths.out, ignore_errors=True)
680     shutil.rmtree(paths.root / "video", ignore_errors=True)
681     # Also drop the shared-storage copy made for external workers.
682     try:
683         storage.delete_prefix(f"{paths.root.name}/input")
684     except Exception:
685         pass