commit 9c7ec98adcef01b6d6b7f768fe7280584d1426ea equwal <truex@equwal.com> 2026-09-20 11:23:46 -0700 Initial commit: SubPlz Web, standalone A web front end that turns an audiobook plus its ebook into split-timed subtitles and a YouTube-ready MP4. Extracted from a subdirectory of a SubPlz checkout and decoupled from it: the alignment backend is now an ordinary dependency located on PATH, not a sibling venv. backend/aligner.py keeps it behind an interface so it can be replaced.
.env.example | 60 +++++ .gitignore | 23 ++ README.md | 355 +++++++++++++++++++++++++++ backend/__init__.py | 0 backend/aligner.py | 297 +++++++++++++++++++++++ backend/api.py | 514 +++++++++++++++++++++++++++++++++++++++ backend/billing.py | 159 ++++++++++++ backend/convert.py | 262 ++++++++++++++++++++ backend/db.py | 158 ++++++++++++ backend/detect.py | 170 +++++++++++++ backend/gen_languages.py | 67 +++++ backend/languages.json | 493 +++++++++++++++++++++++++++++++++++++ backend/languages.py | 63 +++++ backend/main.py | 98 ++++++++ backend/matching.py | 314 ++++++++++++++++++++++++ backend/pricing.py | 130 ++++++++++ backend/queue.py | 102 ++++++++ backend/render.py | 231 ++++++++++++++++++ backend/runner.py | 618 +++++++++++++++++++++++++++++++++++++++++++++++ backend/settings.py | 97 ++++++++ backend/storage.py | 170 +++++++++++++ frontend/app.js | 544 +++++++++++++++++++++++++++++++++++++++++ frontend/index.html | 135 +++++++++++ frontend/style.css | 290 ++++++++++++++++++++++ requirements.txt | 17 ++ run.ps1 | 97 ++++++++ tools/client.py | 259 ++++++++++++++++++++ worker.py | 65 +++++ 28 files changed, 5788 insertions(+)
diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..d346cba --- /dev/null +++ b/.env.example @@ -0,0 +1,60 @@ +# Copy to .env and edit. Every value has a working localhost default, +# so an empty .env is a valid config. + +# --- alignment ------------------------------------------------------------- +# Which backend does the aligning. See backend/aligner.py to add another. +SUBPLZ_WEB_ALIGNER=subplz +# tiny is what upstream recommends for audiobooks: the transcript only needs to +# be good enough to line up against text you already have. +SUBPLZ_WEB_MODEL=tiny +# cpu | cuda +SUBPLZ_WEB_DEVICE=cpu +# 0 = use (CPU count - 1) +SUBPLZ_WEB_THREADS=0 +SUBPLZ_WEB_JOB_TIMEOUT_SECONDS=21600 + +# --- match check ----------------------------------------------------------- +# Transcribe a few short samples on upload and score them against the book, +# so a mismatched pair is caught in seconds instead of after a full run. +SUBPLZ_WEB_MATCH_CHECK=true + +# --- video ----------------------------------------------------------------- +# Render a YouTube-ready MP4: cover image + audio + soft subtitle track. +SUBPLZ_WEB_RENDER_VIDEO=true +SUBPLZ_WEB_VIDEO_WIDTH=1920 +SUBPLZ_WEB_VIDEO_HEIGHT=1080 +# 1 fps is the minimum YouTube accepts and all a still image needs. +SUBPLZ_WEB_VIDEO_FPS=1 +# "auto" probes the H.264 encoders this ffmpeg build can actually run. +# Force one with libx264 / h264_nvenc / h264_amf / h264_qsv / libopenh264. +SUBPLZ_WEB_VIDEO_ENCODER=auto +SUBPLZ_WEB_VIDEO_CRF=28 + +# --- queue ----------------------------------------------------------------- +# memory = run jobs in this process. redis = hand them to worker.py processes. +SUBPLZ_WEB_QUEUE_BACKEND=memory +SUBPLZ_WEB_MAX_CONCURRENT_JOBS=1 +# SUBPLZ_WEB_REDIS_URL=redis://localhost:6379/0 + +# --- storage --------------------------------------------------------------- +# local = under ./data. s3 = a bucket, with presigned download URLs. +SUBPLZ_WEB_STORAGE_BACKEND=local +# SUBPLZ_WEB_S3_BUCKET=subplz-artifacts +# SUBPLZ_WEB_S3_PREFIX=jobs/ +# SUBPLZ_WEB_DOWNLOAD_URL_TTL_SECONDS=3600 + +# --- database -------------------------------------------------------------- +# Empty = SQLite under ./data. Set a Postgres URL for the public release. +# SUBPLZ_WEB_DATABASE_URL=postgresql+psycopg://user:pass@host/subplz + +# --- billing --------------------------------------------------------------- +# Off for localhost. Turning it on requires implementing billing.start_checkout. +SUBPLZ_WEB_BILLING_ENABLED=false +# One free book per rolling 24 hours. +SUBPLZ_WEB_FREE_CONVERSIONS=1 +SUBPLZ_WEB_FREE_WINDOW_HOURS=24 +# Override the plan catalogue (see backend/pricing.py for the shape). +# SUBPLZ_WEB_PLANS_JSON=[...] + +# --- uploads --------------------------------------------------------------- +SUBPLZ_WEB_MAX_UPLOAD_BYTES=2147483648 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..538afc3 --- /dev/null +++ b/.gitignore @@ -0,0 +1,23 @@ +# Runtime state: uploads, artifacts, the job database. +data/ + +# Environments +.venv/ +venv*/ + +# Local config — copy .env.example and edit +.env + +# Python +__pycache__/ +*.py[cod] +*.egg-info/ +.pytest_cache/ + +# Editors / OS +.vscode/ +.idea/ +.DS_Store +Thumbs.db +*.swp +*~ diff --git a/README.md b/README.md new file mode 100644 index 0000000..8e1752f --- /dev/null +++ b/README.md @@ -0,0 +1,355 @@ +# SubPlz Web + +Line an audiobook up with its ebook, sentence by sentence. + +Drop in an audiobook and the ebook it was read from. You get back subtitles +timed to the narration, with the wording taken from your own book rather than +from a machine's guess at what it heard. Two things to do with that: + +- **HoshiReader whispersync** — the `.srt` is the timing file. +- **Subtitled video** — a YouTube-ready MP4 with a selectable caption track. + +One free book per rolling 24 hours; see [Pricing](#pricing). + +--- + +## Quick start + +```powershell +.\run.ps1 +``` + +Opens <http://127.0.0.1:8420>. First run creates `.venv` and installs +everything, including the alignment backend — several minutes, because it +pulls torch. + +Already set up: + +```powershell +.venv\Scripts\python.exe -m uvicorn backend.main:app --port 8420 +``` + +Requires **ffmpeg on PATH** and **Python 3.11** — subplz pins `>=3.10,<3.12`, +so 3.12+ will not work. + +--- + +## About "no Whisper" + +This app never runs `subplz gen`, the mode that transcribes a book from scratch. +Only `subplz sync` is reachable, and `gen` is not exposed anywhere in the API. + +That said, `subplz sync` is itself built on Whisper, and there is no way around +that. It works like this: + +1. A **tiny** Whisper model makes a rough transcript of the audio. +2. That transcript is aligned to *your* ebook text with Needleman–Wunsch. +3. The subtitle text that gets written is **your book's text**, not Whisper's. + +So Whisper is used as a timing device, not as a source of words. Nothing it +mis-hears reaches the `.srt`. If "no Whisper" meant "no model downloads and no +transcription step at all", subplz cannot do that — you would need a different +aligner (aeneas, Montreal Forced Aligner, WhisperX) and a different tool. + +--- + +## Input formats + +**Audio** — `m4b`, `mp3`, `m4a`, `opus`, `flac`, `wav`, `mkv` and friends. +A single file, or **a folder of per-chapter files**: drop all 44 mp3s and they +are merged into one chaptered file, in natural order (`9.mp3` before `10.mp3`). +Files can arrive one drop at a time; the upload starts once both halves are in. + +**Book** — `epub`, `txt`, `srt`, `vtt`, `ass` directly; `fb2`, `fb2.zip`, +`mobi`, `azw3`, `azw` and `prc` are converted on upload. + +Conversion delegates rather than parsing ebook formats by hand +(`backend/convert.py`): + +| From | How | +|------|-----| +| `azw3` / KF8 | the `mobi` package unpacks it straight to epub | +| `mobi` (older) | `mobi` unpacks to HTML, ebooklib rebuilds the epub | +| `fb2`, `fb2.zip` | it is XML; lxml reads it, ebooklib writes the epub | + +`mobi` is the only added dependency; lxml, BeautifulSoup and ebooklib already +ship with subplz. + +**Conversion must preserve chapters**, and that is not a stylistic preference. +See [Will this text even match?](#will-this-text-even-match). + +--- + +## Will this text even match? + +subplz fails late and unhelpfully when the text does not match the audio: it +transcribes the whole book, then reports "the generated transcript and the +provided text file are too different" and writes a `.subfail`. On CPU that is a +wasted hour. + +So the app asks the same question on upload, using subplz's own rule — +transcribe a 60-second sample from the start of up to three chapters, score each +against every chapter of the book with `rapidfuzz.fuzz.ratio`, and compare +against subplz's `SCORE_THRESHOLD = 40`. It costs a few seconds. + +Measured on real files: + +| Pair | Score | Verdict | +|---|---|---| +| Moskva-Petushki audio + its own epub | **74.4** | good | +| Moskva-Petushki audio + an unrelated Russian novel | **40.5** | poor | + +That second row is the important one. An unrelated book in the same language +still clears subplz's threshold of 40, because two prose texts in one language +are roughly 40% similar character-by-character whatever they say. So the bands +here sit well above it: good at ≥60, marginal at ≥48, poor below. A poor score +warns and relabels the button "Start anyway" — it never blocks, because the +check only samples, and it is your book. + +### What actually improves the score + +Two things were tried and measured, and only one of them works. + +**Chapter structure — decisive.** subplz compares the *opening* of each audio +chapter against the *opening* of each text chapter, capped at 2000 characters. +It splits an epub into one text chapter per spine document, but a `.txt` into +exactly one chapter for the whole file. Same book, same audio: + +| Text given to subplz | Score | +|---|---| +| epub, per chapter | **69.1** | +| identical text as one flat file | **38.6** — *below the threshold* | + +A flat text file turns a perfectly good book into a failed run. That is why +`convert.py` emits a chaptered epub rather than plain text, and why the app +warns when a one-document book is paired with multi-chapter audio. + +**Typographic cleanup — no effect, so it was dropped.** Normalising curly +quotes, em dashes, soft hyphens, non-breaking spaces and footnote markers +measured **+0.0** against a real transcript. Joining paragraphs with a space +instead of subplz's empty string was worth +0.3. Both are noise next to +chapterisation, and shipping them would have been dead code. (This is the one +place it matters that `ats/lang.py` implements only Japanese and English — +every other language falls back to `English`, whose `clean()` is nothing but +`.lower()`, so punctuation is compared verbatim. It still does not move the +number.) + +Disable the check with `SUBPLZ_WEB_MATCH_CHECK=false`. + +--- + +## Languages + +97 languages. The catch upstream does not document: subplz splits sentences with +`pysbd`, which supports **23** languages and raises `ValueError` on anything +else — **Portuguese and Finnish included**. subplz can use `stanza` instead when +given `--nlp`, which covers 74 more. + +This app resolves that per request. `backend/languages.json` records which +splitter each language needs, and the aligner adds `--nlp` automatically. You +never see the failure. + +| Language | Splitter | Notes | +|---|---|---| +| Spanish, Russian, Japanese | pysbd | fast path | +| Portuguese, Finnish | stanza | one-off model download on first use | + +subplz also defaults `--language` **and** `--lang` to Japanese. The aligner +always sets both explicitly, so a non-Japanese book cannot silently run as +Japanese. + +Regenerate the registry after upgrading pysbd or stanza: + +```powershell +.venv\Scripts\python.exe -m backend.gen_languages +``` + +--- + +## How it works + +``` +drop files ──► POST /api/uploads stage, pair, convert, detect language + │ + ▼ (draft — correct the language here) + POST /api/jobs/{id}/start entitlement check, enqueue + │ + ▼ + queue ──► runner ──► aligner ──► .srt ──► .mp4 + │ + ▼ + storage + metadata.json +``` + +Upload and start are separate calls on purpose: a wrong language guess costs a +click instead of a re-upload and a wasted multi-hour run. + +### Audio is prepared twice, deliberately + +The aligner is fed a **16 kHz mono** copy; the video keeps the original quality. +That is not tidiness — handing subplz 44.1 kHz stereo crashed ctranslate2 here +(integer divide by zero, part-way through a chapter), reproducibly. Doing the +conversion ourselves also repairs damaged input: the error-tolerant ffmpeg flags +drop corrupt frames instead of letting a single bad chapter abort the whole run. + +### Files + +| Path | Role | +|---|---| +| `backend/api.py` | HTTP routes | +| `backend/aligner.py` | the alignment backend, behind an interface | +| `backend/runner.py` | stages inputs, drives the aligner, collects artifacts | +| `backend/convert.py` | fb2/mobi/azw3 → a chaptered epub | +| `backend/matching.py` | the preflight match score | +| `backend/render.py` | the YouTube MP4 | +| `backend/detect.py` | pairs the dropped files, detects the language | +| `backend/languages.py` | language registry and splitter routing | +| `backend/queue.py` | in-process or Redis dispatch | +| `backend/storage.py` | local disk or S3 | +| `backend/billing.py` | the 24-hour allowance | +| `backend/pricing.py` | the plan catalogue, with the market it was set against | +| `frontend/` | vanilla HTML/CSS/JS, no build step | +| `tools/client.py` | CLI client, and a worked example of the API | +| `worker.py` | standalone worker for the Redis backend | + +### Outputs + +| Artifact | What | +|---|---| +| `<name>.<lang>.srt` | the subtitles | +| `<name>.<lang>.mp4` | cover + audio + soft caption track (`mov_text`) | +| `metadata.json` | language, model, splitter, cue count, timing span | +| `subplz.log` | the full run log — the only way to debug a bad alignment | + +Uploaded media is deleted once a job succeeds. A **failed** job keeps its inputs +so you can fix the language and retry without re-uploading. + +--- + +## Swapping the alignment backend + +subplz is one implementation of `aligner.Aligner`, not a hard dependency. +Everything subplz-specific — its argument names, its Japanese defaults, the +shape of its progress output, where it writes the result — lives in +`SubPlzAligner`. To replace it: + +1. subclass `Aligner` (`build_command`, `progress_reader`, `locate_output`) +2. register it in `ALIGNERS` +3. set `SUBPLZ_WEB_ALIGNER` to its name + +The API, queue, storage and job runner do not change. + +--- + +## Video + +A still cover image at 1 fps, the audio, and the subtitles as a **selectable +track** rather than burned in — so the file stays small, the encode stays fast, +and the viewer can turn captions off. The cover is the largest image in the +epub, or a plain dark card when there is not one. + +The H.264 encoder is **probed, not assumed**: `ffmpeg -encoders` lists encoders +that were compiled in, including hardware ones on machines with no such +hardware, so the app encodes one test frame with each candidate and takes the +first that actually works. On this machine that is `h264_amf`; a build with +`libx264` will prefer that. Set `SUBPLZ_WEB_VIDEO_ENCODER` to force one, or +`SUBPLZ_WEB_RENDER_VIDEO=false` to skip video entirely. + +--- + +## CLI + +```powershell +.venv\Scripts\python.exe tools\client.py submit book.m4b book.epub --wait +.venv\Scripts\python.exe tools\client.py submit .\chapters\ book.epub --wait +.venv\Scripts\python.exe tools\client.py list +.venv\Scripts\python.exe tools\client.py download <job_id> --dir out\ +``` + +Use this rather than `curl` for non-ASCII filenames: curl on a non-UTF-8 console +mangles multipart filenames, and this client does not. + +--- + +## Pricing + +One free book per rolling 24 hours, then paid. The window is rolling rather than +a calendar day: the allowance returns 24 hours after the run that used it. + +| Plan | Price | Per book | +|---|---|---| +| Free | — | 1 per 24h | +| One book | $3.49 | $3.49 | +| 5 books | $12.99 | $2.60 | +| 20 books | $39.99 | $2.00 | +| Unlimited monthly | $14.99 | — | + +Set against the market (2026): the direct competitors are cheap or free — +Voxlight $29.99/year (alignment runs on the user's own Mac), Storyteller and +syncabook free but self-hosted. The adjacent subtitling tools price for +*transcription* and do not transfer: Sonix is $10/hour, so a 10-hour audiobook +would be ~$100, and Happy Scribe's 120-minute $17 tier would not fit one book. +Forced alignment is far cheaper to run than transcription, because the model is +tiny and its output is thrown away. So: per book, priced as an impulse buy, with +the subscription just under Otter ($16.99) and Happy Scribe ($17). + +Every number is an env var — see `backend/pricing.py`. + +**Identity is a cookie, and only a cookie.** Anyone who clears it gets another +free book. That is accepted: hard verification means accounts, email and a +signup wall in front of a tool whose pitch is "drop two files in". The real +protection against abuse is capacity, not identity. + +What is **not** included is a payment provider. `billing.start_checkout` raises +`NotImplementedError` and is the single place Stripe plugs in. Credit +`Account.purchased_credits` from the webhook, never from the success redirect. + +--- + +## Scaling out + +Everything environment-specific is an env var with a localhost default. See +`.env.example`. + +| Concern | localhost | public | +|---|---|---| +| Queue | thread pool in the API process | Redis + `worker.py` | +| Storage | `./data/artifacts` | S3, presigned download URLs | +| Database | SQLite | Postgres | +| Identity | cookie | replace `api.get_account` | +| Billing | off | `SUBPLZ_WEB_BILLING_ENABLED=true` | + +```bash +SUBPLZ_WEB_QUEUE_BACKEND=redis \ +SUBPLZ_WEB_REDIS_URL=redis://redis:6379/0 \ +SUBPLZ_WEB_DATABASE_URL=postgresql+psycopg://user:pass@host/subplz \ +SUBPLZ_WEB_STORAGE_BACKEND=s3 SUBPLZ_WEB_S3_BUCKET=subplz-artifacts \ +SUBPLZ_WEB_DEVICE=cuda \ +python worker.py +``` + +With an external queue the API copies staged uploads into shared storage before +enqueuing, and the worker pulls them down, so the API and workers do not need a +shared filesystem. + +### Before going public + +- Implement `billing.start_checkout` and its webhook. +- Put a reverse proxy in front for TLS and upload limits. +- Add a retention job — audiobooks are large and artifacts are kept forever. +- Run workers on hardware that can take it. Alignment needs roughly 2–3 GB of + RAM; a 1 GB VPS will OOM. A 4h40m Russian audiobook took ~45 minutes on + `tiny`/CPU with 15 threads here, and would take many hours on one vCPU. + +--- + +## Performance + +Alignment is dominated by the Whisper pass. `tiny` on CPU is the slow path; a +GPU with `--device cuda` is far faster. Chaptered files are processed chapter by +chapter, which is also what makes the progress bar meaningful — the runner +counts chapter completions rather than trusting the per-chapter bar, which +restarts at 0% for every chapter. + +Single files longer than about four hours can exhaust RAM, per upstream. Prefer +chaptered `m4b`, or a folder of per-chapter files. diff --git a/backend/__init__.py b/backend/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/backend/aligner.py b/backend/aligner.py new file mode 100644 index 0000000..2f93d9b --- /dev/null +++ b/backend/aligner.py @@ -0,0 +1,297 @@ +"""The alignment backend, behind an interface. + +Everything subplz-specific lives in `SubPlzAligner`: its argument names, its +Japanese defaults, the shape of its progress output, where it writes the result. +The rest of the application talks to the `Aligner` interface, so replacing +subplz means writing one new class and changing SUBPLZ_WEB_ALIGNER - not +touching the API, the queue, storage or the job runner. + +To add a backend: + 1. subclass Aligner + 2. register it in ALIGNERS below + 3. set SUBPLZ_WEB_ALIGNER to its name +""" + +from __future__ import annotations + +import os +import re +from abc import ABC, abstractmethod +from dataclasses import dataclass +from pathlib import Path + +from . import languages +from .settings import settings + + +@dataclass(frozen=True) +class AlignRequest: + """Everything a backend needs to align one book.""" + + audio: Path + text: Path + out_dir: Path + language: str + model: str + device: str + threads: int + # Chapter count of the audio, or 1. Backends may use it for progress. + chapters: int = 1 + + +@dataclass(frozen=True) +class ProgressUpdate: + stage: str + # 0..1 within the alignment run, or None to leave the bar where it is. + fraction: float | None = None + + +class ProgressReader(ABC): + """Per-run state for interpreting a backend's console output.""" + + @abstractmethod + def feed(self, line: str) -> ProgressUpdate | None: + """Interpret one output line. Return None if it says nothing useful.""" + + +class Aligner(ABC): + name: str = "aligner" + output_suffix: str = ".srt" + + # Score at or below which this backend refuses to pair audio with text. + # 0 means the backend has no such notion. + match_threshold: float = 0.0 + + def score_pair(self, audio_text: str, book_text: str) -> float: + """How well a transcript of some audio matches a piece of the book. + + Same arithmetic the backend uses internally, so a preflight number + means the same thing as the one that decides a real run. + """ + return 0.0 + + @abstractmethod + def build_command(self, req: AlignRequest) -> list[str]: + """The subprocess to run.""" + + @abstractmethod + def progress_reader(self, req: AlignRequest) -> ProgressReader: + ... + + @abstractmethod + def locate_output(self, req: AlignRequest) -> Path | None: + """The subtitle file produced, or None if there is not one.""" + + def environment(self) -> dict[str, str]: + """Environment for the subprocess.""" + env = os.environ.copy() + # Backends print emoji; without this a piped stdout dies on a + # cp1252/cp932 console. + env["PYTHONIOENCODING"] = "utf-8" + env["PYTHONUTF8"] = "1" + return env + + def language_note(self, code: str) -> str | None: + """Anything the user should know about this language, or None.""" + return None + + +# --------------------------------------------------------------------------- +# subplz +# --------------------------------------------------------------------------- + +# subplz transcribes one chapter at a time and each chapter gets its own bar +# running 0->100%. Matching only the "Transcribe:" bar keeps the sentence +# splitting and grouping bars from yanking the number around. +_TRANSCRIBE_PCT = re.compile(r"Transcribe:\s*(\d{1,3})%") + +_STAGES: list[tuple[re.Pattern, str]] = [ + (re.compile(r"Starting '"), "Loading audio"), + (re.compile(r"Fuzzy matching chapters"), "Matching chapters"), + (re.compile(r"Splitting transcript into sentences"), "Splitting text into sentences"), + (re.compile(r"Syncing"), "Aligning audio to text"), + (re.compile(r"Grouping based on transcript"), "Grouping subtitle lines"), + (re.compile(r"Writing generated subs"), "Writing subtitles"), +] + +# Transcription dominates the wall clock; the later phases share the tail. +_RUN_START = 0.05 +_TRANSCRIBE_END = 0.80 + +_STAGE_PROGRESS: dict[str, float] = { + "Aligning audio to text": 0.82, + "Grouping subtitle lines": 0.86, + "Writing subtitles": 0.89, +} + + +def _subplz_clean(text: str, lang_code: str) -> str: + """What subplz feeds to fuzz.ratio: lang.normalize(lang.clean(text)). + + Falls back to the library's own implementation when it is importable, so + this cannot drift from the real thing; the inline version is only a + stand-in for when ats is not installed. + """ + try: + from ats.lang import get_lang + + lang = get_lang(lang_code) + return lang.normalize(lang.clean(text)) + except Exception: # noqa: BLE001 - scoring must never break an upload + import unicodedata + + return unicodedata.normalize("NFKD", text.lower()) + + +class _SubPlzProgress(ProgressReader): + def __init__(self, chapters: int): + self.chapters = max(1, chapters) + self.done = 0 + self.last_pct = 0 + self.stage = "Starting subplz" + self.best = _RUN_START + + def feed(self, line: str) -> ProgressUpdate | None: + for pattern, label in _STAGES: + if pattern.search(line): + self.stage = label + break + + fraction = None + m = _TRANSCRIBE_PCT.search(line) + if m: + pct = min(100, max(0, int(m.group(1)))) + # The bar restarting means the previous chapter finished. + if pct < self.last_pct: + self.done = min(self.done + 1, self.chapters - 1) + self.last_pct = pct + frac = min(1.0, (self.done + pct / 100.0) / self.chapters) + fraction = _RUN_START + frac * (_TRANSCRIBE_END - _RUN_START) + self.stage = ( + f"Transcribing chapter {self.done + 1} of {self.chapters}" + if self.chapters > 1 + else "Transcribing audio" + ) + elif self.stage in _STAGE_PROGRESS: + fraction = _STAGE_PROGRESS[self.stage] + + if fraction is not None: + # Progress only ever moves forward. + self.best = max(self.best, fraction) + return ProgressUpdate(stage=self.stage, fraction=self.best) + + +class SubPlzAligner(Aligner): + """kanjieater/SubPlz, driven through its `sync` subcommand. + + `sync` only ever times the text you supply. `gen`, which transcribes a book + from scratch, is deliberately never invoked. + """ + + name = "subplz" + output_suffix = ".srt" + + # subplz/sync.py: SCORE_THRESHOLD = 40. A chapter whose best fuzz.ratio + # never exceeds this is reported as "too different" and left unmatched. + match_threshold = 40.0 + + def score_pair(self, audio_text: str, book_text: str) -> float: + """Reproduces subplz's own match_start() scoring, exactly. + + Two details that are easy to get wrong and both matter a lot: + + * It compares only the **first min(len_a, len_b, 2000) characters**, so + the score is about how chapters *open*, not how similar they are + overall. Front matter or an unread heading at the top of a chapter + sinks it even when the rest is identical. + * Cleaning is language-dependent, and ats/lang.py only implements + Japanese and English - every other language falls back to English, + whose clean() is nothing but .lower(). So for Russian, Spanish, + Portuguese and the rest, punctuation and spacing are compared + verbatim. + """ + from rapidfuzz import fuzz + + a = _subplz_clean(audio_text, self._lang_code) + b = _subplz_clean(book_text, self._lang_code) + # Below this, subplz does not even consider the pair. + if len(a) < 100 or len(b) < 100: + return 0.0 + n = min(len(a), len(b), 2000) + return float(fuzz.ratio(a[:n], b[:n])) + + # Set per request so scoring matches the language subplz will run under. + _lang_code: str = "en" + + def for_language(self, code: str) -> "SubPlzAligner": + clone = SubPlzAligner() + clone._lang_code = code + return clone + + def build_command(self, req: AlignRequest) -> list[str]: + lang = languages.require(req.language) + + cmd = [ + str(settings.subplz_bin), "sync", + # subplz takes either -d, or all three of --audio/--text/ + # --output-dir, and rejects a mix. The explicit form pairs the files + # by name rather than by directory sort order. + "--audio", str(req.audio), + "--text", str(req.text), + "--output-dir", str(req.out_dir), + "--output-format", "srt", + # subplz defaults BOTH of these to Japanese - always set them. + "--language", lang.code, + "--lang", lang.code, + "--lang-ext", lang.code, + "--model", req.model, + "--device", req.device, + "--overwrite", + "--rerun", + "--progress", + "--threads", str(req.threads), + ] + + # pysbd cannot segment this language; make subplz use stanza instead. + if lang.needs_nlp_flag: + cmd.append("--nlp") + + return cmd + + def progress_reader(self, req: AlignRequest) -> ProgressReader: + return _SubPlzProgress(req.chapters) + + def locate_output(self, req: AlignRequest) -> Path | None: + # subplz writes <stem>.<lang-ext>.srt, and the stem is ours. + expected = req.out_dir / f"{req.audio.stem}.{req.language}.srt" + if expected.exists(): + return expected + # Fall back in case upstream changes the convention. + return next(iter(sorted(req.out_dir.glob("*.srt"))), None) + + def language_note(self, code: str) -> str | None: + lang = languages.get(code) + if lang is None or not lang.needs_nlp_flag: + return None + return ( + f"{lang.name} uses the stanza sentence splitter. " + f"The first {lang.name} run downloads a small model." + ) + + +ALIGNERS: dict[str, type[Aligner]] = { + SubPlzAligner.name: SubPlzAligner, +} + + +def get_aligner() -> Aligner: + try: + return ALIGNERS[settings.aligner]() + except KeyError: + known = ", ".join(sorted(ALIGNERS)) + raise RuntimeError( + f"Unknown aligner {settings.aligner!r}. Known backends: {known}" + ) from None + + +aligner: Aligner = get_aligner() diff --git a/backend/api.py b/backend/api.py new file mode 100644 index 0000000..19a2a1e --- /dev/null +++ b/backend/api.py @@ -0,0 +1,514 @@ +"""HTTP API. + +Flow: drop files -> POST /api/uploads (stages, pairs, detects language) -> +POST /api/jobs/{id}/start (entitlement check, enqueue) -> poll GET /api/jobs/{id} +-> download from /api/jobs/{id}/files/{kind}. + +Upload and start are separate so a wrong language guess costs a click rather +than a re-upload and a wasted multi-hour run. +""" + +from __future__ import annotations + +import secrets +import shutil +from pathlib import Path +from typing import Annotated, Literal + +from fastapi import ( + APIRouter, Cookie, Depends, File, HTTPException, Response, UploadFile, +) +from fastapi.responses import FileResponse, RedirectResponse +from pydantic import BaseModel, Field +from sqlalchemy.orm import Session + +from . import billing, convert, detect, languages, matching, pricing +from .aligner import aligner +from .db import Account, Artifact, Job, JobStatus, SessionLocal, new_id, utcnow +from .queue import queue +from .runner import ( + Paths, input_prefix, probe_duration, staged_audio_path, staged_part_path, + staged_text_path, +) +from .settings import settings +from .storage import LocalStorage, storage + +router = APIRouter(prefix="/api") + +DEVICE_COOKIE = "subplz_device" +_CHUNK = 4 * 1024 * 1024 + + +# -------------------------------------------------------------------------- +# session / account +# -------------------------------------------------------------------------- + +def get_session(): + with SessionLocal() as s: + yield s + + +def get_account( + response: Response, + session: Annotated[Session, Depends(get_session)], + subplz_device: Annotated[str | None, Cookie()] = None, +) -> Account: + """Identify the caller. + + A cookie is the whole identity check, on purpose. Anyone who clears it gets + another free book, and that is an accepted cost: hard verification would + mean accounts, email and a signup wall in front of a tool whose pitch is + "drop two files in". The real protection against abuse is capacity - the + queue and per-worker limits - not identity. + + Swapping this for real auth later means changing this one function: + everything downstream just receives an Account. + """ + token = subplz_device + account = None + if token: + account = session.query(Account).filter(Account.device_token == token).first() + + if account is None: + token = secrets.token_urlsafe(24) + account = Account(device_token=token) + session.add(account) + session.commit() + response.set_cookie( + DEVICE_COOKIE, token, + max_age=60 * 60 * 24 * 365, httponly=True, samesite="lax", + ) + return account + + +# -------------------------------------------------------------------------- +# schemas +# -------------------------------------------------------------------------- + +class LanguageOut(BaseModel): + code: str + name: str + splitter: str + # Anything the active backend wants the user to know about this language. + note: str | None = None + + +class DetectionOut(BaseModel): + code: str | None + name: str | None + confidence: float + supported: bool + + +class ArtifactOut(BaseModel): + kind: str + filename: str + size_bytes: int + url: str + + +class JobOut(BaseModel): + id: str + status: str + stage: str + progress: float + language: str + language_name: str + splitter: str + model: str + audio_filename: str + audio_parts: int + text_filename: str + audio_bytes: int + audio_duration_seconds: float | None + error: str | None + created_at: str + artifacts: list[ArtifactOut] = Field(default_factory=list) + + +class UploadOut(BaseModel): + job: JobOut + detected: DetectionOut + # How well the book scores against a sample of the audio, using the + # backend's own rule. None when the check is disabled. + match: dict | None = None + + +class StartIn(BaseModel): + language: str | None = None + model: str | None = None + + +class AccountOut(BaseModel): + id: str + billing_enabled: bool + free_allowance: int + free_window_hours: int + free_tier_summary: str + purchased_credits: int + used: int + remaining: int + allowed: bool + next_free_at: str | None + reason: str + queue_depth: int + + +def _job_out(job: Job, arts: list[Artifact]) -> JobOut: + lang = languages.get(job.language) + return JobOut( + id=job.id, + status=job.status.value, + stage=job.stage, + progress=round(job.progress, 4), + language=job.language, + language_name=lang.name if lang else job.language, + splitter=job.splitter, + model=job.model, + audio_filename=job.audio_filename, + audio_parts=job.audio_parts or 1, + text_filename=job.text_filename, + audio_bytes=job.audio_bytes, + audio_duration_seconds=job.audio_duration_seconds, + error=job.error, + created_at=job.created_at.isoformat(), + artifacts=[ + ArtifactOut( + kind=a.kind, filename=a.filename, size_bytes=a.size_bytes, + url=f"/api/jobs/{job.id}/files/{a.kind}", + ) + for a in arts + ], + ) + + +def _load(session: Session, account: Account, job_id: str) -> Job: + job = session.get(Job, job_id) + if job is None or job.account_id != account.id: + raise HTTPException(404, "Job not found") + return job + + +def _artifacts(session: Session, job_id: str) -> list[Artifact]: + return ( + session.query(Artifact) + .filter(Artifact.job_id == job_id) + .order_by(Artifact.kind) + .all() + ) + + +# -------------------------------------------------------------------------- +# routes +# -------------------------------------------------------------------------- + +@router.get("/languages", response_model=list[LanguageOut]) +def list_languages(): + return [ + LanguageOut( + code=l.code, name=l.name, splitter=l.splitter, + note=aligner.language_note(l.code), + ) + for l in languages.all_languages() + ] + + +@router.get("/account", response_model=AccountOut) +def get_account_info( + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + ent = billing.check(session, account) + return AccountOut( + id=account.id, + billing_enabled=settings.billing_enabled, + free_allowance=ent.free_allowance, + free_window_hours=ent.window_hours, + free_tier_summary=pricing.free_tier_summary(), + purchased_credits=ent.purchased_credits, + used=ent.used, + remaining=ent.remaining, + allowed=ent.allowed, + next_free_at=ent.next_free_at.isoformat() if ent.next_free_at else None, + reason=ent.reason, + queue_depth=queue.depth(), + ) + + +@router.get("/pricing") +def get_pricing(): + """The catalogue, for the paywall. Static - no account needed.""" + return { + "free_tier": pricing.free_tier_summary(), + "billing_enabled": settings.billing_enabled, + "plans": pricing.as_dicts(), + } + + +@router.post("/uploads", response_model=UploadOut) +async def create_upload( + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], + files: Annotated[list[UploadFile], File()], +): + """Stage a dropped pair, work out which is which, and guess the language.""" + names = [f.filename or "" for f in files] + try: + pairing = detect.classify(names) + except detect.DetectionError as exc: + raise HTTPException(400, str(exc)) from exc + + job_id = new_id("job") + paths = Paths.for_job(job_id) + paths.create() + + try: + order = {name: i for i, name in enumerate(pairing.audio_names)} + by_name = {(f.filename or ""): f for f in files} + + text_path = staged_text_path(job_id, pairing.text_name) + await _save(by_name[pairing.text_name], text_path) + + # fb2/mobi/azw3 become something the aligner can read. Do it now, not at + # run time, so a book we cannot open fails while the user is watching. + if convert.needs_conversion(pairing.text_name): + try: + converted = convert.to_readable( + text_path, text_path.with_suffix("") + ) + except convert.ConversionError as exc: + raise HTTPException(400, str(exc)) from exc + if converted != text_path: + text_path.unlink(missing_ok=True) + text_path = converted + + audio_paths: list[Path] = [] + for name in pairing.audio_names: + upload = by_name.get(name) + if upload is None: + raise HTTPException(400, f"Upload did not include {name}.") + dest = ( + staged_part_path(job_id, order[name] + 1, name) + if pairing.is_multipart + else staged_audio_path(job_id, name) + ) + # The staged stem is shared between audio and text, so a single + # audio file must not land on the text file's path. + if dest == text_path: + raise HTTPException( + 400, "The audiobook and the book must be different formats." + ) + await _save(upload, dest) + audio_paths.append(dest) + + if not audio_paths: + raise HTTPException(400, "Upload did not include any audio.") + + detection = detect.detect_language(detect.extract_text_sample(text_path)) + # Parts are merged at run time, so total the durations here. + durations = [probe_duration(p) for p in audio_paths] + duration = sum(d for d in durations if d) if any(durations) else None + + # Does this text actually belong to this audio? Answering now costs a + # few seconds; finding out during the run costs the whole run. + match = matching.check( + audio_paths[0], + text_path, + detection.code if detection.supported else "en", + aligner, + ) + + language = detection.code if detection.supported else "en" + lang = languages.require(language) + + job = Job( + id=job_id, + account_id=account.id, + status=JobStatus.draft, + language=lang.code, + splitter=lang.splitter, + model=settings.model, + audio_filename=pairing.display_name, + text_filename=pairing.text_name, + audio_parts=len(audio_paths), + audio_bytes=sum(p.stat().st_size for p in audio_paths), + audio_duration_seconds=duration, + stage="Ready to start", + ) + session.add(job) + session.commit() + + except HTTPException: + shutil.rmtree(paths.root, ignore_errors=True) + raise + except detect.DetectionError as exc: + shutil.rmtree(paths.root, ignore_errors=True) + raise HTTPException(400, str(exc)) from exc + except Exception as exc: # noqa: BLE001 + shutil.rmtree(paths.root, ignore_errors=True) + raise HTTPException(500, f"Upload failed: {exc}") from exc + + return UploadOut( + job=_job_out(job, []), + detected=DetectionOut( + code=detection.code, name=detection.name, + confidence=round(detection.confidence, 4), + supported=detection.supported, + ), + match=match.as_dict(), + ) + + +async def _save(upload: UploadFile, dest: Path) -> None: + """Stream to disk in chunks - these files run to hundreds of megabytes.""" + dest.parent.mkdir(parents=True, exist_ok=True) + written = 0 + with dest.open("wb") as out: + while chunk := await upload.read(_CHUNK): + written += len(chunk) + if written > settings.max_upload_bytes: + raise HTTPException( + 413, + f"{upload.filename} exceeds the " + f"{settings.max_upload_bytes // (1024**3)} GiB upload limit.", + ) + out.write(chunk) + + +@router.post("/jobs/{job_id}/start", response_model=JobOut) +def start_job( + job_id: str, + body: StartIn, + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + job = _load(session, account, job_id) + # A failed job keeps its staged inputs, so it can be retried in place - + # usually after correcting the language. + if job.status not in (JobStatus.draft, JobStatus.failed): + raise HTTPException(409, f"Job is already {job.status.value}.") + if job.status == JobStatus.failed and not Paths.for_job(job.id).inp.exists(): + raise HTTPException( + 409, "The uploaded files for this job are gone. Upload them again." + ) + + if body.language: + try: + lang = languages.require(body.language) + except languages.UnsupportedLanguage as exc: + raise HTTPException(400, str(exc)) from exc + job.language, job.splitter = lang.code, lang.splitter + if body.model: + job.model = body.model + + ent = billing.check(session, account) + if not ent.allowed: + raise HTTPException(402, ent.reason) + + # With an external queue the worker is probably not this machine, so the + # staged inputs have to go somewhere both sides can reach before enqueuing. + if settings.queue_backend != "memory": + paths = Paths.for_job(job.id) + prefix = input_prefix(job.id) + for local in sorted(paths.inp.rglob("*")): + if local.is_file(): + rel = local.relative_to(paths.inp).as_posix() + storage.put_file(f"{prefix}/{rel}", local) + + billing.consume(session, job) + job.status = JobStatus.queued + job.stage = "Queued" + job.progress = 0.0 + job.error = None + session.commit() + + queue.enqueue(job.id) + return _job_out(job, []) + + +@router.get("/jobs", response_model=list[JobOut]) +def list_jobs( + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + jobs = ( + session.query(Job) + .filter(Job.account_id == account.id) + .order_by(Job.created_at.desc()) + .limit(50) + .all() + ) + return [_job_out(j, _artifacts(session, j.id)) for j in jobs] + + +@router.get("/jobs/{job_id}", response_model=JobOut) +def get_job( + job_id: str, + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + job = _load(session, account, job_id) + return _job_out(job, _artifacts(session, job.id)) + + +@router.post("/jobs/{job_id}/cancel", response_model=JobOut) +def cancel_job( + job_id: str, + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + job = _load(session, account, job_id) + if job.status in (JobStatus.succeeded, JobStatus.failed, JobStatus.canceled): + raise HTTPException(409, f"Job is already {job.status.value}.") + job.status = JobStatus.canceled + job.stage = "Canceled" + job.finished_at = utcnow() + # A canceled job must not eat the free conversion. + billing.refund(session, job) + session.commit() + return _job_out(job, []) + + +@router.delete("/jobs/{job_id}") +def delete_job( + job_id: str, + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + job = _load(session, account, job_id) + if job.status in (JobStatus.queued, JobStatus.running): + raise HTTPException(409, "Cancel the job before deleting it.") + storage.delete_prefix(job.id) + shutil.rmtree(Paths.for_job(job.id).root, ignore_errors=True) + session.delete(job) + session.commit() + return {"deleted": job_id} + + +@router.get("/jobs/{job_id}/files/{kind}") +def download( + job_id: str, + kind: Literal["srt", "video", "metadata", "log"], + account: Annotated[Account, Depends(get_account)], + session: Annotated[Session, Depends(get_session)], +): + job = _load(session, account, job_id) + art = ( + session.query(Artifact) + .filter(Artifact.job_id == job.id, Artifact.kind == kind) + .first() + ) + if art is None: + raise HTTPException(404, f"No {kind} for this job.") + + # S3 hands the browser a presigned URL; local storage serves the file. + url = storage.presigned_url(art.storage_key, art.filename) + if url: + return RedirectResponse(url, status_code=307) + + assert isinstance(storage, LocalStorage) + return FileResponse( + storage.path_for(art.storage_key), + filename=art.filename, + media_type="application/octet-stream", + ) diff --git a/backend/billing.py b/backend/billing.py new file mode 100644 index 0000000..a3bd483 --- /dev/null +++ b/backend/billing.py @@ -0,0 +1,159 @@ +"""Entitlement: one free book per rolling 24 hours, pay for more. + +Disabled on localhost (SUBPLZ_WEB_BILLING_ENABLED=false) so nothing gets in the +way while you use it yourself. The accounting still runs either way - every job +records whether it consumed an allowance - so turning billing on for the public +release does not need a backfill. + +The window is rolling, not a calendar day: the allowance comes back 24 hours +after the run that used it, which avoids a midnight stampede and is easier to +explain than "resets at 00:00 in some timezone". + +Deliberately not included: a payment provider. `start_checkout` is the single +seam where Stripe (or anything else) plugs in. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import datetime, timedelta, timezone + +from sqlalchemy import func +from sqlalchemy.orm import Session + +from .db import Account, Job, JobStatus, utcnow +from .settings import settings + +# Jobs in these states hold an allowance. A failed or cancelled run releases it: +# charging for our own failure is not a business model. +_HOLDING = [JobStatus.queued, JobStatus.running, JobStatus.succeeded] + + +@dataclass(frozen=True) +class Entitlement: + allowed: bool + used: int + free_allowance: int + purchased_credits: int + remaining: int + window_hours: int + # When the next free conversion becomes available, if the window is full. + next_free_at: datetime | None = None + reason: str = "" + + @property + def needs_payment(self) -> bool: + return not self.allowed + + +def _window_start() -> datetime: + return utcnow() - timedelta(hours=settings.free_window_hours) + + +def _jobs_in_window(session: Session, account_id: str) -> int: + """Billable jobs started inside the current window.""" + return ( + session.query(func.count(Job.id)) + .filter( + Job.account_id == account_id, + Job.billed == 1, + Job.status.in_(_HOLDING), + Job.created_at >= _window_start(), + ) + .scalar() + or 0 + ) + + +def _oldest_in_window(session: Session, account_id: str) -> datetime | None: + """The earliest billable job still inside the window. + + Its age is what decides when the allowance frees up again. + """ + return ( + session.query(func.min(Job.created_at)) + .filter( + Job.account_id == account_id, + Job.billed == 1, + Job.status.in_(_HOLDING), + Job.created_at >= _window_start(), + ) + .scalar() + ) + + +def check(session: Session, account: Account) -> Entitlement: + used = _jobs_in_window(session, account.id) + free = settings.free_conversions + purchased = account.purchased_credits + remaining = max(0, free + purchased - used) + window = settings.free_window_hours + + def build(allowed: bool, reason: str = "", when: datetime | None = None): + return Entitlement( + allowed=allowed, used=used, free_allowance=free, + purchased_credits=purchased, remaining=remaining, + window_hours=window, next_free_at=when, reason=reason, + ) + + if not settings.billing_enabled: + return build(True, "billing disabled") + + if remaining > 0: + return build(True) + + oldest = _oldest_in_window(session, account.id) + when = None + if oldest is not None: + if oldest.tzinfo is None: # SQLite hands back naive datetimes + oldest = oldest.replace(tzinfo=timezone.utc) + when = oldest + timedelta(hours=window) + + return build( + False, + reason=( + f"You get {free} free book every {window} hours. " + + ( + f"Your next free conversion unlocks at " + f"{when:%H:%M UTC on %d %b}." + if when + else "Try again later." + ) + + " Add credit to convert one now." + ), + when=when, + ) + + +def consume(session: Session, job: Job) -> None: + """Mark a job as having used an allowance. Called as the job is accepted.""" + job.billed = 1 + session.add(job) + + +def refund(session: Session, job: Job) -> None: + """Release the allowance a failed or cancelled job held.""" + job.billed = 0 + session.add(job) + + +def start_checkout(account: Account, quantity: int = 1) -> str: + """Return a payment URL for `quantity` extra conversions. + + Wire Stripe in here for the public release, e.g.: + + session = stripe.checkout.Session.create( + customer=account.stripe_customer_id, + line_items=[{"price": PRICE_ID, "quantity": quantity}], + mode="payment", + success_url=..., cancel_url=..., + ) + return session.url + + and credit `Account.purchased_credits` from the webhook, not from the + success redirect - the redirect is not a payment confirmation. + """ + raise NotImplementedError( + "No payment provider configured. Implement billing.start_checkout " + "before enabling SUBPLZ_WEB_BILLING_ENABLED." + ) diff --git a/backend/convert.py b/backend/convert.py new file mode 100644 index 0000000..a895eb3 --- /dev/null +++ b/backend/convert.py @@ -0,0 +1,262 @@ +"""Accept fb2, mobi and azw3 by turning them into a chaptered epub. + +Chapters are the whole point, and it took a measurement to see why. subplz +matches audio to text by scoring the *opening* of each audio chapter against +the opening of each text chapter, and it splits an epub into one text chapter +per spine document but a .txt into exactly one chapter for the entire file. + +Measured on Moskva-Petushki (Russian, 4h40m): + + same book, per chapter 69.1 + same book, one flat file 38.6 <- subplz's threshold is 40 + +So converting to plain text would push a perfectly good book *below* the +threshold and produce the "transcript and text are too different" failure. The +structure has to survive the conversion. + +Nothing here parses an ebook format or writes a container by hand: + + azw3 / KF8 -> the `mobi` package unpacks it straight to epub + mobi (older) -> `mobi` unpacks to HTML; ebooklib rebuilds the epub + fb2, fb2.zip -> it is XML; lxml reads it, ebooklib writes the epub + +`mobi` is the only added dependency; lxml, BeautifulSoup and ebooklib already +ship with subplz. +""" + +from __future__ import annotations + +import html +import shutil +import zipfile +from pathlib import Path + +# Accepted and converted on upload. Keep in sync with runner.TEXT_SUFFIXES +# and the frontend accept list. +CONVERTIBLE_SUFFIXES = {".fb2", ".mobi", ".azw", ".azw3", ".prc"} + +_MOBI_SUFFIXES = {".mobi", ".azw", ".azw3", ".prc"} + +# fb2 bodies named this way hold footnotes, not the story. +_FB2_SKIP_BODIES = {"notes", "comments"} + +# Text blocks worth keeping: prose, verse lines, subheadings. +_FB2_BLOCKS = {"p", "v", "subtitle"} + + +class ConversionError(ValueError): + pass + + +def needs_conversion(name: str) -> bool: + lowered = name.lower() + return lowered.endswith(".fb2.zip") or Path(lowered).suffix in CONVERTIBLE_SUFFIXES + + +def to_readable(src: Path, out_stem: Path) -> Path: + """Convert `src` into an epub subplz can align against.""" + name = src.name.lower() + + if name.endswith(".epub"): + return src + + if name.endswith(".fb2.zip"): + return _fb2_to_epub(_unzip_fb2(src), src, out_stem) + + if name.endswith(".fb2"): + return _fb2_to_epub(src.read_bytes(), src, out_stem) + + if Path(name).suffix in _MOBI_SUFFIXES: + return _from_mobi(src, out_stem) + + raise ConversionError(f"Cannot read {src.name}.") + + +# --------------------------------------------------------------------------- +# fb2 +# --------------------------------------------------------------------------- + +def _unzip_fb2(src: Path) -> bytes: + try: + with zipfile.ZipFile(src) as zf: + inner = next((n for n in zf.namelist() if n.lower().endswith(".fb2")), None) + if inner is None: + raise ConversionError(f"{src.name} contains no .fb2 file.") + return zf.read(inner) + except zipfile.BadZipFile as exc: + raise ConversionError(f"{src.name} is not a readable zip archive.") from exc + + +def _fb2_to_epub(raw: bytes, src: Path, out_stem: Path) -> Path: + from lxml import etree + + # recover=True: fb2 in the wild is frequently not well-formed. + root = etree.fromstring(raw, etree.XMLParser(recover=True, huge_tree=True)) + if root is None: + raise ConversionError(f"{src.name} could not be parsed as fb2.") + + def localname(el) -> str | None: + # Comments and processing instructions have a callable .tag, which + # QName rejects. Real fb2 files contain both. + return etree.QName(el).localname if isinstance(el.tag, str) else None + + def text_of(el) -> str: + return " ".join(t.strip() for t in el.itertext() if t and t.strip()) + + title = _first_text(root, localname, "book-title") or src.stem + language = _first_text(root, localname, "lang") or "en" + + chapters: list[tuple[str, list[str]]] = [] + for body in root.iter(): + if localname(body) != "body" or body.get("name") in _FB2_SKIP_BODIES: + continue + + sections = [el for el in body if localname(el) == "section"] + targets = sections or [body] + for i, section in enumerate(targets, start=1): + heading = None + paragraphs: list[str] = [] + last = None + for el in section.iter(): + tag = localname(el) + if tag == "title" and heading is None: + heading = text_of(el) or None + continue + if tag not in _FB2_BLOCKS: + continue + t = text_of(el) + # fb2 nests <p> inside <title>, so a heading is reached twice. + if t and t != last: + paragraphs.append(t) + last = t + if paragraphs: + chapters.append((heading or f"Section {i}", paragraphs)) + + if not chapters: + raise ConversionError(f"No text could be extracted from {src.name}.") + + return write_epub(out_stem.with_suffix(".epub"), title, language, chapters) + + +def _first_text(root, localname, tag: str) -> str | None: + for el in root.iter(): + if localname(el) == tag and el.text and el.text.strip(): + return el.text.strip() + return None + + +# --------------------------------------------------------------------------- +# mobi / azw3 +# --------------------------------------------------------------------------- + +def _from_mobi(src: Path, out_stem: Path) -> Path: + try: + import mobi + except ImportError as exc: # pragma: no cover - dependency is declared + raise ConversionError( + "mobi/azw3 support needs the `mobi` package: pip install mobi" + ) from exc + + tempdir = None + try: + # KF8 (azw3) unpacks straight to epub, chapters and all. + tempdir, produced = mobi.extract(str(src)) + except Exception as exc: # noqa: BLE001 - the library raises bare exceptions + if tempdir: + shutil.rmtree(tempdir, ignore_errors=True) + raise ConversionError( + f"Could not read {src.name}. If it is DRM-protected, it cannot be " + "converted." + ) from exc + + try: + out = Path(produced) + if not out.exists(): + raise ConversionError(f"Nothing could be unpacked from {src.name}.") + + if out.suffix.lower() == ".epub": + dest = out_stem.with_suffix(".epub") + dest.parent.mkdir(parents=True, exist_ok=True) + shutil.copy2(out, dest) + return dest + + pages = [out] if out.suffix.lower() in {".html", ".xhtml", ".htm"} else [] + pages += sorted( + p for p in Path(tempdir).rglob("*") + if p.suffix.lower() in {".html", ".xhtml", ".htm"} and p != out + ) + if not pages: + raise ConversionError(f"No readable text found inside {src.name}.") + + chapters = [] + for i, page in enumerate(pages, start=1): + paragraphs = _html_paragraphs(page) + if paragraphs: + chapters.append((page.stem or f"Section {i}", paragraphs)) + if not chapters: + raise ConversionError(f"No text could be extracted from {src.name}.") + + return write_epub( + out_stem.with_suffix(".epub"), src.stem, "en", chapters + ) + finally: + if tempdir: + shutil.rmtree(tempdir, ignore_errors=True) + + +def _html_paragraphs(path: Path) -> list[str]: + from bs4 import BeautifulSoup + + soup = BeautifulSoup( + path.read_text(encoding="utf-8", errors="replace"), "html.parser" + ) + for bad in soup(["script", "style"]): + bad.decompose() + + blocks = [ + el.get_text(" ", strip=True) + for el in soup.find_all(["p", "h1", "h2", "h3", "h4", "blockquote"]) + ] + blocks = [b for b in blocks if b] + if blocks: + return blocks + # Some MOBI files are one long <div> soup with no paragraph tags. + return [ln.strip() for ln in soup.get_text("\n").splitlines() if ln.strip()] + + +# --------------------------------------------------------------------------- +# epub output +# --------------------------------------------------------------------------- + +def write_epub( + dest: Path, + title: str, + language: str, + chapters: list[tuple[str, list[str]]], +) -> Path: + """Write a chaptered epub with ebooklib - one spine document per chapter.""" + from ebooklib import epub + + book = epub.EpubBook() + book.set_identifier(f"subplz-{abs(hash(title)) % 10**12}") + book.set_title(title) + book.set_language((language or "en")[:8]) + + items = [] + for i, (heading, paragraphs) in enumerate(chapters, start=1): + doc = epub.EpubHtml( + title=heading, file_name=f"chapter{i:04d}.xhtml", lang=language + ) + body = "\n".join(f"<p>{html.escape(p)}</p>" for p in paragraphs) + doc.content = f"<h1>{html.escape(heading)}</h1>\n{body}" + book.add_item(doc) + items.append(doc) + + book.toc = tuple(items) + book.spine = ["nav", *items] + book.add_item(epub.EpubNcx()) + book.add_item(epub.EpubNav()) + + dest.parent.mkdir(parents=True, exist_ok=True) + epub.write_epub(str(dest), book) + return dest diff --git a/backend/db.py b/backend/db.py new file mode 100644 index 0000000..da6f9fe --- /dev/null +++ b/backend/db.py @@ -0,0 +1,158 @@ +"""Persistence. SQLite on localhost, Postgres in production - same models.""" + +import enum +import secrets +from datetime import datetime, timezone + +from sqlalchemy import ( + BigInteger, + DateTime, + Enum, + Float, + ForeignKey, + Integer, + String, + Text, + create_engine, +) +from sqlalchemy.orm import ( + DeclarativeBase, + Mapped, + mapped_column, + relationship, + sessionmaker, +) + +from .settings import settings + + +def utcnow() -> datetime: + return datetime.now(timezone.utc) + + +def new_id(prefix: str) -> str: + return f"{prefix}_{secrets.token_hex(8)}" + + +class Base(DeclarativeBase): + pass + + +class JobStatus(str, enum.Enum): + # Uploaded and analysed, but not started: the user still gets to correct the + # detected language before committing to an hours-long alignment. + draft = "draft" + queued = "queued" + running = "running" + succeeded = "succeeded" + failed = "failed" + canceled = "canceled" + + +class Account(Base): + """One row per identified user. + + On localhost everyone shares a single anonymous account. For the public + release this gains an auth provider id, an email and a Stripe customer id - + billing.py already reads its allowance from here. + """ + + __tablename__ = "accounts" + + id: Mapped[str] = mapped_column( + String(64), primary_key=True, default=lambda: new_id("acct") + ) + # Opaque token the browser stores; becomes a real session subject later. + device_token: Mapped[str] = mapped_column(String(64), unique=True, index=True) + email: Mapped[str | None] = mapped_column(String(320), nullable=True) + stripe_customer_id: Mapped[str | None] = mapped_column(String(64), nullable=True) + # Conversions bought beyond the free allowance. + purchased_credits: Mapped[int] = mapped_column(Integer, default=0) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=utcnow + ) + + jobs: Mapped[list["Job"]] = relationship(back_populates="account") + + +class Job(Base): + __tablename__ = "jobs" + + id: Mapped[str] = mapped_column( + String(64), primary_key=True, default=lambda: new_id("job") + ) + account_id: Mapped[str] = mapped_column(ForeignKey("accounts.id"), index=True) + + status: Mapped[JobStatus] = mapped_column( + Enum(JobStatus), default=JobStatus.queued, index=True + ) + language: Mapped[str] = mapped_column(String(16)) + splitter: Mapped[str] = mapped_column(String(16)) + model: Mapped[str] = mapped_column(String(32)) + + # Display name. For a per-chapter audiobook this is "01.mp3 + 43 more". + audio_filename: Mapped[str] = mapped_column(String(512)) + text_filename: Mapped[str] = mapped_column(String(512)) + # 1 for a single file; higher when the book arrived as per-chapter parts + # that get merged before alignment. + audio_parts: Mapped[int] = mapped_column(Integer, default=1) + audio_bytes: Mapped[int] = mapped_column(BigInteger, default=0) + audio_duration_seconds: Mapped[float | None] = mapped_column(Float, nullable=True) + + progress: Mapped[float] = mapped_column(Float, default=0.0) # 0..1 + stage: Mapped[str] = mapped_column(String(128), default="queued") + error: Mapped[str | None] = mapped_column(Text, nullable=True) + + # 1 once the job has consumed the free/paid allowance. + billed: Mapped[int] = mapped_column(Integer, default=0) + + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=utcnow + ) + started_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) + finished_at: Mapped[datetime | None] = mapped_column( + DateTime(timezone=True), nullable=True + ) + + account: Mapped[Account] = relationship(back_populates="jobs") + artifacts: Mapped[list["Artifact"]] = relationship( + back_populates="job", cascade="all, delete-orphan" + ) + + +class Artifact(Base): + """A downloadable output of a job: the subtitles, metadata and run log.""" + + __tablename__ = "artifacts" + + id: Mapped[str] = mapped_column( + String(64), primary_key=True, default=lambda: new_id("art") + ) + job_id: Mapped[str] = mapped_column(ForeignKey("jobs.id"), index=True) + kind: Mapped[str] = mapped_column(String(32)) # srt | metadata | log + filename: Mapped[str] = mapped_column(String(512)) + # Opaque storage key, not a filesystem path. + storage_key: Mapped[str] = mapped_column(String(1024)) + size_bytes: Mapped[int] = mapped_column(BigInteger, default=0) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), default=utcnow + ) + + job: Mapped[Job] = relationship(back_populates="artifacts") + + +_is_sqlite = settings.resolved_database_url.startswith("sqlite") + +_engine = create_engine( + settings.resolved_database_url, + # check_same_thread only matters for SQLite plus our worker threads. + connect_args={"check_same_thread": False} if _is_sqlite else {}, + pool_pre_ping=True, +) +SessionLocal = sessionmaker(bind=_engine, expire_on_commit=False) + + +def init_db() -> None: + Base.metadata.create_all(_engine) diff --git a/backend/detect.py b/backend/detect.py new file mode 100644 index 0000000..152221c --- /dev/null +++ b/backend/detect.py @@ -0,0 +1,170 @@ +"""Work out what the user dropped on us. + +Two questions, both answered without asking the user anything: + 1. Which file is the audio and which is the book? (by extension) + 2. What language is the book in? (lingua, over text pulled from the epub) + +The detected language is a default, not a verdict - the UI always lets the user +override it, because getting this wrong wastes a long alignment run. +""" + +from __future__ import annotations + +import re +import zipfile +from dataclasses import dataclass +from functools import lru_cache +from pathlib import Path + +from . import languages +from .runner import AUDIO_SUFFIXES, TEXT_SUFFIXES + +# Enough text to be confident without reading a whole book into memory. +_SAMPLE_CHARS = 20_000 + + +class DetectionError(ValueError): + pass + + +@dataclass(frozen=True) +class Pairing: + # One entry for a single-file audiobook, many for a per-chapter set, in + # playback order. + audio_names: list[str] + text_name: str + + @property + def is_multipart(self) -> bool: + return len(self.audio_names) > 1 + + @property + def display_name(self) -> str: + if not self.is_multipart: + return self.audio_names[0] + return f"{self.audio_names[0]} + {len(self.audio_names) - 1} more" + + +def natural_key(name: str) -> tuple: + """Sort key that orders 2.mp3 before 10.mp3. + + Chapter order is playback order, and lexical sorting gets it wrong as soon + as a set passes nine files. + """ + parts = re.split(r"(\d+)", Path(name).stem) + return tuple(int(p) if p.isdigit() else p.lower() for p in parts) + + +def classify(filenames: list[str]) -> Pairing: + """Split dropped filenames into the audio part(s) and the one text file.""" + audio = [n for n in filenames if Path(n).suffix.lower() in AUDIO_SUFFIXES] + text = [n for n in filenames if Path(n).suffix.lower() in TEXT_SUFFIXES] + + if not audio: + raise DetectionError( + "No audio file found. Add an audiobook " + "(m4b, mp3, m4a, opus, flac, wav...)." + ) + if not text: + raise DetectionError( + "No book file found. Add an epub (or a txt/srt/vtt/ass script)." + ) + if len(text) > 1: + raise DetectionError(f"Got {len(text)} text files. Drop one book at a time.") + + # Mixed containers usually mean two different rips got dropped together. + suffixes = {Path(n).suffix.lower() for n in audio} + if len(suffixes) > 1: + raise DetectionError( + "The audio files are not all the same format (" + + ", ".join(sorted(suffixes)) + + "). Drop one audiobook at a time." + ) + + return Pairing( + audio_names=sorted(audio, key=natural_key), + text_name=text[0], + ) + + +def extract_text_sample(path: Path) -> str: + """Pull readable text out of an epub (or plain text file) for detection.""" + suffix = path.suffix.lower() + + if suffix != ".epub": + try: + return path.read_text(encoding="utf-8", errors="replace")[:_SAMPLE_CHARS] + except OSError as exc: + raise DetectionError(f"Could not read {path.name}: {exc}") from exc + + # Read the epub as a zip rather than via ebooklib: much faster, and it + # tolerates the malformed epubs that converted books often are. + try: + from bs4 import BeautifulSoup + + chunks: list[str] = [] + total = 0 + with zipfile.ZipFile(path) as zf: + names = [ + n for n in zf.namelist() + if n.lower().endswith((".xhtml", ".html", ".htm")) + ] + for name in sorted(names): + if total >= _SAMPLE_CHARS: + break + try: + raw = zf.read(name).decode("utf-8", errors="replace") + except (KeyError, OSError): + continue + text = BeautifulSoup(raw, "html.parser").get_text(" ", strip=True) + if text: + chunks.append(text) + total += len(text) + sample = " ".join(chunks)[:_SAMPLE_CHARS] + except zipfile.BadZipFile as exc: + raise DetectionError( + f"{path.name} is not a readable epub (bad zip archive)." + ) from exc + + if not sample.strip(): + raise DetectionError( + f"No text could be read from {path.name}. " + "If it is a scanned/image-only book, it cannot be aligned." + ) + return sample + + +@lru_cache(maxsize=1) +def _detector(): + # Built once and cached: constructing this is the expensive part. + from lingua import LanguageDetectorBuilder + + return LanguageDetectorBuilder.from_all_languages().build() + + +@dataclass(frozen=True) +class Detection: + code: str | None + name: str | None + confidence: float + supported: bool + + +def detect_language(sample: str) -> Detection: + """Best-guess ISO 639-1 code for a block of text.""" + if not sample.strip(): + return Detection(None, None, 0.0, False) + + values = _detector().compute_language_confidence_values(sample) + if not values: + return Detection(None, None, 0.0, False) + + best = values[0] + code = best.language.iso_code_639_1.name.lower() + known = languages.get(code) + return Detection( + code=code, + name=known.name if known else best.language.name.title(), + confidence=float(best.value), + supported=known is not None, + ) diff --git a/backend/gen_languages.py b/backend/gen_languages.py new file mode 100644 index 0000000..d6aa4ed --- /dev/null +++ b/backend/gen_languages.py @@ -0,0 +1,67 @@ +"""Generate languages.json from the actually-installed pysbd + stanza. + +Run once (or after upgrading either package): + python -m backend.gen_languages + +The registry is what makes "works for all languages" concrete: every language +subplz can segment, tagged with the splitter backend it needs. +""" + +import json +from pathlib import Path + +OUT = Path(__file__).with_name("languages.json") + +# Languages where subplz's default splitter (pysbd) works. Fast, no model download. +# Everything else must run with --nlp so subplz uses stanza instead. +def pysbd_languages() -> set[str]: + from pysbd.languages import LANGUAGE_CODES + + return set(LANGUAGE_CODES.keys()) + + +def stanza_languages() -> set[str]: + from stanza.resources.common import list_available_languages + + return set(list_available_languages()) + + +def english_name(code: str) -> str: + import pycountry + + for attr in ("alpha_2", "alpha_3"): + try: + hit = pycountry.languages.get(**{attr: code}) + except (KeyError, LookupError): + hit = None + if hit is not None: + return getattr(hit, "name", code) + return code + + +def build() -> dict: + pysbd = pysbd_languages() + stanza = stanza_languages() + + entries = [] + for code in sorted(pysbd | stanza): + entries.append( + { + "code": code, + "name": english_name(code), + # pysbd wins when available: no model download, much faster start. + "splitter": "pysbd" if code in pysbd else "stanza", + } + ) + return { + "generated_from": {"pysbd": len(pysbd), "stanza": len(stanza)}, + "languages": entries, + } + + +if __name__ == "__main__": + data = build() + OUT.write_text(json.dumps(data, ensure_ascii=False, indent=2), encoding="utf-8") + n = len(data["languages"]) + n_pysbd = sum(1 for e in data["languages"] if e["splitter"] == "pysbd") + print(f"wrote {OUT} - {n} languages ({n_pysbd} pysbd, {n - n_pysbd} stanza)") diff --git a/backend/languages.json b/backend/languages.json new file mode 100644 index 0000000..db458cb --- /dev/null +++ b/backend/languages.json @@ -0,0 +1,493 @@ +{ + "generated_from": { + "pysbd": 23, + "stanza": 95 + }, + "languages": [ + { + "code": "af", + "name": "Afrikaans", + "splitter": "stanza" + }, + { + "code": "am", + "name": "Amharic", + "splitter": "pysbd" + }, + { + "code": "ang", + "name": "Old English (ca. 450-1100)", + "splitter": "stanza" + }, + { + "code": "ar", + "name": "Arabic", + "splitter": "pysbd" + }, + { + "code": "be", + "name": "Belarusian", + "splitter": "stanza" + }, + { + "code": "bg", + "name": "Bulgarian", + "splitter": "pysbd" + }, + { + "code": "bn", + "name": "Bengali", + "splitter": "stanza" + }, + { + "code": "bxr", + "name": "Russia Buriat", + "splitter": "stanza" + }, + { + "code": "ca", + "name": "Catalan", + "splitter": "stanza" + }, + { + "code": "cop", + "name": "Coptic", + "splitter": "stanza" + }, + { + "code": "cs", + "name": "Czech", + "splitter": "stanza" + }, + { + "code": "cu", + "name": "Church Slavic", + "splitter": "stanza" + }, + { + "code": "cy", + "name": "Welsh", + "splitter": "stanza" + }, + { + "code": "da", + "name": "Danish", + "splitter": "pysbd" + }, + { + "code": "de", + "name": "German", + "splitter": "pysbd" + }, + { + "code": "el", + "name": "Modern Greek (1453-)", + "splitter": "pysbd" + }, + { + "code": "en", + "name": "English", + "splitter": "pysbd" + }, + { + "code": "es", + "name": "Spanish", + "splitter": "pysbd" + }, + { + "code": "et", + "name": "Estonian", + "splitter": "stanza" + }, + { + "code": "eu", + "name": "Basque", + "splitter": "stanza" + }, + { + "code": "fa", + "name": "Persian", + "splitter": "pysbd" + }, + { + "code": "fi", + "name": "Finnish", + "splitter": "stanza" + }, + { + "code": "fo", + "name": "Faroese", + "splitter": "stanza" + }, + { + "code": "fr", + "name": "French", + "splitter": "pysbd" + }, + { + "code": "fro", + "name": "Old French (842-ca. 1400)", + "splitter": "stanza" + }, + { + "code": "ga", + "name": "Irish", + "splitter": "stanza" + }, + { + "code": "gd", + "name": "Scottish Gaelic", + "splitter": "stanza" + }, + { + "code": "gl", + "name": "Galician", + "splitter": "stanza" + }, + { + "code": "got", + "name": "Gothic", + "splitter": "stanza" + }, + { + "code": "grc", + "name": "Ancient Greek (to 1453)", + "splitter": "stanza" + }, + { + "code": "gv", + "name": "Manx", + "splitter": "stanza" + }, + { + "code": "hbo", + "name": "Ancient Hebrew", + "splitter": "stanza" + }, + { + "code": "he", + "name": "Hebrew", + "splitter": "stanza" + }, + { + "code": "hi", + "name": "Hindi", + "splitter": "pysbd" + }, + { + "code": "hr", + "name": "Croatian", + "splitter": "stanza" + }, + { + "code": "hsb", + "name": "Upper Sorbian", + "splitter": "stanza" + }, + { + "code": "hu", + "name": "Hungarian", + "splitter": "stanza" + }, + { + "code": "hy", + "name": "Armenian", + "splitter": "pysbd" + }, + { + "code": "hyw", + "name": "Western Armenian", + "splitter": "stanza" + }, + { + "code": "id", + "name": "Indonesian", + "splitter": "stanza" + }, + { + "code": "is", + "name": "Icelandic", + "splitter": "stanza" + }, + { + "code": "it", + "name": "Italian", + "splitter": "pysbd" + }, + { + "code": "ja", + "name": "Japanese", + "splitter": "pysbd" + }, + { + "code": "ka", + "name": "Georgian", + "splitter": "stanza" + }, + { + "code": "kk", + "name": "Kazakh", + "splitter": "pysbd" + }, + { + "code": "kmr", + "name": "Northern Kurdish", + "splitter": "stanza" + }, + { + "code": "ko", + "name": "Korean", + "splitter": "stanza" + }, + { + "code": "kpv", + "name": "Komi-Zyrian", + "splitter": "stanza" + }, + { + "code": "ky", + "name": "Kirghiz", + "splitter": "stanza" + }, + { + "code": "la", + "name": "Latin", + "splitter": "stanza" + }, + { + "code": "lij", + "name": "Ligurian", + "splitter": "stanza" + }, + { + "code": "lt", + "name": "Lithuanian", + "splitter": "stanza" + }, + { + "code": "lv", + "name": "Latvian", + "splitter": "stanza" + }, + { + "code": "lzh", + "name": "Literary Chinese", + "splitter": "stanza" + }, + { + "code": "ml", + "name": "Malayalam", + "splitter": "stanza" + }, + { + "code": "mr", + "name": "Marathi", + "splitter": "pysbd" + }, + { + "code": "mt", + "name": "Maltese", + "splitter": "stanza" + }, + { + "code": "multilingual", + "name": "multilingual", + "splitter": "stanza" + }, + { + "code": "my", + "name": "Burmese", + "splitter": "pysbd" + }, + { + "code": "myv", + "name": "Erzya", + "splitter": "stanza" + }, + { + "code": "nb", + "name": "Norwegian Bokmål", + "splitter": "stanza" + }, + { + "code": "nds", + "name": "Low German", + "splitter": "stanza" + }, + { + "code": "nl", + "name": "Dutch", + "splitter": "pysbd" + }, + { + "code": "nn", + "name": "Norwegian Nynorsk", + "splitter": "stanza" + }, + { + "code": "or", + "name": "Oriya (macrolanguage)", + "splitter": "stanza" + }, + { + "code": "orv", + "name": "Old Russian", + "splitter": "stanza" + }, + { + "code": "ota", + "name": "Ottoman Turkish (1500-1928)", + "splitter": "stanza" + }, + { + "code": "pcm", + "name": "Nigerian Pidgin", + "splitter": "stanza" + }, + { + "code": "pl", + "name": "Polish", + "splitter": "pysbd" + }, + { + "code": "pt", + "name": "Portuguese", + "splitter": "stanza" + }, + { + "code": "qaf", + "name": "qaf", + "splitter": "stanza" + }, + { + "code": "qpm", + "name": "qpm", + "splitter": "stanza" + }, + { + "code": "qtd", + "name": "qtd", + "splitter": "stanza" + }, + { + "code": "ro", + "name": "Romanian", + "splitter": "stanza" + }, + { + "code": "ru", + "name": "Russian", + "splitter": "pysbd" + }, + { + "code": "sa", + "name": "Sanskrit", + "splitter": "stanza" + }, + { + "code": "sd", + "name": "Sindhi", + "splitter": "stanza" + }, + { + "code": "si", + "name": "Sinhala", + "splitter": "stanza" + }, + { + "code": "sk", + "name": "Slovak", + "splitter": "pysbd" + }, + { + "code": "sl", + "name": "Slovenian", + "splitter": "stanza" + }, + { + "code": "sme", + "name": "Northern Sami", + "splitter": "stanza" + }, + { + "code": "sq", + "name": "Albanian", + "splitter": "stanza" + }, + { + "code": "sr", + "name": "Serbian", + "splitter": "stanza" + }, + { + "code": "sv", + "name": "Swedish", + "splitter": "stanza" + }, + { + "code": "ta", + "name": "Tamil", + "splitter": "stanza" + }, + { + "code": "te", + "name": "Telugu", + "splitter": "stanza" + }, + { + "code": "th", + "name": "Thai", + "splitter": "stanza" + }, + { + "code": "tr", + "name": "Turkish", + "splitter": "stanza" + }, + { + "code": "ug", + "name": "Uighur", + "splitter": "stanza" + }, + { + "code": "uk", + "name": "Ukrainian", + "splitter": "stanza" + }, + { + "code": "ur", + "name": "Urdu", + "splitter": "pysbd" + }, + { + "code": "vi", + "name": "Vietnamese", + "splitter": "stanza" + }, + { + "code": "wo", + "name": "Wolof", + "splitter": "stanza" + }, + { + "code": "xcl", + "name": "Classical Armenian", + "splitter": "stanza" + }, + { + "code": "zh", + "name": "Chinese", + "splitter": "pysbd" + }, + { + "code": "zh-hans", + "name": "zh-hans", + "splitter": "stanza" + }, + { + "code": "zh-hant", + "name": "zh-hant", + "splitter": "stanza" + } + ] +} \ No newline at end of file diff --git a/backend/languages.py b/backend/languages.py new file mode 100644 index 0000000..e2f9f5c --- /dev/null +++ b/backend/languages.py @@ -0,0 +1,63 @@ +"""Language registry. + +subplz defaults to Japanese and splits sentences with pysbd, which only knows 23 +languages - pysbd raises ValueError on anything else (Portuguese and Finnish +included). subplz can use stanza instead when passed --nlp, which covers 74 more. + +So: resolve the language here, decide the splitter here, and never let a request +reach subplz with a language its splitter cannot handle. +""" + +import json +from dataclasses import dataclass +from functools import lru_cache +from pathlib import Path + +_REGISTRY = Path(__file__).with_name("languages.json") + + +@dataclass(frozen=True) +class Language: + code: str + name: str + splitter: str # "pysbd" | "stanza" + + @property + def needs_nlp_flag(self) -> bool: + """stanza languages require subplz's --nlp flag (and a one-off model download).""" + return self.splitter == "stanza" + + +@lru_cache(maxsize=1) +def _table() -> dict[str, Language]: + raw = json.loads(_REGISTRY.read_text(encoding="utf-8")) + return {e["code"]: Language(e["code"], e["name"], e["splitter"]) for e in raw["languages"]} + + +def all_languages() -> list[Language]: + # Sort by display name so the dropdown reads naturally. + return sorted(_table().values(), key=lambda l: l.name.lower()) + + +def get(code: str) -> Language | None: + return _table().get((code or "").strip().lower()) + + +def is_supported(code: str) -> bool: + return get(code) is not None + + +class UnsupportedLanguage(ValueError): + def __init__(self, code: str): + super().__init__( + f"Language {code!r} is not supported. subplz can segment " + f"{len(_table())} languages; see GET /api/languages for the list." + ) + self.code = code + + +def require(code: str) -> Language: + lang = get(code) + if lang is None: + raise UnsupportedLanguage(code) + return lang diff --git a/backend/main.py b/backend/main.py new file mode 100644 index 0000000..51a36b5 --- /dev/null +++ b/backend/main.py @@ -0,0 +1,98 @@ +"""Application entrypoint. + + python -m uvicorn backend.main:app --host 127.0.0.1 --port 8420 + +Serves the API and the static frontend from one origin, so localhost needs no +CORS config and no second server. +""" + +from __future__ import annotations + +import logging +from contextlib import asynccontextmanager + +from fastapi import FastAPI +from fastapi.responses import FileResponse +from fastapi.staticfiles import StaticFiles + +from .api import router +from .db import Job, JobStatus, SessionLocal, init_db +from .languages import all_languages +from .queue import queue +from .settings import ROOT, settings + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)-7s %(name)s: %(message)s", +) +log = logging.getLogger("subplz.web") + +FRONTEND = ROOT / "frontend" + + +def _requeue_interrupted() -> None: + """Recover jobs that were mid-flight when the server last stopped. + + An in-process queue dies with the process, so anything left `running` is + orphaned. Put it back in the queue rather than leaving a stuck progress bar. + """ + with SessionLocal() as s: + stale = ( + s.query(Job) + .filter(Job.status.in_([JobStatus.running, JobStatus.queued])) + .all() + ) + for job in stale: + job.status = JobStatus.queued + job.stage = "Queued (resumed after restart)" + job.progress = 0.0 + s.commit() + ids = [j.id for j in stale] + + for job_id in ids: + queue.enqueue(job_id) + if ids: + log.info("re-queued %d interrupted job(s)", len(ids)) + + +@asynccontextmanager +async def lifespan(_: FastAPI): + init_db() + log.info( + "subplz-web ready | %d languages | model=%s device=%s | " + "queue=%s storage=%s | billing=%s", + len(all_languages()), settings.model, settings.device, + settings.queue_backend, settings.storage_backend, + "on" if settings.billing_enabled else "off", + ) + _requeue_interrupted() + yield + queue.shutdown() + + +app = FastAPI( + title="SubPlz Web", + summary="Drag an audiobook and an epub in; get split-timed SRT subtitles out.", + version="1.0.0", + lifespan=lifespan, +) +app.include_router(router) + + +@app.get("/healthz") +def healthz(): + return { + "ok": True, + "languages": len(all_languages()), + "queue_backend": settings.queue_backend, + "queue_depth": queue.depth(), + "storage_backend": settings.storage_backend, + } + + +@app.get("/") +def index(): + return FileResponse(FRONTEND / "index.html") + + +app.mount("/", StaticFiles(directory=FRONTEND), name="static") diff --git a/backend/matching.py b/backend/matching.py new file mode 100644 index 0000000..4400dd7 --- /dev/null +++ b/backend/matching.py @@ -0,0 +1,314 @@ +"""Does this text actually match this audio? + +subplz fails late and unhelpfully when it does not: it transcribes the whole +book first, then reports "the generated transcript and the provided text file +are too different" and writes a .subfail. On CPU that is a wasted hour. + +So we ask the same question up front, cheaply: transcribe a couple of short +samples from the start of a few chapters, and score them against the book with +the backend's own rule. Same arithmetic, same threshold, ~10 seconds instead of +an hour. + +What the number means, and what it does not: a good score says the text lines +up with what is being read. A bad score usually means the wrong book, the wrong +edition (abridged vs full), an audiobook that opens with a publisher +announcement the book does not contain, or a book buried under front matter. +""" + +from __future__ import annotations + +import json +import logging +import subprocess +import zipfile +from dataclasses import dataclass, field +from functools import lru_cache +from pathlib import Path + +from .settings import settings + +log = logging.getLogger(__name__) + +# Scoring compares chapter openings, so a sample from the start of a chapter is +# the relevant thing to transcribe - and a short one is enough. subplz caps the +# comparison at 2000 characters, which is a couple of minutes of speech. +SAMPLE_SECONDS = 60 +MAX_SAMPLES = 3 + + +@dataclass +class ChapterScore: + audio_chapter: int + best_score: float + best_text_chapter: int | None + transcript_head: str = "" + + @property + def matched(self) -> bool: + return self.best_text_chapter is not None + + +@dataclass +class MatchReport: + threshold: float + scores: list[ChapterScore] = field(default_factory=list) + text_chapters: int = 0 + audio_chapters: int = 0 + verdict: str = "unknown" # good | marginal | poor | unknown + summary: str = "" + warnings: list[str] = field(default_factory=list) + skipped: str | None = None + + @property + def best(self) -> float: + return max((s.best_score for s in self.scores), default=0.0) + + @property + def worst(self) -> float: + return min((s.best_score for s in self.scores), default=0.0) + + @property + def matched(self) -> int: + return sum(1 for s in self.scores if s.matched) + + def as_dict(self) -> dict: + return { + "threshold": self.threshold, + "verdict": self.verdict, + "summary": self.summary, + "best": round(self.best, 1), + "worst": round(self.worst, 1), + "sampled": len(self.scores), + "matched": self.matched, + "text_chapters": self.text_chapters, + "audio_chapters": self.audio_chapters, + "warnings": self.warnings, + "skipped": self.skipped, + "scores": [ + { + "audio_chapter": s.audio_chapter, + "score": round(s.best_score, 1), + "text_chapter": s.best_text_chapter, + } + for s in self.scores + ], + } + + +# --------------------------------------------------------------------------- +# book text +# --------------------------------------------------------------------------- + +def book_chapters(text_path: Path) -> list[str]: + """The comparable units, split the way the backend will split them. + + This matters more than it looks. subplz gives an epub one text chapter per + spine document, but a .txt exactly one chapter for the whole file - so a + flat text file offers the matcher a single anchor no matter how many + chapters the audio has. + """ + if text_path.suffix.lower() != ".epub": + return [text_path.read_text(encoding="utf-8", errors="replace")] + + try: + from bs4 import BeautifulSoup + + chapters: list[str] = [] + with zipfile.ZipFile(text_path) as zf: + names = [ + n for n in zf.namelist() + if n.lower().endswith((".xhtml", ".html", ".htm")) + ] + for name in sorted(names): + try: + raw = zf.read(name).decode("utf-8", errors="replace") + except (KeyError, OSError): + continue + soup = BeautifulSoup(raw, "html.parser") + for bad in soup(["script", "style"]): + bad.decompose() + # subplz joins paragraph texts with no separator, so match that. + text = "".join( + p.get_text(" ", strip=True) for p in soup.find_all("p") + ) or soup.get_text(" ", strip=True) + if text.strip(): + chapters.append(text) + return chapters + except zipfile.BadZipFile: + return [] + + +# --------------------------------------------------------------------------- +# audio sampling +# --------------------------------------------------------------------------- + +def audio_chapter_starts(audio: Path) -> list[float]: + """Start time of each chapter, or [0.0] for a flat file.""" + try: + out = subprocess.run( + ["ffprobe", "-v", "error", "-show_chapters", "-print_format", "json", + str(audio)], + capture_output=True, timeout=120, check=True, + ) + chapters = json.loads(out.stdout.decode("utf-8")).get("chapters", []) + starts = [float(c["start_time"]) for c in chapters] + return starts or [0.0] + except (subprocess.SubprocessError, ValueError, OSError, KeyError): + return [0.0] + + +def read_samples(audio: Path, start: float, seconds: int): + """`seconds` of audio from `start`, as the float32 mono 16 kHz the model wants.""" + import numpy as np + + proc = subprocess.run( + [ + "ffmpeg", "-v", "error", + "-max_error_rate", "1.0", "-err_detect", "ignore_err", + "-ss", f"{start:.3f}", "-t", str(seconds), "-i", str(audio), + "-map", "0:a:0", "-f", "s16le", "-acodec", "pcm_s16le", + "-ac", "1", "-ar", "16000", "-", + ], + capture_output=True, timeout=300, + ) + if proc.returncode != 0 or not proc.stdout: + return None + return np.frombuffer(proc.stdout, np.int16).astype(np.float32) / 32768.0 + + +@lru_cache(maxsize=1) +def _model(): + """The same tiny model the alignment itself uses, loaded once.""" + from faster_whisper import WhisperModel + + return WhisperModel(settings.model, device="cpu", compute_type="int8") + + +def transcribe_sample(samples, language: str) -> str: + segments, _ = _model().transcribe( + samples, language=language, beam_size=1, without_timestamps=True + ) + return "".join(seg.text for seg in segments) + + +# --------------------------------------------------------------------------- +# the check +# --------------------------------------------------------------------------- + +def check(audio: Path, text: Path, language: str, aligner) -> MatchReport: + """Score a sample of the audio against the book. Never raises.""" + report = MatchReport(threshold=aligner.match_threshold) + + if not settings.match_check: + report.skipped = "disabled" + return report + + try: + chapters = book_chapters(text) + report.text_chapters = len(chapters) + if not chapters: + report.verdict = "poor" + report.summary = "No readable text could be found in the book." + return report + + starts = audio_chapter_starts(audio) + report.audio_chapters = len(starts) + + if len(chapters) == 1 and len(starts) > 1: + report.warnings.append( + f"The book is one flat document but the audio has " + f"{len(starts)} chapters. Alignment still works, but it has " + f"only one place to anchor - an epub with real chapters aligns " + f"more reliably." + ) + + # Sample from the start of chapters spread across the book, because the + # score is about how chapters open. + picks = _spread(len(starts), MAX_SAMPLES) + scorer = ( + aligner.for_language(language) + if hasattr(aligner, "for_language") else aligner + ) + + for idx in picks: + samples = read_samples(audio, starts[idx], SAMPLE_SECONDS) + if samples is None or len(samples) < 16000: + continue + transcript = transcribe_sample(samples, language) + if len(transcript.strip()) < 40: + continue + + best, best_i = 0.0, None + for ci, chapter in enumerate(chapters): + score = scorer.score_pair(transcript, chapter) + if score > best: + best, best_i = score, ci + + report.scores.append( + ChapterScore( + audio_chapter=idx, + best_score=best, + best_text_chapter=best_i if best > report.threshold else None, + transcript_head=transcript.strip()[:160], + ) + ) + + _verdict(report) + return report + + except Exception as exc: # noqa: BLE001 - a check must never block an upload + log.warning("match check failed: %s", exc, exc_info=True) + report.skipped = str(exc) + return report + + +def _spread(n: int, k: int) -> list[int]: + """Up to k indices spread across range(n), always including the first.""" + if n <= k: + return list(range(n)) + step = n / k + return sorted({min(n - 1, int(i * step)) for i in range(k)}) + + +def _verdict(report: MatchReport) -> None: + """Turn the raw score into advice. + + The bands sit well above the backend's own threshold, deliberately. + Measured on real files: the right book scored 74, while a completely + unrelated Russian novel still scored 40.5 - barely clearing subplz's + threshold of 40. Two prose texts in the same language are roughly 40% + similar character-by-character whatever they say, so "just over the + threshold" means "probably wrong", not "probably fine". + """ + if not report.scores: + report.verdict = "unknown" + report.summary = "Could not sample enough audio to check the match." + return + + t = report.threshold or 40.0 + good_at, weak_at = t * 1.5, t * 1.2 # 60 and 48 for subplz + best = report.best + total = len(report.scores) + strong = sum(1 for s in report.scores if s.best_score >= good_at) + + if best >= good_at: + report.verdict = "good" + report.summary = f"Text and audio line up ({best:.0f}/100 on the best sample)." + if total > 1 and strong < total: + report.warnings.append( + f"Only {strong} of {total} sampled chapters scored well. The " + f"timing may drift in parts of the book." + ) + elif best >= weak_at: + report.verdict = "marginal" + report.summary = ( + f"Weak match ({best:.0f}/100). This may be a different edition or " + f"an abridgement. Alignment can still work, but expect drift." + ) + else: + report.verdict = "poor" + report.summary = ( + f"This text does not look like this audio ({best:.0f}/100). Two " + f"unrelated books in the same language score about this well, so " + f"it is probably the wrong book, edition or abridgement." + ) diff --git a/backend/pricing.py b/backend/pricing.py new file mode 100644 index 0000000..53979e0 --- /dev/null +++ b/backend/pricing.py @@ -0,0 +1,130 @@ +"""What we charge, and why. + +Market as of 2026: + + Direct competitors (audiobook <-> ebook sync) + Voxlight $29.99/year, but alignment runs on the user's own Mac + Spokt freemium, metered sync minutes + Storyteller free, self-hosted; you administer a server + syncabook free CLI + + Adjacent (AI subtitling SaaS, priced for transcription) + Otter $16.99/mo + Happy Scribe $17/mo for 120 AI minutes, $89/mo for 6,000 + Veed $22/mo + Kapwing/Descript $24/mo + Sonix $10/hour pay-as-you-go + +The adjacent tools price per *minute of transcription*, which does not transfer: +a 10-hour audiobook is 600 minutes, so it would cost ~$100 at Sonix's rate and +need Happy Scribe's $89 tier. Forced alignment is much cheaper to run than +transcription - the model is tiny and its output is thrown away, since the +subtitle text comes from the user's own book. + +So we price per book, cheap enough to be an impulse buy, and land the +subscription just under the general subtitling tools: + + free 1 book / 24h + single book $3.49 + 5-book pack $12.99 ($2.60/book) + 20-book pack $39.99 ($2.00/book) + unlimited month $14.99 (under Otter/Happy Scribe, well under Veed/Kapwing) + +Every number is overridable by env var; these are defaults, not decisions cast +in code. +""" + +from __future__ import annotations + +import json +import os +from dataclasses import asdict, dataclass + +from .settings import settings + + +@dataclass(frozen=True) +class Plan: + id: str + name: str + # None for the subscription, which is not a credit pack. + credits: int | None + price_cents: int + currency: str = "usd" + recurring: bool = False + blurb: str = "" + + @property + def price_display(self) -> str: + return f"${self.price_cents / 100:,.2f}" + + @property + def per_book_cents(self) -> int | None: + if not self.credits: + return None + return round(self.price_cents / self.credits) + + +DEFAULT_PLANS: list[Plan] = [ + Plan( + id="single", + name="One book", + credits=1, + price_cents=349, + blurb="One more conversion, whenever you need it.", + ), + Plan( + id="pack5", + name="5 books", + credits=5, + price_cents=1299, + blurb="$2.60 a book. Credits never expire.", + ), + Plan( + id="pack20", + name="20 books", + credits=20, + price_cents=3999, + blurb="$2.00 a book. For working through a series.", + ), + Plan( + id="unlimited", + name="Unlimited monthly", + credits=None, + price_cents=1499, + recurring=True, + blurb="As many books as you like. Cancel any time.", + ), +] + + +def plans() -> list[Plan]: + """The catalogue, overridable with SUBPLZ_WEB_PLANS_JSON.""" + raw = os.environ.get("SUBPLZ_WEB_PLANS_JSON") + if not raw: + return DEFAULT_PLANS + try: + return [Plan(**item) for item in json.loads(raw)] + except (json.JSONDecodeError, TypeError) as exc: + raise RuntimeError(f"SUBPLZ_WEB_PLANS_JSON is not valid: {exc}") from exc + + +def get(plan_id: str) -> Plan | None: + return next((p for p in plans() if p.id == plan_id), None) + + +def as_dicts() -> list[dict]: + out = [] + for p in plans(): + d = asdict(p) + d["price_display"] = p.price_display + d["per_book_cents"] = p.per_book_cents + out.append(d) + return out + + +def free_tier_summary() -> str: + n = settings.free_conversions + hours = settings.free_window_hours + book = "book" if n == 1 else "books" + return f"{n} free {book} every {hours} hours" diff --git a/backend/queue.py b/backend/queue.py new file mode 100644 index 0000000..dc0bf82 --- /dev/null +++ b/backend/queue.py @@ -0,0 +1,102 @@ +"""Job dispatch. + +Localhost runs jobs on a small thread pool inside the API process. The public +deployment sets SUBPLZ_WEB_QUEUE_BACKEND=redis and runs `worker.py` on separate +machines; the API side then only enqueues. Same `enqueue(job_id)` call either way. +""" + +from __future__ import annotations + +import logging +import threading +from abc import ABC, abstractmethod +from concurrent.futures import ThreadPoolExecutor + +from .settings import settings + +log = logging.getLogger(__name__) + + +class JobQueue(ABC): + @abstractmethod + def enqueue(self, job_id: str) -> None: + ... + + @abstractmethod + def depth(self) -> int: + """Jobs waiting or running. Drives the 'N ahead of you' hint in the UI.""" + + def shutdown(self) -> None: + ... + + +class InProcessQueue(JobQueue): + """Thread pool in the API process. Fine for one person on one machine. + + Alignment is CPU-bound and subplz is a subprocess, so the GIL is not the + limit here - max_concurrent_jobs is. + """ + + def __init__(self, workers: int): + self._pool = ThreadPoolExecutor( + max_workers=workers, thread_name_prefix="subplz-job" + ) + self._lock = threading.Lock() + self._pending: set[str] = set() + + def enqueue(self, job_id: str) -> None: + with self._lock: + if job_id in self._pending: + return + self._pending.add(job_id) + self._pool.submit(self._run, job_id) + + def _run(self, job_id: str) -> None: + # Imported here to avoid a circular import at module load. + from .runner import run_job + + try: + run_job(job_id) + except Exception: + log.exception("job %s crashed outside the runner", job_id) + finally: + with self._lock: + self._pending.discard(job_id) + + def depth(self) -> int: + with self._lock: + return len(self._pending) + + def shutdown(self) -> None: + self._pool.shutdown(wait=False, cancel_futures=True) + + +class RedisQueue(JobQueue): + """Public-release backend. Needs `rq` and a reachable Redis.""" + + def __init__(self, url: str): + from redis import Redis # lazy: localhost needs neither package + from rq import Queue as RQQueue + + self._conn = Redis.from_url(url) + self._q = RQQueue("subplz-jobs", connection=self._conn) + + def enqueue(self, job_id: str) -> None: + self._q.enqueue( + "backend.runner.run_job", + job_id, + job_timeout=settings.job_timeout_seconds, + result_ttl=86400, + ) + + def depth(self) -> int: + return self._q.count + self._q.started_job_registry.count + + +def build_queue() -> JobQueue: + if settings.queue_backend == "redis": + return RedisQueue(settings.redis_url) + return InProcessQueue(settings.max_concurrent_jobs) + + +queue: JobQueue = build_queue() diff --git a/backend/render.py b/backend/render.py new file mode 100644 index 0000000..b9172f1 --- /dev/null +++ b/backend/render.py @@ -0,0 +1,231 @@ +"""Render a YouTube-ready MP4: cover image, the audio, and a soft subtitle track. + +The subtitles are a selectable track rather than burned into the picture, so the +file stays small, the encode stays fast, and the viewer can turn captions off. +YouTube reads the track on upload; the .srt is also offered on its own for +people who would rather attach it there. + +The encode is cheap by construction: one still frame per second over a canvas +that is scaled once up front, so ffmpeg never recomputes the scale per frame. +""" + +from __future__ import annotations + +import logging +import subprocess +import zipfile +from functools import lru_cache +from pathlib import Path + +from .settings import settings + +log = logging.getLogger(__name__) + +# Cover art we are willing to pull out of an epub. +_IMAGE_SUFFIXES = {".jpg", ".jpeg", ".png", ".webp"} + +# H.264 encoders in order of preference. ffmpeg builds vary wildly in what they +# include - libx264 is the usual default but is absent from plenty of builds, +# so pick from what is actually compiled in rather than assuming. +_H264_PREFERENCE = [ + "libx264", # best quality per bit, most common + "h264_nvenc", # NVIDIA + "h264_qsv", # Intel Quick Sync + "h264_amf", # AMD + "libopenh264", # software fallback, always safe + "h264_mf", # Windows Media Foundation +] + +# Encoder-specific quality flags. A libx264 -crf/-preset means nothing to the +# others and makes ffmpeg fail outright. +_ENCODER_FLAGS: dict[str, list[str]] = { + "libx264": ["-preset", "veryfast", "-tune", "stillimage"], + "libopenh264": ["-b:v", "1M"], + "h264_nvenc": ["-preset", "p4", "-tune", "ll"], + "h264_qsv": ["-preset", "veryfast"], + "h264_amf": ["-quality", "speed", "-rc", "cqp"], + "h264_mf": [], +} + + +class RenderError(RuntimeError): + pass + + +@lru_cache(maxsize=1) +def available_encoders() -> set[str]: + try: + out = subprocess.run( + ["ffmpeg", "-hide_banner", "-encoders"], + capture_output=True, timeout=60, + ).stdout.decode("utf-8", errors="replace") + except (subprocess.SubprocessError, OSError): + return set() + names = set() + for line in out.splitlines(): + parts = line.split() + # Rows look like " V....D libopenh264 OpenH264 ..." - the name is [1]. + if len(parts) >= 2 and parts[0].startswith("V"): + names.add(parts[1]) + return names + + +def encoder_works(name: str) -> bool: + """Actually encode one frame with `name`. + + Being listed by `ffmpeg -encoders` only means it was compiled in. The + hardware encoders are listed on machines with no such hardware and fail at + runtime, so the only trustworthy check is to run one. + """ + cmd = [ + "ffmpeg", "-hide_banner", "-v", "error", + "-f", "lavfi", "-i", "color=c=black:s=320x240:d=0.1", + "-frames:v", "1", "-c:v", name, + *_ENCODER_FLAGS.get(name, []), + "-pix_fmt", "yuv420p", "-f", "null", "-", + ] + try: + return subprocess.run(cmd, capture_output=True, timeout=60).returncode == 0 + except (subprocess.SubprocessError, OSError): + return False + + +@lru_cache(maxsize=1) +def resolve_encoder() -> str: + """The H.264 encoder to use, honouring an explicit setting if it works.""" + have = available_encoders() + configured = settings.video_encoder + + if configured and configured != "auto": + if encoder_works(configured): + return configured + log.warning( + "video_encoder=%r does not work on this machine; falling back", + configured, + ) + + for name in _H264_PREFERENCE: + if name in have and encoder_works(name): + log.info("using video encoder %s", name) + return name + + raise RenderError( + "No working H.264 encoder found (tried " + + ", ".join(n for n in _H264_PREFERENCE if n in have) + + "). Install an ffmpeg with libx264 or libopenh264, or set " + "SUBPLZ_WEB_RENDER_VIDEO=false." + ) + + +def extract_cover(book: Path, dest_dir: Path) -> Path | None: + """Pull the largest image out of an epub, as a stand-in for cover art. + + "Largest" beats parsing the OPF metadata for this purpose: the cover is + almost always the biggest image, and this still works on the malformed + epubs that converted books often are. + """ + if book.suffix.lower() != ".epub": + return None + try: + with zipfile.ZipFile(book) as zf: + images = [ + info for info in zf.infolist() + if Path(info.filename).suffix.lower() in _IMAGE_SUFFIXES + ] + if not images: + return None + best = max(images, key=lambda i: i.file_size) + if best.file_size < 1024: # a bullet or a rule, not a cover + return None + dest = dest_dir / f"cover{Path(best.filename).suffix.lower()}" + dest.parent.mkdir(parents=True, exist_ok=True) + dest.write_bytes(zf.read(best.filename)) + return dest + except (zipfile.BadZipFile, OSError, KeyError): + return None + + +def build_canvas(cover: Path | None, dest: Path) -> Path: + """Scale the cover onto a fixed canvas once, ahead of the encode.""" + w, h = settings.video_width, settings.video_height + dest.parent.mkdir(parents=True, exist_ok=True) + + if cover is not None and cover.exists(): + cmd = [ + "ffmpeg", "-hide_banner", "-v", "error", "-y", + "-i", str(cover), + "-vf", + f"scale={w}:{h}:force_original_aspect_ratio=decrease," + f"pad={w}:{h}:(ow-iw)/2:(oh-ih)/2:color=black", + "-frames:v", "1", + str(dest), + ] + else: + # No cover: a plain dark card still gives YouTube a valid video stream. + cmd = [ + "ffmpeg", "-hide_banner", "-v", "error", "-y", + "-f", "lavfi", "-i", f"color=c=0x16161a:s={w}x{h}", + "-frames:v", "1", + str(dest), + ] + + proc = subprocess.run(cmd, capture_output=True, timeout=300) + if proc.returncode != 0 or not dest.exists(): + err = proc.stderr.decode("utf-8", errors="replace").strip().splitlines() + raise RenderError( + "could not prepare the cover image: " + + (" | ".join(err[-2:]) if err else "no output") + ) + return dest + + +def render_video( + audio: Path, + subtitles: Path, + canvas: Path, + dest: Path, + duration: float | None = None, +) -> Path: + """Mux canvas + audio + subtitles into an MP4.""" + dest.parent.mkdir(parents=True, exist_ok=True) + + cmd = [ + "ffmpeg", "-hide_banner", "-v", "error", "-y", + "-loop", "1", "-r", str(settings.video_fps), "-i", str(canvas), + "-i", str(audio), + # Declare the format: ffmpeg will not always sniff an srt correctly. + "-f", "srt", "-i", str(subtitles), + ] + if duration: + cmd += ["-t", f"{duration:.3f}"] + else: + cmd += ["-shortest"] + + encoder = resolve_encoder() + cmd += [ + "-map", "0:v", "-map", "1:a", "-map", "2:s", + # Carry chapter marks through; harmless where they are ignored. + "-map_chapters", "1", + "-c:v", encoder, + "-pix_fmt", "yuv420p", + "-c:a", "aac", "-b:a", "128k", + # mov_text is the only subtitle codec MP4 carries. + "-c:s", "mov_text", + # Lets a player start without reading the whole file first. + "-movflags", "+faststart", + ] + cmd += _ENCODER_FLAGS.get(encoder, []) + # -crf is a libx264/libx265 concept; the others have their own rate control. + if encoder == "libx264": + cmd += ["-crf", str(settings.video_crf)] + + cmd.append(str(dest)) + + proc = subprocess.run( + cmd, capture_output=True, timeout=settings.job_timeout_seconds + ) + if proc.returncode != 0 or not dest.exists(): + err = proc.stderr.decode("utf-8", errors="replace").strip().splitlines() + tail = " | ".join(err[-3:]) if err else "no output" + raise RenderError(f"ffmpeg could not build the video: {tail}") + return dest diff --git a/backend/runner.py b/backend/runner.py new file mode 100644 index 0000000..6f2fd2e --- /dev/null +++ b/backend/runner.py @@ -0,0 +1,618 @@ +"""Runs one job: stage inputs, drive the aligner, collect the artifacts. + +Nothing here knows which alignment backend is in use. The command to run, how +to read its progress and where it writes the subtitles all come from +`aligner.Aligner`, so swapping subplz out does not touch this file. + +The backend is always a subprocess. That boundary is also what lets the public +deployment move alignment onto separate GPU workers without code changes. +""" + +from __future__ import annotations + +import json +import logging +import os +import re +import shutil +import subprocess +import time +from dataclasses import dataclass +from datetime import timezone +from pathlib import Path + +from . import languages, render +from .aligner import AlignRequest, aligner +from .db import Artifact, Job, JobStatus, SessionLocal, utcnow +from .settings import settings +from .storage import storage + +log = logging.getLogger(__name__) + +# Audio containers subplz/ffmpeg handle. Checked at upload time. +AUDIO_SUFFIXES = { + ".m4b", ".m4a", ".mp3", ".opus", ".ogg", ".oga", ".flac", ".wav", + ".aac", ".wma", ".mka", ".mkv", ".mp4", ".webm", ".avi", ".mov", +} +# Formats subplz reads directly, plus the ones convert.py turns into epub on +# the way in (fb2/mobi/azw3). +TEXT_SUFFIXES = { + ".epub", ".txt", ".srt", ".vtt", ".ass", + ".fb2", ".fb2.zip", ".mobi", ".azw", ".azw3", ".prc", +} + + +def text_suffix(name: str) -> str: + """Suffix used when staging a text file, keeping the .fb2.zip double.""" + lowered = name.lower() + if lowered.endswith(".fb2.zip"): + return ".fb2.zip" + return Path(lowered).suffix + +# Fixed stem for staged inputs: keeps non-ASCII filenames out of the subprocess +# command line entirely, and makes the output path deterministic. +STAGE_STEM = "source" + + +class JobFailed(RuntimeError): + pass + + +@dataclass +class Paths: + root: Path + inp: Path + out: Path + + @classmethod + def for_job(cls, job_id: str) -> "Paths": + root = settings.data_dir / "work" / job_id + return cls(root=root, inp=root / "input", out=root / "out") + + def create(self) -> None: + self.inp.mkdir(parents=True, exist_ok=True) + self.out.mkdir(parents=True, exist_ok=True) + + +def input_prefix(job_id: str) -> str: + """Storage prefix mirroring a job's staged `input/` tree. + + Only used when the queue is external: the worker that runs the job is + probably not the machine that received the upload. + """ + return f"{job_id}/input" + + +def ensure_inputs(job_id: str) -> None: + """Make sure the staged inputs are on this machine's disk.""" + paths = Paths.for_job(job_id) + paths.create() + if any(p.is_file() for p in paths.inp.rglob("*")): + return # this process received the upload + + prefix = input_prefix(job_id) + keys = storage.list_prefix(prefix) + if not keys: + raise JobFailed( + "the uploaded files for this job are no longer available. " + "Upload them again." + ) + for key in keys: + storage.fetch_to(key, paths.inp / key[len(prefix) + 1 :]) + + +def staged_audio_path(job_id: str, original_name: str) -> Path: + return Paths.for_job(job_id).inp / f"{STAGE_STEM}{Path(original_name).suffix.lower()}" + + +def staged_text_path(job_id: str, original_name: str) -> Path: + return Paths.for_job(job_id).inp / f"{STAGE_STEM}{text_suffix(original_name)}" + + +def staged_part_path(job_id: str, index: int, original_name: str) -> Path: + """Where chapter file `index` (1-based) of a multi-part audiobook is staged. + + Parts get plain numeric names so the ffmpeg concat list never needs quoting + and playback order is unambiguous. + """ + inp = Paths.for_job(job_id).inp + return inp / "parts" / f"{index:04d}{Path(original_name).suffix.lower()}" + + +# Tolerate damaged input rather than aborting the whole book. ffmpeg's default +# max_error_rate (0.667) kills a run over one bad chapter. +_FFMPEG_TOLERANT = [ + "-max_error_rate", "1.0", + "-err_detect", "ignore_err", + "-fflags", "+discardcorrupt", +] + + +def merge_parts(job_id: str, parts: list[Path]) -> Path: + """Join a per-chapter audiobook into one file, one chapter per part. + + Keeping chapter marks matters twice over: subplz processes an m4b chapter by + chapter, and the runner reads chapter completions to drive progress. + """ + paths = Paths.for_job(job_id) + dest = paths.inp / f"{STAGE_STEM}.m4b" + if dest.exists(): + return dest + + durations: list[float] = [] + for p in parts: + d = probe_duration(p) + if d is None: + raise JobFailed(f"could not read the duration of part {p.name}") + durations.append(d) + + scratch = paths.root / "merge" + scratch.mkdir(parents=True, exist_ok=True) + + listing = scratch / "parts.txt" + # Absolute paths: ffmpeg resolves a relative entry against the directory of + # the list file, not the working directory. + listing.write_text( + "".join(f"file '{p.resolve().as_posix()}'\n" for p in parts), + encoding="utf-8", + ) + + # Chapter marks at the part boundaries, in milliseconds. + meta = [";FFMETADATA1"] + start_ms = 0 + for i, seconds in enumerate(durations, start=1): + end_ms = start_ms + int(round(seconds * 1000)) + meta += [ + "[CHAPTER]", + "TIMEBASE=1/1000", + f"START={start_ms}", + f"END={end_ms}", + f"title=Part {i:02d}", + ] + start_ms = end_ms + metadata = scratch / "chapters.txt" + metadata.write_text("\n".join(meta) + "\n", encoding="utf-8") + + cmd = [ + "ffmpeg", "-hide_banner", "-v", "error", "-y", + *_FFMPEG_TOLERANT, + "-f", "concat", "-safe", "0", "-i", str(listing), + "-i", str(metadata), + "-map", "0:a:0", "-map_metadata", "1", "-map_chapters", "1", + # Drop cover art and any stray tracks the parts carry, so the only + # media stream is the audio. (MP4 still writes a bin_data chapter + # track; that one is how the container stores chapters and must stay.) + "-vn", "-sn", + # Keep listenable quality: the aligner downsamples to 16 kHz mono itself + # when it reads the file, but the rendered video uses this same audio, + # and 16 kHz mono would sound awful on YouTube. + "-c:a", "aac", "-b:a", "128k", + str(dest), + ] + proc = subprocess.run(cmd, capture_output=True, timeout=settings.job_timeout_seconds) + if proc.returncode != 0 or not dest.exists(): + err = proc.stderr.decode("utf-8", errors="replace").strip().splitlines() + tail = " | ".join(err[-3:]) if err else "no stderr" + raise JobFailed(f"could not join the {len(parts)} audio parts: {tail}") + + shutil.rmtree(scratch, ignore_errors=True) + return dest + + +def normalize_for_alignment(job_id: str, source: Path) -> Path: + """Produce the 16 kHz mono copy the aligner gets fed. + + Two jobs at once: + + * Format. The acoustic model consumes 16 kHz mono no matter what, and + handing subplz a 44.1 kHz stereo file has crashed ctranslate2 here + (integer divide by zero, part-way through a chapter). Doing the + conversion ourselves, once, keeps the aligner on the input that works. + * Repair. The error-tolerant flags drop corrupt frames rather than letting + ffmpeg abort the run, which is what a damaged chapter would otherwise do. + + The video keeps using `source`, which stays at listenable quality. + """ + dest = Paths.for_job(job_id).inp / f"{STAGE_STEM}.align.m4a" + if dest.exists(): + return dest + + cmd = [ + "ffmpeg", "-hide_banner", "-v", "error", "-y", + *_FFMPEG_TOLERANT, + "-i", str(source), + "-map", "0:a:0", "-map_chapters", "0", + "-vn", "-sn", + "-c:a", "aac", "-b:a", "64k", "-ar", "16000", "-ac", "1", + str(dest), + ] + proc = subprocess.run( + cmd, capture_output=True, timeout=settings.job_timeout_seconds + ) + if proc.returncode != 0 or not dest.exists(): + err = proc.stderr.decode("utf-8", errors="replace").strip().splitlines() + tail = " | ".join(err[-3:]) if err else "no stderr" + raise JobFailed(f"could not prepare the audio for alignment: {tail}") + return dest + + +def probe_chapters(path: Path) -> int: + """Chapter count, or 1 for a flat file. + + subplz processes chaptered m4b files one chapter at a time, so this is what + turns a stream of per-chapter progress bars into an overall percentage. + """ + try: + out = subprocess.run( + ["ffprobe", "-v", "error", "-show_chapters", "-print_format", "json", + str(path)], + capture_output=True, text=True, timeout=120, check=True, + ) + return max(1, len(json.loads(out.stdout).get("chapters", []))) + except (subprocess.SubprocessError, ValueError, OSError): + return 1 + + +def probe_duration(path: Path) -> float | None: + """Audio length in seconds, via ffprobe. Used for progress and metadata.""" + try: + out = subprocess.run( + [ + "ffprobe", "-v", "error", + "-show_entries", "format=duration", + "-of", "default=noprint_wrappers=1:nokey=1", + str(path), + ], + capture_output=True, text=True, timeout=120, check=True, + ) + return float(out.stdout.strip()) + except (subprocess.SubprocessError, ValueError, FileNotFoundError, OSError): + return None + + +def _staged(inp: Path, suffixes: set[str]) -> Path: + """The one staged input with a suffix in `suffixes`.""" + for p in sorted(inp.iterdir()): + if p.is_file() and p.suffix.lower() in suffixes: + return p + raise JobFailed(f"no staged input found with a supported extension in {inp.name}") + + +def build_request(job: Job, paths: Paths, audio: Path, chapters: int) -> AlignRequest: + """Describe the job in backend-neutral terms.""" + languages.require(job.language) # reject an unsupported code before we run + return AlignRequest( + audio=audio, + text=_staged(paths.inp, TEXT_SUFFIXES), + out_dir=paths.out, + language=job.language, + model=job.model, + device=settings.device, + threads=settings.threads or max(1, (os.cpu_count() or 4) - 1), + chapters=chapters, + ) + + +def _set(job_id: str, **fields) -> None: + with SessionLocal() as s: + job = s.get(Job, job_id) + if job is None: + return + for k, v in fields.items(): + setattr(job, k, v) + s.commit() + + +def _is_canceled(job_id: str) -> bool: + with SessionLocal() as s: + job = s.get(Job, job_id) + return job is not None and job.status == JobStatus.canceled + + +def run_job(job_id: str) -> None: + """Execute one job end to end. Never raises; failures land on the Job row.""" + with SessionLocal() as s: + job = s.get(Job, job_id) + if job is None: + return + if job.status == JobStatus.canceled: + return + audio_name, text_name = job.audio_filename, job.text_filename + parts_count = job.audio_parts or 1 + + paths = Paths.for_job(job_id) + log_path = paths.root / "subplz.log" + _set(job_id, status=JobStatus.running, started_at=utcnow(), + stage="Preparing", progress=0.01) + + try: + # No-op when this process received the upload; pulls from shared storage + # when the job was enqueued on another machine. + ensure_inputs(job_id) + + if parts_count > 1: + _set(job_id, stage=f"Joining {parts_count} audio parts", progress=0.02) + parts = sorted(p for p in (paths.inp / "parts").iterdir() if p.is_file()) + if len(parts) != parts_count: + raise JobFailed( + f"expected {parts_count} audio parts but found {len(parts)}" + ) + audio_in = merge_parts(job_id, parts) + else: + audio_in = staged_audio_path(job_id, audio_name) + + if not audio_in.exists(): + raise JobFailed(f"staged audio missing: {audio_in.name}") + + duration = probe_duration(audio_in) + chapters = probe_chapters(audio_in) + + # The aligner gets a normalised 16 kHz mono copy; the video keeps the + # full-quality one. + _set(job_id, stage="Preparing audio", progress=0.03) + align_audio = normalize_for_alignment(job_id, audio_in) + + _set(job_id, audio_duration_seconds=duration, + stage=f"Starting {aligner.name}", progress=0.04) + + with SessionLocal() as s: + job = s.get(Job, job_id) + request = build_request(job, paths, align_audio, chapters) + + returncode = _stream_aligner(job_id, request, log_path) + + if _is_canceled(job_id): + _cleanup_inputs(paths) + return + + if returncode != 0: + raise JobFailed( + f"{aligner.name} exited with code {returncode}. " + f"See the run log for details." + ) + + _set(job_id, stage="Collecting output", progress=0.92) + _collect_artifacts(job_id, paths, log_path, duration, request, audio_in) + + _set(job_id, status=JobStatus.succeeded, stage="Done", progress=1.0, + finished_at=utcnow()) + # Only on success: a failed job keeps its inputs so it can be retried + # without re-uploading hundreds of megabytes. + _cleanup_inputs(paths) + + except Exception as exc: # noqa: BLE001 - surfaced to the user on the Job row + if _is_canceled(job_id): + _cleanup_inputs(paths) + return + # Keep the log even on failure: it is the only way to debug an alignment. + try: + _store_log(job_id, log_path) + except Exception: + pass + _set(job_id, status=JobStatus.failed, error=str(exc)[:4000], + stage="Failed", finished_at=utcnow()) + + +def _stream_aligner(job_id: str, request: AlignRequest, log_path: Path) -> int: + """Run the aligner, mirroring its output to a log file and to job progress. + + The backend decides what its output means; this only moves bytes and + watches for cancellation and the time limit. + """ + cmd = aligner.build_command(request) + reader = aligner.progress_reader(request) + + log_path.parent.mkdir(parents=True, exist_ok=True) + deadline = time.monotonic() + settings.job_timeout_seconds + + with log_path.open("w", encoding="utf-8", errors="replace") as log: + log.write("$ " + " ".join(cmd) + "\n\n") + log.flush() + + proc = subprocess.Popen( + cmd, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + encoding="utf-8", + errors="replace", + bufsize=1, + env=aligner.environment(), + cwd=str(log_path.parent), + ) + + try: + assert proc.stdout is not None + for raw in proc.stdout: + log.write(raw) + log.flush() + + if time.monotonic() > deadline: + proc.kill() + raise JobFailed( + f"job exceeded the {settings.job_timeout_seconds}s time limit" + ) + + if _is_canceled(job_id): + proc.kill() + return proc.wait() + + line = raw.strip() + if not line: + continue + + update = reader.feed(line) + if update is None: + continue + fields: dict = {"stage": update.stage} + if update.fraction is not None: + fields["progress"] = update.fraction + _set(job_id, **fields) + finally: + if proc.poll() is None: + proc.kill() + proc.wait() + + return proc.returncode + + +def _collect_artifacts(job_id: str, paths: Paths, log_path: Path, + duration: float | None, request: AlignRequest, + video_audio: Path) -> None: + with SessionLocal() as s: + job = s.get(Job, job_id) + if job is None: + raise JobFailed("job vanished while collecting output") + language, model = job.language, job.model + audio_name, text_name = job.audio_filename, job.text_filename + splitter = job.splitter + parts_count = job.audio_parts or 1 + + produced = aligner.locate_output(request) + if produced is None or not produced.exists(): + raise JobFailed( + f"{aligner.name} finished but produced no subtitle file. The audio " + "and text may not match, or the language may be wrong for this book." + ) + + # Give the download a filename the user will recognise. A per-chapter + # audiobook has no single audio name worth using, so the book names it. + stem = Path(text_name if parts_count > 1 else audio_name).stem + download_name = f"{stem}.{language}{aligner.output_suffix}" + + cues, first_cue, last_cue = _summarize_srt(produced) + + srt_key = f"{job_id}/{download_name}" + size = storage.put_file(srt_key, produced) + + video_name = video_key = None + video_size = 0 + if settings.render_video: + try: + _set(job_id, stage="Rendering video", progress=0.94) + scratch = paths.root / "video" + cover = render.extract_cover(request.text, scratch) + canvas = render.build_canvas(cover, scratch / "canvas.png") + video_name = f"{stem}.{language}.mp4" + out_video = scratch / video_name + render.render_video( + audio=video_audio, subtitles=produced, canvas=canvas, + dest=out_video, duration=duration, + ) + video_key = f"{job_id}/{video_name}" + video_size = storage.put_file(video_key, out_video) + except Exception as exc: # noqa: BLE001 - never fail a job over the video + # The subtitles are the product, so a failed render is not fatal - + # but it must not be silent either, or it looks like it never ran. + video_name = video_key = None + log.warning("job %s: video render failed: %s", job_id, exc, + exc_info=True) + + metadata = { + "job_id": job_id, + "created_at": utcnow().astimezone(timezone.utc).isoformat(), + "source": { + "audio_filename": audio_name, + "audio_parts": parts_count, + "text_filename": text_name, + "audio_duration_seconds": duration, + }, + "alignment": { + "backend": aligner.name, + "mode": "forced alignment against supplied text (nothing transcribed from scratch)", + "model": model, + "device": settings.device, + "language": language, + "language_name": (languages.get(language).name if languages.get(language) else language), + "sentence_splitter": splitter, + }, + "output": { + "filename": download_name, + "format": "srt", + "cue_count": cues, + "first_cue_start": first_cue, + "last_cue_end": last_cue, + "size_bytes": size, + }, + "video": ( + { + "filename": video_name, + "container": "mp4", + "subtitles": "soft track (mov_text)", + "size_bytes": video_size, + } + if video_name + else None + ), + } + meta_path = paths.root / "metadata.json" + meta_path.write_text(json.dumps(metadata, ensure_ascii=False, indent=2), + encoding="utf-8") + meta_key = f"{job_id}/metadata.json" + meta_size = storage.put_file(meta_key, meta_path) + + rows = [ + Artifact(job_id=job_id, kind="srt", filename=download_name, + storage_key=srt_key, size_bytes=size), + Artifact(job_id=job_id, kind="metadata", filename="metadata.json", + storage_key=meta_key, size_bytes=meta_size), + ] + if video_name and video_key: + rows.append( + Artifact(job_id=job_id, kind="video", filename=video_name, + storage_key=video_key, size_bytes=video_size) + ) + with SessionLocal() as s: + for r in rows: + s.add(r) + s.commit() + + _store_log(job_id, log_path) + + +def _store_log(job_id: str, log_path: Path) -> None: + if not log_path.exists(): + return + with SessionLocal() as s: + existing = ( + s.query(Artifact) + .filter(Artifact.job_id == job_id, Artifact.kind == "log") + .first() + ) + if existing is not None: + return + key = f"{job_id}/subplz.log" + size = storage.put_file(key, log_path) + with SessionLocal() as s: + s.add(Artifact(job_id=job_id, kind="log", filename="subplz.log", + storage_key=key, size_bytes=size)) + s.commit() + + +_TIME = re.compile(r"(\d{2}):(\d{2}):(\d{2}),(\d{3})") + + +def _summarize_srt(path: Path) -> tuple[int, float | None, float | None]: + """Cue count and span - cheap sanity signal that the alignment covered the book.""" + def to_seconds(m: re.Match) -> float: + h, mi, s, ms = (int(g) for g in m.groups()) + return h * 3600 + mi * 60 + s + ms / 1000 + + text = path.read_text(encoding="utf-8", errors="replace") + stamps = list(_TIME.finditer(text)) + count = text.count(" --> ") + if not stamps: + return count, None, None + return count, to_seconds(stamps[0]), to_seconds(stamps[-1]) + + +def _cleanup_inputs(paths: Paths) -> None: + """Drop the uploaded media once the run is over; keep artifacts and the log.""" + shutil.rmtree(paths.inp, ignore_errors=True) + shutil.rmtree(paths.out, ignore_errors=True) + shutil.rmtree(paths.root / "video", ignore_errors=True) + # Also drop the shared-storage copy made for external workers. + try: + storage.delete_prefix(f"{paths.root.name}/input") + except Exception: + pass diff --git a/backend/settings.py b/backend/settings.py new file mode 100644 index 0000000..ed9a9f8 --- /dev/null +++ b/backend/settings.py @@ -0,0 +1,97 @@ +"""Configuration. Every scale-out seam is an env var with a localhost-friendly default.""" + +import shutil +from pathlib import Path +from typing import Literal + +from pydantic_settings import BaseSettings, SettingsConfigDict + +ROOT = Path(__file__).resolve().parent.parent + + +def _find_subplz() -> str: + """Locate the alignment backend on PATH. + + subplz is an ordinary dependency of this project, installed into whatever + environment is running it - not a sibling checkout. Falling back to the bare + name keeps the error useful when it is genuinely missing. + """ + return shutil.which("subplz") or "subplz" + + +class Settings(BaseSettings): + model_config = SettingsConfigDict(env_prefix="SUBPLZ_WEB_", env_file=".env", extra="ignore") + + # --- paths ------------------------------------------------------------- + data_dir: Path = ROOT / "data" + # Override with SUBPLZ_WEB_SUBPLZ_BIN to point at a specific build. + subplz_bin: Path | str = _find_subplz() + + # --- alignment --------------------------------------------------------- + # Which backend does the aligning. See aligner.py to add another. + aligner: str = "subplz" + # "tiny" is what upstream recommends for audiobooks: the transcript only has + # to be good enough to align against text we already have. + model: str = "tiny" + device: Literal["cpu", "cuda"] = "cpu" + # 0 = let the runner pick from CPU count. + threads: int = 0 + # Wall-clock ceiling per job so a wedged run cannot hold a worker forever. + job_timeout_seconds: int = 6 * 60 * 60 + + # --- scale-out seams --------------------------------------------------- + # "memory" runs jobs in a thread in this process (localhost). + # "redis" hands them to external workers (public deployment). + queue_backend: Literal["memory", "redis"] = "memory" + redis_url: str = "redis://localhost:6379/0" + # How many jobs this process will run at once when queue_backend="memory". + max_concurrent_jobs: int = 1 + + # "local" writes under data_dir. "s3" writes to a bucket. + storage_backend: Literal["local", "s3"] = "local" + s3_bucket: str = "" + s3_prefix: str = "jobs/" + # Presigned-URL lifetime for s3 downloads. + download_url_ttl_seconds: int = 3600 + + # SQLite locally; set to a postgresql+psycopg:// URL in production. + database_url: str = "" + + # --- billing ----------------------------------------------------------- + # The public offer: one free book per rolling 24 hours, pay for more. + free_conversions: int = 1 + free_window_hours: int = 24 + # Off on localhost so nothing blocks you; flip on for the public release. + billing_enabled: bool = False + + # --- match check ------------------------------------------------------- + # Transcribe a few short samples on upload and score them against the book, + # so a mismatched pair fails in seconds instead of after a full run. + match_check: bool = True + # Clean the book text before alignment to raise the match score. + text_prep: bool = True + + # --- video ------------------------------------------------------------- + # Render a YouTube-ready MP4 (cover image + audio + soft subtitle track). + render_video: bool = True + video_width: int = 1920 + video_height: int = 1080 + # 1 fps is the minimum YouTube accepts and all a still image needs. + video_fps: int = 1 + # "auto" picks the best H.264 encoder this ffmpeg build actually has. + # Force one with libx264 / h264_nvenc / h264_amf / h264_qsv / libopenh264. + video_encoder: str = "auto" + video_crf: int = 28 + + # --- uploads ----------------------------------------------------------- + max_upload_bytes: int = 2 * 1024 * 1024 * 1024 # 2 GiB + + @property + def resolved_database_url(self) -> str: + if self.database_url: + return self.database_url + return f"sqlite:///{(self.data_dir / 'subplz.db').as_posix()}" + + +settings = Settings() +settings.data_dir.mkdir(parents=True, exist_ok=True) diff --git a/backend/storage.py b/backend/storage.py new file mode 100644 index 0000000..e113adb --- /dev/null +++ b/backend/storage.py @@ -0,0 +1,170 @@ +"""Artifact storage. Local filesystem now, S3 for the public release. + +Everything downstream deals in opaque storage keys, so swapping the backend +touches neither the API nor the job runner. +""" + +import shutil +from abc import ABC, abstractmethod +from pathlib import Path + +from .settings import settings + + +class Storage(ABC): + @abstractmethod + def put_file(self, key: str, source: Path) -> int: + """Store the file at `source` under `key`; return the byte size.""" + + @abstractmethod + def open_stream(self, key: str): + """Return a binary file-like object for `key`.""" + + @abstractmethod + def read_text(self, key: str) -> str: + ... + + @abstractmethod + def delete_prefix(self, prefix: str) -> None: + ... + + @abstractmethod + def list_prefix(self, prefix: str) -> list[str]: + """Every key under `prefix`. Used to move a job's inputs between machines.""" + + def presigned_url(self, key: str, filename: str) -> str | None: + """S3 hands back a direct URL; local storage streams through the app.""" + return None + + def exists(self, key: str) -> bool: + try: + self.open_stream(key).close() + return True + except Exception: + return False + + def fetch_to(self, key: str, dest: Path) -> Path: + """Copy `key` out of storage onto local disk. + + This is what lets a worker on another machine pick up a job whose files + were uploaded to a different API process. + """ + dest.parent.mkdir(parents=True, exist_ok=True) + src = self.open_stream(key) + try: + with dest.open("wb") as out: + while chunk := src.read(8 * 1024 * 1024): + out.write(chunk) + finally: + close = getattr(src, "close", None) + if close: + close() + return dest + + +class LocalStorage(Storage): + def __init__(self, root: Path): + self.root = root.resolve() + self.root.mkdir(parents=True, exist_ok=True) + + def _path(self, key: str) -> Path: + # Keys are app-generated, but refuse traversal regardless. + p = (self.root / key).resolve() + if not str(p).startswith(str(self.root)): + raise ValueError(f"unsafe storage key: {key!r}") + return p + + def path_for(self, key: str) -> Path: + """On-disk location of `key`. Local backend only - used to serve files.""" + return self._path(key) + + def put_file(self, key: str, source: Path) -> int: + dest = self._path(key) + dest.parent.mkdir(parents=True, exist_ok=True) + shutil.copy2(source, dest) + return dest.stat().st_size + + def open_stream(self, key: str): + return self._path(key).open("rb") + + def read_text(self, key: str) -> str: + return self._path(key).read_text(encoding="utf-8") + + def delete_prefix(self, prefix: str) -> None: + target = self._path(prefix) + if target.is_dir(): + shutil.rmtree(target, ignore_errors=True) + elif target.exists(): + target.unlink() + + def list_prefix(self, prefix: str) -> list[str]: + target = self._path(prefix) + if not target.is_dir(): + return [prefix] if target.exists() else [] + return sorted( + p.relative_to(self.root).as_posix() + for p in target.rglob("*") + if p.is_file() + ) + + +class S3Storage(Storage): + """Public-release backend. Needs boto3 and SUBPLZ_WEB_S3_BUCKET.""" + + def __init__(self, bucket: str, prefix: str): + import boto3 # lazy import so localhost never needs boto3 + + self.client = boto3.client("s3") + self.bucket = bucket + self.prefix = prefix.rstrip("/") + "/" if prefix else "" + + def _key(self, key: str) -> str: + return f"{self.prefix}{key}" + + def put_file(self, key: str, source: Path) -> int: + self.client.upload_file(str(source), self.bucket, self._key(key)) + return source.stat().st_size + + def open_stream(self, key: str): + return self.client.get_object(Bucket=self.bucket, Key=self._key(key))["Body"] + + def read_text(self, key: str) -> str: + return self.open_stream(key).read().decode("utf-8") + + def delete_prefix(self, prefix: str) -> None: + paginator = self.client.get_paginator("list_objects_v2") + for page in paginator.paginate(Bucket=self.bucket, Prefix=self._key(prefix)): + keys = [{"Key": o["Key"]} for o in page.get("Contents", [])] + if keys: + self.client.delete_objects(Bucket=self.bucket, Delete={"Objects": keys}) + + def list_prefix(self, prefix: str) -> list[str]: + paginator = self.client.get_paginator("list_objects_v2") + out: list[str] = [] + head = len(self.prefix) + for page in paginator.paginate(Bucket=self.bucket, Prefix=self._key(prefix)): + out += [o["Key"][head:] for o in page.get("Contents", [])] + return sorted(out) + + def presigned_url(self, key: str, filename: str) -> str | None: + disposition = f'attachment; filename="{filename}"' + return self.client.generate_presigned_url( + "get_object", + Params={ + "Bucket": self.bucket, + "Key": self._key(key), + "ResponseContentDisposition": disposition, + }, + ExpiresIn=settings.download_url_ttl_seconds, + ) + + +def get_storage() -> Storage: + if settings.storage_backend == "s3": + if not settings.s3_bucket: + raise RuntimeError("storage_backend=s3 requires SUBPLZ_WEB_S3_BUCKET") + return S3Storage(settings.s3_bucket, settings.s3_prefix) + return LocalStorage(settings.data_dir / "artifacts") + + +storage = get_storage() diff --git a/frontend/app.js b/frontend/app.js new file mode 100644 index 0000000..df8e96d --- /dev/null +++ b/frontend/app.js @@ -0,0 +1,544 @@ +/* SubPlz web UI. No framework, no build step - one file, served as-is. */ + +const $ = (id) => document.getElementById(id); + +const el = { + dropzone: $('dropzone'), picker: $('filepicker'), browse: $('browse'), + uploading: $('uploading'), upbar: $('upbar'), uptext: $('uptext'), + error: $('error'), + confirm: $('confirm'), cAudio: $('c-audio'), cText: $('c-text'), + language: $('language'), detected: $('detected'), splitnote: $('splitnote'), + start: $('start'), discard: $('discard'), eta: $('eta'), + jobsSection: $('jobs-section'), jobs: $('jobs'), quota: $('quota'), + freetier: $('freetier'), staged: $('staged'), dzTitle: $('dz-title'), + match: $('match'), matchBadge: $('match-badge'), + matchSummary: $('match-summary'), matchWarnings: $('match-warnings'), +}; + +let languages = []; +let languagesReady = null; // resolves once the <select> is populated +let draft = null; // the job awaiting confirmation +let pollTimer = null; + +/* ---------------- helpers ---------------- */ + +const fmtBytes = (n) => { + if (!n) return ''; + const u = ['B', 'KB', 'MB', 'GB']; + let i = 0, v = n; + while (v >= 1024 && i < u.length - 1) { v /= 1024; i++; } + return `${v < 10 && i > 0 ? v.toFixed(1) : Math.round(v)} ${u[i]}`; +}; + +const fmtDuration = (s) => { + if (!s && s !== 0) return ''; + const h = Math.floor(s / 3600), m = Math.round((s % 3600) / 60); + return h ? `${h}h ${String(m).padStart(2, '0')}m` : `${m}m`; +}; + +function showError(msg) { + el.error.textContent = msg; + el.error.hidden = false; + el.error.scrollIntoView({ behavior: 'smooth', block: 'nearest' }); +} +const clearError = () => { el.error.hidden = true; }; + +async function api(path, opts = {}) { + const res = await fetch(path, { credentials: 'same-origin', ...opts }); + if (!res.ok) { + let detail = `${res.status} ${res.statusText}`; + try { + const body = await res.json(); + if (body.detail) detail = typeof body.detail === 'string' + ? body.detail : JSON.stringify(body.detail); + } catch { /* non-JSON error body */ } + throw new Error(detail); + } + return res.status === 204 ? null : res.json(); +} + +/* ---------------- setup ---------------- */ + +async function loadLanguages() { + languages = await api('/api/languages'); + el.language.innerHTML = ''; + for (const l of languages) { + const o = document.createElement('option'); + o.value = l.code; + o.textContent = `${l.name} (${l.code})`; + el.language.appendChild(o); + } +} + +/* Setting .value on a <select> silently does nothing when the option is not + there yet, which would leave whatever was selected before on screen while the + badge says something else. Make the mismatch impossible instead. */ +function setLanguage(code) { + el.language.value = code; + if (el.language.value !== code) { + const o = document.createElement('option'); + o.value = code; + o.textContent = code; + el.language.appendChild(o); + el.language.value = code; + } + updateSplitNote(); +} + +async function refreshQuota() { + try { + const a = await api('/api/account'); + if (a.free_tier_summary && el.freetier) { + el.freetier.textContent = a.free_tier_summary.replace(/^./, (c) => c.toUpperCase()) + '.'; + } + if (!a.billing_enabled) { + el.quota.hidden = true; + return; + } + el.quota.hidden = false; + el.quota.classList.toggle('blocked', !a.allowed); + el.quota.textContent = a.allowed + ? `${a.remaining} conversion${a.remaining === 1 ? '' : 's'} left` + : 'Free conversion used — add credit to continue'; + } catch { /* quota is cosmetic; never block the UI on it */ } +} + +/* ---------------- upload ---------------- */ + +/* Files can arrive one drop at a time, so collect them until both halves are + here, then upload. Ignores the cover art and readme files that ride along + inside audiobook folders. */ +const AUDIO_RE = /\.(m4b|m4a|mp3|opus|ogg|oga|flac|wav|aac|wma|mka|mkv|mp4|webm|avi|mov)$/i; +const TEXT_RE = /(\.fb2\.zip|\.(epub|txt|srt|vtt|ass|fb2|mobi|azw|azw3|prc))$/i; + +const staged = { audio: [], text: null }; + +function addFiles(files) { + clearError(); + const list = [...files]; + const audio = list.filter((f) => AUDIO_RE.test(f.name)); + const text = list.filter((f) => TEXT_RE.test(f.name)); + + if (!audio.length && !text.length) { + showError('Nothing usable in that drop — expected an audiobook or an ebook.'); + return; + } + + // Audio accumulates, so a per-chapter set can arrive in batches. Only one + // book makes sense, so a second one replaces the first. + for (const f of audio) { + if (!staged.audio.some((g) => g.name === f.name && g.size === f.size)) { + staged.audio.push(f); + } + } + if (text.length) staged.text = text[text.length - 1]; + + renderStaged(); + if (staged.audio.length && staged.text) { + uploadFiles([...staged.audio, staged.text]); + } +} + +function renderStaged() { + const have = staged.audio.length || staged.text; + el.staged.hidden = !have; + + const a = $('slot-audio'); + a.classList.toggle('filled', staged.audio.length > 0); + a.querySelector('.slot-value').textContent = staged.audio.length + ? (staged.audio.length > 1 + ? `${staged.audio[0].name} + ${staged.audio.length - 1} more` + : staged.audio[0].name) + : 'Waiting for an audio file'; + + const t = $('slot-text'); + t.classList.toggle('filled', !!staged.text); + t.querySelector('.slot-value').textContent = + staged.text ? staged.text.name : 'Waiting for an ebook'; + + // Name the missing half, but only once one half is actually here. + const missing = !have ? null + : !staged.audio.length ? 'audiobook' + : !staged.text ? 'ebook' + : null; + el.dzTitle.innerHTML = missing + ? `Now drop the <em>${missing}</em>` + : 'Drop your audiobook <em>and</em> ebook here'; +} + +function clearStaged() { + staged.audio = []; + staged.text = null; + renderStaged(); +} + +document.querySelectorAll('.slot-x').forEach((btn) => { + btn.addEventListener('click', () => { + if (btn.dataset.slot === 'audio') staged.audio = []; + else staged.text = null; + renderStaged(); + }); +}); + +function uploadFiles(files) { + clearError(); + const list = [...files]; + + el.dropzone.hidden = true; + el.staged.hidden = true; + el.confirm.hidden = true; + el.uploading.hidden = false; + el.upbar.style.width = '0%'; + el.uptext.textContent = 'Starting…'; + + const form = new FormData(); + for (const f of list) form.append('files', f, f.name); + + // XHR rather than fetch: it reports upload progress, and these files are big. + const xhr = new XMLHttpRequest(); + xhr.open('POST', '/api/uploads'); + xhr.withCredentials = true; + + xhr.upload.onprogress = (e) => { + if (!e.lengthComputable) return; + const pct = Math.round((e.loaded / e.total) * 100); + el.upbar.style.width = `${pct}%`; + el.uptext.textContent = pct < 100 + ? `${pct}% — ${fmtBytes(e.loaded)} of ${fmtBytes(e.total)}` + : 'Analysing the book…'; + }; + + xhr.onload = () => { + el.uploading.hidden = true; + if (xhr.status >= 200 && xhr.status < 300) { + try { showConfirm(JSON.parse(xhr.responseText)); } + catch { showError('Server sent an unreadable response.'); resetToDrop(); } + } else { + let msg = `Upload failed (${xhr.status}).`; + try { + const b = JSON.parse(xhr.responseText); + if (b.detail) msg = typeof b.detail === 'string' ? b.detail : msg; + } catch { /* keep the generic message */ } + showError(msg); + resetToDrop(); + } + }; + + xhr.onerror = () => { + el.uploading.hidden = true; + showError('Upload failed: lost connection to the server.'); + resetToDrop(); + }; + + xhr.send(form); +} + +function resetToDrop() { + draft = null; + clearStaged(); + el.dropzone.hidden = false; + el.confirm.hidden = true; + el.uploading.hidden = true; + el.picker.value = ''; +} + +/* ---------------- confirm ---------------- */ + +async function showConfirm(payload) { + // The options must exist before we can select the detected language. + try { await languagesReady; } catch { /* handled at boot */ } + + draft = payload.job; + const d = payload.detected; + + const dur = fmtDuration(draft.audio_duration_seconds); + const bits = [ + draft.audio_parts > 1 ? `${draft.audio_parts} parts` : null, + fmtBytes(draft.audio_bytes), + dur, + ].filter(Boolean); + el.cAudio.innerHTML = + `${escapeHtml(draft.audio_filename)} <span class="meta">${bits.join(' · ')}</span>`; + el.cText.textContent = draft.text_filename; + + setLanguage(draft.language); + + if (d.code && d.supported) { + const pct = Math.round(d.confidence * 100); + el.detected.textContent = `detected ${d.name} · ${pct}%`; + el.detected.classList.toggle('low', d.confidence < 0.7); + el.detected.hidden = false; + } else if (d.code) { + el.detected.textContent = `detected ${d.name} — not supported, pick one`; + el.detected.classList.add('low'); + el.detected.hidden = false; + } else { + el.detected.hidden = true; + } + + renderMatch(payload.match); + + el.eta.textContent = draft.audio_duration_seconds + ? `Alignment usually takes a fraction of the book's length, but on CPU it can approach it. ${dur} of audio — expect a long run.` + : ''; + + el.confirm.hidden = false; + el.start.disabled = false; + el.confirm.scrollIntoView({ behavior: 'smooth', block: 'nearest' }); +} + +/* The preflight score, using the backend's own matching rule. A poor score is + shown but never blocks: the check samples the audio, so it can be wrong, and + it is the user's book. */ +function renderMatch(m) { + if (!m || m.skipped || m.verdict === 'unknown') { + el.match.hidden = true; + el.start.textContent = 'Start alignment'; + return; + } + + el.match.hidden = false; + el.match.className = `match ${m.verdict}`; + el.matchBadge.textContent = + { good: 'match', marginal: 'weak match', poor: 'no match' }[m.verdict] || m.verdict; + el.matchSummary.textContent = m.summary; + + el.matchWarnings.innerHTML = ''; + for (const w of m.warnings || []) { + const li = document.createElement('li'); + li.textContent = w; + el.matchWarnings.appendChild(li); + } + + // Make the user's choice explicit when we expect this to fail. + el.start.textContent = + m.verdict === 'poor' ? 'Start anyway' : 'Start alignment'; +} + +function updateSplitNote() { + // The note comes from the server, so the UI carries no knowledge of which + // alignment backend is running or how it splits sentences. + const lang = languages.find((l) => l.code === el.language.value); + el.splitnote.textContent = (lang && lang.note) || ''; +} + +el.language.addEventListener('change', updateSplitNote); + +el.discard.addEventListener('click', async () => { + if (draft) { try { await api(`/api/jobs/${draft.id}`, { method: 'DELETE' }); } catch {} } + clearError(); + resetToDrop(); + refreshJobs(); +}); + +el.start.addEventListener('click', async () => { + if (!draft) return; + el.start.disabled = true; + clearError(); + try { + await api(`/api/jobs/${draft.id}/start`, { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ language: el.language.value }), + }); + resetToDrop(); + refreshJobs(); + refreshQuota(); + } catch (e) { + showError(e.message); + el.start.disabled = false; + } +}); + +/* ---------------- jobs ---------------- */ + +function escapeHtml(s) { + return String(s).replace(/[&<>"']/g, (c) => ( + { '&': '&', '<': '<', '>': '>', '"': '"', "'": ''' }[c] + )); +} + +function jobCard(j) { + const li = document.createElement('li'); + li.className = 'job'; + + const active = j.status === 'running' || j.status === 'queued'; + const pct = Math.round(j.progress * 100); + + const meta = [ + j.language_name, + j.audio_parts > 1 ? `${j.audio_parts} parts` : null, + fmtDuration(j.audio_duration_seconds), + fmtBytes(j.audio_bytes), + ].filter(Boolean).join(' · '); + + let html = ` + <div class="job-head"> + <div style="min-width:0"> + <div class="job-name">${escapeHtml(j.audio_filename)}</div> + <div class="job-sub">${escapeHtml(meta)}</div> + </div> + <span class="pill ${j.status}">${j.status}</span> + </div>`; + + if (active) { + html += ` + <div class="bar"><div class="bar-fill" style="width:${pct}%"></div></div> + <div class="job-stage"><span>${escapeHtml(j.stage)}</span><span>${pct}%</span></div>`; + } + + if (j.error) html += `<div class="job-err">${escapeHtml(j.error)}</div>`; + + const files = j.artifacts || []; + if (files.length) { + const order = { srt: 0, video: 1, metadata: 2, log: 3 }; + const label = { + srt: '⬇ Subtitles (.srt)', + video: '⬇ Video (.mp4)', + metadata: 'Metadata', + log: 'Run log', + }; + const primary = new Set(['srt', 'video']); + html += '<div class="job-files">' + files + .slice() + .sort((a, b) => (order[a.kind] ?? 9) - (order[b.kind] ?? 9)) + .map((a) => { + const size = a.size_bytes ? ` <span class="dl-size">${fmtBytes(a.size_bytes)}</span>` : ''; + return `<a class="dl ${primary.has(a.kind) ? '' : 'secondary'}" + href="${a.url}" download>${label[a.kind] || a.kind}${size}</a>`; + }) + .join('') + '</div>'; + } + + li.innerHTML = html; + + if (active) { + const cancel = document.createElement('button'); + cancel.className = 'ghost'; + cancel.style.marginTop = '13px'; + cancel.textContent = 'Cancel'; + cancel.onclick = async () => { + cancel.disabled = true; + try { await api(`/api/jobs/${j.id}/cancel`, { method: 'POST' }); } + catch (e) { showError(e.message); } + refreshJobs(); refreshQuota(); + }; + li.appendChild(cancel); + } + return li; +} + +async function refreshJobs() { + let jobs; + try { jobs = await api('/api/jobs'); } + catch { return; } + + // Drafts live in the confirm panel, not the list. + const visible = jobs.filter((j) => j.status !== 'draft'); + el.jobsSection.hidden = visible.length === 0; + el.jobs.innerHTML = ''; + for (const j of visible) el.jobs.appendChild(jobCard(j)); + + const anyActive = visible.some((j) => j.status === 'running' || j.status === 'queued'); + clearTimeout(pollTimer); + if (anyActive) pollTimer = setTimeout(refreshJobs, 1500); +} + +/* ---------------- drag & drop ---------------- */ + +['dragenter', 'dragover'].forEach((ev) => + el.dropzone.addEventListener(ev, (e) => { + e.preventDefault(); + el.dropzone.classList.add('drag'); + })); + +['dragleave', 'drop'].forEach((ev) => + el.dropzone.addEventListener(ev, (e) => { + e.preventDefault(); + if (ev === 'dragleave' && el.dropzone.contains(e.relatedTarget)) return; + el.dropzone.classList.remove('drag'); + })); + +/* Dropping a folder gives directory entries, not files - audiobooks often + arrive as a folder of per-chapter mp3s, so walk it. */ +function readEntries(reader) { + return new Promise((resolve, reject) => reader.readEntries(resolve, reject)); +} + +async function walkEntry(entry, out, depth = 0) { + if (!entry || depth > 4) return; + if (entry.isFile) { + out.push(await new Promise((res, rej) => entry.file(res, rej))); + return; + } + if (entry.isDirectory) { + const reader = entry.createReader(); + // readEntries returns at most 100 per call, so keep going until it is empty. + for (;;) { + const batch = await readEntries(reader); + if (!batch.length) break; + for (const child of batch) await walkEntry(child, out, depth + 1); + } + } +} + +el.dropzone.addEventListener('drop', async (e) => { + const items = e.dataTransfer?.items; + const hasEntries = items && [...items].some( + (i) => i.kind === 'file' && typeof i.webkitGetAsEntry === 'function' + ); + + if (hasEntries) { + const entries = [...items] + .filter((i) => i.kind === 'file') + .map((i) => i.webkitGetAsEntry()) + .filter(Boolean); + if (entries.some((en) => en.isDirectory)) { + el.dropzone.hidden = true; + el.uploading.hidden = false; + el.uptext.textContent = 'Reading folder…'; + const files = []; + try { + for (const en of entries) await walkEntry(en, files); + } catch (err) { + el.uploading.hidden = true; + showError(`Could not read the dropped folder: ${err.message}`); + resetToDrop(); + return; + } + el.uploading.hidden = true; + el.dropzone.hidden = false; + if (files.length) addFiles(files); + return; + } + } + + const files = e.dataTransfer?.files; + if (files?.length) addFiles(files); +}); + +// Dropping anywhere else must not make the browser navigate to the file. +['dragover', 'drop'].forEach((ev) => + window.addEventListener(ev, (e) => e.preventDefault())); + +el.dropzone.addEventListener('click', (e) => { + if (e.target !== el.browse) el.picker.click(); +}); +el.browse.addEventListener('click', (e) => { e.stopPropagation(); el.picker.click(); }); +el.dropzone.addEventListener('keydown', (e) => { + if (e.key === 'Enter' || e.key === ' ') { e.preventDefault(); el.picker.click(); } +}); +el.picker.addEventListener('change', () => { + if (el.picker.files.length) addFiles(el.picker.files); + el.picker.value = ''; // let the same file be chosen again after a removal +}); + +/* ---------------- boot ---------------- */ + +(async function init() { + languagesReady = loadLanguages(); + try { await languagesReady; } + catch (e) { showError(`Could not reach the server: ${e.message}`); } + refreshQuota(); + refreshJobs(); +})(); diff --git a/frontend/index.html b/frontend/index.html new file mode 100644 index 0000000..6db962a --- /dev/null +++ b/frontend/index.html @@ -0,0 +1,135 @@ +<!DOCTYPE html> +<html lang="en"> +<head> +<meta charset="utf-8"> +<meta name="viewport" content="width=device-width, initial-scale=1"> +<title>SubPlz — Audiobook to Subtitles</title> +<link rel="stylesheet" href="/style.css?v=6"> +<link rel="icon" href="data:image/svg+xml,<svg xmlns='http://www.w3.org/2000/svg' viewBox='0 0 100 100'><text y='.9em' font-size='90'>🎧</text></svg>"> +</head> +<body> +<header class="topbar"> + <div class="brand"> + <span class="logo" aria-hidden="true">🎧</span> + <div> + <h1>SubPlz</h1> + <p class="tagline">Line an audiobook up with its ebook, sentence by sentence.</p> + </div> + </div> + <div id="quota" class="quota" hidden></div> +</header> + +<main> + <section class="pitch"> + <p> + Drop in an audiobook and the ebook it was read from. You get back subtitles + timed to the narration, with the wording taken from your own book — not + from a machine's guess at what it heard. + </p> + <ul class="uses"> + <li><strong>HoshiReader whispersync</strong> — a three-file package: audio, book, timings.</li> + <li><strong>Subtitled video</strong> — a YouTube-ready cut with the text on screen.</li> + </ul> + <p class="muted small" id="freetier">One free book every 24 hours.</p> + </section> + + <!-- step 1: drop --> + <section id="dropzone" class="dropzone" tabindex="0" role="button" + aria-label="Drop an audio file and an ebook, or click to choose files"> + <input type="file" id="filepicker" multiple hidden + accept=".m4b,.m4a,.mp3,.opus,.ogg,.oga,.flac,.wav,.aac,.wma,.mka,.mkv,.mp4,.webm,.avi,.mov,.epub,.txt,.srt,.vtt,.ass,.fb2,.zip,.mobi,.azw,.azw3,.prc"> + <div class="dz-inner"> + <div class="dz-icon" aria-hidden="true"> + <svg viewBox="0 0 24 24" fill="none" stroke="currentColor" stroke-width="1.5" + stroke-linecap="round" stroke-linejoin="round"> + <path d="M12 16V4m0 0L7.5 8.5M12 4l4.5 4.5"/> + <path d="M4 15v3a2 2 0 0 0 2 2h12a2 2 0 0 0 2-2v-3"/> + </svg> + </div> + <p class="dz-title" id="dz-title">Drop your audiobook <em>and</em> ebook here</p> + <p class="dz-sub" id="dz-sub">One at a time or both at once — order doesn't matter. + A folder of per-chapter files works too. + Or <button type="button" id="browse" class="linkbtn">browse</button>.</p> + <p class="dz-formats">Audio: m4b · mp3 · m4a · opus · flac · wav · mkv + | Book: epub · fb2 · mobi · azw3 · txt</p> + </div> + </section> + + <!-- what has been dropped so far --> + <section id="staged" class="panel staged" hidden> + <div class="slot" id="slot-audio"> + <span class="slot-tick" aria-hidden="true"></span> + <div class="slot-body"> + <div class="slot-label">Audiobook</div> + <div class="slot-value">Waiting for an audio file</div> + </div> + <button type="button" class="slot-x" data-slot="audio" title="Remove">×</button> + </div> + <div class="slot" id="slot-text"> + <span class="slot-tick" aria-hidden="true"></span> + <div class="slot-body"> + <div class="slot-label">Book</div> + <div class="slot-value">Waiting for an ebook</div> + </div> + <button type="button" class="slot-x" data-slot="text" title="Remove">×</button> + </div> + </section> + + <div id="uploading" class="panel uploading" hidden> + <div class="spinner" aria-hidden="true"></div> + <div class="up-body"> + <p class="up-title">Uploading & analysing…</p> + <div class="bar"><div id="upbar" class="bar-fill"></div></div> + <p id="uptext" class="muted small">0%</p> + </div> + </div> + + <div id="error" class="panel error" hidden role="alert"></div> + + <!-- step 2: confirm --> + <section id="confirm" class="panel confirm" hidden> + <h2>Ready to align</h2> + <dl class="pairing"> + <div><dt>Audiobook</dt><dd id="c-audio"></dd></div> + <div><dt>Book</dt><dd id="c-text"></dd></div> + </dl> + + <div class="langrow"> + <label for="language">Language</label> + <select id="language"></select> + <span id="detected" class="detected"></span> + </div> + <p id="splitnote" class="muted small"></p> + + <div id="match" class="match" hidden> + <div class="match-head"> + <span id="match-badge" class="match-badge"></span> + <span id="match-summary"></span> + </div> + <ul id="match-warnings" class="match-warnings"></ul> + </div> + + <div class="actions"> + <button id="start" class="primary">Start alignment</button> + <button id="discard" class="ghost">Discard</button> + </div> + <p class="muted small" id="eta"></p> + </section> + + <!-- step 3: jobs --> + <section id="jobs-section" hidden> + <h2 class="section-title">Your conversions</h2> + <ul id="jobs" class="jobs"></ul> + </section> +</main> + +<footer> + <p class="muted small"> + Alignment by <a href="https://github.com/kanjieater/SubPlz" target="_blank" rel="noopener">SubPlz</a>. + Every line of text comes from the book you upload, so names and spelling stay right. + </p> +</footer> + +<script src="/app.js?v=6"></script> +</body> +</html> diff --git a/frontend/style.css b/frontend/style.css new file mode 100644 index 0000000..f7ddb2c --- /dev/null +++ b/frontend/style.css @@ -0,0 +1,290 @@ +:root { + --bg: #f7f7f9; + --panel: #ffffff; + --ink: #16161a; + --muted: #6b6b76; + --line: #e3e3e9; + --accent: #4f46e5; + --accent-ink: #ffffff; + --accent-soft: #eef0ff; + --ok: #0f8a5f; + --ok-soft: #e7f6ef; + --err: #c0392b; + --err-soft: #fdecea; + --warn: #9a6700; + --warn-soft: #fff6dd; + --radius: 14px; + --shadow: 0 1px 2px rgba(16, 16, 32, .05), 0 8px 24px rgba(16, 16, 32, .06); +} + +@media (prefers-color-scheme: dark) { + :root { + --bg: #0f1013; + --panel: #17181d; + --ink: #ececf1; + --muted: #9a9aa6; + --line: #2a2b33; + --accent: #7c74ff; + --accent-ink: #0f1013; + --accent-soft: #1e1f33; + --ok: #46d19a; + --ok-soft: #12241d; + --err: #ff8b7a; + --err-soft: #2a1714; + --warn: #e6bb4f; + --warn-soft: #2a2312; + --shadow: 0 1px 2px rgba(0, 0, 0, .4), 0 8px 24px rgba(0, 0, 0, .35); + } +} + +* { box-sizing: border-box; } + +/* An author `display` rule beats the UA stylesheet's [hidden] rule, and several + panels below set display explicitly. Keep [hidden] authoritative. */ +[hidden] { display: none !important; } + +body { + margin: 0; + background: var(--bg); + color: var(--ink); + font: 15px/1.55 ui-sans-serif, system-ui, -apple-system, "Segoe UI", Roboto, + "Helvetica Neue", Arial, "Noto Sans", "Noto Sans JP", sans-serif; + -webkit-font-smoothing: antialiased; +} + +main, .topbar, footer { max-width: 880px; margin: 0 auto; padding: 0 20px; } + +.topbar { + display: flex; align-items: center; justify-content: space-between; + gap: 16px; flex-wrap: wrap; padding-top: 32px; padding-bottom: 20px; +} +.brand { display: flex; align-items: center; gap: 14px; } +.logo { font-size: 34px; line-height: 1; } +h1 { font-size: 22px; margin: 0; letter-spacing: -.01em; } +.tagline { margin: 2px 0 0; color: var(--muted); font-size: 13.5px; } + +.quota { + font-size: 13px; padding: 7px 13px; border-radius: 999px; + background: var(--accent-soft); color: var(--accent); + border: 1px solid transparent; white-space: nowrap; +} +.quota.blocked { background: var(--warn-soft); color: var(--warn); } + +/* ---------- pitch ---------- */ +.pitch { margin: 0 0 22px; } +.pitch p { margin: 0 0 12px; } +.uses { list-style: none; padding: 0; margin: 0 0 12px; display: grid; gap: 7px; } +.uses li { + font-size: 14px; color: var(--muted); + padding-left: 18px; position: relative; +} +.uses li::before { + content: ""; position: absolute; left: 3px; top: 9px; + width: 5px; height: 5px; border-radius: 50%; background: var(--accent); +} +.uses strong { color: var(--ink); font-weight: 600; } + +/* ---------- dropzone ---------- */ +.dropzone { + border: 2px dashed var(--line); + border-radius: var(--radius); + background: var(--panel); + padding: 46px 24px; + text-align: center; + cursor: pointer; + transition: border-color .15s, background .15s, transform .15s; +} +.dropzone:hover, .dropzone:focus-visible { border-color: var(--accent); outline: none; } +.dropzone.drag { + border-color: var(--accent); background: var(--accent-soft); + transform: scale(1.008); +} +.dz-icon { color: var(--accent); } +.dz-icon svg { width: 44px; height: 44px; } +.dz-title { font-size: 17px; font-weight: 600; margin: 14px 0 4px; } +.dz-title em { font-style: normal; text-decoration: underline; text-underline-offset: 3px; } +.dz-sub { margin: 0 0 10px; color: var(--muted); font-size: 14px; } +.dz-formats { margin: 0; color: var(--muted); font-size: 12px; } + +.linkbtn { + background: none; border: none; padding: 0; font: inherit; + color: var(--accent); cursor: pointer; text-decoration: underline; +} + +/* ---------- staged files ---------- */ +.staged { display: grid; gap: 10px; } +.slot { display: flex; align-items: center; gap: 12px; } +.slot-tick { + width: 20px; height: 20px; flex: none; border-radius: 50%; + border: 2px dashed var(--line); +} +.slot.filled .slot-tick { + border: none; background: var(--ok); position: relative; +} +.slot.filled .slot-tick::after { + content: ""; position: absolute; left: 7px; top: 3px; + width: 5px; height: 10px; border: solid var(--panel); + border-width: 0 2px 2px 0; transform: rotate(45deg); +} +.slot-body { flex: 1; min-width: 0; } +.slot-label { + font-size: 11.5px; color: var(--muted); + text-transform: uppercase; letter-spacing: .05em; +} +.slot-value { font-weight: 500; word-break: break-word; } +.slot:not(.filled) .slot-value { font-weight: 400; color: var(--muted); } +.slot-x { + background: none; border: none; color: var(--muted); cursor: pointer; + font-size: 20px; line-height: 1; padding: 0 4px; +} +.slot-x:hover { color: var(--err); } +.slot:not(.filled) .slot-x { visibility: hidden; } + +/* ---------- panels ---------- */ +.panel { + margin-top: 18px; background: var(--panel); border: 1px solid var(--line); + border-radius: var(--radius); padding: 20px; box-shadow: var(--shadow); +} +.panel h2 { margin: 0 0 14px; font-size: 16px; } +.error { border-color: var(--err); background: var(--err-soft); color: var(--err); } + +.uploading { display: flex; gap: 16px; align-items: center; } +.up-body { flex: 1; min-width: 0; } +.up-title { margin: 0 0 8px; font-weight: 600; } + +.spinner { + width: 26px; height: 26px; flex: none; border-radius: 50%; + border: 3px solid var(--line); border-top-color: var(--accent); + animation: spin .8s linear infinite; +} +@keyframes spin { to { transform: rotate(360deg); } } + +.bar { + height: 8px; background: var(--line); border-radius: 999px; overflow: hidden; +} +.bar-fill { + height: 100%; width: 0%; background: var(--accent); + border-radius: 999px; transition: width .35s ease; +} + +/* ---------- confirm ---------- */ +.pairing { margin: 0 0 18px; display: grid; gap: 8px; } +.pairing > div { display: flex; gap: 12px; align-items: baseline; } +.pairing dt { + color: var(--muted); font-size: 12.5px; width: 88px; flex: none; + text-transform: uppercase; letter-spacing: .04em; +} +.pairing dd { + margin: 0; font-weight: 500; word-break: break-word; min-width: 0; +} +.pairing .meta { color: var(--muted); font-weight: 400; font-size: 13px; } + +.langrow { display: flex; align-items: center; gap: 10px; flex-wrap: wrap; } +.langrow label { + font-size: 12.5px; color: var(--muted); + text-transform: uppercase; letter-spacing: .04em; width: 88px; flex: none; +} +select { + font: inherit; padding: 8px 11px; border-radius: 9px; + border: 1px solid var(--line); background: var(--panel); color: var(--ink); + min-width: 230px; +} +.detected { + font-size: 12.5px; padding: 4px 10px; border-radius: 999px; + background: var(--ok-soft); color: var(--ok); +} +.detected.low { background: var(--warn-soft); color: var(--warn); } +#splitnote { margin: 10px 0 0; padding-left: 98px; } + +/* ---------- match check ---------- */ +.match { + margin-top: 16px; padding: 12px 14px; border-radius: 10px; + background: var(--ok-soft); border: 1px solid transparent; +} +.match.marginal { background: var(--warn-soft); } +.match.poor { background: var(--err-soft); } +.match-head { display: flex; align-items: baseline; gap: 10px; flex-wrap: wrap; } +.match-badge { + font-size: 11.5px; font-weight: 700; padding: 3px 9px; border-radius: 999px; + text-transform: uppercase; letter-spacing: .05em; white-space: nowrap; + background: var(--ok); color: var(--panel); +} +.match.marginal .match-badge { background: var(--warn); } +.match.poor .match-badge { background: var(--err); } +.match.good { color: var(--ok); } +.match.marginal { color: var(--warn); } +.match.poor { color: var(--err); } +.match-warnings { + margin: 8px 0 0; padding-left: 18px; font-size: 13px; opacity: .9; +} +.match-warnings:empty { display: none; } + +.actions { display: flex; gap: 10px; margin-top: 20px; flex-wrap: wrap; } + +button.primary, button.ghost, .dl { + font: inherit; font-weight: 500; padding: 9px 17px; + border-radius: 9px; cursor: pointer; border: 1px solid transparent; +} +button.primary { background: var(--accent); color: var(--accent-ink); } +button.primary:hover { filter: brightness(1.08); } +button.primary:disabled { opacity: .55; cursor: not-allowed; filter: none; } +button.ghost { background: transparent; color: var(--muted); border-color: var(--line); } +button.ghost:hover { color: var(--ink); border-color: var(--muted); } + +/* ---------- jobs ---------- */ +.section-title { font-size: 16px; margin: 30px 0 12px; } +.jobs { list-style: none; padding: 0; margin: 0; display: grid; gap: 12px; } + +.job { + background: var(--panel); border: 1px solid var(--line); + border-radius: var(--radius); padding: 16px 18px; box-shadow: var(--shadow); +} +.job-head { + display: flex; align-items: center; gap: 10px; + justify-content: space-between; flex-wrap: wrap; +} +.job-name { font-weight: 600; word-break: break-word; min-width: 0; } +.job-sub { color: var(--muted); font-size: 13px; margin-top: 2px; } + +.pill { + font-size: 11.5px; font-weight: 600; padding: 3px 10px; border-radius: 999px; + text-transform: uppercase; letter-spacing: .05em; white-space: nowrap; +} +.pill.queued, .pill.draft { background: var(--accent-soft); color: var(--accent); } +.pill.running { background: var(--accent-soft); color: var(--accent); } +.pill.succeeded { background: var(--ok-soft); color: var(--ok); } +.pill.failed { background: var(--err-soft); color: var(--err); } +.pill.canceled { background: var(--line); color: var(--muted); } + +.job .bar { margin-top: 12px; } +.job-stage { + display: flex; justify-content: space-between; gap: 10px; + color: var(--muted); font-size: 12.5px; margin-top: 7px; +} +.job-err { + margin-top: 10px; padding: 10px 12px; border-radius: 9px; + background: var(--err-soft); color: var(--err); font-size: 13px; + white-space: pre-wrap; word-break: break-word; +} +.job-files { display: flex; gap: 8px; margin-top: 13px; flex-wrap: wrap; } +.dl { + text-decoration: none; background: var(--accent); color: var(--accent-ink); + display: inline-flex; align-items: center; gap: 7px; font-size: 13.5px; + padding: 8px 14px; +} +.dl:hover { filter: brightness(1.08); } +.dl-size { opacity: .7; font-size: 12px; font-weight: 400; } +.dl.secondary { background: transparent; color: var(--muted); border-color: var(--line); } +.dl.secondary:hover { color: var(--ink); border-color: var(--muted); } + +footer { padding-top: 34px; padding-bottom: 40px; } +footer a { color: var(--accent); } +.muted { color: var(--muted); } +.small { font-size: 12.5px; } + +@media (max-width: 560px) { + .pairing dt, .langrow label { width: 100%; } + .pairing > div { flex-direction: column; gap: 2px; } + #splitnote { padding-left: 0; } + select { min-width: 0; width: 100%; } +} diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..385f677 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,17 @@ +# Web layer only. subplz itself and its ML stack (faster-whisper, stable-ts, +# stanza, torch) come from the parent project: pip install .. in the same venv. +fastapi>=0.115 +uvicorn[standard]>=0.30 +python-multipart>=0.0.9 +sqlalchemy>=2.0 +pydantic-settings>=2.4 + +# fb2 is XML and epub is a zip, but mobi/azw3 need a real parser. +mobi>=0.4 + +# --- public deployment extras (not needed on localhost) --- +# redis>=5.0 # SUBPLZ_WEB_QUEUE_BACKEND=redis +# rq>=1.16 # external worker processes +# boto3>=1.34 # SUBPLZ_WEB_STORAGE_BACKEND=s3 +# psycopg[binary]>=3 # SUBPLZ_WEB_DATABASE_URL=postgresql+psycopg://... +# stripe>=9 # billing.start_checkout diff --git a/run.ps1 b/run.ps1 new file mode 100644 index 0000000..7abc030 --- /dev/null +++ b/run.ps1 @@ -0,0 +1,97 @@ +<# + Start SubPlz Web on localhost. + + .\run.ps1 # start on http://127.0.0.1:8420 + .\run.ps1 -Port 9000 + .\run.ps1 -NoBrowser + + Creates .venv and installs dependencies on first run, including the + alignment backend itself. +#> +[CmdletBinding()] +param( + [int]$Port = 8420, + [string]$BindHost = "127.0.0.1", + [switch]$NoBrowser, + [switch]$Reload +) + +$ErrorActionPreference = "Stop" + +$Root = $PSScriptRoot +$Venv = Join-Path $Root ".venv" +$Python = Join-Path $Venv "Scripts\python.exe" + +# Where to install the alignment backend from. Override to use a fork or a +# local checkout, e.g. -e C:\path\to\SubPlz +$SubPlzSpec = $env:SUBPLZ_INSTALL_SPEC +if (-not $SubPlzSpec) { + $SubPlzSpec = "git+https://github.com/kanjieater/SubPlz.git" +} + +function Write-Step($msg) { Write-Host "==> $msg" -ForegroundColor Cyan } +function Write-Warn($msg) { Write-Host "!! $msg" -ForegroundColor Yellow } + +# --- ffmpeg is not optional --------------------------------------------------- +if (-not (Get-Command ffmpeg -ErrorAction SilentlyContinue)) { + Write-Warn "ffmpeg is not on PATH. It is needed to read audio and build video." + Write-Warn "Install it (winget install Gyan.FFmpeg) and reopen the terminal." +} + +# --- venv --------------------------------------------------------------------- +# subplz pins requires-python >=3.10,<3.12, so 3.11 it is. +if (-not (Test-Path $Python)) { + Write-Step "Creating .venv on Python 3.11" + & py -3.11 --version 2>&1 | Out-Null + if ($LASTEXITCODE -ne 0) { + throw "Python 3.11 not found. Install it, then rerun. (py -0 lists what you have.)" + } + & py -3.11 -m venv $Venv + if ($LASTEXITCODE -ne 0) { throw "venv creation failed" } +} + +# --- dependencies ------------------------------------------------------------- +& $Python -c "import fastapi, uvicorn, sqlalchemy, multipart, pydantic_settings" 2>$null +if ($LASTEXITCODE -ne 0) { + Write-Step "Installing web dependencies" + & $Python -m pip install --upgrade pip + & $Python -m pip install -r (Join-Path $Root "requirements.txt") + if ($LASTEXITCODE -ne 0) { throw "dependency install failed" } +} + +& $Python -c "import subplz" 2>$null +if ($LASTEXITCODE -ne 0) { + Write-Step "Installing the alignment backend ($SubPlzSpec) - pulls torch, several minutes" + & $Python -m pip install $SubPlzSpec + if ($LASTEXITCODE -ne 0) { throw "alignment backend install failed" } +} + +# --- language registry -------------------------------------------------------- +if (-not (Test-Path (Join-Path $Root "backend\languages.json"))) { + Write-Step "Generating the language registry" + Push-Location $Root + try { & $Python -m backend.gen_languages } finally { Pop-Location } +} + +# --- run ---------------------------------------------------------------------- +# The backend prints emoji; without UTF-8 these die on a cp1252/cp932 console. +$env:PYTHONIOENCODING = "utf-8" +$env:PYTHONUTF8 = "1" + +# Put the venv's scripts first, so `subplz` resolves on PATH. +$env:PATH = (Join-Path $Venv "Scripts") + [IO.Path]::PathSeparator + $env:PATH + +$url = "http://${BindHost}:${Port}" +Write-Step "Starting SubPlz Web on $url" +Write-Host " Ctrl+C to stop." -ForegroundColor DarkGray + +if (-not $NoBrowser) { + Start-Job -ScriptBlock { param($u) Start-Sleep -Seconds 3; Start-Process $u } ` + -ArgumentList $url | Out-Null +} + +$uvicornArgs = @("-m", "uvicorn", "backend.main:app", "--host", $BindHost, "--port", $Port) +if ($Reload) { $uvicornArgs += "--reload" } + +Push-Location $Root +try { & $Python @uvicornArgs } finally { Pop-Location } diff --git a/tools/client.py b/tools/client.py new file mode 100644 index 0000000..35c7646 --- /dev/null +++ b/tools/client.py @@ -0,0 +1,259 @@ +"""Command-line client for the SubPlz web API. + +Useful on its own for batch work, and it doubles as a worked example of the +three-call flow the browser uses. + + python tools/client.py submit "book.m4b" "book.epub" [--language ru] [--wait] + python tools/client.py list + python tools/client.py get <job_id> + python tools/client.py download <job_id> [--dir out/] + +Uses only the standard library so it runs anywhere, and sends multipart bodies +itself so non-ASCII filenames survive (curl on a non-UTF-8 console mangles them). +""" + +from __future__ import annotations + +import argparse +import http.client +import json +import mimetypes +import os +import sys +import time +import urllib.parse +import uuid +from pathlib import Path + +DEFAULT_BASE = os.environ.get("SUBPLZ_WEB_URL", "http://127.0.0.1:8420") +COOKIE_FILE = Path.home() / ".subplz_web_cookie" + + +class Client: + def __init__(self, base: str): + parts = urllib.parse.urlsplit(base) + self.host = parts.hostname or "127.0.0.1" + self.port = parts.port or (443 if parts.scheme == "https" else 80) + self.https = parts.scheme == "https" + self.cookie = COOKIE_FILE.read_text().strip() if COOKIE_FILE.exists() else "" + + def _conn(self): + cls = http.client.HTTPSConnection if self.https else http.client.HTTPConnection + return cls(self.host, self.port, timeout=900) + + def _headers(self, extra: dict | None = None) -> dict: + h = {"Accept": "application/json"} + if self.cookie: + h["Cookie"] = self.cookie + h.update(extra or {}) + return h + + def _remember_cookie(self, resp) -> None: + raw = resp.getheader("set-cookie") + if raw: + self.cookie = raw.split(";", 1)[0] + try: + COOKIE_FILE.write_text(self.cookie) + except OSError: + pass + + def request(self, method: str, path: str, body=None, headers=None): + conn = self._conn() + conn.request(method, path, body, self._headers(headers)) + resp = conn.getresponse() + raw = resp.read().decode("utf-8", errors="replace") + self._remember_cookie(resp) + conn.close() + + if resp.status >= 400: + try: + detail = json.loads(raw).get("detail", raw) + except json.JSONDecodeError: + detail = raw[:500] + raise SystemExit(f"error {resp.status}: {detail}") + return json.loads(raw) if raw else None + + def upload(self, audio: list[Path], text: Path) -> dict: + boundary = "----subplz" + uuid.uuid4().hex + body = bytearray() + for path in [*audio, text]: + ctype = mimetypes.guess_type(path.name)[0] or "application/octet-stream" + disposition = ( + f'--{boundary}\r\n' + f'Content-Disposition: form-data; name="files"; ' + f'filename="{path.name}"\r\n' + f"Content-Type: {ctype}\r\n\r\n" + ) + # The header is encoded as UTF-8, which is what browsers send and + # what Starlette decodes - this is the bit curl gets wrong. + body += disposition.encode("utf-8") + body += path.read_bytes() + body += b"\r\n" + body += f"--{boundary}--\r\n".encode() + + return self.request( + "POST", "/api/uploads", bytes(body), + {"Content-Type": f"multipart/form-data; boundary={boundary}"}, + ) + + def start(self, job_id: str, language: str | None) -> dict: + payload = json.dumps({"language": language} if language else {}) + return self.request( + "POST", f"/api/jobs/{job_id}/start", + payload.encode(), {"Content-Type": "application/json"}, + ) + + def get(self, job_id: str) -> dict: + return self.request("GET", f"/api/jobs/{job_id}") + + def list(self) -> list: + return self.request("GET", "/api/jobs") + + def download(self, job_id: str, kind: str, out_dir: Path) -> Path | None: + job = self.get(job_id) + art = next((a for a in job["artifacts"] if a["kind"] == kind), None) + if art is None: + return None + conn = self._conn() + conn.request("GET", art["url"], None, self._headers()) + resp = conn.getresponse() + data = resp.read() + conn.close() + if resp.status >= 400: + raise SystemExit(f"download failed: {resp.status}") + out_dir.mkdir(parents=True, exist_ok=True) + dest = out_dir / art["filename"] + dest.write_bytes(data) + return dest + + +def _fmt(job: dict) -> str: + pct = int(job["progress"] * 100) + return ( + f"{job['id']} {job['status']:<10} {pct:>3}% " + f"{job['stage']:<34} {job['language_name']} {job['audio_filename']}" + ) + + +def wait(client: Client, job_id: str) -> dict: + last = None + while True: + job = client.get(job_id) + line = _fmt(job) + if line != last: + print(line, flush=True) + last = line + if job["status"] in ("succeeded", "failed", "canceled"): + return job + time.sleep(3) + + +def main() -> int: + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--base", default=DEFAULT_BASE) + sub = ap.add_subparsers(dest="cmd", required=True) + + s = sub.add_parser("submit", help="upload a pair and start aligning") + s.add_argument("audio") + s.add_argument("text") + s.add_argument("--language", help="override the detected language") + s.add_argument("--wait", action="store_true") + s.add_argument("--dir", default="subplz-out", help="where to save on --wait") + + sub.add_parser("list", help="list your jobs") + + g = sub.add_parser("get", help="show one job") + g.add_argument("job_id") + + w = sub.add_parser("watch", help="follow a job to completion") + w.add_argument("job_id") + + d = sub.add_parser("download", help="save a finished job's files") + d.add_argument("job_id") + d.add_argument("--dir", default="subplz-out") + + args = ap.parse_args() + client = Client(args.base) + + if args.cmd == "submit": + audio_arg, text = Path(args.audio), Path(args.text) + if not text.exists(): + raise SystemExit(f"no such file: {text}") + + if audio_arg.is_dir(): + # A folder of per-chapter files. Order is settled server-side. + exts = {".m4b", ".m4a", ".mp3", ".opus", ".ogg", ".oga", ".flac", + ".wav", ".aac", ".wma", ".mka", ".mkv", ".mp4", ".webm"} + audio = sorted( + p for p in audio_arg.iterdir() + if p.is_file() and p.suffix.lower() in exts + ) + if not audio: + raise SystemExit(f"no audio files in {audio_arg}") + else: + if not audio_arg.exists(): + raise SystemExit(f"no such file: {audio_arg}") + audio = [audio_arg] + + size_mb = (sum(p.stat().st_size for p in audio) + text.stat().st_size) / 1024**2 + label = f"{len(audio)} audio files" if len(audio) > 1 else audio[0].name + print(f"uploading {label}, {size_mb:.0f} MB ...", flush=True) + res = client.upload(audio, text) + + job, det = res["job"], res["detected"] + print(f"paired: audio={job['audio_filename']} text={job['text_filename']}") + if det["code"]: + print( + f"detected language: {det['name']} ({det['code']}) " + f"{det['confidence'] * 100:.1f}% confident" + + ("" if det["supported"] else " [NOT SUPPORTED - pass --language]") + ) + lang = args.language or job["language"] + print(f"starting with language={lang}") + started = client.start(job["id"], args.language) + print(_fmt(started)) + + if args.wait: + final = wait(client, job["id"]) + if final["status"] == "succeeded": + for kind in ("srt", "metadata"): + dest = client.download(job["id"], kind, Path(args.dir)) + if dest: + print(f"saved {dest}") + return 0 + print(f"job {final['status']}: {final.get('error') or ''}") + return 1 + return 0 + + if args.cmd == "list": + jobs = client.list() + if not jobs: + print("no jobs") + for j in jobs: + print(_fmt(j)) + return 0 + + if args.cmd == "get": + print(json.dumps(client.get(args.job_id), indent=2, ensure_ascii=False)) + return 0 + + if args.cmd == "watch": + final = wait(client, args.job_id) + return 0 if final["status"] == "succeeded" else 1 + + if args.cmd == "download": + any_saved = False + for kind in ("srt", "metadata", "log"): + dest = client.download(args.job_id, kind, Path(args.dir)) + if dest: + print(f"saved {dest}") + any_saved = True + if not any_saved: + print("nothing to download yet") + return 0 + + return 1 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/worker.py b/worker.py new file mode 100644 index 0000000..6a1e0fd --- /dev/null +++ b/worker.py @@ -0,0 +1,65 @@ +"""Standalone job worker for the public deployment. + +On localhost the API process runs jobs itself and you never need this. When you +set SUBPLZ_WEB_QUEUE_BACKEND=redis, the API only enqueues, and these processes +do the work: + + # one per GPU box, or several per box if you have the RAM + SUBPLZ_WEB_QUEUE_BACKEND=redis \ + SUBPLZ_WEB_REDIS_URL=redis://redis:6379/0 \ + SUBPLZ_WEB_DATABASE_URL=postgresql+psycopg://... \ + SUBPLZ_WEB_STORAGE_BACKEND=s3 SUBPLZ_WEB_S3_BUCKET=... \ + SUBPLZ_WEB_DEVICE=cuda SUBPLZ_WEB_MODEL=tiny \ + python worker.py + +Workers and the API must share the database and the storage bucket. They do not +need to share a filesystem: staged uploads are the one thing that is local to +whoever received them, which is why the API writes uploads to shared storage +before enqueuing when the queue backend is redis. +""" + +from __future__ import annotations + +import logging +import sys + +from backend.settings import settings + +logging.basicConfig( + level=logging.INFO, + format="%(asctime)s %(levelname)-7s %(name)s: %(message)s", +) +log = logging.getLogger("subplz.worker") + + +def main() -> int: + if settings.queue_backend != "redis": + log.error( + "queue_backend is %r. worker.py is only for the redis backend - " + "with the memory backend the API process runs jobs itself.", + settings.queue_backend, + ) + return 2 + + try: + from redis import Redis + from rq import Queue, Worker + except ImportError: + log.error("missing dependencies: pip install redis rq") + return 2 + + conn = Redis.from_url(settings.redis_url) + queue = Queue("subplz-jobs", connection=conn) + + log.info( + "worker starting | redis=%s | device=%s model=%s | storage=%s", + settings.redis_url, settings.device, settings.model, + settings.storage_backend, + ) + # burst=False: stay alive and keep taking jobs. + Worker([queue], connection=conn).work(with_scheduler=False) + return 0 + + +if __name__ == "__main__": + sys.exit(main())