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