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