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