inferencelayer 0.2.8

Kortexya's engine-native inference layer — LLM generation + embedding/encoder family on wgpu (WGSL kernels, any adapter) with a pure-Rust CPU fallback
Documentation
//! Realtime input discipline for the duplex loop, extracted so it is testable in isolation
//! and shared with the server binary.
//!
//! The one rule that matters: when the engine falls behind, DROP SILENCE, KEEP SPEECH. Moshi
//! is a time-synchronized dual-stream model — shredding voiced input desynchronizes its
//! timeline and it goes MUTE (measured live: 16.6 s of intelligible speech in, zero audio
//! out, because the all-GPU loop dropped the oldest frames blindly). Silence-first shedding
//! keeps every voiced frame under overload.

/// 24 kHz samples per 80 ms frame.
pub const FRAME_SIZE: usize = 1920;
/// Above this backlog the engine is too far behind — shed down to `KEEP_FRAMES`.
pub const MAX_LAG_FRAMES: usize = 12;
/// Frames kept after a shed (freshest speech + room tone).
pub const KEEP_FRAMES: usize = 9;

/// `rms > 0.01` — speech-level energy.
pub fn is_voiced(frame: &[f32]) -> bool {
    (frame.iter().map(|v| v * v).sum::<f32>() / frame.len().max(1) as f32).sqrt() > 0.01
}

/// If `buf` (contiguous 24 kHz mono) exceeds the lag cap, shed it down to `KEEP_FRAMES` of
/// 1920 samples, SILENT frames first: every voiced frame survives while there are ≤ KEEP of
/// them, and only if voiced speech alone still exceeds the budget do we keep its freshest
/// tail. Returns the number of samples dropped (0 if under the cap). In place.
pub fn shed_silence_first(buf: &mut Vec<f32>) -> usize {
    if buf.len() <= MAX_LAG_FRAMES * FRAME_SIZE {
        return 0;
    }
    let frames: Vec<&[f32]> = buf.chunks_exact(FRAME_SIZE).collect();
    let voiced: Vec<bool> = frames.iter().map(|f| is_voiced(f)).collect();
    let n_voiced = voiced.iter().filter(|v| **v).count();
    let mut kept: Vec<f32> = Vec::with_capacity(KEEP_FRAMES * FRAME_SIZE);
    if n_voiced <= KEEP_FRAMES {
        for (f, v) in frames.iter().zip(&voiced) {
            if *v {
                kept.extend_from_slice(f);
            }
        }
        // top up with the freshest frames so the model still hears the room tone
        let need = KEEP_FRAMES.saturating_sub(kept.len() / FRAME_SIZE);
        let extra: Vec<&&[f32]> = frames
            .iter()
            .zip(&voiced)
            .rev()
            .filter(|(_, v)| !**v)
            .map(|(f, _)| f)
            .take(need)
            .collect();
        for f in extra.iter().rev() {
            kept.extend_from_slice(f);
        }
    } else {
        // keep the freshest KEEP_FRAMES voiced frames, in order
        for (f, v) in frames.iter().zip(&voiced).skip(n_voiced - KEEP_FRAMES) {
            if *v {
                kept.extend_from_slice(f);
            }
        }
    }
    let dropped = buf.len() - kept.len();
    *buf = kept;
    dropped
}

/// Soft-limit a clipping microphone in place: a caller's mic measured 0.999 peak, and clipped
/// square tops are broadband distortion the codec never saw in training. Frames above the
/// ceiling are scaled back under it. Returns true if it acted.
pub fn soft_limit(frame: &mut [f32], ceiling: f32) -> bool {
    let peak = frame.iter().fold(0f32, |m, v| m.max(v.abs()));
    if peak > ceiling {
        let g = ceiling / peak;
        for v in frame.iter_mut() {
            *v *= g;
        }
        true
    } else {
        false
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    fn frame(voiced: bool) -> Vec<f32> {
        vec![if voiced { 0.3 } else { 0.0 }; FRAME_SIZE]
    }

    #[test]
    fn shedding_keeps_every_voiced_frame_when_few() {
        // 20 frames: the first 4 voiced (the user's words), the rest silence
        let mut buf = Vec::new();
        for i in 0..20 {
            buf.extend(frame(i < 4));
        }
        let before_voiced = buf.chunks_exact(FRAME_SIZE).filter(|f| is_voiced(f)).count();
        let dropped = shed_silence_first(&mut buf);
        assert!(dropped > 0, "should have shed (20 > cap)");
        let after = buf.chunks_exact(FRAME_SIZE).collect::<Vec<_>>();
        assert_eq!(after.len(), KEEP_FRAMES, "kept exactly KEEP_FRAMES");
        let after_voiced = after.iter().filter(|f| is_voiced(f)).count();
        assert_eq!(after_voiced, before_voiced, "every voiced frame survived");
    }

    #[test]
    fn shedding_keeps_freshest_speech_when_over_budget() {
        // 16 all-voiced frames: more speech than the budget → keep the freshest KEEP_FRAMES
        let mut buf = Vec::new();
        for _ in 0..16 {
            buf.extend(frame(true));
        }
        let dropped = shed_silence_first(&mut buf);
        assert!(dropped > 0);
        assert_eq!(buf.len(), KEEP_FRAMES * FRAME_SIZE);
    }

    #[test]
    fn no_shed_under_the_cap() {
        let mut buf = frame(true);
        buf.extend(frame(false));
        let len = buf.len();
        assert_eq!(shed_silence_first(&mut buf), 0);
        assert_eq!(buf.len(), len);
    }

    #[test]
    fn soft_limit_scales_clipping_frame() {
        let mut f = vec![0.999f32; 64];
        assert!(soft_limit(&mut f, 0.95));
        assert!(f.iter().all(|v| *v <= 0.951));
        let mut g = vec![0.5f32; 64];
        assert!(!soft_limit(&mut g, 0.95));
    }
}