Skip to main content

active_call/
app.rs

1use crate::{
2    call::{ActiveCallRef, sip::Invitation},
3    callrecord::{
4        CallRecordFormatter, CallRecordManagerBuilder, CallRecordSender, DefaultCallRecordFormatter,
5    },
6    config::Config,
7    locator::RewriteTargetLocator,
8    useragent::{
9        RegisterOption,
10        invitation::{
11            FnCreateInvitationHandler, PendingDialog, PendingDialogGuard,
12            default_create_invite_handler,
13        },
14        public_address::{
15            LearningMessageInspector, SharedPublicAddress, build_contact, build_public_contact_uri,
16            find_local_addr_for_uri,
17        },
18        registration::{RegistrationHandle, UserCredential},
19    },
20};
21
22use crate::media::{cache::set_cache_dir, engine::StreamEngine};
23use anyhow::Result;
24use arc_swap::ArcSwap;
25use async_trait::async_trait;
26use chrono::{DateTime, Local};
27use futures::FutureExt;
28use humantime::parse_duration;
29use rsipstack::rsip::prelude::HeadersExt;
30use rsipstack::rsip::{Accept, Method, typed};
31use rsipstack::transaction::{
32    Endpoint, TransactionReceiver,
33    endpoint::{TargetLocator, TransportEventInspector},
34};
35use rsipstack::transport::transport_layer::DomainResolver;
36use rsipstack::{dialog::dialog_layer::DialogLayer, transaction::endpoint::MessageInspector};
37use rsipstack::{
38    rsip::{Host, HostWithPort},
39    transport::SipAddr,
40};
41use std::future::pending;
42use std::panic::{AssertUnwindSafe, catch_unwind};
43use std::str::FromStr;
44use std::sync::{Arc, RwLock};
45use std::time::Duration;
46use std::{collections::HashMap, net::SocketAddr};
47use std::{collections::HashSet, time::Instant};
48use std::{
49    path::Path,
50    sync::atomic::{AtomicBool, AtomicU64, Ordering},
51};
52use tokio::select;
53use tokio::sync::Mutex;
54use tokio_util::sync::CancellationToken;
55use tracing::{info, warn};
56
57/// Generates a short unique session id for incoming calls, e.g. `s.3f9a2b1c4d5e`,
58/// instead of reusing the raw SIP dialog-id string. Collisions with live
59/// sessions are retried (practically impossible with 48 bits of randomness).
60fn generate_short_session_id(invitation: &Invitation) -> String {
61    loop {
62        let uuid = uuid::Uuid::new_v4().simple().to_string();
63        let session_id = format!("s.{}", &uuid[..12]);
64        if !invitation.session_exists(&session_id) {
65            return session_id;
66        }
67    }
68}
69
70pub struct AppStateInner {
71    pub config: Arc<Config>,
72    pub token: CancellationToken,
73    pub stream_engine: Arc<StreamEngine>,
74    pub callrecord_sender: Option<CallRecordSender>,
75    pub endpoint: Endpoint,
76    pub registration_handles: Mutex<HashMap<String, CancellationToken>>,
77    pub alive_users: Arc<RwLock<HashSet<String>>>,
78    pub dialog_layer: Arc<DialogLayer>,
79    pub create_invitation_handler: Option<FnCreateInvitationHandler>,
80    pub invitation: Invitation,
81    pub routing_state: Arc<crate::call::RoutingState>,
82    pub pending_playbooks: Arc<Mutex<HashMap<String, (String, Instant)>>>,
83    pub learned_public_address: SharedPublicAddress,
84
85    pub active_calls: Arc<std::sync::Mutex<HashMap<String, ActiveCallRef>>>,
86    pub total_calls: AtomicU64,
87    pub total_failed_calls: AtomicU64,
88    pub uptime: DateTime<Local>,
89    pub shutting_down: Arc<AtomicBool>,
90}
91
92pub type AppState = Arc<AppStateInner>;
93
94pub struct AppStateBuilder {
95    pub config: Option<Config>,
96    pub stream_engine: Option<Arc<StreamEngine>>,
97    pub callrecord_sender: Option<CallRecordSender>,
98    pub callrecord_formatter: Option<Arc<dyn CallRecordFormatter>>,
99    pub cancel_token: Option<CancellationToken>,
100    pub create_invitation_handler: Option<FnCreateInvitationHandler>,
101    pub config_path: Option<String>,
102
103    pub message_inspector: Option<Box<dyn MessageInspector>>,
104    pub target_locator: Option<Box<dyn TargetLocator>>,
105    pub transport_inspector: Option<Box<dyn TransportEventInspector>>,
106}
107
108impl AppStateInner {
109    pub fn auto_learn_public_address_enabled(&self) -> bool {
110        self.config.auto_learn_public_address.unwrap_or(false)
111    }
112
113    pub fn get_dump_events_file(&self, session_id: &String) -> String {
114        let recorder_root = self.config.recorder_path();
115        let root = Path::new(&recorder_root);
116        if !root.exists() {
117            match std::fs::create_dir_all(root) {
118                Ok(_) => {
119                    info!("created dump events root: {}", root.to_string_lossy());
120                }
121                Err(e) => {
122                    warn!(
123                        "Failed to create dump events root: {} {}",
124                        e,
125                        root.to_string_lossy()
126                    );
127                }
128            }
129        }
130        root.join(format!("{}.events.jsonl", session_id))
131            .to_string_lossy()
132            .to_string()
133    }
134
135    pub fn get_recorder_file(&self, session_id: &String) -> String {
136        let recorder_root = self.config.recorder_path();
137        let root = Path::new(&recorder_root);
138        if !root.exists() {
139            match std::fs::create_dir_all(root) {
140                Ok(_) => {
141                    info!("created recorder root: {}", root.to_string_lossy());
142                }
143                Err(e) => {
144                    warn!(
145                        "Failed to create recorder root: {} {}",
146                        e,
147                        root.to_string_lossy()
148                    );
149                }
150            }
151        }
152        let desired_ext = self.config.recorder_format().extension();
153        let mut filename = session_id.clone();
154        if !filename
155            .to_lowercase()
156            .ends_with(&format!(".{}", desired_ext.to_lowercase()))
157        {
158            filename = format!("{}.{}", filename, desired_ext);
159        }
160        root.join(filename).to_string_lossy().to_string()
161    }
162
163    pub async fn serve(self: Arc<Self>) -> Result<()> {
164        let incoming_txs = self.endpoint.incoming_transactions()?;
165        let token = self.token.child_token();
166        let endpoint_inner = self.endpoint.inner.clone();
167        let dialog_layer = self.dialog_layer.clone();
168        let app_state_clone = self.clone();
169
170        match self.start_registration().await {
171            Ok(count) => {
172                info!("registration started, count: {}", count);
173            }
174            Err(e) => {
175                warn!("failed to start registration: {:?}", e);
176            }
177        }
178
179        let pending_cleanup_state = self.clone();
180        let pending_cleanup_token = token.clone();
181        crate::spawn(async move {
182            let mut interval = tokio::time::interval(Duration::from_secs(60));
183            let ttl = Duration::from_secs(300);
184            loop {
185                tokio::select! {
186                    _ = pending_cleanup_token.cancelled() => break,
187                    _ = interval.tick() => {
188                        let mut pending = pending_cleanup_state.pending_playbooks.lock().await;
189                        let before = pending.len();
190                        pending.retain(|_, (_, created_at)| created_at.elapsed() < ttl);
191                        let removed = before - pending.len();
192                        if removed > 0 {
193                            info!(removed, remaining = pending.len(), "cleaned up stale pending_playbooks entries");
194                        }
195                    }
196                }
197            }
198        });
199
200        tokio::select! {
201            _ = token.cancelled() => {
202                info!("cancelled");
203            }
204            result = endpoint_inner.serve() => {
205                if let Err(e) = result {
206                    info!("endpoint serve error: {:?}", e);
207                }
208            }
209            result = app_state_clone.process_incoming_request(dialog_layer.clone(), incoming_txs) => {
210                if let Err(e) = result {
211                    info!("process incoming request error: {:?}", e);
212                }
213            },
214        }
215
216        let total_secs = self.config.graceful_shutdown_timeout.unwrap_or(30);
217        let timeout = self
218            .config
219            .graceful_shutdown
220            .map(|_| Duration::from_secs(total_secs));
221
222        match self.stop_registration(timeout).await {
223            Ok(_) => {
224                info!("registration stopped, waiting for clear");
225            }
226            Err(e) => {
227                warn!("failed to stop registration: {:?}", e);
228            }
229        }
230        info!("stopping");
231        Ok(())
232    }
233
234    async fn process_incoming_request(
235        self: Arc<Self>,
236        dialog_layer: Arc<DialogLayer>,
237        mut incoming: TransactionReceiver,
238    ) -> Result<()> {
239        while let Some(mut tx) = incoming.recv().await {
240            let key: &rsipstack::transaction::key::TransactionKey = &tx.key;
241            info!(?key, "received transaction");
242            if tx.original.to_header()?.tag()?.as_ref().is_some() {
243                match dialog_layer.match_dialog(&tx) {
244                    Some(mut d) => {
245                        crate::spawn(async move {
246                            match d.handle(&mut tx).await {
247                                Ok(_) => (),
248                                Err(e) => {
249                                    info!("error handling transaction: {:?}", e);
250                                }
251                            }
252                        });
253                        continue;
254                    }
255                    None => {
256                        info!("dialog not found: {}", tx.original);
257                        match tx
258                            .reply(rsipstack::rsip::StatusCode::CallTransactionDoesNotExist)
259                            .await
260                        {
261                            Ok(_) => (),
262                            Err(e) => {
263                                info!("error replying to request: {:?}", e);
264                            }
265                        }
266                        continue;
267                    }
268                }
269            }
270            // out dialog, new server dialog
271            let (state_sender, state_receiver) = dialog_layer.new_dialog_state_channel();
272            match tx.original.method {
273                rsipstack::rsip::Method::Invite | rsipstack::rsip::Method::Ack => {
274                    // Reject new INVITEs during graceful shutdown
275                    if self.shutting_down.load(Ordering::Relaxed) {
276                        info!(?key, "rejecting INVITE during graceful shutdown");
277                        match tx
278                            .reply_with(
279                                rsipstack::rsip::StatusCode::ServiceUnavailable,
280                                vec![rsipstack::rsip::Header::Other(
281                                    "Reason".into(),
282                                    "SIP;cause=503;text=\"Server shutting down\"".into(),
283                                )],
284                                None,
285                            )
286                            .await
287                        {
288                            Ok(_) => (),
289                            Err(e) => {
290                                info!("error replying to request: {:?}", e);
291                            }
292                        }
293                        continue;
294                    }
295
296                    let invitation_handler = match self.create_invitation_handler {
297                        Some(ref create_invitation_handler) => {
298                            create_invitation_handler(self.config.handler.as_ref()).ok()
299                        }
300                        _ => default_create_invite_handler(
301                            self.config.handler.as_ref(),
302                            Some(self.clone()),
303                        ),
304                    };
305                    let invitation_handler = match invitation_handler {
306                        Some(h) => h,
307                        None => {
308                            info!(?key, "no invite handler configured, rejecting INVITE");
309                            match tx
310                                .reply_with(
311                                    rsipstack::rsip::StatusCode::ServiceUnavailable,
312                                    vec![rsipstack::rsip::Header::Other(
313                                        "Reason".into(),
314                                        "SIP;cause=503;text=\"No invite handler configured\""
315                                            .into(),
316                                    )],
317                                    None,
318                                )
319                                .await
320                            {
321                                Ok(_) => (),
322                                Err(e) => {
323                                    info!("error replying to request: {:?}", e);
324                                }
325                            }
326                            continue;
327                        }
328                    };
329                    let local_addr = tx
330                        .connection
331                        .as_ref()
332                        .map(|connection| connection.get_addr().clone())
333                        .or_else(|| dialog_layer.endpoint.get_addrs().first().cloned());
334                    let contact_username =
335                        tx.original.uri.auth.as_ref().map(|auth| auth.user.as_str());
336                    let contact = local_addr.as_ref().map(|addr| {
337                        build_public_contact_uri(
338                            &self.learned_public_address,
339                            self.auto_learn_public_address_enabled(),
340                            addr,
341                            contact_username,
342                            None,
343                        )
344                    });
345
346                    let dialog = match dialog_layer.get_or_create_server_invite(
347                        &tx,
348                        state_sender,
349                        None,
350                        contact,
351                    ) {
352                        Ok(d) => d,
353                        Err(e) => {
354                            // 481 Dialog/Transaction Does Not Exist
355                            info!("failed to obtain dialog: {:?}", e);
356                            match tx
357                                .reply(rsipstack::rsip::StatusCode::CallTransactionDoesNotExist)
358                                .await
359                            {
360                                Ok(_) => (),
361                                Err(e) => {
362                                    info!("error replying to request: {:?}", e);
363                                }
364                            }
365                            continue;
366                        }
367                    };
368
369                    let dialog_id = dialog.id();
370                    let dialog_id_str = dialog_id.to_string();
371                    // Incoming calls get a short public session id instead of
372                    // the raw dialog-id string; the guard registers the mapping
373                    // below so accept/hangup/message lookups can resolve it.
374                    let session_id = generate_short_session_id(&self.invitation);
375                    let dialog_id_for_cleanup = dialog_id.clone();
376                    let token = self.token.child_token();
377                    let pending_dialog = PendingDialog {
378                        token: token.clone(),
379                        dialog: dialog.clone(),
380                        state_receiver,
381                    };
382
383                    let guard = Arc::new(PendingDialogGuard::new_with_session(
384                        self.invitation.clone(),
385                        dialog_id,
386                        session_id.clone(),
387                        pending_dialog,
388                    ));
389
390                    let accept_timeout = self
391                        .config
392                        .accept_timeout
393                        .as_ref()
394                        .and_then(|t| parse_duration(t).ok())
395                        .unwrap_or_else(|| Duration::from_secs(60));
396
397                    let mut dialog_ref = dialog.clone();
398                    let routing_state = self.routing_state.clone();
399                    let dialog_for_reject = dialog.clone();
400                    let invitation_for_cleanup = self.invitation.clone();
401                    let session_id_for_task = session_id.clone();
402                    crate::spawn(async move {
403                        info!(session_id = session_id_for_task, id = dialog_id_str, "incoming invite task started");
404                        let _pending_guard = guard;
405                        let token_ref = token.clone();
406                        let accept_timeout_sleep = tokio::time::sleep(accept_timeout);
407                        let invite_handler = invitation_handler.on_invite(
408                            session_id_for_task.clone(),
409                            token.clone(),
410                            dialog.clone(),
411                            routing_state,
412                        );
413                        let dialog_handle = dialog_ref.handle(&mut tx);
414                        tokio::pin!(accept_timeout_sleep);
415                        tokio::pin!(invite_handler);
416                        tokio::pin!(dialog_handle);
417
418                        let mut cancel_done = false;
419                        let mut accept_timeout_done = false;
420                        let mut invite_done = false;
421                        loop {
422                            let mut reject_request = None;
423
424                            select! {
425                                _ = token_ref.cancelled(), if !cancel_done => {
426                                    cancel_done = true;
427                                    reject_request = Some((
428                                        rsipstack::rsip::StatusCode::ServiceUnavailable,
429                                        "invite cancelled".to_string(),
430                                        "cancelled",
431                                    ));
432                                }
433                                _ = &mut accept_timeout_sleep, if !accept_timeout_done
434                                    && dialog_for_reject.state().can_cancel() => {
435                                    accept_timeout_done = true;
436                                    reject_request = Some((
437                                        rsipstack::rsip::StatusCode::RequestTimeout,
438                                        "accept timeout".to_string(),
439                                        "accept timeout",
440                                    ));
441                                }
442                                result = &mut invite_handler, if !invite_done => {
443                                    invite_done = true;
444                                    match result {
445                                        Ok(_) => {
446                                            info!(id = dialog_id_str, "invite handler completed");
447                                        }
448                                        Err(e) => {
449                                            info!(id = dialog_id_str, "error handling invite: {:?}", e);
450                                            reject_request = Some((
451                                                rsipstack::rsip::StatusCode::ServiceUnavailable,
452                                                format!("Failed to process invite: {}", e),
453                                                "invite handler error",
454                                            ));
455                                        }
456                                    }
457                                }
458                                result = &mut dialog_handle => {
459                                    match result {
460                                        Ok(_) => {
461                                            info!(id = dialog_id_str, "dialog handling finished");
462                                        }
463                                        Err(e) => {
464                                            info!(
465                                                id = dialog_id_str,
466                                                "dialog handling ended with error: {:?}", e
467                                            );
468                                        }
469                                    }
470                                    if matches!(
471                                        dialog_for_reject.state(),
472                                        rsipstack::dialog::dialog::DialogState::Terminated(_, _)
473                                    ) {
474                                        info!(
475                                            id = dialog_id_str,
476                                            "terminated invite dialog finished, cancelling invite token"
477                                        );
478                                        token_ref.cancel();
479                                    }
480                                    info!(id = dialog_id_str, "incoming invite task finished");
481                                    break;
482                                }
483                            }
484
485                            if let Some((code, reason, source)) = reject_request {
486                                if dialog_for_reject.state().can_cancel() {
487                                    info!(
488                                        id = dialog_id_str,
489                                        ?code,
490                                        %reason,
491                                        source,
492                                        "rejecting invite"
493                                    );
494                                    if let Err(e) =
495                                        dialog_for_reject.reject(Some(code), Some(reason))
496                                    {
497                                        info!(
498                                            id = dialog_id_str,
499                                            "error rejecting invite: {:?}", e
500                                        );
501                                    }
502                                    invitation_for_cleanup.get_pending_call(&dialog_id_for_cleanup);
503                                    invitation_for_cleanup
504                                        .dialog_layer
505                                        .remove_dialog(&dialog_id_for_cleanup);
506                                }
507                            }
508                        }
509                    });
510                }
511                rsipstack::rsip::Method::Options => {
512                    if self.config.enable_options_response.unwrap_or(true) {
513                        info!(?key, "responding to out-of-dialog OPTIONS request");
514                        let allow_header: rsipstack::rsip::Header =
515                            typed::Allow::from(Method::all()).into();
516                        let accept_header =
517                            rsipstack::rsip::Header::Accept(Accept::new("application/sdp"));
518                        match tx
519                            .reply_with(
520                                rsipstack::rsip::StatusCode::OK,
521                                vec![allow_header, accept_header],
522                                None,
523                            )
524                            .await
525                        {
526                            Ok(_) => (),
527                            Err(e) => {
528                                info!("error replying to OPTIONS: {:?}", e);
529                            }
530                        }
531                    } else {
532                        info!(?key, "ignoring out-of-dialog OPTIONS request");
533                    }
534                    continue;
535                }
536                rsipstack::rsip::Method::Refer => {
537                    info!(?key, "ignoring out-of-dialog REFER");
538                    match tx.reply(rsipstack::rsip::StatusCode::BadRequest).await {
539                        Ok(_) => (),
540                        Err(e) => {
541                            info!("error replying to out-of-dialog REFER: {:?}", e);
542                        }
543                    }
544                    continue;
545                }
546                _ => {
547                    info!(?key, "received request: {:?}", tx.original.method);
548                    match tx.reply(rsipstack::rsip::StatusCode::OK).await {
549                        Ok(_) => (),
550                        Err(e) => {
551                            info!("error replying to request: {:?}", e);
552                        }
553                    }
554                }
555            }
556        }
557        Ok(())
558    }
559
560    pub fn stop(&self) {
561        if self.shutting_down.swap(true, Ordering::Relaxed) {
562            return;
563        }
564        info!("stopping, marking as shutting down");
565        self.token.cancel();
566    }
567
568    pub async fn graceful_stop(&self, total_timeout_secs: u64) -> Result<()> {
569        if self.shutting_down.swap(true, Ordering::Relaxed) {
570            return Ok(());
571        }
572
573        info!("graceful stopping, marking as shutting down");
574        let timeout = Duration::from_secs(total_timeout_secs);
575
576        let (reg_result, ()) = tokio::join!(
577            self.stop_registration(Some(timeout)),
578            self.wait_for_active_calls(timeout),
579        );
580        if let Err(e) = reg_result {
581            warn!("stop_registration error: {}", e);
582        }
583        self.token.cancel();
584        Ok(())
585    }
586
587    async fn wait_for_active_calls(&self, timeout: Duration) {
588        let active_calls = self.active_calls.clone();
589        let check_loop = async move {
590            let mut last_count = usize::MAX;
591            loop {
592                let count = active_calls.lock().unwrap().len();
593                if count == 0 {
594                    break;
595                }
596                if count != last_count {
597                    info!(active_calls = count, "waiting for active calls to finish");
598                    last_count = count;
599                }
600                tokio::time::sleep(Duration::from_millis(500)).await;
601            }
602        };
603        match tokio::time::timeout(timeout, check_loop).await {
604            Ok(()) => info!("all active calls finished"),
605            Err(_) => warn!("timed out waiting for active calls to finish, forcing shutdown"),
606        }
607    }
608
609    pub async fn start_registration(&self) -> Result<usize> {
610        let mut count = 0;
611        if let Some(register_users) = &self.config.register_users {
612            for option in register_users.iter() {
613                match self.register(option.clone()).await {
614                    Ok(_) => {
615                        count += 1;
616                    }
617                    Err(e) => {
618                        warn!("failed to register user: {:?} {:?}", e, option);
619                    }
620                }
621            }
622        }
623        Ok(count)
624    }
625
626    pub fn find_credentials_for_callee(&self, callee: &str) -> Option<UserCredential> {
627        let callee_uri = crate::sip_util::ensure_sip_scheme(callee.to_string());
628
629        let parsed_callee = match rsipstack::rsip::Uri::try_from(callee_uri.as_str()) {
630            Ok(uri) => uri,
631            Err(e) => {
632                warn!("failed to parse callee URI: {} {:?}", callee, e);
633                return None;
634            }
635        };
636
637        // Matching is host-kind-homogeneous: a domain callee matches domain
638        // servers, an IP callee matches IP servers.
639        match &parsed_callee.host_with_port.host {
640            rsipstack::rsip::Host::Domain(domain) => {
641                let domain = domain.0.clone();
642                self.find_credentials(
643                    callee,
644                    &format!("domain {}", domain),
645                    |host| matches!(host, rsipstack::rsip::Host::Domain(d) if d.0 == domain),
646                )
647            }
648            rsipstack::rsip::Host::IpAddr(ip) => {
649                let ip = *ip;
650                self.find_credentials(callee, "IP match", move |host| {
651                    matches!(host, rsipstack::rsip::Host::IpAddr(server_ip) if *server_ip == ip)
652                })
653            }
654        }
655    }
656
657    /// Scan registered users for the first server matching `pred` that has a
658    /// credential.
659    fn find_credentials(
660        &self,
661        callee: &str,
662        match_type: &str,
663        pred: impl Fn(&rsipstack::rsip::Host) -> bool,
664    ) -> Option<UserCredential> {
665        let register_users = self.config.register_users.as_ref()?;
666        for option in register_users.iter() {
667            let server = crate::sip_util::ensure_sip_scheme(option.server.clone());
668
669            let parsed_server = match rsipstack::rsip::Uri::try_from(server.as_str()) {
670                Ok(uri) => uri,
671                Err(e) => {
672                    warn!("failed to parse server URI: {} {:?}", option.server, e);
673                    continue;
674                }
675            };
676
677            if pred(&parsed_server.host_with_port.host) {
678                if let Some(cred) = &option.credential {
679                    info!(
680                        callee,
681                        username = cred.username,
682                        server = option.server,
683                        match_type,
684                        "Auto-injecting credentials from registered user for outbound call"
685                    );
686                    return Some(cred.clone());
687                }
688            }
689        }
690        None
691    }
692
693    pub async fn stop_registration(&self, wait_for_clear: Option<Duration>) -> Result<()> {
694        {
695            let mut handles = self.registration_handles.lock().await;
696            for (_, cancel_token) in handles.drain() {
697                cancel_token.cancel();
698            }
699        }
700
701        if let Some(duration) = wait_for_clear {
702            let live_users = self.alive_users.clone();
703            let check_loop = async move {
704                loop {
705                    let is_empty = {
706                        let users = live_users
707                            .read()
708                            .map_err(|_| anyhow::anyhow!("Lock poisoned"))?;
709                        users.is_empty()
710                    };
711                    if is_empty {
712                        break;
713                    }
714                    tokio::time::sleep(Duration::from_millis(50)).await;
715                }
716                Ok::<(), anyhow::Error>(())
717            };
718            match tokio::time::timeout(duration, check_loop).await {
719                Ok(_) => {}
720                Err(e) => {
721                    warn!("failed to wait for clear: {}", e);
722                    return Err(anyhow::anyhow!("failed to wait for clear: {}", e));
723                }
724            }
725        }
726        Ok(())
727    }
728
729    pub async fn register(&self, option: RegisterOption) -> Result<()> {
730        let user = option.aor();
731        let server = crate::sip_util::ensure_sip_scheme(option.server.clone());
732        let sip_server = match rsipstack::rsip::Uri::try_from(server) {
733            Ok(uri) => uri,
734            Err(e) => {
735                warn!("failed to parse server: {} {:?}", e, option.server);
736                return Err(anyhow::anyhow!("failed to parse server: {}", e));
737            }
738        };
739        let cancel_token = self.token.child_token();
740        let credential = option.credential.clone().map(|c| c.into());
741        let registration = rsipstack::dialog::registration::Registration::new(
742            self.endpoint.inner.clone(),
743            credential,
744        );
745        let mut handle = RegistrationHandle {
746            registration,
747            option,
748            cancel_token: cancel_token.clone(),
749            start_time: Instant::now(),
750            last_update: Instant::now(),
751            last_response: None,
752        };
753        self.registration_handles
754            .lock()
755            .await
756            .insert(user.clone(), cancel_token);
757        tracing::debug!(user = user.as_str(), "starting registration task");
758        let alive_users = self.alive_users.clone();
759
760        crate::spawn(async move {
761            handle.start_time = Instant::now();
762            let cancel_token = handle.cancel_token.clone();
763            let addrs = handle.registration.endpoint.get_addrs();
764            let local_bind_addr = if let Some(addr) = find_local_addr_for_uri(&addrs, &sip_server) {
765                addr
766            } else {
767                warn!(
768                    user = user.as_str(),
769                    server = %sip_server,
770                    "failed to get local bind address for registration transport"
771                );
772                alive_users.write().unwrap().remove(&user);
773                return;
774            };
775            let user = handle.option.aor();
776            alive_users.write().unwrap().remove(&user);
777            let mut contact_address = local_bind_addr.addr.clone();
778            let mut contact = build_contact(
779                &local_bind_addr,
780                Some(contact_address.clone()),
781                Some(handle.option.username.as_str()),
782                None,
783            );
784            let mut should_register = true;
785            let mut timer = pending().boxed();
786
787            loop {
788                select! {
789                    _ = cancel_token.cancelled() => {
790                        break;
791                    }
792                    _ = timer.as_mut(), if !should_register => {
793                        should_register = true;
794                        timer = Box::pin(pending());
795                    }
796                    result = handle.do_register(&sip_server, None, &contact), if should_register => {
797                        match result {
798                            Ok((expires, new_addr)) => {
799                                if handle
800                                    .should_retry_registration_now(
801                                        &local_bind_addr,
802                                        &contact_address,
803                                        new_addr.as_ref(),
804                                    )
805                                {
806                                    if let Some(next_contact_address) = new_addr {
807                                        info!(
808                                            user = user.as_str(),
809                                            current_contact = %contact_address,
810                                            next_contact = %next_contact_address,
811                                            "public address changed, retrying registration immediately",
812                                        );
813                                        contact_address = next_contact_address;
814                                        contact = build_contact(
815                                            &local_bind_addr,
816                                            Some(contact_address.clone()),
817                                            Some(handle.option.username.as_str()),
818                                            None,
819                                        );
820                                        continue;
821                                    }
822                                }
823                                info!(
824                                    user = user.as_str(),
825                                    expires = expires,
826                                    contact = %contact_address,
827                                    alive_users = alive_users.read().unwrap().len(),
828                                    "registration refreshed",
829                                );
830                                alive_users.write().unwrap().insert(user.clone());
831                                should_register = false;
832                                timer = Box::pin(tokio::time::sleep(Duration::from_secs(
833                                    (expires * 3 / 4) as u64,
834                                )));
835                            }
836                            Err(e) => {
837                                warn!(
838                                    user = user.as_str(),
839                                    alive_users = alive_users.read().unwrap().len(),
840                                    "registration failed: {:?}", e
841                                );
842                                should_register = false;
843                                timer = Box::pin(tokio::time::sleep(Duration::from_secs(60)));
844                            }
845                        }
846                    }
847                }
848            }
849            handle
850                .do_register(&sip_server, Some(0), &contact)
851                .await
852                .ok();
853            alive_users.write().unwrap().remove(&user);
854        });
855        Ok(())
856    }
857}
858
859impl Drop for AppStateInner {
860    fn drop(&mut self) {
861        self.stop();
862    }
863}
864
865struct SimpleDomainResolver;
866
867#[async_trait]
868impl DomainResolver for SimpleDomainResolver {
869    async fn resolve(&self, target: &SipAddr) -> rsipstack::Result<SipAddr> {
870        match &target.addr.host {
871            Host::Domain(domain) => {
872                let port: u16 = target.addr.port.map(|p| p.value()).unwrap_or(5060);
873                let addr_str = format!("{}:{}", domain, port);
874                match tokio::net::lookup_host(&addr_str).await {
875                    Ok(mut addrs) => {
876                        if let Some(addr) = addrs.next() {
877                            Ok(SipAddr {
878                                r#type: target.r#type,
879                                addr: HostWithPort {
880                                    host: Host::IpAddr(addr.ip()),
881                                    port: Some(rsipstack::rsip::Port(addr.port())),
882                                },
883                            })
884                        } else {
885                            Err(rsipstack::Error::DnsResolutionError(format!(
886                                "no addresses found for {}",
887                                domain
888                            )))
889                        }
890                    }
891                    Err(e) => Err(rsipstack::Error::DnsResolutionError(format!(
892                        "DNS resolution failed for {}: {}",
893                        domain, e
894                    ))),
895                }
896            }
897            _ => Ok(target.clone()),
898        }
899    }
900}
901
902impl AppStateBuilder {
903    pub fn new() -> Self {
904        Self {
905            config: None,
906            stream_engine: None,
907            callrecord_sender: None,
908            callrecord_formatter: None,
909            cancel_token: None,
910            create_invitation_handler: None,
911            config_path: None,
912            message_inspector: None,
913            target_locator: None,
914            transport_inspector: None,
915        }
916    }
917
918    pub fn with_config(mut self, config: Config) -> Self {
919        self.config = Some(config);
920        self
921    }
922
923    pub fn with_stream_engine(mut self, stream_engine: Arc<StreamEngine>) -> Self {
924        self.stream_engine = Some(stream_engine);
925        self
926    }
927
928    pub fn with_callrecord_sender(mut self, sender: CallRecordSender) -> Self {
929        self.callrecord_sender = Some(sender);
930        self
931    }
932
933    pub fn with_cancel_token(mut self, token: CancellationToken) -> Self {
934        self.cancel_token = Some(token);
935        self
936    }
937
938    pub fn with_config_metadata(mut self, path: Option<String>) -> Self {
939        self.config_path = path;
940        self
941    }
942
943    pub fn with_inspector(&mut self, inspector: Box<dyn MessageInspector>) -> &mut Self {
944        self.message_inspector = Some(inspector);
945        self
946    }
947    pub fn with_target_locator(&mut self, locator: Box<dyn TargetLocator>) -> &mut Self {
948        self.target_locator = Some(locator);
949        self
950    }
951
952    pub fn with_transport_inspector(
953        &mut self,
954        inspector: Box<dyn TransportEventInspector>,
955    ) -> &mut Self {
956        self.transport_inspector = Some(inspector);
957        self
958    }
959
960    pub async fn build(self) -> Result<AppState> {
961        let config: Arc<Config> = Arc::new(self.config.unwrap_or_default());
962        let token = self
963            .cancel_token
964            .unwrap_or_else(|| CancellationToken::new());
965        let _ = set_cache_dir(&config.media_cache_path);
966        let local_ip = if !config.addr.is_empty() {
967            std::net::IpAddr::from_str(config.addr.as_str())?
968        } else {
969            crate::net_tool::get_first_non_loopback_interface()?
970        };
971        let transport_layer = match catch_unwind(AssertUnwindSafe(|| {
972            rsipstack::transport::TransportLayer::new(token.clone())
973        })) {
974            Ok(tl) => tl,
975            Err(_) => {
976                warn!(
977                    "failed to initialize default DNS resolver with hickory-resolver, falling back to simple resolver via tokio::net::lookup_host"
978                );
979                rsipstack::transport::TransportLayer::new_with_domain_resolver(
980                    token.clone(),
981                    Box::new(SimpleDomainResolver),
982                )
983            }
984        };
985        let local_addr: SocketAddr = format!("{}:{}", local_ip, config.udp_port).parse()?;
986
987        // Create UDP socket with SO_REUSEPORT for graceful restarts
988        #[cfg(unix)]
989        let std_socket = {
990            use socket2::{Domain, Protocol, SockAddr, Socket, Type};
991
992            let domain = if local_addr.is_ipv4() {
993                Domain::IPV4
994            } else {
995                Domain::IPV6
996            };
997            let socket = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP))
998                .map_err(|err| anyhow::anyhow!("Failed to create UDP socket: {}", err))?;
999
1000            socket
1001                .set_reuse_address(true)
1002                .map_err(|err| anyhow::anyhow!("Failed to set SO_REUSEADDR: {}", err))?;
1003
1004            // SO_REUSEPORT (Linux/BSD)
1005            #[cfg(not(any(target_os = "solaris", target_os = "illumos", target_os = "cygwin")))]
1006            socket
1007                .set_reuse_port(true)
1008                .map_err(|err| anyhow::anyhow!("Failed to set SO_REUSEPORT: {}", err))?;
1009
1010            socket
1011                .bind(&SockAddr::from(local_addr))
1012                .map_err(|err| anyhow::anyhow!("Failed to bind UDP socket: {}", err))?;
1013
1014            let std_socket: std::net::UdpSocket = socket.into();
1015            std_socket
1016        };
1017
1018        #[cfg(not(unix))]
1019        let std_socket = std::net::UdpSocket::bind(local_addr)?;
1020
1021        std_socket.set_nonblocking(true)?;
1022        let tokio_socket = tokio::net::UdpSocket::from_std(std_socket)?;
1023        // Use the actual bound address (important when port=0 lets OS assign a port)
1024        let actual_addr = tokio_socket.local_addr()?;
1025        let bind_addr = rsipstack::transport::SipConnection::resolve_bind_address(actual_addr);
1026        let mut learned_public_address: SharedPublicAddress =
1027            Arc::new(ArcSwap::from_pointee(bind_addr.into()));
1028
1029        let udp_inner = rsipstack::transport::udp::UdpInner {
1030            conn: tokio_socket,
1031            addr: rsipstack::transport::SipAddr {
1032                r#type: Some(rsipstack::rsip::transport::Transport::Udp),
1033                addr: bind_addr.into(),
1034            },
1035        };
1036
1037        let external = config
1038            .external_ip
1039            .as_ref()
1040            .map(|ip| {
1041                format!("{}:{}", ip, actual_addr.port())
1042                    .parse()
1043                    .map_err(|e| anyhow::anyhow!("Failed to parse external address: {}", e))
1044            })
1045            .transpose()?;
1046
1047        let udp_conn = rsipstack::transport::udp::UdpConnection::attach(
1048            udp_inner,
1049            external,
1050            Some(token.child_token()),
1051        )
1052        .await;
1053
1054        info!(
1055            "start useragent, addr: {} (SO_REUSEPORT enabled)",
1056            udp_conn.get_addr()
1057        );
1058
1059        transport_layer.add_transport(udp_conn.into());
1060
1061        // Optional SIP over TLS transport
1062        if let Some(tls_port) = config.tls_port {
1063            let tls_addr: std::net::SocketAddr = format!("{}:{}", local_ip, tls_port).parse()?;
1064            let mut tls_cfg = rsipstack::transport::tls::TlsConfig::default();
1065            if let Some(ref cert_path) = config.tls_cert_file {
1066                tls_cfg.cert = Some(
1067                    std::fs::read(cert_path)
1068                        .map_err(|e| anyhow::anyhow!("tls_cert_file: {}", e))?,
1069                );
1070            }
1071            if let Some(ref key_path) = config.tls_key_file {
1072                tls_cfg.key = Some(
1073                    std::fs::read(key_path).map_err(|e| anyhow::anyhow!("tls_key_file: {}", e))?,
1074                );
1075            }
1076            let external_tls_addr = config
1077                .external_ip
1078                .as_ref()
1079                .and_then(|ip| format!("{}:{}", ip, tls_port).parse().ok());
1080            match rsipstack::transport::tls::TlsListenerConnection::new(
1081                tls_addr,
1082                external_tls_addr,
1083                tls_cfg,
1084            )
1085            .await
1086            {
1087                Ok(tls_conn) => {
1088                    transport_layer.add_transport(tls_conn.into());
1089                    info!("TLS SIP transport started on {}:{}", local_ip, tls_port);
1090                }
1091                Err(e) => {
1092                    return Err(anyhow::anyhow!("Failed to start TLS SIP transport: {}", e));
1093                }
1094            }
1095        }
1096
1097        let endpoint_option = rsipstack::transaction::endpoint::EndpointOption::default();
1098        let mut endpoint_builder = rsipstack::EndpointBuilder::new();
1099        if let Some(ref user_agent) = config.useragent {
1100            endpoint_builder.with_user_agent(user_agent.as_str());
1101        }
1102
1103        let mut endpoint_builder = endpoint_builder
1104            .with_cancel_token(token.child_token())
1105            .with_transport_layer(transport_layer)
1106            .with_option(endpoint_option);
1107
1108        if config.auto_learn_public_address.unwrap_or_default() {
1109            let inspector = LearningMessageInspector::new(bind_addr.into(), self.message_inspector);
1110            learned_public_address = inspector.shared_public_address();
1111            endpoint_builder = endpoint_builder.with_inspector(Box::new(inspector));
1112        } else if let Some(inspector) = self.message_inspector {
1113            endpoint_builder = endpoint_builder.with_inspector(inspector);
1114        }
1115
1116        if let Some(locator) = self.target_locator {
1117            endpoint_builder.with_target_locator(locator);
1118        } else if let Some(ref rules) = config.rewrites {
1119            endpoint_builder
1120                .with_target_locator(Box::new(RewriteTargetLocator::new(rules.clone())));
1121        }
1122
1123        if let Some(inspector) = self.transport_inspector {
1124            endpoint_builder = endpoint_builder.with_transport_inspector(inspector);
1125        }
1126
1127        let endpoint = endpoint_builder.build();
1128        let dialog_layer = Arc::new(DialogLayer::new(endpoint.inner.clone()));
1129
1130        let stream_engine = self.stream_engine.unwrap_or_default();
1131
1132        let callrecord_formatter = if let Some(formatter) = self.callrecord_formatter {
1133            formatter
1134        } else {
1135            let formatter = if let Some(ref callrecord) = config.callrecord {
1136                DefaultCallRecordFormatter::new_with_config(callrecord)
1137            } else {
1138                DefaultCallRecordFormatter::default()
1139            };
1140            Arc::new(formatter)
1141        };
1142
1143        let callrecord_sender = if let Some(sender) = self.callrecord_sender {
1144            Some(sender)
1145        } else if let Some(ref callrecord) = config.callrecord {
1146            let builder = CallRecordManagerBuilder::new()
1147                .with_cancel_token(token.child_token())
1148                .with_config(callrecord.clone())
1149                .with_max_concurrent(32)
1150                .with_formatter(callrecord_formatter.clone());
1151
1152            let mut callrecord_manager = builder.build();
1153            let sender = callrecord_manager.sender.clone();
1154            crate::spawn(async move {
1155                callrecord_manager.serve().await;
1156            });
1157            Some(sender)
1158        } else {
1159            None
1160        };
1161
1162        let app_state = Arc::new(AppStateInner {
1163            config,
1164            token,
1165            stream_engine,
1166            callrecord_sender,
1167            endpoint,
1168            registration_handles: Mutex::new(HashMap::new()),
1169            alive_users: Arc::new(RwLock::new(HashSet::new())),
1170            dialog_layer: dialog_layer.clone(),
1171            create_invitation_handler: self.create_invitation_handler,
1172            invitation: Invitation::new(dialog_layer),
1173            routing_state: Arc::new(crate::call::RoutingState::new()),
1174            pending_playbooks: Arc::new(Mutex::new(HashMap::new())),
1175            learned_public_address,
1176            active_calls: Arc::new(std::sync::Mutex::new(HashMap::new())),
1177            total_calls: AtomicU64::new(0),
1178            total_failed_calls: AtomicU64::new(0),
1179            uptime: Local::now(),
1180            shutting_down: Arc::new(AtomicBool::new(false)),
1181        });
1182
1183        Ok(app_state)
1184    }
1185}