Skip to content
Open
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
59 changes: 51 additions & 8 deletions src/rendezvous_server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use hbb_common::{
register_pk_response::Result::{TOO_FREQUENT, UUID_MISMATCH},
*,
},
tcp::FramedStream,
tcp::{Encrypt, FramedStream},
timeout,
tokio::{
self,
Expand All @@ -31,7 +31,7 @@ use hbb_common::{
AddrMangle, ResultType,
};
use ipnetwork::Ipv4Network;
use sodiumoxide::crypto::sign;
use sodiumoxide::crypto::{box_, sign};
use std::{
collections::HashMap,
net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr},
Expand All @@ -51,7 +51,7 @@ const REG_TIMEOUT: i64 = 30_000;
type TcpStreamSink = SplitSink<Framed<TcpStream, BytesCodec>, Bytes>;
type WsSink = SplitSink<tokio_tungstenite::WebSocketStream<TcpStream>, tungstenite::Message>;
enum Sink {
TcpStream(TcpStreamSink),
TcpStream(TcpStreamSink, Option<Encrypt>),
Ws(WsSink),
}
type Sender = mpsc::UnboundedSender<Data>;
Expand Down Expand Up @@ -852,8 +852,13 @@ impl RendezvousServer {
if let Some(sink) = sink.as_mut() {
if let Ok(bytes) = msg.write_to_bytes() {
match sink {
Sink::TcpStream(s) => {
allow_err!(s.send(Bytes::from(bytes)).await);
Sink::TcpStream(s, encrypt) => {
let bytes = if let Some(encrypt) = encrypt.as_mut() {
Bytes::from(encrypt.enc(&bytes))
} else {
Bytes::from(bytes)
};
allow_err!(s.send(bytes).await);
}
Sink::Ws(ws) => {
allow_err!(ws.send(tungstenite::Message::Binary(bytes)).await);
Expand Down Expand Up @@ -1215,9 +1220,47 @@ impl RendezvousServer {
}
}
} else {
let (a, mut b) = Framed::new(stream, BytesCodec::new()).split();
sink = Some(Sink::TcpStream(a));
while let Ok(Some(Ok(bytes))) = timeout(30_000, b.next()).await {
let (mut a, mut b) = Framed::new(stream, BytesCodec::new()).split();

// Recent logged-in clients secure the TCP rendezvous channel before
// sending requests. Advertise an ephemeral box key signed by the
// server identity so those clients can authenticate this server.
// Logged-out and older clients continue below in plaintext.
let mut handshake_secret = if let Some(server_sk) = self.inner.sk.as_ref() {
let (public_key, secret_key) = box_::gen_keypair();
let mut key_exchange = RendezvousMessage::new();
key_exchange.set_key_exchange(KeyExchange {
keys: vec![Bytes::from(sign::sign(&public_key.0, server_sk))],
..Default::default()
});
a.send(Bytes::from(key_exchange.write_to_bytes()?)).await?;
Some(secret_key)
} else {
None
};
Comment thread
TylonHH marked this conversation as resolved.

sink = Some(Sink::TcpStream(a, None));
let mut receive_encrypt: Option<Encrypt> = None;
while let Ok(Some(Ok(mut bytes))) = timeout(30_000, b.next()).await {
if let Some(secret_key) = handshake_secret.take() {
if let Ok(msg) = RendezvousMessage::parse_from_bytes(&bytes) {
if let Some(rendezvous_message::Union::KeyExchange(exchange)) = msg.union {
if exchange.keys.len() != 2 {
bail!("Invalid secure TCP key exchange");
Comment thread
TylonHH marked this conversation as resolved.
}
let key =
Encrypt::decode(&exchange.keys[1], &exchange.keys[0], &secret_key)?;
receive_encrypt = Some(Encrypt::new(key.clone()));
if let Some(Sink::TcpStream(_, send_encrypt)) = sink.as_mut() {
*send_encrypt = Some(Encrypt::new(key));
}
continue;
}
}
}
Comment thread
TylonHH marked this conversation as resolved.
if let Some(encrypt) = receive_encrypt.as_mut() {
encrypt.dec(&mut bytes)?;
}
if !self.handle_tcp(&bytes, &mut sink, addr, key, ws).await {
break;
}
Expand Down