From fe5c23af3c46d6a6168d9b77bd703065aa65cbc5 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Sat, 22 Aug 2026 10:32:18 +0800
Subject: [PATCH] feat(asr): propagate participant language per session
---
src/main.rs | 138 +++++++++++++++++++++++++++++-----------------
1 files changed, 87 insertions(+), 51 deletions(-)
diff --git a/src/main.rs b/src/main.rs
index af61418..0e179f5 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -357,6 +357,25 @@
})
}
+fn controlled_fixture_probe_binding_decision(
+ probe_sender: &str,
+ current_audio_participant: Option<&str>,
+ expected_participant: Option<&str>,
+) -> Result<(), &'static str> {
+ if !is_bound_user_participant(probe_sender, expected_participant) {
+ return Err("wrong_participant");
+ }
+ let Some(current_audio_participant) = current_audio_participant else {
+ return Err("no_current_participant");
+ };
+ if current_audio_participant != probe_sender
+ || !is_bound_user_participant(current_audio_participant, expected_participant)
+ {
+ return Err("wrong_participant");
+ }
+ Ok(())
+}
+
fn controlled_fixture_visibility_bucket(elapsed: Duration) -> &'static str {
if elapsed <= Duration::from_millis(250) {
"lte_250ms"
@@ -1410,6 +1429,18 @@
}
current_user_participant = Some(participant_for_probe);
}
+ RoomEvent::TrackUnsubscribed {
+ track: RemoteTrack::Audio(_),
+ publication: _,
+ participant,
+ } => {
+ if current_user_participant
+ .as_ref()
+ .is_some_and(|current| current.identity() == participant.identity())
+ {
+ current_user_participant = None;
+ }
+ }
RoomEvent::DataReceived {
payload,
topic: Some(topic),
@@ -1584,16 +1615,25 @@
);
return Some(false);
}
- let decision = observe_controlled_fixture_attributes(
- probe.expires_at,
- &participant.identity().to_string(),
+ let decision = controlled_fixture_probe_binding_decision(
+ probe.sender.as_str(),
+ Some(participant.identity().as_str()),
expected_participant,
- &probe.sequence,
- || participant.attributes(),
- )
- .await;
+ );
+ if let Err(reason) = decision {
+ record_controlled_fixture_attribute_decision(
+ Err(reason),
+ runtime_call_id,
+ runtime_trace_id,
+ &probe.call_id_hash,
+ &probe.call_trace_id_hash,
+ probe.generation,
+ &probe.sequence,
+ );
+ return Some(false);
+ }
let (ack_result, reject_reason, observed) = record_controlled_fixture_attribute_decision(
- decision,
+ Ok(()),
runtime_call_id,
runtime_trace_id,
&probe.call_id_hash,
@@ -4912,35 +4952,42 @@
#[derive(Debug)]
enum PreAudioOrderEvent {
- DataReceived {
- sender: String,
- sequence: String,
- },
- TrackSubscribed {
- participant: String,
- attributes: HashMap<String, String>,
- },
+ DataReceived { sender: String, sequence: String },
+ TrackSubscribed { participant: String },
}
fn drive_pre_audio_order_test_seam(events: &[PreAudioOrderEvent]) -> Vec<&'static str> {
- let mut pending_sequence = None;
+ let call_id = "production-order-call";
+ let trace_id = "production-order-trace";
+ let mut pending_probe = None;
let mut effects = Vec::new();
for event in events {
match event {
- PreAudioOrderEvent::DataReceived { sender, sequence } if sender == "user-1" => {
- pending_sequence = Some(sequence.as_str());
+ PreAudioOrderEvent::DataReceived { sender, sequence } => {
+ let payload = serde_json::to_vec(&json!({
+ "type": CONTROLLED_FIXTURE_PROBE_TOPIC,
+ "protocolVersion": CONTROLLED_FIXTURE_PROTOCOL_VERSION,
+ "callIdHash": sha256_hex(call_id),
+ "callTraceIdHash": sha256_hex(trace_id),
+ "generation": CONTROLLED_FIXTURE_GENERATION,
+ "clientFixtureSequence": sequence,
+ }))
+ .expect("production probe payload");
+ pending_probe = controlled_fixture_probe(
+ &payload,
+ call_id,
+ trace_id,
+ &ParticipantIdentity(sender.clone()),
+ Some("user-1"),
+ );
}
- PreAudioOrderEvent::TrackSubscribed {
- participant,
- attributes,
- } => {
- let pending = pending_sequence.is_some();
- let probe_result = pending_sequence.map(|sequence| {
- classify_controlled_fixture_attributes(
- participant,
+ PreAudioOrderEvent::TrackSubscribed { participant } => {
+ let pending = pending_probe.is_some();
+ let probe_result = pending_probe.as_ref().map(|probe| {
+ controlled_fixture_probe_binding_decision(
+ probe.sender.as_str(),
+ Some(participant),
Some("user-1"),
- attributes,
- sequence,
)
.is_ok()
});
@@ -4958,9 +5005,8 @@
if observer_started {
effects.push("observer_started");
}
- pending_sequence = None;
+ pending_probe = None;
}
- PreAudioOrderEvent::DataReceived { .. } => {}
}
}
effects
@@ -4994,6 +5040,7 @@
"clientFixtureSequence".to_string(),
"fixture-01".to_string(),
),
+ ("language".to_string(), "ja-JP".to_string()),
]);
let mut starts = Vec::new();
for (session_index, sequence) in [(1, "fixture-01"), (2, "fixture-02")] {
@@ -5019,6 +5066,7 @@
let is_in_speech = vad.in_speech;
assert!(!was_in_speech && is_in_speech);
attributes.insert("clientFixtureSequence".to_string(), sequence.to_string());
+ attributes.insert("language".to_string(), "ja-JP".to_string());
let metadata = AudioIngressMetadata::from_participant(&attributes)
.expect("valid participant attributes")
.expect("controlled fixture metadata");
@@ -5033,6 +5081,7 @@
let session_json: serde_json::Value =
serde_json::from_slice(&session_line).expect("session start json");
assert_eq!(sequence, session_json["clientFixtureSequence"]);
+ assert_eq!("ja-JP", session_json["language"]);
starts.push(metadata.client_fixture_sequence);
vad.reset_current_turn();
}
@@ -5125,6 +5174,7 @@
"clientFixtureSequence".to_string(),
"fixture-01".to_string(),
),
+ ("language".to_string(), "ja-JP".to_string()),
]);
let mut upload = None;
let mut last_fixture_sequence = None;
@@ -5161,6 +5211,7 @@
"clientFixtureSequence".to_string(),
"fixture-02".to_string(),
);
+ attrs.insert("language".to_string(), "zh-CN".to_string());
let (was, is, turn) = observe_bound_participant_frame(
"participant-user",
Some("participant-user"),
@@ -5191,6 +5242,8 @@
.expect("second session request");
assert!(first_request.contains("\"clientFixtureSequence\":\"fixture-01\""));
assert!(second_request.contains("\"clientFixtureSequence\":\"fixture-02\""));
+ assert!(first_request.contains("\"language\":\"ja-JP\""));
+ assert!(second_request.contains("\"language\":\"zh-CN\""));
assert!(
first_request.contains("\"audioIngressOriginStatus\":\"controlled_fixture_bound\"")
);
@@ -6174,24 +6227,13 @@
#[test]
fn production_event_order_probe_then_track_publishes_ack_before_observer() {
- let attributes = HashMap::from([
- (
- "inputSourceCategory".to_string(),
- "controlled_fixture".to_string(),
- ),
- (
- "clientFixtureSequence".to_string(),
- "fixture-01".to_string(),
- ),
- ]);
let effects = drive_pre_audio_order_test_seam(&[
PreAudioOrderEvent::DataReceived {
sender: "user-1".to_string(),
- sequence: "fixture-01".to_string(),
+ sequence: "1".to_string(),
},
PreAudioOrderEvent::TrackSubscribed {
participant: "user-1".to_string(),
- attributes,
},
]);
assert_eq!(effects, ["ack_observed", "observer_started"]);
@@ -6199,19 +6241,13 @@
#[test]
fn production_event_order_negative_probe_has_no_observer_or_session_effect() {
- let mut invalid = HashMap::new();
- invalid.insert(
- "inputSourceCategory".to_string(),
- "ordinary_mic".to_string(),
- );
let effects = drive_pre_audio_order_test_seam(&[
PreAudioOrderEvent::DataReceived {
sender: "user-1".to_string(),
- sequence: "fixture-01".to_string(),
+ sequence: "1".to_string(),
},
PreAudioOrderEvent::TrackSubscribed {
- participant: "user-1".to_string(),
- attributes: invalid,
+ participant: "cross-call-user".to_string(),
},
]);
assert!(effects.is_empty());
--
Gitblit v1.9.3