From f5b7e589f2ed275930c1650f0e13412dea332372 Mon Sep 17 00:00:00 2001 From: mayor Date: Sat, 7 Feb 2026 09:17:17 -0600 Subject: [PATCH] fix: implement agent IP detection and remove dead poller code - Implement get_local_ip() using UDP socket trick to resolve local interface address without sending traffic, so heartbeat reports actual agent IP instead of "unknown" - Remove dead poller/scheduler/executor and api_client modules that referenced non-existent modules and were never compiled --- src/api_client.rs | 189 --------------------------- src/poller/executor.rs | 275 ---------------------------------------- src/poller/mod.rs | 5 - src/poller/scheduler.rs | 248 ------------------------------------ src/websocket_client.rs | 17 ++- 5 files changed, 11 insertions(+), 723 deletions(-) delete mode 100644 src/api_client.rs delete mode 100644 src/poller/executor.rs delete mode 100644 src/poller/mod.rs delete mode 100644 src/poller/scheduler.rs diff --git a/src/api_client.rs b/src/api_client.rs deleted file mode 100644 index b419285..0000000 --- a/src/api_client.rs +++ /dev/null @@ -1,189 +0,0 @@ -use crate::config::{AgentConfig, HeartbeatMetadata}; -use crate::metrics::Metric; -use crate::proto::agent; -use prost::Message; -use std::time::Duration; -use thiserror::Error; -use tracing::debug; - -#[derive(Debug, Error)] -pub enum ApiError { - #[error("HTTP request failed: {0}")] - RequestFailed(String), - - #[error("HTTP status error: {0}")] - StatusError(u16), - - #[error("JSON parsing error: {0}")] - JsonError(#[from] std::io::Error), - - #[error("Task join error: {0}")] - JoinError(String), -} - -pub type Result = std::result::Result; - -/// API client for communicating with the Towerops server -#[derive(Clone)] -pub struct ApiClient { - base_url: String, - token: String, -} - -impl ApiClient { - /// Create a new API client - pub fn new(base_url: String, token: String) -> Result { - Ok(Self { base_url, token }) - } - - /// Fetch configuration from the API - pub async fn fetch_config(&self) -> Result { - let url = format!("{}/api/v1/agent/config", self.base_url); - let token = self.token.clone(); - - let config = tokio::task::spawn_blocking(move || { - let response = ureq::get(&url) - .set("Authorization", &format!("Bearer {}", token)) - .timeout(Duration::from_secs(30)) - .call() - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - - let status = response.status(); - debug!("fetch_config response: status={}", status); - - if status != 200 { - return Err(ApiError::StatusError(status)); - } - - let body = response - .into_string() - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - debug!("fetch_config response body: {}", body); - - let config: AgentConfig = serde_json::from_str(&body) - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - }) - .await - .map_err(|e| ApiError::JoinError(e.to_string()))??; - - Ok(config) - } - - /// Submit metrics to the API using Protocol Buffers - pub async fn submit_metrics(&self, metrics: Vec) -> Result<()> { - if metrics.is_empty() { - return Ok(()); - } - - let url = format!("{}/api/v1/agent/metrics", self.base_url); - let token = self.token.clone(); - - tokio::task::spawn_blocking(move || { - // Convert metrics to protobuf - let proto_metrics: Vec = metrics - .into_iter() - .map(|m| convert_metric_to_proto(&m)) - .collect(); - - let batch = agent::MetricBatch { - metrics: proto_metrics, - }; - - // Encode to bytes - let mut buf = Vec::new(); - batch - .encode(&mut buf) - .map_err(|e| ApiError::RequestFailed(format!("Protobuf encoding error: {}", e)))?; - - // Send as protobuf - let response = ureq::post(&url) - .set("Authorization", &format!("Bearer {}", token)) - .set("Content-Type", "application/x-protobuf") - .timeout(Duration::from_secs(30)) - .send_bytes(&buf) - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - - let status = response.status(); - debug!("submit_metrics response: status={}", status); - - if status != 200 { - return Err(ApiError::StatusError(status)); - } - - let body = response - .into_string() - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - debug!("submit_metrics response body: {}", body); - - Ok(()) - }) - .await - .map_err(|e| ApiError::JoinError(e.to_string()))??; - - Ok(()) - } - - /// Send heartbeat to the API - pub async fn heartbeat(&self, metadata: HeartbeatMetadata) -> Result<()> { - let url = format!("{}/api/v1/agent/heartbeat", self.base_url); - let token = self.token.clone(); - - tokio::task::spawn_blocking(move || { - let response = ureq::post(&url) - .set("Authorization", &format!("Bearer {}", token)) - .timeout(Duration::from_secs(30)) - .send_json(&metadata) - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - - let status = response.status(); - debug!("heartbeat response: status={}", status); - - if status != 200 { - return Err(ApiError::StatusError(status)); - } - - let body = response - .into_string() - .map_err(|e| ApiError::RequestFailed(e.to_string()))?; - debug!("heartbeat response body: {}", body); - - Ok(()) - }) - .await - .map_err(|e| ApiError::JoinError(e.to_string()))??; - - Ok(()) - } -} - -/// Convert internal Metric type to protobuf Metric -fn convert_metric_to_proto(metric: &Metric) -> agent::Metric { - use crate::metrics::Metric as M; - - match metric { - M::SensorReading(sr) => agent::Metric { - metric_type: Some(agent::metric::MetricType::SensorReading( - agent::SensorReading { - sensor_id: sr.sensor_id.clone(), - value: sr.value, - status: sr.status.clone(), - timestamp: sr.timestamp.to_unix_timestamp(), - }, - )), - }, - M::InterfaceStat(is) => agent::Metric { - metric_type: Some(agent::metric::MetricType::InterfaceStat( - agent::InterfaceStat { - interface_id: is.interface_id.clone(), - if_in_octets: is.if_in_octets, - if_out_octets: is.if_out_octets, - if_in_errors: is.if_in_errors, - if_out_errors: is.if_out_errors, - if_in_discards: is.if_in_discards, - if_out_discards: is.if_out_discards, - timestamp: is.timestamp.to_unix_timestamp(), - }, - )), - }, - } -} diff --git a/src/poller/executor.rs b/src/poller/executor.rs deleted file mode 100644 index 5aa2162..0000000 --- a/src/poller/executor.rs +++ /dev/null @@ -1,275 +0,0 @@ -use crate::buffer::Storage; -use crate::buffer::StorageError; -use thiserror::Error; - -#[derive(Debug, Error)] -pub enum ExecutorError { - #[error("Storage error: {0}")] - Storage(#[from] StorageError), - #[error("SNMP error: {0}")] - Snmp(#[from] crate::snmp::SnmpError), - #[error("Conversion error: {0}")] - Conversion(String), -} - -pub type Result = std::result::Result; - -use crate::config::EquipmentConfig; -use crate::metrics::{InterfaceStat, Metric, SensorReading}; -use crate::snmp::SnmpClient; - -use crate::metrics::Timestamp; -use log::{error, info, warn}; - -/// Redact SNMP community string for logging, showing first 2 chars only -fn redact_community(community: &str) -> String { - if community.is_empty() { - return "**".to_string(); - } - let visible = std::cmp::min(community.len(), 2); - format!("{}**", &community[..visible]) -} - -/// Executor handles polling individual pieces of equipment -#[derive(Clone)] -pub struct Executor { - snmp_client: SnmpClient, - storage: Storage, -} - -impl Executor { - pub fn new(snmp_client: SnmpClient, storage: Storage) -> Self { - Self { - snmp_client, - storage, - } - } - - /// Poll sensors for a piece of equipment - pub async fn poll_sensors(&self, equipment: &EquipmentConfig) -> Result<()> { - if !equipment.snmp.enabled || equipment.sensors.is_empty() { - return Ok(()); - } - - info!( - "Polling {} sensors for equipment: {} (community: {})", - equipment.sensors.len(), - equipment.name, - redact_community(&equipment.snmp.community) - ); - - for sensor in &equipment.sensors { - match self - .snmp_client - .get( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &sensor.oid, - ) - .await - { - Ok(value) => { - let raw_value = match value.as_f64() { - Some(v) => v, - None => { - warn!( - "Could not convert sensor value to number for sensor {}", - sensor.id - ); - continue; - } - }; - - // Apply divisor if specified - let final_value = if let Some(divisor) = sensor.divisor { - if divisor != 0 { - raw_value / divisor as f64 - } else { - raw_value - } - } else { - raw_value - }; - - let reading = SensorReading { - sensor_id: sensor.id.clone(), - value: final_value, - status: "ok".to_string(), - timestamp: Timestamp::now(), - }; - - if let Err(e) = self.storage.store_metric(&Metric::SensorReading(reading)) { - error!("Failed to store sensor reading: {}", e); - } - } - Err(e) => { - warn!( - "Failed to poll sensor {} for {}: {}", - sensor.oid, equipment.name, e - ); - } - } - } - - Ok(()) - } - - /// Poll interfaces for a piece of equipment - pub async fn poll_interfaces(&self, equipment: &EquipmentConfig) -> Result<()> { - if !equipment.snmp.enabled || equipment.interfaces.is_empty() { - return Ok(()); - } - - info!( - "Polling {} interfaces for equipment: {} (community: {})", - equipment.interfaces.len(), - equipment.name, - redact_community(&equipment.snmp.community) - ); - - for interface in &equipment.interfaces { - // OIDs for interface statistics (from IF-MIB) - let if_in_octets_oid = format!("1.3.6.1.2.1.2.2.1.10.{}", interface.if_index); - let if_out_octets_oid = format!("1.3.6.1.2.1.2.2.1.16.{}", interface.if_index); - let if_in_errors_oid = format!("1.3.6.1.2.1.2.2.1.14.{}", interface.if_index); - let if_out_errors_oid = format!("1.3.6.1.2.1.2.2.1.20.{}", interface.if_index); - 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( - &equipment.ip_address, - &equipment.snmp.community, - &equipment.snmp.version, - equipment.snmp.port, - &if_in_octets_oid, - ) - .await - .unwrap_or(0); - - 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 stat = InterfaceStat { - interface_id: interface.id.clone(), - if_in_octets: in_octets, - if_out_octets: out_octets, - if_in_errors: in_errors, - if_out_errors: out_errors, - if_in_discards: in_discards, - if_out_discards: out_discards, - timestamp: Timestamp::now(), - }; - - if let Err(e) = self.storage.store_metric(&Metric::InterfaceStat(stat)) { - error!("Failed to store interface stat: {}", e); - } - } - - Ok(()) - } - - async fn get_counter( - &self, - ip_address: &str, - community: &str, - version: &str, - port: u16, - oid: &str, - ) -> Result { - let value = self - .snmp_client - .get(ip_address, community, version, port, oid) - .await?; - - value - .as_i64() - .ok_or_else(|| ExecutorError::Conversion("Could not convert value to i64".into())) - } -} - -#[cfg(test)] -mod tests { - use super::*; - - #[test] - fn test_redact_community_normal() { - assert_eq!(redact_community("public"), "pu**"); - } - - #[test] - fn test_redact_community_short() { - assert_eq!(redact_community("ab"), "ab**"); - assert_eq!(redact_community("a"), "a**"); - } - - #[test] - fn test_redact_community_empty() { - assert_eq!(redact_community(""), "**"); - } - - #[test] - fn test_redact_community_three_chars() { - assert_eq!(redact_community("abc"), "ab**"); - } - - #[test] - fn test_redact_community_long() { - assert_eq!(redact_community("mysecretcommunity"), "my**"); - } -} diff --git a/src/poller/mod.rs b/src/poller/mod.rs deleted file mode 100644 index 8c7e04d..0000000 --- a/src/poller/mod.rs +++ /dev/null @@ -1,5 +0,0 @@ -mod executor; -mod scheduler; - -pub use executor::Executor; -pub use scheduler::Scheduler; diff --git a/src/poller/scheduler.rs b/src/poller/scheduler.rs deleted file mode 100644 index 303b883..0000000 --- a/src/poller/scheduler.rs +++ /dev/null @@ -1,248 +0,0 @@ -use crate::api_client::ApiClient; -use crate::buffer::StorageError; -use thiserror::Error; - -#[derive(Debug, Error)] -pub enum SchedulerError { - #[error("API error: {0}")] - Api(#[from] crate::api_client::ApiError), - #[error("Storage error: {0}")] - Storage(#[from] StorageError), - #[error("SNMP error: {0}")] - Snmp(#[from] crate::snmp::SnmpError), - #[error("Executor error: {0}")] - Executor(#[from] super::executor::ExecutorError), -} - -pub type Result = std::result::Result; - -use crate::buffer::Storage; -use crate::config::{AgentConfig, HeartbeatMetadata}; -use crate::poller::Executor; -use crate::snmp::SnmpClient; - -use crate::metrics::Timestamp; -use log::{error, info, warn}; -use std::time::Duration; -use tokio::time::interval; - -/// Main scheduler that orchestrates polling, config refresh, and metrics submission -pub struct Scheduler { - api_client: ApiClient, - storage: Storage, - executor: Executor, - config_refresh_seconds: u64, - current_config: Option, - start_time: Timestamp, -} - -impl Scheduler { - pub fn new( - api_client: ApiClient, - storage: Storage, - snmp_client: SnmpClient, - config_refresh_seconds: u64, - ) -> Self { - let executor = Executor::new(snmp_client, storage.clone()); - - Self { - api_client, - storage, - executor, - config_refresh_seconds, - current_config: None, - start_time: Timestamp::now(), - } - } - - /// Run the main event loop - pub async fn run(&mut self) -> Result<()> { - info!("Starting Towerops agent scheduler"); - - // Fetch initial configuration - if let Err(e) = self.refresh_config().await { - error!("Failed to fetch initial config: {}", e); - } - - let mut config_ticker = interval(Duration::from_secs(self.config_refresh_seconds)); - let mut metrics_ticker = interval(Duration::from_secs(30)); - let mut heartbeat_ticker = interval(Duration::from_secs(60)); - let mut cleanup_ticker = interval(Duration::from_secs(3600)); // Cleanup every hour - let mut poll_ticker = interval(Duration::from_secs(5)); // Check if polling needed every 5s - - loop { - tokio::select! { - _ = config_ticker.tick() => { - if let Err(e) = self.refresh_config().await { - error!("Failed to refresh config: {}", e); - } - } - - _ = metrics_ticker.tick() => { - if let Err(e) = self.flush_metrics().await { - error!("Failed to flush metrics: {}", e); - } - } - - _ = heartbeat_ticker.tick() => { - if let Err(e) = self.send_heartbeat().await { - warn!("Failed to send heartbeat: {}", e); - } - } - - _ = cleanup_ticker.tick() => { - if let Err(e) = self.storage.cleanup_old_metrics() { - error!("Failed to cleanup old metrics: {}", e); - } - } - - _ = poll_ticker.tick() => { - if let Err(e) = self.poll_equipment().await { - error!("Polling error: {}", e); - } - } - } - } - } - - async fn refresh_config(&mut self) -> Result<()> { - info!("Refreshing configuration from API"); - - match self.api_client.fetch_config().await { - Ok(config) => { - info!( - "Configuration updated: {} equipment items", - config.equipment.len() - ); - self.current_config = Some(config); - Ok(()) - } - Err(e) => { - warn!( - "Failed to fetch config, continuing with cached config: {}", - e - ); - Err(e.into()) - } - } - } - - async fn flush_metrics(&self) -> Result<()> { - let pending = self.storage.get_pending_metrics(100)?; - - if pending.is_empty() { - return Ok(()); - } - - info!("Flushing {} pending metrics to API", pending.len()); - - let ids: Vec = pending.iter().map(|(id, _)| *id).collect(); - let metrics: Vec<_> = pending.into_iter().map(|(_, m)| m).collect(); - - match self.api_client.submit_metrics(metrics).await { - Ok(_) => { - self.storage.mark_metrics_sent(&ids)?; - info!("Successfully submitted {} metrics", ids.len()); - Ok(()) - } - Err(e) => { - warn!("Failed to submit metrics, will retry later: {}", e); - Err(e.into()) - } - } - } - - async fn send_heartbeat(&self) -> Result<()> { - let uptime = self.start_time.elapsed_secs() as u64; - - let hostname = hostname::get() - .ok() - .and_then(|h| h.into_string().ok()) - .unwrap_or_else(|| "unknown".to_string()); - - let metadata = HeartbeatMetadata { - version: crate::version::current_version().to_string(), - hostname, - uptime_seconds: uptime, - }; - - self.api_client.heartbeat(metadata).await?; - Ok(()) - } - - async fn poll_equipment(&self) -> Result<()> { - let config = match &self.current_config { - Some(c) => c, - None => return Ok(()), - }; - - let poll_times = self.storage.get_all_last_poll_times()?; - - // Collect equipment that needs polling - let equipment_to_poll: Vec<_> = config - .equipment - .iter() - .filter(|equipment| { - if !equipment.snmp.enabled { - return false; - } - - // Check if it's time to poll this equipment - match poll_times.get(&equipment.id) { - Some(last_poll) => { - let elapsed = last_poll.elapsed_secs() as u64; - elapsed >= equipment.poll_interval_seconds - } - None => true, // Never polled before - } - }) - .cloned() - .collect(); - - if equipment_to_poll.is_empty() { - return Ok(()); - } - - info!( - "Polling {} devices concurrently", - equipment_to_poll.len() - ); - - // Spawn a task for each device - devices poll concurrently - let handles: Vec<_> = equipment_to_poll - .into_iter() - .map(|equipment| { - let executor = self.executor.clone(); - let storage = self.storage.clone(); - - tokio::spawn(async move { - info!("Polling equipment: {}", equipment.name); - - // Poll sensors and interfaces SEQUENTIALLY within each device - // to avoid overwhelming the device with concurrent requests - if let Err(e) = executor.poll_sensors(&equipment).await { - error!("Failed to poll sensors for {}: {}", equipment.name, e); - } - - if let Err(e) = executor.poll_interfaces(&equipment).await { - error!("Failed to poll interfaces for {}: {}", equipment.name, e); - } - - // Update last poll time - if let Err(e) = storage.update_last_poll_time(&equipment.id) { - error!("Failed to update last poll time: {}", e); - } - }) - }) - .collect(); - - // Wait for all device polling tasks to complete - for handle in handles { - if let Err(e) = handle.await { - error!("Device polling task failed: {}", e); - } - } - - Ok(()) - } -} diff --git a/src/websocket_client.rs b/src/websocket_client.rs index 8477046..f999c26 100644 --- a/src/websocket_client.rs +++ b/src/websocket_client.rs @@ -1505,11 +1505,13 @@ fn get_uptime_seconds() -> u64 { 0 } -/// Get local IP address. +/// Get local IP address by connecting a UDP socket to a public address. +/// No data is sent; the OS resolves which local interface would be used. fn get_local_ip() -> Option { - // TODO: Implement IP detection - // Could use local_ip_address crate or parse network interfaces - None + let socket = std::net::UdpSocket::bind("0.0.0.0:0").ok()?; + socket.connect("8.8.8.8:53").ok()?; + let addr = socket.local_addr().ok()?; + Some(addr.ip().to_string()) } #[cfg(test)] @@ -1620,8 +1622,11 @@ mod tests { #[test] fn test_get_local_ip() { let ip = get_local_ip(); - // Currently returns None (not implemented), just verify it's callable - assert!(ip.is_none()); + // Should resolve to a valid local IP via UDP socket trick + assert!(ip.is_some(), "Expected a local IP address"); + let ip_str = ip.unwrap(); + assert!(!ip_str.is_empty()); + assert_ne!(ip_str, "0.0.0.0"); } #[test]