Skip to main content

active_call/call/
sip.rs

1use crate::call::state::{CallProgress, LegShared};
2use crate::callrecord::CallRecordHangupReason;
3use crate::event::EventSender;
4use crate::media::TrackId;
5use crate::media::stream::MediaStream;
6use crate::useragent::invitation::PendingDialog;
7use anyhow::Result;
8use chrono::Utc;
9use rsipstack::dialog::DialogId;
10use rsipstack::dialog::dialog::{
11    Dialog, DialogState, DialogStateReceiver, DialogStateSender, TerminatedReason,
12};
13use rsipstack::dialog::dialog_layer::DialogLayer;
14use rsipstack::dialog::invitation::InviteOption;
15use std::collections::HashMap;
16use std::sync::Arc;
17use tokio_util::sync::CancellationToken;
18use tracing::{info, warn};
19
20/// Remove `id` from the dialog layer and return the dialog, ready to be
21/// hung up. Shared by the dialog guards and `Invitation::hangup`.
22pub(crate) fn remove_dialog(layer: &DialogLayer, id: &DialogId) -> Option<Dialog> {
23    let dialog = layer.get_dialog(id)?;
24    layer.remove_dialog(id);
25    Some(dialog)
26}
27
28pub struct DialogStateReceiverGuard {
29    pub(super) dialog_layer: Arc<DialogLayer>,
30    pub(super) receiver: DialogStateReceiver,
31    pub(super) dialog_id: Option<DialogId>,
32    pub(super) hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
33}
34
35impl DialogStateReceiverGuard {
36    pub fn new(
37        dialog_layer: Arc<DialogLayer>,
38        receiver: DialogStateReceiver,
39        hangup_headers: Option<Vec<rsipstack::rsip::Header>>,
40    ) -> Self {
41        Self {
42            dialog_layer,
43            receiver,
44            dialog_id: None,
45            hangup_headers,
46        }
47    }
48    pub async fn recv(&mut self) -> Option<DialogState> {
49        let state = self.receiver.recv().await;
50        if let Some(ref s) = state {
51            self.dialog_id = Some(s.id().clone());
52        }
53        state
54    }
55
56    fn take_dialog(&mut self) -> Option<Dialog> {
57        let id = self.dialog_id.take()?;
58        info!(%id, "dialog removed on  drop");
59        if let Some(dialog) = remove_dialog(&self.dialog_layer, &id) {
60            return Some(dialog);
61        }
62        // A client dialog is only re-registered under its tag-bearing id once
63        // do_invite's process_invite() fully completes; while still ringing
64        // it's registered under the pre-response id (empty remote tag). Our
65        // tracked `id` already reflects whatever tag the latest DialogState
66        // event carried (e.g. a 183's), so the exact-key lookup above misses
67        // it for any dialog hung up before it's confirmed. Fall back to a
68        // call-id scan, which is unaffected by that id mutation.
69        self.dialog_layer
70            .get_client_dialog_by_call_id(&id.call_id)
71            .into_iter()
72            .next()
73            .map(Dialog::Invite)
74    }
75
76    pub async fn drop_async(&mut self) {
77        if let Some(dialog) = self.take_dialog() {
78            if let Err(e) = dialog.hangup_with_headers(self.hangup_headers.take()).await {
79                warn!(id=%dialog.id(), "error hanging up dialog on drop: {}", e);
80            }
81        }
82    }
83}
84
85impl Drop for DialogStateReceiverGuard {
86    fn drop(&mut self) {
87        if let Some(dialog) = self.take_dialog() {
88            crate::spawn(async move {
89                if let Err(e) = dialog.hangup().await {
90                    warn!(id=%dialog.id(), "error hanging up dialog on drop: {}", e);
91                }
92            });
93        }
94    }
95}
96
97pub(super) struct InviteDialogStates {
98    pub is_client: bool,
99    pub session_id: String,
100    pub track_id: TrackId,
101    pub cancel_token: CancellationToken,
102    pub event_sender: EventSender,
103    /// Lock-free shared state of the leg this dialog belongs to.
104    pub leg: LegShared,
105    pub media_stream: Arc<MediaStream>,
106    pub terminated_reason: Option<TerminatedReason>,
107    pub has_early_media: bool,
108    /// Set once the initial INVITE's answer has been applied via
109    /// `DialogState::Confirmed`. rsipstack reuses the same `Confirmed` event
110    /// for the ACK of a re-INVITE we received (see `handle_reinvite`), but
111    /// in that case the body is the PBX's own locally-generated answer, not
112    /// a remote description — that re-invite was already fully handled via
113    /// `DialogState::Updated`/`handshake()`. Re-applying it here corrupts
114    /// the peer connection's remote SSRC/address with our own values.
115    initial_confirmed: bool,
116    /// Hangup intent carried by this leg (refer legs with `auto_hangup`),
117    /// reported on the leg's TrackEnd so the call actor can hang up.
118    pub hangup_reason: Option<CallRecordHangupReason>,
119}
120
121impl InviteDialogStates {
122    pub(super) fn new(
123        is_client: bool,
124        session_id: String,
125        track_id: TrackId,
126        event_sender: EventSender,
127        media_stream: Arc<MediaStream>,
128        leg: LegShared,
129        cancel_token: CancellationToken,
130        hangup_reason: Option<CallRecordHangupReason>,
131    ) -> Self {
132        Self {
133            is_client,
134            session_id,
135            track_id,
136            cancel_token,
137            event_sender,
138            leg,
139            media_stream,
140            terminated_reason: None,
141            has_early_media: false,
142            initial_confirmed: false,
143            hangup_reason,
144        }
145    }
146}
147
148impl InviteDialogStates {
149    /// Called from `Drop` (synchronous context): everything used here is
150    /// lock-free (ArcSwap rcu/load + broadcast send), so no state or events
151    /// can be lost the way a failed `try_write` used to lose them.
152    pub(super) fn on_terminated(&mut self) {
153        let term = CallProgress::termination(self.terminated_reason.as_ref());
154        let status_code = term.status_code;
155        let reason = term.hangup_reason;
156        self.leg.update_progress(|p| {
157            p.last_status_code = status_code;
158            p.set_hangup_reason(reason.clone());
159        });
160        let progress = self.leg.progress.load_full();
161
162        self.event_sender
163            .send(crate::event::SessionEvent::TrackEnd {
164                track_id: self.track_id.clone(),
165                timestamp: crate::media::get_timestamp(),
166                duration: progress
167                    .answer_time
168                    .map(|t| (Utc::now() - t).num_milliseconds())
169                    .unwrap_or_default() as u64,
170                ssrc: self.leg.ssrc,
171                play_id: None,
172                auto_hangup: self.hangup_reason.clone(),
173            })
174            .ok();
175        let hangup_event = self
176            .leg
177            .build_hangup_event(self.track_id.clone(), Some(term.initiator.to_string()));
178        self.event_sender.send(hangup_event).ok();
179    }
180}
181
182impl Drop for InviteDialogStates {
183    fn drop(&mut self) {
184        self.on_terminated();
185        self.cancel_token.cancel();
186    }
187}
188
189impl DialogStateReceiverGuard {
190    pub(self) async fn dialog_event_loop(&mut self, states: &mut InviteDialogStates) -> Result<()> {
191        while let Some(event) = self.recv().await {
192            match event {
193                DialogState::Calling(dialog_id) => {
194                    info!(session_id=states.session_id, %dialog_id, "dialog calling");
195                    states
196                        .leg
197                        .update_progress(|p| p.session_id = dialog_id.to_string());
198                }
199                DialogState::Trying(_) => {}
200                DialogState::Early(dialog_id, resp) => {
201                    let code = resp.status_code.code();
202                    let body = resp.body();
203                    let answer = String::from_utf8_lossy(body);
204                    let has_sdp = !answer.is_empty();
205                    info!(session_id=states.session_id, %dialog_id, has_sdp=%has_sdp, "dialog early ({}): \n{}", code, answer);
206
207                    states.leg.update_progress(|p| p.on_early(code));
208
209                    if !states.is_client {
210                        continue;
211                    }
212
213                    states
214                        .event_sender
215                        .send(crate::event::SessionEvent::Ringing {
216                            track_id: states.track_id.clone(),
217                            timestamp: crate::media::get_timestamp(),
218                            early_media: has_sdp,
219                            refer: Some(states.leg.is_refer),
220                        })?;
221
222                    if has_sdp {
223                        states.has_early_media = true;
224                        states.leg.update_progress(|p| p.try_set_answer(&answer));
225                        states
226                            .media_stream
227                            .update_remote_description_provisional(
228                                &states.track_id,
229                                &answer.to_string(),
230                            )
231                            .await?;
232                    }
233                }
234                DialogState::Confirmed(dialog_id, msg) => {
235                    info!(session_id=states.session_id, %dialog_id, has_early_media=%states.has_early_media, "dialog confirmed");
236                    states
237                        .leg
238                        .update_progress(|p| p.on_confirmed(dialog_id.to_string()));
239                    if states.is_client && !states.initial_confirmed {
240                        states.initial_confirmed = true;
241                        let answer = String::from_utf8_lossy(msg.body());
242                        let answer = answer.trim();
243                        if !answer.is_empty() {
244                            if states.has_early_media {
245                                info!(
246                                    session_id = states.session_id,
247                                    "updating remote description with final answer after early media (force=true)"
248                                );
249                                // Force update when transitioning from early media (183) to confirmed (200 OK)
250                                // This ensures media parameters are properly updated even if SDP appears similar
251                                if let Err(e) = states
252                                    .media_stream
253                                    .update_remote_description_force(
254                                        &states.track_id,
255                                        &answer.to_string(),
256                                    )
257                                    .await
258                                {
259                                    tracing::warn!(
260                                        session_id = states.session_id,
261                                        "failed to force update remote description on confirmed: {}",
262                                        e
263                                    );
264                                }
265                            } else {
266                                if let Err(e) = states
267                                    .media_stream
268                                    .update_remote_description(
269                                        &states.track_id,
270                                        &answer.to_string(),
271                                    )
272                                    .await
273                                {
274                                    tracing::warn!(
275                                        session_id = states.session_id,
276                                        "failed to update remote description on confirmed: {}",
277                                        e
278                                    );
279                                }
280                            }
281                        }
282                    }
283                }
284                DialogState::Info(dialog_id, req, tx_handle) => {
285                    let body_str = String::from_utf8_lossy(req.body());
286                    info!(session_id=states.session_id, %dialog_id, body=%body_str, "dialog info received");
287                    if body_str.starts_with("Signal=") {
288                        let digit = body_str.trim_start_matches("Signal=").chars().next();
289                        if let Some(digit) = digit {
290                            states.event_sender.send(crate::event::SessionEvent::Dtmf {
291                                track_id: states.track_id.clone(),
292                                timestamp: crate::media::get_timestamp(),
293                                digit: digit.to_string(),
294                                refer: Some(states.leg.is_refer),
295                            })?;
296                        }
297                    }
298                    tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
299                }
300                DialogState::Message(dialog_id, req, tx_handle) => {
301                    let body_str = String::from_utf8_lossy(req.body()).to_string();
302                    let content_type = req.headers.iter().find_map(|h| {
303                        if let rsipstack::rsip::Header::ContentType(content_type) = h {
304                            Some(content_type.value().to_string())
305                        } else {
306                            None
307                        }
308                    });
309                    info!(
310                        session_id=states.session_id,
311                        %dialog_id,
312                        content_type=content_type.as_deref(),
313                        body=%body_str,
314                        "dialog message received"
315                    );
316                    states
317                        .event_sender
318                        .send(crate::event::SessionEvent::Message {
319                            track_id: states.track_id.clone(),
320                            timestamp: crate::media::get_timestamp(),
321                            body: body_str,
322                            content_type,
323                            refer: Some(states.leg.is_refer),
324                        })
325                        .ok();
326                    tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
327                }
328                DialogState::Updated(dialog_id, _req, tx_handle) => {
329                    info!(session_id = states.session_id, %dialog_id, "dialog update received");
330                    let mut answer_sdp = None;
331                    if let Some(sdp_body) = _req.body().get(..) {
332                        let sdp_str = String::from_utf8_lossy(sdp_body);
333                        if !sdp_str.is_empty()
334                            && (_req.method == rsipstack::rsip::Method::Invite
335                                || _req.method == rsipstack::rsip::Method::Update)
336                        {
337                            info!(session_id=states.session_id, %dialog_id, method=%_req.method, "handling re-invite/update offer");
338
339                            // Detect hold state from SDP
340                            let is_on_hold =
341                                crate::media::negotiate::detect_hold_state_from_sdp(&sdp_str);
342                            info!(session_id=states.session_id, %dialog_id, is_on_hold=%is_on_hold, "detected hold state from re-invite SDP");
343
344                            // Update media stream hold state + emit hold event
345                            apply_hold_state(states, is_on_hold).await;
346
347                            match states
348                                .media_stream
349                                .handshake(&states.track_id, sdp_str.to_string(), None)
350                                .await
351                            {
352                                Ok(sdp) => answer_sdp = Some(sdp),
353                                Err(e) => {
354                                    warn!(
355                                        session_id = states.session_id,
356                                        "failed to handle re-invite: {}", e
357                                    );
358                                }
359                            }
360                        } else {
361                            info!(session_id=states.session_id, %dialog_id, "updating remote description:\n{}", sdp_str);
362
363                            // Also check hold state for non-INVITE/UPDATE messages with SDP
364                            let is_on_hold =
365                                crate::media::negotiate::detect_hold_state_from_sdp(&sdp_str);
366                            apply_hold_state(states, is_on_hold).await;
367
368                            states
369                                .media_stream
370                                .update_remote_description(&states.track_id, &sdp_str.to_string())
371                                .await?;
372                        }
373                    }
374
375                    if let Some(sdp) = answer_sdp {
376                        tx_handle
377                            .respond(
378                                rsipstack::rsip::StatusCode::OK,
379                                Some(vec![rsipstack::rsip::Header::ContentType(
380                                    "application/sdp".to_string().into(),
381                                )]),
382                                Some(sdp.into()),
383                            )
384                            .await
385                            .ok();
386                    } else {
387                        tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
388                    }
389                }
390                DialogState::Options(dialog_id, _req, tx_handle) => {
391                    info!(session_id = states.session_id, %dialog_id, "dialog options received");
392                    tx_handle.reply(rsipstack::rsip::StatusCode::OK).await.ok();
393                }
394                DialogState::Refer(dialog_id, req, tx_handle) => {
395                    let refer_to = req
396                        .headers
397                        .iter()
398                        .find_map(|h| {
399                            if let rsipstack::rsip::Header::ReferTo(h) = h {
400                                return Some(h.value().to_string());
401                            }
402                            None
403                        })
404                        .unwrap_or_default();
405                    let referred_by = req.headers.iter().find_map(|h| {
406                        if let rsipstack::rsip::Header::ReferredBy(h) = h {
407                            return Some(h.value().to_string());
408                        }
409                        None
410                    });
411                    info!(session_id = states.session_id, %dialog_id, %refer_to, "received REFER");
412                    tx_handle
413                        .reply(rsipstack::rsip::StatusCode::Other(202, "Accepted".into()))
414                        .await
415                        .ok();
416                    states
417                        .event_sender
418                        .send(crate::event::SessionEvent::TransferRequest {
419                            track_id: states.track_id.clone(),
420                            timestamp: crate::media::get_timestamp(),
421                            refer_to,
422                            referred_by,
423                            refer: Some(states.leg.is_refer),
424                        })
425                        .ok();
426                }
427                DialogState::Terminated(dialog_id, reason) => {
428                    info!(
429                        session_id = states.session_id,
430                        ?dialog_id,
431                        ?reason,
432                        "dialog terminated"
433                    );
434                    states.terminated_reason = Some(reason.clone());
435                    return Ok(());
436                }
437                other_state => {
438                    info!(
439                        session_id = states.session_id,
440                        %other_state,
441                        "dialog received other state"
442                    );
443                }
444            }
445        }
446        Ok(())
447    }
448
449    pub(super) async fn process_dialog(&mut self, mut states: InviteDialogStates) {
450        let token = states.cancel_token.clone();
451        tokio::select! {
452            _ = token.cancelled() => {
453                states.terminated_reason = Some(TerminatedReason::UacCancel);
454            }
455            _ = self.dialog_event_loop(&mut states) => {}
456        };
457
458        // Update hangup headers from the leg extras if available
459        let extras = states.leg.extras.load_full();
460        if let Some(headers) = crate::sip_util::hangup_headers_from_extras(&extras) {
461            match &mut self.hangup_headers {
462                Some(existing) => existing.extend(headers),
463                None => self.hangup_headers = Some(headers),
464            }
465        }
466
467        self.drop_async().await;
468    }
469}
470
471/// Apply a hold/resume transition to the media track and emit the Hold event.
472async fn apply_hold_state(states: &mut InviteDialogStates, is_on_hold: bool) {
473    if is_on_hold {
474        states
475            .media_stream
476            .hold_track(Some(states.track_id.clone()))
477            .await;
478    } else {
479        states
480            .media_stream
481            .resume_track(Some(states.track_id.clone()))
482            .await;
483    }
484    states
485        .event_sender
486        .send(crate::event::SessionEvent::Hold {
487            track_id: states.track_id.clone(),
488            timestamp: crate::media::get_timestamp(),
489            on_hold: is_on_hold,
490            refer: Some(states.leg.is_refer),
491        })
492        .ok();
493}
494
495#[derive(Clone)]
496pub struct Invitation {
497    pub dialog_layer: Arc<DialogLayer>,
498    pub pending_dialogs: Arc<std::sync::Mutex<HashMap<DialogId, PendingDialog>>>,
499    /// Maps the (short) public session id of an incoming call to the DialogId
500    /// captured when the INVITE arrived. Incoming sessions no longer reuse the
501    /// raw dialog-id string, so every session-id based lookup resolves here.
502    sessions: Arc<std::sync::Mutex<HashMap<String, DialogId>>>,
503}
504
505impl Invitation {
506    pub fn new(dialog_layer: Arc<DialogLayer>) -> Self {
507        Self {
508            dialog_layer,
509            pending_dialogs: Arc::new(std::sync::Mutex::new(HashMap::new())),
510            sessions: Arc::new(std::sync::Mutex::new(HashMap::new())),
511        }
512    }
513
514    pub fn register_session(&self, session_id: &str, dialog_id: &DialogId) {
515        self.sessions
516            .lock()
517            .map(|mut ss| ss.insert(session_id.to_string(), dialog_id.clone()))
518            .ok();
519    }
520
521    pub fn unregister_session(&self, session_id: &str) {
522        self.sessions.lock().map(|mut ss| ss.remove(session_id)).ok();
523    }
524
525    pub fn session_exists(&self, session_id: &str) -> bool {
526        self.sessions
527            .lock()
528            .map(|ss| ss.contains_key(session_id))
529            .unwrap_or(false)
530    }
531
532    pub fn add_pending(&self, dialog_id: DialogId, pending: PendingDialog) {
533        self.pending_dialogs
534            .lock()
535            .map(|mut ps| ps.insert(dialog_id, pending))
536            .ok();
537    }
538
539    pub fn get_pending_call(&self, dialog_id: &DialogId) -> Option<PendingDialog> {
540        self.pending_dialogs
541            .lock()
542            .ok()
543            .and_then(|mut ps| ps.remove(dialog_id))
544    }
545
546    pub fn has_pending_call(&self, dialog_id: &DialogId) -> bool {
547        self.pending_dialogs
548            .lock()
549            .ok()
550            .map(|ps| ps.contains_key(dialog_id))
551            .unwrap_or(false)
552    }
553
554    pub fn find_dialog_id_by_session_id(&self, session_id: &str) -> Option<DialogId> {
555        if let Some(id) = self
556            .sessions
557            .lock()
558            .ok()
559            .and_then(|ss| ss.get(session_id).cloned())
560        {
561            return Some(id);
562        }
563        self.pending_dialogs.lock().ok().and_then(|ps| {
564            ps.iter()
565                .find(|(id, _)| id.to_string() == session_id)
566                .map(|(id, _)| id.clone())
567        })
568    }
569
570    /// Reject a pending dialog or hang up an established one.
571    pub async fn hangup(
572        &self,
573        dialog_id: DialogId,
574        code: Option<rsipstack::rsip::StatusCode>,
575        reason: Option<String>,
576    ) -> Result<()> {
577        if let Some(call) = self.get_pending_call(&dialog_id) {
578            call.dialog.reject(code, reason).ok();
579        }
580        if let Some(dialog) = remove_dialog(&self.dialog_layer, &dialog_id) {
581            dialog.hangup().await.ok();
582        }
583        Ok(())
584    }
585
586    pub async fn invite(
587        &self,
588        invite_option: InviteOption,
589        state_sender: DialogStateSender,
590    ) -> Result<(DialogId, Option<Vec<u8>>), rsipstack::Error> {
591        let (dialog, resp) = self
592            .dialog_layer
593            .do_invite(invite_option, state_sender)
594            .await?;
595
596        let offer = match resp {
597            Some(resp) => match resp.status_code.kind() {
598                rsipstack::rsip::StatusCodeKind::Successful => {
599                    let offer = resp.body.clone();
600                    Some(offer)
601                }
602                _ => {
603                    let reason = resp
604                        .reason_phrase()
605                        .unwrap_or(&resp.status_code.to_string())
606                        .to_string();
607                    return Err(rsipstack::Error::DialogError(
608                        reason,
609                        dialog.id(),
610                        resp.status_code,
611                    ));
612                }
613            },
614            None => {
615                return Err(rsipstack::Error::DialogError(
616                    "no response received".to_string(),
617                    dialog.id(),
618                    rsipstack::rsip::StatusCode::NotAcceptableHere,
619                ));
620            }
621        };
622        Ok((dialog.id(), offer))
623    }
624}
625
626#[cfg(test)]
627mod tests {
628    use super::*;
629    use crate::call::state::CallProgress;
630    use crate::media::stream::MediaStreamBuilder;
631
632    // SDP used to simulate an early-media 183 Session Progress response.
633    const EARLY_MEDIA_SDP: &str = "v=0\r\n\
634        o=- 1000 1 IN IP4 192.168.1.100\r\n\
635        s=SIP Call\r\n\
636        t=0 0\r\n\
637        m=audio 10000 RTP/AVP 0\r\n\
638        c=IN IP4 192.168.1.100\r\n\
639        a=rtpmap:0 PCMU/8000\r\n\
640        a=sendrecv\r\n";
641
642    fn make_states(has_early_media: bool) -> InviteDialogStates {
643        let (event_tx, _event_rx) = tokio::sync::broadcast::channel(16);
644        let media_stream = Arc::new(
645            MediaStreamBuilder::new(event_tx.clone())
646                .with_id("test-stream".to_string())
647                .build(),
648        );
649        let cancel_token = CancellationToken::new();
650        let leg = LegShared::new(1000, true, CallProgress::default());
651
652        InviteDialogStates {
653            is_client: true,
654            session_id: "test-session".to_string(),
655            track_id: "test-track".to_string(),
656            cancel_token: cancel_token.clone(),
657            event_sender: event_tx.clone(),
658            leg,
659            media_stream,
660            terminated_reason: None,
661            has_early_media,
662            initial_confirmed: false,
663            hangup_reason: None,
664        }
665    }
666
667    /// Verify that when a 183 Session Progress with SDP arrives (`DialogState::Early`),
668    /// the early SDP is stored in the leg progress so it can serve as a fallback
669    /// when the final 200 OK has an empty body.
670    #[tokio::test]
671    async fn test_early_sdp_stored_in_leg_progress() {
672        let mut states = make_states(false);
673
674        // Simulate DialogState::Early with SDP body (183 Session Progress):
675        // same steps the Early branch performs.
676        let answer = EARLY_MEDIA_SDP.to_string();
677        let has_sdp = !answer.is_empty();
678        if states.is_client && has_sdp {
679            states.has_early_media = true;
680            states.leg.update_progress(|p| p.try_set_answer(&answer));
681        }
682
683        // Assert: early SDP is stored
684        let progress = states.leg.progress.load_full();
685        assert!(
686            progress.answer.is_some(),
687            "leg progress answer should be set after 183 with SDP"
688        );
689        assert_eq!(
690            progress.answer.as_deref().unwrap(),
691            EARLY_MEDIA_SDP,
692            "leg progress answer should contain the early SDP"
693        );
694        assert!(states.has_early_media, "has_early_media should be true");
695    }
696
697    /// Verify that when a 200 OK arrives with an empty body after early media has been
698    /// negotiated, the leg progress retains the early SDP (not overwritten with "").
699    ///
700    /// This is the regression test for the bug where a late 200 OK with empty body would
701    /// cause `SessionEvent::Answer { sdp: "" }` to be emitted, making the answer event
702    /// appear as if no SDP was negotiated.
703    #[tokio::test]
704    async fn test_confirmed_empty_body_keeps_early_sdp() {
705        let mut states = make_states(false);
706
707        // Step 1: simulate 183 with SDP → set has_early_media and progress answer
708        states.has_early_media = true;
709        states
710            .leg
711            .update_progress(|p| p.try_set_answer(EARLY_MEDIA_SDP));
712
713        // Step 2: simulate 200 OK with empty body (Confirmed handler logic)
714        states
715            .leg
716            .update_progress(|p| p.on_confirmed("dialog-1".to_string()));
717        // The Confirmed handler only calls update_remote_description when the body
718        // is non-empty; it does NOT overwrite the progress answer.
719        let confirmed_answer = String::new();
720        assert!(
721            confirmed_answer.trim().is_empty(),
722            "empty body must not be applied"
723        );
724
725        // Assert: leg progress still holds the early SDP
726        let progress = states.leg.progress.load_full();
727        assert!(
728            progress.answer.is_some(),
729            "answer must not be None after 200 OK with empty body"
730        );
731        let stored_answer = progress.answer.as_deref().unwrap();
732        assert!(
733            !stored_answer.is_empty(),
734            "answer must not be empty after 200 OK with empty body"
735        );
736        assert_eq!(
737            stored_answer, EARLY_MEDIA_SDP,
738            "answer should still be the early SDP after 200 OK with empty body"
739        );
740    }
741
742    /// Verify that `create_outgoing_sip_track`'s fallback logic works:
743    /// when the 200 OK body is empty but the leg progress has the early SDP,
744    /// the fallback path is taken and the early SDP is returned (not an empty string).
745    #[tokio::test]
746    async fn test_answer_fallback_to_early_sdp_when_200ok_empty() {
747        // Simulate what the Early (183) handler does: store the early SDP.
748        let states = make_states(true);
749        states
750            .leg
751            .update_progress(|p| p.try_set_answer(EARLY_MEDIA_SDP));
752
753        // Simulate what create_outgoing_sip_track does when 200 OK has empty body.
754        let early = states.leg.progress.load_full().answer.clone();
755        let raw_answer: Option<Vec<u8>> = Some(vec![]); // empty body from 200 OK
756
757        let (answer, already_applied) =
758            crate::call::state::resolve_final_answer(raw_answer, early.as_ref()).unwrap();
759
760        // The answer returned to setup_caller_track (and used in SessionEvent::Answer)
761        // must be the early SDP, not an empty string.
762        assert!(
763            !answer.is_empty(),
764            "Resolved answer must not be empty — should contain the early SDP"
765        );
766        assert_eq!(
767            answer, EARLY_MEDIA_SDP,
768            "Resolved answer should be the early SDP from the 183 handler"
769        );
770        assert!(
771            already_applied,
772            "remote_description_already_applied should be true when using early SDP fallback"
773        );
774    }
775
776    /// Verify the normal case: when 200 OK carries its own SDP body,
777    /// that SDP is used directly (not the early SDP) and remote description
778    /// should be applied.
779    #[tokio::test]
780    async fn test_answer_uses_200ok_sdp_when_present() {
781        const FINAL_SDP: &str = "v=0\r\n\
782            o=- 2000 2 IN IP4 10.0.0.1\r\n\
783            s=SIP Call\r\n\
784            t=0 0\r\n\
785            m=audio 20000 RTP/AVP 0\r\n\
786            c=IN IP4 10.0.0.1\r\n\
787            a=rtpmap:0 PCMU/8000\r\n\
788            a=sendrecv\r\n";
789
790        let states = make_states(true);
791        states
792            .leg
793            .update_progress(|p| p.try_set_answer(EARLY_MEDIA_SDP));
794
795        let early = states.leg.progress.load_full().answer.clone();
796        let (answer, already_applied) = crate::call::state::resolve_final_answer(
797            Some(FINAL_SDP.as_bytes().to_vec()),
798            early.as_ref(),
799        )
800        .unwrap();
801
802        assert_eq!(
803            answer, FINAL_SDP,
804            "When 200 OK has SDP, it should be used (not the early SDP)"
805        );
806        assert!(
807            !already_applied,
808            "remote_description_already_applied should be false when 200 OK has SDP body"
809        );
810    }
811
812    /// Regression: `on_terminated` runs in a synchronous `Drop` context.
813    /// The old implementation used `try_write` and silently dropped the
814    /// status/reason updates AND the TrackEnd/Hangup events when the lock was
815    /// contended. The lock-free (ArcSwap) implementation must always emit both
816    /// events and record the termination, even while another task keeps
817    /// mutating the progress concurrently.
818    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
819    async fn test_on_terminated_always_emits_events_and_records_state() {
820        use crate::event::SessionEvent;
821        use std::sync::atomic::{AtomicBool, Ordering};
822
823        let mut states = make_states(false);
824        states.terminated_reason = Some(TerminatedReason::UacCancel);
825        let leg = states.leg.clone();
826        let event_sender = states.event_sender.clone();
827        let mut event_receiver = event_sender.subscribe();
828
829        // Hammer the progress concurrently, as a busy actor would.
830        let stop = Arc::new(AtomicBool::new(false));
831        let stop2 = stop.clone();
832        let progress = leg.progress.clone();
833        let writer = crate::spawn(async move {
834            while !stop2.load(Ordering::Relaxed) {
835                // Touch unrelated fields, like a busy actor would (never the
836                // termination fields), so writers don't clobber each other.
837                progress.rcu(|p| {
838                    let mut p = CallProgress::clone(p);
839                    p.answer_time.get_or_insert_with(chrono::Utc::now);
840                    p
841                });
842                tokio::task::yield_now().await;
843            }
844        });
845
846        // Give the writer a moment to start contending.
847        tokio::time::sleep(std::time::Duration::from_millis(20)).await;
848
849        // Synchronous drop, exactly like the real dialog task teardown.
850        drop(states);
851
852        stop.store(true, Ordering::Relaxed);
853        let _ = writer.await;
854
855        // Both events must have been emitted despite the concurrent writer.
856        let mut saw_track_end = false;
857        let mut saw_hangup = false;
858        while let Ok(event) = event_receiver.try_recv() {
859            match event {
860                SessionEvent::TrackEnd { .. } => saw_track_end = true,
861                SessionEvent::Hangup { refer, .. } => {
862                    assert_eq!(refer, Some(true), "refer flag comes from the leg");
863                    saw_hangup = true;
864                }
865                _ => {}
866            }
867        }
868        assert!(
869            saw_track_end,
870            "TrackEnd must be emitted from on_terminated even under contention"
871        );
872        assert!(
873            saw_hangup,
874            "Hangup must be emitted from on_terminated even under contention"
875        );
876
877        // The termination must be recorded (487 for UacCancel); the concurrent
878        // writer only ever writes 100, so observing 487 proves the rcu landed.
879        let progress = leg.progress.load_full();
880        assert_eq!(progress.last_status_code, 487);
881        assert_eq!(
882            progress.hangup_reason,
883            Some(crate::callrecord::CallRecordHangupReason::Canceled)
884        );
885    }
886
887    /// A minimal counting track that records how many times its remote
888    /// description was (force-)updated, without touching any real media.
889    struct CountingTrack {
890        id: TrackId,
891        config: crate::media::track::TrackConfig,
892        processor_chain: crate::media::processor::ProcessorChain,
893        updates: std::sync::Arc<std::sync::atomic::AtomicUsize>,
894    }
895
896    #[async_trait::async_trait]
897    impl crate::media::track::Track for CountingTrack {
898        fn ssrc(&self) -> u32 {
899            0
900        }
901        fn id(&self) -> &TrackId {
902            &self.id
903        }
904        fn config(&self) -> &crate::media::track::TrackConfig {
905            &self.config
906        }
907        fn processor_chain(&mut self) -> &mut crate::media::processor::ProcessorChain {
908            &mut self.processor_chain
909        }
910        async fn handshake(
911            &mut self,
912            _offer: String,
913            _timeout: Option<tokio::time::Duration>,
914        ) -> Result<String> {
915            Ok(String::new())
916        }
917        async fn update_remote_description(&mut self, _answer: &String) -> Result<()> {
918            self.updates
919                .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
920            Ok(())
921        }
922        async fn update_remote_description_force(&mut self, _answer: &String) -> Result<()> {
923            self.updates
924                .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
925            Ok(())
926        }
927        async fn start(
928            &mut self,
929            _event_sender: EventSender,
930            _packet_sender: crate::media::track::TrackPacketSender,
931        ) -> Result<()> {
932            Ok(())
933        }
934        async fn stop(&self) -> Result<()> {
935            Ok(())
936        }
937        async fn send_packet(&mut self, _packet: &crate::media::AudioFrame) -> Result<()> {
938            Ok(())
939        }
940    }
941
942    /// Regression: a client dialog must apply its remote SDP answer exactly
943    /// once. The `initial_confirmed` guard exists because the final `Confirmed`
944    /// event fires both when the 200 OK is processed *and* again when the ACK
945    /// for the re-INVITE completes; without the guard the answer would be
946    /// re-applied (and a duplicate `TrackStart`/`SessionEvent::Answer` emitted).
947    #[tokio::test]
948    async fn test_confirmed_applies_remote_answer_only_once() {
949        let mut states = make_states(false);
950
951        let updates = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
952        let track = CountingTrack {
953            id: states.track_id.clone(),
954            config: crate::media::track::TrackConfig::default(),
955            processor_chain: crate::media::processor::ProcessorChain::new(16000),
956            updates: updates.clone(),
957        };
958        states
959            .media_stream
960            .update_track(Box::new(track), None)
961            .await;
962
963        let endpoint = {
964            let mut builder = rsipstack::EndpointBuilder::new();
965            builder.build()
966        };
967        let dialog_layer = Arc::new(DialogLayer::new(endpoint.inner.clone()));
968
969        let (tx, rx) = tokio::sync::mpsc::unbounded_channel::<DialogState>();
970        let mut guard = DialogStateReceiverGuard::new(dialog_layer, rx, None);
971
972        let dialog_id = DialogId {
973            call_id: "test-call-id".to_string(),
974            local_tag: "test-local-tag".to_string(),
975            remote_tag: "test-remote-tag".to_string(),
976        };
977
978        let mut resp = rsipstack::rsip::Response::default();
979        resp.body = EARLY_MEDIA_SDP.as_bytes().to_vec();
980
981        // Two Confirmed events (200 OK + re-INVITE ACK) must only apply once.
982        tx.send(DialogState::Confirmed(dialog_id.clone(), resp.clone()))
983            .unwrap();
984        tx.send(DialogState::Confirmed(dialog_id, resp)).unwrap();
985        drop(tx);
986
987        guard
988            .dialog_event_loop(&mut states)
989            .await
990            .expect("dialog event loop must complete");
991
992        assert_eq!(
993            updates.load(std::sync::atomic::Ordering::SeqCst),
994            1,
995            "remote answer must be applied exactly once, not on the re-INVITE ACK Confirmed event"
996        );
997    }
998
999    #[test]
1000    fn test_session_map_register_find_unregister() {
1001        let endpoint = {
1002            let mut builder = rsipstack::EndpointBuilder::new();
1003            builder.build()
1004        };
1005        let invitation = Invitation::new(Arc::new(DialogLayer::new(endpoint.inner.clone())));
1006
1007        let dialog_id = DialogId {
1008            call_id: "long-call-id-from-carrier".to_string(),
1009            local_tag: String::new(),
1010            remote_tag: "remote-tag".to_string(),
1011        };
1012
1013        assert!(!invitation.session_exists("s.3f9a2b1c4d5e"));
1014        assert!(
1015            invitation
1016                .find_dialog_id_by_session_id("s.3f9a2b1c4d5e")
1017                .is_none()
1018        );
1019
1020        invitation.register_session("s.3f9a2b1c4d5e", &dialog_id);
1021        assert!(invitation.session_exists("s.3f9a2b1c4d5e"));
1022        assert_eq!(
1023            invitation
1024                .find_dialog_id_by_session_id("s.3f9a2b1c4d5e")
1025                .unwrap(),
1026            dialog_id
1027        );
1028
1029        invitation.unregister_session("s.3f9a2b1c4d5e");
1030        assert!(!invitation.session_exists("s.3f9a2b1c4d5e"));
1031        assert!(
1032            invitation
1033                .find_dialog_id_by_session_id("s.3f9a2b1c4d5e")
1034                .is_none()
1035        );
1036        // Unregistering an unknown session is a no-op.
1037        invitation.unregister_session("s.unknown");
1038    }
1039
1040    /// The legacy fallback must keep resolving session ids that match a pending
1041    /// dialog's dialog-id string (e.g. sessions registered before the mapping
1042    /// table existed, or callers still using the raw dialog id).
1043    #[test]
1044    fn test_find_dialog_id_falls_back_to_pending_scan() {
1045        use rsipstack::dialog::dialog::DialogInner;
1046        use rsipstack::dialog::invite_dialog::InviteDialog;
1047        use rsipstack::rsip::typed::{Contact, CSeq, From, To, Via};
1048        use rsipstack::rsip::{Header, Request};
1049        use rsipstack::transaction::key::TransactionRole;
1050
1051        let endpoint = {
1052            let mut builder = rsipstack::EndpointBuilder::new();
1053            builder.build()
1054        };
1055        let invitation = Invitation::new(Arc::new(DialogLayer::new(endpoint.inner.clone())));
1056
1057        let dialog_id = DialogId {
1058            call_id: "legacy-call-id".to_string(),
1059            local_tag: String::new(),
1060            remote_tag: "remote-tag".to_string(),
1061        };
1062
1063        let initial_request = Request {
1064            method: rsipstack::rsip::Method::Invite,
1065            uri: rsipstack::rsip::Uri::try_from("sip:bob@example.com:5060").unwrap(),
1066            headers: vec![
1067                Via::parse("SIP/2.0/UDP alice.example.com:5060;branch=z9hG4bKnashds")
1068                    .unwrap()
1069                    .into(),
1070                CSeq::parse("1 INVITE").unwrap().into(),
1071                From::parse("Alice <sip:alice@example.com>;tag=remote-tag")
1072                    .unwrap()
1073                    .into(),
1074                To::parse("Bob <sip:bob@example.com>").unwrap().into(),
1075                Header::CallId("legacy-call-id".into()),
1076                Contact::parse("<sip:alice@alice.example.com:5060>").unwrap().into(),
1077                Header::MaxForwards("70".into()),
1078            ]
1079            .into(),
1080            version: rsipstack::rsip::Version::V2,
1081            body: vec![],
1082        };
1083
1084        let (state_sender, _state_receiver) = tokio::sync::mpsc::unbounded_channel();
1085        let (tu_sender, _tu_receiver) = tokio::sync::mpsc::unbounded_channel();
1086        let inner = std::sync::Arc::new(
1087            DialogInner::new(
1088                TransactionRole::Server,
1089                dialog_id.clone(),
1090                initial_request,
1091                endpoint.inner.clone(),
1092                state_sender,
1093                None,
1094                None,
1095                tu_sender,
1096            )
1097            .expect("failed to create dialog inner"),
1098        );
1099        let dialog = InviteDialog::from_inner(inner);
1100
1101        assert!(
1102            invitation
1103                .find_dialog_id_by_session_id(&dialog_id.to_string())
1104                .is_none()
1105        );
1106
1107        invitation.add_pending(
1108            dialog_id.clone(),
1109            PendingDialog {
1110                token: tokio_util::sync::CancellationToken::new(),
1111                dialog,
1112                state_receiver: {
1113                    let (_tx, rx) = tokio::sync::mpsc::unbounded_channel();
1114                    rx
1115                },
1116            },
1117        );
1118
1119        assert_eq!(
1120            invitation
1121                .find_dialog_id_by_session_id(&dialog_id.to_string())
1122                .unwrap(),
1123            dialog_id
1124        );
1125    }
1126}