| | |
| | | .await |
| | | { |
| | | Ok(outcome) => { |
| | | let mut published_device_outputs = HashSet::new(); |
| | | for output in &outcome.device_outputs { |
| | | if !should_publish_device_output(&mut published_device_outputs, output) { |
| | | continue; |
| | | } |
| | | if let Err(error) = sink |
| | | .publish_device_output(call_id, trace_id, &turn.turn_id, output) |
| | | .await |
| | |
| | | } |
| | | Some("device_output") => { |
| | | if let Some(output) = event.device_output.as_ref() { |
| | | if !should_publish_device_output(&mut state.published_device_output_ids, output) { |
| | | return Ok(()); |
| | | } |
| | | sink.publish_device_output(call_id, trace_id, &turn.turn_id, output) |
| | | .await?; |
| | | state.device_output_count = state.device_output_count.saturating_add(1); |
| | |
| | | completed: bool, |
| | | audio_chunk_count: u64, |
| | | device_output_count: u64, |
| | | published_device_output_ids: HashSet<String>, |
| | | encoded_audio_buffer: Vec<u8>, |
| | | pcm_stream_decoder: Option<audio::PcmS16leStreamDecoder>, |
| | | pcm_stream_network_chunk_count: u64, |
| | |
| | | completed: false, |
| | | audio_chunk_count: 0, |
| | | device_output_count: 0, |
| | | published_device_output_ids: HashSet::new(), |
| | | encoded_audio_buffer: Vec::new(), |
| | | pcm_stream_decoder: None, |
| | | pcm_stream_network_chunk_count: 0, |
| | |
| | | #[serde(rename = "commandCode")] |
| | | command_code: Option<String>, |
| | | params: Option<serde_json::Value>, |
| | | } |
| | | |
| | | fn should_publish_device_output( |
| | | published_ids: &mut HashSet<String>, |
| | | output: &RuntimeTurnDeviceOutput, |
| | | ) -> bool { |
| | | let Some(command_id) = output |
| | | .command_id |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | else { |
| | | return false; |
| | | }; |
| | | if output |
| | | .command_code |
| | | .as_deref() |
| | | .map(str::trim) |
| | | .filter(|value| !value.is_empty()) |
| | | .is_none() |
| | | { |
| | | return false; |
| | | } |
| | | published_ids.insert(command_id.to_string()) |
| | | } |
| | | |
| | | fn require_safe_segment(value: &str) -> Result<()> { |
| | |
| | | #[cfg(test)] |
| | | mod tests { |
| | | use super::{ |
| | | ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnStreamEvent, RuntimeTurnStreamState, |
| | | RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash, |
| | | ReplyChunkMarker, ReplyChunkMarkerState, RuntimeTurnDeviceOutput, RuntimeTurnStreamEvent, |
| | | RuntimeTurnStreamState, RuntimeTurnStreamTimingPhase, runtime_session_nonce_hash, |
| | | should_publish_device_output, |
| | | }; |
| | | use std::collections::HashSet; |
| | | |
| | | #[test] |
| | | fn reply_chunk_marker_state_emits_turn_first_once_and_later_segment_first_once() { |
| | |
| | | ); |
| | | assert!(state.record_m7().is_none()); |
| | | } |
| | | |
| | | #[test] |
| | | fn device_output_contract_is_reliable_and_deduplicated() { |
| | | let output = RuntimeTurnDeviceOutput { |
| | | command_id: Some("cmd-1".to_string()), |
| | | command_code: Some("custom.app.DeviceLevelChange".to_string()), |
| | | params: Some(serde_json::json!({"level": 1})), |
| | | }; |
| | | let mut published = HashSet::new(); |
| | | assert!(should_publish_device_output(&mut published, &output)); |
| | | assert!(!should_publish_device_output(&mut published, &output)); |
| | | assert_eq!(published.len(), 1); |
| | | } |
| | | |
| | | #[test] |
| | | fn device_output_invalid_or_missing_command_is_fail_closed() { |
| | | for output in [ |
| | | RuntimeTurnDeviceOutput { |
| | | command_id: None, |
| | | command_code: Some("custom.app.DeviceLevelChange".to_string()), |
| | | params: None, |
| | | }, |
| | | RuntimeTurnDeviceOutput { |
| | | command_id: Some("cmd-1".to_string()), |
| | | command_code: None, |
| | | params: None, |
| | | }, |
| | | ] { |
| | | let mut published = HashSet::new(); |
| | | assert!(!should_publish_device_output(&mut published, &output)); |
| | | assert!(published.is_empty()); |
| | | } |
| | | } |
| | | } |