From ec50b7faa5530b87cc6011fd03f5ad17c727c940 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Wed, 14 Jan 2026 19:04:50 -0600 Subject: [PATCH] Fix semaphore error handling and parallelize interface polling MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. Fixed semaphore acquire to handle errors gracefully instead of panicking when permit acquisition fails during concurrent polling 2. Parallelized interface counter polling using tokio::join! to fetch all 6 SNMP counters (in/out octets/errors/discards) concurrently instead of sequentially Performance improvement: Reduces per-interface polling latency from ~30ms (6 × 5ms timeout) to ~5ms (1 parallel batch). --- src/poller/executor.rs | 115 +++++++++++++++++++--------------------- src/poller/scheduler.rs | 8 ++- 2 files changed, 61 insertions(+), 62 deletions(-) diff --git a/src/poller/executor.rs b/src/poller/executor.rs index a77b381..6d8f910 100644 --- a/src/poller/executor.rs +++ b/src/poller/executor.rs @@ -147,72 +147,65 @@ impl Executor { let if_in_discards_oid = format!("1.3.6.1.2.1.2.2.1.13.{}", interface.if_index); let if_out_discards_oid = format!("1.3.6.1.2.1.2.2.1.19.{}", interface.if_index); - // Poll each counter - let in_octets = self - .get_counter( + // Poll all counters in parallel for this interface + let ( + in_octets_result, + out_octets_result, + in_errors_result, + out_errors_result, + in_discards_result, + out_discards_result, + ) = tokio::join!( + self.get_counter( &equipment.ip_address, &equipment.snmp.community, &equipment.snmp.version, equipment.snmp.port, - &if_in_octets_oid, - ) - .await - .unwrap_or(0); + &if_in_octets_oid + ), + self.get_counter( + &equipment.ip_address, + &equipment.snmp.community, + &equipment.snmp.version, + equipment.snmp.port, + &if_out_octets_oid + ), + self.get_counter( + &equipment.ip_address, + &equipment.snmp.community, + &equipment.snmp.version, + equipment.snmp.port, + &if_in_errors_oid + ), + self.get_counter( + &equipment.ip_address, + &equipment.snmp.community, + &equipment.snmp.version, + equipment.snmp.port, + &if_out_errors_oid + ), + self.get_counter( + &equipment.ip_address, + &equipment.snmp.community, + &equipment.snmp.version, + equipment.snmp.port, + &if_in_discards_oid + ), + self.get_counter( + &equipment.ip_address, + &equipment.snmp.community, + &equipment.snmp.version, + equipment.snmp.port, + &if_out_discards_oid + ), + ); - let out_octets = self - .get_counter( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &if_out_octets_oid, - ) - .await - .unwrap_or(0); - - let in_errors = self - .get_counter( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &if_in_errors_oid, - ) - .await - .unwrap_or(0); - - let out_errors = self - .get_counter( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &if_out_errors_oid, - ) - .await - .unwrap_or(0); - - let in_discards = self - .get_counter( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &if_in_discards_oid, - ) - .await - .unwrap_or(0); - - let out_discards = self - .get_counter( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &if_out_discards_oid, - ) - .await - .unwrap_or(0); + let in_octets = in_octets_result.unwrap_or(0); + let out_octets = out_octets_result.unwrap_or(0); + let in_errors = in_errors_result.unwrap_or(0); + let out_errors = out_errors_result.unwrap_or(0); + let in_discards = in_discards_result.unwrap_or(0); + let out_discards = out_discards_result.unwrap_or(0); let stat = InterfaceStat { interface_id: interface.id.clone(), diff --git a/src/poller/scheduler.rs b/src/poller/scheduler.rs index 2112002..518e43b 100644 --- a/src/poller/scheduler.rs +++ b/src/poller/scheduler.rs @@ -275,7 +275,13 @@ impl Scheduler { let task = tokio::spawn(async move { // Acquire permit before polling (limits concurrency) - let _permit = permit.acquire().await.unwrap(); + let _permit = match permit.acquire().await { + Ok(p) => p, + Err(e) => { + error!("Failed to acquire polling permit: {}", e); + return; + } + }; info!("Polling equipment: {}", equipment.name);