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 )
3583 .await;
3584 let (new_id, total_docs, bp_converged) = match reorder_result {
3585 Ok(v) => v,
3586 Err(e) => {
3587 self.delete_output_if_unregistered(output_id, "reorder failure")
3590 .await;
3591 output_cleanup.disarm();
3592 if is_deterministic_source_error(&e) {
3593 self.quarantine_segment(seg_id, &e);
3594 } else if !matches!(&e, Error::IndexClosed) {
3595 self.pause_reorder_retries(seg_id, &e);
3596 }
3597 return Err(e);
3598 }
3599 };
3600
3601 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3607 if let Err(e) = self
3608 .replace_segments(
3609 &[seg_id.to_string()],
3610 new_id,
3611 total_docs,
3612 ReplacementLayout::BpReordered {
3613 converged: ladder_converged,
3614 },
3615 )
3616 .await
3617 {
3618 self.delete_output_if_unregistered(output_id, "replacement failure")
3619 .await;
3620 output_cleanup.disarm();
3621 if !matches!(&e, Error::IndexClosed) {
3622 self.pause_reorder_retries(seg_id, &e);
3623 }
3624 return Err(e);
3625 }
3626 output_cleanup.disarm();
3627 self.clear_reorder_retry(seg_id);
3628
3629 Ok(true)
3630 }
3631
3632 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3639 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3640
3641 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3642 for entry in entries {
3643 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3644 continue;
3645 };
3646 let Some(rest) = filename.strip_prefix("seg_") else {
3647 continue;
3648 };
3649 let Some(hex_id) = rest.get(..32) else {
3650 continue;
3651 };
3652 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3653 continue;
3654 }
3655 orphan_files
3656 .entry(hex_id.to_ascii_lowercase())
3657 .or_default()
3658 .push(entry);
3659 }
3660 }
3661
3662 let mut deleted = 0;
3663 for (hex_id, paths) in &orphan_files {
3664 let deletion_guard = {
3669 let st = self.state.lock().await;
3670 if st.metadata.has_segment(hex_id) {
3671 continue;
3672 }
3673 let Some(guard) = self
3674 .active_operations
3675 .try_register(vec![hex_id.to_string()])
3676 else {
3677 continue;
3678 };
3679 if self.tracker.is_deletion_protected(hex_id) {
3680 drop(guard);
3681 continue;
3682 }
3683 guard
3684 };
3685
3686 let results =
3691 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3692 .await;
3693 let removed = results.into_iter().all(|result| match result {
3694 Ok(()) => true,
3695 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3696 Err(error) => {
3697 log::warn!(
3698 "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
3699 self.schema.index_label(),
3700 hex_id,
3701 error,
3702 );
3703 false
3704 }
3705 });
3706 drop(deletion_guard);
3709 if removed {
3710 deleted += 1;
3711 log::info!(
3712 "[segment_cleanup] index={} swept orphan segment {}",
3713 self.schema.index_label(),
3714 hex_id
3715 );
3716 }
3717 }
3718
3719 Ok(deleted)
3720 }
3721}
3722
3723#[cfg(test)]
3724mod tests {
3725 use super::*;
3726 use std::sync::atomic::{AtomicBool, Ordering};
3727
3728 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3729 let schema = crate::dsl::SchemaBuilder::default().build();
3730 let metadata = IndexMetadata::new(schema.clone());
3731 Arc::new(SegmentManager::new(
3732 Arc::new(crate::directories::RamDirectory::new()),
3733 Arc::new(schema),
3734 metadata,
3735 Box::new(crate::merge::NoMergePolicy),
3736 0,
3737 1,
3738 Arc::new(Semaphore::new(1)),
3739 None,
3740 1024,
3741 Arc::new(ReorderConcurrencyGate::new(1)),
3742 None,
3743 ))
3744 }
3745
3746 #[test]
3747 fn force_merge_planner_pairs_large_and_small_segments() {
3748 let groups = plan_force_merge_groups(
3749 vec![
3750 ("a".into(), 6),
3751 ("b".into(), 6),
3752 ("c".into(), 4),
3753 ("d".into(), 4),
3754 ],
3755 10,
3756 );
3757
3758 assert_eq!(groups.len(), 2);
3759 assert!(groups.iter().all(|group| group.total_docs == 10));
3760 assert!(groups.iter().all(|group| group.segments.len() == 2));
3761 }
3762
3763 #[test]
3764 fn force_merge_planner_leaves_oversized_segments_alone() {
3765 let groups = plan_force_merge_groups(
3766 vec![
3767 ("oversized".into(), 11),
3768 ("small-a".into(), 5),
3769 ("small-b".into(), 5),
3770 ],
3771 10,
3772 );
3773
3774 assert_eq!(groups.len(), 2);
3775 assert_eq!(groups[0].total_docs, 10);
3776 assert_eq!(groups[0].segments.len(), 2);
3777 assert_eq!(groups[1].total_docs, 11);
3778 assert_eq!(groups[1].segments.len(), 1);
3779 }
3780
3781 #[test]
3782 fn force_merge_planner_never_exceeds_segment_format_limit() {
3783 let groups = plan_force_merge_groups(
3784 vec![
3785 ("large-a".into(), 3_000_000_000),
3786 ("large-b".into(), 2_000_000_000),
3787 ],
3788 u64::from(u32::MAX),
3789 );
3790 assert_eq!(groups.len(), 2);
3791 assert!(
3792 groups
3793 .iter()
3794 .all(|group| group.total_docs <= u64::from(u32::MAX))
3795 );
3796 }
3797
3798 #[test]
3799 fn force_merge_hierarchy_has_one_final_bp_pass() {
3800 assert_eq!(force_merge_output_count(1), 0);
3801 assert_eq!(force_merge_output_count(2), 1);
3802 assert_eq!(force_merge_output_count(64), 1);
3803 assert_eq!(force_merge_output_count(65), 2);
3804 assert_eq!(force_merge_output_count(127), 2);
3805 assert_eq!(force_merge_output_count(128), 3);
3806 assert_eq!(force_merge_output_count(1_000), 16);
3807 }
3808
3809 fn expand_force_merge_node(
3810 hierarchy: &ForceMergeHierarchy,
3811 source_count: usize,
3812 node: usize,
3813 sources: &mut Vec<usize>,
3814 ) {
3815 if node < source_count {
3816 sources.push(node);
3817 return;
3818 }
3819
3820 let step_index = node - source_count;
3821 let step = hierarchy
3822 .steps
3823 .get(step_index)
3824 .expect("merge input must refer to an existing source or output");
3825 for &input in &step.inputs {
3826 assert!(
3827 input < node,
3828 "merge step {step_index} refers to a future output node {input}"
3829 );
3830 expand_force_merge_node(hierarchy, source_count, input, sources);
3831 }
3832 }
3833
3834 #[test]
3835 fn force_merge_hierarchy_has_minimal_valid_arity() {
3836 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
3837
3838 for source_count in source_counts {
3839 let hierarchy = plan_force_merge_hierarchy(source_count);
3840 let output_count = hierarchy.steps.len();
3841
3842 assert!(
3843 hierarchy
3844 .steps
3845 .iter()
3846 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
3847 "invalid merge arity for {source_count} sources"
3848 );
3849 assert!(
3850 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
3851 "{output_count} outputs cannot reduce {source_count} sources"
3852 );
3853 assert!(
3854 output_count == 1
3855 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
3856 "{output_count} outputs are not minimal for {source_count} sources"
3857 );
3858 assert_eq!(output_count, force_merge_output_count(source_count));
3859 }
3860 }
3861
3862 #[test]
3863 fn force_merge_hierarchy_preserves_exact_source_order() {
3864 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
3865 let hierarchy = plan_force_merge_hierarchy(source_count);
3866 let mut sources = Vec::with_capacity(source_count);
3867 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
3868 assert_eq!(
3869 sources,
3870 (0..source_count).collect::<Vec<_>>(),
3871 "source order changed for {source_count} sources"
3872 );
3873 }
3874 }
3875
3876 fn force_merge_rewrite_cost(source_count: usize) -> usize {
3877 let hierarchy = plan_force_merge_hierarchy(source_count);
3878 let mut node_weights = vec![1usize; source_count];
3879 let mut rewrite_cost = 0usize;
3880
3881 for (step_index, step) in hierarchy.steps.iter().enumerate() {
3882 let output = source_count + step_index;
3883 let output_weight = step
3884 .inputs
3885 .iter()
3886 .map(|&input| {
3887 assert!(
3888 input < output,
3889 "merge step {step_index} refers to future output {input}"
3890 );
3891 node_weights[input]
3892 })
3893 .sum::<usize>();
3894 rewrite_cost += output_weight;
3895 node_weights.push(output_weight);
3896 }
3897
3898 assert_eq!(node_weights[hierarchy.root], source_count);
3899 rewrite_cost
3900 }
3901
3902 #[test]
3903 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
3904 assert_eq!(force_merge_rewrite_cost(65), 67);
3905 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
3906 }
3907
3908 #[test]
3909 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
3910 assert_eq!(
3911 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
3912 (false, false, 3),
3913 );
3914 assert_eq!(
3915 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
3916 (true, false, 4),
3917 );
3918 assert_eq!(
3919 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
3920 (true, true, 0),
3921 );
3922 }
3923
3924 #[tokio::test]
3925 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
3926 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3927 let field = schema_builder.add_text_field("text", true, true);
3928 let schema = schema_builder.build();
3929 let directory = crate::directories::RamDirectory::new();
3930 let config = crate::index::IndexConfig {
3931 num_indexing_threads: 1,
3932 merge_policy: Box::new(crate::merge::NoMergePolicy),
3933 ..Default::default()
3934 };
3935 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
3936 .await
3937 .unwrap();
3938 for value in ["one", "two", "three"] {
3939 let mut document = crate::dsl::Document::new();
3940 document.add_text(field, value);
3941 writer.add_document(document).unwrap();
3942 writer.commit().await.unwrap();
3943 }
3944
3945 let manager = Arc::clone(writer.segment_manager());
3946 let held_id = manager.get_segment_ids().await.pop().unwrap();
3947 let mut held = Some(
3948 manager
3949 .active_operations
3950 .try_register(vec![held_id])
3951 .unwrap(),
3952 );
3953 let batches = Arc::new(AtomicUsize::new(0));
3954 let batch_count = Arc::clone(&batches);
3955 writer
3956 .force_merge_with_snapshot_refresh(move || {
3957 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
3958 if refresh == 2 {
3961 drop(held.take());
3962 }
3963 std::future::ready(Ok(()))
3964 })
3965 .await
3966 .unwrap();
3967
3968 assert_eq!(manager.get_segment_ids().await.len(), 1);
3969 assert_eq!(
3970 batches.load(Ordering::Relaxed),
3971 4,
3972 "initial/final refreshes plus two replacements are required after the held source releases"
3973 );
3974 }
3975
3976 #[tokio::test]
3977 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
3978 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3979 schema_builder.set_reorder_on_merge(true);
3980 let schema = schema_builder.build();
3981 let mut metadata = IndexMetadata::new(schema.clone());
3982 metadata.add_segment("00000000000000000000000000000001".into(), 1);
3983 metadata.add_segment("00000000000000000000000000000002".into(), 1);
3984
3985 let global_merge_permits = Arc::new(Semaphore::new(1));
3986 let manager = Arc::new(SegmentManager::new(
3987 Arc::new(crate::directories::RamDirectory::new()),
3988 Arc::new(schema),
3989 metadata,
3990 Box::new(crate::merge::NoMergePolicy),
3991 0,
3992 1,
3993 Arc::clone(&global_merge_permits),
3994 None,
3995 1024,
3996 Arc::new(ReorderConcurrencyGate::new(1)),
3997 None,
3998 ));
3999
4000 manager.active_operations.pause_non_indexing();
4005 let force_merge = {
4006 let manager = Arc::clone(&manager);
4007 tokio::spawn(async move { manager.force_merge().await })
4008 };
4009 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4010 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4011 tokio::task::yield_now().await;
4012 }
4013 })
4014 .await
4015 .expect("force merge never reached the paused group claim");
4016
4017 assert_eq!(
4018 global_merge_permits.available_permits(),
4019 1,
4020 "force merge retained global capacity while vector staging blocked group ownership"
4021 );
4022
4023 force_merge.abort();
4024 let _ = force_merge.await;
4025 manager.active_operations.resume_non_indexing();
4026 }
4027
4028 #[test]
4029 fn output_cleanup_guard_runs_during_panic_unwind() {
4030 let cleaned = Arc::new(AtomicBool::new(false));
4031 let cleaned_in_callback = Arc::clone(&cleaned);
4032 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4033 cleaned_in_callback.store(true, Ordering::SeqCst);
4034 });
4035
4036 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4037 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4038 panic!("simulated reorder panic");
4039 }));
4040
4041 assert!(result.is_err());
4042 assert!(
4043 cleaned.load(Ordering::SeqCst),
4044 "partial output cleanup must run during unwind"
4045 );
4046 }
4047
4048 #[test]
4049 fn output_cleanup_guard_disarms_after_commit() {
4050 let cleaned = Arc::new(AtomicBool::new(false));
4051 let cleaned_in_callback = Arc::clone(&cleaned);
4052 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4053 cleaned_in_callback.store(true, Ordering::SeqCst);
4054 });
4055
4056 {
4057 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4058 guard.disarm();
4059 }
4060
4061 assert!(!cleaned.load(Ordering::SeqCst));
4062 }
4063
4064 #[test]
4065 fn test_active_operation_guard_releases_ownership() {
4066 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4067 {
4068 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4069 let snap = active.snapshot();
4070 assert!(snap.contains("a"));
4071 assert!(snap.contains("b"));
4072 }
4073 assert!(active.snapshot().is_empty());
4074 }
4075
4076 #[test]
4077 fn test_non_overlapping_operations_can_run_concurrently() {
4078 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4079 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4080 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4081 let snap = active.snapshot();
4082 assert_eq!(snap.len(), 4);
4083
4084 drop(first);
4085 let snap = active.snapshot();
4086 assert_eq!(snap.len(), 2);
4087 assert!(snap.contains("c"));
4088 assert!(snap.contains("d"));
4089 }
4090
4091 #[test]
4092 fn test_overlapping_operation_is_rejected_until_release() {
4093 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4094 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4095 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4096 drop(first);
4097 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4098 }
4099
4100 #[test]
4101 fn test_active_operation_snapshot() {
4102 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4103 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4104 let snap = active.snapshot();
4105 assert!(snap.contains("x"));
4106 assert!(snap.contains("y"));
4107 assert!(!snap.contains("z"));
4108 }
4109
4110 #[tokio::test]
4111 async fn operation_barrier_ignores_producers_started_after_snapshot() {
4112 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4113 let before_gate = active.try_register(vec!["old".into()]).unwrap();
4114 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4115 assert_eq!(parked_indexing, 0);
4116 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4117
4118 let waiter = {
4119 let active = Arc::clone(&active);
4120 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4121 };
4122 tokio::task::yield_now().await;
4123 assert!(!waiter.is_finished());
4124
4125 drop(before_gate);
4126 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4127 .await
4128 .expect("pre-gate operation barrier was starved by a post-gate producer")
4129 .unwrap();
4130 assert!(active.snapshot().contains("new-flat"));
4131 drop(after_gate);
4132 }
4133
4134 #[tokio::test]
4135 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4136 let manager = lifecycle_test_manager();
4137 manager
4138 .trained
4139 .store(Some(Arc::new(TrainedVectorStructures {
4140 centroids: rustc_hash::FxHashMap::default(),
4141 binary_quantizers: rustc_hash::FxHashMap::default(),
4142 ..Default::default()
4143 })));
4144
4145 let guard = manager.begin_vector_artifact_update().await.unwrap();
4146 assert!(
4147 manager.trained().is_some(),
4148 "search readers keep the last fully validated generation"
4149 );
4150 assert!(
4151 manager.trained_for_segment_build().is_none(),
4152 "new segment producers must stay flat during an artifact update"
4153 );
4154
4155 let detached_transaction_guard = guard.clone();
4156 drop(guard);
4157 assert!(
4158 manager.trained_for_segment_build().is_none(),
4159 "a detached lifecycle transaction must retain the producer gate after request cancellation"
4160 );
4161 drop(detached_transaction_guard);
4162 assert!(manager.trained_for_segment_build().is_some());
4163 }
4164
4165 #[tokio::test]
4166 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4167 let manager = lifecycle_test_manager();
4168 let guard = manager.begin_vector_artifact_update().await.unwrap();
4169 assert!(
4170 manager
4171 .active_operations
4172 .try_register(vec!["merge".into()])
4173 .is_none(),
4174 "ordinary merge/reorder work must not change staged sources"
4175 );
4176 let indexing = manager
4177 .active_operations
4178 .try_register_indexing(vec!["fresh".into()])
4179 .expect("indexing remains available in flat mode");
4180 drop(indexing);
4181
4182 drop(guard);
4183 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4184 assert!(
4185 manager
4186 .active_operations
4187 .try_register(vec!["merge".into()])
4188 .is_some()
4189 );
4190 }
4191
4192 #[tokio::test]
4193 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4194 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4195 let guard = active.try_register(vec!["live".into()]).unwrap();
4196 let cancellation = active.cancellation_flag();
4197 active.stop_accepting();
4198 assert!(cancellation.load(Ordering::Acquire));
4199 assert!(active.try_register(vec!["new".into()]).is_none());
4200
4201 let waiter = {
4202 let active = Arc::clone(&active);
4203 tokio::spawn(async move { active.wait_until_idle().await })
4204 };
4205 tokio::task::yield_now().await;
4206 assert!(!waiter.is_finished());
4207 drop(guard);
4208 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4209 .await
4210 .expect("shutdown waiter missed the final guard notification")
4211 .unwrap();
4212 }
4213
4214 #[tokio::test]
4215 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4216 let manager = lifecycle_test_manager();
4217 let started = Arc::new(Semaphore::new(0));
4218 let release = Arc::new(Semaphore::new(0));
4219 let completed = Arc::new(AtomicBool::new(false));
4220
4221 let request = {
4222 let manager = Arc::clone(&manager);
4223 let started = Arc::clone(&started);
4224 let release = Arc::clone(&release);
4225 let completed = Arc::clone(&completed);
4226 tokio::spawn(async move {
4227 manager
4228 .run_lifecycle_transaction(async move {
4229 started.add_permits(1);
4230 let _permit = release.acquire().await.unwrap();
4231 completed.store(true, Ordering::Release);
4232 Ok(())
4233 })
4234 .await
4235 })
4236 };
4237
4238 let _started = started.acquire().await.unwrap();
4239 request.abort();
4240 assert!(request.await.unwrap_err().is_cancelled());
4241 release.add_permits(1);
4242
4243 manager.begin_shutdown();
4244 tokio::time::timeout(
4245 std::time::Duration::from_secs(1),
4246 manager.wait_for_shutdown(),
4247 )
4248 .await
4249 .expect("shutdown did not drain detached lifecycle transaction");
4250 assert!(completed.load(Ordering::Acquire));
4251 }
4252
4253 #[tokio::test]
4254 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4255 let manager = lifecycle_test_manager();
4256 {
4257 let mut state = manager.state.lock().await;
4258 state.metadata.add_segment_meta(
4259 "eligible".into(),
4260 SegmentMetaInfo {
4261 num_docs: 10,
4262 ancestors: Vec::new(),
4263 generation: 1,
4264 reordered: true,
4265 bp_converged: false,
4266 bp_unconverged_passes: 2,
4267 },
4268 );
4269 state.metadata.add_segment_meta(
4270 "at-limit".into(),
4271 SegmentMetaInfo {
4272 num_docs: 20,
4273 ancestors: Vec::new(),
4274 generation: 1,
4275 reordered: true,
4276 bp_converged: false,
4277 bp_unconverged_passes: 3,
4278 },
4279 );
4280 state.metadata.add_segment_meta(
4281 "carried-debt".into(),
4282 SegmentMetaInfo {
4283 num_docs: 15,
4284 ancestors: Vec::new(),
4285 generation: 2,
4286 reordered: false,
4287 bp_converged: false,
4288 bp_unconverged_passes: 2,
4289 },
4290 );
4291 state.metadata.add_segment_meta(
4292 "carried-debt-at-limit".into(),
4293 SegmentMetaInfo {
4294 num_docs: 25,
4295 ancestors: Vec::new(),
4296 generation: 2,
4297 reordered: false,
4298 bp_converged: false,
4299 bp_unconverged_passes: 3,
4300 },
4301 );
4302 state.metadata.add_segment_meta(
4303 "converged".into(),
4304 SegmentMetaInfo {
4305 num_docs: 30,
4306 ancestors: Vec::new(),
4307 generation: 1,
4308 reordered: true,
4309 bp_converged: true,
4310 bp_unconverged_passes: 0,
4311 },
4312 );
4313 state.metadata.add_segment("fresh".into(), 40);
4314 }
4315
4316 assert_eq!(
4317 manager.unreordered_segments().await,
4318 vec![("fresh".into(), 40)],
4319 "a block-copy output with BP debt is not a fresh first-pass candidate",
4320 );
4321 let mut eligible = manager.unconverged_segments_below(3).await;
4322 eligible.sort_unstable();
4323 assert_eq!(
4324 eligible,
4325 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4326 );
4327 assert!(manager.unconverged_segments_below(0).await.is_empty());
4328 }
4329
4330 #[test]
4331 fn merge_retry_backoff_is_exponential_and_capped() {
4332 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4333 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4334 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4335 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4336 }
4337
4338 #[test]
4339 fn only_deterministic_source_errors_are_quarantined() {
4340 assert!(is_deterministic_source_error(&Error::Corruption(
4341 "bad footer".into()
4342 )));
4343 assert!(is_deterministic_source_error(&Error::Io(
4344 std::io::Error::from(std::io::ErrorKind::NotFound)
4345 )));
4346 assert!(!is_deterministic_source_error(&Error::Io(
4347 std::io::Error::from(std::io::ErrorKind::TimedOut)
4348 )));
4349 assert!(!is_deterministic_source_error(&Error::Io(
4350 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4351 )));
4352 }
4353
4354 #[test]
4355 fn transient_reorder_failure_is_backed_off_until_cleared() {
4356 let manager = lifecycle_test_manager();
4357 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4358 assert!(manager.paused_reorder_segments().contains("source"));
4359 manager.clear_reorder_retry("source");
4360 assert!(!manager.paused_reorder_segments().contains("source"));
4361 }
4362
4363 #[derive(Default)]
4366 struct FailingExistsDirectory(crate::directories::RamDirectory);
4367
4368 #[async_trait::async_trait]
4369 impl crate::directories::Directory for FailingExistsDirectory {
4370 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4371 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4372 }
4373
4374 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4375 self.0.file_size(path).await
4376 }
4377
4378 async fn open_read(
4379 &self,
4380 path: &std::path::Path,
4381 ) -> std::io::Result<crate::directories::FileHandle> {
4382 self.0.open_read(path).await
4383 }
4384
4385 async fn read_range(
4386 &self,
4387 path: &std::path::Path,
4388 range: std::ops::Range<u64>,
4389 ) -> std::io::Result<crate::directories::OwnedBytes> {
4390 self.0.read_range(path, range).await
4391 }
4392
4393 async fn list_files(
4394 &self,
4395 prefix: &std::path::Path,
4396 ) -> std::io::Result<Vec<std::path::PathBuf>> {
4397 self.0.list_files(prefix).await
4398 }
4399
4400 async fn open_lazy(
4401 &self,
4402 path: &std::path::Path,
4403 ) -> std::io::Result<crate::directories::FileHandle> {
4404 self.0.open_lazy(path).await
4405 }
4406 }
4407
4408 #[async_trait::async_trait]
4409 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4410 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4411 self.0.write(path, data).await
4412 }
4413
4414 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4415 self.0.delete(path).await
4416 }
4417
4418 async fn rename(
4419 &self,
4420 from: &std::path::Path,
4421 to: &std::path::Path,
4422 ) -> std::io::Result<()> {
4423 self.0.rename(from, to).await
4424 }
4425
4426 async fn sync(&self) -> std::io::Result<()> {
4427 self.0.sync().await
4428 }
4429
4430 async fn streaming_writer(
4431 &self,
4432 path: &std::path::Path,
4433 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4434 self.0.streaming_writer(path).await
4435 }
4436 }
4437
4438 #[derive(Debug, Clone)]
4439 struct MergeEverythingPolicy;
4440
4441 impl MergePolicy for MergeEverythingPolicy {
4442 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4443 if segments.len() < 2 {
4444 return Vec::new();
4445 }
4446 vec![crate::merge::MergeCandidate {
4447 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4448 }]
4449 }
4450
4451 fn clone_box(&self) -> Box<dyn MergePolicy> {
4452 Box::new(self.clone())
4453 }
4454 }
4455
4456 #[tokio::test]
4457 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4458 let manager = lifecycle_test_manager();
4459 let parked_indexing = manager
4464 .protect_new_segment("00000000000000000000000000000abc".into())
4465 .unwrap();
4466
4467 let error = tokio::time::timeout(
4468 std::time::Duration::from_secs(2),
4469 manager.begin_vector_artifact_update(),
4470 )
4471 .await
4472 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4473 .err()
4474 .expect("an old-generation prepared segment must block artifact replacement")
4475 .to_string();
4476 assert!(error.contains("built but uncommitted"), "{error}");
4477 assert!(
4478 !manager.vector_artifact_update.load(Ordering::Acquire),
4479 "a rejected update must release the producer gate"
4480 );
4481
4482 drop(parked_indexing);
4483
4484 let guard = manager
4485 .begin_vector_artifact_update()
4486 .await
4487 .expect("artifact update should succeed after the pending generation is resolved");
4488 drop(guard);
4489 }
4490
4491 #[tokio::test]
4492 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4493 let manager = lifecycle_test_manager();
4494 let merge_like = manager
4495 .active_operations
4496 .try_register(vec!["merge-source".into()])
4497 .unwrap();
4498
4499 let waiter = {
4500 let manager = Arc::clone(&manager);
4501 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4502 };
4503 for _ in 0..8 {
4504 tokio::task::yield_now().await;
4505 }
4506 assert!(
4507 !waiter.is_finished(),
4508 "artifact update must drain merge/reorder producers that may hold the previous generation"
4509 );
4510
4511 drop(merge_like);
4512 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4513 .await
4514 .expect("artifact update missed the lifecycle guard release")
4515 .unwrap()
4516 .unwrap();
4517 }
4518
4519 #[tokio::test]
4520 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4521 let manager = lifecycle_test_manager();
4522 let release = Arc::new(Semaphore::new(0));
4523 let merge_task = {
4524 let release = Arc::clone(&release);
4525 tokio::spawn(async move {
4526 let _permit = release.acquire().await.unwrap();
4527 })
4528 };
4529 manager.merge_handles.lock().push(merge_task);
4530
4531 let waiter = {
4532 let manager = Arc::clone(&manager);
4533 tokio::spawn(async move { manager.wait_for_all_merges().await })
4534 };
4535 for _ in 0..8 {
4536 tokio::task::yield_now().await;
4537 }
4538 assert!(!waiter.is_finished());
4539 waiter.abort();
4542 let join_error = waiter.await.unwrap_err();
4543 assert!(join_error.is_cancelled());
4544
4545 assert!(
4546 !manager.merge_handles.lock().is_empty(),
4547 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4548 );
4549
4550 release.add_permits(1);
4552 tokio::time::timeout(
4553 std::time::Duration::from_secs(1),
4554 manager.wait_for_all_merges(),
4555 )
4556 .await
4557 .expect("subsequent drain missed the reinserted merge handle");
4558 assert!(manager.merge_handles.lock().is_empty());
4559 }
4560
4561 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4562 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4563 let manager = lifecycle_test_manager();
4564 {
4565 let mut state = manager.state.lock().await;
4566 state
4567 .metadata
4568 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4569 state
4570 .metadata
4571 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4572 }
4573 let reorder_like = manager
4576 .active_operations
4577 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4578 .unwrap();
4579
4580 let force_merge = {
4581 let manager = Arc::clone(&manager);
4582 tokio::spawn(async move { manager.force_merge().await })
4583 };
4584
4585 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4586 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4587 assert!(
4588 retries >= 1,
4589 "force_merge never observed the conflicting owner (retries={retries})"
4590 );
4591 assert!(
4592 retries < 20,
4593 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4594 );
4595
4596 drop(reorder_like);
4597 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4600 .await
4601 .expect("force_merge kept spinning after the conflicting owner released")
4602 .unwrap();
4603 assert!(result.is_err());
4604 }
4605
4606 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4607 async fn force_merge_routes_around_segments_held_by_reorder() {
4608 let manager = lifecycle_test_manager();
4609 {
4610 let mut state = manager.state.lock().await;
4611 state
4612 .metadata
4613 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4614 state
4615 .metadata
4616 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4617 state
4618 .metadata
4619 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4620 }
4621 let _reorder_like = manager
4624 .active_operations
4625 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4626 .unwrap();
4627
4628 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4635 let manager = Arc::clone(&manager);
4636 async move { manager.force_merge().await }
4637 })
4638 .await
4639 .expect("force_merge livelocked on a segment held by an active reorder");
4640 assert!(result.is_err(), "fake segment files must fail the merge");
4641
4642 assert_eq!(
4643 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4644 0,
4645 "batch built from the ownership snapshot must not collide with the held segment"
4646 );
4647 }
4648
4649 #[tokio::test]
4650 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4651 let schema = crate::dsl::SchemaBuilder::default().build();
4652 let mut metadata = IndexMetadata::new(schema.clone());
4653 metadata.add_segment("00000000000000000000000000000001".into(), 10);
4654 metadata.add_segment("00000000000000000000000000000002".into(), 10);
4655 let manager = Arc::new(SegmentManager::new(
4656 Arc::new(FailingExistsDirectory::default()),
4657 Arc::new(schema),
4658 metadata,
4659 Box::new(MergeEverythingPolicy),
4660 0,
4661 1,
4662 Arc::new(Semaphore::new(1)),
4663 None,
4664 1024,
4665 Arc::new(ReorderConcurrencyGate::new(1)),
4666 None,
4667 ));
4668
4669 manager.maybe_merge().await;
4672
4673 tokio::time::timeout(
4674 std::time::Duration::from_secs(5),
4675 manager.wait_for_all_merges(),
4676 )
4677 .await
4678 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4679 assert!(
4680 manager.merge_retry_is_paused(),
4681 "the failed merge should have armed the retry backoff"
4682 );
4683
4684 manager.begin_shutdown();
4686 tokio::time::timeout(
4687 std::time::Duration::from_secs(5),
4688 manager.wait_for_shutdown(),
4689 )
4690 .await
4691 .expect("shutdown did not drain the merge retry wakeup task");
4692 }
4693}