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