use crate::{mq::MqClient, storage::Storage};
use std::sync::Arc;
pub struct BackendParts {
pub namespace: String,
pub storage: Box<dyn Storage>,
pub delivery: Box<dyn MqClient>,
pub committer: Arc<dyn JobCommitter>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct PartitionSequence(pub i64);
#[derive(Clone, Debug)]
pub struct PartitionHeadClaim {
pub topic: String,
pub partition: crate::topic::PartitionId,
pub sequence: PartitionSequence,
pub job_id: crate::JobId,
pub owner: String,
pub lease_epoch: i64,
}
#[async_trait::async_trait]
pub trait JobCommitter: Send + Sync {
async fn commit(
&self,
operations: Vec<crate::storage::StorageOperation>,
queue_name: &str,
payload: &[u8],
) -> anyhow::Result<()>;
async fn claim(&self, job_id: &crate::JobId) -> anyhow::Result<Option<Box<dyn JobLease>>>;
async fn claim_and_fetch(
&self,
job_id: &crate::JobId,
key: &str,
storage: &dyn crate::storage::Storage,
) -> anyhow::Result<Option<ClaimedJob>> {
let Some(lease) = self.claim(job_id).await? else {
return Ok(None);
};
let value = storage.get(key).await?;
Ok(Some(ClaimedJob { lease, value }))
}
async fn stale_leased_job_ids(
&self,
_grace: std::time::Duration,
_limit: usize,
) -> anyhow::Result<Vec<crate::JobId>> {
Ok(Vec::new())
}
async fn register_topics(&self, topics: &[crate::topic::TopicConfig]) -> anyhow::Result<()> {
if topics.is_empty() {
return Ok(());
}
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
async fn commit_partitioned(
&self,
_build_operations: Box<
dyn FnOnce(PartitionSequence) -> anyhow::Result<Vec<crate::storage::StorageOperation>>
+ Send,
>,
_topic: &str,
_partition: crate::topic::PartitionId,
_job_id: &crate::JobId,
_payload_type: &str,
_earliest_run_at: crate::UtcDateTime,
) -> anyhow::Result<PartitionSequence> {
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
async fn claim_partition_head(
&self,
_owner: &str,
_live_workers: &[String],
) -> anyhow::Result<Option<PartitionHeadClaim>> {
Ok(None)
}
async fn peek_assigned_partition_head(
&self,
_owner: &str,
_topic: &str,
_partition: crate::topic::PartitionId,
_lease_epoch: i64,
) -> anyhow::Result<Option<PartitionHeadClaim>> {
Ok(None)
}
async fn partition_backlog_summary(&self) -> anyhow::Result<Vec<(String, u64, f64)>> {
Ok(Vec::new())
}
async fn partition_queue_depth_summary(
&self,
) -> anyhow::Result<Vec<(String, crate::topic::PartitionId, u64)>> {
Ok(Vec::new())
}
async fn admin_peek_partition_head(
&self,
_topic: &str,
_partition: crate::topic::PartitionId,
) -> anyhow::Result<Option<(crate::JobId, PartitionSequence)>> {
Ok(None)
}
async fn admin_force_complete_partition_head(
&self,
_topic: &str,
_partition: crate::topic::PartitionId,
_sequence: PartitionSequence,
_job_id: &crate::JobId,
_operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
async fn heartbeat_worker(&self, _owner: &str) -> anyhow::Result<()> {
Ok(())
}
async fn list_live_workers(&self) -> anyhow::Result<Vec<String>> {
Ok(Vec::new())
}
async fn deregister_worker(&self, _owner: &str) -> anyhow::Result<()> {
Ok(())
}
async fn renew_partition_lease(&self, _claim: &PartitionHeadClaim) -> anyhow::Result<bool> {
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
async fn release_partition_assignment(
&self,
_owner: &str,
_topic: &str,
_partition: crate::topic::PartitionId,
_lease_epoch: i64,
) -> anyhow::Result<bool> {
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
async fn release_partition_head(
&self,
_claim: &PartitionHeadClaim,
_earliest_run_at: crate::UtcDateTime,
_operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
async fn complete_partition_head(
&self,
_claim: &PartitionHeadClaim,
_operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
Err(anyhow::anyhow!(
"topic partitions are not supported by this backend"
))
}
}
#[async_trait::async_trait]
pub trait JobLease: Send + Sync {
async fn finish(self: Box<Self>) -> anyhow::Result<()>;
fn lease_lost(&self) -> tokio::sync::watch::Receiver<()>;
}
pub struct ClaimedJob {
pub lease: Box<dyn JobLease>,
pub value: Option<Vec<u8>>,
}
pub trait Backend: Send {
fn into_parts(self: Box<Self>) -> anyhow::Result<BackendParts>;
}
#[cfg(any(feature = "postgres", feature = "sqlite"))]
fn validate_namespace(namespace: &str) -> anyhow::Result<()> {
if namespace.trim().is_empty() {
return Err(anyhow::anyhow!("backend namespace cannot be empty"));
}
Ok(())
}
#[cfg(feature = "sqlite")]
#[derive(Clone)]
pub struct SqliteBackend {
namespace: String,
storage: crate::storage::Sqlite,
}
#[cfg(feature = "sqlite")]
impl SqliteBackend {
pub fn new(
namespace: impl Into<String>,
storage: crate::storage::Sqlite,
) -> anyhow::Result<Self> {
let namespace = namespace.into();
validate_namespace(&namespace)?;
Ok(Self { namespace, storage })
}
pub async fn connect(namespace: impl Into<String>, url: &str) -> anyhow::Result<Self> {
Self::new(namespace, crate::storage::Sqlite::new(url).await?)
}
pub async fn from_pool(
namespace: impl Into<String>,
pool: sqlx::SqlitePool,
) -> anyhow::Result<Self> {
Self::new(namespace, crate::storage::Sqlite::from_pool(pool).await?)
}
pub fn pool(&self) -> &sqlx::SqlitePool {
self.storage.pool()
}
pub fn namespace(&self) -> &str {
&self.namespace
}
pub async fn enqueue_in(
&self,
transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
message: impl crate::core::JobParameter,
) -> anyhow::Result<crate::JobId> {
let job = crate::bg_job_server_publisher::create_job(
message,
None,
None,
None,
&crate::retry::RetryPolicy::default(),
)?;
let id = job.id.clone();
let routing_key = format!("later-{}", self.namespace);
let command = crate::models::AmqpCommand::ExecuteJob(crate::models::JobAmqp {
payload_type: job.payload_type.clone(),
id: id.clone(),
});
let operations = crate::persist::Persist::job_operations(&routing_key, &job)?;
apply_sqlite_operations(transaction, operations).await?;
insert_sqlite_delivery_intent(
transaction,
&self.namespace,
&routing_key,
&crate::encoder::encode(&command)?,
)
.await?;
Ok(id)
}
pub async fn enqueue_to_partition_in<M>(
&self,
transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
message: M,
) -> anyhow::Result<crate::JobId>
where
M: crate::core::JobParameter + crate::topic::JobPartition + crate::topic::JobTopic,
{
let topic = M::topic_name().to_string();
let key = message.partition_key();
let partition_count =
sqlite_topic_partition_count(transaction, &self.namespace, &topic).await?;
let partition = crate::topic::partition_for_key(&key, partition_count);
let sequence =
sqlite_assign_partition_sequence(transaction, &self.namespace, &topic, partition)
.await?;
let mut job = crate::bg_job_server_publisher::create_job_in_topic(
message,
None,
None,
None,
&crate::retry::RetryPolicy::default(),
Some((topic.clone(), partition.0)),
)?;
job.sequence = Some(sequence.0);
let id = job.id.clone();
let routing_key = format!("later-{}", self.namespace);
let operations = crate::persist::Persist::job_operations(&routing_key, &job)?;
sqlite_insert_partition_job(
transaction,
&self.namespace,
&topic,
partition,
sequence,
&job.id,
&job.payload_type,
chrono::Utc::now(),
)
.await?;
apply_sqlite_operations(transaction, operations).await?;
let wake = crate::encoder::encode(&crate::models::AmqpCommand::PartitionReady {
topic: topic.clone(),
})?;
insert_sqlite_delivery_intent(transaction, &self.namespace, &routing_key, &wake).await?;
Ok(id)
}
}
#[cfg(feature = "sqlite")]
impl Backend for SqliteBackend {
fn into_parts(self: Box<Self>) -> anyhow::Result<BackendParts> {
let delivery: Box<dyn MqClient> = Box::new(crate::mq::sql::SqliteQueue::new(
self.storage.pool().clone(),
self.namespace.clone(),
)?);
let committer: Arc<dyn JobCommitter> = Arc::new(SqliteCommitter {
pool: self.storage.pool().clone(),
namespace: self.namespace.clone(),
});
Ok(BackendParts {
namespace: self.namespace,
storage: Box::new(self.storage),
delivery,
committer,
})
}
}
#[cfg(feature = "postgres")]
#[derive(Clone)]
pub struct PostgresBackend {
namespace: String,
storage: crate::storage::Postgres,
}
#[cfg(feature = "postgres")]
impl PostgresBackend {
pub fn new(
namespace: impl Into<String>,
storage: crate::storage::Postgres,
) -> anyhow::Result<Self> {
let namespace = namespace.into();
validate_namespace(&namespace)?;
Ok(Self { namespace, storage })
}
pub async fn connect(namespace: impl Into<String>, url: &str) -> anyhow::Result<Self> {
Self::new(namespace, crate::storage::Postgres::new(url).await?)
}
pub async fn from_pool(
namespace: impl Into<String>,
pool: sqlx::PgPool,
) -> anyhow::Result<Self> {
Self::new(namespace, crate::storage::Postgres::from_pool(pool).await?)
}
pub fn pool(&self) -> &sqlx::PgPool {
self.storage.pool()
}
pub fn namespace(&self) -> &str {
&self.namespace
}
pub async fn enqueue_in(
&self,
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
message: impl crate::core::JobParameter,
) -> anyhow::Result<crate::JobId> {
let job = crate::bg_job_server_publisher::create_job(
message,
None,
None,
None,
&crate::retry::RetryPolicy::default(),
)?;
let id = job.id.clone();
let routing_key = format!("later-{}", self.namespace);
let command = crate::models::AmqpCommand::ExecuteJob(crate::models::JobAmqp {
payload_type: job.payload_type.clone(),
id: id.clone(),
});
let operations = crate::persist::Persist::job_operations(&routing_key, &job)?;
apply_postgres_operations(transaction, operations).await?;
insert_postgres_delivery_intent(
transaction,
&self.namespace,
&routing_key,
&crate::encoder::encode(&command)?,
)
.await?;
Ok(id)
}
pub async fn enqueue_to_partition_in<M>(
&self,
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
message: M,
) -> anyhow::Result<crate::JobId>
where
M: crate::core::JobParameter + crate::topic::JobPartition + crate::topic::JobTopic,
{
let topic = M::topic_name().to_string();
let key = message.partition_key();
let partition_count =
postgres_topic_partition_count(transaction, &self.namespace, &topic).await?;
let partition = crate::topic::partition_for_key(&key, partition_count);
let sequence =
postgres_assign_partition_sequence(transaction, &self.namespace, &topic, partition)
.await?;
let mut job = crate::bg_job_server_publisher::create_job_in_topic(
message,
None,
None,
None,
&crate::retry::RetryPolicy::default(),
Some((topic.clone(), partition.0)),
)?;
job.sequence = Some(sequence.0);
let id = job.id.clone();
let routing_key = format!("later-{}", self.namespace);
let operations = crate::persist::Persist::job_operations(&routing_key, &job)?;
postgres_insert_partition_job(
transaction,
&self.namespace,
&topic,
partition,
sequence,
&job.id,
&job.payload_type,
chrono::Utc::now(),
)
.await?;
apply_postgres_operations(transaction, operations).await?;
let wake = crate::encoder::encode(&crate::models::AmqpCommand::PartitionReady {
topic: topic.clone(),
})?;
insert_postgres_delivery_intent(transaction, &self.namespace, &routing_key, &wake).await?;
Ok(id)
}
}
#[cfg(feature = "postgres")]
impl Backend for PostgresBackend {
fn into_parts(self: Box<Self>) -> anyhow::Result<BackendParts> {
let delivery: Box<dyn MqClient> = Box::new(crate::mq::sql::PostgresQueue::new(
self.storage.pool().clone(),
self.namespace.clone(),
)?);
let committer: Arc<dyn JobCommitter> = Arc::new(PostgresCommitter {
pool: self.storage.pool().clone(),
namespace: self.namespace.clone(),
});
Ok(BackendParts {
namespace: self.namespace,
storage: Box::new(self.storage),
delivery,
committer,
})
}
}
#[cfg(feature = "sqlite")]
struct SqliteCommitter {
pool: sqlx::SqlitePool,
namespace: String,
}
#[cfg(feature = "sqlite")]
struct SqliteJobLease {
pool: sqlx::SqlitePool,
namespace: String,
job_id: String,
lease_owner: String,
renewal: tokio::task::JoinHandle<()>,
lease_lost: tokio::sync::watch::Receiver<()>,
}
#[cfg(feature = "sqlite")]
impl Drop for SqliteJobLease {
fn drop(&mut self) {
self.renewal.abort();
}
}
#[cfg(feature = "sqlite")]
#[async_trait::async_trait]
impl JobLease for SqliteJobLease {
async fn finish(self: Box<Self>) -> anyhow::Result<()> {
self.renewal.abort();
sqlx::query(
"DELETE FROM later_job_execution \
WHERE namespace = ?1 AND job_id = ?2 AND lease_owner = ?3",
)
.bind(&self.namespace)
.bind(&self.job_id)
.bind(&self.lease_owner)
.execute(&self.pool)
.await?;
Ok(())
}
fn lease_lost(&self) -> tokio::sync::watch::Receiver<()> {
self.lease_lost.clone()
}
}
#[cfg(feature = "sqlite")]
fn spawn_sqlite_job_lease_renewal(
pool: sqlx::SqlitePool,
namespace: String,
job_id: String,
lease_owner: String,
) -> Box<dyn JobLease> {
const LEASE_MILLIS: i64 = 30_000;
let renewal_pool = pool.clone();
let renewal_namespace = namespace.clone();
let renewal_job_id = job_id.clone();
let renewal_owner = lease_owner.clone();
let (lease_lost_tx, lease_lost_rx) = tokio::sync::watch::channel(());
let renewal = tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(10)).await;
let lease_until = chrono::Utc::now()
.timestamp_millis()
.saturating_add(LEASE_MILLIS);
match sqlx::query(
"UPDATE later_job_execution SET lease_until = ?1 \
WHERE namespace = ?2 AND job_id = ?3 AND lease_owner = ?4",
)
.bind(lease_until)
.bind(&renewal_namespace)
.bind(&renewal_job_id)
.bind(&renewal_owner)
.execute(&renewal_pool)
.await
{
Ok(result) if result.rows_affected() == 1 => {}
Ok(_) => {
tracing::warn!(job_id = %renewal_job_id, "SQLite job lease was lost");
let _ = lease_lost_tx.send(());
return;
}
Err(error) => {
tracing::warn!(%error, job_id = %renewal_job_id, "Failed to renew SQLite job lease");
}
}
}
});
Box::new(SqliteJobLease {
pool,
namespace,
job_id,
lease_owner,
lease_lost: lease_lost_rx,
renewal,
})
}
#[cfg(feature = "sqlite")]
#[async_trait::async_trait]
impl JobCommitter for SqliteCommitter {
async fn commit(
&self,
operations: Vec<crate::storage::StorageOperation>,
queue_name: &str,
payload: &[u8],
) -> anyhow::Result<()> {
let mut transaction = self.pool.begin().await?;
apply_sqlite_operations(&mut transaction, operations).await?;
insert_sqlite_delivery_intent(&mut transaction, &self.namespace, queue_name, payload)
.await?;
transaction.commit().await?;
Ok(())
}
async fn claim(&self, job_id: &crate::JobId) -> anyhow::Result<Option<Box<dyn JobLease>>> {
const LEASE_MILLIS: i64 = 30_000;
let now = chrono::Utc::now().timestamp_millis();
let lease_until = now
.checked_add(LEASE_MILLIS)
.ok_or_else(|| anyhow::anyhow!("SQLite job lease is out of range"))?;
let lease_owner = crate::generate_id();
let job_id = job_id.to_string();
let claimed = sqlx::query(
"INSERT INTO later_job_execution \
(namespace, job_id, lease_owner, lease_until) VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(namespace, job_id) DO UPDATE SET \
lease_owner = excluded.lease_owner, lease_until = excluded.lease_until \
WHERE later_job_execution.lease_until <= ?5 \
RETURNING job_id",
)
.bind(&self.namespace)
.bind(&job_id)
.bind(&lease_owner)
.bind(lease_until)
.bind(now)
.fetch_optional(&self.pool)
.await?;
if claimed.is_none() {
return Ok(None);
}
Ok(Some(spawn_sqlite_job_lease_renewal(
self.pool.clone(),
self.namespace.clone(),
job_id,
lease_owner,
)))
}
async fn claim_and_fetch(
&self,
job_id: &crate::JobId,
key: &str,
_storage: &dyn crate::storage::Storage,
) -> anyhow::Result<Option<ClaimedJob>> {
const LEASE_MILLIS: i64 = 30_000;
let now = chrono::Utc::now().timestamp_millis();
let lease_until = now
.checked_add(LEASE_MILLIS)
.ok_or_else(|| anyhow::anyhow!("SQLite job lease is out of range"))?;
let lease_owner = crate::generate_id();
let job_id = job_id.to_string();
let mut transaction = self.pool.begin().await?;
let claimed = sqlx::query(
"INSERT INTO later_job_execution \
(namespace, job_id, lease_owner, lease_until) VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(namespace, job_id) DO UPDATE SET \
lease_owner = excluded.lease_owner, lease_until = excluded.lease_until \
WHERE later_job_execution.lease_until <= ?5 \
RETURNING job_id",
)
.bind(&self.namespace)
.bind(&job_id)
.bind(&lease_owner)
.bind(lease_until)
.bind(now)
.fetch_optional(&mut *transaction)
.await?;
if claimed.is_none() {
transaction.rollback().await?;
return Ok(None);
}
let value: Option<Vec<u8>> = {
use sqlx::Row;
sqlx::query(
"SELECT value FROM later_storage \
WHERE key = ?1 AND (date_expire IS NULL OR date_expire > ?2)",
)
.bind(key)
.bind(now)
.fetch_optional(&mut *transaction)
.await?
.map(|row| row.try_get("value"))
.transpose()?
};
transaction.commit().await?;
Ok(Some(ClaimedJob {
lease: spawn_sqlite_job_lease_renewal(
self.pool.clone(),
self.namespace.clone(),
job_id,
lease_owner,
),
value,
}))
}
async fn stale_leased_job_ids(
&self,
grace: std::time::Duration,
limit: usize,
) -> anyhow::Result<Vec<crate::JobId>> {
let cutoff = chrono::Utc::now().timestamp_millis()
- i64::try_from(grace.as_millis()).unwrap_or(i64::MAX);
let limit = i64::try_from(limit).unwrap_or(i64::MAX);
let job_ids: Vec<(String,)> = sqlx::query_as(
"SELECT job_id FROM later_job_execution \
WHERE namespace = ?1 AND lease_until <= ?2 LIMIT ?3",
)
.bind(&self.namespace)
.bind(cutoff)
.bind(limit)
.fetch_all(&self.pool)
.await?;
Ok(job_ids.into_iter().map(|(id,)| crate::JobId(id)).collect())
}
async fn register_topics(&self, topics: &[crate::topic::TopicConfig]) -> anyhow::Result<()> {
for topic in topics {
let partition_count = i64::from(topic.partition_count());
let existing: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = ?1 AND topic = ?2",
)
.bind(&self.namespace)
.bind(topic.name())
.fetch_optional(&self.pool)
.await?;
match existing {
Some((existing_count,)) if existing_count == partition_count => {}
Some((existing_count,)) => {
return Err(anyhow::anyhow!(
"topic '{}' is already registered with {} partitions, not {}",
topic.name(),
existing_count,
partition_count
));
}
None => {
sqlx::query(
"INSERT INTO later_topic (namespace, topic, partition_count) \
VALUES (?1, ?2, ?3)",
)
.bind(&self.namespace)
.bind(topic.name())
.bind(partition_count)
.execute(&self.pool)
.await?;
}
}
}
Ok(())
}
async fn commit_partitioned(
&self,
build_operations: Box<
dyn FnOnce(PartitionSequence) -> anyhow::Result<Vec<crate::storage::StorageOperation>>
+ Send,
>,
topic: &str,
partition: crate::topic::PartitionId,
job_id: &crate::JobId,
payload_type: &str,
earliest_run_at: crate::UtcDateTime,
) -> anyhow::Result<PartitionSequence> {
let mut transaction = self.pool.begin_with("BEGIN IMMEDIATE").await?;
let sequence =
sqlite_assign_partition_sequence(&mut transaction, &self.namespace, topic, partition)
.await?;
sqlite_insert_partition_job(
&mut transaction,
&self.namespace,
topic,
partition,
sequence,
job_id,
payload_type,
earliest_run_at,
)
.await?;
let operations = build_operations(sequence)?;
apply_sqlite_operations(&mut transaction, operations).await?;
let routing_key = format!("later-{}", self.namespace);
let wake = crate::encoder::encode(&crate::models::AmqpCommand::PartitionReady {
topic: topic.to_string(),
})?;
insert_sqlite_delivery_intent(&mut transaction, &self.namespace, &routing_key, &wake)
.await?;
transaction.commit().await?;
Ok(sequence)
}
async fn claim_partition_head(
&self,
owner: &str,
live_workers: &[String],
) -> anyhow::Result<Option<PartitionHeadClaim>> {
let now = chrono::Utc::now().timestamp_millis();
let lease_until = now + PARTITION_LEASE_MILLIS;
let mut transaction = self.pool.begin_with("BEGIN IMMEDIATE").await?;
let candidates: Vec<(String, i64, i64, String)> = sqlx::query_as(
"SELECT topic, partition_id, sequence, job_id \
FROM later_partition_head \
WHERE namespace = ?1 AND earliest_run_at <= ?2 \
ORDER BY earliest_run_at ASC \
LIMIT 50",
)
.bind(&self.namespace)
.bind(now)
.fetch_all(&mut *transaction)
.await?;
for (topic, partition_id, sequence, job_id) in candidates {
let partition = crate::topic::PartitionId(u32::try_from(partition_id)?);
if crate::topic::rendezvous_owner(&topic, partition, live_workers) != Some(owner) {
continue;
}
let claimed: Option<(i64,)> = sqlx::query_as(
"INSERT INTO later_partition_lease \
(namespace, topic, partition_id, owner, lease_until, lease_epoch) \
VALUES (?1, ?2, ?3, ?4, ?5, 1) \
ON CONFLICT(namespace, topic, partition_id) DO UPDATE SET \
owner = excluded.owner, lease_until = excluded.lease_until, \
lease_epoch = later_partition_lease.lease_epoch + 1 \
WHERE later_partition_lease.lease_until <= ?6 \
RETURNING lease_epoch",
)
.bind(&self.namespace)
.bind(&topic)
.bind(partition_id)
.bind(owner)
.bind(lease_until)
.bind(now)
.fetch_optional(&mut *transaction)
.await?;
if let Some((lease_epoch,)) = claimed {
let current: Option<(i64, String)> = sqlx::query_as(
"SELECT sequence, job_id FROM later_partition_head \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3",
)
.bind(&self.namespace)
.bind(&topic)
.bind(partition_id)
.fetch_optional(&mut *transaction)
.await?;
if current != Some((sequence, job_id.clone())) {
transaction.rollback().await?;
return Ok(None);
}
transaction.commit().await?;
return Ok(Some(PartitionHeadClaim {
topic,
partition,
sequence: PartitionSequence(sequence),
job_id: crate::JobId(job_id),
owner: owner.to_string(),
lease_epoch,
}));
}
}
transaction.rollback().await?;
Ok(None)
}
async fn peek_assigned_partition_head(
&self,
owner: &str,
topic: &str,
partition: crate::topic::PartitionId,
lease_epoch: i64,
) -> anyhow::Result<Option<PartitionHeadClaim>> {
let now = chrono::Utc::now().timestamp_millis();
let row: Option<(i64, String)> = sqlx::query_as(
"SELECT h.sequence, h.job_id \
FROM later_partition_head h \
JOIN later_partition_lease l \
ON l.namespace = h.namespace AND l.topic = h.topic \
AND l.partition_id = h.partition_id \
WHERE h.namespace = ?1 AND h.topic = ?2 AND h.partition_id = ?3 \
AND h.earliest_run_at <= ?4 \
AND l.owner = ?5 AND l.lease_epoch = ?6 AND l.lease_until > ?4",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(now)
.bind(owner)
.bind(lease_epoch)
.fetch_optional(&self.pool)
.await?;
Ok(row.map(|(sequence, job_id)| PartitionHeadClaim {
topic: topic.to_string(),
partition,
sequence: PartitionSequence(sequence),
job_id: crate::JobId(job_id),
owner: owner.to_string(),
lease_epoch,
}))
}
async fn partition_backlog_summary(&self) -> anyhow::Result<Vec<(String, u64, f64)>> {
let now = chrono::Utc::now().timestamp_millis();
let rows: Vec<(String, i64, i64)> = sqlx::query_as(
"SELECT topic, COUNT(*) AS ready_count, MIN(earliest_run_at) AS oldest \
FROM later_partition_head \
WHERE namespace = ?1 AND earliest_run_at <= ?2 \
GROUP BY topic",
)
.bind(&self.namespace)
.bind(now)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(|(topic, ready_count, oldest)| {
let ready_count = u64::try_from(ready_count).unwrap_or(0);
let age_seconds = (now - oldest).max(0) as f64 / 1_000.0;
(topic, ready_count, age_seconds)
})
.collect())
}
async fn partition_queue_depth_summary(
&self,
) -> anyhow::Result<Vec<(String, crate::topic::PartitionId, u64)>> {
let rows: Vec<(String, i64, i64)> = sqlx::query_as(
"SELECT topic, partition_id, COUNT(*) AS depth \
FROM later_partition_job \
WHERE namespace = ?1 \
GROUP BY topic, partition_id",
)
.bind(&self.namespace)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(|(topic, partition_id, depth)| {
(
topic,
crate::topic::PartitionId(u32::try_from(partition_id).unwrap_or(0)),
u64::try_from(depth).unwrap_or(0),
)
})
.collect())
}
async fn admin_peek_partition_head(
&self,
topic: &str,
partition: crate::topic::PartitionId,
) -> anyhow::Result<Option<(crate::JobId, PartitionSequence)>> {
let row: Option<(i64, String)> = sqlx::query_as(
"SELECT sequence, job_id FROM later_partition_head \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.fetch_optional(&self.pool)
.await?;
Ok(row.map(|(sequence, job_id)| (crate::JobId(job_id), PartitionSequence(sequence))))
}
async fn admin_force_complete_partition_head(
&self,
topic: &str,
partition: crate::topic::PartitionId,
sequence: PartitionSequence,
job_id: &crate::JobId,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
let mut transaction = self.pool.begin_with("BEGIN IMMEDIATE").await?;
let deleted = sqlx::query(
"DELETE FROM later_partition_job \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 AND sequence = ?4 \
AND job_id = ?5",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(sequence.0)
.bind(job_id.to_string())
.execute(&mut *transaction)
.await?;
if deleted.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
sqlite_advance_partition_head(&mut transaction, &self.namespace, topic, partition).await?;
let now = chrono::Utc::now().timestamp_millis();
sqlx::query(
"UPDATE later_partition_lease SET lease_epoch = lease_epoch + 1, lease_until = ?4 \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(now)
.execute(&mut *transaction)
.await?;
apply_sqlite_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(true)
}
async fn heartbeat_worker(&self, owner: &str) -> anyhow::Result<()> {
let heartbeat_until = chrono::Utc::now().timestamp_millis() + WORKER_HEARTBEAT_MILLIS;
sqlx::query(
"INSERT INTO later_worker (namespace, owner, heartbeat_until) VALUES (?1, ?2, ?3) \
ON CONFLICT(namespace, owner) DO UPDATE SET heartbeat_until = excluded.heartbeat_until",
)
.bind(&self.namespace)
.bind(owner)
.bind(heartbeat_until)
.execute(&self.pool)
.await?;
Ok(())
}
async fn list_live_workers(&self) -> anyhow::Result<Vec<String>> {
let now = chrono::Utc::now().timestamp_millis();
let rows: Vec<(String,)> = sqlx::query_as(
"SELECT owner FROM later_worker WHERE namespace = ?1 AND heartbeat_until > ?2",
)
.bind(&self.namespace)
.bind(now)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|(owner,)| owner).collect())
}
async fn deregister_worker(&self, owner: &str) -> anyhow::Result<()> {
sqlx::query("DELETE FROM later_worker WHERE namespace = ?1 AND owner = ?2")
.bind(&self.namespace)
.bind(owner)
.execute(&self.pool)
.await?;
Ok(())
}
async fn renew_partition_lease(&self, claim: &PartitionHeadClaim) -> anyhow::Result<bool> {
let now = chrono::Utc::now().timestamp_millis();
let lease_until = chrono::Utc::now().timestamp_millis() + PARTITION_LEASE_MILLIS;
let result = sqlx::query(
"UPDATE later_partition_lease SET lease_until = ?1 \
WHERE namespace = ?2 AND topic = ?3 AND partition_id = ?4 AND owner = ?5 \
AND lease_epoch = ?6 AND lease_until > ?7",
)
.bind(lease_until)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(&claim.owner)
.bind(claim.lease_epoch)
.bind(now)
.execute(&self.pool)
.await?;
Ok(result.rows_affected() == 1)
}
async fn release_partition_assignment(
&self,
owner: &str,
topic: &str,
partition: crate::topic::PartitionId,
lease_epoch: i64,
) -> anyhow::Result<bool> {
let now = chrono::Utc::now().timestamp_millis();
let result = sqlx::query(
"UPDATE later_partition_lease \
SET lease_epoch = lease_epoch + 1, lease_until = ?1 \
WHERE namespace = ?2 AND topic = ?3 AND partition_id = ?4 AND owner = ?5 \
AND lease_epoch = ?6 AND lease_until > ?1",
)
.bind(now)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(owner)
.bind(lease_epoch)
.execute(&self.pool)
.await?;
Ok(result.rows_affected() == 1)
}
async fn release_partition_head(
&self,
claim: &PartitionHeadClaim,
earliest_run_at: crate::UtcDateTime,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
let mut transaction = self.pool.begin().await?;
let now = chrono::Utc::now().timestamp_millis();
let released = sqlx::query(
"UPDATE later_partition_lease SET lease_until = ?1 \
WHERE namespace = ?2 AND topic = ?3 AND partition_id = ?4 AND owner = ?5 \
AND lease_epoch = ?6 AND lease_until > ?7",
)
.bind(now)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(&claim.owner)
.bind(claim.lease_epoch)
.bind(now)
.execute(&mut *transaction)
.await?;
if released.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
let updated = sqlx::query(
"UPDATE later_partition_job SET earliest_run_at = ?1 \
WHERE namespace = ?2 AND topic = ?3 AND partition_id = ?4 AND sequence = ?5 \
AND job_id = ?6",
)
.bind(earliest_run_at.timestamp_millis())
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(claim.sequence.0)
.bind(claim.job_id.to_string())
.execute(&mut *transaction)
.await?;
if updated.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
let head_updated = sqlx::query(
"UPDATE later_partition_head SET earliest_run_at = ?1 \
WHERE namespace = ?2 AND topic = ?3 AND partition_id = ?4 AND sequence = ?5 \
AND job_id = ?6",
)
.bind(earliest_run_at.timestamp_millis())
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(claim.sequence.0)
.bind(claim.job_id.to_string())
.execute(&mut *transaction)
.await?;
if head_updated.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
apply_sqlite_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(true)
}
async fn complete_partition_head(
&self,
claim: &PartitionHeadClaim,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
let mut transaction = self.pool.begin().await?;
let now = chrono::Utc::now().timestamp_millis();
let deleted = sqlx::query(
"DELETE FROM later_partition_job \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 AND sequence = ?4 \
AND job_id = ?5 \
AND EXISTS ( \
SELECT 1 FROM later_partition_lease \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 \
AND owner = ?6 AND lease_epoch = ?7 AND lease_until > ?8 \
)",
)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(claim.sequence.0)
.bind(claim.job_id.to_string())
.bind(&claim.owner)
.bind(claim.lease_epoch)
.bind(now)
.execute(&mut *transaction)
.await?;
if deleted.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
sqlite_advance_partition_head(
&mut transaction,
&self.namespace,
&claim.topic,
claim.partition,
)
.await?;
apply_sqlite_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(true)
}
}
#[cfg(feature = "sqlite")]
const PARTITION_LEASE_MILLIS: i64 = 30_000;
#[cfg(feature = "sqlite")]
const WORKER_HEARTBEAT_MILLIS: i64 = 15_000;
#[cfg(feature = "sqlite")]
async fn sqlite_topic_partition_count(
transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
namespace: &str,
topic: &str,
) -> anyhow::Result<u32> {
let partition_count: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = ?1 AND topic = ?2",
)
.bind(namespace)
.bind(topic)
.fetch_optional(&mut **transaction)
.await?;
let Some((partition_count,)) = partition_count else {
return Err(anyhow::anyhow!("unknown topic '{topic}'"));
};
Ok(u32::try_from(partition_count)?)
}
#[cfg(feature = "sqlite")]
async fn sqlite_assign_partition_sequence(
transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
namespace: &str,
topic: &str,
partition: crate::topic::PartitionId,
) -> anyhow::Result<PartitionSequence> {
let partition_count: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = ?1 AND topic = ?2",
)
.bind(namespace)
.bind(topic)
.fetch_optional(&mut **transaction)
.await?;
let Some((partition_count,)) = partition_count else {
return Err(anyhow::anyhow!("unknown topic '{topic}'"));
};
if i64::from(partition.0) >= partition_count {
return Err(anyhow::anyhow!(
"partition {} is out of range for topic '{topic}' with {partition_count} partitions",
partition.0
));
}
let (sequence,): (i64,) = sqlx::query_as(
"INSERT INTO later_partition_sequence (namespace, topic, partition_id, next_sequence) \
VALUES (?1, ?2, ?3, 2) \
ON CONFLICT(namespace, topic, partition_id) DO UPDATE SET \
next_sequence = later_partition_sequence.next_sequence + 1 \
RETURNING next_sequence - 1",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.fetch_one(&mut **transaction)
.await?;
Ok(PartitionSequence(sequence))
}
#[cfg(feature = "sqlite")]
#[allow(clippy::too_many_arguments)]
async fn sqlite_insert_partition_job(
transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
namespace: &str,
topic: &str,
partition: crate::topic::PartitionId,
sequence: PartitionSequence,
job_id: &crate::JobId,
payload_type: &str,
earliest_run_at: crate::UtcDateTime,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_partition_job \
(namespace, topic, partition_id, sequence, job_id, payload_type, earliest_run_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(sequence.0)
.bind(job_id.to_string())
.bind(payload_type)
.bind(earliest_run_at.timestamp_millis())
.execute(&mut **transaction)
.await?;
sqlx::query(
"INSERT INTO later_partition_head \
(namespace, topic, partition_id, sequence, job_id, earliest_run_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6) \
ON CONFLICT(namespace, topic, partition_id) DO NOTHING",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(sequence.0)
.bind(job_id.to_string())
.bind(earliest_run_at.timestamp_millis())
.execute(&mut **transaction)
.await?;
Ok(())
}
#[cfg(feature = "sqlite")]
async fn sqlite_advance_partition_head(
transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>,
namespace: &str,
topic: &str,
partition: crate::topic::PartitionId,
) -> anyhow::Result<()> {
sqlx::query(
"DELETE FROM later_partition_head \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.execute(&mut **transaction)
.await?;
sqlx::query(
"INSERT INTO later_partition_head \
(namespace, topic, partition_id, sequence, job_id, earliest_run_at) \
SELECT namespace, topic, partition_id, sequence, job_id, earliest_run_at \
FROM later_partition_job \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 \
ORDER BY sequence ASC LIMIT 1",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.execute(&mut **transaction)
.await?;
Ok(())
}
#[cfg(feature = "postgres")]
struct PostgresCommitter {
pool: sqlx::PgPool,
namespace: String,
}
#[cfg(feature = "postgres")]
struct PostgresJobLease {
pool: sqlx::PgPool,
namespace: String,
job_id: String,
lease_owner: String,
renewal: tokio::task::JoinHandle<()>,
lease_lost: tokio::sync::watch::Receiver<()>,
}
#[cfg(feature = "postgres")]
impl Drop for PostgresJobLease {
fn drop(&mut self) {
self.renewal.abort();
}
}
#[cfg(feature = "postgres")]
#[async_trait::async_trait]
impl JobLease for PostgresJobLease {
async fn finish(self: Box<Self>) -> anyhow::Result<()> {
self.renewal.abort();
sqlx::query(
"DELETE FROM later_job_execution \
WHERE namespace = $1 AND job_id = $2 AND lease_owner = $3",
)
.bind(&self.namespace)
.bind(&self.job_id)
.bind(&self.lease_owner)
.execute(&self.pool)
.await?;
Ok(())
}
fn lease_lost(&self) -> tokio::sync::watch::Receiver<()> {
self.lease_lost.clone()
}
}
#[cfg(feature = "postgres")]
fn spawn_postgres_job_lease_renewal(
pool: sqlx::PgPool,
namespace: String,
job_id: String,
lease_owner: String,
) -> Box<dyn JobLease> {
const LEASE_SECONDS: i64 = 30;
let renewal_pool = pool.clone();
let renewal_namespace = namespace.clone();
let renewal_job_id = job_id.clone();
let renewal_owner = lease_owner.clone();
let (lease_lost_tx, lease_lost_rx) = tokio::sync::watch::channel(());
let renewal = tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_secs(10)).await;
match sqlx::query(
"UPDATE later_job_execution SET \
lease_until = NOW() + make_interval(secs => $1) \
WHERE namespace = $2 AND job_id = $3 AND lease_owner = $4",
)
.bind(LEASE_SECONDS)
.bind(&renewal_namespace)
.bind(&renewal_job_id)
.bind(&renewal_owner)
.execute(&renewal_pool)
.await
{
Ok(result) if result.rows_affected() == 1 => {}
Ok(_) => {
tracing::warn!(job_id = %renewal_job_id, "Postgres job lease was lost");
let _ = lease_lost_tx.send(());
return;
}
Err(error) => {
tracing::warn!(%error, job_id = %renewal_job_id, "Failed to renew Postgres job lease");
}
}
}
});
Box::new(PostgresJobLease {
pool,
namespace,
job_id,
lease_owner,
renewal,
lease_lost: lease_lost_rx,
})
}
#[cfg(feature = "postgres")]
#[async_trait::async_trait]
impl JobCommitter for PostgresCommitter {
async fn commit(
&self,
operations: Vec<crate::storage::StorageOperation>,
queue_name: &str,
payload: &[u8],
) -> anyhow::Result<()> {
let mut transaction = self.pool.begin().await?;
apply_postgres_operations(&mut transaction, operations).await?;
insert_postgres_delivery_intent(&mut transaction, &self.namespace, queue_name, payload)
.await?;
transaction.commit().await?;
Ok(())
}
async fn claim(&self, job_id: &crate::JobId) -> anyhow::Result<Option<Box<dyn JobLease>>> {
const LEASE_SECONDS: i64 = 30;
let lease_owner = crate::generate_id();
let job_id = job_id.to_string();
let claimed = sqlx::query(
"INSERT INTO later_job_execution \
(namespace, job_id, lease_owner, lease_until) \
VALUES ($1, $2, $3, NOW() + make_interval(secs => $4)) \
ON CONFLICT(namespace, job_id) DO UPDATE SET \
lease_owner = EXCLUDED.lease_owner, lease_until = EXCLUDED.lease_until \
WHERE later_job_execution.lease_until <= NOW() \
RETURNING job_id",
)
.bind(&self.namespace)
.bind(&job_id)
.bind(&lease_owner)
.bind(LEASE_SECONDS)
.fetch_optional(&self.pool)
.await?;
if claimed.is_none() {
return Ok(None);
}
Ok(Some(spawn_postgres_job_lease_renewal(
self.pool.clone(),
self.namespace.clone(),
job_id,
lease_owner,
)))
}
async fn claim_and_fetch(
&self,
job_id: &crate::JobId,
key: &str,
_storage: &dyn crate::storage::Storage,
) -> anyhow::Result<Option<ClaimedJob>> {
use sqlx::Row;
const LEASE_SECONDS: i64 = 30;
let lease_owner = crate::generate_id();
let job_id = job_id.to_string();
let row = sqlx::query(
"WITH claim AS ( \
INSERT INTO later_job_execution \
(namespace, job_id, lease_owner, lease_until) \
VALUES ($1, $2, $3, NOW() + make_interval(secs => $4)) \
ON CONFLICT(namespace, job_id) DO UPDATE SET \
lease_owner = EXCLUDED.lease_owner, lease_until = EXCLUDED.lease_until \
WHERE later_job_execution.lease_until <= NOW() \
RETURNING job_id \
) \
SELECT claim.job_id AS claimed_job_id, later_storage.value AS value \
FROM (SELECT 1) AS anchor \
LEFT JOIN claim ON true \
LEFT JOIN later_storage \
ON later_storage.key = $5 \
AND (later_storage.date_expire IS NULL OR later_storage.date_expire > NOW())",
)
.bind(&self.namespace)
.bind(&job_id)
.bind(&lease_owner)
.bind(LEASE_SECONDS)
.bind(key)
.fetch_one(&self.pool)
.await?;
let claimed_job_id: Option<String> = row.try_get("claimed_job_id")?;
if claimed_job_id.is_none() {
return Ok(None);
}
let value: Option<Vec<u8>> = row.try_get("value")?;
Ok(Some(ClaimedJob {
lease: spawn_postgres_job_lease_renewal(
self.pool.clone(),
self.namespace.clone(),
job_id,
lease_owner,
),
value,
}))
}
async fn stale_leased_job_ids(
&self,
grace: std::time::Duration,
limit: usize,
) -> anyhow::Result<Vec<crate::JobId>> {
let grace_secs = i64::try_from(grace.as_secs()).unwrap_or(i64::MAX);
let limit = i64::try_from(limit).unwrap_or(i64::MAX);
let job_ids: Vec<(String,)> = sqlx::query_as(
"SELECT job_id FROM later_job_execution \
WHERE namespace = $1 AND lease_until <= NOW() - make_interval(secs => $2) LIMIT $3",
)
.bind(&self.namespace)
.bind(grace_secs)
.bind(limit)
.fetch_all(&self.pool)
.await?;
Ok(job_ids.into_iter().map(|(id,)| crate::JobId(id)).collect())
}
async fn register_topics(&self, topics: &[crate::topic::TopicConfig]) -> anyhow::Result<()> {
for topic in topics {
let partition_count = i64::from(topic.partition_count());
let existing: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = $1 AND topic = $2",
)
.bind(&self.namespace)
.bind(topic.name())
.fetch_optional(&self.pool)
.await?;
match existing {
Some((existing_count,)) if existing_count == partition_count => {}
Some((existing_count,)) => {
return Err(anyhow::anyhow!(
"topic '{}' is already registered with {} partitions, not {}",
topic.name(),
existing_count,
partition_count
));
}
None => {
sqlx::query(
"INSERT INTO later_topic (namespace, topic, partition_count) \
VALUES ($1, $2, $3)",
)
.bind(&self.namespace)
.bind(topic.name())
.bind(partition_count)
.execute(&self.pool)
.await?;
}
}
}
Ok(())
}
async fn commit_partitioned(
&self,
build_operations: Box<
dyn FnOnce(PartitionSequence) -> anyhow::Result<Vec<crate::storage::StorageOperation>>
+ Send,
>,
topic: &str,
partition: crate::topic::PartitionId,
job_id: &crate::JobId,
payload_type: &str,
earliest_run_at: crate::UtcDateTime,
) -> anyhow::Result<PartitionSequence> {
let mut transaction = self.pool.begin().await?;
let sequence =
postgres_assign_partition_sequence(&mut transaction, &self.namespace, topic, partition)
.await?;
postgres_insert_partition_job(
&mut transaction,
&self.namespace,
topic,
partition,
sequence,
job_id,
payload_type,
earliest_run_at,
)
.await?;
let operations = build_operations(sequence)?;
apply_postgres_operations(&mut transaction, operations).await?;
let routing_key = format!("later-{}", self.namespace);
let wake = crate::encoder::encode(&crate::models::AmqpCommand::PartitionReady {
topic: topic.to_string(),
})?;
insert_postgres_delivery_intent(&mut transaction, &self.namespace, &routing_key, &wake)
.await?;
transaction.commit().await?;
Ok(sequence)
}
async fn partition_backlog_summary(&self) -> anyhow::Result<Vec<(String, u64, f64)>> {
let rows: Vec<(String, i64, f64)> = sqlx::query_as(
"SELECT topic, COUNT(*) AS ready_count, \
EXTRACT(EPOCH FROM (NOW() - MIN(earliest_run_at)))::float8 AS oldest_age_seconds \
FROM later_partition_head \
WHERE namespace = $1 AND earliest_run_at <= NOW() \
GROUP BY topic",
)
.bind(&self.namespace)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(|(topic, ready_count, oldest_age_seconds)| {
(
topic,
u64::try_from(ready_count).unwrap_or(0),
oldest_age_seconds.max(0.0),
)
})
.collect())
}
async fn partition_queue_depth_summary(
&self,
) -> anyhow::Result<Vec<(String, crate::topic::PartitionId, u64)>> {
let rows: Vec<(String, i64, i64)> = sqlx::query_as(
"SELECT topic, partition_id, COUNT(*) AS depth \
FROM later_partition_job \
WHERE namespace = $1 \
GROUP BY topic, partition_id",
)
.bind(&self.namespace)
.fetch_all(&self.pool)
.await?;
Ok(rows
.into_iter()
.map(|(topic, partition_id, depth)| {
(
topic,
crate::topic::PartitionId(u32::try_from(partition_id).unwrap_or(0)),
u64::try_from(depth).unwrap_or(0),
)
})
.collect())
}
async fn claim_partition_head(
&self,
owner: &str,
live_workers: &[String],
) -> anyhow::Result<Option<PartitionHeadClaim>> {
let lease_seconds = PARTITION_LEASE_SECONDS;
let mut transaction = self.pool.begin().await?;
let candidates: Vec<(String, i64, i64, String)> = sqlx::query_as(
"SELECT topic, partition_id, sequence, job_id \
FROM later_partition_head \
WHERE namespace = $1 AND earliest_run_at <= NOW() \
ORDER BY earliest_run_at ASC \
LIMIT 50",
)
.bind(&self.namespace)
.fetch_all(&mut *transaction)
.await?;
for (topic, partition_id, sequence, job_id) in candidates {
let partition = crate::topic::PartitionId(u32::try_from(partition_id)?);
if crate::topic::rendezvous_owner(&topic, partition, live_workers) != Some(owner) {
continue;
}
let claimed: Option<(i64,)> = sqlx::query_as(
"INSERT INTO later_partition_lease \
(namespace, topic, partition_id, owner, lease_until, lease_epoch) \
VALUES ($1, $2, $3, $4, NOW() + make_interval(secs => $5), 1) \
ON CONFLICT(namespace, topic, partition_id) DO UPDATE SET \
owner = EXCLUDED.owner, lease_until = EXCLUDED.lease_until, \
lease_epoch = later_partition_lease.lease_epoch + 1 \
WHERE later_partition_lease.lease_until <= NOW() \
RETURNING lease_epoch",
)
.bind(&self.namespace)
.bind(&topic)
.bind(partition_id)
.bind(owner)
.bind(lease_seconds)
.fetch_optional(&mut *transaction)
.await?;
if let Some((lease_epoch,)) = claimed {
let current: Option<(i64, String)> = sqlx::query_as(
"SELECT sequence, job_id FROM later_partition_head \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3",
)
.bind(&self.namespace)
.bind(&topic)
.bind(partition_id)
.fetch_optional(&mut *transaction)
.await?;
if current != Some((sequence, job_id.clone())) {
transaction.rollback().await?;
return Ok(None);
}
transaction.commit().await?;
return Ok(Some(PartitionHeadClaim {
topic,
partition,
sequence: PartitionSequence(sequence),
job_id: crate::JobId(job_id),
owner: owner.to_string(),
lease_epoch,
}));
}
}
transaction.rollback().await?;
Ok(None)
}
async fn peek_assigned_partition_head(
&self,
owner: &str,
topic: &str,
partition: crate::topic::PartitionId,
lease_epoch: i64,
) -> anyhow::Result<Option<PartitionHeadClaim>> {
let row: Option<(i64, String)> = sqlx::query_as(
"SELECT h.sequence, h.job_id \
FROM later_partition_head h \
JOIN later_partition_lease l \
ON l.namespace = h.namespace AND l.topic = h.topic \
AND l.partition_id = h.partition_id \
WHERE h.namespace = $1 AND h.topic = $2 AND h.partition_id = $3 \
AND h.earliest_run_at <= NOW() \
AND l.owner = $4 AND l.lease_epoch = $5 AND l.lease_until > NOW()",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(owner)
.bind(lease_epoch)
.fetch_optional(&self.pool)
.await?;
Ok(row.map(|(sequence, job_id)| PartitionHeadClaim {
topic: topic.to_string(),
partition,
sequence: PartitionSequence(sequence),
job_id: crate::JobId(job_id),
owner: owner.to_string(),
lease_epoch,
}))
}
async fn admin_peek_partition_head(
&self,
topic: &str,
partition: crate::topic::PartitionId,
) -> anyhow::Result<Option<(crate::JobId, PartitionSequence)>> {
let row: Option<(i64, String)> = sqlx::query_as(
"SELECT sequence, job_id FROM later_partition_head \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.fetch_optional(&self.pool)
.await?;
Ok(row.map(|(sequence, job_id)| (crate::JobId(job_id), PartitionSequence(sequence))))
}
async fn admin_force_complete_partition_head(
&self,
topic: &str,
partition: crate::topic::PartitionId,
sequence: PartitionSequence,
job_id: &crate::JobId,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
let mut transaction = self.pool.begin().await?;
let deleted = sqlx::query(
"DELETE FROM later_partition_job \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3 AND sequence = $4 \
AND job_id = $5",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(sequence.0)
.bind(job_id.to_string())
.execute(&mut *transaction)
.await?;
if deleted.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
postgres_advance_partition_head(&mut transaction, &self.namespace, topic, partition)
.await?;
sqlx::query(
"UPDATE later_partition_lease SET lease_epoch = lease_epoch + 1, lease_until = NOW() \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.execute(&mut *transaction)
.await?;
apply_postgres_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(true)
}
async fn heartbeat_worker(&self, owner: &str) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_worker (namespace, owner, heartbeat_until) \
VALUES ($1, $2, NOW() + make_interval(secs => $3)) \
ON CONFLICT(namespace, owner) DO UPDATE SET heartbeat_until = EXCLUDED.heartbeat_until",
)
.bind(&self.namespace)
.bind(owner)
.bind(WORKER_HEARTBEAT_SECONDS)
.execute(&self.pool)
.await?;
Ok(())
}
async fn list_live_workers(&self) -> anyhow::Result<Vec<String>> {
let rows: Vec<(String,)> = sqlx::query_as(
"SELECT owner FROM later_worker WHERE namespace = $1 AND heartbeat_until > NOW()",
)
.bind(&self.namespace)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|(owner,)| owner).collect())
}
async fn deregister_worker(&self, owner: &str) -> anyhow::Result<()> {
sqlx::query("DELETE FROM later_worker WHERE namespace = $1 AND owner = $2")
.bind(&self.namespace)
.bind(owner)
.execute(&self.pool)
.await?;
Ok(())
}
async fn renew_partition_lease(&self, claim: &PartitionHeadClaim) -> anyhow::Result<bool> {
let result = sqlx::query(
"UPDATE later_partition_lease SET lease_until = NOW() + make_interval(secs => $1) \
WHERE namespace = $2 AND topic = $3 AND partition_id = $4 AND owner = $5 \
AND lease_epoch = $6 AND lease_until > NOW()",
)
.bind(PARTITION_LEASE_SECONDS)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(&claim.owner)
.bind(claim.lease_epoch)
.execute(&self.pool)
.await?;
Ok(result.rows_affected() == 1)
}
async fn release_partition_assignment(
&self,
owner: &str,
topic: &str,
partition: crate::topic::PartitionId,
lease_epoch: i64,
) -> anyhow::Result<bool> {
let result = sqlx::query(
"UPDATE later_partition_lease \
SET lease_epoch = lease_epoch + 1, lease_until = NOW() \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3 AND owner = $4 \
AND lease_epoch = $5 AND lease_until > NOW()",
)
.bind(&self.namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(owner)
.bind(lease_epoch)
.execute(&self.pool)
.await?;
Ok(result.rows_affected() == 1)
}
async fn release_partition_head(
&self,
claim: &PartitionHeadClaim,
earliest_run_at: crate::UtcDateTime,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
let mut transaction = self.pool.begin().await?;
let released = sqlx::query(
"UPDATE later_partition_lease SET lease_until = NOW() \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3 AND owner = $4 \
AND lease_epoch = $5 AND lease_until > NOW()",
)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(&claim.owner)
.bind(claim.lease_epoch)
.execute(&mut *transaction)
.await?;
if released.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
let updated = sqlx::query(
"UPDATE later_partition_job SET earliest_run_at = $1 \
WHERE namespace = $2 AND topic = $3 AND partition_id = $4 AND sequence = $5 \
AND job_id = $6",
)
.bind(earliest_run_at)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(claim.sequence.0)
.bind(claim.job_id.to_string())
.execute(&mut *transaction)
.await?;
if updated.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
let head_updated = sqlx::query(
"UPDATE later_partition_head SET earliest_run_at = $1 \
WHERE namespace = $2 AND topic = $3 AND partition_id = $4 AND sequence = $5 \
AND job_id = $6",
)
.bind(earliest_run_at)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(claim.sequence.0)
.bind(claim.job_id.to_string())
.execute(&mut *transaction)
.await?;
if head_updated.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
apply_postgres_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(true)
}
async fn complete_partition_head(
&self,
claim: &PartitionHeadClaim,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<bool> {
let mut transaction = self.pool.begin().await?;
let deleted = sqlx::query(
"DELETE FROM later_partition_job \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3 AND sequence = $4 \
AND job_id = $5 \
AND EXISTS ( \
SELECT 1 FROM later_partition_lease \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3 \
AND owner = $6 AND lease_epoch = $7 AND lease_until > NOW() \
)",
)
.bind(&self.namespace)
.bind(&claim.topic)
.bind(i64::from(claim.partition.0))
.bind(claim.sequence.0)
.bind(claim.job_id.to_string())
.bind(&claim.owner)
.bind(claim.lease_epoch)
.execute(&mut *transaction)
.await?;
if deleted.rows_affected() == 0 {
transaction.rollback().await?;
return Ok(false);
}
postgres_advance_partition_head(
&mut transaction,
&self.namespace,
&claim.topic,
claim.partition,
)
.await?;
apply_postgres_operations(&mut transaction, operations).await?;
transaction.commit().await?;
Ok(true)
}
}
#[cfg(feature = "sqlite")]
async fn insert_sqlite_delivery_intent(
connection: &mut sqlx::SqliteConnection,
namespace: &str,
queue_name: &str,
payload: &[u8],
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_delivery_queue \
(namespace, queue_name, payload, available_at) VALUES (?1, ?2, ?3, ?4)",
)
.bind(namespace)
.bind(queue_name)
.bind(payload)
.bind(chrono::Utc::now().timestamp_millis())
.execute(connection)
.await?;
Ok(())
}
#[cfg(feature = "postgres")]
const PARTITION_LEASE_SECONDS: i64 = 30;
#[cfg(feature = "postgres")]
const WORKER_HEARTBEAT_SECONDS: i64 = 15;
#[cfg(feature = "postgres")]
async fn postgres_topic_partition_count(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
namespace: &str,
topic: &str,
) -> anyhow::Result<u32> {
let partition_count: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = $1 AND topic = $2",
)
.bind(namespace)
.bind(topic)
.fetch_optional(&mut **transaction)
.await?;
let Some((partition_count,)) = partition_count else {
return Err(anyhow::anyhow!("unknown topic '{topic}'"));
};
Ok(u32::try_from(partition_count)?)
}
#[cfg(feature = "postgres")]
async fn postgres_assign_partition_sequence(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
namespace: &str,
topic: &str,
partition: crate::topic::PartitionId,
) -> anyhow::Result<PartitionSequence> {
let sequence: Option<(i64,)> = sqlx::query_as(
"INSERT INTO later_partition_sequence (namespace, topic, partition_id, next_sequence) \
SELECT $1, $2, $3, 2 \
WHERE EXISTS ( \
SELECT 1 FROM later_topic \
WHERE namespace = $1 AND topic = $2 AND partition_count > $3 \
) \
ON CONFLICT(namespace, topic, partition_id) DO UPDATE SET \
next_sequence = later_partition_sequence.next_sequence + 1 \
RETURNING next_sequence - 1",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.fetch_optional(&mut **transaction)
.await?;
let Some((sequence,)) = sequence else {
return Err(anyhow::anyhow!(
"unknown topic '{topic}' or partition {} is out of range",
partition.0
));
};
Ok(PartitionSequence(sequence))
}
#[cfg(feature = "postgres")]
#[allow(clippy::too_many_arguments)]
async fn postgres_insert_partition_job(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
namespace: &str,
topic: &str,
partition: crate::topic::PartitionId,
sequence: PartitionSequence,
job_id: &crate::JobId,
payload_type: &str,
earliest_run_at: crate::UtcDateTime,
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_partition_job \
(namespace, topic, partition_id, sequence, job_id, payload_type, earliest_run_at) \
VALUES ($1, $2, $3, $4, $5, $6, $7)",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(sequence.0)
.bind(job_id.to_string())
.bind(payload_type)
.bind(earliest_run_at)
.execute(&mut **transaction)
.await?;
sqlx::query(
"INSERT INTO later_partition_head \
(namespace, topic, partition_id, sequence, job_id, earliest_run_at) \
VALUES ($1, $2, $3, $4, $5, $6) \
ON CONFLICT(namespace, topic, partition_id) DO NOTHING",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.bind(sequence.0)
.bind(job_id.to_string())
.bind(earliest_run_at)
.execute(&mut **transaction)
.await?;
Ok(())
}
#[cfg(feature = "postgres")]
async fn postgres_advance_partition_head(
transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
namespace: &str,
topic: &str,
partition: crate::topic::PartitionId,
) -> anyhow::Result<()> {
sqlx::query(
"DELETE FROM later_partition_head \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.execute(&mut **transaction)
.await?;
sqlx::query(
"INSERT INTO later_partition_head \
(namespace, topic, partition_id, sequence, job_id, earliest_run_at) \
SELECT namespace, topic, partition_id, sequence, job_id, earliest_run_at \
FROM later_partition_job \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3 \
ORDER BY sequence ASC LIMIT 1",
)
.bind(namespace)
.bind(topic)
.bind(i64::from(partition.0))
.execute(&mut **transaction)
.await?;
Ok(())
}
#[cfg(feature = "postgres")]
async fn insert_postgres_delivery_intent(
connection: &mut sqlx::PgConnection,
namespace: &str,
queue_name: &str,
payload: &[u8],
) -> anyhow::Result<()> {
sqlx::query(
"INSERT INTO later_delivery_queue (namespace, queue_name, payload) \
VALUES ($1, $2, $3)",
)
.bind(namespace)
.bind(queue_name)
.bind(payload)
.execute(connection)
.await?;
Ok(())
}
#[cfg(feature = "sqlite")]
pub(crate) async fn apply_sqlite_operations(
connection: &mut sqlx::SqliteConnection,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<()> {
for operation in operations {
match operation {
crate::storage::StorageOperation::Set { key, value } => {
let now = chrono::Utc::now().timestamp_millis();
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",
)
.bind(key)
.bind(value)
.bind(now)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeAdd { key, value } => {
sqlx::query(
"INSERT INTO later_storage_range (range_key, value) VALUES (?1, ?2) \
ON CONFLICT(range_key, value) DO UPDATE SET date_expire = NULL",
)
.bind(key)
.bind(value)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::Delete { key } => {
sqlx::query("DELETE FROM later_storage WHERE key = ?1")
.bind(key)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::Expire { key, ttl_seconds } => {
let ttl_millis = i64::try_from(ttl_seconds)?
.checked_mul(1_000)
.ok_or_else(|| anyhow::anyhow!("SQLite TTL is too large"))?;
let expires_at = chrono::Utc::now()
.timestamp_millis()
.checked_add(ttl_millis)
.ok_or_else(|| anyhow::anyhow!("SQLite expiry is out of range"))?;
sqlx::query("UPDATE later_storage SET date_expire = ?2 WHERE key = ?1")
.bind(key)
.bind(expires_at)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeRemove { key, value } => {
sqlx::query("DELETE FROM later_storage_range WHERE range_key = ?1 AND value = ?2")
.bind(key)
.bind(value)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeExpire {
key,
value,
ttl_seconds,
} => {
let ttl_millis = i64::try_from(ttl_seconds)?
.checked_mul(1_000)
.ok_or_else(|| anyhow::anyhow!("SQLite range TTL is too large"))?;
let expires_at = chrono::Utc::now()
.timestamp_millis()
.checked_add(ttl_millis)
.ok_or_else(|| anyhow::anyhow!("SQLite range expiry is out of range"))?;
sqlx::query(
"UPDATE later_storage_range SET date_expire = ?3 \
WHERE range_key = ?1 AND value = ?2",
)
.bind(key)
.bind(value)
.bind(expires_at)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeClear { key } => {
sqlx::query("DELETE FROM later_storage_range WHERE range_key = ?1")
.bind(key)
.execute(&mut *connection)
.await?;
}
}
}
Ok(())
}
#[cfg(feature = "postgres")]
pub(crate) async fn apply_postgres_operations(
connection: &mut sqlx::PgConnection,
operations: Vec<crate::storage::StorageOperation>,
) -> anyhow::Result<()> {
if let [crate::storage::StorageOperation::Set { key, value }, crate::storage::StorageOperation::RangeAdd {
key: range_key,
value: range_value,
}] = operations.as_slice()
{
sqlx::query(
"WITH saved AS ( \
INSERT INTO later_storage (key, value, update_count, date_created) \
VALUES ($1, $2, 0, NOW()) \
ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, \
date_updated = NOW(), update_count = later_storage.update_count + 1, \
date_expire = NULL \
RETURNING 1 \
) \
INSERT INTO later_storage_range (range_key, value) \
SELECT $3, $4 FROM saved \
ON CONFLICT (range_key, value) DO UPDATE SET date_expire = NULL",
)
.bind(key)
.bind(value)
.bind(range_key)
.bind(range_value)
.execute(&mut *connection)
.await?;
return Ok(());
}
for operation in operations {
match operation {
crate::storage::StorageOperation::Set { key, value } => {
sqlx::query(
"INSERT INTO later_storage (key, value, update_count, date_created) \
VALUES ($1, $2, 0, NOW()) \
ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, \
date_updated = NOW(), update_count = later_storage.update_count + 1, \
date_expire = NULL",
)
.bind(key)
.bind(value)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeAdd { key, value } => {
sqlx::query(
"INSERT INTO later_storage_range (range_key, value) VALUES ($1, $2) \
ON CONFLICT (range_key, value) DO UPDATE SET date_expire = NULL",
)
.bind(key)
.bind(value)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::Delete { key } => {
sqlx::query("DELETE FROM later_storage WHERE key = $1")
.bind(key)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::Expire { key, ttl_seconds } => {
sqlx::query(
"UPDATE later_storage SET \
date_expire = NOW() + make_interval(secs => $2) WHERE key = $1",
)
.bind(key)
.bind(i64::try_from(ttl_seconds)?)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeRemove { key, value } => {
sqlx::query("DELETE FROM later_storage_range WHERE range_key = $1 AND value = $2")
.bind(key)
.bind(value)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeExpire {
key,
value,
ttl_seconds,
} => {
sqlx::query(
"UPDATE later_storage_range SET \
date_expire = NOW() + make_interval(secs => $3) \
WHERE range_key = $1 AND value = $2",
)
.bind(key)
.bind(value)
.bind(i64::try_from(ttl_seconds)?)
.execute(&mut *connection)
.await?;
}
crate::storage::StorageOperation::RangeClear { key } => {
sqlx::query("DELETE FROM later_storage_range WHERE range_key = $1")
.bind(key)
.execute(&mut *connection)
.await?;
}
}
}
Ok(())
}
#[cfg(all(test, feature = "sqlite"))]
mod sqlite_tests {
use super::{Backend, PartitionSequence, SqliteBackend};
use crate::core::JobParameter;
use serde::{Deserialize, Serialize};
use sqlx::Row;
#[derive(Deserialize, Serialize)]
struct TestJob(String);
impl JobParameter for TestJob {
fn to_bytes(&self) -> anyhow::Result<Vec<u8>> {
Ok(rmp_serde::to_vec(self)?)
}
fn from_bytes(payload: &[u8]) -> Self {
rmp_serde::from_slice(payload).unwrap_or_else(|_| Self(String::new()))
}
fn get_ptype(&self) -> String {
"TestJob".to_owned()
}
}
#[tokio::test]
async fn sqlite_caller_transaction_commits_or_rolls_back_app_data_and_job() -> anyhow::Result<()>
{
let backend = SqliteBackend::connect("atomic-sqlite", "sqlite::memory:").await?;
sqlx::query("CREATE TABLE application_data (value TEXT NOT NULL)")
.execute(backend.pool())
.await?;
let mut rolled_back = backend.pool().begin().await?;
sqlx::query("INSERT INTO application_data (value) VALUES ('rolled-back')")
.execute(&mut *rolled_back)
.await?;
backend
.enqueue_in(&mut rolled_back, TestJob("rolled-back".to_owned()))
.await?;
rolled_back.rollback().await?;
let mut committed = backend.pool().begin().await?;
sqlx::query("INSERT INTO application_data (value) VALUES ('committed')")
.execute(&mut *committed)
.await?;
backend
.enqueue_in(&mut committed, TestJob("committed".to_owned()))
.await?;
committed.commit().await?;
let app_count: i64 = sqlx::query("SELECT COUNT(*) AS count FROM application_data")
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let job_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_storage \
WHERE key LIKE 'later-atomic-sqlite%' AND key NOT LIKE '%-stats-%'",
)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let delivery_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_delivery_queue \
WHERE namespace = 'atomic-sqlite'",
)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
assert_eq!((app_count, job_count, delivery_count), (1, 1, 1));
Ok(())
}
#[derive(Deserialize, Serialize, Clone)]
struct TestPartitionJob(String);
impl JobParameter for TestPartitionJob {
fn to_bytes(&self) -> anyhow::Result<Vec<u8>> {
Ok(rmp_serde::to_vec(self)?)
}
fn from_bytes(payload: &[u8]) -> Self {
rmp_serde::from_slice(payload).unwrap_or_else(|_| Self(String::new()))
}
fn get_ptype(&self) -> String {
"TestPartitionJob".to_owned()
}
}
impl crate::topic::JobPartition for TestPartitionJob {
fn partition_key(&self) -> crate::topic::PartitionKey {
crate::topic::PartitionKey::from("test-partition-job")
}
}
impl crate::topic::JobTopic for TestPartitionJob {
fn topic_name() -> &'static str {
"orders"
}
}
#[tokio::test]
async fn sqlite_partition_caller_transaction_commits_or_rolls_back_app_data_job_and_sequence(
) -> anyhow::Result<()> {
let backend = SqliteBackend::connect("atomic-partition-sqlite", "sqlite::memory:").await?;
Box::new(backend.clone())
.into_parts()?
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 4)?])
.await?;
sqlx::query("CREATE TABLE application_data (value TEXT NOT NULL)")
.execute(backend.pool())
.await?;
let mut rolled_back = backend.pool().begin().await?;
sqlx::query("INSERT INTO application_data (value) VALUES ('rolled-back')")
.execute(&mut *rolled_back)
.await?;
backend
.enqueue_to_partition_in(&mut rolled_back, TestPartitionJob("rolled-back".to_owned()))
.await?;
rolled_back.rollback().await?;
let mut committed = backend.pool().begin().await?;
sqlx::query("INSERT INTO application_data (value) VALUES ('committed')")
.execute(&mut *committed)
.await?;
backend
.enqueue_to_partition_in(&mut committed, TestPartitionJob("committed".to_owned()))
.await?;
committed.commit().await?;
let app_count: i64 = sqlx::query("SELECT COUNT(*) AS count FROM application_data")
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let job_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_storage \
WHERE key LIKE 'later-atomic-partition-sqlite%' AND key NOT LIKE '%-stats-%'",
)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let partition_job_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_partition_job \
WHERE namespace = 'atomic-partition-sqlite'",
)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let partition_head_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_partition_head \
WHERE namespace = 'atomic-partition-sqlite'",
)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let next_sequence: i64 = sqlx::query(
"SELECT next_sequence FROM later_partition_sequence \
WHERE namespace = 'atomic-partition-sqlite'",
)
.fetch_one(backend.pool())
.await?
.try_get("next_sequence")?;
assert_eq!(
(
app_count,
job_count,
partition_job_count,
partition_head_count
),
(1, 1, 1, 1)
);
assert_eq!(
next_sequence, 2,
"the rolled-back enqueue must not leave a gap in the partition sequence"
);
Ok(())
}
#[tokio::test]
async fn sqlite_job_claim_is_exclusive_until_finished() -> anyhow::Result<()> {
let backend = SqliteBackend::connect("lease-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend).into_parts()?;
let job_id = crate::JobId("sqlite-lease-job".to_owned());
let first = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("first SQLite lease was not acquired"))?;
assert!(parts.committer.claim(&job_id).await?.is_none());
first.finish().await?;
let after_finish = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("finished SQLite lease was not released"))?;
after_finish.finish().await?;
Ok(())
}
#[tokio::test]
async fn sqlite_job_lease_loss_is_signaled_after_it_is_stolen() -> anyhow::Result<()> {
let backend = SqliteBackend::connect("lease-lost-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
let job_id = crate::JobId("sqlite-lease-lost-job".to_owned());
let lease = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
let mut lease_lost = lease.lease_lost();
assert!(!lease_lost.has_changed().unwrap_or(true));
sqlx::query(
"UPDATE later_job_execution SET lease_owner = 'someone-else' \
WHERE namespace = 'lease-lost-sqlite' AND job_id = ?1",
)
.bind(job_id.to_string())
.execute(backend.pool())
.await?;
tokio::time::timeout(std::time::Duration::from_secs(15), lease_lost.changed())
.await
.map_err(|_| anyhow::anyhow!("lease_lost did not fire within 15s of being stolen"))??;
Ok(())
}
#[tokio::test]
async fn sqlite_stale_leased_job_ids_finds_leases_past_the_grace_period() -> anyhow::Result<()>
{
let backend = SqliteBackend::connect("stale-lease-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
let grace = std::time::Duration::from_secs(30);
let abandoned_id = crate::JobId("sqlite-abandoned-job".to_owned());
parts.committer.claim(&abandoned_id).await?;
let long_expired = chrono::Utc::now().timestamp_millis() - (grace.as_millis() as i64) * 2;
sqlx::query(
"UPDATE later_job_execution SET lease_until = ?1 \
WHERE namespace = 'stale-lease-sqlite' AND job_id = ?2",
)
.bind(long_expired)
.bind(abandoned_id.to_string())
.execute(backend.pool())
.await?;
let recent_id = crate::JobId("sqlite-recently-expired-job".to_owned());
parts.committer.claim(&recent_id).await?;
let just_expired = chrono::Utc::now().timestamp_millis() - 1_000;
sqlx::query(
"UPDATE later_job_execution SET lease_until = ?1 \
WHERE namespace = 'stale-lease-sqlite' AND job_id = ?2",
)
.bind(just_expired)
.bind(recent_id.to_string())
.execute(backend.pool())
.await?;
let live_id = crate::JobId("sqlite-live-job".to_owned());
parts.committer.claim(&live_id).await?;
let stale = parts.committer.stale_leased_job_ids(grace, 100).await?;
assert_eq!(stale, vec![abandoned_id]);
Ok(())
}
#[tokio::test]
async fn sqlite_redelivery_does_not_wipe_an_abandoned_running_jobs_lease() -> anyhow::Result<()>
{
let backend = SqliteBackend::connect("redelivery-safety-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
let job_id = crate::JobId("sqlite-abandoned-running-job".to_owned());
let committer = parts.committer.clone();
let publisher = crate::BackgroundJobServerPublisher::new(
parts.namespace.clone(),
std::sync::Arc::new(parts.delivery),
parts.storage,
parts.committer,
)
.await?;
publisher
.storage
.save_job(&crate::models::Job {
id: job_id.clone(),
payload_type: "TestJob".to_owned(),
payload: Vec::new(),
config: crate::models::JobConfig::default(),
stage: crate::models::Stage::Running(crate::models::RunningStage {
date: chrono::Utc::now() - chrono::Duration::hours(1),
}),
previous_stages: Vec::new(),
recurring_job_id: None,
topic: None,
partition: None,
sequence: None,
})
.await?;
committer.claim(&job_id).await?;
sqlx::query(
"UPDATE later_job_execution SET lease_until = ?1 \
WHERE namespace = 'redelivery-safety-sqlite' AND job_id = ?2",
)
.bind(chrono::Utc::now().timestamp_millis() - 60_000)
.bind(job_id.to_string())
.execute(backend.pool())
.await?;
struct FakeHandler {
publisher: crate::BackgroundJobServerPublisher,
ctx: (),
}
#[async_trait::async_trait]
impl crate::core::BgJobHandler<()> for FakeHandler {
fn get_ctx(&self) -> &() {
&self.ctx
}
fn get_publisher(&self) -> &crate::BackgroundJobServerPublisher {
&self.publisher
}
async fn dispatch(
&self,
_ptype: String,
_payload: &[u8],
_job_id: crate::JobId,
) -> anyhow::Result<()> {
Err(anyhow::anyhow!(
"a job already Running from a prior attempt must not be dispatched again"
))
}
}
let handler = std::sync::Arc::new(FakeHandler { publisher, ctx: () });
let (tx, _rx) = async_std::channel::unbounded();
crate::commands::handle_amqp_command(
crate::models::AmqpCommand::ExecuteJob(crate::models::JobAmqp {
payload_type: "TestJob".to_owned(),
id: job_id.clone(),
}),
1,
#[cfg(feature = "prometheus")]
"worker-1",
&handler,
&tx,
None,
)
.await?;
let still_present: Option<(String,)> = sqlx::query_as(
"SELECT job_id FROM later_job_execution \
WHERE namespace = 'redelivery-safety-sqlite' AND job_id = ?1",
)
.bind(job_id.to_string())
.fetch_optional(backend.pool())
.await?;
assert!(
still_present.is_some(),
"redelivery of an already-Running job must not finish (delete) its execution lease"
);
let job = handler
.publisher
.storage
.get_job(job_id)
.await?
.expect("job still present");
assert!(matches!(job.stage, crate::models::Stage::Running(_)));
Ok(())
}
#[cfg(feature = "dashboard")]
#[tokio::test]
async fn sqlite_poll_stuck_jobs_repairs_running_index_rows_without_a_lease(
) -> anyhow::Result<()> {
use crate::models::{Job, JobConfig, RunningStage, Stage};
let backend = SqliteBackend::connect("running-index-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
let publisher = crate::BackgroundJobServerPublisher::new(
parts.namespace.clone(),
std::sync::Arc::new(parts.delivery),
parts.storage,
parts.committer,
)
.await?;
let old = chrono::Utc::now() - chrono::Duration::hours(2);
let running_job = |id: &str, topic: Option<&str>| Job {
id: crate::JobId(id.to_owned()),
payload_type: "TestJob".to_owned(),
payload: Vec::new(),
config: JobConfig::default(),
stage: Stage::Running(RunningStage { date: old }),
previous_stages: Vec::new(),
recurring_job_id: None,
topic: topic.map(str::to_owned),
partition: topic.map(|_| 0),
sequence: topic.map(|_| 1),
};
publisher.save(&running_job("lost-lease", None)).await?;
publisher
.save(&running_job("partitioned", Some("topic")))
.await?;
publisher.save(&running_job("orphan", None)).await?;
sqlx::query("DELETE FROM later_storage WHERE key LIKE '%orphan'")
.execute(backend.pool())
.await?;
struct FakeHandler {
publisher: crate::BackgroundJobServerPublisher,
ctx: (),
}
#[async_trait::async_trait]
impl crate::core::BgJobHandler<()> for FakeHandler {
fn get_ctx(&self) -> &() {
&self.ctx
}
fn get_publisher(&self) -> &crate::BackgroundJobServerPublisher {
&self.publisher
}
async fn dispatch(&self, _: String, _: &[u8], _: crate::JobId) -> anyhow::Result<()> {
Ok(())
}
}
let handler = std::sync::Arc::new(FakeHandler { publisher, ctx: () });
crate::commands::handle_poll_stuck_jobs_command(handler.clone()).await?;
let stage_of = |id: &'static str| {
let handler = handler.clone();
async move {
handler
.publisher
.storage
.get_job(crate::JobId(id.to_owned()))
.await
.map(|job| job.map(|job| job.stage.get_name().to_owned()))
}
};
assert_ne!(stage_of("lost-lease").await?.as_deref(), Some("running"));
assert_eq!(stage_of("partitioned").await?.as_deref(), Some("running"));
let remaining: Vec<(String,)> = sqlx::query_as(
"SELECT job_id FROM later_jobs_index WHERE stage = 'running' ORDER BY job_id",
)
.fetch_all(backend.pool())
.await?;
assert_eq!(remaining, vec![("partitioned".to_owned(),)]);
Ok(())
}
#[tokio::test]
async fn sqlite_claim_and_fetch_matches_a_separate_claim_and_get() -> anyhow::Result<()> {
let backend = SqliteBackend::connect("claim-and-fetch-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend).into_parts()?;
let job_id = crate::JobId("sqlite-claim-and-fetch-job".to_owned());
let present_key = "job-present";
let missing_key = "job-missing";
parts.storage.set(present_key, b"job-bytes").await?;
let claimed = parts
.committer
.claim_and_fetch(&job_id, present_key, parts.storage.as_ref())
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
assert_eq!(claimed.value.as_deref(), Some(b"job-bytes".as_slice()));
assert!(parts.committer.claim(&job_id).await?.is_none());
claimed.lease.finish().await?;
let claimed_missing = parts
.committer
.claim_and_fetch(&job_id, missing_key, parts.storage.as_ref())
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
assert!(claimed_missing.value.is_none());
claimed_missing.lease.finish().await?;
let held = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
assert!(parts
.committer
.claim_and_fetch(&job_id, present_key, parts.storage.as_ref())
.await?
.is_none());
held.finish().await?;
Ok(())
}
#[tokio::test]
async fn sqlite_worker_heartbeat_lists_live_workers_until_deregistered() -> anyhow::Result<()> {
let backend = SqliteBackend::connect("worker-heartbeat-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend).into_parts()?;
assert!(parts.committer.list_live_workers().await?.is_empty());
parts.committer.heartbeat_worker("worker-a").await?;
parts.committer.heartbeat_worker("worker-b").await?;
let mut live = parts.committer.list_live_workers().await?;
live.sort();
assert_eq!(live, vec!["worker-a".to_string(), "worker-b".to_string()]);
parts.committer.deregister_worker("worker-a").await?;
assert_eq!(
parts.committer.list_live_workers().await?,
vec!["worker-b".to_string()]
);
Ok(())
}
#[tokio::test]
async fn sqlite_partition_head_moves_to_another_worker_only_after_lease_expiry(
) -> anyhow::Result<()> {
let backend = SqliteBackend::connect("crash-recovery-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
let job_id = crate::JobId("crash-job".to_owned());
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&job_id,
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let claim = parts
.committer
.claim_partition_head("worker-a", &["worker-a".to_string()])
.await?
.ok_or_else(|| anyhow::anyhow!("first worker did not claim the partition head"))?;
assert_eq!(claim.job_id, job_id);
assert_eq!(claim.owner, "worker-a");
assert!(parts
.committer
.claim_partition_head("worker-b", &["worker-b".to_string()])
.await?
.is_none());
sqlx::query(
"UPDATE later_partition_lease SET lease_until = 0 \
WHERE namespace = 'crash-recovery-sqlite' AND topic = 'orders' AND partition_id = 0",
)
.execute(backend.pool())
.await?;
let recovered = parts
.committer
.claim_partition_head("worker-b", &["worker-b".to_string()])
.await?
.ok_or_else(|| {
anyhow::anyhow!("the expired lease was not reclaimed by another worker")
})?;
assert_eq!(recovered.job_id, job_id);
assert_eq!(recovered.sequence.0, claim.sequence.0);
assert_eq!(recovered.owner, "worker-b");
Ok(())
}
#[tokio::test]
async fn sqlite_completing_a_head_keeps_a_valid_assignment() -> anyhow::Result<()> {
let namespace = "partition-epoch-sqlite";
let backend = SqliteBackend::connect(namespace, "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
for id in ["first-job", "second-job"] {
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId(id.to_owned()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
}
let same_owner = "shared-owner";
let first_claim = parts
.committer
.claim_partition_head(same_owner, &[same_owner.to_owned()])
.await?
.ok_or_else(|| anyhow::anyhow!("expected to claim the first job"))?;
assert!(
parts
.committer
.complete_partition_head(&first_claim, Vec::new())
.await?
);
let second_claim = parts
.committer
.peek_assigned_partition_head(
same_owner,
"orders",
crate::topic::PartitionId(0),
first_claim.lease_epoch,
)
.await?
.ok_or_else(|| anyhow::anyhow!("expected to keep the second head assignment"))?;
assert_eq!(
second_claim.lease_epoch, first_claim.lease_epoch,
"a valid assignment must be reused across completions"
);
assert!(
!parts
.committer
.complete_partition_head(&first_claim, Vec::new())
.await?,
"a stale same-owner claim from an earlier job must not complete a later one"
);
assert!(
parts
.committer
.complete_partition_head(&second_claim, Vec::new())
.await?
);
let lease_rows: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM later_partition_lease \
WHERE namespace = ?1 AND topic = 'orders' AND partition_id = 0",
)
.bind(namespace)
.fetch_one(backend.pool())
.await?;
assert_eq!(
lease_rows, 1,
"the lease row must persist across completions, not be deleted"
);
Ok(())
}
#[tokio::test]
async fn sqlite_stale_partition_claim_cannot_complete_a_reclaimed_head() -> anyhow::Result<()> {
let namespace = "partition-fence-sqlite";
let backend = SqliteBackend::connect(namespace, "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
let job_id = crate::JobId("fenced-job".to_owned());
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&job_id,
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let first = parts
.committer
.claim_partition_head("worker-a", &["worker-a".to_owned()])
.await?
.ok_or_else(|| anyhow::anyhow!("first worker did not claim the partition head"))?;
sqlx::query(
"UPDATE later_partition_lease SET lease_until = 0 \
WHERE namespace = ?1 AND topic = 'orders' AND partition_id = 0",
)
.bind(namespace)
.execute(backend.pool())
.await?;
let reclaimed = parts
.committer
.claim_partition_head("worker-b", &["worker-b".to_owned()])
.await?
.ok_or_else(|| anyhow::anyhow!("second worker did not reclaim the partition head"))?;
assert!(reclaimed.lease_epoch > first.lease_epoch);
let stale_release_key = "stale-release-fence";
assert!(
!parts
.committer
.release_partition_head(
&first,
chrono::Utc::now(),
vec![crate::storage::StorageOperation::Set {
key: stale_release_key.to_owned(),
value: b"stale".to_vec(),
}],
)
.await?,
"a stale claim must not save retry state for a reclaimed head"
);
let stale_release_writes: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM later_storage WHERE key = ?1")
.bind(stale_release_key)
.fetch_one(backend.pool())
.await?;
assert_eq!(stale_release_writes, 0);
let stale_completion_key = "stale-completion-fence";
assert!(
!parts
.committer
.complete_partition_head(
&first,
vec![crate::storage::StorageOperation::Set {
key: stale_completion_key.to_owned(),
value: b"stale".to_vec(),
}],
)
.await?,
"a stale claim must not remove a head reclaimed at a newer lease epoch"
);
let stale_completion_writes: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM later_storage WHERE key = ?1")
.bind(stale_completion_key)
.fetch_one(backend.pool())
.await?;
assert_eq!(stale_completion_writes, 0);
assert_eq!(
parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?,
Some((job_id.clone(), first.sequence))
);
assert!(
parts
.committer
.complete_partition_head(&reclaimed, Vec::new())
.await?
);
assert!(parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?
.is_none());
Ok(())
}
#[tokio::test]
async fn sqlite_admin_force_complete_partition_head_unblocks_a_poison_head(
) -> anyhow::Result<()> {
let backend = SqliteBackend::connect("admin-force-fail-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
let poison_id = crate::JobId("poison-job".to_owned());
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&poison_id,
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId("next-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let (peeked_id, sequence) = parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?
.ok_or_else(|| anyhow::anyhow!("expected a head to peek"))?;
assert_eq!(peeked_id, poison_id);
let stale_rejected = parts
.committer
.admin_force_complete_partition_head(
"orders",
crate::topic::PartitionId(0),
PartitionSequence(sequence.0 + 1),
&poison_id,
Vec::new(),
)
.await?;
assert!(!stale_rejected);
let cleared = parts
.committer
.admin_force_complete_partition_head(
"orders",
crate::topic::PartitionId(0),
sequence,
&poison_id,
Vec::new(),
)
.await?;
assert!(cleared);
let lease_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_partition_lease \
WHERE namespace = 'admin-force-fail-sqlite' AND topic = 'orders' AND partition_id = 0",
)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
assert_eq!(lease_count, 0);
let (new_head_id, _) = parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?
.ok_or_else(|| anyhow::anyhow!("expected the next job to become the head"))?;
assert_eq!(new_head_id, crate::JobId("next-job".to_owned()));
let repeated = parts
.committer
.admin_force_complete_partition_head(
"orders",
crate::topic::PartitionId(0),
sequence,
&poison_id,
Vec::new(),
)
.await?;
assert!(!repeated);
Ok(())
}
#[tokio::test]
async fn sqlite_partition_backlog_summary_counts_ready_heads_and_ignores_unready_ones(
) -> anyhow::Result<()> {
let backend = SqliteBackend::connect("backlog-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 2)?])
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId("ready-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now() - chrono::Duration::seconds(30),
)
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(1),
&crate::JobId("future-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now() + chrono::Duration::hours(1),
)
.await?;
let summary = parts.committer.partition_backlog_summary().await?;
assert_eq!(summary.len(), 1);
let (topic, ready_count, oldest_age_seconds) = &summary[0];
assert_eq!(topic, "orders");
assert_eq!(*ready_count, 1);
assert!(
*oldest_age_seconds >= 30.0,
"expected the ready head's age to be at least 30s, got {oldest_age_seconds}"
);
Ok(())
}
#[tokio::test]
async fn sqlite_partition_commit_enqueues_a_wake_up_notification_on_the_sql_queue(
) -> anyhow::Result<()> {
let backend = SqliteBackend::connect("partition-wake-sqlite", "sqlite::memory:").await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId("wake-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let payload: Vec<u8> = sqlx::query(
"SELECT payload FROM later_delivery_queue WHERE namespace = 'partition-wake-sqlite'",
)
.fetch_one(backend.pool())
.await?
.try_get("payload")?;
let command: crate::models::AmqpCommand = crate::encoder::decode(&payload)?;
assert!(
matches!(&command, crate::models::AmqpCommand::PartitionReady { topic } if topic == "orders"),
"expected a PartitionReady notification for topic 'orders', got {command:?}"
);
Ok(())
}
}
#[cfg(all(test, feature = "postgres"))]
mod postgres_tests {
use super::{Backend, PartitionSequence, PostgresBackend};
use crate::core::JobParameter;
use serde::{Deserialize, Serialize};
use sqlx::Row;
#[derive(Deserialize, Serialize)]
struct TestJob(String);
impl JobParameter for TestJob {
fn to_bytes(&self) -> anyhow::Result<Vec<u8>> {
Ok(rmp_serde::to_vec(self)?)
}
fn from_bytes(payload: &[u8]) -> Self {
rmp_serde::from_slice(payload).unwrap_or_else(|_| Self(String::new()))
}
fn get_ptype(&self) -> String {
"TestJob".to_owned()
}
}
#[tokio::test]
async fn postgres_caller_transaction_commits_or_rolls_back_app_data_and_job(
) -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("atomic-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS later_application_data_atomic_test (\
namespace TEXT NOT NULL, value TEXT NOT NULL)",
)
.execute(backend.pool())
.await?;
let mut rolled_back = backend.pool().begin().await?;
sqlx::query(
"INSERT INTO later_application_data_atomic_test (namespace, value) VALUES ($1, 'rolled-back')",
)
.bind(&namespace)
.execute(&mut *rolled_back)
.await?;
backend
.enqueue_in(&mut rolled_back, TestJob("rolled-back".to_owned()))
.await?;
rolled_back.rollback().await?;
let mut committed = backend.pool().begin().await?;
sqlx::query(
"INSERT INTO later_application_data_atomic_test (namespace, value) VALUES ($1, 'committed')",
)
.bind(&namespace)
.execute(&mut *committed)
.await?;
backend
.enqueue_in(&mut committed, TestJob("committed".to_owned()))
.await?;
committed.commit().await?;
let job_key_prefix = format!("later-{namespace}%");
let app_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_application_data_atomic_test WHERE namespace = $1",
)
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let job_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_storage \
WHERE key LIKE $1 AND key NOT LIKE '%-stats-%'",
)
.bind(job_key_prefix)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let delivery_count: i64 =
sqlx::query("SELECT COUNT(*) AS count FROM later_delivery_queue WHERE namespace = $1")
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
assert_eq!((app_count, job_count, delivery_count), (1, 1, 1));
Ok(())
}
#[derive(Deserialize, Serialize, Clone)]
struct TestPartitionJob(String);
impl JobParameter for TestPartitionJob {
fn to_bytes(&self) -> anyhow::Result<Vec<u8>> {
Ok(rmp_serde::to_vec(self)?)
}
fn from_bytes(payload: &[u8]) -> Self {
rmp_serde::from_slice(payload).unwrap_or_else(|_| Self(String::new()))
}
fn get_ptype(&self) -> String {
"TestPartitionJob".to_owned()
}
}
impl crate::topic::JobPartition for TestPartitionJob {
fn partition_key(&self) -> crate::topic::PartitionKey {
crate::topic::PartitionKey::from("test-partition-job")
}
}
impl crate::topic::JobTopic for TestPartitionJob {
fn topic_name() -> &'static str {
"orders"
}
}
#[tokio::test]
async fn postgres_partition_caller_transaction_commits_or_rolls_back_app_data_job_and_sequence(
) -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("atomic-partition-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
Box::new(backend.clone())
.into_parts()?
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 4)?])
.await?;
sqlx::query(
"CREATE TABLE IF NOT EXISTS later_application_data_atomic_partition_test (\
namespace TEXT NOT NULL, value TEXT NOT NULL)",
)
.execute(backend.pool())
.await?;
let mut rolled_back = backend.pool().begin().await?;
sqlx::query(
"INSERT INTO later_application_data_atomic_partition_test (namespace, value) \
VALUES ($1, 'rolled-back')",
)
.bind(&namespace)
.execute(&mut *rolled_back)
.await?;
backend
.enqueue_to_partition_in(&mut rolled_back, TestPartitionJob("rolled-back".to_owned()))
.await?;
rolled_back.rollback().await?;
let mut committed = backend.pool().begin().await?;
sqlx::query(
"INSERT INTO later_application_data_atomic_partition_test (namespace, value) \
VALUES ($1, 'committed')",
)
.bind(&namespace)
.execute(&mut *committed)
.await?;
backend
.enqueue_to_partition_in(&mut committed, TestPartitionJob("committed".to_owned()))
.await?;
committed.commit().await?;
let app_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_application_data_atomic_partition_test \
WHERE namespace = $1",
)
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let job_key_prefix = format!("later-{namespace}%");
let job_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_storage \
WHERE key LIKE $1 AND key NOT LIKE '%-stats-%'",
)
.bind(job_key_prefix)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let partition_job_count: i64 =
sqlx::query("SELECT COUNT(*) AS count FROM later_partition_job WHERE namespace = $1")
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let partition_head_count: i64 =
sqlx::query("SELECT COUNT(*) AS count FROM later_partition_head WHERE namespace = $1")
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
let next_sequence: i64 =
sqlx::query("SELECT next_sequence FROM later_partition_sequence WHERE namespace = $1")
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("next_sequence")?;
assert_eq!(
(
app_count,
job_count,
partition_job_count,
partition_head_count
),
(1, 1, 1, 1)
);
assert_eq!(
next_sequence, 2,
"the rolled-back enqueue must not leave a gap in the partition sequence"
);
Ok(())
}
#[tokio::test]
async fn postgres_job_claim_is_exclusive_until_finished() -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("lease-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace, &connection_string).await?;
let parts = Box::new(backend).into_parts()?;
let job_id = crate::JobId(format!("postgres-lease-job-{}", crate::generate_id()));
let first = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("first Postgres lease was not acquired"))?;
assert!(parts.committer.claim(&job_id).await?.is_none());
first.finish().await?;
let after_finish = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("finished Postgres lease was not released"))?;
after_finish.finish().await?;
Ok(())
}
#[tokio::test]
async fn postgres_job_lease_loss_is_signaled_after_it_is_stolen() -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("lease-lost-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
let job_id = crate::JobId(format!("postgres-lease-lost-job-{}", crate::generate_id()));
let lease = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
let mut lease_lost = lease.lease_lost();
assert!(!lease_lost.has_changed().unwrap_or(true));
sqlx::query(
"UPDATE later_job_execution SET lease_owner = 'someone-else' \
WHERE namespace = $1 AND job_id = $2",
)
.bind(&namespace)
.bind(job_id.to_string())
.execute(backend.pool())
.await?;
tokio::time::timeout(std::time::Duration::from_secs(15), lease_lost.changed())
.await
.map_err(|_| anyhow::anyhow!("lease_lost did not fire within 15s of being stolen"))??;
Ok(())
}
#[tokio::test]
async fn postgres_stale_leased_job_ids_finds_leases_past_the_grace_period() -> anyhow::Result<()>
{
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("stale-lease-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
let grace = std::time::Duration::from_secs(30);
let abandoned_id = crate::JobId(format!("postgres-abandoned-job-{}", crate::generate_id()));
parts.committer.claim(&abandoned_id).await?;
sqlx::query(
"UPDATE later_job_execution \
SET lease_until = NOW() - make_interval(secs => $1) \
WHERE namespace = $2 AND job_id = $3",
)
.bind((grace.as_secs() * 2) as f64)
.bind(&namespace)
.bind(abandoned_id.to_string())
.execute(backend.pool())
.await?;
let recent_id = crate::JobId(format!(
"postgres-recently-expired-job-{}",
crate::generate_id()
));
parts.committer.claim(&recent_id).await?;
sqlx::query(
"UPDATE later_job_execution SET lease_until = NOW() - INTERVAL '1 second' \
WHERE namespace = $1 AND job_id = $2",
)
.bind(&namespace)
.bind(recent_id.to_string())
.execute(backend.pool())
.await?;
let live_id = crate::JobId(format!("postgres-live-job-{}", crate::generate_id()));
parts.committer.claim(&live_id).await?;
let stale = parts.committer.stale_leased_job_ids(grace, 100).await?;
assert_eq!(stale, vec![abandoned_id]);
Ok(())
}
#[tokio::test]
async fn postgres_claim_and_fetch_matches_a_separate_claim_and_get() -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("claim-and-fetch-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace, &connection_string).await?;
let parts = Box::new(backend).into_parts()?;
let job_id = crate::JobId(format!(
"postgres-claim-and-fetch-job-{}",
crate::generate_id()
));
let present_key = format!("job-present-{}", crate::generate_id());
let missing_key = format!("job-missing-{}", crate::generate_id());
parts.storage.set(&present_key, b"job-bytes").await?;
let claimed = parts
.committer
.claim_and_fetch(&job_id, &present_key, parts.storage.as_ref())
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
assert_eq!(claimed.value.as_deref(), Some(b"job-bytes".as_slice()));
assert!(parts.committer.claim(&job_id).await?.is_none());
claimed.lease.finish().await?;
let claimed_missing = parts
.committer
.claim_and_fetch(&job_id, &missing_key, parts.storage.as_ref())
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
assert!(claimed_missing.value.is_none());
claimed_missing.lease.finish().await?;
let held = parts
.committer
.claim(&job_id)
.await?
.ok_or_else(|| anyhow::anyhow!("lease should have been granted"))?;
assert!(parts
.committer
.claim_and_fetch(&job_id, &present_key, parts.storage.as_ref())
.await?
.is_none());
held.finish().await?;
Ok(())
}
#[tokio::test]
async fn postgres_worker_heartbeat_lists_live_workers_until_deregistered() -> anyhow::Result<()>
{
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("worker-heartbeat-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace, &connection_string).await?;
let parts = Box::new(backend).into_parts()?;
assert!(parts.committer.list_live_workers().await?.is_empty());
parts.committer.heartbeat_worker("worker-a").await?;
parts.committer.heartbeat_worker("worker-b").await?;
let mut live = parts.committer.list_live_workers().await?;
live.sort();
assert_eq!(live, vec!["worker-a".to_string(), "worker-b".to_string()]);
parts.committer.deregister_worker("worker-a").await?;
assert_eq!(
parts.committer.list_live_workers().await?,
vec!["worker-b".to_string()]
);
Ok(())
}
#[tokio::test]
async fn postgres_partition_head_moves_to_another_worker_only_after_lease_expiry(
) -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("crash-recovery-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
let job_id = crate::JobId(format!("crash-job-{}", crate::generate_id()));
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&job_id,
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let claim = parts
.committer
.claim_partition_head("worker-a", &["worker-a".to_string()])
.await?
.ok_or_else(|| anyhow::anyhow!("first worker did not claim the partition head"))?;
assert_eq!(claim.job_id, job_id);
assert_eq!(claim.owner, "worker-a");
assert!(parts
.committer
.claim_partition_head("worker-b", &["worker-b".to_string()])
.await?
.is_none());
sqlx::query(
"UPDATE later_partition_lease SET lease_until = NOW() - INTERVAL '1 hour' \
WHERE namespace = $1 AND topic = 'orders' AND partition_id = 0",
)
.bind(&namespace)
.execute(backend.pool())
.await?;
let recovered = parts
.committer
.claim_partition_head("worker-b", &["worker-b".to_string()])
.await?
.ok_or_else(|| {
anyhow::anyhow!("the expired lease was not reclaimed by another worker")
})?;
assert_eq!(recovered.job_id, job_id);
assert_eq!(recovered.sequence.0, claim.sequence.0);
assert_eq!(recovered.owner, "worker-b");
Ok(())
}
#[tokio::test]
async fn postgres_completing_a_head_keeps_a_valid_assignment() -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("partition-epoch-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
for id in ["first-job", "second-job"] {
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId(id.to_owned()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
}
let same_owner = "shared-owner";
let first_claim = parts
.committer
.claim_partition_head(same_owner, &[same_owner.to_owned()])
.await?
.ok_or_else(|| anyhow::anyhow!("expected to claim the first job"))?;
assert!(
parts
.committer
.complete_partition_head(&first_claim, Vec::new())
.await?
);
let second_claim = parts
.committer
.peek_assigned_partition_head(
same_owner,
"orders",
crate::topic::PartitionId(0),
first_claim.lease_epoch,
)
.await?
.ok_or_else(|| anyhow::anyhow!("expected to keep the second head assignment"))?;
assert_eq!(
second_claim.lease_epoch, first_claim.lease_epoch,
"a valid assignment must be reused across completions"
);
assert!(
!parts
.committer
.complete_partition_head(&first_claim, Vec::new())
.await?,
"a stale same-owner claim from an earlier job must not complete a later one"
);
assert!(
parts
.committer
.complete_partition_head(&second_claim, Vec::new())
.await?
);
let lease_rows: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM later_partition_lease \
WHERE namespace = $1 AND topic = 'orders' AND partition_id = 0",
)
.bind(&namespace)
.fetch_one(backend.pool())
.await?;
assert_eq!(
lease_rows, 1,
"the lease row must persist across completions, not be deleted"
);
Ok(())
}
#[tokio::test]
async fn postgres_stale_partition_claim_cannot_complete_a_reclaimed_head() -> anyhow::Result<()>
{
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("partition-fence-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
let job_id = crate::JobId(format!("fenced-job-{}", crate::generate_id()));
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&job_id,
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let first = parts
.committer
.claim_partition_head("worker-a", &["worker-a".to_owned()])
.await?
.ok_or_else(|| anyhow::anyhow!("first worker did not claim the partition head"))?;
sqlx::query(
"UPDATE later_partition_lease SET lease_until = NOW() - INTERVAL '1 hour' \
WHERE namespace = $1 AND topic = 'orders' AND partition_id = 0",
)
.bind(&namespace)
.execute(backend.pool())
.await?;
let reclaimed = parts
.committer
.claim_partition_head("worker-b", &["worker-b".to_owned()])
.await?
.ok_or_else(|| anyhow::anyhow!("second worker did not reclaim the partition head"))?;
assert!(reclaimed.lease_epoch > first.lease_epoch);
let stale_release_key = format!("{namespace}-stale-release-fence");
assert!(
!parts
.committer
.release_partition_head(
&first,
chrono::Utc::now(),
vec![crate::storage::StorageOperation::Set {
key: stale_release_key.clone(),
value: b"stale".to_vec(),
}],
)
.await?,
"a stale claim must not save retry state for a reclaimed head"
);
let stale_release_writes: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM later_storage WHERE key = $1")
.bind(&stale_release_key)
.fetch_one(backend.pool())
.await?;
assert_eq!(stale_release_writes, 0);
let stale_completion_key = format!("{namespace}-stale-completion-fence");
assert!(
!parts
.committer
.complete_partition_head(
&first,
vec![crate::storage::StorageOperation::Set {
key: stale_completion_key.clone(),
value: b"stale".to_vec(),
}],
)
.await?,
"a stale claim must not remove a head reclaimed at a newer lease epoch"
);
let stale_completion_writes: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM later_storage WHERE key = $1")
.bind(&stale_completion_key)
.fetch_one(backend.pool())
.await?;
assert_eq!(stale_completion_writes, 0);
assert_eq!(
parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?,
Some((job_id.clone(), first.sequence))
);
assert!(
parts
.committer
.complete_partition_head(&reclaimed, Vec::new())
.await?
);
assert!(parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?
.is_none());
Ok(())
}
#[tokio::test]
async fn postgres_admin_force_complete_partition_head_unblocks_a_poison_head(
) -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("admin-force-fail-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
let poison_id = crate::JobId("poison-job".to_owned());
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&poison_id,
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId("next-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let (peeked_id, sequence) = parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?
.ok_or_else(|| anyhow::anyhow!("expected a head to peek"))?;
assert_eq!(peeked_id, poison_id);
let stale_rejected = parts
.committer
.admin_force_complete_partition_head(
"orders",
crate::topic::PartitionId(0),
PartitionSequence(sequence.0 + 1),
&poison_id,
Vec::new(),
)
.await?;
assert!(!stale_rejected);
let cleared = parts
.committer
.admin_force_complete_partition_head(
"orders",
crate::topic::PartitionId(0),
sequence,
&poison_id,
Vec::new(),
)
.await?;
assert!(cleared);
let lease_count: i64 = sqlx::query(
"SELECT COUNT(*) AS count FROM later_partition_lease \
WHERE namespace = $1 AND topic = 'orders' AND partition_id = 0",
)
.bind(&namespace)
.fetch_one(backend.pool())
.await?
.try_get("count")?;
assert_eq!(lease_count, 0);
let (new_head_id, _) = parts
.committer
.admin_peek_partition_head("orders", crate::topic::PartitionId(0))
.await?
.ok_or_else(|| anyhow::anyhow!("expected the next job to become the head"))?;
assert_eq!(new_head_id, crate::JobId("next-job".to_owned()));
let repeated = parts
.committer
.admin_force_complete_partition_head(
"orders",
crate::topic::PartitionId(0),
sequence,
&poison_id,
Vec::new(),
)
.await?;
assert!(!repeated);
Ok(())
}
#[tokio::test]
async fn postgres_partition_backlog_summary_counts_ready_heads_and_ignores_unready_ones(
) -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("backlog-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 2)?])
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId("ready-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now() - chrono::Duration::seconds(30),
)
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(1),
&crate::JobId("future-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now() + chrono::Duration::hours(1),
)
.await?;
let summary = parts.committer.partition_backlog_summary().await?;
assert_eq!(summary.len(), 1);
let (topic, ready_count, oldest_age_seconds) = &summary[0];
assert_eq!(topic, "orders");
assert_eq!(*ready_count, 1);
assert!(
*oldest_age_seconds >= 30.0,
"expected the ready head's age to be at least 30s, got {oldest_age_seconds}"
);
Ok(())
}
#[tokio::test]
async fn postgres_partition_commit_enqueues_a_wake_up_notification_on_the_sql_queue(
) -> anyhow::Result<()> {
let connection_string = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let namespace = format!("partition-wake-postgres-{}", crate::generate_id());
let backend = PostgresBackend::connect(namespace.clone(), &connection_string).await?;
let parts = Box::new(backend.clone()).into_parts()?;
parts
.committer
.register_topics(&[crate::topic::TopicConfig::new("orders", 1)?])
.await?;
parts
.committer
.commit_partitioned(
Box::new(|_| Ok(Vec::new())),
"orders",
crate::topic::PartitionId(0),
&crate::JobId("wake-job".to_owned()),
"TestPartitionJob",
chrono::Utc::now(),
)
.await?;
let payload: Vec<u8> =
sqlx::query("SELECT payload FROM later_delivery_queue WHERE namespace = $1")
.bind(namespace)
.fetch_one(backend.pool())
.await?
.try_get("payload")?;
let command: crate::models::AmqpCommand = crate::encoder::decode(&payload)?;
assert!(
matches!(&command, crate::models::AmqpCommand::PartitionReady { topic } if topic == "orders"),
"expected a PartitionReady notification for topic 'orders', got {command:?}"
);
Ok(())
}
}