use super::{TagType, TransactionStatus};
use crate::buffer::inixtime;
use crate::buffer::TimeStatus;
use crate::types::Transaction;
use crate::types::{Metric, Statistics, TXMetric};
use crate::Error;
use sqlx::{Connection, SqliteConnection, SqlitePool};
use tracing::instrument;
use tracing::{debug, error, info, warn};
#[derive(Debug)]
pub(crate) struct MetaTag {
pub key: String,
pub tag: String,
pub value: String,
}
#[cfg(test)]
#[instrument(level = "info", err, skip(pool))]
pub async fn get_name(pool: &SqlitePool, key: &str) -> Result<String, Error> {
let select_sensor = sqlx::query_as("SELECT name FROM sensor WHERE sensor.name = ?");
let mut tx = pool.begin().await?;
let res: (String,) = select_sensor.bind(key).fetch_one(&mut *tx).await?;
tx.commit().await?;
Ok(res.0)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn sensor_id(pool: &SqlitePool, key: &str) -> Result<i64, Error> {
let sensor_by_name =
sqlx::query_as("SELECT s_id FROM sensor WHERE sensor.name = $1 ORDER BY s_id DESC");
let mut tx = pool.begin().await?;
let row: (i64,) = sensor_by_name.bind(key).fetch_one(&mut *tx).await?;
tx.commit().await?;
Ok(row.0)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn count_transactions(pool: &SqlitePool) -> Result<i32, Error> {
let trans_query = sqlx::query_scalar!("SELECT COUNT(*) FROM changes");
let mut tx = pool.begin().await?;
debug!("Counting transactions in database");
let transactions = trans_query.fetch_one(&mut *tx).await?;
tx.commit().await?;
Ok(transactions)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn get_statistics(pool: &SqlitePool) -> Result<Statistics, Error> {
let removed_query =
sqlx::query_scalar!(r#"SELECT count(*) FROM logdata WHERE logdata.status = 'REMOVED';"#);
let timefail_query =
sqlx::query_scalar!(r#"SELECT count(*) FROM logdata WHERE logdata.status = 'TIMEFAIL';"#);
let internal_query = sqlx::query_scalar!(
"SELECT count(*)"
+ " FROM logdata"
+ " JOIN sensor USING(s_id)"
+ " WHERE logdata.status = 'NONE' AND sensor.name LIKE 'modio.%';"
);
let count_query = sqlx::query_scalar!(
"SELECT count(*)"
+ " FROM logdata"
+ " JOIN sensor USING(s_id)"
+ " WHERE logdata.status = 'NONE' AND sensor.name NOT LIKE 'modio.%';"
);
let trans_query = sqlx::query_scalar!(r#"SELECT COUNT(*) FROM changes;"#);
let mut tx = pool.begin().await?;
debug!("get_statistics");
let metrics = count_query
.fetch_one(&mut *tx)
.await
.inspect_err(|e| error!("Failed to count logdata values: err={:?}", e))
.unwrap_or(-1);
let internal = internal_query
.fetch_one(&mut *tx)
.await
.inspect_err(|e| error!("Failed to count internal logdata values: err={:?}", e))
.unwrap_or(-1);
let removed = removed_query
.fetch_one(&mut *tx)
.await
.inspect_err(|e| error!("Failed to count removed logdata values: err={:?}", e))
.unwrap_or(-1);
let timefail = timefail_query
.fetch_one(&mut *tx)
.await
.inspect_err(|e| error!("Failed to count timefail logdata values: err={:?}", e))
.unwrap_or(-1);
let transactions = trans_query
.fetch_one(&mut *tx)
.await
.inspect_err(|e| error!("Failed to count transactions: err={:?}", e))?;
tx.commit().await?;
Ok(Statistics {
metrics,
internal,
removed,
timefail,
transactions,
buffered: 0,
})
}
#[instrument(level = "info", err, skip(pool))]
pub async fn add_sensor(pool: &SqlitePool, key: &str) -> Result<(), Error> {
let mut tx = pool.begin().await?;
debug!("Inserting key into names table.");
sqlx::query!("INSERT OR IGNORE INTO sensor(name) VALUES (?)", key)
.execute(&mut *tx)
.await?;
tx.commit().await?;
Ok(())
}
#[instrument(level = "debug", err, skip(pool, delete_buffer))]
pub async fn persist_data(
pool: &SqlitePool,
delete_buffer: &Vec<(Metric, TimeStatus)>,
) -> Result<(), Error> {
debug!("persist_data_raw, len={}", delete_buffer.len());
let mut pool_tx = pool.begin().await?;
for (metric, tstatus) in delete_buffer {
let status = tstatus.as_str();
let add_logdata = sqlx::query!(
"INSERT INTO logdata (s_id, value, time, status)"
+ " SELECT s_id, $2, $3, $4"
+ " FROM sensor"
+ " WHERE sensor.name = $1",
metric.name,
metric.value,
metric.time,
status,
);
add_logdata
.execute(&mut *pool_tx)
.await
.inspect_err(|e| error!("persist_data_raw. Failed to execute. err={:?}", e))?;
}
pool_tx
.commit()
.await
.inspect_err(|e| error!("persist_data_raw. Failed to commit. err={:?}", e))?;
Ok(())
}
#[instrument(level = "debug", err, skip(pool))]
pub async fn get_batch(pool: &SqlitePool, size: u32) -> Result<Vec<TXMetric>, Error> {
let get_external_batch = sqlx::query_as!(
TXMetric,
"SELECT id, name, value, CAST(time as INTEGER) as time"
+ " FROM logdata"
+ " JOIN sensor USING(s_id)"
+ " WHERE logdata.status = 'NONE' AND sensor.name NOT LIKE 'modio.%'"
+ " ORDER BY ID ASC LIMIT $1",
size
);
let mut tx = pool.begin().await?;
debug!("Retrieving batch of non-modio keys.");
let res = get_external_batch.fetch_all(&mut *tx).await?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "debug", err, skip(pool))]
pub async fn get_internal_batch(pool: &SqlitePool, size: u32) -> Result<Vec<TXMetric>, Error> {
let get_internal_batch = sqlx::query_as!(
TXMetric,
"SELECT id, name, value, CAST(time as INTEGER) as time"
+ " FROM logdata"
+ " JOIN sensor USING(s_id)"
+ " WHERE logdata.status = 'NONE' AND sensor.name LIKE 'modio.%'"
+ " ORDER BY ID ASC LIMIT $1",
size
);
let mut tx = pool.begin().await?;
debug!("Retrieving batch of modio-only keys.");
let res = get_internal_batch
.fetch_all(&mut *tx)
.await
.inspect_err(|e| error!(err=?e, "Failed to fetch internal data batch"))?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "debug", err, skip_all)]
pub async fn drop_batch(pool: &SqlitePool, ids: &[i64]) -> Result<(), Error> {
debug!("Beginning transaction to drop a batch of data.");
let mut tx = pool.begin().await?;
for id in ids {
let query = sqlx::query!(
"UPDATE logdata SET status = 'REMOVED' WHERE id = $1 AND status = 'NONE'",
id
);
query.execute(&mut *tx).await?;
}
debug!("Committing transaction to drop a batch of data.");
tx.commit().await?;
Ok(())
}
#[instrument(level = "info", err, skip(pool))]
pub(crate) async fn fix_timefail(pool: &SqlitePool, adjust: f32) -> Result<u64, Error> {
let fix_timefail = sqlx::query!(
"UPDATE logdata SET time = time + $1, status = 'NONE' WHERE status = 'TIMEFAIL'",
adjust
);
let count = {
let mut tx = pool.begin().await?;
debug!("Fixing timefail for stored data");
let res = fix_timefail.execute(&mut *tx).await?;
tx.commit().await?;
res.rows_affected()
};
Ok(count)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn get_last_datapoint(pool: &SqlitePool, key: &str) -> Result<Metric, Error> {
let last_value = sqlx::query_as!(
Metric,
r#"SELECT name as "name!", value as "value!", time as "time!: f64" FROM logdata JOIN sensor USING(s_id) WHERE sensor.name = $1 ORDER BY time DESC LIMIT 1"#,
key
);
let mut tx = pool.begin().await?;
debug!("Fetching last datapoint from db");
let res = last_value
.fetch_one(&mut *tx)
.await
.inspect_err(|e| error!(err=?e, "Failed to fetch datapoint from db"))?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "info", err, skip_all)]
pub async fn get_latest_logdata(pool: &SqlitePool) -> Result<Vec<Metric>, Error> {
let get_last_all_points = sqlx::query_as!(
Metric,
r#"SELECT name as "name!", value as "value!", MAX(time) as "time!: f64" FROM logdata JOIN sensor USING(s_id) GROUP BY name ORDER BY time ASC"#,
);
let mut tx = pool.begin().await?;
debug!("Fetching all latest datapoints from db");
let res = get_last_all_points.fetch_all(&mut *tx).await?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "debug", err, skip(pool))]
pub(crate) async fn has_transaction(pool: &SqlitePool, token: &str) -> Result<bool, Error> {
let count_query = sqlx::query_scalar!("SELECT count(*) FROM changes WHERE token = $1", token);
let mut tx = pool.begin().await?;
debug!("Checking for transactions for token");
let res = count_query.fetch_one(&mut *tx).await?;
tx.commit().await?;
Ok(res > 0)
}
#[instrument(level = "info", err, skip_all, fields(key))]
pub(crate) async fn transaction_add(
pool: &SqlitePool,
key: &str,
expected: &str,
target: &str,
token: &str,
) -> Result<(), Error> {
let mut tx = pool.begin().await?;
debug!("Writing new transaction to changes table");
sqlx::query!(
"INSERT INTO changes(s_id, token, expected, target)"
+ " SELECT s_id, $2, $3, $4"
+ " FROM sensor"
+ " WHERE sensor.name = $1",
key,
token,
expected,
target,
)
.execute(&mut *tx)
.await?;
tx.commit().await?;
Ok(())
}
#[instrument(level = "info", err, skip(pool))]
pub async fn transaction_get(pool: &SqlitePool, prefix: &str) -> Result<Vec<Transaction>, Error> {
let transaction_pending = sqlx::query!(
"UPDATE changes SET status = 'PENDING'"
+ " FROM changes J1"
+ " JOIN sensor USING(s_id)"
+ " WHERE changes.status = 'NONE' AND sensor.name LIKE $1 || '%'",
prefix
);
let transaction_get = sqlx::query_as!(
Transaction,
"SELECT t_id, name, expected, target, status"
+ " FROM changes"
+ " JOIN sensor USING(s_id)"
+ " WHERE changes.status = 'PENDING' AND sensor.name LIKE $1 || '%'",
prefix
);
let mut tx = pool.begin().await?;
debug!("Marking NEW transactions as pending for prefix={}", prefix);
transaction_pending.execute(&mut *tx).await?;
debug!("Fetching all transactions for prefix");
let res = transaction_get.fetch_all(&mut *tx).await?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn transaction_result_key(
pool: &SqlitePool,
transaction_id: i64,
) -> Result<String, Error> {
let transaction_result_key = sqlx::query!(
"SELECT 'mytemp.internal.change.'||sensor.name as name"
+ " FROM sensor"
+ " JOIN changes USING(s_id)"
+ " WHERE changes.t_id=$1",
transaction_id
);
let change_key = {
let mut tx = pool.begin().await?;
debug!("Fetching transactions for id");
let change_key = transaction_result_key.fetch_one(&mut *tx).await?;
tx.commit().await?;
change_key.name
};
Ok(change_key)
}
#[instrument(level = "info", err, skip(pool))]
async fn transaction_token(pool: &SqlitePool, transaction_id: i64) -> Result<String, Error> {
let query = sqlx::query!(
"SELECT token FROM changes WHERE changes.t_id = $1",
transaction_id
);
let mut tx = pool.begin().await?;
let res = query.fetch_one(&mut *tx).await?;
tx.commit().await?;
Ok(res.token)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn transaction_mark(
pool: &SqlitePool,
transaction_id: i64,
success: TransactionStatus,
timefail: TimeStatus,
) -> Result<u64, Error> {
let transaction_result_str = success.as_str();
let time_status = timefail.as_str();
let time_when = inixtime();
let change_key = transaction_result_key(pool, transaction_id).await?;
add_sensor(pool, &change_key).await?;
let token = transaction_token(pool, transaction_id).await?;
let logval = serde_json::json!({
"clock": time_when,
"id": token,
"status": success.as_str(),
})
.to_string();
let count = {
let mut tx = pool.begin().await?;
let mark_transaction = sqlx::query!(
"UPDATE changes"
+ " SET status = $2"
+ " WHERE changes.t_id = $1 AND changes.status = 'PENDING'",
transaction_id,
transaction_result_str,
);
debug!("Marking transaction from pending");
let res = mark_transaction.execute(&mut *tx).await?;
let count = res.rows_affected();
let add_logdata = sqlx::query!(
"INSERT INTO logdata (s_id, value, time, status)"
+ " SELECT s_id, $2, $3, $4"
+ " FROM sensor"
+ " WHERE sensor.name = $1",
change_key,
logval,
time_when,
time_status,
);
debug!("Adding logdata for reporting");
add_logdata.execute(&mut *tx).await?;
tx.commit().await?;
count
};
Ok(count)
}
#[instrument(level = "info", err, skip(pool))]
pub async fn transaction_fail_pending(pool: &SqlitePool) -> Result<u64, Error> {
let transaction_fail_pending =
sqlx::query!("UPDATE changes SET status = 'FAILED' WHERE status = 'PENDING'");
info!("Failing all pending transactions");
let count = {
let mut tx = pool.begin().await?;
let res = transaction_fail_pending.execute(&mut *tx).await?;
tx.commit().await?;
res.rows_affected()
};
Ok(count)
}
#[instrument(level = "info", err, skip(conn))]
pub async fn delete_old_transactions(conn: &mut SqliteConnection) -> Result<i32, Error> {
let transaction_delete_done = sqlx::query!(
"DELETE FROM changes"
+ " WHERE status in ('FAILED', 'SUCCESS')"
+ " and t_id NOT IN ("
+ "SELECT MAX(t_id) as t_id FROM changes GROUP BY s_id ORDER BY t_id ASC"
+ ")",
);
let count = {
let mut tx = conn.begin().await?;
debug!("delete_old_transactions");
let res = transaction_delete_done.execute(&mut *tx).await?;
tx.commit().await?;
res.rows_affected() as i32
};
Ok(count)
}
#[instrument(level = "info", err, skip(conn))]
pub async fn fail_queued_transactions(conn: &mut SqliteConnection) -> Result<i32, Error> {
let fail_queued_transactions = sqlx::query!(
"\
WITH victims AS ( \
SELECT t_id, s_id, \
ROW_NUMBER() OVER (\
PARTITION BY s_id ORDER BY t_id DESC) AS row \
FROM changes WHERE changes.status='NONE') \
UPDATE changes SET status='FAILED' \
WHERE t_id IN (SELECT t_id FROM victims WHERE victims.row >= 17);"
);
let count = {
let mut tx = conn.begin().await?;
debug!("Failing queued transactions in db");
let res = fail_queued_transactions.execute(&mut *tx).await?;
tx.commit().await?;
res.rows_affected() as i32
};
if count > 0 {
warn!(
"Many unhandled transactions for single keys found, marked {} as failed.",
count
);
};
Ok(count)
}
#[instrument(level = "info", err, skip(conn))]
pub async fn delete_old_logdata(conn: &mut SqliteConnection) -> Result<i32, Error> {
let logdata_delete_done = sqlx::query!(
"\
DELETE FROM logdata \
WHERE status = 'REMOVED' AND id NOT IN (\
SELECT id FROM (\
SELECT id, s_id, MAX(time) as time \
FROM logdata GROUP BY s_id ORDER BY time ASC));",
);
let count = {
let mut tx = conn.begin().await?;
debug!("Deleting old logdata");
let res = logdata_delete_done.execute(&mut *tx).await?;
tx.commit().await?;
res.rows_affected() as i32
};
Ok(count)
}
#[instrument(level = "info", err, skip(conn))]
pub async fn delete_random_data(conn: &mut SqliteConnection) -> Result<i32, Error> {
let count1 = delete_old_logdata(&mut *conn).await?;
let delete_random = sqlx::query!(
"DELETE FROM logdata WHERE ((random() / 9223372036854775808.0 + 1) / 2) < 0.25",
);
let count2 = {
let mut tx = conn.begin().await?;
warn!("Deleting random data");
let res = delete_random.execute(&mut *tx).await?;
tx.commit().await?;
res.rows_affected() as i32
};
let result = count1 + count2;
info!(
"Deleted items, random={}, old={}, total={}.",
count2, count1, result
);
debug!("delete_random_data, vacuuming");
sqlx::query!("VACUUM;").execute(&mut *conn).await?;
Ok(result)
}
#[instrument(level = "info", err, skip(conn))]
pub async fn need_vacuum_or_shrink(conn: &mut SqliteConnection) -> Result<usize, Error> {
debug!("need_vacuum_or_shrink");
const VACUUM_MIN_RATIO: f32 = 0.8;
const VACUUM_MIN_SIZE: f32 = 512.0 * 1024.0;
const MAX_SIZE: f32 = 24.0 * 1024.0 * 1024.0;
let res: (i32,) = sqlx::query_as("pragma page_count")
.fetch_one(&mut *conn)
.await
.inspect_err(|e| error!("Failed pragma page_count: err={:?}", e))?;
#[allow(clippy::cast_precision_loss)]
let pages = res.0 as f32;
let res: (i32,) = sqlx::query_as("pragma freelist_count")
.fetch_one(&mut *conn)
.await
.inspect_err(|e| error!("Failed pragma freelist_count: err={:?}", e))?;
#[allow(clippy::cast_precision_loss)]
let freepages = res.0 as f32;
let res: (i32,) = sqlx::query_as("pragma page_size")
.fetch_one(&mut *conn)
.await
.inspect_err(|e| error!("Failed pragma page_size: err={:?}", e))?;
#[allow(clippy::cast_precision_loss)]
let pagesize: f32 = res.0 as f32;
if freepages * pagesize > VACUUM_MIN_SIZE && freepages > pages * VACUUM_MIN_RATIO {
info!("Vacuuming due to size");
sqlx::query("VACUUM;")
.execute(&mut *conn)
.await
.inspect_err(|e| error!("Failed VACUUM, err={:?}", e))?;
}
if freepages == 0.0 && pages * pagesize >= MAX_SIZE {
info!("DB out of size, removing random data");
delete_random_data(&mut *conn).await?;
}
Ok(0)
}
#[instrument(level = "debug", skip(pool))]
pub async fn get_metadata(pool: &SqlitePool, key: &str) -> Result<Vec<(String, String)>, Error> {
let query = sqlx::query!(
"SELECT sensor.name as key, tag, value"
+ " FROM tag_single"
+ " JOIN sensor USING(s_id)"
+ " WHERE sensor.name = $1",
key
);
let mut tx = pool.begin().await?;
debug!("Fetching all metadata for key={}", key);
let res = query.map(|v| (v.tag, v.value)).fetch_all(&mut *tx).await?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "info", skip(pool))]
pub async fn get_all_metadata_internal(
pool: &SqlitePool,
offset: u32,
) -> Result<Vec<MetaTag>, Error> {
let query = sqlx::query_as!(
MetaTag,
"SELECT sensor.name as key, tag, value"
+ " FROM tag_single"
+ " JOIN sensor USING(s_id)"
+ " WHERE sensor.name LIKE 'modio.%'"
+ " ORDER BY tag_single.tag_id"
+ " LIMIT 100 OFFSET $1",
offset,
);
let mut tx = pool.begin().await?;
debug!("fetching internal metadata batch, offet={offset}");
let pairs = query
.fetch_all(&mut *tx)
.await
.inspect_err(|e| error!(err=?e, "Failed to query db for a batch of internal metadata"))?;
tx.commit().await?;
Ok(pairs)
}
#[instrument(level = "info", skip(pool))]
pub async fn get_all_metadata_customer(
pool: &SqlitePool,
offset: u32,
) -> Result<Vec<MetaTag>, Error> {
let query = sqlx::query_as!(
MetaTag,
"SELECT sensor.name as key, tag, value"
+ " FROM tag_single"
+ " JOIN sensor USING(s_id)"
+ " WHERE sensor.name NOT LIKE 'modio.%'"
+ " ORDER BY tag_single.tag_id"
+ " LIMIT 100 OFFSET $1",
offset
);
let mut tx = pool.begin().await?;
debug!("fetching customer metadata batch, offet={offset}");
let pairs = query
.fetch_all(&mut *tx)
.await
.inspect_err(|e| error!(err=?e, "Failed to query db for a batch of customer metadata"))?;
tx.commit().await?;
Ok(pairs)
}
#[instrument(level = "info", skip(pool))]
pub async fn metadata_replace_tag(
pool: &SqlitePool,
key: &str,
tag: TagType,
value: &str,
) -> Result<(), Error> {
let tag = tag.as_str();
let query = sqlx::query!(
"\
INSERT OR REPLACE INTO tag_single (s_id, tag, value) \
SELECT s_id, $2, $3 \
FROM sensor \
WHERE sensor.name = $1",
key,
tag,
value
);
let mut tx = pool.begin().await?;
debug!(key, tag, value, "Replacing tag in db");
query.execute(&mut *tx).await?;
tx.commit().await?;
Ok(())
}
#[instrument(level = "debug", skip(pool))]
pub async fn metadata_set_tag(
pool: &SqlitePool,
key: &str,
tag: TagType,
value: &str,
) -> Result<(), Error> {
let tag = tag.as_str();
let query = sqlx::query!(
"\
INSERT INTO tag_single (s_id, tag, value) \
SELECT s_id, $2, $3 \
FROM sensor \
WHERE sensor.name = $1",
key,
tag,
value
);
let mut tx = pool.begin().await?;
info!(
"tags: Setting tag='{}' for key='{}' value='{}'",
tag, key, value
);
query.execute(&mut *tx).await?;
tx.commit().await?;
Ok(())
}
#[instrument(level = "debug", skip(pool))]
pub async fn metadata_get_tag(
pool: &SqlitePool,
tag: TagType,
) -> Result<Vec<(String, String)>, Error> {
let tag = tag.as_str();
let query = sqlx::query!(
"\
SELECT sensor.name as key, tag_single.value as value \
FROM tag_single \
JOIN sensor USING(s_id) \
WHERE tag_single.tag = $1 \
",
tag
);
let mut tx = pool.begin().await?;
info!("tags: Retrieving all tags of type: {}", tag);
let res = query
.map(|val| (val.key, val.value))
.fetch_all(&mut *tx)
.await?;
tx.commit().await?;
Ok(res)
}
#[instrument(level = "debug", skip(pool))]
pub async fn metadata_get_single_tag(
pool: &SqlitePool,
key: &str,
tag: TagType,
) -> Result<String, Error> {
let tag = tag.as_str();
let query = sqlx::query_scalar!(
"\
SELECT tag_single.value as value \
FROM tag_single \
JOIN sensor USING(s_id) \
WHERE sensor.name = $1 AND tag_single.tag = $2 \
",
key,
tag
);
let res = {
let mut tx = pool.begin().await?;
debug!("tags: retrieving single tag of type");
let res = query.fetch_one(&mut *tx).await?;
tx.commit().await?;
res
};
Ok(res)
}