xmtp_mls/subscriptions/local_delivery/
types.rs1use 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#[derive(Debug, Clone, Copy)]
13pub(crate) struct LocalDeliveryConfig {
14 pub batch_size: u32,
16 pub max_bytes: u64,
18 pub poll_interval: Duration,
20 pub lease_duration: Duration,
22 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 #[error(transparent)]
73 #[error_code(inherit)]
74 Storage(#[from] StorageError),
75 #[error("The delivered message was not acknowledged")]
77 AcknowledgementRejected,
78 #[error("Delivery stopped after acknowledgement persistence failed")]
80 AcknowledgementFailed,
81 #[error("The queued message belongs to an old delivery selection")]
83 SelectionChanged,
84 #[error("The local delivery configuration is invalid")]
86 InvalidConfiguration,
87 #[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}