|
|
|
@ -1,4 +1,6 @@
|
|
|
|
|
use std::sync::{Arc, Mutex}; |
|
|
|
|
use serde_json::{Result, Value}; |
|
|
|
|
use serde::{Deserialize, Serialize}; |
|
|
|
|
|
|
|
|
|
use rusqlite::{params, Connection}; |
|
|
|
|
|
|
|
|
@ -19,7 +21,10 @@ impl AnalyticUnitService {
|
|
|
|
|
conn.execute( |
|
|
|
|
"CREATE TABLE IF NOT EXISTS analytic_unit ( |
|
|
|
|
id TEXT PRIMARY KEY, |
|
|
|
|
last_detection INTEGER |
|
|
|
|
last_detection INTEGER, |
|
|
|
|
active BOOLEAN, |
|
|
|
|
type INTEGER, |
|
|
|
|
config TEXT |
|
|
|
|
)", |
|
|
|
|
[], |
|
|
|
|
)?; |
|
|
|
@ -30,7 +35,7 @@ impl AnalyticUnitService {
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// TODO: optional id
|
|
|
|
|
pub fn resolve_au(&self, cfg: AnalyticUnitConfig) -> Box<dyn types::AnalyticUnit + Send + Sync> { |
|
|
|
|
pub fn resolve_au(&self, cfg: &AnalyticUnitConfig) -> Box<dyn types::AnalyticUnit + Send + Sync> { |
|
|
|
|
match cfg { |
|
|
|
|
AnalyticUnitConfig::Threshold(c) => Box::new(ThresholdAnalyticUnit::new("1".to_string(), c.clone())), |
|
|
|
|
AnalyticUnitConfig::Pattern(c) => Box::new(PatternAnalyticUnit::new("2".to_string(), c.clone())), |
|
|
|
@ -38,7 +43,8 @@ impl AnalyticUnitService {
|
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pub fn resolve(&self, cfg: AnalyticUnitConfig) -> anyhow::Result<Box<dyn types::AnalyticUnit + Send + Sync>> { |
|
|
|
|
// TODO: get id of analytic_unit which be used also as it's type
|
|
|
|
|
pub fn resolve(&self, cfg: &AnalyticUnitConfig) -> anyhow::Result<Box<dyn types::AnalyticUnit + Send + Sync>> { |
|
|
|
|
let au = self.resolve_au(cfg); |
|
|
|
|
let id = au.as_ref().get_id(); |
|
|
|
|
|
|
|
|
@ -49,16 +55,25 @@ impl AnalyticUnitService {
|
|
|
|
|
let res = stmt.exists(params![id])?; |
|
|
|
|
|
|
|
|
|
if res == false { |
|
|
|
|
let cfg_json = serde_json::to_string(&cfg)?; |
|
|
|
|
conn.execute( |
|
|
|
|
"INSERT INTO analytic_unit (id) VALUES (?1)", |
|
|
|
|
params![id] |
|
|
|
|
"INSERT INTO analytic_unit (id, type, config) VALUES (?1, ?1, ?2)", |
|
|
|
|
params![id, cfg_json] |
|
|
|
|
)?; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
conn.execute( |
|
|
|
|
"UPDATE analytic_unit set active = FALSE where active = TRUE", |
|
|
|
|
params![] |
|
|
|
|
)?; |
|
|
|
|
conn.execute( |
|
|
|
|
"UPDATE analytic_unit set active = TRUE where id = ?1", |
|
|
|
|
params![id] |
|
|
|
|
)?; |
|
|
|
|
|
|
|
|
|
return Ok(au); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// TODO: resolve with saving by id
|
|
|
|
|
pub fn set_last_detection(&self, id: String, last_detection: u64) -> anyhow::Result<()> { |
|
|
|
|
let conn = self.connection.lock().unwrap(); |
|
|
|
|
conn.execute( |
|
|
|
@ -67,4 +82,61 @@ impl AnalyticUnitService {
|
|
|
|
|
)?; |
|
|
|
|
Ok(()) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pub fn get_active(&self) -> anyhow::Result<Box<dyn types::AnalyticUnit + Send + Sync>> { |
|
|
|
|
// TODO: return default when there is no active
|
|
|
|
|
let conn = self.connection.lock().unwrap(); |
|
|
|
|
let mut stmt = conn.prepare( |
|
|
|
|
"SELECT id, type, config from analytic_unit WHERE active = TRUE" |
|
|
|
|
)?; |
|
|
|
|
|
|
|
|
|
let au = stmt.query_row([], |row| { |
|
|
|
|
let c: String = row.get(2)?; |
|
|
|
|
let cfg: AnalyticUnitConfig = serde_json::from_str(&c).unwrap(); |
|
|
|
|
Ok(self.resolve(&cfg)) |
|
|
|
|
})??; |
|
|
|
|
|
|
|
|
|
return Ok(au); |
|
|
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pub fn get_active_config(&self) -> anyhow::Result<AnalyticUnitConfig> { |
|
|
|
|
let exists = { |
|
|
|
|
let conn = self.connection.lock().unwrap(); |
|
|
|
|
let mut stmt = conn.prepare( |
|
|
|
|
"SELECT config from analytic_unit WHERE active = TRUE" |
|
|
|
|
)?; |
|
|
|
|
stmt.exists([])? |
|
|
|
|
}; |
|
|
|
|
|
|
|
|
|
if exists == false { |
|
|
|
|
let c = AnalyticUnitConfig::Pattern(Default::default()); |
|
|
|
|
self.resolve(&c)?; |
|
|
|
|
return Ok(c); |
|
|
|
|
} else { |
|
|
|
|
let conn = self.connection.lock().unwrap(); |
|
|
|
|
let mut stmt = conn.prepare( |
|
|
|
|
"SELECT config from analytic_unit WHERE active = TRUE" |
|
|
|
|
)?; |
|
|
|
|
let acfg = stmt.query_row([], |row| { |
|
|
|
|
let c: String = row.get(0)?; |
|
|
|
|
let cfg = serde_json::from_str(&c).unwrap(); |
|
|
|
|
Ok(cfg) |
|
|
|
|
})?; |
|
|
|
|
return Ok(acfg); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
pub fn update_active_config(&self, cfg: &AnalyticUnitConfig) -> anyhow::Result<()> { |
|
|
|
|
let conn = self.connection.lock().unwrap(); |
|
|
|
|
|
|
|
|
|
let cfg_json = serde_json::to_string(&cfg)?; |
|
|
|
|
|
|
|
|
|
conn.execute( |
|
|
|
|
"UPDATE analytic_unit SET config = ?1 WHERE active = TRUE", |
|
|
|
|
params![cfg_json] |
|
|
|
|
)?; |
|
|
|
|
|
|
|
|
|
return Ok(()); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|