use reliar_inbox::{InboxClaim, InboxMessage, InboxRecordId, InboxScope};
use sqlx::{Postgres, Transaction};
use super::error::PostgresInboxError;
use crate::connection::schema::{restore_search_path, set_search_path};
use super::PostgresInboxStore;
const ADVISORY_LOCK_CLASS: i32 = i32::from_be_bytes(*b"RELI");
pub(super) async fn claim(
store: &PostgresInboxStore,
tx: &mut Transaction<'_, Postgres>,
scope: &InboxScope,
message: InboxMessage<'_>,
) -> Result<InboxClaim, PostgresInboxError> {
let restore = if store.settings.claim_sets_search_path {
Some(set_search_path(tx, &store.settings.schema).await?)
} else {
None
};
let result = claim_locked(tx, scope.as_str(), message).await;
if result.is_ok()
&& let Some(previous) = restore
{
restore_search_path(tx, &previous).await?;
}
result
}
async fn claim_locked(
tx: &mut Transaction<'_, Postgres>,
scope: &str,
message: InboxMessage<'_>,
) -> Result<InboxClaim, PostgresInboxError> {
let message_id = message.id.as_uuid();
let acquired = sqlx::query_scalar!(
r#"SELECT pg_try_advisory_xact_lock($1, hashtext($2 || '/' || $3::uuid::text)) AS "acquired!""#,
ADVISORY_LOCK_CLASS,
scope,
message_id,
)
.fetch_one(&mut **tx)
.await?;
if !acquired {
return Ok(InboxClaim::InProgress);
}
let id = InboxRecordId::new();
if let Some(attempts) = insert_claim_row(tx, id, scope, message).await? {
return Ok(InboxClaim::Claimed {
attempt: claimed_attempt(attempts),
});
}
let row = sqlx::query!(
r#"SELECT id, attempts, completed_at, dead_at FROM inbox WHERE scope = $1 AND message_id = $2"#,
scope,
message_id,
)
.fetch_optional(&mut **tx)
.await?;
let Some(row) = row else {
return upsert_claim_row(tx, InboxRecordId::new(), scope, message).await;
};
if let Some(completed_at) = row.completed_at {
return Ok(InboxClaim::AlreadyCompleted { completed_at });
}
if let Some(dead_at) = row.dead_at {
return Ok(InboxClaim::Dead {
id: InboxRecordId::from_uuid(row.id),
attempts: claimed_attempts_recorded(row.attempts),
dead_at,
});
}
Ok(InboxClaim::Claimed {
attempt: claimed_attempt(row.attempts),
})
}
fn claimed_attempt(attempts: i32) -> u32 {
u32::try_from(attempts)
.unwrap_or(u32::MAX)
.saturating_add(1)
}
fn claimed_attempts_recorded(attempts: i32) -> u32 {
u32::try_from(attempts).unwrap_or(u32::MAX)
}
async fn insert_claim_row(
tx: &mut Transaction<'_, Postgres>,
id: InboxRecordId,
scope: &str,
message: InboxMessage<'_>,
) -> Result<Option<i32>, PostgresInboxError> {
let id = id.as_uuid();
let message_id = message.id.as_uuid();
let message_type = message.message_type.name();
let message_version = i32::from(message.message_type.version());
let conversation_id = message.conversation_id.as_uuid();
let correlation_id = message
.correlation_id
.map(reliar_core::CorrelationId::as_str);
let causation_id = message.causation_id.map(|c| c.as_uuid());
let inserted = 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"#,
id,
scope,
message_id,
message_type,
message_version,
conversation_id,
correlation_id,
causation_id,
)
.fetch_optional(&mut **tx)
.await?;
Ok(inserted)
}
async fn upsert_claim_row(
tx: &mut Transaction<'_, Postgres>,
id: InboxRecordId,
scope: &str,
message: InboxMessage<'_>,
) -> Result<InboxClaim, PostgresInboxError> {
let id = id.as_uuid();
let message_id = message.id.as_uuid();
let message_type = message.message_type.name();
let message_version = i32::from(message.message_type.version());
let conversation_id = message.conversation_id.as_uuid();
let correlation_id = message
.correlation_id
.map(reliar_core::CorrelationId::as_str);
let causation_id = message.causation_id.map(|c| c.as_uuid());
let row = sqlx::query!(
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"#,
id,
scope,
message_id,
message_type,
message_version,
conversation_id,
correlation_id,
causation_id,
)
.fetch_one(&mut **tx)
.await?;
if let Some(completed_at) = row.completed_at {
return Ok(InboxClaim::AlreadyCompleted { completed_at });
}
if let Some(dead_at) = row.dead_at {
return Ok(InboxClaim::Dead {
id: InboxRecordId::from_uuid(row.id),
attempts: claimed_attempts_recorded(row.attempts),
dead_at,
});
}
Ok(InboxClaim::Claimed {
attempt: claimed_attempt(row.attempts),
})
}