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)]
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
46pub(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
67pub(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
95pub trait StreamNotifier: Send + Sync + 'static {
98 fn subscribe(&self) -> StreamNotification;
101
102 fn current_sequence(&self) -> u64;
104
105 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 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}