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 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
135pub 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
202pub 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
242pub 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
277pub 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}