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