forked from torrust/torrust-tracker
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathprocessor.rs
More file actions
132 lines (114 loc) · 5 KB
/
Copy pathprocessor.rs
File metadata and controls
132 lines (114 loc) · 5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
use std::io::Cursor;
use std::net::{IpAddr, SocketAddr};
use std::sync::Arc;
use std::time::Duration;
use aquatic_udp_protocol::Response;
use bittorrent_udp_tracker_core::container::UdpTrackerCoreContainer;
use bittorrent_udp_tracker_core::{self};
use tokio::time::Instant;
use tracing::{instrument, Level};
use super::bound_socket::BoundSocket;
use crate::container::UdpTrackerServerContainer;
use crate::handlers::CookieTimeValues;
use crate::{handlers, statistics, RawRequest};
pub struct Processor {
socket: Arc<BoundSocket>,
udp_tracker_core_container: Arc<UdpTrackerCoreContainer>,
udp_tracker_server_container: Arc<UdpTrackerServerContainer>,
cookie_lifetime: f64,
}
impl Processor {
pub fn new(
socket: Arc<BoundSocket>,
udp_tracker_core_container: Arc<UdpTrackerCoreContainer>,
udp_tracker_server_container: Arc<UdpTrackerServerContainer>,
cookie_lifetime: f64,
) -> Self {
Self {
socket,
udp_tracker_core_container,
udp_tracker_server_container,
cookie_lifetime,
}
}
#[instrument(skip(self, request))]
pub async fn process_request(self, request: RawRequest) {
let from = request.from;
let start_time = Instant::now();
let response = handlers::handle_packet(
request,
self.udp_tracker_core_container.clone(),
self.udp_tracker_server_container.clone(),
self.socket.address(),
CookieTimeValues::new(self.cookie_lifetime),
)
.await;
let elapsed_time = start_time.elapsed();
self.send_response(from, response, elapsed_time).await;
}
#[instrument(skip(self))]
async fn send_response(self, target: SocketAddr, response: Response, req_processing_time: Duration) {
tracing::debug!("send response");
let response_type = match &response {
Response::Connect(_) => "Connect".to_string(),
Response::AnnounceIpv4(_) => "AnnounceIpv4".to_string(),
Response::AnnounceIpv6(_) => "AnnounceIpv6".to_string(),
Response::Scrape(_) => "Scrape".to_string(),
Response::Error(e) => format!("Error: {e:?}"),
};
let udp_response_kind = match &response {
Response::Connect(_) => statistics::event::UdpResponseKind::Connect,
Response::AnnounceIpv4(_) | Response::AnnounceIpv6(_) => statistics::event::UdpResponseKind::Announce,
Response::Scrape(_) => statistics::event::UdpResponseKind::Scrape,
Response::Error(_e) => statistics::event::UdpResponseKind::Error,
};
let mut writer = Cursor::new(Vec::with_capacity(200));
match response.write_bytes(&mut writer) {
Ok(()) => {
let bytes_count = writer.get_ref().len();
let payload = writer.get_ref();
let () = match self.send_packet(&target, payload).await {
Ok(sent_bytes) => {
if tracing::event_enabled!(Level::TRACE) {
tracing::debug!(%bytes_count, %sent_bytes, ?payload, "sent {response_type}");
} else {
tracing::debug!(%bytes_count, %sent_bytes, "sent {response_type}");
}
if let Some(udp_server_stats_event_sender) =
self.udp_tracker_server_container.udp_server_stats_event_sender.as_deref()
{
match target.ip() {
IpAddr::V4(_) => {
udp_server_stats_event_sender
.send_event(statistics::event::Event::Udp4Response {
kind: udp_response_kind,
req_processing_time,
})
.await;
}
IpAddr::V6(_) => {
udp_server_stats_event_sender
.send_event(statistics::event::Event::Udp6Response {
kind: udp_response_kind,
req_processing_time,
})
.await;
}
}
}
}
Err(error) => tracing::warn!(%bytes_count, %error, ?payload, "failed to send"),
};
}
Err(e) => {
tracing::error!(%e, "error");
}
}
}
#[instrument(skip(self))]
async fn send_packet(&self, target: &SocketAddr, payload: &[u8]) -> std::io::Result<usize> {
tracing::trace!("send packet");
// doesn't matter if it reaches or not
self.socket.send_to(payload, target).await
}
}