use crate::topic::{partition_for_key, PartitionId, PartitionKey, TopicConfig};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct LogRecord {
pub topic: String,
pub partition: PartitionId,
pub offset: i64,
pub key: PartitionKey,
pub payload: Vec<u8>,
pub created_at: crate::UtcDateTime,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ConsumerLag {
pub partition: PartitionId,
pub group: String,
pub end_offset: i64,
pub committed_offset: i64,
pub lag: i64,
}
#[async_trait::async_trait]
pub trait LogLagSource: Send + Sync {
async fn consumer_lag(&self, topic: &str) -> anyhow::Result<Vec<ConsumerLag>>;
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct LogPartitionClaim {
pub group: String,
pub topic: String,
pub partition: PartitionId,
pub member_id: String,
pub epoch: i64,
}
const MEMBER_TTL_SECONDS: i64 = 30;
fn validate_member_id(member_id: &str) -> anyhow::Result<()> {
if member_id.trim().is_empty() {
return Err(anyhow::anyhow!("consumer member ID cannot be empty"));
}
Ok(())
}
fn validate_read_limit(limit: u32) -> anyhow::Result<i64> {
if limit == 0 || limit > 10_000 {
return Err(anyhow::anyhow!("log read limit must be from 1 to 10,000"));
}
Ok(i64::from(limit))
}
fn validate_group_name(group: &str) -> anyhow::Result<()> {
if group.trim().is_empty() {
return Err(anyhow::anyhow!("consumer group name cannot be empty"));
}
Ok(())
}
fn validate_cleanup_offset(offset: i64) -> anyhow::Result<()> {
if offset < 1 {
return Err(anyhow::anyhow!("cleanup offset must be positive"));
}
Ok(())
}
#[cfg(feature = "sqlite")]
#[derive(Clone)]
pub struct SqliteLog {
namespace: String,
pool: sqlx::SqlitePool,
}
#[cfg(feature = "sqlite")]
impl SqliteLog {
pub async fn from_pool(
namespace: impl Into<String>,
pool: sqlx::SqlitePool,
) -> anyhow::Result<Self> {
crate::storage::Sqlite::from_pool(pool.clone()).await?;
Ok(Self {
namespace: namespace.into(),
pool,
})
}
pub async fn connect(namespace: impl Into<String>, url: &str) -> anyhow::Result<Self> {
let storage = crate::storage::Sqlite::new(url).await?;
Ok(Self {
namespace: namespace.into(),
pool: storage.pool().clone(),
})
}
pub async fn register_topic(&self, topic: TopicConfig) -> anyhow::Result<()> {
let count = i64::from(topic.partition_count());
let result: Option<(i64,)> = sqlx::query_as(
"INSERT INTO later_topic (namespace, topic, partition_count) VALUES (?1, ?2, ?3) \
ON CONFLICT(namespace, topic) DO UPDATE SET partition_count = later_topic.partition_count \
WHERE later_topic.partition_count = excluded.partition_count RETURNING partition_count",
).bind(&self.namespace).bind(topic.name()).bind(count).fetch_optional(&self.pool).await?;
if result.is_none() {
return Err(anyhow::anyhow!(
"topic '{}' has a different partition count",
topic.name()
));
}
Ok(())
}
pub async fn append(
&self,
topic: &str,
key: PartitionKey,
payload: Vec<u8>,
) -> anyhow::Result<LogRecord> {
use sqlx::Row;
let mut transaction = self.pool.begin().await?;
let count: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = ?1 AND topic = ?2",
)
.bind(&self.namespace)
.bind(topic)
.fetch_optional(&mut *transaction)
.await?;
let Some((count,)) = count else {
return Err(anyhow::anyhow!("unknown topic '{topic}'"));
};
let partition = partition_for_key(&key, u32::try_from(count)?);
let offset: i64 = sqlx::query("INSERT INTO later_log_partition (namespace, topic, partition_id, next_offset) VALUES (?1, ?2, ?3, 2) ON CONFLICT(namespace, topic, partition_id) DO UPDATE SET next_offset = later_log_partition.next_offset + 1 RETURNING next_offset - 1")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).fetch_one(&mut *transaction).await?.try_get(0)?;
let created_at =
chrono::DateTime::from_timestamp_millis(chrono::Utc::now().timestamp_millis())
.ok_or_else(|| anyhow::anyhow!("current timestamp is out of range"))?;
sqlx::query("INSERT INTO later_log_record (namespace, topic, partition_id, record_offset, record_key, payload, created_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).bind(offset).bind(&key.0).bind(&payload).bind(created_at.timestamp_millis()).execute(&mut *transaction).await?;
transaction.commit().await?;
Ok(LogRecord {
topic: topic.to_owned(),
partition,
offset,
key,
payload,
created_at,
})
}
pub async fn read(
&self,
topic: &str,
partition: PartitionId,
offset: i64,
limit: u32,
) -> anyhow::Result<Vec<LogRecord>> {
use sqlx::Row;
let limit = validate_read_limit(limit)?;
let rows = sqlx::query("SELECT record_offset, record_key, payload, created_at FROM later_log_record WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 AND record_offset >= ?4 ORDER BY record_offset LIMIT ?5")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).bind(offset).bind(limit).fetch_all(&self.pool).await?;
rows.into_iter()
.map(|row| {
let millis: i64 = row.try_get("created_at")?;
let created_at = chrono::DateTime::from_timestamp_millis(millis)
.ok_or_else(|| anyhow::anyhow!("stored log timestamp is out of range"))?;
Ok(LogRecord {
topic: topic.to_owned(),
partition,
offset: row.try_get("record_offset")?,
key: PartitionKey(row.try_get("record_key")?),
payload: row.try_get("payload")?,
created_at,
})
})
.collect()
}
pub async fn commit_offset(
&self,
topic: &str,
group: &str,
partition: PartitionId,
next_offset: i64,
) -> anyhow::Result<()> {
validate_group_name(group)?;
if next_offset < 1 {
return Err(anyhow::anyhow!("next offset must be positive"));
}
sqlx::query("INSERT INTO later_log_consumer_offset (namespace, topic, group_name, partition_id, next_offset, committed_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6) ON CONFLICT(namespace, topic, group_name, partition_id) DO UPDATE SET next_offset = excluded.next_offset, committed_at = excluded.committed_at WHERE excluded.next_offset > later_log_consumer_offset.next_offset")
.bind(&self.namespace).bind(topic).bind(group).bind(i64::from(partition.0)).bind(next_offset).bind(chrono::Utc::now().timestamp_millis()).execute(&self.pool).await?;
Ok(())
}
pub async fn committed_offset(
&self,
topic: &str,
group: &str,
partition: PartitionId,
) -> anyhow::Result<i64> {
validate_group_name(group)?;
let row: Option<(i64,)> = sqlx::query_as("SELECT next_offset FROM later_log_consumer_offset WHERE namespace = ?1 AND topic = ?2 AND group_name = ?3 AND partition_id = ?4")
.bind(&self.namespace).bind(topic).bind(group).bind(i64::from(partition.0)).fetch_optional(&self.pool).await?;
Ok(row.map(|(offset,)| offset).unwrap_or(1))
}
pub async fn consumer_lag(&self, topic: &str) -> anyhow::Result<Vec<ConsumerLag>> {
use sqlx::Row;
let rows = sqlx::query(
"SELECT c.partition_id, c.group_name, c.next_offset AS committed_offset, \
p.next_offset AS end_offset \
FROM later_log_consumer_offset c \
JOIN later_log_partition p \
ON p.namespace = c.namespace AND p.topic = c.topic AND p.partition_id = c.partition_id \
WHERE c.namespace = ?1 AND c.topic = ?2",
)
.bind(&self.namespace)
.bind(topic)
.fetch_all(&self.pool)
.await?;
rows.into_iter()
.map(|row| {
let partition_id: i64 = row.try_get("partition_id")?;
let committed_offset: i64 = row.try_get("committed_offset")?;
let end_offset: i64 = row.try_get("end_offset")?;
Ok(ConsumerLag {
partition: PartitionId(u32::try_from(partition_id)?),
group: row.try_get("group_name")?,
end_offset,
committed_offset,
lag: (end_offset - committed_offset).max(0),
})
})
.collect()
}
pub async fn safe_cleanup_offset(
&self,
topic: &str,
partition: PartitionId,
) -> anyhow::Result<Option<i64>> {
let row: (Option<i64>,) = sqlx::query_as("SELECT MIN(next_offset) FROM later_log_consumer_offset WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).fetch_one(&self.pool).await?;
Ok(row.0)
}
pub async fn delete_before(
&self,
topic: &str,
partition: PartitionId,
offset: i64,
) -> anyhow::Result<u64> {
validate_cleanup_offset(offset)?;
let Some(safe_offset) = self.safe_cleanup_offset(topic, partition).await? else {
return Ok(0);
};
if offset > safe_offset {
return Err(anyhow::anyhow!(
"cleanup offset is ahead of an active consumer group"
));
}
Ok(sqlx::query(
"DELETE FROM later_log_record WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3 \
AND record_offset < ?4 \
AND record_offset < COALESCE( \
(SELECT MIN(next_offset) FROM later_log_consumer_offset \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3), 0)",
)
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).bind(offset).execute(&self.pool).await?.rows_affected())
}
pub async fn compact_partition(
&self,
topic: &str,
partition: PartitionId,
) -> anyhow::Result<u64> {
Ok(sqlx::query(
"DELETE FROM later_log_record AS older WHERE older.namespace = ?1 AND older.topic = ?2 AND older.partition_id = ?3 \
AND older.record_offset < COALESCE( \
(SELECT MIN(next_offset) FROM later_log_consumer_offset \
WHERE namespace = ?1 AND topic = ?2 AND partition_id = ?3), 0) \
AND EXISTS (SELECT 1 FROM later_log_record AS newer WHERE newer.namespace = older.namespace AND newer.topic = older.topic AND newer.partition_id = older.partition_id AND newer.record_key = older.record_key AND newer.record_offset > older.record_offset)",
)
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).execute(&self.pool).await?.rows_affected())
}
pub async fn heartbeat(&self, topic: &str, group: &str, member_id: &str) -> anyhow::Result<()> {
validate_group_name(group)?;
validate_member_id(member_id)?;
sqlx::query("INSERT INTO later_log_consumer_member (namespace, topic, group_name, member_id, heartbeat_at) VALUES (?1, ?2, ?3, ?4, ?5) ON CONFLICT(namespace, topic, group_name, member_id) DO UPDATE SET heartbeat_at = excluded.heartbeat_at")
.bind(&self.namespace).bind(topic).bind(group).bind(member_id).bind(chrono::Utc::now().timestamp_millis()).execute(&self.pool).await?;
Ok(())
}
pub async fn claim_partition(
&self,
topic: &str,
group: &str,
member_id: &str,
partition: PartitionId,
) -> anyhow::Result<Option<LogPartitionClaim>> {
use sqlx::Row;
self.heartbeat(topic, group, member_id).await?;
let now = chrono::Utc::now().timestamp_millis();
let members: Vec<(String,)> = sqlx::query_as("SELECT member_id FROM later_log_consumer_member WHERE namespace = ?1 AND topic = ?2 AND group_name = ?3 AND heartbeat_at > ?4")
.bind(&self.namespace).bind(topic).bind(group).bind(now - MEMBER_TTL_SECONDS * 1000).fetch_all(&self.pool).await?;
let members = members
.into_iter()
.map(|(member,)| member)
.collect::<Vec<_>>();
if crate::topic::rendezvous_owner(topic, partition, &members) != Some(member_id) {
return Ok(None);
}
let lease_until = now + MEMBER_TTL_SECONDS * 1000;
let row = sqlx::query("INSERT INTO later_log_consumer_lease (namespace, topic, group_name, partition_id, member_id, lease_until, lease_epoch) VALUES (?1, ?2, ?3, ?4, ?5, ?6, 1) ON CONFLICT(namespace, topic, group_name, partition_id) DO UPDATE SET member_id = excluded.member_id, lease_until = excluded.lease_until, lease_epoch = later_log_consumer_lease.lease_epoch + 1 WHERE later_log_consumer_lease.lease_until <= ?7 RETURNING lease_epoch")
.bind(&self.namespace).bind(topic).bind(group).bind(i64::from(partition.0)).bind(member_id).bind(lease_until).bind(now).fetch_optional(&self.pool).await?;
let Some(row) = row else {
return Ok(None);
};
Ok(Some(LogPartitionClaim {
group: group.to_owned(),
topic: topic.to_owned(),
partition,
member_id: member_id.to_owned(),
epoch: row.try_get(0)?,
}))
}
pub async fn renew_claim(&self, claim: &LogPartitionClaim) -> anyhow::Result<bool> {
let now = chrono::Utc::now().timestamp_millis();
let result = sqlx::query("UPDATE later_log_consumer_lease SET lease_until = ?1 WHERE namespace = ?2 AND topic = ?3 AND group_name = ?4 AND partition_id = ?5 AND member_id = ?6 AND lease_epoch = ?7 AND lease_until > ?8")
.bind(now + MEMBER_TTL_SECONDS * 1000).bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(&claim.member_id).bind(claim.epoch).bind(now).execute(&self.pool).await?;
Ok(result.rows_affected() == 1)
}
pub async fn release_claim(&self, claim: &LogPartitionClaim) -> anyhow::Result<bool> {
let result = sqlx::query("DELETE FROM later_log_consumer_lease WHERE namespace = ?1 AND topic = ?2 AND group_name = ?3 AND partition_id = ?4 AND member_id = ?5 AND lease_epoch = ?6")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(&claim.member_id).bind(claim.epoch).execute(&self.pool).await?;
Ok(result.rows_affected() == 1)
}
pub async fn leave_group(
&self,
topic: &str,
group: &str,
member_id: &str,
) -> anyhow::Result<()> {
validate_group_name(group)?;
validate_member_id(member_id)?;
let mut transaction = self.pool.begin().await?;
sqlx::query("DELETE FROM later_log_consumer_member WHERE namespace = ?1 AND topic = ?2 AND group_name = ?3 AND member_id = ?4")
.bind(&self.namespace).bind(topic).bind(group).bind(member_id).execute(&mut *transaction).await?;
sqlx::query("DELETE FROM later_log_consumer_lease WHERE namespace = ?1 AND topic = ?2 AND group_name = ?3 AND member_id = ?4")
.bind(&self.namespace).bind(topic).bind(group).bind(member_id).execute(&mut *transaction).await?;
transaction.commit().await?;
Ok(())
}
pub async fn commit_claimed_offset(
&self,
claim: &LogPartitionClaim,
next_offset: i64,
) -> anyhow::Result<bool> {
if next_offset < 1 {
return Err(anyhow::anyhow!("next offset must be positive"));
}
let now = chrono::Utc::now().timestamp_millis();
let mut transaction = self.pool.begin().await?;
let owned: Option<(i64,)> = sqlx::query_as("SELECT lease_epoch FROM later_log_consumer_lease WHERE namespace = ?1 AND topic = ?2 AND group_name = ?3 AND partition_id = ?4 AND member_id = ?5 AND lease_epoch = ?6 AND lease_until > ?7")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(&claim.member_id).bind(claim.epoch).bind(now).fetch_optional(&mut *transaction).await?;
if owned.is_none() {
transaction.rollback().await?;
return Ok(false);
}
sqlx::query("INSERT INTO later_log_consumer_offset (namespace, topic, group_name, partition_id, next_offset, committed_at) VALUES (?1, ?2, ?3, ?4, ?5, ?6) ON CONFLICT(namespace, topic, group_name, partition_id) DO UPDATE SET next_offset = excluded.next_offset, committed_at = excluded.committed_at WHERE excluded.next_offset > later_log_consumer_offset.next_offset")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(next_offset).bind(now).execute(&mut *transaction).await?;
transaction.commit().await?;
Ok(true)
}
}
#[cfg(feature = "postgres")]
#[derive(Clone)]
pub struct PostgresLog {
namespace: String,
pool: sqlx::PgPool,
}
#[cfg(feature = "postgres")]
impl PostgresLog {
pub async fn from_pool(
namespace: impl Into<String>,
pool: sqlx::PgPool,
) -> anyhow::Result<Self> {
crate::storage::Postgres::from_pool(pool.clone()).await?;
Ok(Self {
namespace: namespace.into(),
pool,
})
}
pub async fn connect(namespace: impl Into<String>, url: &str) -> anyhow::Result<Self> {
let storage = crate::storage::Postgres::new(url).await?;
Ok(Self {
namespace: namespace.into(),
pool: storage.pool().clone(),
})
}
pub async fn register_topic(&self, topic: TopicConfig) -> anyhow::Result<()> {
let count = i64::from(topic.partition_count());
let result: Option<(i64,)> = sqlx::query_as(
"INSERT INTO later_topic (namespace, topic, partition_count) VALUES ($1, $2, $3) \
ON CONFLICT(namespace, topic) DO UPDATE SET partition_count = later_topic.partition_count \
WHERE later_topic.partition_count = excluded.partition_count RETURNING partition_count",
).bind(&self.namespace).bind(topic.name()).bind(count).fetch_optional(&self.pool).await?;
if result.is_none() {
return Err(anyhow::anyhow!(
"topic '{}' has a different partition count",
topic.name()
));
}
Ok(())
}
pub async fn append(
&self,
topic: &str,
key: PartitionKey,
payload: Vec<u8>,
) -> anyhow::Result<LogRecord> {
use sqlx::Row;
let mut transaction = self.pool.begin().await?;
let count: Option<(i64,)> = sqlx::query_as(
"SELECT partition_count FROM later_topic WHERE namespace = $1 AND topic = $2",
)
.bind(&self.namespace)
.bind(topic)
.fetch_optional(&mut *transaction)
.await?;
let Some((count,)) = count else {
return Err(anyhow::anyhow!("unknown topic '{topic}'"));
};
let partition = partition_for_key(&key, u32::try_from(count)?);
let offset: i64 = sqlx::query("INSERT INTO later_log_partition (namespace, topic, partition_id, next_offset) VALUES ($1, $2, $3, 2) ON CONFLICT(namespace, topic, partition_id) DO UPDATE SET next_offset = later_log_partition.next_offset + 1 RETURNING next_offset - 1")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).fetch_one(&mut *transaction).await?.try_get(0)?;
let created_at: crate::UtcDateTime = sqlx::query("INSERT INTO later_log_record (namespace, topic, partition_id, record_offset, record_key, payload, created_at) VALUES ($1, $2, $3, $4, $5, $6, NOW()) RETURNING created_at")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).bind(offset).bind(&key.0).bind(&payload).fetch_one(&mut *transaction).await?.try_get("created_at")?;
transaction.commit().await?;
Ok(LogRecord {
topic: topic.to_owned(),
partition,
offset,
key,
payload,
created_at,
})
}
pub async fn read(
&self,
topic: &str,
partition: PartitionId,
offset: i64,
limit: u32,
) -> anyhow::Result<Vec<LogRecord>> {
use sqlx::Row;
let limit = validate_read_limit(limit)?;
let rows = sqlx::query("SELECT record_offset, record_key, payload, created_at FROM later_log_record WHERE namespace = $1 AND topic = $2 AND partition_id = $3 AND record_offset >= $4 ORDER BY record_offset LIMIT $5")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).bind(offset).bind(limit).fetch_all(&self.pool).await?;
rows.into_iter()
.map(|row| {
Ok(LogRecord {
topic: topic.to_owned(),
partition,
offset: row.try_get("record_offset")?,
key: PartitionKey(row.try_get("record_key")?),
payload: row.try_get("payload")?,
created_at: row.try_get("created_at")?,
})
})
.collect()
}
pub async fn commit_offset(
&self,
topic: &str,
group: &str,
partition: PartitionId,
next_offset: i64,
) -> anyhow::Result<()> {
validate_group_name(group)?;
if next_offset < 1 {
return Err(anyhow::anyhow!("next offset must be positive"));
}
sqlx::query("INSERT INTO later_log_consumer_offset (namespace, topic, group_name, partition_id, next_offset) VALUES ($1, $2, $3, $4, $5) ON CONFLICT(namespace, topic, group_name, partition_id) DO UPDATE SET next_offset = excluded.next_offset, committed_at = NOW() WHERE excluded.next_offset > later_log_consumer_offset.next_offset")
.bind(&self.namespace).bind(topic).bind(group).bind(i64::from(partition.0)).bind(next_offset).execute(&self.pool).await?;
Ok(())
}
pub async fn committed_offset(
&self,
topic: &str,
group: &str,
partition: PartitionId,
) -> anyhow::Result<i64> {
validate_group_name(group)?;
let row: Option<(i64,)> = sqlx::query_as("SELECT next_offset FROM later_log_consumer_offset WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND partition_id = $4")
.bind(&self.namespace).bind(topic).bind(group).bind(i64::from(partition.0)).fetch_optional(&self.pool).await?;
Ok(row.map(|(offset,)| offset).unwrap_or(1))
}
pub async fn consumer_lag(&self, topic: &str) -> anyhow::Result<Vec<ConsumerLag>> {
use sqlx::Row;
let rows = sqlx::query(
"SELECT c.partition_id, c.group_name, c.next_offset AS committed_offset, \
p.next_offset AS end_offset \
FROM later_log_consumer_offset c \
JOIN later_log_partition p \
ON p.namespace = c.namespace AND p.topic = c.topic AND p.partition_id = c.partition_id \
WHERE c.namespace = $1 AND c.topic = $2",
)
.bind(&self.namespace)
.bind(topic)
.fetch_all(&self.pool)
.await?;
rows.into_iter()
.map(|row| {
let partition_id: i64 = row.try_get("partition_id")?;
let committed_offset: i64 = row.try_get("committed_offset")?;
let end_offset: i64 = row.try_get("end_offset")?;
Ok(ConsumerLag {
partition: PartitionId(u32::try_from(partition_id)?),
group: row.try_get("group_name")?,
end_offset,
committed_offset,
lag: (end_offset - committed_offset).max(0),
})
})
.collect()
}
pub async fn safe_cleanup_offset(
&self,
topic: &str,
partition: PartitionId,
) -> anyhow::Result<Option<i64>> {
let row: (Option<i64>,) = sqlx::query_as("SELECT MIN(next_offset) FROM later_log_consumer_offset WHERE namespace = $1 AND topic = $2 AND partition_id = $3")
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).fetch_one(&self.pool).await?;
Ok(row.0)
}
pub async fn delete_before(
&self,
topic: &str,
partition: PartitionId,
offset: i64,
) -> anyhow::Result<u64> {
validate_cleanup_offset(offset)?;
let Some(safe_offset) = self.safe_cleanup_offset(topic, partition).await? else {
return Ok(0);
};
if offset > safe_offset {
return Err(anyhow::anyhow!(
"cleanup offset is ahead of an active consumer group"
));
}
Ok(sqlx::query(
"DELETE FROM later_log_record WHERE namespace = $1 AND topic = $2 AND partition_id = $3 \
AND record_offset < $4 \
AND record_offset < COALESCE( \
(SELECT MIN(next_offset) FROM later_log_consumer_offset \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3), 0)",
)
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).bind(offset).execute(&self.pool).await?.rows_affected())
}
pub async fn compact_partition(
&self,
topic: &str,
partition: PartitionId,
) -> anyhow::Result<u64> {
Ok(sqlx::query(
"DELETE FROM later_log_record AS older WHERE older.namespace = $1 AND older.topic = $2 AND older.partition_id = $3 \
AND older.record_offset < COALESCE( \
(SELECT MIN(next_offset) FROM later_log_consumer_offset \
WHERE namespace = $1 AND topic = $2 AND partition_id = $3), 0) \
AND EXISTS (SELECT 1 FROM later_log_record AS newer WHERE newer.namespace = older.namespace AND newer.topic = older.topic AND newer.partition_id = older.partition_id AND newer.record_key = older.record_key AND newer.record_offset > older.record_offset)",
)
.bind(&self.namespace).bind(topic).bind(i64::from(partition.0)).execute(&self.pool).await?.rows_affected())
}
pub async fn heartbeat(&self, topic: &str, group: &str, member_id: &str) -> anyhow::Result<()> {
validate_group_name(group)?;
validate_member_id(member_id)?;
sqlx::query("INSERT INTO later_log_consumer_member (namespace, topic, group_name, member_id) VALUES ($1, $2, $3, $4) ON CONFLICT(namespace, topic, group_name, member_id) DO UPDATE SET heartbeat_at = NOW()")
.bind(&self.namespace).bind(topic).bind(group).bind(member_id).execute(&self.pool).await?;
Ok(())
}
pub async fn claim_partition(
&self,
topic: &str,
group: &str,
member_id: &str,
partition: PartitionId,
) -> anyhow::Result<Option<LogPartitionClaim>> {
use sqlx::Row;
self.heartbeat(topic, group, member_id).await?;
let members: Vec<(String,)> = sqlx::query_as("SELECT member_id FROM later_log_consumer_member WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND heartbeat_at > NOW() - INTERVAL '30 seconds'")
.bind(&self.namespace).bind(topic).bind(group).fetch_all(&self.pool).await?;
let members = members
.into_iter()
.map(|(member,)| member)
.collect::<Vec<_>>();
if crate::topic::rendezvous_owner(topic, partition, &members) != Some(member_id) {
return Ok(None);
}
let row = sqlx::query("INSERT INTO later_log_consumer_lease (namespace, topic, group_name, partition_id, member_id, lease_until, lease_epoch) VALUES ($1, $2, $3, $4, $5, NOW() + INTERVAL '30 seconds', 1) ON CONFLICT(namespace, topic, group_name, partition_id) DO UPDATE SET member_id = excluded.member_id, lease_until = excluded.lease_until, lease_epoch = later_log_consumer_lease.lease_epoch + 1 WHERE later_log_consumer_lease.lease_until <= NOW() RETURNING lease_epoch")
.bind(&self.namespace).bind(topic).bind(group).bind(i64::from(partition.0)).bind(member_id).fetch_optional(&self.pool).await?;
let Some(row) = row else {
return Ok(None);
};
Ok(Some(LogPartitionClaim {
group: group.to_owned(),
topic: topic.to_owned(),
partition,
member_id: member_id.to_owned(),
epoch: row.try_get(0)?,
}))
}
pub async fn renew_claim(&self, claim: &LogPartitionClaim) -> anyhow::Result<bool> {
let result = sqlx::query("UPDATE later_log_consumer_lease SET lease_until = NOW() + INTERVAL '30 seconds' WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND partition_id = $4 AND member_id = $5 AND lease_epoch = $6 AND lease_until > NOW()")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(&claim.member_id).bind(claim.epoch).execute(&self.pool).await?;
Ok(result.rows_affected() == 1)
}
pub async fn release_claim(&self, claim: &LogPartitionClaim) -> anyhow::Result<bool> {
let result = sqlx::query("DELETE FROM later_log_consumer_lease WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND partition_id = $4 AND member_id = $5 AND lease_epoch = $6")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(&claim.member_id).bind(claim.epoch).execute(&self.pool).await?;
Ok(result.rows_affected() == 1)
}
pub async fn leave_group(
&self,
topic: &str,
group: &str,
member_id: &str,
) -> anyhow::Result<()> {
validate_group_name(group)?;
validate_member_id(member_id)?;
let mut transaction = self.pool.begin().await?;
sqlx::query("DELETE FROM later_log_consumer_member WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND member_id = $4")
.bind(&self.namespace).bind(topic).bind(group).bind(member_id).execute(&mut *transaction).await?;
sqlx::query("DELETE FROM later_log_consumer_lease WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND member_id = $4")
.bind(&self.namespace).bind(topic).bind(group).bind(member_id).execute(&mut *transaction).await?;
transaction.commit().await?;
Ok(())
}
pub async fn commit_claimed_offset(
&self,
claim: &LogPartitionClaim,
next_offset: i64,
) -> anyhow::Result<bool> {
if next_offset < 1 {
return Err(anyhow::anyhow!("next offset must be positive"));
}
let mut transaction = self.pool.begin().await?;
let owned: Option<(i64,)> = sqlx::query_as("SELECT lease_epoch FROM later_log_consumer_lease WHERE namespace = $1 AND topic = $2 AND group_name = $3 AND partition_id = $4 AND member_id = $5 AND lease_epoch = $6 AND lease_until > NOW()")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(&claim.member_id).bind(claim.epoch).fetch_optional(&mut *transaction).await?;
if owned.is_none() {
transaction.rollback().await?;
return Ok(false);
}
sqlx::query("INSERT INTO later_log_consumer_offset (namespace, topic, group_name, partition_id, next_offset) VALUES ($1, $2, $3, $4, $5) ON CONFLICT(namespace, topic, group_name, partition_id) DO UPDATE SET next_offset = excluded.next_offset, committed_at = NOW() WHERE excluded.next_offset > later_log_consumer_offset.next_offset")
.bind(&self.namespace).bind(&claim.topic).bind(&claim.group).bind(i64::from(claim.partition.0)).bind(next_offset).execute(&mut *transaction).await?;
transaction.commit().await?;
Ok(true)
}
}
#[cfg(feature = "sqlite")]
#[async_trait::async_trait]
impl LogLagSource for SqliteLog {
async fn consumer_lag(&self, topic: &str) -> anyhow::Result<Vec<ConsumerLag>> {
self.consumer_lag(topic).await
}
}
#[cfg(feature = "postgres")]
#[async_trait::async_trait]
impl LogLagSource for PostgresLog {
async fn consumer_lag(&self, topic: &str) -> anyhow::Result<Vec<ConsumerLag>> {
self.consumer_lag(topic).await
}
}
#[cfg(all(test, feature = "sqlite"))]
mod sqlite_tests {
use super::*;
#[tokio::test]
async fn cleanup_with_no_consumer_groups_deletes_nothing() -> anyhow::Result<()> {
let log = SqliteLog::connect("log-test", "sqlite::memory:").await?;
log.register_topic(TopicConfig::new("orphan", 1)?).await?;
log.append("orphan", PartitionKey::from("a"), b"1".to_vec())
.await?;
log.append("orphan", PartitionKey::from("a"), b"2".to_vec())
.await?;
assert_eq!(
log.safe_cleanup_offset("orphan", PartitionId(0)).await?,
None
);
assert_eq!(log.compact_partition("orphan", PartitionId(0)).await?, 0);
assert_eq!(
log.read("orphan", PartitionId(0), 1, 10).await?.len(),
2,
"compact_partition must not have deleted anything with no consumer groups"
);
Ok(())
}
#[tokio::test]
async fn normal_topic_retains_every_record_until_explicitly_trimmed() -> anyhow::Result<()> {
let log = SqliteLog::connect("log-test", "sqlite::memory:").await?;
log.register_topic(TopicConfig::new("events", 1)?).await?;
let mut records = Vec::new();
for index in 1..=5 {
records.push(
log.append(
"events",
PartitionKey::from(format!("event-{index}")),
format!("payload-{index}").into_bytes(),
)
.await?,
);
}
assert_eq!(
records.iter().map(|r| r.offset).collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5]
);
log.commit_offset("events", "audit", PartitionId(0), 6)
.await?;
log.commit_offset("events", "search", PartitionId(0), 3)
.await?;
assert_eq!(log.compact_partition("events", PartitionId(0)).await?, 0);
assert_eq!(
log.read("events", PartitionId(0), 1, 10).await?,
records.clone()
);
assert_eq!(
log.safe_cleanup_offset("events", PartitionId(0)).await?,
Some(3)
);
assert!(log
.delete_before("events", PartitionId(0), 4)
.await
.is_err());
assert_eq!(log.delete_before("events", PartitionId(0), 3).await?, 2);
assert_eq!(
log.read("events", PartitionId(0), 1, 10).await?,
records[2..].to_vec()
);
log.commit_offset("events", "search", PartitionId(0), 6)
.await?;
assert_eq!(log.delete_before("events", PartitionId(0), 6).await?, 3);
assert_eq!(log.read("events", PartitionId(0), 1, 10).await?, Vec::new());
Ok(())
}
#[tokio::test]
async fn consumer_lag_reports_each_groups_own_distance_from_the_end() -> anyhow::Result<()> {
let log = SqliteLog::connect("log-test", "sqlite::memory:").await?;
log.register_topic(TopicConfig::new("events", 1)?).await?;
assert_eq!(log.consumer_lag("events").await?, Vec::new());
for index in 1..=5 {
log.append(
"events",
PartitionKey::from(format!("event-{index}")),
format!("payload-{index}").into_bytes(),
)
.await?;
}
log.commit_offset("events", "caught-up", PartitionId(0), 6)
.await?;
log.commit_offset("events", "behind", PartitionId(0), 3)
.await?;
let mut lag = log.consumer_lag("events").await?;
lag.sort_by(|a, b| a.group.cmp(&b.group));
assert_eq!(
lag,
vec![
ConsumerLag {
partition: PartitionId(0),
group: "behind".to_string(),
end_offset: 6,
committed_offset: 3,
lag: 3,
},
ConsumerLag {
partition: PartitionId(0),
group: "caught-up".to_string(),
end_offset: 6,
committed_offset: 6,
lag: 0,
},
]
);
Ok(())
}
#[tokio::test]
async fn compacted_topic_keeps_only_the_latest_record_per_key() -> anyhow::Result<()> {
let log = SqliteLog::connect("log-test", "sqlite::memory:").await?;
log.register_topic(TopicConfig::new("orders", 1)?).await?;
let order_one = PartitionKey::from("order-1");
let order_two = PartitionKey::from("order-2");
let order_three = PartitionKey::from("order-3");
log.append("orders", order_one.clone(), b"v1".to_vec())
.await?; log.append("orders", order_two.clone(), b"v1".to_vec())
.await?; log.append("orders", order_one.clone(), b"v2".to_vec())
.await?; log.append("orders", order_three.clone(), b"only".to_vec())
.await?; let order_one_latest = log
.append("orders", order_one.clone(), b"v3".to_vec())
.await?; let order_two_latest = log
.append("orders", order_two.clone(), b"v2".to_vec())
.await?;
log.commit_offset("orders", "projector", PartitionId(0), 7)
.await?;
log.commit_offset("orders", "slow-reader", PartitionId(0), 4)
.await?;
assert_eq!(
log.safe_cleanup_offset("orders", PartitionId(0)).await?,
Some(4)
);
assert_eq!(log.compact_partition("orders", PartitionId(0)).await?, 3);
let order_three_record = log
.read("orders", PartitionId(0), 1, 1)
.await?
.into_iter()
.next()
.ok_or_else(|| anyhow::anyhow!("order-3's only record should survive compaction"))?;
assert_eq!(order_three_record.key, order_three);
assert_eq!(order_three_record.payload, b"only");
assert_eq!(log.compact_partition("orders", PartitionId(0)).await?, 0);
log.commit_offset("orders", "slow-reader", PartitionId(0), 7)
.await?;
assert_eq!(log.compact_partition("orders", PartitionId(0)).await?, 0);
assert_eq!(
log.read("orders", PartitionId(0), 1, 10).await?,
vec![order_three_record, order_one_latest, order_two_latest]
);
Ok(())
}
#[tokio::test]
async fn append_and_replay_keep_partition_offsets() -> anyhow::Result<()> {
let log = SqliteLog::connect("log-test", "sqlite::memory:").await?;
log.register_topic(TopicConfig::new("orders", 1)?).await?;
let first = log
.append("orders", PartitionKey::from("one"), b"first".to_vec())
.await?;
let second = log
.append("orders", PartitionKey::from("two"), b"second".to_vec())
.await?;
assert_eq!((first.offset, second.offset), (1, 2));
assert_eq!(
log.read("orders", PartitionId(0), 1, 10).await?,
vec![first.clone(), second.clone()]
);
assert_eq!(
log.committed_offset("orders", "billing", PartitionId(0))
.await?,
1
);
log.commit_offset("orders", "billing", PartitionId(0), 3)
.await?;
log.commit_offset("orders", "billing", PartitionId(0), 2)
.await?;
assert_eq!(
log.committed_offset("orders", "billing", PartitionId(0))
.await?,
3
);
let replacement = log
.append("orders", PartitionKey::from("one"), b"replacement".to_vec())
.await?;
log.commit_offset("orders", "billing", PartitionId(0), 4)
.await?;
assert_eq!(log.compact_partition("orders", PartitionId(0)).await?, 1);
assert_eq!(
log.read("orders", PartitionId(0), 1, 10).await?,
vec![second.clone(), replacement]
);
let claim = log
.claim_partition("orders", "billing", "member-a", PartitionId(0))
.await?
.ok_or_else(|| anyhow::anyhow!("member-a did not claim its partition"))?;
assert!(log.renew_claim(&claim).await?);
assert!(log.commit_claimed_offset(&claim, 4).await?);
assert!(log
.claim_partition("orders", "billing", "member-b", PartitionId(0))
.await?
.is_none());
log.leave_group("orders", "billing", "member-a").await?;
assert!(log
.claim_partition("orders", "billing", "member-b", PartitionId(0))
.await?
.is_some());
Ok(())
}
}
#[cfg(all(test, feature = "postgres"))]
mod postgres_tests {
use super::*;
#[tokio::test]
async fn append_and_replay_keep_partition_offsets() -> anyhow::Result<()> {
let url = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let log = PostgresLog::connect(format!("log-test-{}", crate::generate_id()), &url).await?;
log.register_topic(TopicConfig::new("orders", 1)?).await?;
let first = log
.append("orders", PartitionKey::from("one"), b"first".to_vec())
.await?;
let second = log
.append("orders", PartitionKey::from("two"), b"second".to_vec())
.await?;
assert_eq!((first.offset, second.offset), (1, 2));
assert_eq!(
log.read("orders", PartitionId(0), 1, 10).await?,
vec![first.clone(), second.clone()]
);
assert_eq!(
log.committed_offset("orders", "billing", PartitionId(0))
.await?,
1
);
log.commit_offset("orders", "billing", PartitionId(0), 3)
.await?;
log.commit_offset("orders", "billing", PartitionId(0), 2)
.await?;
assert_eq!(
log.committed_offset("orders", "billing", PartitionId(0))
.await?,
3
);
let replacement = log
.append("orders", PartitionKey::from("one"), b"replacement".to_vec())
.await?;
log.commit_offset("orders", "billing", PartitionId(0), 4)
.await?;
assert_eq!(log.compact_partition("orders", PartitionId(0)).await?, 1);
assert_eq!(
log.read("orders", PartitionId(0), 1, 10).await?,
vec![second.clone(), replacement]
);
let claim = log
.claim_partition("orders", "billing", "member-a", PartitionId(0))
.await?
.ok_or_else(|| anyhow::anyhow!("member-a did not claim its partition"))?;
assert!(log.renew_claim(&claim).await?);
assert!(log.commit_claimed_offset(&claim, 4).await?);
assert!(log
.claim_partition("orders", "billing", "member-b", PartitionId(0))
.await?
.is_none());
log.leave_group("orders", "billing", "member-a").await?;
assert!(log
.claim_partition("orders", "billing", "member-b", PartitionId(0))
.await?
.is_some());
Ok(())
}
#[tokio::test]
async fn consumer_lag_reports_each_groups_own_distance_from_the_end() -> anyhow::Result<()> {
let url = std::env::var("LATER_POSTGRES_TEST_URL")
.unwrap_or_else(|_| "postgres://test:test@127.0.0.1:55432/later_test".to_owned());
let log = PostgresLog::connect(format!("log-test-{}", crate::generate_id()), &url).await?;
log.register_topic(TopicConfig::new("events", 1)?).await?;
assert_eq!(log.consumer_lag("events").await?, Vec::new());
for index in 1..=5 {
log.append(
"events",
PartitionKey::from(format!("event-{index}")),
format!("payload-{index}").into_bytes(),
)
.await?;
}
log.commit_offset("events", "caught-up", PartitionId(0), 6)
.await?;
log.commit_offset("events", "behind", PartitionId(0), 3)
.await?;
let mut lag = log.consumer_lag("events").await?;
lag.sort_by(|a, b| a.group.cmp(&b.group));
assert_eq!(
lag,
vec![
ConsumerLag {
partition: PartitionId(0),
group: "behind".to_string(),
end_offset: 6,
committed_offset: 3,
lag: 3,
},
ConsumerLag {
partition: PartitionId(0),
group: "caught-up".to_string(),
end_offset: 6,
committed_offset: 6,
lag: 0,
},
]
);
Ok(())
}
}