Command Lifecycles
Status: Current Last updated: 2026-09-15 19:40 EDT
End-to-end sequence diagrams showing how jobs flow through the system, from CLI invocation to output files. Every batchalign command now fits one of the explicit workflow families surfaced in the new contributor-facing architecture: per-file transform, cross-file batch transform, reference projection, composite workflow, or media analysis. For per-command option-driven flowcharts, see Command Flowcharts.
Contributor rule of thumb: if you are adding new command semantics, start in
crates/batchalign/src/commands/ and then jump to the owning module
(compare.rs, benchmark.rs, transcribe/, fa/, morphosyntax/, etc.).
runner/ owns lifecycle and queueing; dispatch/ should remain thin.
Workflow Families Overview
| Workflow Family | Commands | Parallelism | Key shape |
|---|---|---|---|
| Per-file transform | align, transcribe, transcribe_s, morphotag | Concurrent files (semaphore-bounded by num_workers) | One file in, one primary output out |
| Cross-file batch transform | utseg, translate, coref | Cross-file batching: pool utterances, group by language, dispatch languages concurrently, chunk large language groups across multiple workers | Two-level parallelism: cross-language × intra-language chunking (up to max_workers_per_key per language) |
| Reference projection | compare | Concurrent files, but with two primary CHAT inputs per file | Main+gold comparison bundle plus AST-first materializers |
| Composite workflow | benchmark | Concurrent files (semaphore-bounded by num_workers) | Transcribe first, then compare via typed command composition |
| Media analysis V2 | opensmile, avqi | Concurrent files (semaphore-bounded by num_workers) | Rust prepares audio, sends typed execute_v2 requests, Python returns raw analysis payloads |
All workflow families are server-side orchestrated or Rust-owned at the request boundary.
Parallelism Model
All per-file dispatch shapes (align, transcribe, benchmark,
morphotag, opensmile, avqi) process files concurrently using
supervised tokio::spawn tasks bounded by a
tokio::sync::Semaphore(num_workers). The number of workers is auto-tuned
based on available memory and CPU cores, or set explicitly with --workers N.
Each file opens its first durable attempt before setup work such as input
reads, media resolution, or conversion so early failures are visible in the
attempt history rather than only in terminal file state. Job-level media
prevalidation now uses an explicit file_setup attempt for the same reason.
Per-file dispatch code now routes the common processing/retry/completion
sequence through a shared FileRunTracker helper instead of open-coding those
store mutations in every pipeline. Supervised file tasks also report explicit
FileTaskOutcome values back to the runner so the runner does not need to
infer task success by rereading shared file state after the task exits. Media
preflight failures now use that same lifecycle boundary via an explicit setup
failure path instead of a one-off runner helper. The runner-owned lifecycle
labels are now also typed through FileStage rather than repeated as ad hoc
strings in each dispatch module. The same typed label vocabulary now flows
through the shared internal progress channel used by FA and transcribe
pipelines. The API now exposes that state in two parallel fields:
progress_stage for stable client logic and progress_label as the derived
operator-facing display string.
Batched text commands (utseg, translate, coref) take a
different approach: they pool all utterances from all files, group them by
per-item language, and dispatch with two levels of bounded parallelism. At
the outer level, language groups run concurrently but bounded by a semaphore
(max_total_workers / max_workers_per_key concurrent groups) to prevent
exceeding the global worker cap (morphosyntax/batch.rs). At the inner level,
each language group’s infer_batch call (morphosyntax/worker.rs) splits
large batches into chunks across up to max_workers_per_key workers of the
same language. When a language group finishes and its workers return to the
pool, the next queued group starts. compare
does not use this pooled-text shape any more; it is its own reference
projection workflow because it needs both a main transcript and a gold
companion per file. benchmark is a composite workflow that composes
transcribe and compare rather than inventing its own third orchestration style.
| Shape | File-level parallelism | Within-file parallelism |
|---|---|---|
| Per-file transform | Supervised tasks + Semaphore(N) | Single worker call or Rust composition per file |
| Reference projection | Supervised tasks + Semaphore(N) | Main+gold comparison bundle plus materializers per file |
| Composite workflow | Supervised tasks + Semaphore(N) | Transcribe then compare per file |
| Cross-file batch transform | N/A (single batch) | One GPU batch call covers all files |
| Media analysis V2 | Supervised tasks + Semaphore(N) | Single worker call per file |
Scenario 1: align: 3 files, 2 workers
Forced alignment is the most complex dispatch shape. Each file has its own audio, so files are processed sequentially. Within each file, utterances are grouped into time windows and batched to the worker.
sequenceDiagram
participant CLI
participant Server
participant JobStore
participant Pool as WorkerPool
participant W1 as Worker 1
participant W2 as Worker 2
participant Cache
CLI->>Server: POST /jobs (align, 3 files, lang=eng)
Server->>JobStore: Create job (Queued)
Server-->>CLI: 202 Accepted {job_id}
Note over Server: Runner picks up job
Server->>JobStore: Memory gate check
JobStore-->>Server: OK (idle worker bypass)
Server->>JobStore: Mark Running
Server->>Pool: pre_scale(num_workers=2)
Pool->>W1: spawn (command=align, lang=eng)
Pool->>W2: spawn (command=align, lang=eng)
loop For each file (concurrent, bounded by num_workers semaphore)
Server->>Server: resolve audio via media_mappings
Server->>Server: ensure_wav(): convert mp4→wav if needed (cached)
Server->>Server: parse_lenient()
alt Complete reusable %wor timing
Note over Server: Cheap rerun path:<br/>verify main↔%wor mapping,<br/>rehydrate main-tier word bullets,<br/>refresh utterance bullets / %wor
Server->>JobStore: Mark file Done
else Normal align path
Note over Server: UTR pre-pass (detect-and-skip)
Server->>Server: count_utterance_timing()
alt Untimed utterances exist AND utr_engine configured
Server->>Pool: checkout infer:asr worker
Server->>W1: execute_v2(task="asr", typed_input)
W1-->>Server: typed raw ASR result
Server->>Server: inject_utr_timing(): exact subsequence fast path, else global DP
Note over Server: Untimed utterances get<br/>utterance-level bullets from ASR
else All utterances timed OR no utr_engine
Note over Server: Skip UTR (use existing bullets<br/>or interpolation fallback)
end
Server->>Server: pre-validate (MainTierValid)
Server->>Server: group_utterances() → time windows
Server->>Cache: batch lookup (BLAKE3 keys: words + audio identity + window)
Cache-->>Server: hits[] + misses[]
alt Cache misses exist
Server->>Pool: checkout worker
Pool-->>Server: CheckedOutWorker (RAII)
Server->>W1: execute_v2(task="fa", prepared_audio + prepared_text)
Note over W1: Read prepared artifacts,<br/>run FA model
W1-->>Server: typed raw FA timings
Note over Server: Drop CheckedOutWorker → returns to pool
end
Server->>Server: parse_fa_response(): DP-align model output to transcript
Server->>Server: apply_fa_results(): inject timings + postprocess
Note over Server: Timing: chunk-relative → file-absolute ms<br/>Generate %wor tier<br/>Monotonicity check (E362)<br/>Same-speaker overlap (E704)
Server->>Cache: store new entries
Server->>Server: post-validate → serialize CHAT
Server->>JobStore: Mark file Done
end
end
Server->>JobStore: Mark job Completed
CLI->>Server: GET /jobs/{id}/results
Server-->>CLI: output files
Walkthrough
- CLI discovers
.chafiles in the input directory (sorted largest-first), submits them as a single job viaPOST /jobs. - Server creates the job in
Queuedstate and returns immediately. - The runner checks the memory gate, if an idle worker already exists
for
(align, eng), the memory check is bypassed entirely. - Pre-scaling spawns 2 worker processes to avoid sequential spawn overhead. Workers load the FA model (Whisper or Wave2Vec) at startup.
- Files are processed concurrently, bounded by
num_workersvia atokio::sync::Semaphore. For each file, the server resolves the audio file by walking the parent directory or using media mappings for matching.wav/.mp3/.mp4files. If the resolved file is MP4 (or another container format),ensure_wavconverts it to WAV via ffmpeg and caches the result at~/.batchalign3/media_cache/(see Media Conversion). 5b. Cheap rerun path: After parsing, the server first checks whether the file already has complete reusable%wortiming. If main↔%woralignment is clean and every mapped%worword is timed, the server copies that timing back to main-tier words, refreshes utterance bullets, resolves adjacent-utterance end overlap (--end-overlap-policy), and only THEN optionally regenerates%worfrom that resolved state; FA is skipped for the file entirely. 5c. UTR pre-pass (detect-and-skip): If the file is not already fully reusable, the server callscount_utterance_timing(). If untimed utterances exist and a UTR engine is configured (--utr, the default), it runs a Rust-owned UTR backend on the full audio, theninject_utr_timing()first tries a cheap exact-subsequence match and falls back to one global Hirschberg DP when the transcript/ASR match is missing or ambiguous. For--utr-engine rev, the server uses the shared Rust Rev.AI client directly. For worker-backed engines such as Whisper, the server still uses the worker ASR task. If all utterances are already timed, UTR is skipped entirely. If no UTR engine is configured (--no-utr), untimed utterances fall back to proportional interpolation. The updated CHAT text (with recovered bullets) is then used for FA grouping. 5c. The server groups utterances into time windows (max 20s for Whisper, 15s for Wave2Vec). - Cache lookup uses BLAKE3 hashes of (words + audio identity + time window
- gap-healing policy + engine). Cache hits skip worker IPC entirely.
- Cache misses are sent to a checked-out worker via typed
execute_v2requests. TheCheckedOutWorkerRAII guard returns the worker to the pool on drop. - The Rust server DP-aligns model timestamps to transcript words
(Hirschberg algorithm), converts chunk-relative times to file-absolute
milliseconds, generates
%wortiers, and runs monotonicity/overlap checks. - New results are cached for future reuse.
- The CLI polls for results and writes output files.
Scenario 2: morphotag: 2 multilingual files, per-file dispatch
Morphotag processes files concurrently, bounded by num_workers.
Within each file, utterances are analyzed independently (with optional
cache hits) and results are injected back into the AST.
sequenceDiagram
participant CLI
participant Server
participant Pool as WorkerPool
participant W as Worker
participant Cache
CLI->>Server: POST /jobs (morphotag, 2 files)
Server-->>CLI: 202 Accepted {job_id}
Note over Server: Runner: mark Running, select dispatch_morphotag_job
loop For each file (concurrent, bounded by num_workers semaphore)
Server->>Server: parse_lenient()
alt @Options: CA in header
Server->>Server: serialize parsed file as-is (no %mor/%gra added)
Server->>JobStore: Mark file Done
else
Server->>Server: clear existing %mor/%gra
Server->>Server: collect_payloads()
Server->>Cache: batch lookup all utterances
Cache-->>Server: hits[] + misses[]
Server->>Server: inject cache hits immediately
alt Misses exist
Server->>Pool: checkout worker
Pool-->>Server: CheckedOutWorker
Server->>W: execute_v2(task="morphosyntax", misses batch)
W-->>Server: UD results
Note over Server: Worker returned to pool
end
Server->>Server: inject_results() → insert %mor/%gra tiers
Server->>Server: validate alignment → serialize CHAT
Server->>JobStore: Mark file Done
end
end
Server->>JobStore: Mark job Completed
CLI->>Server: GET /jobs/{id}/results
Server-->>CLI: output files
Walkthrough
- CLI submits CHAT files for morphosyntactic enrichment.
- The runner selects the per-file dispatch path (
dispatch_morphotag_job). - Files are processed concurrently, bounded by
num_workers. This prevents the BA2 over-parallelism crash mode while maximizing throughput on multi-core hosts. - For each file, the server parses the transcript. If the parsed
header declares
@Options: CA, the file is serialized back as-is, no%mor/%gratiers are added or removed, and no provenance comment is injected (mirroringalign’s@Options: NoAlignpass-through). Otherwise the server clears any stale morphology and collects payloads (word lists + language metadata). - Cache lookup checks all utterances in the file at once. BLAKE3 keys include (words + language + terminator + special forms + engine version).
- Cache misses are sent to a checked-out worker in a single batch. The worker runs the Stanza NLP pipeline for the appropriate language(s).
- Results (both from cache and worker) are injected back into the file’s
AST, inserting new
%morand%gratiers. - The file is validated (ensuring morphology matches the main tier) and serialized back to CHAT.
- Each file’s result is written to disk immediately as it finishes, allowing for incremental progress visibility on large corpora.
- This per-file shape replaces the previous complex cross-file windowing logic, providing better reliability and simpler progress reporting.
Scenario 2b: compare: 1 main file + 1 gold companion
Compare is the reference-projection shape. It pairs each primary transcript with
a FILE.gold.cha companion, morphotags only the main side, and materializes one
or more outputs from a typed comparison bundle.
sequenceDiagram
participant CLI
participant Server
participant Pool as WorkerPool
participant W as Worker
participant Cmp as compare()
CLI->>Server: POST /jobs (compare, 1 main file)
Server-->>CLI: 202 Accepted {job_id}
Server->>Server: Resolve FILE.gold.cha companion
alt Missing gold companion
Server->>Server: Mark file Error
else Gold companion present
Server->>Pool: checkout worker
Pool-->>Server: CheckedOutWorker
Server->>W: execute_v2(task="morphosyntax", main transcript only)
W-->>Server: typed morphosyntax result
Note over Server: Worker returned to pool
Server->>Server: PostValidated::into_judged_document() -> AST_main
Server->>Server: parse_lenient(raw gold) -> AST_gold
Server->>Cmp: compare(AST_main, AST_gold)
Note over Cmp: conform -> per-gold window search -> local DP<br/>main view + gold view + structural word matches + metrics
Cmp-->>Server: ComparisonBundle
Server->>Server: project_gold_structurally() on gold AST
Note over Server: Exact matches copy %mor / %gra / %wor;<br/>unsafe partial projection stays conservative
Server->>Server: build typed %xsrep / %xsmor models\nlower once to gold AST tiers
opt Internal benchmark/main path
Server->>Server: MainAnnotatedCompareMaterializer<br/>reuse typed tier models on main AST
end
Server->>Server: CompareMetricsCsvTable -> csv crate -> .compare.csv
Server->>Server: validate -> serialize
Server-->>CLI: output .cha + .compare.csv
end
Walkthrough
- The CLI submits only primary
.chainputs; the gold companion is resolved by the compare planner/dispatch layer. - BA3 runs morphosyntax on the main transcript only. The gold transcript stays raw during artifact construction so deletions retain reference-side shape instead of picking up invented tags.
compare()performs BA2-style per-gold-utterance window selection and local DP, then returns aComparisonBundlecontaining main-anchored tokens, gold-anchored tokens, structural gold↔main word matches, and aggregate metrics.- The released materializer projects onto the gold/reference AST and injects
%xsrep/%xsmorthere. The main-annotated materializer still exists for internal benchmark-style flows, but it is no longer the compare command surface. %xsrep/%xsmorare emitted from typed compare-tier models, not from raw string hacking, and.compare.csvis rendered from the same bundle through a structured table model. That keeps transcript annotations and metric output in lockstep.
Scenario 3: transcribe: 1 file, audio to CHAT
Transcription creates CHAT from scratch rather than modifying existing files. It has the longest pipeline: ASR → post-processing → CHAT assembly → optional follow-up commands.
Speaker label handling: convert_asr_response() always uses speaker
labels from the ASR engine when present (matching BA2’s process_generation()
which unconditionally reads utterance["speaker"]). The --diarization flag
only controls whether a dedicated Pyannote/NeMo stage runs, it does not
suppress ASR-provided labels. This means batchalign3 transcribe (without
--diarization) still produces multi-speaker output when Rev.AI returns
speaker-labeled monologues. When --diarization enabled is explicitly
requested, BA3 runs the dedicated speaker stage even on top of Rev-labeled
output and applies its evidence before utterance segmentation.
Rev.AI skip_postprocessing: For English only,
skip_postprocessing=true is sent to Rev.AI (matching BA2), so BA3’s own
pre-CHAT utterance model handles segmentation from raw output. In --lang auto
mode, the Rust server first runs Rev.AI language ID. If that resolves to a
supported language such as English before submission, the request path becomes
the same as explicit --lang eng. If language ID fails or returns an unmapped
code, BA3 keeps a true Rev auto request instead; downstream processing may
still later resolve the output to English, but the provider request was not the
same as --lang eng.
sequenceDiagram
participant CLI
participant Server
participant Pool as WorkerPool
participant W as Worker
CLI->>Server: POST /jobs (transcribe, 1 audio file)
Server-->>CLI: 202 Accepted {job_id}
Note over Server: Runner: mark Running
Server->>Server: resolve audio path
Server->>Server: ensure_wav(): convert mp4→wav if needed (cached)
alt Rev.AI engine selected
Server->>Server: derive evidence key and acquire per-key lease
alt valid durable evidence hit
Server->>Server: replay provider-shaped transcript evidence
else typed cache miss
Server->>Server: authorize one Rev.AI request
Server->>Server: optional language ID, submit, poll, validate
Server->>Server: durably commit evidence before continuing
end
else worker-backed ASR engine
Server->>Pool: checkout worker
Pool-->>Server: CheckedOutWorker
Server->>W: execute_v2(task="asr", typed_input)
Note over W: Run local or provider-backed ASR model
W-->>Server: typed raw ASR result
Note over Server: Worker returned to pool
end
Note over Server: convert_asr_response(): ALWAYS groups<br/>tokens by speaker label (no flag gating)
opt --diarization enabled
Server->>Pool: checkout worker
Pool-->>Server: CheckedOutWorker
Server->>W: execute_v2(task="speaker", prepared_audio)
Note over W: Run diarization model<br/>(pyannoteAI default, local alternatives explicit)
W-->>Server: typed raw speaker result (speaker_result)
Note over Server: Worker returned to pool
opt --debug-dir configured
Server->>Server: Write exact same-job canonical turns<br/>with typed backend provenance
end
end
Note over Server: Rust ASR normalization over typed monologues
Server->>Server: 1. Compound merging (adjacent subword tokens)
Server->>Server: 2. Timed word extraction (seconds to ms)
Server->>Server: 2d. Cantonese normalization, once per monologue (lang=yue only)
Server->>Server: 3. Multi-word splitting (timestamp interpolation)
Server->>Server: 4. Number expansion (digits to words)
Server->>Server: 5. Long-turn splitting (chunk at >300 words)
Server->>Server: 5b. Long-pause fallback splitting
opt dedicated speaker segments present
Server->>Server: project_speakers_onto_chunks():<br/>assign timed words by summed overlap,<br/>split at speaker changes
end
opt language has BA2 utterance model (eng/zho/yue)
Server->>W: execute_v2(task="utseg", prepared word batch)
Note over W: BA2-style model returns typed boundary assignments
W-->>Server: typed assignments
Note over Server: Apply assignments before CHAT build
end
Server->>Server: 6. Retokenization (punctuation fallback / cleanup)
Server->>Server: 7. Disfluency cleanup + retrace detection
Server->>Server: build_chat(): ChatFile AST
Note over Server: Generate headers: @Languages, @Participants,<br/>@ID (PAR/INV/CHI/MOT...), @Media<br/>Build utterances with %wor tiers<br/>Speaker codes from ASR labels used directly
opt with_utseg=true (default)
Server->>Server: process_utseg_with_evidence(): re-segment utterance boundaries
end
opt with_morphosyntax=true (default: false)
Server->>Server: process_morphosyntax(): add %mor/%gra tiers
end
Server->>Server: validate, serialize, .cha output
Server-->>CLI: output .cha file
Walkthrough
- Rev.AI evidence resolution: For Rev.AI-backed transcription, the server
hashes the complete provider-visible inference media plus all
inference-affecting request settings, then acquires a per-key singleflight
lease. A valid hit replays provider-shaped transcript evidence without a
network request. A typed miss is the only state allowed to authorize
language ID and paid submission; validated evidence must be durably committed
before the pipeline continues. For English,
skip_postprocessing=trueis sent so BA3’s own pre-CHAT utterance model handles segmentation. In--lang automode there are two real branches:- Language ID succeeds and maps cleanly: BA3 collapses to a resolved
language before submission. If it resolves to
eng, the Rev request path is the same as explicit--lang eng. - Language ID fails or returns an unmapped code: BA3 submits a true Rev
auto request. Downstream code may still later resolve the transcript to
English for segmentation and CHAT headers, but provider-side options such
as
speakers_countandskip_postprocessingwere not the explicit-English ones. The former parallel pre-submission path is currently disabled because it submitted work before cache lookup and bypassed typed miss authorization. This protects correctness and billing at the cost of lower cold-cache batch throughput until a cache-aware typed parallel planner is implemented.
- Language ID succeeds and maps cleanly: BA3 collapses to a resolved
language before submission. If it resolves to
- The worker or Rust-owned Rev path returns a typed ASR response. BA3
preserves both a flattened token view and provider-shaped monologues so
later stages can keep punctuation and speaker boundaries instead of trying to
re-infer them from plain text. The inference boundary has zero CHAT
awareness.
2b. Speaker label handling:
convert_asr_response()always groups tokens by their speaker labels when present. There is nouse_speaker_labelsparameter, this matches BA2’s unconditional speaker reading. The--diarizationflag only gates the dedicated speaker stage (step 2c), not the use of ASR-provided labels. 2c. Dedicated diarization (optional): If--diarization enabledis set, the server dispatchesexecute_v2(task="speaker"). The typed backend is pyannoteAI Precision-2 by default, with local Pyannote and NeMo alternatives. ASR-provided labels are read first, but the explicit dedicated result is authoritative. 2d. Same-job turn retention (optional): If--debug-diris configured, BA3 writes the exact dedicated segments before they can be discarded. The typed label-coordinate map is shared with CHAT projection, provenance is derived fromSpeakerBackendV2, and an enabled write failure fails the file. This is an interim debug/research artifact, not the final durable evidence sidecar. - All post-processing happens in Rust (
batchalign), not Python. The normalization stages inprepare_asr_chunks()are:- Compound merging, joins adjacent subword tokens
- Timed word extraction, seconds to milliseconds, filter pauses
2d. Cantonese normalization (lang=yue only), simplified to traditional via
ferrous-opencc+ domain replacements (pure Rust), run once over the whole monologue before anything splits it - Multi-word splitting, split space-separated tokens, interpolate timestamps
- Number expansion, digits to spelled-out words (language-aware)
- Long-turn splitting, chunk monologues at >300 words 5b. Long-pause fallback splitting, split strongly separated runs when provider punctuation is missing
- Speaker projection and pre-CHAT utterance segmentation: Dedicated
segments are first projected onto timed ASR words by greatest summed overlap,
and prepared chunks are split at label changes. For supported languages (
eng,zho,yue), BA3 now calls the BA2 utterance model at this seam through the V2utsegworker task. Python returns typed word-group assignments, Rust applies them to the prepared ASR chunks, and only then does punctuation retokenization run as fallback/cleanup. This seam runs on the effective resolved language seen by post-processing, so both Revauto -> resolved engand Revauto -> true auto request -> later resolved engcan eventually reach the English utterance model. Only the first branch is provider-request equivalent to explicit--lang eng. - CHAT assembly (
build_chat) creates a completeChatFileAST with proper headers (participant codes derived from speaker indices: PAR, INV, CHI, MOT, etc.) and utterances with%wortiming tiers. Speaker-safe chunks already carry the labels from the dedicated projection when it ran. - Optional follow-up commands (utseg defaults on, morphotag defaults off) are chained automatically, reusing the same worker pool.
Scenario 4: Server Startup & Lazy Capability Detection
At startup the server recovers persisted state and begins accepting jobs. Capability detection is lazy: there is no probe worker at startup. Instead, capabilities are detected from the first real worker spawn for each profile.
sequenceDiagram
participant Server
participant Pool as WorkerPool
participant W as First Worker
participant DB as SQLite
Server->>DB: Mark queued/running jobs interrupted
Server->>DB: Prune expired entries
Server->>DB: Load jobs and reconcile runtime state
Note over Server,DB: Requeue resumable work; promote all-terminal jobs to final state; persist canonical status and cleared leases
Server->>DB: Load persisted jobs, init utterance cache
Note over Server: Server ready, accepting requests<br/>(capabilities not yet known)
Note over Server: First job arrives (e.g., morphotag)
Server->>Pool: checkout worker (morphotag, eng)
Pool->>W: python -m batchalign.worker --task morphosyntax --lang eng
W-->>Pool: {"ready": true, "pid": N}
Server->>W: capabilities()
Note over W: Import-probe each InferTask:<br/>stanza → Morphosyntax, Utseg, Coref ✓<br/>googletrans → Translate ✓<br/>torch+torchaudio → FA ✓<br/>whisper or Rev key → ASR ✓<br/>parselmouth+torchaudio → AVQI ✓<br/>(no opensmile) → OpenSMILE ✗
W-->>Server: CapabilitiesResponse {infer_tasks, engine_versions, commands=[]}
Server->>Pool: record_capabilities(): admit the report once, store it per worker key
Note over Server: One availability rule (command_supported, primary infer task only):<br/>morphotag needs Morphosyntax ✓<br/>utseg needs Utseg ✓<br/>translate needs Translate ✓<br/>coref needs Coref ✓<br/>align needs FA ✓ (FA engine name read at dispatch, after FA loads)<br/>opensmile needs OpenSMILE ✗ → excluded
Server->>Server: Build final capabilities list
Note over Server: /health now advertises:<br/>commands: [morphotag, utseg, translate, coref, align, transcribe, ...]<br/>infer_tasks: [Morphosyntax, Utseg, Translate, Coref, FA, ASR, ...]
Note over Server: Worker stays in pool for actual job work
Walkthrough
- DB recovery: Any jobs left in
QueuedorRunningstate from a previous crash are first markedInterrupted, and expired entries are pruned. - Runtime reconciliation:
JobStore::load_from_db()rebuilds each job and then uses theJobrecovery transition to choose a canonical state: resumable files are re-queued, while all-terminal jobs are promoted toCompletedorFailed. The reconciled status and cleared lease metadata are written back to SQLite so memory and persistence agree. - The server begins accepting requests immediately. Capabilities are not yet known, they are populated lazily.
- When the first job arrives, the server spawns a real worker for the
requested command. During this first worker’s startup, the server calls
capabilities()which import-probes eachInferTask: for each task, the worker tries to import the required Python packages (e.g.,stanzafor Morphosyntax,torch+torchaudiofor FA). If imports succeed, the task is reported as available. Itsengine_versionsentry isnull, except forced alignment’s, which names the FA engine once an FA model has loaded. - The worker stays in the pool for actual job work, it is not shut down after capability detection.
WorkerPool::record_capabilities()admits the report once (WorkerEngineReports::admit: every advertised task needs an entry, only forced alignment’s may be a name, and no entry may name an unadvertised task) and stores the outcome per worker key: the admitted report, or the refusal, which/healthreturns inworker_capability_admissions.WorkerCapabilitySnapshot::detected()then derives the released command surface withcapability::command_supported: a command is advertised when the worker supports its primary infer task, whether or not the model behind it has loaded. Engine names are not consulted. At dispatch the same step runs again on a post-load report, and the forced-alignment dispatch arm (onlyalign, whose cache rows are namespaced by the FA engine) additionally reads the FA engine from that report withFaCacheNamespace::from_loaded, refusing a job whose worker still names no FA engine after the load.- The
/healthendpoint advertises the validated capability set once it is known. The CLI checks this before submitting jobs, if a required command is missing, it errors immediately rather than queueing a job that will fail.
Cross-Cutting Concerns
CHAT Ownership Boundary
In all scenarios above, the Rust server owns the full CHAT lifecycle: parsing, AST manipulation, validation, caching, and serialization. Python workers receive extracted data (word lists, audio paths) and return raw ML output. No CHAT text crosses the IPC boundary.
Cache Behavior
Cache checks happen before any worker IPC. A fully-cached file (e.g., re-running morphotag on unchanged input) completes without touching a Python worker at all. Cache keys include the engine version, so model upgrades automatically invalidate stale entries.
Error Boundaries
- Worker crash: Pool detects exit, decrements worker count, spawns replacement on next checkout.
- Retryable errors (FA timeout, Rev.AI throttle): Exponential backoff
with configurable
max_attempts. - Terminal errors (parse failure, validation rejection): File marked
Error, job continues processing remaining files. - Memory pressure: Job re-queued with backoff rather than OOM-killed.
Worker Pool Mechanics
Workers are keyed by (CommandName, LanguageCode3). The pool uses
Mutex<VecDeque> for the idle queue and tokio::sync::Semaphore for
availability. CheckedOutWorker is an RAII guard that returns the worker
to the pool on drop, no manual checkin needed.
This page last changed: 2026-09-16 (commit 197c81e6). The whole book last changed: 2026-09-16 (commit 34d249d8).