Skip to main content

kcode_speaker_v3_analysis/
lib.rs

1#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
2use kcode_speaker_v3_gemini_analysis::GeminiFeatureProgress;
3pub use kcode_speaker_v3_llm_protocol::{
4    GEMINI_FEATURE_PROMPT_ONE, GEMINI_FEATURE_PROMPT_ONE_REVISION, GEMINI_FEATURE_PROMPT_REVISIONS,
5    GEMINI_FEATURE_PROMPT_THREE, GEMINI_FEATURE_PROMPT_THREE_REVISION, GEMINI_FEATURE_PROMPT_TWO,
6    GEMINI_FEATURE_PROMPT_TWO_REVISION, GEMINI_TRANSCRIPT_PROMPT,
7    GEMINI_TRANSCRIPT_PROMPT_REVISION, GPT_STRUCTURING_PROMPT, GPT_STRUCTURING_PROMPT_REVISION,
8    SpeakerFeatureEvidence,
9};
10pub use kcode_speaker_v3_schema::{
11    FEATURE_NAMES, FEATURE_SCHEMA_REVISION, FeatureVector24, LocalSpeakerLabel,
12    MAX_AUDIO_DURATION_MS, OGG_MEDIA_TYPE, OggAudioMetadata, StructuredAnalysis, StructuredSpeaker,
13    ValidationError, VocalGenderPresentation,
14};
15use serde::{Deserialize, Serialize};
16use std::{collections::BTreeSet, error::Error, fmt};
17#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
18use std::{future::Future, pin::Pin};
19const GEMINI_MODEL_ID: &str = "gemini-3.1-pro-preview";
20const TERRA_MODEL_ID: &str = "gpt-5.6-terra";
21
22#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
23pub struct GeminiCohort {
24    pub model_id: String,
25    pub transcript_prompt_revision: String,
26    pub feature_prompt_revisions: [String; 3],
27    pub feature_schema_revision: String,
28}
29impl GeminiCohort {
30    pub fn new(model_id: impl Into<String>) -> Self {
31        Self {
32            model_id: model_id.into(),
33            transcript_prompt_revision: GEMINI_TRANSCRIPT_PROMPT_REVISION.into(),
34            feature_prompt_revisions: GEMINI_FEATURE_PROMPT_REVISIONS.map(str::to_owned),
35            feature_schema_revision: FEATURE_SCHEMA_REVISION.into(),
36        }
37    }
38    pub fn validate(&self) -> Result<(), ValidationError> {
39        validate_text(&self.model_id, "gemini_model_id")?;
40        validate_text(
41            &self.transcript_prompt_revision,
42            "transcript_prompt_revision",
43        )?;
44        for value in &self.feature_prompt_revisions {
45            validate_text(value, "feature_prompt_revision")?;
46        }
47        validate_text(&self.feature_schema_revision, "feature_schema_revision")
48    }
49}
50#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
51pub struct StructurerProvenance {
52    pub model_id: String,
53    pub prompt_revision: String,
54}
55impl StructurerProvenance {
56    pub fn new(model_id: impl Into<String>) -> Self {
57        Self {
58            model_id: model_id.into(),
59            prompt_revision: GPT_STRUCTURING_PROMPT_REVISION.into(),
60        }
61    }
62    pub fn validate(&self) -> Result<(), ValidationError> {
63        validate_text(&self.model_id, "structurer_model_id")?;
64        validate_text(&self.prompt_revision, "structurer_prompt_revision")
65    }
66}
67#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
68pub struct AnalysisEnvelope {
69    pub audio: OggAudioMetadata,
70    pub analysis: StructuredAnalysis,
71    pub gemini: GeminiCohort,
72    pub structurer: StructurerProvenance,
73}
74impl AnalysisEnvelope {
75    pub fn validate(&self) -> Result<(), ValidationError> {
76        self.audio.validate()?;
77        self.analysis.validate()?;
78        self.gemini.validate()?;
79        self.structurer.validate()
80    }
81}
82#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
83pub struct ExecutedAnalysis {
84    pub envelope: AnalysisEnvelope,
85    pub label_extractor: StructurerProvenance,
86}
87#[derive(Debug, Clone, PartialEq, Serialize, Default)]
88pub struct RecoveredAnalysisArtifacts {
89    pub transcript: Option<String>,
90    pub speaker_labels: Option<Vec<LocalSpeakerLabel>>,
91    pub feature_evidence: Vec<SpeakerFeatureEvidence>,
92    pub executed_analysis: Option<ExecutedAnalysis>,
93}
94#[derive(Debug, Clone, PartialEq, Serialize)]
95pub enum AnalysisArtifact {
96    GeminiTranscript(String),
97    TerraSpeakerLabels(Vec<LocalSpeakerLabel>),
98    GeminiFeatureEvidence(SpeakerFeatureEvidence),
99    ExecutedAnalysis(Box<ExecutedAnalysis>),
100}
101#[derive(Debug, Clone, PartialEq, Eq)]
102pub enum AnalysisStage {
103    Transcript,
104    SpeakerLabels,
105    SpeakerFeatures,
106    Structuring,
107}
108#[derive(Debug, Clone, PartialEq, Eq)]
109pub enum AnalysisJob {
110    Transcript,
111    SpeakerLabels,
112    SpeakerFeature {
113        speaker: LocalSpeakerLabel,
114        packet: u8,
115    },
116    Structuring,
117}
118#[derive(Debug, Clone, PartialEq, Eq)]
119pub enum AnalysisProgress {
120    JobStarted { sequence: u64, job: AnalysisJob },
121    JobSucceeded { sequence: u64 },
122    JobFailed { sequence: u64, error: String },
123    StageCompleted { stage: AnalysisStage },
124}
125#[derive(Debug, Clone, PartialEq, Eq)]
126pub enum AnalysisError {
127    Input(String),
128    Progress(String),
129    GeminiTranscript(String),
130    TerraLabels(String),
131    GeminiCache(String),
132    GeminiFeature {
133        speaker: LocalSpeakerLabel,
134        packet: u8,
135        message: String,
136    },
137    TerraStructuring(String),
138    TranscriptMismatch,
139    SpeakerSetMismatch,
140}
141impl fmt::Display for AnalysisError {
142    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
143        match self {
144            Self::Input(v) => write!(f, "invalid input: {v}"),
145            Self::Progress(v) => write!(f, "progress reporting failed: {v}"),
146            Self::GeminiTranscript(v) => write!(f, "Gemini transcript failed: {v}"),
147            Self::TerraLabels(v) => write!(f, "Terra speaker-label extraction failed: {v}"),
148            Self::GeminiCache(v) => write!(f, "Gemini feature cache creation failed: {v}"),
149            Self::GeminiFeature {
150                speaker,
151                packet,
152                message,
153            } => write!(
154                f,
155                "Gemini feature call failed for {speaker}, packet {packet}: {message}"
156            ),
157            Self::TerraStructuring(v) => write!(f, "Terra final structuring failed: {v}"),
158            Self::TranscriptMismatch => f.write_str("Terra returned a different transcript"),
159            Self::SpeakerSetMismatch => f.write_str("Terra returned a different speaker set"),
160        }
161    }
162}
163impl Error for AnalysisError {}
164#[derive(Debug, Clone, PartialEq, Eq)]
165pub enum ResumableAnalysisError {
166    Analysis(AnalysisError),
167    Recovered(String),
168    Artifact(String),
169}
170impl fmt::Display for ResumableAnalysisError {
171    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
172        match self {
173            Self::Analysis(error) => write!(f, "{error}"),
174            Self::Recovered(v) => write!(f, "invalid recovered artifacts: {v}"),
175            Self::Artifact(v) => write!(f, "artifact reporting failed: {v}"),
176        }
177    }
178}
179impl Error for ResumableAnalysisError {
180    fn source(&self) -> Option<&(dyn Error + 'static)> {
181        match self {
182            Self::Analysis(error) => Some(error),
183            Self::Recovered(_) | Self::Artifact(_) => None,
184        }
185    }
186}
187impl From<AnalysisError> for ResumableAnalysisError {
188    fn from(value: AnalysisError) -> Self {
189        Self::Analysis(value)
190    }
191}
192fn validate_text(value: &str, field: &'static str) -> Result<(), ValidationError> {
193    (!value.trim().is_empty())
194        .then_some(())
195        .ok_or(ValidationError::Blank(field))
196}
197
198#[cfg(any(feature = "providers", feature = "adapter-providers"))]
199pub struct Analyzer {
200    operations: ProviderOperations,
201}
202#[cfg(any(feature = "providers", feature = "adapter-providers"))]
203impl Analyzer {
204    #[cfg(feature = "providers")]
205    pub fn new(
206        gemini: kcode_gemini_3_1_pro::Gemini31Pro,
207        terra: kcode_codex_terra::CodexTerra,
208    ) -> Self {
209        Self {
210            operations: ProviderOperations {
211                gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis::new(gemini),
212                terra: kcode_speaker_v3_terra_analysis::TerraAnalysis::new(terra),
213            },
214        }
215    }
216    #[cfg(feature = "adapter-providers")]
217    pub fn from_codex_adapter(
218        gemini: kcode_gemini_3_1_pro::Gemini31Pro,
219        adapter: kcode_k1_codex_adapter::Adapter,
220    ) -> Self {
221        Self {
222            operations: ProviderOperations {
223                gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis::new(gemini),
224                terra: kcode_speaker_v3_terra_analysis::TerraAnalysis::from_codex_adapter(adapter),
225            },
226        }
227    }
228    pub async fn analyze_ogg_with_progress<F>(
229        &self,
230        bytes: &[u8],
231        report: F,
232    ) -> Result<ExecutedAnalysis, AnalysisError>
233    where
234        F: FnMut(AnalysisProgress) -> Result<(), String>,
235    {
236        execute_strict(&self.operations, bytes, report).await
237    }
238    pub async fn analyze_ogg(
239        &self,
240        bytes: &[u8],
241        duration_ms: u64,
242        filename: Option<String>,
243    ) -> Result<ExecutedAnalysis, AnalysisError> {
244        execute_legacy(&self.operations, bytes, duration_ms, filename).await
245    }
246    pub async fn analyze_ogg_resumable<P, A>(
247        &self,
248        bytes: &[u8],
249        recovered: RecoveredAnalysisArtifacts,
250        report: P,
251        emit: A,
252    ) -> Result<ExecutedAnalysis, ResumableAnalysisError>
253    where
254        P: FnMut(AnalysisProgress) -> Result<(), String>,
255        A: FnMut(AnalysisArtifact) -> Result<(), String>,
256    {
257        execute_resume_strict(&self.operations, bytes, recovered, report, emit).await
258    }
259}
260
261#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
262type AnalysisFuture<'a, T> = Pin<Box<dyn Future<Output = T> + 'a>>;
263#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
264trait AnalysisOperations: Sync {
265    fn transcript<'a>(
266        &'a self,
267        audio: &'a [u8],
268    ) -> AnalysisFuture<'a, Result<String, AnalysisError>>;
269    fn speaker_labels<'a>(
270        &'a self,
271        transcript: &'a str,
272    ) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>>;
273    fn feature_evidence_with_progress<'a>(
274        &'a self,
275        audio: &'a [u8],
276        transcript: &'a str,
277        labels: &'a [LocalSpeakerLabel],
278        report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
279    ) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>>;
280    fn structured_analysis<'a>(
281        &'a self,
282        transcript: &'a str,
283        evidence: Vec<SpeakerFeatureEvidence>,
284    ) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>>;
285}
286#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
287struct ProviderOperations {
288    gemini: kcode_speaker_v3_gemini_analysis::GeminiAnalysis,
289    terra: kcode_speaker_v3_terra_analysis::TerraAnalysis,
290}
291#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
292impl AnalysisOperations for ProviderOperations {
293    fn transcript<'a>(
294        &'a self,
295        audio: &'a [u8],
296    ) -> AnalysisFuture<'a, Result<String, AnalysisError>> {
297        Box::pin(async move {
298            self.gemini
299                .transcript(audio)
300                .await
301                .map_err(|error| match error {
302                    kcode_speaker_v3_gemini_analysis::GeminiTranscriptError::Provider(v)
303                    | kcode_speaker_v3_gemini_analysis::GeminiTranscriptError::Protocol(v) => {
304                        AnalysisError::GeminiTranscript(v)
305                    }
306                })
307        })
308    }
309    fn speaker_labels<'a>(
310        &'a self,
311        transcript: &'a str,
312    ) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>> {
313        Box::pin(async move {
314            self.terra
315                .speaker_labels(transcript)
316                .await
317                .map_err(|error| match error {
318                    kcode_speaker_v3_terra_analysis::TerraAnalysisError::Protocol(v)
319                    | kcode_speaker_v3_terra_analysis::TerraAnalysisError::Provider(v) => {
320                        AnalysisError::TerraLabels(v)
321                    }
322                })
323        })
324    }
325    fn feature_evidence_with_progress<'a>(
326        &'a self,
327        audio: &'a [u8],
328        transcript: &'a str,
329        labels: &'a [LocalSpeakerLabel],
330        report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
331    ) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>> {
332        Box::pin(async move {
333            self.gemini
334                .feature_evidence_with_progress(audio, transcript, labels, report)
335                .await
336                .map_err(|error| match error {
337                    kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Cache(v) => {
338                        AnalysisError::GeminiCache(v)
339                    }
340                    kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Progress(v) => {
341                        AnalysisError::Progress(v)
342                    }
343                    kcode_speaker_v3_gemini_analysis::GeminiFeatureError::Feature {
344                        speaker,
345                        packet,
346                        message,
347                    } => AnalysisError::GeminiFeature {
348                        speaker,
349                        packet,
350                        message,
351                    },
352                })
353        })
354    }
355    fn structured_analysis<'a>(
356        &'a self,
357        transcript: &'a str,
358        evidence: Vec<SpeakerFeatureEvidence>,
359    ) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>> {
360        Box::pin(async move {
361            self.terra
362                .structured_analysis(transcript, evidence)
363                .await
364                .map_err(|error| match error {
365                    kcode_speaker_v3_terra_analysis::TerraAnalysisError::Protocol(v)
366                    | kcode_speaker_v3_terra_analysis::TerraAnalysisError::Provider(v) => {
367                        AnalysisError::TerraStructuring(v)
368                    }
369                })
370        })
371    }
372}
373
374#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
375async fn execute_strict<O: AnalysisOperations, P: FnMut(AnalysisProgress) -> Result<(), String>>(
376    operations: &O,
377    bytes: &[u8],
378    report: P,
379) -> Result<ExecutedAnalysis, AnalysisError> {
380    let audio =
381        OggAudioMetadata::from_ogg_bytes(bytes).map_err(|e| AnalysisError::Input(e.to_string()))?;
382    execute_admitted(
383        operations,
384        bytes,
385        audio,
386        RecoveredAnalysisArtifacts::default(),
387        report,
388        |_| Ok(()),
389        false,
390    )
391    .await
392    .map_err(non_resumable_error)
393}
394#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
395async fn execute_legacy<O: AnalysisOperations>(
396    operations: &O,
397    bytes: &[u8],
398    duration_ms: u64,
399    filename: Option<String>,
400) -> Result<ExecutedAnalysis, AnalysisError> {
401    let audio = OggAudioMetadata::from_bytes(bytes, duration_ms, filename)
402        .map_err(|e| AnalysisError::Input(e.to_string()))?;
403    execute_admitted(
404        operations,
405        bytes,
406        audio,
407        RecoveredAnalysisArtifacts::default(),
408        |_| Ok(()),
409        |_| Ok(()),
410        false,
411    )
412    .await
413    .map_err(non_resumable_error)
414}
415#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
416async fn execute_resume_strict<
417    O: AnalysisOperations,
418    P: FnMut(AnalysisProgress) -> Result<(), String>,
419    A: FnMut(AnalysisArtifact) -> Result<(), String>,
420>(
421    operations: &O,
422    bytes: &[u8],
423    recovered: RecoveredAnalysisArtifacts,
424    report: P,
425    emit: A,
426) -> Result<ExecutedAnalysis, ResumableAnalysisError> {
427    let audio =
428        OggAudioMetadata::from_ogg_bytes(bytes).map_err(|e| AnalysisError::Input(e.to_string()))?;
429    execute_admitted(operations, bytes, audio, recovered, report, emit, true).await
430}
431
432#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
433async fn execute_admitted<O, P, A>(
434    operations: &O,
435    bytes: &[u8],
436    audio: OggAudioMetadata,
437    mut recovered: RecoveredAnalysisArtifacts,
438    mut report: P,
439    mut emit: A,
440    resumable: bool,
441) -> Result<ExecutedAnalysis, ResumableAnalysisError>
442where
443    O: AnalysisOperations,
444    P: FnMut(AnalysisProgress) -> Result<(), String>,
445    A: FnMut(AnalysisArtifact) -> Result<(), String>,
446{
447    let (recovered_speakers, final_speakers) = if resumable {
448        validate_recovered(&recovered, &audio)?
449    } else {
450        (BTreeSet::new(), None)
451    };
452    let transcript = if let Some(value) = recovered.transcript.take() {
453        value
454    } else if let Some(executed) = recovered.executed_analysis.as_ref() {
455        executed.envelope.analysis.transcript.clone()
456    } else {
457        progress(
458            &mut report,
459            AnalysisProgress::JobStarted {
460                sequence: 1,
461                job: AnalysisJob::Transcript,
462            },
463        )?;
464        let value = match operations.transcript(bytes).await {
465            Ok(v) => v,
466            Err(e) => {
467                failed(&mut report, 1, &e)?;
468                return Err(e.into());
469            }
470        };
471        if let Err(error) = validate_generated_transcript(&value, resumable) {
472            failed(&mut report, 1, &error)?;
473            return Err(error.into());
474        }
475        progress(&mut report, AnalysisProgress::JobSucceeded { sequence: 1 })?;
476        if resumable {
477            emit(AnalysisArtifact::GeminiTranscript(value.clone()))
478                .map_err(ResumableAnalysisError::Artifact)?;
479        }
480        progress(
481            &mut report,
482            AnalysisProgress::StageCompleted {
483                stage: AnalysisStage::Transcript,
484            },
485        )?;
486        value
487    };
488    let labels = if let Some(value) = recovered.speaker_labels.take() {
489        value
490    } else {
491        progress(
492            &mut report,
493            AnalysisProgress::JobStarted {
494                sequence: 2,
495                job: AnalysisJob::SpeakerLabels,
496            },
497        )?;
498        let value = match operations.speaker_labels(&transcript).await {
499            Ok(v) => v,
500            Err(e) => {
501                failed(&mut report, 2, &e)?;
502                return Err(e.into());
503            }
504        };
505        if let Err(error) = validate_generated_labels(
506            &value,
507            &recovered_speakers,
508            final_speakers.as_ref(),
509            resumable,
510        ) {
511            failed(&mut report, 2, &error)?;
512            return Err(error);
513        }
514        progress(&mut report, AnalysisProgress::JobSucceeded { sequence: 2 })?;
515        if resumable {
516            emit(AnalysisArtifact::TerraSpeakerLabels(value.clone()))
517                .map_err(ResumableAnalysisError::Artifact)?;
518        }
519        progress(
520            &mut report,
521            AnalysisProgress::StageCompleted {
522                stage: AnalysisStage::SpeakerLabels,
523            },
524        )?;
525        value
526    };
527    let label_set = labels.iter().copied().collect::<BTreeSet<_>>();
528    let sequence = structuring_sequence(labels.len())?;
529    let missing = labels
530        .iter()
531        .copied()
532        .filter(|label| !recovered_speakers.contains(label))
533        .collect::<Vec<_>>();
534    let targets = if resumable { missing } else { labels.clone() };
535    let produced = if !resumable || !targets.is_empty() {
536        let mut feature_report = |event| report(map_feature_progress(&labels, event)?);
537        let value = operations
538            .feature_evidence_with_progress(bytes, &transcript, &targets, &mut feature_report)
539            .await?;
540        if resumable {
541            let returned =
542                validate_evidence(&value).map_err(|message| feature_error(targets[0], message))?;
543            let expected = targets.iter().copied().collect::<BTreeSet<_>>();
544            if returned != expected {
545                return Err(feature_error(
546                    targets[0],
547                    "generated feature speaker set differs from requested labels".into(),
548                )
549                .into());
550            }
551            for target in &targets {
552                let item = value
553                    .iter()
554                    .find(|item| item.speaker() == *target)
555                    .expect("validated feature set");
556                emit(AnalysisArtifact::GeminiFeatureEvidence(item.clone()))
557                    .map_err(ResumableAnalysisError::Artifact)?;
558            }
559        }
560        progress(
561            &mut report,
562            AnalysisProgress::StageCompleted {
563                stage: AnalysisStage::SpeakerFeatures,
564            },
565        )?;
566        value
567    } else {
568        Vec::new()
569    };
570    let evidence = if resumable {
571        recovered.feature_evidence.extend(produced);
572        let complete = validate_evidence(&recovered.feature_evidence)
573            .map_err(ResumableAnalysisError::Recovered)?;
574        if complete != label_set {
575            return Err(ResumableAnalysisError::Recovered(
576                "feature evidence does not cover the complete label set".into(),
577            ));
578        }
579        labels
580            .iter()
581            .map(|label| {
582                recovered
583                    .feature_evidence
584                    .iter()
585                    .find(|item| item.speaker() == *label)
586                    .expect("validated complete evidence")
587                    .clone()
588            })
589            .collect()
590    } else {
591        produced
592    };
593    if let Some(executed) = recovered.executed_analysis.take() {
594        return Ok(executed);
595    }
596    progress(
597        &mut report,
598        AnalysisProgress::JobStarted {
599            sequence,
600            job: AnalysisJob::Structuring,
601        },
602    )?;
603    let analysis = match operations.structured_analysis(&transcript, evidence).await {
604        Ok(v) => v,
605        Err(e) => {
606            failed(&mut report, sequence, &e)?;
607            return Err(e.into());
608        }
609    };
610    progress(&mut report, AnalysisProgress::JobSucceeded { sequence })?;
611    if analysis.transcript != transcript {
612        return Err(AnalysisError::TranscriptMismatch.into());
613    }
614    let returned = analysis
615        .speakers
616        .iter()
617        .map(|speaker| speaker.speaker)
618        .collect::<BTreeSet<_>>();
619    if returned != label_set {
620        return Err(AnalysisError::SpeakerSetMismatch.into());
621    }
622    let envelope = AnalysisEnvelope {
623        audio,
624        analysis,
625        gemini: GeminiCohort::new(GEMINI_MODEL_ID),
626        structurer: StructurerProvenance::new(TERRA_MODEL_ID),
627    };
628    envelope
629        .validate()
630        .map_err(|e| AnalysisError::TerraStructuring(e.to_string()))?;
631    let label_extractor = StructurerProvenance {
632        model_id: TERRA_MODEL_ID.into(),
633        prompt_revision: kcode_speaker_v3_llm_protocol::TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
634    };
635    label_extractor
636        .validate()
637        .map_err(|e| AnalysisError::TerraLabels(e.to_string()))?;
638    let executed = ExecutedAnalysis {
639        envelope,
640        label_extractor,
641    };
642    if resumable {
643        emit(AnalysisArtifact::ExecutedAnalysis(Box::new(
644            executed.clone(),
645        )))
646        .map_err(ResumableAnalysisError::Artifact)?;
647    }
648    progress(
649        &mut report,
650        AnalysisProgress::StageCompleted {
651            stage: AnalysisStage::Structuring,
652        },
653    )?;
654    Ok(executed)
655}
656
657#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
658fn validate_generated_transcript(value: &str, resumable: bool) -> Result<(), AnalysisError> {
659    if !resumable {
660        return Ok(());
661    }
662    validate_text(value, "transcript")
663        .map_err(|error| AnalysisError::GeminiTranscript(error.to_string()))
664}
665#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
666fn validate_generated_labels(
667    labels: &[LocalSpeakerLabel],
668    recovered: &BTreeSet<LocalSpeakerLabel>,
669    final_speakers: Option<&BTreeSet<LocalSpeakerLabel>>,
670    resumable: bool,
671) -> Result<(), ResumableAnalysisError> {
672    if !resumable {
673        return Ok(());
674    }
675    let labels = validate_labels(labels).map_err(AnalysisError::TerraLabels)?;
676    if !recovered.is_subset(&labels) {
677        return Err(ResumableAnalysisError::Recovered(
678            "feature evidence contains a speaker absent from generated labels".into(),
679        ));
680    }
681    if final_speakers.is_some_and(|expected| expected != &labels) {
682        return Err(ResumableAnalysisError::Recovered(
683            "generated labels differ from the recovered final speaker set".into(),
684        ));
685    }
686    Ok(())
687}
688#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
689fn validate_labels(labels: &[LocalSpeakerLabel]) -> Result<BTreeSet<LocalSpeakerLabel>, String> {
690    let set = labels.iter().copied().collect::<BTreeSet<_>>();
691    if set.len() != labels.len() {
692        return Err("speaker labels contain a duplicate".into());
693    }
694    Ok(set)
695}
696#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
697fn validate_evidence(
698    values: &[SpeakerFeatureEvidence],
699) -> Result<BTreeSet<LocalSpeakerLabel>, String> {
700    let mut set = BTreeSet::new();
701    for value in values {
702        let packets = value.packets();
703        SpeakerFeatureEvidence::new(
704            value.speaker(),
705            packets[0].into(),
706            packets[1].into(),
707            packets[2].into(),
708        )
709        .map_err(|e| e.to_string())?;
710        if !set.insert(value.speaker()) {
711            return Err(format!(
712                "duplicate feature evidence for {}",
713                value.speaker()
714            ));
715        }
716    }
717    Ok(set)
718}
719#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
720fn validate_recovered(
721    value: &RecoveredAnalysisArtifacts,
722    audio: &OggAudioMetadata,
723) -> Result<
724    (
725        BTreeSet<LocalSpeakerLabel>,
726        Option<BTreeSet<LocalSpeakerLabel>>,
727    ),
728    ResumableAnalysisError,
729> {
730    if let Some(transcript) = &value.transcript {
731        validate_text(transcript, "transcript")
732            .map_err(|error| ResumableAnalysisError::Recovered(error.to_string()))?;
733    }
734    let evidence =
735        validate_evidence(&value.feature_evidence).map_err(ResumableAnalysisError::Recovered)?;
736    let labels = value
737        .speaker_labels
738        .as_deref()
739        .map(validate_labels)
740        .transpose()
741        .map_err(ResumableAnalysisError::Recovered)?;
742    if labels
743        .as_ref()
744        .is_some_and(|labels| !evidence.is_subset(labels))
745    {
746        return Err(ResumableAnalysisError::Recovered(
747            "feature evidence contains a speaker absent from labels".into(),
748        ));
749    }
750    let final_speakers = value
751        .executed_analysis
752        .as_ref()
753        .map(|executed| {
754            validate_recovered_final(
755                executed,
756                audio,
757                value.transcript.as_deref(),
758                labels.as_ref(),
759                &evidence,
760            )
761        })
762        .transpose()?;
763    Ok((evidence, final_speakers))
764}
765#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
766fn validate_recovered_final(
767    executed: &ExecutedAnalysis,
768    audio: &OggAudioMetadata,
769    transcript: Option<&str>,
770    labels: Option<&BTreeSet<LocalSpeakerLabel>>,
771    evidence: &BTreeSet<LocalSpeakerLabel>,
772) -> Result<BTreeSet<LocalSpeakerLabel>, ResumableAnalysisError> {
773    executed
774        .envelope
775        .validate()
776        .map_err(|error| ResumableAnalysisError::Recovered(error.to_string()))?;
777    executed
778        .label_extractor
779        .validate()
780        .map_err(|error| ResumableAnalysisError::Recovered(error.to_string()))?;
781    if executed.envelope.audio != *audio {
782        return Err(ResumableAnalysisError::Recovered(
783            "final audio metadata differs from admitted audio".into(),
784        ));
785    }
786    if transcript.is_some_and(|value| value != executed.envelope.analysis.transcript.as_str()) {
787        return Err(ResumableAnalysisError::Recovered(
788            "transcript differs from the recovered final transcript".into(),
789        ));
790    }
791    let speakers = executed
792        .envelope
793        .analysis
794        .speakers
795        .iter()
796        .map(|speaker| speaker.speaker)
797        .collect::<BTreeSet<_>>();
798    if labels.is_some_and(|labels| labels != &speakers) {
799        return Err(ResumableAnalysisError::Recovered(
800            "labels differ from the recovered final speaker set".into(),
801        ));
802    }
803    if !evidence.is_subset(&speakers) {
804        return Err(ResumableAnalysisError::Recovered(
805            "feature evidence contains a speaker absent from the recovered final".into(),
806        ));
807    }
808    Ok(speakers)
809}
810#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
811fn non_resumable_error(error: ResumableAnalysisError) -> AnalysisError {
812    match error {
813        ResumableAnalysisError::Analysis(error) => error,
814        ResumableAnalysisError::Recovered(_) | ResumableAnalysisError::Artifact(_) => {
815            unreachable!("non-resumable execution created a resumable-only error")
816        }
817    }
818}
819#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
820fn feature_error(speaker: LocalSpeakerLabel, message: String) -> AnalysisError {
821    AnalysisError::GeminiFeature {
822        speaker,
823        packet: 1,
824        message,
825    }
826}
827#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
828fn progress<P: FnMut(AnalysisProgress) -> Result<(), String>>(
829    report: &mut P,
830    event: AnalysisProgress,
831) -> Result<(), AnalysisError> {
832    report(event).map_err(AnalysisError::Progress)
833}
834#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
835fn failed<P: FnMut(AnalysisProgress) -> Result<(), String>, E: fmt::Display>(
836    report: &mut P,
837    sequence: u64,
838    error: &E,
839) -> Result<(), AnalysisError> {
840    progress(
841        report,
842        AnalysisProgress::JobFailed {
843            sequence,
844            error: error.to_string(),
845        },
846    )
847}
848#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
849fn structuring_sequence(count: usize) -> Result<u64, AnalysisError> {
850    u64::try_from(count)
851        .ok()
852        .and_then(|v| v.checked_mul(3))
853        .and_then(|v| v.checked_add(3))
854        .ok_or_else(sequence_error)
855}
856#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
857fn sequence_error() -> AnalysisError {
858    AnalysisError::Progress("analysis job sequence overflow".into())
859}
860#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
861fn feature_sequence(
862    labels: &[LocalSpeakerLabel],
863    speaker: LocalSpeakerLabel,
864    packet: u8,
865) -> Result<u64, String> {
866    if !(1..=3).contains(&packet) {
867        return Err(format!(
868            "Gemini reported invalid feature packet {packet} for {speaker}"
869        ));
870    }
871    let index = labels
872        .iter()
873        .position(|v| *v == speaker)
874        .ok_or_else(|| format!("Gemini reported an unknown feature speaker {speaker}"))?;
875    u64::try_from(index)
876        .ok()
877        .and_then(|v| v.checked_mul(3))
878        .and_then(|v| v.checked_add(u64::from(packet)))
879        .and_then(|v| v.checked_add(2))
880        .ok_or_else(|| sequence_error().to_string())
881}
882#[cfg(any(feature = "providers", feature = "adapter-providers", test))]
883fn map_feature_progress(
884    labels: &[LocalSpeakerLabel],
885    event: GeminiFeatureProgress,
886) -> Result<AnalysisProgress, String> {
887    match event {
888        GeminiFeatureProgress::Started { speaker, packet } => Ok(AnalysisProgress::JobStarted {
889            sequence: feature_sequence(labels, speaker, packet)?,
890            job: AnalysisJob::SpeakerFeature { speaker, packet },
891        }),
892        GeminiFeatureProgress::Succeeded { speaker, packet } => {
893            Ok(AnalysisProgress::JobSucceeded {
894                sequence: feature_sequence(labels, speaker, packet)?,
895            })
896        }
897        GeminiFeatureProgress::Failed {
898            speaker,
899            packet,
900            error,
901        } => {
902            let error = AnalysisError::GeminiFeature {
903                speaker,
904                packet,
905                message: error,
906            }
907            .to_string();
908            Ok(AnalysisProgress::JobFailed {
909                sequence: feature_sequence(labels, speaker, packet)?,
910                error,
911            })
912        }
913    }
914}
915
916#[cfg(test)]
917mod tests {
918    use super::*;
919    use futures::{executor::block_on, future::poll_fn, join};
920    use std::{
921        sync::{
922            Arc, Mutex,
923            atomic::{AtomicBool, Ordering},
924        },
925        task::Poll,
926    };
927    const TRANSCRIPT: &str = "[high] Speaker 1: exact transcript";
928    struct State {
929        labels: Vec<LocalSpeakerLabel>,
930        analysis: StructuredAnalysis,
931        fail: Option<&'static str>,
932        calls: Mutex<Vec<&'static str>>,
933        targets: Mutex<Vec<LocalSpeakerLabel>>,
934        final_evidence: Mutex<Vec<SpeakerFeatureEvidence>>,
935        wait: Option<Arc<AtomicBool>>,
936        mark: Option<Arc<AtomicBool>>,
937    }
938    #[derive(Clone)]
939    struct Fake(Arc<State>);
940    impl Fake {
941        fn new(count: u32) -> Self {
942            Self::configured(count, analysis(TRANSCRIPT, count), None, None, None)
943        }
944        fn configured(
945            count: u32,
946            result: StructuredAnalysis,
947            fail: Option<&'static str>,
948            wait: Option<Arc<AtomicBool>>,
949            mark: Option<Arc<AtomicBool>>,
950        ) -> Self {
951            Self(Arc::new(State {
952                labels: (1..=count).map(label).collect(),
953                analysis: result,
954                fail,
955                calls: Mutex::new(Vec::new()),
956                targets: Mutex::new(Vec::new()),
957                final_evidence: Mutex::new(Vec::new()),
958                wait,
959                mark,
960            }))
961        }
962        fn calls(&self) -> Vec<&'static str> {
963            self.0.calls.lock().unwrap().clone()
964        }
965    }
966    impl AnalysisOperations for Fake {
967        fn transcript<'a>(
968            &'a self,
969            _: &'a [u8],
970        ) -> AnalysisFuture<'a, Result<String, AnalysisError>> {
971            Box::pin(async move {
972                self.0.calls.lock().unwrap().push("transcript");
973                if let Some(gate) = &self.0.wait {
974                    poll_fn(|cx| {
975                        if gate.load(Ordering::SeqCst) {
976                            Poll::Ready(())
977                        } else {
978                            cx.waker().wake_by_ref();
979                            Poll::Pending
980                        }
981                    })
982                    .await;
983                }
984                if self.0.fail == Some("transcript") {
985                    Err(AnalysisError::GeminiTranscript("failed".into()))
986                } else {
987                    Ok(TRANSCRIPT.into())
988                }
989            })
990        }
991        fn speaker_labels<'a>(
992            &'a self,
993            _: &'a str,
994        ) -> AnalysisFuture<'a, Result<Vec<LocalSpeakerLabel>, AnalysisError>> {
995            Box::pin(async move {
996                self.0.calls.lock().unwrap().push("labels");
997                if self.0.fail == Some("labels") {
998                    Err(AnalysisError::TerraLabels("failed".into()))
999                } else {
1000                    Ok(self.0.labels.clone())
1001                }
1002            })
1003        }
1004        fn feature_evidence_with_progress<'a>(
1005            &'a self,
1006            _: &'a [u8],
1007            _: &'a str,
1008            labels: &'a [LocalSpeakerLabel],
1009            report: &'a mut dyn FnMut(GeminiFeatureProgress) -> Result<(), String>,
1010        ) -> AnalysisFuture<'a, Result<Vec<SpeakerFeatureEvidence>, AnalysisError>> {
1011            Box::pin(async move {
1012                self.0.calls.lock().unwrap().push("features");
1013                *self.0.targets.lock().unwrap() = labels.to_vec();
1014                if self.0.fail == Some("features") {
1015                    return Err(AnalysisError::GeminiCache("failed".into()));
1016                }
1017                for &speaker in labels {
1018                    for packet in 1..=3 {
1019                        report(GeminiFeatureProgress::Started { speaker, packet })
1020                            .map_err(AnalysisError::Progress)?;
1021                        report(GeminiFeatureProgress::Succeeded { speaker, packet })
1022                            .map_err(AnalysisError::Progress)?;
1023                    }
1024                }
1025                Ok(labels.iter().copied().map(evidence).collect())
1026            })
1027        }
1028        fn structured_analysis<'a>(
1029            &'a self,
1030            _: &'a str,
1031            evidence: Vec<SpeakerFeatureEvidence>,
1032        ) -> AnalysisFuture<'a, Result<StructuredAnalysis, AnalysisError>> {
1033            Box::pin(async move {
1034                self.0.calls.lock().unwrap().push("final");
1035                *self.0.final_evidence.lock().unwrap() = evidence;
1036                if self.0.fail == Some("final") {
1037                    return Err(AnalysisError::TerraStructuring("failed".into()));
1038                }
1039                if let Some(mark) = &self.0.mark {
1040                    mark.store(true, Ordering::SeqCst);
1041                }
1042                Ok(self.0.analysis.clone())
1043            })
1044        }
1045    }
1046    fn label(n: u32) -> LocalSpeakerLabel {
1047        LocalSpeakerLabel::new(n).unwrap()
1048    }
1049    fn evidence(speaker: LocalSpeakerLabel) -> SpeakerFeatureEvidence {
1050        SpeakerFeatureEvidence::new(
1051            speaker,
1052            format!("{speaker}-1"),
1053            format!("{speaker}-2"),
1054            format!("{speaker}-3"),
1055        )
1056        .unwrap()
1057    }
1058    fn analysis(transcript: &str, count: u32) -> StructuredAnalysis {
1059        StructuredAnalysis {
1060            transcript: transcript.into(),
1061            speakers: (1..=count)
1062                .map(|n| StructuredSpeaker {
1063                    speaker: label(n),
1064                    language: "English".into(),
1065                    features: FeatureVector24::default(),
1066                    features_usable_for_training: false,
1067                })
1068                .collect(),
1069        }
1070    }
1071    fn ogg() -> Vec<u8> {
1072        let mut value = vec![0; 28];
1073        value[..4].copy_from_slice(b"OggS");
1074        value[26] = 1;
1075        value
1076    }
1077    fn audio(bytes: &[u8]) -> OggAudioMetadata {
1078        OggAudioMetadata::from_bytes(bytes, 1, None).unwrap()
1079    }
1080    fn executed(bytes: &[u8], count: u32) -> ExecutedAnalysis {
1081        ExecutedAnalysis {
1082            envelope: AnalysisEnvelope {
1083                audio: audio(bytes),
1084                analysis: analysis(TRANSCRIPT, count),
1085                gemini: GeminiCohort::new(GEMINI_MODEL_ID),
1086                structurer: StructurerProvenance::new(TERRA_MODEL_ID),
1087            },
1088            label_extractor: StructurerProvenance {
1089                model_id: TERRA_MODEL_ID.into(),
1090                prompt_revision:
1091                    kcode_speaker_v3_llm_protocol::TERRA_SPEAKER_LABELS_PROMPT_REVISION.into(),
1092            },
1093        }
1094    }
1095    fn resume(
1096        fake: &Fake,
1097        recovered: RecoveredAnalysisArtifacts,
1098    ) -> (
1099        Result<ExecutedAnalysis, ResumableAnalysisError>,
1100        Vec<AnalysisProgress>,
1101        Vec<AnalysisArtifact>,
1102    ) {
1103        let bytes = ogg();
1104        let mut progress = Vec::new();
1105        let mut artifacts = Vec::new();
1106        let result = block_on(execute_admitted(
1107            fake,
1108            &bytes,
1109            audio(&bytes),
1110            recovered,
1111            |v| {
1112                progress.push(v);
1113                Ok(())
1114            },
1115            |v| {
1116                artifacts.push(v);
1117                Ok(())
1118            },
1119            true,
1120        ));
1121        (result, progress, artifacts)
1122    }
1123    fn recovered(count: u32) -> RecoveredAnalysisArtifacts {
1124        RecoveredAnalysisArtifacts {
1125            transcript: Some(TRANSCRIPT.into()),
1126            speaker_labels: Some((1..=count).map(label).collect()),
1127            feature_evidence: (1..=count).map(|n| evidence(label(n))).collect(),
1128            ..Default::default()
1129        }
1130    }
1131
1132    #[test]
1133    fn all_missing_matches_old_flow_and_legacy_contract() {
1134        let bytes = ogg();
1135        let old = Fake::new(1);
1136        let mut old_events = Vec::new();
1137        let old_result = block_on(execute_admitted(
1138            &old,
1139            &bytes,
1140            audio(&bytes),
1141            RecoveredAnalysisArtifacts::default(),
1142            |v| {
1143                old_events.push(v);
1144                Ok(())
1145            },
1146            |_| Ok(()),
1147            false,
1148        ))
1149        .map_err(non_resumable_error);
1150        let new = Fake::new(1);
1151        let (new_result, new_events, artifacts) =
1152            resume(&new, RecoveredAnalysisArtifacts::default());
1153        assert_eq!(old_result, new_result.map_err(non_resumable_error));
1154        assert_eq!(old_events, new_events);
1155        assert_eq!(old.calls(), new.calls());
1156        assert!(matches!(
1157            &artifacts[..],
1158            [
1159                AnalysisArtifact::GeminiTranscript(_),
1160                AnalysisArtifact::TerraSpeakerLabels(_),
1161                AnalysisArtifact::GeminiFeatureEvidence(_),
1162                AnalysisArtifact::ExecutedAnalysis(_)
1163            ]
1164        ));
1165        let starts = new_events
1166            .iter()
1167            .filter_map(|v| {
1168                if let AnalysisProgress::JobStarted { sequence, .. } = v {
1169                    Some(*sequence)
1170                } else {
1171                    None
1172                }
1173            })
1174            .collect::<Vec<_>>();
1175        assert_eq!(starts, [1, 2, 3, 4, 5, 6]);
1176        let legacy = Fake::new(0);
1177        let result = block_on(execute_legacy(
1178            &legacy,
1179            &bytes,
1180            1234,
1181            Some("voice.ogg".into()),
1182        ))
1183        .unwrap();
1184        assert_eq!(result.envelope.audio.filename(), Some("voice.ogg"));
1185        assert_eq!(
1186            legacy.calls(),
1187            ["transcript", "labels", "features", "final"]
1188        );
1189        let invalid = Fake::new(1);
1190        let mut reports = 0;
1191        assert!(matches!(
1192            block_on(execute_strict(&invalid, b"bad", |_| {
1193                reports += 1;
1194                Ok(())
1195            })),
1196            Err(AnalysisError::Input(_))
1197        ));
1198        assert_eq!(reports, 0);
1199        assert!(invalid.calls().is_empty());
1200    }
1201
1202    #[test]
1203    fn recovered_and_partial_stages_are_selective_and_preserved() {
1204        let full = Fake::new(1);
1205        let (result, events, artifacts) = resume(&full, recovered(1));
1206        result.unwrap();
1207        assert_eq!(full.calls(), ["final"]);
1208        assert!(matches!(
1209            &artifacts[..],
1210            [AnalysisArtifact::ExecutedAnalysis(_)]
1211        ));
1212        assert_eq!(events.len(), 3);
1213        let transcript = Fake::new(1);
1214        let (_, _, artifacts) = resume(
1215            &transcript,
1216            RecoveredAnalysisArtifacts {
1217                transcript: Some(TRANSCRIPT.into()),
1218                ..Default::default()
1219            },
1220        );
1221        assert_eq!(transcript.calls(), ["labels", "features", "final"]);
1222        assert_eq!(artifacts.len(), 3);
1223        let labels = Fake::new(1);
1224        let (_, _, artifacts) = resume(
1225            &labels,
1226            RecoveredAnalysisArtifacts {
1227                speaker_labels: Some(vec![label(1)]),
1228                ..Default::default()
1229            },
1230        );
1231        assert_eq!(labels.calls(), ["transcript", "features", "final"]);
1232        assert!(matches!(
1233            &artifacts[..],
1234            [
1235                AnalysisArtifact::GeminiTranscript(_),
1236                AnalysisArtifact::GeminiFeatureEvidence(_),
1237                AnalysisArtifact::ExecutedAnalysis(_)
1238            ]
1239        ));
1240        let packet = evidence(label(1));
1241        let downstream = Fake::new(1);
1242        let (_, _, artifacts) = resume(
1243            &downstream,
1244            RecoveredAnalysisArtifacts {
1245                feature_evidence: vec![packet.clone()],
1246                ..Default::default()
1247            },
1248        );
1249        assert_eq!(downstream.calls(), ["transcript", "labels", "final"]);
1250        assert_eq!(*downstream.0.final_evidence.lock().unwrap(), [packet]);
1251        assert_eq!(artifacts.len(), 3);
1252        let first = evidence(label(1));
1253        let partial = Fake::new(2);
1254        let (_, _, artifacts) = resume(
1255            &partial,
1256            RecoveredAnalysisArtifacts {
1257                transcript: Some(TRANSCRIPT.into()),
1258                speaker_labels: Some(vec![label(1), label(2)]),
1259                feature_evidence: vec![first.clone()],
1260                ..Default::default()
1261            },
1262        );
1263        assert_eq!(partial.calls(), ["features", "final"]);
1264        assert_eq!(*partial.0.targets.lock().unwrap(), [label(2)]);
1265        assert!(
1266            matches!(&artifacts[..], [AnalysisArtifact::GeminiFeatureEvidence(v), AnalysisArtifact::ExecutedAnalysis(_)] if v.speaker() == label(2))
1267        );
1268        assert_eq!(partial.0.final_evidence.lock().unwrap()[0], first);
1269    }
1270
1271    #[test]
1272    fn recovered_final_repairs_missing_stages_without_structuring() {
1273        let bytes = ogg();
1274        let final_value = executed(&bytes, 2);
1275        let fake = Fake::new(2);
1276        let (result, events, artifacts) = resume(
1277            &fake,
1278            RecoveredAnalysisArtifacts {
1279                feature_evidence: vec![evidence(label(1))],
1280                executed_analysis: Some(final_value.clone()),
1281                ..Default::default()
1282            },
1283        );
1284        assert_eq!(result.unwrap(), final_value);
1285        assert_eq!(fake.calls(), ["labels", "features"]);
1286        assert_eq!(*fake.0.targets.lock().unwrap(), [label(2)]);
1287        assert!(
1288            matches!(&artifacts[..], [AnalysisArtifact::TerraSpeakerLabels(_), AnalysisArtifact::GeminiFeatureEvidence(value)] if value.speaker() == label(2))
1289        );
1290        assert!(!events.iter().any(|event| matches!(
1291            event,
1292            AnalysisProgress::JobStarted {
1293                job: AnalysisJob::Structuring,
1294                ..
1295            } | AnalysisProgress::StageCompleted {
1296                stage: AnalysisStage::Structuring
1297            }
1298        )));
1299        let mismatch = Fake::new(2);
1300        let (_, _, artifacts) = resume(
1301            &mismatch,
1302            RecoveredAnalysisArtifacts {
1303                executed_analysis: Some(executed(&bytes, 1)),
1304                ..Default::default()
1305            },
1306        );
1307        assert_eq!(mismatch.calls(), ["labels"]);
1308        assert!(artifacts.is_empty());
1309    }
1310
1311    #[test]
1312    fn recovered_final_supplies_transcript_and_is_returned_exactly() {
1313        let bytes = ogg();
1314        let final_value = executed(&bytes, 1);
1315        let fake = Fake::configured(1, analysis(TRANSCRIPT, 1), Some("transcript"), None, None);
1316        let (result, events, artifacts) = resume(
1317            &fake,
1318            RecoveredAnalysisArtifacts {
1319                speaker_labels: Some(vec![label(1)]),
1320                feature_evidence: vec![evidence(label(1))],
1321                executed_analysis: Some(final_value.clone()),
1322                ..Default::default()
1323            },
1324        );
1325        assert_eq!(result.unwrap(), final_value);
1326        assert!(fake.calls().is_empty());
1327        assert!(events.is_empty());
1328        assert!(artifacts.is_empty());
1329    }
1330
1331    #[test]
1332    fn final_recovery_conflicts_are_effect_free() {
1333        let bytes = ogg();
1334        let valid = executed(&bytes, 1);
1335        let mut wrong_audio = valid.clone();
1336        wrong_audio.envelope.audio = OggAudioMetadata::from_bytes(&bytes, 2, None).unwrap();
1337        let mut bad_envelope = valid.clone();
1338        bad_envelope.envelope.gemini.model_id.clear();
1339        let mut bad_extractor = valid.clone();
1340        bad_extractor.label_extractor.model_id.clear();
1341        let bad = [
1342            RecoveredAnalysisArtifacts {
1343                executed_analysis: Some(wrong_audio),
1344                ..Default::default()
1345            },
1346            RecoveredAnalysisArtifacts {
1347                transcript: Some("different".into()),
1348                executed_analysis: Some(valid.clone()),
1349                ..Default::default()
1350            },
1351            RecoveredAnalysisArtifacts {
1352                speaker_labels: Some(vec![label(2)]),
1353                executed_analysis: Some(valid.clone()),
1354                ..Default::default()
1355            },
1356            RecoveredAnalysisArtifacts {
1357                feature_evidence: vec![evidence(label(2))],
1358                executed_analysis: Some(valid.clone()),
1359                ..Default::default()
1360            },
1361            RecoveredAnalysisArtifacts {
1362                executed_analysis: Some(bad_envelope),
1363                ..Default::default()
1364            },
1365            RecoveredAnalysisArtifacts {
1366                executed_analysis: Some(bad_extractor),
1367                ..Default::default()
1368            },
1369        ];
1370        for value in bad {
1371            let fake = Fake::new(1);
1372            let (result, events, artifacts) = resume(&fake, value);
1373            assert!(matches!(result, Err(ResumableAnalysisError::Recovered(_))));
1374            assert!(fake.calls().is_empty());
1375            assert!(events.is_empty());
1376            assert!(artifacts.is_empty());
1377        }
1378    }
1379
1380    #[test]
1381    fn recovered_empty_labels_are_complete_and_zero_speaker_final_works() {
1382        let fake = Fake::new(0);
1383        let (result, events, artifacts) = resume(
1384            &fake,
1385            RecoveredAnalysisArtifacts {
1386                transcript: Some(TRANSCRIPT.into()),
1387                speaker_labels: Some(vec![]),
1388                ..Default::default()
1389            },
1390        );
1391        result.unwrap();
1392        assert_eq!(fake.calls(), ["final"]);
1393        assert!(matches!(
1394            &artifacts[..],
1395            [AnalysisArtifact::ExecutedAnalysis(_)]
1396        ));
1397        assert!(!events.iter().any(|event| matches!(
1398            event,
1399            AnalysisProgress::StageCompleted {
1400                stage: AnalysisStage::SpeakerFeatures
1401            }
1402        )));
1403        let bytes = ogg();
1404        let final_value = executed(&bytes, 0);
1405        let recovered = Fake::new(0);
1406        let (result, events, artifacts) = resume(
1407            &recovered,
1408            RecoveredAnalysisArtifacts {
1409                speaker_labels: Some(vec![]),
1410                executed_analysis: Some(final_value.clone()),
1411                ..Default::default()
1412            },
1413        );
1414        assert_eq!(result.unwrap(), final_value);
1415        assert!(recovered.calls().is_empty());
1416        assert!(events.is_empty());
1417        assert!(artifacts.is_empty());
1418    }
1419
1420    #[test]
1421    fn artifact_failure_fences_the_next_stage() {
1422        let fake = Fake::new(1);
1423        let bytes = ogg();
1424        let mut events = Vec::new();
1425        let result = block_on(execute_admitted(
1426            &fake,
1427            &bytes,
1428            audio(&bytes),
1429            RecoveredAnalysisArtifacts::default(),
1430            |v| {
1431                events.push(v);
1432                Ok(())
1433            },
1434            |_| Err("durability unavailable".into()),
1435            true,
1436        ));
1437        assert_eq!(
1438            result,
1439            Err(ResumableAnalysisError::Artifact(
1440                "durability unavailable".into()
1441            ))
1442        );
1443        assert_eq!(fake.calls(), ["transcript"]);
1444        assert!(
1445            !events
1446                .iter()
1447                .any(|v| matches!(v, AnalysisProgress::StageCompleted { .. }))
1448        );
1449    }
1450
1451    #[test]
1452    fn malformed_recovery_is_effect_free_and_generated_final_conflicts_fail() {
1453        let bad = [
1454            RecoveredAnalysisArtifacts {
1455                transcript: Some(" ".into()),
1456                ..Default::default()
1457            },
1458            RecoveredAnalysisArtifacts {
1459                speaker_labels: Some(vec![label(1), label(1)]),
1460                ..Default::default()
1461            },
1462            RecoveredAnalysisArtifacts {
1463                speaker_labels: Some(vec![label(1)]),
1464                feature_evidence: vec![evidence(label(2))],
1465                ..Default::default()
1466            },
1467            RecoveredAnalysisArtifacts {
1468                feature_evidence: vec![evidence(label(1)), evidence(label(1))],
1469                ..Default::default()
1470            },
1471        ];
1472        for value in bad {
1473            let fake = Fake::new(1);
1474            assert!(matches!(
1475                resume(&fake, value).0,
1476                Err(ResumableAnalysisError::Recovered(_))
1477            ));
1478            assert!(fake.calls().is_empty());
1479        }
1480        let transcript = Fake::configured(1, analysis("different", 1), None, None, None);
1481        assert_eq!(
1482            resume(&transcript, recovered(1)).0,
1483            Err(ResumableAnalysisError::Analysis(
1484                AnalysisError::TranscriptMismatch
1485            ))
1486        );
1487        let speakers = Fake::configured(1, analysis(TRANSCRIPT, 2), None, None, None);
1488        assert_eq!(
1489            resume(&speakers, recovered(1)).0,
1490            Err(ResumableAnalysisError::Analysis(
1491                AnalysisError::SpeakerSetMismatch
1492            ))
1493        );
1494    }
1495
1496    #[test]
1497    fn failures_do_not_retry_and_analyses_remain_concurrent() {
1498        for (stage, calls) in [
1499            ("transcript", vec!["transcript"]),
1500            ("labels", vec!["transcript", "labels"]),
1501            ("features", vec!["transcript", "labels", "features"]),
1502            ("final", vec!["transcript", "labels", "features", "final"]),
1503        ] {
1504            let fake = Fake::configured(1, analysis(TRANSCRIPT, 1), Some(stage), None, None);
1505            assert!(
1506                resume(&fake, RecoveredAnalysisArtifacts::default())
1507                    .0
1508                    .is_err()
1509            );
1510            assert_eq!(fake.calls(), calls);
1511        }
1512        let done = Arc::new(AtomicBool::new(false));
1513        let slow = Fake::configured(1, analysis(TRANSCRIPT, 1), None, Some(done.clone()), None);
1514        let fast = Fake::configured(1, analysis(TRANSCRIPT, 1), None, None, Some(done.clone()));
1515        let bytes = ogg();
1516        let (a, b) = block_on(async {
1517            join!(
1518                execute_admitted(
1519                    &slow,
1520                    &bytes,
1521                    audio(&bytes),
1522                    RecoveredAnalysisArtifacts::default(),
1523                    |_| Ok(()),
1524                    |_| Ok(()),
1525                    true
1526                ),
1527                execute_admitted(
1528                    &fast,
1529                    &bytes,
1530                    audio(&bytes),
1531                    RecoveredAnalysisArtifacts::default(),
1532                    |_| Ok(()),
1533                    |_| Ok(()),
1534                    true
1535                )
1536            )
1537        });
1538        a.unwrap();
1539        b.unwrap();
1540        assert!(done.load(Ordering::SeqCst));
1541    }
1542
1543    #[test]
1544    fn public_contract_and_provider_operations_compile() {
1545        let cohort = GeminiCohort::new("gemini");
1546        cohort.validate().unwrap();
1547        StructurerProvenance::new("terra").validate().unwrap();
1548        fn require<O: AnalysisOperations>() {}
1549        require::<ProviderOperations>();
1550        let wrapped: ResumableAnalysisError = AnalysisError::Input("bad".into()).into();
1551        assert!(matches!(
1552            wrapped,
1553            ResumableAnalysisError::Analysis(AnalysisError::Input(_))
1554        ));
1555    }
1556    #[cfg(feature = "providers")]
1557    #[test]
1558    fn legacy_constructor_compiles() {
1559        let _: fn(kcode_gemini_3_1_pro::Gemini31Pro, kcode_codex_terra::CodexTerra) -> Analyzer =
1560            Analyzer::new;
1561    }
1562    #[cfg(feature = "adapter-providers")]
1563    #[test]
1564    fn adapter_constructor_compiles() {
1565        let _: fn(kcode_gemini_3_1_pro::Gemini31Pro, kcode_k1_codex_adapter::Adapter) -> Analyzer =
1566            Analyzer::from_codex_adapter;
1567    }
1568}