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