Skip to main content

cortexkit_bus_naming/
limits.rs

1use std::error::Error;
2use std::fmt;
3use std::time::Duration;
4
5use crate::names::AccountNames;
6
7pub const HOUR: Duration = Duration::from_secs(60 * 60);
8pub const MIB: u64 = 1024 * 1024;
9pub const GIB: u64 = 1024 * MIB;
10
11#[derive(Debug, Clone, Copy, PartialEq, Eq)]
12pub enum StreamKind {
13    Room,
14    Wake,
15    Peer,
16    Effect,
17    EffectDead,
18    Event,
19}
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22pub enum DiscardPolicy {
23    Old,
24    New,
25}
26
27#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct StreamSpec {
29    pub kind: StreamKind,
30    pub name: String,
31    pub subjects: Vec<String>,
32    pub max_age: Duration,
33    pub max_bytes: u64,
34    pub discard: DiscardPolicy,
35    pub work_queue: bool,
36    /// The per-subject message cap (`max_msgs_per_subject`), or `None` to set
37    /// none. Only the event stream has one, so one noisy module or event cannot
38    /// evict every other subject's history; the five older streams leave it
39    /// unset and their emitted configuration is unchanged.
40    pub max_msgs_per_subject: Option<i64>,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
44pub struct ConsumerSpec {
45    pub durable: String,
46    pub stream: String,
47    pub filter_subjects: Vec<String>,
48}
49
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub enum LimitError {
52    InvalidSubject {
53        subject: String,
54        reason: &'static str,
55    },
56    DuplicateStream {
57        stream: String,
58    },
59    OverlappingStreamFilters {
60        left_stream: String,
61        left_filter: String,
62        right_stream: String,
63        right_filter: String,
64        involves_work_queue: bool,
65    },
66    MissingConsumerStream {
67        durable: String,
68        stream: String,
69    },
70    ConsumerFilterNotSubset {
71        durable: String,
72        filter: String,
73        stream: String,
74        bindings: Vec<String>,
75    },
76    StoreRetentionTooShort {
77        stream: String,
78        max_age: Duration,
79        store_retention: Duration,
80    },
81}
82
83impl fmt::Display for LimitError {
84    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
85        match self {
86            Self::InvalidSubject { subject, reason } => {
87                write!(f, "invalid subject filter {subject:?}: {reason}")
88            }
89            Self::DuplicateStream { stream } => write!(f, "duplicate stream {stream:?}"),
90            Self::OverlappingStreamFilters {
91                left_stream,
92                left_filter,
93                right_stream,
94                right_filter,
95                involves_work_queue,
96            } => {
97                let queue = if *involves_work_queue {
98                    "work-queue "
99                } else {
100                    ""
101                };
102                write!(
103                    f,
104                    "{queue}stream filter overlap: {left_stream} {left_filter:?} overlaps {right_stream} {right_filter:?}"
105                )
106            }
107            Self::MissingConsumerStream { durable, stream } => {
108                write!(f, "consumer {durable:?} names unknown stream {stream:?}")
109            }
110            Self::ConsumerFilterNotSubset {
111                durable,
112                filter,
113                stream,
114                bindings,
115            } => write!(
116                f,
117                "consumer {durable:?} filter {filter:?} is not a subset of stream {stream} bindings {bindings:?}"
118            ),
119            Self::StoreRetentionTooShort {
120                stream,
121                max_age,
122                store_retention,
123            } => write!(
124                f,
125                "stream {stream} max-age {}s exceeds declared store retention {}s",
126                max_age.as_secs(),
127                store_retention.as_secs()
128            ),
129        }
130    }
131}
132
133impl Error for LimitError {}
134
135/// Returns the six normative stream configurations for an account.
136pub fn shipped_streams(account: &AccountNames) -> Vec<StreamSpec> {
137    let names = account.streams();
138    vec![
139        StreamSpec {
140            kind: StreamKind::Room,
141            name: names.room.clone(),
142            subjects: vec![account.room_binding()],
143            max_age: HOUR * 24,
144            max_bytes: GIB,
145            discard: DiscardPolicy::Old,
146            work_queue: false,
147            max_msgs_per_subject: None,
148        },
149        StreamSpec {
150            kind: StreamKind::Wake,
151            name: names.wake.clone(),
152            subjects: vec![account.wake_binding()],
153            max_age: HOUR * 24,
154            max_bytes: GIB,
155            discard: DiscardPolicy::Old,
156            work_queue: false,
157            max_msgs_per_subject: None,
158        },
159        StreamSpec {
160            kind: StreamKind::Peer,
161            name: names.peer.clone(),
162            subjects: vec![account.peer_binding()],
163            max_age: HOUR * 24,
164            max_bytes: GIB,
165            discard: DiscardPolicy::Old,
166            work_queue: false,
167            max_msgs_per_subject: None,
168        },
169        StreamSpec {
170            kind: StreamKind::Effect,
171            name: names.effect.clone(),
172            subjects: vec![account.effect_binding()],
173            max_age: HOUR * 24,
174            max_bytes: 256 * MIB,
175            discard: DiscardPolicy::New,
176            work_queue: true,
177            max_msgs_per_subject: None,
178        },
179        StreamSpec {
180            kind: StreamKind::EffectDead,
181            name: names.effect_dead.clone(),
182            subjects: vec![account.effect_dead()],
183            max_age: HOUR * 24 * 7,
184            max_bytes: 64 * MIB,
185            discard: DiscardPolicy::Old,
186            work_queue: false,
187            max_msgs_per_subject: None,
188        },
189        StreamSpec {
190            kind: StreamKind::Event,
191            name: names.event.clone(),
192            subjects: vec![account.event_binding()],
193            max_age: HOUR * 24 * 7,
194            max_bytes: GIB,
195            discard: DiscardPolicy::Old,
196            work_queue: false,
197            max_msgs_per_subject: Some(10_000),
198        },
199    ]
200}
201
202/// Rejects duplicate streams, malformed filters, and every pair of overlapping filters.
203pub fn validate_streams(streams: &[StreamSpec]) -> Result<(), LimitError> {
204    for (index, stream) in streams.iter().enumerate() {
205        if streams[..index]
206            .iter()
207            .any(|other| other.name == stream.name)
208        {
209            return Err(LimitError::DuplicateStream {
210                stream: stream.name.clone(),
211            });
212        }
213        for subject in &stream.subjects {
214            SubjectPattern::parse(subject)?;
215        }
216    }
217
218    for left_index in 0..streams.len() {
219        for right_index in (left_index + 1)..streams.len() {
220            let left = &streams[left_index];
221            let right = &streams[right_index];
222            for left_filter in &left.subjects {
223                let left_pattern = SubjectPattern::parse(left_filter)?;
224                for right_filter in &right.subjects {
225                    let right_pattern = SubjectPattern::parse(right_filter)?;
226                    if left_pattern.overlaps(&right_pattern) {
227                        return Err(LimitError::OverlappingStreamFilters {
228                            left_stream: left.name.clone(),
229                            left_filter: left_filter.clone(),
230                            right_stream: right.name.clone(),
231                            right_filter: right_filter.clone(),
232                            involves_work_queue: left.work_queue || right.work_queue,
233                        });
234                    }
235                }
236            }
237        }
238    }
239    Ok(())
240}
241
242/// Rejects a consumer filter unless every subject it can match is in its stream.
243pub fn validate_consumer(
244    consumer: &ConsumerSpec,
245    streams: &[StreamSpec],
246) -> Result<(), LimitError> {
247    let stream = streams
248        .iter()
249        .find(|stream| stream.name == consumer.stream)
250        .ok_or_else(|| LimitError::MissingConsumerStream {
251            durable: consumer.durable.clone(),
252            stream: consumer.stream.clone(),
253        })?;
254
255    let bindings = stream
256        .subjects
257        .iter()
258        .map(|subject| SubjectPattern::parse(subject))
259        .collect::<Result<Vec<_>, _>>()?;
260    for filter in &consumer.filter_subjects {
261        let filter_pattern = SubjectPattern::parse(filter)?;
262        if !bindings
263            .iter()
264            .any(|binding| filter_pattern.is_subset_of(binding))
265        {
266            return Err(LimitError::ConsumerFilterNotSubset {
267                durable: consumer.durable.clone(),
268                filter: filter.clone(),
269                stream: stream.name.clone(),
270                bindings: stream.subjects.clone(),
271            });
272        }
273    }
274    Ok(())
275}
276
277/// Enforces that the owning store keeps a body at least as long as its stream.
278///
279/// A module declaring a `resolved` event calls this against the event stream:
280/// the event carries only an id, so the module must keep the record for at
281/// least as long as the stream keeps the event.
282pub fn validate_store_retention(
283    stream: &StreamSpec,
284    store_retention: Duration,
285) -> Result<(), LimitError> {
286    if stream.max_age > store_retention {
287        Err(LimitError::StoreRetentionTooShort {
288            stream: stream.name.clone(),
289            max_age: stream.max_age,
290            store_retention,
291        })
292    } else {
293        Ok(())
294    }
295}
296
297#[derive(Debug)]
298struct SubjectPattern<'a> {
299    tokens: Vec<&'a str>,
300}
301
302impl<'a> SubjectPattern<'a> {
303    fn parse(raw: &'a str) -> Result<Self, LimitError> {
304        if raw.is_empty() {
305            return Err(invalid_subject(raw, "subject must not be empty"));
306        }
307        let tokens = raw.split('.').collect::<Vec<_>>();
308        if tokens.iter().any(|token| token.is_empty()) {
309            return Err(invalid_subject(raw, "subject contains an empty token"));
310        }
311        for (index, token) in tokens.iter().enumerate() {
312            if token.bytes().any(|byte| byte.is_ascii_whitespace()) {
313                return Err(invalid_subject(
314                    raw,
315                    "subject tokens cannot contain whitespace",
316                ));
317            }
318            if (token.contains('*') && *token != "*") || (token.contains('>') && *token != ">") {
319                return Err(invalid_subject(
320                    raw,
321                    "wildcards must occupy a whole subject token",
322                ));
323            }
324            if *token == ">" && index + 1 != tokens.len() {
325                return Err(invalid_subject(
326                    raw,
327                    "the > wildcard must be the final token",
328                ));
329            }
330        }
331        Ok(Self { tokens })
332    }
333
334    fn overlaps(&self, other: &Self) -> bool {
335        patterns_overlap(&self.tokens, &other.tokens)
336    }
337
338    fn is_subset_of(&self, other: &Self) -> bool {
339        pattern_subset(&self.tokens, &other.tokens)
340    }
341}
342
343fn invalid_subject(subject: &str, reason: &'static str) -> LimitError {
344    LimitError::InvalidSubject {
345        subject: subject.to_owned(),
346        reason,
347    }
348}
349
350fn patterns_overlap(left: &[&str], right: &[&str]) -> bool {
351    match (left.first(), right.first()) {
352        (None, None) => true,
353        (None, Some(_)) | (Some(_), None) => false,
354        (Some(&">"), Some(_)) | (Some(_), Some(&">")) => true,
355        (Some(left_head), Some(right_head))
356            if left_head == right_head || *left_head == "*" || *right_head == "*" =>
357        {
358            patterns_overlap(&left[1..], &right[1..])
359        }
360        _ => false,
361    }
362}
363
364fn pattern_subset(candidate: &[&str], binding: &[&str]) -> bool {
365    match (candidate.first(), binding.first()) {
366        (None, None) => true,
367        (None, Some(_)) | (Some(_), None) => false,
368        (Some(_), Some(&">")) => true,
369        (Some(&">"), Some(_)) => false,
370        (Some(candidate_head), Some(binding_head)) => {
371            let head_is_subset = binding_head == candidate_head || *binding_head == "*";
372            head_is_subset && pattern_subset(&candidate[1..], &binding[1..])
373        }
374    }
375}