diff --git a/.github/workflows/core.yml b/.github/workflows/core.yml index 7e8f4ee3..b3e6b1ed 100644 --- a/.github/workflows/core.yml +++ b/.github/workflows/core.yml @@ -165,7 +165,10 @@ jobs: - uses: taiki-e/install-action@v2 if: ${{ contains(matrix.OS, 'ubuntu') }} with: - tool: cargo-zigbuild + # v0.23.3 emits -mcpu=generic+v6+strict_align for + # arm-unknown-linux-musleabi, which zig 0.16.0 rejects; + # unpin only together with a zig bump. + tool: cargo-zigbuild@0.23.2 - name: Build if: ${{ !contains(matrix.TARGET, 'mips') }} diff --git a/easytier-core/src/foundation/stats.rs b/easytier-core/src/foundation/stats.rs index fcbd25cb..228fcfa5 100644 --- a/easytier-core/src/foundation/stats.rs +++ b/easytier-core/src/foundation/stats.rs @@ -169,6 +169,15 @@ pub enum MetricName { /// Traffic packets forwarded for foreign network, forward TrafficPacketsForeignForwardForwarded, + /// Bytes accepted from one VPN portal client into the mesh + VpnPortalClientBytesTx, + /// Bytes delivered from the mesh to one VPN portal client + VpnPortalClientBytesRx, + /// Packets accepted from one VPN portal client into the mesh + VpnPortalClientPacketsTx, + /// Packets delivered from the mesh to one VPN portal client + VpnPortalClientPacketsRx, + /// UDP broadcast relay packets captured from the raw socket UdpBroadcastRelayPacketsCaptured, /// UDP broadcast relay packets ignored before forwarding @@ -260,6 +269,15 @@ impl fmt::Display for MetricName { write!(f, "traffic_packets_foreign_forward_forwarded") } + MetricName::VpnPortalClientBytesTx => write!(f, "vpn_portal_client_bytes_tx"), + MetricName::VpnPortalClientBytesRx => write!(f, "vpn_portal_client_bytes_rx"), + MetricName::VpnPortalClientPacketsTx => { + write!(f, "vpn_portal_client_packets_tx") + } + MetricName::VpnPortalClientPacketsRx => { + write!(f, "vpn_portal_client_packets_rx") + } + MetricName::UdpBroadcastRelayPacketsCaptured => { write!(f, "udp_broadcast_relay_packets_captured") } @@ -314,6 +332,8 @@ pub enum LabelType { DstIp(String), /// Mapped Dst Ip MappedDstIp(String), + /// Stable VPN portal client name + VpnPortalClient(String), } impl fmt::Display for LabelType { @@ -333,6 +353,7 @@ impl fmt::Display for LabelType { LabelType::Status(status) => write!(f, "status={}", status), LabelType::DstIp(ip) => write!(f, "dst_ip={}", ip), LabelType::MappedDstIp(ip) => write!(f, "mapped_dst_ip={}", ip), + LabelType::VpnPortalClient(client) => write!(f, "vpn_portal_client={}", client), } } } @@ -354,6 +375,7 @@ impl LabelType { LabelType::Status(_) => "status", LabelType::DstIp(_) => "dst_ip", LabelType::MappedDstIp(_) => "mapped_dst_ip", + LabelType::VpnPortalClient(_) => "vpn_portal_client", } } @@ -373,6 +395,7 @@ impl LabelType { LabelType::Status(status) => status.clone(), LabelType::DstIp(ip) => ip.clone(), LabelType::MappedDstIp(ip) => ip.clone(), + LabelType::VpnPortalClient(client) => client.clone(), } } } diff --git a/easytier-core/src/gateway/vpn_portal/runtime.rs b/easytier-core/src/gateway/vpn_portal/runtime.rs index d2a69567..d4ce3dce 100644 --- a/easytier-core/src/gateway/vpn_portal/runtime.rs +++ b/easytier-core/src/gateway/vpn_portal/runtime.rs @@ -22,6 +22,7 @@ use tokio_util::sync::CancellationToken; use crate::{ config::runtime::{CoreInstanceRuntimeConfig, CoreRuntimeConfigStore}, events::{CoreEvent, CoreEventSink}, + foundation::stats::{CounterHandle, LabelSet, LabelType, MetricName, StatsManager}, peers::{ attached::{AttachedPeerConfig, AttachedPeerRuntime}, peer_manager::PeerManagerCore, @@ -150,6 +151,38 @@ impl Default for ClientStatus { } } +#[derive(Clone)] +struct PortalClientTrafficMetrics { + upload_bytes: CounterHandle, + upload_packets: CounterHandle, + download_bytes: CounterHandle, + download_packets: CounterHandle, +} + +impl PortalClientTrafficMetrics { + fn new(stats: &StatsManager, network_name: &str, client_name: &str) -> Self { + let labels = LabelSet::new() + .with_label_type(LabelType::NetworkName(network_name.to_owned())) + .with_label_type(LabelType::VpnPortalClient(client_name.to_owned())); + Self { + upload_bytes: stats.get_counter(MetricName::VpnPortalClientBytesTx, labels.clone()), + upload_packets: stats.get_counter(MetricName::VpnPortalClientPacketsTx, labels.clone()), + download_bytes: stats.get_counter(MetricName::VpnPortalClientBytesRx, labels.clone()), + download_packets: stats.get_counter(MetricName::VpnPortalClientPacketsRx, labels), + } + } + + fn record_upload(&self, bytes: usize) { + self.upload_bytes.add(bytes as u64); + self.upload_packets.inc(); + } + + fn record_download(&self, bytes: usize) { + self.download_bytes.add(bytes as u64); + self.download_packets.inc(); + } +} + struct PortalRuntime { cancel: CancellationToken, tasks: JoinSet<()>, @@ -165,6 +198,7 @@ pub struct PortalModule { events: Arc, statuses: Arc>>, session_locks: Arc>>>>, + traffic_metrics: Arc>>, runtime: Mutex>, } @@ -179,6 +213,25 @@ impl PortalModule { if let Some(config) = config.as_ref() { validate_clients(config, runtime_config.snapshot().as_ref())?; } + let network_name = runtime_config + .snapshot() + .peer + .runtime + .network_identity + .network_name + .clone(); + let stats = peer_manager.stats_manager(); + let traffic_metrics = config + .as_ref() + .into_iter() + .flat_map(|config| &config.clients) + .map(|client| { + ( + client.name.clone(), + PortalClientTrafficMetrics::new(&stats, &network_name, &client.name), + ) + }) + .collect(); Ok(Arc::new(Self { operation: Mutex::new(()), peer_manager, @@ -188,6 +241,7 @@ impl PortalModule { events, statuses: Arc::new(RwLock::new(BTreeMap::new())), session_locks: Arc::new(RwLock::new(BTreeMap::new())), + traffic_metrics: Arc::new(StdRwLock::new(traffic_metrics)), runtime: Mutex::new(None), })) } @@ -249,11 +303,30 @@ impl PortalModule { }; let applied_clients = candidate.clients; + { + let network_name = self + .runtime_config + .snapshot() + .peer + .runtime + .network_identity + .network_name + .clone(); + let stats = self.peer_manager.stats_manager(); + let mut traffic_metrics = self.traffic_metrics.write().unwrap(); + traffic_metrics.retain(|name, _| applied.contains(name)); + for name in &applied { + traffic_metrics.entry(name.clone()).or_insert_with(|| { + PortalClientTrafficMetrics::new(&stats, &network_name, name) + }); + } + } + { let mut statuses = self.statuses.write().await; statuses.retain(|name, _| applied.contains(name)); - for name in applied { - statuses.entry(name).or_default(); + for name in &applied { + statuses.entry(name.clone()).or_default(); } } { @@ -320,6 +393,7 @@ impl PortalModule { self.config.clone().expect("checked above"), self.statuses.clone(), self.session_locks.clone(), + self.traffic_metrics.clone(), self.events.clone(), cancel.clone(), start_signal.clone(), @@ -347,6 +421,7 @@ impl PortalModule { config: Arc>, statuses: Arc>>, session_locks: Arc>>>>, + traffic_metrics: Arc>>, events: Arc, cancel: CancellationToken, start_signal: CancellationToken, @@ -371,6 +446,7 @@ impl PortalModule { config.clone(), statuses.clone(), session_locks.clone(), + traffic_metrics.clone(), events.clone(), cancel.clone(), )); @@ -397,6 +473,7 @@ impl PortalModule { config: Arc>, statuses: Arc>>, session_locks: Arc>>>>, + traffic_metrics: Arc>>, events: Arc, cancel: CancellationToken, ) { @@ -419,6 +496,12 @@ impl PortalModule { _ = cancel.cancelled() => return, guard = session_lock.lock() => guard, }; + let Some(traffic) = traffic_metrics.read().unwrap().get(&client.name).cloned() else { + // The client was removed by a concurrent config update while this + // session was waiting for the session lock. + tracing::warn!(client = %client.name, "VPN portal client removed before session started"); + return; + }; let generation = { let mut statuses = statuses.write().await; let status = statuses.entry(client.name.clone()).or_default(); @@ -482,6 +565,7 @@ impl PortalModule { let attached = attached.clone(); let name = client.name.clone(); let virtual_ip = client.virtual_ip.address(); + let traffic = traffic.clone(); tokio::spawn(async move { while let Some(payload) = client_stream.recv().await { if !has_ipv4_source(&payload, virtual_ip) { @@ -492,6 +576,7 @@ impl PortalModule { tracing::debug!(?error, client = %name, "attached peer send failed"); break; } + traffic.record_upload(payload.len()); } }) }; @@ -500,9 +585,11 @@ impl PortalModule { tokio::spawn(async move { while let Some(packet) = attached.recv_packet().await { let payload = packet.payload().to_vec(); + let bytes = payload.len(); if client_sink.send(payload).await.is_err() { break; } + traffic.record_download(bytes); } }) }; @@ -1051,6 +1138,24 @@ mod tests { } } + fn traffic_metrics( + peer_manager: &PeerManagerCore, + clients: &[&str], + ) -> Arc>> { + let stats = peer_manager.stats_manager(); + Arc::new(StdRwLock::new( + clients + .iter() + .map(|name| { + ( + (*name).to_owned(), + PortalClientTrafficMetrics::new(&stats, "portal-test", name), + ) + }) + .collect(), + )) + } + fn raw_ipv4(source: Ipv4Addr, destination: Ipv4Addr) -> Vec { let mut packet = vec![0u8; 28]; packet[0] = 0x45; @@ -1077,6 +1182,56 @@ mod tests { )); assert!(!has_ipv4_source(&[0u8; 8], assigned)); } + + #[tokio::test] + async fn portal_client_traffic_accumulates_across_sessions() { + let stats = StatsManager::new(); + let first = PortalClientTrafficMetrics::new(&stats, "portal-test", "client-a"); + first.record_upload(80); + first.record_download(120); + drop(first); + + let reconnected = PortalClientTrafficMetrics::new(&stats, "portal-test", "client-a"); + reconnected.record_upload(20); + reconnected.record_download(30); + + let labels = LabelSet::new() + .with_label_type(LabelType::NetworkName("portal-test".to_owned())) + .with_label_type(LabelType::VpnPortalClient("client-a".to_owned())); + assert_eq!( + stats + .get_metric(MetricName::VpnPortalClientBytesTx, &labels) + .unwrap() + .value, + 100 + ); + assert_eq!( + stats + .get_metric(MetricName::VpnPortalClientPacketsTx, &labels) + .unwrap() + .value, + 2 + ); + assert_eq!( + stats + .get_metric(MetricName::VpnPortalClientBytesRx, &labels) + .unwrap() + .value, + 150 + ); + assert_eq!( + stats + .get_metric(MetricName::VpnPortalClientPacketsRx, &labels) + .unwrap() + .value, + 2 + ); + assert!(stats.export_prometheus().contains( + "vpn_portal_client_bytes_tx{network_name=\"portal-test\",vpn_portal_client=\"client-a\"} 100" + )); + stats.stop_cleanup_task().await; + } + fn network_runtime() -> (Arc, CoreRuntimeConfigStore) { network_runtime_with_secure_mode(false) } @@ -1252,6 +1407,7 @@ mod tests { Arc::new(StdRwLock::new(config)), statuses.clone(), session_locks, + traffic_metrics(&peer_manager, &["alice"]), Arc::new(()), cancel, )); @@ -1300,6 +1456,21 @@ mod tests { .await .expect("mesh packet was not delivered before the first client packet"); assert_eq!(&outbound[16..20], virtual_ip.octets().as_slice()); + let labels = LabelSet::new() + .with_label_type(LabelType::NetworkName("portal-test".to_owned())) + .with_label_type(LabelType::VpnPortalClient("alice".to_owned())); + assert!( + peer_manager + .stats_manager() + .get_metric(MetricName::VpnPortalClientBytesRx, &labels) + .is_some_and(|metric| metric.value >= mesh_packet.len() as u64) + ); + assert!( + peer_manager + .stats_manager() + .get_metric(MetricName::VpnPortalClientPacketsRx, &labels) + .is_some_and(|metric| metric.value >= 1) + ); endpoint_sender.send("portal://roamed".to_owned()).unwrap(); tokio::time::timeout(std::time::Duration::from_secs(5), async { @@ -1361,6 +1532,7 @@ mod tests { Arc::new(StdRwLock::new(config)), statuses.clone(), session_locks, + traffic_metrics(&peer_manager, &["alice"]), events.clone(), cancel.clone(), )); @@ -1467,14 +1639,13 @@ mod tests { Arc::new(StdRwLock::new(config)), statuses.clone(), session_locks, + traffic_metrics(&peer_manager, &["alice"]), events.clone(), CancellationToken::new(), )); - to_runtime - .send(raw_ipv4(virtual_ip, Ipv4Addr::new(10, 82, 0, 1))) - .await - .unwrap(); + let client_packet = raw_ipv4(virtual_ip, Ipv4Addr::new(10, 82, 0, 1)); + to_runtime.send(client_packet.clone()).await.unwrap(); let attached_peer_id = tokio::time::timeout(std::time::Duration::from_secs(5), async { loop { let status = statuses.read().await.get("alice").cloned().unwrap(); @@ -1517,6 +1688,26 @@ mod tests { } task.await.unwrap(); + let labels = LabelSet::new() + .with_label_type(LabelType::NetworkName("portal-test".to_owned())) + .with_label_type(LabelType::VpnPortalClient("alice".to_owned())); + assert_eq!( + peer_manager + .stats_manager() + .get_metric(MetricName::VpnPortalClientBytesTx, &labels) + .unwrap() + .value, + client_packet.len() as u64 + ); + assert_eq!( + peer_manager + .stats_manager() + .get_metric(MetricName::VpnPortalClientPacketsTx, &labels) + .unwrap() + .value, + 1 + ); + let status = statuses.read().await.get("alice").cloned().unwrap(); assert_eq!(status.state, PortalClientState::Offline); assert!(status.peer_id.is_none());