forked from torrust/torrust-tracker
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathserver.rs
More file actions
88 lines (72 loc) · 2.79 KB
/
Copy pathserver.rs
File metadata and controls
88 lines (72 loc) · 2.79 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
use std::io::Cursor;
use std::net::SocketAddr;
use std::sync::Arc;
use aquatic_udp_protocol::Response;
use log::{debug, error, info};
use tokio::net::UdpSocket;
use crate::tracker;
use crate::udp::handlers::handle_packet;
use crate::udp::MAX_PACKET_SIZE;
pub struct Udp {
socket: Arc<UdpSocket>,
tracker: Arc<tracker::Tracker>,
}
impl Udp {
/// # Errors
///
/// Will return `Err` unable to bind to the supplied `bind_address`.
pub async fn new(tracker: Arc<tracker::Tracker>, bind_address: &str) -> tokio::io::Result<Udp> {
let socket = UdpSocket::bind(bind_address).await?;
Ok(Udp {
socket: Arc::new(socket),
tracker,
})
}
/// # Panics
///
/// It would panic if unable to resolve the `local_addr` from the supplied ´socket´.
pub async fn start(&self) {
loop {
let mut data = [0; MAX_PACKET_SIZE];
let socket = self.socket.clone();
let tracker = self.tracker.clone();
tokio::select! {
_ = tokio::signal::ctrl_c() => {
info!("Stopping UDP server: {}..", socket.local_addr().unwrap());
break;
}
Ok((valid_bytes, remote_addr)) = socket.recv_from(&mut data) => {
let payload = data[..valid_bytes].to_vec();
info!("Received {} bytes", payload.len());
debug!("From: {}", &remote_addr);
debug!("Payload: {:?}", payload);
let response = handle_packet(remote_addr, payload, tracker).await;
Udp::send_response(socket, remote_addr, response).await;
}
}
}
}
async fn send_response(socket: Arc<UdpSocket>, remote_addr: SocketAddr, response: Response) {
let buffer = vec![0u8; MAX_PACKET_SIZE];
let mut cursor = Cursor::new(buffer);
match response.write(&mut cursor) {
Ok(_) => {
#[allow(clippy::cast_possible_truncation)]
let position = cursor.position() as usize;
let inner = cursor.get_ref();
info!("Sending {} bytes ...", &inner[..position].len());
debug!("To: {:?}", &remote_addr);
debug!("Payload: {:?}", &inner[..position]);
Udp::send_packet(socket, &remote_addr, &inner[..position]).await;
info!("{} bytes sent", &inner[..position].len());
}
Err(_) => {
error!("could not write response to bytes.");
}
}
}
async fn send_packet(socket: Arc<UdpSocket>, remote_addr: &SocketAddr, payload: &[u8]) {
// doesn't matter if it reaches or not
drop(socket.send_to(payload, remote_addr).await);
}
}