From 8673dd9db0467f9ce9b930316839de752fc70385 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Sat, 08 Aug 2026 04:37:06 +0800
Subject: [PATCH] feat: pair helper device output downlink

---
 helper-paired-manifest.json |   28 ++++++++++++++
 src/main.rs                 |   72 +++++++++++++++++++++++++++++++++++-
 2 files changed, 98 insertions(+), 2 deletions(-)

diff --git a/helper-paired-manifest.json b/helper-paired-manifest.json
new file mode 100644
index 0000000..64327c5
--- /dev/null
+++ b/helper-paired-manifest.json
@@ -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"
+  }
+}
diff --git a/src/main.rs b/src/main.rs
index 611993a..e76d222 100644
--- a/src/main.rs
+++ b/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());
+        }
+    }
 }

--
Gitblit v1.9.3