use crate::context::DjogiContext;
use crate::error::DbError;
use crate::{DjogiError, HeerId};
use time::OffsetDateTime;
pub const MAX_RETRY_COUNT: i32 = 10;
#[derive(Debug, Clone)]
pub struct OutboxRow {
pub id: HeerId,
pub row_id: String,
pub action: String,
pub payload: serde_json::Value,
pub created_at: OffsetDateTime,
}
fn validate_table_ident(name: &str) -> Result<(), DjogiError> {
use crate::ident::IdentError;
crate::ident::check_user_supplied_ident(name, false).map_err(|e| {
let msg = match e {
IdentError::Empty | IdentError::TooLong { .. } => {
format!("invalid outbox table name {name:?}: must be 1–63 bytes")
}
IdentError::BadFirst { .. } => format!(
"invalid outbox table name {name:?}: first character must be an ASCII letter or underscore"
),
IdentError::BadByte { byte, .. } => format!(
"invalid outbox table name {name:?}: contains disallowed character '{}'",
byte as char
),
IdentError::Reserved => unreachable!("check_user_supplied_ident(reserved=false) cannot return Reserved"),
IdentError::ReservedDjogiPrefix => format!(
"invalid outbox table name {name:?}: starts with the framework-reserved \
`__djogi_` prefix; choose a different name"
),
};
DjogiError::Db(DbError::other(msg))
})
}
pub async fn claim_pending(
ctx: &mut DjogiContext,
outbox_table: &str,
batch_size: u32,
lease_duration: time::Duration,
) -> Result<Vec<OutboxRow>, DjogiError> {
validate_table_ident(outbox_table)?;
let lease_secs = lease_duration.whole_seconds();
let sql = format!(
"UPDATE {outbox_table} \
SET state = 'processing', leased_until = now() + make_interval(secs => {lease_secs}) \
WHERE id IN ( \
SELECT id FROM {outbox_table} \
WHERE state = 'pending' \
AND (leased_until IS NULL OR leased_until <= now()) \
ORDER BY created_at \
LIMIT $1 \
FOR UPDATE SKIP LOCKED \
) \
RETURNING id, row_id::text, action, payload, created_at"
);
let batch_size_i64 = batch_size as i64;
let params: &[&(dyn postgres_types::ToSql + Sync)] = &[&batch_size_i64];
let rows = ctx.query_all(&sql, params).await?;
let mut result = Vec::with_capacity(rows.len());
for row in rows {
let id_raw: i64 = row.try_get(0).map_err(|e| {
DjogiError::Db(DbError::other(format!("claim_pending: decode id: {e}")))
})?;
let row_id: String = row.try_get(1).map_err(|e| {
DjogiError::Db(DbError::other(format!("claim_pending: decode row_id: {e}")))
})?;
let action: String = row.try_get(2).map_err(|e| {
DjogiError::Db(DbError::other(format!("claim_pending: decode action: {e}")))
})?;
let payload: serde_json::Value = row.try_get(3).map_err(|e| {
DjogiError::Db(DbError::other(format!(
"claim_pending: decode payload: {e}"
)))
})?;
let created_at: OffsetDateTime = row.try_get(4).map_err(|e| {
DjogiError::Db(DbError::other(format!(
"claim_pending: decode created_at: {e}"
)))
})?;
let id = HeerId::from_i64(id_raw).map_err(|e| {
DjogiError::Db(DbError::other(format!(
"claim_pending: invalid HeerId value {id_raw}: {e}"
)))
})?;
result.push(OutboxRow {
id,
row_id,
action,
payload,
created_at,
});
}
Ok(result)
}
pub async fn mark_published(
ctx: &mut DjogiContext,
outbox_table: &str,
row_id: HeerId,
) -> Result<(), DjogiError> {
validate_table_ident(outbox_table)?;
let sql = format!(
"UPDATE {outbox_table} \
SET state = 'published' \
WHERE id = $1 AND state = 'processing'"
);
let id_raw = row_id.as_i64();
let params: &[&(dyn postgres_types::ToSql + Sync)] = &[&id_raw];
ctx.execute(&sql, params).await?;
Ok(())
}
pub async fn mark_failed(
ctx: &mut DjogiContext,
outbox_table: &str,
row_id: HeerId,
error_message: &str,
retryable: bool,
) -> Result<(), DjogiError> {
validate_table_ident(outbox_table)?;
let fetch_sql =
format!("SELECT retry_count FROM {outbox_table} WHERE id = $1 AND state = 'processing'");
let id_raw = row_id.as_i64();
let fetch_params: &[&(dyn postgres_types::ToSql + Sync)] = &[&id_raw];
let maybe_row = ctx.query_opt(&fetch_sql, fetch_params).await?;
let current_retry: i32 = match maybe_row {
Some(row) => row.try_get(0).map_err(|e| {
DjogiError::Db(DbError::other(format!(
"mark_failed: decode retry_count: {e}"
)))
})?,
None => {
return Ok(());
}
};
let transition_to_pending = retryable && current_retry < MAX_RETRY_COUNT;
if transition_to_pending {
let new_retry = current_retry.saturating_add(1);
let shift = new_retry.clamp(0, 10) as u32;
let backoff_secs: i64 = (1i64 << shift).min(1024);
let sql = format!(
"UPDATE {outbox_table} \
SET state = 'pending', retry_count = retry_count + 1, \
leased_until = now() + make_interval(secs => {backoff_secs}), \
failed_reason = $2 \
WHERE id = $1 AND state = 'processing'"
);
let params: &[&(dyn postgres_types::ToSql + Sync)] = &[&id_raw, &error_message];
ctx.execute(&sql, params).await?;
} else {
let sql = format!(
"UPDATE {outbox_table} \
SET state = 'failed', failed_reason = $2 \
WHERE id = $1 AND state = 'processing'"
);
let params: &[&(dyn postgres_types::ToSql + Sync)] = &[&id_raw, &error_message];
ctx.execute(&sql, params).await?;
}
Ok(())
}
pub async fn recover_stale(
ctx: &mut DjogiContext,
outbox_table: &str,
stale_threshold: time::Duration,
) -> Result<u64, DjogiError> {
validate_table_ident(outbox_table)?;
let threshold_secs = stale_threshold.whole_seconds();
let sql = format!(
"UPDATE {outbox_table} \
SET state = 'pending', leased_until = NULL \
WHERE state = 'processing' AND leased_until < now() - make_interval(secs => {threshold_secs})"
);
let params: &[&(dyn postgres_types::ToSql + Sync)] = &[];
let count = ctx.execute(&sql, params).await?;
Ok(count)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn valid_idents_pass() {
assert!(validate_table_ident("worker_outbox").is_ok());
assert!(validate_table_ident("accounts_outbox").is_ok());
assert!(validate_table_ident("orders_outbox_v2").is_ok());
assert!(validate_table_ident("_private").is_ok());
assert!(validate_table_ident("a").is_ok());
}
#[test]
fn empty_ident_rejected() {
assert!(validate_table_ident("").is_err());
}
#[test]
fn ident_over_63_bytes_rejected() {
let long = "a".repeat(64);
assert!(validate_table_ident(&long).is_err());
}
#[test]
fn ident_exactly_63_bytes_ok() {
let ok = "a".repeat(63);
assert!(validate_table_ident(&ok).is_ok());
}
#[test]
fn digit_first_byte_rejected() {
assert!(validate_table_ident("1outbox").is_err());
}
#[test]
fn hyphen_rejected() {
assert!(validate_table_ident("worker-outbox").is_err());
}
#[test]
fn dot_rejected() {
assert!(validate_table_ident("schema.table").is_err());
}
#[test]
fn space_rejected() {
assert!(validate_table_ident("worker outbox").is_err());
}
#[test]
fn djogi_reserved_prefix_rejected() {
let err = validate_table_ident("__djogi_outbox").expect_err("must reject reserved prefix");
let msg = err.to_string();
assert!(msg.contains("`__djogi_` prefix"), "got: {msg}");
assert!(validate_table_ident("__djogi_").is_err());
assert!(validate_table_ident("__myadopter_outbox").is_ok());
assert!(validate_table_ident("_djogi_outbox").is_ok());
}
}