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 futures::{
FutureExt,
future::{BoxFuture, join_all},
};
#[cfg(any(feature = "providers", test))]
use kcode_speaker_v3_llm_protocol::{
FEATURE_PACKETS, FeaturePacket, GeminiRequestPart, SpeakerFeatureEvidence,
TERRA_SPEAKER_LABELS_PROMPT_REVISION, TerraFinalInput, TerraSpeakerLabelsInput, ToolDefinition,
decode_record_speaker_analysis_arguments, decode_record_speaker_labels_arguments,
extract_gemini_text, gemini_feature_cached_prefix, gemini_feature_suffix,
gemini_transcript_request, record_speaker_analysis_tool, record_speaker_labels_tool,
};
#[cfg(any(feature = "providers", test))]
use serde_json::Value;
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 AnalysisError {
Input(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::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 {
backend: ProviderBackend,
}
#[cfg(feature = "providers")]
impl Analyzer {
pub fn new(
gemini: kcode_gemini_3_1_pro::Gemini31Pro,
terra: kcode_codex_terra::CodexTerra,
) -> Self {
Self {
backend: ProviderBackend { gemini, terra },
}
}
pub async fn analyze_ogg(
&self,
bytes: &[u8],
duration_ms: u64,
filename: Option<String>,
) -> Result<ExecutedAnalysis, AnalysisError> {
execute(&self.backend, 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))]
trait AnalysisBackend: Sync {
type Cache: Send + Sync;
fn gemini_transcript<'a>(&'a self, audio: &'a [u8]) -> BoxFuture<'a, Result<Value, String>>;
fn terra_labels<'a>(
&'a self,
input: String,
tool: ToolDefinition,
) -> BoxFuture<'a, Result<Value, String>>;
fn gemini_cache<'a>(
&'a self,
audio: &'a [u8],
transcript: &'a str,
) -> BoxFuture<'a, Result<Self::Cache, String>>;
fn gemini_feature<'a>(
&'a self,
cache: &'a Self::Cache,
speaker: LocalSpeakerLabel,
packet: FeaturePacket,
) -> BoxFuture<'a, Result<Value, String>>;
fn terra_final<'a>(
&'a self,
input: String,
tool: ToolDefinition,
) -> BoxFuture<'a, Result<Value, String>>;
}
#[cfg(any(feature = "providers", test))]
async fn execute<B: AnalysisBackend>(
backend: &B,
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()))?;
let transcript_response = backend
.gemini_transcript(bytes)
.await
.map_err(AnalysisError::GeminiTranscript)?;
let transcript = extract_gemini_text(&transcript_response)
.map_err(|error| AnalysisError::GeminiTranscript(error.to_string()))?;
let labels_input = TerraSpeakerLabelsInput::new(transcript.clone())
.map_err(|error| AnalysisError::TerraLabels(error.to_string()))?;
let labels_arguments = backend
.terra_labels(labels_input.render(), record_speaker_labels_tool())
.await
.map_err(AnalysisError::TerraLabels)?;
let labels = decode_record_speaker_labels_arguments(&labels_arguments)
.map_err(|error| AnalysisError::TerraLabels(error.to_string()))?;
let evidence = collect_feature_evidence(backend, bytes, &transcript, &labels).await?;
let final_input = TerraFinalInput::new(transcript.clone(), evidence)
.map_err(|error| AnalysisError::TerraStructuring(error.to_string()))?;
let final_arguments = backend
.terra_final(final_input.render(), record_speaker_analysis_tool())
.await
.map_err(AnalysisError::TerraStructuring)?;
let analysis = decode_record_speaker_analysis_arguments(&final_arguments)
.map_err(|error| AnalysisError::TerraStructuring(error.to_string()))?;
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: TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
};
label_extractor
.validate()
.map_err(|error| AnalysisError::TerraLabels(error.to_string()))?;
Ok(ExecutedAnalysis {
envelope,
label_extractor,
})
}
#[cfg(any(feature = "providers", test))]
async fn collect_feature_evidence<B: AnalysisBackend>(
backend: &B,
bytes: &[u8],
transcript: &str,
labels: &[LocalSpeakerLabel],
) -> Result<Vec<SpeakerFeatureEvidence>, AnalysisError> {
if labels.is_empty() {
return Ok(Vec::new());
}
let cache = backend
.gemini_cache(bytes, transcript)
.await
.map_err(AnalysisError::GeminiCache)?;
let jobs = labels
.iter()
.copied()
.flat_map(|speaker| {
FEATURE_PACKETS
.into_iter()
.map(move |packet| (speaker, packet))
})
.collect::<Vec<_>>();
let responses = join_all(
jobs.iter()
.map(|(speaker, packet)| backend.gemini_feature(&cache, *speaker, *packet)),
)
.await;
let mut texts = Vec::with_capacity(responses.len());
for ((speaker, packet), response) in jobs.into_iter().zip(responses) {
let response = response.map_err(|message| AnalysisError::GeminiFeature {
speaker,
packet: packet.index() as u8 + 1,
message,
})?;
let text =
extract_gemini_text(&response).map_err(|error| AnalysisError::GeminiFeature {
speaker,
packet: packet.index() as u8 + 1,
message: error.to_string(),
})?;
texts.push(text);
}
let mut evidence = Vec::with_capacity(labels.len());
for (speaker, packets) in labels.iter().copied().zip(texts.chunks_exact(3)) {
evidence.push(
SpeakerFeatureEvidence::new(
speaker,
packets[0].clone(),
packets[1].clone(),
packets[2].clone(),
)
.map_err(|error| AnalysisError::GeminiFeature {
speaker,
packet: 1,
message: error.to_string(),
})?,
);
}
Ok(evidence)
}
#[cfg(any(feature = "providers", test))]
#[derive(Clone)]
struct ProviderBackend {
gemini: kcode_gemini_3_1_pro::Gemini31Pro,
terra: kcode_codex_terra::CodexTerra,
}
#[cfg(any(feature = "providers", test))]
impl AnalysisBackend for ProviderBackend {
type Cache = kcode_gemini_3_1_pro::CachedPrefix;
fn gemini_transcript<'a>(&'a self, audio: &'a [u8]) -> BoxFuture<'a, Result<Value, String>> {
async move {
let contents = vec![gemini_content(gemini_transcript_request(audio))];
self.gemini
.generate(contents, None, None)
.await
.map(|generation| generation.response)
.map_err(|error| error.to_string())
}
.boxed()
}
fn terra_labels<'a>(
&'a self,
input: String,
tool: ToolDefinition,
) -> BoxFuture<'a, Result<Value, String>> {
async move { run_terra(&self.terra, input, tool).await }.boxed()
}
fn gemini_cache<'a>(
&'a self,
audio: &'a [u8],
transcript: &'a str,
) -> BoxFuture<'a, Result<Self::Cache, String>> {
async move {
let contents = vec![gemini_content(gemini_feature_cached_prefix(
audio, transcript,
))];
self.gemini
.create_cached_prefix(contents, None, std::time::Duration::from_secs(60 * 60))
.await
.map_err(|error| error.to_string())
}
.boxed()
}
fn gemini_feature<'a>(
&'a self,
cache: &'a Self::Cache,
speaker: LocalSpeakerLabel,
packet: FeaturePacket,
) -> BoxFuture<'a, Result<Value, String>> {
async move {
let contents = vec![kcode_gemini_3_1_pro::Content {
parts: vec![kcode_gemini_3_1_pro::Part::Text(gemini_feature_suffix(
packet, speaker,
))],
}];
self.gemini
.generate(contents, Some(cache), None)
.await
.map(|generation| generation.response)
.map_err(|error| error.to_string())
}
.boxed()
}
fn terra_final<'a>(
&'a self,
input: String,
tool: ToolDefinition,
) -> BoxFuture<'a, Result<Value, String>> {
async move { run_terra(&self.terra, input, tool).await }.boxed()
}
}
#[cfg(any(feature = "providers", test))]
fn gemini_content<'a>(
parts: impl IntoIterator<Item = GeminiRequestPart<'a>>,
) -> kcode_gemini_3_1_pro::Content {
kcode_gemini_3_1_pro::Content {
parts: parts
.into_iter()
.map(|part| match part {
GeminiRequestPart::Audio { media_type, bytes } => {
kcode_gemini_3_1_pro::Part::InlineData {
mime_type: media_type.into(),
bytes: bytes.to_vec(),
}
}
GeminiRequestPart::Text(text) => kcode_gemini_3_1_pro::Part::Text(text),
})
.collect(),
}
}
#[cfg(any(feature = "providers", test))]
async fn run_terra(
terra: &kcode_codex_terra::CodexTerra,
input: String,
tool: ToolDefinition,
) -> Result<Value, String> {
terra
.run(kcode_codex_terra::ToolRun {
input,
tool_name: tool.name.into(),
tool_description: tool.description.into(),
input_schema: tool.input_schema,
})
.await
.map(|result| result.arguments)
.map_err(|error| error.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
use futures::{executor::block_on, future::poll_fn, join};
use serde_json::{Map, json};
use std::{
sync::{
Arc, Mutex,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
task::Poll,
};
#[derive(Clone)]
struct FakeBackend {
state: Arc<FakeState>,
}
struct FakeState {
transcript: Value,
labels: Value,
final_arguments: Value,
transcript_error: Option<String>,
labels_error: Option<String>,
cache_error: Option<String>,
feature_error: Option<(LocalSpeakerLabel, FeaturePacket)>,
final_error: Option<String>,
wait_for: Option<Arc<AtomicBool>>,
mark_complete: Option<Arc<AtomicBool>>,
transcript_calls: AtomicUsize,
labels_calls: AtomicUsize,
cache_calls: AtomicUsize,
feature_calls: AtomicUsize,
final_calls: AtomicUsize,
active_features: AtomicUsize,
maximum_active_features: AtomicUsize,
targets: Mutex<Vec<(LocalSpeakerLabel, FeaturePacket)>>,
}
impl FakeBackend {
fn successful(speaker_count: u32) -> Self {
let labels = (1..=speaker_count)
.map(|number| format!("Speaker {number}"))
.collect::<Vec<_>>();
Self {
state: Arc::new(FakeState {
transcript: gemini_text("[high] Speaker 1: exact transcript"),
labels: json!({ "speakers": labels }),
final_arguments: final_arguments(
"[high] Speaker 1: exact transcript",
speaker_count,
),
transcript_error: None,
labels_error: None,
cache_error: None,
feature_error: None,
final_error: None,
wait_for: None,
mark_complete: None,
transcript_calls: AtomicUsize::new(0),
labels_calls: AtomicUsize::new(0),
cache_calls: AtomicUsize::new(0),
feature_calls: AtomicUsize::new(0),
final_calls: AtomicUsize::new(0),
active_features: AtomicUsize::new(0),
maximum_active_features: AtomicUsize::new(0),
targets: Mutex::new(Vec::new()),
}),
}
}
fn with_state(&self, update: impl FnOnce(&FakeState) -> FakeState) -> Self {
Self {
state: Arc::new(update(&self.state)),
}
}
}
impl AnalysisBackend for FakeBackend {
type Cache = ();
fn gemini_transcript<'a>(
&'a self,
_audio: &'a [u8],
) -> BoxFuture<'a, Result<Value, String>> {
async move {
self.state.transcript_calls.fetch_add(1, Ordering::SeqCst);
if let Some(wait_for) = &self.state.wait_for {
poll_fn(|context| {
if wait_for.load(Ordering::SeqCst) {
Poll::Ready(())
} else {
context.waker().wake_by_ref();
Poll::Pending
}
})
.await;
}
match &self.state.transcript_error {
Some(error) => Err(error.clone()),
None => Ok(self.state.transcript.clone()),
}
}
.boxed()
}
fn terra_labels<'a>(
&'a self,
_input: String,
_tool: ToolDefinition,
) -> BoxFuture<'a, Result<Value, String>> {
async move {
self.state.labels_calls.fetch_add(1, Ordering::SeqCst);
match &self.state.labels_error {
Some(error) => Err(error.clone()),
None => Ok(self.state.labels.clone()),
}
}
.boxed()
}
fn gemini_cache<'a>(
&'a self,
_audio: &'a [u8],
_transcript: &'a str,
) -> BoxFuture<'a, Result<Self::Cache, String>> {
async move {
self.state.cache_calls.fetch_add(1, Ordering::SeqCst);
match &self.state.cache_error {
Some(error) => Err(error.clone()),
None => Ok(()),
}
}
.boxed()
}
fn gemini_feature<'a>(
&'a self,
_cache: &'a Self::Cache,
speaker: LocalSpeakerLabel,
packet: FeaturePacket,
) -> BoxFuture<'a, Result<Value, String>> {
async move {
self.state.feature_calls.fetch_add(1, Ordering::SeqCst);
self.state.targets.lock().unwrap().push((speaker, packet));
let active = self.state.active_features.fetch_add(1, Ordering::SeqCst) + 1;
self.state
.maximum_active_features
.fetch_max(active, Ordering::SeqCst);
let mut yielded = false;
poll_fn(|context| {
if yielded {
Poll::Ready(())
} else {
yielded = true;
context.waker().wake_by_ref();
Poll::Pending
}
})
.await;
self.state.active_features.fetch_sub(1, Ordering::SeqCst);
if self.state.feature_error == Some((speaker, packet)) {
Err("feature failure".into())
} else {
Ok(gemini_text(&format!(
"{speaker} packet {}",
packet.index() + 1
)))
}
}
.boxed()
}
fn terra_final<'a>(
&'a self,
_input: String,
_tool: ToolDefinition,
) -> BoxFuture<'a, Result<Value, String>> {
async move {
self.state.final_calls.fetch_add(1, Ordering::SeqCst);
match &self.state.final_error {
Some(error) => Err(error.clone()),
None => {
if let Some(mark_complete) = &self.state.mark_complete {
mark_complete.store(true, Ordering::SeqCst);
}
Ok(self.state.final_arguments.clone())
}
}
}
.boxed()
}
}
fn copy_state(state: &FakeState) -> FakeState {
FakeState {
transcript: state.transcript.clone(),
labels: state.labels.clone(),
final_arguments: state.final_arguments.clone(),
transcript_error: state.transcript_error.clone(),
labels_error: state.labels_error.clone(),
cache_error: state.cache_error.clone(),
feature_error: state.feature_error,
final_error: state.final_error.clone(),
wait_for: state.wait_for.clone(),
mark_complete: state.mark_complete.clone(),
transcript_calls: AtomicUsize::new(0),
labels_calls: AtomicUsize::new(0),
cache_calls: AtomicUsize::new(0),
feature_calls: AtomicUsize::new(0),
final_calls: AtomicUsize::new(0),
active_features: AtomicUsize::new(0),
maximum_active_features: AtomicUsize::new(0),
targets: Mutex::new(Vec::new()),
}
}
fn gemini_text(text: &str) -> Value {
json!({
"candidates": [{
"content": {
"parts": [{ "text": text }]
}
}]
})
}
fn null_features() -> Map<String, Value> {
FEATURE_NAMES
.into_iter()
.map(|name| (name.into(), Value::Null))
.collect()
}
fn final_arguments(transcript: &str, speaker_count: u32) -> Value {
let speakers = (1..=speaker_count)
.map(|number| {
json!({
"speaker": format!("Speaker {number}"),
"language": "English",
"features": null_features(),
"features_usable_for_training": false
})
})
.collect::<Vec<_>>();
json!({
"transcript": transcript,
"speakers": speakers
})
}
fn 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
}
#[test]
fn zero_one_and_many_speaker_workflows_have_exact_calls() {
for speaker_count in [0, 1, 40] {
let backend = FakeBackend::successful(speaker_count);
let result = block_on(execute(&backend, &ogg(), 1, Some("voice.ogg".into()))).unwrap();
assert_eq!(
result.envelope.analysis.speakers.len(),
speaker_count as usize
);
assert_eq!(
result.label_extractor,
StructurerProvenance {
model_id: TERRA_MODEL_ID.into(),
prompt_revision: TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
}
);
assert_eq!(backend.state.transcript_calls.load(Ordering::SeqCst), 1);
assert_eq!(backend.state.labels_calls.load(Ordering::SeqCst), 1);
assert_eq!(
backend.state.cache_calls.load(Ordering::SeqCst),
usize::from(speaker_count > 0)
);
assert_eq!(
backend.state.feature_calls.load(Ordering::SeqCst),
speaker_count as usize * 3
);
assert_eq!(backend.state.final_calls.load(Ordering::SeqCst), 1);
}
}
#[test]
fn terra_labels_control_targets_and_features_overlap() {
let backend = FakeBackend::successful(3);
block_on(execute(&backend, &ogg(), 1, None)).unwrap();
let expected = (1..=3)
.flat_map(|number| {
FEATURE_PACKETS
.into_iter()
.map(move |packet| (LocalSpeakerLabel::new(number).unwrap(), packet))
})
.collect::<Vec<_>>();
assert_eq!(*backend.state.targets.lock().unwrap(), expected);
assert!(backend.state.maximum_active_features.load(Ordering::SeqCst) > 1);
}
#[test]
fn input_and_each_provider_stage_fail_without_retry() {
let backend = FakeBackend::successful(1);
assert!(matches!(
block_on(execute(&backend, b"bad", 1, None)),
Err(AnalysisError::Input(_))
));
assert_eq!(backend.state.transcript_calls.load(Ordering::SeqCst), 0);
let transcript_failure = backend.with_state(|state| {
let mut state = copy_state(state);
state.transcript_error = Some("transcript".into());
state
});
assert!(matches!(
block_on(execute(&transcript_failure, &ogg(), 1, None)),
Err(AnalysisError::GeminiTranscript(_))
));
assert_eq!(
transcript_failure
.state
.transcript_calls
.load(Ordering::SeqCst),
1
);
assert_eq!(
transcript_failure.state.labels_calls.load(Ordering::SeqCst),
0
);
let labels_failure = backend.with_state(|state| {
let mut state = copy_state(state);
state.labels_error = Some("labels".into());
state
});
assert!(matches!(
block_on(execute(&labels_failure, &ogg(), 1, None)),
Err(AnalysisError::TerraLabels(_))
));
assert_eq!(labels_failure.state.labels_calls.load(Ordering::SeqCst), 1);
assert_eq!(labels_failure.state.cache_calls.load(Ordering::SeqCst), 0);
let cache_failure = backend.with_state(|state| {
let mut state = copy_state(state);
state.cache_error = Some("cache".into());
state
});
assert!(matches!(
block_on(execute(&cache_failure, &ogg(), 1, None)),
Err(AnalysisError::GeminiCache(_))
));
assert_eq!(cache_failure.state.cache_calls.load(Ordering::SeqCst), 1);
assert_eq!(cache_failure.state.feature_calls.load(Ordering::SeqCst), 0);
let feature_failure = backend.with_state(|state| {
let mut state = copy_state(state);
state.feature_error = Some((LocalSpeakerLabel::new(1).unwrap(), FeaturePacket::Two));
state
});
assert!(matches!(
block_on(execute(&feature_failure, &ogg(), 1, None)),
Err(AnalysisError::GeminiFeature { packet: 2, .. })
));
assert_eq!(
feature_failure.state.feature_calls.load(Ordering::SeqCst),
3
);
assert_eq!(feature_failure.state.final_calls.load(Ordering::SeqCst), 0);
let final_failure = backend.with_state(|state| {
let mut state = copy_state(state);
state.final_error = Some("final".into());
state
});
assert!(matches!(
block_on(execute(&final_failure, &ogg(), 1, None)),
Err(AnalysisError::TerraStructuring(_))
));
assert_eq!(final_failure.state.final_calls.load(Ordering::SeqCst), 1);
}
#[test]
fn protocol_failures_and_cross_stage_mismatches_are_rejected() {
let backend = FakeBackend::successful(1);
let invalid_transcript = backend.with_state(|state| {
let mut state = copy_state(state);
state.transcript = json!({});
state
});
assert!(matches!(
block_on(execute(&invalid_transcript, &ogg(), 1, None)),
Err(AnalysisError::GeminiTranscript(_))
));
let invalid_labels = backend.with_state(|state| {
let mut state = copy_state(state);
state.labels = json!({ "speakers": ["Unknown"] });
state
});
assert!(matches!(
block_on(execute(&invalid_labels, &ogg(), 1, None)),
Err(AnalysisError::TerraLabels(_))
));
let invalid_final = backend.with_state(|state| {
let mut state = copy_state(state);
state.final_arguments = json!({});
state
});
assert!(matches!(
block_on(execute(&invalid_final, &ogg(), 1, None)),
Err(AnalysisError::TerraStructuring(_))
));
let transcript_mismatch = backend.with_state(|state| {
let mut state = copy_state(state);
state.final_arguments = final_arguments("different", 1);
state
});
assert_eq!(
block_on(execute(&transcript_mismatch, &ogg(), 1, None)),
Err(AnalysisError::TranscriptMismatch)
);
let speaker_mismatch = backend.with_state(|state| {
let mut state = copy_state(state);
state.final_arguments = final_arguments("[high] Speaker 1: exact transcript", 2);
state
});
assert_eq!(
block_on(execute(&speaker_mismatch, &ogg(), 1, None)),
Err(AnalysisError::SpeakerSetMismatch)
);
}
#[test]
fn a_blocked_analysis_does_not_block_an_unrelated_analysis() {
let completed = Arc::new(AtomicBool::new(false));
let fast = FakeBackend::successful(0).with_state(|state| {
let mut state = copy_state(state);
state.mark_complete = Some(completed.clone());
state
});
let slow = FakeBackend::successful(0).with_state(|state| {
let mut state = copy_state(state);
state.wait_for = Some(completed.clone());
state
});
let slow_audio = ogg();
let fast_audio = ogg();
let (slow_result, fast_result) = block_on(async {
join!(
execute(&slow, &slow_audio, 1, None),
execute(&fast, &fast_audio, 1, None)
)
});
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 concrete_provider_backend_compiles() {
fn require_backend<B: AnalysisBackend>() {}
require_backend::<ProviderBackend>();
}
}