use serde::{Deserialize, Serialize};
use std::{collections::BTreeSet, error::Error, fmt};
pub use kcode_speaker_v3_llm_protocol::{
GEMINI_FEATURE_PROMPT_ONE, GEMINI_FEATURE_PROMPT_ONE_REVISION, GEMINI_FEATURE_PROMPT_REVISIONS,
GEMINI_FEATURE_PROMPT_THREE, GEMINI_FEATURE_PROMPT_THREE_REVISION, GEMINI_FEATURE_PROMPT_TWO,
GEMINI_FEATURE_PROMPT_TWO_REVISION, GEMINI_TRANSCRIPT_PROMPT,
GEMINI_TRANSCRIPT_PROMPT_REVISION, GPT_STRUCTURING_PROMPT, GPT_STRUCTURING_PROMPT_REVISION,
};
pub use kcode_speaker_v3_schema::{
FEATURE_NAMES, FEATURE_SCHEMA_REVISION, FeatureVector24, LocalSpeakerLabel,
MAX_AUDIO_DURATION_MS, OGG_MEDIA_TYPE, OggAudioMetadata, StructuredAnalysis, StructuredSpeaker,
ValidationError, VocalGenderPresentation,
};
#[cfg(any(feature = "providers", test))]
use kcode_speaker_v3_gemini_analysis::GeminiFeatureProgress;
#[cfg(any(feature = "providers", test))]
use kcode_speaker_v3_llm_protocol::SpeakerFeatureEvidence;
#[cfg(any(feature = "providers", test))]
use std::{future::Future, pin::Pin};
const GEMINI_MODEL_ID: &str = "gemini-3.1-pro-preview";
const TERRA_MODEL_ID: &str = "gpt-5.6-terra";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct GeminiCohort {
pub model_id: String,
pub transcript_prompt_revision: String,
pub feature_prompt_revisions: [String; 3],
pub feature_schema_revision: String,
}
impl GeminiCohort {
pub fn new(model_id: impl Into<String>) -> Self {
Self {
model_id: model_id.into(),
transcript_prompt_revision: GEMINI_TRANSCRIPT_PROMPT_REVISION.into(),
feature_prompt_revisions: GEMINI_FEATURE_PROMPT_REVISIONS.map(str::to_owned),
feature_schema_revision: FEATURE_SCHEMA_REVISION.into(),
}
}
pub fn validate(&self) -> Result<(), ValidationError> {
validate_text(&self.model_id, "gemini_model_id")?;
validate_text(
&self.transcript_prompt_revision,
"transcript_prompt_revision",
)?;
for revision in &self.feature_prompt_revisions {
validate_text(revision, "feature_prompt_revision")?;
}
validate_text(&self.feature_schema_revision, "feature_schema_revision")
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct StructurerProvenance {
pub model_id: String,
pub prompt_revision: String,
}
impl StructurerProvenance {
pub fn new(model_id: impl Into<String>) -> Self {
Self {
model_id: model_id.into(),
prompt_revision: GPT_STRUCTURING_PROMPT_REVISION.into(),
}
}
pub fn validate(&self) -> Result<(), ValidationError> {
validate_text(&self.model_id, "structurer_model_id")?;
validate_text(&self.prompt_revision, "structurer_prompt_revision")
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AnalysisEnvelope {
pub audio: OggAudioMetadata,
pub analysis: StructuredAnalysis,
pub gemini: GeminiCohort,
pub structurer: StructurerProvenance,
}
impl AnalysisEnvelope {
pub fn validate(&self) -> Result<(), ValidationError> {
self.audio.validate()?;
self.analysis.validate()?;
self.gemini.validate()?;
self.structurer.validate()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AnalysisStage {
Transcript,
SpeakerLabels,
SpeakerFeatures,
Structuring,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AnalysisJob {
Transcript,
SpeakerLabels,
SpeakerFeature {
speaker: LocalSpeakerLabel,
packet: u8,
},
Structuring,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AnalysisProgress {
JobStarted { sequence: u64, job: AnalysisJob },
JobSucceeded { sequence: u64 },
JobFailed { sequence: u64, error: String },
StageCompleted { stage: AnalysisStage },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum AnalysisError {
Input(String),
Progress(String),
GeminiTranscript(String),
TerraLabels(String),
GeminiCache(String),
GeminiFeature {
speaker: LocalSpeakerLabel,
packet: u8,
message: String,
},
TerraStructuring(String),
TranscriptMismatch,
SpeakerSetMismatch,
}
impl fmt::Display for AnalysisError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Input(message) => write!(formatter, "invalid input: {message}"),
Self::Progress(message) => write!(formatter, "progress reporting failed: {message}"),
Self::GeminiTranscript(message) => {
write!(formatter, "Gemini transcript failed: {message}")
}
Self::TerraLabels(message) => {
write!(
formatter,
"Terra speaker-label extraction failed: {message}"
)
}
Self::GeminiCache(message) => {
write!(formatter, "Gemini feature cache creation failed: {message}")
}
Self::GeminiFeature {
speaker,
packet,
message,
} => write!(
formatter,
"Gemini feature call failed for {speaker}, packet {packet}: {message}"
),
Self::TerraStructuring(message) => {
write!(formatter, "Terra final structuring failed: {message}")
}
Self::TranscriptMismatch => {
formatter.write_str("Terra returned a different transcript")
}
Self::SpeakerSetMismatch => {
formatter.write_str("Terra returned a different speaker set")
}
}
}
}
impl Error for AnalysisError {}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct ExecutedAnalysis {
pub envelope: AnalysisEnvelope,
pub label_extractor: StructurerProvenance,
}
#[cfg(feature = "providers")]
pub struct Analyzer {
operations: ProviderOperations,
}
#[cfg(feature = "providers")]
impl Analyzer {
pub fn new(
gemini: kcode_gemini_3_1_pro::Gemini31Pro,
terra: kcode_codex_terra::CodexTerra,
) -> Self {
Self {
operations: ProviderOperations {
gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis::new(gemini),
terra: kcode_speaker_v3_terra_analysis::TerraAnalysis::new(terra),
},
}
}
pub async fn analyze_ogg_with_progress<F>(
&self,
bytes: &[u8],
report: F,
) -> Result<ExecutedAnalysis, AnalysisError>
where
F: FnMut(AnalysisProgress) -> Result<(), String>,
{
execute_strict(&self.operations, bytes, report).await
}
pub async fn analyze_ogg(
&self,
bytes: &[u8],
duration_ms: u64,
filename: Option<String>,
) -> Result<ExecutedAnalysis, AnalysisError> {
execute_legacy(&self.operations, bytes, duration_ms, filename).await
}
}
fn validate_text(value: &str, field: &'static str) -> Result<(), ValidationError> {
(!value.trim().is_empty())
.then_some(())
.ok_or(ValidationError::Blank(field))
}
#[cfg(any(feature = "providers", test))]
type AnalysisFuture<'a, T> = Pin<Box<dyn Future<Output = T> + 'a>>;
#[cfg(any(feature = "providers", test))]
trait AnalysisOperations: Sync {
fn transcript<'a>(
&'a self,
audio: &'a [u8],
) -> AnalysisFuture<'a, Result<String, AnalysisError>>;
fn speaker_labels<'a>(
&'a self,
transcript: &'a str,
) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>>;
fn feature_evidence_with_progress<'a>(
&'a self,
audio: &'a [u8],
transcript: &'a str,
labels: &'a [LocalSpeakerLabel],
report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>>;
fn structured_analysis<'a>(
&'a self,
transcript: &'a str,
evidence: Vec<SpeakerFeatureEvidence>,
) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>>;
}
#[cfg(any(feature = "providers", test))]
struct ProviderOperations {
gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis,
terra: kcode_speaker_v3_terra_analysis::TerraAnalysis,
}
#[cfg(any(feature = "providers", test))]
impl AnalysisOperations for ProviderOperations {
fn transcript<'a>(
&'a self,
audio: &'a [u8],
) -> AnalysisFuture<'a, Result<String, AnalysisError>> {
Box::pin(async move {
self.gemini.transcript(audio).await.map_err(|error| match error {
kcode_speaker_v3_gemini_analysis::GeminiTranscriptError::Provider(message)
| kcode_speaker_v3_gemini_analysis::GeminiTranscriptError::Protocol(message) => {
AnalysisError::GeminiTranscript(message)
}
})
})
}
fn speaker_labels<'a>(
&'a self,
transcript: &'a str,
) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>> {
Box::pin(async move {
self.terra
.speaker_labels(transcript)
.await
.map_err(|error| match error {
kcode_speaker_v3_terra_analysis::TerraAnalysisError::Protocol(message)
| kcode_speaker_v3_terra_analysis::TerraAnalysisError::Provider(message) => {
AnalysisError::TerraLabels(message)
}
})
})
}
fn feature_evidence_with_progress<'a>(
&'a self,
audio: &'a [u8],
transcript: &'a str,
labels: &'a [LocalSpeakerLabel],
report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>> {
Box::pin(async move {
self.gemini
.feature_evidence_with_progress(audio, transcript, labels, report)
.await
.map_err(|error| match error {
kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Cache(message) => {
AnalysisError::GeminiCache(message)
}
kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Progress(message) => {
AnalysisError::Progress(message)
}
kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Feature {
speaker,
packet,
message,
} => AnalysisError::GeminiFeature {
speaker,
packet,
message,
},
})
})
}
fn structured_analysis<'a>(
&'a self,
transcript: &'a str,
evidence: Vec<SpeakerFeatureEvidence>,
) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>> {
Box::pin(async move {
self.terra
.structured_analysis(transcript, evidence)
.await
.map_err(|error| match error {
kcode_speaker_v3_terra_analysis::TerraAnalysisError::Protocol(message)
| kcode_speaker_v3_terra_analysis::TerraAnalysisError::Provider(message) => {
AnalysisError::TerraStructuring(message)
}
})
})
}
}
#[cfg(any(feature = "providers", test))]
async fn execute_strict<O, F>(
operations: &O,
bytes: &[u8],
report: F,
) -> Result<ExecutedAnalysis, AnalysisError>
where
O: AnalysisOperations,
F: FnMut(AnalysisProgress) -> Result<(), String>,
{
let audio = OggAudioMetadata::from_ogg_bytes(bytes)
.map_err(|error| AnalysisError::Input(error.to_string()))?;
execute_admitted(operations, bytes, audio, report).await
}
#[cfg(any(feature = "providers", test))]
async fn execute_legacy<O: AnalysisOperations>(
operations: &O,
bytes: &[u8],
duration_ms: u64,
filename: Option<String>,
) -> Result<ExecutedAnalysis, AnalysisError> {
let audio = OggAudioMetadata::from_bytes(bytes, duration_ms, filename)
.map_err(|error| AnalysisError::Input(error.to_string()))?;
execute_admitted(operations, bytes, audio, |_| Ok(())).await
}
#[cfg(any(feature = "providers", test))]
async fn execute_admitted<O, F>(
operations: &O,
bytes: &[u8],
audio: OggAudioMetadata,
mut report: F,
) -> Result<ExecutedAnalysis, AnalysisError>
where
O: AnalysisOperations,
F: FnMut(AnalysisProgress) -> Result<(), String>,
{
report_progress(
&mut report,
AnalysisProgress::JobStarted {
sequence: 1,
job: AnalysisJob::Transcript,
},
)?;
let transcript = match operations.transcript(bytes).await {
Ok(transcript) => transcript,
Err(error) => {
report_leaf_failure(&mut report, 1, &error)?;
return Err(error);
}
};
report_progress(&mut report, AnalysisProgress::JobSucceeded { sequence: 1 })?;
report_progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Transcript,
},
)?;
report_progress(
&mut report,
AnalysisProgress::JobStarted {
sequence: 2,
job: AnalysisJob::SpeakerLabels,
},
)?;
let labels = match operations.speaker_labels(&transcript).await {
Ok(labels) => labels,
Err(error) => {
report_leaf_failure(&mut report, 2, &error)?;
return Err(error);
}
};
report_progress(&mut report, AnalysisProgress::JobSucceeded { sequence: 2 })?;
report_progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerLabels,
},
)?;
let structuring_sequence = structuring_sequence(labels.len())?;
let evidence = {
let mut feature_report = |progress| {
let progress = map_feature_progress(&labels, progress)?;
report(progress)
};
operations
.feature_evidence_with_progress(bytes, &transcript, &labels, &mut feature_report)
.await?
};
report_progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerFeatures,
},
)?;
report_progress(
&mut report,
AnalysisProgress::JobStarted {
sequence: structuring_sequence,
job: AnalysisJob::Structuring,
},
)?;
let analysis = match operations.structured_analysis(&transcript, evidence).await {
Ok(analysis) => analysis,
Err(error) => {
report_leaf_failure(&mut report, structuring_sequence, &error)?;
return Err(error);
}
};
report_progress(
&mut report,
AnalysisProgress::JobSucceeded {
sequence: structuring_sequence,
},
)?;
if analysis.transcript != transcript {
return Err(AnalysisError::TranscriptMismatch);
}
let expected_speakers = labels.iter().copied().collect::<BTreeSet<_>>();
let returned_speakers = analysis
.speakers
.iter()
.map(|speaker| speaker.speaker)
.collect::<BTreeSet<_>>();
if expected_speakers != returned_speakers {
return Err(AnalysisError::SpeakerSetMismatch);
}
let envelope = AnalysisEnvelope {
audio,
analysis,
gemini: GeminiCohort::new(GEMINI_MODEL_ID),
structurer: StructurerProvenance::new(TERRA_MODEL_ID),
};
envelope
.validate()
.map_err(|error| AnalysisError::TerraStructuring(error.to_string()))?;
let label_extractor = StructurerProvenance {
model_id: TERRA_MODEL_ID.into(),
prompt_revision: kcode_speaker_v3_llm_protocol::TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
};
label_extractor
.validate()
.map_err(|error| AnalysisError::TerraLabels(error.to_string()))?;
let executed = ExecutedAnalysis {
envelope,
label_extractor,
};
report_progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring,
},
)?;
Ok(executed)
}
#[cfg(any(feature = "providers", test))]
fn report_progress<F>(report: &mut F, progress: AnalysisProgress) -> Result<(), AnalysisError>
where
F: FnMut(AnalysisProgress) -> Result<(), String>,
{
report(progress).map_err(AnalysisError::Progress)
}
#[cfg(any(feature = "providers", test))]
fn report_leaf_failure<F>(
report: &mut F,
sequence: u64,
error: &AnalysisError,
) -> Result<(), AnalysisError>
where
F: FnMut(AnalysisProgress) -> Result<(), String>,
{
report_progress(
report,
AnalysisProgress::JobFailed {
sequence,
error: error.to_string(),
},
)
}
#[cfg(any(feature = "providers", test))]
fn structuring_sequence(label_count: usize) -> Result<u64, AnalysisError> {
u64::try_from(label_count)
.ok()
.and_then(|count| count.checked_mul(3))
.and_then(|count| count.checked_add(3))
.ok_or_else(sequence_error)
}
#[cfg(any(feature = "providers", test))]
fn feature_sequence(
labels: &[LocalSpeakerLabel],
speaker: LocalSpeakerLabel,
packet: u8,
) -> Result<u64, String> {
if !(1..=3).contains(&packet) {
return Err(format!(
"Gemini reported invalid feature packet {packet} for {speaker}"
));
}
let label_index = labels
.iter()
.position(|candidate| *candidate == speaker)
.ok_or_else(|| format!("Gemini reported an unknown feature speaker {speaker}"))?;
u64::try_from(label_index)
.ok()
.and_then(|index| index.checked_mul(3))
.and_then(|index| index.checked_add(u64::from(packet)))
.and_then(|index| index.checked_add(2))
.ok_or_else(|| sequence_error().to_string())
}
#[cfg(any(feature = "providers", test))]
fn sequence_error() -> AnalysisError {
AnalysisError::Progress("analysis job sequence overflow".into())
}
#[cfg(any(feature = "providers", test))]
fn map_feature_progress(
labels: &[LocalSpeakerLabel],
progress: GeminiFeatureProgress,
) -> Result<AnalysisProgress, String> {
match progress {
GeminiFeatureProgress::Started { speaker, packet } => Ok(AnalysisProgress::JobStarted {
sequence: feature_sequence(labels, speaker, packet)?,
job: AnalysisJob::SpeakerFeature { speaker, packet },
}),
GeminiFeatureProgress::Succeeded { speaker, packet } => {
Ok(AnalysisProgress::JobSucceeded {
sequence: feature_sequence(labels, speaker, packet)?,
})
}
GeminiFeatureProgress::Failed {
speaker,
packet,
error,
} => {
let displayed_error = AnalysisError::GeminiFeature {
speaker,
packet,
message: error,
}
.to_string();
Ok(AnalysisProgress::JobFailed {
sequence: feature_sequence(labels, speaker, packet)?,
error: displayed_error,
})
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::{executor::block_on, future::poll_fn, join};
use std::{
sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
task::Poll,
time::Instant,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum FailureStage {
Transcript,
Labels,
Cache,
Features,
Final,
}
#[derive(Clone)]
struct FakeConfig {
transcript: String,
labels: Vec<LocalSpeakerLabel>,
analysis: StructuredAnalysis,
failure: Option<FailureStage>,
feature_terminal_order: Option<Vec<(LocalSpeakerLabel, u8)>>,
feature_failures: Vec<(LocalSpeakerLabel, u8, String)>,
wait_for: Option<Arc<AtomicBool>>,
mark_complete: Option<Arc<AtomicBool>>,
}
struct FakeState {
config: FakeConfig,
transcript_calls: AtomicUsize,
label_calls: AtomicUsize,
feature_calls: AtomicUsize,
final_calls: AtomicUsize,
calls: Mutex<Vec<&'static str>>,
}
#[derive(Clone)]
struct FakeOperations {
state: Arc<FakeState>,
}
impl FakeOperations {
fn successful(speaker_count: u32) -> Self {
let transcript = "[high] Speaker 1: exact transcript".to_owned();
Self::from_config(FakeConfig {
labels: (1..=speaker_count).map(label).collect(),
analysis: structured_analysis(&transcript, speaker_count),
transcript,
failure: None,
feature_terminal_order: None,
feature_failures: Vec::new(),
wait_for: None,
mark_complete: None,
})
}
fn from_config(config: FakeConfig) -> Self {
Self {
state: Arc::new(FakeState {
config,
transcript_calls: AtomicUsize::new(0),
label_calls: AtomicUsize::new(0),
feature_calls: AtomicUsize::new(0),
final_calls: AtomicUsize::new(0),
calls: Mutex::new(Vec::new()),
}),
}
}
fn with_config(&self, update: impl FnOnce(&mut FakeConfig)) -> Self {
let mut config = self.state.config.clone();
update(&mut config);
Self::from_config(config)
}
fn feature_failure(&self, speaker: LocalSpeakerLabel, packet: u8) -> Option<String> {
if self.state.config.failure == Some(FailureStage::Features)
&& speaker == label(1)
&& packet == 2
{
return Some("features".into());
}
self.state
.config
.feature_failures
.iter()
.find(|(candidate, candidate_packet, _)| {
*candidate == speaker && *candidate_packet == packet
})
.map(|(_, _, message)| message.clone())
}
}
impl AnalysisOperations for FakeOperations {
fn transcript<'a>(
&'a self,
_audio: &'a [u8],
) -> AnalysisFuture<'a, Result<String, AnalysisError>> {
Box::pin(async move {
self.state.transcript_calls.fetch_add(1, Ordering::SeqCst);
self.state.calls.lock().unwrap().push("transcript");
if let Some(wait_for) = &self.state.config.wait_for {
poll_fn(|context| {
if wait_for.load(Ordering::SeqCst) {
Poll::Ready(())
} else {
context.waker().wake_by_ref();
Poll::Pending
}
})
.await;
}
if self.state.config.failure == Some(FailureStage::Transcript) {
return Err(AnalysisError::GeminiTranscript("transcript".into()));
}
Ok(self.state.config.transcript.clone())
})
}
fn speaker_labels<'a>(
&'a self,
_transcript: &'a str,
) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>> {
Box::pin(async move {
self.state.label_calls.fetch_add(1, Ordering::SeqCst);
self.state.calls.lock().unwrap().push("labels");
if self.state.config.failure == Some(FailureStage::Labels) {
return Err(AnalysisError::TerraLabels("labels".into()));
}
Ok(self.state.config.labels.clone())
})
}
fn feature_evidence_with_progress<'a>(
&'a self,
_audio: &'a [u8],
_transcript: &'a str,
labels: &'a [LocalSpeakerLabel],
report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>> {
Box::pin(async move {
self.state.feature_calls.fetch_add(1, Ordering::SeqCst);
self.state.calls.lock().unwrap().push("features");
if self.state.config.failure == Some(FailureStage::Cache) {
return Err(AnalysisError::GeminiCache("cache".into()));
}
let jobs = labels
.iter()
.copied()
.flat_map(|speaker| (1..=3).map(move |packet| (speaker, packet)))
.collect::<Vec<_>>();
for &(speaker, packet) in &jobs {
report(GeminiFeatureProgress::Started { speaker, packet })
.map_err(AnalysisError::Progress)?;
}
let terminal_order = self
.state
.config
.feature_terminal_order
.clone()
.unwrap_or_else(|| jobs.clone());
for (speaker, packet) in terminal_order {
if let Some(error) = self.feature_failure(speaker, packet) {
report(GeminiFeatureProgress::Failed {
speaker,
packet,
error,
})
.map_err(AnalysisError::Progress)?;
} else {
report(GeminiFeatureProgress::Succeeded { speaker, packet })
.map_err(AnalysisError::Progress)?;
}
}
if let Some((speaker, packet, message)) =
jobs.iter().find_map(|&(speaker, packet)| {
self.feature_failure(speaker, packet)
.map(|message| (speaker, packet, message))
})
{
return Err(AnalysisError::GeminiFeature {
speaker,
packet,
message,
});
}
labels
.iter()
.copied()
.map(|speaker| {
SpeakerFeatureEvidence::new(
speaker,
format!("{speaker} packet 1"),
format!("{speaker} packet 2"),
format!("{speaker} packet 3"),
)
.map_err(|error| AnalysisError::GeminiFeature {
speaker,
packet: 1,
message: error.to_string(),
})
})
.collect()
})
}
fn structured_analysis<'a>(
&'a self,
_transcript: &'a str,
_evidence: Vec<SpeakerFeatureEvidence>,
) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>> {
Box::pin(async move {
self.state.final_calls.fetch_add(1, Ordering::SeqCst);
self.state.calls.lock().unwrap().push("final");
if self.state.config.failure == Some(FailureStage::Final) {
return Err(AnalysisError::TerraStructuring("final".into()));
}
if let Some(mark_complete) = &self.state.config.mark_complete {
mark_complete.store(true, Ordering::SeqCst);
}
Ok(self.state.config.analysis.clone())
})
}
}
fn label(number: u32) -> LocalSpeakerLabel {
LocalSpeakerLabel::new(number).unwrap()
}
fn structured_analysis(transcript: &str, speaker_count: u32) -> StructuredAnalysis {
StructuredAnalysis {
transcript: transcript.into(),
speakers: (1..=speaker_count)
.map(|number| StructuredSpeaker {
speaker: label(number),
language: "English".into(),
features: FeatureVector24::default(),
features_usable_for_training: false,
})
.collect(),
}
}
fn legacy_ogg() -> Vec<u8> {
let mut bytes = vec![0; 28];
bytes[..4].copy_from_slice(b"OggS");
bytes[4] = 0;
bytes[26] = 1;
bytes[27] = 0;
bytes
}
fn strict_ogg(samples: u64) -> Vec<u8> {
let pre_skip = 312_u16;
let mut head = b"OpusHead".to_vec();
head.push(1);
head.push(1);
head.extend_from_slice(&pre_skip.to_le_bytes());
head.extend_from_slice(&48_000_u32.to_le_bytes());
head.extend_from_slice(&0_i16.to_le_bytes());
head.push(0);
let mut tags = b"OpusTags".to_vec();
tags.extend_from_slice(&0_u32.to_le_bytes());
tags.extend_from_slice(&0_u32.to_le_bytes());
let mut bytes = ogg_page(0x02, 0, 0, &head);
bytes.extend_from_slice(&ogg_page(0, 0, 1, &tags));
bytes.extend_from_slice(&ogg_page(
0x04,
u64::from(pre_skip) + samples,
2,
&[0xf8, 0xff, 0xfe],
));
bytes
}
fn ogg_page(header_type: u8, granule: u64, sequence: u32, payload: &[u8]) -> Vec<u8> {
let payload_length = u8::try_from(payload.len()).unwrap();
let mut page = Vec::with_capacity(28 + payload.len());
page.extend_from_slice(b"OggS");
page.push(0);
page.push(header_type);
page.extend_from_slice(&granule.to_le_bytes());
page.extend_from_slice(&0x534b_5633_u32.to_le_bytes());
page.extend_from_slice(&sequence.to_le_bytes());
page.extend_from_slice(&0_u32.to_le_bytes());
page.push(1);
page.push(payload_length);
page.extend_from_slice(payload);
let checksum = ogg_crc(&page);
page[22..26].copy_from_slice(&checksum.to_le_bytes());
page
}
fn ogg_crc(bytes: &[u8]) -> u32 {
let mut crc = 0_u32;
for &byte in bytes {
crc ^= u32::from(byte) << 24;
for _ in 0..8 {
crc = if crc & 0x8000_0000 == 0 {
crc << 1
} else {
(crc << 1) ^ 0x04c1_1db7
};
}
}
crc
}
fn run_with_progress(
operations: &FakeOperations,
) -> (
Result<ExecutedAnalysis, AnalysisError>,
Vec<AnalysisProgress>,
) {
let audio = strict_ogg(48);
let mut events = Vec::new();
let result = block_on(execute_strict(operations, &audio, |event| {
events.push(event);
Ok(())
}));
(result, events)
}
#[test]
fn strict_admission_precedes_reports_and_operations() {
for bytes in [b"bad".to_vec(), strict_ogg(7_200_001)] {
let operations = FakeOperations::successful(1);
let reports = AtomicUsize::new(0);
assert!(matches!(
block_on(execute_strict(&operations, &bytes, |_| {
reports.fetch_add(1, Ordering::SeqCst);
Ok(())
})),
Err(AnalysisError::Input(_))
));
assert_eq!(reports.load(Ordering::SeqCst), 0);
assert!(operations.state.calls.lock().unwrap().is_empty());
}
}
#[test]
fn zero_label_sequences_are_exact() {
let operations = FakeOperations::successful(0);
let (result, events) = run_with_progress(&operations);
result.unwrap();
assert_eq!(
events,
vec![
AnalysisProgress::JobStarted {
sequence: 1,
job: AnalysisJob::Transcript,
},
AnalysisProgress::JobSucceeded { sequence: 1 },
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Transcript,
},
AnalysisProgress::JobStarted {
sequence: 2,
job: AnalysisJob::SpeakerLabels,
},
AnalysisProgress::JobSucceeded { sequence: 2 },
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerLabels,
},
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerFeatures,
},
AnalysisProgress::JobStarted {
sequence: 3,
job: AnalysisJob::Structuring,
},
AnalysisProgress::JobSucceeded { sequence: 3 },
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring,
},
]
);
}
#[test]
fn multiple_label_sequences_are_exact() {
let operations = FakeOperations::successful(2);
let (result, events) = run_with_progress(&operations);
result.unwrap();
let mut expected = vec![
AnalysisProgress::JobStarted {
sequence: 1,
job: AnalysisJob::Transcript,
},
AnalysisProgress::JobSucceeded { sequence: 1 },
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Transcript,
},
AnalysisProgress::JobStarted {
sequence: 2,
job: AnalysisJob::SpeakerLabels,
},
AnalysisProgress::JobSucceeded { sequence: 2 },
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerLabels,
},
];
for (sequence, speaker, packet) in [
(3, label(1), 1),
(4, label(1), 2),
(5, label(1), 3),
(6, label(2), 1),
(7, label(2), 2),
(8, label(2), 3),
] {
expected.push(AnalysisProgress::JobStarted {
sequence,
job: AnalysisJob::SpeakerFeature { speaker, packet },
});
}
for sequence in 3..=8 {
expected.push(AnalysisProgress::JobSucceeded { sequence });
}
expected.extend([
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerFeatures,
},
AnalysisProgress::JobStarted {
sequence: 9,
job: AnalysisJob::Structuring,
},
AnalysisProgress::JobSucceeded { sequence: 9 },
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring,
},
]);
assert_eq!(events, expected);
}
#[test]
fn feature_terminals_preserve_actual_completion_order() {
let operations = FakeOperations::successful(2).with_config(|config| {
config.feature_terminal_order = Some(vec![
(label(2), 2),
(label(1), 3),
(label(2), 1),
(label(1), 1),
(label(2), 3),
(label(1), 2),
]);
});
let (result, events) = run_with_progress(&operations);
result.unwrap();
let terminals = events
.iter()
.filter_map(|event| match event {
AnalysisProgress::JobSucceeded { sequence } if (3..9).contains(sequence) => {
Some(*sequence)
}
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(terminals, [7, 5, 6, 3, 8, 4]);
}
#[test]
fn multiple_feature_failures_are_all_reported_without_duplicate() {
let operations = FakeOperations::successful(2).with_config(|config| {
config.feature_terminal_order = Some(vec![
(label(2), 1),
(label(1), 2),
(label(1), 1),
(label(1), 3),
(label(2), 2),
(label(2), 3),
]);
config.feature_failures = vec![
(label(1), 2, "first deterministic failure".into()),
(label(2), 1, "first completed failure".into()),
];
});
let (result, events) = run_with_progress(&operations);
assert_eq!(
result,
Err(AnalysisError::GeminiFeature {
speaker: label(1),
packet: 2,
message: "first deterministic failure".into(),
})
);
let failures = events
.iter()
.filter_map(|event| match event {
AnalysisProgress::JobFailed { sequence, error } => Some((*sequence, error.clone())),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(
failures,
vec![
(
6,
AnalysisError::GeminiFeature {
speaker: label(2),
packet: 1,
message: "first completed failure".into(),
}
.to_string(),
),
(
4,
AnalysisError::GeminiFeature {
speaker: label(1),
packet: 2,
message: "first deterministic failure".into(),
}
.to_string(),
),
]
);
assert!(!events.contains(&AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerFeatures,
}));
assert_eq!(
*operations.state.calls.lock().unwrap(),
["transcript", "labels", "features"]
);
}
#[test]
fn reporter_failure_supersedes_and_stops_orchestration() {
let operations = FakeOperations::successful(1);
let audio = strict_ogg(48);
let mut events = Vec::new();
let result = block_on(execute_strict(&operations, &audio, |event| {
events.push(event.clone());
if event
== (AnalysisProgress::JobStarted {
sequence: 3,
job: AnalysisJob::SpeakerFeature {
speaker: label(1),
packet: 1,
},
})
{
Err("reporter closed".into())
} else {
Ok(())
}
}));
assert_eq!(
result,
Err(AnalysisError::Progress("reporter closed".into()))
);
assert_eq!(
*operations.state.calls.lock().unwrap(),
["transcript", "labels", "features"]
);
assert_eq!(
events.last(),
Some(&AnalysisProgress::JobStarted {
sequence: 3,
job: AnalysisJob::SpeakerFeature {
speaker: label(1),
packet: 1,
},
})
);
}
#[test]
fn provider_stages_fail_without_retry_and_leaf_failures_are_exact() {
for (stage, expected_error, expected_calls) in [
(
FailureStage::Transcript,
AnalysisError::GeminiTranscript("transcript".into()),
vec!["transcript"],
),
(
FailureStage::Labels,
AnalysisError::TerraLabels("labels".into()),
vec!["transcript", "labels"],
),
(
FailureStage::Cache,
AnalysisError::GeminiCache("cache".into()),
vec!["transcript", "labels", "features"],
),
(
FailureStage::Features,
AnalysisError::GeminiFeature {
speaker: label(1),
packet: 2,
message: "features".into(),
},
vec!["transcript", "labels", "features"],
),
(
FailureStage::Final,
AnalysisError::TerraStructuring("final".into()),
vec!["transcript", "labels", "features", "final"],
),
] {
let operations =
FakeOperations::successful(1).with_config(|config| config.failure = Some(stage));
let (result, events) = run_with_progress(&operations);
assert_eq!(result, Err(expected_error.clone()));
assert_eq!(*operations.state.calls.lock().unwrap(), expected_calls);
let matching_failures = events
.iter()
.filter(|event| {
matches!(
event,
AnalysisProgress::JobFailed { error, .. }
if error == &expected_error.to_string()
)
})
.count();
if stage == FailureStage::Cache {
assert_eq!(matching_failures, 0);
} else {
assert_eq!(matching_failures, 1);
}
}
}
#[test]
fn cross_stage_transcript_and_speaker_mismatches_are_rejected() {
let transcript = FakeOperations::successful(1).with_config(|config| {
config.analysis = structured_analysis("different", 1);
});
let (result, events) = run_with_progress(&transcript);
assert_eq!(result, Err(AnalysisError::TranscriptMismatch));
assert!(events.contains(&AnalysisProgress::JobSucceeded { sequence: 6 }));
assert!(!events.contains(&AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring,
}));
let speakers = FakeOperations::successful(1).with_config(|config| {
config.analysis = structured_analysis(&config.transcript, 2);
});
let (result, events) = run_with_progress(&speakers);
assert_eq!(result, Err(AnalysisError::SpeakerSetMismatch));
assert!(events.contains(&AnalysisProgress::JobSucceeded { sequence: 6 }));
assert!(!events.contains(&AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring,
}));
}
#[test]
fn legacy_api_retains_metadata_and_previous_stage_order() {
for speaker_count in [0, 1, 40] {
let operations = FakeOperations::successful(speaker_count);
let result = block_on(execute_legacy(
&operations,
&legacy_ogg(),
1234,
Some("voice.ogg".into()),
))
.unwrap();
assert_eq!(result.envelope.audio.duration_ms(), 1234);
assert_eq!(result.envelope.audio.filename(), Some("voice.ogg"));
assert_eq!(
result.envelope.analysis.speakers.len(),
speaker_count as usize
);
assert_eq!(
*operations.state.calls.lock().unwrap(),
["transcript", "labels", "features", "final"]
);
assert_eq!(
result.label_extractor,
StructurerProvenance {
model_id: TERRA_MODEL_ID.into(),
prompt_revision:
kcode_speaker_v3_llm_protocol::TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
}
);
}
let operations = FakeOperations::successful(1);
assert!(matches!(
block_on(execute_legacy(&operations, b"bad", 1, None)),
Err(AnalysisError::Input(_))
));
assert!(operations.state.calls.lock().unwrap().is_empty());
}
#[test]
fn a_blocked_analysis_does_not_block_an_unrelated_analysis() {
let completed = Arc::new(AtomicBool::new(false));
let fast = FakeOperations::successful(0).with_config(|config| {
config.mark_complete = Some(completed.clone());
});
let slow = FakeOperations::successful(0).with_config(|config| {
config.wait_for = Some(completed.clone());
});
let slow_audio = strict_ogg(48);
let fast_audio = strict_ogg(48);
let (slow_result, fast_result) = block_on(async {
join!(
execute_strict(&slow, &slow_audio, |_| Ok(())),
execute_strict(&fast, &fast_audio, |_| Ok(()))
)
});
slow_result.unwrap();
fast_result.unwrap();
assert!(completed.load(Ordering::SeqCst));
}
#[test]
fn provenance_preserves_the_previous_public_contract() {
let cohort = GeminiCohort::new("gemini-model");
assert_eq!(
cohort.feature_prompt_revisions,
GEMINI_FEATURE_PROMPT_REVISIONS.map(str::to_owned)
);
cohort.validate().unwrap();
StructurerProvenance::new("gpt-5.6").validate().unwrap();
assert_eq!(
GeminiCohort::new(" ").validate(),
Err(ValidationError::Blank("gemini_model_id"))
);
}
#[test]
fn reference_scale_local_orchestration_completes_within_envelope() {
let started = Instant::now();
let operations = FakeOperations::successful(1000);
let (result, _) = run_with_progress(&operations);
assert_eq!(result.unwrap().envelope.analysis.speakers.len(), 1000);
assert!(started.elapsed().as_secs() < 10);
}
#[test]
fn concrete_provider_operations_compile() {
fn require_operations<O: AnalysisOperations>() {}
require_operations::<ProviderOperations>();
}
}