Recently Written · git

subplz-web

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

Log | Files | Refs


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     )