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