Skip to main content

kcode_audio_session_view/
lib.rs

1#![forbid(unsafe_code)]
2
3//! Deterministic adaptation of Audio Ingress and History Handoff projections.
4
5use kcode_audio_history_handoff::{PieceProjection, RecordingProjection};
6use kcode_audio_ingress::{
7    ConfirmationState, RecordingState, RecordingStatus, Step, StepState, TranscriptionStatus,
8};
9use serde_json::Value;
10use uuid::Uuid;
11
12/// Established facade view of one audio recording.
13#[derive(Clone, Debug)]
14pub struct Recording {
15    pub id: Uuid,
16    pub sha256: String,
17    pub original_filename: String,
18    pub content_type: &'static str,
19    pub size_bytes: u64,
20    pub source_created_at: String,
21    pub received_at: String,
22    pub updated_at: String,
23    pub status: String,
24    pub transcription_model: String,
25    pub reconciliation_model: String,
26    pub reconciliation_reasoning: String,
27    pub transcription_status: Option<Value>,
28    pub attempt_count: i64,
29    pub next_attempt_at: Option<String>,
30    pub last_error: Option<String>,
31    pub speaker_review: Option<SpeakerReview>,
32    pub transcript_piece_count: usize,
33    pub completed_piece_count: usize,
34}
35
36/// Bounded review summary for one classifier-aware recording.
37#[derive(Clone, Debug)]
38pub struct SpeakerReview {
39    pub confirmation_state: ConfirmationState,
40    pub observation_count: usize,
41    pub signed_off_chunk_count: usize,
42}
43
44/// One deterministic transcript piece and its Session History lifecycle.
45#[derive(Clone, Debug)]
46pub struct IngressPiece {
47    pub id: String,
48    pub recording_id: Uuid,
49    pub sha256: String,
50    pub original_filename: String,
51    pub source_created_at: String,
52    pub piece_index: u32,
53    pub piece_count: u32,
54    pub transcript_text: String,
55    pub estimated_tokens: u64,
56    pub phase: String,
57    pub provenance_id: Option<String>,
58    pub state: Value,
59    pub version: i64,
60    pub ingress_failure_count: i64,
61    pub ingress_failures: Value,
62    pub created_at: String,
63    pub updated_at: String,
64}
65
66/// Complete deterministic facade view for one recording and its projected pieces.
67#[derive(Clone, Debug)]
68pub struct RecordingView {
69    pub recording: Recording,
70    pub pieces: Vec<IngressPiece>,
71}
72
73/// Adapts one Audio Ingress status and its matching History Handoff projection.
74pub fn render(recording: RecordingStatus, projection: RecordingProjection) -> RecordingView {
75    let RecordingProjection {
76        recording_id: _,
77        pieces: projected_pieces,
78    } = projection;
79    let pieces = projected_pieces
80        .into_iter()
81        .map(|piece| adapt_piece(&recording, piece))
82        .collect::<Vec<_>>();
83    let recording = adapt_recording(recording, &pieces);
84    RecordingView { recording, pieces }
85}
86
87fn adapt_recording(recording: RecordingStatus, pieces: &[IngressPiece]) -> Recording {
88    let speaker_review = recording
89        .correction_packet
90        .as_ref()
91        .map(|packet| SpeakerReview {
92            confirmation_state: packet.confirmation_state,
93            observation_count: packet
94                .chunks
95                .iter()
96                .map(|chunk| chunk.observations.len())
97                .sum(),
98            signed_off_chunk_count: packet
99                .chunks
100                .iter()
101                .filter(|chunk| chunk.signed_off)
102                .count(),
103        });
104    let awaiting_speaker_review = speaker_review
105        .as_ref()
106        .is_some_and(|review| review.confirmation_state != ConfirmationState::Confirmed);
107    let (mut status, transcription_status, attempt_count, last_error) = match recording.state {
108        RecordingState::Queued => ("uploaded".into(), None, 0, None),
109        RecordingState::Processing { attempt, progress } => (
110            processing_stage(&progress).into(),
111            serde_json::to_value(progress).ok(),
112            i64::from(attempt),
113            None,
114        ),
115        RecordingState::AwaitingReview => ("speaker_review".into(), None, 0, None),
116        RecordingState::Complete { .. } if awaiting_speaker_review => {
117            ("speaker_review".into(), None, 0, None)
118        }
119        RecordingState::Complete { .. } => ("ready_for_ingress".into(), None, 0, None),
120        RecordingState::Failed {
121            attempts, error, ..
122        } => ("failed".into(), None, i64::from(attempts), Some(error)),
123    };
124
125    if !pieces.is_empty() {
126        status = if pieces.iter().all(|piece| piece.phase == "complete") {
127            "complete".into()
128        } else if pieces.iter().any(|piece| piece.phase == "ingress_failed") {
129            "ingress_failed".into()
130        } else if pieces
131            .iter()
132            .any(|piece| piece.phase == "ingress_in_progress")
133        {
134            "ingressing".into()
135        } else {
136            "ready_for_ingress".into()
137        };
138    }
139    let completed_piece_count = pieces
140        .iter()
141        .filter(|piece| piece.phase == "complete")
142        .count();
143
144    Recording {
145        id: recording.id,
146        sha256: recording.sha256,
147        original_filename: recording.original_filename,
148        content_type: "audio/wav",
149        size_bytes: recording.size_bytes,
150        source_created_at: recording.recorded_at.to_rfc3339(),
151        received_at: recording.received_at.to_rfc3339(),
152        updated_at: recording.received_at.to_rfc3339(),
153        status,
154        transcription_model: recording.transcription_model,
155        reconciliation_model: recording.reconciliation_model,
156        reconciliation_reasoning: recording.reconciliation_reasoning,
157        transcription_status,
158        attempt_count,
159        next_attempt_at: None,
160        last_error,
161        speaker_review,
162        transcript_piece_count: pieces.len(),
163        completed_piece_count,
164    }
165}
166
167fn adapt_piece(recording: &RecordingStatus, piece: PieceProjection) -> IngressPiece {
168    let record = piece.record;
169    IngressPiece {
170        id: record.id,
171        recording_id: recording.id,
172        sha256: recording.sha256.clone(),
173        original_filename: recording.original_filename.clone(),
174        source_created_at: recording.recorded_at.to_rfc3339(),
175        piece_index: piece.piece_index,
176        piece_count: piece.piece_count,
177        transcript_text: piece.transcript_text,
178        estimated_tokens: piece.estimated_tokens,
179        phase: record.phase,
180        provenance_id: record.provenance_id,
181        state: record.state,
182        version: record.version,
183        ingress_failure_count: record.ingress_failure_count,
184        ingress_failures: record.ingress_failures,
185        created_at: record.started_at,
186        updated_at: record.updated_at,
187    }
188}
189
190fn processing_stage(status: &TranscriptionStatus) -> &'static str {
191    if status.steps.len() == 1 && status.steps[0].step == Step::ReconcileTranscript {
192        return "reconciling";
193    }
194    let plan_complete = status
195        .steps
196        .iter()
197        .any(|entry| entry.step == Step::PlanChunks && entry.state == StepState::Completed);
198    if !plan_complete {
199        return "chunking";
200    }
201
202    let chunks_complete = status
203        .steps
204        .iter()
205        .filter(|entry| matches!(entry.step, Step::TranscribeChunk { .. }))
206        .all(|entry| entry.state == StepState::Completed);
207    if !chunks_complete {
208        return "transcribing";
209    }
210
211    let analyses_complete = status
212        .steps
213        .iter()
214        .any(|entry| entry.step == Step::AnalyzeSpeakers && entry.state == StepState::Completed);
215    if !analyses_complete {
216        return "analyzing_speakers";
217    }
218
219    "reconciling"
220}
221
222#[cfg(test)]
223mod tests {
224    use chrono::{DateTime, Utc};
225    use kcode_audio_history_handoff::{PieceProjection, RecordingProjection};
226    use kcode_audio_ingress::{
227        ConfirmationState, CorrectionChunk, CorrectionObservation, CorrectionPacket, JobState,
228        ObservationKey, ParsedChunk, RecordingState, RecordingStatus, Step, StepState, StepStatus,
229        TranscriptionStatus,
230    };
231    use kcode_session_history::SessionRecord;
232    use serde_json::json;
233    use uuid::Uuid;
234
235    use super::render;
236
237    fn timestamp(value: &str) -> DateTime<Utc> {
238        DateTime::parse_from_rfc3339(value)
239            .unwrap()
240            .with_timezone(&Utc)
241    }
242
243    fn recording(state: RecordingState) -> RecordingStatus {
244        RecordingStatus {
245            id: Uuid::parse_str("b067a460-69bb-4c49-a9d7-037715d3137d").unwrap(),
246            user_id: "user".into(),
247            sha256: "ab".repeat(32),
248            original_filename: "meeting.final.WAV".into(),
249            size_bytes: 42,
250            recorded_at: timestamp("2026-01-02T03:04:05Z"),
251            received_at: timestamp("2026-01-02T03:05:06Z"),
252            transcription_model: "transcription-model".into(),
253            reconciliation_model: "reconciliation-model".into(),
254            reconciliation_reasoning: "xhigh".into(),
255            state,
256            correction_packet: None,
257        }
258    }
259
260    fn with_speaker_packet(
261        mut recording: RecordingStatus,
262        confirmation_state: ConfirmationState,
263    ) -> RecordingStatus {
264        let observations = (0..2)
265            .map(|speaker_ordinal| CorrectionObservation {
266                local_label: format!("Speaker {speaker_ordinal}"),
267                speaker_ordinal,
268                observation_key: ObservationKey {
269                    object_id: format!("observation-{speaker_ordinal}"),
270                    piece_index: speaker_ordinal,
271                },
272                candidate: None,
273                resolution: None,
274            })
275            .collect();
276        recording.correction_packet = Some(CorrectionPacket {
277            recording_id: recording.id,
278            user_id: recording.user_id.clone(),
279            sha256: recording.sha256.clone(),
280            original_filename: recording.original_filename.clone(),
281            size_bytes: recording.size_bytes,
282            recorded_at: recording.recorded_at,
283            chunk_count: 1,
284            chunks: vec![CorrectionChunk {
285                chunk_index: 0,
286                chunk_count: 1,
287                audio_start_ms: 0,
288                audio_end_ms: 1_000,
289                raw_gemini_response: "raw".into(),
290                parsed: ParsedChunk {
291                    clip_valid: true,
292                    clip_validity_reason: None,
293                    speakers: Vec::new(),
294                },
295                observations,
296                signed_off: confirmation_state == ConfirmationState::Confirmed,
297            }],
298            confirmation_state,
299        });
300        recording
301    }
302
303    fn projection(recording_id: Uuid, phases: &[&str]) -> RecordingProjection {
304        let piece_count = u32::try_from(phases.len()).unwrap();
305        RecordingProjection {
306            recording_id,
307            pieces: phases
308                .iter()
309                .enumerate()
310                .map(|(index, phase)| PieceProjection {
311                    piece_index: u32::try_from(index).unwrap(),
312                    piece_count,
313                    transcript_text: format!("Transcript piece {index}"),
314                    estimated_tokens: 10 + u64::try_from(index).unwrap(),
315                    record: SessionRecord {
316                        id: format!("session-{index}"),
317                        phase: (*phase).into(),
318                        started_at: format!("2026-01-02T03:0{index}:00Z"),
319                        updated_at: format!("2026-01-02T04:0{index}:00Z"),
320                        state: json!({"piece":index}),
321                        provenance_id: Some(format!("provenance-{index}")),
322                        version: 7 + i64::try_from(index).unwrap(),
323                        last_user_message_at: None,
324                        ended_at: None,
325                        ingress_failure_count: 2 + i64::try_from(index).unwrap(),
326                        ingress_failures: json!([{"piece":index,"message":"retained"}]),
327                        ingress_next_attempt_at: None,
328                        summary: false,
329                    },
330                })
331                .collect(),
332        }
333    }
334
335    fn empty_projection(recording_id: Uuid) -> RecordingProjection {
336        projection(recording_id, &[])
337    }
338
339    #[test]
340    fn queued_failed_complete_and_speaker_review_states_are_preserved() {
341        let queued = recording(RecordingState::Queued);
342        let queued_id = queued.id;
343        let queued = render(queued, empty_projection(queued_id)).recording;
344        assert_eq!(queued.status, "uploaded");
345        assert_eq!(queued.attempt_count, 0);
346        assert!(queued.last_error.is_none());
347
348        let failed = recording(RecordingState::Failed {
349            attempts: 3,
350            error: "provider stopped".into(),
351            retryable: true,
352        });
353        let failed_id = failed.id;
354        let failed = render(failed, empty_projection(failed_id)).recording;
355        assert_eq!(failed.status, "failed");
356        assert_eq!(failed.attempt_count, 3);
357        assert_eq!(failed.last_error.as_deref(), Some("provider stopped"));
358
359        let complete = recording(RecordingState::Complete {
360            transcript: "Transcript".into(),
361        });
362        let complete_id = complete.id;
363        let complete = render(complete, projection(complete_id, &["complete"])).recording;
364        assert_eq!(complete.status, "complete");
365
366        let awaiting = with_speaker_packet(
367            recording(RecordingState::Complete {
368                transcript: "Transcript".into(),
369            }),
370            ConfirmationState::Unconfirmed,
371        );
372        let awaiting_id = awaiting.id;
373        let awaiting = render(awaiting, empty_projection(awaiting_id)).recording;
374        assert_eq!(awaiting.status, "speaker_review");
375        let review = awaiting.speaker_review.unwrap();
376        assert_eq!(review.confirmation_state, ConfirmationState::Unconfirmed);
377        assert_eq!(review.observation_count, 2);
378
379        let confirmed = with_speaker_packet(
380            recording(RecordingState::Complete {
381                transcript: "Transcript".into(),
382            }),
383            ConfirmationState::Confirmed,
384        );
385        let confirmed_id = confirmed.id;
386        let confirmed = render(confirmed, empty_projection(confirmed_id)).recording;
387        assert_eq!(confirmed.status, "ready_for_ingress");
388    }
389
390    #[test]
391    fn piece_phase_priority_and_completed_counts_are_exact() {
392        let cases = [
393            (vec!["complete", "complete"], "complete", 2),
394            (
395                vec!["complete", "ingress_in_progress", "ingress_failed"],
396                "ingress_failed",
397                1,
398            ),
399            (
400                vec!["complete", "ingress_pending", "ingress_in_progress"],
401                "ingressing",
402                1,
403            ),
404            (vec!["complete", "ingress_pending"], "ready_for_ingress", 1),
405        ];
406
407        for (phases, expected_status, expected_completed) in cases {
408            let recording = recording(RecordingState::Complete {
409                transcript: "Transcript".into(),
410            });
411            let recording_id = recording.id;
412            let view = render(recording, projection(recording_id, &phases));
413            assert_eq!(view.recording.status, expected_status);
414            assert_eq!(view.recording.completed_piece_count, expected_completed);
415            assert_eq!(view.recording.transcript_piece_count, phases.len());
416        }
417    }
418
419    #[test]
420    fn recording_and_piece_fields_are_adapted_without_drift() {
421        let recording = recording(RecordingState::Complete {
422            transcript: "Transcript".into(),
423        });
424        let recording_id = recording.id;
425        let projection_id = Uuid::parse_str("2456325d-b033-4caf-92f0-e044080ed9b8").unwrap();
426        let view = render(recording, projection(projection_id, &["ingress_failed"]));
427
428        let adapted = &view.recording;
429        assert_eq!(adapted.id, recording_id);
430        assert_eq!(adapted.sha256, "ab".repeat(32));
431        assert_eq!(adapted.original_filename, "meeting.final.WAV");
432        assert_eq!(adapted.content_type, "audio/wav");
433        assert_eq!(adapted.size_bytes, 42);
434        assert_eq!(adapted.source_created_at, "2026-01-02T03:04:05+00:00");
435        assert_eq!(adapted.received_at, "2026-01-02T03:05:06+00:00");
436        assert_eq!(adapted.updated_at, adapted.received_at);
437        assert_eq!(adapted.status, "ingress_failed");
438        assert_eq!(adapted.transcription_model, "transcription-model");
439        assert_eq!(adapted.reconciliation_model, "reconciliation-model");
440        assert_eq!(adapted.reconciliation_reasoning, "xhigh");
441        assert!(adapted.transcription_status.is_none());
442        assert_eq!(adapted.attempt_count, 0);
443        assert!(adapted.next_attempt_at.is_none());
444        assert!(adapted.last_error.is_none());
445        assert!(adapted.speaker_review.is_none());
446        assert_eq!(adapted.transcript_piece_count, 1);
447        assert_eq!(adapted.completed_piece_count, 0);
448
449        let piece = &view.pieces[0];
450        assert_eq!(piece.id, "session-0");
451        assert_eq!(piece.recording_id, recording_id);
452        assert_eq!(piece.sha256, "ab".repeat(32));
453        assert_eq!(piece.original_filename, "meeting.final.WAV");
454        assert_eq!(piece.source_created_at, "2026-01-02T03:04:05+00:00");
455        assert_eq!(piece.piece_index, 0);
456        assert_eq!(piece.piece_count, 1);
457        assert_eq!(piece.transcript_text, "Transcript piece 0");
458        assert_eq!(piece.estimated_tokens, 10);
459        assert_eq!(piece.phase, "ingress_failed");
460        assert_eq!(piece.provenance_id.as_deref(), Some("provenance-0"));
461        assert_eq!(piece.state, json!({"piece":0}));
462        assert_eq!(piece.version, 7);
463        assert_eq!(piece.ingress_failure_count, 2);
464        assert_eq!(
465            piece.ingress_failures,
466            json!([{"piece":0,"message":"retained"}])
467        );
468        assert_eq!(piece.created_at, "2026-01-02T03:00:00Z");
469        assert_eq!(piece.updated_at, "2026-01-02T04:00:00Z");
470    }
471
472    fn step(step: Step, state: StepState) -> StepStatus {
473        StepStatus {
474            step,
475            state,
476            attempts: 1,
477            retry_after: None,
478            error: None,
479        }
480    }
481
482    fn progress(steps: Vec<StepStatus>) -> TranscriptionStatus {
483        TranscriptionStatus {
484            state: JobState::Running,
485            steps,
486            transcript: None,
487            correction_packet: None,
488        }
489    }
490
491    #[test]
492    fn processing_stage_mapping_and_status_json_are_preserved() {
493        let cases = [
494            (
495                "reconciling",
496                progress(vec![step(Step::ReconcileTranscript, StepState::Running)]),
497            ),
498            (
499                "chunking",
500                progress(vec![step(Step::PlanChunks, StepState::Running)]),
501            ),
502            (
503                "transcribing",
504                progress(vec![
505                    step(Step::PlanChunks, StepState::Completed),
506                    step(
507                        Step::TranscribeChunk { index: 0, total: 1 },
508                        StepState::Running,
509                    ),
510                ]),
511            ),
512            (
513                "analyzing_speakers",
514                progress(vec![
515                    step(Step::PlanChunks, StepState::Completed),
516                    step(
517                        Step::TranscribeChunk { index: 0, total: 1 },
518                        StepState::Completed,
519                    ),
520                    step(Step::AnalyzeSpeakers, StepState::Running),
521                ]),
522            ),
523            (
524                "reconciling",
525                progress(vec![
526                    step(Step::PlanChunks, StepState::Completed),
527                    step(
528                        Step::TranscribeChunk { index: 0, total: 1 },
529                        StepState::Completed,
530                    ),
531                    step(Step::AnalyzeSpeakers, StepState::Completed),
532                ]),
533            ),
534        ];
535
536        for (expected_stage, progress) in cases {
537            let expected_json = serde_json::to_value(&progress).unwrap();
538            let recording = recording(RecordingState::Processing {
539                attempt: 4,
540                progress,
541            });
542            let recording_id = recording.id;
543            let adapted = render(recording, empty_projection(recording_id)).recording;
544            assert_eq!(adapted.status, expected_stage);
545            assert_eq!(adapted.attempt_count, 4);
546            assert_eq!(adapted.transcription_status, Some(expected_json));
547        }
548    }
549}