Skip to main content

active_call/
app.rs

1use crate::{
2    call::{ActiveCallRef, sip::Invitation},
3    callrecord::{
4        CallRecordFormatter, CallRecordManagerBuilder, CallRecordSender, DefaultCallRecordFormatter,
5    },
6    config::Config,
7    locator::RewriteTargetLocator,
8    useragent::{
9        RegisterOption,
10        invitation::{
11            FnCreateInvitationHandler, PendingDialog, PendingDialogGuard,
12            default_create_invite_handler,
13        },
14        public_address::{
15            LearningMessageInspector, SharedPublicAddress, build_contact, build_public_contact_uri,
16            find_local_addr_for_uri,
17        },
18        registration::{RegistrationHandle, UserCredential},
19    },
20};
21
22use crate::media::{cache::set_cache_dir, engine::StreamEngine};
23use anyhow::Result;
24use arc_swap::ArcSwap;
25use chrono::{DateTime, Local};
26use futures::FutureExt;
27use humantime::parse_duration;
28use rsipstack::rsip::prelude::HeadersExt;
29use rsipstack::rsip::{Accept, Method, typed};
30use rsipstack::transaction::{
31    Endpoint, TransactionReceiver,
32    endpoint::{TargetLocator, TransportEventInspector},
33};
34use rsipstack::{dialog::dialog_layer::DialogLayer, transaction::endpoint::MessageInspector};
35use 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        let total_secs = self.config.graceful_shutdown_timeout.unwrap_or(30);
197        let timeout = self
198            .config
199            .graceful_shutdown
200            .map(|_| Duration::from_secs(total_secs));
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.on_invite(
382                            dialog_id_str.clone(),
383                            token.clone(),
384                            dialog.clone(),
385                            routing_state,
386                        );
387                        let dialog_handle = dialog_ref.handle(&mut tx);
388                        tokio::pin!(accept_timeout_sleep);
389                        tokio::pin!(invite_handler);
390                        tokio::pin!(dialog_handle);
391
392                        let mut cancel_done = false;
393                        let mut accept_timeout_done = false;
394                        let mut invite_done = false;
395                        loop {
396                            let mut reject_request = None;
397
398                            select! {
399                                _ = token_ref.cancelled(), if !cancel_done => {
400                                    cancel_done = true;
401                                    reject_request = Some((
402                                        rsipstack::rsip::StatusCode::ServiceUnavailable,
403                                        "invite cancelled".to_string(),
404                                        "cancelled",
405                                    ));
406                                }
407                                _ = &mut accept_timeout_sleep, if !accept_timeout_done
408                                    && dialog_for_reject.state().can_cancel() => {
409                                    accept_timeout_done = true;
410                                    reject_request = Some((
411                                        rsipstack::rsip::StatusCode::RequestTimeout,
412                                        "accept timeout".to_string(),
413                                        "accept timeout",
414                                    ));
415                                }
416                                result = &mut invite_handler, if !invite_done => {
417                                    invite_done = true;
418                                    match result {
419                                        Ok(_) => {
420                                            info!(id = dialog_id_str, "invite handler completed");
421                                        }
422                                        Err(e) => {
423                                            info!(id = dialog_id_str, "error handling invite: {:?}", e);
424                                            reject_request = Some((
425                                                rsipstack::rsip::StatusCode::ServiceUnavailable,
426                                                format!("Failed to process invite: {}", e),
427                                                "invite handler error",
428                                            ));
429                                        }
430                                    }
431                                }
432                                result = &mut dialog_handle => {
433                                    match result {
434                                        Ok(_) => {
435                                            info!(id = dialog_id_str, "dialog handling finished");
436                                        }
437                                        Err(e) => {
438                                            info!(
439                                                id = dialog_id_str,
440                                                "dialog handling ended with error: {:?}", e
441                                            );
442                                        }
443                                    }
444                                    if matches!(
445                                        dialog_for_reject.state(),
446                                        rsipstack::dialog::dialog::DialogState::Terminated(_, _)
447                                    ) {
448                                        info!(
449                                            id = dialog_id_str,
450                                            "terminated invite dialog finished, cancelling invite token"
451                                        );
452                                        token_ref.cancel();
453                                    }
454                                    info!(id = dialog_id_str, "incoming invite task finished");
455                                    break;
456                                }
457                            }
458
459                            if let Some((code, reason, source)) = reject_request {
460                                if dialog_for_reject.state().can_cancel() {
461                                    info!(
462                                        id = dialog_id_str,
463                                        ?code,
464                                        %reason,
465                                        source,
466                                        "rejecting invite"
467                                    );
468                                    if let Err(e) =
469                                        dialog_for_reject.reject(Some(code), Some(reason))
470                                    {
471                                        info!(
472                                            id = dialog_id_str,
473                                            "error rejecting invite: {:?}", e
474                                        );
475                                    }
476                                    invitation_for_cleanup.get_pending_call(&dialog_id_for_cleanup);
477                                    invitation_for_cleanup
478                                        .dialog_layer
479                                        .remove_dialog(&dialog_id_for_cleanup);
480                                }
481                            }
482                        }
483                    });
484                }
485                rsipstack::rsip::Method::Options => {
486                    if self.config.enable_options_response.unwrap_or(true) {
487                        info!(?key, "responding to out-of-dialog OPTIONS request");
488                        let allow_header: rsipstack::rsip::Header =
489                            typed::Allow::from(Method::all()).into();
490                        let accept_header =
491                            rsipstack::rsip::Header::Accept(Accept::new("application/sdp"));
492                        match tx
493                            .reply_with(
494                                rsipstack::rsip::StatusCode::OK,
495                                vec![allow_header, accept_header],
496                                None,
497                            )
498                            .await
499                        {
500                            Ok(_) => (),
501                            Err(e) => {
502                                info!("error replying to OPTIONS: {:?}", e);
503                            }
504                        }
505                    } else {
506                        info!(?key, "ignoring out-of-dialog OPTIONS request");
507                    }
508                    continue;
509                }
510                rsipstack::rsip::Method::Refer => {
511                    info!(?key, "ignoring out-of-dialog REFER");
512                    match tx.reply(rsipstack::rsip::StatusCode::BadRequest).await {
513                        Ok(_) => (),
514                        Err(e) => {
515                            info!("error replying to out-of-dialog REFER: {:?}", e);
516                        }
517                    }
518                    continue;
519                }
520                _ => {
521                    info!(?key, "received request: {:?}", tx.original.method);
522                    match tx.reply(rsipstack::rsip::StatusCode::OK).await {
523                        Ok(_) => (),
524                        Err(e) => {
525                            info!("error replying to request: {:?}", e);
526                        }
527                    }
528                }
529            }
530        }
531        Ok(())
532    }
533
534    pub fn stop(&self) {
535        if self.shutting_down.swap(true, Ordering::Relaxed) {
536            return;
537        }
538        info!("stopping, marking as shutting down");
539        self.token.cancel();
540    }
541
542    pub async fn graceful_stop(&self, total_timeout_secs: u64) -> Result<()> {
543        if self.shutting_down.swap(true, Ordering::Relaxed) {
544            return Ok(());
545        }
546
547        info!("graceful stopping, marking as shutting down");
548        let timeout = Duration::from_secs(total_timeout_secs);
549
550        let (reg_result, ()) = tokio::join!(
551            self.stop_registration(Some(timeout)),
552            self.wait_for_active_calls(timeout),
553        );
554        if let Err(e) = reg_result {
555            warn!("stop_registration error: {}", e);
556        }
557        self.token.cancel();
558        Ok(())
559    }
560
561    async fn wait_for_active_calls(&self, timeout: Duration) {
562        let active_calls = self.active_calls.clone();
563        let check_loop = async move {
564            let mut last_count = usize::MAX;
565            loop {
566                let count = active_calls.lock().unwrap().len();
567                if count == 0 {
568                    break;
569                }
570                if count != last_count {
571                    info!(active_calls = count, "waiting for active calls to finish");
572                    last_count = count;
573                }
574                tokio::time::sleep(Duration::from_millis(500)).await;
575            }
576        };
577        match tokio::time::timeout(timeout, check_loop).await {
578            Ok(()) => info!("all active calls finished"),
579            Err(_) => warn!("timed out waiting for active calls to finish, forcing shutdown"),
580        }
581    }
582
583    pub async fn start_registration(&self) -> Result<usize> {
584        let mut count = 0;
585        if let Some(register_users) = &self.config.register_users {
586            for option in register_users.iter() {
587                match self.register(option.clone()).await {
588                    Ok(_) => {
589                        count += 1;
590                    }
591                    Err(e) => {
592                        warn!("failed to register user: {:?} {:?}", e, option);
593                    }
594                }
595            }
596        }
597        Ok(count)
598    }
599
600    pub fn find_credentials_for_callee(&self, callee: &str) -> Option<UserCredential> {
601        let callee_uri = callee
602            .strip_prefix("sip:")
603            .or_else(|| callee.strip_prefix("sips:"))
604            .unwrap_or(callee);
605        let callee_uri = if !callee_uri.starts_with("sip:") && !callee_uri.starts_with("sips:") {
606            format!("sip:{}", callee_uri)
607        } else {
608            callee_uri.to_string()
609        };
610
611        let parsed_callee = match rsipstack::rsip::Uri::try_from(callee_uri.as_str()) {
612            Ok(uri) => uri,
613            Err(e) => {
614                warn!("failed to parse callee URI: {} {:?}", callee, e);
615                return None;
616            }
617        };
618
619        let callee_host = match &parsed_callee.host_with_port.host {
620            rsipstack::rsip::Host::Domain(domain) => domain.to_string(),
621            rsipstack::rsip::Host::IpAddr(ip) => return self.find_credentials_by_ip(ip),
622        };
623
624        // Look through registered users to find one matching this domain
625        if let Some(register_users) = &self.config.register_users {
626            for option in register_users.iter() {
627                let mut server = option.server.clone();
628                if !server.starts_with("sip:") && !server.starts_with("sips:") {
629                    server = format!("sip:{}", server);
630                }
631
632                let parsed_server = match rsipstack::rsip::Uri::try_from(server.as_str()) {
633                    Ok(uri) => uri,
634                    Err(e) => {
635                        warn!("failed to parse server URI: {} {:?}", option.server, e);
636                        continue;
637                    }
638                };
639
640                let server_host = match &parsed_server.host_with_port.host {
641                    rsipstack::rsip::Host::Domain(domain) => domain.to_string(),
642                    rsipstack::rsip::Host::IpAddr(ip) => {
643                        // Compare IP addresses
644                        if let rsipstack::rsip::Host::IpAddr(callee_ip) =
645                            &parsed_callee.host_with_port.host
646                        {
647                            if ip == callee_ip {
648                                if let Some(cred) = &option.credential {
649                                    info!(
650                                        callee,
651                                        username = cred.username,
652                                        server = option.server,
653                                        "Auto-injecting credentials from registered user for outbound call (IP match)"
654                                    );
655                                    return Some(cred.clone());
656                                }
657                            }
658                        }
659                        continue;
660                    }
661                };
662
663                if server_host == callee_host {
664                    if let Some(cred) = &option.credential {
665                        info!(
666                            callee,
667                            username = cred.username,
668                            server = option.server,
669                            "Auto-injecting credentials from registered user for outbound call"
670                        );
671                        return Some(cred.clone());
672                    }
673                }
674            }
675        }
676
677        None
678    }
679
680    /// Helper function to find credentials by IP address
681    fn find_credentials_by_ip(
682        &self,
683        callee_ip: &std::net::IpAddr,
684    ) -> Option<crate::useragent::registration::UserCredential> {
685        if let Some(register_users) = &self.config.register_users {
686            for option in register_users.iter() {
687                let mut server = option.server.clone();
688                if !server.starts_with("sip:") && !server.starts_with("sips:") {
689                    server = format!("sip:{}", server);
690                }
691
692                if let Ok(parsed_server) = rsipstack::rsip::Uri::try_from(server.as_str()) {
693                    if let rsipstack::rsip::Host::IpAddr(server_ip) =
694                        &parsed_server.host_with_port.host
695                    {
696                        if server_ip == callee_ip {
697                            if let Some(cred) = &option.credential {
698                                info!(
699                                    callee_ip = %callee_ip,
700                                    username = cred.username,
701                                    server = option.server,
702                                    "Auto-injecting credentials from registered user for outbound call (IP match)"
703                                );
704                                return Some(cred.clone());
705                            }
706                        }
707                    }
708                }
709            }
710        }
711        None
712    }
713
714    pub async fn stop_registration(&self, wait_for_clear: Option<Duration>) -> Result<()> {
715        {
716            let mut handles = self.registration_handles.lock().await;
717            for (_, cancel_token) in handles.drain() {
718                cancel_token.cancel();
719            }
720        }
721
722        if let Some(duration) = wait_for_clear {
723            let live_users = self.alive_users.clone();
724            let check_loop = async move {
725                loop {
726                    let is_empty = {
727                        let users = live_users
728                            .read()
729                            .map_err(|_| anyhow::anyhow!("Lock poisoned"))?;
730                        users.is_empty()
731                    };
732                    if is_empty {
733                        break;
734                    }
735                    tokio::time::sleep(Duration::from_millis(50)).await;
736                }
737                Ok::<(), anyhow::Error>(())
738            };
739            match tokio::time::timeout(duration, check_loop).await {
740                Ok(_) => {}
741                Err(e) => {
742                    warn!("failed to wait for clear: {}", e);
743                    return Err(anyhow::anyhow!("failed to wait for clear: {}", e));
744                }
745            }
746        }
747        Ok(())
748    }
749
750    pub async fn register(&self, option: RegisterOption) -> Result<()> {
751        let user = option.aor();
752        let mut server = option.server.clone();
753        if !server.starts_with("sip:") && !server.starts_with("sips:") {
754            server = format!("sip:{}", server);
755        }
756        let sip_server = match rsipstack::rsip::Uri::try_from(server) {
757            Ok(uri) => uri,
758            Err(e) => {
759                warn!("failed to parse server: {} {:?}", e, option.server);
760                return Err(anyhow::anyhow!("failed to parse server: {}", e));
761            }
762        };
763        let cancel_token = self.token.child_token();
764        let credential = option.credential.clone().map(|c| c.into());
765        let registration = rsipstack::dialog::registration::Registration::new(
766            self.endpoint.inner.clone(),
767            credential,
768        );
769        let mut handle = RegistrationHandle {
770            registration,
771            option,
772            cancel_token: cancel_token.clone(),
773            start_time: Instant::now(),
774            last_update: Instant::now(),
775            last_response: None,
776        };
777        self.registration_handles
778            .lock()
779            .await
780            .insert(user.clone(), cancel_token);
781        tracing::debug!(user = user.as_str(), "starting registration task");
782        let alive_users = self.alive_users.clone();
783
784        crate::spawn(async move {
785            handle.start_time = Instant::now();
786            let cancel_token = handle.cancel_token.clone();
787            let addrs = handle.registration.endpoint.get_addrs();
788            let local_bind_addr = if let Some(addr) = find_local_addr_for_uri(&addrs, &sip_server) {
789                addr
790            } else {
791                warn!(
792                    user = user.as_str(),
793                    server = %sip_server,
794                    "failed to get local bind address for registration transport"
795                );
796                alive_users.write().unwrap().remove(&user);
797                return;
798            };
799            let user = handle.option.aor();
800            alive_users.write().unwrap().remove(&user);
801            let mut contact_address = local_bind_addr.addr.clone();
802            let mut contact = build_contact(
803                &local_bind_addr,
804                Some(contact_address.clone()),
805                Some(handle.option.username.as_str()),
806                None,
807            );
808            let mut should_register = true;
809            let mut timer = pending().boxed();
810
811            loop {
812                select! {
813                    _ = cancel_token.cancelled() => {
814                        break;
815                    }
816                    _ = timer.as_mut(), if !should_register => {
817                        should_register = true;
818                        timer = Box::pin(pending());
819                    }
820                    result = handle.do_register(&sip_server, None, &contact), if should_register => {
821                        match result {
822                            Ok((expires, new_addr)) => {
823                                if handle
824                                    .should_retry_registration_now(
825                                        &local_bind_addr,
826                                        &contact_address,
827                                        new_addr.as_ref(),
828                                    )
829                                {
830                                    if let Some(next_contact_address) = new_addr {
831                                        info!(
832                                            user = user.as_str(),
833                                            current_contact = %contact_address,
834                                            next_contact = %next_contact_address,
835                                            "public address changed, retrying registration immediately",
836                                        );
837                                        contact_address = next_contact_address;
838                                        contact = build_contact(
839                                            &local_bind_addr,
840                                            Some(contact_address.clone()),
841                                            Some(handle.option.username.as_str()),
842                                            None,
843                                        );
844                                        continue;
845                                    }
846                                }
847                                info!(
848                                    user = user.as_str(),
849                                    expires = expires,
850                                    contact = %contact_address,
851                                    alive_users = alive_users.read().unwrap().len(),
852                                    "registration refreshed",
853                                );
854                                alive_users.write().unwrap().insert(user.clone());
855                                should_register = false;
856                                timer = Box::pin(tokio::time::sleep(Duration::from_secs(
857                                    (expires * 3 / 4) as u64,
858                                )));
859                            }
860                            Err(e) => {
861                                warn!(
862                                    user = user.as_str(),
863                                    alive_users = alive_users.read().unwrap().len(),
864                                    "registration failed: {:?}", e
865                                );
866                                should_register = false;
867                                timer = Box::pin(tokio::time::sleep(Duration::from_secs(60)));
868                            }
869                        }
870                    }
871                }
872            }
873            handle
874                .do_register(&sip_server, Some(0), &contact)
875                .await
876                .ok();
877            alive_users.write().unwrap().remove(&user);
878        });
879        Ok(())
880    }
881}
882
883impl Drop for AppStateInner {
884    fn drop(&mut self) {
885        self.stop();
886    }
887}
888
889impl AppStateBuilder {
890    pub fn new() -> Self {
891        Self {
892            config: None,
893            stream_engine: None,
894            callrecord_sender: None,
895            callrecord_formatter: None,
896            cancel_token: None,
897            create_invitation_handler: None,
898            config_path: None,
899            message_inspector: None,
900            target_locator: None,
901            transport_inspector: None,
902        }
903    }
904
905    pub fn with_config(mut self, config: Config) -> Self {
906        self.config = Some(config);
907        self
908    }
909
910    pub fn with_stream_engine(mut self, stream_engine: Arc<StreamEngine>) -> Self {
911        self.stream_engine = Some(stream_engine);
912        self
913    }
914
915    pub fn with_callrecord_sender(mut self, sender: CallRecordSender) -> Self {
916        self.callrecord_sender = Some(sender);
917        self
918    }
919
920    pub fn with_cancel_token(mut self, token: CancellationToken) -> Self {
921        self.cancel_token = Some(token);
922        self
923    }
924
925    pub fn with_config_metadata(mut self, path: Option<String>) -> Self {
926        self.config_path = path;
927        self
928    }
929
930    pub fn with_inspector(&mut self, inspector: Box<dyn MessageInspector>) -> &mut Self {
931        self.message_inspector = Some(inspector);
932        self
933    }
934    pub fn with_target_locator(&mut self, locator: Box<dyn TargetLocator>) -> &mut Self {
935        self.target_locator = Some(locator);
936        self
937    }
938
939    pub fn with_transport_inspector(
940        &mut self,
941        inspector: Box<dyn TransportEventInspector>,
942    ) -> &mut Self {
943        self.transport_inspector = Some(inspector);
944        self
945    }
946
947    pub async fn build(self) -> Result<AppState> {
948        let config: Arc<Config> = Arc::new(self.config.unwrap_or_default());
949        let token = self
950            .cancel_token
951            .unwrap_or_else(|| CancellationToken::new());
952        let _ = set_cache_dir(&config.media_cache_path);
953        let local_ip = if !config.addr.is_empty() {
954            std::net::IpAddr::from_str(config.addr.as_str())?
955        } else {
956            crate::net_tool::get_first_non_loopback_interface()?
957        };
958        let transport_layer = rsipstack::transport::TransportLayer::new(token.clone());
959        let local_addr: SocketAddr = format!("{}:{}", local_ip, config.udp_port).parse()?;
960
961        // Create UDP socket with SO_REUSEPORT for graceful restarts
962        #[cfg(unix)]
963        let std_socket = {
964            use socket2::{Domain, Protocol, SockAddr, Socket, Type};
965
966            let domain = if local_addr.is_ipv4() {
967                Domain::IPV4
968            } else {
969                Domain::IPV6
970            };
971            let socket = Socket::new(domain, Type::DGRAM, Some(Protocol::UDP))
972                .map_err(|err| anyhow::anyhow!("Failed to create UDP socket: {}", err))?;
973
974            socket
975                .set_reuse_address(true)
976                .map_err(|err| anyhow::anyhow!("Failed to set SO_REUSEADDR: {}", err))?;
977
978            // SO_REUSEPORT (Linux/BSD)
979            #[cfg(not(any(target_os = "solaris", target_os = "illumos", target_os = "cygwin")))]
980            socket
981                .set_reuse_port(true)
982                .map_err(|err| anyhow::anyhow!("Failed to set SO_REUSEPORT: {}", err))?;
983
984            socket
985                .bind(&SockAddr::from(local_addr))
986                .map_err(|err| anyhow::anyhow!("Failed to bind UDP socket: {}", err))?;
987
988            let std_socket: std::net::UdpSocket = socket.into();
989            std_socket
990        };
991
992        #[cfg(not(unix))]
993        let std_socket = std::net::UdpSocket::bind(local_addr)?;
994
995        std_socket.set_nonblocking(true)?;
996        let tokio_socket = tokio::net::UdpSocket::from_std(std_socket)?;
997        // Use the actual bound address (important when port=0 lets OS assign a port)
998        let actual_addr = tokio_socket.local_addr()?;
999        let bind_addr = rsipstack::transport::SipConnection::resolve_bind_address(actual_addr);
1000        let mut learned_public_address: SharedPublicAddress =
1001            Arc::new(ArcSwap::from_pointee(bind_addr.into()));
1002
1003        let udp_inner = rsipstack::transport::udp::UdpInner {
1004            conn: tokio_socket,
1005            addr: rsipstack::transport::SipAddr {
1006                r#type: Some(rsipstack::rsip::transport::Transport::Udp),
1007                addr: bind_addr.into(),
1008            },
1009        };
1010
1011        let external = config
1012            .external_ip
1013            .as_ref()
1014            .map(|ip| {
1015                format!("{}:{}", ip, actual_addr.port())
1016                    .parse()
1017                    .map_err(|e| anyhow::anyhow!("Failed to parse external address: {}", e))
1018            })
1019            .transpose()?;
1020
1021        let udp_conn = rsipstack::transport::udp::UdpConnection::attach(
1022            udp_inner,
1023            external,
1024            Some(token.child_token()),
1025        )
1026        .await;
1027
1028        info!(
1029            "start useragent, addr: {} (SO_REUSEPORT enabled)",
1030            udp_conn.get_addr()
1031        );
1032
1033        transport_layer.add_transport(udp_conn.into());
1034
1035        // Optional SIP over TLS transport
1036        if let Some(tls_port) = config.tls_port {
1037            let tls_addr: std::net::SocketAddr = format!("{}:{}", local_ip, tls_port).parse()?;
1038            let tls_sip_addr = rsipstack::transport::SipAddr {
1039                r#type: Some(rsipstack::rsip::transport::Transport::Tls),
1040                addr: tls_addr.into(),
1041            };
1042            let mut tls_cfg = rsipstack::transport::tls::TlsConfig::default();
1043            if let Some(ref cert_path) = config.tls_cert_file {
1044                tls_cfg.cert = Some(
1045                    std::fs::read(cert_path)
1046                        .map_err(|e| anyhow::anyhow!("tls_cert_file: {}", e))?,
1047                );
1048            }
1049            if let Some(ref key_path) = config.tls_key_file {
1050                tls_cfg.key = Some(
1051                    std::fs::read(key_path).map_err(|e| anyhow::anyhow!("tls_key_file: {}", e))?,
1052                );
1053            }
1054            let external_tls_addr = config
1055                .external_ip
1056                .as_ref()
1057                .and_then(|ip| format!("{}:{}", ip, tls_port).parse().ok());
1058            match rsipstack::transport::tls::TlsListenerConnection::new(
1059                tls_sip_addr,
1060                external_tls_addr,
1061                tls_cfg,
1062            )
1063            .await
1064            {
1065                Ok(tls_conn) => {
1066                    transport_layer.add_transport(tls_conn.into());
1067                    info!("TLS SIP transport started on {}:{}", local_ip, tls_port);
1068                }
1069                Err(e) => {
1070                    return Err(anyhow::anyhow!("Failed to start TLS SIP transport: {}", e));
1071                }
1072            }
1073        }
1074
1075        let endpoint_option = rsipstack::transaction::endpoint::EndpointOption::default();
1076        let mut endpoint_builder = rsipstack::EndpointBuilder::new();
1077        if let Some(ref user_agent) = config.useragent {
1078            endpoint_builder.with_user_agent(user_agent.as_str());
1079        }
1080
1081        let mut endpoint_builder = endpoint_builder
1082            .with_cancel_token(token.child_token())
1083            .with_transport_layer(transport_layer)
1084            .with_option(endpoint_option);
1085
1086        if config.auto_learn_public_address.unwrap_or_default() {
1087            let inspector = LearningMessageInspector::new(bind_addr.into(), self.message_inspector);
1088            learned_public_address = inspector.shared_public_address();
1089            endpoint_builder = endpoint_builder.with_inspector(Box::new(inspector));
1090        } else if let Some(inspector) = self.message_inspector {
1091            endpoint_builder = endpoint_builder.with_inspector(inspector);
1092        }
1093
1094        if let Some(locator) = self.target_locator {
1095            endpoint_builder.with_target_locator(locator);
1096        } else if let Some(ref rules) = config.rewrites {
1097            endpoint_builder
1098                .with_target_locator(Box::new(RewriteTargetLocator::new(rules.clone())));
1099        }
1100
1101        if let Some(inspector) = self.transport_inspector {
1102            endpoint_builder = endpoint_builder.with_transport_inspector(inspector);
1103        }
1104
1105        let endpoint = endpoint_builder.build();
1106        let dialog_layer = Arc::new(DialogLayer::new(endpoint.inner.clone()));
1107
1108        let stream_engine = self.stream_engine.unwrap_or_default();
1109
1110        let callrecord_formatter = if let Some(formatter) = self.callrecord_formatter {
1111            formatter
1112        } else {
1113            let formatter = if let Some(ref callrecord) = config.callrecord {
1114                DefaultCallRecordFormatter::new_with_config(callrecord)
1115            } else {
1116                DefaultCallRecordFormatter::default()
1117            };
1118            Arc::new(formatter)
1119        };
1120
1121        let callrecord_sender = if let Some(sender) = self.callrecord_sender {
1122            Some(sender)
1123        } else if let Some(ref callrecord) = config.callrecord {
1124            let builder = CallRecordManagerBuilder::new()
1125                .with_cancel_token(token.child_token())
1126                .with_config(callrecord.clone())
1127                .with_max_concurrent(32)
1128                .with_formatter(callrecord_formatter.clone());
1129
1130            let mut callrecord_manager = builder.build();
1131            let sender = callrecord_manager.sender.clone();
1132            crate::spawn(async move {
1133                callrecord_manager.serve().await;
1134            });
1135            Some(sender)
1136        } else {
1137            None
1138        };
1139
1140        let app_state = Arc::new(AppStateInner {
1141            config,
1142            token,
1143            stream_engine,
1144            callrecord_sender,
1145            endpoint,
1146            registration_handles: Mutex::new(HashMap::new()),
1147            alive_users: Arc::new(RwLock::new(HashSet::new())),
1148            dialog_layer: dialog_layer.clone(),
1149            create_invitation_handler: self.create_invitation_handler,
1150            invitation: Invitation::new(dialog_layer),
1151            routing_state: Arc::new(crate::call::RoutingState::new()),
1152            pending_playbooks: Arc::new(Mutex::new(HashMap::new())),
1153            learned_public_address,
1154            active_calls: Arc::new(std::sync::Mutex::new(HashMap::new())),
1155            total_calls: AtomicU64::new(0),
1156            total_failed_calls: AtomicU64::new(0),
1157            uptime: Local::now(),
1158            shutting_down: Arc::new(AtomicBool::new(false)),
1159        });
1160
1161        Ok(app_state)
1162    }
1163}