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