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