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