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