use reliar_core::Serializer;
use reliar_outbox::{AcquireRequest, AcquiredBatch, PoisonedRow};
use super::error::PostgresOutboxError;
use crate::records::{RawRow, decode_row};
use super::PostgresOutboxStore;
pub(super) async fn acquire<Ser: Serializer + Send + Sync + 'static>(
store: &PostgresOutboxStore<Ser>,
request: AcquireRequest,
) -> Result<AcquiredBatch, PostgresOutboxError> {
let batch_size = i64::from(request.batch_size);
let lease_ms = i64::try_from(request.lease.as_millis()).unwrap_or(i64::MAX);
let worker = request.worker.as_str();
let rows = if store.settings.statement_timeout.is_zero() {
claim_rows(&store.pool, batch_size, worker, lease_ms)
.await
.map_err(|e| store.map_err(e))?
} else {
let mut tx = store.pool.begin().await.map_err(|e| store.map_err(e))?;
let timeout_ms = i64::try_from(store.settings.statement_timeout.as_millis())
.unwrap_or(i64::MAX)
.to_string();
sqlx::query_scalar!(
"SELECT set_config('statement_timeout', $1, true)",
timeout_ms
)
.fetch_one(&mut *tx)
.await
.map_err(|e| store.map_err(e))?;
let rows = claim_rows(&mut *tx, batch_size, worker, lease_ms)
.await
.map_err(|e| store.map_err(e))?;
tx.commit().await.map_err(|e| store.map_err(e))?;
rows
};
let mut records = Vec::with_capacity(rows.len());
let mut poisoned = Vec::new();
let mut poisoned_ids = Vec::new();
let mut poisoned_errors = Vec::new();
for raw in rows {
match decode_row(raw) {
Ok(record) => records.push(record),
Err(err) => {
poisoned_ids.push(err.id.as_uuid());
poisoned_errors.push(crate::records::truncate_last_error(err.detail.clone()));
poisoned.push(PoisonedRow::new(err.id, err.message_id, err.detail));
}
}
}
if !poisoned_ids.is_empty() {
let undecodable =
crate::records::encode_dead_reason(reliar_outbox::DeadReason::Undecodable);
let sweep_result = if store.settings.statement_timeout.is_zero() {
poison_sweep_rows(
&store.pool,
&poisoned_ids,
&poisoned_errors,
worker,
undecodable,
)
.await
} else {
async {
let mut tx = store.pool.begin().await?;
store.set_local_timeout_raw(&mut tx).await?;
poison_sweep_rows(
&mut *tx,
&poisoned_ids,
&poisoned_errors,
worker,
undecodable,
)
.await?;
tx.commit().await
}
.await
};
if let Err(err) = sweep_result {
tracing::warn!(
target: "reliar.outbox.acquire",
worker_id = %worker,
poisoned_count = poisoned_ids.len(),
error = %store.map_err(err),
"poison sweep failed; the claimed batch is returned and the undecodable rows \
stay leased until their lease lapses"
);
}
}
Ok(AcquiredBatch::new(records, poisoned))
}
async fn poison_sweep_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
poisoned_ids: &[uuid::Uuid],
poisoned_errors: &[String],
worker: &str,
undecodable: &str,
) -> Result<(), sqlx::Error> {
sqlx::query!(
r#"UPDATE outbox o
SET dead_at = now(),
dead_reason = $4,
last_error = f.err,
locked_by = NULL,
locked_until = NULL,
updated_at = now()
FROM UNNEST($1::uuid[], $2::text[]) AS f(id, err)
WHERE o.id = f.id AND o.locked_by = $3"#,
poisoned_ids,
poisoned_errors,
worker,
undecodable,
)
.execute(executor)
.await?;
Ok(())
}
async fn claim_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
batch_size: i64,
worker: &str,
lease_ms: i64,
) -> Result<Vec<RawRow>, sqlx::Error> {
sqlx::query_as!(
RawRow,
r#"WITH claimed AS (
SELECT id FROM outbox
WHERE published_at IS NULL AND dead_at IS NULL
AND available_at <= now()
AND (locked_until IS NULL OR locked_until < now())
AND (expires_at IS NULL OR expires_at > now())
ORDER BY available_at, sequence
LIMIT $1
FOR UPDATE SKIP LOCKED
)
UPDATE outbox o
SET locked_by = $2,
locked_until = now() + ($3::bigint * interval '1 millisecond'),
available_at = now() + ($3::bigint * interval '1 millisecond'),
updated_at = now()
FROM claimed
WHERE o.id = claimed.id
RETURNING o.id, o.message_id, o.message_type, o.message_version,
o.correlation_id, o.conversation_id, o.causation_id, o.request_id,
o.content_type, o.payload, o.tenant_id, o.expires_at, o.ordering_key,
o.metadata, o.headers, o.metadata_version,
o.created_at, o.available_at,
o.attempts, o.locked_by, o.locked_until,
o.published_at, o.dead_at, o.dead_reason, o.last_error"#,
batch_size,
worker,
lease_ms,
)
.fetch_all(executor)
.await
}