Skip to main content

exoware_server/
stream.rs

1//! Live stream coordination for `log.stream.v1`.
2//!
3//! A [`StreamNotifier`] tracks the highest published batch sequence and wakes
4//! subscribers. Each subscriber pulls batches from the log at its own pace.
5//! A small ordered lookahead overlaps log reads without allowing an unbounded
6//! per-subscriber backlog.
7//!
8//! `StreamNotifier` is an in-process coordination primitive. Split deployments
9//! need a separate remote notification path that advances a local notifier after
10//! the query worker can serve the announced batches.
11
12use 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
23/// `ErrorInfo.domain` used for all stream-service errors.
24pub const STREAM_ERROR_DOMAIN: &str = "log.stream";
25/// `ErrorInfo.reason` when a `since_sequence_number` or `Get(seq)` references a
26/// batch that has been pruned from the log.
27pub const REASON_BATCH_EVICTED: &str = "BATCH_EVICTED";
28/// `ErrorInfo.reason` when a `Get(seq)` references a sequence number greater
29/// than any that has ever been issued.
30pub const REASON_BATCH_NOT_FOUND: &str = "BATCH_NOT_FOUND";
31/// Metadata key on `BATCH_EVICTED` errors carrying the lowest retained seq.
32pub const METADATA_OLDEST_RETAINED: &str = "oldest_retained";
33
34/// One validated and compiled stream selector.
35#[derive(Clone, Debug)]
36pub struct CompiledSelector {
37    prefix: Prefix,
38    regex: Regex,
39}
40
41impl CompiledSelector {
42    /// The raw prefix used to narrow keys before evaluating the payload regex.
43    pub fn prefix(&self) -> &[u8] {
44        self.prefix.as_bytes()
45    }
46
47    /// Return whether this selector accepts the key.
48    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/// A stream filter rejected during validation or compilation.
56///
57/// Messages describe the invalid selector or value filter and are safe to
58/// surface to clients.
59#[derive(Clone, Debug, thiserror::Error)]
60#[error("{0}")]
61pub struct InvalidFilter(String);
62
63/// Validated stream matchers compiled for repeated evaluation.
64#[derive(Clone, Debug)]
65pub struct CompiledMatchers {
66    keys: Vec<CompiledSelector>,
67    values: Option<CompiledFilters>,
68}
69
70impl CompiledMatchers {
71    /// Validate and compile a `StreamFilter`.
72    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    /// Return whether both predicates accept an already-resolved entry.
91    ///
92    /// The key predicate runs first, so key-excluded entries avoid the
93    /// value-filter scan. Implementations that resolve values lazily must call
94    /// [`Self::matches_key`] before resolving the value and
95    /// [`Self::matches_value`] afterward.
96    pub fn matches(&self, key: &[u8], value: &[u8]) -> bool {
97        self.matches_key(key) && self.matches_value(value)
98    }
99
100    /// Return whether any selector matches the key.
101    ///
102    /// Selector regular expressions evaluate bytes after the matching prefix.
103    pub fn matches_key(&self, key: &[u8]) -> bool {
104        self.keys.iter().any(|matcher| matcher.matches_key(key))
105    }
106
107    /// The compiled selectors in request order.
108    pub fn selectors(&self) -> &[CompiledSelector] {
109        &self.keys
110    }
111
112    /// Return whether the configured value filters match the value.
113    ///
114    /// Values pass when no value filter is configured.
115    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
122/// Validate and compile a `StreamFilter`. Shared between replay and live
123/// delivery so both paths match identically and regexes are compiled once per
124/// subscribe.
125pub(crate) fn compile_matchers(filter: &StreamFilter) -> Result<CompiledMatchers, ConnectError> {
126    CompiledMatchers::compile(filter).map_err(|e| ConnectError::invalid_argument(e.to_string()))
127}
128
129/// Apply compiled matchers while consuming the source entries.
130pub(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
142// TODO (#56): Add a separate remote stream notification abstraction for split deployments.
143/// In-process notification capability for stream subscribers.
144pub trait StreamNotifier: Send + Sync + 'static {
145    /// Atomically snapshot the visible batch frontier and return a notifier
146    /// that wakes when the frontier may have advanced.
147    fn subscribe(&self) -> StreamNotification;
148
149    /// Highest batch sequence currently visible to live subscribers.
150    fn current_sequence(&self) -> u64;
151
152    /// Announce that batches through `seq` may now be visible.
153    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    /// Announce a newly committed batch sequence to subscribers.
170    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    // Verbatim copy of the pre-refactor apply_filter loop. The differential
304    // test below pins the rewrite to the old first-match-wins semantics.
305    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            // Matches both prefix-1 selectors and must appear exactly once.
355            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            // Duplicate entry preserves multiplicity.
364            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}