cai
2026-07-08 e1612a0151ef28c4e02c0f1843dd725a92d4c268
src/main.rs
@@ -12,7 +12,7 @@
};
use anyhow::{Context, Result, anyhow};
use audio::{AudioDiagnostics, PcmFrame, load_pre_recorded_frames};
use audio::{AudioDiagnostics, load_pre_recorded_frames};
use base64::{Engine as _, engine::general_purpose};
use futures_util::StreamExt;
use libwebrtc::{
@@ -34,8 +34,12 @@
use tokio::{sync::mpsc::UnboundedReceiver, task::JoinHandle};
use tracing::{info, warn};
const TARGET_SAMPLE_RATE_HZ: u32 = 48_000;
const TARGET_NUM_CHANNELS: u16 = 1;
const USER_AUDIO_SAMPLE_RATE_HZ: u32 = 48_000;
const USER_AUDIO_NUM_CHANNELS: u16 = 1;
const DEFAULT_BOT_AUDIO_PROFILE: &str = "pcm-16k";
const LIVEKIT_48K_SAMPLE_RATE_HZ: u32 = 48_000;
const PCM_16K_SAMPLE_RATE_HZ: u32 = 16_000;
const BOT_NUM_CHANNELS: u16 = 1;
const TRACK_NAME: &str = "bot-main-audio";
#[tokio::main(flavor = "multi_thread")]
@@ -58,8 +62,8 @@
        &http,
        config.greeting_audio_file.as_deref(),
        config.greeting_audio_url.as_deref(),
        TARGET_SAMPLE_RATE_HZ,
        TARGET_NUM_CHANNELS,
        config.bot_audio_profile.sample_rate_hz,
        config.bot_audio_profile.num_channels,
        config.audio_debug_dump_dir.as_deref(),
        &config.call_id,
        "greeting",
@@ -81,6 +85,9 @@
        room_alias = %redact(&config.room_id),
        participant_alias = %redact(&config.bot_participant_identity),
        greeting_source = %config.greeting_source,
        bot_audio_profile = %config.bot_audio_profile.profile,
        bot_sample_rate_hz = config.bot_audio_profile.sample_rate_hz,
        bot_num_channels = config.bot_audio_profile.num_channels,
        "combrabo voice runtime helper connected"
    );
    emit_activity(
@@ -129,8 +136,9 @@
        &config.call_id,
        &config.trace_id,
        TRACK_NAME,
        TARGET_SAMPLE_RATE_HZ,
        u32::from(TARGET_NUM_CHANNELS),
        config.bot_audio_profile.profile.clone(),
        config.bot_audio_profile.sample_rate_hz,
        u32::from(config.bot_audio_profile.num_channels),
        config.user_participant_identity.clone(),
    )
    .await?;
@@ -269,6 +277,54 @@
    simple_vad_gate_until_greeting_done: bool,
    simple_vad_post_greeting_delay_ms: u64,
    simple_vad_config: SimpleVadConfig,
    bot_audio_profile: BotAudioProfile,
}
#[derive(Clone)]
struct BotAudioProfile {
    profile: String,
    sample_rate_hz: u32,
    num_channels: u16,
}
impl BotAudioProfile {
    fn from_env() -> Result<Self> {
        let profile = env::var("CV_BOT_AUDIO_PROFILE")
            .unwrap_or_else(|_| DEFAULT_BOT_AUDIO_PROFILE.to_string())
            .trim()
            .to_ascii_lowercase();
        match profile.as_str() {
            "livekit-48k" | "48k" => Ok(Self {
                profile: "livekit-48k".to_string(),
                sample_rate_hz: LIVEKIT_48K_SAMPLE_RATE_HZ,
                num_channels: BOT_NUM_CHANNELS,
            }),
            "pcm-16k" | "16k" => Ok(Self {
                profile: "pcm-16k".to_string(),
                sample_rate_hz: PCM_16K_SAMPLE_RATE_HZ,
                num_channels: BOT_NUM_CHANNELS,
            }),
            "custom" => {
                let sample_rate_hz = u32_env("CV_BOT_SAMPLE_RATE_HZ", LIVEKIT_48K_SAMPLE_RATE_HZ);
                let num_channels = u16_env("CV_BOT_NUM_CHANNELS", BOT_NUM_CHANNELS);
                if sample_rate_hz == 0 {
                    return Err(anyhow!("CV_BOT_SAMPLE_RATE_HZ must be positive"));
                }
                if num_channels == 0 {
                    return Err(anyhow!("CV_BOT_NUM_CHANNELS must be positive"));
                }
                Ok(Self {
                    profile,
                    sample_rate_hz,
                    num_channels,
                })
            }
            _ => Err(anyhow!(
                "unsupported CV_BOT_AUDIO_PROFILE {}; expected livekit-48k, pcm-16k or custom",
                profile
            )),
        }
    }
}
#[derive(Clone)]
@@ -346,6 +402,7 @@
            simple_vad_gate_until_greeting_done: bool_env("CV_VAD_GATE_UNTIL_GREETING_DONE", true),
            simple_vad_post_greeting_delay_ms: u64_env("CV_VAD_POST_GREETING_DELAY_MS", 800),
            simple_vad_config: SimpleVadConfig::from_env(),
            bot_audio_profile: BotAudioProfile::from_env()?,
        })
    }
}
@@ -489,6 +546,13 @@
    env::var(key)
        .ok()
        .and_then(|value| value.trim().parse::<u32>().ok())
        .unwrap_or(default_value)
}
fn u16_env(key: &str, default_value: u16) -> u16 {
    env::var(key)
        .ok()
        .and_then(|value| value.trim().parse::<u16>().ok())
        .unwrap_or(default_value)
}
@@ -930,8 +994,8 @@
    audio::write_pcm_wav(
        &output_path,
        &turn.samples,
        TARGET_SAMPLE_RATE_HZ,
        TARGET_NUM_CHANNELS,
        USER_AUDIO_SAMPLE_RATE_HZ,
        USER_AUDIO_NUM_CHANNELS,
    )?;
    let byte_size = fs::metadata(&output_path)
        .context("failed to stat turn artifact")?
@@ -959,8 +1023,8 @@
            artifact_type: "local_file".to_string(),
            path_ref: path_ref.to_string(),
            format: "wav".to_string(),
            sample_rate: TARGET_SAMPLE_RATE_HZ,
            channels: u32::from(TARGET_NUM_CHANNELS),
            sample_rate: USER_AUDIO_SAMPLE_RATE_HZ,
            channels: u32::from(USER_AUDIO_NUM_CHANNELS),
            duration_ms: turn.duration_ms,
            byte_size,
        },
@@ -1343,52 +1407,59 @@
        .trim()
        .to_ascii_lowercase();
    let frames = if format == "pcm_s16le" {
        let frame_alignment_bytes = pcm_s16le_frame_alignment_bytes(
            audio_chunk
                .channels
                .unwrap_or(u32::from(TARGET_NUM_CHANNELS)),
        )?;
        let (aligned_payload, dropped_tail_bytes) = take_aligned_pcm_payload(
            &mut state.pcm_audio_buffer,
            &payload,
            frame_alignment_bytes,
            audio_chunk.last.unwrap_or(false),
        );
        if dropped_tail_bytes > 0 {
        let sample_rate = audio_chunk.sample_rate.unwrap_or(sink.sample_rate_hz);
        let channels = audio_chunk.channels.unwrap_or(u32::from(sink.num_channels));
        if state.pcm_stream_decoder.is_none() {
            state.pcm_stream_decoder = Some(audio::PcmS16leStreamDecoder::new(
                sample_rate,
                channels,
                sink.sample_rate_hz,
                sink.num_channels,
            )?);
        }
        state.pcm_stream_network_chunk_count =
            state.pcm_stream_network_chunk_count.saturating_add(1);
        let stream_result = state
            .pcm_stream_decoder
            .as_mut()
            .expect("pcm stream decoder initialized")
            .push_bytes(
                &payload,
                sample_rate,
                channels,
                audio_chunk.last.unwrap_or(false),
            )?;
        if stream_result.dropped_tail_bytes > 0 {
            warn!(
                call_id = %call_id,
                trace_id = %trace_id,
                turn_id = %turn.turn_id,
                chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(),
                dropped_tail_bytes,
                dropped_tail_bytes = stream_result.dropped_tail_bytes,
                "runtime helper stream_audio_pcm_unaligned_tail_dropped"
            );
        }
        if aligned_payload.is_empty() {
        if stream_result.frames.is_empty() {
            warn!(
                call_id = %call_id,
                trace_id = %trace_id,
                turn_id = %turn.turn_id,
                chunk_seq = audio_chunk.chunk_seq.unwrap_or_default(),
                buffered_bytes = state.pcm_audio_buffer.len(),
                "runtime helper stream_audio_pcm_waiting_for_sample_boundary"
                buffered_source_bytes = stream_result.buffered_source_bytes,
                buffered_source_samples = stream_result.buffered_source_samples,
                network_chunk_count = state.pcm_stream_network_chunk_count,
                "runtime helper stream_audio_pcm_waiting_for_20ms_frame"
            );
            return Ok(0);
        }
        pcm_s16le_payload_to_frames(
            &aligned_payload,
            audio_chunk.sample_rate.unwrap_or(TARGET_SAMPLE_RATE_HZ),
            audio_chunk
                .channels
                .unwrap_or(u32::from(TARGET_NUM_CHANNELS)),
        )?
        stream_result.frames
    } else if matches!(format.as_str(), "mp3" | "mpeg" | "wav") {
        state.encoded_audio_buffer.extend_from_slice(&payload);
        match audio::decode_audio_bytes_to_frames(
            &state.encoded_audio_buffer,
            "stream_chunk",
            TARGET_SAMPLE_RATE_HZ,
            TARGET_NUM_CHANNELS,
            sink.sample_rate_hz,
            sink.num_channels,
            bridge_config.audio_debug_dump_dir.as_deref(),
            call_id,
            &format!("stream-reply-{}", turn.turn_id),
@@ -1473,6 +1544,34 @@
        sleep_until(pacing_started_at + Duration::from_millis(((index + 1) as u64) * 20)).await;
    }
    if audio_chunk.last.unwrap_or(false) {
        let mut debug_source_path = None;
        let mut debug_pcm_wav_path = None;
        let mut debug_pcm_wav_size_bytes = None;
        if format == "pcm_s16le" {
            if let (Some(debug_dump_dir), Some(decoder)) = (
                bridge_config.audio_debug_dump_dir.as_deref(),
                state.pcm_stream_decoder.as_ref(),
            ) {
                match decoder.write_debug_dump(
                    debug_dump_dir,
                    call_id,
                    &format!("stream-reply-{}", turn.turn_id),
                ) {
                    Ok(debug_dump) => {
                        debug_source_path = debug_dump.debug_source_path;
                        debug_pcm_wav_path = debug_dump.debug_pcm_wav_path;
                        debug_pcm_wav_size_bytes = debug_dump.debug_pcm_wav_size_bytes;
                    }
                    Err(error) => warn!(
                        call_id = %call_id,
                        trace_id = %trace_id,
                        turn_id = %turn.turn_id,
                        error = %safe_error(&error.to_string()),
                        "runtime helper stream_audio_pcm_debug_dump_failed"
                    ),
                }
            }
        }
        emit_activity(
            call_id,
            trace_id,
@@ -1485,94 +1584,20 @@
                "replyPlaybackMode": state.reply_playback_mode.as_str(),
                "format": format.as_str(),
                "chunkSeq": audio_chunk.chunk_seq,
                "networkChunkCount": state.pcm_stream_network_chunk_count,
                "debugSourcePath": debug_source_path,
                "debugPcmWavPath": debug_pcm_wav_path,
                "debugPcmWavSizeBytes": debug_pcm_wav_size_bytes,
                "sourceSampleRate": audio_chunk.sample_rate,
                "sourceChannels": audio_chunk.channels,
                "targetAudioProfile": sink.profile.as_str(),
                "targetSampleRate": sink.sample_rate_hz,
                "targetChannels": sink.num_channels,
                "replyTotalAfterVadEndMs": turn_pipeline_started_at.elapsed().as_millis() as u64,
            }),
        );
    }
    Ok(frames.len())
}
fn pcm_s16le_payload_to_frames(
    payload: &[u8],
    sample_rate: u32,
    channels: u32,
) -> Result<Vec<PcmFrame>> {
    audio::pcm_s16le_bytes_to_frames(
        payload,
        sample_rate,
        channels,
        TARGET_SAMPLE_RATE_HZ,
        TARGET_NUM_CHANNELS,
    )
}
fn pcm_s16le_frame_alignment_bytes(channels: u32) -> Result<usize> {
    if channels == 0 {
        return Err(anyhow!("pcm_s16le channel count is zero"));
    }
    let channel_count = usize::try_from(channels)
        .map_err(|_| anyhow!("unsupported pcm_s16le channel count {channels}"))?;
    Ok(2 * channel_count)
}
fn take_aligned_pcm_payload(
    buffer: &mut Vec<u8>,
    payload: &[u8],
    frame_alignment_bytes: usize,
    last: bool,
) -> (Vec<u8>, usize) {
    let alignment = frame_alignment_bytes.max(2);
    buffer.extend_from_slice(payload);
    let aligned_len = buffer.len() - buffer.len() % alignment;
    let aligned_payload = if aligned_len == 0 {
        Vec::new()
    } else {
        buffer.drain(..aligned_len).collect()
    };
    let dropped_tail_bytes = if last && !buffer.is_empty() {
        let dropped = buffer.len();
        buffer.clear();
        dropped
    } else {
        0
    };
    (aligned_payload, dropped_tail_bytes)
}
#[cfg(test)]
mod pcm_stream_tests {
    use super::take_aligned_pcm_payload;
    #[test]
    fn take_aligned_pcm_payload_buffers_split_sample_bytes() {
        let mut buffer = Vec::new();
        let (first, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x01], 2, false);
        assert!(first.is_empty());
        assert_eq!(0, dropped);
        assert_eq!(vec![0x01], buffer);
        let (second, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x02, 0x03], 2, false);
        assert_eq!(vec![0x01, 0x02], second);
        assert_eq!(0, dropped);
        assert_eq!(vec![0x03], buffer);
        let (third, dropped) = take_aligned_pcm_payload(&mut buffer, &[0x04], 2, true);
        assert_eq!(vec![0x03, 0x04], third);
        assert_eq!(0, dropped);
        assert!(buffer.is_empty());
    }
    #[test]
    fn take_aligned_pcm_payload_drops_final_half_sample() {
        let mut buffer = Vec::new();
        let (payload, dropped) =
            take_aligned_pcm_payload(&mut buffer, &[0x01, 0x02, 0x03], 2, true);
        assert_eq!(vec![0x01, 0x02], payload);
        assert_eq!(1, dropped);
        assert!(buffer.is_empty());
    }
}
fn trim_ascii_whitespace(value: &[u8]) -> &[u8] {
@@ -1606,8 +1631,8 @@
            artifact_type: "local_file".to_string(),
            path_ref: path_ref.to_string(),
            format: "wav".to_string(),
            sample_rate: TARGET_SAMPLE_RATE_HZ,
            channels: u32::from(TARGET_NUM_CHANNELS),
            sample_rate: USER_AUDIO_SAMPLE_RATE_HZ,
            channels: u32::from(USER_AUDIO_NUM_CHANNELS),
            duration_ms: turn.duration_ms,
            byte_size,
        },
@@ -1770,8 +1795,8 @@
        http,
        Some(&audio_path_string),
        None,
        TARGET_SAMPLE_RATE_HZ,
        TARGET_NUM_CHANNELS,
        sink.sample_rate_hz,
        sink.num_channels,
        bridge_config.audio_debug_dump_dir.as_deref(),
        call_id,
        &reply_debug_label,
@@ -1950,7 +1975,8 @@
    audio_chunk_count: u64,
    device_output_count: u64,
    encoded_audio_buffer: Vec<u8>,
    pcm_audio_buffer: Vec<u8>,
    pcm_stream_decoder: Option<audio::PcmS16leStreamDecoder>,
    pcm_stream_network_chunk_count: u64,
}
impl Default for RuntimeTurnStreamState {
@@ -1964,7 +1990,8 @@
            audio_chunk_count: 0,
            device_output_count: 0,
            encoded_audio_buffer: Vec::new(),
            pcm_audio_buffer: Vec::new(),
            pcm_stream_decoder: None,
            pcm_stream_network_chunk_count: 0,
        }
    }
}
@@ -2113,8 +2140,8 @@
    tokio::spawn(async move {
        let mut stream = NativeAudioStream::new(
            track.rtc_track(),
            TARGET_SAMPLE_RATE_HZ as i32,
            i32::from(TARGET_NUM_CHANNELS),
            USER_AUDIO_SAMPLE_RATE_HZ as i32,
            i32::from(USER_AUDIO_NUM_CHANNELS),
        );
        let started_at = Instant::now();
        let mut frame_count: u64 = 0;
@@ -2601,6 +2628,9 @@
    rtc_source: NativeAudioSource,
    track: LocalAudioTrack,
    device_output_destination_identity: Option<String>,
    profile: String,
    sample_rate_hz: u32,
    num_channels: u16,
}
impl BotAudioOutputSink {
@@ -2611,6 +2641,7 @@
        call_id: &str,
        trace_id: &str,
        track_name: &str,
        profile: String,
        sample_rate: u32,
        num_channels: u32,
        device_output_destination_identity: Option<String>,
@@ -2639,11 +2670,14 @@
                    "failed to publish bot audio track in room {room_alias} for participant {participant_alias}: {error}"
                )
            })?;
        let num_channels_u16 = u16::try_from(num_channels)
            .map_err(|_| anyhow!("unsupported bot audio channel count {num_channels}"))?;
        info!(
            room_alias = %room_alias,
            participant_alias = %participant_alias,
            track_name = %track_name,
            bot_audio_profile = %profile,
            sample_rate,
            num_channels,
            "runtime helper published bot audio track"
@@ -2658,6 +2692,7 @@
            None,
            json!({
                "trackName": track_name,
                "audioProfile": profile,
                "sampleRate": sample_rate,
                "numChannels": num_channels,
            }),
@@ -2668,6 +2703,9 @@
            rtc_source,
            track,
            device_output_destination_identity,
            profile,
            sample_rate_hz: sample_rate,
            num_channels: num_channels_u16,
        })
    }