"""Fail-closed quality runtime helpers for the BlueMagpie-TTS Space. The module deliberately has no import-time model downloads. Both Whisper and ECAPA are supplied through small injectable boundaries so the online Space can load the real models lazily while unit tests remain deterministic and offline. """ from __future__ import annotations import hashlib import json import math import operator import threading from dataclasses import dataclass, replace from math import gcd from typing import Any, Callable, Sequence import numpy as np import torch import torch.nn.functional as torch_functional from production import AsrComparison, CandidateSequenceSelection, compare_asr_text WHISPER_MODEL_ID = "openai/whisper-large-v3-turbo" WHISPER_REVISION = "41f01f3fe87f28c78e2fbf8b568835947dd65ed9" VERIFICATION_WHISPER_MODEL_ID = "openai/whisper-large-v3" VERIFICATION_WHISPER_REVISION = "06f233fe06e710322aca913c1bc4249a0d71fce1" WHISPER_ATTENTION_IMPLEMENTATION = "eager" WHISPER_RETURN_ATTENTION_MASK = True WHISPER_SAMPLE_RATE = 16_000 WHISPER_MAX_SEGMENT_SECONDS = 28.0 ADAPTIVE_CASCADE_STAGE_LIMITS = (1, 5, 10, 15, 20) REQUEST_SEED_LIMIT = 2_147_483_648 ACTIVE_VOICE_TOP_DB = 35.0 ACTIVE_VOICE_FRAME_MS = 25.0 ACTIVE_VOICE_HOP_MS = 10.0 ACTIVE_VOICE_MIN_RMS = 1.0e-4 RELEASE_SPEAKER_TRIGGER_SECONDS = 1.48 SEQUENCE_FALLBACK_MAX_LOCAL_BOUNDARY_SPEAKER_DROP = 0.15 SEQUENCE_FALLBACK_SPEAKER_WEIGHT = 0.05 SEQUENCE_FALLBACK_BOUNDARY_WEIGHT = 0.10 CASCADE_EVIDENCE_SCHEMA_VERSION = 1 CASCADE_EVIDENCE_LOG_PREFIX = "[BlueMagpie] cascade evidence " CASCADE_EVIDENCE_MAX_ATTEMPTS = ADAPTIVE_CASCADE_STAGE_LIMITS[-1] CASCADE_EVIDENCE_MAX_LOCAL_RESULTS = 20 CASCADE_EVIDENCE_MAX_REASONS = 8 CASCADE_EVIDENCE_MAX_SEQUENCE_PATHS = 3 _CASCADE_EVIDENCE_OUTCOMES = frozenset( {"returned", "no_qualified_candidate", "final_output_rejected"} ) _CASCADE_EVIDENCE_SELECTION_MODES = frozenset( {"whole_trajectory", "sequence_dp"} ) _CASCADE_EVIDENCE_REJECTION_CODES = frozenset( { "boundary_speaker_drop", "empty_trajectory", "invalid_audio_duration", "invalid_chunk_artifacts", "invalid_gate_config", "invalid_joined_verification", "malformed_observation", "missing_pace_evidence", "missing_speaker_evidence", "nonfinite_score", "nonfinite_trajectory_score", "pace_too_fast", "semantic_gate", "speaker_similarity", "truncated", } ) @dataclass(frozen=True) class GenerationPolicy: """Candidate-specific endpoint duration estimate used by the Space.""" name: str cjk_cps: float ascii_cps: float hard_stop_margin_steps: int BASE_GENERATION_POLICY = GenerationPolicy( name="base", cjk_cps=5.2, ascii_cps=4.6, hard_stop_margin_steps=1, ) SAFE_DURATION_GENERATION_POLICY = GenerationPolicy( name="safe_duration", cjk_cps=4.6, ascii_cps=4.0, hard_stop_margin_steps=1, ) def generation_policy_for_candidate_offset(candidate_offset: int) -> GenerationPolicy: """Map candidate zero to base and every retry to the safe estimate. The policy changes only the native-duration endpoint estimate. It is not a minimum-length policy and therefore never holds the generation loop open to enforce playback pace. """ if isinstance(candidate_offset, (bool, np.bool_)): raise ValueError("candidate_offset must be a non-negative integer") try: offset = operator.index(candidate_offset) except (TypeError, ValueError, OverflowError) as error: raise ValueError("candidate_offset must be a non-negative integer") from error if offset < 0: raise ValueError("candidate_offset must be a non-negative integer") return BASE_GENERATION_POLICY if offset == 0 else SAFE_DURATION_GENERATION_POLICY def resolve_request_seed( request_seed: int | None, random_seed_factory: Callable[[int], int], ) -> int: """Return a validated root seed, drawing randomness only for ``None``.""" candidate = ( random_seed_factory(REQUEST_SEED_LIMIT) if request_seed is None else request_seed ) if isinstance(candidate, (bool, np.bool_)): raise ValueError(f"request_seed must be an integer in [0, {REQUEST_SEED_LIMIT})") try: # ``operator.index`` semantics reject floats and numeric strings while # accepting Python and NumPy integer scalars. seed = operator.index(candidate) except (AttributeError, TypeError, ValueError, OverflowError) as error: raise ValueError( f"request_seed must be an integer in [0, {REQUEST_SEED_LIMIT})" ) from error seed = int(seed) if not 0 <= seed < REQUEST_SEED_LIMIT: raise ValueError(f"request_seed must be an integer in [0, {REQUEST_SEED_LIMIT})") return seed def _finite_float( value: Any, *, minimum: float | None = None, maximum: float | None = None, ) -> float | None: if isinstance(value, (bool, np.bool_)): return None try: result = float(value) except (TypeError, ValueError, OverflowError): return None if not math.isfinite(result): return None if minimum is not None and result < minimum: return None if maximum is not None and result > maximum: return None return result def release_speaker_measurement_required(active_duration_seconds: float) -> bool: """Trigger online ECAPA slightly before the external 1.50 s hard gate. The 20 ms margin covers a bounded RMS-frame shift caused by WAV serialization while leaving the published external eligibility threshold unchanged. """ duration = _finite_float(active_duration_seconds, minimum=0.0) return duration is not None and duration >= RELEASE_SPEAKER_TRIGGER_SECONDS def _mono_audio(audio: np.ndarray | Sequence[float]) -> np.ndarray: """Return contiguous mono float32 audio, rejecting ambiguous/bad inputs.""" waveform = np.asarray(audio) if waveform.ndim == 1: pass elif waveform.ndim == 2: first, second = waveform.shape if first <= 8 and second > first: waveform = waveform.mean(axis=0) elif second <= 8 and first > second: waveform = waveform.mean(axis=1) else: raise ValueError("2-D audio must have an identifiable channel axis (at most 8 channels)") else: raise ValueError("audio must be a one- or two-dimensional array") waveform = np.asarray(waveform, dtype=np.float32).reshape(-1) if waveform.size == 0: raise ValueError("audio is empty") if not np.isfinite(waveform).all(): raise ValueError("audio contains non-finite samples") return np.ascontiguousarray(waveform) def _resample_audio(audio: np.ndarray, sample_rate: int, target_rate: int) -> np.ndarray: try: source_rate = int(sample_rate) destination_rate = int(target_rate) except (TypeError, ValueError, OverflowError) as error: raise ValueError("sample rates must be positive integers") from error if source_rate <= 0 or destination_rate <= 0: raise ValueError("sample rates must be positive integers") if source_rate == destination_rate: return np.ascontiguousarray(audio, dtype=np.float32) # SpeechBrain already depends on SciPy. ``resample_poly`` avoids an # undeclared optional ``librosa`` resampler dependency in the Space image. from scipy.signal import resample_poly common_divisor = gcd(source_rate, destination_rate) output = resample_poly( np.asarray(audio, dtype=np.float32), destination_rate // common_divisor, source_rate // common_divisor, ) output = np.asarray(output, dtype=np.float32).reshape(-1) if output.size == 0 or not np.isfinite(output).all(): raise ValueError("resampling produced invalid audio") return np.ascontiguousarray(output) def _trim_active_speech( audio: np.ndarray, *, top_db: float = 35.0, frame_length: int = 512, hop_length: int = 128, ) -> np.ndarray: threshold_db = _finite_float(top_db, minimum=0.0) if threshold_db is None: raise ValueError("top_db must be finite and non-negative") if float(np.max(np.abs(audio))) <= 1.0e-7: raise ValueError("audio contains no active speech") import librosa intervals = librosa.effects.split( audio, top_db=threshold_db, frame_length=max(32, int(frame_length)), hop_length=max(1, int(hop_length)), ) if intervals.size == 0: raise ValueError("audio contains no active speech") start = int(intervals[0, 0]) stop = int(intervals[-1, 1]) active = np.asarray(audio[start:stop], dtype=np.float32) if active.size == 0 or float(np.max(np.abs(active))) <= 1.0e-7: raise ValueError("audio contains no active speech") return np.ascontiguousarray(active) def trim_release_speaker_activity( audio: np.ndarray | Sequence[float], sample_rate: int, *, top_db: float = 35.0, margin_seconds: float = 0.05, ) -> np.ndarray: """Trim speaker audio with the independent release-gate contract. This deliberately mirrors ``evaluate_space_profile.trim_speaker_activity``: 25 ms RMS frames, 10 ms hop, peak-minus-35 dB with a 1e-4 floor, and a 50 ms outer margin. It remains separate from the interval-union duration used for pace and speaker eligibility. """ signal = _mono_audio(audio) try: rate = operator.index(sample_rate) except (TypeError, ValueError, OverflowError) as error: raise ValueError("sample rate must be positive") from error relative_db = _finite_float(top_db, minimum=0.0) margin_duration = _finite_float(margin_seconds, minimum=0.0) if rate <= 0 or relative_db is None or margin_duration is None: raise ValueError("release speaker trim settings are invalid") frame = max(160, int(round(0.025 * rate))) hop = max(80, int(round(0.010 * rate))) if signal.size < frame: active = signal else: starts = np.arange(0, signal.size - frame + 1, hop, dtype=np.int64) rms = np.asarray( [ float( np.sqrt( np.mean( np.square(signal[start : start + frame], dtype=np.float64) ) ) ) for start in starts ], dtype=np.float64, ) peak = float(rms.max(initial=0.0)) if peak <= 0.0: active = signal else: threshold = max(1.0e-4, peak * (10.0 ** (-relative_db / 20.0))) active_frames = np.flatnonzero(rms >= threshold) if active_frames.size == 0: active = signal else: margin = max(0, int(round(margin_duration * rate))) begin = max(0, int(starts[int(active_frames[0])]) - margin) end = min( signal.size, int(starts[int(active_frames[-1])]) + frame + margin, ) active = signal[begin:end] if end > begin else signal active = np.ascontiguousarray(active, dtype=np.float32) if active.size == 0 or float(np.max(np.abs(active))) <= 1.0e-7: raise ValueError("audio contains no active speech") return active def active_voiced_intervals( audio: np.ndarray | Sequence[float], sample_rate: int, *, top_db: float = ACTIVE_VOICE_TOP_DB, frame_ms: float = ACTIVE_VOICE_FRAME_MS, hop_ms: float = ACTIVE_VOICE_HOP_MS, min_rms: float = ACTIVE_VOICE_MIN_RMS, ) -> tuple[tuple[int, int], ...]: """Return the deterministic union of active RMS-frame intervals. This is the same 25 ms / 10 ms, peak-minus-35 dB, 1e-4 floor contract used by the independent hosted evaluator. Unlike first-to-last trimming, the interval union excludes internal punctuation and joining pauses from both online pace evidence and the speaker-gate duration threshold. """ signal = _mono_audio(audio) if isinstance(sample_rate, (bool, np.bool_)): raise ValueError("sample rate must be positive") try: source_rate = operator.index(sample_rate) except (TypeError, ValueError, OverflowError) as error: raise ValueError("sample rate must be positive") from error if source_rate <= 0: raise ValueError("sample rate must be positive") frame_duration = _finite_float(frame_ms, minimum=0.0) hop_duration = _finite_float(hop_ms, minimum=0.0) rms_floor = _finite_float(min_rms, minimum=0.0) relative_db = _finite_float(top_db, minimum=0.0) if ( frame_duration is None or frame_duration <= 0.0 or hop_duration is None or hop_duration <= 0.0 or rms_floor is None or rms_floor <= 0.0 or relative_db is None ): raise ValueError("active-voice detector settings are invalid") frame = max(1, int(round(frame_duration * source_rate / 1000.0))) hop = max(1, int(round(hop_duration * source_rate / 1000.0))) if signal.size <= frame: starts = np.asarray([0], dtype=np.int64) else: starts = np.arange(0, signal.size - frame + 1, hop, dtype=np.int64) final_start = signal.size - frame if int(starts[-1]) != final_start: starts = np.append(starts, final_start) rms = np.asarray( [ float( np.sqrt( np.mean( np.square( signal[int(start) : int(start) + frame], dtype=np.float64, ) ) ) ) for start in starts ], dtype=np.float64, ) peak = float(rms.max(initial=0.0)) threshold = max(rms_floor, peak * 10.0 ** (-relative_db / 20.0)) active_starts = starts[rms >= threshold] intervals: list[list[int]] = [] for raw_start in active_starts: start = int(raw_start) end = min(signal.size, start + frame) if intervals and start <= intervals[-1][1]: intervals[-1][1] = max(intervals[-1][1], end) else: intervals.append([start, end]) return tuple((start, end) for start, end in intervals) def active_voiced_duration_seconds( audio: np.ndarray | Sequence[float], sample_rate: int, **detector_kwargs: Any, ) -> float: """Measure active interval-union duration under the hosted gate contract.""" intervals = active_voiced_intervals(audio, sample_rate, **detector_kwargs) active_samples = sum(end - start for start, end in intervals) return active_samples / float(operator.index(sample_rate)) @dataclass(frozen=True) class WhisperRuntime: """Loaded processor/model pair for deterministic Whisper transcription.""" processor: Any model: Any device: torch.device dtype: torch.dtype def _load_whisper_runtime( model_id: str, revision: str, *, device: str | torch.device | None = None, processor_factory: Any | None = None, model_factory: Any | None = None, ) -> WhisperRuntime: """Load one exact ASR revision through the shared deterministic contract. Factory injection exists for offline tests. The default imports ``transformers`` only when this function is first called. """ if processor_factory is None or model_factory is None: from transformers import AutoModelForSpeechSeq2Seq, AutoProcessor processor_factory = processor_factory or AutoProcessor model_factory = model_factory or AutoModelForSpeechSeq2Seq selected_device = torch.device( device if device is not None else ("cuda" if torch.cuda.is_available() else "cpu") ) dtype = torch.float16 if selected_device.type == "cuda" else torch.float32 processor = processor_factory.from_pretrained( model_id, revision=revision, ) model = model_factory.from_pretrained( model_id, revision=revision, attn_implementation=WHISPER_ATTENTION_IMPLEMENTATION, torch_dtype=dtype, low_cpu_mem_usage=True, use_safetensors=True, ) model = model.to(selected_device) model.eval() return WhisperRuntime( processor=processor, model=model, device=selected_device, dtype=dtype, ) def load_pinned_whisper_runtime( *, device: str | torch.device | None = None, processor_factory: Any | None = None, model_factory: Any | None = None, ) -> WhisperRuntime: """Load the pinned turbo ASR used for local candidate screening.""" return _load_whisper_runtime( WHISPER_MODEL_ID, WHISPER_REVISION, device=device, processor_factory=processor_factory, model_factory=model_factory, ) def load_pinned_verification_whisper_runtime( *, device: str | torch.device | None = None, processor_factory: Any | None = None, model_factory: Any | None = None, ) -> WhisperRuntime: """Load the independent full large-v3 ASR used only for whole outputs.""" return _load_whisper_runtime( VERIFICATION_WHISPER_MODEL_ID, VERIFICATION_WHISPER_REVISION, device=device, processor_factory=processor_factory, model_factory=model_factory, ) class LazyWhisperASR: """Thread-safe one-shot lazy loader with an injectable runtime factory.""" def __init__(self, runtime_loader: Callable[[], WhisperRuntime] | None = None) -> None: self._runtime_loader = runtime_loader or load_pinned_whisper_runtime self._runtime: WhisperRuntime | None = None self._lock = threading.Lock() def get_runtime(self) -> WhisperRuntime: runtime = self._runtime if runtime is not None: return runtime with self._lock: if self._runtime is None: self._runtime = self._runtime_loader() return self._runtime _DEFAULT_WHISPER = LazyWhisperASR() _DEFAULT_VERIFICATION_WHISPER = LazyWhisperASR( load_pinned_verification_whisper_runtime ) def _split_whisper_audio( waveform: np.ndarray, *, sample_rate: int = WHISPER_SAMPLE_RATE, max_segment_seconds: float = WHISPER_MAX_SEGMENT_SECONDS, boundary_search_seconds: float = 1.5, ) -> tuple[np.ndarray, ...]: """Split long audio at real pauses while staying below Whisper's cap. Equal-duration cuts can land inside a word even when the production waveform already contains punctuation pauses. Search the complete safe 6-to-28-second window for a sustained low-energy run and choose the latest one; only fall back to a local low-energy cut when no such pause exists. """ maximum_seconds = _finite_float(max_segment_seconds, minimum=1.0) search_seconds = _finite_float(boundary_search_seconds, minimum=0.0) if maximum_seconds is None or search_seconds is None or sample_rate <= 0: raise ValueError("invalid Whisper segmentation settings") maximum_samples = max(1, int(round(maximum_seconds * sample_rate))) if waveform.size <= maximum_samples: return (waveform,) search_samples = int(round(search_seconds * sample_rate)) probe_radius = max(1, int(round(0.02 * sample_rate))) probe_hop = max(1, int(round(0.02 * sample_rate))) minimum_segment_samples = min( maximum_samples // 2, max(probe_radius * 2, int(round(6.0 * sample_rate))), ) minimum_quiet_run_samples = max(1, int(round(0.08 * sample_rate))) boundaries = [0] while waveform.size - boundaries[-1] > maximum_samples: segment_start = boundaries[-1] hard_limit = min(waveform.size, segment_start + maximum_samples) lower = segment_start + minimum_segment_samples upper = hard_limit - probe_radius probes = np.arange(lower, upper + 1, probe_hop, dtype=np.int64) rms = np.asarray( [ float( np.sqrt( np.mean( np.square( waveform[index - probe_radius : index + probe_radius], dtype=np.float64, ) ) ) ) for index in probes ], dtype=np.float64, ) segment_peak = float(rms.max(initial=0.0)) quiet_threshold = max(1.0e-4, segment_peak * 10.0 ** (-35.0 / 20.0)) quiet_positions = probes[rms < quiet_threshold] quiet_runs: list[tuple[int, int]] = [] for position in quiet_positions.tolist(): if quiet_runs and position <= quiet_runs[-1][1] + probe_hop: quiet_runs[-1] = (quiet_runs[-1][0], position) else: quiet_runs.append((position, position)) sustained = [ run for run in quiet_runs if run[1] - run[0] + 2 * probe_radius >= minimum_quiet_run_samples ] if sustained: begin, end = sustained[-1] boundary = (begin + end) // 2 else: fallback_lower = max(lower, hard_limit - search_samples) fallback_mask = probes >= fallback_lower fallback_probes = probes[fallback_mask] fallback_rms = rms[fallback_mask] if fallback_probes.size: boundary = int(fallback_probes[int(np.argmin(fallback_rms))]) else: boundary = hard_limit boundary = min(hard_limit, max(segment_start + 1, int(boundary))) boundaries.append(boundary) boundaries.append(waveform.size) segments = tuple( np.ascontiguousarray(waveform[start:stop], dtype=np.float32) for start, stop in zip(boundaries, boundaries[1:]) ) if ( not segments or any(segment.size == 0 for segment in segments) or any(segment.size > int(round(30.0 * sample_rate)) for segment in segments) ): raise ValueError("failed to split audio within Whisper's segment limit") return segments @torch.inference_mode() def transcribe_whisper( audio: np.ndarray | Sequence[float], sample_rate: int, *, lazy_asr: LazyWhisperASR | None = None, runtime: WhisperRuntime | None = None, language: str = "zh", task: str = "transcribe", max_new_tokens: int = 128, ) -> str: """Transcribe ndarray audio using deterministic decoding. ``runtime`` and ``lazy_asr`` are mutually exclusive injection points. An empty decoded string is returned as-is; the semantic verifier will reject it rather than accepting an arbitrary TTS fallback. """ if runtime is not None and lazy_asr is not None: raise ValueError("pass either runtime or lazy_asr, not both") token_limit = int(max_new_tokens) if token_limit <= 0: raise ValueError("max_new_tokens must be positive") waveform = _resample_audio(_mono_audio(audio), int(sample_rate), WHISPER_SAMPLE_RATE) segments = _split_whisper_audio(waveform) selected_runtime = runtime or (lazy_asr or _DEFAULT_WHISPER).get_runtime() processor_input: np.ndarray | list[np.ndarray] processor_input = segments[0] if len(segments) == 1 else list(segments) processor_output = selected_runtime.processor( processor_input, sampling_rate=WHISPER_SAMPLE_RATE, return_tensors="pt", return_attention_mask=WHISPER_RETURN_ATTENTION_MASK, ) features = processor_output.input_features.to( device=selected_runtime.device, dtype=selected_runtime.dtype, ) attention_mask = processor_output.attention_mask.to(device=selected_runtime.device) token_ids = selected_runtime.model.generate( features, attention_mask=attention_mask, language=language, task=task, do_sample=False, num_beams=1, max_new_tokens=token_limit, ) decoded = selected_runtime.processor.batch_decode(token_ids, skip_special_tokens=True) if not decoded or len(decoded) != len(segments): return "" return " ".join(str(text).strip() for text in decoded if str(text).strip()) def transcribe_verification_whisper( audio: np.ndarray | Sequence[float], sample_rate: int, *, lazy_asr: LazyWhisperASR | None = None, runtime: WhisperRuntime | None = None, language: str = "zh", task: str = "transcribe", max_new_tokens: int = 128, ) -> str: """Transcribe with the separately pinned full large-v3 final verifier.""" selected_lazy = lazy_asr if runtime is None and selected_lazy is None: selected_lazy = _DEFAULT_VERIFICATION_WHISPER return transcribe_whisper( audio, sample_rate, lazy_asr=selected_lazy, runtime=runtime, language=language, task=task, max_new_tokens=max_new_tokens, ) @dataclass(frozen=True) class PreparedCandidateAudio: """Validated candidate waveform and its ASR transcript.""" waveform: np.ndarray duration_seconds: float transcript_text: str def prepare_candidate_audio( audio: np.ndarray | Sequence[float], sample_rate: int, *, transcriber: Callable[[np.ndarray, int], str] | None = None, ) -> PreparedCandidateAudio | None: """Prepare one candidate, returning ``None`` for candidate-data errors. Invalid/empty/non-finite audio and a transcriber's ``ValueError`` describe an unusable candidate, not a service outage. They therefore become a normal gate rejection so the cascade can try the next seed. Runtime and I/O failures deliberately propagate and abort the request fail-closed. """ try: selected_sample_rate = int(sample_rate) if selected_sample_rate <= 0: raise ValueError("sample_rate must be positive") waveform = _mono_audio(audio) transcript = (transcriber or transcribe_whisper)( waveform, selected_sample_rate, ) if not isinstance(transcript, str): raise ValueError("ASR transcript must be a string") except (TypeError, ValueError, OverflowError): return None return PreparedCandidateAudio( waveform=waveform, duration_seconds=waveform.size / float(selected_sample_rate), transcript_text=transcript.strip(), ) @torch.inference_mode() def _encode_speaker_segments( segments: Sequence[np.ndarray], encoder: Any, *, device: str | torch.device, ) -> np.ndarray: """Encode a variable-length segment batch in one ECAPA forward pass.""" waveforms = tuple(np.asarray(segment, dtype=np.float32).reshape(-1) for segment in segments) if not waveforms or any( waveform.size == 0 or not np.isfinite(waveform).all() for waveform in waveforms ): raise ValueError("speaker segments must be non-empty and finite") maximum_length = max(waveform.size for waveform in waveforms) batch = np.zeros((len(waveforms), maximum_length), dtype=np.float32) relative_lengths = np.empty(len(waveforms), dtype=np.float32) for index, waveform in enumerate(waveforms): batch[index, : waveform.size] = waveform relative_lengths[index] = waveform.size / float(maximum_length) selected_device = torch.device(device) tensor = torch.from_numpy(batch).to(selected_device) wav_lens = torch.from_numpy(relative_lengths).to(selected_device) embeddings = encoder.encode_batch(tensor, wav_lens=wav_lens) embeddings = torch.as_tensor(embeddings).detach().float() if embeddings.ndim == 0 or embeddings.shape[0] != len(waveforms): raise ValueError("speaker encoder returned an invalid batch size") embeddings = embeddings.reshape(len(waveforms), -1) if embeddings.shape[1] == 0 or not torch.isfinite(embeddings).all(): raise ValueError("speaker encoder returned an invalid embedding") norms = torch.linalg.vector_norm(embeddings, dim=1) if not torch.isfinite(norms).all() or bool(torch.any(norms <= 1.0e-8)): raise ValueError("speaker encoder returned a zero-norm embedding") normalized = torch_functional.normalize(embeddings, dim=1).cpu().numpy().astype(np.float32) if not np.isfinite(normalized).all(): raise ValueError("speaker encoder returned a non-finite embedding") return normalized @torch.inference_mode() def speaker_embedding_from_audio( audio: np.ndarray | Sequence[float], sample_rate: int, encoder: Any, *, device: str | torch.device = "cpu", target_sample_rate: int = 16_000, active_top_db: float = 35.0, ) -> np.ndarray: """Extract one normalized ECAPA embedding from active ndarray speech.""" waveform = _mono_audio(audio) waveform = _resample_audio(waveform, int(sample_rate), int(target_sample_rate)) waveform = _trim_active_speech(waveform, top_db=active_top_db) return _encode_speaker_segments((waveform,), encoder, device=device)[0] def cosine_similarity(left: np.ndarray | Sequence[float], right: np.ndarray | Sequence[float]) -> float: """Return a finite cosine similarity, raising on unusable embeddings.""" left_array = np.asarray(left, dtype=np.float64).reshape(-1) right_array = np.asarray(right, dtype=np.float64).reshape(-1) if left_array.size == 0 or left_array.shape != right_array.shape: raise ValueError("speaker embeddings must have equal non-empty shapes") if not np.isfinite(left_array).all() or not np.isfinite(right_array).all(): raise ValueError("speaker embeddings must be finite") denominator = float(np.linalg.norm(left_array) * np.linalg.norm(right_array)) if not math.isfinite(denominator) or denominator <= 1.0e-12: raise ValueError("speaker embeddings must have non-zero norm") similarity = float(np.dot(left_array, right_array) / denominator) if not math.isfinite(similarity): raise ValueError("speaker cosine similarity is non-finite") return float(np.clip(similarity, -1.0, 1.0)) @dataclass(frozen=True) class SpeakerEvidence: similarity: float begin_similarity: float end_similarity: float boundary_drop: float active_duration_seconds: float speaker_embedding: np.ndarray active_rms_db: float def speaker_evidence_from_audio( audio: np.ndarray | Sequence[float], sample_rate: int, encoder: Any, anchor_embedding: np.ndarray | Sequence[float], *, device: str | torch.device = "cpu", edge_seconds: float = 1.5, whole_window_seconds: float = 3.0, whole_max_windows: int = 4, active_top_db: float = 35.0, ) -> SpeakerEvidence: """Measure whole/begin/end anchor similarity on trimmed active speech.""" edge_duration = _finite_float(edge_seconds, minimum=0.01) window_duration = _finite_float(whole_window_seconds, minimum=0.01) try: max_windows = int(whole_max_windows) except (TypeError, ValueError, OverflowError): max_windows = 0 if edge_duration is None or window_duration is None or max_windows <= 0: raise ValueError("speaker window settings must be finite and positive") waveform = _resample_audio(_mono_audio(audio), int(sample_rate), 16_000) active_duration_seconds = active_voiced_duration_seconds( waveform, 16_000, top_db=active_top_db, ) active = _trim_active_speech(waveform, top_db=active_top_db) edge_samples = max(1, int(round(edge_duration * 16_000))) whole_window_samples = max(1, int(round(window_duration * 16_000))) begin = active[:edge_samples] end = active[-edge_samples:] if active.size <= whole_window_samples: whole_segments = [active] else: starts = np.linspace( 0, active.size - whole_window_samples, num=max_windows, ).round().astype(int) whole_segments = [ active[start : start + whole_window_samples] for start in dict.fromkeys(starts.tolist()) ] embeddings = _encode_speaker_segments( (*whole_segments, begin, end), encoder, device=device, ) whole_embeddings = embeddings[: len(whole_segments)] whole_embedding = np.mean(whole_embeddings, axis=0, dtype=np.float64) whole_norm = float(np.linalg.norm(whole_embedding)) if not math.isfinite(whole_norm) or whole_norm <= 1.0e-8: raise ValueError("speaker windows produced a zero-norm embedding") whole_embedding = np.asarray(whole_embedding / whole_norm, dtype=np.float32) begin_embedding, end_embedding = embeddings[-2:] similarity = cosine_similarity(whole_embedding, anchor_embedding) begin_similarity = cosine_similarity(begin_embedding, anchor_embedding) end_similarity = cosine_similarity(end_embedding, anchor_embedding) boundary_drop = max(0.0, begin_similarity - end_similarity) active_rms = float(np.sqrt(np.mean(np.square(active, dtype=np.float64)))) if not math.isfinite(active_rms) or active_rms <= 1.0e-8: raise ValueError("active speech has invalid RMS") return SpeakerEvidence( similarity=similarity, begin_similarity=begin_similarity, end_similarity=end_similarity, boundary_drop=boundary_drop, active_duration_seconds=active_duration_seconds, speaker_embedding=whole_embedding.copy(), active_rms_db=20.0 * math.log10(active_rms), ) def release_speaker_evidence_from_audio( audio: np.ndarray | Sequence[float], sample_rate: int, encoder: Any, anchor_embedding: np.ndarray | Sequence[float], *, device: str | torch.device = "cpu", active_top_db: float = 35.0, ) -> SpeakerEvidence: """Measure the exact full/third speaker evidence used for promotion. Candidate-local scoring keeps the bounded window metric for latency. A returned waveform is additionally checked with this independent-aligned metric so a low online boundary drop cannot hide a release-gate failure. """ try: rate = operator.index(sample_rate) except (TypeError, ValueError, OverflowError) as error: raise ValueError("sample rate must be positive") from error if rate <= 0: raise ValueError("sample rate must be positive") waveform = _mono_audio(audio) active_duration_seconds = active_voiced_duration_seconds( waveform, rate, top_db=active_top_db, ) active = trim_release_speaker_activity( waveform, rate, top_db=active_top_db, ) source_segments = (active, *tuple(np.array_split(active, 3))) if any(segment.size == 0 for segment in source_segments): raise ValueError("release speaker segments must be non-empty") import librosa segments_16k: list[np.ndarray] = [] for segment in source_segments: resampled = ( segment if rate == 16_000 else librosa.resample(segment, orig_sr=rate, target_sr=16_000) ) segments_16k.append(np.ascontiguousarray(resampled, dtype=np.float32)) # Encode one segment at a time to match the independent verifier rather # than allowing padding/batching to perturb short boundary embeddings. selected_device = torch.device(device) encoded: list[np.ndarray] = [] for segment in segments_16k: tensor = torch.from_numpy(segment).unsqueeze(0).to(selected_device) embedding = ( torch.as_tensor(encoder.encode_batch(tensor)) .detach() .float() .reshape(-1) ) if embedding.numel() == 0 or not bool(torch.isfinite(embedding).all()): raise ValueError("speaker encoder returned an invalid embedding") norm = torch.linalg.vector_norm(embedding) if not bool(torch.isfinite(norm)) or float(norm) <= 1.0e-8: raise ValueError("speaker encoder returned a zero-norm embedding") encoded.append((embedding / norm).cpu().numpy().astype(np.float32)) embeddings = np.stack(encoded, axis=0) whole_embedding = embeddings[0] begin_embedding = embeddings[1] end_embedding = embeddings[-1] similarity = cosine_similarity(whole_embedding, anchor_embedding) begin_similarity = cosine_similarity(begin_embedding, anchor_embedding) end_similarity = cosine_similarity(end_embedding, anchor_embedding) active_rms = float(np.sqrt(np.mean(np.square(active, dtype=np.float64)))) if not math.isfinite(active_rms) or active_rms <= 1.0e-8: raise ValueError("active speech has invalid RMS") return SpeakerEvidence( similarity=similarity, begin_similarity=begin_similarity, end_similarity=end_similarity, boundary_drop=max(0.0, begin_similarity - end_similarity), active_duration_seconds=active_duration_seconds, speaker_embedding=whole_embedding.copy(), active_rms_db=20.0 * math.log10(active_rms), ) def active_audio_rms_db(audio: np.ndarray | Sequence[float], *, top_db: float = 35.0) -> float: """Measure finite RMS dB on the active region of candidate audio.""" active = _trim_active_speech(_mono_audio(audio), top_db=top_db) rms = float(np.sqrt(np.mean(np.square(active, dtype=np.float64)))) if not math.isfinite(rms) or rms <= 1.0e-8: raise ValueError("active speech has invalid RMS") return 20.0 * math.log10(rms) @dataclass(frozen=True) class CandidateObservation: target_text: str transcript_text: str audio_duration_seconds: float speaker_similarity: float | None = None begin_speaker_similarity: float | None = None end_speaker_similarity: float | None = None pace_cps: float | None = None truncated: bool = False @dataclass(frozen=True) class CandidateGateResult: passed: bool comparison: AsrComparison speaker_gate_applied: bool speaker_similarity: float | None boundary_speaker_drop: float | None pace_cps: float | None score: float rejection_reasons: tuple[str, ...] @dataclass(frozen=True) class ChunkCandidateArtifact: """Acoustic evidence retained for sequence-level candidate selection.""" speaker_embedding: np.ndarray | None = None rms_db: float | None = None median_f0_hz: float | None = None @dataclass(frozen=True) class CandidateGateEvidence: """Bounded, content-free diagnostic snapshot of one hard-gate result.""" passed: bool target_units: int hypothesis_units: int edit_distance: int cer: float | None prefix_cer: float | None suffix_cer: float | None prefix_deletions: int suffix_deletions: int extra_tail_units: int speaker_gate_applied: bool speaker_similarity: float | None boundary_speaker_drop: float | None pace_cps: float | None score: float | None rejection_reasons: tuple[str, ...] @dataclass(frozen=True) class TrajectoryGateEvidence: """Content-free evidence for a joined, sequence, or final waveform gate.""" passed: bool result_count: int score: float | None rejection_reasons: tuple[str, ...] result: CandidateGateEvidence | None @dataclass(frozen=True) class CandidateAttemptEvidence: """Diagnostics for one generated trajectory without waveform/text payloads.""" candidate_index: int seed: int trajectory_passed: bool trajectory_score: float | None trajectory_rejection_reasons: tuple[str, ...] local_result_count: int local_results: tuple[CandidateGateEvidence, ...] joined_output: TrajectoryGateEvidence | None = None @dataclass(frozen=True) class SequencePathEvidence: """Final-verifier evidence for one ranked DP path.""" rank: int chunk_candidate_indices: tuple[int, ...] chunk_seeds: tuple[int, ...] final_output: TrajectoryGateEvidence @dataclass(frozen=True) class SequenceSearchEvidence: """Finite-lattice diagnostics explaining why sequence DP can or cannot run.""" eligible_candidate_counts: tuple[int, ...] finite_transition_counts: tuple[int, ...] ranked_path_count: int checked_paths: tuple[SequencePathEvidence, ...] = () @dataclass(frozen=True) class CascadeDiagnostics: """Bounded diagnostics retained on both success and fail-closed outcomes.""" attempts: tuple[CandidateAttemptEvidence, ...] = () sequence_search: SequenceSearchEvidence | None = None def _evidence_float(value: Any) -> float | None: return _finite_float(value) def _evidence_int(value: Any, *, maximum: int = 1_000_000) -> int: if isinstance(value, (bool, np.bool_)): return 0 try: integer = operator.index(value) except (TypeError, ValueError, OverflowError): return 0 return min(max(0, int(integer)), maximum) def _sanitize_rejection_reason(reason: Any) -> str: """Return an allow-listed reason code without forwarding arbitrary text.""" if not isinstance(reason, str) or len(reason) > 96: return "unknown_rejection" tokens = reason.split(":") if ( not tokens or len(tokens) > 3 or tokens[-1] not in _CASCADE_EVIDENCE_REJECTION_CODES ): return "unknown_rejection" for token in tokens[:-1]: if token == "joined_output": continue if token.startswith("chunk_"): chunk_index = token[6:] if ( chunk_index.isdigit() and len(chunk_index) <= 2 and int(chunk_index) < CASCADE_EVIDENCE_MAX_LOCAL_RESULTS ): continue return "unknown_rejection" return ":".join(tokens) def _bounded_rejection_reasons(reasons: Any) -> tuple[str, ...]: try: values = tuple(reasons) except TypeError: values = () return tuple( _sanitize_rejection_reason(reason) for reason in values[:CASCADE_EVIDENCE_MAX_REASONS] ) def candidate_gate_evidence(result: CandidateGateResult) -> CandidateGateEvidence: """Project a gate result to finite scalar/count evidence only.""" if not isinstance(result, CandidateGateResult): raise TypeError("result must be a CandidateGateResult") comparison = result.comparison return CandidateGateEvidence( passed=result.passed is True, target_units=_evidence_int(len(comparison.target_text)), hypothesis_units=_evidence_int(len(comparison.transcript_text)), edit_distance=_evidence_int(comparison.edit_distance), cer=_evidence_float(comparison.cer), prefix_cer=_evidence_float(comparison.prefix_cer), suffix_cer=_evidence_float(comparison.suffix_cer), prefix_deletions=_evidence_int(comparison.prefix_deletions), suffix_deletions=_evidence_int(comparison.suffix_deletions), extra_tail_units=_evidence_int(comparison.extra_tail_units), speaker_gate_applied=result.speaker_gate_applied is True, speaker_similarity=_evidence_float(result.speaker_similarity), boundary_speaker_drop=_evidence_float(result.boundary_speaker_drop), pace_cps=_evidence_float(result.pace_cps), score=_evidence_float(result.score), rejection_reasons=_bounded_rejection_reasons(result.rejection_reasons), ) def verify_candidate( observation: CandidateObservation, *, locale: str = "zh-TW", short_text_units: int = 6, short_text_max_cer: float = 0.0, max_cer: float = 0.20, prefix_units: int = 6, suffix_units: int = 6, max_prefix_cer: float = 0.0, max_suffix_cer: float = 0.0, max_extra_tail_units: int = 0, speaker_gate_enabled: bool = True, short_audio_seconds: float = 1.5, min_speaker_similarity: float = 0.10, max_boundary_speaker_drop: float = 0.03, max_pace_cps: float | None = None, speaker_weight: float = 0.05, boundary_weight: float = 0.10, ) -> CandidateGateResult: """Apply strict semantic and duration-aware speaker gates to a candidate.""" duration = _finite_float(observation.audio_duration_seconds, minimum=0.0) short_duration_limit = _finite_float(short_audio_seconds, minimum=0.0) general_cer_limit = _finite_float(max_cer, minimum=0.0) exact_cer_limit = _finite_float(short_text_max_cer, minimum=0.0) min_similarity = _finite_float(min_speaker_similarity, minimum=-1.0, maximum=1.0) max_boundary = _finite_float(max_boundary_speaker_drop, minimum=0.0) max_pace = None if max_pace_cps is None else _finite_float(max_pace_cps, minimum=0.0) speaker_cost_weight = _finite_float(speaker_weight, minimum=0.0) boundary_cost_weight = _finite_float(boundary_weight, minimum=0.0) try: short_unit_limit = max(0, int(short_text_units)) except (TypeError, ValueError, OverflowError): short_unit_limit = -1 speaker_enabled = isinstance(speaker_gate_enabled, (bool, np.bool_)) if speaker_enabled: speaker_enabled = bool(speaker_gate_enabled) # First normalize with a permissive finite limit to determine target units. preliminary = compare_asr_text( observation.target_text, observation.transcript_text, locale=locale, prefix_units=prefix_units, suffix_units=suffix_units, max_cer=general_cer_limit if general_cer_limit is not None else math.nan, max_prefix_cer=max_prefix_cer, max_suffix_cer=max_suffix_cer, max_extra_tail_units=max_extra_tail_units, ) selected_cer_limit = general_cer_limit if short_unit_limit >= 0 and len(preliminary.target_text) <= short_unit_limit: selected_cer_limit = exact_cer_limit comparison = compare_asr_text( observation.target_text, observation.transcript_text, locale=locale, prefix_units=prefix_units, suffix_units=suffix_units, max_cer=selected_cer_limit if selected_cer_limit is not None else math.nan, max_prefix_cer=max_prefix_cer, max_suffix_cer=max_suffix_cer, max_extra_tail_units=max_extra_tail_units, ) reasons: list[str] = [] if duration is None or duration <= 0.0: reasons.append("invalid_audio_duration") if observation.truncated is not False: reasons.append("truncated") if not comparison.passed: reasons.append("semantic_gate") pace = None if max_pace_cps is not None: pace = _finite_float(observation.pace_cps, minimum=0.0) if max_pace is None: reasons.append("invalid_gate_config") elif pace is None: reasons.append("missing_pace_evidence") elif pace > max_pace: reasons.append("pace_too_fast") valid_common_config = bool( isinstance(speaker_gate_enabled, (bool, np.bool_)) and short_duration_limit is not None and general_cer_limit is not None and exact_cer_limit is not None and short_unit_limit >= 0 and ( not speaker_enabled or all( value is not None for value in ( min_similarity, max_boundary, speaker_cost_weight, boundary_cost_weight, ) ) ) ) if not valid_common_config: reasons.append("invalid_gate_config") speaker_gate_applied = bool( speaker_enabled and duration is not None and short_duration_limit is not None and duration >= short_duration_limit ) similarity: float | None = None boundary_drop: float | None = None if speaker_gate_applied: similarity = _finite_float( observation.speaker_similarity, minimum=-1.0, maximum=1.0, ) begin_similarity = _finite_float( observation.begin_speaker_similarity, minimum=-1.0, maximum=1.0, ) end_similarity = _finite_float( observation.end_speaker_similarity, minimum=-1.0, maximum=1.0, ) if similarity is None or begin_similarity is None or end_similarity is None: reasons.append("missing_speaker_evidence") else: boundary_drop = max(0.0, begin_similarity - end_similarity) if min_similarity is None or similarity < min_similarity: reasons.append("speaker_similarity") if max_boundary is None or boundary_drop > max_boundary: reasons.append("boundary_speaker_drop") score = math.inf if not reasons: score = comparison.cer if speaker_gate_applied: assert similarity is not None and boundary_drop is not None assert speaker_cost_weight is not None and boundary_cost_weight is not None score += speaker_cost_weight * (1.0 - similarity) score += boundary_cost_weight * boundary_drop if not math.isfinite(score) or score < 0.0: reasons.append("nonfinite_score") score = math.inf return CandidateGateResult( passed=not reasons, comparison=comparison, speaker_gate_applied=speaker_gate_applied, speaker_similarity=similarity, boundary_speaker_drop=boundary_drop, pace_cps=pace, score=score, rejection_reasons=tuple(reasons), ) @dataclass(frozen=True) class TrajectoryGateResult: passed: bool candidate_results: tuple[CandidateGateResult, ...] score: float rejection_reasons: tuple[str, ...] chunk_artifacts: tuple[ChunkCandidateArtifact, ...] = () @dataclass(frozen=True) class CandidateVerification: """Selection result plus content-free joined-output evidence.""" verification: TrajectoryGateResult joined_output: TrajectoryGateEvidence | None = None def trajectory_gate_evidence( verification: TrajectoryGateResult, ) -> TrajectoryGateEvidence: """Project a joined/final verification without retaining recognized text.""" if not isinstance(verification, TrajectoryGateResult): raise TypeError("verification must be a TrajectoryGateResult") sole_result = ( verification.candidate_results[0] if len(verification.candidate_results) == 1 else None ) result = ( candidate_gate_evidence(sole_result) if isinstance(sole_result, CandidateGateResult) else None ) return TrajectoryGateEvidence( passed=verification.passed is True, result_count=_evidence_int(len(verification.candidate_results)), score=_evidence_float(verification.score), rejection_reasons=_bounded_rejection_reasons( verification.rejection_reasons ), result=result, ) def exact_waveform_sha256( audio: np.ndarray | Sequence[float], sample_rate: int, ) -> str: """Hash exact canonical float32 samples together with their sample rate.""" if isinstance(sample_rate, (bool, np.bool_)): raise ValueError("sample_rate must be a positive integer") try: rate = operator.index(sample_rate) except (TypeError, ValueError, OverflowError) as error: raise ValueError("sample_rate must be a positive integer") from error if rate <= 0: raise ValueError("sample_rate must be a positive integer") waveform = _mono_audio(audio) canonical = np.ascontiguousarray(waveform, dtype=np.dtype(" None: self._entries: dict[ tuple[str, int, str, str], TrajectoryGateResult, ] = {} @property def entry_count(self) -> int: return len(self._entries) def verify( self, audio: np.ndarray | Sequence[float], sample_rate: int, target_text: str, verifier_profile: str, verifier: Callable[[np.ndarray, int, str], TrajectoryGateResult], ) -> TrajectoryGateResult: if not isinstance(target_text, str) or not target_text: raise ValueError("target_text must be a non-empty string") if not isinstance(verifier_profile, str) or not verifier_profile: raise ValueError("verifier_profile must be a non-empty string") if not callable(verifier): raise ValueError("verifier must be callable") waveform = _mono_audio(audio) if isinstance(sample_rate, (bool, np.bool_)): raise ValueError("sample_rate must be a positive integer") try: rate = operator.index(sample_rate) except (TypeError, ValueError, OverflowError) as error: raise ValueError("sample_rate must be a positive integer") from error if rate <= 0: raise ValueError("sample_rate must be a positive integer") waveform_hash = exact_waveform_sha256(waveform, rate) key = (waveform_hash, rate, target_text, verifier_profile) cached = self._entries.get(key) if cached is not None: return cached verification = verifier(waveform, rate, target_text) if not isinstance(verification, TrajectoryGateResult): raise RuntimeError("whole-waveform verifier returned an invalid result") self._entries[key] = verification return verification def verify_trajectory( observations: Sequence[CandidateObservation], *, chunk_artifacts: Sequence[ChunkCandidateArtifact] = (), **candidate_gate_kwargs: Any, ) -> TrajectoryGateResult: """Require every chunk in a non-empty trajectory to pass all hard gates.""" try: candidates = tuple(observations) except TypeError: candidates = () if not candidates: return TrajectoryGateResult(False, (), math.inf, ("empty_trajectory",)) try: artifacts = tuple(chunk_artifacts) except TypeError: artifacts = () if artifacts and ( len(artifacts) != len(candidates) or any(not isinstance(artifact, ChunkCandidateArtifact) for artifact in artifacts) ): return TrajectoryGateResult( False, (), math.inf, ("invalid_chunk_artifacts",), ) results: list[CandidateGateResult] = [] rejection_reasons: list[str] = [] for index, observation in enumerate(candidates): try: result = verify_candidate(observation, **candidate_gate_kwargs) except (TypeError, ValueError, OverflowError): # A malformed observation must reject the entire trajectory. comparison = compare_asr_text("", "") result = CandidateGateResult( passed=False, comparison=comparison, speaker_gate_applied=False, speaker_similarity=None, boundary_speaker_drop=None, pace_cps=None, score=math.inf, rejection_reasons=("malformed_observation",), ) results.append(result) rejection_reasons.extend(f"chunk_{index}:{reason}" for reason in result.rejection_reasons) passed = bool(results) and all(result.passed for result in results) score = sum(result.score for result in results) if passed else math.inf if not math.isfinite(score): passed = False score = math.inf if not rejection_reasons: rejection_reasons.append("nonfinite_trajectory_score") return TrajectoryGateResult( passed=passed, candidate_results=tuple(results), score=score, rejection_reasons=tuple(rejection_reasons), chunk_artifacts=artifacts, ) def qualify_trajectory_with_joined_output( local_verification: TrajectoryGateResult, joined_verification: TrajectoryGateResult, ) -> TrajectoryGateResult: """Require joined-output safety without discarding DP-local evidence. ``candidate_results`` and ``chunk_artifacts`` always stay local to the generated chunks. This lets sequence DP reuse individually safe chunks when RMS matching, fades, pauses or crossfade make the same-seed joined waveform fail its whole-output gate. """ if not isinstance(local_verification, TrajectoryGateResult): raise TypeError("local_verification must be a TrajectoryGateResult") if local_verification.passed is not True or not math.isfinite( local_verification.score ): return local_verification if not isinstance(joined_verification, TrajectoryGateResult): raise TypeError("joined_verification must be a TrajectoryGateResult") joined_passed = bool( joined_verification.passed is True and math.isfinite(joined_verification.score) and len(joined_verification.candidate_results) == 1 ) if joined_passed: return local_verification joined_reasons = joined_verification.rejection_reasons or ( "invalid_joined_verification", ) return TrajectoryGateResult( passed=False, candidate_results=local_verification.candidate_results, score=math.inf, rejection_reasons=tuple( f"joined_output:{reason}" for reason in joined_reasons ), chunk_artifacts=local_verification.chunk_artifacts, ) class NoQualifiedCandidateError(RuntimeError): """Raised when the full adaptive cascade has no verified trajectory.""" def __init__( self, message: str, *, diagnostics: CascadeDiagnostics | None = None, ) -> None: super().__init__(message) self.diagnostics = diagnostics or CascadeDiagnostics() class FinalOutputRejectedError(RuntimeError): """Raised when post-join output fails the final whole-waveform gate.""" def require_verified_final_output( verification: TrajectoryGateResult, ) -> TrajectoryGateResult: """Return verified final evidence or reject without an audio fallback.""" if not isinstance(verification, TrajectoryGateResult): raise FinalOutputRejectedError("final verifier returned an invalid result") if ( verification.passed is not True or not verification.candidate_results or not all(result.passed for result in verification.candidate_results) or not math.isfinite(verification.score) ): reasons = ",".join(verification.rejection_reasons) or "unsafe_final_output" raise FinalOutputRejectedError(f"final output rejected: {reasons}") return verification @dataclass(frozen=True) class CascadeResult: trajectory: Any verification: TrajectoryGateResult seed: int | None candidate_index: int | None attempted_seeds: tuple[int, ...] chunk_candidate_indices: tuple[int, ...] = () chunk_seeds: tuple[int, ...] = () selection_mode: str = "whole_trajectory" sequence_path_rank: int | None = None sequence_paths_checked: int = 0 diagnostics: CascadeDiagnostics = CascadeDiagnostics() def _candidate_gate_evidence_payload( evidence: CandidateGateEvidence, ) -> dict[str, Any]: return { "passed": evidence.passed, "target_units": evidence.target_units, "hypothesis_units": evidence.hypothesis_units, "edit_distance": evidence.edit_distance, "cer": evidence.cer, "prefix_cer": evidence.prefix_cer, "suffix_cer": evidence.suffix_cer, "prefix_deletions": evidence.prefix_deletions, "suffix_deletions": evidence.suffix_deletions, "extra_tail_units": evidence.extra_tail_units, "speaker_gate_applied": evidence.speaker_gate_applied, "speaker_similarity": evidence.speaker_similarity, "boundary_speaker_drop": evidence.boundary_speaker_drop, "pace_cps": evidence.pace_cps, "score": evidence.score, "reasons": list(evidence.rejection_reasons), } def _trajectory_gate_evidence_payload( evidence: TrajectoryGateEvidence | None, ) -> dict[str, Any] | None: if evidence is None: return None return { "passed": evidence.passed, "result_count": evidence.result_count, "score": evidence.score, "reasons": list(evidence.rejection_reasons), "result": ( None if evidence.result is None else _candidate_gate_evidence_payload(evidence.result) ), } def _candidate_attempt_evidence_payload( evidence: CandidateAttemptEvidence, ) -> dict[str, Any]: return { "candidate_index": evidence.candidate_index, "seed": evidence.seed, "policy": generation_policy_for_candidate_offset( evidence.candidate_index ).name, "trajectory_passed": evidence.trajectory_passed, "trajectory_score": evidence.trajectory_score, "trajectory_reasons": list( evidence.trajectory_rejection_reasons ), "local_result_count": evidence.local_result_count, "local_results": [ _candidate_gate_evidence_payload(result) for result in evidence.local_results ], "joined_output": _trajectory_gate_evidence_payload( evidence.joined_output ), } def _sequence_search_evidence_payload( evidence: SequenceSearchEvidence | None, ) -> dict[str, Any] | None: if evidence is None: return None eligible_counts = evidence.eligible_candidate_counts[ :CASCADE_EVIDENCE_MAX_LOCAL_RESULTS ] transition_counts = evidence.finite_transition_counts[ : max(0, CASCADE_EVIDENCE_MAX_LOCAL_RESULTS - 1) ] return { "row_count": len(evidence.eligible_candidate_counts), "eligible_candidate_counts": list(eligible_counts), "zero_eligible_rows": [ index for index, count in enumerate(eligible_counts) if count == 0 ], "transition_boundary_count": len(evidence.finite_transition_counts), "finite_transition_counts": list(transition_counts), "zero_transition_edges": [ index for index, count in enumerate(transition_counts) if count == 0 ], "ranked_path_count": evidence.ranked_path_count, "checked_path_count": len(evidence.checked_paths), "checked_paths": [ { "rank": path.rank, "chunk_candidates": list( path.chunk_candidate_indices[ :CASCADE_EVIDENCE_MAX_LOCAL_RESULTS ] ), "chunk_seeds": list( path.chunk_seeds[:CASCADE_EVIDENCE_MAX_LOCAL_RESULTS] ), "final_output": _trajectory_gate_evidence_payload( path.final_output ), } for path in evidence.checked_paths[ :CASCADE_EVIDENCE_MAX_SEQUENCE_PATHS ] ], } def format_cascade_evidence_log( diagnostics: CascadeDiagnostics, *, outcome: str, candidate_limit: int, selection: CascadeResult | None = None, final_output: TrajectoryGateEvidence | None = None, ) -> str: """Return one canonical JSON log line containing only bounded evidence.""" if not isinstance(diagnostics, CascadeDiagnostics): raise TypeError("diagnostics must be CascadeDiagnostics") if outcome not in _CASCADE_EVIDENCE_OUTCOMES: raise ValueError("invalid cascade evidence outcome") if isinstance(candidate_limit, (bool, np.bool_)): raise ValueError("candidate_limit must be an integer between 1 and 20") try: limit = operator.index(candidate_limit) except (TypeError, ValueError, OverflowError) as error: raise ValueError( "candidate_limit must be an integer between 1 and 20" ) from error if not 1 <= limit <= CASCADE_EVIDENCE_MAX_ATTEMPTS: raise ValueError("candidate_limit must be an integer between 1 and 20") selected: dict[str, Any] | None = None if selection is not None: if not isinstance(selection, CascadeResult): raise TypeError("selection must be a CascadeResult") mode = ( selection.selection_mode if selection.selection_mode in _CASCADE_EVIDENCE_SELECTION_MODES else None ) selected = { "selection": mode, "candidate_index": selection.candidate_index, "seed": selection.seed, "chunk_candidates": list( selection.chunk_candidate_indices[ :CASCADE_EVIDENCE_MAX_LOCAL_RESULTS ] ), "chunk_seeds": list( selection.chunk_seeds[:CASCADE_EVIDENCE_MAX_LOCAL_RESULTS] ), "sequence_rank": selection.sequence_path_rank, "sequence_paths_checked": selection.sequence_paths_checked, } attempts = diagnostics.attempts[:CASCADE_EVIDENCE_MAX_ATTEMPTS] payload = { "schema_version": CASCADE_EVIDENCE_SCHEMA_VERSION, "outcome": outcome, "candidate_limit": int(limit), "attempt_count": len(diagnostics.attempts), "attempts": [ _candidate_attempt_evidence_payload(attempt) for attempt in attempts ], "sequence_search": _sequence_search_evidence_payload( diagnostics.sequence_search ), "selection": selected, "final_output": _trajectory_gate_evidence_payload(final_output), } return CASCADE_EVIDENCE_LOG_PREFIX + json.dumps( payload, allow_nan=False, ensure_ascii=True, separators=(",", ":"), sort_keys=True, ) def select_k_candidate_sequences( local_scores: Sequence[Sequence[float]], transition_scores: Sequence[Sequence[Sequence[float]]] = (), *, max_paths: int = 3, ) -> tuple[CandidateSequenceSelection, ...]: """Return up to three distinct finite paths in stable cost order. Each DP state retains only its ``max_paths`` best prefixes. This keeps the search bounded at ``O(N K² max_paths)`` while still producing exact k-best paths for the requested small bound. Single-chunk requests intentionally return no sequence fallback. """ if isinstance(max_paths, (bool, np.bool_)): raise ValueError("max_paths must be an integer between 1 and 3") try: path_limit = operator.index(max_paths) except (TypeError, ValueError, OverflowError) as error: raise ValueError("max_paths must be an integer between 1 and 3") from error if not 1 <= path_limit <= 3: raise ValueError("max_paths must be an integer between 1 and 3") try: raw_local = [list(row) for row in local_scores] except TypeError: return () if len(raw_local) <= 1 or any(not row for row in raw_local): return () safe_local = [ [ score if score is not None else math.inf for score in (_finite_float(value, minimum=0.0) for value in row) ] for row in raw_local ] try: raw_transitions = [ [list(row) for row in matrix] for matrix in transition_scores ] except TypeError: return () if len(raw_transitions) != len(safe_local) - 1: return () safe_transitions: list[list[list[float]]] = [] for step, matrix in enumerate(raw_transitions): previous_count = len(safe_local[step]) current_count = len(safe_local[step + 1]) if len(matrix) != previous_count or any( len(row) != current_count for row in matrix ): return () safe_transitions.append( [ [ score if score is not None else math.inf for score in ( _finite_float(value, minimum=0.0) for value in row ) ] for row in matrix ] ) # One list of (cost, path) prefixes for each current candidate position. states: list[list[tuple[float, tuple[int, ...]]]] = [] for candidate_index, score in enumerate(safe_local[0]): states.append( [(score, (candidate_index,))] if math.isfinite(score) else [] ) for step in range(1, len(safe_local)): next_states: list[list[tuple[float, tuple[int, ...]]]] = [] for current_index, local_score in enumerate(safe_local[step]): options: dict[tuple[int, ...], float] = {} if math.isfinite(local_score): for previous_index, prefixes in enumerate(states): edge_score = safe_transitions[step - 1][previous_index][ current_index ] if not math.isfinite(edge_score): continue for previous_score, prefix in prefixes: total = previous_score + edge_score + local_score path = prefix + (current_index,) if math.isfinite(total): old_score = options.get(path, math.inf) if total < old_score: options[path] = total ranked = sorted( ((score, path) for path, score in options.items()), key=lambda item: (item[0], item[1]), )[:path_limit] next_states.append(ranked) states = next_states complete: dict[tuple[int, ...], float] = {} for prefixes in states: for score, path in prefixes: old_score = complete.get(path, math.inf) if score < old_score: complete[path] = score ranked_complete = sorted( ((score, path) for path, score in complete.items()), key=lambda item: (item[0], item[1]), )[:path_limit] return tuple( CandidateSequenceSelection(candidate_indices=path, total_score=score) for score, path in ranked_complete ) def candidate_chunk_transition_score( previous_result: CandidateGateResult, previous_artifact: ChunkCandidateArtifact, current_result: CandidateGateResult, current_artifact: ChunkCandidateArtifact, *, speaker_weight: float = 1.0, rms_db_weight: float = 0.05, median_f0_weight: float = 0.10, ) -> float: """Return a finite adjacent-chunk cost or ``inf`` for an unsafe edge.""" if previous_result.passed is not True or current_result.passed is not True: return math.inf speaker_w = _finite_float(speaker_weight, minimum=0.0) rms_w = _finite_float(rms_db_weight, minimum=0.0) f0_w = _finite_float(median_f0_weight, minimum=0.0) previous_rms = _finite_float(previous_artifact.rms_db) current_rms = _finite_float(current_artifact.rms_db) if None in (speaker_w, rms_w, f0_w, previous_rms, current_rms): return math.inf previous_embedding = previous_artifact.speaker_embedding current_embedding = current_artifact.speaker_embedding speaker_cost = 0.0 if previous_result.speaker_gate_applied and previous_embedding is None: return math.inf if current_result.speaker_gate_applied and current_embedding is None: return math.inf if previous_embedding is not None and current_embedding is not None: try: speaker_cost = 1.0 - cosine_similarity(previous_embedding, current_embedding) except ValueError: return math.inf assert previous_rms is not None and current_rms is not None rms_cost = abs(previous_rms - current_rms) f0_cost = 0.0 previous_f0 = previous_artifact.median_f0_hz current_f0 = current_artifact.median_f0_hz if previous_f0 is not None or current_f0 is not None: previous_pitch = _finite_float(previous_f0, minimum=1.0) current_pitch = _finite_float(current_f0, minimum=1.0) if previous_pitch is None or current_pitch is None: return math.inf f0_cost = abs(math.log2(current_pitch / previous_pitch)) assert speaker_w is not None and rms_w is not None and f0_w is not None score = speaker_w * speaker_cost + rms_w * rms_cost + f0_w * f0_cost return score if math.isfinite(score) and score >= 0.0 else math.inf def candidate_limit_for_chunk_budget( chunk_count: int, *, max_candidates: int = 20, max_generated_chunks: int = 20, total_text_units: int | None = None, max_generated_text_units: int | None = None, ) -> int: """Return a cap bounded by chunk count and optional text-generation work.""" if isinstance(chunk_count, (bool, np.bool_)): raise ValueError("chunk_count must be a positive integer") try: chunks = int(chunk_count) candidates = int(max_candidates) generated_chunks = int(max_generated_chunks) except (TypeError, ValueError, OverflowError) as error: raise ValueError("candidate budget values must be integers") from error if chunks <= 0: raise ValueError("chunk_count must be a positive integer") if candidates <= 0 or candidates > 20: raise ValueError("max_candidates must be between 1 and 20") if generated_chunks <= 0: raise ValueError("max_generated_chunks must be positive") if chunks > generated_chunks: raise ValueError("one trajectory exceeds the generated-chunk budget") limit = min(candidates, generated_chunks // chunks) if total_text_units is None and max_generated_text_units is None: return limit if total_text_units is None or max_generated_text_units is None: raise ValueError("text-unit budget fields must be provided together") if isinstance(total_text_units, (bool, np.bool_)) or isinstance( max_generated_text_units, (bool, np.bool_), ): raise ValueError("text-unit budget values must be integers") try: units = int(total_text_units) generated_units = int(max_generated_text_units) except (TypeError, ValueError, OverflowError) as error: raise ValueError("text-unit budget values must be integers") from error if units <= 0 or generated_units <= 0: raise ValueError("text-unit budget values must be positive") if units > generated_units: raise ValueError("one trajectory exceeds the generated-text-unit budget") return min(limit, generated_units // units) @dataclass(frozen=True) class _VerifiedTrajectoryCandidate: candidate_index: int seed: int trajectory: Any verification: TrajectoryGateResult joined_output: TrajectoryGateEvidence | None = None def _unwrap_candidate_verification( value: Any, ) -> tuple[TrajectoryGateResult, TrajectoryGateEvidence | None]: if isinstance(value, CandidateVerification): verification = value.verification joined_output = value.joined_output else: verification = value joined_output = None if not isinstance(verification, TrajectoryGateResult): raise TypeError("candidate_verifier must return TrajectoryGateResult") if joined_output is not None and not isinstance( joined_output, TrajectoryGateEvidence, ): raise TypeError("candidate joined-output evidence is invalid") return verification, joined_output def _candidate_attempt_evidence( candidate: _VerifiedTrajectoryCandidate, ) -> CandidateAttemptEvidence: verification = candidate.verification local_results = verification.candidate_results[ :CASCADE_EVIDENCE_MAX_LOCAL_RESULTS ] return CandidateAttemptEvidence( candidate_index=_evidence_int(candidate.candidate_index), seed=_evidence_int( candidate.seed, maximum=REQUEST_SEED_LIMIT + CASCADE_EVIDENCE_MAX_ATTEMPTS, ), trajectory_passed=verification.passed is True, trajectory_score=_evidence_float(verification.score), trajectory_rejection_reasons=_bounded_rejection_reasons( verification.rejection_reasons ), local_result_count=_evidence_int(len(verification.candidate_results)), local_results=tuple( candidate_gate_evidence(result) for result in local_results if isinstance(result, CandidateGateResult) ), joined_output=candidate.joined_output, ) def _cascade_diagnostics( candidates: Sequence[_VerifiedTrajectoryCandidate], sequence_search: SequenceSearchEvidence | None = None, ) -> CascadeDiagnostics: return CascadeDiagnostics( attempts=tuple( _candidate_attempt_evidence(candidate) for candidate in tuple(candidates)[:CASCADE_EVIDENCE_MAX_ATTEMPTS] ), sequence_search=sequence_search, ) def _whole_trajectory_result( candidate: _VerifiedTrajectoryCandidate, attempted_seeds: Sequence[int], chunk_count: int, *, diagnostics: CascadeDiagnostics, ) -> CascadeResult: return CascadeResult( trajectory=candidate.trajectory, verification=candidate.verification, seed=candidate.seed, candidate_index=candidate.candidate_index, attempted_seeds=tuple(attempted_seeds), chunk_candidate_indices=(candidate.candidate_index,) * chunk_count, chunk_seeds=(candidate.seed,) * chunk_count, selection_mode="whole_trajectory", diagnostics=diagnostics, ) @dataclass(frozen=True) class _SequenceFallbackSearchResult: results: tuple[CascadeResult, ...] evidence: SequenceSearchEvidence def _sequence_fallback_search( candidates: Sequence[_VerifiedTrajectoryCandidate], attempted_seeds: Sequence[int], chunk_count: int, *, max_paths: int, max_local_boundary_speaker_drop: float | None = None, ) -> _SequenceFallbackSearchResult: """Rank mixed-seed paths using safe or narrowly recoverable local chunks. A boundary-only local rejection may be made eligible up to the supplied fallback cap. This is intentionally an internal DP representation: the caller must still verify the exactly assembled mixed path with the stricter joined/final gate before any audio can be returned. """ if not candidates: return _SequenceFallbackSearchResult( results=(), evidence=SequenceSearchEvidence( eligible_candidate_counts=(0,) * max(0, chunk_count), finite_transition_counts=(0,) * max(0, chunk_count - 1), ranked_path_count=0, ), ) if chunk_count <= 1: eligible = 0 if chunk_count == 1: for candidate in candidates: results = candidate.verification.candidate_results if len(results) != 1: continue result = results[0] if ( isinstance(result, CandidateGateResult) and result.passed is True and math.isfinite(result.score) and result.score >= 0.0 ): eligible += 1 return _SequenceFallbackSearchResult( results=(), evidence=SequenceSearchEvidence( eligible_candidate_counts=((eligible,) if chunk_count == 1 else ()), finite_transition_counts=(), ranked_path_count=0, ), ) local_scores: list[list[float]] = [[] for _ in range(chunk_count)] usable: list[bool] = [] sequence_results: list[tuple[CandidateGateResult, ...]] = [] for candidate in candidates: verification = candidate.verification try: trajectory_length = len(candidate.trajectory) except TypeError: trajectory_length = -1 candidate_usable = bool( trajectory_length == chunk_count and len(verification.candidate_results) == chunk_count and len(verification.chunk_artifacts) == chunk_count ) usable.append(candidate_usable) adjusted_results: list[CandidateGateResult] = [] for chunk_index in range(chunk_count): score = math.inf if candidate_usable: result = verification.candidate_results[chunk_index] artifact = verification.chunk_artifacts[chunk_index] adjusted_result = _sequence_fallback_candidate_result( result, max_local_boundary_speaker_drop=max_local_boundary_speaker_drop, ) adjusted_results.append(adjusted_result) if ( adjusted_result.passed and math.isfinite(adjusted_result.score) and adjusted_result.score >= 0.0 ): speaker_artifact_valid = True if adjusted_result.speaker_gate_applied: try: speaker_artifact_valid = bool( artifact.speaker_embedding is not None and cosine_similarity( artifact.speaker_embedding, artifact.speaker_embedding, ) >= 1.0 - 1.0e-6 ) except ValueError: speaker_artifact_valid = False if speaker_artifact_valid: score = adjusted_result.score local_scores[chunk_index].append(score) sequence_results.append(tuple(adjusted_results)) transitions: list[list[list[float]]] = [] for chunk_index in range(1, chunk_count): matrix: list[list[float]] = [] for previous_position, previous_candidate in enumerate(candidates): row: list[float] = [] for current_position, current_candidate in enumerate(candidates): score = math.inf if usable[previous_position] and usable[current_position]: score = candidate_chunk_transition_score( sequence_results[previous_position][chunk_index - 1], previous_candidate.verification.chunk_artifacts[chunk_index - 1], sequence_results[current_position][chunk_index], current_candidate.verification.chunk_artifacts[chunk_index], ) row.append(score) matrix.append(row) transitions.append(matrix) selections = select_k_candidate_sequences( local_scores, transitions, max_paths=max_paths, ) output: list[CascadeResult] = [] for rank, selection in enumerate(selections, 1): selected_candidates = tuple( candidates[position] for position in selection.candidate_indices ) selected_trajectory = tuple( candidate.trajectory[chunk_index] for chunk_index, candidate in enumerate(selected_candidates) ) selected_results = tuple( sequence_results[position][chunk_index] for chunk_index, position in enumerate(selection.candidate_indices) ) selected_artifacts = tuple( candidate.verification.chunk_artifacts[chunk_index] for chunk_index, candidate in enumerate(selected_candidates) ) verification = TrajectoryGateResult( passed=True, candidate_results=selected_results, score=selection.total_score, rejection_reasons=(), chunk_artifacts=selected_artifacts, ) output.append( CascadeResult( trajectory=selected_trajectory, verification=verification, seed=None, candidate_index=None, attempted_seeds=tuple(attempted_seeds), chunk_candidate_indices=tuple( candidate.candidate_index for candidate in selected_candidates ), chunk_seeds=tuple( candidate.seed for candidate in selected_candidates ), selection_mode="sequence_dp", sequence_path_rank=rank, ) ) results = tuple(output) evidence = SequenceSearchEvidence( eligible_candidate_counts=tuple( sum(math.isfinite(score) for score in row) for row in local_scores ), finite_transition_counts=tuple( sum( math.isfinite(score) for row in matrix for score in row ) for matrix in transitions ), ranked_path_count=len(results), ) return _SequenceFallbackSearchResult(results=results, evidence=evidence) def _sequence_fallback_results( candidates: Sequence[_VerifiedTrajectoryCandidate], attempted_seeds: Sequence[int], chunk_count: int, *, max_paths: int, max_local_boundary_speaker_drop: float | None = None, ) -> tuple[CascadeResult, ...]: """Compatibility wrapper returning only ranked sequence results.""" return _sequence_fallback_search( candidates, attempted_seeds, chunk_count, max_paths=max_paths, max_local_boundary_speaker_drop=max_local_boundary_speaker_drop, ).results def _sequence_fallback_candidate_result( result: CandidateGateResult, *, max_local_boundary_speaker_drop: float | None, ) -> CandidateGateResult: """Return a DP-safe view of one local result or an unchanged hard reject.""" if result.passed is True and math.isfinite(result.score) and result.score >= 0.0: return result limit = ( None if max_local_boundary_speaker_drop is None else _finite_float( max_local_boundary_speaker_drop, minimum=0.0, maximum=1.0, ) ) if ( limit is None or result.rejection_reasons != ("boundary_speaker_drop",) or result.speaker_gate_applied is not True or result.comparison.passed is not True ): return result similarity = _finite_float( result.speaker_similarity, minimum=-1.0, maximum=1.0, ) boundary_drop = _finite_float(result.boundary_speaker_drop, minimum=0.0) cer = _finite_float(result.comparison.cer, minimum=0.0) if ( similarity is None or boundary_drop is None or boundary_drop > limit + 1.0e-12 or cer is None ): return result score = ( cer + SEQUENCE_FALLBACK_SPEAKER_WEIGHT * (1.0 - similarity) + SEQUENCE_FALLBACK_BOUNDARY_WEIGHT * boundary_drop ) if not math.isfinite(score) or score < 0.0: return result return replace( result, passed=True, score=score, rejection_reasons=(), ) def _preferred_speaker_verification( verification: TrajectoryGateResult, *, min_similarity: float, max_boundary_drop: float, ) -> bool: if verification.passed is not True or not math.isfinite(verification.score): return False for result in verification.candidate_results: if not result.speaker_gate_applied: continue similarity = _finite_float(result.speaker_similarity, minimum=-1.0, maximum=1.0) boundary_drop = _finite_float(result.boundary_speaker_drop, minimum=0.0) if ( similarity is None or boundary_drop is None or similarity < min_similarity or boundary_drop > max_boundary_drop ): return False return True def run_adaptive_cascade( chunks: Sequence[str], root_seed: int, candidate_generator: Callable[[tuple[str, ...], int], Any], candidate_verifier: Callable[ [Any, tuple[str, ...], int], TrajectoryGateResult | CandidateVerification, ], *, initial_candidates: int = 1, max_candidates: int = 5, preferred_min_speaker_similarity: float = 0.25, preferred_max_boundary_speaker_drop: float = 0.05, sequence_final_verifier: ( Callable[[CascadeResult, tuple[str, ...]], TrajectoryGateResult] | None ) = None, max_sequence_paths: int = 3, sequence_fallback_max_local_boundary_speaker_drop: float | None = None, ) -> CascadeResult: """Run a deterministic 1-to-5-to-10-to-15-to-20 fail-closed cascade. The generator is called once per trajectory with ``root_seed + offset``. It receives all chunks in one call, making the shared per-trajectory seed contract explicit and preventing accidental per-chunk seed drift. """ try: chunk_tuple = tuple(str(chunk) for chunk in chunks) base_seed = int(root_seed) first_stage = int(initial_candidates) limit = int(max_candidates) except (TypeError, ValueError, OverflowError) as error: raise ValueError("invalid adaptive cascade arguments") from error if not chunk_tuple or any(not chunk for chunk in chunk_tuple): raise ValueError("adaptive cascade requires non-empty text chunks") if first_stage != 1: raise ValueError("online adaptive cascade must start with exactly one candidate") if limit < first_stage or limit > ADAPTIVE_CASCADE_STAGE_LIMITS[-1]: raise ValueError("adaptive cascade supports between 1 and 20 candidates") if isinstance(max_sequence_paths, (bool, np.bool_)): raise ValueError("max_sequence_paths must be an integer between 1 and 3") try: sequence_path_limit = operator.index(max_sequence_paths) except (TypeError, ValueError, OverflowError) as error: raise ValueError( "max_sequence_paths must be an integer between 1 and 3" ) from error if not 1 <= sequence_path_limit <= 3: raise ValueError("max_sequence_paths must be an integer between 1 and 3") if sequence_final_verifier is not None and not callable(sequence_final_verifier): raise ValueError("sequence_final_verifier must be callable") sequence_fallback_boundary = None if sequence_fallback_max_local_boundary_speaker_drop is not None: sequence_fallback_boundary = _finite_float( sequence_fallback_max_local_boundary_speaker_drop, minimum=0.0, maximum=1.0, ) if sequence_fallback_boundary is None: raise ValueError("sequence fallback boundary threshold must be finite") if sequence_final_verifier is None: raise ValueError( "sequence fallback boundary relaxation requires a final verifier" ) preferred_similarity = _finite_float( preferred_min_speaker_similarity, minimum=-1.0, maximum=1.0, ) preferred_boundary = _finite_float( preferred_max_boundary_speaker_drop, minimum=0.0, ) if preferred_similarity is None or preferred_boundary is None: raise ValueError("preferred speaker thresholds must be finite") attempted_seeds: list[int] = [] candidates: list[_VerifiedTrajectoryCandidate] = [] sequence_search_evidence: SequenceSearchEvidence | None = None first_seed = base_seed first_trajectory = candidate_generator(chunk_tuple, first_seed) first_verification, first_joined_output = _unwrap_candidate_verification( candidate_verifier(first_trajectory, chunk_tuple, first_seed) ) attempted_seeds.append(first_seed) first_candidate = _VerifiedTrajectoryCandidate( candidate_index=0, seed=first_seed, trajectory=first_trajectory, verification=first_verification, joined_output=first_joined_output, ) candidates.append(first_candidate) if ( first_verification.passed and math.isfinite(first_verification.score) and ( limit == 1 or _preferred_speaker_verification( first_verification, min_similarity=preferred_similarity, max_boundary_drop=preferred_boundary, ) ) ): return _whole_trajectory_result( first_candidate, attempted_seeds, len(chunk_tuple), diagnostics=_cascade_diagnostics(candidates), ) stages = list( dict.fromkeys( min(stage_limit, limit) for stage_limit in ADAPTIVE_CASCADE_STAGE_LIMITS[1:] ) ) next_candidate = first_stage for stage_size in stages: for candidate_index in range(next_candidate, stage_size): seed = base_seed + candidate_index trajectory = candidate_generator(chunk_tuple, seed) verification, joined_output = _unwrap_candidate_verification( candidate_verifier(trajectory, chunk_tuple, seed) ) attempted_seeds.append(seed) candidates.append( _VerifiedTrajectoryCandidate( candidate_index=candidate_index, seed=seed, trajectory=trajectory, verification=verification, joined_output=joined_output, ) ) qualified = [ candidate for candidate in candidates if candidate.verification.passed and math.isfinite(candidate.verification.score) ] final_stage = stage_size == limit if qualified: preferred_qualified = [ candidate for candidate in qualified if _preferred_speaker_verification( candidate.verification, min_similarity=preferred_similarity, max_boundary_drop=preferred_boundary, ) ] selectable = qualified if final_stage else preferred_qualified if selectable: selected_whole = min( selectable, key=lambda candidate: ( candidate.verification.score, candidate.candidate_index, ), ) return _whole_trajectory_result( selected_whole, attempted_seeds, len(chunk_tuple), diagnostics=_cascade_diagnostics(candidates), ) next_candidate = stage_size continue # With a final-aware callback, defer DP until all whole-trajectory # candidates in the request budget have been exhausted. This keeps # whole trajectories globally preferred and bounds joined checks to # at most ``max_sequence_paths`` once per request. if sequence_final_verifier is not None and not final_stage: next_candidate = stage_size continue sequence_search = _sequence_fallback_search( candidates, attempted_seeds, len(chunk_tuple), max_paths=(sequence_path_limit if sequence_final_verifier else 1), max_local_boundary_speaker_drop=sequence_fallback_boundary, ) sequence_results = sequence_search.results sequence_search_evidence = sequence_search.evidence if sequence_final_verifier is None: sequence_result = sequence_results[0] if sequence_results else None if sequence_result is not None and ( final_stage or _preferred_speaker_verification( sequence_result.verification, min_similarity=preferred_similarity, max_boundary_drop=preferred_boundary, ) ): return replace( sequence_result, diagnostics=_cascade_diagnostics( candidates, sequence_search_evidence, ), ) else: checked_paths: list[SequencePathEvidence] = [] for checked_count, sequence_result in enumerate(sequence_results, 1): try: final_verification = sequence_final_verifier( sequence_result, chunk_tuple, ) except Exception as error: raise RuntimeError( "sequence final verification failed; refusing unverified audio" ) from error if not isinstance(final_verification, TrajectoryGateResult): raise RuntimeError( "sequence final verifier returned an invalid result" ) checked_paths.append( SequencePathEvidence( rank=_evidence_int(sequence_result.sequence_path_rank), chunk_candidate_indices=tuple( _evidence_int(index) for index in sequence_result.chunk_candidate_indices ), chunk_seeds=tuple( _evidence_int( seed, maximum=( REQUEST_SEED_LIMIT + CASCADE_EVIDENCE_MAX_ATTEMPTS ), ) for seed in sequence_result.chunk_seeds ), final_output=trajectory_gate_evidence( final_verification ), ) ) sequence_search_evidence = replace( sequence_search.evidence, checked_paths=tuple(checked_paths), ) if ( final_verification.passed is True and math.isfinite(final_verification.score) and len(final_verification.candidate_results) == 1 and isinstance( final_verification.candidate_results[0], CandidateGateResult, ) and final_verification.candidate_results[0].passed is True and math.isfinite( final_verification.candidate_results[0].score ) and not final_verification.rejection_reasons ): return CascadeResult( trajectory=sequence_result.trajectory, verification=sequence_result.verification, seed=sequence_result.seed, candidate_index=sequence_result.candidate_index, attempted_seeds=sequence_result.attempted_seeds, chunk_candidate_indices=( sequence_result.chunk_candidate_indices ), chunk_seeds=sequence_result.chunk_seeds, selection_mode=sequence_result.selection_mode, sequence_path_rank=sequence_result.sequence_path_rank, sequence_paths_checked=checked_count, diagnostics=_cascade_diagnostics( candidates, sequence_search_evidence, ), ) next_candidate = stage_size raise NoQualifiedCandidateError( f"no verified TTS trajectory after {len(attempted_seeds)} candidates", diagnostics=_cascade_diagnostics( candidates, sequence_search_evidence, ), )