cai
2026-08-08 8673dd9db0467f9ce9b930316839de752fc70385
feat: pair helper device output downlink
1 files modified
1 files added
100 ■■■■■ changed files
helper-paired-manifest.json 28 ●●●●● patch | view | raw | blame | history
src/main.rs 72 ●●●●● patch | view | raw | blame | history
helper-paired-manifest.json
New file
@@ -0,0 +1,28 @@
{
  "schemaVersion": "1",
  "repository": "lm-livekit-helper",
  "baseCommit": "b0a4e2ea93fafc4d65d474f3ada616a0ad5e19d5",
  "rollbackCommit": "0cec0c4ffed92f277778adbbd14eea3e5aa562b6",
  "source": {
    "lmDocHelperCommit": "18008c8c88ba5144b2d728a6f09ce7a46e55e41d",
    "lmDocHelperTree": "3462a8af48ca986d4f9f0ea4e5babc551b6f0766",
    "manifestCommit": "39dd2f93680a1799845b32e0d24eb9e6ee724954"
  },
  "pairedJava": {
    "commit": "cdbd85605212297149b434b2cd3154c55f8559d3",
    "tree": "b085cd0a56ce5f77cbcb143a417134248f83a98e"
  },
  "downlinkContract": {
    "inputEvent": "device_output",
    "topic": "device_output",
    "reliable": true,
    "dedupeKey": "commandId",
    "invalidInput": "fail_closed"
  },
  "verification": {
    "status": "NEEDS_CI_VALIDATION",
    "localCargo": "WEBRTC_SYS_BUILD_TIMEOUT",
    "linuxArtifact": "PENDING_CI",
    "scope": "src/main.rs only; no dependency, protocol, config, or runtime changes"
  }
}
src/main.rs
@@ -985,7 +985,11 @@
            .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
@@ -1593,6 +1597,9 @@
        }
        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);
@@ -2421,6 +2428,7 @@
    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,
@@ -2438,6 +2446,7 @@
            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,
@@ -2828,6 +2837,30 @@
    #[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<()> {
@@ -3894,9 +3927,11 @@
#[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() {
@@ -3993,4 +4028,37 @@
        );
        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());
        }
    }
}