diff --git a/easytier-contrib/easytier-ohrs/src/exports/runtime_api.rs b/easytier-contrib/easytier-ohrs/src/exports/runtime_api.rs index e28368e0..f1792a03 100644 --- a/easytier-contrib/easytier-ohrs/src/exports/runtime_api.rs +++ b/easytier-contrib/easytier-ohrs/src/exports/runtime_api.rs @@ -118,11 +118,7 @@ pub(crate) fn set_tun_fd( }) } -pub(crate) fn get_runtime_snapshot() -> RuntimeAggregateState { - get_runtime_snapshot_inner() -} - -pub(crate) fn get_runtime_snapshot_inner() -> RuntimeAggregateState { +pub(crate) fn collect_runtime_state() -> RuntimeAggregateState { let infos = match ASYNC_RUNTIME.block_on(INSTANCE_MANAGER.collect_network_infos()) { Ok(infos) => infos, Err(err) => { diff --git a/easytier-contrib/easytier-ohrs/src/kernel_bridge.rs b/easytier-contrib/easytier-ohrs/src/kernel_bridge.rs index c415e930..f01868b0 100644 --- a/easytier-contrib/easytier-ohrs/src/kernel_bridge.rs +++ b/easytier-contrib/easytier-ohrs/src/kernel_bridge.rs @@ -3,6 +3,4 @@ mod routing; mod socket_server; pub(crate) use routing::aggregate_requested_tun_routes; -pub use socket_server::{ - set_snapshot_broadcast_enabled, start_local_socket_server, stop_local_socket_server, -}; +pub use socket_server::{start_local_socket_server, stop_local_socket_server}; diff --git a/easytier-contrib/easytier-ohrs/src/kernel_bridge/protocol.rs b/easytier-contrib/easytier-ohrs/src/kernel_bridge/protocol.rs index 5d503232..e6f3eea1 100644 --- a/easytier-contrib/easytier-ohrs/src/kernel_bridge/protocol.rs +++ b/easytier-contrib/easytier-ohrs/src/kernel_bridge/protocol.rs @@ -32,6 +32,13 @@ pub(crate) fn send_local_socket_message( Ok(()) } +fn shrink_clients_if_sparse(clients: &mut Vec) { + let sparse_limit = clients.len().saturating_mul(2).max(4); + if clients.capacity() > sparse_limit { + clients.shrink_to_fit(); + } +} + pub(crate) fn broadcast_local_socket_message( clients: &mut Vec, message_type: &str, @@ -45,6 +52,7 @@ pub(crate) fn broadcast_local_socket_message( active_clients.push(client); } } + shrink_clients_if_sparse(&mut active_clients); *clients = active_clients; delivered } @@ -79,6 +87,7 @@ pub(crate) fn broadcast_local_socket_json_payload_message( active_clients.push(client); } } + shrink_clients_if_sparse(&mut active_clients); *clients = active_clients; delivered } diff --git a/easytier-contrib/easytier-ohrs/src/kernel_bridge/socket_server.rs b/easytier-contrib/easytier-ohrs/src/kernel_bridge/socket_server.rs index a992901c..26d898fb 100644 --- a/easytier-contrib/easytier-ohrs/src/kernel_bridge/socket_server.rs +++ b/easytier-contrib/easytier-ohrs/src/kernel_bridge/socket_server.rs @@ -1,13 +1,20 @@ use super::protocol::{ TunRequestPayload, broadcast_local_socket_json_payload_message, broadcast_local_socket_message, }; -use crate::INSTANCE_MANAGER; +use crate::collect_runtime_state_inner; use crate::config::repository::kernel_socket_path; -use crate::get_runtime_snapshot_inner; use crate::kernel_bridge::routing::aggregate_tun_routes; +use crate::runtime::state::runtime_state::{ + PeerConnInfo as RuntimePeerConnInfo, RuntimeAggregateState, peer_conn_to_view, +}; +use crate::{ASYNC_RUNTIME, INSTANCE_MANAGER}; use easytier::common::global_ctx::{EventBusSubscriber, GlobalCtxEvent}; +use easytier::proto::api::instance::ListPeerRequest; +use easytier::proto::rpc_types::controller::BaseController; use once_cell::sync::Lazy; +use serde::Serialize; use std::collections::{HashMap, HashSet}; +use std::hash::Hash; use std::io::ErrorKind; use std::os::unix::net::{UnixListener, UnixStream}; use std::path::PathBuf; @@ -23,13 +30,75 @@ struct LocalSocketState { } static LOCAL_SOCKET_STATE: Lazy>> = Lazy::new(|| Mutex::new(None)); -static SNAPSHOT_BROADCAST_ENABLED: AtomicBool = AtomicBool::new(true); const SOCKET_TICK_INTERVAL: Duration = Duration::from_millis(250); +const TRAFFIC_STATS_INTERVAL: Duration = Duration::from_secs(1); +const INSTANCE_POLL_INTERVAL: Duration = Duration::from_secs(1); const TUN_FAST_CHECK_WINDOW: Duration = Duration::from_secs(8); const EVENT_RECEIVER_SYNC_INTERVAL: Duration = Duration::from_secs(1); -pub fn set_snapshot_broadcast_enabled(enabled: bool) { - SNAPSHOT_BROADCAST_ENABLED.store(enabled, Ordering::Relaxed); +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +struct TrafficStatsPayload { + instances: Vec, +} + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +struct InstanceTrafficStats { + config_id: String, + instance_id: String, + rx_bytes: i64, + tx_bytes: i64, + peers: Vec, +} + +#[derive(Debug, Clone, Serialize)] +#[serde(rename_all = "camelCase")] +struct PeerTrafficStats { + peer_id: i64, + rx_bytes: i64, + tx_bytes: i64, + total_bytes: i64, + latency_us: i64, + loss_rate: f64, +} + +struct PendingPeerEvent { + event: &'static str, + instance_id: String, + peer_id: i64, + conn: Option, +} + +#[derive(Default)] +struct DrainedKernelEvents { + tun_refresh: bool, + topology_lost: bool, + peer_events: Vec, +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct RuntimePeerEventPayload { + event: &'static str, + config_id: String, + instance_id: String, + peer_id: i64, + conn: Option, +} + +fn shrink_hash_map_if_sparse(map: &mut HashMap) { + let sparse_limit = map.len().saturating_mul(2).max(8); + if map.capacity() > sparse_limit { + map.shrink_to_fit(); + } +} + +fn shrink_hash_set_if_sparse(set: &mut HashSet) { + let sparse_limit = set.len().saturating_mul(2).max(8); + if set.capacity() > sparse_limit { + set.shrink_to_fit(); + } } fn sync_tun_event_receivers(receivers: &mut HashMap) { @@ -44,36 +113,65 @@ fn sync_tun_event_receivers(receivers: &mut HashMap) } } receivers.retain(|instance_id, _| active_instance_ids.contains(instance_id)); + shrink_hash_map_if_sparse(receivers); } fn event_needs_tun_refresh(event: &GlobalCtxEvent) -> bool { matches!( event, - GlobalCtxEvent::DhcpIpv4Changed(_, _) - | GlobalCtxEvent::DhcpIpv4Conflicted(_) - | GlobalCtxEvent::PublicIpv6Changed(_, _) - | GlobalCtxEvent::PublicIpv6RoutesUpdated(_, _) - | GlobalCtxEvent::ProxyCidrsUpdated(_, _) - | GlobalCtxEvent::ConfigPatched(_) - | GlobalCtxEvent::PeerAdded(_) - | GlobalCtxEvent::PeerRemoved(_) - | GlobalCtxEvent::PeerConnAdded(_) - | GlobalCtxEvent::PeerConnRemoved(_) + GlobalCtxEvent::DhcpIpv4Changed(_, _) | GlobalCtxEvent::ProxyCidrsUpdated(_, _) ) } -fn drain_tun_refresh_events(receivers: &mut HashMap) -> bool { - let mut refresh_needed = false; +fn drain_kernel_events(receivers: &mut HashMap) -> DrainedKernelEvents { + let mut drained = DrainedKernelEvents::default(); let mut closed_receivers = Vec::new(); for (instance_id, receiver) in receivers.iter_mut() { loop { match receiver.try_recv() { Ok(event) => { - refresh_needed = event_needs_tun_refresh(&event) || refresh_needed; + drained.tun_refresh = event_needs_tun_refresh(&event) || drained.tun_refresh; + match event { + GlobalCtxEvent::PeerAdded(peer_id) => { + drained.peer_events.push(PendingPeerEvent { + event: "peer_added", + instance_id: instance_id.clone(), + peer_id: peer_id as i64, + conn: None, + }); + } + GlobalCtxEvent::PeerRemoved(peer_id) => { + drained.peer_events.push(PendingPeerEvent { + event: "peer_removed", + instance_id: instance_id.clone(), + peer_id: peer_id as i64, + conn: None, + }); + } + GlobalCtxEvent::PeerConnAdded(conn_info) => { + let peer_id = conn_info.peer_id as i64; + drained.peer_events.push(PendingPeerEvent { + event: "peer_conn_added", + instance_id: instance_id.clone(), + peer_id, + conn: Some(peer_conn_to_view(conn_info)), + }); + } + GlobalCtxEvent::PeerConnRemoved(conn_info) => { + let peer_id = conn_info.peer_id as i64; + drained.peer_events.push(PendingPeerEvent { + event: "peer_conn_removed", + instance_id: instance_id.clone(), + peer_id, + conn: Some(peer_conn_to_view(conn_info)), + }); + } + _ => {} + } } Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break, Err(tokio::sync::broadcast::error::TryRecvError::Lagged(_)) => { - refresh_needed = true; + drained.topology_lost = true; continue; } Err(tokio::sync::broadcast::error::TryRecvError::Closed) => { @@ -86,7 +184,124 @@ fn drain_tun_refresh_events(receivers: &mut HashMap) for instance_id in closed_receivers { receivers.remove(&instance_id); } - refresh_needed + drained +} + +fn broadcast_runtime_peer_events( + clients: &mut Vec, + peer_events: Vec, +) { + for event in peer_events { + let payload = RuntimePeerEventPayload { + event: event.event, + config_id: event.instance_id.clone(), + instance_id: event.instance_id, + peer_id: event.peer_id, + conn: event.conn, + }; + match serde_json::to_string(&payload) { + Ok(json) => { + let _ = broadcast_local_socket_json_payload_message( + clients, + "runtime_peer_event", + &json, + ); + } + Err(err) => { + ohrs_log_error!("[Rust] serialize runtime peer event failed: {}", err); + } + } + } +} + +fn tun_candidate_ids(snapshot: &RuntimeAggregateState) -> HashSet { + snapshot + .instances + .iter() + .filter(|instance| instance.running && instance.tun_required) + .map(|instance| instance.instance_id.clone()) + .collect() +} + +fn collect_traffic_stats() -> TrafficStatsPayload { + let services = INSTANCE_MANAGER + .iter() + .filter_map(|instance| { + instance + .value() + .get_api_service() + .map(|api_service| (instance.key().to_string(), api_service)) + }) + .collect::>(); + + let instances = ASYNC_RUNTIME.block_on(async { + let mut instances = Vec::new(); + for (instance_id, api_service) in services { + let peers = match api_service + .get_peer_manage_service() + .list_peer(BaseController::default(), ListPeerRequest::default()) + .await + { + Ok(response) => response.peer_infos, + Err(err) => { + ohrs_log_debug!( + "[Rust] collect traffic stats list_peer failed instance={}: {}", + instance_id, + err + ); + continue; + } + }; + + let mut instance_rx_bytes = 0i64; + let mut instance_tx_bytes = 0i64; + let mut peer_stats = Vec::with_capacity(peers.len()); + + for peer in peers { + let mut peer_rx_bytes = 0i64; + let mut peer_tx_bytes = 0i64; + let mut latency_us = i64::MAX; + let mut loss_rate = 0f64; + + for conn in peer.conns { + if let Some(stats) = conn.stats { + let rx_bytes = stats.rx_bytes as i64; + let tx_bytes = stats.tx_bytes as i64; + peer_rx_bytes += rx_bytes; + peer_tx_bytes += tx_bytes; + latency_us = latency_us.min(stats.latency_us as i64); + } + loss_rate = loss_rate.max(conn.loss_rate as f64); + } + + instance_rx_bytes += peer_rx_bytes; + instance_tx_bytes += peer_tx_bytes; + peer_stats.push(PeerTrafficStats { + peer_id: peer.peer_id as i64, + rx_bytes: peer_rx_bytes, + tx_bytes: peer_tx_bytes, + total_bytes: peer_rx_bytes + peer_tx_bytes, + latency_us: if latency_us == i64::MAX { + -1 + } else { + latency_us + }, + loss_rate, + }); + } + + instances.push(InstanceTrafficStats { + config_id: instance_id.clone(), + instance_id, + rx_bytes: instance_rx_bytes, + tx_bytes: instance_tx_bytes, + peers: peer_stats, + }); + } + instances + }); + + TrafficStatsPayload { instances } } pub fn start_local_socket_server() -> bool { @@ -131,21 +346,25 @@ pub fn start_local_socket_server() -> bool { let stop_flag = std::sync::Arc::new(AtomicBool::new(false)); let worker_stop_flag = stop_flag.clone(); let worker = thread::spawn(move || { - let mut last_snapshot_json = String::new(); + let mut last_topology_json = String::new(); let mut delivered_tun_requests = HashSet::new(); let mut last_tun_route_signatures = HashMap::::new(); let mut tun_fast_until = Instant::now() + TUN_FAST_CHECK_WINDOW; let mut tun_bootstrap_done = false; let mut last_event_receiver_sync_at: Option = None; + let mut last_traffic_stats_at: Option = None; + let mut last_instance_poll_at: Option = None; let mut tun_event_receivers = HashMap::::new(); let mut clients = Vec::::new(); while !worker_stop_flag.load(Ordering::Relaxed) { + let mut full_topology_dirty = false; let mut accepted_client = false; loop { match listener.accept() { Ok((stream, _addr)) => { accepted_client = true; + full_topology_dirty = true; clients.push(stream); tun_fast_until = Instant::now() + TUN_FAST_CHECK_WINDOW; tun_bootstrap_done = false; @@ -158,15 +377,21 @@ pub fn start_local_socket_server() -> bool { } } - let snapshot_enabled = SNAPSHOT_BROADCAST_ENABLED.load(Ordering::Relaxed); if clients.is_empty() { - if !last_snapshot_json.is_empty() { - last_snapshot_json.clear(); + if !last_topology_json.is_empty() { + last_topology_json.clear(); + last_topology_json.shrink_to_fit(); } delivered_tun_requests.clear(); + shrink_hash_set_if_sparse(&mut delivered_tun_requests); last_tun_route_signatures.clear(); + shrink_hash_map_if_sparse(&mut last_tun_route_signatures); tun_event_receivers.clear(); + shrink_hash_map_if_sparse(&mut tun_event_receivers); + clients.shrink_to_fit(); last_event_receiver_sync_at = None; + last_traffic_stats_at = None; + last_instance_poll_at = None; tun_bootstrap_done = false; thread::sleep(SOCKET_TICK_INTERVAL); continue; @@ -181,115 +406,164 @@ pub fn start_local_socket_server() -> bool { sync_tun_event_receivers(&mut tun_event_receivers); last_event_receiver_sync_at = Some(now); } - if drain_tun_refresh_events(&mut tun_event_receivers) { + let drained_events = drain_kernel_events(&mut tun_event_receivers); + let tun_refresh = drained_events.tun_refresh; + let topology_lost = drained_events.topology_lost; + let peer_events = drained_events.peer_events; + if topology_lost { + full_topology_dirty = true; + } + if tun_refresh { tun_bootstrap_done = false; tun_fast_until = now + TUN_FAST_CHECK_WINDOW; } - let should_collect_snapshot = snapshot_enabled - || accepted_client - || (!tun_bootstrap_done && now < tun_fast_until); - if !should_collect_snapshot { - if !last_snapshot_json.is_empty() { - last_snapshot_json.clear(); + if !peer_events.is_empty() { + broadcast_runtime_peer_events(&mut clients, peer_events); + } + let should_collect_traffic_stats = last_traffic_stats_at + .map(|last| now.duration_since(last) >= TRAFFIC_STATS_INTERVAL) + .unwrap_or(true); + if should_collect_traffic_stats { + last_traffic_stats_at = Some(now); + match serde_json::to_string(&collect_traffic_stats()) { + Ok(json) => { + let _ = broadcast_local_socket_json_payload_message( + &mut clients, + "traffic_stats", + &json, + ); + } + Err(err) => { + ohrs_log_error!("[Rust] serialize traffic stats failed: {}", err); + } } + } + let should_poll_instance = last_instance_poll_at + .map(|last| now.duration_since(last) >= INSTANCE_POLL_INTERVAL) + .unwrap_or(true); + let should_collect_topology = accepted_client + || full_topology_dirty + || tun_refresh + || should_poll_instance + || (!tun_bootstrap_done && now < tun_fast_until); + if !should_collect_topology { thread::sleep(SOCKET_TICK_INTERVAL); continue; } - let snapshot = get_runtime_snapshot_inner(); - if snapshot_enabled { - let snapshot_json = match serde_json::to_string(&snapshot) { - Ok(json) => json, - Err(err) => { - ohrs_log_error!("[Rust] serialize runtime snapshot failed: {}", err); - thread::sleep(SOCKET_TICK_INTERVAL); - continue; + let snapshot = collect_runtime_state_inner(); + last_instance_poll_at = Some(now); + match serde_json::to_string(&snapshot) { + Ok(json) => { + if accepted_client || full_topology_dirty || json != last_topology_json { + let _ = broadcast_local_socket_json_payload_message( + &mut clients, + "runtime_topology", + &json, + ); + last_topology_json = json; } - }; - - if accepted_client || snapshot_json != last_snapshot_json { - let _ = broadcast_local_socket_json_payload_message( - &mut clients, - "runtime_snapshot", - &snapshot_json, - ); - last_snapshot_json = snapshot_json; } - } else if !last_snapshot_json.is_empty() { - last_snapshot_json.clear(); + Err(err) => { + ohrs_log_error!("[Rust] serialize runtime topology failed: {}", err); + } } + let active_tun_candidate_ids = tun_candidate_ids(&snapshot); + delivered_tun_requests + .retain(|instance_id| active_tun_candidate_ids.contains(instance_id)); + last_tun_route_signatures + .retain(|instance_id, _| active_tun_candidate_ids.contains(instance_id)); + shrink_hash_set_if_sparse(&mut delivered_tun_requests); + shrink_hash_map_if_sparse(&mut last_tun_route_signatures); + let has_undelivered_tun_candidate = active_tun_candidate_ids + .iter() + .any(|instance_id| !delivered_tun_requests.contains(instance_id)); + let should_evaluate_tun = + tun_refresh || !tun_bootstrap_done || has_undelivered_tun_candidate; let mut saw_running_instance = false; let mut saw_tun_candidate = false; - for instance in snapshot.instances.iter() { - if instance.running { - saw_running_instance = true; - } - if instance.running && instance.tun_required { - saw_tun_candidate = true; - let virtual_ipv4 = instance - .my_node_info - .as_ref() - .and_then(|info| info.virtual_ipv4.clone()); - let virtual_ipv4_cidr = instance - .my_node_info - .as_ref() - .and_then(|info| info.virtual_ipv4_cidr.clone()); - if clients.is_empty() { - continue; + if should_evaluate_tun { + for instance in snapshot.instances.iter() { + if instance.running { + saw_running_instance = true; } - if virtual_ipv4.is_none() || virtual_ipv4_cidr.is_none() { - continue; - } - let aggregated_routes = aggregate_tun_routes(instance); - let route_signature = serde_json::to_string(&( - &virtual_ipv4, - &virtual_ipv4_cidr, - &aggregated_routes, - instance.magic_dns_enabled, - instance.need_exit_node, - )) - .unwrap_or_else(|_| "[]".to_string()); - let should_send = accepted_client - || !delivered_tun_requests.contains(&instance.instance_id) - || last_tun_route_signatures - .get(&instance.instance_id) - .map(|value| value != &route_signature) - .unwrap_or(true); - if !should_send { - continue; - } - let payload = TunRequestPayload { - config_id: instance.config_id.clone(), - instance_id: instance.instance_id.clone(), - display_name: instance.display_name.clone(), - virtual_ipv4, - virtual_ipv4_cidr, - aggregated_routes, - magic_dns_enabled: instance.magic_dns_enabled, - need_exit_node: instance.need_exit_node, - }; - let payload_json = match serde_json::to_string(&payload) { - Ok(json) => json, - Err(err) => { - ohrs_log_error!("[Rust] serialize tun request failed: {}", err); + if instance.running && instance.tun_required { + saw_tun_candidate = true; + let virtual_ipv4 = instance + .my_node_info + .as_ref() + .and_then(|info| info.virtual_ipv4.clone()); + let virtual_ipv4_cidr = instance + .my_node_info + .as_ref() + .and_then(|info| info.virtual_ipv4_cidr.clone()); + if clients.is_empty() { continue; } - }; - if broadcast_local_socket_message(&mut clients, "tun_request", &payload_json) { - delivered_tun_requests.insert(instance.instance_id.clone()); - last_tun_route_signatures - .insert(instance.instance_id.clone(), route_signature); + if virtual_ipv4.is_none() || virtual_ipv4_cidr.is_none() { + continue; + } + let aggregated_routes = aggregate_tun_routes(instance); + let route_signature = serde_json::to_string(&( + &virtual_ipv4, + &virtual_ipv4_cidr, + &aggregated_routes, + instance.magic_dns_enabled, + instance.need_exit_node, + )) + .unwrap_or_else(|_| "[]".to_string()); + let allow_route_signature_refresh = tun_refresh || !tun_bootstrap_done; + let should_send = !delivered_tun_requests.contains(&instance.instance_id) + || (allow_route_signature_refresh + && last_tun_route_signatures + .get(&instance.instance_id) + .map(|value| value != &route_signature) + .unwrap_or(true)); + if !should_send { + continue; + } + let payload = TunRequestPayload { + config_id: instance.config_id.clone(), + instance_id: instance.instance_id.clone(), + display_name: instance.display_name.clone(), + virtual_ipv4, + virtual_ipv4_cidr, + aggregated_routes, + magic_dns_enabled: instance.magic_dns_enabled, + need_exit_node: instance.need_exit_node, + }; + let payload_json = match serde_json::to_string(&payload) { + Ok(json) => json, + Err(err) => { + ohrs_log_error!("[Rust] serialize tun request failed: {}", err); + continue; + } + }; + if broadcast_local_socket_message( + &mut clients, + "tun_request", + &payload_json, + ) { + delivered_tun_requests.insert(instance.instance_id.clone()); + last_tun_route_signatures + .insert(instance.instance_id.clone(), route_signature); + } + } + } + } else { + for instance in snapshot.instances.iter() { + if instance.running { + saw_running_instance = true; + } + if instance.running && instance.tun_required { + saw_tun_candidate = true; } - } else { - delivered_tun_requests.remove(&instance.instance_id); - last_tun_route_signatures.remove(&instance.instance_id); } } - if !snapshot_enabled - && (!delivered_tun_requests.is_empty() - || (saw_running_instance && !saw_tun_candidate) - || now >= tun_fast_until) + if !delivered_tun_requests.is_empty() + || (saw_running_instance && !saw_tun_candidate) + || now >= tun_fast_until { tun_bootstrap_done = true; } diff --git a/easytier-contrib/easytier-ohrs/src/lib.rs b/easytier-contrib/easytier-ohrs/src/lib.rs index ad222865..3fa5f7fd 100644 --- a/easytier-contrib/easytier-ohrs/src/lib.rs +++ b/easytier-contrib/easytier-ohrs/src/lib.rs @@ -63,7 +63,7 @@ use easytier::proto::api::manage::NetworkConfig; use easytier::proto::api::manage::NetworkingMethod; use easytier::web_client::{WebClient, WebClientHooks, run_web_client}; use kernel_bridge::{ - set_snapshot_broadcast_enabled, start_local_socket_server as start_local_socket_server_inner, + start_local_socket_server as start_local_socket_server_inner, stop_local_socket_server as stop_local_socket_server_inner, }; use napi_derive_ohos::napi; @@ -517,18 +517,8 @@ mod tests { } } -#[napi] -pub fn get_runtime_snapshot() -> RuntimeAggregateState { - exports::runtime_api::get_runtime_snapshot() -} - -#[napi] -pub fn set_kernel_snapshot_enabled(enabled: bool) { - set_snapshot_broadcast_enabled(enabled); -} - -pub(crate) fn get_runtime_snapshot_inner() -> RuntimeAggregateState { - exports::runtime_api::get_runtime_snapshot_inner() +pub(crate) fn collect_runtime_state_inner() -> RuntimeAggregateState { + exports::runtime_api::collect_runtime_state() } #[napi] diff --git a/easytier-contrib/easytier-ohrs/src/runtime/state/runtime_state.rs b/easytier-contrib/easytier-ohrs/src/runtime/state/runtime_state.rs index 948b7b59..baff8f1b 100644 --- a/easytier-contrib/easytier-ohrs/src/runtime/state/runtime_state.rs +++ b/easytier-contrib/easytier-ohrs/src/runtime/state/runtime_state.rs @@ -324,7 +324,7 @@ fn route_to_view(route: api::instance::Route) -> RouteView { } } -fn peer_conn_to_view(conn: api::instance::PeerConnInfo) -> PeerConnInfo { +pub(crate) fn peer_conn_to_view(conn: api::instance::PeerConnInfo) -> PeerConnInfo { let stats = conn.stats.map(|stats| PeerConnStats { rx_bytes: stats.rx_bytes as i64, tx_bytes: stats.tx_bytes as i64,