use crate::common;
use reliar_core::{Classify, CorrelationId, Envelope, FailureKind, Message, MessageType};
use reliar_inbox::{InboxClaim, InboxFailure, InboxMessage, InboxScope, InboxState, InboxStore};
use reliar_store_postgres::{PostgresInboxSettings, PostgresInboxStore};
#[derive(Debug)]
struct HandlerFailed;
impl std::fmt::Display for HandlerFailed {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "handler failed, by design")
}
}
impl std::error::Error for HandlerFailed {}
#[derive(serde::Serialize, serde::Deserialize)]
struct OrderCreated;
impl Message for OrderCreated {
const TYPE: &'static str = "orders.created";
const VERSION: u16 = 1;
}
async fn surrogate_id_and_the_unique_key() {
let pool = common::fresh_db().await;
let store = PostgresInboxStore::connect(pool.clone(), PostgresInboxSettings::default())
.await
.unwrap();
let envelope = Envelope::builder(OrderCreated).build();
let message = InboxMessage::from_envelope(&envelope);
let scope_a = InboxScope::new("consumer-a").unwrap();
let scope_b = InboxScope::new("consumer-b").unwrap();
let mut tx_a = pool.begin().await.unwrap();
store.claim(&mut tx_a, &scope_a, message).await.unwrap();
tx_a.commit().await.unwrap();
let mut tx_b = pool.begin().await.unwrap();
store.claim(&mut tx_b, &scope_b, message).await.unwrap();
tx_b.commit().await.unwrap();
let record_a = store.find(&scope_a, envelope.id).await.unwrap().unwrap();
let record_b = store.find(&scope_b, envelope.id).await.unwrap().unwrap();
assert_ne!(
record_a.id, record_b.id,
"two scopes claiming the same message get distinct row ids"
);
let err = sqlx::query!(
"INSERT INTO inbox (id, scope, message_id, message_type, message_version, \
conversation_id) \
VALUES ($1, $2, $3, 'orders.created', 1, $3)",
uuid::Uuid::now_v7(),
scope_a.as_str(),
envelope.id.as_uuid(),
)
.execute(&pool)
.await
.unwrap_err();
assert!(
matches!(&err, sqlx::Error::Database(db) if db.code().as_deref() == Some("23505")),
"expected a unique-violation SQLSTATE 23505, got {err:?}"
);
let err = sqlx::query!(
"INSERT INTO inbox (id, scope, message_id, message_type, message_version, \
conversation_id, completed_at, dead_at) \
VALUES ($1, 'consumer-c', $2, 'orders.created', 1, $2, now(), now())",
uuid::Uuid::now_v7(),
uuid::Uuid::now_v7(),
)
.execute(&pool)
.await
.unwrap_err();
assert!(
matches!(&err, sqlx::Error::Database(db) if db.constraint() == Some("ck_inbox_terminal")),
"expected ck_inbox_terminal to reject completed_at and dead_at both set, got {err:?}"
);
}
async fn trace_fields_persisted() {
let pool = common::fresh_db().await;
let store = PostgresInboxStore::connect(pool.clone(), PostgresInboxSettings::default())
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let bare = Envelope::builder(OrderCreated).build();
let mut tx = pool.begin().await.unwrap();
store
.claim(&mut tx, &scope, InboxMessage::from_envelope(&bare))
.await
.unwrap();
tx.commit().await.unwrap();
let record = store.find(&scope, bare.id).await.unwrap().unwrap();
assert_eq!(record.message_type.name(), "orders.created");
assert_eq!(record.message_type.version(), 1);
assert_eq!(
record.conversation_id,
bare.metadata.correlation.conversation_id
);
assert!(record.correlation_id.is_none());
assert!(record.causation_id.is_none());
let correlation_id = CorrelationId::parse("checkout-42").unwrap();
let cause = reliar_core::MessageId::new();
let envelope = Envelope::builder(OrderCreated)
.correlation_id(correlation_id.clone())
.causation(cause)
.build();
let mut tx = pool.begin().await.unwrap();
store
.claim(&mut tx, &scope, InboxMessage::from_envelope(&envelope))
.await
.unwrap();
tx.commit().await.unwrap();
let record = store.find(&scope, envelope.id).await.unwrap().unwrap();
assert_eq!(record.correlation_id, Some(correlation_id));
assert_eq!(record.causation_id, Some(cause));
assert_ne!(record.causation_id, Some(record.message_id));
}
async fn fail_creates_a_fully_populated_row() {
let pool = common::fresh_db().await;
let store = PostgresInboxStore::connect(pool.clone(), PostgresInboxSettings::default())
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let envelope = Envelope::builder(OrderCreated).build();
let message = InboxMessage::from_envelope(&envelope);
{
let mut tx = pool.begin().await.unwrap();
store.claim(&mut tx, &scope, message).await.unwrap();
}
store.fail(&scope, message, &HandlerFailed).await.unwrap();
let record = store.find(&scope, envelope.id).await.unwrap().unwrap();
assert_eq!(record.message_type, envelope.message_type);
assert_eq!(
record.conversation_id,
envelope.metadata.correlation.conversation_id
);
assert_eq!(record.attempts, 1);
}
async fn bounded_retries_in_sql() {
let pool = common::fresh_db().await;
let settings = PostgresInboxSettings::default().max_attempts(3);
let store = PostgresInboxStore::connect(pool.clone(), settings)
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let envelope = Envelope::builder(OrderCreated).build();
let message = InboxMessage::from_envelope(&envelope);
let first = store.fail(&scope, message, &HandlerFailed).await.unwrap();
assert_eq!(first, InboxFailure::Recorded { attempts: 1 });
let second = store.fail(&scope, message, &HandlerFailed).await.unwrap();
assert_eq!(second, InboxFailure::Recorded { attempts: 2 });
let third = store.fail(&scope, message, &HandlerFailed).await.unwrap();
let dead_at = match third {
InboxFailure::Dead {
attempts: 3,
dead_at,
..
} => dead_at,
other => panic!("expected Dead{{attempts: 3}}, got {other:?}"),
};
let mut tx = pool.begin().await.unwrap();
let claim = store.claim(&mut tx, &scope, message).await.unwrap();
match claim {
InboxClaim::Dead {
attempts: 3,
dead_at: claimed_dead_at,
..
} => assert_eq!(claimed_dead_at, dead_at),
other => panic!("expected InboxClaim::Dead, got {other:?}"),
}
tx.rollback().await.unwrap();
let settings_one = PostgresInboxSettings::default().max_attempts(1);
let store_one = PostgresInboxStore::connect(pool.clone(), settings_one)
.await
.unwrap();
let envelope2 = Envelope::builder(OrderCreated).build();
let message2 = InboxMessage::from_envelope(&envelope2);
let outcome = store_one
.fail(&scope, message2, &HandlerFailed)
.await
.unwrap();
assert!(matches!(outcome, InboxFailure::Dead { attempts: 1, .. }));
}
async fn the_dead_transition_is_atomic() {
const N: u32 = 10;
let pool = common::fresh_db().await;
let settings = PostgresInboxSettings::default().max_attempts(N);
let store = PostgresInboxStore::connect(pool.clone(), settings)
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let envelope = Envelope::builder(OrderCreated).build();
let message_id = envelope.id;
let conversation_id = envelope.metadata.correlation.conversation_id;
let message_type: &'static MessageType = Box::leak(Box::new(envelope.message_type.clone()));
let message = InboxMessage::new(message_id, message_type).conversation(conversation_id);
let mut tasks = Vec::new();
for _ in 0..N {
let store = store.clone();
let scope = scope.clone();
tasks.push(tokio::spawn(async move {
store.fail(&scope, message, &HandlerFailed).await.unwrap()
}));
}
let mut dead_count = 0;
let mut dead_at_values = std::collections::HashSet::new();
for task in tasks {
if let InboxFailure::Dead { dead_at, .. } = task.await.unwrap() {
dead_count += 1;
dead_at_values.insert(dead_at);
}
}
assert_eq!(dead_count, 1, "exactly one caller must observe Dead");
assert_eq!(dead_at_values.len(), 1);
let record = store.find(&scope, message_id).await.unwrap().unwrap();
assert_eq!(record.attempts, N);
assert!(record.dead_at.is_some());
}
async fn fail_on_an_already_dead_row_keeps_the_original_dead_at() {
let pool = common::fresh_db().await;
let settings = PostgresInboxSettings::default().max_attempts(1);
let store = PostgresInboxStore::connect(pool.clone(), settings)
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let envelope = Envelope::builder(OrderCreated).build();
let message = InboxMessage::from_envelope(&envelope);
let first = store.fail(&scope, message, &HandlerFailed).await.unwrap();
let dead_at = match first {
InboxFailure::Dead { dead_at, .. } => dead_at,
other => panic!("expected Dead, got {other:?}"),
};
let second = store.fail(&scope, message, &HandlerFailed).await.unwrap();
match second {
InboxFailure::Dead {
attempts: 2,
dead_at: second_dead_at,
..
} => assert_eq!(
second_dead_at, dead_at,
"dead_at must not move on a later fail"
),
other => panic!("expected Dead {{ attempts: 2, .. }}, got {other:?}"),
}
let record = store.find(&scope, envelope.id).await.unwrap().unwrap();
assert_eq!(record.attempts, 2);
assert_eq!(record.dead_at, Some(dead_at));
}
async fn complete_on_a_dead_row_is_not_claimed() {
let pool = common::fresh_db().await;
let settings = PostgresInboxSettings::default().max_attempts(1);
let store = PostgresInboxStore::connect(pool.clone(), settings)
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let envelope = Envelope::builder(OrderCreated).build();
let message = InboxMessage::from_envelope(&envelope);
store.fail(&scope, message, &HandlerFailed).await.unwrap();
let mut tx = pool.begin().await.unwrap();
let err = store
.complete(&mut tx, &scope, envelope.id)
.await
.unwrap_err();
assert_eq!(err.kind(), FailureKind::Permanent);
tx.rollback().await.unwrap();
}
async fn max_attempts_zero_is_rejected_at_connect() {
let pool = common::fresh_db().await;
let settings = PostgresInboxSettings::default().max_attempts(0);
let err = PostgresInboxStore::connect(pool, settings)
.await
.unwrap_err();
assert_eq!(err.kind(), FailureKind::Permanent);
assert!(err.to_string().contains("max_attempts"));
}
async fn state_against_real_rows() {
use reliar_inbox::InboxDeadLetters;
let pool = common::fresh_db().await;
let settings = PostgresInboxSettings::default().max_attempts(1);
let store = PostgresInboxStore::connect(pool.clone(), settings)
.await
.unwrap();
let scope = InboxScope::new("orders-projection").unwrap();
let claimed_envelope = Envelope::builder(OrderCreated).build();
let mut tx = pool.begin().await.unwrap();
store
.claim(
&mut tx,
&scope,
InboxMessage::from_envelope(&claimed_envelope),
)
.await
.unwrap();
tx.commit().await.unwrap();
let claimed = store
.find(&scope, claimed_envelope.id)
.await
.unwrap()
.unwrap();
assert_eq!(claimed.state(), InboxState::Claimed);
let retrying_settings = PostgresInboxSettings::default().max_attempts(5);
let retrying_store = PostgresInboxStore::connect(pool.clone(), retrying_settings)
.await
.unwrap();
let retrying_envelope = Envelope::builder(OrderCreated).build();
let retrying_message = InboxMessage::from_envelope(&retrying_envelope);
retrying_store
.fail(&scope, retrying_message, &HandlerFailed)
.await
.unwrap();
let retrying = retrying_store
.find(&scope, retrying_envelope.id)
.await
.unwrap()
.unwrap();
assert_eq!(retrying.state(), InboxState::Retrying);
let completed_envelope = Envelope::builder(OrderCreated).build();
let mut tx = pool.begin().await.unwrap();
store
.claim(
&mut tx,
&scope,
InboxMessage::from_envelope(&completed_envelope),
)
.await
.unwrap();
store
.complete(&mut tx, &scope, completed_envelope.id)
.await
.unwrap();
tx.commit().await.unwrap();
let completed = store
.find(&scope, completed_envelope.id)
.await
.unwrap()
.unwrap();
assert_eq!(completed.state(), InboxState::Completed);
let dead_envelope = Envelope::builder(OrderCreated).build();
let dead_message = InboxMessage::from_envelope(&dead_envelope);
store
.fail(&scope, dead_message, &HandlerFailed)
.await
.unwrap();
let dead = store.find(&scope, dead_envelope.id).await.unwrap().unwrap();
assert_eq!(dead.state(), InboxState::Dead);
store.retry_dead(&[dead.id]).await.unwrap();
let retried = store.find(&scope, dead_envelope.id).await.unwrap().unwrap();
assert_eq!(retried.state(), InboxState::Claimed);
}
pub(crate) fn trials(rt: &'static tokio::runtime::Runtime) -> Vec<libtest_mimic::Trial> {
vec![
libtest_mimic::Trial::test("inbox_dead::surrogate_id_and_the_unique_key", move || {
rt.block_on(surrogate_id_and_the_unique_key());
Ok(())
}),
libtest_mimic::Trial::test("inbox_dead::trace_fields_persisted", move || {
rt.block_on(trace_fields_persisted());
Ok(())
}),
libtest_mimic::Trial::test(
"inbox_dead::fail_creates_a_fully_populated_row",
move || {
rt.block_on(fail_creates_a_fully_populated_row());
Ok(())
},
),
libtest_mimic::Trial::test("inbox_dead::bounded_retries_in_sql", move || {
rt.block_on(bounded_retries_in_sql());
Ok(())
}),
libtest_mimic::Trial::test("inbox_dead::the_dead_transition_is_atomic", move || {
rt.block_on(the_dead_transition_is_atomic());
Ok(())
}),
libtest_mimic::Trial::test(
"inbox_dead::fail_on_an_already_dead_row_keeps_the_original_dead_at",
move || {
rt.block_on(fail_on_an_already_dead_row_keeps_the_original_dead_at());
Ok(())
},
),
libtest_mimic::Trial::test(
"inbox_dead::complete_on_a_dead_row_is_not_claimed",
move || {
rt.block_on(complete_on_a_dead_row_is_not_claimed());
Ok(())
},
),
libtest_mimic::Trial::test(
"inbox_dead::max_attempts_zero_is_rejected_at_connect",
move || {
rt.block_on(max_attempts_zero_is_rejected_at_connect());
Ok(())
},
),
libtest_mimic::Trial::test("inbox_dead::state_against_real_rows", move || {
rt.block_on(state_against_real_rows());
Ok(())
}),
]
}