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