From 10c0b59bab3107745d461dc9a3d2a654ed8208c3 Mon Sep 17 00:00:00 2001
From: cai <cai@nbcai.cc>
Date: Tue, 11 Aug 2026 21:36:12 +0800
Subject: [PATCH] test(helper): bind probe order to production spawn gate
---
src/main.rs | 84 +++++++++++++++++++++++++++++++++---------
1 files changed, 66 insertions(+), 18 deletions(-)
diff --git a/src/main.rs b/src/main.rs
index 43d1737..056ed15 100644
--- a/src/main.rs
+++ b/src/main.rs
@@ -177,6 +177,21 @@
!pending_probe || probe_result == Some(true)
}
+fn start_observer_after_controlled_fixture_probe<F>(
+ pending_probe: bool,
+ probe_result: Option<bool>,
+ spawn: F,
+) -> bool
+where
+ F: FnOnce(),
+{
+ if !controlled_fixture_observer_gate(pending_probe, probe_result) {
+ return false;
+ }
+ spawn();
+ true
+}
+
#[tokio::main(flavor = "multi_thread")]
async fn main() -> Result<()> {
init_tracing();
@@ -930,7 +945,28 @@
user_participant_identity.as_deref(),
)
.await;
- if !controlled_fixture_observer_gate(pending_sequence.is_some(), probe_result) {
+ let observer_started = start_observer_after_controlled_fixture_probe(
+ pending_sequence.is_some(),
+ probe_result,
+ || {
+ spawn_user_audio_frame_observer(
+ track,
+ call_id.clone(),
+ trace_id.clone(),
+ participant_alias,
+ track_sid_alias,
+ simple_vad_enabled,
+ simple_vad_config.clone(),
+ vad_enabled_gate.clone(),
+ turn_bridge_config.clone(),
+ http.clone(),
+ sink.clone(),
+ user_participant_identity.clone(),
+ participant,
+ );
+ },
+ );
+ if !observer_started {
warn!(
call_id = %call_id,
trace_id = %trace_id,
@@ -938,21 +974,6 @@
);
continue;
}
- spawn_user_audio_frame_observer(
- track,
- call_id.clone(),
- trace_id.clone(),
- participant_alias,
- track_sid_alias,
- simple_vad_enabled,
- simple_vad_config.clone(),
- vad_enabled_gate.clone(),
- turn_bridge_config.clone(),
- http.clone(),
- sink.clone(),
- user_participant_identity.clone(),
- participant,
- );
current_user_participant = Some(participant_for_probe);
}
RoomEvent::DataReceived {
@@ -4391,10 +4412,18 @@
)
.is_ok()
});
- if controlled_fixture_observer_gate(pending, probe_result) {
+ let mut ack_observed = false;
+ let mut observer_started = false;
+ start_observer_after_controlled_fixture_probe(pending, probe_result, || {
if pending && probe_result == Some(true) {
- effects.push("ack_observed");
+ ack_observed = true;
}
+ observer_started = true;
+ });
+ if ack_observed {
+ effects.push("ack_observed");
+ }
+ if observer_started {
effects.push("observer_started");
}
pending_sequence = None;
@@ -5113,4 +5142,23 @@
]);
assert!(effects.is_empty());
}
+
+ #[test]
+ fn production_audio_branch_orders_probe_before_spawn_callsite() {
+ let source = include_str!("main.rs");
+ let branch = source
+ .find("RoomEvent::TrackSubscribed {\n track: RemoteTrack::Audio")
+ .expect("audio TrackSubscribed production branch");
+ let branch_source = &source[branch..];
+ let probe = branch_source
+ .find("let probe_result = process_controlled_fixture_probe")
+ .expect("probe must be processed in audio branch");
+ let spawn = branch_source
+ .find("start_observer_after_controlled_fixture_probe")
+ .expect("spawn must use shared order entry");
+ assert!(
+ probe < spawn,
+ "probe must precede shared observer spawn entry"
+ );
+ }
}
--
Gitblit v1.9.3