Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 0 additions & 49 deletions src/rendezvous_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -357,7 +357,7 @@
bytes: &BytesMut,
addr: SocketAddr,
socket: &mut FramedSocket,
key: &str,

Check warning on line 360 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Compile Check

unused variable: `key`

Check warning on line 360 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Run Tests

unused variable: `key`
) -> ResultType<()> {
if let Ok(msg_in) = RendezvousMessage::parse_from_bytes(bytes) {
match msg_in.union {
Expand Down Expand Up @@ -411,24 +411,9 @@
}
}
}
Some(rendezvous_message::Union::PunchHoleRequest(ph)) => {
if self.pm.is_in_memory(&ph.id).await {
self.handle_udp_punch_hole_request(addr, ph, key).await?;
} else {
// not in memory, fetch from db with spawn in case blocking me
let mut me = self.clone();
let key = key.to_owned();
tokio::spawn(async move {
allow_err!(me.handle_udp_punch_hole_request(addr, ph, &key).await);
});
}
}
Some(rendezvous_message::Union::PunchHoleSent(phs)) => {
self.handle_hole_sent(phs, addr, Some(socket)).await?;
}
Some(rendezvous_message::Union::LocalAddr(la)) => {
self.handle_local_addr(la, addr, Some(socket)).await?;
}
Some(rendezvous_message::Union::ConfigureUpdate(mut cu)) => {
if try_into_v4(addr).ip().is_loopback() && cu.serial > self.inner.serial {
let mut inner: Inner = (*self.inner).clone();
Expand Down Expand Up @@ -1070,40 +1055,6 @@
Ok(())
}

#[inline]
async fn handle_udp_punch_hole_request(
&mut self,
addr: SocketAddr,
ph: PunchHoleRequest,
key: &str,
) -> ResultType<()> {
let (msg, to_addr) = self.handle_punch_hole_request(addr, ph, key, false).await?;
match to_addr {
Some(to_addr) => {
// Check if target is a WS/TCP peer with persistent connection
let addr_v4 = try_into_v4(to_addr);
let sink_arc = self.ws_map.lock().await.get(&addr_v4).cloned();
if let Some(sink_arc) = sink_arc {
let mut ws_sink = sink_arc.lock().await;
Self::send_to_sink(&mut *ws_sink, msg).await;
} else {
let sink_arc = self.tcp_map.lock().await.get(&addr_v4).cloned();
if let Some(sink_arc) = sink_arc {
let mut tcp_sink = sink_arc.lock().await;
Self::send_to_sink(&mut *tcp_sink, msg).await;
} else {
self.tx.send(Data::Msg(msg.into(), to_addr))?;
}
}
}
None => {
// Error response goes back to the UDP requester
self.tx.send(Data::Msg(msg.into(), addr))?;
}
}
Ok(())
}

fn make_register_pk_response(result: register_pk_response::Result) -> RendezvousMessage {
let mut msg_out = RendezvousMessage::new();
msg_out.set_register_pk_response(RegisterPkResponse {
Expand Down Expand Up @@ -1380,7 +1331,7 @@
let arg = fds.next();
if let Some("-") = arg { lock.clear(); }
else {
let mut start = arg.and_then(|x| x.parse::<usize>().ok()).unwrap_or(0);

Check warning on line 1334 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Compile Check

variable does not need to be mutable

Check warning on line 1334 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Run Tests

variable does not need to be mutable
let mut page_size = fds.next().and_then(|x| x.parse::<usize>().ok()).unwrap_or(10);
if page_size == 0 { page_size = 10; }
for (_, e) in lock.iter().enumerate().skip(start).take(page_size) {
Expand Down Expand Up @@ -1570,12 +1521,12 @@
async fn handle_listener_inner(
&mut self,
stream: TcpStream,
mut addr: SocketAddr,

Check warning on line 1524 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Compile Check

variable does not need to be mutable

Check warning on line 1524 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Run Tests

variable does not need to be mutable
key: &str,
ws: bool,
) -> ResultType<()> {
let mut sink;
let mut forwarded_ip: Option<String> = None;

Check warning on line 1529 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Compile Check

value assigned to `forwarded_ip` is never read

Check warning on line 1529 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Run Tests

value assigned to `forwarded_ip` is never read
if ws {
use tokio_tungstenite::tungstenite::handshake::server::{Request, Response};
let forwarded_ip_cell: std::sync::Arc<std::sync::Mutex<Option<String>>> =
Expand Down Expand Up @@ -1855,7 +1806,7 @@
}

#[inline]
async fn send_rk_res(

Check warning on line 1809 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Compile Check

function `send_rk_res` is never used

Check warning on line 1809 in src/rendezvous_server.rs

View workflow job for this annotation

GitHub Actions / Run Tests

function `send_rk_res` is never used
socket: &mut FramedSocket,
addr: SocketAddr,
res: register_pk_response::Result,
Expand Down
Loading