Skip to content
Merged
Show file tree
Hide file tree
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
2 changes: 1 addition & 1 deletion crates/netscli-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ pub use dns::{resolve_a, resolve_aaaa};
pub use inspect::{InspectEngine, InspectResult};
#[cfg(feature = "mdns")]
pub use mdns::{MdnsEngine, MdnsService, COMMON_SERVICE_TYPES};
pub use ops::{resolve_host_ip, Ops, OpsConfig, PingSummary};
pub use ops::{resolve_host_ip, Ops, OpsConfig, PingSummary, MAX_CONCURRENCY};
pub use oui::lookup_vendor;
pub use pcap::{
PcapCancelToken, PcapConfig, PcapEngine, PcapPacketSummary, PcapParseResult, PcapResult,
Expand Down
2 changes: 1 addition & 1 deletion crates/netscli-core/src/ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,5 +5,5 @@ mod pcap;
mod scan;
mod validation;

pub use config::{Ops, OpsConfig};
pub use config::{Ops, OpsConfig, MAX_CONCURRENCY};
pub use host::{resolve_host_ip, resolve_host_ip_with_timeout, PingSummary};
14 changes: 10 additions & 4 deletions crates/netscli-core/src/ops/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,14 +27,20 @@ pub struct Ops {
pub(super) cfg: OpsConfig,
}

/// Upper bound on in-flight probes.
///
/// Exported so callers clamp to the same value rather than guessing: the MCP
/// server used to clamp to 4096 here, which meant 1025..=4096 was accepted at
/// the API boundary and then silently reduced by `Ops::new` (C-10).
pub const MAX_CONCURRENCY: usize = 1024;

impl Ops {
pub fn new(cfg: OpsConfig) -> Self {
let mut cfg = cfg;
// Clamp [1, 1024]. Lower bound prevents semaphore deadlock; upper
// bound matches the MCP server's existing per-call clamp and stops
// users from triggering kernel ephemeral-port exhaustion on aggressive
// Lower bound prevents semaphore deadlock; the upper bound stops users
// from triggering kernel ephemeral-port exhaustion on aggressive
// --concurrency values.
cfg.concurrency = cfg.concurrency.clamp(1, 1024);
cfg.concurrency = cfg.concurrency.clamp(1, MAX_CONCURRENCY);
cfg.scan_timeout_ms = cfg.scan_timeout_ms.max(1);
cfg.ping_timeout_ms = cfg.ping_timeout_ms.max(1);
cfg.dns_timeout_ms = cfg.dns_timeout_ms.max(1);
Expand Down
76 changes: 23 additions & 53 deletions crates/netscli-core/src/stats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,14 @@ impl NetworkMonitor {

let (rx, tx, _) = sum_traffic(&networks, None);
let now = Instant::now();
let inactive = now - Self::ACTIVE_HOLD - Self::ACTIVE_HOLD;
// `Instant - Duration` panics on underflow, and within ~360ms of the
// process's monotonic-clock origin `now` is smaller than the two
// holds we subtract (C-01). `checked_sub` degrades to "as early as
// representable", which reads as inactive exactly as intended.
let inactive = now
.checked_sub(Self::ACTIVE_HOLD)
.and_then(|t| t.checked_sub(Self::ACTIVE_HOLD))
.unwrap_or(now);

Self {
state: Mutex::new(MonitorState {
Expand All @@ -73,7 +80,14 @@ impl NetworkMonitor {
let (rx, tx, _) = sum_traffic(&s.networks, s.selected_interface.as_deref());

let now = Instant::now();
let inactive = now - Self::ACTIVE_HOLD - Self::ACTIVE_HOLD;
// `Instant - Duration` panics on underflow, and within ~360ms of the
// process's monotonic-clock origin `now` is smaller than the two
// holds we subtract (C-01). `checked_sub` degrades to "as early as
// representable", which reads as inactive exactly as intended.
let inactive = now
.checked_sub(Self::ACTIVE_HOLD)
.and_then(|t| t.checked_sub(Self::ACTIVE_HOLD))
.unwrap_or(now);

s.last_update = now;
s.last_stats = (rx, tx);
Expand Down Expand Up @@ -190,13 +204,15 @@ fn sum_traffic(networks: &Networks, selected: Option<&str>) -> (u64, u64, bool)
return (0, 0, false);
}

let mut rx = 0;
let mut tx = 0;
let mut rx: u64 = 0;
let mut tx: u64 = 0;
let mut any = false;
for (_, network) in networks {
any = true;
rx += network.total_received();
tx += network.total_transmitted();
// Saturating, not `+=`: these are per-interface byte counters summed
// across every adapter, and a debug build panics on overflow (C-02).
rx = rx.saturating_add(network.total_received());
tx = tx.saturating_add(network.total_transmitted());
}
(rx, tx, any)
}
Expand All @@ -216,50 +232,4 @@ fn bytes_to_mbps(bytes: u64, elapsed: Duration) -> f64 {
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn bytes_to_mbps_typical() {
// 1_000_000 bytes in 1 second = 8 Mbps
assert!((bytes_to_mbps(1_000_000, Duration::from_secs(1)) - 8.0).abs() < 1e-6);
}

#[test]
fn bytes_to_mbps_zero_elapsed_safe() {
assert_eq!(bytes_to_mbps(1_000_000, Duration::from_secs(0)), 0.0);
}

#[test]
fn bytes_to_mbps_reports_true_value_for_gbe() {
// 1.25 GB in 1s = 10 Gbps = 10_000 Mbps. Previously this was
// silently clamped to 999.99 which hid gigabit+ link activity.
let actual = bytes_to_mbps(1_250_000_000, Duration::from_secs(1));
assert!(
(actual - 10_000.0).abs() < 1.0,
"expected ~10_000 Mbps, got {actual}"
);
}

#[test]
fn new_monitor_reports_available_on_refresh() {
// Smoke test: constructing and calling get_stats should never panic
// and should report `available=true` if any network interface exists
// on the test host (all CI runners do).
let monitor = NetworkMonitor::new();
let stats = monitor.get_stats();
// Not asserting `available` since some sandboxed CI may lack interfaces;
// just ensuring the call path is sound.
assert!(stats.upload_mbps >= 0.0 && stats.download_mbps >= 0.0);
}

#[test]
fn selected_interface_round_trips() {
let monitor = NetworkMonitor::new();
assert_eq!(monitor.selected_interface(), None);
monitor.set_interface(Some("lo0".to_string()));
assert_eq!(monitor.selected_interface(), Some("lo0".to_string()));
monitor.set_interface(None);
assert_eq!(monitor.selected_interface(), None);
}
}
mod tests;
78 changes: 78 additions & 0 deletions crates/netscli-core/src/stats/tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
//! Tests for the network traffic monitor.
//!
//! Split from `stats.rs` to keep it under the maintainability cap.

use super::*;

// C-01: `now - ACTIVE_HOLD - ACTIVE_HOLD` panics on underflow, which is
// reachable within ~360ms of the monotonic clock's origin. There is no
// portable way to construct a near-zero `Instant`, so this pins the
// saturating arithmetic itself rather than the call site.
#[test]
fn instant_underflow_saturates_instead_of_panicking() {
let origin = Instant::now();
let hold = NetworkMonitor::ACTIVE_HOLD;

// Mirrors the expression in `new()` and `refresh()`.
let inactive = origin
.checked_sub(hold)
.and_then(|t| t.checked_sub(hold))
.unwrap_or(origin);

// Either it stepped back by two holds, or it clamped at the origin —
// never a panic, and never a time in the future.
assert!(inactive <= origin);
}

// C-02: these are per-interface byte counters summed across every
// adapter, and `+=` panics on overflow in a debug build.
#[test]
fn traffic_accumulation_saturates_on_overflow() {
let mut rx: u64 = u64::MAX - 1;
rx = rx.saturating_add(1_000_000);
assert_eq!(rx, u64::MAX);
}

#[test]
fn bytes_to_mbps_typical() {
// 1_000_000 bytes in 1 second = 8 Mbps
assert!((bytes_to_mbps(1_000_000, Duration::from_secs(1)) - 8.0).abs() < 1e-6);
}

#[test]
fn bytes_to_mbps_zero_elapsed_safe() {
assert_eq!(bytes_to_mbps(1_000_000, Duration::from_secs(0)), 0.0);
}

#[test]
fn bytes_to_mbps_reports_true_value_for_gbe() {
// 1.25 GB in 1s = 10 Gbps = 10_000 Mbps. Previously this was
// silently clamped to 999.99 which hid gigabit+ link activity.
let actual = bytes_to_mbps(1_250_000_000, Duration::from_secs(1));
assert!(
(actual - 10_000.0).abs() < 1.0,
"expected ~10_000 Mbps, got {actual}"
);
}

#[test]
fn new_monitor_reports_available_on_refresh() {
// Smoke test: constructing and calling get_stats should never panic
// and should report `available=true` if any network interface exists
// on the test host (all CI runners do).
let monitor = NetworkMonitor::new();
let stats = monitor.get_stats();
// Not asserting `available` since some sandboxed CI may lack interfaces;
// just ensuring the call path is sound.
assert!(stats.upload_mbps >= 0.0 && stats.download_mbps >= 0.0);
}

#[test]
fn selected_interface_round_trips() {
let monitor = NetworkMonitor::new();
assert_eq!(monitor.selected_interface(), None);
monitor.set_interface(Some("lo0".to_string()));
assert_eq!(monitor.selected_interface(), Some("lo0".to_string()));
monitor.set_interface(None);
assert_eq!(monitor.selected_interface(), None);
}
27 changes: 26 additions & 1 deletion crates/netscli-core/tests/config_safety.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,32 @@ fn ops_new_clamps_concurrency_upper_bound() {
};

let ops = Ops::new(cfg);
assert_eq!(ops.config().concurrency, 1024);
assert_eq!(ops.config().concurrency, netscli_core::MAX_CONCURRENCY);
}

// C-10: the MCP server clamped to 4096 while this clamps to 1024, so any
// value in 1025..=4096 was accepted at the API boundary and then silently
// reduced -- and a comment claimed the two bounds matched. Both now derive
// from this constant, so a change to one cannot drift from the other.
#[test]
fn max_concurrency_is_the_value_ops_actually_enforces() {
let cfg = OpsConfig {
concurrency: netscli_core::MAX_CONCURRENCY + 1,
..Default::default()
};
let ops = Ops::new(cfg);
assert_eq!(ops.config().concurrency, netscli_core::MAX_CONCURRENCY);

// And a value just inside the bound survives untouched, which is what
// makes the constant meaningful rather than merely an upper limit.
let cfg = OpsConfig {
concurrency: netscli_core::MAX_CONCURRENCY,
..Default::default()
};
assert_eq!(
Ops::new(cfg).config().concurrency,
netscli_core::MAX_CONCURRENCY
);
}

#[tokio::test]
Expand Down
8 changes: 7 additions & 1 deletion crates/netscli-mcp/src/server/schemas.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,15 @@ use super::errors::RpcError;

const MAX_SUBNET_ADDRESSES: u64 = 1 << 16; // /16

/// Clamp to the same ceiling `Ops` enforces.
///
/// This used to clamp to 4096 while `Ops::new` re-clamped to 1024, so any
/// value in 1025..=4096 was accepted here and then silently reduced, and a
/// comment in `ops/config.rs` claimed the two bounds matched (C-10).
/// Deferring to the core constant makes that true by construction.
pub(super) fn clamp_concurrency(max_concurrent: Option<usize>, default: usize) -> usize {
let c = max_concurrent.unwrap_or(default);
c.clamp(1, 4096)
c.clamp(1, netscli_core::MAX_CONCURRENCY)
}

pub(super) fn clamp_timeout_ms(timeout_ms: Option<u64>, default: u64) -> u64 {
Expand Down
58 changes: 53 additions & 5 deletions scripts/generate-oui.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,14 +33,18 @@ fn canon_prefix(hex: &str) -> Option<String> {
Some(format!("{}:{}:{}", &p[0..2], &p[2..4], &p[4..6]))
}

// `error_for_status` matters more than it looks here (A-09): IEEE rate-limits
// unknown user agents, and without this a 403 body was handed straight to the
// CSV parser, which found no header row and returned an empty map -- silently
// producing a gutted dataset rather than failing.
async fn fetch_url(client: &Client, url: &str) -> Result<String, Box<dyn std::error::Error>> {
let response = client.get(url).send().await?;
let response = client.get(url).send().await?.error_for_status()?;
let text = response.text().await?;
Ok(text)
}

async fn fetch_gz(client: &Client, url: &str) -> Result<Vec<u8>, Box<dyn std::error::Error>> {
let response = client.get(url).send().await?;
let response = client.get(url).send().await?.error_for_status()?;
let bytes = response.bytes().await?;
Ok(bytes.to_vec())
}
Expand All @@ -63,7 +67,12 @@ fn parse_oui_csv(csv: &str) -> HashMap<String, String> {
.iter()
.position(|c| c.to_lowercase().contains("organization"));

// A missing header row used to yield an empty map indistinguishable from
// "this file legitimately had no rows" (A-09). Warn loudly; the
// minimum-entry floor in main() is what actually stops the bad write.
if p_idx.is_none() || o_idx.is_none() {
eprintln!(" WARNING: no assignment/organization header found");
eprintln!(" The response was probably an error page, not CSV.");
return map;
}

Expand Down Expand Up @@ -132,8 +141,13 @@ fn parse_wireshark_manuf(
.replace(':', "")
.to_uppercase();

if hex.len() >= 6 {
let canon = format!("{}:{}:{}", &hex[0..2], &hex[2..4], &hex[4..6]);
// Use the shared canonicaliser rather than slicing `hex` directly
// (C-36). Unlike the CSV path, `hex` here has only had colons removed
// — it is not hex-filtered — so `&hex[0..2]` panicked on any entry
// whose first token began with a multi-byte character. `canon_prefix`
// filters to ASCII hex digits first, which also drops the malformed
// rows that would otherwise have produced junk vendor keys.
if let Some(canon) = canon_prefix(&hex) {
map.entry(canon).or_insert_with(|| vendor);
}
}
Expand All @@ -143,7 +157,12 @@ fn parse_wireshark_manuf(

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let client = Client::builder().timeout(FETCH_TIMEOUT).build()?;
// IEEE rate-limits unknown user agents, which is what made the silent
// 403-parsed-as-CSV failure reachable in the first place.
let client = Client::builder()
.timeout(FETCH_TIMEOUT)
.user_agent(concat!("netscli-oui-generator/", env!("CARGO_PKG_VERSION")))
.build()?;
println!("Fetching OUI data from IEEE...");
let mut merged: HashMap<String, String> = HashMap::new();

Expand Down Expand Up @@ -185,6 +204,35 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
.join("netscli-core")
.join("data")
.join("oui.min.json.gz");
// Refuse to overwrite a good dataset with a substantially smaller one
// (A-09). The output used to be written unconditionally, so a throttled
// or reshaped upstream produced a near-empty vendor database that looked
// like a successful run and would ship in the next release.
//
// 90% rather than "any shrinkage": entries do legitimately disappear when
// registrations lapse, and a hard equality check would fail every run.
if let Ok(existing) = std::fs::read(&out_path) {
let mut decoder = flate2::read::GzDecoder::new(&existing[..]);
let mut previous = String::new();
if std::io::Read::read_to_string(&mut decoder, &mut previous).is_ok() {
if let Ok(Value::Object(old)) = serde_json::from_str::<Value>(&previous) {
let floor = old.len() * 9 / 10;
if json_map.len() < floor {
return Err(format!(
"refusing to write: {} prefixes is below 90% of the existing {} \
(floor {}). The upstream fetch was probably throttled or \
reshaped. Re-run, or delete {} to override.",
json_map.len(),
old.len(),
floor,
out_path.display()
)
.into());
}
}
}
}

std::fs::create_dir_all(out_path.parent().unwrap_or(repo_root))?;
let file = File::create(&out_path)?;
let mut encoder = GzEncoder::new(file, Compression::default());
Expand Down
Loading