diff --git a/Cargo.lock b/Cargo.lock index dab85cd6..9381ec07 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -257,6 +257,12 @@ version = "0.7.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50" +[[package]] +name = "ascii" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d92bec98840b8f03a5ff5413de5293bfcd8bf96467cf5452609f939ec6f5de16" + [[package]] name = "async-broadcast" version = "0.7.2" @@ -1248,6 +1254,12 @@ dependencies = [ "windows-targets 0.52.6", ] +[[package]] +name = "chunked_transfer" +version = "1.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e4de3bc4ea267985becf712dc6d9eed8b04c953b3fcfb339ebc87acd9804901" + [[package]] name = "cidr" version = "0.3.1" @@ -2088,6 +2100,16 @@ dependencies = [ "dirs-sys 0.5.0", ] +[[package]] +name = "dirs-next" +version = "2.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b98cf8ebf19c3d1b223e151f99a4f9f0690dca41414773390fc824184ac833e1" +dependencies = [ + "cfg-if", + "dirs-sys-next", +] + [[package]] name = "dirs-sys" version = "0.3.7" @@ -2111,6 +2133,17 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "dirs-sys-next" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4ebda144c4fe02d1f7ea1a7d9641b6fc6b580adcfa024ae48797ecdeb6825b4d" +dependencies = [ + "libc", + "redox_users 0.4.5", + "winapi", +] + [[package]] name = "dispatch2" version = "0.3.1" @@ -2289,6 +2322,7 @@ dependencies = [ "hickory-resolver", "hickory-server", "hmac", + "hotpath", "http", "http_req", "humansize", @@ -2630,6 +2664,12 @@ version = "1.2.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4ef6b89e5b37196644d8796de5268852ff179b44e96276cf4290264843743bb7" +[[package]] +name = "encode_unicode" +version = "1.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34aa73646ffb006b8f5147f3dc182bd4bcb190227ce861fc4a4844bf8e3cb2c0" + [[package]] name = "encoding" version = "0.2.33" @@ -3888,6 +3928,61 @@ dependencies = [ "windows-sys 0.52.0", ] +[[package]] +name = "hotpath" +version = "0.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc2c28b1fa962e433f800ed1ea0bf53dc028d3745cf2acec6cfd28b65ac96afa" +dependencies = [ + "arc-swap", + "cfg-if", + "crossbeam-channel", + "flate2", + "flume 0.12.0", + "futures-util", + "hdrhistogram", + "hotpath-macros", + "hotpath-meta", + "libc", + "object", + "parking_lot", + "pin-project-lite", + "prettytable-rs", + "quanta", + "regex", + "rustc-demangle", + "serde", + "serde_json", + "tiny_http", + "tokio", +] + +[[package]] +name = "hotpath-macros" +version = "0.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a585238d8daf746e27df0f24d1bbdcd2410e9febff63f9a0173f90d7e71c50f6" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.117", +] + +[[package]] +name = "hotpath-macros-meta" +version = "0.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "309f63c2f755dead454dd4b3ea8ab5c947f14f8ea435fbcd37fa820e17290e80" + +[[package]] +name = "hotpath-meta" +version = "0.18.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "68faa91a9e1114dff668cd90560f332da6bbde40dae37ec28ea1c43ca5ce3be3" +dependencies = [ + "hotpath-macros-meta", +] + [[package]] name = "html5ever" version = "0.29.1" @@ -4495,6 +4590,17 @@ dependencies = [ "once_cell", ] +[[package]] +name = "is-terminal" +version = "0.4.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" +dependencies = [ + "hermit-abi", + "libc", + "windows-sys 0.61.2", +] + [[package]] name = "is-wsl" version = "0.4.0" @@ -5613,7 +5719,7 @@ version = "0.7.6" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "680998035259dcfcafe653688bf2aa6d3e2dc05e98be6ab46afb089dc84f1df8" dependencies = [ - "proc-macro-crate 2.0.0", + "proc-macro-crate 3.5.0", "proc-macro2", "quote", "syn 2.0.117", @@ -5845,6 +5951,15 @@ dependencies = [ "objc2-foundation", ] +[[package]] +name = "object" +version = "0.36.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "62948e14d923ea95ea2c7c86c71013138b66525b86bdc08d2dcc262bdb497b87" +dependencies = [ + "memchr", +] + [[package]] name = "once_cell" version = "1.21.3" @@ -6715,6 +6830,19 @@ dependencies = [ "syn 2.0.117", ] +[[package]] +name = "prettytable-rs" +version = "0.10.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eea25e07510aa6ab6547308ebe3c036016d162b8da920dbb079e3ba8acf3d95a" +dependencies = [ + "encode_unicode", + "is-terminal", + "lazy_static", + "term", + "unicode-width 0.1.11", +] + [[package]] name = "primeorder" version = "0.13.6" @@ -7016,6 +7144,21 @@ version = "0.1.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b5a041e753da8b807c9255f28de81879c78c876392ff2469cde94799b2896b9d" +[[package]] +name = "quanta" +version = "0.12.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f3ab5a9d756f0d97bdc89019bd2e4ea098cf9cde50ee7564dde6b81ccc8f06c7" +dependencies = [ + "crossbeam-utils", + "libc", + "once_cell", + "raw-cpuid", + "wasi 0.11.0+wasi-snapshot-preview1", + "web-sys", + "winapi", +] + [[package]] name = "quick-error" version = "2.0.1" @@ -7269,6 +7412,15 @@ dependencies = [ "rand_core 0.5.1", ] +[[package]] +name = "raw-cpuid" +version = "11.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "498cd0dc59d73224351ee52a95fee0f1a617a2eae0e7d9d720cc622c73a54186" +dependencies = [ + "bitflags 2.8.0", +] + [[package]] name = "raw-window-handle" version = "0.6.2" @@ -7750,6 +7902,12 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "rustc-demangle" +version = "0.1.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b50b8869d9fc858ce7266cce0194bd74df58b9d0e3f6df3a9fc8eb470d95c09d" + [[package]] name = "rustc-hash" version = "2.1.0" @@ -9674,6 +9832,17 @@ dependencies = [ "utf-8", ] +[[package]] +name = "term" +version = "0.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c59df8ac95d96ff9bede18eb7300b0fda5e5d8d90960e76f8e14ae765eedbf1f" +dependencies = [ + "dirs-next", + "rustversion", + "winapi", +] + [[package]] name = "terminal_size" version = "0.4.1" @@ -9829,6 +9998,18 @@ version = "1.0.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "825f6c8a18bc36d56a62f66af7296385b628c9c5543a8663d4c217fc920bfefd" +[[package]] +name = "tiny_http" +version = "0.12.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "389915df6413a2e74fb181895f933386023c71110878cd0825588928e64cdc82" +dependencies = [ + "ascii", + "chunked_transfer", + "httpdate", + "log", +] + [[package]] name = "tinystr" version = "0.7.6" diff --git a/easytier/Cargo.toml b/easytier/Cargo.toml index 8c83246a..c61cf02b 100644 --- a/easytier/Cargo.toml +++ b/easytier/Cargo.toml @@ -52,6 +52,7 @@ toml = "0.8.12" chrono = { version = "0.4.37", features = ["serde"] } guarden = "0.2" +hotpath = { version = "0.18", default-features = false, optional = true } delegate = "0.13.5" @@ -401,6 +402,15 @@ jemalloc-prof = [ "jemalloc-sys/stats", ] tracing = ["tokio/tracing", "dep:console-subscriber"] +hotpath = [ + "dep:hotpath", + "hotpath/hotpath", + "hotpath/tokio", + "hotpath/parking_lot", + "hotpath/flume", +] +hotpath-cpu = ["hotpath", "hotpath/hotpath-cpu"] +hotpath-alloc = ["hotpath", "hotpath/hotpath-alloc"] magic-dns = ["dep:hickory-client", "dep:hickory-server"] faketcp = ["dep:flume"] zstd = ["dep:zstd"] diff --git a/easytier/src/easytier-core.rs b/easytier/src/easytier-core.rs index 03cf30b4..1bc5d47a 100644 --- a/easytier/src/easytier-core.rs +++ b/easytier/src/easytier-core.rs @@ -1,5 +1,11 @@ use easytier::core; +#[cfg(all( + feature = "hotpath-alloc", + any(feature = "jemalloc", feature = "mimalloc") +))] +compile_error!("feature `hotpath-alloc` cannot be enabled together with `jemalloc` or `mimalloc`"); + #[cfg(all(feature = "mimalloc", not(feature = "jemalloc")))] use mimalloc::MiMalloc; @@ -24,6 +30,16 @@ pub static malloc_conf: &[u8] = b"retain:false\0"; rust_i18n::i18n!("locales", fallback = "en"); #[tokio::main(flavor = "current_thread")] +#[cfg_attr( + all( + feature = "hotpath", + not(all( + feature = "hotpath-alloc", + any(feature = "jemalloc", feature = "mimalloc") + )) + ), + hotpath::main +)] async fn main() -> std::process::ExitCode { core::main().await } diff --git a/easytier/src/hotpath_off.rs b/easytier/src/hotpath_off.rs new file mode 100644 index 00000000..0ad97020 --- /dev/null +++ b/easytier/src/hotpath_off.rs @@ -0,0 +1,40 @@ +//! No-op stand-in for the `hotpath` macros used by this crate, selected when +//! the `hotpath` feature is disabled. +//! +//! Keeping `hotpath` as an optional dependency means default builds do not pull +//! the profiler (or any of its transitive dependencies) into the dependency +//! graph. These macros expand to their input unchanged, mirroring `hotpath`'s +//! own disabled mode so call sites compile identically with or without the +//! feature. +//! +//! The macros are `#[macro_export]`-ed so that `lib.rs`' `extern crate self as +//! hotpath` alias exposes them through the same `hotpath::...` paths used when +//! the feature is enabled. + +/// No-op mirroring `hotpath::channel!`: returns the channel expression +/// unchanged (dropping any optional trailing `label`/`log`/`capacity` args). +#[doc(hidden)] +#[macro_export] +macro_rules! channel { + ($expr:expr $(, $($rest:tt)*)?) => { + $expr + }; +} + +/// No-op mirroring `hotpath::mutex!`: returns the expression unchanged. +#[doc(hidden)] +#[macro_export] +macro_rules! mutex { + ($expr:expr $(, $($rest:tt)*)?) => { + $expr + }; +} + +/// No-op mirroring `hotpath::rw_lock!`: returns the expression unchanged. +#[doc(hidden)] +#[macro_export] +macro_rules! rw_lock { + ($expr:expr $(, $($rest:tt)*)?) => { + $expr + }; +} diff --git a/easytier/src/lib.rs b/easytier/src/lib.rs index 92ef893a..2d8f6c7b 100644 --- a/easytier/src/lib.rs +++ b/easytier/src/lib.rs @@ -5,6 +5,15 @@ use std::io; use clap::Command; use clap_complete::{Generator, Shell}; +// When the `hotpath` feature is off, alias the current crate as `hotpath` so +// call sites keep using `hotpath::...` paths, and provide a local no-op shim +// for the profiling macros. This keeps `hotpath` an optional dependency: the +// profiler is absent from the dependency graph entirely in default builds. +#[cfg(not(feature = "hotpath"))] +extern crate self as hotpath; +#[cfg(not(feature = "hotpath"))] +mod hotpath_off; + mod arch; mod gateway; pub mod instance; diff --git a/easytier/src/peers/mod.rs b/easytier/src/peers/mod.rs index c94a65de..debf7e9d 100644 --- a/easytier/src/peers/mod.rs +++ b/easytier/src/peers/mod.rs @@ -59,7 +59,7 @@ type BoxNicPacketFilter = Box; pub type PacketRecvChan = tokio::sync::mpsc::Sender; pub type PacketRecvChanReceiver = tokio::sync::mpsc::Receiver; pub fn create_packet_recv_chan() -> (PacketRecvChan, PacketRecvChanReceiver) { - tokio::sync::mpsc::channel(128) + hotpath::channel!(tokio::sync::mpsc::channel(128)) } pub async fn recv_packet_from_chan( packet_recv_chan_receiver: &mut PacketRecvChanReceiver, diff --git a/easytier/src/peers/peer.rs b/easytier/src/peers/peer.rs index 63af93f8..fd543512 100644 --- a/easytier/src/peers/peer.rs +++ b/easytier/src/peers/peer.rs @@ -207,6 +207,7 @@ impl Peer { .map(|conn| conn.clone()) } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "Peer"))] pub async fn send_msg(&self, msg: ZCPacket) -> Result<(), Error> { let Some(conn) = self.select_conn().await else { return Err(Error::PeerNoConnectionError(self.peer_node_id)); diff --git a/easytier/src/peers/peer_conn.rs b/easytier/src/peers/peer_conn.rs index 90374286..2e640547 100644 --- a/easytier/src/peers/peer_conn.rs +++ b/easytier/src/peers/peer_conn.rs @@ -10,6 +10,15 @@ use std::{ }, }; +#[cfg(feature = "hotpath")] +use hotpath::wrap::std::sync::Mutex as StdMutex; +#[cfg(feature = "hotpath")] +use hotpath::wrap::tokio::sync::Mutex; +#[cfg(not(feature = "hotpath"))] +use std::sync::Mutex as StdMutex; +#[cfg(not(feature = "hotpath"))] +use tokio::sync::Mutex; + use base64::Engine as _; use base64::engine::general_purpose::STANDARD as BASE64_STANDARD; use guarden::guard; @@ -17,7 +26,7 @@ use hmac::Mac; use prost::Message; use tokio::{ - sync::{Mutex, broadcast}, + sync::broadcast, task::JoinSet, time::{Duration, timeout}, }; @@ -98,7 +107,7 @@ struct PeerSessionTunnelFilter { enabled: bool, my_peer_id: Arc>, peer_id: Arc>>, - session: Arc>>>, + session: Arc>>>, } impl PeerSessionTunnelFilter { @@ -107,7 +116,7 @@ impl PeerSessionTunnelFilter { enabled, my_peer_id: Arc::new(AtomicCell::new(PeerId::default())), peer_id: Arc::new(AtomicCell::new(None)), - session: Arc::new(std::sync::Mutex::new(None)), + session: Arc::new(hotpath::mutex!(std::sync::Mutex::new(None))), } } @@ -116,7 +125,7 @@ impl PeerSessionTunnelFilter { enabled, my_peer_id: Arc::new(AtomicCell::new(my_peer_id)), peer_id: Arc::new(AtomicCell::new(None)), - session: Arc::new(std::sync::Mutex::new(None)), + session: Arc::new(hotpath::mutex!(std::sync::Mutex::new(None))), } } @@ -379,11 +388,12 @@ impl PeerConn { session_filter, noise_handshake_result: None, - tunnel: Arc::new(Mutex::new(Box::new( + tunnel: Arc::new(hotpath::mutex!(tokio::sync::Mutex::new(Box::new( guard!([mut mpsc_tunnel] mpsc_tunnel.close()), - ))), + ) + as Box))), sink, - recv: Mutex::new(Some(recv)), + recv: hotpath::mutex!(tokio::sync::Mutex::new(Some(recv))), tunnel_info, tasks: JoinSet::new(), @@ -1460,6 +1470,7 @@ impl PeerConn { }); } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerConn"))] pub async fn send_msg(&self, msg: ZCPacket) -> Result<(), Error> { Ok(self.sink.send(msg).await?) } diff --git a/easytier/src/peers/peer_manager.rs b/easytier/src/peers/peer_manager.rs index 17144d4c..f7028962 100644 --- a/easytier/src/peers/peer_manager.rs +++ b/easytier/src/peers/peer_manager.rs @@ -10,11 +10,12 @@ use std::{ time::{Duration, Instant, SystemTime}, }; +#[cfg(feature = "hotpath")] +use hotpath::wrap::tokio::sync::{Mutex, RwLock}; +#[cfg(not(feature = "hotpath"))] +use tokio::sync::{Mutex, RwLock}; use tokio::{ - sync::{ - Mutex, RwLock, - mpsc::{self, UnboundedReceiver, UnboundedSender}, - }, + sync::mpsc::{self, UnboundedReceiver, UnboundedSender}, task::JoinSet, }; @@ -277,8 +278,8 @@ impl PeerManager { let rpc_tspt = Arc::new(RpcTransport { my_peer_id, peers: Arc::downgrade(&peers), - foreign_peers: Mutex::new(None), - packet_recv: Mutex::new(peer_rpc_tspt_recv), + foreign_peers: hotpath::mutex!(tokio::sync::Mutex::new(None)), + packet_recv: hotpath::mutex!(tokio::sync::Mutex::new(peer_rpc_tspt_recv)), peer_rpc_tspt_sender, encryptor: encryptor.clone(), is_secure_mode_enabled, @@ -410,17 +411,21 @@ impl PeerManager { global_ctx, nic_channel, - tasks: Mutex::new(JoinSet::new()), + tasks: hotpath::mutex!(tokio::sync::Mutex::new(JoinSet::new())), - packet_recv: Arc::new(Mutex::new(Some(packet_recv))), + packet_recv: Arc::new(hotpath::mutex!(tokio::sync::Mutex::new(Some(packet_recv)))), peers, peer_rpc_mgr, peer_rpc_tspt: rpc_tspt, - peer_packet_process_pipeline: Arc::new(RwLock::new(Vec::new())), - nic_packet_process_pipeline: Arc::new(RwLock::new(Vec::new())), + peer_packet_process_pipeline: Arc::new(hotpath::rw_lock!(tokio::sync::RwLock::new( + Vec::new() + ))), + nic_packet_process_pipeline: Arc::new(hotpath::rw_lock!(tokio::sync::RwLock::new( + Vec::new() + ))), route_algo_inst, @@ -431,7 +436,7 @@ impl PeerManager { encryptor, data_compress_algo, - exit_nodes: RwLock::new(exit_nodes), + exit_nodes: hotpath::rw_lock!(tokio::sync::RwLock::new(exit_nodes)), reserved_my_peer_id_map: DashMap::new(), recent_have_traffic: Arc::new(DashMap::new()), @@ -1438,6 +1443,7 @@ impl PeerManager { self.get_route().get_foreign_network_summary().await } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerManager"))] async fn run_nic_packet_process_pipeline(&self, data: &mut ZCPacket) -> bool { // Enforce ACL for outbound (NIC-originated) packets. If ACL denies, stop processing. if !self.global_ctx.get_acl_filter().process_packet_with_acl( @@ -1523,6 +1529,7 @@ impl PeerManager { result } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerManager"))] async fn send_msg_internal( peers: &Arc, foreign_network_client: &Arc, diff --git a/easytier/src/peers/peer_map.rs b/easytier/src/peers/peer_map.rs index 46e55aaf..a9d32676 100644 --- a/easytier/src/peers/peer_map.rs +++ b/easytier/src/peers/peer_map.rs @@ -132,6 +132,7 @@ impl PeerMap { peer_id == self.my_peer_id || self.peer_map.contains_key(&peer_id) } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerMap"))] pub async fn send_msg_directly(&self, msg: ZCPacket, dst_peer_id: PeerId) -> Result<(), Error> { if dst_peer_id == self.my_peer_id { let packet_send = self.packet_send.clone(); @@ -163,6 +164,7 @@ impl PeerMap { Ok(()) } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerMap"))] pub async fn get_gateway_peer_id( &self, dst_peer_id: PeerId, diff --git a/easytier/src/peers/peer_ospf_route.rs b/easytier/src/peers/peer_ospf_route.rs index b790aa0b..7c5916d5 100644 --- a/easytier/src/peers/peer_ospf_route.rs +++ b/easytier/src/peers/peer_ospf_route.rs @@ -1393,6 +1393,7 @@ impl RouteTable { } } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "RouteTable"))] fn get_next_hop(&self, dst_peer_id: PeerId) -> Option { if self.suppressed_peer_ids.contains_key(&dst_peer_id) { return None; @@ -1400,6 +1401,7 @@ impl RouteTable { self.get_topology_next_hop(dst_peer_id) } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "RouteTable"))] fn get_topology_next_hop(&self, dst_peer_id: PeerId) -> Option { let cur_version = self.next_hop_map_version.get(); self.next_hop_map.get(&dst_peer_id).and_then(|x| { diff --git a/easytier/src/peers/peer_session.rs b/easytier/src/peers/peer_session.rs index 34599ead..e15a1710 100644 --- a/easytier/src/peers/peer_session.rs +++ b/easytier/src/peers/peer_session.rs @@ -337,6 +337,7 @@ impl PeerSession { } } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerSession"))] pub fn encrypt_payload( &self, sender_peer_id: PeerId, @@ -350,6 +351,7 @@ impl PeerSession { .encrypt_payload(Self::dir_for_sender(sender_peer_id, receiver_peer_id), pkt) } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "PeerSession"))] pub fn decrypt_payload( &self, sender_peer_id: PeerId, diff --git a/easytier/src/peers/secure_datagram.rs b/easytier/src/peers/secure_datagram.rs index f4722e44..cf89acdf 100644 --- a/easytier/src/peers/secure_datagram.rs +++ b/easytier/src/peers/secure_datagram.rs @@ -701,6 +701,10 @@ impl SecureDatagramSession { false } + #[cfg_attr( + feature = "hotpath", + hotpath::measure(impl_type = "SecureDatagramSession") + )] pub fn encrypt_payload( &self, dir: SecureDatagramDirection, @@ -719,6 +723,10 @@ impl SecureDatagramSession { Ok(()) } + #[cfg_attr( + feature = "hotpath", + hotpath::measure(impl_type = "SecureDatagramSession") + )] pub fn decrypt_payload( &self, dir: SecureDatagramDirection, diff --git a/easytier/src/tunnel/fake_tcp/stack.rs b/easytier/src/tunnel/fake_tcp/stack.rs index a7f1b779..3e6a4056 100644 --- a/easytier/src/tunnel/fake_tcp/stack.rs +++ b/easytier/src/tunnel/fake_tcp/stack.rs @@ -159,7 +159,7 @@ impl Socket { ack: Option, state: State, ) -> (Socket, flume::Sender) { - let (incoming_tx, incoming_rx) = flume::bounded(MPMC_BUFFER_LEN); + let (incoming_tx, incoming_rx) = hotpath::channel!(flume::bounded(MPMC_BUFFER_LEN)); ( Socket { diff --git a/easytier/src/tunnel/mpsc.rs b/easytier/src/tunnel/mpsc.rs index e15231ae..fb4f68c9 100644 --- a/easytier/src/tunnel/mpsc.rs +++ b/easytier/src/tunnel/mpsc.rs @@ -43,7 +43,7 @@ pub struct MpscTunnel { impl MpscTunnel { pub fn new(tunnel: T, send_timeout: Option) -> Self { - let (tx, mut rx) = channel(32); + let (tx, mut rx) = hotpath::channel!(channel(32)); let (stream, mut sink) = tunnel.split(); let task = tokio::spawn(async move { @@ -66,6 +66,7 @@ impl MpscTunnel { } } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "MpscTunnel"))] async fn forward_one_round( rx: &mut Receiver, sink: &mut Pin>, @@ -79,6 +80,7 @@ impl MpscTunnel { } } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "MpscTunnel"))] async fn forward_one_round_no_timeout( rx: &mut Receiver, sink: &mut Pin>, @@ -96,6 +98,7 @@ impl MpscTunnel { sink.flush().await } + #[cfg_attr(feature = "hotpath", hotpath::measure(impl_type = "MpscTunnel"))] async fn forward_one_round_with_timeout( rx: &mut Receiver, sink: &mut Pin>, diff --git a/easytier/src/tunnel/ring.rs b/easytier/src/tunnel/ring.rs index 2573bdcc..43f8a783 100644 --- a/easytier/src/tunnel/ring.rs +++ b/easytier/src/tunnel/ring.rs @@ -196,7 +196,8 @@ pub struct RingTunnelListener { impl RingTunnelListener { pub fn new(key: url::Url) -> Self { - let (conn_sender, conn_receiver) = tokio::sync::mpsc::unbounded_channel(); + let (conn_sender, conn_receiver) = + hotpath::channel!(tokio::sync::mpsc::unbounded_channel()); RingTunnelListener { listener_addr: key, conn_sender, diff --git a/easytier/src/tunnel/udp.rs b/easytier/src/tunnel/udp.rs index b7ebd979..41f5a52d 100644 --- a/easytier/src/tunnel/udp.rs +++ b/easytier/src/tunnel/udp.rs @@ -572,8 +572,9 @@ pub struct UdpTunnelListener { impl UdpTunnelListener { pub fn new(addr: url::Url) -> Self { - let (close_event_send, close_event_recv) = tokio::sync::mpsc::unbounded_channel(); - let (conn_send, conn_recv) = tokio::sync::mpsc::channel(100); + let (close_event_send, close_event_recv) = + hotpath::channel!(tokio::sync::mpsc::unbounded_channel()); + let (conn_send, conn_recv) = hotpath::channel!(tokio::sync::mpsc::channel(100)); Self { addr: addr.clone(), socket: None, @@ -784,7 +785,8 @@ impl UdpTunnelConnector { "udp build tunnel for connector" ); - let (close_event_sender, mut close_event_recv) = tokio::sync::mpsc::unbounded_channel(); + let (close_event_sender, mut close_event_recv) = + hotpath::channel!(tokio::sync::mpsc::unbounded_channel()); let ring_recv = RingStream::new(ring_for_send_udp.clone()); let ring_sender = RingSink::new(ring_for_recv_udp.clone()); diff --git a/easytier/src/tunnel/wireguard.rs b/easytier/src/tunnel/wireguard.rs index e0711f2e..443fe2f2 100644 --- a/easytier/src/tunnel/wireguard.rs +++ b/easytier/src/tunnel/wireguard.rs @@ -468,7 +468,7 @@ pub struct WgTunnelListener { impl WgTunnelListener { pub fn new(addr: url::Url, config: WgConfig) -> Self { - let (conn_send, conn_recv) = tokio::sync::mpsc::unbounded_channel(); + let (conn_send, conn_recv) = hotpath::channel!(tokio::sync::mpsc::unbounded_channel()); WgTunnelListener { addr, config,