first commit
This commit is contained in:
@@ -0,0 +1,341 @@
|
||||
use std::{path::Path, sync::{Arc, Mutex}};
|
||||
use anyhow::{Context, Result};
|
||||
use chrono::{DateTime, Duration, Utc};
|
||||
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},
|
||||
queries,
|
||||
};
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Db {
|
||||
conn: Arc<Mutex<Connection>>,
|
||||
}
|
||||
|
||||
impl Db {
|
||||
pub fn open(path: &Path) -> Result<Self> {
|
||||
let conn = Connection::open(path)
|
||||
.with_context(|| format!("cannot open SQLite database {}", path.display()))?;
|
||||
conn.busy_timeout(std::time::Duration::from_secs(5))?;
|
||||
conn.execute_batch(queries::INIT_SCHEMA)?;
|
||||
Ok(Self { conn: Arc::new(Mutex::new(conn)) })
|
||||
}
|
||||
|
||||
fn lock(&self) -> Result<std::sync::MutexGuard<'_, Connection>> {
|
||||
self.conn.lock().map_err(|_| anyhow::anyhow!("database mutex poisoned"))
|
||||
}
|
||||
|
||||
fn from_json<T: DeserializeOwned>(payload: String) -> Result<T> {
|
||||
Ok(serde_json::from_str(&payload)?)
|
||||
}
|
||||
|
||||
fn to_json<T: Serialize>(value: &T) -> Result<String> {
|
||||
Ok(serde_json::to_string(value)?)
|
||||
}
|
||||
|
||||
pub fn count_devices(&self) -> Result<u64> {
|
||||
let conn = self.lock()?;
|
||||
let count: i64 = conn.query_row(queries::COUNT_DEVICES, [], |row| row.get(0))?;
|
||||
Ok(count.max(0) as u64)
|
||||
}
|
||||
|
||||
pub fn save_device(&self, device: &Device) -> Result<()> {
|
||||
let payload = Self::to_json(device)?;
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::UPSERT_DEVICE,
|
||||
params![device.id, device.mac, device.name, device.ip, device.simulated as i64, payload, device.updated_at.to_rfc3339()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn list_devices(&self) -> Result<Vec<Device>> {
|
||||
let conn = self.lock()?;
|
||||
let mut stmt = conn.prepare(queries::LIST_DEVICES)?;
|
||||
let payloads = stmt.query_map([], |row| row.get::<_, String>(0))?
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
payloads.into_iter().map(Self::from_json).collect()
|
||||
}
|
||||
|
||||
pub fn get_device(&self, id: &str) -> Result<Option<Device>> {
|
||||
let conn = self.lock()?;
|
||||
let payload: Option<String> = conn.query_row(queries::GET_DEVICE_BY_ID, [id], |row| row.get(0)).optional()?;
|
||||
payload.map(Self::from_json).transpose()
|
||||
}
|
||||
|
||||
pub fn get_device_by_mac(&self, mac: &str) -> Result<Option<Device>> {
|
||||
let conn = self.lock()?;
|
||||
let payload: Option<String> = conn.query_row(queries::GET_DEVICE_BY_MAC, [mac], |row| row.get(0)).optional()?;
|
||||
payload.map(Self::from_json).transpose()
|
||||
}
|
||||
|
||||
pub fn delete_device(&self, id: &str) -> Result<bool> {
|
||||
let mut conn = self.lock()?;
|
||||
let tx = conn.transaction()?;
|
||||
tx.execute(queries::DELETE_DEVICE_READINGS, [id])?;
|
||||
tx.execute(queries::DELETE_ZONES_BY_DEVICE_ID, [id])?;
|
||||
let changed = tx.execute(queries::DELETE_DEVICE, [id])? > 0;
|
||||
tx.commit()?;
|
||||
Ok(changed)
|
||||
}
|
||||
|
||||
pub fn save_zone(&self, zone: &Zone) -> Result<()> {
|
||||
let payload = Self::to_json(zone)?;
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::UPSERT_ZONE,
|
||||
params![zone.id, payload, zone.updated_at.to_rfc3339()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn list_zones(&self) -> Result<Vec<Zone>> {
|
||||
self.list_payloads(queries::LIST_ZONES)
|
||||
}
|
||||
|
||||
pub fn get_zone(&self, id: &str) -> Result<Option<Zone>> {
|
||||
self.get_payload(queries::GET_ZONE, id)
|
||||
}
|
||||
|
||||
pub fn delete_zone(&self, id: &str) -> Result<bool> {
|
||||
let mut conn = self.lock()?;
|
||||
let tx = conn.transaction()?;
|
||||
tx.execute(queries::DELETE_SCHEDULES_BY_ZONE_ID, [id])?;
|
||||
let changed = tx.execute(queries::DELETE_ZONE, [id])? > 0;
|
||||
tx.commit()?;
|
||||
Ok(changed)
|
||||
}
|
||||
|
||||
pub fn save_schedule(&self, schedule: &Schedule) -> Result<()> {
|
||||
let payload = Self::to_json(schedule)?;
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::UPSERT_SCHEDULE,
|
||||
params![schedule.id, schedule.zone_id, payload, schedule.updated_at.to_rfc3339()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn list_schedules(&self) -> Result<Vec<Schedule>> {
|
||||
self.list_payloads(queries::LIST_SCHEDULES)
|
||||
}
|
||||
|
||||
pub fn get_schedule(&self, id: &str) -> Result<Option<Schedule>> {
|
||||
self.get_payload(queries::GET_SCHEDULE, id)
|
||||
}
|
||||
|
||||
pub fn delete_schedule(&self, id: &str) -> Result<bool> {
|
||||
self.delete_by_id("schedules", id)
|
||||
}
|
||||
|
||||
pub fn save_automation(&self, item: &Automation) -> Result<()> {
|
||||
let payload = Self::to_json(item)?;
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::UPSERT_AUTOMATION,
|
||||
params![item.id, payload, item.updated_at.to_rfc3339()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn list_automations(&self) -> Result<Vec<Automation>> {
|
||||
self.list_payloads(queries::LIST_AUTOMATIONS)
|
||||
}
|
||||
|
||||
pub fn get_automation(&self, id: &str) -> Result<Option<Automation>> {
|
||||
self.get_payload(queries::GET_AUTOMATION, id)
|
||||
}
|
||||
|
||||
pub fn delete_automation(&self, id: &str) -> Result<bool> {
|
||||
self.delete_by_id("automations", id)
|
||||
}
|
||||
|
||||
fn list_payloads<T: DeserializeOwned>(&self, sql: &str) -> Result<Vec<T>> {
|
||||
let conn = self.lock()?;
|
||||
let mut stmt = conn.prepare(sql)?;
|
||||
let payloads = stmt.query_map([], |row| row.get::<_, String>(0))?
|
||||
.collect::<std::result::Result<Vec<_>, _>>()?;
|
||||
payloads.into_iter().map(Self::from_json).collect()
|
||||
}
|
||||
|
||||
fn get_payload<T: DeserializeOwned>(&self, sql: &str, id: &str) -> Result<Option<T>> {
|
||||
let conn = self.lock()?;
|
||||
let payload: Option<String> = conn.query_row(sql, [id], |row| row.get(0)).optional()?;
|
||||
payload.map(Self::from_json).transpose()
|
||||
}
|
||||
|
||||
fn delete_by_id(&self, table: &str, id: &str) -> Result<bool> {
|
||||
let sql = match table {
|
||||
"schedules" => queries::DELETE_SCHEDULE,
|
||||
"automations" => queries::DELETE_AUTOMATION,
|
||||
_ => anyhow::bail!("unsupported table"),
|
||||
};
|
||||
let conn = self.lock()?;
|
||||
Ok(conn.execute(sql, [id])? > 0)
|
||||
}
|
||||
|
||||
pub fn add_reading(&self, reading: &Reading) -> Result<i64> {
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::INSERT_READING,
|
||||
params![reading.device_id, reading.timestamp.to_rfc3339(), reading.indoor_temperature,
|
||||
reading.outdoor_temperature, reading.target_temperature, reading.power as i64, reading.source],
|
||||
)?;
|
||||
Ok(conn.last_insert_rowid())
|
||||
}
|
||||
|
||||
pub fn list_readings(&self, device_id: Option<&str>, since: DateTime<Utc>, limit: u32) -> Result<Vec<Reading>> {
|
||||
let conn = self.lock()?;
|
||||
let limit = limit.clamp(1, 5000) as i64;
|
||||
let mut rows_out = Vec::new();
|
||||
if let Some(device_id) = device_id {
|
||||
let mut stmt = conn.prepare(queries::LIST_READINGS_BY_DEVICE)?;
|
||||
let rows = stmt.query_map(params![device_id, since.to_rfc3339(), limit], Self::map_reading)?;
|
||||
for row in rows { rows_out.push(row?); }
|
||||
} else {
|
||||
let mut stmt = conn.prepare(queries::LIST_READINGS_ALL)?;
|
||||
let rows = stmt.query_map(params![since.to_rfc3339(), 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 {
|
||||
id: row.get(0)?,
|
||||
device_id: row.get(1)?,
|
||||
timestamp: DateTime::parse_from_rfc3339(×tamp)
|
||||
.map(|v| v.with_timezone(&Utc))
|
||||
.unwrap_or_else(|_| Utc::now()),
|
||||
indoor_temperature: row.get(3)?,
|
||||
outdoor_temperature: row.get(4)?,
|
||||
target_temperature: row.get(5)?,
|
||||
power: row.get::<_, i64>(6)? != 0,
|
||||
source: row.get(7)?,
|
||||
})
|
||||
}
|
||||
|
||||
pub fn prune_readings(&self, retention_days: i64) -> Result<u64> {
|
||||
let before = Utc::now() - Duration::days(retention_days.max(1));
|
||||
let conn = self.lock()?;
|
||||
Ok(conn.execute(queries::PRUNE_READINGS, [before.to_rfc3339()])? as u64)
|
||||
}
|
||||
|
||||
pub fn log_event(&self, level: &str, kind: &str, message: &str, metadata: &Value) -> Result<i64> {
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::INSERT_EVENT,
|
||||
params![Utc::now().to_rfc3339(), level, kind, message, serde_json::to_string(metadata)?],
|
||||
)?;
|
||||
Ok(conn.last_insert_rowid())
|
||||
}
|
||||
|
||||
pub fn list_events(&self, limit: u32) -> Result<Vec<EventLog>> {
|
||||
let conn = self.lock()?;
|
||||
let mut stmt = conn.prepare(queries::LIST_EVENTS)?;
|
||||
let rows = stmt.query_map([limit.clamp(1, 1000) as i64], |row| {
|
||||
let ts: String = row.get(1)?;
|
||||
let metadata: String = row.get(5)?;
|
||||
Ok(EventLog {
|
||||
id: row.get(0)?,
|
||||
timestamp: DateTime::parse_from_rfc3339(&ts).map(|v| v.with_timezone(&Utc)).unwrap_or_else(|_| Utc::now()),
|
||||
level: row.get(2)?,
|
||||
kind: row.get(3)?,
|
||||
message: row.get(4)?,
|
||||
metadata: serde_json::from_str(&metadata).unwrap_or(Value::Null),
|
||||
})
|
||||
})?;
|
||||
rows.collect::<std::result::Result<Vec<_>, _>>().map_err(Into::into)
|
||||
}
|
||||
|
||||
pub fn list_api_tokens(&self) -> Result<Vec<ApiTokenInfo>> {
|
||||
let conn = self.lock()?;
|
||||
let mut stmt = conn.prepare(queries::LIST_API_TOKENS)?;
|
||||
let rows = stmt.query_map([], |row| {
|
||||
let created_at: String = row.get(3)?;
|
||||
Ok(ApiTokenInfo {
|
||||
id: row.get(0)?,
|
||||
name: row.get(1)?,
|
||||
token_prefix: row.get(2)?,
|
||||
created_at: DateTime::parse_from_rfc3339(&created_at)
|
||||
.map(|value| value.with_timezone(&Utc))
|
||||
.unwrap_or_else(|_| Utc::now()),
|
||||
})
|
||||
})?;
|
||||
rows.collect::<std::result::Result<Vec<_>, _>>().map_err(Into::into)
|
||||
}
|
||||
|
||||
pub fn save_api_token(&self, token: &ApiTokenInfo, token_hash: &str) -> Result<()> {
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::INSERT_API_TOKEN,
|
||||
params![token.id, token.name, token_hash, token.token_prefix, token.created_at.to_rfc3339()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn api_token_exists(&self, token_hash: &str) -> Result<bool> {
|
||||
let conn = self.lock()?;
|
||||
let found: Option<i64> = conn.query_row(
|
||||
queries::API_TOKEN_EXISTS,
|
||||
[token_hash],
|
||||
|row| row.get(0),
|
||||
).optional()?;
|
||||
Ok(found.is_some())
|
||||
}
|
||||
|
||||
pub fn delete_api_token(&self, id: &str) -> Result<bool> {
|
||||
let conn = self.lock()?;
|
||||
Ok(conn.execute(queries::DELETE_API_TOKEN, [id])? > 0)
|
||||
}
|
||||
|
||||
pub fn load_runtime_settings(&self) -> Result<Option<RuntimeSettings>> {
|
||||
let conn = self.lock()?;
|
||||
let value: Option<String> = conn.query_row(queries::LOAD_RUNTIME_SETTINGS, [], |row| row.get(0)).optional()?;
|
||||
value.map(Self::from_json).transpose()
|
||||
}
|
||||
|
||||
pub fn save_runtime_settings(&self, settings: &RuntimeSettings) -> Result<()> {
|
||||
let value = Self::to_json(settings)?;
|
||||
let conn = self.lock()?;
|
||||
conn.execute(
|
||||
queries::UPSERT_RUNTIME_SETTINGS,
|
||||
params![value, Utc::now().to_rfc3339()],
|
||||
)?;
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use crate::models::{ApiTokenInfo, Device};
|
||||
|
||||
#[test]
|
||||
fn sqlite_round_trip() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let db = Db::open(&dir.path().join("test.db")).unwrap();
|
||||
let device = Device::simulated_default();
|
||||
db.save_device(&device).unwrap();
|
||||
let loaded = db.get_device(&device.id).unwrap().unwrap();
|
||||
assert_eq!(loaded.mac, device.mac);
|
||||
assert_eq!(db.list_devices().unwrap().len(), 1);
|
||||
db.log_event("info", "test", "ok", &serde_json::json!({"a":1})).unwrap();
|
||||
assert_eq!(db.list_events(10).unwrap().len(), 1);
|
||||
|
||||
let access_token = ApiTokenInfo {
|
||||
id: "token-1".into(),
|
||||
name: "Home Assistant".into(),
|
||||
token_prefix: "gree_controller_test...".into(),
|
||||
created_at: Utc::now(),
|
||||
};
|
||||
db.save_api_token(&access_token, "test-hash").unwrap();
|
||||
assert!(db.api_token_exists("test-hash").unwrap());
|
||||
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());
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user