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