use core::fmt;
use reliar_core::{ConversationId, CorrelationId, MessageId, MessageType};
use time::OffsetDateTime;
use crate::message::InboxMessage;
use crate::record_id::InboxRecordId;
use crate::scope::InboxScope;
pub(crate) const MAX_ERROR_LEN: usize = 2048;
pub(crate) const TRUNCATION_MARKER: &str = "…[truncated]";
pub(crate) fn truncate_error(error: impl Into<String>) -> String {
let error = error.into();
if error.len() <= MAX_ERROR_LEN {
return error;
}
let budget = MAX_ERROR_LEN.saturating_sub(TRUNCATION_MARKER.len());
let mut end = budget.min(error.len());
while end > 0 && !error.is_char_boundary(end) {
end -= 1;
}
let mut truncated = String::with_capacity(end + TRUNCATION_MARKER.len());
truncated.push_str(&error[..end]);
truncated.push_str(TRUNCATION_MARKER);
truncated
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum InboxState {
Claimed,
Retrying,
Completed,
Dead,
}
impl fmt::Display for InboxState {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(match self {
Self::Claimed => "claimed",
Self::Retrying => "retrying",
Self::Completed => "completed",
Self::Dead => "dead",
})
}
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct InboxRecord {
pub id: InboxRecordId,
pub scope: InboxScope,
pub message_id: MessageId,
pub message_type: MessageType,
pub conversation_id: ConversationId,
pub correlation_id: Option<CorrelationId>,
pub causation_id: Option<MessageId>,
pub received_at: OffsetDateTime,
pub updated_at: OffsetDateTime,
pub completed_at: Option<OffsetDateTime>,
pub dead_at: Option<OffsetDateTime>,
pub attempts: u32,
pub last_error: Option<String>,
}
impl InboxRecord {
pub fn builder(
id: InboxRecordId,
scope: InboxScope,
message: InboxMessage<'_>,
received_at: OffsetDateTime,
) -> InboxRecordBuilder {
InboxRecordBuilder::new(id, scope, message, received_at)
}
#[must_use]
pub const fn state(&self) -> InboxState {
if self.completed_at.is_some() {
InboxState::Completed
} else if self.dead_at.is_some() {
InboxState::Dead
} else if self.attempts > 0 {
InboxState::Retrying
} else {
InboxState::Claimed
}
}
}
#[must_use]
#[derive(Debug)]
pub struct InboxRecordBuilder {
id: InboxRecordId,
scope: InboxScope,
message_id: MessageId,
message_type: MessageType,
conversation_id: ConversationId,
correlation_id: Option<CorrelationId>,
causation_id: Option<MessageId>,
received_at: OffsetDateTime,
updated_at: OffsetDateTime,
completed_at: Option<OffsetDateTime>,
dead_at: Option<OffsetDateTime>,
attempts: u32,
last_error: Option<String>,
}
impl InboxRecordBuilder {
fn new(
id: InboxRecordId,
scope: InboxScope,
message: InboxMessage<'_>,
received_at: OffsetDateTime,
) -> Self {
Self {
id,
scope,
message_id: message.id,
message_type: message.message_type.clone(),
conversation_id: message.conversation_id,
correlation_id: message.correlation_id.cloned(),
causation_id: message.causation_id,
received_at,
updated_at: received_at,
completed_at: None,
dead_at: None,
attempts: 0,
last_error: None,
}
}
pub const fn updated_at(mut self, updated_at: OffsetDateTime) -> Self {
self.updated_at = updated_at;
self
}
pub const fn completed_at(mut self, completed_at: Option<OffsetDateTime>) -> Self {
self.completed_at = completed_at;
self
}
pub const fn dead_at(mut self, dead_at: Option<OffsetDateTime>) -> Self {
self.dead_at = dead_at;
self
}
pub const fn attempts(mut self, attempts: u32) -> Self {
self.attempts = attempts;
self
}
pub fn last_error(mut self, error: Option<String>) -> Self {
self.last_error = error.map(truncate_error);
self
}
#[must_use]
pub fn build(self) -> InboxRecord {
InboxRecord {
id: self.id,
scope: self.scope,
message_id: self.message_id,
message_type: self.message_type,
conversation_id: self.conversation_id,
correlation_id: self.correlation_id,
causation_id: self.causation_id,
received_at: self.received_at,
updated_at: self.updated_at,
completed_at: self.completed_at,
dead_at: self.dead_at,
attempts: self.attempts,
last_error: self.last_error,
}
}
}