pub(in crate::inbox) struct TryAdvisoryLockParams<'a> {
pub(in crate::inbox) class: i32,
pub(in crate::inbox) scope: &'a str,
pub(in crate::inbox) message_id: uuid::Uuid,
}
pub(in crate::inbox) async fn try_advisory_lock<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: TryAdvisoryLockParams<'_>,
) -> Result<bool, sqlx::Error> {
sqlx::query_scalar!(
r#"SELECT pg_try_advisory_xact_lock($1, hashtext($2 || '/' || $3::uuid::text)) AS "acquired!""#,
params.class,
params.scope,
params.message_id,
)
.fetch_one(executor)
.await
}
pub(in crate::inbox) struct ClaimRowParams<'a> {
pub(in crate::inbox) id: uuid::Uuid,
pub(in crate::inbox) scope: &'a str,
pub(in crate::inbox) message_id: uuid::Uuid,
pub(in crate::inbox) message_type: &'a str,
pub(in crate::inbox) message_version: i32,
pub(in crate::inbox) conversation_id: uuid::Uuid,
pub(in crate::inbox) correlation_id: Option<&'a str>,
pub(in crate::inbox) causation_id: Option<uuid::Uuid>,
}
pub(in crate::inbox) struct ClaimStateRow {
pub(in crate::inbox) id: uuid::Uuid,
pub(in crate::inbox) attempts: i32,
pub(in crate::inbox) completed_at: Option<time::OffsetDateTime>,
pub(in crate::inbox) dead_at: Option<time::OffsetDateTime>,
}
pub(in crate::inbox) async fn insert_claim_row<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: ClaimRowParams<'_>,
) -> Result<Option<i32>, sqlx::Error> {
sqlx::query_scalar!(
r#"INSERT INTO inbox (id, scope, message_id, message_type, message_version,
conversation_id, correlation_id, causation_id)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (scope, message_id) DO NOTHING
RETURNING attempts"#,
params.id,
params.scope,
params.message_id,
params.message_type,
params.message_version,
params.conversation_id,
params.correlation_id,
params.causation_id,
)
.fetch_optional(executor)
.await
}
pub(in crate::inbox) struct SelectClaimStateParams<'a> {
pub(in crate::inbox) scope: &'a str,
pub(in crate::inbox) message_id: uuid::Uuid,
}
pub(in crate::inbox) async fn select_claim_state<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: SelectClaimStateParams<'_>,
) -> Result<Option<ClaimStateRow>, sqlx::Error> {
sqlx::query_as!(
ClaimStateRow,
r#"SELECT id, attempts, completed_at, dead_at FROM inbox WHERE scope = $1 AND message_id = $2"#,
params.scope,
params.message_id,
)
.fetch_optional(executor)
.await
}
pub(in crate::inbox) async fn upsert_claim_row<'e>(
executor: impl sqlx::PgExecutor<'e>,
params: ClaimRowParams<'_>,
) -> Result<ClaimStateRow, sqlx::Error> {
sqlx::query_as!(
ClaimStateRow,
r#"INSERT INTO inbox (id, scope, message_id, message_type, message_version,
conversation_id, correlation_id, causation_id)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (scope, message_id) DO UPDATE
SET updated_at = inbox.updated_at
RETURNING id, attempts, completed_at, dead_at"#,
params.id,
params.scope,
params.message_id,
params.message_type,
params.message_version,
params.conversation_id,
params.correlation_id,
params.causation_id,
)
.fetch_one(executor)
.await
}