Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 35 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,41 @@ metadata and the backend fallback mirror it.
- License notice: commercial use is free under the AGPL; the paid licence is for closed-source use, with Pro plans linked (#2578)

### Fixed

- Keep saved legacy Spanish pronunciation scopes active without guessing ambiguous or unrelated language codes (#2585)

- Preserve subtitle edits made during transcription, keep dub lock waits off the event loop, report committed tracks as complete after late cancellation, and prevent cancelled ingest history from referencing deleted files (#2585)
- Dub publication keeps file and database work off the event loop, preserves source metadata, waits safely on cancellation, and restores audio after save failures (#2585)
- Dubbing finishes when quality-check annotations arrive during assembly and clears measurements of replaced audio while still protecting subtitle edits (#2585)

- Saving a voice design skips cold engine loading and downloads, including when a warm engine unloads during the save (#2583) — thanks @simoncheese!

- Electron streaming previews drain the final PCM chunk and crossfade recovered chunks only while audio overlaps (#2518) — thanks @rudycelekli!
- Allow application-data relocation into existing empty folders without removing files added during copying (#2521) — thanks @rudycelekli!
- Preserve models folder names containing comment characters across restart (#2519) — thanks @rudycelekli!
- Re-render cached longform audio when its synthesis language changes (#2524) — thanks @rudycelekli!
- Reusing Clone takes restores their saved WAV precision and mastering controls (#2526) — thanks @rudycelekli!
- Keep Electron remote connection and WebSocket-ticket deadlines active while reading response bodies (#2527) — thanks @rudycelekli!
- Preserve CRLF and CR metadata paragraphs in longform audio exports (#2528) — thanks @rudycelekli!
- Cancelled dictation starts cannot replace the next session after a delayed connection ticket arrives (#2533) — thanks @rudycelekli!
- Electron remote WebSockets retain the selected backend path prefix and path-bound tickets (#2537) — thanks @rudycelekli!
- Refresh longform audio after changing a voice reference and keep prior clips usable until active renders finish (#2535) — thanks @rudycelekli!
- Retire cancelled longform streams and failed setup jobs while preserving resume checkpoints (#2536) — thanks @rudycelekli!
- Buffered dictation utterances receive distinct saved-history IDs so deleting one preserves the others (#2538) — thanks @rudycelekli!
- Compare voices warns when a generated preview omits speech, while keeping the surviving audio playable (#2548) — thanks @rudycelekli!
- Pronunciation scopes match language picker names and ISO codes without confusing Spanish and Estonian (#2542) — thanks @rudycelekli!
- Keep batch retries and deletion from racing over active job files (#2547) — thanks @rudycelekli!
- Dictionary backups preserve duplicate-entry pronunciation order across preview, synthesis and restore (#2552) — thanks @rudycelekli!
- Gallery trimming no longer stalls waiting for audio metadata before decoding (#2558) — thanks @rudycelekli!
- Native exports preserve existing files when a replacement write fails and retry temporary file locks (#2560) — thanks @rudycelekli!
- Runtime setup keeps concurrent package download progress separate across equivalent package spellings (#2562) — thanks @rudycelekli!
- Honor storage scan budgets in large flat directories (#2564) — thanks @rudycelekli!
- Clean visual-context frame directories after worker completion (#2566) — thanks @rudycelekli!
- Preserve concurrent partial MCP binding edits (#2568) — thanks @rudycelekli!
- Preserve CPU forced-alignment fallback when loading the MPS aligner fails (#2570) — thanks @rudycelekli!
- Keep untimed transcript segments alongside precise word-timed speech (#2572) — thanks @rudycelekli!
- Score dub quality against the selected track’s saved language text, reject stale checks, and preserve audio and subtitle edits when generation conflicts (#2574) — thanks @rudycelekli!
- Accept Japanese Han letters in translation and refinement script checks (#2576) — thanks @rudycelekli!
- Elevated Windows app removal stops before deleting data and points to a normal PowerShell window or Settings (#2578)
- Contributor audits inspect committed files and exclude submodules, while still stopping on failed file attribution (#2556)

Expand Down
6 changes: 4 additions & 2 deletions backend/api/routers/archetypes.py
Original file line number Diff line number Diff line change
Expand Up @@ -306,7 +306,7 @@ def _is_unusable_audio(audio_tensor) -> bool:
return flatness is not None and flatness < _DEGENERATE_FLATNESS


async def _render_archetype_wav(a: dict, out_path: Path) -> None:
async def _render_archetype_wav(a: dict, out_path: Path, *, allow_model_load: bool = True) -> None:
"""Render an archetype's sample script to ``out_path`` using the live engine.

Reuses generation.py's inference primitives so there is exactly one TTS code
Expand All @@ -321,7 +321,9 @@ async def _render_archetype_wav(a: dict, out_path: Path) -> None:
_safe_torchaudio_save,
)

model = await get_model()
# Save-time samples are optional: an unload after the residency probe must
# not turn persistence into a cold load. Explicit previews retain loading.
model = await get_model() if allow_model_load else await get_model(allow_load=False)
language = a["language"]
if language in (None, "", "Auto"):
language = None
Expand Down
107 changes: 76 additions & 31 deletions backend/api/routers/audiobook.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
"""

import asyncio
import anyio
import json
import logging
import os
Expand Down Expand Up @@ -587,8 +588,11 @@ def _build_synth(
from services.tts_backend import OmniVoiceBackend, active_backend_id, get_backend_class

opts = opts or ExpressiveOptions()
from core.voice_reference_snapshots import VoiceReferenceSnapshot, voice_file_lock

cache: dict = {}
token_cache: dict = {}
references = VoiceReferenceSnapshot()

def resolve(voice_id):
# Translate the span token ([voice:NAME] / exact id / None) to a profile
Expand All @@ -598,9 +602,13 @@ def resolve(voice_id):
if voice_id not in token_cache:
token_cache[voice_id] = _map_span_voice(voice_id, default_voice, voice_map)
key = token_cache[voice_id]
if key not in cache:
cache[key] = _resolve_voice(key)
return cache[key]
# Resolve and claim custody atomically with profile file deletion.
# The closure keeps the snapshot alive for any pending render worker.
with voice_file_lock:
if key not in cache:
cache[key] = _resolve_voice(key)
references.retain(cache[key]["ref_audio"])
return cache[key]

engine_id = active_backend_id()
cls = get_backend_class(engine_id)
Expand Down Expand Up @@ -778,6 +786,12 @@ def _render_chapter_cached(chapter, synth, sr, engine_id, resolve, cache_dir, le
seg_extra_sig = f"{lex_sig}\x00{expr_sig}" if expr_sig else lex_sig
if vmap_sig:
seg_extra_sig = f"{seg_extra_sig}\x00{vmap_sig}"
# Language reaches the engine even when normalization leaves text unchanged.
# Partition both layers; unknown-language legacy audio cannot satisfy an
# explicit language. Autodetect keeps its released cache derivation.
if language:
sig["\x00language"] = language
seg_extra_sig = f"{seg_extra_sig}\x00language={json.dumps(language)}"
marking = will_mark()
if marking:
# Provenance-marked chapters cache under their own key (#1169): a
Expand Down Expand Up @@ -808,7 +822,7 @@ def _render_chapter_cached(chapter, synth, sr, engine_id, resolve, cache_dir, le
inputs: dict = {
"sample rate": sr, "engine": engine_id, "normalized text": spans_tuples,
"pronunciation lexicon": lex_sig, "expressive settings": expr_sig,
"voice map": vmap_sig, "watermark": marking,
"voice map": vmap_sig, "watermark": marking, "language": language,
}
for k, v in resolved.items():
label = f"voice {re.sub(r'[^A-Za-z0-9_-]', '', k)[:40] or '(default)'}"
Expand Down Expand Up @@ -1116,35 +1130,41 @@ async def _render_longform_sse(

def _emit(payload: dict) -> str:
if job_store is not None:
if payload.get("type") == "error":
try:
if (job_store.get(job_id) or {}).get("status") in ("pending", "running"):
job_store.mark_failed(job_id, payload.get("error") or "render failed")
except Exception:
pass # terminal setup errors must still reach the client
try:
job_store.append_event(job_id, json.dumps(payload))
except Exception:
pass # best-effort job history; never block the stream
return f"data: {json.dumps(payload)}\n\n"

if not plan.chapters:
yield _emit({"type": "error", "error": "nothing to render (no chapters)"})
return
ffmpeg = find_ffmpeg()
if not ffmpeg:
yield _emit({"type": "error", "error": "ffmpeg not available; the output needs it"})
return

# Confined work dir (job_id is already token-sanitized above; work_dir adds
# the basename + realpath barrier so CodeQL sees a clean path).
work = longform_resume.work_dir(job_type, job_id)
if work is None:
yield _emit({"type": "error", "error": "invalid job id"})
return
os.makedirs(work, exist_ok=True)
# Chapter WAVs are content-addressed in a shared cache so a re-run (after a
# failure or interruption) reuses what already rendered — only the
# missing/changed chapters synthesize again (resume). Shared across both
# front doors: an identical chapter renders once.
cache_dir = os.path.join(OUTPUTS_DIR, LONGFORM_CACHE_SUBDIR)
os.makedirs(cache_dir, exist_ok=True)
prune_cache_dir(cache_dir) # bound disk before this job adds its chapters
try:
if not plan.chapters:
yield _emit({"type": "error", "error": "nothing to render (no chapters)"})
return
ffmpeg = find_ffmpeg()
if not ffmpeg:
yield _emit({"type": "error", "error": "ffmpeg not available; the output needs it"})
return

# Confined work dir (job_id is already token-sanitized above; work_dir adds
# the basename + realpath barrier so CodeQL sees a clean path).
work = longform_resume.work_dir(job_type, job_id)
if work is None:
yield _emit({"type": "error", "error": "invalid job id"})
return
os.makedirs(work, exist_ok=True)
# Chapter WAVs are content-addressed in a shared cache so a re-run (after a
# failure or interruption) reuses what already rendered — only the
# missing/changed chapters synthesize again (resume). Shared across both
# front doors: an identical chapter renders once.
cache_dir = os.path.join(OUTPUTS_DIR, LONGFORM_CACHE_SUBDIR)
os.makedirs(cache_dir, exist_ok=True)
prune_cache_dir(cache_dir) # bound disk before this job adds its chapters
resolved_lang = _resolve_default_language(language, default_voice)
operation = "audiobook" if job_type == "audiobook" else "longform"
decision = gpu_gateway.decide(operation)
Expand Down Expand Up @@ -1285,7 +1305,7 @@ def _emit(payload: dict) -> str:

yield _emit({"type": "assembling"})
meta_path = os.path.join(work, "chapters.ffmeta")
with open(meta_path, "w", encoding="utf-8") as f:
with open(meta_path, "w", encoding="utf-8", newline="") as f:
f.write(build_ffmetadata(chapters_meta, global_meta=metadata))
concat_path = os.path.join(work, "concat.txt")
with open(concat_path, "w", encoding="utf-8") as f:
Expand Down Expand Up @@ -1356,6 +1376,16 @@ def _emit(payload: dict) -> str:
"measured_i": measured.input_i if measured else None,
}
yield _emit(done)
except (asyncio.CancelledError, GeneratorExit):
# Transport cancellation/iterator closure bypass Exception; keep the
# checkpoint, but do not leave this finished response recorded as live.
if job_store is not None:
try:
if (job_store.get(job_id) or {}).get("status") in ("pending", "running"):
job_store.mark_cancelled(job_id)
except Exception:
pass # job history is best-effort, including during shutdown
raise
except Exception as e: # surface, don't 500 the stream
logger.exception("[%s] longform render failed", job_id)
if job_store is not None:
Expand All @@ -1369,8 +1399,10 @@ def _emit(payload: dict) -> str:

async def _public_longform_stream(plan, **render_kwargs):
"""Keep generator diagnostics local if setup fails before its own guard."""
stream = None
try:
async for event in _render_longform_sse(plan, **render_kwargs):
stream = _render_longform_sse(plan, **render_kwargs)
async for event in stream:
yield event
except asyncio.CancelledError:
raise
Expand All @@ -1385,6 +1417,19 @@ async def _public_longform_stream(plan, **render_kwargs):
)
yield f"data: {json.dumps({'type': 'error', 'error': error})}\n\n"

finally:
if stream is not None:
await stream.aclose()

class _ClosingLongformResponse(StreamingResponse):
async def __call__(self, scope, receive, send):
try:
await super().__call__(scope, receive, send)
finally:
# Older ASGI disconnects and outer AnyIO scopes can cancel cleanup.
with anyio.CancelScope(shield=True):
await self.body_iterator.aclose()


@router.post("/audiobook")
async def audiobook_synthesize(req: AudiobookRequest, request: Request = None):
Expand All @@ -1393,7 +1438,7 @@ async def audiobook_synthesize(req: AudiobookRequest, request: Request = None):
# `request` is injected by FastAPI on the HTTP path (the default only applies
# to a direct in-process call, e.g. a unit test); its disconnect poll is what
# lets Stop cancel the render mid-book (#1216).
return StreamingResponse(
return _ClosingLongformResponse(
_public_longform_stream(
plan, default_voice=req.default_voice, language=req.language,
fmt=req.format, bitrate=req.bitrate,
Expand Down Expand Up @@ -1458,7 +1503,7 @@ async def longform_render(req: LongformRenderRequest, request: Request = None):
if spans:
chapters.append(Chapter(title=c.title or f"Chapter {i + 1}", spans=spans))
plan = AudiobookPlan(chapters=chapters)
return StreamingResponse(
return _ClosingLongformResponse(
_public_longform_stream(
plan, default_voice=req.default_voice, language=req.language,
fmt=req.format, bitrate=req.bitrate,
Expand Down Expand Up @@ -1553,7 +1598,7 @@ def retire_checkpoint():
# job id), so the already-rendered chapters still hit instantly — only the
# unrendered ones synthesize. Using a fresh id means the request's job_id
# never names a work dir / output file (defence-in-depth path-injection).
return StreamingResponse(
return _ClosingLongformResponse(
_public_longform_stream(
plan, default_voice=p.get("default_voice"), language=p.get("language"),
fmt=p.get("fmt", "m4b"), bitrate=p.get("bitrate", "128k"),
Expand Down
51 changes: 50 additions & 1 deletion backend/api/routers/batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@
import uuid
import time
import asyncio
import contextlib
import threading
import logging
from typing import Optional, List

Expand Down Expand Up @@ -55,9 +57,25 @@
_queue: asyncio.Queue = None # Lazily initialised
_worker_task: asyncio.Task = None # Background consumer
_processing_job_ids: set[str] = set()
_mutating_job_ids: set[str] = set()
_job_mutation_lock = threading.Lock()
_jobs: dict = {} # job_id → status dict


@contextlib.contextmanager
def _job_mutation(job_id: str):
"""Reserve retry/deletion across awaits and sync endpoint threads."""
with _job_mutation_lock:
if job_id in _mutating_job_ids:
raise HTTPException(409, "A batch job action is already in progress")
_mutating_job_ids.add(job_id)
try:
yield
finally:
with _job_mutation_lock:
_mutating_job_ids.discard(job_id)


class BatchJobStatus(BaseModel):
id: str
status: str # "queued" | "running" | "done" | "failed" | "cancelled"
Expand Down Expand Up @@ -1031,6 +1049,12 @@ def cancel_batch_job(job_id: str):

@router.post("/batch/jobs/{job_id}/retry")
async def retry_batch_job(job_id: str):
"""Retry a terminal job once while protecting its input and output files."""
with _job_mutation(job_id):
return await _retry_batch_job(job_id)


async def _retry_batch_job(job_id: str):
"""Retry a terminal job using its original app-owned upload and settings."""
job = _jobs.get(job_id)
if not job:
Expand Down Expand Up @@ -1081,7 +1105,24 @@ async def retry_batch_job(job_id: str):
raise HTTPException(status_code=400, detail="Invalid batch job path")
try:
if os.path.isdir(output_dir):
await asyncio.to_thread(shutil.rmtree, output_dir)
# Cancelling the request cannot stop a filesystem worker. Keep
# custody until it settles so a second action cannot remove or
# recreate the directory while the first cleanup is still running.
cleanup = asyncio.get_running_loop().run_in_executor(None, shutil.rmtree, output_dir)
cancellation = None
while not cleanup.done():
try:
await asyncio.shield(cleanup)
except asyncio.CancelledError as exc:
cancellation = exc
except OSError:
if cancellation is None:
raise
if cancellation is not None:
with contextlib.suppress(OSError):
cleanup.result()
raise cancellation
cleanup.result()
except OSError as exc:
raise HTTPException(
status_code=500,
Expand Down Expand Up @@ -1114,10 +1155,18 @@ async def retry_batch_job(job_id: str):

@router.delete("/batch/jobs/{job_id}")
def delete_batch_job(job_id: str):
"""Delete a settled job without racing retry admission or an active worker."""
with _job_mutation(job_id):
return _delete_batch_job(job_id)


def _delete_batch_job(job_id: str):
"""Delete a batch job record and every app-owned input/output file."""
job = _jobs.get(job_id)
if not job:
raise HTTPException(404, "Job not found")
if job.get("status") in ("queued", "running") or job_id in _processing_job_ids:
raise HTTPException(409, "Cancel the batch job and wait for it to stop before deleting")
if job.get("video_path"):
try:
unlink_if_present(job["video_path"])
Expand Down
Loading
Loading