1use std::sync::atomic::{AtomicU64, Ordering};
13use std::sync::Arc;
14
15use connectrpc::ConnectError;
16use exoware_sdk::common::Entry;
17use exoware_sdk::keys::Prefix;
18use exoware_sdk::selector::compile_payload_regex;
19use exoware_sdk::stream_filter::{validate_filter, CompiledFilters, StreamFilter};
20use regex::bytes::Regex;
21use tokio::sync::Notify;
22
23pub const STREAM_ERROR_DOMAIN: &str = "log.stream";
25pub const REASON_BATCH_EVICTED: &str = "BATCH_EVICTED";
28pub const REASON_BATCH_NOT_FOUND: &str = "BATCH_NOT_FOUND";
31pub const METADATA_OLDEST_RETAINED: &str = "oldest_retained";
33
34#[derive(Clone, Debug)]
36pub struct CompiledSelector {
37 prefix: Prefix,
38 regex: Regex,
39}
40
41impl CompiledSelector {
42 pub fn prefix(&self) -> &[u8] {
44 self.prefix.as_bytes()
45 }
46
47 pub fn matches_key(&self, key: &[u8]) -> bool {
49 self.prefix
50 .strip_slice(key)
51 .is_some_and(|payload| self.regex.is_match(payload))
52 }
53}
54
55#[derive(Clone, Debug, thiserror::Error)]
60#[error("{0}")]
61pub struct InvalidFilter(String);
62
63#[derive(Clone, Debug)]
65pub struct CompiledMatchers {
66 keys: Vec<CompiledSelector>,
67 values: Option<CompiledFilters>,
68}
69
70impl CompiledMatchers {
71 pub fn compile(filter: &StreamFilter) -> Result<Self, InvalidFilter> {
73 validate_filter(filter).map_err(|e| InvalidFilter(e.to_string()))?;
74 let keys = filter
75 .selectors
76 .iter()
77 .map(|mk| {
78 let regex = compile_payload_regex(&mk.payload_regex)
79 .map_err(|e| InvalidFilter(e.to_string()))?;
80 let prefix =
81 Prefix::new(mk.prefix.clone()).map_err(|e| InvalidFilter(e.to_string()))?;
82 Ok(CompiledSelector { prefix, regex })
83 })
84 .collect::<Result<Vec<_>, InvalidFilter>>()?;
85 let values = CompiledFilters::compile(&filter.value_filters)
86 .map_err(|e| InvalidFilter(format!("invalid value_filter: {e}")))?;
87 Ok(Self { keys, values })
88 }
89
90 pub fn matches(&self, key: &[u8], value: &[u8]) -> bool {
97 self.matches_key(key) && self.matches_value(value)
98 }
99
100 pub fn matches_key(&self, key: &[u8]) -> bool {
104 self.keys.iter().any(|matcher| matcher.matches_key(key))
105 }
106
107 pub fn selectors(&self) -> &[CompiledSelector] {
109 &self.keys
110 }
111
112 pub fn matches_value(&self, value: &[u8]) -> bool {
116 self.values
117 .as_ref()
118 .is_none_or(|matcher| matcher.matches(value))
119 }
120}
121
122pub(crate) fn compile_matchers(filter: &StreamFilter) -> Result<CompiledMatchers, ConnectError> {
126 CompiledMatchers::compile(filter).map_err(|e| ConnectError::invalid_argument(e.to_string()))
127}
128
129pub(crate) fn apply_filter(matchers: &CompiledMatchers, kvs: Vec<Entry>) -> Vec<Entry> {
131 kvs.into_iter()
132 .filter(|kv| matchers.matches(kv.key.as_ref(), kv.value.as_ref()))
133 .collect()
134}
135
136#[derive(Clone)]
137pub struct StreamNotification {
138 pub current_sequence: u64,
139 pub notify: Arc<Notify>,
140}
141
142pub trait StreamNotifier: Send + Sync + 'static {
145 fn subscribe(&self) -> StreamNotification;
148
149 fn current_sequence(&self) -> u64;
151
152 fn advance(&self, seq: u64);
154}
155
156pub struct StreamHub {
157 published_sequence: AtomicU64,
158 notify: Arc<Notify>,
159}
160
161impl StreamHub {
162 pub fn new(initial_sequence: u64) -> Self {
163 Self {
164 published_sequence: AtomicU64::new(initial_sequence),
165 notify: Arc::new(Notify::new()),
166 }
167 }
168
169 pub fn publish(&self, seq: u64) {
171 self.advance(seq);
172 }
173}
174
175impl StreamNotifier for StreamHub {
176 fn subscribe(&self) -> StreamNotification {
177 StreamNotification {
178 current_sequence: self.published_sequence.load(Ordering::Acquire),
179 notify: self.notify.clone(),
180 }
181 }
182
183 fn current_sequence(&self) -> u64 {
184 self.published_sequence.load(Ordering::Acquire)
185 }
186
187 fn advance(&self, seq: u64) {
188 self.published_sequence.fetch_max(seq, Ordering::SeqCst);
189 self.notify.notify_waiters();
190 }
191}
192
193#[cfg(test)]
194mod tests {
195 use super::*;
196 use crate::{Filter, Selector, StreamFilter, Utf8};
197 use bytes::Bytes;
198 use exoware_sdk::keys::Key;
199
200 fn filter(prefix: u8, regex: &str) -> StreamFilter {
201 StreamFilter {
202 selectors: vec![Selector {
203 prefix: Bytes::from(vec![prefix]),
204 payload_regex: Utf8::from(regex),
205 }],
206 value_filters: vec![],
207 }
208 }
209
210 fn filter_with_values(prefix: u8, regex: &str, value_filters: Vec<Filter>) -> StreamFilter {
211 StreamFilter {
212 selectors: vec![Selector {
213 prefix: Bytes::from(vec![prefix]),
214 payload_regex: Utf8::from(regex),
215 }],
216 value_filters,
217 }
218 }
219
220 fn key(family: u8, payload: &[u8]) -> Key {
221 Prefix::from_byte(family).encode(payload).unwrap()
222 }
223
224 fn kv(family: u8, payload: &[u8], value: &'static [u8]) -> Entry {
225 Entry {
226 key: key(family, payload).to_vec(),
227 value: Bytes::from_static(value),
228 ..Default::default()
229 }
230 }
231
232 #[test]
233 fn publish_sequence_is_monotonic() {
234 let hub = StreamHub::new(7);
235 assert_eq!(hub.current_sequence(), 7);
236 hub.publish(3);
237 assert_eq!(hub.current_sequence(), 7);
238 hub.publish(9);
239 assert_eq!(hub.current_sequence(), 9);
240 }
241
242 #[test]
243 fn subscribe_snapshots_current_sequence() {
244 let hub = StreamHub::new(11);
245 let subscription = hub.subscribe();
246 assert_eq!(subscription.current_sequence, 11);
247 }
248
249 #[test]
250 fn apply_filter_still_selects_matching_entries() {
251 let matchers = compile_matchers(&filter(1, "(?s).*")).unwrap();
252 let kvs = vec![kv(1, b"hit", b"v1"), kv(2, b"miss", b"v2")];
253 let entries = apply_filter(&matchers, kvs);
254 assert_eq!(entries.len(), 1);
255 assert_eq!(entries[0].value.as_ref(), b"v1");
256 }
257
258 #[test]
259 fn subscribe_rejects_invalid_filter() {
260 let bad = StreamFilter {
261 selectors: vec![],
262 value_filters: vec![],
263 };
264 assert!(compile_matchers(&bad).is_err());
265 }
266
267 #[test]
268 fn value_filter_intersects_with_key_filter() {
269 let matchers = compile_matchers(&filter_with_values(
270 1,
271 "(?s).*",
272 vec![Filter::Regex("^keep$".into())],
273 ))
274 .unwrap();
275 let kvs = vec![kv(1, b"a", b"keep"), kv(1, b"b", b"drop")];
276 let entries = apply_filter(&matchers, kvs);
277 assert_eq!(entries.len(), 1);
278 assert_eq!(entries[0].value.as_ref(), b"keep");
279 }
280
281 #[test]
282 fn value_filter_exact_match() {
283 let matchers = compile_matchers(&filter_with_values(
284 1,
285 "(?s).*",
286 vec![Filter::Exact(Bytes::from_static(b"target"))],
287 ))
288 .unwrap();
289 let kvs = vec![kv(1, b"a", b"target"), kv(1, b"b", b"other")];
290 let entries = apply_filter(&matchers, kvs);
291 assert_eq!(entries.len(), 1);
292 assert_eq!(entries[0].value.as_ref(), b"target");
293 }
294
295 #[test]
296 fn value_filter_empty_accepts_all_matching_keys() {
297 let matchers = compile_matchers(&filter(1, "(?s).*")).unwrap();
298 let kvs = vec![kv(1, b"a", b"one"), kv(1, b"b", b"two")];
299 let entries = apply_filter(&matchers, kvs);
300 assert_eq!(entries.len(), 2);
301 }
302
303 fn apply_filter_reference(matchers: &CompiledMatchers, kvs: &[Entry]) -> Vec<Entry> {
306 let mut out = Vec::with_capacity(kvs.len());
307 'outer: for kv in kvs {
308 let v = kv.value.as_ref();
309 let value_ok = matchers.values.as_ref().is_none_or(|m| m.matches(v));
310 if !value_ok {
311 continue;
312 }
313 for matcher in &matchers.keys {
314 let Some(payload) = matcher.prefix.strip_slice(&kv.key) else {
315 continue;
316 };
317 if matcher.regex.is_match(payload) {
318 out.push(kv.clone());
319 continue 'outer;
320 }
321 }
322 }
323 out
324 }
325
326 #[test]
327 fn apply_filter_matches_reference_implementation() {
328 let overlapping = StreamFilter {
329 selectors: vec![
330 Selector {
331 prefix: Bytes::from_static(&[1]),
332 payload_regex: Utf8::from("^a"),
333 },
334 Selector {
335 prefix: Bytes::from_static(&[1]),
336 payload_regex: Utf8::from("b$"),
337 },
338 Selector {
339 prefix: Bytes::from_static(&[2]),
340 payload_regex: Utf8::from("(?s).*"),
341 },
342 ],
343 value_filters: vec![],
344 };
345 let with_values = StreamFilter {
346 value_filters: vec![
347 Filter::Prefix(Bytes::from_static(b"keep")),
348 Filter::Exact(Bytes::from_static(b"exact")),
349 Filter::Regex("^regex$".into()),
350 ],
351 ..overlapping.clone()
352 };
353 let entries = vec![
354 kv(1, b"ab", b"keep-anything"),
356 kv(1, b"aX", b"exact"),
357 kv(1, b"Xb", b"regex"),
358 kv(1, b"XX", b"keep"),
359 kv(1, b"", b"keep"),
360 kv(2, b"", b"exact"),
361 kv(2, b"anything", b"drop"),
362 kv(3, b"ab", b"keep"),
363 kv(1, b"ab", b"keep-anything"),
365 kv(1, b"ab", b"nope"),
366 ];
367
368 for filter in [overlapping, with_values] {
369 let matchers = compile_matchers(&filter).unwrap();
370 let expected = apply_filter_reference(&matchers, &entries);
371 let actual = apply_filter(&matchers, entries.clone());
372 assert_eq!(actual, expected);
373 }
374 }
375
376 #[test]
377 fn public_compile_builds_usable_matchers_and_rejects_invalid_filters() {
378 let matchers = CompiledMatchers::compile(&filter_with_values(
379 1,
380 "^hit$",
381 vec![Filter::Exact(Bytes::from_static(b"v"))],
382 ))
383 .expect("compile");
384 assert!(matchers.matches(key(1, b"hit").as_ref(), b"v"));
385 assert!(!matchers.matches(key(1, b"miss").as_ref(), b"v"));
386 assert!(!matchers.matches(key(1, b"hit").as_ref(), b"other"));
387
388 let invalid = StreamFilter {
389 selectors: vec![],
390 value_filters: vec![],
391 };
392 assert!(CompiledMatchers::compile(&invalid).is_err());
393 }
394
395 #[test]
396 fn matcher_methods_preserve_key_and_value_semantics() {
397 let matchers = compile_matchers(&filter_with_values(
398 1,
399 "^hit$",
400 vec![
401 Filter::Prefix(Bytes::from_static(b"keep")),
402 Filter::Exact(Bytes::from_static(b"exact")),
403 Filter::Regex("^regex$".into()),
404 ],
405 ))
406 .unwrap();
407
408 assert!(matchers.matches_key(key(1, b"hit").as_ref()));
409 assert!(!matchers.matches_key(key(1, b"miss").as_ref()));
410 assert!(!matchers.matches_key(key(2, b"hit").as_ref()));
411 assert!(matchers.matches_value(b"keep-suffix"));
412 assert!(matchers.matches_value(b"exact"));
413 assert!(matchers.matches_value(b"regex"));
414 assert!(!matchers.matches_value(b"drop"));
415
416 let without_value_filters = compile_matchers(&filter(1, "(?s).*")).unwrap();
417 assert!(without_value_filters.matches_value(b"anything"));
418 }
419
420 #[test]
421 fn compiled_selectors_expose_prefixes_and_preserve_match_semantics() {
422 let filter = StreamFilter {
423 selectors: vec![
424 Selector {
425 prefix: Bytes::from_static(&[1]),
426 payload_regex: Utf8::from("^hit$"),
427 },
428 Selector {
429 prefix: Bytes::from_static(&[1, 2]),
430 payload_regex: Utf8::from("suffix"),
431 },
432 ],
433 value_filters: vec![],
434 };
435 let matchers = CompiledMatchers::compile(&filter).expect("compile");
436 let selectors = matchers.selectors();
437
438 assert_eq!(selectors.len(), 2);
439 assert_eq!(selectors[0].prefix(), &[1]);
440 assert!(selectors[0].matches_key(&[1, b'h', b'i', b't']));
441 assert!(!selectors[0].matches_key(&[1, b'm', b'i', b's', b's']));
442 assert_eq!(selectors[1].prefix(), &[1, 2]);
443 assert!(selectors[1].matches_key(&[1, 2, b'x', b's', b'u', b'f', b'f', b'i', b'x']));
444 assert!(!selectors[1].matches_key(&[1, b'x', b's', b'u', b'f', b'f', b'i', b'x']));
445 }
446}