diff --git a/src/common.rs b/src/common.rs index 57e91db85..648bc6b5c 100644 --- a/src/common.rs +++ b/src/common.rs @@ -1167,13 +1167,6 @@ pub fn get_ipv6_punch_enabled() -> bool { ) } -pub fn get_port_forward_mux_enabled() -> bool { - config::option2bool( - keys::OPTION_ENABLE_PORT_FORWARD_MUX, - &get_local_option(keys::OPTION_ENABLE_PORT_FORWARD_MUX), - ) -} - pub fn get_local_option(key: &str) -> String { let v = LocalConfig::get_option(key); if key == keys::OPTION_ENABLE_UDP_PUNCH || key == keys::OPTION_ENABLE_IPV6_PUNCH { diff --git a/src/port_forward.rs b/src/port_forward.rs index a11aff445..b5f8348a8 100644 --- a/src/port_forward.rs +++ b/src/port_forward.rs @@ -82,9 +82,6 @@ pub async fn listen( remote_port: i32, ) -> ResultType<()> { let listener = tcp::new_listener(format!("127.0.0.1:{}", port), true).await?; - // One tunnel per mapping: every accept here goes to the one target the - // peer authenticated, and dropping it at the end closes that tunnel. - let tunnel = Tunnel::new(); let addr = listener.local_addr()?; log::info!("listening on port {:?}", addr); let is_rdp = port == 0; @@ -92,151 +89,76 @@ pub async fn listen( run_rdp(addr.port(), &rdp_display_name(&lc, &id)); } let mut ui_receiver = ui_receiver; + // One tunnel per mapping; the listener drops it on its way out, and that + // ends the tunnel. + let tunnel = Tunnel::new(); loop { tokio::select! { - // `addr` above the loop is the listener's own address; `run_rdp` needs - // that port. The accepted peer's address gets its own name so it can - // never shadow it. - Ok((forward, peer_addr)) = listener.accept() => { - log::debug!("new connection from {:?}", peer_addr); - match tunnel.claim() { + Ok((forward, addr)) = listener.accept() => { + log::info!("new connection from {:?}", addr); + // A multiplexed mapping takes the connection on its tunnel, or + // probes for one on its first accept. Everything else, the + // setting off or a peer without the feature, is the raw pipe + // below, as it always was. + let claim = if mux_enabled() { tunnel.claim() } else { Claim::Legacy }; + match claim { Claim::Muxed(handle) => { if let Err(e) = handle.open(&remote_host, remote_port, forward, Vec::new()) { - log::debug!("cannot open channel for {:?}: {}", peer_addr, e); + log::debug!("cannot open channel for {:?}: {}", addr, e); } + continue; } - // The claiming accept negotiates: it asks for the tunnel, and - // the peer's answer fixes this listener's mode until it closes. Claim::Claimed => { - let interface = interface.with_port_forward(login_target( - &remote_host, - remote_port, - crate::common::get_port_forward_mux_enabled(), - )); - let mut forward = Framed::new(forward, BytesCodec::new()); - let mut close_port_forward = false; - match connect_and_login(&id, &password, &mut ui_receiver, interface.clone(), &mut forward, key, token, is_rdp, &mut close_port_forward).await { - Ok(Some(outcome)) if outcome.mux => { - let handle = tunnel.set_muxed(outcome.stream, interface.clone()); - if !outcome.local_eof { - let (socket, prebuf) = take_socket(forward, outcome.prebuf); - if let Err(e) = handle.open(&remote_host, remote_port, socket, prebuf) { - log::debug!("cannot open channel for {:?}: {}", peer_addr, e); - } - } - } - Ok(Some(outcome)) => { - tunnel.set_legacy(); - if outcome.local_eof { - log::debug!("legacy peer and local {:?} already gone", peer_addr); - } else { - run_legacy(outcome, forward, peer_addr, interface.clone()); - } - } - _ if close_port_forward => { - tunnel.set_failed(); - break; - } - Err(err) => { - tunnel.set_failed(); - interface.on_establish_connection_error(err.to_string()); - } - _ => tunnel.set_failed(), + if establish_tunnel(&tunnel, &id, &password, &mut ui_receiver, &interface, forward, addr, key, token, is_rdp, &remote_host, remote_port).await { + break; } + continue; } - // A `Legacy` listener stays legacy until it closes: every accept - // logs in on its own, asks for no tunnel, and takes the raw pipe - // whatever the peer reports. Re-adding the mapping is how a user - // picks up an upgraded peer; nothing switches modes underneath - // live connections. - Claim::Legacy => { - let interface = interface.with_port_forward(login_target(&remote_host, remote_port, false)); - let mut forward = Framed::new(forward, BytesCodec::new()); - let mut close_port_forward = false; - match connect_and_login(&id, &password, &mut ui_receiver, interface.clone(), &mut forward, key, token, is_rdp, &mut close_port_forward).await { - Ok(Some(outcome)) if outcome.local_eof => { - log::debug!("legacy peer and local {:?} already gone", peer_addr); + Claim::Legacy => {} + } + // The target rides with this accept's interface clone, so two + // mappings logging in at once cannot overwrite each other's. + let interface = interface.with_port_forward(login_target(&remote_host, remote_port, false)); + let id = id.clone(); + let password = password.clone(); + let mut forward = Framed::new(forward, BytesCodec::new()); + let mut close_port_forward = false; + match connect_and_login(&id, &password, &mut ui_receiver, interface.clone(), &mut forward, key, token, is_rdp, &mut close_port_forward).await { + Ok(Some(stream)) => { + let interface = interface.clone(); + tokio::spawn(async move { + if let Err(err) = run_forward(forward, stream).await { + interface.msgbox("error", "Error", &err.to_string(), ""); } - Ok(Some(outcome)) => run_legacy(outcome, forward, peer_addr, interface.clone()), - _ if close_port_forward => break, - Err(err) => interface.on_establish_connection_error(err.to_string()), - _ => {} - } + log::info!("connection from {:?} closed", addr); + }); } + _ if close_port_forward => { + break; + } + Err(err) => { + interface.on_establish_connection_error(err.to_string()); + } + _ => {} + } + } + d = ui_receiver.recv() => { + match d { + Some(Data::Close) => { + break; + } + Some(Data::NewRDP) => { + println!("receive run_rdp from ui_receiver"); + run_rdp(addr.port(), &rdp_display_name(&lc, &id)); + } + _ => {} } } - d = ui_receiver.recv() => if on_ui_command(d, addr.port(), &lc, &id) { - break; - }, } } Ok(()) } -/// Commands the window sends its listener. `true` means stop listening: the -/// window is closing, or its sender is gone. -fn on_ui_command(d: Option, port: u16, lc: &Arc>, id: &str) -> bool { - match d { - Some(Data::Close) | None => true, - Some(Data::NewRDP) => { - run_rdp(port, &rdp_display_name(lc, id)); - false - } - _ => false, - } -} - -/// Today's raw pipe, for peers without multiplexing. -fn run_legacy( - outcome: LoginOutcome, - forward: Framed, - addr: std::net::SocketAddr, - interface: impl Interface, -) { - let mut stream = outcome.stream; - let prebuf = outcome.prebuf; - tokio::spawn(async move { - stream.set_raw(); - if !prebuf.is_empty() { - allow_err!(stream.send_bytes(prebuf.into()).await); - } - if let Err(err) = run_forward(forward, stream).await { - interface.msgbox("error", "Error", &err.to_string(), ""); - } - log::info!("connection from {:?} closed", addr); - }); -} - -pub(crate) struct LoginOutcome { - stream: Stream, - mux: bool, - prebuf: Vec, - local_eof: bool, -} - -/// The target this accept's login asks for. It travels with the interface -/// clone rather than through the shared `LoginConfigHandler`, so mappings -/// logging in at the same time cannot overwrite each other's target. -fn login_target(host: &str, port: i32, multiplex: bool) -> PortForward { - PortForward { - host: host.to_owned(), - port, - multiplex, - ..Default::default() - } -} - -fn peer_supports_mux(pi: &PeerInfo) -> bool { - pi.features.as_ref().map(|f| f.port_forward_mux).unwrap_or(false) -} - -/// `into_inner()` would drop bytes the codec pulled but never yielded. -fn take_socket(forward: Framed, mut prebuf: Vec) -> (TcpStream, Vec) { - let parts = forward.into_parts(); - prebuf.extend_from_slice(&parts.read_buf); - (parts.io, prebuf) -} - async fn connect_and_login( id: &str, password: &str, @@ -247,6 +169,166 @@ async fn connect_and_login( token: &str, is_rdp: bool, close_port_forward: &mut bool, +) -> ResultType> { + let conn_type = if is_rdp { + ConnType::RDP + } else { + ConnType::PORT_FORWARD + }; + let ((mut stream, direct, _pk, _kcp, _stream_type), (feedback, rendezvous_server)) = + Client::start(id, key, token, conn_type, interface.clone()).await?; + interface.update_direct(Some(direct)); + if !stream.is_secured() && !crate::common::is_direct_ip_access(id) { + if !confirm_insecure_connection(&interface, ui_receiver).await { + *close_port_forward = true; + return Ok(None); + } + } + let mut buffer = Vec::new(); + let mut received = false; + + let _keep_it = hc_connection(feedback, rendezvous_server, token).await; + + loop { + tokio::select! { + res = timeout(READ_TIMEOUT, stream.next()) => match res { + Err(_) => { + bail!("Timeout"); + } + Ok(Some(Ok(bytes))) => { + if !received { + received = true; + interface.update_received(true); + } + let msg_in = Message::parse_from_bytes(&bytes)?; + match msg_in.union { + Some(message::Union::Hash(hash)) => { + if !interface.handle_hash(password, hash, &mut stream).await { + return Ok(None); + } + } + Some(message::Union::LoginResponse(lr)) => match lr.union { + Some(login_response::Union::Error(err)) => { + if !interface.handle_login_error(&err) { + return Ok(None); + } + } + Some(login_response::Union::PeerInfo(pi)) => { + interface.handle_peer_info(pi); + break; + } + _ => {} + } + Some(message::Union::TestDelay(t)) => { + interface.handle_test_delay(t, &mut stream).await; + } + _ => {} + } + } + Ok(Some(Err(err))) => { + bail!("Connection closed: {}", err); + } + _ => { + bail!("Reset by the peer"); + } + }, + d = ui_receiver.recv() => { + match d { + Some(Data::Login((os_username, os_password, password, remember))) => { + interface.handle_login_from_ui(os_username, os_password, password, remember, &mut stream).await; + } + Some(Data::Message(msg)) => { + allow_err!(stream.send(&msg).await); + } + _ => {} + } + }, + res = forward.next() => { + if let Some(Ok(bytes)) = res { + buffer.extend(bytes); + } else { + return Ok(None); + } + }, + } + } + stream.set_raw(); + if !buffer.is_empty() { + allow_err!(stream.send_bytes(buffer.into()).await); + } + Ok(Some(stream)) +} + +/// The first accept of a multiplexed mapping. It logs in asking for the +/// tunnel, and the peer's answer fixes this listener's mode until it closes: +/// a peer with the feature gets a tunnel every later accept joins, one +/// without gets today's raw pipe for this connection and `Legacy` for the +/// rest. Re-adding the mapping is how a user picks up an upgraded peer; +/// nothing switches modes underneath live connections. Returns `true` when +/// the listener should stop. +async fn establish_tunnel( + tunnel: &Tunnel, + id: &str, + password: &str, + ui_receiver: &mut mpsc::UnboundedReceiver, + interface: &impl Interface, + forward: TcpStream, + addr: std::net::SocketAddr, + key: &str, + token: &str, + is_rdp: bool, + remote_host: &str, + remote_port: i32, +) -> bool { + let interface = interface.with_port_forward(login_target(remote_host, remote_port, true)); + let mut forward = Framed::new(forward, BytesCodec::new()); + let mut close_port_forward = false; + match connect_and_login_mux(id, password, ui_receiver, interface.clone(), &mut forward, key, token, is_rdp, &mut close_port_forward).await { + Ok(Some(outcome)) if outcome.mux => { + let handle = tunnel.set_muxed(outcome.stream, interface.clone()); + if !outcome.local_eof { + let (socket, prebuf) = take_socket(forward, outcome.prebuf); + if let Err(e) = handle.open(remote_host, remote_port, socket, prebuf) { + log::debug!("cannot open channel for {:?}: {}", addr, e); + } + } + } + Ok(Some(outcome)) => { + tunnel.set_legacy(); + if outcome.local_eof { + log::debug!("legacy peer and local {:?} already gone", addr); + } else { + run_legacy(outcome, forward, addr, interface.clone()); + } + } + _ if close_port_forward => { + tunnel.set_failed(); + return true; + } + Err(err) => { + tunnel.set_failed(); + interface.on_establish_connection_error(err.to_string()); + } + _ => tunnel.set_failed(), + } + false +} + +/// `connect_and_login` for a mapping that wants the tunnel: the login asks +/// for it, the pre-read stops at one window rather than growing without +/// bound, and a local EOF no longer ends the login, since the tunnel may +/// still be wanted. It reports what the peer answered rather than a raw +/// stream, because the caller's next step depends on it. +async fn connect_and_login_mux( + id: &str, + password: &str, + ui_receiver: &mut mpsc::UnboundedReceiver, + interface: impl Interface, + forward: &mut Framed, + key: &str, + token: &str, + is_rdp: bool, + close_port_forward: &mut bool, ) -> ResultType> { let conn_type = if is_rdp { ConnType::RDP @@ -348,74 +430,64 @@ async fn connect_and_login( })) } - -/// A mapping's login is built from the window's shared handler: -/// `create_login_msg` reads `port_forward` and `handle_login_from_ui` reads -/// `hash`. Mappings log in concurrently, so each fills them and sends under -/// the window's turn lock, or one login carried another mapping's target or -/// answered another's challenge. -async fn login_with_hash( - interface: &impl Interface, - password: &str, - hash: Hash, - remote_host: &str, - remote_port: i32, - stream: &mut Stream, -) -> bool { - let lc = interface.get_lch(); - let turn = lc.read().unwrap().port_forward_login_turn.clone(); - let _turn = turn.lock().await; - lc.write().unwrap().port_forward = (remote_host.to_owned(), remote_port); - interface.handle_hash(password, hash, stream).await -} - -type UiLogin = (String, String, String, bool); - -/// This connection's `Hash`. The window's password prompt is broadcast to -/// every mapping and can reach this one first, so a password typed while -/// the `Hash` was on its way is kept and answers it now, rather than being -/// dropped in the hope that the mapping which prompted has already stored -/// it in the shared handler. -async fn hash_arrived( - interface: &impl Interface, - password: &str, - hash: Hash, - pending_login: Option, - remote_host: &str, - remote_port: i32, - stream: &mut Stream, -) -> bool { - match pending_login { - Some(login) => { - login_from_ui(interface, &hash, login, remote_host, remote_port, stream).await; - true - } - None => login_with_hash(interface, password, hash, remote_host, remote_port, stream).await, - } -} - -/// The window's password prompt is broadcast to every mapping; this one -/// answers it with its own challenge. -async fn login_from_ui( - interface: &impl Interface, - hash: &Hash, - login: UiLogin, - remote_host: &str, - remote_port: i32, - stream: &mut Stream, +/// Today's raw pipe, for peers without multiplexing. +fn run_legacy( + outcome: LoginOutcome, + forward: Framed, + addr: std::net::SocketAddr, + interface: impl Interface, ) { - let lc = interface.get_lch(); - let turn = lc.read().unwrap().port_forward_login_turn.clone(); - let _turn = turn.lock().await; - { - let mut lc = lc.write().unwrap(); - lc.port_forward = (remote_host.to_owned(), remote_port); - lc.set_hash(hash.clone()); + let mut stream = outcome.stream; + let prebuf = outcome.prebuf; + tokio::spawn(async move { + stream.set_raw(); + if !prebuf.is_empty() { + allow_err!(stream.send_bytes(prebuf.into()).await); + } + if let Err(err) = run_forward(forward, stream).await { + interface.msgbox("error", "Error", &err.to_string(), ""); + } + log::info!("connection from {:?} closed", addr); + }); +} + +struct LoginOutcome { + stream: Stream, + mux: bool, + prebuf: Vec, + local_eof: bool, +} + +/// The target this accept's login asks for. It travels with the interface +/// clone rather than through the shared `LoginConfigHandler`, so mappings +/// logging in at the same time cannot overwrite each other's target. +fn login_target(host: &str, port: i32, multiplex: bool) -> PortForward { + PortForward { + host: host.to_owned(), + port, + multiplex, + ..Default::default() } - let (os_username, os_password, password, remember) = login; - interface - .handle_login_from_ui(os_username, os_password, password, remember, stream) - .await; +} + +fn peer_supports_mux(pi: &PeerInfo) -> bool { + pi.features.as_ref().map(|f| f.port_forward_mux).unwrap_or(false) +} + +/// `into_inner()` would drop bytes the codec pulled but never yielded. +fn take_socket(forward: Framed, mut prebuf: Vec) -> (TcpStream, Vec) { + let parts = forward.into_parts(); + prebuf.extend_from_slice(&parts.read_buf); + (parts.io, prebuf) +} + +/// The controlling side's `enable-port-forward-mux`: on unless set to `N`. +fn mux_enabled() -> bool { + use hbb_common::config::{keys, option2bool, LocalConfig}; + option2bool( + keys::OPTION_ENABLE_PORT_FORWARD_MUX, + &LocalConfig::get_option(keys::OPTION_ENABLE_PORT_FORWARD_MUX), + ) } async fn run_forward(forward: Framed, stream: Stream) -> ResultType<()> {