use crate::records::RawRow;
pub(in crate::outbox) struct ClaimRowsParams<'a> {
pub(in crate::outbox) batch_size: i64,
pub(in crate::outbox) worker: &'a str,
pub(in crate::outbox) lease_ms: i64,
}
pub(in crate::outbox) async fn claim_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: ClaimRowsParams<'_>,
) -> 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 (expires_at IS NULL OR expires_at > now())
ORDER BY available_at, id
LIMIT $1
FOR UPDATE SKIP LOCKED
)
UPDATE outbox o
SET locked_by = $2,
available_at = now() + ($3::bigint * interval '1 millisecond'),
claim_token = uuidv7(),
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.claim_token,
o.published_at, o.dead_at, o.dead_reason, o.last_error"#,
params.batch_size,
params.worker,
params.lease_ms,
)
.fetch_all(executor)
.await
}
pub(in crate::outbox) struct PoisonSweepRowsParams<'a> {
pub(in crate::outbox) ids: &'a [uuid::Uuid],
pub(in crate::outbox) tokens: &'a [Option<uuid::Uuid>],
pub(in crate::outbox) errors: &'a [String],
pub(in crate::outbox) dead_reason: &'a str,
}
pub(in crate::outbox) async fn poison_sweep_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: PoisonSweepRowsParams<'_>,
) -> Result<(), sqlx::Error> {
sqlx::query!(
r#"UPDATE outbox o
SET dead_at = now(),
dead_reason = $4,
last_error = f.err,
locked_by = NULL,
claim_token = NULL,
updated_at = now()
FROM UNNEST($1::uuid[], $2::uuid[], $3::text[]) AS f(id, token, err)
WHERE o.id = f.id AND o.claim_token = f.token"#,
params.ids,
params.tokens as &[Option<uuid::Uuid>],
params.errors,
params.dead_reason,
)
.execute(executor)
.await?;
Ok(())
}