Skip to main content

xmtp_mls/subscriptions/local_delivery/
types.rs

1use thiserror::Error;
2use xmtp_common::{ErrorCode, RetryableError, time::Duration};
3use xmtp_db::StorageError;
4
5pub use xmtp_db::delivery::DeliveryFilter as LocalDeliveryFilter;
6
7use crate::subscriptions::policy::StreamPolicy;
8
9const LEASE_RENEWAL_DIVISOR: u32 = 3;
10
11/// Bounds local reads and maintains exclusive default-consumer ownership.
12#[derive(Debug, Clone, Copy)]
13pub(crate) struct LocalDeliveryConfig {
14    /// Maximum candidates fetched in one local read.
15    pub batch_size: u32,
16    /// Maximum candidate bytes, checked before message blobs are loaded.
17    pub max_bytes: u64,
18    /// Maximum wait before checking cross-process writes and missed local wakes.
19    pub poll_interval: Duration,
20    /// Time before another process can take over an unrenewed default consumer.
21    pub lease_duration: Duration,
22    /// Renewal interval; must be shorter than the lease duration.
23    pub renew_interval: Duration,
24}
25
26impl Default for LocalDeliveryConfig {
27    fn default() -> Self {
28        Self::from(&StreamPolicy::default())
29    }
30}
31
32impl From<&StreamPolicy> for LocalDeliveryConfig {
33    fn from(settings: &StreamPolicy) -> Self {
34        Self {
35            batch_size: settings.max_local_read_rows,
36            max_bytes: settings.max_local_read_bytes,
37            poll_interval: settings.active_database_poll_interval,
38            lease_duration: settings.default_consumer_lease_duration,
39            renew_interval: settings.default_consumer_lease_duration / LEASE_RENEWAL_DIVISOR,
40        }
41    }
42}
43
44impl LocalDeliveryConfig {
45    pub(super) fn validate(self) -> Result<(), LocalDeliveryError> {
46        if self.batch_size == 0
47            || self.max_bytes == 0
48            || self.poll_interval.is_zero()
49            || self.renew_interval.is_zero()
50            || self.renew_interval >= self.lease_duration
51        {
52            return Err(LocalDeliveryError::InvalidConfiguration);
53        }
54        self.lease_until(0)?;
55        Ok(())
56    }
57
58    pub(super) fn lease_until(self, now: i64) -> Result<i64, LocalDeliveryError> {
59        now.checked_add(self.lease_duration_ns()?)
60            .ok_or(LocalDeliveryError::InvalidConfiguration)
61    }
62
63    pub(super) fn lease_duration_ns(self) -> Result<i64, LocalDeliveryError> {
64        i64::try_from(self.lease_duration.as_nanos())
65            .map_err(|_| LocalDeliveryError::InvalidConfiguration)
66    }
67}
68
69#[derive(Debug, Error, ErrorCode)]
70pub enum LocalDeliveryError {
71    /// Database receipt, lease, cursor, or acknowledgement failure. May be retryable.
72    #[error(transparent)]
73    #[error_code(inherit)]
74    Storage(#[from] StorageError),
75    /// The callback failed or its token was dropped before acknowledgement. Not retryable.
76    #[error("The delivered message was not acknowledged")]
77    AcknowledgementRejected,
78    /// A previous acknowledgement write failed. Reopen to retry delivery. Not retryable.
79    #[error("Delivery stopped after acknowledgement persistence failed")]
80    AcknowledgementFailed,
81    /// Scope, filters, or retained content changed before dispatch. Reselect without acknowledgement.
82    #[error("The queued message belongs to an old delivery selection")]
83    SelectionChanged,
84    /// The batch or timing settings cannot maintain a valid consumer lease. Not retryable.
85    #[error("The local delivery configuration is invalid")]
86    InvalidConfiguration,
87    /// This reader has been closed. Not retryable.
88    #[error("The local message reader is closed")]
89    Closed,
90}
91
92impl RetryableError for LocalDeliveryError {
93    fn is_retryable(&self) -> bool {
94        match self {
95            Self::Storage(error) => error.is_retryable(),
96            _ => false,
97        }
98    }
99}