Skip to main content

vox_rtc_server/
session.rs

1use crate::error::{Result, VoxRtcError};
2use crate::socket::RawSocketChannel;
3use crate::types::*;
4use serde_json::Value;
5use std::ops::ControlFlow;
6use std::sync::{Arc, Mutex};
7use tokio::sync::broadcast::error::RecvError;
8use uuid::Uuid;
9use tokio::task::JoinHandle;
10use tokio::time::{Duration, timeout};
11
12#[derive(Clone)]
13pub struct VoxRtcControlSession {
14    channel: RawSocketChannel,
15    session_id: String,
16    channel_name: String,
17    join_timeout: Duration,
18    response_generation: Arc<Mutex<ResponseGeneration>>,
19}
20
21#[derive(Default)]
22struct ResponseGeneration {
23    counter: u64,
24    id: Option<String>,
25}
26
27pub struct Listener {
28    handle: JoinHandle<()>,
29}
30
31impl Drop for Listener {
32    fn drop(&mut self) {
33        self.handle.abort();
34    }
35}
36
37impl VoxRtcControlSession {
38    pub(crate) fn new(
39        channel: RawSocketChannel,
40        session_id: String,
41        join_timeout: Duration,
42    ) -> Self {
43        let channel_name = format!("/rtc/{session_id}");
44        Self {
45            channel,
46            session_id,
47            channel_name,
48            join_timeout,
49            response_generation: Arc::new(Mutex::new(ResponseGeneration::default())),
50        }
51    }
52
53    pub fn session_id(&self) -> &str {
54        &self.session_id
55    }
56
57    pub fn channel_name(&self) -> &str {
58        &self.channel_name
59    }
60
61    pub async fn join(&self) -> Result<()> {
62        let mut states = self.channel.subscribe_state();
63        self.channel.join().await?;
64        let channel_name = self.channel.name().to_owned();
65        let channel = self.channel.clone();
66        timeout(self.join_timeout, async move {
67            loop {
68                let state = *states.borrow_and_update();
69                match state {
70                    ChannelState::Joined => return Ok(()),
71                    ChannelState::Closed | ChannelState::Declined => {
72                        let reason = join_decline_reason(channel.decline_reason().await);
73                        return Err(VoxRtcError::JoinFailed {
74                            channel: channel_name,
75                            state: format!("{state:?}"),
76                            reason,
77                        });
78                    }
79                    _ => {}
80                }
81                if states.changed().await.is_err() {
82                    return Err(VoxRtcError::Disconnected);
83                }
84            }
85        })
86        .await
87        .map_err(|_| VoxRtcError::JoinTimeout(self.channel_name.clone()))?
88    }
89
90    pub async fn close(&self) -> Result<()> {
91        self.channel.leave().await
92    }
93
94    pub fn on_event<F>(&self, handler: F) -> Listener
95    where
96        F: Fn(WireEvent) + Send + Sync + 'static,
97    {
98        let mut messages = self.channel.subscribe_messages();
99        let session_id = self.session_id.clone();
100        let channel_name = self.channel_name.clone();
101        Listener {
102            handle: tokio::spawn(async move {
103                loop {
104                    match next_message(messages.recv().await) {
105                        ControlFlow::Break(()) => break,
106                        ControlFlow::Continue(None) => continue,
107                        ControlFlow::Continue(Some((event, payload))) => handler(WireEvent {
108                            r#type: event,
109                            data: payload,
110                            session_id: session_id.clone(),
111                            channel_name: channel_name.clone(),
112                        }),
113                    }
114                }
115            }),
116        }
117    }
118
119    pub fn on<F>(&self, event_name: impl Into<String>, handler: F) -> Listener
120    where
121        F: Fn(EventData) + Send + Sync + 'static,
122    {
123        let event_name = event_name.into();
124        let mut messages = self.channel.subscribe_messages();
125        Listener {
126            handle: tokio::spawn(async move {
127                loop {
128                    match next_message(messages.recv().await) {
129                        ControlFlow::Break(()) => break,
130                        ControlFlow::Continue(None) => continue,
131                        ControlFlow::Continue(Some((event, payload))) => {
132                            if event == event_name {
133                                handler(payload);
134                            }
135                        }
136                    }
137                }
138            }),
139        }
140    }
141
142    pub fn on_session_attached<F>(&self, handler: F) -> Listener
143    where
144        F: Fn(SessionAttachedEvent) + Send + Sync + 'static,
145    {
146        let session_id = self.session_id.clone();
147        let channel_name = self.channel_name.clone();
148        self.on(EVENT_RTC_SESSION_ATTACHED, move |payload| {
149            handler(SessionAttachedEvent {
150                session_id: session_id.clone(),
151                channel_name: channel_name.clone(),
152                data: payload,
153            })
154        })
155    }
156
157    pub fn on_session_created<F>(&self, handler: F) -> Listener
158    where
159        F: Fn(SessionCreatedEvent) + Send + Sync + 'static,
160    {
161        let session_id = self.session_id.clone();
162        let channel_name = self.channel_name.clone();
163        self.on(EVENT_SESSION_CREATED, move |payload| {
164            let session = payload.get("session").and_then(Value::as_object).cloned();
165            handler(SessionCreatedEvent {
166                session_id: session_id.clone(),
167                channel_name: channel_name.clone(),
168                data: payload,
169                session,
170            });
171        })
172    }
173
174    pub fn on_transcript<F>(&self, handler: F) -> Listener
175    where
176        F: Fn(TranscriptEvent) + Send + Sync + 'static,
177    {
178        let session_id = self.session_id.clone();
179        let channel_name = self.channel_name.clone();
180        self.on(EVENT_TRANSCRIPT_COMPLETED, move |payload| {
181            handler(TranscriptEvent {
182                session_id: session_id.clone(),
183                channel_name: channel_name.clone(),
184                transcript: required_string(&payload, "transcript", ""),
185                language: optional_string(&payload, "language"),
186                start_ms: optional_number(&payload, "start_ms"),
187                end_ms: optional_number(&payload, "end_ms"),
188                eou_probability: optional_number(&payload, "eou_probability"),
189                topics: optional_string_vec(&payload, "topics"),
190                entities: transcript_entities(&payload),
191                words: transcript_words(&payload),
192                speech_context: payload
193                    .get("speech_context")
194                    .cloned()
195                    .and_then(|value| serde_json::from_value::<SpeechContext>(value).ok())
196                    .filter(SpeechContext::is_valid),
197                data: payload,
198            });
199        })
200    }
201
202    pub fn on_turn_state_changed<F>(&self, handler: F) -> Listener
203    where
204        F: Fn(TurnStateEvent) + Send + Sync + 'static,
205    {
206        let session_id = self.session_id.clone();
207        let channel_name = self.channel_name.clone();
208        self.on(EVENT_TURN_STATE_CHANGED, move |payload| {
209            handler(TurnStateEvent {
210                session_id: session_id.clone(),
211                channel_name: channel_name.clone(),
212                state: required_string(&payload, "state", "unknown"),
213                previous_state: optional_string(&payload, "previous_state"),
214                data: payload,
215            });
216        })
217    }
218
219    pub fn on_speech_started<F>(&self, handler: F) -> Listener
220    where
221        F: Fn(SpeechStartedEvent) + Send + Sync + 'static,
222    {
223        let session_id = self.session_id.clone();
224        let channel_name = self.channel_name.clone();
225        self.on(EVENT_SPEECH_STARTED, move |payload| {
226            handler(SpeechStartedEvent {
227                session_id: session_id.clone(),
228                channel_name: channel_name.clone(),
229                timestamp_ms: optional_number(&payload, "timestamp_ms"),
230                data: payload,
231            });
232        })
233    }
234
235    pub fn on_speech_stopped<F>(&self, handler: F) -> Listener
236    where
237        F: Fn(SpeechStoppedEvent) + Send + Sync + 'static,
238    {
239        let session_id = self.session_id.clone();
240        let channel_name = self.channel_name.clone();
241        self.on(EVENT_SPEECH_STOPPED, move |payload| {
242            handler(SpeechStoppedEvent {
243                session_id: session_id.clone(),
244                channel_name: channel_name.clone(),
245                timestamp_ms: optional_number(&payload, "timestamp_ms"),
246                data: payload,
247            });
248        })
249    }
250
251    pub fn on_transcript_delta<F>(&self, handler: F) -> Listener
252    where
253        F: Fn(TranscriptDeltaEvent) + Send + Sync + 'static,
254    {
255        let session_id = self.session_id.clone();
256        let channel_name = self.channel_name.clone();
257        self.on(EVENT_TRANSCRIPT_DELTA, move |payload| {
258            handler(TranscriptDeltaEvent {
259                session_id: session_id.clone(),
260                channel_name: channel_name.clone(),
261                delta: required_string(&payload, "delta", ""),
262                start_ms: optional_number(&payload, "start_ms"),
263                end_ms: optional_number(&payload, "end_ms"),
264                data: payload,
265            });
266        })
267    }
268
269    pub fn on_turn_eou_predicted<F>(&self, handler: F) -> Listener
270    where
271        F: Fn(TurnEouPredictedEvent) + Send + Sync + 'static,
272    {
273        let session_id = self.session_id.clone();
274        let channel_name = self.channel_name.clone();
275        self.on(EVENT_TURN_EOU_PREDICTED, move |payload| {
276            handler(TurnEouPredictedEvent {
277                session_id: session_id.clone(),
278                channel_name: channel_name.clone(),
279                probability: optional_number(&payload, "probability"),
280                threshold: optional_number(&payload, "threshold"),
281                delay_ms: optional_number(&payload, "delay_ms"),
282                start_ms: optional_number(&payload, "start_ms"),
283                end_ms: optional_number(&payload, "end_ms"),
284                decision: optional_string(&payload, "decision"),
285                action: optional_string(&payload, "action"),
286                turn_detector: optional_string(&payload, "turn_detector"),
287                data: payload,
288            });
289        })
290    }
291
292    pub fn on_response_created<F>(&self, handler: F) -> Listener
293    where
294        F: Fn(ResponseEvent) + Send + Sync + 'static,
295    {
296        self.on_response_event(EVENT_RESPONSE_CREATED, handler)
297    }
298
299    pub fn on_response_committed<F>(&self, handler: F) -> Listener
300    where
301        F: Fn(ResponseEvent) + Send + Sync + 'static,
302    {
303        self.on_response_event(EVENT_RESPONSE_COMMITTED, handler)
304    }
305
306    pub fn on_response_done<F>(&self, handler: F) -> Listener
307    where
308        F: Fn(ResponseEvent) + Send + Sync + 'static,
309    {
310        self.on_response_event(EVENT_RESPONSE_DONE, handler)
311    }
312
313    pub fn on_response_cancelled<F>(&self, handler: F) -> Listener
314    where
315        F: Fn(ResponseEvent) + Send + Sync + 'static,
316    {
317        self.on_response_event(EVENT_RESPONSE_CANCELLED, handler)
318    }
319
320    pub fn on_response_audio_clear<F>(&self, handler: F) -> Listener
321    where
322        F: Fn(ResponseEvent) + Send + Sync + 'static,
323    {
324        self.on_response_event(EVENT_RESPONSE_AUDIO_CLEAR, handler)
325    }
326
327    fn on_response_event<F>(&self, event_name: &'static str, handler: F) -> Listener
328    where
329        F: Fn(ResponseEvent) + Send + Sync + 'static,
330    {
331        let session_id = self.session_id.clone();
332        let channel_name = self.channel_name.clone();
333        self.on(event_name, move |payload| {
334            handler(response_event(payload, &session_id, &channel_name));
335        })
336    }
337
338    pub fn on_interruption_detected<F>(&self, handler: F) -> Listener
339    where
340        F: Fn(InterruptionEvent) + Send + Sync + 'static,
341    {
342        self.on_interruption_event(EVENT_INTERRUPTION_DETECTED, handler)
343    }
344
345    pub fn on_interruption_false_positive<F>(&self, handler: F) -> Listener
346    where
347        F: Fn(InterruptionEvent) + Send + Sync + 'static,
348    {
349        self.on_interruption_event(EVENT_INTERRUPTION_FALSE_POSITIVE, handler)
350    }
351
352    fn on_interruption_event<F>(&self, event_name: &'static str, handler: F) -> Listener
353    where
354        F: Fn(InterruptionEvent) + Send + Sync + 'static,
355    {
356        let session_id = self.session_id.clone();
357        let channel_name = self.channel_name.clone();
358        self.on(event_name, move |payload| {
359            handler(InterruptionEvent {
360                response: response_event(payload.clone(), &session_id, &channel_name),
361                vad_active_ms: optional_number(&payload, "vad_active_ms"),
362                partial_transcript: optional_string(&payload, "partial_transcript"),
363                reason: optional_nonempty_string(&payload, "reason"),
364            });
365        })
366    }
367
368    pub fn on_browser_event<F>(&self, handler: F) -> Listener
369    where
370        F: Fn(BrowserEvent) + Send + Sync + 'static,
371    {
372        let session_id = self.session_id.clone();
373        let channel_name = self.channel_name.clone();
374        self.on(EVENT_BROWSER_EVENT, move |payload| {
375            handler(BrowserEvent {
376                session_id: session_id.clone(),
377                channel_name: channel_name.clone(),
378                event: required_string(&payload, "event", ""),
379                payload: payload.get("payload").cloned().unwrap_or(Value::Null),
380                data: payload,
381            });
382        })
383    }
384
385    pub fn on_close<F>(&self, handler: F) -> Listener
386    where
387        F: Fn(CloseEvent) + Send + Sync + 'static,
388    {
389        let session_id = self.session_id.clone();
390        let channel_name = self.channel_name.clone();
391        self.on(EVENT_RTC_CLIENT_DISCONNECTED, move |payload| {
392            handler(CloseEvent {
393                session_id: session_id.clone(),
394                channel_name: channel_name.clone(),
395                reason: required_string(&payload, "reason", "unknown"),
396                connection_state: optional_string(&payload, "connection_state"),
397                ice_connection_state: optional_string(&payload, "ice_connection_state"),
398                data_channel_state: optional_string(&payload, "data_channel_state"),
399                data: payload,
400            });
401        })
402    }
403
404    pub fn on_error<F>(&self, handler: F) -> Listener
405    where
406        F: Fn(ErrorEvent) + Send + Sync + 'static,
407    {
408        let session_id = self.session_id.clone();
409        let channel_name = self.channel_name.clone();
410        self.on(EVENT_ERROR, move |payload| {
411            handler(ErrorEvent {
412                session_id: session_id.clone(),
413                channel_name: channel_name.clone(),
414                message: optional_string(&payload, "message"),
415                code: optional_nonempty_string(&payload, "code"),
416                recoverable: recoverable_flag(&payload),
417                generation_id: optional_nonempty_string(&payload, "generation_id"),
418                data: payload,
419            });
420        })
421    }
422
423    pub fn on_signaling_error<F>(&self, handler: F) -> Listener
424    where
425        F: Fn(SignalingErrorEvent) + Send + Sync + 'static,
426    {
427        let session_id = self.session_id.clone();
428        let channel_name = self.channel_name.clone();
429        self.on(EVENT_RTC_SIGNALING_ERROR, move |payload| {
430            handler(SignalingErrorEvent {
431                session_id: session_id.clone(),
432                channel_name: channel_name.clone(),
433                message: optional_string(&payload, "message"),
434                generation: optional_i64(&payload, "generation"),
435                data: payload,
436            });
437        })
438    }
439
440    pub async fn send_control(&self, event: &str, payload: EventData) -> Result<()> {
441        self.channel.send_message(event, payload).await
442    }
443
444    pub async fn configure(&self, config: SessionConfig) -> Result<()> {
445        let mut payload = EventData::new();
446        payload.insert(
447            "session".to_owned(),
448            Value::Object(session_config_payload(config)),
449        );
450        self.send_control("session.update", payload).await
451    }
452
453    pub async fn start_response(&self, options: Option<ResponseOptions>) -> Result<()> {
454        let (_, payload) = self.start_payload(options);
455        self.send_control("response.start", payload).await
456    }
457
458    pub async fn start_response_and_wait(
459        &self,
460        options: Option<ResponseOptions>,
461        wait_timeout: Duration,
462    ) -> Result<StartAck> {
463        let (generation_id, payload) = self.start_payload(options);
464        let mut messages = self.channel.subscribe_messages();
465        self.send_control("response.start", payload).await?;
466        timeout(wait_timeout, async move {
467            loop {
468                match next_message(messages.recv().await) {
469                    ControlFlow::Break(()) => return Err(VoxRtcError::ChannelClosed),
470                    ControlFlow::Continue(None) => continue,
471                    ControlFlow::Continue(Some((event, data))) => {
472                        if optional_nonempty_string(&data, "generation_id").as_deref()
473                            != Some(generation_id.as_str())
474                        {
475                            continue;
476                        }
477                        if event == EVENT_RESPONSE_CREATED {
478                            return Ok(StartAck {
479                                accepted: true,
480                                generation_id: generation_id.clone(),
481                                response_id: optional_string(&data, "response_id"),
482                                output: response_output(&data),
483                                error_code: None,
484                                error_message: None,
485                                recoverable: true,
486                            });
487                        }
488                        if event == EVENT_ERROR {
489                            return Ok(StartAck {
490                                accepted: false,
491                                generation_id: generation_id.clone(),
492                                response_id: optional_string(&data, "response_id"),
493                                output: None,
494                                error_code: optional_nonempty_string(&data, "code"),
495                                error_message: optional_string(&data, "message"),
496                                recoverable: recoverable_flag(&data),
497                            });
498                        }
499                    }
500                }
501            }
502        })
503        .await
504        .map_err(|_| VoxRtcError::Timeout("response.start acknowledgement"))?
505    }
506
507    pub async fn append_response_text(
508        &self,
509        delta: impl Into<String>,
510        options: Option<ResponseOptions>,
511    ) -> Result<()> {
512        let explicit = explicit_generation(&options);
513        let mut payload = response_options_payload(options);
514        payload.insert("delta".to_owned(), Value::String(delta.into()));
515        self.thread_generation(&mut payload, explicit);
516        self.send_control("response.delta", payload).await
517    }
518
519    pub async fn commit_response(&self, options: Option<ResponseOptions>) -> Result<()> {
520        let explicit = explicit_generation(&options);
521        let mut payload = EventData::new();
522        self.thread_generation(&mut payload, explicit);
523        self.send_control("response.commit", payload).await
524    }
525
526    pub async fn cancel_response(&self, options: Option<ResponseOptions>) -> Result<()> {
527        let explicit = explicit_generation(&options);
528        let mut payload = EventData::new();
529        self.thread_generation(&mut payload, explicit);
530        self.clear_response_generation();
531        self.send_control("response.cancel", payload).await
532    }
533
534    pub async fn replace_response_text(
535        &self,
536        text: impl Into<String>,
537        options: Option<ResponseOptions>,
538    ) -> Result<()> {
539        self.clear_response_generation();
540        let explicit = explicit_generation(&options);
541        let mut payload = response_options_payload(options);
542        payload.insert("text".to_owned(), Value::String(text.into()));
543        if let Some(generation_id) = explicit {
544            payload.insert("generation_id".to_owned(), Value::String(generation_id));
545        }
546        self.send_control("response.replace_text", payload).await
547    }
548
549    pub async fn send_text_response(
550        &self,
551        text: impl Into<String>,
552        options: Option<ResponseOptions>,
553        cancel_first: bool,
554    ) -> Result<()> {
555        let text = text.into();
556        if cancel_first {
557            return self.replace_response_text(text, options).await;
558        }
559        self.start_response(options.clone()).await?;
560        self.append_response_text(text, options.clone()).await?;
561        self.commit_response(options).await
562    }
563
564    pub async fn send_client_event(&self, envelope: ClientEventEnvelope) -> Result<()> {
565        let mut payload = EventData::new();
566        payload.insert("event".to_owned(), Value::String(envelope.event));
567        payload.insert("payload".to_owned(), envelope.payload);
568        self.send_control(EVENT_CLIENT_EVENT, payload).await
569    }
570
571    fn start_payload(&self, options: Option<ResponseOptions>) -> (String, EventData) {
572        let explicit = explicit_generation(&options);
573        let mut payload = response_options_payload(options);
574        let generation_id = match explicit {
575            Some(id) => self.set_response_generation(id),
576            None => self.next_response_generation(),
577        };
578        payload.insert(
579            "generation_id".to_owned(),
580            Value::String(generation_id.clone()),
581        );
582        (generation_id, payload)
583    }
584
585    fn next_response_generation(&self) -> String {
586        let mut state = self
587            .response_generation
588            .lock()
589            .expect("response generation mutex poisoned");
590        state.counter += 1;
591        let generation_id = format!("generation_{}_{}", state.counter, Uuid::new_v4());
592        state.id = Some(generation_id.clone());
593        generation_id
594    }
595
596    fn set_response_generation(&self, generation_id: String) -> String {
597        let mut state = self
598            .response_generation
599            .lock()
600            .expect("response generation mutex poisoned");
601        state.counter += 1;
602        state.id = Some(generation_id.clone());
603        generation_id
604    }
605
606    fn thread_generation(&self, payload: &mut EventData, explicit: Option<String>) {
607        match explicit {
608            Some(generation_id) => {
609                payload.insert("generation_id".to_owned(), Value::String(generation_id));
610            }
611            None => self.add_response_generation(payload),
612        }
613    }
614
615    fn add_response_generation(&self, payload: &mut EventData) {
616        let state = self
617            .response_generation
618            .lock()
619            .expect("response generation mutex poisoned");
620        if let Some(generation_id) = &state.id {
621            payload.insert(
622                "generation_id".to_owned(),
623                Value::String(generation_id.clone()),
624            );
625        }
626    }
627
628    fn clear_response_generation(&self) {
629        self.response_generation
630            .lock()
631            .expect("response generation mutex poisoned")
632            .id = None;
633    }
634}
635
636fn next_message(
637    result: std::result::Result<(String, EventData), RecvError>,
638) -> ControlFlow<(), Option<(String, EventData)>> {
639    match result {
640        Ok(message) => ControlFlow::Continue(Some(message)),
641        Err(RecvError::Lagged(_)) => ControlFlow::Continue(None),
642        Err(RecvError::Closed) => ControlFlow::Break(()),
643    }
644}
645
646fn join_decline_reason(reason: Option<EventData>) -> Option<String> {
647    let reason = reason?;
648    for key in ["message", "reason", "error"] {
649        if let Some(value) = reason.get(key).and_then(Value::as_str)
650            && !value.is_empty()
651        {
652            return Some(value.to_owned());
653        }
654    }
655    if reason.is_empty() {
656        None
657    } else {
658        Some(Value::Object(reason).to_string())
659    }
660}
661
662fn insert_opt(session: &mut EventData, key: &str, value: Option<String>) {
663    if let Some(value) = value {
664        session.insert(key.to_owned(), Value::String(value));
665    }
666}
667
668fn session_config_payload(config: SessionConfig) -> EventData {
669    let mut session = config.extra;
670    insert_opt(&mut session, "stt_model", config.stt_model);
671    insert_opt(&mut session, "tts_model", config.tts_model);
672    insert_opt(&mut session, "voice", config.voice);
673    insert_opt(&mut session, "turn_profile", config.turn_profile);
674    insert_opt(&mut session, "vad_backend", config.vad_backend);
675    insert_opt(&mut session, "turn_detector", config.turn_detector);
676    if let Some(enabled) = config.speech_context {
677        session.insert("speech_context".to_owned(), Value::Bool(enabled));
678    }
679    session
680}
681
682fn response_options_payload(options: Option<ResponseOptions>) -> EventData {
683    let mut payload = EventData::new();
684    if let Some(options) = options {
685        if let Some(allow) = options.allow_interruptions {
686            payload.insert("allow_interruptions".to_owned(), Value::Bool(allow));
687        }
688        if let Some(output) = options.output {
689            payload.insert(
690                "output".to_owned(),
691                Value::Object(response_output_options_payload(output)),
692            );
693        }
694    }
695    payload
696}
697
698fn response_output_options_payload(output: ResponseOutputOptions) -> EventData {
699    let mut payload = EventData::new();
700    insert_opt(&mut payload, "model", output.model);
701    insert_opt(&mut payload, "voice", output.voice);
702    insert_opt(&mut payload, "language", output.language);
703    if let Some(speed) = output.speed
704        && let Some(number) = serde_json::Number::from_f64(speed)
705    {
706        payload.insert("speed".to_owned(), Value::Number(number));
707    }
708    if let Some(params) = output.params {
709        payload.insert("params".to_owned(), Value::Object(params));
710    }
711    payload
712}
713
714fn explicit_generation(options: &Option<ResponseOptions>) -> Option<String> {
715    options
716        .as_ref()
717        .and_then(|options| options.generation_id.clone())
718        .filter(|id| !id.is_empty())
719}
720
721fn response_event(payload: EventData, session_id: &str, channel_name: &str) -> ResponseEvent {
722    let output = response_output(&payload);
723    ResponseEvent {
724        session_id: session_id.to_owned(),
725        channel_name: channel_name.to_owned(),
726        response_id: optional_string(&payload, "response_id"),
727        generation_id: optional_nonempty_string(&payload, "generation_id"),
728        output,
729        data: payload,
730    }
731}
732
733fn response_output(payload: &EventData) -> Option<ResponseOutput> {
734    let output = payload.get("output")?.as_object()?;
735    let model = optional_nonempty_string(output, "model")?;
736    let language = optional_nonempty_string(output, "language")?;
737    let speed = optional_number(output, "speed")?;
738    let params = output.get("params")?.as_object()?.clone();
739    Some(ResponseOutput {
740        model,
741        voice: optional_nonempty_string(output, "voice"),
742        language,
743        speed,
744        params,
745    })
746}
747
748#[cfg(test)]
749mod tests {
750    use super::*;
751    use crate::socket::test_channel;
752    use serde_json::json;
753    use tokio::sync::broadcast;
754    use tokio::sync::mpsc;
755
756    async fn session() -> (VoxRtcControlSession, broadcast::Sender<(String, EventData)>) {
757        let (channel, sender) = test_channel().await;
758        let session =
759            VoxRtcControlSession::new(channel, "sess-1".to_owned(), Duration::from_secs(1));
760        (session, sender)
761    }
762
763    fn payload(value: Value) -> EventData {
764        value.as_object().cloned().expect("object payload")
765    }
766
767    #[test]
768    fn join_decline_reason_prefers_structured_message_fields() {
769        assert_eq!(
770            join_decline_reason(Some(payload(json!({ "message": "expired" })))),
771            Some("expired".to_owned())
772        );
773        assert_eq!(
774            join_decline_reason(Some(payload(json!({ "reason": "missing" })))),
775            Some("missing".to_owned())
776        );
777        assert_eq!(
778            join_decline_reason(Some(payload(json!({ "channel": "/rtc/abc" })))),
779            Some(r#"{"channel":"/rtc/abc"}"#.to_owned())
780        );
781        assert_eq!(join_decline_reason(Some(EventData::new())), None);
782        assert_eq!(join_decline_reason(None), None);
783    }
784
785    async fn recv<T>(rx: &mut mpsc::UnboundedReceiver<T>) -> T {
786        timeout(Duration::from_secs(1), rx.recv())
787            .await
788            .expect("handler fired within timeout")
789            .expect("handler produced an event")
790    }
791
792    #[test]
793    fn next_message_classifies_lag_close_and_ok() {
794        assert!(matches!(
795            next_message(Ok(("e".to_owned(), EventData::new()))),
796            ControlFlow::Continue(Some(_))
797        ));
798        assert!(matches!(
799            next_message(Err(RecvError::Lagged(7))),
800            ControlFlow::Continue(None)
801        ));
802        assert!(matches!(
803            next_message(Err(RecvError::Closed)),
804            ControlFlow::Break(())
805        ));
806    }
807
808    #[test]
809    fn session_config_serializes_explicit_false_speech_context() {
810        let payload = session_config_payload(SessionConfig {
811            speech_context: Some(false),
812            ..Default::default()
813        });
814        assert_eq!(payload.get("speech_context"), Some(&Value::Bool(false)));
815    }
816
817    #[tokio::test]
818    async fn response_commands_share_one_generation_id() {
819        let (session, _) = session().await;
820        let generation_id = session.next_response_generation();
821        let mut delta = payload(json!({ "delta": "hello" }));
822        session.add_response_generation(&mut delta);
823        let mut commit = EventData::new();
824        session.add_response_generation(&mut commit);
825
826        assert_eq!(
827            delta.get("generation_id"),
828            Some(&Value::String(generation_id.clone()))
829        );
830        assert_eq!(
831            commit.get("generation_id"),
832            Some(&Value::String(generation_id))
833        );
834    }
835
836    #[tokio::test]
837    async fn on_error_parses_typed_fields() {
838        let (session, sender) = session().await;
839        let (tx, mut rx) = mpsc::unbounded_channel();
840        let _listener = session.on_error(move |event| {
841            tx.send(event).unwrap();
842        });
843        sender
844            .send((
845                EVENT_ERROR.to_owned(),
846                payload(json!({
847                    "message": "cannot start now",
848                    "code": ERROR_CODE_SESSION_FAILED,
849                    "recoverable": false,
850                    "generation_id": "gen-9"
851                })),
852            ))
853            .unwrap();
854        let event = recv(&mut rx).await;
855        assert_eq!(event.message.as_deref(), Some("cannot start now"));
856        assert_eq!(event.code.as_deref(), Some(ERROR_CODE_SESSION_FAILED));
857        assert!(!event.recoverable);
858        assert_eq!(event.generation_id.as_deref(), Some("gen-9"));
859    }
860
861    #[tokio::test]
862    async fn on_error_defaults_missing_recoverable_to_true() {
863        let (session, sender) = session().await;
864        let (tx, mut rx) = mpsc::unbounded_channel();
865        let _listener = session.on_error(move |event| {
866            tx.send(event).unwrap();
867        });
868        sender
869            .send((
870                EVENT_ERROR.to_owned(),
871                payload(json!({ "message": "legacy server", "code": "" })),
872            ))
873            .unwrap();
874        let event = recv(&mut rx).await;
875        assert!(event.recoverable);
876        assert_eq!(event.code, None);
877        assert_eq!(event.generation_id, None);
878    }
879
880    #[tokio::test]
881    async fn start_payload_uses_explicit_generation_id() {
882        let (session, _) = session().await;
883        let options = ResponseOptions {
884            allow_interruptions: Some(false),
885            generation_id: Some("gen-7".to_owned()),
886            ..Default::default()
887        };
888        let (generation_id, start) = session.start_payload(Some(options));
889        assert_eq!(generation_id, "gen-7");
890        assert_eq!(
891            start.get("generation_id"),
892            Some(&Value::String("gen-7".to_owned()))
893        );
894        assert_eq!(start.get("allow_interruptions"), Some(&Value::Bool(false)));
895
896        let mut commit = EventData::new();
897        session.thread_generation(&mut commit, None);
898        assert_eq!(
899            commit.get("generation_id"),
900            Some(&Value::String("gen-7".to_owned()))
901        );
902    }
903
904    #[tokio::test]
905    async fn start_payload_generates_generation_id_when_absent() {
906        let (session, _) = session().await;
907        let (generation_id, start) = session.start_payload(None);
908        assert!(generation_id.starts_with("generation_1_"));
909        assert!(generation_id.len() > "generation_1_".len());
910        assert_eq!(
911            start.get("generation_id"),
912            Some(&Value::String(generation_id))
913        );
914    }
915
916    #[tokio::test]
917    async fn start_payload_serializes_response_output() {
918        let (session, _) = session().await;
919        let options = ResponseOptions {
920            generation_id: Some("gen-output".to_owned()),
921            output: Some(ResponseOutputOptions {
922                model: Some("qwen3-tts:0.6b-clone".to_owned()),
923                voice: Some("samantha".to_owned()),
924                language: Some("fr".to_owned()),
925                speed: Some(0.9),
926                params: Some(payload(json!({ "temperature": 0.7 }))),
927            }),
928            ..Default::default()
929        };
930
931        let (_, start) = session.start_payload(Some(options));
932
933        assert_eq!(
934            start.get("output"),
935            Some(&json!({
936                "model": "qwen3-tts:0.6b-clone",
937                "voice": "samantha",
938                "language": "fr",
939                "speed": 0.9,
940                "params": { "temperature": 0.7 }
941            }))
942        );
943    }
944
945    #[tokio::test]
946    async fn explicit_generation_id_overrides_tracked_one() {
947        let (session, _) = session().await;
948        let tracked = session.next_response_generation();
949        let mut delta = payload(json!({ "delta": "hi" }));
950        session.thread_generation(&mut delta, Some("gen-42".to_owned()));
951        assert_eq!(
952            delta.get("generation_id"),
953            Some(&Value::String("gen-42".to_owned()))
954        );
955        assert_ne!(tracked, "gen-42");
956    }
957
958    #[tokio::test]
959    async fn response_events_expose_generation_id() {
960        let (session, sender) = session().await;
961        let (tx, mut rx) = mpsc::unbounded_channel();
962        let _listener = session.on_response_created(move |event| {
963            tx.send(event).unwrap();
964        });
965        sender
966            .send((
967                EVENT_RESPONSE_CREATED.to_owned(),
968                payload(json!({ "response_id": "resp-1", "generation_id": "gen-1" })),
969            ))
970            .unwrap();
971        let event = recv(&mut rx).await;
972        assert_eq!(event.response_id.as_deref(), Some("resp-1"));
973        assert_eq!(event.generation_id.as_deref(), Some("gen-1"));
974    }
975
976    #[tokio::test]
977    async fn audio_clear_and_interruption_expose_generation_id() {
978        let (session, sender) = session().await;
979        let (clear_tx, mut clear_rx) = mpsc::unbounded_channel();
980        let _clear = session.on_response_audio_clear(move |event| {
981            clear_tx.send(event).unwrap();
982        });
983        let (int_tx, mut int_rx) = mpsc::unbounded_channel();
984        let _interruption = session.on_interruption_detected(move |event| {
985            int_tx.send(event).unwrap();
986        });
987        sender
988            .send((
989                EVENT_RESPONSE_AUDIO_CLEAR.to_owned(),
990                payload(json!({ "response_id": "resp-2", "generation_id": "gen-2" })),
991            ))
992            .unwrap();
993        sender
994            .send((
995                EVENT_INTERRUPTION_DETECTED.to_owned(),
996                payload(json!({
997                    "response_id": "resp-2",
998                    "generation_id": "gen-2",
999                    "vad_active_ms": 250
1000                })),
1001            ))
1002            .unwrap();
1003        let clear = recv(&mut clear_rx).await;
1004        assert_eq!(clear.generation_id.as_deref(), Some("gen-2"));
1005        let interruption = recv(&mut int_rx).await;
1006        assert_eq!(interruption.response.generation_id.as_deref(), Some("gen-2"));
1007        assert_eq!(interruption.vad_active_ms, Some(250.0));
1008    }
1009
1010    #[tokio::test]
1011    async fn on_signaling_error_parses_message_and_generation() {
1012        let (session, sender) = session().await;
1013        let (tx, mut rx) = mpsc::unbounded_channel();
1014        let _listener = session.on_signaling_error(move |event| {
1015            tx.send(event).unwrap();
1016        });
1017        sender
1018            .send((
1019                EVENT_RTC_SIGNALING_ERROR.to_owned(),
1020                payload(json!({
1021                    "message": "setLocalDescription failed",
1022                    "generation": 3
1023                })),
1024            ))
1025            .unwrap();
1026        let event = recv(&mut rx).await;
1027        assert_eq!(event.message.as_deref(), Some("setLocalDescription failed"));
1028        assert_eq!(event.generation, Some(3));
1029    }
1030
1031    #[tokio::test]
1032    async fn on_signaling_error_leaves_generation_none_when_absent() {
1033        let (session, sender) = session().await;
1034        let (tx, mut rx) = mpsc::unbounded_channel();
1035        let _listener = session.on_signaling_error(move |event| {
1036            tx.send(event).unwrap();
1037        });
1038        sender
1039            .send((
1040                EVENT_RTC_SIGNALING_ERROR.to_owned(),
1041                payload(json!({ "message": "RTC signaling failed" })),
1042            ))
1043            .unwrap();
1044        let event = recv(&mut rx).await;
1045        assert_eq!(event.message.as_deref(), Some("RTC signaling failed"));
1046        assert_eq!(event.generation, None);
1047    }
1048
1049    #[tokio::test]
1050    async fn on_transcript_exposes_entities_and_words() {
1051        let (session, sender) = session().await;
1052        let (tx, mut rx) = mpsc::unbounded_channel();
1053        let _listener = session.on_transcript(move |event| {
1054            tx.send(event).unwrap();
1055        });
1056        sender
1057            .send((
1058                EVENT_TRANSCRIPT_COMPLETED.to_owned(),
1059                payload(json!({
1060                    "transcript": "call Ada",
1061                    "entities": [
1062                        { "type": "PRODUCT", "text": "Ada", "start_char": 5, "end_char": 8 }
1063                    ],
1064                    "words": [
1065                        { "word": "call", "start_ms": 0, "end_ms": 300 },
1066                        { "word": "Ada", "start_ms": 300, "end_ms": 600, "confidence": 0.91 }
1067                    ],
1068                    "speech_context": serde_json::from_str::<Value>(include_str!(
1069                        "../../../fixtures/speech-context-v2.json"
1070                    )).unwrap()
1071                })),
1072            ))
1073            .unwrap();
1074        let event = recv(&mut rx).await;
1075        assert_eq!(
1076            event.entities,
1077            vec![TranscriptEntity {
1078                r#type: "PRODUCT".to_owned(),
1079                text: "Ada".to_owned(),
1080                start_char: 5,
1081                end_char: 8,
1082            }]
1083        );
1084        assert_eq!(event.words.len(), 2);
1085        assert_eq!(event.words[0].word, "call");
1086        assert_eq!(event.words[0].start_ms, 0.0);
1087        assert_eq!(event.words[0].confidence, None);
1088        assert_eq!(event.words[1].confidence, Some(0.91));
1089        let context = event.speech_context.expect("speech context");
1090        assert_eq!(context.schema_version, 2);
1091        assert_eq!(context.status, SpeechContextStatus::Complete);
1092        assert_eq!(
1093            context.emotions.as_deref(),
1094            Some(
1095                &[SpeechContextSpan {
1096                    label: "surprised".to_owned(),
1097                    start_ms: 0,
1098                    end_ms: 2500,
1099                }][..]
1100            )
1101        );
1102        let sounds = context.sounds.expect("sound spans");
1103        assert_eq!(sounds.len(), 2);
1104        assert_eq!(sounds[0].span.label, "fireworks");
1105        assert_eq!(sounds[0].score, 0.42);
1106    }
1107
1108    #[tokio::test]
1109    async fn on_transcript_defaults_entities_and_words_to_empty() {
1110        let (session, sender) = session().await;
1111        let (tx, mut rx) = mpsc::unbounded_channel();
1112        let _listener = session.on_transcript(move |event| {
1113            tx.send(event).unwrap();
1114        });
1115        sender
1116            .send((
1117                EVENT_TRANSCRIPT_COMPLETED.to_owned(),
1118                payload(json!({ "transcript": "hello" })),
1119            ))
1120            .unwrap();
1121        let event = recv(&mut rx).await;
1122        assert!(event.entities.is_empty());
1123        assert!(event.words.is_empty());
1124        assert!(event.speech_context.is_none());
1125    }
1126
1127    #[tokio::test]
1128    async fn on_transcript_preserves_text_but_rejects_malformed_speech_context() {
1129        let (session, sender) = session().await;
1130        let (tx, mut rx) = mpsc::unbounded_channel();
1131        let _listener = session.on_transcript(move |event| {
1132            tx.send(event).unwrap();
1133        });
1134        sender
1135            .send((
1136                EVENT_TRANSCRIPT_COMPLETED.to_owned(),
1137                payload(json!({
1138                    "transcript": "still delivered",
1139                    "speech_context": {
1140                        "schema_version": 2,
1141                        "status": "complete",
1142                        "emotions": [],
1143                        "vocal": [],
1144                        "sounds": [
1145                            {
1146                                "label": "fireworks",
1147                                "start_ms": 0,
1148                                "end_ms": 960,
1149                                "score": 1.1
1150                            }
1151                        ]
1152                    }
1153                })),
1154            ))
1155            .unwrap();
1156        let event = recv(&mut rx).await;
1157        assert_eq!(event.transcript, "still delivered");
1158        assert!(event.speech_context.is_none());
1159    }
1160
1161    #[tokio::test]
1162    async fn interruption_events_expose_reason() {
1163        let (session, sender) = session().await;
1164        let (det_tx, mut det_rx) = mpsc::unbounded_channel();
1165        let _detected = session.on_interruption_detected(move |event| {
1166            det_tx.send(event).unwrap();
1167        });
1168        let (fp_tx, mut fp_rx) = mpsc::unbounded_channel();
1169        let _false_positive = session.on_interruption_false_positive(move |event| {
1170            fp_tx.send(event).unwrap();
1171        });
1172        sender
1173            .send((
1174                EVENT_INTERRUPTION_DETECTED.to_owned(),
1175                payload(json!({
1176                    "response_id": "resp-3",
1177                    "generation_id": "gen-3",
1178                    "reason": "speech_overlap"
1179                })),
1180            ))
1181            .unwrap();
1182        sender
1183            .send((
1184                EVENT_INTERRUPTION_FALSE_POSITIVE.to_owned(),
1185                payload(json!({ "response_id": "resp-3", "reason": "backchannel" })),
1186            ))
1187            .unwrap();
1188        let detected = recv(&mut det_rx).await;
1189        assert_eq!(detected.reason.as_deref(), Some("speech_overlap"));
1190        let false_positive = recv(&mut fp_rx).await;
1191        assert_eq!(false_positive.reason.as_deref(), Some("backchannel"));
1192    }
1193
1194    #[tokio::test]
1195    async fn start_response_and_wait_resolves_on_matching_created() {
1196        let (session, sender) = session().await;
1197        let options = ResponseOptions {
1198            generation_id: Some("gen-ack".to_owned()),
1199            ..Default::default()
1200        };
1201        tokio::spawn(async move {
1202            tokio::time::sleep(Duration::from_millis(100)).await;
1203            sender
1204                .send((
1205                    EVENT_RESPONSE_CREATED.to_owned(),
1206                    payload(json!({ "response_id": "resp-other", "generation_id": "gen-other" })),
1207                ))
1208                .unwrap();
1209            sender
1210                .send((
1211                    EVENT_RESPONSE_CREATED.to_owned(),
1212                    payload(json!({
1213                        "response_id": "resp-9",
1214                        "generation_id": "gen-ack",
1215                        "output": {
1216                            "model": "qwen3-tts:0.6b-clone",
1217                            "voice": "samantha",
1218                            "language": "fr",
1219                            "speed": 0.9,
1220                            "params": { "temperature": 0.7 }
1221                        }
1222                    })),
1223                ))
1224                .unwrap();
1225        });
1226        let ack = session
1227            .start_response_and_wait(Some(options), Duration::from_secs(2))
1228            .await
1229            .expect("ack within timeout");
1230        assert!(ack.accepted);
1231        assert_eq!(ack.generation_id, "gen-ack");
1232        assert_eq!(ack.response_id.as_deref(), Some("resp-9"));
1233        assert_eq!(
1234            ack.output,
1235            Some(ResponseOutput {
1236                model: "qwen3-tts:0.6b-clone".to_owned(),
1237                voice: Some("samantha".to_owned()),
1238                language: "fr".to_owned(),
1239                speed: 0.9,
1240                params: payload(json!({ "temperature": 0.7 })),
1241            })
1242        );
1243        assert!(ack.recoverable);
1244        assert_eq!(ack.error_code, None);
1245    }
1246
1247    #[tokio::test]
1248    async fn start_response_and_wait_surfaces_typed_rejection() {
1249        let (session, sender) = session().await;
1250        let options = ResponseOptions {
1251            generation_id: Some("gen-rejected".to_owned()),
1252            ..Default::default()
1253        };
1254        tokio::spawn(async move {
1255            tokio::time::sleep(Duration::from_millis(100)).await;
1256            sender
1257                .send((
1258                    EVENT_ERROR.to_owned(),
1259                    payload(json!({
1260                        "message": "busy",
1261                        "code": ERROR_CODE_RESPONSE_ALREADY_ACTIVE,
1262                        "recoverable": true,
1263                        "generation_id": "gen-rejected"
1264                    })),
1265                ))
1266                .unwrap();
1267        });
1268        let ack = session
1269            .start_response_and_wait(Some(options), Duration::from_secs(2))
1270            .await
1271            .expect("rejection within timeout");
1272        assert!(!ack.accepted);
1273        assert_eq!(ack.generation_id, "gen-rejected");
1274        assert_eq!(
1275            ack.error_code.as_deref(),
1276            Some(ERROR_CODE_RESPONSE_ALREADY_ACTIVE)
1277        );
1278        assert_eq!(ack.error_message.as_deref(), Some("busy"));
1279        assert!(ack.recoverable);
1280    }
1281
1282    #[tokio::test]
1283    async fn start_response_and_wait_times_out_without_ack() {
1284        let (session, _sender) = session().await;
1285        let error = session
1286            .start_response_and_wait(None, Duration::from_millis(100))
1287            .await
1288            .expect_err("no ack must time out");
1289        assert!(matches!(error, VoxRtcError::Timeout(_)));
1290    }
1291
1292    #[tokio::test]
1293    async fn on_speech_started_fires_with_timestamp() {
1294        let (session, sender) = session().await;
1295        let (tx, mut rx) = mpsc::unbounded_channel();
1296        let _listener = session.on_speech_started(move |event| {
1297            tx.send(event).unwrap();
1298        });
1299        sender
1300            .send((
1301                EVENT_SPEECH_STARTED.to_owned(),
1302                payload(json!({ "session_id": "sess-1", "timestamp_ms": 1234 })),
1303            ))
1304            .unwrap();
1305        let event = recv(&mut rx).await;
1306        assert_eq!(event.session_id, "sess-1");
1307        assert_eq!(event.channel_name, "/rtc/sess-1");
1308        assert_eq!(event.timestamp_ms, Some(1234.0));
1309    }
1310
1311    #[tokio::test]
1312    async fn on_speech_stopped_fires_with_timestamp() {
1313        let (session, sender) = session().await;
1314        let (tx, mut rx) = mpsc::unbounded_channel();
1315        let _listener = session.on_speech_stopped(move |event| {
1316            tx.send(event).unwrap();
1317        });
1318        sender
1319            .send((
1320                EVENT_SPEECH_STOPPED.to_owned(),
1321                payload(json!({ "timestamp_ms": 5678 })),
1322            ))
1323            .unwrap();
1324        let event = recv(&mut rx).await;
1325        assert_eq!(event.timestamp_ms, Some(5678.0));
1326    }
1327
1328    #[tokio::test]
1329    async fn on_transcript_delta_fires_with_fields() {
1330        let (session, sender) = session().await;
1331        let (tx, mut rx) = mpsc::unbounded_channel();
1332        let _listener = session.on_transcript_delta(move |event| {
1333            tx.send(event).unwrap();
1334        });
1335        sender
1336            .send((
1337                EVENT_TRANSCRIPT_DELTA.to_owned(),
1338                payload(json!({ "delta": "hel", "start_ms": 10, "end_ms": 20 })),
1339            ))
1340            .unwrap();
1341        let event = recv(&mut rx).await;
1342        assert_eq!(event.delta, "hel");
1343        assert_eq!(event.start_ms, Some(10.0));
1344        assert_eq!(event.end_ms, Some(20.0));
1345    }
1346
1347    #[tokio::test]
1348    async fn on_turn_eou_predicted_fires_with_fields() {
1349        let (session, sender) = session().await;
1350        let (tx, mut rx) = mpsc::unbounded_channel();
1351        let _listener = session.on_turn_eou_predicted(move |event| {
1352            tx.send(event).unwrap();
1353        });
1354        sender
1355            .send((
1356                EVENT_TURN_EOU_PREDICTED.to_owned(),
1357                payload(json!({
1358                    "probability": 0.82,
1359                    "threshold": 0.5,
1360                    "delay_ms": 120,
1361                    "start_ms": 0,
1362                    "end_ms": 300,
1363                    "decision": "end",
1364                    "action": "commit",
1365                    "turn_detector": "smart"
1366                })),
1367            ))
1368            .unwrap();
1369        let event = recv(&mut rx).await;
1370        assert_eq!(event.probability, Some(0.82));
1371        assert_eq!(event.threshold, Some(0.5));
1372        assert_eq!(event.delay_ms, Some(120.0));
1373        assert_eq!(event.start_ms, Some(0.0));
1374        assert_eq!(event.end_ms, Some(300.0));
1375        assert_eq!(event.decision.as_deref(), Some("end"));
1376        assert_eq!(event.action.as_deref(), Some("commit"));
1377        assert_eq!(event.turn_detector.as_deref(), Some("smart"));
1378    }
1379
1380    #[tokio::test]
1381    async fn handler_survives_a_lagged_broadcast() {
1382        let (session, sender) = session().await;
1383        let (tx, mut rx) = mpsc::unbounded_channel();
1384        let _listener = session.on_speech_started(move |event| {
1385            tx.send(event.timestamp_ms).unwrap();
1386        });
1387
1388        for index in 0..2100u32 {
1389            let _ = sender.send((
1390                EVENT_SPEECH_STARTED.to_owned(),
1391                payload(json!({ "timestamp_ms": index })),
1392            ));
1393        }
1394        let _ = sender.send((
1395            EVENT_SPEECH_STARTED.to_owned(),
1396            payload(json!({ "timestamp_ms": 9999 })),
1397        ));
1398
1399        let mut saw_final = false;
1400        while let Ok(Some(value)) = timeout(Duration::from_secs(1), rx.recv()).await {
1401            if value == Some(9999.0) {
1402                saw_final = true;
1403                break;
1404            }
1405        }
1406        assert!(
1407            saw_final,
1408            "loop must keep delivering events after a broadcast lag"
1409        );
1410    }
1411}