hermes_core/index/
reader.rs1use 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
19struct SearcherState<D: DirectoryWriter + 'static> {
25 searcher: Arc<Searcher<D>>,
26 segment_ids: Vec<String>,
27}
28
29struct 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: Arc<Schema>,
43 segment_manager: Arc<crate::merge::SegmentManager<D>>,
45 state: ArcSwap<SearcherState<D>>,
47 resources: SearcherResources,
49 last_reload_check: RwLock<std::time::Instant>,
51 reload_check_interval: std::time::Duration,
53 reloading: AtomicBool,
55}
56
57impl<D: DirectoryWriter + 'static> IndexReader<D> {
58 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 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 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 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 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 pub fn set_reload_interval(&mut self, interval: std::time::Duration) {
139 self.reload_check_interval = interval;
140 }
141
142 pub async fn searcher(&self) -> Result<Arc<Searcher<D>>> {
147 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 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 self.do_reload_check().await?;
163 }
164 }
166
167 Ok(Arc::clone(&self.state.load().searcher))
169 }
170
171 async fn do_reload_check(&self) -> Result<()> {
173 *self.last_reload_check.write() = std::time::Instant::now();
174
175 let new_segment_ids = self.segment_manager.get_segment_ids().await;
177
178 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 pub async fn reload(&self) -> Result<()> {
205 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 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 async fn reload_with_segments(&self, new_segment_ids: Vec<String>) -> Result<()> {
243 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 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 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 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}