backend/api.py (29878 bytes)
1 """HTTP API. 2 3 Flow: drop files -> POST /api/uploads (stages, pairs, detects language) -> 4 POST /api/jobs/{id}/start (entitlement check, enqueue) -> poll GET /api/jobs/{id} 5 -> download from /api/jobs/{id}/files/{kind}. 6 7 Upload and start are separate so a wrong language guess costs a click rather 8 than a re-upload and a wasted multi-hour run. 9 10 Around that: /api/auth/* (email sign-in links) and /api/billing/* (Stripe 11 checkout, its return trip and its webhook). 12 """ 13 14 from __future__ import annotations 15 16 import shutil 17 from pathlib import Path 18 from typing import Annotated, Literal 19 20 from fastapi import ( 21 APIRouter, Cookie, Depends, File, HTTPException, Request, Response, 22 UploadFile, 23 ) 24 from fastapi.responses import FileResponse, RedirectResponse 25 from pydantic import BaseModel, Field 26 from sqlalchemy.orm import Session 27 28 from . import ( 29 accounts, auth, billing, convert, detect, languages, mailer, matching, 30 payments, pricing, 31 ) 32 from .aligner import aligner 33 from .db import Account, Artifact, Job, JobStatus, SessionLocal, new_id, utcnow 34 from .queue import queue 35 from .runner import ( 36 Paths, input_prefix, probe_duration, staged_audio_path, staged_cover_path, 37 staged_part_path, staged_text_path, 38 ) 39 from .settings import settings 40 from .storage import LocalStorage, storage 41 42 router = APIRouter(prefix="/api") 43 44 DEVICE_COOKIE = "subplz_device" 45 _CHUNK = 4 * 1024 * 1024 46 47 48 # -------------------------------------------------------------------------- 49 # session / account 50 # -------------------------------------------------------------------------- 51 52 def get_session(): 53 with SessionLocal() as s: 54 yield s 55 56 57 def get_account( 58 response: Response, 59 session: Annotated[Session, Depends(get_session)], 60 subplz_device: Annotated[str | None, Cookie()] = None, 61 ) -> Account: 62 """Identify the caller. 63 64 A cookie is the whole identity check for the free tier, on purpose. Anyone 65 who clears it gets another free book, and that is an accepted cost: a 66 signup wall does not belong in front of a tool whose pitch is "drop two 67 files in". The real protection against abuse is capacity - the queue and 68 per-worker limits - not identity. 69 70 Signing in does not replace the cookie, it re-points it: the cookie then 71 names the signed-in account, on every device that has signed in. 72 """ 73 account = None 74 if subplz_device: 75 account = ( 76 session.query(Account) 77 .filter(Account.device_token == subplz_device) 78 .first() 79 ) 80 81 if account is None: 82 account = accounts.new_account(session) 83 session.commit() 84 _set_identity(response, account) 85 return account 86 87 # This device's anonymous row was folded into a real account while no 88 # browser was attached (a payment landing by webhook). Follow it. 89 survivor = accounts.resolve(session, account) 90 if survivor.id != account.id: 91 _set_identity(response, survivor) 92 return survivor 93 94 95 def _set_identity(response: Response, account: Account) -> None: 96 response.set_cookie( 97 DEVICE_COOKIE, account.device_token, 98 max_age=60 * 60 * 24 * 365, httponly=True, samesite="lax", 99 secure=settings.cookie_secure, 100 ) 101 102 103 # -------------------------------------------------------------------------- 104 # schemas 105 # -------------------------------------------------------------------------- 106 107 class LanguageOut(BaseModel): 108 code: str 109 name: str 110 splitter: str 111 # Anything the active backend wants the user to know about this language. 112 note: str | None = None 113 114 115 class DetectionOut(BaseModel): 116 code: str | None 117 name: str | None 118 confidence: float 119 supported: bool 120 121 122 class ArtifactOut(BaseModel): 123 kind: str 124 filename: str 125 size_bytes: int 126 url: str 127 128 129 class JobOut(BaseModel): 130 id: str 131 status: str 132 stage: str 133 progress: float 134 language: str 135 language_name: str 136 splitter: str 137 model: str 138 audio_filename: str 139 audio_parts: int 140 cover_filename: str | None 141 text_filename: str 142 audio_bytes: int 143 audio_duration_seconds: float | None 144 error: str | None 145 created_at: str 146 local: bool = False 147 tier: str = billing.FREE 148 artifacts: list[ArtifactOut] = Field(default_factory=list) 149 150 151 class UploadOut(BaseModel): 152 job: JobOut 153 detected: DetectionOut 154 # Set when an image was dropped but the visitor is not signed in. 155 cover_requires_sign_in: bool = False 156 # How well the book scores against a sample of the audio, using the 157 # backend's own rule. None when the check is disabled. 158 match: dict | None = None 159 160 161 class StartIn(BaseModel): 162 language: str | None = None 163 model: str | None = None 164 165 166 class AccountOut(BaseModel): 167 id: str 168 signed_in: bool 169 email: str | None 170 billing_enabled: bool 171 # Whether the server can actually take money / send sign-in email yet. 172 payments_available: bool 173 email_sign_in_available: bool 174 free_tier_summary: str 175 # The cloud tier: conversions on this server's hardware. 176 cloud_available: bool 177 credits: int 178 subscribed: bool 179 subscription_ends: str | None 180 cloud_allowed: bool 181 queue_depth: int 182 183 184 class LocalJobIn(BaseModel): 185 """A conversion about to run in the visitor's browser.""" 186 audio_filename: str = Field(max_length=512) 187 audio_parts: int = Field(default=1, ge=1, le=2000) 188 audio_bytes: int = Field(default=0, ge=0) 189 audio_duration_seconds: float | None = None 190 text_filename: str = Field(max_length=512) 191 language: str = Field(max_length=16) 192 193 194 class LocalFinishIn(BaseModel): 195 # A twenty-hour book is a few megabytes of subtitles. 196 srt: str = Field(max_length=8_000_000) 197 filename: str = Field(max_length=512) 198 metadata: dict = Field(default_factory=dict) 199 200 201 class LocalFailIn(BaseModel): 202 error: str = Field(default="", max_length=4000) 203 204 205 class EmailIn(BaseModel): 206 email: str 207 208 209 class TokenIn(BaseModel): 210 token: str 211 212 213 class CheckoutIn(BaseModel): 214 plan_id: str 215 216 217 def _job_out(job: Job, arts: list[Artifact]) -> JobOut: 218 lang = languages.get(job.language) 219 return JobOut( 220 id=job.id, 221 status=job.status.value, 222 stage=job.stage, 223 progress=round(job.progress, 4), 224 language=job.language, 225 language_name=lang.name if lang else job.language, 226 splitter=job.splitter, 227 model=job.model, 228 audio_filename=job.audio_filename, 229 audio_parts=job.audio_parts or 1, 230 cover_filename=job.cover_filename, 231 text_filename=job.text_filename, 232 audio_bytes=job.audio_bytes, 233 audio_duration_seconds=job.audio_duration_seconds, 234 error=job.error, 235 created_at=job.created_at.isoformat(), 236 local=bool(job.local), 237 tier=job.tier or billing.FREE, 238 artifacts=[ 239 ArtifactOut( 240 kind=a.kind, filename=a.filename, size_bytes=a.size_bytes, 241 url=f"/api/jobs/{job.id}/files/{a.kind}", 242 ) 243 for a in arts 244 ], 245 ) 246 247 248 def _load(session: Session, account: Account, job_id: str) -> Job: 249 job = session.get(Job, job_id) 250 if job is None or job.account_id != account.id: 251 raise HTTPException(404, "Job not found") 252 return job 253 254 255 def _artifacts(session: Session, job_id: str) -> list[Artifact]: 256 return ( 257 session.query(Artifact) 258 .filter(Artifact.job_id == job_id) 259 .order_by(Artifact.kind) 260 .all() 261 ) 262 263 264 # -------------------------------------------------------------------------- 265 # routes 266 # -------------------------------------------------------------------------- 267 268 @router.get("/languages", response_model=list[LanguageOut]) 269 def list_languages(): 270 return [ 271 LanguageOut( 272 code=l.code, name=l.name, splitter=l.splitter, 273 note=aligner.language_note(l.code), 274 ) 275 for l in languages.all_languages() 276 ] 277 278 279 @router.get("/account", response_model=AccountOut) 280 def get_account_info( 281 account: Annotated[Account, Depends(get_account)], 282 session: Annotated[Session, Depends(get_session)], 283 ): 284 return _account_out(session, account) 285 286 287 def _account_out(session: Session, account: Account) -> AccountOut: 288 ent = billing.check(account) 289 return AccountOut( 290 id=account.id, 291 signed_in=account.signed_in, 292 email=account.email, 293 billing_enabled=settings.billing_enabled, 294 payments_available=settings.payments_configured, 295 email_sign_in_available=settings.sign_in_available, 296 free_tier_summary=pricing.free_tier_summary(), 297 credits=ent.credits, 298 subscribed=ent.subscribed, 299 subscription_ends=( 300 ent.subscription_ends.isoformat() if ent.subscription_ends else None 301 ), 302 cloud_available=settings.cloud_enabled, 303 cloud_allowed=ent.cloud_allowed or not settings.billing_enabled, 304 queue_depth=queue.depth(), 305 ) 306 307 308 @router.get("/pricing") 309 def get_pricing(): 310 """The catalogue, for the paywall. Static - no account needed.""" 311 return { 312 "free_tier": pricing.free_tier_summary(), 313 "billing_enabled": settings.billing_enabled, 314 "payments_available": settings.payments_configured, 315 "tiers": pricing.TIER_OUTPUTS, 316 "contact_email": settings.contact_email, 317 "input_retention_hours": settings.input_retention_hours, 318 "artifact_retention_days": settings.artifact_retention_days, 319 "plans": pricing.as_dicts(), 320 } 321 322 323 # -------------------------------------------------------------------------- 324 # sign-in 325 # -------------------------------------------------------------------------- 326 327 @router.post("/auth/request") 328 def request_sign_in( 329 body: EmailIn, 330 account: Annotated[Account, Depends(get_account)], 331 session: Annotated[Session, Depends(get_session)], 332 ): 333 """Mail a one-time sign-in link.""" 334 try: 335 email = accounts.normalize_email(body.email) 336 except accounts.InvalidEmail as exc: 337 raise HTTPException(400, str(exc)) from exc 338 339 # A server that cannot send mail has no safe way to do this: the only 340 # fallback is showing the link, which would sign anyone in as anyone. 341 if not settings.sign_in_available: 342 raise HTTPException( 343 503, "Email sign-in is not set up on this server yet." 344 ) 345 346 try: 347 link = auth.issue(session, account, email) 348 except auth.TooManyRequests as exc: 349 raise HTTPException(429, str(exc)) from exc 350 351 try: 352 sent = mailer.send_login_link(email, link) 353 except mailer.MailError as exc: 354 raise HTTPException(502, str(exc)) from exc 355 356 out = {"sent": sent, "email": email} 357 if not sent and settings.dev_login_links: 358 out["dev_link"] = link 359 return out 360 361 362 @router.post("/auth/verify", response_model=AccountOut) 363 def verify_sign_in( 364 body: TokenIn, 365 response: Response, 366 account: Annotated[Account, Depends(get_account)], 367 session: Annotated[Session, Depends(get_session)], 368 ): 369 owner = auth.redeem(session, body.token, account) 370 if owner is None: 371 raise HTTPException( 372 400, "That sign-in link has expired or was already used. " 373 "Request a new one." 374 ) 375 _set_identity(response, owner) 376 return _account_out(session, owner) 377 378 379 @router.post("/auth/signout") 380 def sign_out(response: Response): 381 """Forget this browser. The account and everything in it stay put.""" 382 response.delete_cookie(DEVICE_COOKIE) 383 return {"signed_out": True} 384 385 386 # -------------------------------------------------------------------------- 387 # billing 388 # -------------------------------------------------------------------------- 389 390 @router.post("/billing/checkout") 391 def create_checkout( 392 body: CheckoutIn, 393 account: Annotated[Account, Depends(get_account)], 394 session: Annotated[Session, Depends(get_session)], 395 ): 396 plan = pricing.get(body.plan_id) 397 if plan is None: 398 raise HTTPException(404, "No such plan.") 399 try: 400 url = payments.start_checkout(session, account, plan) 401 except payments.PaymentsUnavailable as exc: 402 raise HTTPException(503, str(exc)) from exc 403 except payments.PaymentError as exc: 404 raise HTTPException(400, str(exc)) from exc 405 return {"url": url} 406 407 408 @router.get("/billing/return") 409 def checkout_return( 410 session_id: str, 411 session: Annotated[Session, Depends(get_session)], 412 ): 413 """Where Stripe sends the browser after paying. 414 415 Fulfils from here as well as from the webhook, so credits are there by the 416 time the page loads rather than whenever the webhook gets round to it. 417 Identity needs no handling: if paying folded this device into an existing 418 account, get_account re-points the cookie on the very next request. 419 """ 420 try: 421 paid = payments.fulfil_by_id(session, session_id) is not None 422 except payments.PaymentsUnavailable: 423 paid = False 424 state = "paid" if paid else "pending" 425 return RedirectResponse(f"/?checkout={state}", status_code=303) 426 427 428 @router.post("/billing/webhook") 429 async def stripe_webhook( 430 request: Request, 431 session: Annotated[Session, Depends(get_session)], 432 ): 433 payload = await request.body() 434 try: 435 event = payments.verify_webhook( 436 payload, request.headers.get("stripe-signature") 437 ) 438 except payments.PaymentsUnavailable as exc: 439 raise HTTPException(503, str(exc)) from exc 440 except payments.PaymentError as exc: 441 raise HTTPException(400, str(exc)) from exc 442 payments.handle_event(session, event) 443 return {"received": True} 444 445 446 @router.post("/billing/portal") 447 def billing_portal( 448 account: Annotated[Account, Depends(get_account)], 449 ): 450 try: 451 return {"url": payments.portal_url(account)} 452 except payments.PaymentsUnavailable as exc: 453 raise HTTPException(503, str(exc)) from exc 454 except payments.PaymentError as exc: 455 raise HTTPException(400, str(exc)) from exc 456 457 458 @router.post("/uploads", response_model=UploadOut) 459 async def create_upload( 460 account: Annotated[Account, Depends(get_account)], 461 session: Annotated[Session, Depends(get_session)], 462 files: Annotated[list[UploadFile], File()], 463 ): 464 """Stage a dropped pair, work out which is which, and guess the language.""" 465 names = [f.filename or "" for f in files] 466 try: 467 pairing = detect.classify(names) 468 except detect.DetectionError as exc: 469 raise HTTPException(400, str(exc)) from exc 470 471 job_id = new_id("job") 472 paths = Paths.for_job(job_id) 473 paths.create() 474 475 try: 476 order = {name: i for i, name in enumerate(pairing.audio_names)} 477 by_name = {(f.filename or ""): f for f in files} 478 479 text_path = staged_text_path(job_id, pairing.text_name) 480 await _save(by_name[pairing.text_name], text_path) 481 482 # fb2/mobi/azw3 become something the aligner can read. Do it now, not at 483 # run time, so a book we cannot open fails while the user is watching. 484 if convert.needs_conversion(pairing.text_name): 485 try: 486 converted = convert.to_readable( 487 text_path, text_path.with_suffix("") 488 ) 489 except convert.ConversionError as exc: 490 raise HTTPException(400, str(exc)) from exc 491 if converted != text_path: 492 text_path.unlink(missing_ok=True) 493 text_path = converted 494 495 audio_paths: list[Path] = [] 496 for name in pairing.audio_names: 497 upload = by_name.get(name) 498 if upload is None: 499 raise HTTPException(400, f"Upload did not include {name}.") 500 dest = ( 501 staged_part_path(job_id, order[name] + 1, name) 502 if pairing.is_multipart 503 else staged_audio_path(job_id, name) 504 ) 505 # The staged stem is shared between audio and text, so a single 506 # audio file must not land on the text file's path. 507 if dest == text_path: 508 raise HTTPException( 509 400, "The audiobook and the book must be different formats." 510 ) 511 await _save(upload, dest) 512 audio_paths.append(dest) 513 514 if not audio_paths: 515 raise HTTPException(400, "Upload did not include any audio.") 516 517 # Custom cover art is a signed-in feature. Say so plainly rather than 518 # accepting the file and quietly ignoring it: the UI labels the drop 519 # zone before anyone picks an image, and this is the backstop. 520 cover_name = None 521 cover_blocked = False 522 if pairing.cover_name: 523 if account.signed_in: 524 upload = by_name.get(pairing.cover_name) 525 if upload is not None: 526 await _save(upload, staged_cover_path(job_id, pairing.cover_name)) 527 cover_name = pairing.cover_name 528 else: 529 cover_blocked = True 530 531 detection = detect.detect_language(detect.extract_text_sample(text_path)) 532 # Parts are merged at run time, so total the durations here. 533 durations = [probe_duration(p) for p in audio_paths] 534 duration = sum(d for d in durations if d) if any(durations) else None 535 536 # Does this text actually belong to this audio? Answering now costs a 537 # few seconds; finding out during the run costs the whole run. 538 match = matching.check( 539 audio_paths[0], 540 text_path, 541 detection.code if detection.supported else "en", 542 aligner, 543 ) 544 545 language = detection.code if detection.supported else "en" 546 lang = languages.require(language) 547 548 job = Job( 549 id=job_id, 550 account_id=account.id, 551 status=JobStatus.draft, 552 language=lang.code, 553 splitter=lang.splitter, 554 model=settings.model, 555 audio_filename=pairing.display_name, 556 text_filename=pairing.text_name, 557 audio_parts=len(audio_paths), 558 cover_filename=cover_name, 559 audio_bytes=sum(p.stat().st_size for p in audio_paths), 560 audio_duration_seconds=duration, 561 stage="Ready to start", 562 ) 563 session.add(job) 564 session.commit() 565 566 except HTTPException: 567 shutil.rmtree(paths.root, ignore_errors=True) 568 raise 569 except detect.DetectionError as exc: 570 shutil.rmtree(paths.root, ignore_errors=True) 571 raise HTTPException(400, str(exc)) from exc 572 except Exception as exc: # noqa: BLE001 573 shutil.rmtree(paths.root, ignore_errors=True) 574 raise HTTPException(500, f"Upload failed: {exc}") from exc 575 576 return UploadOut( 577 job=_job_out(job, []), 578 detected=DetectionOut( 579 code=detection.code, name=detection.name, 580 confidence=round(detection.confidence, 4), 581 supported=detection.supported, 582 ), 583 match=match.as_dict(), 584 cover_requires_sign_in=cover_blocked, 585 ) 586 587 588 async def _save(upload: UploadFile, dest: Path) -> None: 589 """Stream to disk in chunks - these files run to hundreds of megabytes.""" 590 dest.parent.mkdir(parents=True, exist_ok=True) 591 written = 0 592 with dest.open("wb") as out: 593 while chunk := await upload.read(_CHUNK): 594 written += len(chunk) 595 if written > settings.max_upload_bytes: 596 raise HTTPException( 597 413, 598 f"{upload.filename} exceeds the " 599 f"{settings.max_upload_bytes // (1024**3)} GiB upload limit.", 600 ) 601 out.write(chunk) 602 603 604 @router.post("/jobs/{job_id}/start", response_model=JobOut) 605 def start_job( 606 job_id: str, 607 body: StartIn, 608 account: Annotated[Account, Depends(get_account)], 609 session: Annotated[Session, Depends(get_session)], 610 ): 611 job = _load(session, account, job_id) 612 # A failed job keeps its staged inputs, so it can be retried in place - 613 # usually after correcting the language. 614 if job.status not in (JobStatus.draft, JobStatus.failed): 615 raise HTTPException(409, f"Job is already {job.status.value}.") 616 if job.status == JobStatus.failed and not Paths.for_job(job.id).inp.exists(): 617 raise HTTPException( 618 409, "The uploaded files for this job are gone. Upload them again." 619 ) 620 621 if body.language: 622 try: 623 lang = languages.require(body.language) 624 except languages.UnsupportedLanguage as exc: 625 raise HTTPException(400, str(exc)) from exc 626 job.language, job.splitter = lang.code, lang.splitter 627 if body.model: 628 job.model = body.model 629 630 try: 631 billing.authorize_start(session, account, job) 632 except billing.PaymentRequired as exc: 633 session.rollback() 634 raise HTTPException(402, str(exc)) from exc 635 636 # With an external queue the worker is probably not this machine, so the 637 # staged inputs have to go somewhere both sides can reach before enqueuing. 638 if settings.queue_backend != "memory": 639 paths = Paths.for_job(job.id) 640 prefix = input_prefix(job.id) 641 for local in sorted(paths.inp.rglob("*")): 642 if local.is_file(): 643 rel = local.relative_to(paths.inp).as_posix() 644 storage.put_file(f"{prefix}/{rel}", local) 645 646 job.status = JobStatus.queued 647 job.stage = "Queued" 648 job.progress = 0.0 649 job.error = None 650 session.commit() 651 652 queue.enqueue(job.id) 653 return _job_out(job, []) 654 655 656 # -------------------------------------------------------------------------- 657 # jobs that run in the browser 658 # -------------------------------------------------------------------------- 659 # 660 # The audio never reaches us. The server's part is to say whether the job may 661 # start (the free window, or a credit), and to keep the finished subtitles so 662 # they are still there on another device. Nothing here can be enforced against 663 # someone who edits the page's JavaScript, and that is accepted: see billing.py. 664 665 _LOCAL_STALE_HOURS = 48 666 667 668 @router.post("/local/jobs", response_model=JobOut) 669 def start_local_job( 670 body: LocalJobIn, 671 account: Annotated[Account, Depends(get_account)], 672 session: Annotated[Session, Depends(get_session)], 673 ): 674 # The same book started again - a closed tab, a reload - carries on under 675 # the job it already paid for rather than being charged a second time. 676 running = ( 677 session.query(Job) 678 .filter( 679 Job.account_id == account.id, Job.local == 1, 680 Job.status == JobStatus.running, 681 Job.audio_filename == body.audio_filename, 682 Job.audio_bytes == body.audio_bytes, 683 ) 684 .order_by(Job.created_at.desc()) 685 .first() 686 ) 687 if running is not None: 688 running.language = body.language 689 session.commit() 690 return _job_out(running, []) 691 692 job = Job( 693 account_id=account.id, status=JobStatus.running, local=1, 694 language=body.language, splitter="browser", model="whisper-tiny", 695 audio_filename=body.audio_filename, text_filename=body.text_filename, 696 audio_parts=body.audio_parts, audio_bytes=body.audio_bytes, 697 audio_duration_seconds=body.audio_duration_seconds, 698 stage="Running in your browser", started_at=utcnow(), 699 ) 700 session.add(job) 701 try: 702 billing.authorize_start(session, account, job) 703 except billing.PaymentRequired as exc: 704 session.rollback() 705 raise HTTPException(402, str(exc)) from exc 706 session.commit() 707 return _job_out(job, []) 708 709 710 @router.post("/local/jobs/{job_id}/finish", response_model=JobOut) 711 def finish_local_job( 712 job_id: str, 713 body: LocalFinishIn, 714 account: Annotated[Account, Depends(get_account)], 715 session: Annotated[Session, Depends(get_session)], 716 ): 717 job = _load(session, account, job_id) 718 if not job.local or job.status != JobStatus.running: 719 raise HTTPException(409, f"Job is already {job.status.value}.") 720 721 import json 722 import tempfile 723 724 name = Path(body.filename).name or "subtitles.srt" 725 with tempfile.TemporaryDirectory() as tmp: 726 for kind, filename, content in ( 727 ("srt", name, body.srt), 728 ("metadata", "metadata.json", 729 json.dumps({"job_id": job.id, **body.metadata}, ensure_ascii=False, indent=2)), 730 ): 731 src = Path(tmp) / filename 732 # Bytes, not text: on Windows write_text would turn every line ending into CRLF. 733 src.write_bytes(content.encode("utf-8")) 734 key = f"{job.id}/{filename}" 735 size = storage.put_file(key, src) 736 session.add(Artifact(job_id=job.id, kind=kind, filename=filename, 737 storage_key=key, size_bytes=size)) 738 739 job.status = JobStatus.succeeded 740 job.stage, job.progress, job.finished_at = "Done", 1.0, utcnow() 741 session.commit() 742 return _job_out(job, _artifacts(session, job.id)) 743 744 745 @router.post("/local/jobs/{job_id}/fail", response_model=JobOut) 746 def fail_local_job( 747 job_id: str, 748 body: LocalFailIn, 749 account: Annotated[Account, Depends(get_account)], 750 session: Annotated[Session, Depends(get_session)], 751 ): 752 job = _load(session, account, job_id) 753 if not job.local or job.status != JobStatus.running: 754 raise HTTPException(409, f"Job is already {job.status.value}.") 755 job.status = JobStatus.failed 756 job.stage, job.error, job.finished_at = "Failed", body.error or "Failed in the browser.", utcnow() 757 # Whatever went wrong, they got nothing: the free slot or the credit goes back. 758 billing.refund(session, job) 759 session.commit() 760 return _job_out(job, []) 761 762 763 def expire_stale_local_jobs(session: Session) -> int: 764 """A tab that was closed for good never reports back. Release what it held.""" 765 from datetime import timedelta 766 767 stale = ( 768 session.query(Job) 769 .filter(Job.local == 1, Job.status == JobStatus.running, 770 Job.created_at < utcnow() - timedelta(hours=_LOCAL_STALE_HOURS)) 771 .all() 772 ) 773 for job in stale: 774 job.status, job.stage, job.finished_at = JobStatus.canceled, "Abandoned", utcnow() 775 billing.refund(session, job) 776 session.commit() 777 return len(stale) 778 779 780 @router.post("/convert") 781 async def convert_book(file: Annotated[UploadFile, File()]): 782 """mobi / azw3 to epub. The one thing the browser cannot do for itself: 783 those formats need a real parser, and a book is small enough to send.""" 784 name = Path(file.filename or "book").name 785 if not convert.needs_conversion(name): 786 raise HTTPException(400, "That format does not need converting.") 787 import tempfile 788 789 with tempfile.TemporaryDirectory() as tmp: 790 src = Path(tmp) / name 791 await _save(file, src) 792 try: 793 out = convert.to_readable(src, Path(tmp) / "converted") 794 except convert.ConversionError as exc: 795 raise HTTPException(400, str(exc)) from exc 796 data = out.read_bytes() 797 return Response(data, media_type="application/epub+zip", 798 headers={"X-Filename": Path(name).stem + ".epub"}) 799 800 801 @router.get("/jobs", response_model=list[JobOut]) 802 def list_jobs( 803 account: Annotated[Account, Depends(get_account)], 804 session: Annotated[Session, Depends(get_session)], 805 ): 806 jobs = ( 807 session.query(Job) 808 .filter(Job.account_id == account.id) 809 .order_by(Job.created_at.desc()) 810 .limit(50) 811 .all() 812 ) 813 return [_job_out(j, _artifacts(session, j.id)) for j in jobs] 814 815 816 @router.get("/jobs/{job_id}", response_model=JobOut) 817 def get_job( 818 job_id: str, 819 account: Annotated[Account, Depends(get_account)], 820 session: Annotated[Session, Depends(get_session)], 821 ): 822 job = _load(session, account, job_id) 823 return _job_out(job, _artifacts(session, job.id)) 824 825 826 @router.post("/jobs/{job_id}/cancel", response_model=JobOut) 827 def cancel_job( 828 job_id: str, 829 account: Annotated[Account, Depends(get_account)], 830 session: Annotated[Session, Depends(get_session)], 831 ): 832 job = _load(session, account, job_id) 833 if job.status in (JobStatus.succeeded, JobStatus.failed, JobStatus.canceled): 834 raise HTTPException(409, f"Job is already {job.status.value}.") 835 job.status = JobStatus.canceled 836 job.stage = "Canceled" 837 job.finished_at = utcnow() 838 # A canceled job must not eat the free conversion. 839 billing.refund(session, job) 840 session.commit() 841 return _job_out(job, []) 842 843 844 @router.delete("/jobs/{job_id}") 845 def delete_job( 846 job_id: str, 847 account: Annotated[Account, Depends(get_account)], 848 session: Annotated[Session, Depends(get_session)], 849 ): 850 job = _load(session, account, job_id) 851 if job.status in (JobStatus.queued, JobStatus.running): 852 raise HTTPException(409, "Cancel the job before deleting it.") 853 storage.delete_prefix(job.id) 854 shutil.rmtree(Paths.for_job(job.id).root, ignore_errors=True) 855 session.delete(job) 856 session.commit() 857 return {"deleted": job_id} 858 859 860 @router.get("/jobs/{job_id}/files/{kind}") 861 def download( 862 job_id: str, 863 kind: Literal["srt", "video", "video_embedded", "metadata", "log"], 864 account: Annotated[Account, Depends(get_account)], 865 session: Annotated[Session, Depends(get_session)], 866 ): 867 job = _load(session, account, job_id) 868 art = ( 869 session.query(Artifact) 870 .filter(Artifact.job_id == job.id, Artifact.kind == kind) 871 .first() 872 ) 873 if art is None: 874 raise HTTPException(404, f"No {kind} for this job.") 875 876 # S3 hands the browser a presigned URL; local storage serves the file. 877 url = storage.presigned_url(art.storage_key, art.filename) 878 if url: 879 return RedirectResponse(url, status_code=307) 880 881 assert isinstance(storage, LocalStorage) 882 return FileResponse( 883 storage.path_for(art.storage_key), 884 filename=art.filename, 885 media_type="application/octet-stream", 886 )