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}