1#![forbid(unsafe_code)]
2
3use 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#[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#[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#[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#[derive(Clone, Debug)]
68pub struct RecordingView {
69 pub recording: Recording,
70 pub pieces: Vec<IngressPiece>,
71}
72
73pub 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}