v0.4.4
This commit is contained in:
@@ -5,7 +5,7 @@ use rusqlite::{params, Connection, OptionalExtension};
|
||||
use serde::{de::DeserializeOwned, Serialize};
|
||||
use serde_json::Value;
|
||||
use crate::{
|
||||
models::{ApiTokenInfo, Automation, Device, EventLog, Reading, RuntimeSettings, Schedule, Zone, ZoneReading},
|
||||
models::{ApiTokenInfo, Automation, Device, EventLog, HaReading, Reading, RuntimeSettings, Schedule, Zone, ZoneReading},
|
||||
queries,
|
||||
};
|
||||
|
||||
@@ -219,6 +219,23 @@ impl Db {
|
||||
Ok(rows_out)
|
||||
}
|
||||
|
||||
pub fn list_device_history(&self, device_id: Option<&str>, since: DateTime<Utc>, bucket_seconds: i64, limit: u32) -> Result<Vec<Reading>> {
|
||||
let conn = self.lock()?;
|
||||
let bucket_seconds = bucket_seconds.max(1);
|
||||
let limit = limit.clamp(1, 20_000) as i64;
|
||||
let mut rows_out = Vec::new();
|
||||
if let Some(device_id) = device_id {
|
||||
let mut stmt = conn.prepare(queries::LIST_DEVICE_HISTORY_BY_DEVICE_BUCKETED)?;
|
||||
let rows = stmt.query_map(params![device_id, since.to_rfc3339(), bucket_seconds, limit], Self::map_reading)?;
|
||||
for row in rows { rows_out.push(row?); }
|
||||
} else {
|
||||
let mut stmt = conn.prepare(queries::LIST_DEVICE_HISTORY_ALL_BUCKETED)?;
|
||||
let rows = stmt.query_map(params![since.to_rfc3339(), bucket_seconds, limit], Self::map_reading)?;
|
||||
for row in rows { rows_out.push(row?); }
|
||||
}
|
||||
Ok(rows_out)
|
||||
}
|
||||
|
||||
fn map_reading(row: &rusqlite::Row<'_>) -> rusqlite::Result<Reading> {
|
||||
let timestamp: String = row.get(2)?;
|
||||
Ok(Reading {
|
||||
@@ -240,7 +257,8 @@ impl Db {
|
||||
let conn = self.lock()?;
|
||||
let device = conn.execute(queries::PRUNE_READINGS, [before.to_rfc3339()])? as u64;
|
||||
let zone = conn.execute(queries::PRUNE_ZONE_READINGS, [before.to_rfc3339()])? as u64;
|
||||
Ok(device + zone)
|
||||
let ha = conn.execute(queries::PRUNE_HA_READINGS, [before.to_rfc3339()])? as u64;
|
||||
Ok(device + zone + ha)
|
||||
}
|
||||
|
||||
pub fn add_zone_reading_if_due(&self, reading: &ZoneReading, min_interval_seconds: i64) -> Result<bool> {
|
||||
@@ -311,6 +329,60 @@ impl Db {
|
||||
})
|
||||
}
|
||||
|
||||
pub fn add_ha_reading_if_due(&self, reading: &HaReading, min_interval_seconds: i64) -> Result<bool> {
|
||||
let cutoff = reading.timestamp.clone() - Duration::seconds(min_interval_seconds.max(1));
|
||||
let conn = self.lock()?;
|
||||
let changed = conn.execute(
|
||||
queries::INSERT_HA_READING_IF_DUE,
|
||||
params![
|
||||
reading.entity_id,
|
||||
reading.zone_id,
|
||||
reading.kind,
|
||||
reading.timestamp.to_rfc3339(),
|
||||
reading.temperature,
|
||||
cutoff.to_rfc3339(),
|
||||
],
|
||||
)?;
|
||||
Ok(changed > 0)
|
||||
}
|
||||
|
||||
pub fn list_ha_history(&self, entity_id: Option<&str>, since: DateTime<Utc>, bucket_seconds: i64, limit: u32) -> Result<Vec<HaReading>> {
|
||||
let conn = self.lock()?;
|
||||
let bucket_seconds = bucket_seconds.max(1);
|
||||
let limit = limit.clamp(1, 20_000) as i64;
|
||||
let mut rows_out = Vec::new();
|
||||
if let Some(entity_id) = entity_id {
|
||||
let mut stmt = conn.prepare(queries::LIST_HA_HISTORY_BY_ENTITY_BUCKETED)?;
|
||||
let rows = stmt.query_map(params![entity_id, since.to_rfc3339(), bucket_seconds, limit], Self::map_ha_reading)?;
|
||||
for row in rows { rows_out.push(row?); }
|
||||
} else {
|
||||
let mut stmt = conn.prepare(queries::LIST_HA_HISTORY_ALL_BUCKETED)?;
|
||||
let rows = stmt.query_map(params![since.to_rfc3339(), bucket_seconds, limit], Self::map_ha_reading)?;
|
||||
for row in rows { rows_out.push(row?); }
|
||||
}
|
||||
Ok(rows_out)
|
||||
}
|
||||
|
||||
fn map_ha_reading(row: &rusqlite::Row<'_>) -> rusqlite::Result<HaReading> {
|
||||
let timestamp: String = row.get(4)?;
|
||||
Ok(HaReading {
|
||||
id: row.get(0)?,
|
||||
entity_id: row.get(1)?,
|
||||
zone_id: row.get(2)?,
|
||||
kind: row.get(3)?,
|
||||
timestamp: DateTime::parse_from_rfc3339(×tamp)
|
||||
.map(|value| value.with_timezone(&Utc))
|
||||
.unwrap_or_else(|_| Utc::now()),
|
||||
temperature: row.get(5)?,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn history_counts(&self) -> Result<(u64, u64, u64)> {
|
||||
let conn = self.lock()?;
|
||||
let (device, zone, ha): (i64, i64, i64) = conn.query_row(queries::HISTORY_COUNTS, [], |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)))?;
|
||||
Ok((device.max(0) as u64, zone.max(0) as u64, ha.max(0) as u64))
|
||||
}
|
||||
|
||||
pub fn log_event(&self, level: &str, kind: &str, message: &str, metadata: &Value) -> Result<i64> {
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
@@ -399,7 +471,7 @@ impl Db {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::models::{ApiTokenInfo, Device};
|
||||
use crate::models::{ApiTokenInfo, Device, HaReading, Reading};
|
||||
|
||||
#[test]
|
||||
fn sqlite_round_trip() {
|
||||
@@ -424,5 +496,19 @@ mod tests {
|
||||
assert_eq!(db.list_api_tokens().unwrap().len(), 1);
|
||||
assert!(db.delete_api_token(&access_token.id).unwrap());
|
||||
assert!(!db.api_token_exists("test-hash").unwrap());
|
||||
|
||||
let now = Utc::now();
|
||||
db.add_reading(&Reading {
|
||||
id: 0, device_id: device.id.clone(), timestamp: now.clone(), indoor_temperature: Some(22.5),
|
||||
outdoor_temperature: Some(31.0), target_temperature: 23.0, power: true, source: "gree".into(),
|
||||
}).unwrap();
|
||||
assert_eq!(db.list_device_history(Some(&device.id), now.clone() - Duration::minutes(1), 30, 100).unwrap().len(), 1);
|
||||
|
||||
db.add_ha_reading_if_due(&HaReading {
|
||||
id: 0, entity_id: "sensor.room".into(), zone_id: Some("zone-room".into()), kind: "room".into(),
|
||||
timestamp: now.clone(), temperature: 22.1,
|
||||
}, 15).unwrap();
|
||||
assert_eq!(db.list_ha_history(Some("sensor.room"), now.clone() - Duration::minutes(1), 30, 100).unwrap().len(), 1);
|
||||
assert_eq!(db.history_counts().unwrap(), (1, 0, 1));
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user