use crate::UtcDateTime;
pub mod memory;
#[cfg(feature = "postgres")]
pub mod postgres;
#[cfg(feature = "redis")]
pub mod redis;
#[cfg(feature = "sqlite")]
pub mod sqlite;
mod job_index;
mod range;
#[cfg(feature = "postgres")]
pub use crate::storage::postgres::Postgres;
#[cfg(feature = "redis")]
pub use crate::storage::redis::Redis;
#[cfg(feature = "sqlite")]
pub use crate::storage::sqlite::Sqlite;
pub use job_index::{
JobIndexChange, JobIndexMetricsBatch, JobIndexPage, JobIndexRow, JobIndexWaitStats, QueueSample,
};
pub use range::{RangeItem, RangeOrder, RangePage, StorageOperation};
#[cfg(test)]
mod tests;
#[async_trait::async_trait]
pub trait Storage: Sync + Send {
async fn get(&self, key: &str) -> anyhow::Result<Option<Vec<u8>>>;
async fn apply(&self, operations: Vec<StorageOperation>) -> anyhow::Result<()>;
async fn range_page(
&self,
key: &str,
cursor: Option<i64>,
limit: usize,
order: RangeOrder,
) -> anyhow::Result<RangePage>;
async fn range_count(&self, key: &str) -> anyhow::Result<usize>;
async fn set(&self, key: &str, value: &[u8]) -> anyhow::Result<()> {
self.apply(vec![StorageOperation::Set {
key: key.to_string(),
value: value.to_vec(),
}])
.await
}
async fn del(&self, key: &str) -> anyhow::Result<()> {
self.apply(vec![StorageOperation::Delete {
key: key.to_string(),
}])
.await
}
async fn exist(&self, key: &str) -> anyhow::Result<bool> {
Ok(self.get(key).await?.is_some())
}
async fn expire(&self, key: &str, ttl_seconds: usize) -> anyhow::Result<()> {
self.apply(vec![StorageOperation::Expire {
key: key.to_string(),
ttl_seconds,
}])
.await
}
async fn set_if_absent(&self, key: &str, value: &[u8]) -> anyhow::Result<bool> {
if self.exist(key).await? {
return Ok(false);
}
self.set(key, value).await?;
Ok(true)
}
async fn total_db_size_bytes(&self) -> anyhow::Result<Option<u64>> {
Ok(None)
}
async fn sweep_expired(&self, _limit: usize) -> anyhow::Result<usize> {
Ok(0)
}
async fn checkpoint_wal(&self) -> anyhow::Result<()> {
Ok(())
}
async fn vacuum(&self) -> anyhow::Result<()> {
Ok(())
}
async fn job_index_upsert(&self, namespace: &str, row: JobIndexRow) -> anyhow::Result<()>;
async fn apply_with_job_index(
&self,
operations: Vec<StorageOperation>,
namespace: &str,
change: JobIndexChange,
) -> anyhow::Result<()> {
self.apply(operations).await?;
let result = match change {
JobIndexChange::Upsert(row) => self.job_index_upsert(namespace, row).await,
JobIndexChange::Remove(job_id) => self.job_index_remove(namespace, &job_id).await,
};
if let Err(error) = result {
tracing::warn!(%error, "Could not update dashboard job index");
}
Ok(())
}
async fn job_index_remove(&self, _namespace: &str, _job_id: &str) -> anyhow::Result<()> {
Ok(())
}
fn queue_backs_enqueued_stage(&self) -> bool {
false
}
async fn queue_sample_record(
&self,
_namespace: &str,
_at: UtcDateTime,
_queued: i64,
_partitioned: i64,
) -> anyhow::Result<()> {
Ok(())
}
async fn queue_samples_since(
&self,
_namespace: &str,
_since: UtcDateTime,
) -> anyhow::Result<Vec<QueueSample>> {
Ok(Vec::new())
}
async fn purge_unused_dashboard_rows(
&self,
_namespace: &str,
_limit: usize,
) -> anyhow::Result<usize> {
Ok(0)
}
async fn job_index_stage_counts(
&self,
namespace: &str,
) -> anyhow::Result<std::collections::HashMap<String, usize>>;
async fn job_index_list_by_stage(
&self,
namespace: &str,
stage: &str,
cursor: Option<i64>,
limit: usize,
) -> anyhow::Result<JobIndexPage>;
async fn job_index_list_by_partition(
&self,
namespace: &str,
topic: &str,
partition: u32,
cursor: Option<i64>,
limit: usize,
) -> anyhow::Result<JobIndexPage>;
async fn job_index_partition_neighbor(
&self,
namespace: &str,
topic: &str,
partition: u32,
sequence: i64,
older: bool,
) -> anyhow::Result<Option<JobIndexRow>>;
async fn job_index_list_continuations(
&self,
namespace: &str,
parent_job_id: &str,
limit: usize,
) -> anyhow::Result<Vec<JobIndexRow>>;
async fn job_index_truncate_partition(
&self,
namespace: &str,
topic: &str,
partition: u32,
keep_last: usize,
) -> anyhow::Result<usize>;
async fn job_index_record_transition(
&self,
namespace: &str,
stage: &str,
at: UtcDateTime,
) -> anyhow::Result<()>;
async fn job_index_record_metrics_batch(
&self,
namespace: &str,
batch: &JobIndexMetricsBatch,
) -> anyhow::Result<()>;
async fn job_index_reconcile_stage_counts(&self, _namespace: &str) -> anyhow::Result<()> {
Ok(())
}
async fn job_index_recent_transition_count(
&self,
namespace: &str,
stage: &str,
now: UtcDateTime,
) -> anyhow::Result<usize>;
async fn job_index_record_wait(
&self,
namespace: &str,
mode: &str,
wait_ms: i64,
at: UtcDateTime,
) -> anyhow::Result<()>;
async fn job_index_recent_wait_stats(
&self,
namespace: &str,
mode: &str,
now: UtcDateTime,
) -> anyhow::Result<JobIndexWaitStats>;
async fn job_index_sweep_expired(
&self,
namespace: &str,
now: UtcDateTime,
limit: usize,
) -> anyhow::Result<usize>;
}