How VoiceStudio's Batch Queue Handles Progress, Folder Watching, and Pinned-Voice Routing with Incompatible Engines
VoiceStudio's batch queue uses an async folder watcher to detect new jobs, creates prioritized Task objects, and automatically falls back to compatible engines or per-segment generation when a pinned voice engine lacks native batch support.
The batch queue in the debpalash/VoiceStudio repository is the core orchestration layer for large-scale TTS and dubbing workflows. It continuously monitors a progress/ directory for incoming job folders, transforms each discovery into a trackable Task, and intelligently routes work to voice engines—gracefully handling cases where a user-pinned engine cannot process batch operations.
Folder Watcher Architecture
The queue relies on a filesystem watcher to trigger job processing without explicit API calls. This design enables integration with external systems that simply drop job manifests into a monitored directory.
Async Watcher Implementation
In api/routers/batch.py, the watcher runs as a background task during FastAPI startup:
# api/routers/batch.py (simplified)
async def _watch_progress_folder():
async for event in awatch(PROGRESS_ROOT):
if event.is_dir_created:
await _enqueue_job(event.path)
The watcher uses an asyncio-compatible observer (such as watchdog with async support) to detect directory creation events in real time.
Job Enqueueing
Once a folder appears, _enqueue_job loads the manifest and creates a prioritized task:
async def _enqueue_job(job_path: Path):
job = await _load_job_manifest(job_path)
task = Task(
job_id=job.id,
priority=PriorityClass.BATCH,
payload=job
)
task_store.create(task, now=time.time())
By default, batch tasks receive PriorityClass.BATCH, which ranks below PriorityClass.INTERACTIVE. This ensures that on-demand requests are prioritized over background workloads.
Task Priority and Scheduling
The TaskStore manages queue ordering based on priority classes. The implementation in api/routers/batch.py → TaskStore.create enforces strict precedence rules.
Verified behavior from tests/test_worker_scheduler.py:
- Batch tasks enter the queue with status
QUEUED - Interactive tasks always outrank batch tasks when workers poll for work
- Priority inversion is prevented through timestamp-based tiebreaking
This guarantees responsive user-facing generation even under heavy batch load.
Pinned-Voice Routing and Engine Resolution
Users can pin a job to a specific voice engine via the engine_id field in the job manifest. The resolution logic in _resolve_engine balances user preference against technical capability.
Engine Registry Lookup
The global _REGISTRY in api/engines/_registry.py contains all installed voice backends. Resolution follows this order:
- Check if
engine_idis specified and registered - Verify
engine.supports_batchproperty - If incompatible, fall back through the chain
async def _resolve_engine(task: Task) -> Engine:
pinned = task.payload.get("engine_id")
if pinned and pinned in _REGISTRY:
engine = _REGISTRY[pinned]
if engine.supports_batch:
return engine
# Incompatible → proceed to fallback
return _REGISTRY.get(DEFAULT_BATCH_ENGINE) or _fallback_engine()
Incompatible-Engine Fallback Strategies
When a pinned engine cannot handle batch operations, the queue implements a two-tier fallback in _run_batch_pipeline (located in api/routers/batch.py):
Tier 1: Native Batch Engine
If DEFAULT_BATCH_ENGINE exists in _REGISTRY and reports supports_batch=True, the job is transparently redirected:
engine = _REGISTRY.get(pinned_id)
if engine and not engine.supports_batch:
logger.warning("Pinned engine %s lacks batch support – using fallback.", pinned_id)
engine = _REGISTRY.get(DEFAULT_BATCH_ENGINE)
This fallback is exercised in tests/test_text_normalization_routes.py, which verifies that batch width is respected (e.g., 8 segments then 2 segments for a 10-segment job).
Tier 2: Per-Segment Generation
If no native batch engine is available, the pipeline decomposes the job into individual generation calls:
async def _run_batch_pipeline(job_id: str, job: Job):
engine = await _resolve_engine(task)
if engine.supports_batch:
await engine.generate_batch(job.segments, **job.opts)
else:
# Per-segment fallback
for seg in job.segments:
await engine.generate(seg, **job.opts)
Performance implication: Per-segment generation incurs higher overhead but guarantees completion. The tests/test_perf_operation_budgets.py suite validates that native batch engines never trigger this slower path unexpectedly.
Post-Processing and Watermarking
After audio generation completes, the queue handles:
- File placement: Results written to
job_path / "dubbed_<lang>.wav" - Watermark injection: If the
markingflag is set,api/routers/watermark.pyprocesses the audio
The watermark integration is verified in tests/test_synthetic_audio_watermark_1169.py, confirming that batch-dubbed tracks carry correct synthetic audio markers.
Complete Workflow Example
Client: Creating a Batch Job
import httpx, pathlib, json
payload = {
"text": ["Hello", "world", "..."],
"engine_id": "my-pinned-engine", # optional; may be incompatible
"options": {"duration": [1.0, 1.0, 1.0]}
}
job_dir = pathlib.Path("/tmp/progress") / "job123"
job_dir.mkdir(parents=True)
(job_dir / "manifest.json").write_text(json.dumps(payload))
# Server's folder watcher detects job123 automatically
Server: Fallback Handling
# Excerpt from api/routers/batch.py
engine = _REGISTRY.get(pinned_id)
if engine and not engine.supports_batch:
logger.warning("Pinned engine %s lacks batch support – using fallback.", pinned_id)
engine = _REGISTRY.get(DEFAULT_BATCH_ENGINE)
Key Source Files
| File | Purpose |
|---|---|
api/routers/batch.py |
Main batch-queue implementation, folder watcher, task enqueuing, engine resolution, fallback logic |
api/engines/_registry.py |
Global voice-engine registry (_REGISTRY) |
tests/test_worker_scheduler.py |
Priority queue behavior and task state verification |
tests/test_text_normalization_routes.py |
Pinned-engine selection and native-batch fallback testing |
tests/test_perf_operation_budgets.py |
Validates native-batch operation budgets |
tests/test_synthetic_audio_watermark_1169.py |
Confirms watermark application on batch output |
Summary
- Folder watching uses an async filesystem observer in
api/routers/batch.pyto trigger job ingestion without explicit API calls - Task priority separates interactive and batch workloads, with
INTERACTIVEalways outrankingBATCH - Pinned-voice routing respects user preference through
engine_idmetadata, with capability checks against_REGISTRY - Incompatible-engine fallback first attempts
DEFAULT_BATCH_ENGINE, then degrades to per-segment generation if necessary - Graceful degradation ensures job completion regardless of engine compatibility, verified across multiple test suites
Frequently Asked Questions
What triggers the batch queue to start processing a job?
The folder watcher in api/routers/batch.py detects new directory creation events in the progress/ folder. When a job folder appears, _enqueue_job is invoked automatically—no HTTP request required. This enables drop-in integration with external job submission systems.
How does the queue handle a pinned engine that cannot process batches?
The _resolve_engine function checks engine.supports_batch. If false, it logs a warning and falls back to DEFAULT_BATCH_ENGINE. If no batch-capable engine exists, _run_batch_pipeline switches to per-segment generation. Both paths guarantee job completion.
Why are there two priority classes for tasks?
INTERACTIVE and BATCH priority classes ensure responsive user-facing generation. Tests in tests/test_worker_scheduler.py confirm that interactive tasks are always serviced before batch tasks, preventing latency spikes for on-demand requests during heavy background workloads.
Where is the fallback behavior tested?
Three test files cover fallback scenarios: tests/test_text_normalization_routes.py verifies engine selection and native-batch width, tests/test_perf_operation_budgets.py ensures budgeted operation counts, and tests/test_worker_scheduler.py validates priority ordering throughout the fallback chain.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →