Skip to main content

weavatrix_search_vector/mutable/
rebuild.rs

1use super::MutableVectorIndex;
2use crate::error::SearchError;
3use crate::hnsw::VectorIndex;
4use std::sync::Arc;
5
6impl MutableVectorIndex {
7    /// Rebuilds only the mutable delta into its own HNSW graph.
8    ///
9    /// This bounds exact staged scanning without rebuilding the immutable
10    /// base. Concurrent mutations are retried and never lost.
11    ///
12    /// # Errors
13    ///
14    /// Returns a build error or [`SearchError::MutationConflict`] after three
15    /// conflicting attempts.
16    pub fn seal_delta(&self) -> Result<(), SearchError> {
17        for _ in 0..3 {
18            let (generation, sealed, pending, deleted) = {
19                let state = self
20                    .state
21                    .read()
22                    .unwrap_or_else(std::sync::PoisonError::into_inner);
23                (
24                    state.generation,
25                    state.sealed.clone(),
26                    state.pending.clone(),
27                    state.deleted.clone(),
28                )
29            };
30            if pending.is_empty() {
31                return Ok(());
32            }
33            let mut owned = Vec::new();
34            owned
35                .try_reserve_exact(
36                    sealed
37                        .as_ref()
38                        .map_or(0, |index| index.len())
39                        .saturating_add(pending.len()),
40                )
41                .map_err(|_| SearchError::AllocationFailed)?;
42            if let Some(index) = sealed {
43                for key in index.keys() {
44                    if deleted.contains(&key) || pending.contains_key(&key) {
45                        continue;
46                    }
47                    owned.push((
48                        key,
49                        index
50                            .vector(key)
51                            .ok_or(SearchError::MissingKey(key))?
52                            .to_vec(),
53                    ));
54                }
55            }
56            owned.extend(pending.iter().map(|(key, vector)| (*key, vector.clone())));
57            owned.sort_unstable_by_key(|record| record.0);
58            let borrowed = owned
59                .iter()
60                .map(|(key, vector)| (*key, vector.as_slice()))
61                .collect::<Vec<_>>();
62            let rebuilt = Arc::new(VectorIndex::build(self.config.clone(), &borrowed)?);
63            let mut state = self
64                .state
65                .write()
66                .unwrap_or_else(std::sync::PoisonError::into_inner);
67            if state.generation != generation {
68                continue;
69            }
70            state.sealed = Some(rebuilt);
71            state.pending.clear();
72            state.generation = state.generation.wrapping_add(1);
73            return Ok(());
74        }
75        Err(SearchError::MutationConflict)
76    }
77
78    /// Deterministically rebuilds the immutable base. Concurrent mutations are
79    /// detected and retried without losing updates.
80    ///
81    /// # Errors
82    ///
83    /// Returns a build error or [`SearchError::MutationConflict`] after three
84    /// conflicting rebuilds.
85    pub fn compact(&self) -> Result<(), SearchError> {
86        for _ in 0..3 {
87            let (generation, base, sealed, pending, deleted) = {
88                let state = self
89                    .state
90                    .read()
91                    .unwrap_or_else(std::sync::PoisonError::into_inner);
92                (
93                    state.generation,
94                    Arc::clone(&state.base),
95                    state.sealed.clone(),
96                    state.pending.clone(),
97                    state.deleted.clone(),
98                )
99            };
100            let mut owned = Vec::new();
101            owned
102                .try_reserve_exact(
103                    base.len()
104                        .saturating_sub(deleted.len())
105                        .saturating_add(sealed.as_ref().map_or(0, |index| index.len()))
106                        .saturating_add(pending.len()),
107                )
108                .map_err(|_| SearchError::AllocationFailed)?;
109            for key in base.keys() {
110                if deleted.contains(&key)
111                    || pending.contains_key(&key)
112                    || sealed
113                        .as_ref()
114                        .is_some_and(|index| index.vector(key).is_some())
115                {
116                    continue;
117                }
118                let vector = base.vector(key).ok_or(SearchError::MissingKey(key))?;
119                owned.push((key, vector.to_vec()));
120            }
121            if let Some(index) = sealed {
122                for key in index.keys() {
123                    if deleted.contains(&key) || pending.contains_key(&key) {
124                        continue;
125                    }
126                    let vector = index.vector(key).ok_or(SearchError::MissingKey(key))?;
127                    owned.push((key, vector.to_vec()));
128                }
129            }
130            owned.extend(pending.iter().map(|(key, vector)| (*key, vector.clone())));
131            owned.sort_unstable_by_key(|record| record.0);
132            let borrowed = owned
133                .iter()
134                .map(|(key, vector)| (*key, vector.as_slice()))
135                .collect::<Vec<_>>();
136            let rebuilt = Arc::new(VectorIndex::build(self.config.clone(), &borrowed)?);
137            let mut state = self
138                .state
139                .write()
140                .unwrap_or_else(std::sync::PoisonError::into_inner);
141            if state.generation != generation {
142                continue;
143            }
144            state.base = rebuilt;
145            state.sealed = None;
146            state.pending.clear();
147            state.deleted.clear();
148            return Ok(());
149        }
150        Err(SearchError::MutationConflict)
151    }
152}