Fix semaphore error handling and parallelize interface polling
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).
This commit is contained in:
parent
6462bdbf81
commit
ec50b7faa5
2 changed files with 61 additions and 62 deletions
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue