Skip to main content

hermes_core/index/
reader.rs

1//! IndexReader - manages Searcher with reload policy (native only)
2//!
3//! The IndexReader periodically reloads its Searcher to pick up new segments.
4//! Uses SegmentManager as authoritative source for segment state.
5
6use std::sync::Arc;
7use std::sync::atomic::{AtomicBool, Ordering};
8
9use arc_swap::ArcSwap;
10use parking_lot::RwLock;
11
12use crate::directories::DirectoryWriter;
13use crate::dsl::Schema;
14use crate::error::Result;
15
16use super::Searcher;
17use super::searcher::SearcherResources;
18
19/// IndexReader - manages Searcher with reload policy
20///
21/// The IndexReader periodically reloads its Searcher to pick up new segments.
22/// Uses SegmentManager as authoritative source for segment state (avoids race conditions).
23/// Combined searcher + segment IDs, swapped atomically via ArcSwap (wait-free reads).
24struct SearcherState<D: DirectoryWriter + 'static> {
25    searcher: Arc<Searcher<D>>,
26    segment_ids: Vec<String>,
27}
28
29/// Cancellation-safe ownership of the reload flag. Async reload checks may be
30/// dropped at any await point; resetting manually only on normal return leaves
31/// every future reload disabled after request cancellation or panic.
32struct ReloadGuard<'a>(&'a AtomicBool);
33
34impl Drop for ReloadGuard<'_> {
35    fn drop(&mut self) {
36        self.0.store(false, Ordering::Release);
37    }
38}
39
40pub struct IndexReader<D: DirectoryWriter + 'static> {
41    /// Schema
42    schema: Arc<Schema>,
43    /// Segment manager - authoritative source for segments
44    segment_manager: Arc<crate::merge::SegmentManager<D>>,
45    /// Current searcher + segment IDs (ArcSwap for wait-free reads)
46    state: ArcSwap<SearcherState<D>>,
47    /// Cache and CPU policy preserved across every searcher reload.
48    resources: SearcherResources,
49    /// Last reload check time
50    last_reload_check: RwLock<std::time::Instant>,
51    /// Reload check interval (default 1 second)
52    reload_check_interval: std::time::Duration,
53    /// Guard against concurrent reloads
54    reloading: AtomicBool,
55}
56
57impl<D: DirectoryWriter + 'static> IndexReader<D> {
58    /// Create a new IndexReader from a segment manager
59    ///
60    /// Centroids are loaded dynamically from metadata on each reload,
61    /// so the reader always picks up centroids trained after Index::create().
62    pub async fn from_segment_manager(
63        schema: Arc<Schema>,
64        segment_manager: Arc<crate::merge::SegmentManager<D>>,
65        term_cache_blocks: usize,
66        reload_interval_ms: u64,
67    ) -> Result<Self> {
68        const STANDALONE_STORE_CACHE_BYTES: usize = 32 * 1024 * 1024;
69        let resources = SearcherResources::new(
70            term_cache_blocks,
71            STANDALONE_STORE_CACHE_BYTES,
72            crate::default_search_threads(),
73            4,
74        )?;
75        Self::from_segment_manager_with_resources(
76            schema,
77            segment_manager,
78            reload_interval_ms,
79            resources,
80        )
81        .await
82    }
83
84    /// Internal constructor used by `Index` to preserve its configured cache
85    /// and search CPU policy across reader reloads.
86    pub(crate) async fn from_segment_manager_with_resources(
87        schema: Arc<Schema>,
88        segment_manager: Arc<crate::merge::SegmentManager<D>>,
89        reload_interval_ms: u64,
90        resources: SearcherResources,
91    ) -> Result<Self> {
92        // Get initial segment IDs
93        let initial_segment_ids = segment_manager.get_segment_ids().await;
94
95        let reader = Self::create_reader(&schema, &segment_manager, resources.clone()).await?;
96
97        Ok(Self {
98            schema,
99            segment_manager,
100            state: ArcSwap::from_pointee(SearcherState {
101                searcher: Arc::new(reader),
102                segment_ids: initial_segment_ids,
103            }),
104            resources,
105            last_reload_check: RwLock::new(std::time::Instant::now()),
106            reload_check_interval: std::time::Duration::from_millis(reload_interval_ms),
107            reloading: AtomicBool::new(false),
108        })
109    }
110
111    /// Create a new reader with fresh snapshot from segment manager
112    ///
113    /// Captures segment IDs and their trained vector generation together.
114    async fn create_reader(
115        schema: &Arc<Schema>,
116        segment_manager: &Arc<crate::merge::SegmentManager<D>>,
117        resources: SearcherResources,
118    ) -> Result<Searcher<D>> {
119        let snapshot = segment_manager.acquire_snapshot().await;
120
121        // The snapshot carries the exact trained-artifact generation paired
122        // with its segment IDs. This remains correct across atomic retrains.
123        let trained = snapshot
124            .trained_vectors()
125            .unwrap_or_else(|| Arc::new(crate::segment::TrainedVectorStructures::default()));
126
127        Searcher::from_snapshot(
128            segment_manager.directory(),
129            Arc::clone(schema),
130            snapshot,
131            trained,
132            resources,
133        )
134        .await
135    }
136
137    /// Set reload check interval
138    pub fn set_reload_interval(&mut self, interval: std::time::Duration) {
139        self.reload_check_interval = interval;
140    }
141
142    /// Get current searcher (reloads only if segments changed)
143    ///
144    /// Wait-free read path via ArcSwap::load(). Reload checks are guarded
145    /// by an AtomicBool to prevent concurrent reloads.
146    pub async fn searcher(&self) -> Result<Arc<Searcher<D>>> {
147        // Check if we should check for segment changes
148        let should_check = {
149            let last = self.last_reload_check.read();
150            last.elapsed() >= self.reload_check_interval
151        };
152
153        if should_check {
154            // Try to acquire the reload guard (non-blocking)
155            if self
156                .reloading
157                .compare_exchange(false, true, Ordering::Acquire, Ordering::Relaxed)
158                .is_ok()
159            {
160                let _reload_guard = ReloadGuard(&self.reloading);
161                // We won the race — do the reload check
162                self.do_reload_check().await?;
163            }
164            // Otherwise another reload is in progress — just return current searcher
165        }
166
167        // Wait-free load (no lock contention with reloads)
168        Ok(Arc::clone(&self.state.load().searcher))
169    }
170
171    /// Actual reload check (called under the `reloading` guard)
172    async fn do_reload_check(&self) -> Result<()> {
173        *self.last_reload_check.write() = std::time::Instant::now();
174
175        // Get current segment IDs from segment manager
176        let new_segment_ids = self.segment_manager.get_segment_ids().await;
177
178        // Check if segments actually changed (wait-free read)
179        let segments_changed = {
180            let state = self.state.load();
181            state.segment_ids != new_segment_ids
182        };
183
184        if segments_changed {
185            let old_count = self.state.load().segment_ids.len();
186            let new_count = new_segment_ids.len();
187            log::info!(
188                "[index_reload] index={} old_count={} new_count={}",
189                self.schema.index_label(),
190                old_count,
191                new_count
192            );
193            self.reload_with_segments(new_segment_ids).await?;
194        }
195        Ok(())
196    }
197
198    /// Force reload reader with fresh snapshot.
199    ///
200    /// Waits for any in-progress reload (from `searcher()`) to finish, then
201    /// performs its own reload with the latest segment IDs. This guarantees
202    /// the reload actually happens — unlike `searcher()` which silently skips
203    /// if another reload is in progress.
204    pub async fn reload(&self) -> Result<()> {
205        // Wait for any in-progress reload to finish, then acquire the guard.
206        // This is critical: a concurrent do_reload_check() may have started
207        // before a commit, so its reload won't see the new segments.
208        loop {
209            if self
210                .reloading
211                .compare_exchange(false, true, Ordering::Acquire, Ordering::Relaxed)
212                .is_ok()
213            {
214                break;
215            }
216            tokio::task::yield_now().await;
217        }
218        let _reload_guard = ReloadGuard(&self.reloading);
219        let new_segment_ids = self.segment_manager.get_segment_ids().await;
220
221        // Fast path: skip reload if segments haven't changed
222        let segments_changed = {
223            let state = self.state.load();
224            state.segment_ids != new_segment_ids
225        };
226
227        if segments_changed {
228            self.reload_with_segments(new_segment_ids).await
229        } else {
230            log::debug!(
231                "[reload] index={} segments unchanged, skipping",
232                self.schema.index_label()
233            );
234            Ok(())
235        }
236    }
237
238    /// Internal reload with specific segment IDs.
239    /// Reuses existing segment readers for unchanged segments (avoids re-opening
240    /// mmaps, fast fields, sparse indexes, etc.).
241    /// Atomic swap via ArcSwap::store (wait-free for readers).
242    async fn reload_with_segments(&self, new_segment_ids: Vec<String>) -> Result<()> {
243        // Collect existing segment readers for reuse
244        let existing_segments: Vec<Arc<crate::segment::SegmentReader>> =
245            self.state.load().searcher.segment_readers().to_vec();
246
247        let snapshot = self.segment_manager.acquire_snapshot().await;
248
249        // Use the trained-artifact generation captured with this snapshot.
250        let trained = snapshot
251            .trained_vectors()
252            .unwrap_or_else(|| Arc::new(crate::segment::TrainedVectorStructures::default()));
253
254        let new_reader = Searcher::from_snapshot_reuse(
255            self.segment_manager.directory(),
256            Arc::clone(&self.schema),
257            snapshot,
258            trained,
259            self.resources.clone(),
260            &existing_segments,
261        )
262        .await?;
263
264        // Atomic swap — readers see old or new state, never a torn read
265        self.state.store(Arc::new(SearcherState {
266            searcher: Arc::new(new_reader),
267            segment_ids: new_segment_ids,
268        }));
269
270        Ok(())
271    }
272
273    /// Get schema
274    pub fn schema(&self) -> &Schema {
275        &self.schema
276    }
277}
278
279#[cfg(test)]
280mod tests {
281    use super::*;
282
283    #[test]
284    fn reload_guard_releases_flag_on_unwind() {
285        let reloading = AtomicBool::new(true);
286        let result = std::panic::catch_unwind(|| {
287            let _guard = ReloadGuard(&reloading);
288            panic!("cancel reload");
289        });
290        assert!(result.is_err());
291        assert!(!reloading.load(Ordering::Acquire));
292    }
293}