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 then pulls batches from the log at its own
5//! pace, so live delivery is naturally paced by client reads instead of an
6//! internal 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#[derive(Clone)]
35pub(crate) struct CompiledKeyMatcher {
36    prefix: Prefix,
37    regex: Regex,
38}
39
40#[derive(Clone)]
41pub(crate) struct CompiledMatchers {
42    pub keys: Vec<CompiledKeyMatcher>,
43    pub values: Option<CompiledFilters>,
44}
45
46/// Validate and compile a `StreamFilter`. Shared between replay and live
47/// delivery so both paths match identically and regexes are compiled once per
48/// subscribe.
49pub(crate) fn compile_matchers(filter: &StreamFilter) -> Result<CompiledMatchers, ConnectError> {
50    validate_filter(filter).map_err(|e| ConnectError::invalid_argument(e.to_string()))?;
51    let keys = filter
52        .selectors
53        .iter()
54        .map(|mk| {
55            let regex = compile_payload_regex(&mk.payload_regex)
56                .map_err(|e| ConnectError::invalid_argument(e.to_string()))?;
57            let prefix = Prefix::new(mk.prefix.clone())
58                .map_err(|e| ConnectError::invalid_argument(e.to_string()))?;
59            Ok(CompiledKeyMatcher { prefix, regex })
60        })
61        .collect::<Result<Vec<_>, ConnectError>>()?;
62    let values = CompiledFilters::compile(&filter.value_filters)
63        .map_err(|e| ConnectError::invalid_argument(format!("invalid value_filter: {e}")))?;
64    Ok(CompiledMatchers { keys, values })
65}
66
67/// Apply a compiled filter to a batch. First-match-wins per entry.
68pub(crate) fn apply_filter(matchers: &CompiledMatchers, kvs: &[Entry]) -> Vec<Entry> {
69    let mut out = Vec::with_capacity(kvs.len());
70    'outer: for kv in kvs {
71        let v = kv.value.as_ref();
72        let value_ok = matchers.values.as_ref().is_none_or(|m| m.matches(v));
73        if !value_ok {
74            continue;
75        }
76        for matcher in &matchers.keys {
77            let Some(payload) = matcher.prefix.strip_slice(&kv.key) else {
78                continue;
79            };
80            if matcher.regex.is_match(payload) {
81                out.push(kv.clone());
82                continue 'outer;
83            }
84        }
85    }
86    out
87}
88
89#[derive(Clone)]
90pub struct StreamNotification {
91    pub current_sequence: u64,
92    pub notify: Arc<Notify>,
93}
94
95// TODO (#56): Add a separate remote stream notification abstraction for split deployments.
96/// In-process notification capability for stream subscribers.
97pub trait StreamNotifier: Send + Sync + 'static {
98    /// Atomically snapshot the visible batch frontier and return a notifier
99    /// that wakes when the frontier may have advanced.
100    fn subscribe(&self) -> StreamNotification;
101
102    /// Highest batch sequence currently visible to live subscribers.
103    fn current_sequence(&self) -> u64;
104
105    /// Announce that batches through `seq` may now be visible.
106    fn advance(&self, seq: u64);
107}
108
109pub struct StreamHub {
110    published_sequence: AtomicU64,
111    notify: Arc<Notify>,
112}
113
114impl StreamHub {
115    pub fn new(initial_sequence: u64) -> Self {
116        Self {
117            published_sequence: AtomicU64::new(initial_sequence),
118            notify: Arc::new(Notify::new()),
119        }
120    }
121
122    /// Announce a newly committed batch sequence to subscribers.
123    pub fn publish(&self, seq: u64) {
124        self.advance(seq);
125    }
126}
127
128impl StreamNotifier for StreamHub {
129    fn subscribe(&self) -> StreamNotification {
130        StreamNotification {
131            current_sequence: self.published_sequence.load(Ordering::Acquire),
132            notify: self.notify.clone(),
133        }
134    }
135
136    fn current_sequence(&self) -> u64 {
137        self.published_sequence.load(Ordering::Acquire)
138    }
139
140    fn advance(&self, seq: u64) {
141        self.published_sequence.fetch_max(seq, Ordering::SeqCst);
142        self.notify.notify_waiters();
143    }
144}
145
146#[cfg(test)]
147mod tests {
148    use super::*;
149    use bytes::Bytes;
150    use exoware_sdk::keys::Key;
151    use exoware_sdk::kv_codec::Utf8;
152    use exoware_sdk::selector::Selector;
153    use exoware_sdk::stream_filter::Filter;
154
155    fn filter(prefix: u8, regex: &str) -> StreamFilter {
156        StreamFilter {
157            selectors: vec![Selector {
158                prefix: Bytes::from(vec![prefix]),
159                payload_regex: Utf8::from(regex),
160            }],
161            value_filters: vec![],
162        }
163    }
164
165    fn filter_with_values(prefix: u8, regex: &str, value_filters: Vec<Filter>) -> StreamFilter {
166        StreamFilter {
167            selectors: vec![Selector {
168                prefix: Bytes::from(vec![prefix]),
169                payload_regex: Utf8::from(regex),
170            }],
171            value_filters,
172        }
173    }
174
175    fn key(family: u8, payload: &[u8]) -> Key {
176        Prefix::from_byte(family).encode(payload).unwrap()
177    }
178
179    fn kv(family: u8, payload: &[u8], value: &'static [u8]) -> Entry {
180        Entry {
181            key: key(family, payload).to_vec(),
182            value: Bytes::from_static(value),
183            ..Default::default()
184        }
185    }
186
187    #[test]
188    fn publish_sequence_is_monotonic() {
189        let hub = StreamHub::new(7);
190        assert_eq!(hub.current_sequence(), 7);
191        hub.publish(3);
192        assert_eq!(hub.current_sequence(), 7);
193        hub.publish(9);
194        assert_eq!(hub.current_sequence(), 9);
195    }
196
197    #[test]
198    fn subscribe_snapshots_current_sequence() {
199        let hub = StreamHub::new(11);
200        let subscription = hub.subscribe();
201        assert_eq!(subscription.current_sequence, 11);
202    }
203
204    #[test]
205    fn apply_filter_still_selects_matching_entries() {
206        let matchers = compile_matchers(&filter(1, "(?s).*")).unwrap();
207        let kvs = vec![kv(1, b"hit", b"v1"), kv(2, b"miss", b"v2")];
208        let entries = apply_filter(&matchers, &kvs);
209        assert_eq!(entries.len(), 1);
210        assert_eq!(entries[0].value.as_ref(), b"v1");
211    }
212
213    #[test]
214    fn subscribe_rejects_invalid_filter() {
215        let bad = StreamFilter {
216            selectors: vec![],
217            value_filters: vec![],
218        };
219        assert!(compile_matchers(&bad).is_err());
220    }
221
222    #[test]
223    fn value_filter_intersects_with_key_filter() {
224        let matchers = compile_matchers(&filter_with_values(
225            1,
226            "(?s).*",
227            vec![Filter::Regex("^keep$".into())],
228        ))
229        .unwrap();
230        let kvs = vec![kv(1, b"a", b"keep"), kv(1, b"b", b"drop")];
231        let entries = apply_filter(&matchers, &kvs);
232        assert_eq!(entries.len(), 1);
233        assert_eq!(entries[0].value.as_ref(), b"keep");
234    }
235
236    #[test]
237    fn value_filter_exact_match() {
238        let matchers = compile_matchers(&filter_with_values(
239            1,
240            "(?s).*",
241            vec![Filter::Exact(Bytes::from_static(b"target"))],
242        ))
243        .unwrap();
244        let kvs = vec![kv(1, b"a", b"target"), kv(1, b"b", b"other")];
245        let entries = apply_filter(&matchers, &kvs);
246        assert_eq!(entries.len(), 1);
247        assert_eq!(entries[0].value.as_ref(), b"target");
248    }
249
250    #[test]
251    fn value_filter_empty_accepts_all_matching_keys() {
252        let matchers = compile_matchers(&filter(1, "(?s).*")).unwrap();
253        let kvs = vec![kv(1, b"a", b"one"), kv(1, b"b", b"two")];
254        let entries = apply_filter(&matchers, &kvs);
255        assert_eq!(entries.len(), 2);
256    }
257}