1use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwapOption;
49use tokio::sync::Mutex as AsyncMutex;
50use tokio::sync::{Notify, Semaphore};
51use tokio::task::JoinHandle;
52
53use crate::directories::DirectoryWriter;
54use crate::error::{Error, Result};
55use crate::index::{IndexMetadata, SegmentMetaInfo};
56use crate::segment::{
57 SegmentFiles, SegmentId, SegmentMeta, SegmentSnapshot, SegmentTracker, TrainedVectorStructures,
58};
59#[cfg(feature = "native")]
60use crate::segment::{SegmentMerger, SegmentReader};
61
62use super::{MergePolicy, SegmentInfo};
63
64struct ActiveOperationState {
74 segment_ids: HashSet<String>,
75 operation_tokens: HashSet<u64>,
76 next_operation_token: u64,
77 accepting: bool,
78}
79
80struct ActiveSegmentOperations {
81 inner: parking_lot::Mutex<ActiveOperationState>,
82 idle: Notify,
83 shutdown: Notify,
84}
85
86impl ActiveSegmentOperations {
87 fn new() -> Self {
88 Self {
89 inner: parking_lot::Mutex::new(ActiveOperationState {
90 segment_ids: HashSet::new(),
91 operation_tokens: HashSet::new(),
92 next_operation_token: 0,
93 accepting: true,
94 }),
95 idle: Notify::new(),
96 shutdown: Notify::new(),
97 }
98 }
99
100 fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
103 let mut inner = self.inner.lock();
104 if !inner.accepting {
105 log::debug!("[segment_lifecycle] rejected operation during shutdown");
106 return None;
107 }
108 for id in &segment_ids {
110 if inner.segment_ids.contains(id) {
111 log::debug!(
112 "[segment_lifecycle] rejected: {} overlaps with an active operation ({} active IDs)",
113 id,
114 inner.segment_ids.len()
115 );
116 return None;
117 }
118 }
119 log::debug!(
120 "[segment_lifecycle] registered {} IDs (total active: {})",
121 segment_ids.len(),
122 inner.segment_ids.len() + segment_ids.len()
123 );
124 let operation_token = inner.next_operation_token;
125 let next_operation_token = operation_token.checked_add(1)?;
126 for id in &segment_ids {
127 inner.segment_ids.insert(id.clone());
128 }
129 inner.next_operation_token = next_operation_token;
130 inner.operation_tokens.insert(operation_token);
131 Some(SegmentOperationGuard {
132 active_operations: Arc::clone(self),
133 segment_ids,
134 operation_token,
135 })
136 }
137
138 fn snapshot(&self) -> HashSet<String> {
140 self.inner.lock().segment_ids.clone()
141 }
142
143 fn operation_tokens_snapshot(&self) -> HashSet<u64> {
148 self.inner.lock().operation_tokens.clone()
149 }
150
151 fn stop_accepting(&self) {
154 let mut inner = self.inner.lock();
155 inner.accepting = false;
156 self.shutdown.notify_waiters();
157 if inner.segment_ids.is_empty() {
158 self.idle.notify_waiters();
159 }
160 }
161
162 fn is_accepting(&self) -> bool {
163 self.inner.lock().accepting
164 }
165
166 async fn wait_until_idle(&self) {
170 loop {
171 let notified = self.idle.notified();
172 if self.inner.lock().segment_ids.is_empty() {
173 return;
174 }
175 notified.await;
176 }
177 }
178
179 async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
180 while !operations.is_empty() {
181 let notified = self.idle.notified();
182 if self.inner.lock().operation_tokens.is_disjoint(operations) {
183 return;
184 }
185 notified.await;
186 }
187 }
188
189 async fn wait_for_shutdown(&self) {
192 loop {
193 let notified = self.shutdown.notified();
194 if !self.inner.lock().accepting {
195 return;
196 }
197 notified.await;
198 }
199 }
200}
201
202pub(crate) struct SegmentOperationGuard {
206 active_operations: Arc<ActiveSegmentOperations>,
207 segment_ids: Vec<String>,
208 operation_token: u64,
209}
210
211impl Drop for SegmentOperationGuard {
212 fn drop(&mut self) {
213 let mut inner = self.active_operations.inner.lock();
214 for id in &self.segment_ids {
215 inner.segment_ids.remove(id);
216 }
217 inner.operation_tokens.remove(&self.operation_token);
218 self.active_operations.idle.notify_waiters();
221 if inner.segment_ids.is_empty() {
222 debug_assert!(inner.operation_tokens.is_empty());
223 }
224 }
225}
226
227struct VectorArtifactUpdateLease {
235 updating: Arc<AtomicBool>,
236}
237
238impl Drop for VectorArtifactUpdateLease {
239 fn drop(&mut self) {
240 self.updating.store(false, Ordering::Release);
241 }
242}
243
244#[derive(Clone)]
245pub(crate) struct VectorArtifactUpdateGuard {
246 _lease: Arc<VectorArtifactUpdateLease>,
247}
248
249static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
253
254const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
255const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
256
257#[derive(Default)]
258struct MergeRetryState {
259 retry_after: Option<std::time::Instant>,
260 consecutive_failures: u32,
261}
262
263fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
264 let shift = consecutive_failures.saturating_sub(1).min(16);
265 MERGE_RETRY_BASE_DELAY
266 .checked_mul(1u32 << shift)
267 .unwrap_or(MERGE_RETRY_MAX_DELAY)
268 .min(MERGE_RETRY_MAX_DELAY)
269}
270
271fn try_spawn_lifecycle<F>(
278 handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
279 runtime: &tokio::runtime::Handle,
280 future: F,
281) -> bool
282where
283 F: std::future::Future<Output = ()> + Send + 'static,
284{
285 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
286 let mut handles = handles.lock();
287 handles.retain(|handle| !handle.is_finished());
288 handles.push(runtime.spawn(future));
289 }))
290 .is_ok()
291}
292
293struct OutputCleanupGuard {
301 segment_id: SegmentId,
302 cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
303}
304
305impl OutputCleanupGuard {
306 fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
307 Self {
308 segment_id,
309 cleanup: Some(cleanup),
310 }
311 }
312
313 fn disarm(&mut self) {
314 self.cleanup = None;
315 }
316}
317
318impl Drop for OutputCleanupGuard {
319 fn drop(&mut self) {
320 if let Some(cleanup) = self.cleanup.take() {
321 cleanup(self.segment_id);
322 }
323 }
324}
325
326struct ManagerState {
328 metadata: IndexMetadata,
329 merge_policy: Box<dyn MergePolicy>,
330}
331
332#[cfg(feature = "native")]
333struct MergeTaskError {
334 error: Error,
335 unavailable_segments: Vec<String>,
336}
337
338#[cfg(feature = "native")]
339impl MergeTaskError {
340 fn source(segment_id: String, error: Error) -> Self {
341 Self {
342 error,
343 unavailable_segments: vec![segment_id],
344 }
345 }
346
347 fn sources(segment_ids: Vec<String>, error: Error) -> Self {
348 Self {
349 error,
350 unavailable_segments: segment_ids,
351 }
352 }
353}
354
355#[cfg(feature = "native")]
356impl From<Error> for MergeTaskError {
357 fn from(error: Error) -> Self {
358 Self {
359 error,
360 unavailable_segments: Vec::new(),
361 }
362 }
363}
364
365#[cfg(feature = "native")]
366fn is_deterministic_source_error(error: &Error) -> bool {
367 matches!(error, Error::Corruption(_) | Error::Serialization(_))
368 || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
369}
370
371#[cfg(feature = "native")]
372fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
373 if is_deterministic_source_error(&error) {
374 MergeTaskError::source(segment_id, error)
375 } else {
376 MergeTaskError::from(error)
380 }
381}
382
383#[cfg(feature = "native")]
384type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
385
386pub struct SegmentManager<D: DirectoryWriter + 'static> {
390 state: Arc<AsyncMutex<ManagerState>>,
392
393 active_operations: Arc<ActiveSegmentOperations>,
395
396 quarantined_segments: parking_lot::Mutex<HashSet<String>>,
401
402 merge_retry: parking_lot::Mutex<MergeRetryState>,
405
406 reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
410
411 merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
413
414 global_merge_wakeup_pending: AtomicBool,
418
419 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
423
424 trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
428
429 vector_artifact_update: Arc<AtomicBool>,
433
434 tracker: Arc<SegmentTracker>,
436
437 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
439
440 directory: Arc<D>,
442 schema: Arc<crate::dsl::Schema>,
444 term_cache_blocks: usize,
446 merge_permits: Arc<Semaphore>,
450 global_merge_permits: Arc<Semaphore>,
452 reorder_permits: Arc<Semaphore>,
456 reorder_on_merge: bool,
461 merge_bp_time_budget: Option<std::time::Duration>,
465 bp_memory_budget_bytes: usize,
468 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
471}
472
473impl<D: DirectoryWriter + 'static> SegmentManager<D> {
474 #[allow(clippy::too_many_arguments)]
476 pub fn new(
477 directory: Arc<D>,
478 schema: Arc<crate::dsl::Schema>,
479 metadata: IndexMetadata,
480 merge_policy: Box<dyn MergePolicy>,
481 term_cache_blocks: usize,
482 max_concurrent_merges: usize,
483 global_merge_permits: Arc<Semaphore>,
484 merge_bp_time_budget: Option<std::time::Duration>,
485 bp_memory_budget_bytes: usize,
486 reorder_permits: Arc<Semaphore>,
487 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
488 ) -> Self {
489 let reorder_on_merge = schema.reorder_on_merge();
492 if reorder_on_merge {
493 log::info!("[merge] reorder-on-merge enabled by index schema");
494 }
495
496 let tracker = Arc::new(SegmentTracker::new());
497 for seg_id in metadata.segment_metas.keys() {
498 tracker.register(seg_id);
499 }
500
501 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
502 Arc::new(parking_lot::Mutex::new(Vec::new()));
503 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
504 let dir = Arc::clone(&directory);
505 let tracker = Arc::clone(&tracker);
506 let lifecycle_handles = Arc::clone(&lifecycle_handles);
507 Arc::new(move |segment_ids| {
508 let Ok(handle) = tokio::runtime::Handle::try_current() else {
511 tracker.complete_deletion(&segment_ids);
514 return;
515 };
516 let dir = Arc::clone(&dir);
517 let task_tracker = Arc::clone(&tracker);
518 let cleanup_ids = segment_ids.clone();
519 let future = async move {
520 for &segment_id in &segment_ids {
521 log::info!(
522 "[segment_cleanup] deleting deferred segment {}",
523 segment_id.to_hex()
524 );
525 if let Err(error) =
526 crate::segment::delete_segment(dir.as_ref(), segment_id).await
527 {
528 log::warn!(
529 "[segment_cleanup] deferred delete failed for {}: {}",
530 segment_id.to_hex(),
531 error,
532 );
533 }
534 }
535 task_tracker.complete_deletion(&segment_ids);
536 };
537 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
538 tracker.complete_deletion(&cleanup_ids);
542 log::warn!(
543 "[segment_cleanup] runtime rejected deferred deletion; files will be swept later"
544 );
545 }
546 })
547 };
548
549 Self {
550 state: Arc::new(AsyncMutex::new(ManagerState {
551 metadata,
552 merge_policy,
553 })),
554 active_operations: Arc::new(ActiveSegmentOperations::new()),
555 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
556 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
557 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
558 merge_handles: parking_lot::Mutex::new(Vec::new()),
559 global_merge_wakeup_pending: AtomicBool::new(false),
560 lifecycle_handles,
561 trained: Arc::new(ArcSwapOption::new(None)),
562 vector_artifact_update: Arc::new(AtomicBool::new(false)),
563 tracker,
564 delete_fn,
565 directory,
566 schema,
567 term_cache_blocks,
568 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
569 global_merge_permits,
570 reorder_permits,
571 reorder_on_merge,
572 merge_bp_time_budget,
573 bp_memory_budget_bytes,
574 background_reorder_pool,
575 }
576 }
577
578 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
582 if let Some(pool) = &self.background_reorder_pool {
583 return Arc::clone(pool);
584 }
585 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
586 let threads = (num_cpus::get() / 2).max(1);
587 log::info!(
588 "[merge] process-wide background CPU pool: {} thread(s)",
589 threads
590 );
591 Arc::new(
592 rayon::ThreadPoolBuilder::new()
593 .num_threads(threads)
594 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
595 .build()
596 .expect("failed to build background CPU pool"),
597 )
598 }))
599 }
600
601 pub fn begin_shutdown(&self) {
605 self.active_operations.stop_accepting();
606 }
607
608 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
616 where
617 T: Send + 'static,
618 F: std::future::Future<Output = Result<T>> + Send + 'static,
619 {
620 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
621 let future = async move {
622 let result = transaction.await;
623 let _ = result_tx.send(result);
624 };
625 let runtime = tokio::runtime::Handle::current();
626 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
627 return Err(Error::Internal(
628 "runtime rejected lifecycle metadata transaction".into(),
629 ));
630 }
631 result_rx.await.map_err(|_| {
632 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
633 })?
634 }
635
636 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
638 let manager = Arc::clone(self);
639 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
640 let Ok(handle) = tokio::runtime::Handle::try_current() else {
641 log::warn!(
642 "[segment_cleanup] runtime unavailable; partial output {} will be swept on startup",
643 segment_id.to_hex(),
644 );
645 return;
646 };
647
648 let cleanup_manager = Arc::clone(&manager);
649 let future = async move {
650 cleanup_manager
651 .delete_output_if_unregistered(segment_id, "task unwind")
652 .await;
653 };
654 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
655 log::warn!(
656 "[segment_cleanup] runtime rejected output cleanup; {} will be swept on startup",
657 segment_id.to_hex(),
658 );
659 }
660 });
661
662 OutputCleanupGuard::new(output_id, cleanup)
663 }
664
665 pub(crate) fn schedule_unpublished_segment_cleanup(
670 self: &Arc<Self>,
671 output_id: SegmentId,
672 operation: SegmentOperationGuard,
673 runtime: tokio::runtime::Handle,
674 ) {
675 let manager = Arc::clone(self);
676 let output_hex = output_id.to_hex();
677 let future = async move {
678 manager
679 .delete_output_if_unregistered(output_id, "indexing abort or failure")
680 .await;
681 drop(operation);
682 };
683 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
684 log::warn!(
687 "[segment_cleanup] runtime unavailable; indexing output {} will be swept on startup",
688 output_hex,
689 );
690 }
691 }
692
693 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
699 match self
700 .active_operations
701 .try_register(vec![segment_id.clone()])
702 {
703 Some(operation) => Ok(operation),
704 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
705 None => Err(Error::Corruption(format!(
706 "new segment ID {} is already owned by an active operation",
707 segment_id
708 ))),
709 }
710 }
711
712 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
716 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
717 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
718 })?;
719 let files = SegmentFiles::new(id.0);
720
721 for path in files.mandatory_paths() {
722 if !self.directory.exists(path).await.map_err(Error::Io)? {
723 return Err(Error::Corruption(format!(
724 "segment {} cannot be published: mandatory file {:?} is missing",
725 segment_id, path
726 )));
727 }
728 }
729
730 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
731 Error::Corruption(format!(
732 "segment {} cannot be published: missing/unreadable {:?}: {}",
733 segment_id, files.meta, e
734 ))
735 })?;
736 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
737 Error::Corruption(format!(
738 "segment {} cannot be published: failed reading {:?}: {}",
739 segment_id, files.meta, e
740 ))
741 })?;
742 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
743 Error::Corruption(format!(
744 "segment {} cannot be published: invalid {:?}: {}",
745 segment_id, files.meta, e
746 ))
747 })?;
748
749 if meta.id != id.0 || meta.num_docs != expected_docs {
750 return Err(Error::Corruption(format!(
751 "segment {} cannot be published: metadata identity/docs mismatch \
752 (id={:032x}, docs={}, expected_docs={})",
753 segment_id, meta.id, meta.num_docs, expected_docs
754 )));
755 }
756
757 Ok(())
758 }
759
760 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
761 let inserted = self
762 .quarantined_segments
763 .lock()
764 .insert(segment_id.to_string());
765 if inserted {
766 log::error!(
767 "[merge] quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
768 It remains metadata-live for explicit repair but is excluded from merges until restart",
769 segment_id,
770 error,
771 );
772 }
773 }
774
775 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
776 let mut retry = self.merge_retry.lock();
777 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
778 let delay = merge_retry_delay(retry.consecutive_failures);
779 retry.retry_after = std::time::Instant::now().checked_add(delay);
780 log::warn!(
781 "[merge] pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
782 delay.as_secs_f64(),
783 retry.consecutive_failures,
784 error,
785 );
786 delay
787 }
788
789 fn clear_merge_retry_backoff(&self) {
790 *self.merge_retry.lock() = MergeRetryState::default();
791 }
792
793 fn merge_retry_is_paused(&self) -> bool {
794 let mut retry = self.merge_retry.lock();
795 match retry.retry_after {
796 Some(deadline) if deadline > std::time::Instant::now() => true,
797 Some(_) => {
798 retry.retry_after = None;
799 false
800 }
801 None => false,
802 }
803 }
804
805 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
806 let mut retries = self.reorder_retries.lock();
807 let retry = retries.entry(segment_id.to_string()).or_default();
808 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
809 let delay = merge_retry_delay(retry.consecutive_failures);
810 retry.retry_after = std::time::Instant::now().checked_add(delay);
811 log::warn!(
812 "[reorder] pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
813 segment_id,
814 delay.as_secs_f64(),
815 retry.consecutive_failures,
816 error,
817 );
818 }
819
820 fn clear_reorder_retry(&self, segment_id: &str) {
821 self.reorder_retries.lock().remove(segment_id);
822 }
823
824 fn paused_reorder_segments(&self) -> HashSet<String> {
825 let now = std::time::Instant::now();
826 let mut retries = self.reorder_retries.lock();
827 let mut paused = HashSet::new();
828 for (segment_id, retry) in retries.iter_mut() {
829 match retry.retry_after {
830 Some(deadline) if deadline > now => {
831 paused.insert(segment_id.clone());
832 }
833 Some(_) => retry.retry_after = None,
834 None => {}
835 }
836 }
837 paused
838 }
839
840 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
844 if self
845 .global_merge_wakeup_pending
846 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
847 .is_err()
848 {
849 return;
850 }
851
852 let manager = Arc::clone(self);
853 let future = async move {
854 let capacity = tokio::select! {
855 biased;
856 () = manager.active_operations.wait_for_shutdown() => None,
857 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
858 };
859
860 manager
861 .global_merge_wakeup_pending
862 .store(false, Ordering::Release);
863 if let Some(permit) = capacity {
864 drop(permit);
868 manager.maybe_merge().await;
869 }
870 };
871 let runtime = tokio::runtime::Handle::current();
872 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
873 self.global_merge_wakeup_pending
874 .store(false, Ordering::Release);
875 log::warn!("[merge] runtime rejected global-capacity wakeup task");
876 }
877 }
878
879 #[cfg(test)]
880 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
881 self.quarantined_segments.lock().contains(segment_id)
882 }
883
884 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
889 let output_hex = output_id.to_hex();
890 {
891 let st = self.state.lock().await;
892 if st.metadata.has_segment(&output_hex) {
893 return;
894 }
895 }
896
897 log::info!(
901 "[segment_cleanup] deleting uncommitted output {} after {}",
902 output_hex,
903 reason,
904 );
905 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
906 {
907 log::warn!(
908 "[segment_cleanup] failed deleting uncommitted output {}: {}",
909 output_hex,
910 error,
911 );
912 }
913 }
914
915 pub async fn get_segment_ids(&self) -> Vec<String> {
921 self.state.lock().await.metadata.segment_ids()
922 }
923
924 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
926 self.trained.load_full()
927 }
928
929 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
936 if self.vector_artifact_update.load(Ordering::Acquire) {
937 return None;
938 }
939 let trained = self.trained.load_full();
940 if self.vector_artifact_update.load(Ordering::Acquire) {
941 None
942 } else {
943 trained
944 }
945 }
946
947 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
955 self.vector_artifact_update
956 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
957 .map_err(|_| {
958 Error::Internal("a trained-vector artifact update is already in progress".into())
959 })?;
960 let guard = VectorArtifactUpdateGuard {
961 _lease: Arc::new(VectorArtifactUpdateLease {
962 updating: Arc::clone(&self.vector_artifact_update),
963 }),
964 };
965 let preexisting = self.active_operations.operation_tokens_snapshot();
966 self.active_operations
967 .wait_until_operations_finish(&preexisting)
968 .await;
969 Ok(guard)
970 }
971
972 pub async fn load_and_publish_trained(&self) {
975 if let Err(error) = self.try_load_and_publish_trained().await {
976 self.trained.store(None);
977 log::error!("[trained] refusing to publish trained artifacts: {error}");
978 }
979 }
980
981 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
984 let vector_fields = {
986 let st = self.state.lock().await;
987 st.metadata.vector_fields.clone()
988 };
989 let trained = IndexMetadata::try_load_trained_from_fields(
991 &vector_fields,
992 self.schema.as_ref(),
993 self.directory.as_ref(),
994 )
995 .await?
996 .map(Arc::new);
997 self.trained.store(trained);
1001 Ok(())
1002 }
1003
1004 pub(crate) async fn update_vector_metadata_and_publish<F>(
1011 self: &Arc<Self>,
1012 artifact_update: &VectorArtifactUpdateGuard,
1013 update: F,
1014 ) -> Result<()>
1015 where
1016 F: FnOnce(&mut IndexMetadata),
1017 {
1018 let mut st = Arc::clone(&self.state).lock_owned().await;
1019 let mut next = st.metadata.clone();
1020 update(&mut next);
1021
1022 let next_trained = IndexMetadata::try_load_trained_from_fields(
1023 &next.vector_fields,
1024 self.schema.as_ref(),
1025 self.directory.as_ref(),
1026 )
1027 .await?
1028 .map(Arc::new);
1029
1030 let directory = Arc::clone(&self.directory);
1031 let trained = Arc::clone(&self.trained);
1032 let artifact_update = artifact_update.clone();
1036 self.run_lifecycle_transaction(async move {
1037 let _artifact_update = artifact_update;
1038 next.save(directory.as_ref()).await?;
1039 st.metadata = next;
1040 trained.store(next_trained);
1041 Ok(())
1042 })
1043 .await
1044 }
1045
1046 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1048 where
1049 F: FnOnce(&IndexMetadata) -> R,
1050 {
1051 let st = self.state.lock().await;
1052 f(&st.metadata)
1053 }
1054
1055 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1057 where
1058 F: FnOnce(&mut IndexMetadata),
1059 {
1060 let mut st = Arc::clone(&self.state).lock_owned().await;
1061 let mut next = st.metadata.clone();
1062 f(&mut next);
1063 let directory = Arc::clone(&self.directory);
1064 self.run_lifecycle_transaction(async move {
1065 next.save(directory.as_ref()).await?;
1066 st.metadata = next;
1067 Ok(())
1068 })
1069 .await
1070 }
1071
1072 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1075 let acquired = {
1076 let st = self.state.lock().await;
1077 let segment_ids = st.metadata.segment_ids();
1078 self.tracker.acquire(&segment_ids)
1079 };
1080
1081 SegmentSnapshot::with_delete_fn(
1082 Arc::clone(&self.tracker),
1083 acquired,
1084 Arc::clone(&self.delete_fn),
1085 )
1086 }
1087
1088 pub fn tracker(&self) -> Arc<SegmentTracker> {
1090 Arc::clone(&self.tracker)
1091 }
1092
1093 pub fn directory(&self) -> Arc<D> {
1095 Arc::clone(&self.directory)
1096 }
1097}
1098
1099#[cfg(feature = "native")]
1104impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1105 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1107 for (segment_id, num_docs) in new_segments {
1110 self.validate_completed_segment(segment_id, *num_docs)
1111 .await?;
1112 }
1113
1114 let mut st = Arc::clone(&self.state).lock_owned().await;
1115 let mut next = st.metadata.clone();
1116 let mut added = Vec::new();
1117 for (segment_id, num_docs) in new_segments {
1118 if !next.has_segment(segment_id) {
1119 next.add_segment(segment_id.clone(), *num_docs);
1120 added.push(segment_id.clone());
1121 }
1122 }
1123
1124 let directory = Arc::clone(&self.directory);
1130 let tracker = Arc::clone(&self.tracker);
1131 self.run_lifecycle_transaction(async move {
1132 next.save(directory.as_ref()).await?;
1133 for segment_id in &added {
1134 tracker.register(segment_id);
1135 }
1136 st.metadata = next;
1137 Ok(())
1138 })
1139 .await
1140 }
1141
1142 pub async fn maybe_merge(self: &Arc<Self>) {
1153 if !self.active_operations.is_accepting() {
1154 log::debug!("[maybe_merge] manager is shutting down, skipping");
1155 return;
1156 }
1157 if self.merge_retry_is_paused() {
1158 log::debug!("[maybe_merge] retry backoff active, skipping");
1159 return;
1160 }
1161
1162 {
1165 let mut handles = self.merge_handles.lock();
1166 handles.retain(|h| !h.is_finished());
1167 }
1168 let local_slots = self.merge_permits.available_permits();
1169 let global_slots = self.global_merge_permits.available_permits();
1170 let slots_available = local_slots.min(global_slots);
1171
1172 let new_handles = {
1176 let st = self.state.lock().await;
1177 let quarantined = self.quarantined_segments.lock().clone();
1178 let active_ids = self.active_operations.snapshot();
1179
1180 let segments: Vec<SegmentInfo> = st
1183 .metadata
1184 .segment_metas
1185 .iter()
1186 .filter(|(id, _)| {
1187 !self.tracker.is_pending_deletion(id)
1188 && !active_ids.contains(*id)
1189 && !quarantined.contains(*id)
1190 })
1191 .map(|(id, info)| SegmentInfo {
1192 id: id.clone(),
1193 num_docs: info.num_docs,
1194 })
1195 .collect();
1196
1197 log::debug!("[maybe_merge] {} eligible segments", segments.len());
1198
1199 let candidates = st.merge_policy.find_merges(&segments);
1200
1201 if candidates.is_empty() {
1202 return;
1203 }
1204
1205 if slots_available == 0 {
1209 if local_slots > 0 && global_slots == 0 {
1210 self.schedule_global_merge_wakeup();
1211 }
1212 log::debug!("[maybe_merge] at max concurrent merges, skipping");
1213 return;
1214 }
1215
1216 log::debug!(
1217 "[maybe_merge] {} merge candidates, {} slots available",
1218 candidates.len(),
1219 slots_available
1220 );
1221
1222 let mut handles = Vec::new();
1223 for c in candidates {
1224 if handles.len() >= slots_available {
1225 break;
1226 }
1227 if let Some(h) = self.spawn_merge(c.segment_ids) {
1228 handles.push(h);
1229 }
1230 }
1231 handles
1232 };
1234
1235 if !new_handles.is_empty() {
1236 self.merge_handles.lock().extend(new_handles);
1240 }
1241 }
1242
1243 fn spawn_merge(self: &Arc<Self>, segment_ids_to_merge: Vec<String>) -> Option<JoinHandle<()>> {
1252 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1253 Ok(permit) => permit,
1254 Err(_) => {
1255 log::debug!("[spawn_merge] skipped: global merge capacity is full");
1256 self.schedule_global_merge_wakeup();
1257 return None;
1258 }
1259 };
1260 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1261 Ok(permit) => permit,
1262 Err(_) => {
1263 log::debug!("[spawn_merge] skipped: no merge permit available");
1264 return None;
1265 }
1266 };
1267 let output_id = SegmentId::new();
1268 let output_hex = output_id.to_hex();
1269
1270 let mut all_ids = segment_ids_to_merge.clone();
1271 all_ids.push(output_hex);
1272
1273 let guard = match self.active_operations.try_register(all_ids) {
1274 Some(g) => g,
1275 None => {
1276 log::debug!("[spawn_merge] skipped: segments overlap with an active operation");
1277 return None;
1278 }
1279 };
1280
1281 let sm = Arc::clone(self);
1282 let ids = segment_ids_to_merge;
1283
1284 Some(tokio::spawn(async move {
1285 let mut output_cleanup = sm.output_cleanup_guard(output_id);
1286 let mut reevaluate = false;
1287 let mut retry_delay = None;
1288
1289 let trained_snap = sm.trained_for_segment_build();
1290 let granularity = sm.merge_granularity(&ids).await;
1291 let result = Self::do_merge(
1292 sm.directory.as_ref(),
1293 &sm.schema,
1294 &ids,
1295 output_id,
1296 sm.term_cache_blocks,
1297 trained_snap.as_deref(),
1298 sm.reorder_on_merge,
1299 granularity,
1300 sm.merge_bp_time_budget,
1301 sm.bp_memory_budget_bytes,
1302 Arc::clone(&sm.reorder_permits),
1303 Some(sm.background_cpu_pool()),
1304 )
1305 .await;
1306
1307 match result {
1308 Ok((new_id, doc_count, bp_converged)) => {
1309 match sm
1310 .replace_segments(
1311 &ids,
1312 new_id,
1313 doc_count,
1314 sm.reorder_on_merge,
1315 bp_converged,
1316 )
1317 .await
1318 {
1319 Ok(()) => {
1320 output_cleanup.disarm();
1321 sm.clear_merge_retry_backoff();
1322 reevaluate = true;
1323 }
1324 Err(e) => {
1325 sm.delete_output_if_unregistered(output_id, "replacement failure")
1326 .await;
1327 output_cleanup.disarm();
1328 retry_delay = Some(sm.pause_merge_retries(&e));
1329 log::error!("[merge] failed to publish merged segment: {}", e);
1330 }
1331 }
1332 }
1333 Err(MergeTaskError {
1334 error,
1335 unavailable_segments,
1336 }) => {
1337 log::error!(
1338 "[merge] background merge failed for segments {:?}: {}",
1339 ids,
1340 error
1341 );
1342 if !unavailable_segments.is_empty() {
1343 for segment_id in &unavailable_segments {
1344 sm.quarantine_segment(segment_id, &error);
1345 }
1346 reevaluate = true;
1350 } else {
1351 retry_delay = Some(sm.pause_merge_retries(&error));
1352 }
1353 sm.delete_output_if_unregistered(output_id, "merge failure")
1354 .await;
1355 output_cleanup.disarm();
1356 }
1357 }
1358 drop(guard);
1361 drop(merge_permit);
1363 drop(global_merge_permit);
1364
1365 if reevaluate {
1366 sm.maybe_merge().await;
1367 } else if let Some(retry_delay) = retry_delay {
1368 tokio::select! {
1372 () = tokio::time::sleep(retry_delay) => {
1373 sm.maybe_merge().await;
1374 }
1375 () = sm.active_operations.wait_for_shutdown() => {}
1376 }
1377 }
1378 }))
1379 }
1380
1381 async fn replace_segments(
1385 self: &Arc<Self>,
1386 old_ids: &[String],
1387 new_id: String,
1388 doc_count: u32,
1389 reordered: bool,
1390 bp_converged: bool,
1391 ) -> Result<()> {
1392 self.validate_completed_segment(&new_id, doc_count).await?;
1395 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
1396 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
1397 })?;
1398 let output_reader = SegmentReader::open(
1399 self.directory.as_ref(),
1400 output_id,
1401 Arc::clone(&self.schema),
1402 self.term_cache_blocks,
1403 )
1404 .await
1405 .map_err(|error| match error {
1406 Error::Io(_) | Error::IndexClosed => error,
1410 error => Error::Corruption(format!(
1411 "replacement segment {new_id} failed full reader validation: {error}"
1412 )),
1413 })?;
1414 if output_reader.num_docs() != doc_count {
1415 return Err(Error::Corruption(format!(
1416 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
1417 output_reader.num_docs(),
1418 )));
1419 }
1420 drop(output_reader);
1421
1422 let mut st = Arc::clone(&self.state).lock_owned().await;
1423 let missing: Vec<&String> = old_ids
1427 .iter()
1428 .filter(|id| !st.metadata.has_segment(id))
1429 .collect();
1430 if !missing.is_empty() {
1431 return Err(Error::Corruption(format!(
1432 "replace_segments: source segment(s) {:?} not in metadata — \
1433 refusing to add output {} (would duplicate documents)",
1434 missing, new_id
1435 )));
1436 }
1437
1438 let parent_generation = old_ids
1439 .iter()
1440 .filter_map(|id| st.metadata.segment_metas.get(id))
1441 .map(|info| info.generation)
1442 .max()
1443 .unwrap_or(0)
1444 .checked_add(1)
1445 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
1446 let parent_unconverged_passes = old_ids
1447 .iter()
1448 .filter_map(|id| st.metadata.segment_metas.get(id))
1449 .map(|info| info.bp_unconverged_passes)
1450 .max()
1451 .unwrap_or(0);
1452 let bp_unconverged_passes = if reordered && !bp_converged {
1453 parent_unconverged_passes.saturating_add(1)
1454 } else {
1455 0
1456 };
1457 let retired_ids = old_ids.to_vec();
1458 let mut next = st.metadata.clone();
1459 for id in old_ids {
1460 next.remove_segment(id);
1461 }
1462 next.add_segment_meta(
1463 new_id.clone(),
1464 SegmentMetaInfo {
1465 num_docs: doc_count,
1466 ancestors: retired_ids.clone(),
1467 generation: parent_generation,
1468 reordered,
1469 bp_converged,
1470 bp_unconverged_passes,
1471 },
1472 );
1473
1474 let directory = Arc::clone(&self.directory);
1475 let tracker = Arc::clone(&self.tracker);
1476 self.run_lifecycle_transaction(async move {
1477 next.save(directory.as_ref()).await?;
1480 tracker.register(&new_id);
1481 st.metadata = next;
1482
1483 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
1487 drop(st);
1488 for &segment_id in &ready_to_delete {
1489 if let Err(error) =
1490 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1491 {
1492 log::warn!(
1493 "[segment_cleanup] immediate delete failed for {}: {}",
1494 segment_id.to_hex(),
1495 error,
1496 );
1497 }
1498 }
1499 tracker.complete_deletion(&ready_to_delete);
1500 Ok(())
1501 })
1502 .await
1503 }
1504
1505 #[allow(clippy::too_many_arguments)]
1510 async fn do_merge(
1511 directory: &D,
1512 schema: &Arc<crate::dsl::Schema>,
1513 segment_ids_to_merge: &[String],
1514 output_segment_id: SegmentId,
1515 term_cache_blocks: usize,
1516 trained: Option<&TrainedVectorStructures>,
1517 reorder_bmp: bool,
1518 granularity: crate::segment::reorder::BpGranularity,
1519 merge_bp_time_budget: Option<std::time::Duration>,
1520 bp_memory_budget_bytes: usize,
1521 reorder_permits: Arc<Semaphore>,
1522 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
1523 ) -> MergeTaskResult<(String, u32, bool)> {
1524 let output_hex = output_segment_id.to_hex();
1525 let load_start = std::time::Instant::now();
1526
1527 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
1528 for id_str in segment_ids_to_merge {
1529 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
1530 MergeTaskError::source(
1531 id_str.clone(),
1532 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
1533 )
1534 })?;
1535 segment_ids.push(id);
1536 }
1537
1538 let mut unavailable_sources = Vec::new();
1543 let mut missing_files = Vec::new();
1544 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
1545 let files = SegmentFiles::new(id.0);
1546 let mut source_unavailable = false;
1547 for path in files.mandatory_paths() {
1548 let exists = directory
1549 .exists(path)
1550 .await
1551 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
1552 if !exists {
1553 source_unavailable = true;
1554 missing_files.push(format!("{}:{:?}", id_str, path));
1555 }
1556 }
1557 if source_unavailable {
1558 unavailable_sources.push(id_str.clone());
1559 }
1560 }
1561 if !unavailable_sources.is_empty() {
1562 return Err(MergeTaskError::sources(
1563 unavailable_sources,
1564 Error::Corruption(format!(
1565 "merge sources are missing mandatory files: {}",
1566 missing_files.join(", ")
1567 )),
1568 ));
1569 }
1570
1571 let schema_arc = Arc::clone(schema);
1572 let futures: Vec<_> = segment_ids
1573 .iter()
1574 .map(|&sid| {
1575 let sch = Arc::clone(&schema_arc);
1576 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
1577 })
1578 .collect();
1579
1580 let results = futures::future::join_all(futures).await;
1581 let mut readers = Vec::with_capacity(results.len());
1582 let mut total_docs = 0u64;
1583 for (i, result) in results.into_iter().enumerate() {
1584 match result {
1585 Ok(r) => {
1586 total_docs += r.meta().num_docs as u64;
1587 readers.push(r);
1588 }
1589 Err(e) => {
1590 log::error!(
1591 "[merge] Failed to open segment {}: {:?}",
1592 segment_ids_to_merge[i],
1593 e
1594 );
1595 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
1596 }
1597 }
1598 }
1599 if total_docs > u32::MAX as u64 {
1600 return Err(Error::Internal(format!(
1601 "Merged segment doc count ({}) exceeds u32::MAX",
1602 total_docs
1603 ))
1604 .into());
1605 }
1606
1607 for (i, reader) in readers.iter().enumerate() {
1611 let meta_docs = reader.meta().num_docs;
1612 let store_docs = reader.store().num_docs();
1613 if store_docs != meta_docs {
1614 return Err(MergeTaskError::source(
1615 segment_ids_to_merge[i].clone(),
1616 Error::Corruption(format!(
1617 "pre-merge validation: segment {} store has {} docs but meta says {}",
1618 segment_ids_to_merge[i], store_docs, meta_docs
1619 )),
1620 ));
1621 }
1622 }
1623
1624 log::info!(
1625 "[merge] loaded {} segment readers in {:.1}s",
1626 readers.len(),
1627 load_start.elapsed().as_secs_f64()
1628 );
1629
1630 let merger = SegmentMerger::new(Arc::clone(schema))
1631 .with_bmp_reorder(reorder_bmp)
1632 .with_granularity(granularity)
1633 .with_bp_budget(crate::segment::BpBudget {
1634 min_partition_docs: None,
1635 time_budget: merge_bp_time_budget,
1636 })
1637 .with_bp_memory_budget(bp_memory_budget_bytes)
1638 .with_reorder_permits(reorder_permits)
1639 .with_background_pool(bg_cpu_pool);
1640
1641 log::info!(
1642 "[merge] {} segments -> {} (trained={})",
1643 segment_ids_to_merge.len(),
1644 output_hex,
1645 trained.map_or(0, |t| t.centroids.len()),
1646 );
1647
1648 let (_merged_meta, merge_stats) = merger
1649 .merge(directory, &readers, output_segment_id, trained)
1650 .await
1651 .map_err(|error| {
1652 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
1653 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
1659 } else {
1660 MergeTaskError::from(error)
1661 }
1662 })?;
1663 let bp_converged = merge_stats.bp_converged;
1664 if !bp_converged {
1665 log::info!(
1666 "[merge] merge-time BP hit its wall-clock budget — output marked unconverged; \
1667 the background optimizer deepens it later",
1668 );
1669 }
1670
1671 log::info!(
1672 "[merge] total wall-clock: {:.1}s ({} segments, {} docs)",
1673 load_start.elapsed().as_secs_f64(),
1674 readers.len(),
1675 total_docs,
1676 );
1677
1678 Ok((output_hex, total_docs as u32, bp_converged))
1679 }
1680
1681 pub async fn abort_merges(&self) {
1689 loop {
1690 let handles: Vec<JoinHandle<()>> = { std::mem::take(&mut *self.merge_handles.lock()) };
1691 if handles.is_empty() {
1692 return;
1693 }
1694 for handle in handles {
1695 if let Err(error) = handle.await
1696 && error.is_panic()
1697 {
1698 log::error!("[merge] background task panicked while draining: {}", error);
1699 }
1700 }
1701 }
1702 }
1703
1704 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
1706 let handles: Vec<JoinHandle<()>> = { std::mem::take(&mut *self.merge_handles.lock()) };
1707 for h in handles {
1708 let _ = h.await;
1709 }
1710 }
1711
1712 pub async fn wait_for_all_merges(self: &Arc<Self>) {
1718 loop {
1719 let handles: Vec<JoinHandle<()>> = { std::mem::take(&mut *self.merge_handles.lock()) };
1720 if handles.is_empty() {
1721 break;
1722 }
1723 for h in handles {
1724 let _ = h.await;
1725 }
1726 }
1727 }
1728
1729 pub async fn wait_for_shutdown(self: &Arc<Self>) {
1734 self.wait_for_all_merges().await;
1735 self.active_operations.wait_until_idle().await;
1736 loop {
1737 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
1738 if handles.is_empty() {
1739 break;
1740 }
1741 for handle in handles {
1742 if let Err(error) = handle.await
1743 && error.is_panic()
1744 {
1745 log::error!("[segment_cleanup] task panicked while draining: {}", error);
1746 }
1747 }
1748 }
1749 }
1750
1751 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
1760 const FORCE_MERGE_BATCH: usize = 64;
1761
1762 let max_segment_docs = {
1763 let st = self.state.lock().await;
1764 st.merge_policy.max_segment_docs()
1765 };
1766
1767 self.wait_for_all_merges().await;
1770
1771 loop {
1772 if !self.active_operations.is_accepting() {
1773 return Err(Error::IndexClosed);
1774 }
1775 let mut segments: Vec<(String, u32)> = {
1777 let st = self.state.lock().await;
1778 st.metadata
1779 .segment_metas
1780 .iter()
1781 .map(|(id, info)| (id.clone(), info.num_docs))
1782 .collect()
1783 };
1784
1785 if segments.len() < 2 {
1786 return Ok(());
1787 }
1788
1789 segments.sort_by_key(|(_, docs)| *docs);
1790
1791 let max_docs = max_segment_docs.map(|m| m as u64).unwrap_or(u64::MAX);
1793 let mut batch = Vec::new();
1794 let mut batch_docs = 0u64;
1795
1796 for (id, docs) in &segments {
1797 if batch.len() >= FORCE_MERGE_BATCH {
1798 break;
1799 }
1800 let next_total = batch_docs + *docs as u64;
1801 if next_total > max_docs && !batch.is_empty() {
1802 break;
1803 }
1804 batch.push(id.clone());
1805 batch_docs += *docs as u64;
1806 }
1807
1808 if batch.len() < 2 {
1809 return Ok(());
1810 }
1811
1812 log::info!(
1813 "[force_merge] merging batch of {} segments ({} docs)",
1814 batch.len(),
1815 batch_docs
1816 );
1817
1818 let _global_merge_permit = tokio::select! {
1819 biased;
1820 () = self.active_operations.wait_for_shutdown() => {
1821 return Err(Error::IndexClosed);
1822 }
1823 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
1824 permit.map_err(|_| {
1825 Error::Internal("global background merge scheduler is closed".into())
1826 })?
1827 }
1828 };
1829
1830 let output_id = SegmentId::new();
1831 let output_hex = output_id.to_hex();
1832
1833 let mut all_ids = batch.clone();
1836 all_ids.push(output_hex);
1837 let guard = {
1838 let st = self.state.lock().await;
1839 batch
1840 .iter()
1841 .all(|id| st.metadata.has_segment(id))
1842 .then(|| self.active_operations.try_register(all_ids))
1843 .flatten()
1844 };
1845 let _guard = match guard {
1846 Some(g) => g,
1847 None if !self.active_operations.is_accepting() => {
1848 return Err(Error::IndexClosed);
1849 }
1850 None => {
1851 self.wait_for_merging_thread().await;
1853 continue;
1854 }
1855 };
1856 let mut output_cleanup = self.output_cleanup_guard(output_id);
1857
1858 let trained_snap = self.trained_for_segment_build();
1859 let granularity = self.merge_granularity(&batch).await;
1860 let merge_result = Self::do_merge(
1861 self.directory.as_ref(),
1862 &self.schema,
1863 &batch,
1864 output_id,
1865 self.term_cache_blocks,
1866 trained_snap.as_deref(),
1867 self.reorder_on_merge,
1868 granularity,
1869 self.merge_bp_time_budget,
1870 self.bp_memory_budget_bytes,
1871 Arc::clone(&self.reorder_permits),
1872 Some(self.background_cpu_pool()),
1873 )
1874 .await;
1875 let (new_segment_id, total_docs, bp_converged) = match merge_result {
1876 Ok(v) => v,
1877 Err(MergeTaskError {
1878 error,
1879 unavailable_segments,
1880 }) => {
1881 for segment_id in &unavailable_segments {
1882 self.quarantine_segment(segment_id, &error);
1883 }
1884 self.delete_output_if_unregistered(output_id, "force-merge failure")
1885 .await;
1886 output_cleanup.disarm();
1887 return Err(error);
1888 }
1889 };
1890
1891 if let Err(e) = self
1892 .replace_segments(
1893 &batch,
1894 new_segment_id,
1895 total_docs,
1896 self.reorder_on_merge,
1897 bp_converged,
1898 )
1899 .await
1900 {
1901 self.delete_output_if_unregistered(output_id, "replacement failure")
1902 .await;
1903 output_cleanup.disarm();
1904 return Err(e);
1905 }
1906 output_cleanup.disarm();
1907
1908 }
1910 }
1911
1912 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
1919 self.wait_for_all_merges().await;
1920 let segment_ids = self.get_segment_ids().await;
1921
1922 if segment_ids.is_empty() {
1923 log::info!("[reorder] no segments to reorder");
1924 return Ok(());
1925 }
1926
1927 log::info!("[reorder] reordering {} segments", segment_ids.len());
1928
1929 for seg_id in segment_ids {
1930 match self
1931 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
1932 .await
1933 {
1934 Ok(true) => {}
1935 Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
1936 Err(e) => return Err(e),
1937 }
1938 }
1939
1940 log::info!("[reorder] all segments reordered");
1941 Ok(())
1942 }
1943
1944 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
1949 self.unreordered_segments()
1950 .await
1951 .into_iter()
1952 .map(|(id, _)| id)
1953 .collect()
1954 }
1955
1956 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
1959 let quarantined = self.quarantined_segments.lock().clone();
1960 let paused = self.paused_reorder_segments();
1961 let st = self.state.lock().await;
1962 let active_ids = self.active_operations.snapshot();
1963 st.metadata
1964 .segment_metas
1965 .iter()
1966 .filter(|(id, info)| {
1967 !info.reordered
1968 && !active_ids.contains(*id)
1969 && !quarantined.contains(*id)
1970 && !paused.contains(*id)
1971 })
1972 .map(|(id, info)| (id.clone(), info.num_docs))
1973 .collect()
1974 }
1975
1976 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
1980 self.unconverged_segments_below(u32::MAX)
1981 .await
1982 .into_iter()
1983 .map(|(id, docs, _)| (id, docs))
1984 .collect()
1985 }
1986
1987 pub async fn unconverged_segments_below(
1990 &self,
1991 max_unconverged_passes: u32,
1992 ) -> Vec<(String, u32, u32)> {
1993 let quarantined = self.quarantined_segments.lock().clone();
1994 let paused = self.paused_reorder_segments();
1995 let st = self.state.lock().await;
1996 let active_ids = self.active_operations.snapshot();
1997 st.metadata
1998 .segment_metas
1999 .iter()
2000 .filter(|(id, info)| {
2001 info.reordered
2002 && !info.bp_converged
2003 && info.bp_unconverged_passes < max_unconverged_passes
2004 && !active_ids.contains(*id)
2005 && !quarantined.contains(*id)
2006 && !paused.contains(*id)
2007 })
2008 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
2009 .collect()
2010 }
2011
2012 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
2022 let st = self.state.lock().await;
2023 let deepening = ids.iter().any(|id| {
2024 st.metadata
2025 .segment_metas
2026 .get(id)
2027 .is_some_and(|info| info.reordered && !info.bp_converged)
2028 });
2029 drop(st);
2030 if deepening {
2031 log::info!(
2032 "[reorder] source segment(s) unconverged — forcing record-level BP (deepening pass)",
2033 );
2034 crate::segment::reorder::BpGranularity::Records
2035 } else {
2036 crate::segment::reorder::BpGranularity::Auto
2037 }
2038 }
2039
2040 pub async fn reorder_single_segment(
2045 self: &Arc<Self>,
2046 seg_id: &str,
2047 rayon_pool: Option<Arc<rayon::ThreadPool>>,
2048 bp_budget: crate::segment::BpBudget,
2049 ) -> Result<bool> {
2050 let source_id = SegmentId::from_hex(seg_id)
2051 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
2052 if self.quarantined_segments.lock().contains(seg_id) {
2053 return Err(Error::Corruption(format!(
2054 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
2055 seg_id
2056 )));
2057 }
2058
2059 let _reorder_permit = tokio::select! {
2064 biased;
2065 () = self.active_operations.wait_for_shutdown() => {
2066 return Err(Error::IndexClosed);
2067 }
2068 permit = Arc::clone(&self.reorder_permits).acquire_owned() => {
2069 permit.map_err(|_| {
2070 Error::Internal("background reorder scheduler is closed".into())
2071 })?
2072 }
2073 };
2074
2075 let output_id = SegmentId::new();
2076 let output_hex = output_id.to_hex();
2077 let source_ids = [seg_id.to_string()];
2078 let granularity = self.merge_granularity(&source_ids).await;
2079
2080 let all_ids = vec![seg_id.to_string(), output_hex];
2086 let (_guard, source_docs) = {
2087 let st = self.state.lock().await;
2088 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
2089 log::info!(
2090 "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
2091 seg_id
2092 );
2093 self.clear_reorder_retry(seg_id);
2094 return Ok(false);
2095 };
2096
2097 match self.active_operations.try_register(all_ids) {
2098 Some(guard) => (guard, source_meta.num_docs),
2099 None if !self.active_operations.is_accepting() => {
2100 return Err(Error::IndexClosed);
2101 }
2102 None => {
2103 log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
2104 return Ok(false);
2105 }
2106 }
2107 };
2108
2109 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
2114 if is_deterministic_source_error(&error) {
2115 self.quarantine_segment(seg_id, &error);
2116 } else if !matches!(&error, Error::IndexClosed) {
2117 self.pause_reorder_retries(seg_id, &error);
2118 }
2119 return Err(error);
2120 }
2121
2122 let mut output_cleanup = self.output_cleanup_guard(output_id);
2123
2124 let reorder_result = crate::segment::reorder::reorder_segment(
2125 self.directory.as_ref(),
2126 &self.schema,
2127 source_id,
2128 output_id,
2129 self.term_cache_blocks,
2130 self.bp_memory_budget_bytes,
2131 bp_budget,
2132 granularity,
2133 rayon_pool,
2134 )
2135 .await;
2136 let (new_id, total_docs, bp_converged) = match reorder_result {
2137 Ok(v) => v,
2138 Err(e) => {
2139 self.delete_output_if_unregistered(output_id, "reorder failure")
2142 .await;
2143 output_cleanup.disarm();
2144 if is_deterministic_source_error(&e) {
2145 self.quarantine_segment(seg_id, &e);
2146 } else if !matches!(&e, Error::IndexClosed) {
2147 self.pause_reorder_retries(seg_id, &e);
2148 }
2149 return Err(e);
2150 }
2151 };
2152
2153 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
2159 if let Err(e) = self
2160 .replace_segments(
2161 &[seg_id.to_string()],
2162 new_id,
2163 total_docs,
2164 true,
2165 ladder_converged,
2166 )
2167 .await
2168 {
2169 self.delete_output_if_unregistered(output_id, "replacement failure")
2170 .await;
2171 output_cleanup.disarm();
2172 if !matches!(&e, Error::IndexClosed) {
2173 self.pause_reorder_retries(seg_id, &e);
2174 }
2175 return Err(e);
2176 }
2177 output_cleanup.disarm();
2178 self.clear_reorder_retry(seg_id);
2179
2180 Ok(true)
2181 }
2182
2183 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
2190 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
2191
2192 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
2193 for entry in entries {
2194 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
2195 continue;
2196 };
2197 let Some(rest) = filename.strip_prefix("seg_") else {
2198 continue;
2199 };
2200 let Some(hex_id) = rest.get(..32) else {
2201 continue;
2202 };
2203 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
2204 continue;
2205 }
2206 orphan_files
2207 .entry(hex_id.to_ascii_lowercase())
2208 .or_default()
2209 .push(entry);
2210 }
2211 }
2212
2213 let mut deleted = 0;
2214 for (hex_id, paths) in &orphan_files {
2215 let deletion_guard = {
2220 let st = self.state.lock().await;
2221 if st.metadata.has_segment(hex_id) {
2222 continue;
2223 }
2224 let Some(guard) = self
2225 .active_operations
2226 .try_register(vec![hex_id.to_string()])
2227 else {
2228 continue;
2229 };
2230 if self.tracker.is_deletion_protected(hex_id) {
2231 drop(guard);
2232 continue;
2233 }
2234 guard
2235 };
2236
2237 let results =
2242 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
2243 .await;
2244 let removed = results.into_iter().all(|result| match result {
2245 Ok(()) => true,
2246 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
2247 Err(error) => {
2248 log::warn!(
2249 "[segment_cleanup] failed sweeping orphan segment {}: {}",
2250 hex_id,
2251 error,
2252 );
2253 false
2254 }
2255 });
2256 drop(deletion_guard);
2259 if removed {
2260 deleted += 1;
2261 log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
2262 }
2263 }
2264
2265 Ok(deleted)
2266 }
2267}
2268
2269#[cfg(test)]
2270mod tests {
2271 use super::*;
2272 use std::sync::atomic::{AtomicBool, Ordering};
2273
2274 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
2275 let schema = crate::dsl::SchemaBuilder::default().build();
2276 let metadata = IndexMetadata::new(schema.clone());
2277 Arc::new(SegmentManager::new(
2278 Arc::new(crate::directories::RamDirectory::new()),
2279 Arc::new(schema),
2280 metadata,
2281 Box::new(crate::merge::NoMergePolicy),
2282 0,
2283 1,
2284 Arc::new(Semaphore::new(1)),
2285 None,
2286 1024,
2287 Arc::new(Semaphore::new(1)),
2288 None,
2289 ))
2290 }
2291
2292 #[test]
2293 fn output_cleanup_guard_runs_during_panic_unwind() {
2294 let cleaned = Arc::new(AtomicBool::new(false));
2295 let cleaned_in_callback = Arc::clone(&cleaned);
2296 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
2297 cleaned_in_callback.store(true, Ordering::SeqCst);
2298 });
2299
2300 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
2301 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
2302 panic!("simulated reorder panic");
2303 }));
2304
2305 assert!(result.is_err());
2306 assert!(
2307 cleaned.load(Ordering::SeqCst),
2308 "partial output cleanup must run during unwind"
2309 );
2310 }
2311
2312 #[test]
2313 fn output_cleanup_guard_disarms_after_commit() {
2314 let cleaned = Arc::new(AtomicBool::new(false));
2315 let cleaned_in_callback = Arc::clone(&cleaned);
2316 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
2317 cleaned_in_callback.store(true, Ordering::SeqCst);
2318 });
2319
2320 {
2321 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
2322 guard.disarm();
2323 }
2324
2325 assert!(!cleaned.load(Ordering::SeqCst));
2326 }
2327
2328 #[test]
2329 fn test_active_operation_guard_releases_ownership() {
2330 let active = Arc::new(ActiveSegmentOperations::new());
2331 {
2332 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
2333 let snap = active.snapshot();
2334 assert!(snap.contains("a"));
2335 assert!(snap.contains("b"));
2336 }
2337 assert!(active.snapshot().is_empty());
2338 }
2339
2340 #[test]
2341 fn test_non_overlapping_operations_can_run_concurrently() {
2342 let active = Arc::new(ActiveSegmentOperations::new());
2343 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
2344 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
2345 let snap = active.snapshot();
2346 assert_eq!(snap.len(), 4);
2347
2348 drop(first);
2349 let snap = active.snapshot();
2350 assert_eq!(snap.len(), 2);
2351 assert!(snap.contains("c"));
2352 assert!(snap.contains("d"));
2353 }
2354
2355 #[test]
2356 fn test_overlapping_operation_is_rejected_until_release() {
2357 let active = Arc::new(ActiveSegmentOperations::new());
2358 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
2359 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
2360 drop(first);
2361 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
2362 }
2363
2364 #[test]
2365 fn test_active_operation_snapshot() {
2366 let active = Arc::new(ActiveSegmentOperations::new());
2367 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
2368 let snap = active.snapshot();
2369 assert!(snap.contains("x"));
2370 assert!(snap.contains("y"));
2371 assert!(!snap.contains("z"));
2372 }
2373
2374 #[tokio::test]
2375 async fn operation_barrier_ignores_producers_started_after_snapshot() {
2376 let active = Arc::new(ActiveSegmentOperations::new());
2377 let before_gate = active.try_register(vec!["old".into()]).unwrap();
2378 let barrier = active.operation_tokens_snapshot();
2379 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
2380
2381 let waiter = {
2382 let active = Arc::clone(&active);
2383 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
2384 };
2385 tokio::task::yield_now().await;
2386 assert!(!waiter.is_finished());
2387
2388 drop(before_gate);
2389 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
2390 .await
2391 .expect("pre-gate operation barrier was starved by a post-gate producer")
2392 .unwrap();
2393 assert!(active.snapshot().contains("new-flat"));
2394 drop(after_gate);
2395 }
2396
2397 #[tokio::test]
2398 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
2399 let manager = lifecycle_test_manager();
2400 manager
2401 .trained
2402 .store(Some(Arc::new(TrainedVectorStructures {
2403 centroids: rustc_hash::FxHashMap::default(),
2404 codebooks: rustc_hash::FxHashMap::default(),
2405 })));
2406
2407 let guard = manager.begin_vector_artifact_update().await.unwrap();
2408 assert!(
2409 manager.trained().is_some(),
2410 "search readers keep the last fully validated generation"
2411 );
2412 assert!(
2413 manager.trained_for_segment_build().is_none(),
2414 "new segment producers must stay flat during an artifact update"
2415 );
2416
2417 let detached_transaction_guard = guard.clone();
2418 drop(guard);
2419 assert!(
2420 manager.trained_for_segment_build().is_none(),
2421 "a detached lifecycle transaction must retain the producer gate after request cancellation"
2422 );
2423 drop(detached_transaction_guard);
2424 assert!(manager.trained_for_segment_build().is_some());
2425 }
2426
2427 #[tokio::test]
2428 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
2429 let active = Arc::new(ActiveSegmentOperations::new());
2430 let guard = active.try_register(vec!["live".into()]).unwrap();
2431 active.stop_accepting();
2432 assert!(active.try_register(vec!["new".into()]).is_none());
2433
2434 let waiter = {
2435 let active = Arc::clone(&active);
2436 tokio::spawn(async move { active.wait_until_idle().await })
2437 };
2438 tokio::task::yield_now().await;
2439 assert!(!waiter.is_finished());
2440 drop(guard);
2441 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
2442 .await
2443 .expect("shutdown waiter missed the final guard notification")
2444 .unwrap();
2445 }
2446
2447 #[tokio::test]
2448 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
2449 let manager = lifecycle_test_manager();
2450 let started = Arc::new(Semaphore::new(0));
2451 let release = Arc::new(Semaphore::new(0));
2452 let completed = Arc::new(AtomicBool::new(false));
2453
2454 let request = {
2455 let manager = Arc::clone(&manager);
2456 let started = Arc::clone(&started);
2457 let release = Arc::clone(&release);
2458 let completed = Arc::clone(&completed);
2459 tokio::spawn(async move {
2460 manager
2461 .run_lifecycle_transaction(async move {
2462 started.add_permits(1);
2463 let _permit = release.acquire().await.unwrap();
2464 completed.store(true, Ordering::Release);
2465 Ok(())
2466 })
2467 .await
2468 })
2469 };
2470
2471 let _started = started.acquire().await.unwrap();
2472 request.abort();
2473 assert!(request.await.unwrap_err().is_cancelled());
2474 release.add_permits(1);
2475
2476 manager.begin_shutdown();
2477 tokio::time::timeout(
2478 std::time::Duration::from_secs(1),
2479 manager.wait_for_shutdown(),
2480 )
2481 .await
2482 .expect("shutdown did not drain detached lifecycle transaction");
2483 assert!(completed.load(Ordering::Acquire));
2484 }
2485
2486 #[tokio::test]
2487 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
2488 let manager = lifecycle_test_manager();
2489 {
2490 let mut state = manager.state.lock().await;
2491 state.metadata.add_segment_meta(
2492 "eligible".into(),
2493 SegmentMetaInfo {
2494 num_docs: 10,
2495 ancestors: Vec::new(),
2496 generation: 1,
2497 reordered: true,
2498 bp_converged: false,
2499 bp_unconverged_passes: 2,
2500 },
2501 );
2502 state.metadata.add_segment_meta(
2503 "at-limit".into(),
2504 SegmentMetaInfo {
2505 num_docs: 20,
2506 ancestors: Vec::new(),
2507 generation: 1,
2508 reordered: true,
2509 bp_converged: false,
2510 bp_unconverged_passes: 3,
2511 },
2512 );
2513 state.metadata.add_segment_meta(
2514 "converged".into(),
2515 SegmentMetaInfo {
2516 num_docs: 30,
2517 ancestors: Vec::new(),
2518 generation: 1,
2519 reordered: true,
2520 bp_converged: true,
2521 bp_unconverged_passes: 0,
2522 },
2523 );
2524 state.metadata.add_segment("fresh".into(), 40);
2525 }
2526
2527 assert_eq!(
2528 manager.unconverged_segments_below(3).await,
2529 vec![("eligible".into(), 10, 2)]
2530 );
2531 assert!(manager.unconverged_segments_below(0).await.is_empty());
2532 }
2533
2534 #[test]
2535 fn merge_retry_backoff_is_exponential_and_capped() {
2536 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
2537 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
2538 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
2539 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
2540 }
2541
2542 #[test]
2543 fn only_deterministic_source_errors_are_quarantined() {
2544 assert!(is_deterministic_source_error(&Error::Corruption(
2545 "bad footer".into()
2546 )));
2547 assert!(is_deterministic_source_error(&Error::Io(
2548 std::io::Error::from(std::io::ErrorKind::NotFound)
2549 )));
2550 assert!(!is_deterministic_source_error(&Error::Io(
2551 std::io::Error::from(std::io::ErrorKind::TimedOut)
2552 )));
2553 assert!(!is_deterministic_source_error(&Error::Io(
2554 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
2555 )));
2556 }
2557
2558 #[test]
2559 fn transient_reorder_failure_is_backed_off_until_cleared() {
2560 let manager = lifecycle_test_manager();
2561 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
2562 assert!(manager.paused_reorder_segments().contains("source"));
2563 manager.clear_reorder_retry("source");
2564 assert!(!manager.paused_reorder_segments().contains("source"));
2565 }
2566}