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] old_count={} new_count={}",
189                old_count,
190                new_count
191            );
192            self.reload_with_segments(new_segment_ids).await?;
193        }
194        Ok(())
195    }
196
197    /// Force reload reader with fresh snapshot.
198    ///
199    /// Waits for any in-progress reload (from `searcher()`) to finish, then
200    /// performs its own reload with the latest segment IDs. This guarantees
201    /// the reload actually happens — unlike `searcher()` which silently skips
202    /// if another reload is in progress.
203    pub async fn reload(&self) -> Result<()> {
204        // Wait for any in-progress reload to finish, then acquire the guard.
205        // This is critical: a concurrent do_reload_check() may have started
206        // before a commit, so its reload won't see the new segments.
207        loop {
208            if self
209                .reloading
210                .compare_exchange(false, true, Ordering::Acquire, Ordering::Relaxed)
211                .is_ok()
212            {
213                break;
214            }
215            tokio::task::yield_now().await;
216        }
217        let _reload_guard = ReloadGuard(&self.reloading);
218        let new_segment_ids = self.segment_manager.get_segment_ids().await;
219
220        // Fast path: skip reload if segments haven't changed
221        let segments_changed = {
222            let state = self.state.load();
223            state.segment_ids != new_segment_ids
224        };
225
226        if segments_changed {
227            self.reload_with_segments(new_segment_ids).await
228        } else {
229            log::debug!("[reload] segments unchanged, skipping");
230            Ok(())
231        }
232    }
233
234    /// Internal reload with specific segment IDs.
235    /// Reuses existing segment readers for unchanged segments (avoids re-opening
236    /// mmaps, fast fields, sparse indexes, etc.).
237    /// Atomic swap via ArcSwap::store (wait-free for readers).
238    async fn reload_with_segments(&self, new_segment_ids: Vec<String>) -> Result<()> {
239        // Collect existing segment readers for reuse
240        let existing_segments: Vec<Arc<crate::segment::SegmentReader>> =
241            self.state.load().searcher.segment_readers().to_vec();
242
243        let snapshot = self.segment_manager.acquire_snapshot().await;
244
245        // Use the trained-artifact generation captured with this snapshot.
246        let trained = snapshot
247            .trained_vectors()
248            .unwrap_or_else(|| Arc::new(crate::segment::TrainedVectorStructures::default()));
249
250        let new_reader = Searcher::from_snapshot_reuse(
251            self.segment_manager.directory(),
252            Arc::clone(&self.schema),
253            snapshot,
254            trained,
255            self.resources.clone(),
256            &existing_segments,
257        )
258        .await?;
259
260        // Atomic swap — readers see old or new state, never a torn read
261        self.state.store(Arc::new(SearcherState {
262            searcher: Arc::new(new_reader),
263            segment_ids: new_segment_ids,
264        }));
265
266        Ok(())
267    }
268
269    /// Get schema
270    pub fn schema(&self) -> &Schema {
271        &self.schema
272    }
273}
274
275#[cfg(test)]
276mod tests {
277    use super::*;
278
279    #[test]
280    fn reload_guard_releases_flag_on_unwind() {
281        let reloading = AtomicBool::new(true);
282        let result = std::panic::catch_unwind(|| {
283            let _guard = ReloadGuard(&reloading);
284            panic!("cancel reload");
285        });
286        assert!(result.is_err());
287        assert!(!reloading.load(Ordering::Acquire));
288    }
289}