* refactor(launchpad): quieter, borderless design refresh The launchpad carried decoration from an earlier direction: icon chips, corner-hung count badges, a permanently visible filled arrow, uppercase mono card titles, and a dotted stipple divider — plus a frame that had been invisible since the app-wide border tokens were zeroed. Rework it around what the borderless direction actually implies: - Feature tiles get a whisper-faint surface instead of a dead frame, and read as three bands (bare glyph + count / title + arrow / description). `--card-hue` is spent sparingly — the glyph at rest, the surface, count and arrow only once raised. Titles move to sans sentence case; counts are plain tabular numerals. Lift softened 4px -> 2px, coloured glow -> neutral shadow, plus an explicit focus ring and a staggered entrance. - Hero drops the boxed "646" pill and the filled A/B-Compare button for quiet type, with a hairline standing in for the separation. - Section labels trade the dotted stipple for a single fading hairline; rows are transparent until hover and reveal "Open" on hover/focus (it stays in the DOM, so AT and keyboard always reach it). - Hero, tiles, recent files, callout and project lists now share one 1180px column — previously only the top half was capped, so lists ran edge-to-edge on a wide display while the deck stayed centred. Two bugs found and fixed while doing it: - Buttons that had `border border-solid border-transparent` removed fell back to the UA default border and rendered a visible 1px outline. They now carry `border-0` explicitly. - `.lp-animate` used `animation-fill-mode: both`, so after the entrance it kept owning `transform` — and animation-origin declarations outrank normal ones, which silently killed the card hover lift. Now `backwards`, which still holds the from-state through the stagger delay. Also drops CSS the page has not rendered since #904: the cursor-spotlight layer, the breath ring, and the per-card waveform strip. Verified with headless renders at 1600/1280/940 and the empty state. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(dictation): decode Wayland portal signals and show the capture pill The GlobalShortcuts portal declares Activated/Deactivated as (o session, s shortcut_id, t timestamp, a{sv} options). We decoded the timestamp as u32, so zbus rejected every signal with Signature mismatch: got `(osta{sv})`, expected `(osua{sv})` and the press was dropped as an invalid signal. Registration succeeded and the desktop even reported the bound chord back, so the hotkey looked wired up while doing nothing at all — on every Wayland compositor, for the whole life of the feature (#1490). Decode the 64-bit timestamp, and keep the 32-bit spelling as a fallback so a non-conforming portal degrades to working rather than to silence. With presses arriving, the second half of the failure showed: nothing had shown the widget window since it became a hidden recorder host, so a capture ran with no pill on screen — and a mic or Accessibility failure rendered into a window nobody could see. Add show_dictation_pill, which bottom-centres the capsule on the monitor under the pointer and shows it without taking focus (Windows keeps SW_SHOWNOACTIVATE so paste still lands in the user's document), and call it from the widget for every state but idle. Wayland denies clients their own placement, so the compositor picks the spot there; the pill still appears. dispatch_dictation_capture now logs whether a press was emitted or queued — a press that reaches Rust and produces nothing was otherwise indistinguishable from one the compositor never delivered. Tests: portal signals decode at both timestamp widths (the 64-bit case fails before this change with the exact production error); pill placement centres, respects a second monitor's origin, and clamps rather than going off-screen; the widget shows for a state needing the user, stays hidden while idle, and never shows for a press that arrives while dictation is disabled. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * chore: sync in-progress workspace changes Uncommitted work already in the tree, checkpointed so the branch matches the local machine: - Remote GPU workers: join-from-the-app flow, one-time secrets, QR join codes, a Compute control in the status bar, and the device-list Workers panel (#1516) - Model Catalogue workspace, with Settings pointing at it - Settings sidebar search and keyboard navigation - Demo assets for dubbing, dictation and voice design, plus the scripts that render them - Backend: validation-error handling, ASR request-path degradation, and the accompanying tests - CHANGELOG entries for the above and for the Wayland dictation fix Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * fix(tests): follow Engines to the Model Catalogue, and green the sweep - test_supertonic3 asserted the license gate points at "Settings" while the engine now names Model Catalogue → Engines, which is where the accept button actually lives. The assertion follows the move; what it pins is unchanged — the hint must name a place the user can reach it. - Carries the CJK allowlist entries for the rendered dub bundle (#1517) and the regenerated route snapshot for /workers/agent (#1516), both of which this branch inherits from the workspace sync. - docs/install/linux.md: the dictation capsule is bottom-anchored everywhere except Wayland, where the protocol gives applications no say in their placement. Documented rather than left as a surprise (CodeRabbit). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * ci: stop a flaky dependency fetch from failing green runs en-core-web-sm resolves to a direct GitHub release URL, and github.com intermittently answers `http2 error: refused stream before processing any application logic`. uv's own three retries all land within the same few seconds and fail together, so the whole job dies on a dependency that has nothing to do with the change under test — it cost #1518 and #1517 an otherwise-green run tonight. Two changes: back off between whole `uv sync` attempts, which is what actually clears it, and pass --no-sync to the pytest steps. `uv run` re-resolves the environment before running, so every test step was a fresh chance to hit the same fetch even though the install step had already synced — that is exactly how #1518 failed, in the isolated backend/tests step, with all 5467 tests already passed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * ci: one retry seam for every uv sync, not just the job that failed last en-core-web-sm resolves to a direct GitHub *release* URL rather than a package index, and github.com intermittently answers `http2 error: refused stream before processing any application logic`. uv's own retries all land inside the same ~10 seconds and fail together, so a job dies on a dependency unrelated to the change under test. Tonight that cost four otherwise-green runs across #1515, #1517 and #1518 — and the first fix only covered the Tests job, so the next failure simply moved to Smoke (Linux), which syncs separately. The fetch is per-job, so the fix has to be per-job: scripts/uv-sync-retry.sh backs off between whole attempts (15s, 45s, 90s) and every workflow that syncs now goes through it — ci.yml (tests + the platform matrix), release.yml, security.yml, evals.yml. It still fails loudly after four attempts, so a genuinely broken lockfile is not disguised as a flake. The Tests job also lacked the UV_HTTP_TIMEOUT / UV_HTTP_RETRIES the smoke matrix has always set, which is part of why it was the one that kept dying; it has them now. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * test(ci): pin the Intel-Mac contract by intent, not by command spelling test_ci_verifies_intel_mac_as_the_documented_remote_only_host asserted the literal line `run: uv sync --extra pockettts`, so routing every sync through scripts/uv-sync-retry.sh read as a broken Intel-Mac contract. The contract it exists to protect is that the pockettts extra installs ONLY on backend_supported legs — which the regex now pins, while leaving how the sync is invoked free to change. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> * ci: keep every uv run out of the resolver, and bound the retry budget CodeRabbit, #1517: - `uv run` re-resolves before running, so the smoke suite, the worker-artifact tests, the release test run and the eval run were each a fresh chance to hit the flaky direct-URL fetch outside the retry loop. All of them pass --no-sync now; the environment is already synced by the step that owns the retries. security.yml's `uv run --with pip-audit` is deliberately left alone — it layers an ephemeral package rather than running the project's own tests. - The retry count multiplied uv's own budget (UV_HTTP_RETRIES=5 with a 120 s timeout on the smoke matrix). Three attempts and 60 s of total backoff outlast the refusals actually observed while staying well inside the jobs' timeout-minutes. - The Intel-Mac contract test pinned the smoke command literally too, so --no-sync tripped it exactly like the sync line did. Same fix: assert the contract (smoke runs only on backend_supported legs), not its spelling. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
621 lines
26 KiB
Python
621 lines
26 KiB
Python
"""Batch dubbing queue — POST videos with settings, process sequentially.
|
|
|
|
This is a lightweight batch orchestrator. Each job is a dub project that
|
|
runs through the same ingest→transcribe→translate→generate pipeline as
|
|
a manual dub, but driven by the queue instead of the UI.
|
|
|
|
The queue is in-memory (lives for the process lifetime). Jobs persist to
|
|
the SQLite `jobs` table for history, but the queue itself restarts empty
|
|
on backend restart — intentional, since GPU jobs can't be safely resumed.
|
|
"""
|
|
import os
|
|
import uuid
|
|
import time
|
|
import asyncio
|
|
import logging
|
|
from typing import Optional, List
|
|
|
|
from fastapi import APIRouter, File, UploadFile, HTTPException, Form
|
|
from pydantic import BaseModel
|
|
|
|
from core.config import DATA_DIR
|
|
from core import failure
|
|
from core.logging_utils import log_safe
|
|
from core.file_cleanup import FileCleanupError, unlink_if_present
|
|
|
|
router = APIRouter()
|
|
logger = logging.getLogger("omnivoice.batch")
|
|
|
|
# ── In-memory queue ─────────────────────────────────────────────────────
|
|
|
|
_queue: asyncio.Queue = None # Lazily initialised
|
|
_worker_task: asyncio.Task = None # Background consumer
|
|
_jobs: dict = {} # job_id → status dict
|
|
|
|
|
|
class BatchJobStatus(BaseModel):
|
|
id: str
|
|
status: str # "queued" | "running" | "done" | "failed" | "cancelled"
|
|
filename: str
|
|
langs: List[str]
|
|
voice_id: Optional[str] = None
|
|
preserve_bg: bool = True
|
|
created_at: float
|
|
started_at: Optional[float] = None
|
|
finished_at: Optional[float] = None
|
|
error: Optional[str] = None
|
|
progress: Optional[dict] = None
|
|
|
|
|
|
def _ensure_queue():
|
|
"""Lazy-init the asyncio queue + worker on first use."""
|
|
global _queue, _worker_task
|
|
if _queue is None:
|
|
_queue = asyncio.Queue()
|
|
_worker_task = asyncio.ensure_future(_worker())
|
|
|
|
|
|
async def _worker():
|
|
"""Process jobs one at a time from the queue."""
|
|
while True:
|
|
job_id = await _queue.get()
|
|
job = _jobs.get(job_id)
|
|
if not job or job["status"] == "cancelled":
|
|
_queue.task_done()
|
|
continue
|
|
|
|
job["status"] = "running"
|
|
job["started_at"] = time.time()
|
|
logger.info("Batch job %s starting: %s", job_id, job["filename"])
|
|
|
|
try:
|
|
await _run_batch_pipeline(job_id, job)
|
|
if job["status"] != "cancelled":
|
|
job["status"] = "done"
|
|
job["finished_at"] = time.time()
|
|
logger.info(
|
|
"Batch job %s completed in %.1fs",
|
|
job_id, job["finished_at"] - job["started_at"],
|
|
)
|
|
except asyncio.CancelledError:
|
|
# Task cancellation always means SHUTDOWN: the job-level cancel
|
|
# endpoint only flips job["status"] — nothing ever cancels this
|
|
# task to abort a single job. Swallowing the CancelledError here
|
|
# made the worker unkillable (the while-loop re-entered
|
|
# _queue.get() and event-loop teardown hung forever in
|
|
# _cancel_all_tasks waiting on a task that never finishes). Mark
|
|
# the in-flight job, then let the cancellation propagate.
|
|
job["status"] = "cancelled"
|
|
job["finished_at"] = time.time()
|
|
raise
|
|
except Exception as e:
|
|
job["status"] = "failed"
|
|
# plan-04 (#131): guaranteed non-empty, structured reason.
|
|
job["error"] = failure.build_failure(e, stage="batch", include_diagnostic=False)["reason"]
|
|
job["finished_at"] = time.time()
|
|
logger.error("Batch job %s failed: %s", job_id, e, exc_info=True)
|
|
finally:
|
|
_queue.task_done()
|
|
|
|
|
|
def _set_progress(job, stage, percent=0, **extra):
|
|
"""Update a job's progress dict."""
|
|
job["progress"] = {"stage": stage, "percent": percent, **extra}
|
|
|
|
|
|
async def _run_batch_pipeline(job_id: str, job: dict):
|
|
"""Full batch dub pipeline: extract → transcribe → translate → generate → mix → export."""
|
|
import subprocess
|
|
|
|
loop = asyncio.get_running_loop()
|
|
video_path = job["video_path"]
|
|
langs = job["langs"]
|
|
batch_dir = os.path.join(DATA_DIR, "batch", job_id)
|
|
os.makedirs(batch_dir, exist_ok=True)
|
|
|
|
# ── 1. Extract audio ──────────────────────────────────────────────
|
|
_set_progress(job, "extract", 0)
|
|
audio_path = os.path.join(batch_dir, "audio.wav")
|
|
|
|
from services.ffmpeg_utils import bed_mix_filter, find_ffmpeg
|
|
ffmpeg = find_ffmpeg()
|
|
|
|
def _extract():
|
|
subprocess.run(
|
|
[ffmpeg, "-y", "-i", video_path,
|
|
"-vn", "-acodec", "pcm_s16le", "-ar", "22050", "-ac", "1",
|
|
audio_path],
|
|
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
|
|
timeout=300, check=True,
|
|
)
|
|
# Get duration
|
|
result = subprocess.run(
|
|
[ffmpeg, "-i", audio_path],
|
|
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
|
|
timeout=30,
|
|
)
|
|
import re
|
|
match = re.search(r"Duration: (\d+):(\d+):(\d+)\.(\d+)", result.stderr.decode("utf-8", errors="replace"))
|
|
if match:
|
|
h, m, s, cs = match.groups()
|
|
return int(h) * 3600 + int(m) * 60 + int(s) + int(cs) / 100
|
|
return 0.0
|
|
|
|
duration = await loop.run_in_executor(None, _extract)
|
|
job["duration"] = duration
|
|
_set_progress(job, "extract", 100)
|
|
|
|
if job["status"] == "cancelled":
|
|
return
|
|
|
|
# ── 2. Transcribe ─────────────────────────────────────────────────
|
|
_set_progress(job, "transcribe", 0)
|
|
|
|
from services.asr_backend import load_active_asr_backend
|
|
from services.model_manager import _gpu_pool, _cpu_pool, run_on_gpu_pool_guarded
|
|
from services.segmentation import (
|
|
segment_transcript, assign_speakers_heuristic,
|
|
)
|
|
|
|
def _transcribe():
|
|
# `load_*`, not `get_*`: the plain selector returns engines whose
|
|
# shallow probe passed but whose deep import chain is broken, failing
|
|
# the whole batch job at `.transcribe()` instead of degrading (#1185).
|
|
backend = load_active_asr_backend()
|
|
result = backend.transcribe(audio_path, word_timestamps=True)
|
|
detected_lang = result.get("language", "en")
|
|
segments = segment_transcript(result, duration=duration)
|
|
segments = assign_speakers_heuristic(segments)
|
|
for i, s in enumerate(segments):
|
|
s["id"] = f"s{i:05x}"
|
|
s.setdefault("text_original", s.get("text", ""))
|
|
try:
|
|
backend.unload()
|
|
except Exception:
|
|
pass
|
|
return segments, detected_lang
|
|
|
|
# Bound the batch transcribe (#730) so a wedged whisperx/CTranslate2 call
|
|
# can't hold its GPU-pool worker forever and starve the rest of the backend
|
|
# ("can't reach backend"); run_transcribe_guarded also resets the pool on
|
|
# timeout to restore capacity.
|
|
from services.asr_backend import run_transcribe_guarded
|
|
segments, source_lang = await run_transcribe_guarded(_gpu_pool, _transcribe, what="Batch")
|
|
source_lang = (source_lang or "en").split("_")[0][:2].lower()
|
|
job["segments"] = segments
|
|
job["source_lang"] = source_lang
|
|
_set_progress(job, "transcribe", 100, segments_count=len(segments))
|
|
|
|
if job["status"] == "cancelled" or not segments:
|
|
if not segments:
|
|
job["error"] = "Transcription produced no segments"
|
|
job["status"] = "failed"
|
|
return
|
|
|
|
# ── Engine resolution (issue #312 class) ────────────────────────────
|
|
# Batch used to hardcode VoiceStudio via get_model() regardless of the
|
|
# engine selected in Model Catalogue → Engines. require_cloning only when a
|
|
# specific voice is pinned (job["voice_id"]) — an unpinned job is fine on
|
|
# any active engine. Resolved ONCE for the whole job (every language
|
|
# below shares the same active engine); an uncaught ValueError here
|
|
# propagates to _worker()'s existing except-Exception handling, which
|
|
# already records a structured job failure via core.failure.build_failure.
|
|
from services.tts_backend import resolve_generation_backend
|
|
backend = await resolve_generation_backend(
|
|
require_cloning=bool(job.get("voice_id")),
|
|
cloning_purpose="this batch job's pinned voice",
|
|
)
|
|
sr = backend.sample_rate
|
|
|
|
# ── 3. Translate + Generate per language ───────────────────────────
|
|
total_langs = len(langs)
|
|
outputs = {}
|
|
|
|
for lang_idx, target_lang in enumerate(langs):
|
|
if job["status"] == "cancelled":
|
|
return
|
|
|
|
# ── 3a. Translate ─────────────────────────────────────────────
|
|
_set_progress(
|
|
job, "translate",
|
|
percent=int((lang_idx / total_langs) * 100),
|
|
current_lang=target_lang,
|
|
)
|
|
|
|
translated_segments = list(segments) # copy
|
|
if target_lang != source_lang:
|
|
try:
|
|
def _translate_batch(segs, src, tgt):
|
|
"""Translate segment texts via Google Translate."""
|
|
from deep_translator import GoogleTranslator
|
|
TRANSLATE_CODES = {
|
|
"en": "en", "es": "es", "fr": "fr", "de": "de",
|
|
"it": "it", "pt": "pt", "ru": "ru", "ja": "ja",
|
|
"ko": "ko", "zh": "zh-CN", "ar": "ar", "hi": "hi",
|
|
"tr": "tr", "pl": "pl", "nl": "nl", "sv": "sv",
|
|
}
|
|
src_code = TRANSLATE_CODES.get(src, src) or "auto"
|
|
tgt_code = TRANSLATE_CODES.get(tgt, tgt)
|
|
translator = GoogleTranslator(source=src_code, target=tgt_code)
|
|
out = []
|
|
for s in segs:
|
|
s_copy = dict(s)
|
|
text = s.get("text", "").strip()
|
|
if text:
|
|
try:
|
|
s_copy["text"] = translator.translate(text) or text
|
|
except Exception as e:
|
|
logger.warning("Translate seg failed: %s", e)
|
|
out.append(s_copy)
|
|
return out
|
|
|
|
translated_segments = await loop.run_in_executor(
|
|
_cpu_pool, _translate_batch,
|
|
segments, source_lang, target_lang,
|
|
)
|
|
except ImportError:
|
|
logger.warning("deep_translator not installed, skipping translation for %s", target_lang)
|
|
except Exception as e:
|
|
logger.warning("Translation failed for %s: %s, using original", target_lang, e)
|
|
translated_segments = segments
|
|
|
|
if job["status"] == "cancelled":
|
|
return
|
|
|
|
# ── 3b. Generate TTS ──────────────────────────────────────────
|
|
_set_progress(
|
|
job, "generate",
|
|
percent=int((lang_idx / total_langs) * 100),
|
|
current_lang=target_lang,
|
|
current_segment=0,
|
|
total_segments=len(translated_segments),
|
|
)
|
|
|
|
from services.audio_dsp import apply_mastering, normalize_audio
|
|
from services.audio_io import atomic_save_wav
|
|
import torch
|
|
|
|
total_samples = int(duration * sr)
|
|
full_audio = torch.zeros(1, total_samples)
|
|
total_segs = len(translated_segments)
|
|
|
|
for i, seg in enumerate(translated_segments):
|
|
if job["status"] == "cancelled":
|
|
return
|
|
|
|
_set_progress(
|
|
job, "generate",
|
|
percent=int(((lang_idx + (i / total_segs)) / total_langs) * 100),
|
|
current_lang=target_lang,
|
|
current_segment=i + 1,
|
|
total_segments=total_segs,
|
|
)
|
|
|
|
seg_start = seg.get("start", 0)
|
|
seg_end = seg.get("end", 0)
|
|
seg_duration = seg_end - seg_start
|
|
seg_text = seg.get("text", "").strip()
|
|
|
|
if seg_duration <= 0.05 or not seg_text:
|
|
continue
|
|
|
|
def _gen(text=seg_text, lang=target_lang, dur=seg_duration):
|
|
# Normalize once at the segment's text→engine choke point —
|
|
# the same pre-pass as /generate and dub_generate's _gen.
|
|
# `lang` is the job's target language code. Pref-gated,
|
|
# idempotent, never raises.
|
|
from services.text_normalization import normalize_for_tts
|
|
text = normalize_for_tts(text, lang)
|
|
|
|
ref_audio = None
|
|
ref_text = None
|
|
|
|
# Use voice_id if provided
|
|
if job.get("voice_id"):
|
|
from core.db import db_conn
|
|
from core.config import VOICES_DIR as _VD
|
|
with db_conn() as conn:
|
|
row = conn.execute(
|
|
"SELECT * FROM voice_profiles WHERE id=?",
|
|
(job["voice_id"],),
|
|
).fetchone()
|
|
if row:
|
|
if row["is_locked"] and row["locked_audio_path"]:
|
|
ref_audio = os.path.join(_VD, row["locked_audio_path"])
|
|
elif row["ref_audio_path"]:
|
|
ref_audio = os.path.join(_VD, row["ref_audio_path"])
|
|
ref_text = row.get("ref_text")
|
|
|
|
try:
|
|
audio_out = backend.generate(
|
|
text=text, language=lang,
|
|
ref_audio=ref_audio, ref_text=ref_text,
|
|
duration=dur, num_step=16,
|
|
guidance_scale=2.0, speed=1.0,
|
|
denoise=True, postprocess_output=True,
|
|
)
|
|
if not getattr(backend, "applies_own_mastering", False):
|
|
audio_out = apply_mastering(audio_out, sample_rate=sr)
|
|
return normalize_audio(audio_out, target_dBFS=-2.0)
|
|
except Exception as e:
|
|
logger.warning("TTS failed for seg %d (lang=%s): %s", i, lang, e)
|
|
# #1190: the silence still stands in for the segment (one
|
|
# bad line shouldn't bin an otherwise good dub), but it is
|
|
# no longer INVISIBLE — the job carries a warning the UI /
|
|
# API consumer can see instead of shipping a
|
|
# finished-looking track with unexplained silence.
|
|
job.setdefault("warnings", []).append(
|
|
f"Segment {i + 1} of the {lang} track failed to "
|
|
f"synthesize and was left silent: {e}"
|
|
)
|
|
return torch.zeros(1, int(dur * sr))
|
|
|
|
try:
|
|
# Bounded + pool-reset on hang so a wedged batch segment can't
|
|
# starve the GPU pool and brick the backend (#730 class).
|
|
# Budget is the shared length-scaled one (#1190): a long segment
|
|
# on CPU-class hardware no longer dies on the flat 300s.
|
|
from services.model_manager import generate_timeout_s
|
|
audio_tensor = await run_on_gpu_pool_guarded(
|
|
_gen, what="Batch generate",
|
|
timeout=generate_timeout_s(seg_text),
|
|
)
|
|
|
|
# Fit to slot
|
|
target_samples_seg = int(seg_duration * sr)
|
|
current_samples = audio_tensor.shape[-1]
|
|
if target_samples_seg > current_samples:
|
|
audio_tensor = torch.nn.functional.pad(
|
|
audio_tensor, (0, target_samples_seg - current_samples)
|
|
)
|
|
elif current_samples > target_samples_seg:
|
|
audio_tensor = audio_tensor[..., :target_samples_seg]
|
|
|
|
# Crossfade
|
|
fade_samples = int(0.015 * sr)
|
|
wl = audio_tensor.shape[-1]
|
|
if wl > fade_samples * 2:
|
|
ramp_up = torch.linspace(0, 1, fade_samples)
|
|
ramp_down = torch.linspace(1, 0, fade_samples)
|
|
audio_tensor[0, :fade_samples] *= ramp_up
|
|
audio_tensor[0, -fade_samples:] *= ramp_down
|
|
|
|
s_idx = int(seg_start * sr)
|
|
e_idx = min(s_idx + wl, total_samples)
|
|
full_audio[:, s_idx:e_idx] += audio_tensor[:, :e_idx - s_idx]
|
|
|
|
except TimeoutError as e:
|
|
# #1190/#1202: a GPU timeout (or a saturated pool) used to be
|
|
# swallowed into a silent gap in the dubbed track — the user got
|
|
# a finished-looking video with missing speech and no warning,
|
|
# and on a 1-worker host the abandoned job made every later
|
|
# segment likelier to time out too (the "22-chunk batch dies at
|
|
# chunk 3" cascade). Fail the job loudly instead: _worker()'s
|
|
# except-Exception handler records a structured failure the UI
|
|
# surfaces. Non-timeout per-segment errors keep the old
|
|
# degrade-to-gap behaviour, but are now recorded on the job.
|
|
logger.error("Batch TTS seg %d timed out — failing the job: %s", i, e)
|
|
raise RuntimeError(
|
|
f"Segment {i + 1} of the {target_lang} track did not "
|
|
f"render, so the dubbed track would have shipped with a "
|
|
f"silent gap. {e}"
|
|
) from e
|
|
except Exception as e:
|
|
logger.warning("Batch TTS seg %d failed: %s", i, e)
|
|
job.setdefault("warnings", []).append(
|
|
f"Segment {i + 1} of the {target_lang} track failed and was "
|
|
f"left silent: {e}"
|
|
)
|
|
|
|
# ── 3c. Save dubbed audio track ───────────────────────────────
|
|
# Invisible provenance mark on the assembled track (#1169), tensor
|
|
# stage, before the WAV write / aac mux — batch dubs used to ship
|
|
# unmarked while the interactive dub pipeline marked every segment.
|
|
# One whole-track embed (chunked internally, #1045) is equivalent to
|
|
# dub_generate's per-segment marks: the 16-bit message repeats
|
|
# throughout. Runs in the GPU pool like generate's finalize; never
|
|
# raises (degrades to unmarked on failure, same as every producer).
|
|
# Dispatched to the dedicated watermark pool, not the GPU pool (#1190):
|
|
# AudioSeal embedding is CPU work that holds no VRAM, and a whole-track
|
|
# embed is long enough that occupying a GPU worker with it stalled the
|
|
# next language's segments on 1-worker hosts.
|
|
from services.watermark import mark_synthetic
|
|
from services.model_manager import get_watermark_pool
|
|
import functools
|
|
full_audio = await loop.run_in_executor(
|
|
get_watermark_pool(),
|
|
functools.partial(mark_synthetic, full_audio, sr,
|
|
context="batch.dub_track"),
|
|
)
|
|
|
|
# Same assembly pattern as dub_generate.py:390 — `full_audio` is a
|
|
# zero-init tensor that gets +='d from torch.cat-style slices, so
|
|
# it can land non-contiguous + out-of-range. Go through the
|
|
# audited + atomic helper to defend against #48 silent corruption
|
|
# and partial-write truncation simultaneously.
|
|
track_path = os.path.join(batch_dir, f"dubbed_{target_lang}.wav")
|
|
atomic_save_wav(track_path, full_audio, sr)
|
|
|
|
# ── 3d. Mix with original video ───────────────────────────────
|
|
_set_progress(
|
|
job, "mix",
|
|
percent=int(((lang_idx + 0.8) / total_langs) * 100),
|
|
current_lang=target_lang,
|
|
)
|
|
|
|
output_path = os.path.join(batch_dir, f"output_{target_lang}.mp4")
|
|
|
|
def _mix(bg=job.get("preserve_bg", True)):
|
|
if bg:
|
|
# Mix dubbed audio with original background
|
|
subprocess.run(
|
|
[ffmpeg, "-y",
|
|
"-i", video_path,
|
|
"-i", track_path,
|
|
"-filter_complex",
|
|
bed_mix_filter("0:a", "1:a", out="out", duration="first"),
|
|
"-map", "0:v", "-map", "[out]",
|
|
"-c:v", "copy", "-c:a", "aac", "-b:a", "192k",
|
|
"-shortest", output_path],
|
|
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
|
|
timeout=600, check=True,
|
|
)
|
|
else:
|
|
# Replace audio entirely
|
|
subprocess.run(
|
|
[ffmpeg, "-y",
|
|
"-i", video_path,
|
|
"-i", track_path,
|
|
"-map", "0:v", "-map", "1:a",
|
|
"-c:v", "copy", "-c:a", "aac", "-b:a", "192k",
|
|
"-shortest", output_path],
|
|
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
|
|
timeout=600, check=True,
|
|
)
|
|
|
|
await loop.run_in_executor(None, _mix)
|
|
outputs[target_lang] = output_path
|
|
|
|
job["outputs"] = outputs
|
|
_set_progress(job, "done", 100)
|
|
|
|
|
|
# ── Endpoints ───────────────────────────────────────────────────────────
|
|
|
|
@router.post("/batch/enqueue")
|
|
async def enqueue_batch_job(
|
|
video: UploadFile = File(...),
|
|
langs: str = Form("es"), # comma-separated lang codes
|
|
voice_id: Optional[str] = Form(None),
|
|
preserve_bg: bool = Form(True),
|
|
):
|
|
"""Enqueue a video for batch dubbing.
|
|
|
|
The video is saved to disk and a job is added to the queue.
|
|
Returns the job ID for status polling.
|
|
"""
|
|
_ensure_queue()
|
|
|
|
job_id = str(uuid.uuid4())[:12]
|
|
lang_list = [l.strip() for l in langs.split(",") if l.strip()]
|
|
if not lang_list:
|
|
raise HTTPException(400, "At least one target language is required")
|
|
|
|
# TTS-only install: no ASR model on disk → typed 409 with a download CTA
|
|
# now, instead of accepting the job and having the transcribe stage
|
|
# silently auto-download multi-GB whisper weights (or fail) in the worker.
|
|
from services.asr_backend import asr_model_missing_detail, asr_model_missing_error
|
|
missing = await asyncio.to_thread(asr_model_missing_error)
|
|
if missing is not None:
|
|
raise HTTPException(409, {**missing, "message": asr_model_missing_detail(missing)})
|
|
|
|
# Save the uploaded video
|
|
batch_dir = os.path.join(DATA_DIR, "batch")
|
|
os.makedirs(batch_dir, exist_ok=True)
|
|
ext = os.path.splitext(video.filename or "video.mp4")[1] or ".mp4"
|
|
video_path = os.path.join(batch_dir, f"{job_id}{ext}")
|
|
|
|
with open(video_path, "wb") as f:
|
|
content = await video.read()
|
|
f.write(content)
|
|
|
|
job = {
|
|
"id": job_id,
|
|
"status": "queued",
|
|
"filename": video.filename or f"{job_id}{ext}",
|
|
"video_path": video_path,
|
|
"langs": lang_list,
|
|
"voice_id": voice_id,
|
|
"preserve_bg": preserve_bg,
|
|
"created_at": time.time(),
|
|
"started_at": None,
|
|
"finished_at": None,
|
|
"error": None,
|
|
"progress": None,
|
|
}
|
|
_jobs[job_id] = job
|
|
await _queue.put(job_id)
|
|
|
|
logger.info(
|
|
"Batch job %s enqueued (%d target languages)",
|
|
log_safe(job_id), len(lang_list),
|
|
)
|
|
return {"job_id": job_id, "status": "queued", "queue_position": _queue.qsize()}
|
|
|
|
|
|
@router.get("/batch/jobs")
|
|
def list_batch_jobs(status: Optional[str] = None, limit: int = 50):
|
|
"""List batch jobs, optionally filtered by status."""
|
|
jobs = list(_jobs.values())
|
|
if status:
|
|
if status == "active":
|
|
jobs = [j for j in jobs if j["status"] in ("queued", "running")]
|
|
else:
|
|
jobs = [j for j in jobs if j["status"] == status]
|
|
jobs.sort(key=lambda j: j["created_at"], reverse=True)
|
|
return jobs[:limit]
|
|
|
|
|
|
@router.get("/batch/jobs/{job_id}")
|
|
def get_batch_job(job_id: str):
|
|
"""Get the status of a specific batch job."""
|
|
job = _jobs.get(job_id)
|
|
if not job:
|
|
raise HTTPException(404, "Job not found")
|
|
return job
|
|
|
|
|
|
@router.post("/batch/jobs/{job_id}/cancel")
|
|
def cancel_batch_job(job_id: str):
|
|
"""Cancel a queued or running batch job."""
|
|
job = _jobs.get(job_id)
|
|
if not job:
|
|
raise HTTPException(404, "Job not found")
|
|
if job["status"] in ("done", "failed", "cancelled"):
|
|
return {"already": job["status"]}
|
|
job["status"] = "cancelled"
|
|
job["finished_at"] = time.time()
|
|
return {"cancelled": True}
|
|
|
|
|
|
@router.delete("/batch/jobs/{job_id}")
|
|
def delete_batch_job(job_id: str):
|
|
"""Delete a batch job record and its video file."""
|
|
job = _jobs.get(job_id)
|
|
if not job:
|
|
raise HTTPException(404, "Job not found")
|
|
if job.get("video_path"):
|
|
try:
|
|
unlink_if_present(job["video_path"])
|
|
except FileCleanupError as exc:
|
|
raise HTTPException(
|
|
status_code=500,
|
|
detail="Could not delete the batch video file. Close any app using it and retry.",
|
|
) from exc
|
|
_jobs.pop(job_id, None)
|
|
return {"deleted": True}
|
|
|
|
|
|
@router.get("/batch/download/{job_id}/{lang}")
|
|
def download_batch_output(job_id: str, lang: str):
|
|
"""Download a completed batch job's output video for a given language."""
|
|
from fastapi.responses import FileResponse
|
|
|
|
job = _jobs.get(job_id)
|
|
if not job:
|
|
raise HTTPException(404, "Job not found")
|
|
if job["status"] != "done":
|
|
raise HTTPException(400, f"Job is {job['status']}, not done")
|
|
|
|
outputs = job.get("outputs", {})
|
|
path = outputs.get(lang)
|
|
if not path or not os.path.exists(path):
|
|
raise HTTPException(404, f"No output for language '{lang}'")
|
|
|
|
filename = f"{os.path.splitext(job['filename'])[0]}_{lang}.mp4"
|
|
return FileResponse(
|
|
path,
|
|
media_type="video/mp4",
|
|
filename=filename,
|
|
)
|