#![deny(missing_docs)]
#[macro_use]
extern crate diesel;
#[macro_use]
extern crate diesel_migrations;
use tracing::{debug, error, info, warn};
use diesel::prelude::*;
use diesel::{insert_into, sql_query};
use metrics::{GaugeValue, Key, KeyName, SetRecorderError, SharedString, Unit};
use diesel_migrations::{EmbeddedMigrations, MigrationHarness};
use std::sync::Mutex;
use std::{
collections::{HashMap, VecDeque},
path::{Path, PathBuf},
sync::mpsc::{Receiver, RecvTimeoutError, SyncSender},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
use thiserror::Error;
const FLUSH_QUEUE_LIMIT: usize = 1000;
const BACKGROUND_CHANNEL_LIMIT: usize = 8000;
const SQLITE_DEFAULT_MAX_VARIABLES: usize = 999;
const METRIC_FIELDS_PER_ROW: usize = 3;
const INSERT_BATCH_SIZE: usize = SQLITE_DEFAULT_MAX_VARIABLES / METRIC_FIELDS_PER_ROW;
const QUEUE_HARD_LIMIT: usize = 100_000;
const RECONNECT_AFTER_FAILURES: u64 = 3;
const RECONNECT_BACKOFF: Duration = Duration::from_secs(30);
const ERROR_LOG_INTERVAL: Duration = Duration::from_secs(60);
struct LogThrottle {
interval: Duration,
last_logged: Option<Instant>,
suppressed: u64,
}
impl LogThrottle {
const fn new(interval: Duration) -> Self {
LogThrottle {
interval,
last_logged: None,
suppressed: 0,
}
}
fn allow(&mut self) -> Option<u64> {
let now = Instant::now();
let due = match self.last_logged {
Some(last) => now.duration_since(last) >= self.interval,
None => true,
};
if due {
self.last_logged = Some(now);
Some(std::mem::take(&mut self.suppressed))
} else {
self.suppressed += 1;
None
}
}
fn log_if_due(&mut self, emit: impl FnOnce(u64)) {
if let Some(suppressed) = self.allow() {
emit(suppressed);
}
}
}
#[derive(Debug, Error)]
pub enum MetricsError {
#[error("Database error: {0}")]
DbConnectionError(#[from] ConnectionError),
#[error("Migration error: {0}")]
MigrationError(Box<dyn std::error::Error + Send + Sync>),
#[error("Error querying DB: {0}")]
QueryError(#[from] diesel::result::Error),
#[error("Invalid database path")]
InvalidDatabasePath,
#[cfg(feature = "csv")]
#[error("IO Error: {0}")]
IoError(#[from] std::io::Error),
#[cfg(feature = "csv")]
#[error("CSV Error: {0}")]
CsvError(#[from] csv::Error),
#[error("Database has no metrics stored in it")]
EmptyDatabase,
#[error("Metric key {0} not found in database")]
KeyNotFound(String),
#[error("Exporter task has been stopped or crashed")]
ExporterUnavailable,
#[error("Session for signpost `{0}` has zero duration")]
ZeroLengthSession(String),
#[error("No metrics recorded for `{0}` in requested session")]
NoMetricsForKey(String),
}
impl MetricsError {
fn is_malformed_db(&self) -> bool {
self.to_string().contains("malformed")
}
}
pub type Result<T, E = MetricsError> = std::result::Result<T, E>;
mod maintenance;
mod metrics_db;
mod models;
mod recorder;
mod schema;
use crate::metrics_db::query;
pub use metrics_db::{MetricsDb, Session};
pub use models::{Metric, MetricKey, NewMetric};
pub(crate) const MIGRATIONS: EmbeddedMigrations = embed_migrations!();
#[derive(QueryableByName)]
struct PragmaCheckResult {
#[diesel(sql_type = diesel::sql_types::Text)]
#[diesel(column_name = quick_check)]
result: String,
}
fn remove_db_files(path: &Path) {
let db_path = PathBuf::from(path);
for suffix in &["", "-wal", "-shm"] {
let mut file_path = db_path.clone().into_os_string();
file_path.push(suffix);
let file_path = PathBuf::from(file_path);
if file_path.exists() {
if let Err(e) = std::fs::remove_file(&file_path) {
error!("Failed to remove {}: {}", file_path.display(), e);
} else {
info!("Removed corrupt database file: {}", file_path.display());
}
}
}
}
fn setup_db<P: AsRef<Path>>(path: P) -> Result<SqliteConnection> {
let url = path
.as_ref()
.to_str()
.ok_or(MetricsError::InvalidDatabasePath)?;
let mut db = SqliteConnection::establish(url)?;
sql_query("PRAGMA journal_mode=WAL;").execute(&mut db)?;
sql_query("PRAGMA busy_timeout = 5000;").execute(&mut db)?;
sql_query("PRAGMA journal_size_limit = 4194304;").execute(&mut db)?;
db.run_pending_migrations(MIGRATIONS)
.map_err(MetricsError::MigrationError)?;
let check: String = sql_query("PRAGMA quick_check;")
.get_result::<PragmaCheckResult>(&mut db)?
.result;
if check != "ok" {
return Err(MetricsError::QueryError(
diesel::result::Error::DatabaseError(
diesel::result::DatabaseErrorKind::Unknown,
Box::new(format!("database disk image is malformed: {check}")),
),
));
}
Ok(db)
}
fn setup_db_or_reset<P: AsRef<Path>>(path: P) -> Result<(SqliteConnection, bool)> {
let path = path.as_ref();
match setup_db(path) {
Ok(db) => Ok((db, false)),
Err(err) if err.is_malformed_db() => {
warn!(
"Database is malformed, removing and recreating: {}",
path.display()
);
remove_db_files(path);
setup_db(path).map(|db| (db, true))
}
Err(err) => Err(err),
}
}
enum RegisterType {
Counter,
Gauge,
Histogram,
}
enum Event {
Stop,
DescribeKey(RegisterType, KeyName, Option<Unit>, SharedString),
IncrementCounter(Duration, Key, u64),
AbsoluteCounter(Duration, Key, u64),
UpdateGauge(Duration, Key, GaugeValue),
UpdateHistogram(Duration, Key, f64),
SetHousekeeping {
retention_period: Option<Duration>,
housekeeping_period: Option<Duration>,
record_limit: Option<usize>,
},
RequestSummaryFromSignpost {
signpost_key: String,
keys: Vec<String>,
tx: tokio::sync::oneshot::Sender<Result<HashMap<String, f64>>>,
},
}
pub struct SqliteExporterHandle {
sender: SyncSender<Event>,
}
impl SqliteExporterHandle {
pub fn request_average_metrics(
&self,
from_signpost: &str,
with_keys: &[&str],
) -> Result<HashMap<String, f64>> {
let (tx, rx) = tokio::sync::oneshot::channel();
self.sender
.send(Event::RequestSummaryFromSignpost {
signpost_key: from_signpost.to_string(),
keys: with_keys.iter().map(|s| s.to_string()).collect(),
tx,
})
.map_err(|_| MetricsError::ExporterUnavailable)?;
match rx.blocking_recv() {
Ok(metrics) => Ok(metrics?),
Err(_) => Err(MetricsError::ExporterUnavailable),
}
}
}
pub struct SqliteExporter {
thread: Option<JoinHandle<()>>,
sender: SyncSender<Event>,
send_error_throttle: Mutex<LogThrottle>,
}
struct InnerState {
db: SqliteConnection,
db_path: PathBuf,
last_housekeeping: Instant,
housekeeping: Option<Duration>,
retention: Option<Duration>,
default_retention: Option<Duration>,
record_limit: Option<usize>,
inserted_since_housekeeping: usize,
maintenance: Option<maintenance::Maintenance>,
startup_cleanup_pending: bool,
last_maintenance_step: Instant,
last_vacuum_attempt: Option<Instant>,
flush_duration: Duration,
last_flush: Instant,
last_values: HashMap<Key, f64>,
counters: HashMap<Key, u64>,
key_ids: HashMap<String, i64>,
queue: VecDeque<NewMetric>,
consecutive_flush_failures: u64,
last_reconnect: Option<Instant>,
}
impl InnerState {
fn new(flush_duration: Duration, db: SqliteConnection, db_path: PathBuf) -> Self {
InnerState {
db,
db_path,
last_housekeeping: Instant::now(),
housekeeping: None,
retention: None,
default_retention: None,
record_limit: None,
inserted_since_housekeeping: 0,
maintenance: None,
startup_cleanup_pending: false,
last_maintenance_step: Instant::now(),
last_vacuum_attempt: None,
flush_duration,
last_flush: Instant::now(),
last_values: HashMap::new(),
counters: HashMap::new(),
key_ids: HashMap::new(),
queue: VecDeque::with_capacity(FLUSH_QUEUE_LIMIT),
consecutive_flush_failures: 0,
last_reconnect: None,
}
}
fn set_housekeeping(
&mut self,
retention: Option<Duration>,
housekeeping_duration: Option<Duration>,
record_limit: Option<usize>,
) {
self.retention = retention.or(self.default_retention);
self.housekeeping = housekeeping_duration;
self.last_housekeeping = Instant::now();
self.record_limit = record_limit;
if !self.startup_cleanup_pending {
self.maintenance = self.maintenance.as_ref().and_then(|_| {
housekeeping_duration
.map(|_| maintenance::Maintenance::new(self.retention, self.record_limit))
});
}
self.inserted_since_housekeeping = 0;
}
fn should_housekeep(&self) -> bool {
match self.housekeeping {
Some(duration) => {
let row_trigger = self.record_limit.map_or(100_000, |limit| {
(limit / 4).clamp(FLUSH_QUEUE_LIMIT, 100_000)
});
self.last_housekeeping.elapsed() >= duration
|| ((self.retention.is_some() || self.record_limit.is_some())
&& self.inserted_since_housekeeping >= row_trigger
&& self.last_housekeeping.elapsed() >= Duration::from_secs(1))
}
None => false,
}
}
fn housekeep(&mut self) -> Result<(), diesel::result::Error> {
let result = self.housekeep_step();
self.last_maintenance_step = Instant::now();
result
}
fn housekeep_step(&mut self) -> Result<(), diesel::result::Error> {
if self.maintenance.is_none() {
self.maintenance = Some(maintenance::Maintenance::new(
self.retention,
self.record_limit,
));
self.last_housekeeping = Instant::now();
self.inserted_since_housekeeping = 0;
}
if self.maintenance.as_mut().unwrap().step(&mut self.db)? {
self.maintenance = None;
self.startup_cleanup_pending = false;
let reclaim = maintenance::reclaim(&mut self.db, &mut self.last_vacuum_attempt);
let checkpoint = maintenance::checkpoint(&mut self.db);
reclaim?;
checkpoint?;
}
Ok(())
}
fn should_flush(&self) -> bool {
if self.last_flush.elapsed() > self.flush_duration {
true
} else if self.queue.len() >= FLUSH_QUEUE_LIMIT {
debug!("Flushing due to queue size ({} items)", self.queue.len());
true
} else {
false
}
}
fn flush(&mut self) -> Result<(), diesel::result::Error> {
if self.queue.is_empty() {
self.last_flush = Instant::now();
return Ok(());
}
let (front, back) = self.queue.as_slices();
match Self::insert_metrics(&mut self.db, [front, back]) {
Ok(()) => {
self.inserted_since_housekeeping = self
.inserted_since_housekeeping
.saturating_add(self.queue.len());
self.queue.clear();
self.last_flush = Instant::now();
self.consecutive_flush_failures = 0;
Ok(())
}
Err(e) => {
self.consecutive_flush_failures += 1;
self.enforce_queue_cap();
if Self::is_connection_fatal(&e)
|| self.consecutive_flush_failures >= RECONNECT_AFTER_FAILURES
{
self.reconnect();
}
Err(e)
}
}
}
fn insert_metrics<'a, S>(
db: &mut SqliteConnection,
slabs: S,
) -> Result<(), diesel::result::Error>
where
S: IntoIterator<Item = &'a [NewMetric]>,
{
use crate::schema::metrics::dsl::metrics;
db.transaction::<_, diesel::result::Error, _>(|db| {
let chunk_size = INSERT_BATCH_SIZE.max(1);
for slab in slabs {
for chunk in slab.chunks(chunk_size) {
insert_into(metrics).values(chunk).execute(db)?;
}
}
Ok(())
})
}
fn enforce_queue_cap(&mut self) {
if self.queue.len() > QUEUE_HARD_LIMIT {
let overflow = self.queue.len() - QUEUE_HARD_LIMIT;
self.queue.drain(..overflow);
warn!(
"metrics-sqlite queue exceeded {} items while flushing kept failing, dropped {} oldest metrics",
QUEUE_HARD_LIMIT, overflow
);
}
}
fn is_connection_fatal(e: &diesel::result::Error) -> bool {
matches!(e, diesel::result::Error::BrokenTransactionManager)
}
fn reconnect(&mut self) {
if let Some(last) = self.last_reconnect {
if last.elapsed() < RECONNECT_BACKOFF {
return;
}
}
self.last_reconnect = Some(Instant::now());
warn!(
"metrics-sqlite database connection is broken, reconnecting to {}",
self.db_path.display()
);
match setup_db_or_reset(&self.db_path) {
Ok((db, was_reset)) => {
self.db = db;
self.key_ids.clear();
self.consecutive_flush_failures = 0;
if was_reset {
let dropped = self.queue.len();
if dropped > 0 {
warn!(
"metrics-sqlite database was recreated; dropping {dropped} queued metrics with stale key ids"
);
self.queue.clear();
}
}
info!("metrics-sqlite database connection re-established");
}
Err(e) => {
error!("metrics-sqlite failed to reconnect to database: {:?}", e);
}
}
}
fn queue_metric(&mut self, timestamp: Duration, key: &str, value: f64) -> Result<()> {
let metric_key_id = match self.key_ids.get(key) {
Some(key) => *key,
None => {
debug!("Looking up {}", key);
let key_id = MetricKey::key_by_name(key, &mut self.db)?.id;
self.key_ids.insert(key.to_string(), key_id);
key_id
}
};
let metric = NewMetric {
timestamp: timestamp.as_secs_f64(),
metric_key_id,
value: value as _,
};
self.queue.push_back(metric);
Ok(())
}
pub fn metrics_summary_for_signpost_and_keys(
&mut self,
signpost: String,
metrics: Vec<String>,
) -> Result<HashMap<String, f64>> {
query::metrics_summary_for_signpost_and_keys(&mut self.db, &signpost, metrics)
}
}
fn run_worker(
db: SqliteConnection,
db_path: PathBuf,
receiver: Receiver<Event>,
flush_duration: Duration,
keep_duration: Option<Duration>,
) -> JoinHandle<()> {
thread::Builder::new()
.name("metrics-sqlite: worker".to_string())
.spawn(move || {
let mut state = InnerState::new(flush_duration, db, db_path);
state.default_retention = keep_duration;
state.retention = keep_duration;
if keep_duration.is_some() {
state.maintenance = Some(maintenance::Maintenance::new(keep_duration, None));
state.startup_cleanup_pending = true;
}
let mut flush_error_throttle = LogThrottle::new(ERROR_LOG_INTERVAL);
let mut queue_error_throttle = LogThrottle::new(ERROR_LOG_INTERVAL);
info!("SQLite worker started");
loop {
let time_based_flush = state.last_flush.elapsed() >= flush_duration;
let mut should_flush = false;
let mut should_exit = false;
let wait = if state.maintenance.is_some() {
flush_duration.min(maintenance::STEP_INTERVAL)
} else {
flush_duration
};
match receiver.recv_timeout(wait) {
Ok(Event::Stop) => {
info!("Stopping SQLiteExporter worker, flushing & exiting");
should_flush = true;
should_exit = true;
}
Ok(Event::SetHousekeeping {
retention_period,
housekeeping_period,
record_limit,
}) => {
state.set_housekeeping(retention_period, housekeeping_period, record_limit);
}
Ok(Event::DescribeKey(_key_type, key, unit, desc)) => {
info!("Describing key {:?}", key);
if let Err(e) = MetricKey::create_or_update(
key.as_str(),
unit,
Some(desc.as_ref()),
&mut state.db,
) {
error!("Failed to create key entry: {:?}", e);
}
}
Ok(Event::IncrementCounter(timestamp, key, value)) => {
let key_name = key.name();
let entry = state.counters.entry(key.clone()).or_insert(0);
let value = {
*entry += value;
*entry
};
if let Err(e) = state.queue_metric(timestamp, key_name, value as _) {
queue_error_throttle.log_if_due(|suppressed| {
if suppressed > 0 {
error!(
"Error queueing metric: {:?} ({} similar errors suppressed in the last {}s)",
e,
suppressed,
ERROR_LOG_INTERVAL.as_secs()
);
} else {
error!("Error queueing metric: {:?}", e);
}
});
}
should_flush = state.should_flush();
}
Ok(Event::AbsoluteCounter(timestamp, key, value)) => {
let key_name = key.name();
state.counters.insert(key.clone(), value);
if let Err(e) = state.queue_metric(timestamp, key_name, value as _) {
queue_error_throttle.log_if_due(|suppressed| {
if suppressed > 0 {
error!(
"Error queueing metric: {:?} ({} similar errors suppressed in the last {}s)",
e,
suppressed,
ERROR_LOG_INTERVAL.as_secs()
);
} else {
error!("Error queueing metric: {:?}", e);
}
});
}
should_flush = state.should_flush();
}
Ok(Event::UpdateGauge(timestamp, key, value)) => {
let key_name = key.name();
let entry = state.last_values.entry(key.clone()).or_insert(0.0);
let value = match value {
GaugeValue::Absolute(v) => {
*entry = v;
*entry
}
GaugeValue::Increment(v) => {
*entry += v;
*entry
}
GaugeValue::Decrement(v) => {
*entry -= v;
*entry
}
};
if let Err(e) = state.queue_metric(timestamp, key_name, value) {
queue_error_throttle.log_if_due(|suppressed| {
if suppressed > 0 {
error!(
"Error queueing metric: {:?} ({} similar errors suppressed in the last {}s)",
e,
suppressed,
ERROR_LOG_INTERVAL.as_secs()
);
} else {
error!("Error queueing metric: {:?}", e);
}
});
}
should_flush = state.should_flush();
}
Ok(Event::UpdateHistogram(timestamp, key, value)) => {
let key_name = key.name();
if let Err(e) = state.queue_metric(timestamp, key_name, value) {
queue_error_throttle.log_if_due(|suppressed| {
if suppressed > 0 {
error!(
"Error queueing metric: {:?} ({} similar errors suppressed in the last {}s)",
e,
suppressed,
ERROR_LOG_INTERVAL.as_secs()
);
} else {
error!("Error queueing metric: {:?}", e);
}
});
}
should_flush = state.should_flush();
}
Ok(Event::RequestSummaryFromSignpost {
signpost_key,
keys,
tx,
}) => {
match state.flush() {
Ok(()) => match state
.metrics_summary_for_signpost_and_keys(signpost_key, keys)
{
Ok(metrics) => {
if tx.send(Ok(metrics)).is_err() {
error!(
"Failed to respond with metrics results, discarding"
);
}
}
Err(e) => {
if let Err(e) = tx.send(Err(e)) {
error!(
"Failed to respond with metrics error result, discarding: {e:?}"
);
}
}
},
Err(e) => {
let err = MetricsError::from(e);
error!(
"Failed to flush pending metrics before summary request: {err:?}"
);
if let Err(send_err) = tx.send(Err(err)) {
error!(
"Failed to respond with metrics flush error result, discarding: {send_err:?}"
);
}
}
}
}
Err(RecvTimeoutError::Timeout) => {
should_flush = state.should_flush();
}
Err(RecvTimeoutError::Disconnected) => {
warn!("SQLiteExporter channel disconnected, exiting worker");
should_flush = true;
should_exit = true;
}
}
if time_based_flush || should_flush {
if time_based_flush {
debug!("Flushing due to elapsed time ({}s)", flush_duration.as_secs());
}
if let Err(e) = state.flush() {
if let Some(suppressed) = flush_error_throttle.allow() {
if suppressed > 0 {
error!(
"Error flushing metrics: {} ({} similar errors suppressed in the last {}s)",
e,
suppressed,
ERROR_LOG_INTERVAL.as_secs()
);
} else {
error!("Error flushing metrics: {}", e);
}
}
}
}
if should_exit {
let _ = maintenance::checkpoint(&mut state.db);
break;
}
if (state.maintenance.is_some() || state.should_housekeep())
&& state.last_maintenance_step.elapsed() >= maintenance::STEP_INTERVAL
{
if let Err(e) = state.housekeep() {
error!("Failed running house keeping: {:?}", e);
state.maintenance = None;
state.startup_cleanup_pending = false;
state.last_housekeeping = Instant::now();
state.inserted_since_housekeeping = 0;
}
}
}
})
.unwrap()
}
impl SqliteExporter {
pub fn new<P: AsRef<Path>>(
flush_interval: Duration,
keep_duration: Option<Duration>,
path: P,
) -> Result<Self> {
let path = path.as_ref().to_path_buf();
let (db, _was_reset) = setup_db_or_reset(&path)?;
let (sender, receiver) = std::sync::mpsc::sync_channel(BACKGROUND_CHANNEL_LIMIT);
let thread = run_worker(db, path, receiver, flush_interval, keep_duration);
let exporter = SqliteExporter {
thread: Some(thread),
sender,
send_error_throttle: Mutex::new(LogThrottle::new(ERROR_LOG_INTERVAL)),
};
Ok(exporter)
}
pub fn set_periodic_housekeeping(
&self,
periodic_duration: Option<Duration>,
retention: Option<Duration>,
record_limit: Option<usize>,
) {
if let Err(e) = self.sender.send(Event::SetHousekeeping {
retention_period: retention,
housekeeping_period: periodic_duration,
record_limit,
}) {
error!("Failed to set house keeping settings: {:?}", e);
}
}
fn log_send_failure(&self, context: &str, err: &dyn std::fmt::Debug) {
if let Ok(mut throttle) = self.send_error_throttle.lock() {
if let Some(suppressed) = throttle.allow() {
if suppressed > 0 {
error!(
"Error sending metric {} to SQLite worker: {:?} ({} similar errors suppressed in the last {}s)",
context,
err,
suppressed,
ERROR_LOG_INTERVAL.as_secs()
);
} else {
error!(
"Error sending metric {} to SQLite worker: {:?}",
context, err
);
}
}
}
}
pub fn install(self) -> Result<SqliteExporterHandle, SetRecorderError<Self>> {
let handle = SqliteExporterHandle {
sender: self.sender.clone(),
};
metrics::set_global_recorder(self)?;
Ok(handle)
}
}
impl Drop for SqliteExporter {
fn drop(&mut self) {
let _ = self.sender.send(Event::Stop);
let _ = self.thread.take().unwrap().join();
}
}
#[cfg(test)]
mod tests {
use crate::{
InnerState, LogThrottle, NewMetric, QUEUE_HARD_LIMIT, SqliteExporter, setup_db_or_reset,
};
use std::time::{Duration, Instant};
fn test_state() -> (InnerState, tempfile::TempDir) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("metrics.db");
let (db, _was_reset) = setup_db_or_reset(&path).unwrap();
(InnerState::new(Duration::from_secs(5), db, path), dir)
}
#[test]
fn configuring_housekeeping_preserves_pending_startup_cleanup() {
let (mut state, _dir) = test_state();
state.default_retention = Some(Duration::from_secs(60));
state.startup_cleanup_pending = true;
state.maintenance = Some(crate::maintenance::Maintenance::new(
state.default_retention,
None,
));
state.queue_metric(Duration::ZERO, "old", 1.0).unwrap();
state.flush().unwrap();
state.set_housekeeping(None, Some(Duration::from_secs(1800)), None);
assert!(state.maintenance.is_some());
while state.maintenance.is_some() {
state.housekeep().unwrap();
}
use diesel::prelude::*;
assert_eq!(
crate::schema::metrics::table
.count()
.get_result::<i64>(&mut state.db)
.unwrap(),
0
);
}
#[test]
fn disabling_periodic_cleanup_cancels_remaining_batches() {
let (mut state, _dir) = test_state();
for _ in 0..2500 {
state.queue_metric(Duration::ZERO, "old", 1.0).unwrap();
}
state.flush().unwrap();
state.set_housekeeping(
Some(Duration::from_secs(60)),
Some(Duration::from_secs(1)),
None,
);
state.housekeep().unwrap();
assert!(state.maintenance.is_some());
state.set_housekeeping(Some(Duration::from_secs(60)), None, None);
assert!(state.maintenance.is_none());
assert!(!state.should_housekeep());
use diesel::prelude::*;
assert_eq!(
crate::schema::metrics::table
.count()
.get_result::<i64>(&mut state.db)
.unwrap(),
1500
);
}
#[test]
fn disabling_periodic_cleanup_preserves_startup_policy() {
let (mut state, _dir) = test_state();
state.default_retention = Some(Duration::from_secs(60));
state.startup_cleanup_pending = true;
state.maintenance = Some(crate::maintenance::Maintenance::new(
state.default_retention,
None,
));
state.queue_metric(Duration::ZERO, "old", 1.0).unwrap();
state
.queue_metric(
std::time::SystemTime::UNIX_EPOCH.elapsed().unwrap(),
"recent",
1.0,
)
.unwrap();
state.flush().unwrap();
state.set_housekeeping(None, None, Some(0));
while state.maintenance.is_some() {
state.housekeep().unwrap();
}
assert!(!state.startup_cleanup_pending);
assert!(!state.should_housekeep());
use diesel::prelude::*;
assert_eq!(
crate::schema::metrics::table
.count()
.get_result::<i64>(&mut state.db)
.unwrap(),
1
);
}
#[test]
fn maintenance_interval_starts_after_sqlite_lock_wait() {
use diesel::{prelude::*, sql_query};
let (mut state, _dir) = test_state();
state.queue_metric(Duration::ZERO, "old", 1.0).unwrap();
state.flush().unwrap();
state.set_housekeeping(
Some(Duration::from_secs(60)),
Some(Duration::from_secs(1)),
None,
);
let mut blocker = crate::setup_db(&state.db_path).unwrap();
sql_query("BEGIN IMMEDIATE").execute(&mut blocker).unwrap();
let unlocker = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(100));
let released = Instant::now();
sql_query("COMMIT").execute(&mut blocker).unwrap();
released
});
state.housekeep().unwrap();
let released = unlocker.join().unwrap();
assert!(
state.last_maintenance_step >= released,
"the interval must start after waiting for SQLite, not before"
);
}
#[test]
fn failed_vacuum_is_throttled_and_does_not_skip_checkpoint() {
use diesel::{prelude::*, sql_query};
let (mut state, _dir) = test_state();
sql_query("CREATE TABLE ballast (data BLOB)")
.execute(&mut state.db)
.unwrap();
sql_query("WITH RECURSIVE n(x) AS (VALUES(1) UNION ALL SELECT x+1 FROM n WHERE x<1100) INSERT INTO ballast SELECT zeroblob(65536) FROM n")
.execute(&mut state.db).unwrap();
sql_query("DELETE FROM ballast")
.execute(&mut state.db)
.unwrap();
let wal = state.db_path.with_file_name("metrics.db-wal");
assert!(std::fs::metadata(&wal).unwrap().len() > 0);
sql_query("PRAGMA query_only = ON")
.execute(&mut state.db)
.unwrap();
assert!(state.housekeep().is_err());
let attempt = state
.last_vacuum_attempt
.expect("failed VACUUM records its attempt");
assert_eq!(
std::fs::metadata(&wal).unwrap().len(),
0,
"checkpoint must run even when reclamation fails"
);
state.housekeep().unwrap();
assert_eq!(state.last_vacuum_attempt, Some(attempt));
state.last_vacuum_attempt = Some(Instant::now() - crate::maintenance::VACUUM_INTERVAL);
assert!(
state.housekeep().is_err(),
"retry after the cooldown expires"
);
assert!(state.last_vacuum_attempt.unwrap() > attempt);
}
#[test]
fn housekeeping_inherits_retention_and_counts_only_successful_writes() {
let (mut state, _dir) = test_state();
let default = Duration::from_secs(60);
state.default_retention = Some(default);
state.set_housekeeping(None, Some(Duration::from_secs(1800)), Some(1000));
assert_eq!(state.retention, Some(default));
state.last_housekeeping = Instant::now() - Duration::from_secs(2);
for _ in 0..1000 {
state.queue_metric(Duration::ZERO, "test", 1.0).unwrap();
}
assert!(!state.should_housekeep());
state.flush().unwrap();
assert!(state.should_housekeep());
state.housekeep().unwrap();
while state.maintenance.is_some() {
state.housekeep().unwrap();
}
use diesel::prelude::*;
assert_eq!(
crate::schema::metrics::table
.count()
.get_result::<i64>(&mut state.db)
.unwrap(),
0
);
assert!(!state.should_housekeep());
state.set_housekeeping(Some(Duration::from_secs(120)), None, Some(1000));
assert_eq!(state.retention, Some(Duration::from_secs(120)));
state.inserted_since_housekeeping = 100_000;
assert!(
!state.should_housekeep(),
"disabled housekeeping ignores row trigger"
);
}
#[test]
fn enforce_queue_cap_drops_oldest_when_over_limit() {
let (mut state, _dir) = test_state();
let total = QUEUE_HARD_LIMIT + 250;
for i in 0..total {
state.queue.push_back(NewMetric {
timestamp: 0.0,
metric_key_id: 1,
value: i as f64,
});
}
state.enforce_queue_cap();
assert_eq!(state.queue.len(), QUEUE_HARD_LIMIT);
assert_eq!(state.queue.front().unwrap().value, 250.0);
assert_eq!(state.queue.back().unwrap().value, (total - 1) as f64);
state.enforce_queue_cap();
assert_eq!(state.queue.len(), QUEUE_HARD_LIMIT);
}
#[test]
fn failed_flush_leaves_queue_intact() {
use diesel::connection::SimpleConnection;
let (mut state, _dir) = test_state();
state
.db
.batch_execute("DROP TABLE metrics")
.expect("setup: drop metrics table");
for i in 0..5_000 {
state.queue.push_back(NewMetric {
timestamp: 0.0,
metric_key_id: 1,
value: i as f64,
});
}
let before = state.queue.len();
assert!(state.flush().is_err(), "flush should fail");
assert_eq!(state.queue.len(), before);
assert_eq!(state.queue.front().unwrap().value, 0.0);
assert_eq!(state.queue.back().unwrap().value, 4_999.0);
assert_eq!(state.consecutive_flush_failures, 1);
assert_eq!(state.inserted_since_housekeeping, 0);
}
#[test]
fn reconnect_drops_queue_when_database_is_recreated() {
use crate::setup_db;
use diesel::connection::SimpleConnection;
use std::io::{Seek, SeekFrom, Write};
let (mut state, _dir) = test_state();
for i in 0..100 {
state.queue.push_back(NewMetric {
timestamp: 0.0,
metric_key_id: 1,
value: i as f64,
});
}
state
.db
.batch_execute("PRAGMA wal_checkpoint(TRUNCATE)")
.expect("setup: wal checkpoint");
state.db = setup_db(":memory:").expect("setup: in-memory placeholder");
let mut f = std::fs::OpenOptions::new()
.write(true)
.open(&state.db_path)
.expect("setup: open db file");
f.seek(SeekFrom::Start(100))
.expect("setup: seek past header");
f.write_all(&[0xffu8; 16 * 1024])
.expect("setup: write garbage pages");
f.sync_all().expect("setup: sync garbage to disk");
drop(f);
state.reconnect();
assert!(
state.queue.is_empty(),
"queue should be cleared after reset"
);
assert_eq!(state.consecutive_flush_failures, 0);
assert!(state.flush().is_ok());
}
#[test]
fn reconnect_rebuilds_connection_with_backoff() {
let (mut state, _dir) = test_state();
state.consecutive_flush_failures = 5;
state.reconnect();
assert_eq!(state.consecutive_flush_failures, 0);
let first_attempt = state.last_reconnect.expect("reconnect should run");
assert!(state.flush().is_ok());
state.consecutive_flush_failures = 5;
state.reconnect();
assert_eq!(state.consecutive_flush_failures, 5);
assert_eq!(state.last_reconnect, Some(first_attempt));
}
#[test]
fn log_throttle_suppresses_and_reports_count() {
let mut throttle = LogThrottle::new(Duration::from_millis(50));
assert_eq!(throttle.allow(), Some(0));
for _ in 0..7 {
assert_eq!(throttle.allow(), None);
}
std::thread::sleep(Duration::from_millis(60));
assert_eq!(throttle.allow(), Some(7));
std::thread::sleep(Duration::from_millis(60));
assert_eq!(throttle.allow(), Some(0));
}
#[test]
fn test_threading() {
use std::thread;
let dir = tempfile::tempdir().unwrap();
let exporter = std::sync::Arc::new(
SqliteExporter::new(
Duration::from_millis(500),
None,
dir.path().join("metrics.db"),
)
.unwrap(),
);
let joins: Vec<thread::JoinHandle<()>> = (0..5)
.map(|_| {
let exporter = exporter.clone();
thread::spawn(move || {
metrics::with_local_recorder(exporter.as_ref(), || {
let start = Instant::now();
loop {
metrics::gauge!("rate").set(1.0);
metrics::counter!("hits").increment(1);
metrics::histogram!("histogram").record(5.0);
if start.elapsed().as_secs() >= 5 {
break;
}
}
});
})
})
.collect();
for j in joins {
j.join().unwrap();
}
}
}