From 666ef4268864f905267c84cdbf76ddbfe149a963 Mon Sep 17 00:00:00 2001 From: Gvidas Juknevicius Date: Mon, 7 Sep 2026 16:41:26 +0300 Subject: [PATCH] Add new tables, add automatic data cleanup with configurable duration, move company business logic into seperate folder --- .../2026-08-20-121835-0000_initial/up.sql | 8 +- .../down.sql | 1 + .../up.sql | 19 ++ .../down.sql | 1 + .../up.sql | 6 + src/bes_wifi/mod.rs | 261 ++++++++++++++++++ src/config/cli.rs | 3 + src/db/mod.rs | 3 +- src/db/models.rs | 52 ---- src/db/models/antenna.rs | 20 ++ src/db/models/client_readings.rs | 83 ++++++ src/db/models/client_statistics.rs | 49 ++++ src/db/models/mod.rs | 9 + src/db/models/wara.rs | 18 ++ src/db/schema.rs | 65 ++++- src/main.rs | 193 +------------ 16 files changed, 537 insertions(+), 254 deletions(-) create mode 100644 migrations/2026-09-07-062247-0000_add_client_statistics_table/down.sql create mode 100644 migrations/2026-09-07-062247-0000_add_client_statistics_table/up.sql create mode 100644 migrations/2026-09-07-073346-0000_add_wara_system_settings_table/down.sql create mode 100644 migrations/2026-09-07-073346-0000_add_wara_system_settings_table/up.sql create mode 100644 src/bes_wifi/mod.rs delete mode 100644 src/db/models.rs create mode 100644 src/db/models/antenna.rs create mode 100644 src/db/models/client_readings.rs create mode 100644 src/db/models/client_statistics.rs create mode 100644 src/db/models/mod.rs create mode 100644 src/db/models/wara.rs diff --git a/migrations/2026-08-20-121835-0000_initial/up.sql b/migrations/2026-08-20-121835-0000_initial/up.sql index 3cbbd25..71a2c31 100644 --- a/migrations/2026-08-20-121835-0000_initial/up.sql +++ b/migrations/2026-08-20-121835-0000_initial/up.sql @@ -9,16 +9,18 @@ CREATE TABLE IF NOT EXISTS antennas ( CREATE TABLE IF NOT EXISTS client_readings ( id BIGINT AUTO_INCREMENT PRIMARY KEY, - ip VARCHAR(15) NOT NULL, + + ip VARCHAR(15), mac VARCHAR(17) NOT NULL, antenna_ip VARCHAR(15) NOT NULL, - db_reading INT NOT NULL, read_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + + db_reading INT NOT NULL, tx_ccq INT, rx_ccq INT, radio_name VARCHAR(255), tx_rate VARCHAR(50), rx_rate VARCHAR(50), - + CONSTRAINT fk_client_readings_antenna FOREIGN KEY (antenna_ip) REFERENCES antennas(ip) ); diff --git a/migrations/2026-09-07-062247-0000_add_client_statistics_table/down.sql b/migrations/2026-09-07-062247-0000_add_client_statistics_table/down.sql new file mode 100644 index 0000000..a0d484b --- /dev/null +++ b/migrations/2026-09-07-062247-0000_add_client_statistics_table/down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS client_statistics; diff --git a/migrations/2026-09-07-062247-0000_add_client_statistics_table/up.sql b/migrations/2026-09-07-062247-0000_add_client_statistics_table/up.sql new file mode 100644 index 0000000..fd5178b --- /dev/null +++ b/migrations/2026-09-07-062247-0000_add_client_statistics_table/up.sql @@ -0,0 +1,19 @@ + +CREATE TABLE IF NOT EXISTS client_statistics ( + id BIGINT AUTO_INCREMENT PRIMARY KEY, + + ip VARCHAR(15), + mac VARCHAR(17) NOT NULL, + antenna_ip VARCHAR(15) NOT NULL, + read_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + + db_reading INT NOT NULL, + tx_ccq INT, + rx_ccq INT, + radio_name VARCHAR(255), + tx_rate VARCHAR(50), + rx_rate VARCHAR(50), + + CONSTRAINT uq_client_statistics_mac UNIQUE (mac), + CONSTRAINT fk_client_statistics_antenna FOREIGN KEY (antenna_ip) REFERENCES antennas(ip) +); diff --git a/migrations/2026-09-07-073346-0000_add_wara_system_settings_table/down.sql b/migrations/2026-09-07-073346-0000_add_wara_system_settings_table/down.sql new file mode 100644 index 0000000..6feca8f --- /dev/null +++ b/migrations/2026-09-07-073346-0000_add_wara_system_settings_table/down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS wara; diff --git a/migrations/2026-09-07-073346-0000_add_wara_system_settings_table/up.sql b/migrations/2026-09-07-073346-0000_add_wara_system_settings_table/up.sql new file mode 100644 index 0000000..417e889 --- /dev/null +++ b/migrations/2026-09-07-073346-0000_add_wara_system_settings_table/up.sql @@ -0,0 +1,6 @@ + +CREATE TABLE IF NOT EXISTS wara ( + id BIGINT AUTO_INCREMENT PRIMARY KEY, + + last_client_reading_cleanup TIMESTAMP +); diff --git a/src/bes_wifi/mod.rs b/src/bes_wifi/mod.rs new file mode 100644 index 0000000..0856b8a --- /dev/null +++ b/src/bes_wifi/mod.rs @@ -0,0 +1,261 @@ +use std::collections::HashMap; +use std::ops::Deref; +use std::str::FromStr; + +use chrono::{Duration, NaiveDateTime, Utc}; +use diesel::query_builder::AsQuery; +use diesel::query_dsl::methods::{FilterDsl, SelectDsl}; +use diesel::{ExpressionMethods, OptionalExtension, QueryResult, RunQueryDsl, sql_query}; + +use crate::config::{Config, ConfigHost}; +use crate::db::DbConnection; +use crate::db::models::*; + +pub async fn run_bes_wifi_collection( + interval: &tokio::time::Duration, + cleanup_interval: &tokio::time::Duration, + cfg: &Config, +) -> anyhow::Result<()> { + let mut conn = crate::db::init_db(&cfg.db_url); + + for host in &cfg.hosts { + use crate::db::schema::antennas; + diesel::insert_into(antennas::table) + .values(&NewAntenna { + name: host.name.clone(), + ip: host.ip.to_string(), + error: None, + }) + .on_conflict(diesel::dsl::DuplicatedKeys) + .do_update() + .set(( + antennas::name.eq(host.name.clone()), + antennas::error.eq(&None as &Option), + )) + .execute(&mut conn) + .unwrap(); + } + + let v6_hosts = cfg + .hosts + .iter() + .filter(|v| v.tags.contains(&"v6".to_string())) + .cloned() + .collect::>(); + + let v7_hosts = cfg + .hosts + .iter() + .filter(|v| v.tags.contains(&"v7".to_string())) + .cloned() + .collect::>(); + + let v6_cmd = "/interface/wireless/registration-table/print"; + let v7_cmd = "/interface/wifi/registration-table/print"; + + let mut interval = tokio::time::interval(*interval); + + let mut last_cleanup: Option = None; + + loop { + interval.tick().await; + + { + if let Some(last) = last_cleanup { + let now = Utc::now().naive_utc(); + if now - last > chrono::Duration::from_std(cleanup_interval.clone())? { + cleanup_old_client_readings(&mut conn)?; + } + } else { + use crate::db::schema::wara; + + let row = wara::table.first::(&mut conn).optional()?; + if let Some(row) = row { + last_cleanup = row.last_client_reading_cleanup; + } else { + diesel::insert_into(wara::table) + .values(NewWara { + last_client_reading_cleanup: None, + }) + .execute(&mut conn)?; + } + } + } + + let (v6_res, v7_res) = futures::join!( + crate::mt_commander::MtCommander::run_command_on_hosts(&v6_hosts, v6_cmd), + crate::mt_commander::MtCommander::run_command_on_hosts(&v7_hosts, v7_cmd), + ); + let mut records = Vec::new(); + + for (antenna, res) in v6_res?.iter().chain(v7_res?.iter()) { + match res { + Ok(vals) if antenna.tags.contains(&"v7".to_string()) => { + for val in vals { + records.push(NewClientReadingTimed { + ip: Default::default(), + mac: get_val_from_cmd_output(antenna, val, "mac-address") + .cloned() + .unwrap_or_default(), + antenna_ip: antenna.ip.to_string(), + db_reading: get_val_from_cmd_output(antenna, val, "signal") + .map(|db| parse_int::parse::(db).expect("Unparseble int")) + .unwrap_or_default(), + radio_name: Default::default(), + rx_ccq: Default::default(), + tx_ccq: Default::default(), + tx_rate: get_val_from_cmd_output(antenna, val, "tx-rate") + .cloned() + .ok(), + rx_rate: get_val_from_cmd_output(antenna, val, "rx-rate") + .cloned() + .ok(), + read_at: Utc::now().naive_utc(), + }); + } + } + Ok(vals) if antenna.tags.contains(&"v6".to_string()) => { + for val in vals { + records.push(NewClientReadingTimed { + ip: get_val_from_cmd_output(antenna, val, "last-ip") + .cloned() + .ok(), + mac: get_val_from_cmd_output(antenna, val, "mac-address") + .cloned() + .unwrap_or_default(), + antenna_ip: antenna.ip.to_string(), + db_reading: get_val_from_cmd_output(antenna, val, "signal-strength") + .cloned() + .map(|mut db| { + if db.contains('@') { + db = db.split('@').nth(0).unwrap().to_string(); + } + parse_int::parse::(&db).expect("Unparseble int") + }) + .unwrap_or_default(), + radio_name: get_val_from_cmd_output(antenna, val, "radio-name") + .cloned() + .ok(), + rx_ccq: get_val_from_cmd_output(antenna, val, "rx-ccq") + .map(|db| parse_int::parse::(db).expect("Unparseble int")) + .ok(), + tx_ccq: get_val_from_cmd_output(antenna, val, "tx-ccq") + .map(|db| parse_int::parse::(db).expect("Unparseble int")) + .ok(), + tx_rate: get_val_from_cmd_output(antenna, val, "tx-rate") + .cloned() + .ok(), + rx_rate: get_val_from_cmd_output(antenna, val, "rx-rate") + .cloned() + .ok(), + read_at: Utc::now().naive_utc(), + }); + } + } + Ok(_) => unreachable!(), + Err(e) => { + use crate::db::schema::antennas; + diesel::update(antennas::table.filter(antennas::ip.eq(antenna.ip.to_string()))) + .set(antennas::error.eq(e.to_string())) + .execute(&mut conn)?; + } + } + } + + log::debug!("inserting: {records:?}"); + diesel::insert_into(crate::db::schema::client_readings::table) + .values(&records) + .execute(&mut conn) + .unwrap(); + upsert_client_statistics(&mut conn, &records)?; + } +} + +fn get_val_from_cmd_output<'a, T>( + antenna: &ConfigHost, + command_output: &'a HashMap>, + key: &str, +) -> anyhow::Result<&'a T> { + match command_output.get(key) { + Some(Some(val)) => Ok(val), + Some(None) => { + log::warn!( + "[{} | {}] Value of '{key}' was none", + antenna.ip, + antenna.name + ); + anyhow::bail!( + "[{} | {}] Value of '{key}' was none", + antenna.ip, + antenna.name + ) + } + None => { + log::warn!( + "[{} | {}] Unable to find '{key}' in command output", + antenna.ip, + antenna.name + ); + anyhow::bail!( + "[{} | {}] Unable to find '{key}' in command output", + antenna.ip, + antenna.name + ) + } + } +} + +pub fn upsert_client_statistics( + conn: &mut crate::db::DbConnection, + readings: &Vec, +) -> anyhow::Result<()> { + use crate::db::schema::client_statistics; + for reading in readings.into_iter() { + let existing_id = client_statistics::table + .filter(client_statistics::mac.eq(&reading.mac)) + .select(client_statistics::id) + .first::(conn) + .optional()?; + + match existing_id { + Some(id) => { + diesel::update(client_statistics::table.filter(client_statistics::id.eq(id))) + .set(( + client_statistics::ip.eq(&reading.ip), + client_statistics::antenna_ip.eq(&reading.antenna_ip), + client_statistics::db_reading.eq(reading.db_reading), + client_statistics::read_at.eq(reading.read_at), + client_statistics::tx_ccq.eq(reading.tx_ccq), + client_statistics::rx_ccq.eq(reading.rx_ccq), + client_statistics::radio_name.eq(&reading.radio_name), + client_statistics::tx_rate.eq(&reading.tx_rate), + client_statistics::rx_rate.eq(&reading.rx_rate), + )) + .execute(conn)?; + } + + None => { + let val: NewClientStatisticWithTime = reading.clone().into(); + diesel::insert_into(client_statistics::table) + .values(val) + .execute(conn)?; + } + } + } + Ok(()) +} + +pub fn cleanup_old_client_readings(conn: &mut DbConnection) -> anyhow::Result<()> { + use crate::db::schema::{client_readings, wara}; + + let cutoff = Utc::now().naive_utc() - Duration::days(7); + + let count = diesel::delete(client_readings::table.filter(client_readings::read_at.lt(cutoff))) + .execute(conn)?; + log::info!("Deleted {count} old 'client_readings' entries"); + + diesel::update(wara::table.filter(wara::id.eq(1))) + .set(wara::last_client_reading_cleanup.eq(Some(Utc::now().naive_utc()))) + .execute(conn)?; + Ok(()) +} diff --git a/src/config/cli.rs b/src/config/cli.rs index 8eba7ff..170d4e8 100644 --- a/src/config/cli.rs +++ b/src/config/cli.rs @@ -94,6 +94,9 @@ pub enum CliCommand { BesWifi { #[arg(long, short = 'i', value_parser = humantime::parse_duration, default_value="5m")] interval: std::time::Duration, + + #[arg(long, short = 'C', value_parser = humantime::parse_duration, default_value="1w")] + cleanup_interval: std::time::Duration, }, #[default] diff --git a/src/db/mod.rs b/src/db/mod.rs index 7075e78..85fb18d 100644 --- a/src/db/mod.rs +++ b/src/db/mod.rs @@ -1,4 +1,5 @@ -use diesel::MysqlConnection as DbConnection; +pub use diesel::MysqlConnection as DbConnection; +pub use diesel::mysql::Mysql as DbBackend; use diesel::prelude::*; use diesel_migrations::{EmbeddedMigrations, MigrationHarness, embed_migrations}; diff --git a/src/db/models.rs b/src/db/models.rs deleted file mode 100644 index 651b019..0000000 --- a/src/db/models.rs +++ /dev/null @@ -1,52 +0,0 @@ -#![allow(dead_code)] -use chrono::NaiveDateTime; -use diesel::prelude::*; - -#[derive(Debug, Queryable, Selectable)] -#[diesel(table_name = crate::db::schema::antennas)] -#[diesel(check_for_backend(diesel::mysql::Mysql))] -pub struct Antenna { - pub id: i64, - pub name: String, - pub ip: String, - pub error: Option, -} - -#[derive(Debug, Insertable)] -#[diesel(table_name = crate::db::schema::antennas)] -pub struct NewAntenna { - pub name: String, - pub ip: String, - pub error: Option, -} - -#[derive(Debug, Queryable, Selectable)] -#[diesel(table_name = crate::db::schema::client_readings)] -#[diesel(check_for_backend(diesel::mysql::Mysql))] -pub struct ClientReading { - pub id: i64, - pub ip: String, - pub mac: String, - pub antenna_ip: String, - pub db_reading: i32, - pub read_at: NaiveDateTime, - pub tx_ccq: i32, - pub rx_ccq: i32, - pub radio_name: String, - pub tx_rate: String, - pub rx_rate: String, -} - -#[derive(Debug, Insertable)] -#[diesel(table_name = crate::db::schema::client_readings)] -pub struct NewClientReading { - pub ip: String, - pub mac: String, - pub antenna_ip: String, - pub db_reading: i32, - pub tx_ccq: i32, - pub rx_ccq: i32, - pub radio_name: String, - pub tx_rate: String, - pub rx_rate: String, -} diff --git a/src/db/models/antenna.rs b/src/db/models/antenna.rs new file mode 100644 index 0000000..2e47afa --- /dev/null +++ b/src/db/models/antenna.rs @@ -0,0 +1,20 @@ +#![allow(dead_code)] +use diesel::prelude::*; + +#[derive(Debug, Queryable, Selectable)] +#[diesel(table_name = crate::db::schema::antennas)] +#[diesel(check_for_backend(crate::db::DbBackend))] +pub struct Antenna { + pub id: i64, + pub name: String, + pub ip: String, + pub error: Option, +} + +#[derive(Debug, Insertable)] +#[diesel(table_name = crate::db::schema::antennas)] +pub struct NewAntenna { + pub name: String, + pub ip: String, + pub error: Option, +} diff --git a/src/db/models/client_readings.rs b/src/db/models/client_readings.rs new file mode 100644 index 0000000..01a6ca5 --- /dev/null +++ b/src/db/models/client_readings.rs @@ -0,0 +1,83 @@ +#![allow(dead_code)] +use chrono::NaiveDateTime; +use diesel::prelude::*; + +use crate::db::models::NewClientStatisticWithTime; + +#[derive(Debug, Clone, Queryable, Selectable)] +#[diesel(table_name = crate::db::schema::client_readings)] +#[diesel(check_for_backend(crate::db::DbBackend))] +pub struct ClientReading { + pub id: i64, + pub ip: Option, + pub mac: String, + pub antenna_ip: String, + pub db_reading: i32, + pub read_at: NaiveDateTime, + pub tx_ccq: Option, + pub rx_ccq: Option, + pub radio_name: Option, + pub tx_rate: Option, + pub rx_rate: Option, +} + +#[derive(Debug, Clone, Insertable)] +#[diesel(table_name = crate::db::schema::client_readings)] +pub struct NewClientReading { + pub ip: Option, + pub mac: String, + pub antenna_ip: String, + pub db_reading: i32, + pub tx_ccq: Option, + pub rx_ccq: Option, + pub radio_name: Option, + pub tx_rate: Option, + pub rx_rate: Option, +} +#[derive(Debug, Clone, Insertable)] +#[diesel(table_name = crate::db::schema::client_readings)] +pub struct NewClientReadingTimed { + pub ip: Option, + pub mac: String, + pub antenna_ip: String, + pub db_reading: i32, + pub read_at: NaiveDateTime, + pub tx_ccq: Option, + pub rx_ccq: Option, + pub radio_name: Option, + pub tx_rate: Option, + pub rx_rate: Option, +} + +impl Into for NewClientReadingTimed { + fn into(self) -> NewClientReading { + NewClientReading { + ip: self.ip, + mac: self.mac, + antenna_ip: self.antenna_ip, + db_reading: self.db_reading, + tx_ccq: self.tx_ccq, + rx_ccq: self.rx_ccq, + radio_name: self.radio_name, + tx_rate: self.tx_rate, + rx_rate: self.rx_rate, + } + } +} + +impl Into for NewClientReadingTimed { + fn into(self) -> NewClientStatisticWithTime { + NewClientStatisticWithTime { + ip: self.ip, + mac: self.mac, + antenna_ip: self.antenna_ip, + db_reading: self.db_reading, + read_at: self.read_at, + tx_ccq: self.tx_ccq, + rx_ccq: self.rx_ccq, + radio_name: self.radio_name, + tx_rate: self.tx_rate, + rx_rate: self.rx_rate, + } + } +} diff --git a/src/db/models/client_statistics.rs b/src/db/models/client_statistics.rs new file mode 100644 index 0000000..4d6e37e --- /dev/null +++ b/src/db/models/client_statistics.rs @@ -0,0 +1,49 @@ +#![allow(dead_code)] +use chrono::NaiveDateTime; +use diesel::prelude::*; + +#[derive(Debug, Clone, Queryable, Selectable)] +#[diesel(table_name = crate::db::schema::client_statistics)] +#[diesel(check_for_backend(crate::db::DbBackend))] +pub struct ClientStatistics { + pub id: i64, + pub ip: Option, + pub mac: String, + pub antenna_ip: String, + pub db_reading: i32, + pub read_at: NaiveDateTime, + pub tx_ccq: Option, + pub rx_ccq: Option, + pub radio_name: Option, + pub tx_rate: Option, + pub rx_rate: Option, +} + +#[derive(Debug, Clone, Insertable)] +#[diesel(table_name = crate::db::schema::client_statistics)] +pub struct NewClientStatistic { + pub ip: Option, + pub mac: String, + pub antenna_ip: String, + pub db_reading: i32, + pub tx_ccq: Option, + pub rx_ccq: Option, + pub radio_name: Option, + pub tx_rate: Option, + pub rx_rate: Option, +} + +#[derive(Debug, Clone, Insertable)] +#[diesel(table_name = crate::db::schema::client_statistics)] +pub struct NewClientStatisticWithTime { + pub ip: Option, + pub mac: String, + pub antenna_ip: String, + pub db_reading: i32, + pub read_at: NaiveDateTime, + pub tx_ccq: Option, + pub rx_ccq: Option, + pub radio_name: Option, + pub tx_rate: Option, + pub rx_rate: Option, +} diff --git a/src/db/models/mod.rs b/src/db/models/mod.rs new file mode 100644 index 0000000..db559e6 --- /dev/null +++ b/src/db/models/mod.rs @@ -0,0 +1,9 @@ +mod antenna; +mod client_readings; +mod client_statistics; +mod wara; + +pub use antenna::*; +pub use client_readings::*; +pub use client_statistics::*; +pub use wara::*; diff --git a/src/db/models/wara.rs b/src/db/models/wara.rs new file mode 100644 index 0000000..f1dba00 --- /dev/null +++ b/src/db/models/wara.rs @@ -0,0 +1,18 @@ +#![allow(dead_code)] +use chrono::NaiveDateTime; +use diesel::prelude::*; + +#[derive(Debug, Queryable, Selectable)] +#[diesel(table_name = crate::db::schema::wara)] +#[diesel(check_for_backend(crate::db::DbBackend))] +pub struct Wara { + pub id: i64, + + pub last_client_reading_cleanup: Option, +} + +#[derive(Debug, Insertable)] +#[diesel(table_name = crate::db::schema::wara)] +pub struct NewWara { + pub last_client_reading_cleanup: Option, +} diff --git a/src/db/schema.rs b/src/db/schema.rs index c77f463..510f639 100644 --- a/src/db/schema.rs +++ b/src/db/schema.rs @@ -1,28 +1,65 @@ -// Written manually because diesel cant fucking read the db properly for some reason +// @generated automatically by Diesel CLI. diesel::table! { antennas (id) { - id -> Int8, + id -> Bigint, name -> Text, - ip -> Text, + #[max_length = 15] + ip -> Varchar, error -> Nullable, + last_error_t -> Nullable, } } diesel::table! { client_readings (id) { - id -> Int8, - ip -> Text, - mac -> Text, - antenna_ip -> Text, - db_reading -> Int4, + id -> Bigint, + #[max_length = 15] + ip -> Nullable, + #[max_length = 17] + mac -> Varchar, + #[max_length = 15] + antenna_ip -> Varchar, read_at -> Timestamp, - tx_ccq -> Int4, - rx_ccq -> Int4, - radio_name -> Varchar, - tx_rate -> Varchar, - rx_rate -> Varchar, + db_reading -> Integer, + tx_ccq -> Nullable, + rx_ccq -> Nullable, + #[max_length = 255] + radio_name -> Nullable, + #[max_length = 50] + tx_rate -> Nullable, + #[max_length = 50] + rx_rate -> Nullable, } } -diesel::allow_tables_to_appear_in_same_query!(antennas, client_readings,); +diesel::table! { + client_statistics (id) { + id -> Bigint, + #[max_length = 15] + ip -> Nullable, + #[max_length = 17] + mac -> Varchar, + #[max_length = 15] + antenna_ip -> Varchar, + read_at -> Timestamp, + db_reading -> Integer, + tx_ccq -> Nullable, + rx_ccq -> Nullable, + #[max_length = 255] + radio_name -> Nullable, + #[max_length = 50] + tx_rate -> Nullable, + #[max_length = 50] + rx_rate -> Nullable, + } +} + +diesel::table! { + wara (id) { + id -> Bigint, + last_client_reading_cleanup -> Nullable, + } +} + +diesel::allow_tables_to_appear_in_same_query!(antennas, client_readings, client_statistics, wara,); diff --git a/src/main.rs b/src/main.rs index d4a24f4..4f3a325 100644 --- a/src/main.rs +++ b/src/main.rs @@ -13,22 +13,16 @@ use std::{ str::FromStr, }; -use diesel::{ExpressionMethods, RunQueryDsl, query_dsl::methods::FilterDsl}; use log::LevelFilter; use serde_json::json; -use tokio::time::{self, Duration}; +use tokio::time::Duration; -use crate::{ - config::{ - Config, ConfigHost, - cli::{CliCommand, CliDeviceCommand, OutputFormat}, - }, - db::{ - models::{NewAntenna, NewClientReading}, - schema::antennas, - }, +use crate::config::{ + Config, ConfigHost, + cli::{CliCommand, CliDeviceCommand, OutputFormat}, }; +mod bes_wifi; mod config; mod db; mod mt_commander; @@ -55,179 +49,10 @@ async fn main() -> anyhow::Result<()> { mt_commander::MtCommander::init(cfg.clone())?; match &cfg.cli.command { - CliCommand::BesWifi { interval } => { - let mut conn = db::init_db(&cfg.db_url); - - for host in &cfg.hosts { - use crate::db::schema::antennas; - diesel::insert_into(antennas::table) - .values(&NewAntenna { - name: host.name.clone(), - ip: host.ip.to_string(), - error: None, - }) - .on_conflict(diesel::dsl::DuplicatedKeys) - .do_update() - .set(( - antennas::name.eq(host.name.clone()), - antennas::error.eq(&None as &Option), - )) - .execute(&mut conn) - .unwrap(); - } - - let v6_hosts = cfg - .hosts - .iter() - .filter(|v| v.tags.contains(&"v6".to_string())) - .cloned() - .collect::>(); - - let v7_hosts = cfg - .hosts - .iter() - .filter(|v| v.tags.contains(&"v7".to_string())) - .cloned() - .collect::>(); - - let v6_cmd = "/interface/wireless/registration-table/print"; - let v7_cmd = "/interface/wifi/registration-table/print"; - - let mut interval = time::interval(*interval); - - loop { - interval.tick().await; - - let (v6_res, v7_res) = futures::join!( - mt_commander::MtCommander::run_command_on_hosts(&v6_hosts, v6_cmd), - mt_commander::MtCommander::run_command_on_hosts(&v7_hosts, v7_cmd), - ); - let mut records = Vec::new(); - - fn get_val_from_cmd_output<'a, T>( - antenna: &ConfigHost, - command_output: &'a HashMap>, - key: &str, - ) -> anyhow::Result<&'a T> { - match command_output.get(key) { - Some(Some(val)) => Ok(val), - Some(None) => { - log::warn!( - "[{} | {}] Value of '{key}' was none", - antenna.ip, - antenna.name - ); - anyhow::bail!( - "[{} | {}] Value of '{key}' was none", - antenna.ip, - antenna.name - ) - } - None => { - log::warn!( - "[{} | {}] Unable to find '{key}' in command output", - antenna.ip, - antenna.name - ); - anyhow::bail!( - "[{} | {}] Unable to find '{key}' in command output", - antenna.ip, - antenna.name - ) - } - } - } - - for (antenna, res) in v6_res?.iter().chain(v7_res?.iter()) { - match res { - Ok(vals) if antenna.tags.contains(&"v7".to_string()) => { - for val in vals { - records.push(NewClientReading { - ip: Default::default(), - mac: get_val_from_cmd_output(antenna, val, "mac-address") - .cloned() - .unwrap_or_default(), - antenna_ip: antenna.ip.to_string(), - db_reading: get_val_from_cmd_output(antenna, val, "signal") - .map(|db| { - parse_int::parse::(db).expect("Unparseble int") - }) - .unwrap_or_default(), - radio_name: Default::default(), - rx_ccq: Default::default(), - tx_ccq: Default::default(), - tx_rate: get_val_from_cmd_output(antenna, val, "tx-rate") - .cloned() - .unwrap_or_default(), - rx_rate: get_val_from_cmd_output(antenna, val, "rx-rate") - .cloned() - .unwrap_or_default(), - }); - } - } - Ok(vals) if antenna.tags.contains(&"v6".to_string()) => { - for val in vals { - records.push(NewClientReading { - ip: get_val_from_cmd_output(antenna, val, "last-ip") - .cloned() - .unwrap_or_default(), - mac: get_val_from_cmd_output(antenna, val, "mac-address") - .cloned() - .unwrap_or_default(), - antenna_ip: antenna.ip.to_string(), - db_reading: get_val_from_cmd_output( - antenna, - val, - "signal-strength", - ) - .cloned() - .map(|mut db| { - if db.contains('@') { - db = db.split('@').nth(0).unwrap().to_string(); - } - parse_int::parse::(&db).expect("Unparseble int") - }) - .unwrap_or_default(), - radio_name: get_val_from_cmd_output(antenna, val, "radio-name") - .cloned() - .unwrap_or_default(), - rx_ccq: get_val_from_cmd_output(antenna, val, "rx-ccq") - .map(|db| { - parse_int::parse::(db).expect("Unparseble int") - }) - .unwrap_or_default(), - tx_ccq: get_val_from_cmd_output(antenna, val, "tx-ccq") - .map(|db| { - parse_int::parse::(db).expect("Unparseble int") - }) - .unwrap_or_default(), - tx_rate: get_val_from_cmd_output(antenna, val, "tx-rate") - .cloned() - .unwrap_or_default(), - rx_rate: get_val_from_cmd_output(antenna, val, "rx-rate") - .cloned() - .unwrap_or_default(), - }); - } - } - Ok(_) => unreachable!(), - Err(e) => { - diesel::update( - antennas::table.filter(antennas::ip.eq(antenna.ip.to_string())), - ) - .set(antennas::error.eq(e.to_string())) - .execute(&mut conn)?; - } - } - } - - log::debug!("inserting: {records:?}"); - diesel::insert_into(crate::db::schema::client_readings::table) - .values(&records) - .execute(&mut conn) - .unwrap(); - } - } + CliCommand::BesWifi { + interval, + cleanup_interval, + } => bes_wifi::run_bes_wifi_collection(interval, cleanup_interval, &cfg).await?, CliCommand::Run { all, devices,