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