use crate::common;
use std::time::Duration;
use reliar_outbox::{
AcquireRequest, Classify, FailedRecord, FailureKind, FailureOutcome, OutboxStore, WorkerId,
};
use reliar_store_postgres::{PostgresOutboxError, PostgresOutboxSettings, PostgresOutboxStore};
async fn a_generous_non_zero_timeout_behaves_identically_to_zero_for_acquire_complete_and_fail() {
let pool = common::fresh_db().await;
let settings = PostgresOutboxSettings::default().statement_timeout(Duration::from_secs(30));
let store = PostgresOutboxStore::with_settings(pool.clone(), settings);
let envelopes = common::seed(&store, &pool, 2).await;
let worker = WorkerId::generate();
let batch = store
.acquire(AcquireRequest::new(worker.clone()).batch_size(10))
.await
.unwrap();
assert_eq!(batch.records.len(), 2, "acquire behaves the same wrapped");
let ids: Vec<_> = batch.records.iter().map(|r| r.envelope.id).collect();
assert!(ids.contains(&envelopes[0].id));
assert!(ids.contains(&envelopes[1].id));
let affected = store
.complete(&worker, &[batch.records[0].record_ref()])
.await
.unwrap();
assert_eq!(affected, 1, "complete behaves the same wrapped");
let affected = store
.fail(
&worker,
&[FailedRecord::new(
batch.records[1].record_ref(),
"transient, retry shortly",
FailureOutcome::Retry {
delay: Duration::from_millis(1),
},
)],
)
.await
.unwrap();
assert_eq!(affected, 1, "fail behaves the same wrapped");
}
async fn a_too_short_timeout_yields_a_clean_transient_error_instead_of_hanging() {
let pool = common::fresh_db().await;
let store = PostgresOutboxStore::new(pool.clone());
let envelopes = common::seed(&store, &pool, 1).await;
let worker = WorkerId::generate();
let batch = store
.acquire(AcquireRequest::new(worker.clone()))
.await
.unwrap();
let record_ref = batch.records[0].record_ref();
assert_eq!(batch.records[0].envelope.id, envelopes[0].id);
let mut holder = pool.begin().await.unwrap();
sqlx::query("SELECT id FROM outbox WHERE id = $1 FOR UPDATE")
.bind(record_ref.id.as_uuid())
.fetch_all(&mut *holder)
.await
.unwrap();
let settings = PostgresOutboxSettings::default().statement_timeout(Duration::from_millis(50));
let timeout_store = PostgresOutboxStore::with_settings(pool.clone(), settings);
let result = tokio::time::timeout(
Duration::from_secs(5),
timeout_store.complete(&worker, &[record_ref]),
)
.await
.expect("must return well within the outer bound, not hang");
let err = result.expect_err("a canceled statement must surface as a typed error");
assert!(
matches!(err, PostgresOutboxError::Database { .. }),
"expected Database, got {err:?}"
);
assert_eq!(
err.kind(),
FailureKind::Transient,
"a canceled statement is safe to retry"
);
holder.rollback().await.unwrap();
}
pub(crate) fn trials(rt: &'static tokio::runtime::Runtime) -> Vec<libtest_mimic::Trial> {
vec![
libtest_mimic::Trial::test(
"outbox_statement_timeout::a_generous_non_zero_timeout_behaves_identically_to_zero_for_acquire_complete_and_fail",
move || {
rt.block_on(a_generous_non_zero_timeout_behaves_identically_to_zero_for_acquire_complete_and_fail());
Ok(())
},
),
libtest_mimic::Trial::test(
"outbox_statement_timeout::a_too_short_timeout_yields_a_clean_transient_error_instead_of_hanging",
move || {
rt.block_on(
a_too_short_timeout_yields_a_clean_transient_error_instead_of_hanging(),
);
Ok(())
},
),
]
}