Skip to main content

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