1use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwapOption;
49use tokio::sync::Mutex as AsyncMutex;
50use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore};
51use tokio::task::JoinHandle;
52
53use crate::directories::DirectoryWriter;
54use crate::error::{Error, Result};
55use crate::index::{IndexMetadata, ReorderConcurrencyGate, ReorderPriority, 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
64const FORCE_MERGE_MAX_FAN_IN: usize = 64;
65
66#[derive(Debug)]
67struct ForceMergeGroup {
68 segments: Vec<(String, u32)>,
69 total_docs: u64,
70}
71
72fn plan_force_merge_groups(
81 mut segments: Vec<(String, u32)>,
82 max_docs: u64,
83) -> Vec<ForceMergeGroup> {
84 segments.sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
85 right_docs
86 .cmp(left_docs)
87 .then_with(|| left_id.cmp(right_id))
88 });
89
90 let mut groups: Vec<ForceMergeGroup> = Vec::new();
91 for segment in segments {
92 let docs = u64::from(segment.1);
93 let best_group = groups
94 .iter()
95 .enumerate()
96 .filter_map(|(index, group)| {
97 group
98 .total_docs
99 .checked_add(docs)
100 .filter(|&total| total <= max_docs)
101 .map(|_| (index, group.total_docs))
102 })
103 .max_by_key(|&(index, used)| (used, std::cmp::Reverse(index)))
104 .map(|(index, _)| index);
105
106 if let Some(index) = best_group {
107 groups[index].total_docs += docs;
108 groups[index].segments.push(segment);
109 } else {
110 groups.push(ForceMergeGroup {
111 segments: vec![segment],
112 total_docs: docs,
113 });
114 }
115 }
116
117 for group in &mut groups {
121 group
122 .segments
123 .sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
124 left_docs
125 .cmp(right_docs)
126 .then_with(|| left_id.cmp(right_id))
127 });
128 }
129
130 groups.sort_unstable_by(|left, right| {
133 left.total_docs
134 .cmp(&right.total_docs)
135 .then_with(|| left.segments[0].0.cmp(&right.segments[0].0))
136 });
137 groups
138}
139
140fn force_merge_output_count(source_count: usize) -> usize {
141 if source_count < 2 {
142 return 0;
143 }
144 (source_count - 1).div_ceil(FORCE_MERGE_MAX_FAN_IN - 1)
145}
146
147#[derive(Debug)]
148struct ForceMergeStep {
149 inputs: Vec<usize>,
151}
152
153#[derive(Debug)]
154struct ForceMergeHierarchy {
155 steps: Vec<ForceMergeStep>,
156 root: usize,
157}
158
159fn plan_force_merge_hierarchy(source_count: usize) -> ForceMergeHierarchy {
164 debug_assert!(source_count >= 2);
165 let internal_count = force_merge_output_count(source_count);
166 let max_leaves = 1usize
167 .checked_add(internal_count.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
168 .expect("force-merge hierarchy size exceeds usize");
169 let deficit = max_leaves - source_count;
170 debug_assert!(deficit < FORCE_MERGE_MAX_FAN_IN - 1);
171
172 fn build(
173 leaf_start: usize,
174 leaf_count: usize,
175 internal_count: usize,
176 deficit: usize,
177 source_count: usize,
178 steps: &mut Vec<ForceMergeStep>,
179 ) -> usize {
180 debug_assert!(internal_count > 0);
181 if internal_count == 1 {
182 let arity = FORCE_MERGE_MAX_FAN_IN - deficit;
183 debug_assert_eq!(leaf_count, arity);
184 debug_assert!((2..=FORCE_MERGE_MAX_FAN_IN).contains(&arity));
185 let output = source_count + steps.len();
186 steps.push(ForceMergeStep {
187 inputs: (leaf_start..leaf_start + arity).collect(),
188 });
189 return output;
190 }
191
192 let child_internal_total = internal_count - 1;
197 let base = child_internal_total / FORCE_MERGE_MAX_FAN_IN;
198 let extra = child_internal_total % FORCE_MERGE_MAX_FAN_IN;
199 let mut child_internal = vec![base; FORCE_MERGE_MAX_FAN_IN];
200 for count in &mut child_internal[..extra] {
201 *count += 1;
202 }
203 let partial_child = (deficit > 0).then(|| {
204 child_internal
205 .iter()
206 .position(|&count| count > 0)
207 .expect("a non-root partial node requires an internal child")
208 });
209
210 let mut cursor = leaf_start;
211 let mut inputs = Vec::with_capacity(FORCE_MERGE_MAX_FAN_IN);
212 for (child, &child_internals) in child_internal.iter().enumerate() {
213 if child_internals == 0 {
214 inputs.push(cursor);
215 cursor += 1;
216 continue;
217 }
218 let child_deficit = usize::from(partial_child == Some(child)) * deficit;
219 let child_leaves = 1usize
220 .checked_add(child_internals.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
221 .and_then(|maximum| maximum.checked_sub(child_deficit))
222 .expect("force-merge child size exceeds usize");
223 inputs.push(build(
224 cursor,
225 child_leaves,
226 child_internals,
227 child_deficit,
228 source_count,
229 steps,
230 ));
231 cursor += child_leaves;
232 }
233 debug_assert_eq!(cursor, leaf_start + leaf_count);
234 let output = source_count + steps.len();
235 steps.push(ForceMergeStep { inputs });
236 output
237 }
238
239 let mut steps = Vec::with_capacity(internal_count);
240 let root = build(
241 0,
242 source_count,
243 internal_count,
244 deficit,
245 source_count,
246 &mut steps,
247 );
248 debug_assert_eq!(steps.len(), internal_count);
249 ForceMergeHierarchy { steps, root }
250}
251
252struct ActiveOperationState {
262 segment_ids: HashSet<String>,
263 operation_tokens: HashSet<u64>,
264 indexing_tokens: HashSet<u64>,
270 next_operation_token: u64,
271 accepting: bool,
272 non_indexing_paused: bool,
276}
277
278struct ActiveSegmentOperations {
279 inner: parking_lot::Mutex<ActiveOperationState>,
280 idle: Notify,
281 shutdown: Notify,
282 shutdown_requested: Arc<AtomicBool>,
283 index_label: Arc<str>,
286}
287
288impl ActiveSegmentOperations {
289 fn new(index_label: Arc<str>) -> Self {
290 Self {
291 inner: parking_lot::Mutex::new(ActiveOperationState {
292 segment_ids: HashSet::new(),
293 operation_tokens: HashSet::new(),
294 indexing_tokens: HashSet::new(),
295 next_operation_token: 0,
296 accepting: true,
297 non_indexing_paused: false,
298 }),
299 idle: Notify::new(),
300 shutdown: Notify::new(),
301 shutdown_requested: Arc::new(AtomicBool::new(false)),
302 index_label,
303 }
304 }
305
306 fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
310 self.try_register_kind(segment_ids, false, false)
311 }
312
313 fn try_register_indexing(
316 self: &Arc<Self>,
317 segment_ids: Vec<String>,
318 ) -> Option<SegmentOperationGuard> {
319 self.try_register_kind(segment_ids, true, false)
320 }
321
322 fn try_register_vector_update(
325 self: &Arc<Self>,
326 segment_ids: Vec<String>,
327 ) -> Option<SegmentOperationGuard> {
328 self.try_register_kind(segment_ids, false, true)
329 }
330
331 fn try_register_kind(
332 self: &Arc<Self>,
333 segment_ids: Vec<String>,
334 indexing: bool,
335 vector_update: bool,
336 ) -> Option<SegmentOperationGuard> {
337 let mut inner = self.inner.lock();
338 if !inner.accepting {
339 log::debug!(
340 "[segment_lifecycle] index={} rejected operation during shutdown",
341 self.index_label
342 );
343 return None;
344 }
345 if !indexing && !vector_update && inner.non_indexing_paused {
346 log::debug!(
347 "[segment_lifecycle] index={} deferred operation during dense vector retraining",
348 self.index_label
349 );
350 return None;
351 }
352 for id in &segment_ids {
354 if inner.segment_ids.contains(id) {
355 log::debug!(
356 "[segment_lifecycle] index={} rejected: {} overlaps with an active operation ({} active IDs)",
357 self.index_label,
358 id,
359 inner.segment_ids.len()
360 );
361 return None;
362 }
363 }
364 log::debug!(
365 "[segment_lifecycle] index={} registered {} IDs (total active: {})",
366 self.index_label,
367 segment_ids.len(),
368 inner.segment_ids.len() + segment_ids.len()
369 );
370 let operation_token = inner.next_operation_token;
371 let next_operation_token = operation_token.checked_add(1)?;
372 for id in &segment_ids {
373 inner.segment_ids.insert(id.clone());
374 }
375 inner.next_operation_token = next_operation_token;
376 inner.operation_tokens.insert(operation_token);
377 if indexing {
378 inner.indexing_tokens.insert(operation_token);
379 }
380 Some(SegmentOperationGuard {
381 active_operations: Arc::clone(self),
382 segment_ids,
383 operation_token,
384 })
385 }
386
387 fn snapshot(&self) -> HashSet<String> {
389 self.inner.lock().segment_ids.clone()
390 }
391
392 fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
403 let inner = self.inner.lock();
404 let tokens = inner
405 .operation_tokens
406 .difference(&inner.indexing_tokens)
407 .copied()
408 .collect();
409 (tokens, inner.indexing_tokens.len())
410 }
411
412 fn stop_accepting(&self) {
415 self.shutdown_requested.store(true, Ordering::Release);
416 let mut inner = self.inner.lock();
417 inner.accepting = false;
418 self.shutdown.notify_waiters();
419 if inner.segment_ids.is_empty() {
420 self.idle.notify_waiters();
421 }
422 }
423
424 fn pause_non_indexing(&self) {
425 self.inner.lock().non_indexing_paused = true;
426 }
427
428 fn resume_non_indexing(&self) {
429 self.inner.lock().non_indexing_paused = false;
430 self.idle.notify_waiters();
431 }
432
433 fn is_accepting(&self) -> bool {
434 self.inner.lock().accepting
435 }
436
437 fn cancellation_flag(&self) -> Arc<AtomicBool> {
438 Arc::clone(&self.shutdown_requested)
439 }
440
441 async fn wait_until_idle(&self) {
445 loop {
446 let notified = self.idle.notified();
447 if self.inner.lock().segment_ids.is_empty() {
448 return;
449 }
450 notified.await;
451 }
452 }
453
454 async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
455 while !operations.is_empty() {
456 let notified = self.idle.notified();
457 if self.inner.lock().operation_tokens.is_disjoint(operations) {
458 return;
459 }
460 notified.await;
461 }
462 }
463
464 async fn wait_for_shutdown(&self) {
467 loop {
468 let notified = self.shutdown.notified();
469 if !self.inner.lock().accepting {
470 return;
471 }
472 notified.await;
473 }
474 }
475}
476
477pub(crate) struct SegmentOperationGuard {
481 active_operations: Arc<ActiveSegmentOperations>,
482 segment_ids: Vec<String>,
483 operation_token: u64,
484}
485
486impl Drop for SegmentOperationGuard {
487 fn drop(&mut self) {
488 let mut inner = self.active_operations.inner.lock();
489 for id in &self.segment_ids {
490 inner.segment_ids.remove(id);
491 }
492 inner.operation_tokens.remove(&self.operation_token);
493 inner.indexing_tokens.remove(&self.operation_token);
494 self.active_operations.idle.notify_waiters();
497 if inner.segment_ids.is_empty() {
498 debug_assert!(inner.operation_tokens.is_empty());
499 }
500 }
501}
502
503struct VectorArtifactUpdateLease {
510 updating: Arc<AtomicBool>,
511 active_operations: Arc<ActiveSegmentOperations>,
512}
513
514impl Drop for VectorArtifactUpdateLease {
515 fn drop(&mut self) {
516 self.updating.store(false, Ordering::Release);
517 self.active_operations.resume_non_indexing();
518 }
519}
520
521#[derive(Clone)]
522pub(crate) struct VectorArtifactUpdateGuard {
523 _lease: Arc<VectorArtifactUpdateLease>,
524}
525
526static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
530
531const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
532const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
533
534#[derive(Default)]
535struct MergeRetryState {
536 retry_after: Option<std::time::Instant>,
537 consecutive_failures: u32,
538}
539
540fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
541 let shift = consecutive_failures.saturating_sub(1).min(16);
542 MERGE_RETRY_BASE_DELAY
543 .checked_mul(1u32 << shift)
544 .unwrap_or(MERGE_RETRY_MAX_DELAY)
545 .min(MERGE_RETRY_MAX_DELAY)
546}
547
548struct DrainedMergeHandles<'a> {
557 shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
558 drained: Vec<JoinHandle<()>>,
559}
560
561impl<'a> DrainedMergeHandles<'a> {
562 fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
563 let drained = std::mem::take(&mut *shared.lock());
564 Self { shared, drained }
565 }
566
567 fn is_empty(&self) -> bool {
568 self.drained.is_empty()
569 }
570
571 async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
575 let handle = self.drained.last_mut()?;
576 let result = handle.await;
577 self.drained.pop();
578 Some(result)
579 }
580}
581
582impl Drop for DrainedMergeHandles<'_> {
583 fn drop(&mut self) {
584 if !self.drained.is_empty() {
585 self.shared.lock().append(&mut self.drained);
586 }
587 }
588}
589
590fn try_spawn_lifecycle<F>(
597 handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
598 runtime: &tokio::runtime::Handle,
599 future: F,
600) -> bool
601where
602 F: std::future::Future<Output = ()> + Send + 'static,
603{
604 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
605 let mut handles = handles.lock();
606 handles.retain(|handle| !handle.is_finished());
607 handles.push(runtime.spawn(future));
608 }))
609 .is_ok()
610}
611
612struct OutputCleanupGuard {
620 segment_id: SegmentId,
621 cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
622}
623
624impl OutputCleanupGuard {
625 fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
626 Self {
627 segment_id,
628 cleanup: Some(cleanup),
629 }
630 }
631
632 fn disarm(&mut self) {
633 self.cleanup = None;
634 }
635}
636
637impl Drop for OutputCleanupGuard {
638 fn drop(&mut self) {
639 if let Some(cleanup) = self.cleanup.take() {
640 cleanup(self.segment_id);
641 }
642 }
643}
644
645struct ManagerState {
647 metadata: IndexMetadata,
648 merge_policy: Box<dyn MergePolicy>,
649}
650
651type ReplacementRefresh = Arc<
652 dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>
653 + Send
654 + Sync,
655>;
656
657async fn refresh_replacement_topology(refresh: Option<ReplacementRefresh>, index_label: &str) {
661 let Some(refresh) = refresh else {
662 return;
663 };
664 let mut last_error = None;
665 for attempt in 0..3 {
666 match refresh().await {
667 Ok(()) => return,
668 Err(error) => {
669 last_error = Some(error);
670 if attempt < 2 {
671 tokio::time::sleep(std::time::Duration::from_secs(1 << attempt)).await;
672 }
673 }
674 }
675 }
676 if let Some(error) = last_error {
677 log::warn!(
678 "[segment_lifecycle] index={index_label} replacement topology refresh failed after 3 attempts: {}",
679 error,
680 );
681 }
682}
683
684#[cfg(feature = "native")]
685struct MergeTaskError {
686 error: Error,
687 unavailable_segments: Vec<String>,
688}
689
690#[cfg(feature = "native")]
691impl MergeTaskError {
692 fn source(segment_id: String, error: Error) -> Self {
693 Self {
694 error,
695 unavailable_segments: vec![segment_id],
696 }
697 }
698
699 fn sources(segment_ids: Vec<String>, error: Error) -> Self {
700 Self {
701 error,
702 unavailable_segments: segment_ids,
703 }
704 }
705}
706
707#[cfg(feature = "native")]
708impl From<Error> for MergeTaskError {
709 fn from(error: Error) -> Self {
710 Self {
711 error,
712 unavailable_segments: Vec::new(),
713 }
714 }
715}
716
717#[cfg(feature = "native")]
718fn is_deterministic_source_error(error: &Error) -> bool {
719 matches!(error, Error::Corruption(_) | Error::Serialization(_))
720 || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
721}
722
723#[cfg(feature = "native")]
724fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
725 if is_deterministic_source_error(&error) {
726 MergeTaskError::source(segment_id, error)
727 } else {
728 MergeTaskError::from(error)
732 }
733}
734
735#[cfg(feature = "native")]
736type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
737
738#[derive(Clone, Copy)]
739enum ReplacementLayout {
740 BlockCopy,
744 BpReordered { converged: bool },
746 PreserveSingleSource,
749}
750
751fn replacement_bp_state(
752 parent_has_debt: bool,
753 parent_unconverged_passes: u32,
754 layout: ReplacementLayout,
755) -> (bool, bool, u32) {
756 match layout {
757 ReplacementLayout::BlockCopy => (
758 false,
759 !parent_has_debt,
760 if parent_has_debt {
761 parent_unconverged_passes
762 } else {
763 0
764 },
765 ),
766 ReplacementLayout::BpReordered { converged } => (
767 true,
768 converged,
769 if converged {
770 0
771 } else {
772 parent_unconverged_passes.saturating_add(1)
773 },
774 ),
775 ReplacementLayout::PreserveSingleSource => {
776 unreachable!("preserved layouts retain the complete source metadata")
777 }
778 }
779}
780
781#[derive(Clone, Copy, Debug, Eq, PartialEq)]
782enum VectorSegmentRewriteOutcome {
783 Rewritten,
784 AlreadyCurrent,
785 SourceGone,
786 Conflict,
787 Deferred,
788}
789
790pub(crate) struct StagedVectorSegment {
793 source_id: String,
794 output_id: SegmentId,
795 doc_count: u32,
796 _operation: SegmentOperationGuard,
797 cleanup: OutputCleanupGuard,
798}
799
800pub struct SegmentManager<D: DirectoryWriter + 'static> {
804 state: Arc<AsyncMutex<ManagerState>>,
806
807 active_operations: Arc<ActiveSegmentOperations>,
809
810 quarantined_segments: parking_lot::Mutex<HashSet<String>>,
815
816 merge_retry: parking_lot::Mutex<MergeRetryState>,
819
820 reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
824
825 merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
827
828 global_merge_wakeup_pending: AtomicBool,
832
833 force_merge_active: AtomicUsize,
838
839 #[cfg(test)]
842 force_merge_conflict_retries: std::sync::atomic::AtomicU64,
843
844 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
848
849 trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
853
854 vector_artifact_update: Arc<AtomicBool>,
858
859 tracker: Arc<SegmentTracker>,
861
862 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
864
865 directory: Arc<D>,
867 schema: Arc<crate::dsl::Schema>,
869 term_cache_blocks: usize,
871 merge_permits: Arc<Semaphore>,
875 global_merge_permits: Arc<Semaphore>,
877 reorder_permits: Arc<ReorderConcurrencyGate>,
881 reorder_on_merge: bool,
886 merge_bp_time_budget: Option<std::time::Duration>,
890 bp_memory_budget_bytes: usize,
893 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
896 replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
900}
901
902struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
903
904impl Drop for ForceMergeActivityGuard<'_> {
905 fn drop(&mut self) {
906 self.0.fetch_sub(1, Ordering::AcqRel);
907 }
908}
909
910impl<D: DirectoryWriter + 'static> SegmentManager<D> {
911 #[allow(clippy::too_many_arguments)]
913 pub fn new(
914 directory: Arc<D>,
915 schema: Arc<crate::dsl::Schema>,
916 metadata: IndexMetadata,
917 merge_policy: Box<dyn MergePolicy>,
918 term_cache_blocks: usize,
919 max_concurrent_merges: usize,
920 global_merge_permits: Arc<Semaphore>,
921 merge_bp_time_budget: Option<std::time::Duration>,
922 bp_memory_budget_bytes: usize,
923 reorder_permits: Arc<ReorderConcurrencyGate>,
924 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
925 ) -> Self {
926 let reorder_on_merge = schema.reorder_on_merge();
929 if reorder_on_merge {
930 log::info!(
931 "[merge] index={} reorder-on-merge enabled by index schema",
932 schema.index_label()
933 );
934 }
935
936 let tracker = Arc::new(SegmentTracker::new());
937 for seg_id in metadata.segment_metas.keys() {
938 tracker.register(seg_id);
939 }
940
941 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
942 Arc::new(parking_lot::Mutex::new(Vec::new()));
943 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
944 let dir = Arc::clone(&directory);
945 let tracker = Arc::clone(&tracker);
946 let lifecycle_handles = Arc::clone(&lifecycle_handles);
947 let cleanup_index_label: Arc<str> = schema.index_label().into();
948 Arc::new(move |segment_ids| {
949 let Ok(handle) = tokio::runtime::Handle::try_current() else {
952 tracker.complete_deletion(&segment_ids);
955 return;
956 };
957 let dir = Arc::clone(&dir);
958 let task_tracker = Arc::clone(&tracker);
959 let task_index_label = Arc::clone(&cleanup_index_label);
960 let cleanup_ids = segment_ids.clone();
961 let future = async move {
962 for &segment_id in &segment_ids {
963 log::info!(
964 "[segment_cleanup] index={} deleting deferred segment {}",
965 task_index_label,
966 segment_id.to_hex()
967 );
968 if let Err(error) =
969 crate::segment::delete_segment(dir.as_ref(), segment_id).await
970 {
971 log::warn!(
972 "[segment_cleanup] index={} deferred delete failed for {}: {}",
973 task_index_label,
974 segment_id.to_hex(),
975 error,
976 );
977 }
978 }
979 task_tracker.complete_deletion(&segment_ids);
980 };
981 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
982 tracker.complete_deletion(&cleanup_ids);
986 log::warn!(
987 "[segment_cleanup] index={} runtime rejected deferred deletion; files will be swept later",
988 cleanup_index_label
989 );
990 }
991 })
992 };
993
994 Self {
995 state: Arc::new(AsyncMutex::new(ManagerState {
996 metadata,
997 merge_policy,
998 })),
999 active_operations: Arc::new(ActiveSegmentOperations::new(schema.index_label().into())),
1000 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
1001 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
1002 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
1003 merge_handles: parking_lot::Mutex::new(Vec::new()),
1004 global_merge_wakeup_pending: AtomicBool::new(false),
1005 force_merge_active: AtomicUsize::new(0),
1006 #[cfg(test)]
1007 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
1008 lifecycle_handles,
1009 trained: Arc::new(ArcSwapOption::new(None)),
1010 vector_artifact_update: Arc::new(AtomicBool::new(false)),
1011 tracker,
1012 delete_fn,
1013 directory,
1014 schema,
1015 term_cache_blocks,
1016 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
1017 global_merge_permits,
1018 reorder_permits,
1019 reorder_on_merge,
1020 merge_bp_time_budget,
1021 bp_memory_budget_bytes,
1022 background_reorder_pool,
1023 replacement_refresh: parking_lot::RwLock::new(None),
1024 }
1025 }
1026
1027 pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1028 where
1029 F: Fn() -> Fut + Send + Sync + 'static,
1030 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1031 {
1032 *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1033 }
1034
1035 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1041 if let Some(pool) = &self.background_reorder_pool {
1042 return Arc::clone(pool);
1043 }
1044 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1045 let threads = (num_cpus::get() / 2).max(1);
1046 log::info!(
1047 "[merge] process-wide background CPU pool: {} thread(s)",
1048 threads
1049 );
1050 Arc::new(
1051 rayon::ThreadPoolBuilder::new()
1052 .num_threads(threads)
1053 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
1054 .build()
1055 .expect("failed to build background CPU pool"),
1056 )
1057 }))
1058 }
1059
1060 pub fn begin_shutdown(&self) {
1064 self.active_operations.stop_accepting();
1065 }
1066
1067 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1075 where
1076 T: Send + 'static,
1077 F: std::future::Future<Output = Result<T>> + Send + 'static,
1078 {
1079 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1080 let future = async move {
1081 let result = transaction.await;
1082 let _ = result_tx.send(result);
1083 };
1084 let runtime = tokio::runtime::Handle::current();
1085 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1086 return Err(Error::Internal(
1087 "runtime rejected lifecycle metadata transaction".into(),
1088 ));
1089 }
1090 result_rx.await.map_err(|_| {
1091 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1092 })?
1093 }
1094
1095 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1097 let manager = Arc::clone(self);
1098 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1099 let Ok(handle) = tokio::runtime::Handle::try_current() else {
1100 log::warn!(
1101 "[segment_cleanup] index={} runtime unavailable; partial output {} will be swept on startup",
1102 manager.schema.index_label(),
1103 segment_id.to_hex(),
1104 );
1105 return;
1106 };
1107
1108 let cleanup_manager = Arc::clone(&manager);
1109 let future = async move {
1110 cleanup_manager
1111 .delete_output_if_unregistered(segment_id, "task unwind")
1112 .await;
1113 };
1114 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1115 log::warn!(
1116 "[segment_cleanup] index={} runtime rejected output cleanup; {} will be swept on startup",
1117 manager.schema.index_label(),
1118 segment_id.to_hex(),
1119 );
1120 }
1121 });
1122
1123 OutputCleanupGuard::new(output_id, cleanup)
1124 }
1125
1126 pub(crate) fn schedule_unpublished_segment_cleanup(
1131 self: &Arc<Self>,
1132 output_id: SegmentId,
1133 operation: SegmentOperationGuard,
1134 runtime: tokio::runtime::Handle,
1135 ) {
1136 let manager = Arc::clone(self);
1137 let output_hex = output_id.to_hex();
1138 let future = async move {
1139 manager
1140 .delete_output_if_unregistered(output_id, "indexing abort or failure")
1141 .await;
1142 drop(operation);
1143 };
1144 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1145 log::warn!(
1148 "[segment_cleanup] index={} runtime unavailable; indexing output {} will be swept on startup",
1149 self.schema.index_label(),
1150 output_hex,
1151 );
1152 }
1153 }
1154
1155 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1161 match self
1162 .active_operations
1163 .try_register_indexing(vec![segment_id.clone()])
1164 {
1165 Some(operation) => Ok(operation),
1166 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1167 None => Err(Error::Corruption(format!(
1168 "new segment ID {} is already owned by an active operation",
1169 segment_id
1170 ))),
1171 }
1172 }
1173
1174 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1178 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1179 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1180 })?;
1181 let files = SegmentFiles::new(id.0);
1182
1183 for path in files.mandatory_paths() {
1184 if !self.directory.exists(path).await.map_err(Error::Io)? {
1185 return Err(Error::Corruption(format!(
1186 "segment {} cannot be published: mandatory file {:?} is missing",
1187 segment_id, path
1188 )));
1189 }
1190 }
1191
1192 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1193 Error::Corruption(format!(
1194 "segment {} cannot be published: missing/unreadable {:?}: {}",
1195 segment_id, files.meta, e
1196 ))
1197 })?;
1198 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1199 Error::Corruption(format!(
1200 "segment {} cannot be published: failed reading {:?}: {}",
1201 segment_id, files.meta, e
1202 ))
1203 })?;
1204 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1205 Error::Corruption(format!(
1206 "segment {} cannot be published: invalid {:?}: {}",
1207 segment_id, files.meta, e
1208 ))
1209 })?;
1210
1211 if meta.id != id.0 || meta.num_docs != expected_docs {
1212 return Err(Error::Corruption(format!(
1213 "segment {} cannot be published: metadata identity/docs mismatch \
1214 (id={:032x}, docs={}, expected_docs={})",
1215 segment_id, meta.id, meta.num_docs, expected_docs
1216 )));
1217 }
1218
1219 Ok(())
1220 }
1221
1222 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1223 let inserted = self
1224 .quarantined_segments
1225 .lock()
1226 .insert(segment_id.to_string());
1227 if inserted {
1228 log::error!(
1229 "[merge] index={} quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1230 It remains metadata-live for explicit repair but is excluded from merges until restart",
1231 self.schema.index_label(),
1232 segment_id,
1233 error,
1234 );
1235 }
1236 }
1237
1238 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1239 let mut retry = self.merge_retry.lock();
1240 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1241 let delay = merge_retry_delay(retry.consecutive_failures);
1242 retry.retry_after = std::time::Instant::now().checked_add(delay);
1243 log::warn!(
1244 "[merge] index={} pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1245 self.schema.index_label(),
1246 delay.as_secs_f64(),
1247 retry.consecutive_failures,
1248 error,
1249 );
1250 delay
1251 }
1252
1253 fn clear_merge_retry_backoff(&self) {
1254 *self.merge_retry.lock() = MergeRetryState::default();
1255 }
1256
1257 fn merge_retry_is_paused(&self) -> bool {
1258 let mut retry = self.merge_retry.lock();
1259 match retry.retry_after {
1260 Some(deadline) if deadline > std::time::Instant::now() => true,
1261 Some(_) => {
1262 retry.retry_after = None;
1263 false
1264 }
1265 None => false,
1266 }
1267 }
1268
1269 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1270 let mut retries = self.reorder_retries.lock();
1271 let retry = retries.entry(segment_id.to_string()).or_default();
1272 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1273 let delay = merge_retry_delay(retry.consecutive_failures);
1274 retry.retry_after = std::time::Instant::now().checked_add(delay);
1275 log::warn!(
1276 "[reorder] index={} pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1277 self.schema.index_label(),
1278 segment_id,
1279 delay.as_secs_f64(),
1280 retry.consecutive_failures,
1281 error,
1282 );
1283 }
1284
1285 fn clear_reorder_retry(&self, segment_id: &str) {
1286 self.reorder_retries.lock().remove(segment_id);
1287 }
1288
1289 fn paused_reorder_segments(&self) -> HashSet<String> {
1290 let now = std::time::Instant::now();
1291 let mut retries = self.reorder_retries.lock();
1292 let mut paused = HashSet::new();
1293 for (segment_id, retry) in retries.iter_mut() {
1294 match retry.retry_after {
1295 Some(deadline) if deadline > now => {
1296 paused.insert(segment_id.clone());
1297 }
1298 Some(_) => retry.retry_after = None,
1299 None => {}
1300 }
1301 }
1302 paused
1303 }
1304
1305 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1309 if self
1310 .global_merge_wakeup_pending
1311 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1312 .is_err()
1313 {
1314 return;
1315 }
1316
1317 let manager = Arc::clone(self);
1318 let future = async move {
1319 let capacity = tokio::select! {
1320 biased;
1321 () = manager.active_operations.wait_for_shutdown() => None,
1322 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1323 };
1324
1325 manager
1326 .global_merge_wakeup_pending
1327 .store(false, Ordering::Release);
1328 if let Some(permit) = capacity {
1329 drop(permit);
1333 manager.maybe_merge().await;
1334 }
1335 };
1336 let runtime = tokio::runtime::Handle::current();
1337 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1338 self.global_merge_wakeup_pending
1339 .store(false, Ordering::Release);
1340 log::warn!(
1341 "[merge] index={} runtime rejected global-capacity wakeup task",
1342 self.schema.index_label()
1343 );
1344 }
1345 }
1346
1347 #[cfg(test)]
1348 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1349 self.quarantined_segments.lock().contains(segment_id)
1350 }
1351
1352 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1357 let output_hex = output_id.to_hex();
1358 {
1359 let st = self.state.lock().await;
1360 if st.metadata.has_segment(&output_hex) {
1361 return;
1362 }
1363 }
1364
1365 log::info!(
1369 "[segment_cleanup] index={} deleting uncommitted output {} after {}",
1370 self.schema.index_label(),
1371 output_hex,
1372 reason,
1373 );
1374 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1375 {
1376 log::warn!(
1377 "[segment_cleanup] index={} failed deleting uncommitted output {}: {}",
1378 self.schema.index_label(),
1379 output_hex,
1380 error,
1381 );
1382 }
1383 }
1384
1385 pub async fn get_segment_ids(&self) -> Vec<String> {
1391 self.state.lock().await.metadata.segment_ids()
1392 }
1393
1394 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1396 self.trained.load_full()
1397 }
1398
1399 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1406 if self.vector_artifact_update.load(Ordering::Acquire) {
1407 return None;
1408 }
1409 let trained = self.trained.load_full();
1410 if self.vector_artifact_update.load(Ordering::Acquire) {
1411 None
1412 } else {
1413 trained
1414 }
1415 }
1416
1417 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1436 self.vector_artifact_update
1437 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1438 .map_err(|_| {
1439 Error::Internal("a trained-vector artifact update is already in progress".into())
1440 })?;
1441 self.active_operations.pause_non_indexing();
1442 let guard = VectorArtifactUpdateGuard {
1443 _lease: Arc::new(VectorArtifactUpdateLease {
1444 updating: Arc::clone(&self.vector_artifact_update),
1445 active_operations: Arc::clone(&self.active_operations),
1446 }),
1447 };
1448 let (preexisting, parked_indexing) =
1449 self.active_operations.draining_operation_tokens_snapshot();
1450 if parked_indexing > 0 {
1451 return Err(Error::Internal(format!(
1452 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1453 segment(s) are built but uncommitted; commit or abort the pending \
1454 generation and retry"
1455 )));
1456 }
1457 self.active_operations
1458 .wait_until_operations_finish(&preexisting)
1459 .await;
1460 Ok(guard)
1461 }
1462
1463 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1466 let vector_fields = {
1468 let st = self.state.lock().await;
1469 st.metadata.vector_fields.clone()
1470 };
1471 let trained = IndexMetadata::try_load_trained_from_fields(
1473 &vector_fields,
1474 self.schema.as_ref(),
1475 self.directory.as_ref(),
1476 )
1477 .await?
1478 .map(Arc::new);
1479 self.trained.store(trained);
1483 Ok(())
1484 }
1485
1486 pub(crate) async fn publish_vector_generation(
1493 self: &Arc<Self>,
1494 artifact_update: &VectorArtifactUpdateGuard,
1495 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1496 next_trained: Arc<TrainedVectorStructures>,
1497 mut staged: Vec<StagedVectorSegment>,
1498 ) -> Result<()> {
1499 if !self.vector_artifact_update.load(Ordering::Acquire) {
1500 return Err(Error::Internal(
1501 "vector generation publication lost its exclusive update lease".into(),
1502 ));
1503 }
1504
1505 for replacement in &staged {
1506 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1507 .await?;
1508 }
1509
1510 let mut st = Arc::clone(&self.state).lock_owned().await;
1511 let mut next = st.metadata.clone();
1512 next.vector_fields = vector_fields;
1513 next.refresh_total_vectors();
1514
1515 for replacement in &staged {
1516 let source_info = next
1517 .segment_metas
1518 .remove(&replacement.source_id)
1519 .ok_or_else(|| {
1520 Error::Corruption(format!(
1521 "vector generation source {} disappeared before publication",
1522 replacement.source_id,
1523 ))
1524 })?;
1525 let output_hex = replacement.output_id.to_hex();
1526 if next.segment_metas.contains_key(&output_hex) {
1527 return Err(Error::Corruption(format!(
1528 "vector generation output {output_hex} is already metadata-live"
1529 )));
1530 }
1531 next.add_segment_meta(output_hex, source_info);
1534 }
1535
1536 let directory = Arc::clone(&self.directory);
1537 let trained = Arc::clone(&self.trained);
1538 let tracker = Arc::clone(&self.tracker);
1539 let replacement_refresh = self.replacement_refresh.read().clone();
1540 let artifact_update = artifact_update.clone();
1544 let index_label = self.schema.index_label().to_owned();
1545 self.run_lifecycle_transaction(async move {
1546 let _artifact_update = artifact_update;
1547 next.save(directory.as_ref()).await?;
1548
1549 for replacement in &staged {
1550 tracker.register(&replacement.output_id.to_hex());
1551 }
1552 st.metadata = next;
1553 trained.store(Some(next_trained));
1554
1555 for replacement in &mut staged {
1558 replacement.cleanup.disarm();
1559 }
1560 let retired = staged
1561 .iter()
1562 .map(|replacement| replacement.source_id.clone())
1563 .collect::<Vec<_>>();
1564 let ready_to_delete = tracker.mark_for_deletion(&retired);
1565 drop(st);
1566 for &segment_id in &ready_to_delete {
1567 if let Err(error) =
1568 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1569 {
1570 log::warn!(
1571 "[segment_cleanup] index={index_label} immediate dense-vector generation delete failed for {}: {}",
1572 segment_id.to_hex(),
1573 error,
1574 );
1575 }
1576 }
1577 tracker.complete_deletion(&ready_to_delete);
1578 refresh_replacement_topology(replacement_refresh, &index_label).await;
1579 Ok(())
1580 })
1581 .await
1582 }
1583
1584 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1586 where
1587 F: FnOnce(&IndexMetadata) -> R,
1588 {
1589 let st = self.state.lock().await;
1590 f(&st.metadata)
1591 }
1592
1593 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1595 where
1596 F: FnOnce(&mut IndexMetadata),
1597 {
1598 let mut st = Arc::clone(&self.state).lock_owned().await;
1599 let mut next = st.metadata.clone();
1600 f(&mut next);
1601 let directory = Arc::clone(&self.directory);
1602 self.run_lifecycle_transaction(async move {
1603 next.save(directory.as_ref()).await?;
1604 st.metadata = next;
1605 Ok(())
1606 })
1607 .await
1608 }
1609
1610 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1613 let (acquired, trained) = {
1614 let st = self.state.lock().await;
1615 let segment_ids = st.metadata.segment_ids();
1616 (self.tracker.acquire(&segment_ids), self.trained.load_full())
1617 };
1618
1619 SegmentSnapshot::with_generation(
1620 Arc::clone(&self.tracker),
1621 acquired,
1622 trained,
1623 Arc::clone(&self.delete_fn),
1624 )
1625 }
1626
1627 pub fn tracker(&self) -> Arc<SegmentTracker> {
1629 Arc::clone(&self.tracker)
1630 }
1631
1632 pub fn directory(&self) -> Arc<D> {
1634 Arc::clone(&self.directory)
1635 }
1636}
1637
1638#[cfg(feature = "native")]
1643impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1644 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1646 for (segment_id, num_docs) in new_segments {
1649 self.validate_completed_segment(segment_id, *num_docs)
1650 .await?;
1651 }
1652
1653 let mut st = Arc::clone(&self.state).lock_owned().await;
1654 let mut next = st.metadata.clone();
1655 let mut added = Vec::new();
1656 for (segment_id, num_docs) in new_segments {
1657 if !next.has_segment(segment_id) {
1658 next.add_segment(segment_id.clone(), *num_docs);
1659 added.push(segment_id.clone());
1660 }
1661 }
1662
1663 let directory = Arc::clone(&self.directory);
1669 let tracker = Arc::clone(&self.tracker);
1670 self.run_lifecycle_transaction(async move {
1671 next.save(directory.as_ref()).await?;
1672 for segment_id in &added {
1673 tracker.register(segment_id);
1674 }
1675 st.metadata = next;
1676 Ok(())
1677 })
1678 .await
1679 }
1680
1681 pub async fn maybe_merge(self: &Arc<Self>) {
1692 if !self.active_operations.is_accepting() {
1693 log::debug!(
1694 "[maybe_merge] index={} manager is shutting down, skipping",
1695 self.schema.index_label()
1696 );
1697 return;
1698 }
1699 if self.merge_retry_is_paused() {
1700 log::debug!(
1701 "[maybe_merge] index={} retry backoff active, skipping",
1702 self.schema.index_label()
1703 );
1704 return;
1705 }
1706
1707 {
1710 let mut handles = self.merge_handles.lock();
1711 handles.retain(|h| !h.is_finished());
1712 }
1713 let local_slots = self.merge_permits.available_permits();
1714 let global_slots = self.global_merge_permits.available_permits();
1715 let slots_available = local_slots.min(global_slots);
1716
1717 {
1721 let st = self.state.lock().await;
1722 let quarantined = self.quarantined_segments.lock().clone();
1723 let active_ids = self.active_operations.snapshot();
1724
1725 let live_segments: Vec<SegmentInfo> = st
1730 .metadata
1731 .segment_metas
1732 .iter()
1733 .filter(|(id, _)| {
1734 !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1735 })
1736 .map(|(id, info)| SegmentInfo {
1737 id: id.clone(),
1738 num_docs: info.num_docs,
1739 })
1740 .collect();
1741 let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1742
1743 let segments: Vec<SegmentInfo> = live_segments
1746 .iter()
1747 .filter(|segment| !active_ids.contains(&segment.id))
1748 .cloned()
1749 .collect();
1750
1751 log::debug!(
1752 "[maybe_merge] index={} {} eligible segments",
1753 self.schema.index_label(),
1754 segments.len()
1755 );
1756
1757 let candidates = st.merge_policy.find_merges(&segments);
1758
1759 if candidates.is_empty() {
1760 return;
1761 }
1762
1763 if slots_available == 0 {
1767 if local_slots > 0 && global_slots == 0 {
1768 self.schedule_global_merge_wakeup();
1769 }
1770 log::debug!(
1771 "[maybe_merge] index={} at max concurrent merges, skipping",
1772 self.schema.index_label()
1773 );
1774 return;
1775 }
1776
1777 log::debug!(
1778 "[maybe_merge] index={} {} merge candidates, {} slots available",
1779 self.schema.index_label(),
1780 candidates.len(),
1781 slots_available
1782 );
1783
1784 let mut handles = Vec::new();
1785 for c in candidates {
1786 if handles.len() >= slots_available {
1787 break;
1788 }
1789 let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1795 if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1796 handles.push(h);
1797 }
1798 }
1799 if !handles.is_empty() {
1800 if severe_backlog && self.reorder_on_merge {
1801 log::info!(
1802 "[maybe_merge] index={} severe backlog: {} live segments; started {} fast \
1803 block-copy merge(s), deferring BP to the optimizer",
1804 self.schema.index_label(),
1805 live_segments.len(),
1806 handles.len(),
1807 );
1808 }
1809 self.merge_handles.lock().extend(handles);
1814 }
1815 }
1816 }
1817
1818 fn spawn_merge(
1827 self: &Arc<Self>,
1828 segment_ids_to_merge: Vec<String>,
1829 reorder_bmp: bool,
1830 ) -> Option<JoinHandle<()>> {
1831 if self.force_merge_active.load(Ordering::Acquire) > 0 {
1832 log::debug!(
1833 "[spawn_merge] index={} skipped: explicit force merge has priority",
1834 self.schema.index_label()
1835 );
1836 return None;
1837 }
1838 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1839 Ok(permit) => permit,
1840 Err(_) => {
1841 log::debug!(
1842 "[spawn_merge] index={} skipped: global merge capacity is full",
1843 self.schema.index_label()
1844 );
1845 self.schedule_global_merge_wakeup();
1846 return None;
1847 }
1848 };
1849 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1850 Ok(permit) => permit,
1851 Err(_) => {
1852 log::debug!(
1853 "[spawn_merge] index={} skipped: no merge permit available",
1854 self.schema.index_label()
1855 );
1856 return None;
1857 }
1858 };
1859 let output_id = SegmentId::new();
1860 let output_hex = output_id.to_hex();
1861
1862 let mut all_ids = segment_ids_to_merge.clone();
1863 all_ids.push(output_hex);
1864
1865 let guard = match self.active_operations.try_register(all_ids) {
1866 Some(g) => g,
1867 None => {
1868 log::debug!(
1869 "[spawn_merge] index={} skipped: segments overlap with an active operation",
1870 self.schema.index_label()
1871 );
1872 return None;
1873 }
1874 };
1875
1876 let sm = Arc::clone(self);
1877 let ids = segment_ids_to_merge;
1878
1879 let index_label = self.schema.index_label().to_owned();
1880 Some(tokio::spawn(async move {
1881 let mut reevaluate = false;
1882 let mut retry_delay = None;
1883
1884 let result = sm
1885 .merge_and_replace_registered(
1886 &ids,
1887 output_id,
1888 reorder_bmp,
1889 ReorderPriority::AutomaticMerge,
1890 )
1891 .await;
1892
1893 match result {
1894 Ok(_) => {
1895 sm.clear_merge_retry_backoff();
1896 reevaluate = true;
1897 }
1898 Err(MergeTaskError {
1899 error: Error::IndexClosed,
1900 ..
1901 }) => {
1902 log::debug!(
1903 "[merge] index={index_label} background merge for segments {:?} cancelled during shutdown",
1904 ids,
1905 );
1906 }
1907 Err(MergeTaskError {
1908 error,
1909 unavailable_segments,
1910 }) => {
1911 log::error!(
1912 "[merge] index={index_label} background merge failed for segments {:?}: {}",
1913 ids,
1914 error
1915 );
1916 if !unavailable_segments.is_empty() {
1917 reevaluate = true;
1921 } else {
1922 retry_delay = Some(sm.pause_merge_retries(&error));
1923 }
1924 }
1925 }
1926 drop(guard);
1929 drop(merge_permit);
1931 drop(global_merge_permit);
1932
1933 if reevaluate {
1934 sm.maybe_merge().await;
1935 } else if let Some(retry_delay) = retry_delay {
1936 sm.schedule_merge_retry_wakeup(retry_delay);
1943 }
1944 }))
1945 }
1946
1947 async fn merge_and_replace_registered(
1954 self: &Arc<Self>,
1955 ids: &[String],
1956 output_id: SegmentId,
1957 reorder_bmp: bool,
1958 priority: ReorderPriority,
1959 ) -> MergeTaskResult<(String, u32, bool)> {
1960 let mut output_cleanup = self.output_cleanup_guard(output_id);
1961 let trained = self.trained_for_segment_build();
1962 let granularity = if reorder_bmp {
1963 self.merge_granularity(ids).await
1964 } else {
1965 crate::segment::reorder::BpGranularity::Auto
1966 };
1967 let result = Self::do_merge(
1968 self.directory.as_ref(),
1969 &self.schema,
1970 ids,
1971 output_id,
1972 self.term_cache_blocks,
1973 trained.as_deref(),
1974 reorder_bmp,
1975 granularity,
1976 self.merge_bp_time_budget,
1977 self.bp_memory_budget_bytes,
1978 Arc::clone(&self.reorder_permits),
1979 priority,
1980 self.active_operations.cancellation_flag(),
1981 Some(self.background_cpu_pool()),
1982 )
1983 .await;
1984
1985 let (new_id, doc_count, bp_converged) = match result {
1986 Ok(value) => value,
1987 Err(error) => {
1988 for segment_id in &error.unavailable_segments {
1989 self.quarantine_segment(segment_id, &error.error);
1990 }
1991 self.delete_output_if_unregistered(output_id, "merge failure")
1992 .await;
1993 output_cleanup.disarm();
1994 return Err(error);
1995 }
1996 };
1997
1998 let layout = if reorder_bmp {
1999 ReplacementLayout::BpReordered {
2000 converged: bp_converged,
2001 }
2002 } else {
2003 ReplacementLayout::BlockCopy
2004 };
2005 if let Err(error) = self
2006 .replace_segments(ids, new_id.clone(), doc_count, layout)
2007 .await
2008 {
2009 self.delete_output_if_unregistered(output_id, "replacement failure")
2010 .await;
2011 output_cleanup.disarm();
2012 return Err(MergeTaskError::from(error));
2013 }
2014 output_cleanup.disarm();
2015 Ok((new_id, doc_count, bp_converged))
2016 }
2017
2018 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
2022 let manager = Arc::clone(self);
2023 let future = async move {
2024 tokio::select! {
2025 () = tokio::time::sleep(retry_delay) => {
2026 manager.maybe_merge().await;
2027 }
2028 () = manager.active_operations.wait_for_shutdown() => {}
2029 }
2030 };
2031 let runtime = tokio::runtime::Handle::current();
2032 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2033 log::warn!(
2034 "[merge] index={} runtime rejected merge-retry wakeup task; eligible segments may stay \
2035 unmerged until the next commit re-runs merge policy evaluation",
2036 self.schema.index_label()
2037 );
2038 }
2039 }
2040
2041 async fn replace_segments(
2045 self: &Arc<Self>,
2046 old_ids: &[String],
2047 new_id: String,
2048 doc_count: u32,
2049 layout: ReplacementLayout,
2050 ) -> Result<()> {
2051 self.validate_completed_segment(&new_id, doc_count).await?;
2054 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
2055 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
2056 })?;
2057 let output_reader = SegmentReader::open(
2058 self.directory.as_ref(),
2059 output_id,
2060 Arc::clone(&self.schema),
2061 self.term_cache_blocks,
2062 )
2063 .await
2064 .map_err(|error| match error {
2065 Error::Io(_) | Error::IndexClosed => error,
2069 error => Error::Corruption(format!(
2070 "replacement segment {new_id} failed full reader validation: {error}"
2071 )),
2072 })?;
2073 if output_reader.num_docs() != doc_count {
2074 return Err(Error::Corruption(format!(
2075 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2076 output_reader.num_docs(),
2077 )));
2078 }
2079 drop(output_reader);
2080
2081 let mut st = Arc::clone(&self.state).lock_owned().await;
2082 let missing: Vec<&String> = old_ids
2086 .iter()
2087 .filter(|id| !st.metadata.has_segment(id))
2088 .collect();
2089 if !missing.is_empty() {
2090 return Err(Error::Corruption(format!(
2091 "replace_segments: source segment(s) {:?} not in metadata — \
2092 refusing to add output {} (would duplicate documents)",
2093 missing, new_id
2094 )));
2095 }
2096
2097 let replacement_info = match layout {
2098 ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2099 let generation = old_ids
2100 .iter()
2101 .filter_map(|id| st.metadata.segment_metas.get(id))
2102 .map(|info| info.generation)
2103 .max()
2104 .unwrap_or(0)
2105 .checked_add(1)
2106 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2107 let parent_unconverged_passes = old_ids
2108 .iter()
2109 .filter_map(|id| st.metadata.segment_metas.get(id))
2110 .map(|info| info.bp_unconverged_passes)
2111 .max()
2112 .unwrap_or(0);
2113 let parent_has_debt = old_ids
2114 .iter()
2115 .filter_map(|id| st.metadata.segment_metas.get(id))
2116 .any(|info| !info.bp_converged);
2117 let (reordered, bp_converged, bp_unconverged_passes) =
2118 replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2119 SegmentMetaInfo {
2120 num_docs: doc_count,
2121 ancestors: old_ids.to_vec(),
2122 generation,
2123 reordered,
2124 bp_converged,
2125 bp_unconverged_passes,
2126 }
2127 }
2128 ReplacementLayout::PreserveSingleSource => {
2129 let [source_id] = old_ids else {
2130 return Err(Error::Internal(
2131 "layout-preserving replacement requires exactly one source".into(),
2132 ));
2133 };
2134 let mut source = st
2135 .metadata
2136 .segment_metas
2137 .get(source_id)
2138 .cloned()
2139 .ok_or_else(|| {
2140 Error::Corruption(format!(
2141 "layout-preserving replacement source {source_id} disappeared"
2142 ))
2143 })?;
2144 source.num_docs = doc_count;
2145 source
2146 }
2147 };
2148 let retired_ids = old_ids.to_vec();
2149 let mut next = st.metadata.clone();
2150 for id in old_ids {
2151 next.remove_segment(id);
2152 }
2153 next.add_segment_meta(new_id.clone(), replacement_info);
2154
2155 let directory = Arc::clone(&self.directory);
2156 let tracker = Arc::clone(&self.tracker);
2157 let replacement_refresh = self.replacement_refresh.read().clone();
2158 let index_label = self.schema.index_label().to_owned();
2159 self.run_lifecycle_transaction(async move {
2160 next.save(directory.as_ref()).await?;
2163 tracker.register(&new_id);
2164 st.metadata = next;
2165
2166 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2170 drop(st);
2171 for &segment_id in &ready_to_delete {
2172 if let Err(error) =
2173 crate::segment::delete_segment(directory.as_ref(), segment_id).await
2174 {
2175 log::warn!(
2176 "[segment_cleanup] index={index_label} immediate delete failed for {}: {}",
2177 segment_id.to_hex(),
2178 error,
2179 );
2180 }
2181 }
2182 tracker.complete_deletion(&ready_to_delete);
2183 refresh_replacement_topology(replacement_refresh, &index_label).await;
2184 Ok(())
2185 })
2186 .await
2187 }
2188
2189 #[allow(clippy::too_many_arguments)]
2194 async fn do_merge(
2195 directory: &D,
2196 schema: &Arc<crate::dsl::Schema>,
2197 segment_ids_to_merge: &[String],
2198 output_segment_id: SegmentId,
2199 term_cache_blocks: usize,
2200 trained: Option<&TrainedVectorStructures>,
2201 reorder_bmp: bool,
2202 granularity: crate::segment::reorder::BpGranularity,
2203 merge_bp_time_budget: Option<std::time::Duration>,
2204 bp_memory_budget_bytes: usize,
2205 reorder_permits: Arc<ReorderConcurrencyGate>,
2206 reorder_priority: ReorderPriority,
2207 cancellation: Arc<AtomicBool>,
2208 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2209 ) -> MergeTaskResult<(String, u32, bool)> {
2210 let output_hex = output_segment_id.to_hex();
2211 let load_start = std::time::Instant::now();
2212
2213 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2214 for id_str in segment_ids_to_merge {
2215 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2216 MergeTaskError::source(
2217 id_str.clone(),
2218 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2219 )
2220 })?;
2221 segment_ids.push(id);
2222 }
2223
2224 let mut unavailable_sources = Vec::new();
2229 let mut missing_files = Vec::new();
2230 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2231 let files = SegmentFiles::new(id.0);
2232 let mut source_unavailable = false;
2233 for path in files.mandatory_paths() {
2234 let exists = directory
2235 .exists(path)
2236 .await
2237 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2238 if !exists {
2239 source_unavailable = true;
2240 missing_files.push(format!("{}:{:?}", id_str, path));
2241 }
2242 }
2243 if source_unavailable {
2244 unavailable_sources.push(id_str.clone());
2245 }
2246 }
2247 if !unavailable_sources.is_empty() {
2248 return Err(MergeTaskError::sources(
2249 unavailable_sources,
2250 Error::Corruption(format!(
2251 "merge sources are missing mandatory files: {}",
2252 missing_files.join(", ")
2253 )),
2254 ));
2255 }
2256
2257 let schema_arc = Arc::clone(schema);
2258 let futures: Vec<_> = segment_ids
2259 .iter()
2260 .map(|&sid| {
2261 let sch = Arc::clone(&schema_arc);
2262 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
2263 })
2264 .collect();
2265
2266 let results = futures::future::join_all(futures).await;
2267 let mut readers = Vec::with_capacity(results.len());
2268 let mut total_docs = 0u64;
2269 for (i, result) in results.into_iter().enumerate() {
2270 match result {
2271 Ok(r) => {
2272 total_docs += r.meta().num_docs as u64;
2273 readers.push(r);
2274 }
2275 Err(e) => {
2276 log::error!(
2277 "[merge] index={} Failed to open segment {}: {:?}",
2278 schema.index_label(),
2279 segment_ids_to_merge[i],
2280 e
2281 );
2282 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2283 }
2284 }
2285 }
2286 if total_docs > u32::MAX as u64 {
2287 return Err(Error::Internal(format!(
2288 "Merged segment doc count ({}) exceeds u32::MAX",
2289 total_docs
2290 ))
2291 .into());
2292 }
2293
2294 for (i, reader) in readers.iter().enumerate() {
2298 let meta_docs = reader.meta().num_docs;
2299 let store_docs = reader.store().num_docs();
2300 if store_docs != meta_docs {
2301 return Err(MergeTaskError::source(
2302 segment_ids_to_merge[i].clone(),
2303 Error::Corruption(format!(
2304 "pre-merge validation: segment {} store has {} docs but meta says {}",
2305 segment_ids_to_merge[i], store_docs, meta_docs
2306 )),
2307 ));
2308 }
2309 }
2310
2311 log::info!(
2312 "[merge] index={} loaded {} segment readers in {:.1}s",
2313 schema.index_label(),
2314 readers.len(),
2315 load_start.elapsed().as_secs_f64()
2316 );
2317
2318 let merger = SegmentMerger::new(Arc::clone(schema))
2319 .with_bmp_reorder(reorder_bmp)
2320 .with_granularity(granularity)
2321 .with_bp_budget(crate::segment::BpBudget {
2322 min_partition_docs: None,
2323 time_budget: merge_bp_time_budget,
2324 })
2325 .with_cancellation(cancellation)
2326 .with_bp_memory_budget(bp_memory_budget_bytes)
2327 .with_reorder_permits(reorder_permits)
2328 .with_reorder_priority(reorder_priority)
2329 .with_background_pool(bg_cpu_pool);
2330
2331 log::info!(
2332 "[merge] index={} {} segments -> {} (trained={})",
2333 schema.index_label(),
2334 segment_ids_to_merge.len(),
2335 output_hex,
2336 trained.map_or(0, |t| t.centroids.len()),
2337 );
2338
2339 let (_merged_meta, merge_stats) = merger
2340 .merge(directory, &readers, output_segment_id, trained)
2341 .await
2342 .map_err(|error| {
2343 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2344 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2350 } else {
2351 MergeTaskError::from(error)
2352 }
2353 })?;
2354 let bp_converged = merge_stats.bp_converged;
2355 if !bp_converged {
2356 log::info!(
2357 "[merge] index={} merge-time BP hit its wall-clock budget — output marked unconverged; \
2358 the background optimizer deepens it later",
2359 schema.index_label(),
2360 );
2361 }
2362
2363 log::info!(
2364 "[merge] index={} total wall-clock: {:.1}s ({} segments, {} docs)",
2365 schema.index_label(),
2366 load_start.elapsed().as_secs_f64(),
2367 readers.len(),
2368 total_docs,
2369 );
2370
2371 Ok((output_hex, total_docs as u32, bp_converged))
2372 }
2373
2374 pub async fn abort_merges(&self) {
2384 loop {
2385 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2386 if handles.is_empty() {
2387 return;
2388 }
2389 while let Some(result) = handles.join_next().await {
2390 if let Err(error) = result
2391 && error.is_panic()
2392 {
2393 log::error!(
2394 "[merge] index={} background task panicked while draining: {}",
2395 self.schema.index_label(),
2396 error
2397 );
2398 }
2399 }
2400 }
2401 }
2402
2403 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2408 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2409 while handles.join_next().await.is_some() {}
2410 }
2411
2412 pub async fn wait_for_all_merges(self: &Arc<Self>) {
2421 loop {
2422 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2423 if handles.is_empty() {
2424 break;
2425 }
2426 while handles.join_next().await.is_some() {}
2427 }
2428 }
2429
2430 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2435 self.wait_for_all_merges().await;
2436 self.active_operations.wait_until_idle().await;
2437 loop {
2438 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2439 if handles.is_empty() {
2440 break;
2441 }
2442 for handle in handles {
2443 if let Err(error) = handle.await
2444 && error.is_panic()
2445 {
2446 log::error!(
2447 "[segment_cleanup] index={} task panicked while draining: {}",
2448 self.schema.index_label(),
2449 error
2450 );
2451 }
2452 }
2453 }
2454 }
2455
2456 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2465 self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2466 .await
2467 }
2468
2469 pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2479 self: &Arc<Self>,
2480 mut refresh_snapshots: F,
2481 ) -> Result<()>
2482 where
2483 F: FnMut() -> Fut,
2484 Fut: std::future::Future<Output = Result<()>>,
2485 {
2486 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2491 std::time::Duration::from_millis(100);
2492 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2496
2497 let (_force_merge_activity, max_segment_docs) = {
2498 let st = self.state.lock().await;
2499 self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2503 (
2504 ForceMergeActivityGuard(&self.force_merge_active),
2505 st.merge_policy.max_segment_docs(),
2506 )
2507 };
2508
2509 let background_merges = self
2512 .merge_handles
2513 .lock()
2514 .iter()
2515 .filter(|handle| !handle.is_finished())
2516 .count();
2517 if background_merges > 0 {
2518 log::info!(
2519 "[force_merge] index={} waiting for {} in-flight background merge(s) before planning",
2520 self.schema.index_label(),
2521 background_merges,
2522 );
2523 }
2524 let drain_start = std::time::Instant::now();
2525 self.wait_for_all_merges().await;
2526 if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2527 log::info!(
2528 "[force_merge] index={} drained background merges in {:.1}s",
2529 self.schema.index_label(),
2530 drain_start.elapsed().as_secs_f64(),
2531 );
2532 }
2533
2534 let refresh_start = std::time::Instant::now();
2539 refresh_snapshots().await?;
2540 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2541 log::info!(
2542 "[force_merge] index={} initial snapshot refresh took {:.1}s",
2543 self.schema.index_label(),
2544 refresh_start.elapsed().as_secs_f64(),
2545 );
2546 }
2547
2548 let mut completed_outputs = HashSet::new();
2552 let mut logged_held_wait = false;
2554
2555 loop {
2556 if !self.active_operations.is_accepting() {
2557 return Err(Error::IndexClosed);
2558 }
2559
2560 let segments: Vec<(String, u32)> = {
2561 let st = self.state.lock().await;
2562 st.metadata
2563 .segment_metas
2564 .iter()
2565 .filter(|(id, _)| !completed_outputs.contains(*id))
2566 .map(|(id, info)| (id.clone(), info.num_docs))
2567 .collect()
2568 };
2569
2570 let active_ids = self.active_operations.snapshot();
2575 let held = segments
2576 .iter()
2577 .filter(|(id, _)| active_ids.contains(id))
2578 .count();
2579 let free_segments: Vec<_> = segments
2580 .into_iter()
2581 .filter(|(id, _)| !active_ids.contains(id))
2582 .collect();
2583 let max_docs = max_segment_docs
2587 .map(u64::from)
2588 .unwrap_or(u64::from(u32::MAX));
2589 let next_group = plan_force_merge_groups(free_segments, max_docs)
2590 .into_iter()
2591 .find(|group| group.segments.len() >= 2);
2592
2593 let Some(group) = next_group else {
2594 if held == 0 {
2595 if !completed_outputs.is_empty() {
2596 completed_outputs.clear();
2602 continue;
2603 }
2604 refresh_snapshots().await?;
2611 return Ok(());
2612 }
2613 if !logged_held_wait {
2614 log::info!(
2615 "[force_merge] index={} waiting: {} segment(s) held by active \
2616 merge/reorder operations, no free group can merge",
2617 self.schema.index_label(),
2618 held
2619 );
2620 logged_held_wait = true;
2621 } else {
2622 log::debug!(
2623 "[force_merge] index={} still waiting on {} held segment(s)",
2624 self.schema.index_label(),
2625 held
2626 );
2627 }
2628 #[cfg(test)]
2629 self.force_merge_conflict_retries
2630 .fetch_add(1, Ordering::Relaxed);
2631 tokio::select! {
2632 biased;
2633 () = self.active_operations.wait_for_shutdown() => {
2634 return Err(Error::IndexClosed);
2635 }
2636 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2637 }
2638 continue;
2639 };
2640 logged_held_wait = false;
2641
2642 let hierarchy = plan_force_merge_hierarchy(group.segments.len());
2646 let output_ids: Vec<_> = (0..hierarchy.steps.len())
2647 .map(|_| SegmentId::new())
2648 .collect();
2649 let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
2650 let mut all_ids = source_ids.clone();
2651 all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
2652 let group_guard = {
2653 let st = self.state.lock().await;
2654 source_ids
2655 .iter()
2656 .all(|id| st.metadata.has_segment(id))
2657 .then(|| self.active_operations.try_register(all_ids))
2658 .flatten()
2659 };
2660 let _group_guard = match group_guard {
2661 Some(guard) => guard,
2662 None if !self.active_operations.is_accepting() => {
2663 return Err(Error::IndexClosed);
2664 }
2665 None => {
2666 #[cfg(test)]
2667 self.force_merge_conflict_retries
2668 .fetch_add(1, Ordering::Relaxed);
2669 log::debug!(
2670 "[force_merge] index={} group lost a registration race, replanning",
2671 self.schema.index_label()
2672 );
2673 let had_tracked_merges = !self.merge_handles.lock().is_empty();
2674 self.wait_for_merging_thread().await;
2675 if !had_tracked_merges {
2676 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2677 }
2678 continue;
2679 }
2680 };
2681
2682 log::info!(
2683 "[force_merge] index={} planned final group: {} segments, {} docs, {} merge pass(es)",
2684 self.schema.index_label(),
2685 group.segments.len(),
2686 group.total_docs,
2687 output_ids.len(),
2688 );
2689
2690 let group_global_merge_permit = if self.reorder_on_merge {
2700 let capacity_start = std::time::Instant::now();
2701 let permit = tokio::select! {
2702 biased;
2703 () = self.active_operations.wait_for_shutdown() => {
2704 return Err(Error::IndexClosed);
2705 }
2706 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2707 permit.map_err(|_| {
2708 Error::Internal(
2709 "global background merge scheduler is closed".into(),
2710 )
2711 })?
2712 }
2713 };
2714 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2715 log::info!(
2716 "[force_merge] index={} waited {:.1}s for foreground global merge capacity",
2717 self.schema.index_label(),
2718 capacity_start.elapsed().as_secs_f64(),
2719 );
2720 }
2721 Some(permit)
2722 } else {
2723 None
2724 };
2725 let _foreground_reorder = if self.reorder_on_merge {
2726 log::info!(
2727 "[force_merge] index={} prioritizing BP capacity ({} total pass slot(s))",
2728 self.schema.index_label(),
2729 self.reorder_permits.limit(),
2730 );
2731 let admission_start = std::time::Instant::now();
2732 let guard = Arc::clone(&self.reorder_permits)
2733 .begin_foreground()
2734 .await
2735 .map_err(|_| {
2736 Error::Internal("background reorder scheduler is closed".into())
2737 })?;
2738 if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
2739 log::info!(
2740 "[force_merge] index={} acquired foreground BP capacity in {:.1}s",
2741 self.schema.index_label(),
2742 admission_start.elapsed().as_secs_f64(),
2743 );
2744 }
2745 Some(guard)
2746 } else {
2747 None
2748 };
2749
2750 let source_count = group.segments.len();
2751 let mut nodes: Vec<Option<(String, u32)>> =
2752 group.segments.into_iter().map(Some).collect();
2753 nodes.resize_with(source_count + hierarchy.steps.len(), || None);
2754 for (step_index, step) in hierarchy.steps.iter().enumerate() {
2755 let final_pass = step_index + 1 == hierarchy.steps.len();
2756 let mut batch_entries = Vec::with_capacity(step.inputs.len());
2757 for &node in &step.inputs {
2758 let entry = nodes
2759 .get_mut(node)
2760 .and_then(Option::take)
2761 .expect("force-merge hierarchy must reference an available node");
2762 batch_entries.push(entry);
2763 }
2764 let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
2765 let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
2766 let output_id = output_ids[step_index];
2767
2768 let capacity_start = std::time::Instant::now();
2769 let step_global_merge_permit = if group_global_merge_permit.is_none() {
2770 Some(tokio::select! {
2771 biased;
2772 () = self.active_operations.wait_for_shutdown() => {
2773 return Err(Error::IndexClosed);
2774 }
2775 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2776 permit.map_err(|_| {
2777 Error::Internal(
2778 "global background merge scheduler is closed".into(),
2779 )
2780 })?
2781 }
2782 })
2783 } else {
2784 None
2785 };
2786 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2787 log::info!(
2788 "[force_merge] index={} waited {:.1}s for global merge capacity",
2789 self.schema.index_label(),
2790 capacity_start.elapsed().as_secs_f64(),
2791 );
2792 }
2793
2794 let reorder_bmp = final_pass && self.reorder_on_merge;
2798 log::info!(
2799 "[force_merge] index={} {} pass: {} segments ({} docs, bp={})",
2800 self.schema.index_label(),
2801 if final_pass {
2802 "final"
2803 } else {
2804 "fan-in reduction"
2805 },
2806 batch.len(),
2807 batch_docs,
2808 reorder_bmp,
2809 );
2810 let (new_segment_id, total_docs, _) = self
2811 .merge_and_replace_registered(
2812 &batch,
2813 output_id,
2814 reorder_bmp,
2815 ReorderPriority::Foreground,
2816 )
2817 .await
2818 .map_err(|error| error.error)?;
2819 drop(step_global_merge_permit);
2820
2821 let refresh_start = std::time::Instant::now();
2824 refresh_snapshots().await?;
2825 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2826 log::info!(
2827 "[force_merge] index={} post-replacement snapshot refresh took {:.1}s",
2828 self.schema.index_label(),
2829 refresh_start.elapsed().as_secs_f64(),
2830 );
2831 }
2832
2833 let output_node = source_count + step_index;
2834 debug_assert!(nodes[output_node].is_none());
2835 nodes[output_node] = Some((new_segment_id, total_docs));
2836 }
2837 let (root_id, _) = nodes[hierarchy.root]
2838 .take()
2839 .expect("force-merge hierarchy must produce its root");
2840 debug_assert!(nodes.into_iter().all(|node| node.is_none()));
2841 completed_outputs.insert(root_id);
2842 }
2843 }
2844
2845 fn segment_needs_vector_rewrite(
2846 &self,
2847 reader: &SegmentReader,
2848 field_ids: &[u32],
2849 trained: &TrainedVectorStructures,
2850 rewrite_existing: bool,
2851 ) -> Result<bool> {
2852 for &field_id in field_ids {
2853 let flat = reader.flat_vectors().get(&field_id);
2854 let ann = reader.vector_indexes().get(&field_id);
2855 if ann.is_some() && flat.is_none() {
2856 return Err(Error::Corruption(format!(
2857 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2858 reader.meta().id,
2859 )));
2860 }
2861
2862 let Some(flat) = flat else {
2863 continue;
2864 };
2865 if flat.num_vectors == 0 {
2866 continue;
2867 }
2868 if rewrite_existing {
2869 return Ok(true);
2870 }
2871 let field = crate::dsl::Field(field_id);
2872 let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2873 Error::Corruption(format!(
2874 "segment {:032x} references unknown vector field {field_id}",
2875 reader.meta().id,
2876 ))
2877 })?;
2878 let current = match entry.field_type {
2879 crate::dsl::FieldType::DenseVector
2883 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2884 config.index_type == crate::dsl::VectorIndexType::Tq
2885 }) =>
2886 {
2887 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
2888 }
2889 crate::dsl::FieldType::DenseVector
2890 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2891 config.index_type == crate::dsl::VectorIndexType::IvfTq
2892 }) =>
2893 {
2894 let config = entry
2895 .dense_vector_config
2896 .as_ref()
2897 .expect("matched IVF-TQ configuration");
2898 match (ann, trained.centroids.get(&field_id)) {
2899 (
2900 Some(crate::segment::VectorIndex::IvfTq { index, .. }),
2901 Some(centroids),
2902 ) => {
2903 let header = index.get().header();
2904 crate::structures::is_ivf_tq_cosine_generation(centroids.version)
2905 && crate::structures::is_ivf_tq_cosine_generation(
2906 header.quantizer_version,
2907 )
2908 && header.dim == config.dim
2909 && header.num_clusters == centroids.num_clusters
2910 && header.quantizer_version == centroids.version
2911 && header.codebook_version
2912 == crate::structures::vector::quantization::tq_expected_fingerprint(
2913 config.dim,
2914 )
2915 && header.routing == config.ivf_routing
2916 }
2917 _ => false,
2918 }
2919 }
2920 crate::dsl::FieldType::DenseVector => false,
2923 crate::dsl::FieldType::BinaryDenseVector => {
2924 matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2925 }
2926 _ => false,
2927 };
2928 if !current {
2929 return Ok(true);
2930 }
2931 }
2932 Ok(false)
2933 }
2934
2935 async fn acquire_vector_rewrite_capacity(
2936 &self,
2937 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2938 let global = tokio::select! {
2939 biased;
2940 () = self.active_operations.wait_for_shutdown() => {
2941 return Err(Error::IndexClosed);
2942 }
2943 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2944 permit.map_err(|_| Error::Internal(
2945 "global background merge scheduler is closed".into()
2946 ))?
2947 }
2948 };
2949 let local = tokio::select! {
2950 biased;
2951 () = self.active_operations.wait_for_shutdown() => {
2952 return Err(Error::IndexClosed);
2953 }
2954 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2955 permit.map_err(|_| Error::Internal(
2956 "background merge scheduler is closed".into()
2957 ))?
2958 }
2959 };
2960 Ok((global, local))
2961 }
2962
2963 async fn build_vector_replacement(
2964 self: &Arc<Self>,
2965 segment_id: &str,
2966 source_id: SegmentId,
2967 output_id: SegmentId,
2968 trained: &TrainedVectorStructures,
2969 failure_context: &'static str,
2970 ) -> Result<(String, u32, OutputCleanupGuard)> {
2971 let mut cleanup = self.output_cleanup_guard(output_id);
2972 match crate::segment::reorder::rewrite_vector_segment(
2973 self.directory.as_ref(),
2974 &self.schema,
2975 source_id,
2976 output_id,
2977 self.term_cache_blocks,
2978 trained,
2979 Some(self.background_cpu_pool()),
2980 )
2981 .await
2982 {
2983 Ok((new_id, doc_count)) => {
2984 self.validate_completed_segment(&new_id, doc_count).await?;
2985 Ok((new_id, doc_count, cleanup))
2986 }
2987 Err(error) => {
2988 self.delete_output_if_unregistered(output_id, failure_context)
2989 .await;
2990 cleanup.disarm();
2991 if is_deterministic_source_error(&error) {
2992 self.quarantine_segment(segment_id, &error);
2993 }
2994 Err(error)
2995 }
2996 }
2997 }
2998
2999 pub(crate) async fn stage_vector_generation(
3003 self: &Arc<Self>,
3004 _artifact_update: &VectorArtifactUpdateGuard,
3005 segment_ids: &[String],
3006 field_ids: &[u32],
3007 trained: Arc<TrainedVectorStructures>,
3008 rewrite_existing: bool,
3009 ) -> Result<Vec<StagedVectorSegment>> {
3010 if !self.vector_artifact_update.load(Ordering::Acquire) {
3011 return Err(Error::Internal(
3012 "cannot stage a vector generation without an exclusive update lease".into(),
3013 ));
3014 }
3015
3016 let mut staged = Vec::new();
3017 for segment_id in segment_ids {
3018 if self.quarantined_segments.lock().contains(segment_id) {
3019 return Err(Error::Corruption(format!(
3020 "segment {segment_id} is quarantined after a deterministic source failure"
3021 )));
3022 }
3023 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3024 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3025 })?;
3026
3027 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3030
3031 let output_id = SegmentId::new();
3032 let output_hex = output_id.to_hex();
3033 let operation = {
3034 let st = self.state.lock().await;
3035 if !st.metadata.has_segment(segment_id) {
3036 return Err(Error::Corruption(format!(
3037 "vector generation source {segment_id} disappeared while lifecycle work was paused"
3038 )));
3039 }
3040 self.active_operations
3041 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
3042 }
3043 .ok_or_else(|| {
3044 if self.active_operations.is_accepting() {
3045 Error::Internal(format!(
3046 "vector generation could not claim stable source {segment_id}"
3047 ))
3048 } else {
3049 Error::IndexClosed
3050 }
3051 })?;
3052
3053 let reader = SegmentReader::open(
3054 self.directory.as_ref(),
3055 source_id,
3056 Arc::clone(&self.schema),
3057 self.term_cache_blocks,
3058 )
3059 .await?;
3060 if !self.segment_needs_vector_rewrite(
3061 &reader,
3062 field_ids,
3063 trained.as_ref(),
3064 rewrite_existing,
3065 )? {
3066 continue;
3067 }
3068 drop(reader);
3069
3070 let (new_id, doc_count, cleanup) = self
3071 .build_vector_replacement(
3072 segment_id,
3073 source_id,
3074 output_id,
3075 trained.as_ref(),
3076 "vector generation staging failure",
3077 )
3078 .await?;
3079 debug_assert_eq!(new_id, output_hex);
3080 let output_reader = SegmentReader::open(
3081 self.directory.as_ref(),
3082 output_id,
3083 Arc::clone(&self.schema),
3084 self.term_cache_blocks,
3085 )
3086 .await?;
3087 if self.segment_needs_vector_rewrite(
3088 &output_reader,
3089 field_ids,
3090 trained.as_ref(),
3091 false,
3092 )? {
3093 return Err(Error::Corruption(format!(
3094 "staged vector segment {new_id} does not match its candidate codebook generation"
3095 )));
3096 }
3097
3098 staged.push(StagedVectorSegment {
3099 source_id: segment_id.clone(),
3100 output_id,
3101 doc_count,
3102 _operation: operation,
3103 cleanup,
3104 });
3105 }
3106 Ok(staged)
3107 }
3108
3109 async fn rewrite_vector_segment_once(
3110 self: &Arc<Self>,
3111 segment_id: &str,
3112 field_ids: &[u32],
3113 ) -> Result<VectorSegmentRewriteOutcome> {
3114 if self.quarantined_segments.lock().contains(segment_id) {
3115 return Err(Error::Corruption(format!(
3116 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3117 )));
3118 }
3119 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3120 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3121 })?;
3122
3123 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3128
3129 let output_id = SegmentId::new();
3130 let output_hex = output_id.to_hex();
3131 let all_ids = vec![segment_id.to_owned(), output_hex];
3132 let operation = {
3133 let st = self.state.lock().await;
3134 if !st.metadata.has_segment(segment_id) {
3135 return Ok(VectorSegmentRewriteOutcome::SourceGone);
3136 }
3137 self.active_operations.try_register(all_ids)
3138 };
3139 let _operation = match operation {
3140 Some(operation) => operation,
3141 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3142 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3143 };
3144
3145 let Some(trained) = self.trained_for_segment_build() else {
3146 return Ok(VectorSegmentRewriteOutcome::Deferred);
3147 };
3148
3149 let reader = SegmentReader::open(
3150 self.directory.as_ref(),
3151 source_id,
3152 Arc::clone(&self.schema),
3153 self.term_cache_blocks,
3154 )
3155 .await?;
3156 if !self.segment_needs_vector_rewrite(&reader, field_ids, trained.as_ref(), false)? {
3157 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3158 }
3159 drop(reader);
3160
3161 let (new_id, doc_count, mut output_cleanup) = self
3162 .build_vector_replacement(
3163 segment_id,
3164 source_id,
3165 output_id,
3166 trained.as_ref(),
3167 "vector rewrite failure",
3168 )
3169 .await?;
3170
3171 if let Err(error) = self
3172 .replace_segments(
3173 &[segment_id.to_owned()],
3174 new_id,
3175 doc_count,
3176 ReplacementLayout::PreserveSingleSource,
3177 )
3178 .await
3179 {
3180 self.delete_output_if_unregistered(output_id, "vector replacement failure")
3181 .await;
3182 output_cleanup.disarm();
3183 return Err(error);
3184 }
3185 output_cleanup.disarm();
3186 Ok(VectorSegmentRewriteOutcome::Rewritten)
3187 }
3188
3189 pub(crate) async fn rewrite_vector_segments(
3194 self: &Arc<Self>,
3195 field_ids: &[u32],
3196 ) -> Result<usize> {
3197 if field_ids.is_empty() {
3198 return Ok(0);
3199 }
3200 let mut rewritten = 0usize;
3201 loop {
3202 let segment_ids = self.get_segment_ids().await;
3203 let mut conflicted = false;
3204 let mut changed = false;
3205 for segment_id in segment_ids {
3206 match self
3207 .rewrite_vector_segment_once(&segment_id, field_ids)
3208 .await?
3209 {
3210 VectorSegmentRewriteOutcome::Rewritten => {
3211 rewritten += 1;
3212 changed = true;
3213 }
3214 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3215 VectorSegmentRewriteOutcome::Deferred => {
3216 return Err(Error::Internal(
3217 "ANN finalization lost the published trained generation".into(),
3218 ));
3219 }
3220 VectorSegmentRewriteOutcome::AlreadyCurrent
3221 | VectorSegmentRewriteOutcome::SourceGone => {}
3222 }
3223 }
3224 if !conflicted && !changed {
3225 log::info!(
3226 "[dense_vector_rewrite] index={} ANN finalization complete ({} segment(s) rewritten)",
3227 self.schema.index_label(),
3228 rewritten,
3229 );
3230 return Ok(rewritten);
3231 }
3232 tokio::select! {
3233 biased;
3234 () = self.active_operations.wait_for_shutdown() => {
3235 return Err(Error::IndexClosed);
3236 }
3237 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3238 }
3239 }
3240 }
3241
3242 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3247 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3248 return;
3249 }
3250 let manager = Arc::clone(self);
3251 let future = async move {
3252 let field_ids = manager
3253 .read_metadata(|metadata| {
3254 metadata
3255 .vector_fields
3256 .keys()
3257 .filter(|field_id| metadata.is_field_built(**field_id))
3258 .copied()
3259 .collect::<Vec<_>>()
3260 })
3261 .await;
3262 for segment_id in segment_ids {
3263 loop {
3264 match manager
3265 .rewrite_vector_segment_once(&segment_id, &field_ids)
3266 .await
3267 {
3268 Ok(VectorSegmentRewriteOutcome::Conflict) => {
3269 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3270 }
3271 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3272 Ok(_) => break,
3273 Err(error) => {
3274 log::error!(
3275 "[dense_vector_rewrite] index={} failed to upgrade newly committed segment {}: {}",
3276 manager.schema.index_label(),
3277 segment_id,
3278 error,
3279 );
3280 break;
3281 }
3282 }
3283 }
3284 }
3285 };
3286 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3287 log::warn!(
3288 "[dense_vector_rewrite] index={} runtime unavailable; newly committed flat segment upgrade deferred",
3289 self.schema.index_label()
3290 );
3291 return;
3292 };
3293 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3294 log::warn!(
3295 "[dense_vector_rewrite] index={} runtime rejected newly committed flat segment upgrade",
3296 self.schema.index_label()
3297 );
3298 }
3299 }
3300
3301 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3308 self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3309 .await
3310 }
3311
3312 pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3315 self: &Arc<Self>,
3316 mut refresh_snapshots: F,
3317 ) -> Result<()>
3318 where
3319 F: FnMut() -> Fut,
3320 Fut: std::future::Future<Output = Result<()>>,
3321 {
3322 self.wait_for_all_merges().await;
3323 refresh_snapshots().await?;
3324 let segment_ids = self.get_segment_ids().await;
3325
3326 if segment_ids.is_empty() {
3327 log::info!(
3328 "[reorder] index={} no segments to reorder",
3329 self.schema.index_label()
3330 );
3331 return Ok(());
3332 }
3333
3334 log::info!(
3335 "[reorder] index={} reordering {} segments",
3336 self.schema.index_label(),
3337 segment_ids.len()
3338 );
3339
3340 for seg_id in segment_ids {
3341 match self
3342 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3343 .await
3344 {
3345 Ok(true) => refresh_snapshots().await?,
3346 Ok(false) => log::warn!(
3347 "[reorder] index={} segment {} skipped (in merge)",
3348 self.schema.index_label(),
3349 seg_id
3350 ),
3351 Err(e) => return Err(e),
3352 }
3353 }
3354
3355 refresh_snapshots().await?;
3358 log::info!(
3359 "[reorder] index={} all segments reordered",
3360 self.schema.index_label()
3361 );
3362 Ok(())
3363 }
3364
3365 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3370 self.unreordered_segments()
3371 .await
3372 .into_iter()
3373 .map(|(id, _)| id)
3374 .collect()
3375 }
3376
3377 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3380 let quarantined = self.quarantined_segments.lock().clone();
3381 let paused = self.paused_reorder_segments();
3382 let st = self.state.lock().await;
3383 let active_ids = self.active_operations.snapshot();
3384 st.metadata
3385 .segment_metas
3386 .iter()
3387 .filter(|(id, info)| {
3388 !info.reordered
3389 && info.bp_converged
3390 && !active_ids.contains(*id)
3391 && !quarantined.contains(*id)
3392 && !paused.contains(*id)
3393 })
3394 .map(|(id, info)| (id.clone(), info.num_docs))
3395 .collect()
3396 }
3397
3398 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3402 self.unconverged_segments_below(u32::MAX)
3403 .await
3404 .into_iter()
3405 .map(|(id, docs, _)| (id, docs))
3406 .collect()
3407 }
3408
3409 pub async fn unconverged_segments_below(
3412 &self,
3413 max_unconverged_passes: u32,
3414 ) -> Vec<(String, u32, u32)> {
3415 let quarantined = self.quarantined_segments.lock().clone();
3416 let paused = self.paused_reorder_segments();
3417 let st = self.state.lock().await;
3418 let active_ids = self.active_operations.snapshot();
3419 st.metadata
3420 .segment_metas
3421 .iter()
3422 .filter(|(id, info)| {
3423 !info.bp_converged
3424 && info.bp_unconverged_passes < max_unconverged_passes
3425 && !active_ids.contains(*id)
3426 && !quarantined.contains(*id)
3427 && !paused.contains(*id)
3428 })
3429 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3430 .collect()
3431 }
3432
3433 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3443 let st = self.state.lock().await;
3444 let deepening = ids.iter().any(|id| {
3445 st.metadata
3446 .segment_metas
3447 .get(id)
3448 .is_some_and(|info| !info.bp_converged)
3449 });
3450 drop(st);
3451 if deepening {
3452 log::info!(
3453 "[reorder] index={} source BP lineage unconverged — forcing record-level BP (deepening pass)",
3454 self.schema.index_label(),
3455 );
3456 crate::segment::reorder::BpGranularity::Records
3457 } else {
3458 crate::segment::reorder::BpGranularity::Auto
3459 }
3460 }
3461
3462 pub async fn reorder_single_segment(
3467 self: &Arc<Self>,
3468 seg_id: &str,
3469 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3470 bp_budget: crate::segment::BpBudget,
3471 ) -> Result<bool> {
3472 let source_id = SegmentId::from_hex(seg_id)
3473 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3474 if self.quarantined_segments.lock().contains(seg_id) {
3475 return Err(Error::Corruption(format!(
3476 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3477 seg_id
3478 )));
3479 }
3480 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3481 log::debug!(
3482 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3483 self.schema.index_label(),
3484 seg_id,
3485 );
3486 return Ok(false);
3487 }
3488
3489 let reorder_gate = Arc::clone(&self.reorder_permits);
3494 let _reorder_permit = tokio::select! {
3495 biased;
3496 () = self.active_operations.wait_for_shutdown() => {
3497 return Err(Error::IndexClosed);
3498 }
3499 permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3500 permit.map_err(|_| {
3501 Error::Internal("background reorder scheduler is closed".into())
3502 })?
3503 }
3504 };
3505
3506 let output_id = SegmentId::new();
3507 let output_hex = output_id.to_hex();
3508 let source_ids = [seg_id.to_string()];
3509 let granularity = self.merge_granularity(&source_ids).await;
3510
3511 let all_ids = vec![seg_id.to_string(), output_hex];
3517 let (_guard, source_docs) = {
3518 let st = self.state.lock().await;
3519 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3523 log::debug!(
3524 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3525 self.schema.index_label(),
3526 seg_id,
3527 );
3528 return Ok(false);
3529 }
3530 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3531 log::info!(
3532 "[optimizer] index={} segment {} no longer in metadata (merged away), skipping reorder",
3533 self.schema.index_label(),
3534 seg_id
3535 );
3536 self.clear_reorder_retry(seg_id);
3537 return Ok(false);
3538 };
3539
3540 match self.active_operations.try_register(all_ids) {
3541 Some(guard) => (guard, source_meta.num_docs),
3542 None if !self.active_operations.is_accepting() => {
3543 return Err(Error::IndexClosed);
3544 }
3545 None => {
3546 log::debug!(
3547 "[optimizer] index={} segment {} in active merge, skipping",
3548 self.schema.index_label(),
3549 seg_id
3550 );
3551 return Ok(false);
3552 }
3553 }
3554 };
3555
3556 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3561 if is_deterministic_source_error(&error) {
3562 self.quarantine_segment(seg_id, &error);
3563 } else if !matches!(&error, Error::IndexClosed) {
3564 self.pause_reorder_retries(seg_id, &error);
3565 }
3566 return Err(error);
3567 }
3568
3569 let mut output_cleanup = self.output_cleanup_guard(output_id);
3570
3571 let reorder_result = crate::segment::reorder::reorder_segment(
3572 self.directory.as_ref(),
3573 &self.schema,
3574 source_id,
3575 output_id,
3576 self.term_cache_blocks,
3577 self.bp_memory_budget_bytes,
3578 bp_budget,
3579 granularity,
3580 rayon_pool,
3581 Some(self.active_operations.cancellation_flag()),
3582 self.trained_for_segment_build(),
3583 )
3584 .await;
3585 let (new_id, total_docs, bp_converged) = match reorder_result {
3586 Ok(v) => v,
3587 Err(e) => {
3588 self.delete_output_if_unregistered(output_id, "reorder failure")
3591 .await;
3592 output_cleanup.disarm();
3593 if is_deterministic_source_error(&e) {
3594 self.quarantine_segment(seg_id, &e);
3595 } else if !matches!(&e, Error::IndexClosed) {
3596 self.pause_reorder_retries(seg_id, &e);
3597 }
3598 return Err(e);
3599 }
3600 };
3601
3602 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3608 if let Err(e) = self
3609 .replace_segments(
3610 &[seg_id.to_string()],
3611 new_id,
3612 total_docs,
3613 ReplacementLayout::BpReordered {
3614 converged: ladder_converged,
3615 },
3616 )
3617 .await
3618 {
3619 self.delete_output_if_unregistered(output_id, "replacement failure")
3620 .await;
3621 output_cleanup.disarm();
3622 if !matches!(&e, Error::IndexClosed) {
3623 self.pause_reorder_retries(seg_id, &e);
3624 }
3625 return Err(e);
3626 }
3627 output_cleanup.disarm();
3628 self.clear_reorder_retry(seg_id);
3629
3630 Ok(true)
3631 }
3632
3633 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3640 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3641
3642 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3643 for entry in entries {
3644 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3645 continue;
3646 };
3647 let Some(rest) = filename.strip_prefix("seg_") else {
3648 continue;
3649 };
3650 let Some(hex_id) = rest.get(..32) else {
3651 continue;
3652 };
3653 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3654 continue;
3655 }
3656 orphan_files
3657 .entry(hex_id.to_ascii_lowercase())
3658 .or_default()
3659 .push(entry);
3660 }
3661 }
3662
3663 let mut deleted = 0;
3664 for (hex_id, paths) in &orphan_files {
3665 let deletion_guard = {
3670 let st = self.state.lock().await;
3671 if st.metadata.has_segment(hex_id) {
3672 continue;
3673 }
3674 let Some(guard) = self
3675 .active_operations
3676 .try_register(vec![hex_id.to_string()])
3677 else {
3678 continue;
3679 };
3680 if self.tracker.is_deletion_protected(hex_id) {
3681 drop(guard);
3682 continue;
3683 }
3684 guard
3685 };
3686
3687 let results =
3692 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3693 .await;
3694 let removed = results.into_iter().all(|result| match result {
3695 Ok(()) => true,
3696 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3697 Err(error) => {
3698 log::warn!(
3699 "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
3700 self.schema.index_label(),
3701 hex_id,
3702 error,
3703 );
3704 false
3705 }
3706 });
3707 drop(deletion_guard);
3710 if removed {
3711 deleted += 1;
3712 log::info!(
3713 "[segment_cleanup] index={} swept orphan segment {}",
3714 self.schema.index_label(),
3715 hex_id
3716 );
3717 }
3718 }
3719
3720 Ok(deleted)
3721 }
3722}
3723
3724#[cfg(test)]
3725mod tests {
3726 use super::*;
3727 use std::sync::atomic::{AtomicBool, Ordering};
3728
3729 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3730 let schema = crate::dsl::SchemaBuilder::default().build();
3731 let metadata = IndexMetadata::new(schema.clone());
3732 Arc::new(SegmentManager::new(
3733 Arc::new(crate::directories::RamDirectory::new()),
3734 Arc::new(schema),
3735 metadata,
3736 Box::new(crate::merge::NoMergePolicy),
3737 0,
3738 1,
3739 Arc::new(Semaphore::new(1)),
3740 None,
3741 1024,
3742 Arc::new(ReorderConcurrencyGate::new(1)),
3743 None,
3744 ))
3745 }
3746
3747 #[test]
3748 fn force_merge_planner_pairs_large_and_small_segments() {
3749 let groups = plan_force_merge_groups(
3750 vec![
3751 ("a".into(), 6),
3752 ("b".into(), 6),
3753 ("c".into(), 4),
3754 ("d".into(), 4),
3755 ],
3756 10,
3757 );
3758
3759 assert_eq!(groups.len(), 2);
3760 assert!(groups.iter().all(|group| group.total_docs == 10));
3761 assert!(groups.iter().all(|group| group.segments.len() == 2));
3762 }
3763
3764 #[test]
3765 fn force_merge_planner_leaves_oversized_segments_alone() {
3766 let groups = plan_force_merge_groups(
3767 vec![
3768 ("oversized".into(), 11),
3769 ("small-a".into(), 5),
3770 ("small-b".into(), 5),
3771 ],
3772 10,
3773 );
3774
3775 assert_eq!(groups.len(), 2);
3776 assert_eq!(groups[0].total_docs, 10);
3777 assert_eq!(groups[0].segments.len(), 2);
3778 assert_eq!(groups[1].total_docs, 11);
3779 assert_eq!(groups[1].segments.len(), 1);
3780 }
3781
3782 #[test]
3783 fn force_merge_planner_never_exceeds_segment_format_limit() {
3784 let groups = plan_force_merge_groups(
3785 vec![
3786 ("large-a".into(), 3_000_000_000),
3787 ("large-b".into(), 2_000_000_000),
3788 ],
3789 u64::from(u32::MAX),
3790 );
3791 assert_eq!(groups.len(), 2);
3792 assert!(
3793 groups
3794 .iter()
3795 .all(|group| group.total_docs <= u64::from(u32::MAX))
3796 );
3797 }
3798
3799 #[test]
3800 fn force_merge_hierarchy_has_one_final_bp_pass() {
3801 assert_eq!(force_merge_output_count(1), 0);
3802 assert_eq!(force_merge_output_count(2), 1);
3803 assert_eq!(force_merge_output_count(64), 1);
3804 assert_eq!(force_merge_output_count(65), 2);
3805 assert_eq!(force_merge_output_count(127), 2);
3806 assert_eq!(force_merge_output_count(128), 3);
3807 assert_eq!(force_merge_output_count(1_000), 16);
3808 }
3809
3810 fn expand_force_merge_node(
3811 hierarchy: &ForceMergeHierarchy,
3812 source_count: usize,
3813 node: usize,
3814 sources: &mut Vec<usize>,
3815 ) {
3816 if node < source_count {
3817 sources.push(node);
3818 return;
3819 }
3820
3821 let step_index = node - source_count;
3822 let step = hierarchy
3823 .steps
3824 .get(step_index)
3825 .expect("merge input must refer to an existing source or output");
3826 for &input in &step.inputs {
3827 assert!(
3828 input < node,
3829 "merge step {step_index} refers to a future output node {input}"
3830 );
3831 expand_force_merge_node(hierarchy, source_count, input, sources);
3832 }
3833 }
3834
3835 #[test]
3836 fn force_merge_hierarchy_has_minimal_valid_arity() {
3837 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
3838
3839 for source_count in source_counts {
3840 let hierarchy = plan_force_merge_hierarchy(source_count);
3841 let output_count = hierarchy.steps.len();
3842
3843 assert!(
3844 hierarchy
3845 .steps
3846 .iter()
3847 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
3848 "invalid merge arity for {source_count} sources"
3849 );
3850 assert!(
3851 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
3852 "{output_count} outputs cannot reduce {source_count} sources"
3853 );
3854 assert!(
3855 output_count == 1
3856 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
3857 "{output_count} outputs are not minimal for {source_count} sources"
3858 );
3859 assert_eq!(output_count, force_merge_output_count(source_count));
3860 }
3861 }
3862
3863 #[test]
3864 fn force_merge_hierarchy_preserves_exact_source_order() {
3865 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
3866 let hierarchy = plan_force_merge_hierarchy(source_count);
3867 let mut sources = Vec::with_capacity(source_count);
3868 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
3869 assert_eq!(
3870 sources,
3871 (0..source_count).collect::<Vec<_>>(),
3872 "source order changed for {source_count} sources"
3873 );
3874 }
3875 }
3876
3877 fn force_merge_rewrite_cost(source_count: usize) -> usize {
3878 let hierarchy = plan_force_merge_hierarchy(source_count);
3879 let mut node_weights = vec![1usize; source_count];
3880 let mut rewrite_cost = 0usize;
3881
3882 for (step_index, step) in hierarchy.steps.iter().enumerate() {
3883 let output = source_count + step_index;
3884 let output_weight = step
3885 .inputs
3886 .iter()
3887 .map(|&input| {
3888 assert!(
3889 input < output,
3890 "merge step {step_index} refers to future output {input}"
3891 );
3892 node_weights[input]
3893 })
3894 .sum::<usize>();
3895 rewrite_cost += output_weight;
3896 node_weights.push(output_weight);
3897 }
3898
3899 assert_eq!(node_weights[hierarchy.root], source_count);
3900 rewrite_cost
3901 }
3902
3903 #[test]
3904 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
3905 assert_eq!(force_merge_rewrite_cost(65), 67);
3906 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
3907 }
3908
3909 #[test]
3910 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
3911 assert_eq!(
3912 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
3913 (false, false, 3),
3914 );
3915 assert_eq!(
3916 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
3917 (true, false, 4),
3918 );
3919 assert_eq!(
3920 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
3921 (true, true, 0),
3922 );
3923 }
3924
3925 #[tokio::test]
3926 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
3927 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3928 let field = schema_builder.add_text_field("text", true, true);
3929 let schema = schema_builder.build();
3930 let directory = crate::directories::RamDirectory::new();
3931 let config = crate::index::IndexConfig {
3932 num_indexing_threads: 1,
3933 merge_policy: Box::new(crate::merge::NoMergePolicy),
3934 ..Default::default()
3935 };
3936 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
3937 .await
3938 .unwrap();
3939 for value in ["one", "two", "three"] {
3940 let mut document = crate::dsl::Document::new();
3941 document.add_text(field, value);
3942 writer.add_document(document).unwrap();
3943 writer.commit().await.unwrap();
3944 }
3945
3946 let manager = Arc::clone(writer.segment_manager());
3947 let held_id = manager.get_segment_ids().await.pop().unwrap();
3948 let mut held = Some(
3949 manager
3950 .active_operations
3951 .try_register(vec![held_id])
3952 .unwrap(),
3953 );
3954 let batches = Arc::new(AtomicUsize::new(0));
3955 let batch_count = Arc::clone(&batches);
3956 writer
3957 .force_merge_with_snapshot_refresh(move || {
3958 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
3959 if refresh == 2 {
3962 drop(held.take());
3963 }
3964 std::future::ready(Ok(()))
3965 })
3966 .await
3967 .unwrap();
3968
3969 assert_eq!(manager.get_segment_ids().await.len(), 1);
3970 assert_eq!(
3971 batches.load(Ordering::Relaxed),
3972 4,
3973 "initial/final refreshes plus two replacements are required after the held source releases"
3974 );
3975 }
3976
3977 #[tokio::test]
3978 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
3979 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3980 schema_builder.set_reorder_on_merge(true);
3981 let schema = schema_builder.build();
3982 let mut metadata = IndexMetadata::new(schema.clone());
3983 metadata.add_segment("00000000000000000000000000000001".into(), 1);
3984 metadata.add_segment("00000000000000000000000000000002".into(), 1);
3985
3986 let global_merge_permits = Arc::new(Semaphore::new(1));
3987 let manager = Arc::new(SegmentManager::new(
3988 Arc::new(crate::directories::RamDirectory::new()),
3989 Arc::new(schema),
3990 metadata,
3991 Box::new(crate::merge::NoMergePolicy),
3992 0,
3993 1,
3994 Arc::clone(&global_merge_permits),
3995 None,
3996 1024,
3997 Arc::new(ReorderConcurrencyGate::new(1)),
3998 None,
3999 ));
4000
4001 manager.active_operations.pause_non_indexing();
4006 let force_merge = {
4007 let manager = Arc::clone(&manager);
4008 tokio::spawn(async move { manager.force_merge().await })
4009 };
4010 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4011 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4012 tokio::task::yield_now().await;
4013 }
4014 })
4015 .await
4016 .expect("force merge never reached the paused group claim");
4017
4018 assert_eq!(
4019 global_merge_permits.available_permits(),
4020 1,
4021 "force merge retained global capacity while vector staging blocked group ownership"
4022 );
4023
4024 force_merge.abort();
4025 let _ = force_merge.await;
4026 manager.active_operations.resume_non_indexing();
4027 }
4028
4029 #[test]
4030 fn output_cleanup_guard_runs_during_panic_unwind() {
4031 let cleaned = Arc::new(AtomicBool::new(false));
4032 let cleaned_in_callback = Arc::clone(&cleaned);
4033 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4034 cleaned_in_callback.store(true, Ordering::SeqCst);
4035 });
4036
4037 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4038 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4039 panic!("simulated reorder panic");
4040 }));
4041
4042 assert!(result.is_err());
4043 assert!(
4044 cleaned.load(Ordering::SeqCst),
4045 "partial output cleanup must run during unwind"
4046 );
4047 }
4048
4049 #[test]
4050 fn output_cleanup_guard_disarms_after_commit() {
4051 let cleaned = Arc::new(AtomicBool::new(false));
4052 let cleaned_in_callback = Arc::clone(&cleaned);
4053 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4054 cleaned_in_callback.store(true, Ordering::SeqCst);
4055 });
4056
4057 {
4058 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4059 guard.disarm();
4060 }
4061
4062 assert!(!cleaned.load(Ordering::SeqCst));
4063 }
4064
4065 #[test]
4066 fn test_active_operation_guard_releases_ownership() {
4067 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4068 {
4069 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4070 let snap = active.snapshot();
4071 assert!(snap.contains("a"));
4072 assert!(snap.contains("b"));
4073 }
4074 assert!(active.snapshot().is_empty());
4075 }
4076
4077 #[test]
4078 fn test_non_overlapping_operations_can_run_concurrently() {
4079 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4080 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4081 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4082 let snap = active.snapshot();
4083 assert_eq!(snap.len(), 4);
4084
4085 drop(first);
4086 let snap = active.snapshot();
4087 assert_eq!(snap.len(), 2);
4088 assert!(snap.contains("c"));
4089 assert!(snap.contains("d"));
4090 }
4091
4092 #[test]
4093 fn test_overlapping_operation_is_rejected_until_release() {
4094 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4095 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4096 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4097 drop(first);
4098 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4099 }
4100
4101 #[test]
4102 fn test_active_operation_snapshot() {
4103 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4104 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4105 let snap = active.snapshot();
4106 assert!(snap.contains("x"));
4107 assert!(snap.contains("y"));
4108 assert!(!snap.contains("z"));
4109 }
4110
4111 #[tokio::test]
4112 async fn operation_barrier_ignores_producers_started_after_snapshot() {
4113 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4114 let before_gate = active.try_register(vec!["old".into()]).unwrap();
4115 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4116 assert_eq!(parked_indexing, 0);
4117 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4118
4119 let waiter = {
4120 let active = Arc::clone(&active);
4121 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4122 };
4123 tokio::task::yield_now().await;
4124 assert!(!waiter.is_finished());
4125
4126 drop(before_gate);
4127 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4128 .await
4129 .expect("pre-gate operation barrier was starved by a post-gate producer")
4130 .unwrap();
4131 assert!(active.snapshot().contains("new-flat"));
4132 drop(after_gate);
4133 }
4134
4135 #[tokio::test]
4136 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4137 let manager = lifecycle_test_manager();
4138 manager
4139 .trained
4140 .store(Some(Arc::new(TrainedVectorStructures {
4141 centroids: rustc_hash::FxHashMap::default(),
4142 binary_quantizers: rustc_hash::FxHashMap::default(),
4143 ..Default::default()
4144 })));
4145
4146 let guard = manager.begin_vector_artifact_update().await.unwrap();
4147 assert!(
4148 manager.trained().is_some(),
4149 "search readers keep the last fully validated generation"
4150 );
4151 assert!(
4152 manager.trained_for_segment_build().is_none(),
4153 "new segment producers must stay flat during an artifact update"
4154 );
4155
4156 let detached_transaction_guard = guard.clone();
4157 drop(guard);
4158 assert!(
4159 manager.trained_for_segment_build().is_none(),
4160 "a detached lifecycle transaction must retain the producer gate after request cancellation"
4161 );
4162 drop(detached_transaction_guard);
4163 assert!(manager.trained_for_segment_build().is_some());
4164 }
4165
4166 #[tokio::test]
4167 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4168 let manager = lifecycle_test_manager();
4169 let guard = manager.begin_vector_artifact_update().await.unwrap();
4170 assert!(
4171 manager
4172 .active_operations
4173 .try_register(vec!["merge".into()])
4174 .is_none(),
4175 "ordinary merge/reorder work must not change staged sources"
4176 );
4177 let indexing = manager
4178 .active_operations
4179 .try_register_indexing(vec!["fresh".into()])
4180 .expect("indexing remains available in flat mode");
4181 drop(indexing);
4182
4183 drop(guard);
4184 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4185 assert!(
4186 manager
4187 .active_operations
4188 .try_register(vec!["merge".into()])
4189 .is_some()
4190 );
4191 }
4192
4193 #[tokio::test]
4194 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4195 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4196 let guard = active.try_register(vec!["live".into()]).unwrap();
4197 let cancellation = active.cancellation_flag();
4198 active.stop_accepting();
4199 assert!(cancellation.load(Ordering::Acquire));
4200 assert!(active.try_register(vec!["new".into()]).is_none());
4201
4202 let waiter = {
4203 let active = Arc::clone(&active);
4204 tokio::spawn(async move { active.wait_until_idle().await })
4205 };
4206 tokio::task::yield_now().await;
4207 assert!(!waiter.is_finished());
4208 drop(guard);
4209 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4210 .await
4211 .expect("shutdown waiter missed the final guard notification")
4212 .unwrap();
4213 }
4214
4215 #[tokio::test]
4216 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4217 let manager = lifecycle_test_manager();
4218 let started = Arc::new(Semaphore::new(0));
4219 let release = Arc::new(Semaphore::new(0));
4220 let completed = Arc::new(AtomicBool::new(false));
4221
4222 let request = {
4223 let manager = Arc::clone(&manager);
4224 let started = Arc::clone(&started);
4225 let release = Arc::clone(&release);
4226 let completed = Arc::clone(&completed);
4227 tokio::spawn(async move {
4228 manager
4229 .run_lifecycle_transaction(async move {
4230 started.add_permits(1);
4231 let _permit = release.acquire().await.unwrap();
4232 completed.store(true, Ordering::Release);
4233 Ok(())
4234 })
4235 .await
4236 })
4237 };
4238
4239 let _started = started.acquire().await.unwrap();
4240 request.abort();
4241 assert!(request.await.unwrap_err().is_cancelled());
4242 release.add_permits(1);
4243
4244 manager.begin_shutdown();
4245 tokio::time::timeout(
4246 std::time::Duration::from_secs(1),
4247 manager.wait_for_shutdown(),
4248 )
4249 .await
4250 .expect("shutdown did not drain detached lifecycle transaction");
4251 assert!(completed.load(Ordering::Acquire));
4252 }
4253
4254 #[tokio::test]
4255 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4256 let manager = lifecycle_test_manager();
4257 {
4258 let mut state = manager.state.lock().await;
4259 state.metadata.add_segment_meta(
4260 "eligible".into(),
4261 SegmentMetaInfo {
4262 num_docs: 10,
4263 ancestors: Vec::new(),
4264 generation: 1,
4265 reordered: true,
4266 bp_converged: false,
4267 bp_unconverged_passes: 2,
4268 },
4269 );
4270 state.metadata.add_segment_meta(
4271 "at-limit".into(),
4272 SegmentMetaInfo {
4273 num_docs: 20,
4274 ancestors: Vec::new(),
4275 generation: 1,
4276 reordered: true,
4277 bp_converged: false,
4278 bp_unconverged_passes: 3,
4279 },
4280 );
4281 state.metadata.add_segment_meta(
4282 "carried-debt".into(),
4283 SegmentMetaInfo {
4284 num_docs: 15,
4285 ancestors: Vec::new(),
4286 generation: 2,
4287 reordered: false,
4288 bp_converged: false,
4289 bp_unconverged_passes: 2,
4290 },
4291 );
4292 state.metadata.add_segment_meta(
4293 "carried-debt-at-limit".into(),
4294 SegmentMetaInfo {
4295 num_docs: 25,
4296 ancestors: Vec::new(),
4297 generation: 2,
4298 reordered: false,
4299 bp_converged: false,
4300 bp_unconverged_passes: 3,
4301 },
4302 );
4303 state.metadata.add_segment_meta(
4304 "converged".into(),
4305 SegmentMetaInfo {
4306 num_docs: 30,
4307 ancestors: Vec::new(),
4308 generation: 1,
4309 reordered: true,
4310 bp_converged: true,
4311 bp_unconverged_passes: 0,
4312 },
4313 );
4314 state.metadata.add_segment("fresh".into(), 40);
4315 }
4316
4317 assert_eq!(
4318 manager.unreordered_segments().await,
4319 vec![("fresh".into(), 40)],
4320 "a block-copy output with BP debt is not a fresh first-pass candidate",
4321 );
4322 let mut eligible = manager.unconverged_segments_below(3).await;
4323 eligible.sort_unstable();
4324 assert_eq!(
4325 eligible,
4326 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4327 );
4328 assert!(manager.unconverged_segments_below(0).await.is_empty());
4329 }
4330
4331 #[test]
4332 fn merge_retry_backoff_is_exponential_and_capped() {
4333 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4334 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4335 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4336 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4337 }
4338
4339 #[test]
4340 fn only_deterministic_source_errors_are_quarantined() {
4341 assert!(is_deterministic_source_error(&Error::Corruption(
4342 "bad footer".into()
4343 )));
4344 assert!(is_deterministic_source_error(&Error::Io(
4345 std::io::Error::from(std::io::ErrorKind::NotFound)
4346 )));
4347 assert!(!is_deterministic_source_error(&Error::Io(
4348 std::io::Error::from(std::io::ErrorKind::TimedOut)
4349 )));
4350 assert!(!is_deterministic_source_error(&Error::Io(
4351 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4352 )));
4353 }
4354
4355 #[test]
4356 fn transient_reorder_failure_is_backed_off_until_cleared() {
4357 let manager = lifecycle_test_manager();
4358 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4359 assert!(manager.paused_reorder_segments().contains("source"));
4360 manager.clear_reorder_retry("source");
4361 assert!(!manager.paused_reorder_segments().contains("source"));
4362 }
4363
4364 #[derive(Default)]
4367 struct FailingExistsDirectory(crate::directories::RamDirectory);
4368
4369 #[async_trait::async_trait]
4370 impl crate::directories::Directory for FailingExistsDirectory {
4371 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4372 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4373 }
4374
4375 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4376 self.0.file_size(path).await
4377 }
4378
4379 async fn open_read(
4380 &self,
4381 path: &std::path::Path,
4382 ) -> std::io::Result<crate::directories::FileHandle> {
4383 self.0.open_read(path).await
4384 }
4385
4386 async fn read_range(
4387 &self,
4388 path: &std::path::Path,
4389 range: std::ops::Range<u64>,
4390 ) -> std::io::Result<crate::directories::OwnedBytes> {
4391 self.0.read_range(path, range).await
4392 }
4393
4394 async fn list_files(
4395 &self,
4396 prefix: &std::path::Path,
4397 ) -> std::io::Result<Vec<std::path::PathBuf>> {
4398 self.0.list_files(prefix).await
4399 }
4400
4401 async fn open_lazy(
4402 &self,
4403 path: &std::path::Path,
4404 ) -> std::io::Result<crate::directories::FileHandle> {
4405 self.0.open_lazy(path).await
4406 }
4407 }
4408
4409 #[async_trait::async_trait]
4410 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4411 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4412 self.0.write(path, data).await
4413 }
4414
4415 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4416 self.0.delete(path).await
4417 }
4418
4419 async fn rename(
4420 &self,
4421 from: &std::path::Path,
4422 to: &std::path::Path,
4423 ) -> std::io::Result<()> {
4424 self.0.rename(from, to).await
4425 }
4426
4427 async fn sync(&self) -> std::io::Result<()> {
4428 self.0.sync().await
4429 }
4430
4431 async fn streaming_writer(
4432 &self,
4433 path: &std::path::Path,
4434 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4435 self.0.streaming_writer(path).await
4436 }
4437 }
4438
4439 #[derive(Debug, Clone)]
4440 struct MergeEverythingPolicy;
4441
4442 impl MergePolicy for MergeEverythingPolicy {
4443 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4444 if segments.len() < 2 {
4445 return Vec::new();
4446 }
4447 vec![crate::merge::MergeCandidate {
4448 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4449 }]
4450 }
4451
4452 fn clone_box(&self) -> Box<dyn MergePolicy> {
4453 Box::new(self.clone())
4454 }
4455 }
4456
4457 #[tokio::test]
4458 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4459 let manager = lifecycle_test_manager();
4460 let parked_indexing = manager
4465 .protect_new_segment("00000000000000000000000000000abc".into())
4466 .unwrap();
4467
4468 let error = tokio::time::timeout(
4469 std::time::Duration::from_secs(2),
4470 manager.begin_vector_artifact_update(),
4471 )
4472 .await
4473 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4474 .err()
4475 .expect("an old-generation prepared segment must block artifact replacement")
4476 .to_string();
4477 assert!(error.contains("built but uncommitted"), "{error}");
4478 assert!(
4479 !manager.vector_artifact_update.load(Ordering::Acquire),
4480 "a rejected update must release the producer gate"
4481 );
4482
4483 drop(parked_indexing);
4484
4485 let guard = manager
4486 .begin_vector_artifact_update()
4487 .await
4488 .expect("artifact update should succeed after the pending generation is resolved");
4489 drop(guard);
4490 }
4491
4492 #[tokio::test]
4493 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4494 let manager = lifecycle_test_manager();
4495 let merge_like = manager
4496 .active_operations
4497 .try_register(vec!["merge-source".into()])
4498 .unwrap();
4499
4500 let waiter = {
4501 let manager = Arc::clone(&manager);
4502 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4503 };
4504 for _ in 0..8 {
4505 tokio::task::yield_now().await;
4506 }
4507 assert!(
4508 !waiter.is_finished(),
4509 "artifact update must drain merge/reorder producers that may hold the previous generation"
4510 );
4511
4512 drop(merge_like);
4513 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4514 .await
4515 .expect("artifact update missed the lifecycle guard release")
4516 .unwrap()
4517 .unwrap();
4518 }
4519
4520 #[tokio::test]
4521 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4522 let manager = lifecycle_test_manager();
4523 let release = Arc::new(Semaphore::new(0));
4524 let merge_task = {
4525 let release = Arc::clone(&release);
4526 tokio::spawn(async move {
4527 let _permit = release.acquire().await.unwrap();
4528 })
4529 };
4530 manager.merge_handles.lock().push(merge_task);
4531
4532 let waiter = {
4533 let manager = Arc::clone(&manager);
4534 tokio::spawn(async move { manager.wait_for_all_merges().await })
4535 };
4536 for _ in 0..8 {
4537 tokio::task::yield_now().await;
4538 }
4539 assert!(!waiter.is_finished());
4540 waiter.abort();
4543 let join_error = waiter.await.unwrap_err();
4544 assert!(join_error.is_cancelled());
4545
4546 assert!(
4547 !manager.merge_handles.lock().is_empty(),
4548 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4549 );
4550
4551 release.add_permits(1);
4553 tokio::time::timeout(
4554 std::time::Duration::from_secs(1),
4555 manager.wait_for_all_merges(),
4556 )
4557 .await
4558 .expect("subsequent drain missed the reinserted merge handle");
4559 assert!(manager.merge_handles.lock().is_empty());
4560 }
4561
4562 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4563 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4564 let manager = lifecycle_test_manager();
4565 {
4566 let mut state = manager.state.lock().await;
4567 state
4568 .metadata
4569 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4570 state
4571 .metadata
4572 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4573 }
4574 let reorder_like = manager
4577 .active_operations
4578 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4579 .unwrap();
4580
4581 let force_merge = {
4582 let manager = Arc::clone(&manager);
4583 tokio::spawn(async move { manager.force_merge().await })
4584 };
4585
4586 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4587 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4588 assert!(
4589 retries >= 1,
4590 "force_merge never observed the conflicting owner (retries={retries})"
4591 );
4592 assert!(
4593 retries < 20,
4594 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4595 );
4596
4597 drop(reorder_like);
4598 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4601 .await
4602 .expect("force_merge kept spinning after the conflicting owner released")
4603 .unwrap();
4604 assert!(result.is_err());
4605 }
4606
4607 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4608 async fn force_merge_routes_around_segments_held_by_reorder() {
4609 let manager = lifecycle_test_manager();
4610 {
4611 let mut state = manager.state.lock().await;
4612 state
4613 .metadata
4614 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4615 state
4616 .metadata
4617 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4618 state
4619 .metadata
4620 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4621 }
4622 let _reorder_like = manager
4625 .active_operations
4626 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4627 .unwrap();
4628
4629 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4636 let manager = Arc::clone(&manager);
4637 async move { manager.force_merge().await }
4638 })
4639 .await
4640 .expect("force_merge livelocked on a segment held by an active reorder");
4641 assert!(result.is_err(), "fake segment files must fail the merge");
4642
4643 assert_eq!(
4644 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4645 0,
4646 "batch built from the ownership snapshot must not collide with the held segment"
4647 );
4648 }
4649
4650 #[tokio::test]
4651 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4652 let schema = crate::dsl::SchemaBuilder::default().build();
4653 let mut metadata = IndexMetadata::new(schema.clone());
4654 metadata.add_segment("00000000000000000000000000000001".into(), 10);
4655 metadata.add_segment("00000000000000000000000000000002".into(), 10);
4656 let manager = Arc::new(SegmentManager::new(
4657 Arc::new(FailingExistsDirectory::default()),
4658 Arc::new(schema),
4659 metadata,
4660 Box::new(MergeEverythingPolicy),
4661 0,
4662 1,
4663 Arc::new(Semaphore::new(1)),
4664 None,
4665 1024,
4666 Arc::new(ReorderConcurrencyGate::new(1)),
4667 None,
4668 ));
4669
4670 manager.maybe_merge().await;
4673
4674 tokio::time::timeout(
4675 std::time::Duration::from_secs(5),
4676 manager.wait_for_all_merges(),
4677 )
4678 .await
4679 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4680 assert!(
4681 manager.merge_retry_is_paused(),
4682 "the failed merge should have armed the retry backoff"
4683 );
4684
4685 manager.begin_shutdown();
4687 tokio::time::timeout(
4688 std::time::Duration::from_secs(5),
4689 manager.wait_for_shutdown(),
4690 )
4691 .await
4692 .expect("shutdown did not drain the merge retry wakeup task");
4693 }
4694}