pub(in crate::outbox) struct CompleteRowsParams<'a> {
pub(in crate::outbox) ids: &'a [uuid::Uuid],
pub(in crate::outbox) tokens: &'a [Option<uuid::Uuid>],
}
pub(in crate::outbox) async fn complete_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: CompleteRowsParams<'_>,
) -> Result<Vec<uuid::Uuid>, sqlx::Error> {
let rows = sqlx::query!(
r#"UPDATE outbox o
SET published_at = now(),
attempts = o.attempts + 1,
locked_by = NULL,
claim_token = NULL,
updated_at = now()
FROM UNNEST($1::uuid[], $2::uuid[]) AS f(id, token)
WHERE o.id = f.id AND o.claim_token = f.token
RETURNING o.id"#,
params.ids,
params.tokens as &[Option<uuid::Uuid>],
)
.fetch_all(executor)
.await?;
Ok(rows.into_iter().map(|r| r.id).collect())
}
pub(in crate::outbox) struct ReleaseRowsParams<'a> {
pub(in crate::outbox) ids: &'a [uuid::Uuid],
pub(in crate::outbox) tokens: &'a [Option<uuid::Uuid>],
}
pub(in crate::outbox) async fn release_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: ReleaseRowsParams<'_>,
) -> Result<Vec<uuid::Uuid>, sqlx::Error> {
let rows = sqlx::query!(
r#"UPDATE outbox o
SET locked_by = NULL,
claim_token = NULL,
available_at = now(),
updated_at = now()
FROM UNNEST($1::uuid[], $2::uuid[]) AS f(id, token)
WHERE o.id = f.id AND o.claim_token = f.token
RETURNING o.id"#,
params.ids,
params.tokens as &[Option<uuid::Uuid>],
)
.fetch_all(executor)
.await?;
Ok(rows.into_iter().map(|r| r.id).collect())
}
pub(in crate::outbox) struct ExtendLeaseRowsParams<'a> {
pub(in crate::outbox) ids: &'a [uuid::Uuid],
pub(in crate::outbox) tokens: &'a [Option<uuid::Uuid>],
pub(in crate::outbox) lease_ms: i64,
}
pub(in crate::outbox) async fn extend_lease_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: ExtendLeaseRowsParams<'_>,
) -> Result<Vec<uuid::Uuid>, sqlx::Error> {
let rows = sqlx::query!(
r#"UPDATE outbox o
SET available_at = now() + ($3::bigint * interval '1 millisecond'),
updated_at = now()
FROM UNNEST($1::uuid[], $2::uuid[]) AS f(id, token)
WHERE o.id = f.id AND o.claim_token = f.token
RETURNING o.id"#,
params.ids,
params.tokens as &[Option<uuid::Uuid>],
params.lease_ms,
)
.fetch_all(executor)
.await?;
Ok(rows.into_iter().map(|r| r.id).collect())
}
pub(in crate::outbox) struct FailRetryRowsParams<'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) delays_ms: &'a [i64],
}
pub(in crate::outbox) async fn fail_retry_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: FailRetryRowsParams<'_>,
) -> Result<Vec<uuid::Uuid>, sqlx::Error> {
let rows = sqlx::query!(
r#"UPDATE outbox o
SET attempts = o.attempts + 1,
last_error = f.err,
locked_by = NULL,
claim_token = NULL,
available_at = now() + (f.delay_ms * interval '1 millisecond'),
updated_at = now()
FROM UNNEST($1::uuid[], $2::uuid[], $3::text[], $4::bigint[]) AS f(id, token, err, delay_ms)
WHERE o.id = f.id AND o.claim_token = f.token
RETURNING o.id"#,
params.ids,
params.tokens as &[Option<uuid::Uuid>],
params.errors,
params.delays_ms,
)
.fetch_all(executor)
.await?;
Ok(rows.into_iter().map(|r| r.id).collect())
}
pub(in crate::outbox) struct FailDeadRowsParams<'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) reasons: &'a [&'a str],
}
pub(in crate::outbox) async fn fail_dead_rows<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: FailDeadRowsParams<'_>,
) -> Result<Vec<uuid::Uuid>, sqlx::Error> {
let rows = sqlx::query!(
r#"UPDATE outbox o
SET attempts = o.attempts + 1,
last_error = f.err,
dead_at = now(),
dead_reason = f.reason,
locked_by = NULL,
claim_token = NULL,
updated_at = now()
FROM UNNEST($1::uuid[], $2::uuid[], $3::text[], $4::text[]) AS f(id, token, err, reason)
WHERE o.id = f.id AND o.claim_token = f.token
RETURNING o.id"#,
params.ids,
params.tokens as &[Option<uuid::Uuid>],
params.errors,
params.reasons as &[&str],
)
.fetch_all(executor)
.await?;
Ok(rows.into_iter().map(|r| r.id).collect())
}