#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
use kcode_speaker_v3_gemini_analysis::GeminiFeatureProgress;
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,
SpeakerFeatureEvidence,
};
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,
};
use serde::{Deserialize, Serialize};
use std::{collections::BTreeSet, error::Error, fmt};
#[cfg(any(feature = "providers", feature = "adapter-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 value in &self.feature_prompt_revisions {
validate_text(value, "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, Serialize, Deserialize)]
pub struct ExecutedAnalysis {
pub envelope: AnalysisEnvelope,
pub label_extractor: StructurerProvenance,
}
#[derive(Debug, Clone, PartialEq, Serialize, Default)]
pub struct RecoveredAnalysisArtifacts {
pub transcript: Option<String>,
pub speaker_labels: Option<Vec<LocalSpeakerLabel>>,
pub feature_evidence: Vec<SpeakerFeatureEvidence>,
pub executed_analysis: Option<ExecutedAnalysis>,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub enum AnalysisArtifact {
GeminiTranscript(String),
TerraSpeakerLabels(Vec<LocalSpeakerLabel>),
GeminiFeatureEvidence(SpeakerFeatureEvidence),
ExecutedAnalysis(Box<ExecutedAnalysis>),
}
#[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, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Input(v) => write!(f, "invalid input: {v}"),
Self::Progress(v) => write!(f, "progress reporting failed: {v}"),
Self::GeminiTranscript(v) => write!(f, "Gemini transcript failed: {v}"),
Self::TerraLabels(v) => write!(f, "Terra speaker-label extraction failed: {v}"),
Self::GeminiCache(v) => write!(f, "Gemini feature cache creation failed: {v}"),
Self::GeminiFeature {
speaker,
packet,
message,
} => write!(
f,
"Gemini feature call failed for {speaker}, packet {packet}: {message}"
),
Self::TerraStructuring(v) => write!(f, "Terra final structuring failed: {v}"),
Self::TranscriptMismatch => f.write_str("Terra returned a different transcript"),
Self::SpeakerSetMismatch => f.write_str("Terra returned a different speaker set"),
}
}
}
impl Error for AnalysisError {}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ResumableAnalysisError {
Analysis(AnalysisError),
Recovered(String),
Artifact(String),
}
impl fmt::Display for ResumableAnalysisError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Analysis(error) => write!(f, "{error}"),
Self::Recovered(v) => write!(f, "invalid recovered artifacts: {v}"),
Self::Artifact(v) => write!(f, "artifact reporting failed: {v}"),
}
}
}
impl Error for ResumableAnalysisError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Analysis(error) => Some(error),
Self::Recovered(_) | Self::Artifact(_) => None,
}
}
}
impl From<AnalysisError> for ResumableAnalysisError {
fn from(value: AnalysisError) -> Self {
Self::Analysis(value)
}
}
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", feature = "adapter-providers"))]
pub struct Analyzer {
operations: ProviderOperations,
}
#[cfg(any(feature = "providers", feature = "adapter-providers"))]
impl Analyzer {
#[cfg(feature = "providers")]
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),
},
}
}
#[cfg(feature = "adapter-providers")]
pub fn from_codex_adapter(
gemini: kcode_gemini_3_1_pro::Gemini31Pro,
adapter: kcode_k1_codex_adapter::Adapter,
) -> Self {
Self {
operations: ProviderOperations {
gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis::new(gemini),
terra: kcode_speaker_v3_terra_analysis::TerraAnalysis::from_codex_adapter(adapter),
},
}
}
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
}
pub async fn analyze_ogg_resumable<P, A>(
&self,
bytes: &[u8],
recovered: RecoveredAnalysisArtifacts,
report: P,
emit: A,
) -> Result<ExecutedAnalysis, ResumableAnalysisError>
where
P: FnMut(AnalysisProgress) -> Result<(), String>,
A: FnMut(AnalysisArtifact) -> Result<(), String>,
{
execute_resume_strict(&self.operations, bytes, recovered, report, emit).await
}
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
type AnalysisFuture<'a, T> = Pin<Box<dyn Future<Output = T> + 'a>>;
#[cfg(any(feature = "providers", feature = "adapter-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", feature = "adapter-providers", test))]
struct ProviderOperations {
gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis,
terra: kcode_speaker_v3_terra_analysis::TerraAnalysis,
}
#[cfg(any(feature = "providers", feature = "adapter-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(v)
| kcode_speaker_v3_gemini_analysis::GeminiTranscriptError::Protocol(v) => {
AnalysisError::GeminiTranscript(v)
}
})
})
}
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(v)
| kcode_speaker_v3_terra_analysis::TerraAnalysisError::Provider(v) => {
AnalysisError::TerraLabels(v)
}
})
})
}
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(v) => {
AnalysisError::GeminiCache(v)
}
kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Progress(v) => {
AnalysisError::Progress(v)
}
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(v)
| kcode_speaker_v3_terra_analysis::TerraAnalysisError::Provider(v) => {
AnalysisError::TerraStructuring(v)
}
})
})
}
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
async fn execute_strict<O: AnalysisOperations, P: FnMut(AnalysisProgress) -> Result<(), String>>(
operations: &O,
bytes: &[u8],
report: P,
) -> Result<ExecutedAnalysis, AnalysisError> {
let audio =
OggAudioMetadata::from_ogg_bytes(bytes).map_err(|e| AnalysisError::Input(e.to_string()))?;
execute_admitted(
operations,
bytes,
audio,
RecoveredAnalysisArtifacts::default(),
report,
|_| Ok(()),
false,
)
.await
.map_err(non_resumable_error)
}
#[cfg(any(feature = "providers", feature = "adapter-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(|e| AnalysisError::Input(e.to_string()))?;
execute_admitted(
operations,
bytes,
audio,
RecoveredAnalysisArtifacts::default(),
|_| Ok(()),
|_| Ok(()),
false,
)
.await
.map_err(non_resumable_error)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
async fn execute_resume_strict<
O: AnalysisOperations,
P: FnMut(AnalysisProgress) -> Result<(), String>,
A: FnMut(AnalysisArtifact) -> Result<(), String>,
>(
operations: &O,
bytes: &[u8],
recovered: RecoveredAnalysisArtifacts,
report: P,
emit: A,
) -> Result<ExecutedAnalysis, ResumableAnalysisError> {
let audio =
OggAudioMetadata::from_ogg_bytes(bytes).map_err(|e| AnalysisError::Input(e.to_string()))?;
execute_admitted(operations, bytes, audio, recovered, report, emit, true).await
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
async fn execute_admitted<O, P, A>(
operations: &O,
bytes: &[u8],
audio: OggAudioMetadata,
mut recovered: RecoveredAnalysisArtifacts,
mut report: P,
mut emit: A,
resumable: bool,
) -> Result<ExecutedAnalysis, ResumableAnalysisError>
where
O: AnalysisOperations,
P: FnMut(AnalysisProgress) -> Result<(), String>,
A: FnMut(AnalysisArtifact) -> Result<(), String>,
{
let (recovered_speakers, final_speakers) = if resumable {
validate_recovered(&recovered, &audio)?
} else {
(BTreeSet::new(), None)
};
let transcript = if let Some(value) = recovered.transcript.take() {
value
} else if let Some(executed) = recovered.executed_analysis.as_ref() {
executed.envelope.analysis.transcript.clone()
} else {
progress(
&mut report,
AnalysisProgress::JobStarted {
sequence: 1,
job: AnalysisJob::Transcript,
},
)?;
let value = match operations.transcript(bytes).await {
Ok(v) => v,
Err(e) => {
failed(&mut report, 1, &e)?;
return Err(e.into());
}
};
if let Err(error) = validate_generated_transcript(&value, resumable) {
failed(&mut report, 1, &error)?;
return Err(error.into());
}
progress(&mut report, AnalysisProgress::JobSucceeded { sequence: 1 })?;
if resumable {
emit(AnalysisArtifact::GeminiTranscript(value.clone()))
.map_err(ResumableAnalysisError::Artifact)?;
}
progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Transcript,
},
)?;
value
};
let labels = if let Some(value) = recovered.speaker_labels.take() {
value
} else {
progress(
&mut report,
AnalysisProgress::JobStarted {
sequence: 2,
job: AnalysisJob::SpeakerLabels,
},
)?;
let value = match operations.speaker_labels(&transcript).await {
Ok(v) => v,
Err(e) => {
failed(&mut report, 2, &e)?;
return Err(e.into());
}
};
if let Err(error) = validate_generated_labels(
&value,
&recovered_speakers,
final_speakers.as_ref(),
resumable,
) {
failed(&mut report, 2, &error)?;
return Err(error);
}
progress(&mut report, AnalysisProgress::JobSucceeded { sequence: 2 })?;
if resumable {
emit(AnalysisArtifact::TerraSpeakerLabels(value.clone()))
.map_err(ResumableAnalysisError::Artifact)?;
}
progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerLabels,
},
)?;
value
};
let label_set = labels.iter().copied().collect::<BTreeSet<_>>();
let sequence = structuring_sequence(labels.len())?;
let missing = labels
.iter()
.copied()
.filter(|label| !recovered_speakers.contains(label))
.collect::<Vec<_>>();
let targets = if resumable { missing } else { labels.clone() };
let produced = if !resumable || !targets.is_empty() {
let mut feature_report = |event| report(map_feature_progress(&labels, event)?);
let value = operations
.feature_evidence_with_progress(bytes, &transcript, &targets, &mut feature_report)
.await?;
if resumable {
let returned =
validate_evidence(&value).map_err(|message| feature_error(targets[0], message))?;
let expected = targets.iter().copied().collect::<BTreeSet<_>>();
if returned != expected {
return Err(feature_error(
targets[0],
"generated feature speaker set differs from requested labels".into(),
)
.into());
}
for target in &targets {
let item = value
.iter()
.find(|item| item.speaker() == *target)
.expect("validated feature set");
emit(AnalysisArtifact::GeminiFeatureEvidence(item.clone()))
.map_err(ResumableAnalysisError::Artifact)?;
}
}
progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerFeatures,
},
)?;
value
} else {
Vec::new()
};
let evidence = if resumable {
recovered.feature_evidence.extend(produced);
let complete = validate_evidence(&recovered.feature_evidence)
.map_err(ResumableAnalysisError::Recovered)?;
if complete != label_set {
return Err(ResumableAnalysisError::Recovered(
"feature evidence does not cover the complete label set".into(),
));
}
labels
.iter()
.map(|label| {
recovered
.feature_evidence
.iter()
.find(|item| item.speaker() == *label)
.expect("validated complete evidence")
.clone()
})
.collect()
} else {
produced
};
if let Some(executed) = recovered.executed_analysis.take() {
return Ok(executed);
}
progress(
&mut report,
AnalysisProgress::JobStarted {
sequence,
job: AnalysisJob::Structuring,
},
)?;
let analysis = match operations.structured_analysis(&transcript, evidence).await {
Ok(v) => v,
Err(e) => {
failed(&mut report, sequence, &e)?;
return Err(e.into());
}
};
progress(&mut report, AnalysisProgress::JobSucceeded { sequence })?;
if analysis.transcript != transcript {
return Err(AnalysisError::TranscriptMismatch.into());
}
let returned = analysis
.speakers
.iter()
.map(|speaker| speaker.speaker)
.collect::<BTreeSet<_>>();
if returned != label_set {
return Err(AnalysisError::SpeakerSetMismatch.into());
}
let envelope = AnalysisEnvelope {
audio,
analysis,
gemini: GeminiCohort::new(GEMINI_MODEL_ID),
structurer: StructurerProvenance::new(TERRA_MODEL_ID),
};
envelope
.validate()
.map_err(|e| AnalysisError::TerraStructuring(e.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(|e| AnalysisError::TerraLabels(e.to_string()))?;
let executed = ExecutedAnalysis {
envelope,
label_extractor,
};
if resumable {
emit(AnalysisArtifact::ExecutedAnalysis(Box::new(
executed.clone(),
)))
.map_err(ResumableAnalysisError::Artifact)?;
}
progress(
&mut report,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring,
},
)?;
Ok(executed)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn validate_generated_transcript(value: &str, resumable: bool) -> Result<(), AnalysisError> {
if !resumable {
return Ok(());
}
validate_text(value, "transcript")
.map_err(|error| AnalysisError::GeminiTranscript(error.to_string()))
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn validate_generated_labels(
labels: &[LocalSpeakerLabel],
recovered: &BTreeSet<LocalSpeakerLabel>,
final_speakers: Option<&BTreeSet<LocalSpeakerLabel>>,
resumable: bool,
) -> Result<(), ResumableAnalysisError> {
if !resumable {
return Ok(());
}
let labels = validate_labels(labels).map_err(AnalysisError::TerraLabels)?;
if !recovered.is_subset(&labels) {
return Err(ResumableAnalysisError::Recovered(
"feature evidence contains a speaker absent from generated labels".into(),
));
}
if final_speakers.is_some_and(|expected| expected != &labels) {
return Err(ResumableAnalysisError::Recovered(
"generated labels differ from the recovered final speaker set".into(),
));
}
Ok(())
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn validate_labels(labels: &[LocalSpeakerLabel]) -> Result<BTreeSet<LocalSpeakerLabel>, String> {
let set = labels.iter().copied().collect::<BTreeSet<_>>();
if set.len() != labels.len() {
return Err("speaker labels contain a duplicate".into());
}
Ok(set)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn validate_evidence(
values: &[SpeakerFeatureEvidence],
) -> Result<BTreeSet<LocalSpeakerLabel>, String> {
let mut set = BTreeSet::new();
for value in values {
let packets = value.packets();
SpeakerFeatureEvidence::new(
value.speaker(),
packets[0].into(),
packets[1].into(),
packets[2].into(),
)
.map_err(|e| e.to_string())?;
if !set.insert(value.speaker()) {
return Err(format!(
"duplicate feature evidence for {}",
value.speaker()
));
}
}
Ok(set)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn validate_recovered(
value: &RecoveredAnalysisArtifacts,
audio: &OggAudioMetadata,
) -> Result<
(
BTreeSet<LocalSpeakerLabel>,
Option<BTreeSet<LocalSpeakerLabel>>,
),
ResumableAnalysisError,
> {
if let Some(transcript) = &value.transcript {
validate_text(transcript, "transcript")
.map_err(|error| ResumableAnalysisError::Recovered(error.to_string()))?;
}
let evidence =
validate_evidence(&value.feature_evidence).map_err(ResumableAnalysisError::Recovered)?;
let labels = value
.speaker_labels
.as_deref()
.map(validate_labels)
.transpose()
.map_err(ResumableAnalysisError::Recovered)?;
if labels
.as_ref()
.is_some_and(|labels| !evidence.is_subset(labels))
{
return Err(ResumableAnalysisError::Recovered(
"feature evidence contains a speaker absent from labels".into(),
));
}
let final_speakers = value
.executed_analysis
.as_ref()
.map(|executed| {
validate_recovered_final(
executed,
audio,
value.transcript.as_deref(),
labels.as_ref(),
&evidence,
)
})
.transpose()?;
Ok((evidence, final_speakers))
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn validate_recovered_final(
executed: &ExecutedAnalysis,
audio: &OggAudioMetadata,
transcript: Option<&str>,
labels: Option<&BTreeSet<LocalSpeakerLabel>>,
evidence: &BTreeSet<LocalSpeakerLabel>,
) -> Result<BTreeSet<LocalSpeakerLabel>, ResumableAnalysisError> {
executed
.envelope
.validate()
.map_err(|error| ResumableAnalysisError::Recovered(error.to_string()))?;
executed
.label_extractor
.validate()
.map_err(|error| ResumableAnalysisError::Recovered(error.to_string()))?;
if executed.envelope.audio != *audio {
return Err(ResumableAnalysisError::Recovered(
"final audio metadata differs from admitted audio".into(),
));
}
if transcript.is_some_and(|value| value != executed.envelope.analysis.transcript.as_str()) {
return Err(ResumableAnalysisError::Recovered(
"transcript differs from the recovered final transcript".into(),
));
}
let speakers = executed
.envelope
.analysis
.speakers
.iter()
.map(|speaker| speaker.speaker)
.collect::<BTreeSet<_>>();
if labels.is_some_and(|labels| labels != &speakers) {
return Err(ResumableAnalysisError::Recovered(
"labels differ from the recovered final speaker set".into(),
));
}
if !evidence.is_subset(&speakers) {
return Err(ResumableAnalysisError::Recovered(
"feature evidence contains a speaker absent from the recovered final".into(),
));
}
Ok(speakers)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn non_resumable_error(error: ResumableAnalysisError) -> AnalysisError {
match error {
ResumableAnalysisError::Analysis(error) => error,
ResumableAnalysisError::Recovered(_) | ResumableAnalysisError::Artifact(_) => {
unreachable!("non-resumable execution created a resumable-only error")
}
}
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn feature_error(speaker: LocalSpeakerLabel, message: String) -> AnalysisError {
AnalysisError::GeminiFeature {
speaker,
packet: 1,
message,
}
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn progress<P: FnMut(AnalysisProgress) -> Result<(), String>>(
report: &mut P,
event: AnalysisProgress,
) -> Result<(), AnalysisError> {
report(event).map_err(AnalysisError::Progress)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn failed<P: FnMut(AnalysisProgress) -> Result<(), String>, E: fmt::Display>(
report: &mut P,
sequence: u64,
error: &E,
) -> Result<(), AnalysisError> {
progress(
report,
AnalysisProgress::JobFailed {
sequence,
error: error.to_string(),
},
)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn structuring_sequence(count: usize) -> Result<u64, AnalysisError> {
u64::try_from(count)
.ok()
.and_then(|v| v.checked_mul(3))
.and_then(|v| v.checked_add(3))
.ok_or_else(sequence_error)
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn sequence_error() -> AnalysisError {
AnalysisError::Progress("analysis job sequence overflow".into())
}
#[cfg(any(feature = "providers", feature = "adapter-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 index = labels
.iter()
.position(|v| *v == speaker)
.ok_or_else(|| format!("Gemini reported an unknown feature speaker {speaker}"))?;
u64::try_from(index)
.ok()
.and_then(|v| v.checked_mul(3))
.and_then(|v| v.checked_add(u64::from(packet)))
.and_then(|v| v.checked_add(2))
.ok_or_else(|| sequence_error().to_string())
}
#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
fn map_feature_progress(
labels: &[LocalSpeakerLabel],
event: GeminiFeatureProgress,
) -> Result<AnalysisProgress, String> {
match event {
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 error = AnalysisError::GeminiFeature {
speaker,
packet,
message: error,
}
.to_string();
Ok(AnalysisProgress::JobFailed {
sequence: feature_sequence(labels, speaker, packet)?,
error,
})
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use futures::{executor::block_on, future::poll_fn, join};
use std::{
sync::{
Arc, Mutex,
atomic::{AtomicBool, Ordering},
},
task::Poll,
};
const TRANSCRIPT: &str = "[high] Speaker 1: exact transcript";
struct State {
labels: Vec<LocalSpeakerLabel>,
analysis: StructuredAnalysis,
fail: Option<&'static str>,
calls: Mutex<Vec<&'static str>>,
targets: Mutex<Vec<LocalSpeakerLabel>>,
final_evidence: Mutex<Vec<SpeakerFeatureEvidence>>,
wait: Option<Arc<AtomicBool>>,
mark: Option<Arc<AtomicBool>>,
}
#[derive(Clone)]
struct Fake(Arc<State>);
impl Fake {
fn new(count: u32) -> Self {
Self::configured(count, analysis(TRANSCRIPT, count), None, None, None)
}
fn configured(
count: u32,
result: StructuredAnalysis,
fail: Option<&'static str>,
wait: Option<Arc<AtomicBool>>,
mark: Option<Arc<AtomicBool>>,
) -> Self {
Self(Arc::new(State {
labels: (1..=count).map(label).collect(),
analysis: result,
fail,
calls: Mutex::new(Vec::new()),
targets: Mutex::new(Vec::new()),
final_evidence: Mutex::new(Vec::new()),
wait,
mark,
}))
}
fn calls(&self) -> Vec<&'static str> {
self.0.calls.lock().unwrap().clone()
}
}
impl AnalysisOperations for Fake {
fn transcript<'a>(
&'a self,
_: &'a [u8],
) -> AnalysisFuture<'a, Result<String, AnalysisError>> {
Box::pin(async move {
self.0.calls.lock().unwrap().push("transcript");
if let Some(gate) = &self.0.wait {
poll_fn(|cx| {
if gate.load(Ordering::SeqCst) {
Poll::Ready(())
} else {
cx.waker().wake_by_ref();
Poll::Pending
}
})
.await;
}
if self.0.fail == Some("transcript") {
Err(AnalysisError::GeminiTranscript("failed".into()))
} else {
Ok(TRANSCRIPT.into())
}
})
}
fn speaker_labels<'a>(
&'a self,
_: &'a str,
) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>> {
Box::pin(async move {
self.0.calls.lock().unwrap().push("labels");
if self.0.fail == Some("labels") {
Err(AnalysisError::TerraLabels("failed".into()))
} else {
Ok(self.0.labels.clone())
}
})
}
fn feature_evidence_with_progress<'a>(
&'a self,
_: &'a [u8],
_: &'a str,
labels: &'a [LocalSpeakerLabel],
report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>> {
Box::pin(async move {
self.0.calls.lock().unwrap().push("features");
*self.0.targets.lock().unwrap() = labels.to_vec();
if self.0.fail == Some("features") {
return Err(AnalysisError::GeminiCache("failed".into()));
}
for &speaker in labels {
for packet in 1..=3 {
report(GeminiFeatureProgress::Started { speaker, packet })
.map_err(AnalysisError::Progress)?;
report(GeminiFeatureProgress::Succeeded { speaker, packet })
.map_err(AnalysisError::Progress)?;
}
}
Ok(labels.iter().copied().map(evidence).collect())
})
}
fn structured_analysis<'a>(
&'a self,
_: &'a str,
evidence: Vec<SpeakerFeatureEvidence>,
) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>> {
Box::pin(async move {
self.0.calls.lock().unwrap().push("final");
*self.0.final_evidence.lock().unwrap() = evidence;
if self.0.fail == Some("final") {
return Err(AnalysisError::TerraStructuring("failed".into()));
}
if let Some(mark) = &self.0.mark {
mark.store(true, Ordering::SeqCst);
}
Ok(self.0.analysis.clone())
})
}
}
fn label(n: u32) -> LocalSpeakerLabel {
LocalSpeakerLabel::new(n).unwrap()
}
fn evidence(speaker: LocalSpeakerLabel) -> SpeakerFeatureEvidence {
SpeakerFeatureEvidence::new(
speaker,
format!("{speaker}-1"),
format!("{speaker}-2"),
format!("{speaker}-3"),
)
.unwrap()
}
fn analysis(transcript: &str, count: u32) -> StructuredAnalysis {
StructuredAnalysis {
transcript: transcript.into(),
speakers: (1..=count)
.map(|n| StructuredSpeaker {
speaker: label(n),
language: "English".into(),
features: FeatureVector24::default(),
features_usable_for_training: false,
})
.collect(),
}
}
fn ogg() -> Vec<u8> {
let mut value = vec![0; 28];
value[..4].copy_from_slice(b"OggS");
value[26] = 1;
value
}
fn audio(bytes: &[u8]) -> OggAudioMetadata {
OggAudioMetadata::from_bytes(bytes, 1, None).unwrap()
}
fn executed(bytes: &[u8], count: u32) -> ExecutedAnalysis {
ExecutedAnalysis {
envelope: AnalysisEnvelope {
audio: audio(bytes),
analysis: analysis(TRANSCRIPT, count),
gemini: GeminiCohort::new(GEMINI_MODEL_ID),
structurer: StructurerProvenance::new(TERRA_MODEL_ID),
},
label_extractor: StructurerProvenance {
model_id: TERRA_MODEL_ID.into(),
prompt_revision:
kcode_speaker_v3_llm_protocol::TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
},
}
}
fn resume(
fake: &Fake,
recovered: RecoveredAnalysisArtifacts,
) -> (
Result<ExecutedAnalysis, ResumableAnalysisError>,
Vec<AnalysisProgress>,
Vec<AnalysisArtifact>,
) {
let bytes = ogg();
let mut progress = Vec::new();
let mut artifacts = Vec::new();
let result = block_on(execute_admitted(
fake,
&bytes,
audio(&bytes),
recovered,
|v| {
progress.push(v);
Ok(())
},
|v| {
artifacts.push(v);
Ok(())
},
true,
));
(result, progress, artifacts)
}
fn recovered(count: u32) -> RecoveredAnalysisArtifacts {
RecoveredAnalysisArtifacts {
transcript: Some(TRANSCRIPT.into()),
speaker_labels: Some((1..=count).map(label).collect()),
feature_evidence: (1..=count).map(|n| evidence(label(n))).collect(),
..Default::default()
}
}
#[test]
fn all_missing_matches_old_flow_and_legacy_contract() {
let bytes = ogg();
let old = Fake::new(1);
let mut old_events = Vec::new();
let old_result = block_on(execute_admitted(
&old,
&bytes,
audio(&bytes),
RecoveredAnalysisArtifacts::default(),
|v| {
old_events.push(v);
Ok(())
},
|_| Ok(()),
false,
))
.map_err(non_resumable_error);
let new = Fake::new(1);
let (new_result, new_events, artifacts) =
resume(&new, RecoveredAnalysisArtifacts::default());
assert_eq!(old_result, new_result.map_err(non_resumable_error));
assert_eq!(old_events, new_events);
assert_eq!(old.calls(), new.calls());
assert!(matches!(
&artifacts[..],
[
AnalysisArtifact::GeminiTranscript(_),
AnalysisArtifact::TerraSpeakerLabels(_),
AnalysisArtifact::GeminiFeatureEvidence(_),
AnalysisArtifact::ExecutedAnalysis(_)
]
));
let starts = new_events
.iter()
.filter_map(|v| {
if let AnalysisProgress::JobStarted { sequence, .. } = v {
Some(*sequence)
} else {
None
}
})
.collect::<Vec<_>>();
assert_eq!(starts, [1, 2, 3, 4, 5, 6]);
let legacy = Fake::new(0);
let result = block_on(execute_legacy(
&legacy,
&bytes,
1234,
Some("voice.ogg".into()),
))
.unwrap();
assert_eq!(result.envelope.audio.filename(), Some("voice.ogg"));
assert_eq!(
legacy.calls(),
["transcript", "labels", "features", "final"]
);
let invalid = Fake::new(1);
let mut reports = 0;
assert!(matches!(
block_on(execute_strict(&invalid, b"bad", |_| {
reports += 1;
Ok(())
})),
Err(AnalysisError::Input(_))
));
assert_eq!(reports, 0);
assert!(invalid.calls().is_empty());
}
#[test]
fn recovered_and_partial_stages_are_selective_and_preserved() {
let full = Fake::new(1);
let (result, events, artifacts) = resume(&full, recovered(1));
result.unwrap();
assert_eq!(full.calls(), ["final"]);
assert!(matches!(
&artifacts[..],
[AnalysisArtifact::ExecutedAnalysis(_)]
));
assert_eq!(events.len(), 3);
let transcript = Fake::new(1);
let (_, _, artifacts) = resume(
&transcript,
RecoveredAnalysisArtifacts {
transcript: Some(TRANSCRIPT.into()),
..Default::default()
},
);
assert_eq!(transcript.calls(), ["labels", "features", "final"]);
assert_eq!(artifacts.len(), 3);
let labels = Fake::new(1);
let (_, _, artifacts) = resume(
&labels,
RecoveredAnalysisArtifacts {
speaker_labels: Some(vec![label(1)]),
..Default::default()
},
);
assert_eq!(labels.calls(), ["transcript", "features", "final"]);
assert!(matches!(
&artifacts[..],
[
AnalysisArtifact::GeminiTranscript(_),
AnalysisArtifact::GeminiFeatureEvidence(_),
AnalysisArtifact::ExecutedAnalysis(_)
]
));
let packet = evidence(label(1));
let downstream = Fake::new(1);
let (_, _, artifacts) = resume(
&downstream,
RecoveredAnalysisArtifacts {
feature_evidence: vec![packet.clone()],
..Default::default()
},
);
assert_eq!(downstream.calls(), ["transcript", "labels", "final"]);
assert_eq!(*downstream.0.final_evidence.lock().unwrap(), [packet]);
assert_eq!(artifacts.len(), 3);
let first = evidence(label(1));
let partial = Fake::new(2);
let (_, _, artifacts) = resume(
&partial,
RecoveredAnalysisArtifacts {
transcript: Some(TRANSCRIPT.into()),
speaker_labels: Some(vec![label(1), label(2)]),
feature_evidence: vec![first.clone()],
..Default::default()
},
);
assert_eq!(partial.calls(), ["features", "final"]);
assert_eq!(*partial.0.targets.lock().unwrap(), [label(2)]);
assert!(
matches!(&artifacts[..], [AnalysisArtifact::GeminiFeatureEvidence(v), AnalysisArtifact::ExecutedAnalysis(_)] if v.speaker() == label(2))
);
assert_eq!(partial.0.final_evidence.lock().unwrap()[0], first);
}
#[test]
fn recovered_final_repairs_missing_stages_without_structuring() {
let bytes = ogg();
let final_value = executed(&bytes, 2);
let fake = Fake::new(2);
let (result, events, artifacts) = resume(
&fake,
RecoveredAnalysisArtifacts {
feature_evidence: vec![evidence(label(1))],
executed_analysis: Some(final_value.clone()),
..Default::default()
},
);
assert_eq!(result.unwrap(), final_value);
assert_eq!(fake.calls(), ["labels", "features"]);
assert_eq!(*fake.0.targets.lock().unwrap(), [label(2)]);
assert!(
matches!(&artifacts[..], [AnalysisArtifact::TerraSpeakerLabels(_), AnalysisArtifact::GeminiFeatureEvidence(value)] if value.speaker() == label(2))
);
assert!(!events.iter().any(|event| matches!(
event,
AnalysisProgress::JobStarted {
job: AnalysisJob::Structuring,
..
} | AnalysisProgress::StageCompleted {
stage: AnalysisStage::Structuring
}
)));
let mismatch = Fake::new(2);
let (_, _, artifacts) = resume(
&mismatch,
RecoveredAnalysisArtifacts {
executed_analysis: Some(executed(&bytes, 1)),
..Default::default()
},
);
assert_eq!(mismatch.calls(), ["labels"]);
assert!(artifacts.is_empty());
}
#[test]
fn recovered_final_supplies_transcript_and_is_returned_exactly() {
let bytes = ogg();
let final_value = executed(&bytes, 1);
let fake = Fake::configured(1, analysis(TRANSCRIPT, 1), Some("transcript"), None, None);
let (result, events, artifacts) = resume(
&fake,
RecoveredAnalysisArtifacts {
speaker_labels: Some(vec![label(1)]),
feature_evidence: vec![evidence(label(1))],
executed_analysis: Some(final_value.clone()),
..Default::default()
},
);
assert_eq!(result.unwrap(), final_value);
assert!(fake.calls().is_empty());
assert!(events.is_empty());
assert!(artifacts.is_empty());
}
#[test]
fn final_recovery_conflicts_are_effect_free() {
let bytes = ogg();
let valid = executed(&bytes, 1);
let mut wrong_audio = valid.clone();
wrong_audio.envelope.audio = OggAudioMetadata::from_bytes(&bytes, 2, None).unwrap();
let mut bad_envelope = valid.clone();
bad_envelope.envelope.gemini.model_id.clear();
let mut bad_extractor = valid.clone();
bad_extractor.label_extractor.model_id.clear();
let bad = [
RecoveredAnalysisArtifacts {
executed_analysis: Some(wrong_audio),
..Default::default()
},
RecoveredAnalysisArtifacts {
transcript: Some("different".into()),
executed_analysis: Some(valid.clone()),
..Default::default()
},
RecoveredAnalysisArtifacts {
speaker_labels: Some(vec![label(2)]),
executed_analysis: Some(valid.clone()),
..Default::default()
},
RecoveredAnalysisArtifacts {
feature_evidence: vec![evidence(label(2))],
executed_analysis: Some(valid.clone()),
..Default::default()
},
RecoveredAnalysisArtifacts {
executed_analysis: Some(bad_envelope),
..Default::default()
},
RecoveredAnalysisArtifacts {
executed_analysis: Some(bad_extractor),
..Default::default()
},
];
for value in bad {
let fake = Fake::new(1);
let (result, events, artifacts) = resume(&fake, value);
assert!(matches!(result, Err(ResumableAnalysisError::Recovered(_))));
assert!(fake.calls().is_empty());
assert!(events.is_empty());
assert!(artifacts.is_empty());
}
}
#[test]
fn recovered_empty_labels_are_complete_and_zero_speaker_final_works() {
let fake = Fake::new(0);
let (result, events, artifacts) = resume(
&fake,
RecoveredAnalysisArtifacts {
transcript: Some(TRANSCRIPT.into()),
speaker_labels: Some(vec![]),
..Default::default()
},
);
result.unwrap();
assert_eq!(fake.calls(), ["final"]);
assert!(matches!(
&artifacts[..],
[AnalysisArtifact::ExecutedAnalysis(_)]
));
assert!(!events.iter().any(|event| matches!(
event,
AnalysisProgress::StageCompleted {
stage: AnalysisStage::SpeakerFeatures
}
)));
let bytes = ogg();
let final_value = executed(&bytes, 0);
let recovered = Fake::new(0);
let (result, events, artifacts) = resume(
&recovered,
RecoveredAnalysisArtifacts {
speaker_labels: Some(vec![]),
executed_analysis: Some(final_value.clone()),
..Default::default()
},
);
assert_eq!(result.unwrap(), final_value);
assert!(recovered.calls().is_empty());
assert!(events.is_empty());
assert!(artifacts.is_empty());
}
#[test]
fn artifact_failure_fences_the_next_stage() {
let fake = Fake::new(1);
let bytes = ogg();
let mut events = Vec::new();
let result = block_on(execute_admitted(
&fake,
&bytes,
audio(&bytes),
RecoveredAnalysisArtifacts::default(),
|v| {
events.push(v);
Ok(())
},
|_| Err("durability unavailable".into()),
true,
));
assert_eq!(
result,
Err(ResumableAnalysisError::Artifact(
"durability unavailable".into()
))
);
assert_eq!(fake.calls(), ["transcript"]);
assert!(
!events
.iter()
.any(|v| matches!(v, AnalysisProgress::StageCompleted { .. }))
);
}
#[test]
fn malformed_recovery_is_effect_free_and_generated_final_conflicts_fail() {
let bad = [
RecoveredAnalysisArtifacts {
transcript: Some(" ".into()),
..Default::default()
},
RecoveredAnalysisArtifacts {
speaker_labels: Some(vec![label(1), label(1)]),
..Default::default()
},
RecoveredAnalysisArtifacts {
speaker_labels: Some(vec![label(1)]),
feature_evidence: vec![evidence(label(2))],
..Default::default()
},
RecoveredAnalysisArtifacts {
feature_evidence: vec![evidence(label(1)), evidence(label(1))],
..Default::default()
},
];
for value in bad {
let fake = Fake::new(1);
assert!(matches!(
resume(&fake, value).0,
Err(ResumableAnalysisError::Recovered(_))
));
assert!(fake.calls().is_empty());
}
let transcript = Fake::configured(1, analysis("different", 1), None, None, None);
assert_eq!(
resume(&transcript, recovered(1)).0,
Err(ResumableAnalysisError::Analysis(
AnalysisError::TranscriptMismatch
))
);
let speakers = Fake::configured(1, analysis(TRANSCRIPT, 2), None, None, None);
assert_eq!(
resume(&speakers, recovered(1)).0,
Err(ResumableAnalysisError::Analysis(
AnalysisError::SpeakerSetMismatch
))
);
}
#[test]
fn failures_do_not_retry_and_analyses_remain_concurrent() {
for (stage, calls) in [
("transcript", vec!["transcript"]),
("labels", vec!["transcript", "labels"]),
("features", vec!["transcript", "labels", "features"]),
("final", vec!["transcript", "labels", "features", "final"]),
] {
let fake = Fake::configured(1, analysis(TRANSCRIPT, 1), Some(stage), None, None);
assert!(
resume(&fake, RecoveredAnalysisArtifacts::default())
.0
.is_err()
);
assert_eq!(fake.calls(), calls);
}
let done = Arc::new(AtomicBool::new(false));
let slow = Fake::configured(1, analysis(TRANSCRIPT, 1), None, Some(done.clone()), None);
let fast = Fake::configured(1, analysis(TRANSCRIPT, 1), None, None, Some(done.clone()));
let bytes = ogg();
let (a, b) = block_on(async {
join!(
execute_admitted(
&slow,
&bytes,
audio(&bytes),
RecoveredAnalysisArtifacts::default(),
|_| Ok(()),
|_| Ok(()),
true
),
execute_admitted(
&fast,
&bytes,
audio(&bytes),
RecoveredAnalysisArtifacts::default(),
|_| Ok(()),
|_| Ok(()),
true
)
)
});
a.unwrap();
b.unwrap();
assert!(done.load(Ordering::SeqCst));
}
#[test]
fn public_contract_and_provider_operations_compile() {
let cohort = GeminiCohort::new("gemini");
cohort.validate().unwrap();
StructurerProvenance::new("terra").validate().unwrap();
fn require<O: AnalysisOperations>() {}
require::<ProviderOperations>();
let wrapped: ResumableAnalysisError = AnalysisError::Input("bad".into()).into();
assert!(matches!(
wrapped,
ResumableAnalysisError::Analysis(AnalysisError::Input(_))
));
}
#[cfg(feature = "providers")]
#[test]
fn legacy_constructor_compiles() {
let _: fn(kcode_gemini_3_1_pro::Gemini31Pro, kcode_codex_terra::CodexTerra) -> Analyzer =
Analyzer::new;
}
#[cfg(feature = "adapter-providers")]
#[test]
fn adapter_constructor_compiles() {
let _: fn(kcode_gemini_3_1_pro::Gemini31Pro, kcode_k1_codex_adapter::Adapter) -> Analyzer =
Analyzer::from_codex_adapter;
}
}