aster_forge_cache/
bloom.rs

1//! Concurrent Bloom-filter primitives for cache-aside existence checks.
2
3use std::sync::Arc;
4
5use bloomfilter::Bloom;
6use parking_lot::{Mutex, RwLock};
7
8/// Configuration used to size a [`BloomFilter`].
9#[derive(Debug, Clone, Copy, PartialEq)]
10pub struct BloomConfig {
11    /// Expected number of stored keys before reserve capacity is applied.
12    pub expected_items: usize,
13    /// Desired false-positive probability in the range `(0, 1)`.
14    pub false_positive_rate: f64,
15}
16
17impl BloomConfig {
18    /// Creates a Bloom-filter configuration.
19    #[must_use]
20    pub const fn new(expected_items: usize, false_positive_rate: f64) -> Self {
21        Self {
22            expected_items,
23            false_positive_rate,
24        }
25    }
26}
27
28/// Bloom-filter construction or rebuild failure.
29#[derive(Debug, thiserror::Error)]
30pub enum BloomError {
31    /// The underlying Bloom filter rejected its capacity or probability.
32    #[error("invalid Bloom filter configuration: {0}")]
33    InvalidConfiguration(String),
34    /// A second rebuild was started before the active rebuild completed.
35    #[error("Bloom filter rebuild already in progress")]
36    RebuildInProgress,
37}
38
39/// Concurrent Bloom filter with atomic rebuild support.
40///
41/// Inserts made while a rebuild is active are recorded and applied to the new
42/// filter immediately before it replaces the old filter.
43pub struct BloomFilter {
44    inner: RwLock<Bloom<str>>,
45    rebuild_buffer: Mutex<Option<Vec<String>>>,
46}
47
48impl BloomFilter {
49    /// Creates a Bloom filter from the expected item count and false-positive rate.
50    ///
51    /// # Errors
52    ///
53    /// Returns [`BloomError::InvalidConfiguration`] when the item count or false-positive rate is
54    /// rejected by the underlying Bloom filter.
55    pub fn new(config: BloomConfig) -> Result<Self, BloomError> {
56        Ok(Self {
57            inner: RwLock::new(build_filter(config)?),
58            rebuild_buffer: Mutex::new(None),
59        })
60    }
61
62    /// Returns whether the key may be present.
63    pub fn contains(&self, key: &str) -> bool {
64        self.inner.read().check(key)
65    }
66
67    /// Inserts a key, including it in any active rebuild.
68    pub fn insert(&self, key: &str) {
69        let mut rebuild_buffer = self.rebuild_buffer.lock();
70        self.inner.write().set(key);
71        if let Some(pending) = rebuild_buffer.as_mut() {
72            pending.push(key.to_string());
73        }
74    }
75
76    /// Inserts several keys into the current filter.
77    pub fn insert_many<'a>(&self, keys: impl IntoIterator<Item = &'a str>) {
78        let mut rebuild_buffer = self.rebuild_buffer.lock();
79        let mut filter = self.inner.write();
80        for key in keys {
81            filter.set(key);
82            if let Some(pending) = rebuild_buffer.as_mut() {
83                pending.push(key.to_string());
84            }
85        }
86    }
87
88    /// Replaces the current filter with an empty filter using new sizing.
89    ///
90    /// # Errors
91    ///
92    /// Returns [`BloomError::InvalidConfiguration`] for invalid sizing, or
93    /// [`BloomError::RebuildInProgress`] while an atomic rebuild is active.
94    pub fn clear(&self, config: BloomConfig) -> Result<(), BloomError> {
95        let replacement = build_filter(config)?;
96        let rebuild_buffer = self.rebuild_buffer.lock();
97        if rebuild_buffer.is_some() {
98            return Err(BloomError::RebuildInProgress);
99        }
100        *self.inner.write() = replacement;
101        Ok(())
102    }
103
104    /// Starts an atomic rebuild session.
105    ///
106    /// Feed streamed batches into the returned session and call
107    /// [`BloomRebuild::commit`] after the source has completed. Dropping the
108    /// session leaves the previous filter active and stops buffering inserts.
109    ///
110    /// # Errors
111    ///
112    /// Returns [`BloomError::InvalidConfiguration`] for invalid sizing, or
113    /// [`BloomError::RebuildInProgress`] when another rebuild already owns the session.
114    pub fn start_rebuild(
115        self: &Arc<Self>,
116        config: BloomConfig,
117    ) -> Result<BloomRebuild, BloomError> {
118        let replacement = build_filter(config)?;
119        let mut rebuild_buffer = self.rebuild_buffer.lock();
120        if rebuild_buffer.is_some() {
121            return Err(BloomError::RebuildInProgress);
122        }
123        *rebuild_buffer = Some(Vec::new());
124        Ok(BloomRebuild {
125            owner: Arc::clone(self),
126            replacement: Some(replacement),
127            loaded: 0,
128        })
129    }
130}
131
132/// In-progress atomic rebuild created by [`BloomFilter::start_rebuild`].
133pub struct BloomRebuild {
134    owner: Arc<BloomFilter>,
135    replacement: Option<Bloom<str>>,
136    loaded: usize,
137}
138
139impl BloomRebuild {
140    /// Adds a streamed batch to the replacement filter.
141    pub fn insert_many<'a>(&mut self, keys: impl IntoIterator<Item = &'a str>) {
142        let Some(replacement) = self.replacement.as_mut() else {
143            return;
144        };
145        for key in keys {
146            replacement.set(key);
147            self.loaded = self.loaded.saturating_add(1);
148        }
149    }
150
151    /// Atomically installs the rebuilt filter and returns the number of keys
152    /// loaded from the source plus concurrent inserts captured during rebuild.
153    #[must_use]
154    pub fn commit(mut self) -> usize {
155        let Some(mut replacement) = self.replacement.take() else {
156            return self.loaded;
157        };
158        let mut rebuild_buffer = self.owner.rebuild_buffer.lock();
159        let buffered = rebuild_buffer.take().unwrap_or_default();
160        for key in &buffered {
161            replacement.set(key);
162        }
163        *self.owner.inner.write() = replacement;
164        self.loaded.saturating_add(buffered.len())
165    }
166}
167
168impl Drop for BloomRebuild {
169    fn drop(&mut self) {
170        if self.replacement.is_some() {
171            self.owner.rebuild_buffer.lock().take();
172        }
173    }
174}
175
176fn build_filter(config: BloomConfig) -> Result<Bloom<str>, BloomError> {
177    Bloom::new_for_fp_rate(
178        reserved_capacity(config.expected_items),
179        config.false_positive_rate,
180    )
181    .map_err(|error| BloomError::InvalidConfiguration(error.to_string()))
182}
183
184fn reserved_capacity(count: usize) -> usize {
185    let reserve = if count < 5_000 {
186        count / 2
187    } else if count < 100_000 {
188        count / 5
189    } else {
190        (count / 10).min(1_000_000)
191    };
192    count.saturating_add(reserve.max(1_000))
193}
194
195#[cfg(test)]
196mod tests {
197    use super::*;
198
199    fn config(expected_items: usize) -> BloomConfig {
200        BloomConfig::new(expected_items, 0.001)
201    }
202
203    #[test]
204    fn insert_and_clear_update_membership() {
205        let filter = BloomFilter::new(config(100)).expect("valid Bloom config");
206        assert!(!filter.contains("missing"));
207
208        filter.insert("present");
209        assert!(filter.contains("present"));
210
211        filter.clear(config(1_000)).expect("valid Bloom config");
212        assert!(!filter.contains("present"));
213    }
214
215    #[test]
216    fn rebuild_replaces_old_keys_and_keeps_concurrent_inserts() {
217        let filter = Arc::new(BloomFilter::new(config(100)).expect("valid Bloom config"));
218        filter.insert("old");
219
220        let mut rebuild = filter
221            .start_rebuild(config(2))
222            .expect("first rebuild should start");
223        rebuild.insert_many(["new-a", "new-b"]);
224        filter.insert("concurrent");
225
226        assert_eq!(rebuild.commit(), 3);
227        assert!(!filter.contains("old"));
228        assert!(filter.contains("new-a"));
229        assert!(filter.contains("new-b"));
230        assert!(filter.contains("concurrent"));
231    }
232
233    #[test]
234    fn dropped_rebuild_keeps_old_filter_and_releases_session() {
235        let filter = Arc::new(BloomFilter::new(config(100)).expect("valid Bloom config"));
236        filter.insert("old");
237
238        let rebuild = filter
239            .start_rebuild(config(1))
240            .expect("first rebuild should start");
241        filter.insert("concurrent");
242        drop(rebuild);
243
244        assert!(filter.contains("old"));
245        assert!(filter.contains("concurrent"));
246        assert!(filter.start_rebuild(config(1)).is_ok());
247    }
248
249    #[test]
250    fn overlapping_rebuilds_are_rejected() {
251        let filter = Arc::new(BloomFilter::new(config(100)).expect("valid Bloom config"));
252        let _active = filter
253            .start_rebuild(config(1))
254            .expect("first rebuild should start");
255
256        assert!(matches!(
257            filter.start_rebuild(config(1)),
258            Err(BloomError::RebuildInProgress)
259        ));
260    }
261
262    #[test]
263    fn clear_during_rebuild_is_rejected_without_discarding_the_session() {
264        let filter = Arc::new(BloomFilter::new(config(100)).expect("valid Bloom config"));
265        let mut rebuild = filter
266            .start_rebuild(config(1))
267            .expect("first rebuild should start");
268        rebuild.insert_many(["rebuilt"]);
269
270        assert!(matches!(
271            filter.clear(config(1_000)),
272            Err(BloomError::RebuildInProgress)
273        ));
274        assert_eq!(rebuild.commit(), 1);
275        assert!(filter.contains("rebuilt"));
276    }
277
278    #[test]
279    fn bulk_insert_preserves_all_keys() {
280        let filter = BloomFilter::new(config(100)).expect("valid Bloom config");
281        let keys: Vec<String> = (0..100).map(|index| format!("key-{index}")).collect();
282        filter.insert_many(keys.iter().map(String::as_str));
283
284        assert!(keys.iter().all(|key| filter.contains(key)));
285    }
286}