use super::{
JobIndexPage, JobIndexRow, JobIndexWaitStats, RangeItem, RangeOrder, RangePage, Storage,
StorageOperation,
};
use crate::UtcDateTime;
use sqlx::{
sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions},
Row, SqlitePool,
};
use std::{str::FromStr, time::Duration};
const METRIC_BUCKET_MS: i64 = 60_000;
const QUEUE_SAMPLE_RETENTION_MINUTES: i64 = 7 * 24 * 60;
fn minute_bucket(at: UtcDateTime) -> i64 {
at.timestamp_millis().div_euclid(METRIC_BUCKET_MS)
}
fn row_to_job_index(row: &sqlx::sqlite::SqliteRow) -> Result<JobIndexRow, sqlx::Error> {
use chrono::TimeZone;
let stage_date_ms: i64 = row.try_get("stage_date_ms")?;
let created_at_ms: i64 = row.try_get("created_at_ms")?;
let date_expire_ms: Option<i64> = row.try_get("date_expire")?;
Ok(JobIndexRow {
job_id: row.try_get("job_id")?,
payload_type: row.try_get("payload_type")?,
stage: row.try_get("stage")?,
stage_date: chrono::Utc.timestamp_millis_opt(stage_date_ms).unwrap(),
revision: row.try_get("revision")?,
created_at: chrono::Utc.timestamp_millis_opt(created_at_ms).unwrap(),
wait_ms: row.try_get("wait_ms")?,
wait_mode: row.try_get("wait_mode")?,
topic: row.try_get("topic")?,
partition: row.try_get("partition_id")?,
sequence: row.try_get("sequence")?,
parent_job_id: row.try_get("parent_job_id")?,
date_expire: date_expire_ms.map(|ms| chrono::Utc.timestamp_millis_opt(ms).unwrap()),
})
}
#[derive(Clone)]
pub struct Sqlite {
pool: SqlitePool,
wal_check_not_before: std::sync::Arc<std::sync::atomic::AtomicI64>,
}
#[derive(Debug, PartialEq, Eq)]
pub(crate) enum WalCheckpoint {
BelowThreshold,
Truncated,
GaveUp,
}
const WAL_CHECKPOINT_FRAMES: i64 = 32_768;
const WAL_CHECK_INTERVAL_MS: i64 = 5_000;
const WAL_TRUNCATE_WAIT_MS: i64 = 15_000;
const WAL_RETRY_AFTER_BUSY_MS: i64 = 15_000;
const LIVE_STAGES: &str = "'delayed', 'waiting', 'enqueued', 'running', 'requeued'";
const COUNTED_LIVE_STAGES: &str = "'delayed', 'waiting', 'running', 'requeued'";
const FINISHED_STAGES: &str = "'success', 'failed'";
const MIGRATIONS_TABLE: &str = "later_schema_migrations";
async fn run_migrations(pool: &SqlitePool) -> anyhow::Result<()> {
sqlx::query(&format!(
"CREATE TABLE IF NOT EXISTS {MIGRATIONS_TABLE} ( \
version BIGINT PRIMARY KEY, \
description TEXT NOT NULL, \
checksum BLOB NOT NULL, \
applied_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP \
)"
))
.execute(pool)
.await?;
let migrator = sqlx::migrate!("./migrations/sqlite");
for migration in migrator.iter() {
if migration.migration_type.is_down_migration() {
continue;
}
let existing: Option<Vec<u8>> = sqlx::query_scalar(&format!(
"SELECT checksum FROM {MIGRATIONS_TABLE} WHERE version = ?1"
))
.bind(migration.version)
.fetch_optional(pool)
.await?;
match existing {
Some(checksum) if checksum == migration.checksum.as_ref() => continue,
Some(_) => {
anyhow::bail!(
"Later migration {} ({}) has changed since it was applied - \
migrations must never be edited after release",
migration.version,
migration.description,
);
}
None => {}
}
let mut tx = pool.begin().await?;
sqlx::raw_sql(&migration.sql).execute(&mut *tx).await?;
sqlx::query(&format!(
"INSERT INTO {MIGRATIONS_TABLE} (version, description, checksum) VALUES (?1, ?2, ?3)"
))
.bind(migration.version)
.bind(migration.description.as_ref())
.bind(migration.checksum.as_ref())
.execute(&mut *tx)
.await?;
tx.commit().await?;
}
Ok(())
}
impl Sqlite {
pub async fn new(connection_string: &str) -> anyhow::Result<Self> {
let options = SqliteConnectOptions::from_str(connection_string)?
.create_if_missing(true)
.foreign_keys(true)
.busy_timeout(Duration::from_secs(5))
.journal_mode(SqliteJournalMode::Wal);
let max_connections = if connection_string.contains(":memory:") {
1
} else {
5
};
let pool = SqlitePoolOptions::new()
.max_connections(max_connections)
.connect_with(options)
.await?;
run_migrations(&pool).await?;
Ok(Self::from_parts(pool))
}
pub async fn from_pool(pool: SqlitePool) -> anyhow::Result<Self> {
run_migrations(&pool).await?;
Ok(Self::from_parts(pool))
}
async fn list_queued_jobs(
&self,
queue_name: &str,
cursor: Option<i64>,
limit: usize,
) -> anyhow::Result<JobIndexPage> {
use chrono::TimeZone;
let rows = sqlx::query(
"SELECT sequence, payload, available_at FROM later_delivery_queue \
WHERE sequence < ?1 AND queue_name = ?2 AND lease_until IS NULL \
ORDER BY sequence DESC LIMIT ?3",
)
.bind(cursor.unwrap_or(i64::MAX))
.bind(queue_name)
.bind(i64::try_from(limit)?)
.fetch_all(&self.pool)
.await?;
let mut items = Vec::with_capacity(rows.len());
for row in &rows {
let payload: Vec<u8> = row.try_get("payload")?;
let Ok(crate::models::AmqpCommand::ExecuteJob(job)) =
crate::encoder::decode::<crate::models::AmqpCommand>(&payload)
else {
continue;
};
let queued_at = chrono::Utc
.timestamp_millis_opt(row.try_get("available_at")?)
.single()
.unwrap_or_else(chrono::Utc::now);
items.push(JobIndexRow {
job_id: job.id.to_string(),
payload_type: job.payload_type,
stage: "enqueued".to_string(),
stage_date: queued_at,
revision: 0,
created_at: queued_at,
wait_ms: None,
wait_mode: None,
topic: None,
partition: None,
sequence: None,
parent_job_id: None,
date_expire: None,
});
}
let next_cursor = if rows.len() == limit {
match rows.last() {
Some(row) => Some(row.try_get::<i64, _>("sequence")?),
None => None,
}
} else {
None
};
Ok(JobIndexPage { items, next_cursor })
}
fn from_parts(pool: SqlitePool) -> Self {
Self {
pool,
wal_check_not_before: std::sync::Arc::new(std::sync::atomic::AtomicI64::new(0)),
}
}
pub(crate) async fn checkpoint_wal_above(
&self,
frames: i64,
truncate_wait_ms: i64,
) -> anyhow::Result<WalCheckpoint> {
let mut connection = self.pool.acquire().await?;
let (_, log_frames, _): (i64, i64, i64) = sqlx::query_as("PRAGMA wal_checkpoint(PASSIVE)")
.fetch_one(&mut *connection)
.await?;
if log_frames < frames {
return Ok(WalCheckpoint::BelowThreshold);
}
let original_wait: i64 = sqlx::query_scalar("PRAGMA busy_timeout")
.fetch_one(&mut *connection)
.await?;
sqlx::query(&format!("PRAGMA busy_timeout = {truncate_wait_ms}"))
.execute(&mut *connection)
.await?;
let truncated: Result<(i64, i64, i64), _> =
sqlx::query_as("PRAGMA wal_checkpoint(TRUNCATE)")
.fetch_one(&mut *connection)
.await;
sqlx::query(&format!("PRAGMA busy_timeout = {original_wait}"))
.execute(&mut *connection)
.await?;
let (busy, ..) = truncated?;
Ok(if busy == 0 {
WalCheckpoint::Truncated
} else {
WalCheckpoint::GaveUp
})
}
pub fn pool(&self) -> &SqlitePool {
&self.pool
}
fn now_millis() -> i64 {
chrono::Utc::now().timestamp_millis()
}
}
async fn upsert_job_index<'e, E>(
executor: E,
namespace: &str,
row: &JobIndexRow,
) -> anyhow::Result<()>
where
E: sqlx::Executor<'e, Database = sqlx::Sqlite>,
{
sqlx::query(
"INSERT INTO later_jobs_index \
(namespace, job_id, payload_type, stage, stage_date_ms, revision, \
created_at_ms, wait_ms, wait_mode, topic, partition_id, sequence, \
parent_job_id, date_expire) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14) \
ON CONFLICT(namespace, job_id) DO UPDATE SET \
payload_type = excluded.payload_type, \
stage = excluded.stage, \
stage_date_ms = excluded.stage_date_ms, \
revision = excluded.revision, \
wait_ms = COALESCE(excluded.wait_ms, later_jobs_index.wait_ms), \
wait_mode = COALESCE(excluded.wait_mode, later_jobs_index.wait_mode), \
topic = excluded.topic, \
partition_id = excluded.partition_id, \
sequence = excluded.sequence, \
parent_job_id = COALESCE(excluded.parent_job_id, later_jobs_index.parent_job_id), \
date_expire = excluded.date_expire \
WHERE excluded.revision >= later_jobs_index.revision",
)
.bind(namespace)
.bind(&row.job_id)
.bind(&row.payload_type)
.bind(&row.stage)
.bind(row.stage_date.timestamp_millis())
.bind(row.revision)
.bind(row.created_at.timestamp_millis())
.bind(row.wait_ms)
.bind(&row.wait_mode)
.bind(&row.topic)
.bind(row.partition)
.bind(row.sequence)
.bind(&row.parent_job_id)
.bind(row.date_expire.map(|d| d.timestamp_millis()))
.execute(executor)
.await?;
Ok(())
}
#[async_trait::async_trait]
impl Storage for Sqlite {
async fn get(&self, key: &str) -> anyhow::Result<Option<Vec<u8>>> {
sqlx::query(
"SELECT value FROM later_storage \
WHERE key = ?1 AND (date_expire IS NULL OR date_expire > ?2)",
)
.bind(key)
.bind(Self::now_millis())
.fetch_optional(&self.pool)
.await?
.map(|row| row.try_get("value"))
.transpose()
.map_err(anyhow::Error::from)
}
async fn apply(&self, operations: Vec<StorageOperation>) -> anyhow::Result<()> {
let mut transaction = self.pool.begin().await?;
crate::backend::apply_sqlite_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(())
}
async fn set_if_absent(&self, key: &str, value: &[u8]) -> anyhow::Result<bool> {
use sqlx::Row;
let now = Self::now_millis();
let row = sqlx::query(
"INSERT INTO later_storage (key, value, counter_value, date_created) \
VALUES (?1, ?2, NULL, ?3) \
ON CONFLICT(key) DO UPDATE SET value = excluded.value, counter_value = NULL, \
date_updated = excluded.date_created, date_expire = NULL \
WHERE later_storage.date_expire IS NOT NULL AND later_storage.date_expire <= ?4 \
RETURNING 1 AS created",
)
.bind(key)
.bind(value)
.bind(now)
.bind(now)
.fetch_optional(&self.pool)
.await?;
Ok(row.is_some_and(|row| row.get::<i64, _>("created") == 1))
}
async fn range_page(
&self,
key: &str,
cursor: Option<i64>,
limit: usize,
order: RangeOrder,
) -> anyhow::Result<RangePage> {
let limit_i64 = i64::try_from(limit)?;
let rows = match (order, cursor) {
(RangeOrder::OldestFirst, Some(cursor)) => {
sqlx::query(
"SELECT sequence, value FROM later_storage_range \
WHERE range_key = ?1 AND sequence > ?2 \
AND (date_expire IS NULL OR date_expire > ?3) \
ORDER BY sequence ASC LIMIT ?4",
)
.bind(key)
.bind(cursor)
.bind(Self::now_millis())
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
(RangeOrder::OldestFirst, None) => {
sqlx::query(
"SELECT sequence, value FROM later_storage_range \
WHERE range_key = ?1 AND (date_expire IS NULL OR date_expire > ?2) \
ORDER BY sequence ASC LIMIT ?3",
)
.bind(key)
.bind(Self::now_millis())
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
(RangeOrder::NewestFirst, Some(cursor)) => {
sqlx::query(
"SELECT sequence, value FROM later_storage_range \
WHERE range_key = ?1 AND sequence < ?2 \
AND (date_expire IS NULL OR date_expire > ?3) \
ORDER BY sequence DESC LIMIT ?4",
)
.bind(key)
.bind(cursor)
.bind(Self::now_millis())
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
(RangeOrder::NewestFirst, None) => {
sqlx::query(
"SELECT sequence, value FROM later_storage_range \
WHERE range_key = ?1 AND (date_expire IS NULL OR date_expire > ?2) \
ORDER BY sequence DESC LIMIT ?3",
)
.bind(key)
.bind(Self::now_millis())
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
};
let items = rows
.into_iter()
.map(|row| {
Ok(RangeItem {
cursor: row.try_get("sequence")?,
value: row.try_get("value")?,
})
})
.collect::<Result<Vec<_>, sqlx::Error>>()?;
let next_cursor = (items.len() == limit)
.then(|| items.last().map(|item| item.cursor))
.flatten();
Ok(RangePage { items, next_cursor })
}
async fn range_count(&self, key: &str) -> anyhow::Result<usize> {
let now = Self::now_millis();
let count = sqlx::query(
"SELECT COUNT(*) AS item_count FROM later_storage_range \
WHERE range_key = ?1 AND (date_expire IS NULL OR date_expire > ?2)",
)
.bind(key)
.bind(now)
.fetch_one(&self.pool)
.await?
.try_get::<i64, _>("item_count")?;
Ok(usize::try_from(count)?)
}
async fn total_db_size_bytes(&self) -> anyhow::Result<Option<u64>> {
let size: i64 = sqlx::query_scalar(
"SELECT (SELECT page_count FROM pragma_page_count()) \
* (SELECT page_size FROM pragma_page_size())",
)
.fetch_one(&self.pool)
.await?;
Ok(Some(u64::try_from(size)?))
}
async fn sweep_expired(&self, limit: usize) -> anyhow::Result<usize> {
let now = Self::now_millis();
let limit = i64::try_from(limit).unwrap_or(i64::MAX);
let storage_removed = sqlx::query(
"DELETE FROM later_storage WHERE key IN ( \
SELECT key FROM later_storage \
WHERE date_expire IS NOT NULL AND date_expire <= ?1 LIMIT ?2 \
)",
)
.bind(now)
.bind(limit)
.execute(&self.pool)
.await?
.rows_affected();
let remaining = limit.saturating_sub(i64::try_from(storage_removed).unwrap_or(i64::MAX));
let range_removed = if remaining > 0 {
sqlx::query(
"DELETE FROM later_storage_range WHERE sequence IN ( \
SELECT sequence FROM later_storage_range \
WHERE date_expire IS NOT NULL AND date_expire <= ?1 LIMIT ?2 \
)",
)
.bind(now)
.bind(remaining)
.execute(&self.pool)
.await?
.rows_affected()
} else {
0
};
Ok(usize::try_from(storage_removed + range_removed)?)
}
async fn checkpoint_wal(&self) -> anyhow::Result<()> {
use std::sync::atomic::Ordering;
let now = Self::now_millis();
if now < self.wal_check_not_before.load(Ordering::Relaxed) {
return Ok(());
}
self.wal_check_not_before
.store(now + WAL_CHECK_INTERVAL_MS, Ordering::Relaxed);
match self
.checkpoint_wal_above(WAL_CHECKPOINT_FRAMES, WAL_TRUNCATE_WAIT_MS)
.await?
{
WalCheckpoint::BelowThreshold => {}
WalCheckpoint::Truncated => tracing::info!("Truncated the SQLite write-ahead log"),
WalCheckpoint::GaveUp => {
tracing::debug!(
"A reader kept the SQLite write-ahead log from truncating; backing off"
);
self.wal_check_not_before
.store(now + WAL_RETRY_AFTER_BUSY_MS, Ordering::Relaxed);
}
}
Ok(())
}
async fn vacuum(&self) -> anyhow::Result<()> {
sqlx::query("VACUUM").execute(&self.pool).await?;
Ok(())
}
async fn job_index_upsert(&self, namespace: &str, row: JobIndexRow) -> anyhow::Result<()> {
upsert_job_index(&self.pool, namespace, &row).await
}
async fn apply_with_job_index(
&self,
operations: Vec<StorageOperation>,
namespace: &str,
change: crate::storage::JobIndexChange,
) -> anyhow::Result<()> {
let mut transaction = self.pool.begin().await?;
crate::backend::apply_sqlite_operations(&mut transaction, operations).await?;
match change {
crate::storage::JobIndexChange::Upsert(row) => {
upsert_job_index(&mut *transaction, namespace, &row).await?
}
crate::storage::JobIndexChange::Remove(job_id) => {
sqlx::query("DELETE FROM later_jobs_index WHERE namespace = ?1 AND job_id = ?2")
.bind(namespace)
.bind(job_id)
.execute(&mut *transaction)
.await?;
}
}
transaction.commit().await?;
Ok(())
}
async fn job_index_remove(&self, namespace: &str, job_id: &str) -> anyhow::Result<()> {
sqlx::query("DELETE FROM later_jobs_index WHERE namespace = ?1 AND job_id = ?2")
.bind(namespace)
.bind(job_id)
.execute(&self.pool)
.await?;
Ok(())
}
async fn queue_sample_record(
&self,
namespace: &str,
at: UtcDateTime,
queued: i64,
partitioned: i64,
) -> anyhow::Result<()> {
let bucket = minute_bucket(at);
let mut tx = self.pool.begin().await?;
sqlx::query(
"INSERT INTO later_queue_samples (namespace, minute_bucket, queued, partitioned) \
VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT (namespace, minute_bucket) \
DO UPDATE SET queued = excluded.queued, partitioned = excluded.partitioned",
)
.bind(namespace)
.bind(bucket)
.bind(queued)
.bind(partitioned)
.execute(&mut *tx)
.await?;
sqlx::query("DELETE FROM later_queue_samples WHERE namespace = ?1 AND minute_bucket < ?2")
.bind(namespace)
.bind(bucket - QUEUE_SAMPLE_RETENTION_MINUTES)
.execute(&mut *tx)
.await?;
tx.commit().await?;
Ok(())
}
async fn queue_samples_since(
&self,
namespace: &str,
since: UtcDateTime,
) -> anyhow::Result<Vec<crate::storage::QueueSample>> {
use chrono::TimeZone;
let rows = sqlx::query(
"SELECT minute_bucket, queued, partitioned FROM later_queue_samples \
WHERE namespace = ?1 AND minute_bucket >= ?2 ORDER BY minute_bucket",
)
.bind(namespace)
.bind(minute_bucket(since))
.fetch_all(&self.pool)
.await?;
rows.iter()
.map(|row| {
let bucket: i64 = row.try_get("minute_bucket")?;
Ok(crate::storage::QueueSample {
at: chrono::Utc
.timestamp_millis_opt(bucket * METRIC_BUCKET_MS)
.single()
.ok_or_else(|| anyhow::anyhow!("queue sample time out of range"))?,
queued: row.try_get("queued")?,
partitioned: row.try_get("partitioned")?,
})
})
.collect()
}
fn queue_backs_enqueued_stage(&self) -> bool {
true
}
async fn purge_unused_dashboard_rows(
&self,
namespace: &str,
limit: usize,
) -> anyhow::Result<usize> {
let mut budget = i64::try_from(limit)?;
let mut removed = 0usize;
let exact = [
"all-jobs",
"stage-enqueued-jobs",
"stage-running-jobs",
"stage-waiting-jobs",
];
for suffix in exact {
if budget == 0 {
break;
}
let done = sqlx::query(
"DELETE FROM later_storage_range WHERE sequence IN ( \
SELECT sequence FROM later_storage_range WHERE range_key = ?1 LIMIT ?2)",
)
.bind(format!("{namespace}-{suffix}"))
.bind(budget)
.execute(&self.pool)
.await?
.rows_affected();
budget -= i64::try_from(done)?;
removed += usize::try_from(done)?;
}
for prefix in ["stats-job-in-", "stats-partition-jobs-"] {
if budget == 0 {
break;
}
let from = format!("{namespace}-{prefix}");
let mut to = from.clone();
to.pop();
to.push('.');
let done = sqlx::query(
"DELETE FROM later_storage_range WHERE sequence IN ( \
SELECT sequence FROM later_storage_range \
WHERE range_key >= ?1 AND range_key < ?2 LIMIT ?3)",
)
.bind(from)
.bind(to)
.bind(budget)
.execute(&self.pool)
.await?
.rows_affected();
budget -= i64::try_from(done)?;
removed += usize::try_from(done)?;
}
if budget > 0 {
let done = sqlx::query(
"DELETE FROM later_jobs_index WHERE rowid IN ( \
SELECT rowid FROM later_jobs_index \
WHERE namespace = ?1 AND stage = 'enqueued' AND topic IS NULL LIMIT ?2)",
)
.bind(namespace)
.bind(budget)
.execute(&self.pool)
.await?
.rows_affected();
removed += usize::try_from(done)?;
}
Ok(removed)
}
async fn job_index_stage_counts(
&self,
namespace: &str,
) -> anyhow::Result<std::collections::HashMap<String, usize>> {
let now = Self::now_millis();
let rows = sqlx::query(&format!(
"SELECT stage, count AS item_count FROM later_stage_counts \
WHERE namespace = ?1 AND count > 0 AND stage IN ({COUNTED_LIVE_STAGES}) \
UNION ALL \
SELECT 'enqueued', SUM(ready) FROM later_queue_depth \
WHERE queue_name = ?1 AND ready > 0 HAVING SUM(ready) > 0 \
UNION ALL \
SELECT stage, COUNT(*) AS item_count FROM later_jobs_index \
WHERE namespace = ?1 AND stage IN ({FINISHED_STAGES}) \
AND (date_expire IS NULL OR date_expire > ?2) \
GROUP BY stage"
))
.bind(namespace)
.bind(now)
.fetch_all(&self.pool)
.await?;
rows.into_iter()
.map(|row| {
let stage: String = row.try_get(0)?;
let count: i64 = row.try_get(1)?;
Ok((stage, usize::try_from(count)?))
})
.collect()
}
async fn job_index_reconcile_stage_counts(&self, namespace: &str) -> anyhow::Result<()> {
let mut tx = self.pool.begin().await?;
sqlx::query("DELETE FROM later_stage_counts WHERE namespace = ?1")
.bind(namespace)
.execute(&mut *tx)
.await?;
sqlx::query(&format!(
"INSERT INTO later_stage_counts (namespace, stage, count) \
SELECT namespace, stage, COUNT(*) FROM later_jobs_index \
WHERE namespace = ?1 AND stage IN ({LIVE_STAGES}) GROUP BY stage"
))
.bind(namespace)
.execute(&mut *tx)
.await?;
sqlx::query("DELETE FROM later_queue_depth WHERE queue_name = ?1")
.bind(namespace)
.execute(&mut *tx)
.await?;
sqlx::query(
"INSERT INTO later_queue_depth (namespace, queue_name, ready) \
SELECT namespace, queue_name, COUNT(*) FROM later_delivery_queue \
WHERE queue_name = ?1 AND lease_until IS NULL GROUP BY namespace, queue_name",
)
.bind(namespace)
.execute(&mut *tx)
.await?;
tx.commit().await?;
Ok(())
}
async fn job_index_list_by_stage(
&self,
namespace: &str,
stage: &str,
cursor: Option<i64>,
limit: usize,
) -> anyhow::Result<JobIndexPage> {
if stage == "enqueued" {
return self.list_queued_jobs(namespace, cursor, limit).await;
}
let now = Self::now_millis();
let limit_i64 = i64::try_from(limit)?;
let rows = match cursor {
Some(cursor) => {
sqlx::query(
"SELECT * FROM later_jobs_index \
WHERE namespace = ?1 AND stage = ?2 AND stage_date_ms < ?3 \
AND (date_expire IS NULL OR date_expire > ?4) \
ORDER BY stage_date_ms DESC, job_id DESC LIMIT ?5",
)
.bind(namespace)
.bind(stage)
.bind(cursor)
.bind(now)
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
None => {
sqlx::query(
"SELECT * FROM later_jobs_index \
WHERE namespace = ?1 AND stage = ?2 \
AND (date_expire IS NULL OR date_expire > ?3) \
ORDER BY stage_date_ms DESC, job_id DESC LIMIT ?4",
)
.bind(namespace)
.bind(stage)
.bind(now)
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
};
let items = rows
.iter()
.map(row_to_job_index)
.collect::<Result<Vec<_>, _>>()?;
let next_cursor = (items.len() == limit)
.then(|| items.last().map(|item| item.stage_date.timestamp_millis()))
.flatten();
Ok(JobIndexPage { items, next_cursor })
}
async fn job_index_list_by_partition(
&self,
namespace: &str,
topic: &str,
partition: u32,
cursor: Option<i64>,
limit: usize,
) -> anyhow::Result<JobIndexPage> {
let limit_i64 = i64::try_from(limit)?;
let partition_id = i64::from(partition);
let rows = match cursor {
Some(cursor) => {
sqlx::query(
"SELECT * FROM later_jobs_index \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 \
AND sequence < ?4 \
ORDER BY sequence DESC LIMIT ?5",
)
.bind(namespace)
.bind(topic)
.bind(partition_id)
.bind(cursor)
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
None => {
sqlx::query(
"SELECT * FROM later_jobs_index \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 \
ORDER BY sequence DESC LIMIT ?4",
)
.bind(namespace)
.bind(topic)
.bind(partition_id)
.bind(limit_i64)
.fetch_all(&self.pool)
.await?
}
};
let items = rows
.iter()
.map(row_to_job_index)
.collect::<Result<Vec<_>, _>>()?;
let next_cursor = (items.len() == limit)
.then(|| items.last().and_then(|item| item.sequence))
.flatten();
Ok(JobIndexPage { items, next_cursor })
}
async fn job_index_partition_neighbor(
&self,
namespace: &str,
topic: &str,
partition: u32,
sequence: i64,
older: bool,
) -> anyhow::Result<Option<JobIndexRow>> {
let target = if older { sequence - 1 } else { sequence + 1 };
let row = sqlx::query(
"SELECT * FROM later_jobs_index \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 AND sequence = ?4",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition))
.bind(target)
.fetch_optional(&self.pool)
.await?;
Ok(row.as_ref().map(row_to_job_index).transpose()?)
}
async fn job_index_list_continuations(
&self,
namespace: &str,
parent_job_id: &str,
limit: usize,
) -> anyhow::Result<Vec<JobIndexRow>> {
let rows = sqlx::query(
"SELECT * FROM later_jobs_index \
WHERE namespace = ?1 AND parent_job_id = ?2 \
ORDER BY stage_date_ms ASC LIMIT ?3",
)
.bind(namespace)
.bind(parent_job_id)
.bind(i64::try_from(limit)?)
.fetch_all(&self.pool)
.await?;
rows.iter()
.map(row_to_job_index)
.collect::<Result<_, _>>()
.map_err(anyhow::Error::from)
}
async fn job_index_truncate_partition(
&self,
namespace: &str,
topic: &str,
partition: u32,
keep_last: usize,
) -> anyhow::Result<usize> {
let partition_id = i64::from(partition);
let keep_last = i64::try_from(keep_last)?;
let result = sqlx::query(
"DELETE FROM later_jobs_index WHERE rowid IN ( \
SELECT rowid FROM later_jobs_index \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 \
ORDER BY sequence DESC \
LIMIT -1 OFFSET ?4 \
)",
)
.bind(namespace)
.bind(topic)
.bind(partition_id)
.bind(keep_last)
.execute(&self.pool)
.await?;
Ok(usize::try_from(result.rows_affected())?)
}
async fn job_index_record_transition(
&self,
namespace: &str,
stage: &str,
at: UtcDateTime,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_job_transition_metrics (namespace, stage, minute_bucket, count) \
VALUES (?1, ?2, ?3, 1) \
ON CONFLICT(namespace, stage, minute_bucket) \
DO UPDATE SET count = count + 1",
)
.bind(namespace)
.bind(stage)
.bind(minute_bucket(at))
.execute(&self.pool)
.await?;
Ok(())
}
async fn job_index_record_metrics_batch(
&self,
namespace: &str,
batch: &crate::storage::JobIndexMetricsBatch,
) -> anyhow::Result<()> {
let mut tx = self.pool.begin().await?;
for (stage, at, count) in &batch.transitions {
sqlx::query(
"INSERT INTO later_job_transition_metrics (namespace, stage, minute_bucket, count) \
VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(namespace, stage, minute_bucket) \
DO UPDATE SET count = count + ?4",
)
.bind(namespace)
.bind(stage)
.bind(minute_bucket(*at))
.bind(count)
.execute(&mut *tx)
.await?;
}
for (mode, at, count, sum_ms, max_ms) in &batch.waits {
sqlx::query(
"INSERT INTO later_job_wait_metrics \
(namespace, mode, minute_bucket, count, sum_ms, max_ms) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6) \
ON CONFLICT(namespace, mode, minute_bucket) DO UPDATE SET \
count = count + ?4, \
sum_ms = sum_ms + ?5, \
max_ms = MAX(max_ms, ?6)",
)
.bind(namespace)
.bind(mode)
.bind(minute_bucket(*at))
.bind(count)
.bind(sum_ms)
.bind(max_ms)
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
Ok(())
}
async fn job_index_recent_transition_count(
&self,
namespace: &str,
stage: &str,
now: UtcDateTime,
) -> anyhow::Result<usize> {
let count: i64 = sqlx::query_scalar(
"SELECT COALESCE(SUM(count), 0) FROM later_job_transition_metrics \
WHERE namespace = ?1 AND stage = ?2 AND minute_bucket >= ?3",
)
.bind(namespace)
.bind(stage)
.bind(minute_bucket(now) - 1)
.fetch_one(&self.pool)
.await?;
Ok(usize::try_from(count)?)
}
async fn job_index_record_wait(
&self,
namespace: &str,
mode: &str,
wait_ms: i64,
at: UtcDateTime,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_job_wait_metrics \
(namespace, mode, minute_bucket, count, sum_ms, max_ms) \
VALUES (?1, ?2, ?3, 1, ?4, ?4) \
ON CONFLICT(namespace, mode, minute_bucket) DO UPDATE SET \
count = count + 1, \
sum_ms = sum_ms + ?4, \
max_ms = MAX(max_ms, ?4)",
)
.bind(namespace)
.bind(mode)
.bind(minute_bucket(at))
.bind(wait_ms)
.execute(&self.pool)
.await?;
Ok(())
}
async fn job_index_recent_wait_stats(
&self,
namespace: &str,
mode: &str,
now: UtcDateTime,
) -> anyhow::Result<JobIndexWaitStats> {
let row = sqlx::query(
"SELECT COALESCE(SUM(count), 0) AS c, COALESCE(SUM(sum_ms), 0) AS s, \
COALESCE(MAX(max_ms), 0) AS m \
FROM later_job_wait_metrics \
WHERE namespace = ?1 AND mode = ?2 AND minute_bucket >= ?3",
)
.bind(namespace)
.bind(mode)
.bind(minute_bucket(now) - 1)
.fetch_one(&self.pool)
.await?;
let count: i64 = row.try_get("c")?;
let sum_ms: i64 = row.try_get("s")?;
let max_ms: i64 = row.try_get("m")?;
let count = usize::try_from(count)?;
Ok(JobIndexWaitStats {
count,
avg_ms: if count == 0 {
0
} else {
sum_ms / i64::try_from(count)?
},
max_ms,
})
}
async fn job_index_sweep_expired(
&self,
namespace: &str,
now: UtcDateTime,
limit: usize,
) -> anyhow::Result<usize> {
let now_ms = now.timestamp_millis();
let limit_i64 = i64::try_from(limit)?;
let result = sqlx::query(
"DELETE FROM later_jobs_index WHERE rowid IN ( \
SELECT rowid FROM later_jobs_index \
WHERE namespace = ?1 AND date_expire IS NOT NULL AND date_expire <= ?2 \
LIMIT ?3 \
)",
)
.bind(namespace)
.bind(now_ms)
.bind(limit_i64)
.execute(&self.pool)
.await?;
let stale_bucket = minute_bucket(now) - 5;
sqlx::query(
"DELETE FROM later_job_transition_metrics \
WHERE namespace = ?1 AND minute_bucket < ?2",
)
.bind(namespace)
.bind(stale_bucket)
.execute(&self.pool)
.await?;
sqlx::query(
"DELETE FROM later_job_wait_metrics WHERE namespace = ?1 AND minute_bucket < ?2",
)
.bind(namespace)
.bind(stale_bucket)
.execute(&self.pool)
.await?;
Ok(usize::try_from(result.rows_affected())?)
}
}
#[cfg(test)]
mod stage_count_tests {
use super::*;
use crate::storage::Storage;
fn row(id: &str, stage: &str, date_expire: Option<UtcDateTime>) -> crate::storage::JobIndexRow {
let now = chrono::Utc::now();
crate::storage::JobIndexRow {
job_id: id.to_string(),
payload_type: "t".to_string(),
stage: stage.to_string(),
stage_date: now,
revision: 0,
created_at: now,
wait_ms: None,
wait_mode: None,
topic: None,
partition: None,
sequence: None,
parent_job_id: None,
date_expire,
}
}
#[tokio::test]
async fn stage_counts_skip_expired_terminal_rows_and_count_live_ones() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("c.db").display()
))
.await?;
let past = chrono::Utc::now() - chrono::Duration::hours(1);
let future = chrono::Utc::now() + chrono::Duration::hours(1);
for (id, stage, expire) in [
("a", "delayed", None),
("b", "delayed", None),
("c", "running", None),
("d", "success", Some(future)),
("e", "success", Some(past)),
("f", "failed", None),
] {
storage
.job_index_upsert("ns", row(id, stage, expire))
.await?;
}
storage.job_index_reconcile_stage_counts("ns").await?;
let counts = storage.job_index_stage_counts("ns").await?;
assert_eq!(counts.get("delayed"), Some(&2));
assert_eq!(counts.get("running"), Some(&1));
assert_eq!(
counts.get("success"),
Some(&1),
"expired success row must not count"
);
assert_eq!(counts.get("failed"), Some(&1));
Ok(())
}
async fn count_of(storage: &Sqlite, stage: &str) -> usize {
storage
.job_index_stage_counts("ns")
.await
.expect("stage counts")
.get(stage)
.copied()
.unwrap_or(0)
}
#[tokio::test]
async fn counters_follow_the_index_rows_exactly() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("t.db").display()
))
.await?;
let mut a = row("a", "delayed", None);
storage.job_index_upsert("ns", a.clone()).await?;
storage
.job_index_upsert("ns", row("b", "delayed", None))
.await?;
assert_eq!(count_of(&storage, "delayed").await, 2);
for _ in 0..5 {
storage.job_index_upsert("ns", a.clone()).await?;
}
assert_eq!(
count_of(&storage, "delayed").await,
2,
"same-stage re-saves change nothing"
);
a.stage = "running".to_string();
a.revision = 1;
storage.job_index_upsert("ns", a.clone()).await?;
assert_eq!(
(
count_of(&storage, "delayed").await,
count_of(&storage, "running").await
),
(1, 1)
);
let mut stale = row("a", "delayed", None);
stale.revision = 0;
storage.job_index_upsert("ns", stale).await?;
assert_eq!(
(
count_of(&storage, "delayed").await,
count_of(&storage, "running").await
),
(1, 1),
"a stale write must not move a counter"
);
a.stage = "success".to_string();
a.revision = 2;
a.date_expire = Some(chrono::Utc::now() + chrono::Duration::hours(1));
storage.job_index_upsert("ns", a).await?;
assert_eq!(
(
count_of(&storage, "running").await,
count_of(&storage, "success").await
),
(0, 1)
);
sqlx::query("DELETE FROM later_jobs_index WHERE namespace = 'ns' AND job_id = 'b'")
.execute(storage.pool())
.await?;
assert_eq!(count_of(&storage, "delayed").await, 0);
storage
.job_index_upsert("other", row("z", "delayed", None))
.await?;
assert_eq!(count_of(&storage, "delayed").await, 0);
Ok(())
}
#[tokio::test]
async fn a_recount_replaces_whatever_the_counters_drifted_to() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("r.db").display()
))
.await?;
storage
.job_index_upsert("ns", row("a", "delayed", None))
.await?;
sqlx::query("UPDATE later_stage_counts SET count = 999")
.execute(storage.pool())
.await?;
sqlx::query(
"INSERT INTO later_stage_counts (namespace, stage, count) VALUES ('ns', 'waiting', 40)",
)
.execute(storage.pool())
.await?;
assert_eq!(count_of(&storage, "delayed").await, 999);
storage.job_index_reconcile_stage_counts("ns").await?;
assert_eq!(count_of(&storage, "delayed").await, 1);
assert_eq!(count_of(&storage, "waiting").await, 0);
Ok(())
}
#[tokio::test]
#[ignore = "seeds 1.5M rows; run on demand"]
async fn stage_counts_stay_fast_at_a_million_and_a_half_jobs() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("big.db").display()
))
.await?;
let now = chrono::Utc::now().timestamp_millis();
sqlx::query(
"WITH RECURSIVE n(i) AS (SELECT 1 UNION ALL SELECT i + 1 FROM n WHERE i < 1500000) \
INSERT INTO later_jobs_index \
(namespace, job_id, payload_type, stage, stage_date_ms, revision, created_at_ms) \
SELECT 'ns', printf('job-%08d', i), 't', 'delayed', ?1 + i, 0, ?1 FROM n",
)
.bind(now)
.execute(storage.pool())
.await?;
storage.job_index_reconcile_stage_counts("ns").await?;
let started = std::time::Instant::now();
let counts = storage.job_index_stage_counts("ns").await?;
let counter_read = started.elapsed();
assert_eq!(counts.get("delayed"), Some(&1_500_000));
let started = std::time::Instant::now();
let scanned: i64 = sqlx::query_scalar(&format!(
"SELECT COUNT(*) FROM later_jobs_index WHERE namespace = 'ns' AND stage IN ({LIVE_STAGES})"
))
.fetch_one(storage.pool())
.await?;
let scan = started.elapsed();
assert_eq!(scanned, 1_500_000);
println!("counter read: {counter_read:?} index scan: {scan:?}");
assert!(
counter_read < std::time::Duration::from_millis(20),
"counter read took {counter_read:?}"
);
Ok(())
}
fn queued_job_message(id: &str, payload_type: &str) -> Vec<u8> {
crate::encoder::encode(&crate::models::AmqpCommand::ExecuteJob(
crate::models::JobAmqp {
payload_type: payload_type.to_string(),
id: crate::JobId(id.to_string()),
},
))
.expect("encode queue message")
}
#[tokio::test]
async fn queued_jobs_are_counted_and_listed_from_the_delivery_queue() -> anyhow::Result<()> {
use crate::mq::MqClient;
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("q.db").display()
))
.await?;
let queue = crate::mq::sql::SqliteQueue::new(storage.pool().clone(), "mq")?;
let publisher = queue.new_publisher("jobs").await?;
for number in 0..30 {
publisher
.publish(&queued_job_message(&format!("job-{number:02}"), "work"))
.await?;
}
publisher
.publish(&crate::encoder::encode(
&crate::models::AmqpCommand::PartitionReady {
topic: "t".to_string(),
},
)?)
.await?;
assert_eq!(
count_of_queue(&storage, "jobs").await,
31,
"publish adds to the counter"
);
let mut consumer = queue.new_consumer("jobs", 1).await?;
let first = consumer.next().await.expect("a message")?;
let claimed = sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM later_delivery_queue WHERE lease_until IS NOT NULL",
)
.fetch_one(storage.pool())
.await? as usize;
assert_eq!(count_of_queue(&storage, "jobs").await, 31 - claimed);
first.nack_requeue().await?;
let after_nack = count_of_queue(&storage, "jobs").await;
let still_claimed = sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM later_delivery_queue WHERE lease_until IS NOT NULL",
)
.fetch_one(storage.pool())
.await? as usize;
assert_eq!(after_nack, 31 - still_claimed);
drop(consumer);
sqlx::query("UPDATE later_delivery_queue SET lease_owner = NULL, lease_until = NULL")
.execute(storage.pool())
.await?;
sqlx::query("DELETE FROM later_queue_depth")
.execute(storage.pool())
.await?;
storage.job_index_reconcile_stage_counts("jobs").await?;
assert_eq!(count_of_queue(&storage, "jobs").await, 31);
let page = storage
.job_index_list_by_stage("jobs", "enqueued", None, 10)
.await?;
let ids: Vec<_> = page.items.iter().map(|row| row.job_id.as_str()).collect();
assert_eq!(
ids,
vec![
"job-29", "job-28", "job-27", "job-26", "job-25", "job-24", "job-23", "job-22",
"job-21"
],
"the PartitionReady message (newest) is skipped, jobs are newest first"
);
assert!(page.next_cursor.is_some());
let second = storage
.job_index_list_by_stage("jobs", "enqueued", page.next_cursor, 10)
.await?;
assert_eq!(
second.items.first().map(|row| row.job_id.as_str()),
Some("job-20")
);
assert!(page
.items
.iter()
.all(|row| row.stage == "enqueued" && row.payload_type == "work"));
Ok(())
}
async fn count_of_queue(storage: &Sqlite, queue: &str) -> usize {
storage
.job_index_stage_counts(queue)
.await
.expect("stage counts")
.get("enqueued")
.copied()
.unwrap_or(0)
}
#[tokio::test]
async fn the_purge_removes_only_rows_nothing_reads_any_more() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("p.db").display()
))
.await?;
let range = |key: &str, count: usize| {
(0..count)
.map(|n| StorageOperation::RangeAdd {
key: key.to_string(),
value: format!("v{n}").into_bytes(),
})
.collect::<Vec<_>>()
};
let mut operations = Vec::new();
for (key, count) in [
("ns-all-jobs", 7),
("ns-stage-enqueued-jobs", 5),
("ns-stage-running-jobs", 3),
("ns-stats-job-in-success", 4),
("ns-stats-partition-jobs-analytics-0", 6),
] {
operations.extend(range(key, count));
}
for (key, count) in [
("ns-stage-delayed-jobs", 2),
("ns-stage-requeued-jobs", 2),
("ns-stats-active-workers", 2),
("ns-all-recurring-jobs", 2),
("ns-job-abc-next", 2),
("other-all-jobs", 3),
] {
operations.extend(range(key, count));
}
storage.apply(operations).await?;
let queued = row("plain", "enqueued", None);
let mut partitioned = row("topic-job", "enqueued", None);
partitioned.topic = Some("t".to_string());
for index_row in [queued, partitioned, row("run", "running", None)] {
storage.job_index_upsert("ns", index_row).await?;
}
let mut removed = 0;
loop {
let batch = storage.purge_unused_dashboard_rows("ns", 4).await?;
if batch == 0 {
break;
}
removed += batch;
}
assert_eq!(
removed,
7 + 5 + 3 + 4 + 6 + 1,
"every dead range row and the queued index row"
);
let remaining: Vec<(String, i64)> = sqlx::query_as(
"SELECT range_key, COUNT(*) FROM later_storage_range GROUP BY range_key ORDER BY range_key",
)
.fetch_all(storage.pool())
.await?;
assert_eq!(
remaining,
vec![
("ns-all-recurring-jobs".to_string(), 2),
("ns-job-abc-next".to_string(), 2),
("ns-stage-delayed-jobs".to_string(), 2),
("ns-stage-requeued-jobs".to_string(), 2),
("ns-stats-active-workers".to_string(), 2),
("other-all-jobs".to_string(), 3),
]
);
let index_left: Vec<String> =
sqlx::query_scalar("SELECT job_id FROM later_jobs_index ORDER BY job_id")
.fetch_all(storage.pool())
.await?;
assert_eq!(index_left, vec!["run".to_string(), "topic-job".to_string()]);
Ok(())
}
#[tokio::test]
async fn live_stage_count_is_answered_from_the_index_alone() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let storage = Sqlite::new(&format!(
"sqlite://{}",
directory.path().join("p.db").display()
))
.await?;
let plan = sqlx::query(&format!(
"EXPLAIN QUERY PLAN SELECT stage, COUNT(*) FROM later_jobs_index \
WHERE namespace = ?1 AND stage IN ({LIVE_STAGES}) GROUP BY stage"
))
.bind("ns")
.fetch_all(storage.pool())
.await?;
let details: Vec<String> = plan
.iter()
.map(|r| r.try_get::<String, _>("detail"))
.collect::<Result<_, _>>()?;
assert!(
details
.iter()
.any(|d| d.contains("COVERING INDEX later_jobs_index_stage_idx")),
"live stage count must not read the table: {details:?}"
);
Ok(())
}
}
#[cfg(test)]
mod wal_tests {
use super::*;
use crate::storage::{Storage, StorageOperation};
async fn write_some(storage: &Sqlite, count: usize) -> anyhow::Result<()> {
for number in 0..count {
storage
.apply(vec![StorageOperation::Set {
key: format!("wal-key-{number}"),
value: vec![7; 2_000],
}])
.await?;
}
Ok(())
}
fn wal_len(path: &std::path::Path) -> u64 {
std::fs::metadata(path.with_extension("db-wal"))
.map(|metadata| metadata.len())
.unwrap_or(0)
}
#[tokio::test]
async fn checkpoint_truncates_a_large_log_and_leaves_a_small_one() -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let path = directory.path().join("wal.db");
let storage = Sqlite::new(&format!("sqlite://{}", path.display())).await?;
write_some(&storage, 40).await?;
assert!(wal_len(&path) > 0, "writes should have left a log behind");
assert_eq!(
storage.checkpoint_wal_above(1_000_000, 2_000).await?,
WalCheckpoint::BelowThreshold
);
assert!(wal_len(&path) > 0, "a small log is left alone");
assert_eq!(
storage.checkpoint_wal_above(1, 1_000).await?,
WalCheckpoint::Truncated
);
assert_eq!(wal_len(&path), 0, "a large log is cut back to nothing");
Ok(())
}
#[tokio::test]
async fn checkpoint_gives_up_quickly_behind_a_reader_instead_of_stalling_writers(
) -> anyhow::Result<()> {
let directory = tempfile::tempdir()?;
let path = directory.path().join("wal.db");
let storage = Sqlite::new(&format!("sqlite://{}", path.display())).await?;
write_some(&storage, 10).await?;
let mut reader = storage.pool().begin().await?;
let _: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM later_storage")
.fetch_one(&mut *reader)
.await?;
write_some(&storage, 10).await?;
let started = std::time::Instant::now();
assert_eq!(
storage.checkpoint_wal_above(1, 1_000).await?,
WalCheckpoint::GaveUp
);
assert!(
started.elapsed() < std::time::Duration::from_secs(6),
"gave up after {:?}",
started.elapsed()
);
let wait: i64 = sqlx::query_scalar("PRAGMA busy_timeout")
.fetch_one(storage.pool())
.await?;
assert!(
wait == 5_000 || wait == 30_000,
"busy_timeout left at {wait}"
);
reader.rollback().await?;
assert_eq!(
storage.checkpoint_wal_above(1, 1_000).await?,
WalCheckpoint::Truncated
);
assert_eq!(wal_len(&path), 0);
Ok(())
}
}