1use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwap;
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 PublishedIndexGeneration, SegmentFiles, SegmentId, SegmentMeta, SegmentSnapshot,
58 SegmentTracker, TrainedVectorStructures,
59};
60#[cfg(feature = "native")]
61use crate::segment::{SegmentMerger, SegmentReader};
62
63use super::{MergePolicy, SegmentInfo};
64
65const FORCE_MERGE_MAX_FAN_IN: usize = 64;
66
67#[derive(Debug)]
68struct ForceMergeGroup {
69 segments: Vec<(String, u32)>,
70 total_docs: u64,
71}
72
73fn plan_force_merge_groups(
82 mut segments: Vec<(String, u32)>,
83 max_docs: u64,
84) -> Vec<ForceMergeGroup> {
85 segments.sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
86 right_docs
87 .cmp(left_docs)
88 .then_with(|| left_id.cmp(right_id))
89 });
90
91 let mut groups: Vec<ForceMergeGroup> = Vec::new();
92 for segment in segments {
93 let docs = u64::from(segment.1);
94 let best_group = groups
95 .iter()
96 .enumerate()
97 .filter_map(|(index, group)| {
98 group
99 .total_docs
100 .checked_add(docs)
101 .filter(|&total| total <= max_docs)
102 .map(|_| (index, group.total_docs))
103 })
104 .max_by_key(|&(index, used)| (used, std::cmp::Reverse(index)))
105 .map(|(index, _)| index);
106
107 if let Some(index) = best_group {
108 groups[index].total_docs += docs;
109 groups[index].segments.push(segment);
110 } else {
111 groups.push(ForceMergeGroup {
112 segments: vec![segment],
113 total_docs: docs,
114 });
115 }
116 }
117
118 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 Compacted,
746 BpReordered { converged: bool },
748 MaintenanceOnly,
750 PreserveSingleSource,
753}
754
755fn replacement_bp_state(
756 parent_has_debt: bool,
757 parent_unconverged_passes: u32,
758 layout: ReplacementLayout,
759) -> (bool, bool, u32) {
760 match layout {
761 ReplacementLayout::Compacted => {
762 unreachable!("compaction preserves its single source lineage")
763 }
764 ReplacementLayout::BlockCopy => (
765 false,
766 !parent_has_debt,
767 if parent_has_debt {
768 parent_unconverged_passes
769 } else {
770 0
771 },
772 ),
773 ReplacementLayout::BpReordered { converged } => (
774 true,
775 converged,
776 if converged {
777 0
778 } else {
779 parent_unconverged_passes.saturating_add(1)
780 },
781 ),
782 ReplacementLayout::PreserveSingleSource | ReplacementLayout::MaintenanceOnly => {
783 unreachable!("preserved layouts retain the complete source metadata")
784 }
785 }
786}
787
788fn replacement_seismic_state<'a>(
791 parents: impl Iterator<Item = &'a SegmentMetaInfo>,
792 pending_terms: u32,
793 layout: ReplacementLayout,
794) -> (u32, u32) {
795 if pending_terms == 0 {
796 return (0, 0);
797 }
798 let mut parent_count = 0usize;
799 let mut parent_pending_terms = 0u64;
800 let mut passes = 0u32;
801 let mut stalls = 0u32;
802 for parent in parents {
803 parent_count += 1;
804 parent_pending_terms += u64::from(parent.seismic_pending_terms);
805 passes = passes.max(parent.seismic_maintenance_passes);
806 stalls = stalls.max(parent.seismic_no_progress_passes);
807 }
808 let made_progress = u64::from(pending_terms) < parent_pending_terms;
811 if parent_count > 1 || made_progress {
812 stalls = 0;
813 }
814 if matches!(
815 layout,
816 ReplacementLayout::BpReordered { .. } | ReplacementLayout::MaintenanceOnly
817 ) {
818 passes = passes.saturating_add(1);
819 if parent_count == 1 && !made_progress {
820 stalls = stalls.saturating_add(1);
821 }
822 }
823 (passes, stalls)
824}
825
826#[derive(Clone, Copy, Debug, Eq, PartialEq)]
827enum VectorSegmentRewriteOutcome {
828 Rewritten,
829 AlreadyCurrent,
830 SourceGone,
831 Conflict,
832 Deferred,
833}
834
835pub(crate) struct StagedVectorSegment {
838 source_id: String,
839 output_id: SegmentId,
840 doc_count: u32,
841 _operation: SegmentOperationGuard,
842 cleanup: OutputCleanupGuard,
843}
844
845pub struct SegmentManager<D: DirectoryWriter + 'static> {
849 state: Arc<AsyncMutex<ManagerState>>,
851
852 active_operations: Arc<ActiveSegmentOperations>,
854
855 quarantined_segments: parking_lot::Mutex<HashSet<String>>,
860
861 merge_retry: parking_lot::Mutex<MergeRetryState>,
864
865 reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
869
870 merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
872
873 global_merge_wakeup_pending: AtomicBool,
877
878 force_merge_active: AtomicUsize,
883
884 #[cfg(test)]
887 force_merge_conflict_retries: std::sync::atomic::AtomicU64,
888
889 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
893
894 published_generation: Arc<ArcSwap<PublishedIndexGeneration>>,
898
899 vector_artifact_update: Arc<AtomicBool>,
903
904 tracker: Arc<SegmentTracker>,
906
907 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
909
910 directory: Arc<D>,
912 schema: Arc<crate::dsl::Schema>,
914 optimization: crate::structures::IndexOptimization,
916 posting_codec: crate::structures::PostingCodec,
918 term_dict_block_size: crate::structures::SSTableBlockSize,
919 term_cache_blocks: usize,
921 term_cache_budget_bytes: Option<usize>,
922 merge_permits: Arc<Semaphore>,
926 global_merge_permits: Arc<Semaphore>,
928 reorder_permits: Arc<ReorderConcurrencyGate>,
932 reorder_on_merge: bool,
937 merge_bp_time_budget: Option<std::time::Duration>,
941 bp_memory_budget_bytes: usize,
944 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
947 replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
951}
952
953struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
954
955impl Drop for ForceMergeActivityGuard<'_> {
956 fn drop(&mut self) {
957 self.0.fetch_sub(1, Ordering::AcqRel);
958 }
959}
960
961impl<D: DirectoryWriter + 'static> SegmentManager<D> {
962 #[allow(clippy::too_many_arguments)]
964 pub fn new(
965 directory: Arc<D>,
966 schema: Arc<crate::dsl::Schema>,
967 metadata: IndexMetadata,
968 merge_policy: Box<dyn MergePolicy>,
969 term_cache_blocks: usize,
970 max_concurrent_merges: usize,
971 global_merge_permits: Arc<Semaphore>,
972 merge_bp_time_budget: Option<std::time::Duration>,
973 bp_memory_budget_bytes: usize,
974 reorder_permits: Arc<ReorderConcurrencyGate>,
975 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
976 ) -> Self {
977 let reorder_on_merge = schema.reorder_on_merge();
980 if reorder_on_merge {
981 log::info!(
982 "[merge] index={} reorder-on-merge enabled by index schema",
983 schema.index_label()
984 );
985 }
986
987 let tracker = Arc::new(SegmentTracker::new());
988 for seg_id in metadata.owned_ids() {
989 tracker.register(&seg_id);
990 }
991
992 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
993 Arc::new(parking_lot::Mutex::new(Vec::new()));
994 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
995 let dir = Arc::clone(&directory);
996 let tracker = Arc::clone(&tracker);
997 let lifecycle_handles = Arc::clone(&lifecycle_handles);
998 let cleanup_index_label: Arc<str> = schema.index_label().into();
999 Arc::new(move |segment_ids| {
1000 let Ok(handle) = tokio::runtime::Handle::try_current() else {
1003 tracker.complete_deletion(&segment_ids);
1006 return;
1007 };
1008 let dir = Arc::clone(&dir);
1009 let task_tracker = Arc::clone(&tracker);
1010 let task_index_label = Arc::clone(&cleanup_index_label);
1011 let cleanup_ids = segment_ids.clone();
1012 let future = async move {
1013 for &segment_id in &segment_ids {
1014 log::info!(
1015 "[segment_cleanup] index={} deleting deferred segment {}",
1016 task_index_label,
1017 segment_id.to_hex()
1018 );
1019 if let Err(error) =
1020 crate::segment::delete_segment(dir.as_ref(), segment_id).await
1021 {
1022 log::warn!(
1023 "[segment_cleanup] index={} deferred delete failed for {}: {}",
1024 task_index_label,
1025 segment_id.to_hex(),
1026 error,
1027 );
1028 }
1029 }
1030 task_tracker.complete_deletion(&segment_ids);
1031 };
1032 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
1033 tracker.complete_deletion(&cleanup_ids);
1037 log::warn!(
1038 "[segment_cleanup] index={} runtime rejected deferred deletion; files will be swept later",
1039 cleanup_index_label
1040 );
1041 }
1042 })
1043 };
1044
1045 let initial_generation = Arc::new(PublishedIndexGeneration {
1046 publication_id: metadata.publication_generation,
1047 schema: Arc::clone(&schema),
1048 trained_vectors: None,
1049 });
1050 Self {
1051 state: Arc::new(AsyncMutex::new(ManagerState {
1052 metadata,
1053 merge_policy,
1054 })),
1055 active_operations: Arc::new(ActiveSegmentOperations::new(schema.index_label().into())),
1056 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
1057 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
1058 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
1059 merge_handles: parking_lot::Mutex::new(Vec::new()),
1060 global_merge_wakeup_pending: AtomicBool::new(false),
1061 force_merge_active: AtomicUsize::new(0),
1062 #[cfg(test)]
1063 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
1064 lifecycle_handles,
1065 published_generation: Arc::new(ArcSwap::new(initial_generation)),
1066 vector_artifact_update: Arc::new(AtomicBool::new(false)),
1067 tracker,
1068 delete_fn,
1069 directory,
1070 schema,
1071 optimization: crate::structures::IndexOptimization::default(),
1072 posting_codec: crate::structures::PostingCodec::default(),
1073 term_dict_block_size: crate::structures::SSTableBlockSize::default(),
1074 term_cache_blocks,
1075 term_cache_budget_bytes: None,
1076 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
1077 global_merge_permits,
1078 reorder_permits,
1079 reorder_on_merge,
1080 merge_bp_time_budget,
1081 bp_memory_budget_bytes,
1082 background_reorder_pool,
1083 replacement_refresh: parking_lot::RwLock::new(None),
1084 }
1085 }
1086
1087 pub fn with_term_cache_budget(mut self, bytes: Option<usize>) -> Self {
1090 self.term_cache_budget_bytes = bytes;
1091 self
1092 }
1093
1094 pub fn with_term_dict_block_size(mut self, size: crate::structures::SSTableBlockSize) -> Self {
1096 self.term_dict_block_size = size;
1097 self
1098 }
1099
1100 pub fn with_posting_config(
1102 mut self,
1103 optimization: crate::structures::IndexOptimization,
1104 posting_codec: crate::structures::PostingCodec,
1105 ) -> Self {
1106 self.optimization = optimization;
1107 self.posting_codec = posting_codec;
1108 self
1109 }
1110
1111 pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1112 where
1113 F: Fn() -> Fut + Send + Sync + 'static,
1114 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1115 {
1116 *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1117 }
1118
1119 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1125 if let Some(pool) = &self.background_reorder_pool {
1126 return Arc::clone(pool);
1127 }
1128 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1129 let threads = (num_cpus::get() / 2).max(1);
1130 log::info!(
1131 "[merge] process-wide background CPU pool: {} thread(s)",
1132 threads
1133 );
1134 Arc::new(
1135 rayon::ThreadPoolBuilder::new()
1136 .num_threads(threads)
1137 .thread_name(|i| format!("summa-bg-cpu-{}", i))
1138 .build()
1139 .expect("failed to build background CPU pool"),
1140 )
1141 }))
1142 }
1143
1144 pub fn begin_shutdown(&self) {
1148 self.active_operations.stop_accepting();
1149 }
1150
1151 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1159 where
1160 T: Send + 'static,
1161 F: std::future::Future<Output = Result<T>> + Send + 'static,
1162 {
1163 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1164 let future = async move {
1165 let result = transaction.await;
1166 let _ = result_tx.send(result);
1167 };
1168 let runtime = tokio::runtime::Handle::current();
1169 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1170 return Err(Error::Internal(
1171 "runtime rejected lifecycle metadata transaction".into(),
1172 ));
1173 }
1174 result_rx.await.map_err(|_| {
1175 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1176 })?
1177 }
1178
1179 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1181 let manager = Arc::clone(self);
1182 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1183 let Ok(handle) = tokio::runtime::Handle::try_current() else {
1184 log::warn!(
1185 "[segment_cleanup] index={} runtime unavailable; partial output {} will be swept on startup",
1186 manager.schema.index_label(),
1187 segment_id.to_hex(),
1188 );
1189 return;
1190 };
1191
1192 let cleanup_manager = Arc::clone(&manager);
1193 let future = async move {
1194 cleanup_manager
1195 .delete_output_if_unregistered(segment_id, "task unwind")
1196 .await;
1197 };
1198 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1199 log::warn!(
1200 "[segment_cleanup] index={} runtime rejected output cleanup; {} will be swept on startup",
1201 manager.schema.index_label(),
1202 segment_id.to_hex(),
1203 );
1204 }
1205 });
1206
1207 OutputCleanupGuard::new(output_id, cleanup)
1208 }
1209
1210 pub(crate) fn schedule_unpublished_segment_cleanup(
1215 self: &Arc<Self>,
1216 output_id: SegmentId,
1217 operation: SegmentOperationGuard,
1218 runtime: tokio::runtime::Handle,
1219 ) {
1220 let manager = Arc::clone(self);
1221 let output_hex = output_id.to_hex();
1222 let future = async move {
1223 manager
1224 .delete_output_if_unregistered(output_id, "indexing abort or failure")
1225 .await;
1226 drop(operation);
1227 };
1228 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1229 log::warn!(
1232 "[segment_cleanup] index={} runtime unavailable; indexing output {} will be swept on startup",
1233 self.schema.index_label(),
1234 output_hex,
1235 );
1236 }
1237 }
1238
1239 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1245 match self
1246 .active_operations
1247 .try_register_indexing(vec![segment_id.clone()])
1248 {
1249 Some(operation) => Ok(operation),
1250 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1251 None => Err(Error::Corruption(format!(
1252 "new segment ID {} is already owned by an active operation",
1253 segment_id
1254 ))),
1255 }
1256 }
1257
1258 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1262 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1263 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1264 })?;
1265 let files = SegmentFiles::new(id.0);
1266
1267 for path in files.mandatory_paths() {
1268 if !self.directory.exists(path).await.map_err(Error::Io)? {
1269 return Err(Error::Corruption(format!(
1270 "segment {} cannot be published: mandatory file {:?} is missing",
1271 segment_id, path
1272 )));
1273 }
1274 }
1275
1276 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1277 Error::Corruption(format!(
1278 "segment {} cannot be published: missing/unreadable {:?}: {}",
1279 segment_id, files.meta, e
1280 ))
1281 })?;
1282 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1283 Error::Corruption(format!(
1284 "segment {} cannot be published: failed reading {:?}: {}",
1285 segment_id, files.meta, e
1286 ))
1287 })?;
1288 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1289 Error::Corruption(format!(
1290 "segment {} cannot be published: invalid {:?}: {}",
1291 segment_id, files.meta, e
1292 ))
1293 })?;
1294
1295 if meta.id != id.0 || meta.num_docs != expected_docs {
1296 return Err(Error::Corruption(format!(
1297 "segment {} cannot be published: metadata identity/docs mismatch \
1298 (id={:032x}, docs={}, expected_docs={})",
1299 segment_id, meta.id, meta.num_docs, expected_docs
1300 )));
1301 }
1302
1303 Ok(())
1304 }
1305
1306 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1307 let inserted = self
1308 .quarantined_segments
1309 .lock()
1310 .insert(segment_id.to_string());
1311 if inserted {
1312 log::error!(
1313 "[merge] index={} quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1314 It remains metadata-live for explicit repair but is excluded from merges until restart",
1315 self.schema.index_label(),
1316 segment_id,
1317 error,
1318 );
1319 }
1320 }
1321
1322 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1323 let mut retry = self.merge_retry.lock();
1324 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1325 let delay = merge_retry_delay(retry.consecutive_failures);
1326 retry.retry_after = std::time::Instant::now().checked_add(delay);
1327 log::warn!(
1328 "[merge] index={} pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1329 self.schema.index_label(),
1330 delay.as_secs_f64(),
1331 retry.consecutive_failures,
1332 error,
1333 );
1334 delay
1335 }
1336
1337 fn clear_merge_retry_backoff(&self) {
1338 *self.merge_retry.lock() = MergeRetryState::default();
1339 }
1340
1341 fn merge_retry_is_paused(&self) -> bool {
1342 let mut retry = self.merge_retry.lock();
1343 match retry.retry_after {
1344 Some(deadline) if deadline > std::time::Instant::now() => true,
1345 Some(_) => {
1346 retry.retry_after = None;
1347 false
1348 }
1349 None => false,
1350 }
1351 }
1352
1353 async fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1354 let st = self.state.lock().await;
1357 if !st.metadata.has_segment(segment_id) {
1358 return;
1359 }
1360 let mut retries = self.reorder_retries.lock();
1361 let retry = retries.entry(segment_id.to_string()).or_default();
1362 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1363 let delay = merge_retry_delay(retry.consecutive_failures);
1364 retry.retry_after = std::time::Instant::now().checked_add(delay);
1365 log::warn!(
1366 "[reorder] index={} pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1367 self.schema.index_label(),
1368 segment_id,
1369 delay.as_secs_f64(),
1370 retry.consecutive_failures,
1371 error,
1372 );
1373 }
1374
1375 fn clear_reorder_retry(&self, segment_id: &str) {
1376 self.reorder_retries.lock().remove(segment_id);
1377 }
1378
1379 fn retire_reorder_retries(&self, retired: &[String]) {
1381 let mut retries = self.reorder_retries.lock();
1382 for id in retired {
1383 retries.remove(id);
1384 }
1385 }
1386
1387 fn paused_reorder_segments(&self) -> HashSet<String> {
1388 let now = std::time::Instant::now();
1389 let mut retries = self.reorder_retries.lock();
1390 let mut paused = HashSet::new();
1391 for (segment_id, retry) in retries.iter_mut() {
1392 match retry.retry_after {
1393 Some(deadline) if deadline > now => {
1394 paused.insert(segment_id.clone());
1395 }
1396 Some(_) => retry.retry_after = None,
1397 None => {}
1398 }
1399 }
1400 paused
1401 }
1402
1403 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1407 if self
1408 .global_merge_wakeup_pending
1409 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1410 .is_err()
1411 {
1412 return;
1413 }
1414
1415 let manager = Arc::clone(self);
1416 let future = async move {
1417 let capacity = tokio::select! {
1418 biased;
1419 () = manager.active_operations.wait_for_shutdown() => None,
1420 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1421 };
1422
1423 manager
1424 .global_merge_wakeup_pending
1425 .store(false, Ordering::Release);
1426 if let Some(permit) = capacity {
1427 drop(permit);
1431 manager.maybe_merge().await;
1432 }
1433 };
1434 let runtime = tokio::runtime::Handle::current();
1435 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1436 self.global_merge_wakeup_pending
1437 .store(false, Ordering::Release);
1438 log::warn!(
1439 "[merge] index={} runtime rejected global-capacity wakeup task",
1440 self.schema.index_label()
1441 );
1442 }
1443 }
1444
1445 #[cfg(test)]
1446 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1447 self.quarantined_segments.lock().contains(segment_id)
1448 }
1449
1450 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1455 let output_hex = output_id.to_hex();
1456 {
1457 let st = self.state.lock().await;
1458 if st.metadata.owns_id(&output_hex) {
1459 return;
1460 }
1461 }
1462
1463 log::info!(
1467 "[segment_cleanup] index={} deleting uncommitted output {} after {}",
1468 self.schema.index_label(),
1469 output_hex,
1470 reason,
1471 );
1472 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1473 {
1474 log::warn!(
1475 "[segment_cleanup] index={} failed deleting uncommitted output {}: {}",
1476 self.schema.index_label(),
1477 output_hex,
1478 error,
1479 );
1480 }
1481 }
1482
1483 pub async fn get_segment_ids(&self) -> Vec<String> {
1489 self.state.lock().await.metadata.segment_ids()
1490 }
1491
1492 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1494 self.published_generation.load().trained_vectors.clone()
1495 }
1496
1497 pub(crate) fn published_generation(&self) -> Arc<PublishedIndexGeneration> {
1500 self.published_generation.load_full()
1501 }
1502
1503 pub(crate) fn publication_id(&self) -> u64 {
1504 self.published_generation.load().publication_id
1505 }
1506
1507 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1514 if self.vector_artifact_update.load(Ordering::Acquire) {
1515 return None;
1516 }
1517 let trained = self.published_generation.load().trained_vectors.clone();
1518 if self.vector_artifact_update.load(Ordering::Acquire) {
1519 None
1520 } else {
1521 trained
1522 }
1523 }
1524
1525 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1544 self.vector_artifact_update
1545 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1546 .map_err(|_| {
1547 Error::Internal("a trained-vector artifact update is already in progress".into())
1548 })?;
1549 self.active_operations.pause_non_indexing();
1550 let guard = VectorArtifactUpdateGuard {
1551 _lease: Arc::new(VectorArtifactUpdateLease {
1552 updating: Arc::clone(&self.vector_artifact_update),
1553 active_operations: Arc::clone(&self.active_operations),
1554 }),
1555 };
1556 let (preexisting, parked_indexing) =
1557 self.active_operations.draining_operation_tokens_snapshot();
1558 if parked_indexing > 0 {
1559 return Err(Error::Internal(format!(
1560 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1561 segment(s) are built but uncommitted; commit or abort the pending \
1562 generation and retry"
1563 )));
1564 }
1565 self.active_operations
1566 .wait_until_operations_finish(&preexisting)
1567 .await;
1568 Ok(guard)
1569 }
1570
1571 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1574 let (vector_fields, schema, publication_id) = {
1576 let st = self.state.lock().await;
1577 (
1578 st.metadata.vector_fields.clone(),
1579 Arc::new(st.metadata.schema.clone()),
1580 st.metadata.publication_generation,
1581 )
1582 };
1583 let trained = IndexMetadata::try_load_trained_from_fields(
1585 &vector_fields,
1586 schema.as_ref(),
1587 self.directory.as_ref(),
1588 )
1589 .await?
1590 .map(Arc::new);
1591 self.published_generation
1595 .store(Arc::new(PublishedIndexGeneration {
1596 publication_id,
1597 schema,
1598 trained_vectors: trained,
1599 }));
1600 Ok(())
1601 }
1602
1603 pub(crate) async fn publish_vector_generation(
1610 self: &Arc<Self>,
1611 artifact_update: &VectorArtifactUpdateGuard,
1612 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1613 next_trained: Arc<TrainedVectorStructures>,
1614 staged: Vec<StagedVectorSegment>,
1615 ) -> Result<()> {
1616 let schema = self.published_generation().schema.clone();
1617 self.publish_vector_generation_with_schema(
1618 artifact_update,
1619 schema,
1620 vector_fields,
1621 Some(next_trained),
1622 staged,
1623 )
1624 .await
1625 }
1626
1627 pub(crate) async fn publish_vector_generation_with_schema(
1628 self: &Arc<Self>,
1629 artifact_update: &VectorArtifactUpdateGuard,
1630 schema: Arc<crate::dsl::Schema>,
1631 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1632 next_trained: Option<Arc<TrainedVectorStructures>>,
1633 mut staged: Vec<StagedVectorSegment>,
1634 ) -> Result<()> {
1635 if !self.vector_artifact_update.load(Ordering::Acquire) {
1636 return Err(Error::Internal(
1637 "vector generation publication lost its exclusive update lease".into(),
1638 ));
1639 }
1640
1641 for replacement in &staged {
1642 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1643 .await?;
1644 }
1645
1646 let mut st = Arc::clone(&self.state).lock_owned().await;
1647 let mut next = st.metadata.clone();
1648 next.schema = (*schema).clone();
1649 next.vector_fields = vector_fields;
1650 next.refresh_total_vectors();
1651
1652 for replacement in &staged {
1653 let source_info = next
1654 .segment_metas
1655 .remove(&replacement.source_id)
1656 .ok_or_else(|| {
1657 Error::Corruption(format!(
1658 "vector generation source {} disappeared before publication",
1659 replacement.source_id,
1660 ))
1661 })?;
1662 let output_hex = replacement.output_id.to_hex();
1663 if next.segment_metas.contains_key(&output_hex) {
1664 return Err(Error::Corruption(format!(
1665 "vector generation output {output_hex} is already metadata-live"
1666 )));
1667 }
1668 next.add_segment_meta(output_hex, source_info);
1671 }
1672
1673 let directory = Arc::clone(&self.directory);
1674 let published_generation = Arc::clone(&self.published_generation);
1675 let tracker = Arc::clone(&self.tracker);
1676 let replacement_refresh = self.replacement_refresh.read().clone();
1677 let manager = Arc::clone(self);
1678 let artifact_update = artifact_update.clone();
1682 let index_label = self.schema.index_label().to_owned();
1683 next.publication_generation =
1684 next.publication_generation.checked_add(1).ok_or_else(|| {
1685 Error::Corruption("vector publication generation exhausted u64".into())
1686 })?;
1687 let next_schema = schema;
1688 let next_publication_id = next.publication_generation;
1689 self.run_lifecycle_transaction(async move {
1690 let _artifact_update = artifact_update;
1691 next.save(directory.as_ref()).await?;
1692
1693 for replacement in &staged {
1694 tracker.register(&replacement.output_id.to_hex());
1695 }
1696 st.metadata = next;
1697 published_generation.store(Arc::new(PublishedIndexGeneration {
1698 publication_id: next_publication_id,
1699 schema: next_schema,
1700 trained_vectors: next_trained,
1701 }));
1702
1703 for replacement in &mut staged {
1706 replacement.cleanup.disarm();
1707 }
1708 let retired = staged
1709 .iter()
1710 .map(|replacement| replacement.source_id.clone())
1711 .collect::<Vec<_>>();
1712 manager.retire_reorder_retries(&retired);
1713 let ready_to_delete = tracker.mark_for_deletion(&retired);
1714 drop(st);
1715 for &segment_id in &ready_to_delete {
1716 if let Err(error) =
1717 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1718 {
1719 log::warn!(
1720 "[segment_cleanup] index={index_label} immediate dense-vector generation delete failed for {}: {}",
1721 segment_id.to_hex(),
1722 error,
1723 );
1724 }
1725 }
1726 tracker.complete_deletion(&ready_to_delete);
1727 refresh_replacement_topology(replacement_refresh, &index_label).await;
1728 Ok(())
1729 })
1730 .await
1731 }
1732
1733 pub(crate) async fn publish_vector_schema_only(
1737 self: &Arc<Self>,
1738 artifact_update: &VectorArtifactUpdateGuard,
1739 schema: Arc<crate::dsl::Schema>,
1740 ) -> Result<()> {
1741 let vector_fields = self
1742 .read_metadata(|metadata| metadata.vector_fields.clone())
1743 .await;
1744 let trained = self.published_generation().trained_vectors.clone();
1745 self.publish_vector_generation_with_schema(
1746 artifact_update,
1747 schema,
1748 vector_fields,
1749 trained,
1750 Vec::new(),
1751 )
1752 .await
1753 }
1754
1755 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1757 where
1758 F: FnOnce(&IndexMetadata) -> R,
1759 {
1760 let st = self.state.lock().await;
1761 f(&st.metadata)
1762 }
1763
1764 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1766 where
1767 F: FnOnce(&mut IndexMetadata),
1768 {
1769 let mut st = Arc::clone(&self.state).lock_owned().await;
1770 let mut next = st.metadata.clone();
1771 f(&mut next);
1772 let directory = Arc::clone(&self.directory);
1773 self.run_lifecycle_transaction(async move {
1774 next.save(directory.as_ref()).await?;
1775 st.metadata = next;
1776 Ok(())
1777 })
1778 .await
1779 }
1780
1781 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1784 let st = self.state.lock().await;
1785 let acquired = self.tracker.acquire(&st.metadata.segment_ids());
1786 SegmentSnapshot::with_generation(
1787 Arc::clone(&self.tracker),
1788 acquired,
1789 self.published_generation.load_full(),
1790 Arc::clone(&self.delete_fn),
1791 )
1792 .with_deletions(&st.metadata)
1793 }
1794
1795 pub fn tracker(&self) -> Arc<SegmentTracker> {
1797 Arc::clone(&self.tracker)
1798 }
1799
1800 pub fn directory(&self) -> Arc<D> {
1802 Arc::clone(&self.directory)
1803 }
1804}
1805
1806#[cfg(feature = "native")]
1811impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1812 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1814 for (segment_id, num_docs) in new_segments {
1817 self.validate_completed_segment(segment_id, *num_docs)
1818 .await?;
1819 }
1820
1821 let mut st = Arc::clone(&self.state).lock_owned().await;
1822 let mut next = st.metadata.clone();
1823 let mut added = Vec::new();
1824 for (segment_id, num_docs) in new_segments {
1825 if !next.has_segment(segment_id) {
1826 next.add_segment(segment_id.clone(), *num_docs);
1827 added.push(segment_id.clone());
1828 }
1829 }
1830
1831 let directory = Arc::clone(&self.directory);
1837 let tracker = Arc::clone(&self.tracker);
1838 self.run_lifecycle_transaction(async move {
1839 next.save(directory.as_ref()).await?;
1840 for segment_id in &added {
1841 tracker.register(segment_id);
1842 }
1843 st.metadata = next;
1844 Ok(())
1845 })
1846 .await
1847 }
1848
1849 pub async fn maybe_merge(self: &Arc<Self>) {
1860 if !self.active_operations.is_accepting() {
1861 log::debug!(
1862 "[maybe_merge] index={} manager is shutting down, skipping",
1863 self.schema.index_label()
1864 );
1865 return;
1866 }
1867 if self.merge_retry_is_paused() {
1868 log::debug!(
1869 "[maybe_merge] index={} retry backoff active, skipping",
1870 self.schema.index_label()
1871 );
1872 return;
1873 }
1874
1875 {
1878 let mut handles = self.merge_handles.lock();
1879 handles.retain(|h| !h.is_finished());
1880 }
1881 let local_slots = self.merge_permits.available_permits();
1882 let global_slots = self.global_merge_permits.available_permits();
1883 let slots_available = local_slots.min(global_slots);
1884
1885 {
1889 let st = self.state.lock().await;
1890 let quarantined = self.quarantined_segments.lock().clone();
1891 let active_ids = self.active_operations.snapshot();
1892
1893 let live_segments: Vec<SegmentInfo> = st
1898 .metadata
1899 .segment_metas
1900 .iter()
1901 .filter(|(id, _)| {
1902 !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1903 })
1904 .map(|(id, info)| SegmentInfo {
1905 id: id.clone(),
1906 num_docs: info.num_docs,
1907 })
1908 .collect();
1909 let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1910
1911 let segments: Vec<SegmentInfo> = live_segments
1914 .iter()
1915 .filter(|segment| !active_ids.contains(&segment.id))
1916 .cloned()
1917 .collect();
1918
1919 log::debug!(
1920 "[maybe_merge] index={} {} eligible segments",
1921 self.schema.index_label(),
1922 segments.len()
1923 );
1924
1925 let candidates = st.merge_policy.find_merges(&segments);
1926
1927 if candidates.is_empty() {
1928 return;
1929 }
1930
1931 if slots_available == 0 {
1935 if local_slots > 0 && global_slots == 0 {
1936 self.schedule_global_merge_wakeup();
1937 }
1938 log::debug!(
1939 "[maybe_merge] index={} at max concurrent merges, skipping",
1940 self.schema.index_label()
1941 );
1942 return;
1943 }
1944
1945 log::debug!(
1946 "[maybe_merge] index={} {} merge candidates, {} slots available",
1947 self.schema.index_label(),
1948 candidates.len(),
1949 slots_available
1950 );
1951
1952 let mut handles = Vec::new();
1953 for c in candidates {
1954 if handles.len() >= slots_available {
1955 break;
1956 }
1957 let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1963 if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1964 handles.push(h);
1965 }
1966 }
1967 if !handles.is_empty() {
1968 if severe_backlog && self.reorder_on_merge {
1969 log::info!(
1970 "[maybe_merge] index={} severe backlog: {} live segments; started {} fast \
1971 block-copy merge(s), deferring BP to the optimizer",
1972 self.schema.index_label(),
1973 live_segments.len(),
1974 handles.len(),
1975 );
1976 }
1977 self.merge_handles.lock().extend(handles);
1982 }
1983 }
1984 }
1985
1986 fn spawn_merge(
1995 self: &Arc<Self>,
1996 segment_ids_to_merge: Vec<String>,
1997 reorder_bmp: bool,
1998 ) -> Option<JoinHandle<()>> {
1999 if self.force_merge_active.load(Ordering::Acquire) > 0 {
2000 log::debug!(
2001 "[spawn_merge] index={} skipped: explicit force merge has priority",
2002 self.schema.index_label()
2003 );
2004 return None;
2005 }
2006 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
2007 Ok(permit) => permit,
2008 Err(_) => {
2009 log::debug!(
2010 "[spawn_merge] index={} skipped: global merge capacity is full",
2011 self.schema.index_label()
2012 );
2013 self.schedule_global_merge_wakeup();
2014 return None;
2015 }
2016 };
2017 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
2018 Ok(permit) => permit,
2019 Err(_) => {
2020 log::debug!(
2021 "[spawn_merge] index={} skipped: no merge permit available",
2022 self.schema.index_label()
2023 );
2024 return None;
2025 }
2026 };
2027 let output_id = SegmentId::new();
2028 let output_hex = output_id.to_hex();
2029
2030 let mut all_ids = segment_ids_to_merge.clone();
2031 all_ids.push(output_hex);
2032
2033 let guard = match self.active_operations.try_register(all_ids) {
2034 Some(g) => g,
2035 None => {
2036 log::debug!(
2037 "[spawn_merge] index={} skipped: segments overlap with an active operation",
2038 self.schema.index_label()
2039 );
2040 return None;
2041 }
2042 };
2043
2044 let sm = Arc::clone(self);
2045 let ids = segment_ids_to_merge;
2046
2047 let index_label = self.schema.index_label().to_owned();
2048 Some(tokio::spawn(async move {
2049 let mut reevaluate = false;
2050 let mut retry_delay = None;
2051
2052 let result = sm
2053 .merge_and_replace_registered(
2054 &ids,
2055 output_id,
2056 reorder_bmp,
2057 ReorderPriority::AutomaticMerge,
2058 Arc::new((guard, merge_permit, global_merge_permit)),
2059 )
2060 .await;
2061
2062 match result {
2063 Ok(_) => {
2064 sm.clear_merge_retry_backoff();
2065 reevaluate = true;
2066 }
2067 Err(MergeTaskError {
2068 error: Error::IndexClosed,
2069 ..
2070 }) => {
2071 log::debug!(
2072 "[merge] index={index_label} background merge for segments {:?} cancelled during shutdown",
2073 ids,
2074 );
2075 }
2076 Err(MergeTaskError {
2077 error,
2078 unavailable_segments,
2079 }) => {
2080 log::error!(
2081 "[merge] index={index_label} background merge failed for segments {:?}: {}",
2082 ids,
2083 error
2084 );
2085 if !unavailable_segments.is_empty() {
2086 reevaluate = true;
2090 } else {
2091 retry_delay = Some(sm.pause_merge_retries(&error));
2092 }
2093 }
2094 }
2095 if reevaluate {
2099 sm.maybe_merge().await;
2100 } else if let Some(retry_delay) = retry_delay {
2101 sm.schedule_merge_retry_wakeup(retry_delay);
2108 }
2109 }))
2110 }
2111
2112 async fn merge_and_replace_registered(
2119 self: &Arc<Self>,
2120 ids: &[String],
2121 output_id: SegmentId,
2122 reorder_bmp: bool,
2123 priority: ReorderPriority,
2124 ownership: Arc<dyn Send + Sync>,
2125 ) -> MergeTaskResult<(String, u32, bool)> {
2126 let manager = Arc::clone(self);
2127 let ids = ids.to_vec();
2128 self.run_lifecycle_transaction(async move {
2131 let _ownership = ownership;
2132 Ok(manager
2133 .merge_and_replace_owned(&ids, output_id, reorder_bmp, priority)
2134 .await)
2135 })
2136 .await
2137 .map_err(MergeTaskError::from)?
2138 }
2139
2140 async fn merge_and_replace_owned(
2141 self: &Arc<Self>,
2142 ids: &[String],
2143 output_id: SegmentId,
2144 reorder_bmp: bool,
2145 priority: ReorderPriority,
2146 ) -> MergeTaskResult<(String, u32, bool)> {
2147 let mut output_cleanup = self.output_cleanup_guard(output_id);
2148 let generation = self.published_generation();
2149 let trained = self.trained_for_segment_build();
2150 let granularity = if reorder_bmp {
2151 self.merge_granularity(ids).await
2152 } else {
2153 crate::segment::reorder::BpGranularity::Auto
2154 };
2155 let result = Self::do_merge(
2156 self.directory.as_ref(),
2157 &generation.schema,
2158 ids,
2159 output_id,
2160 self.term_cache_blocks,
2161 self.term_cache_budget_bytes,
2162 self.optimization,
2163 self.posting_codec,
2164 self.term_dict_block_size,
2165 trained.as_deref(),
2166 reorder_bmp,
2167 granularity,
2168 self.merge_bp_time_budget,
2169 self.bp_memory_budget_bytes,
2170 Arc::clone(&self.reorder_permits),
2171 priority,
2172 self.active_operations.cancellation_flag(),
2173 Some(self.background_cpu_pool()),
2174 )
2175 .await;
2176
2177 let (new_id, doc_count, bp_converged) = match result {
2178 Ok(value) => value,
2179 Err(error) => {
2180 for segment_id in &error.unavailable_segments {
2181 self.quarantine_segment(segment_id, &error.error);
2182 }
2183 self.delete_output_if_unregistered(output_id, "merge failure")
2184 .await;
2185 output_cleanup.disarm();
2186 return Err(error);
2187 }
2188 };
2189
2190 let layout = if reorder_bmp {
2191 ReplacementLayout::BpReordered {
2192 converged: bp_converged,
2193 }
2194 } else {
2195 ReplacementLayout::BlockCopy
2196 };
2197 if let Err(error) = self
2198 .replace_segments(ids, new_id.clone(), doc_count, layout, None)
2199 .await
2200 {
2201 self.delete_output_if_unregistered(output_id, "replacement failure")
2202 .await;
2203 output_cleanup.disarm();
2204 return Err(MergeTaskError::from(error));
2205 }
2206 output_cleanup.disarm();
2207 Ok((new_id, doc_count, bp_converged))
2208 }
2209
2210 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
2214 let manager = Arc::clone(self);
2215 let future = async move {
2216 tokio::select! {
2217 () = tokio::time::sleep(retry_delay) => {
2218 manager.maybe_merge().await;
2219 }
2220 () = manager.active_operations.wait_for_shutdown() => {}
2221 }
2222 };
2223 let runtime = tokio::runtime::Handle::current();
2224 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2225 log::warn!(
2226 "[merge] index={} runtime rejected merge-retry wakeup task; eligible segments may stay \
2227 unmerged until the next commit re-runs merge policy evaluation",
2228 self.schema.index_label()
2229 );
2230 }
2231 }
2232
2233 async fn replace_segments(
2237 self: &Arc<Self>,
2238 old_ids: &[String],
2239 new_id: String,
2240 doc_count: u32,
2241 layout: ReplacementLayout,
2242 expected_deletions: Option<&HashMap<String, Option<crate::segment::DeletionMeta>>>,
2243 ) -> Result<()> {
2244 self.validate_completed_segment(&new_id, doc_count).await?;
2247 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
2248 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
2249 })?;
2250 let output_reader = SegmentReader::open_with_term_cache_budget(
2251 self.directory.as_ref(),
2252 output_id,
2253 self.published_generation().schema.clone(),
2254 self.term_cache_blocks,
2255 self.term_cache_budget_bytes,
2256 )
2257 .await
2258 .map_err(|error| match error {
2259 Error::Io(_) | Error::IndexClosed => error,
2263 error => Error::Corruption(format!(
2264 "replacement segment {new_id} failed full reader validation: {error}"
2265 )),
2266 })?;
2267 if output_reader.num_docs() != doc_count {
2268 return Err(Error::Corruption(format!(
2269 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2270 output_reader.num_docs(),
2271 )));
2272 }
2273 let has_surviving_bmp = !output_reader.bmp_indexes().is_empty();
2274 let ann_fragmented = output_reader
2275 .vector_indexes()
2276 .iter()
2277 .any(|(&field, index)| {
2278 matches!(
2279 index,
2280 crate::segment::VectorIndex::BinaryIvf(_)
2281 | crate::segment::VectorIndex::ScannBinary(_)
2282 ) && output_reader
2283 .ann_health(crate::dsl::Field(field))
2284 .is_some_and(|health| health.fragmentation() > 1.0)
2285 });
2286 let seismic_pending_terms =
2287 output_reader
2288 .seismic_indexes()
2289 .values()
2290 .try_fold(0u32, |total, index| {
2291 total.checked_add(index.pending_terms()).ok_or_else(|| {
2292 Error::Corruption("Seismic maintenance debt exceeds u32".into())
2293 })
2294 })?;
2295 drop(output_reader);
2296
2297 let mut st = Arc::clone(&self.state).lock_owned().await;
2298 let missing: Vec<&String> = old_ids
2302 .iter()
2303 .filter(|id| !st.metadata.has_segment(id))
2304 .collect();
2305 if !missing.is_empty() {
2306 return Err(Error::Corruption(format!(
2307 "replace_segments: source segment(s) {:?} not in metadata — \
2308 refusing to add output {} (would duplicate documents)",
2309 missing, new_id
2310 )));
2311 }
2312
2313 if let Some(expected) = expected_deletions {
2314 if old_ids
2315 .iter()
2316 .any(|id| expected.get(id) != Some(&st.metadata.segment_metas[id].deletions))
2317 {
2318 return Err(Error::Internal(
2319 "row visibility changed during compaction; retry with a fresh snapshot".into(),
2320 ));
2321 }
2322 } else if matches!(layout, ReplacementLayout::Compacted) {
2323 return Err(Error::Internal(
2324 "compaction requires a visibility snapshot".into(),
2325 ));
2326 }
2327 let mut replacement_info = match layout {
2328 ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2329 let generation = old_ids
2330 .iter()
2331 .filter_map(|id| st.metadata.segment_metas.get(id))
2332 .map(|info| info.generation)
2333 .max()
2334 .unwrap_or(0)
2335 .checked_add(1)
2336 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2337 let parent_unconverged_passes = old_ids
2338 .iter()
2339 .filter_map(|id| st.metadata.segment_metas.get(id))
2340 .map(|info| info.bp_unconverged_passes)
2341 .max()
2342 .unwrap_or(0);
2343 let parent_has_debt = old_ids
2344 .iter()
2345 .filter_map(|id| st.metadata.segment_metas.get(id))
2346 .any(|info| !info.bp_converged);
2347 let (reordered, bp_converged, bp_unconverged_passes) =
2348 replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2349 SegmentMetaInfo {
2350 deletions: None,
2351 num_docs: doc_count,
2352 ancestors: old_ids.to_vec(),
2353 generation,
2354 reordered,
2355 bp_converged,
2356 bp_unconverged_passes,
2357 seismic_pending_terms,
2358 seismic_maintenance_passes: 0,
2359 seismic_no_progress_passes: 0,
2360 ann_fragmented,
2361 }
2362 }
2363 ReplacementLayout::PreserveSingleSource
2364 | ReplacementLayout::Compacted
2365 | ReplacementLayout::MaintenanceOnly => {
2366 let [source_id] = old_ids else {
2367 return Err(Error::Internal(
2368 "layout-preserving replacement requires exactly one source".into(),
2369 ));
2370 };
2371 let mut source = st
2372 .metadata
2373 .segment_metas
2374 .get(source_id)
2375 .cloned()
2376 .ok_or_else(|| {
2377 Error::Corruption(format!(
2378 "layout-preserving replacement source {source_id} disappeared"
2379 ))
2380 })?;
2381 source.num_docs = doc_count;
2382 if matches!(
2383 layout,
2384 ReplacementLayout::Compacted | ReplacementLayout::MaintenanceOnly
2385 ) {
2386 source.ancestors = old_ids.to_vec();
2387 source.generation = source.generation.checked_add(1).ok_or_else(|| {
2388 Error::Corruption("maintenance generation overflow".into())
2389 })?;
2390 }
2391 if matches!(layout, ReplacementLayout::Compacted) {
2392 source.deletions = None;
2393 if has_surviving_bmp && source.reordered {
2394 source.bp_converged = false;
2395 }
2396 }
2397 source
2398 }
2399 };
2400 replacement_info.seismic_pending_terms = seismic_pending_terms;
2401 replacement_info.ann_fragmented = ann_fragmented;
2402 let (passes, no_progress_passes) = replacement_seismic_state(
2403 old_ids.iter().map(|id| &st.metadata.segment_metas[id]),
2404 seismic_pending_terms,
2405 layout,
2406 );
2407 replacement_info.seismic_maintenance_passes = passes;
2408 replacement_info.seismic_no_progress_passes = no_progress_passes;
2409 let mut retired_ids = old_ids.to_vec();
2413 let mut deletion_output = None;
2414 let has_deletions = old_ids
2415 .iter()
2416 .any(|id| st.metadata.segment_metas[id].deletions.is_some());
2417 if matches!(layout, ReplacementLayout::Compacted) {
2418 for id in old_ids {
2419 if let Some(deletion) = &st.metadata.segment_metas[id].deletions {
2420 retired_ids.push(deletion.id.clone());
2421 }
2422 }
2423 } else if has_deletions && old_ids.len() == 1 {
2424 let source = &st.metadata.segment_metas[&old_ids[0]];
2425 if source.num_docs != doc_count {
2426 return Err(Error::Corruption(
2427 "layout-preserving replacement changed physical row count".into(),
2428 ));
2429 }
2430 replacement_info.deletions = source.deletions.clone();
2433 } else if has_deletions {
2434 let mut alive = crate::query::DocBitset::all(doc_count);
2435 let mut offset = 0u32;
2436 for id in old_ids {
2437 let info = &st.metadata.segment_metas[id];
2438 if let Some(deletion) = &info.deletions {
2439 let source = deletion
2440 .load(self.directory.as_ref(), info.num_docs)
2441 .await?;
2442 crate::segment::deletion::append_dead_rows(
2443 &mut alive,
2444 &source,
2445 info.num_docs,
2446 offset,
2447 )?;
2448 retired_ids.push(deletion.id.clone());
2449 }
2450 offset = offset
2451 .checked_add(info.num_docs)
2452 .ok_or_else(|| Error::Corruption("merge row count overflow".into()))?;
2453 }
2454 if offset != doc_count {
2455 return Err(Error::Corruption(
2456 "layout-preserving merge changed physical row count".into(),
2457 ));
2458 }
2459 let deletion_id = SegmentId::new();
2460 let claim = self.protect_new_segment(deletion_id.to_hex())?;
2461 deletion_output = Some((claim, self.output_cleanup_guard(deletion_id)));
2462 replacement_info.deletions = Some(
2463 crate::segment::deletion::write(
2464 self.directory.as_ref(),
2465 deletion_id,
2466 doc_count,
2467 &alive,
2468 )
2469 .await?,
2470 );
2471 }
2472 let mut next = st.metadata.clone();
2473 for id in old_ids {
2474 next.remove_segment(id);
2475 }
2476 next.add_segment_meta(new_id.clone(), replacement_info);
2477 retired_ids.sort_unstable();
2478 retired_ids.dedup();
2479 retired_ids.retain(|id| !next.owns_id(id));
2480
2481 let directory = Arc::clone(&self.directory);
2482 let tracker = Arc::clone(&self.tracker);
2483 let replacement_refresh = self.replacement_refresh.read().clone();
2484 let manager = Arc::clone(self);
2485 let index_label = self.schema.index_label().to_owned();
2486 self.run_lifecycle_transaction(async move {
2487 next.save(directory.as_ref()).await?;
2490 tracker.register(&new_id);
2491 if let Some(deletion) = &next.segment_metas[&new_id].deletions {
2492 tracker.register(&deletion.id);
2493 }
2494 if let Some((_, cleanup)) = deletion_output.as_mut() {
2495 cleanup.disarm();
2496 }
2497 st.metadata = next;
2498
2499 manager.retire_reorder_retries(&retired_ids);
2500
2501 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2505 drop(st);
2506 for &segment_id in &ready_to_delete {
2507 if let Err(error) =
2508 crate::segment::delete_segment(directory.as_ref(), segment_id).await
2509 {
2510 log::warn!(
2511 "[segment_cleanup] index={index_label} immediate delete failed for {}: {}",
2512 segment_id.to_hex(),
2513 error,
2514 );
2515 }
2516 }
2517 tracker.complete_deletion(&ready_to_delete);
2518 refresh_replacement_topology(replacement_refresh, &index_label).await;
2519 Ok(())
2520 })
2521 .await
2522 }
2523
2524 #[allow(clippy::too_many_arguments)]
2529 async fn do_merge(
2530 directory: &D,
2531 schema: &Arc<crate::dsl::Schema>,
2532 segment_ids_to_merge: &[String],
2533 output_segment_id: SegmentId,
2534 term_cache_blocks: usize,
2535 term_cache_budget_bytes: Option<usize>,
2536 optimization: crate::structures::IndexOptimization,
2537 posting_codec: crate::structures::PostingCodec,
2538 term_dict_block_size: crate::structures::SSTableBlockSize,
2539 trained: Option<&TrainedVectorStructures>,
2540 reorder_bmp: bool,
2541 granularity: crate::segment::reorder::BpGranularity,
2542 merge_bp_time_budget: Option<std::time::Duration>,
2543 bp_memory_budget_bytes: usize,
2544 reorder_permits: Arc<ReorderConcurrencyGate>,
2545 reorder_priority: ReorderPriority,
2546 cancellation: Arc<AtomicBool>,
2547 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2548 ) -> MergeTaskResult<(String, u32, bool)> {
2549 let output_hex = output_segment_id.to_hex();
2550 let load_start = std::time::Instant::now();
2551
2552 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2553 for id_str in segment_ids_to_merge {
2554 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2555 MergeTaskError::source(
2556 id_str.clone(),
2557 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2558 )
2559 })?;
2560 segment_ids.push(id);
2561 }
2562
2563 let mut unavailable_sources = Vec::new();
2568 let mut missing_files = Vec::new();
2569 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2570 let files = SegmentFiles::new(id.0);
2571 let mut source_unavailable = false;
2572 for path in files.mandatory_paths() {
2573 let exists = directory
2574 .exists(path)
2575 .await
2576 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2577 if !exists {
2578 source_unavailable = true;
2579 missing_files.push(format!("{}:{:?}", id_str, path));
2580 }
2581 }
2582 if source_unavailable {
2583 unavailable_sources.push(id_str.clone());
2584 }
2585 }
2586 if !unavailable_sources.is_empty() {
2587 return Err(MergeTaskError::sources(
2588 unavailable_sources,
2589 Error::Corruption(format!(
2590 "merge sources are missing mandatory files: {}",
2591 missing_files.join(", ")
2592 )),
2593 ));
2594 }
2595
2596 let schema_arc = Arc::clone(schema);
2597 let futures: Vec<_> = segment_ids
2598 .iter()
2599 .map(|&sid| {
2600 let sch = Arc::clone(&schema_arc);
2601 async move {
2602 SegmentReader::open_with_term_cache_budget(
2603 directory,
2604 sid,
2605 sch,
2606 term_cache_blocks,
2607 term_cache_budget_bytes,
2608 )
2609 .await
2610 }
2611 })
2612 .collect();
2613
2614 let results = futures::future::join_all(futures).await;
2615 let mut readers = Vec::with_capacity(results.len());
2616 let mut total_docs = 0u64;
2617 for (i, result) in results.into_iter().enumerate() {
2618 match result {
2619 Ok(r) => {
2620 total_docs += r.meta().num_docs as u64;
2621 readers.push(r);
2622 }
2623 Err(e) => {
2624 log::error!(
2625 "[merge] index={} Failed to open segment {}: {:?}",
2626 schema.index_label(),
2627 segment_ids_to_merge[i],
2628 e
2629 );
2630 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2631 }
2632 }
2633 }
2634 if total_docs > u32::MAX as u64 {
2635 return Err(Error::Internal(format!(
2636 "Merged segment doc count ({}) exceeds u32::MAX",
2637 total_docs
2638 ))
2639 .into());
2640 }
2641
2642 for (i, reader) in readers.iter().enumerate() {
2646 let meta_docs = reader.meta().num_docs;
2647 let store_docs = reader.store().num_docs();
2648 if store_docs != meta_docs {
2649 return Err(MergeTaskError::source(
2650 segment_ids_to_merge[i].clone(),
2651 Error::Corruption(format!(
2652 "pre-merge validation: segment {} store has {} docs but meta says {}",
2653 segment_ids_to_merge[i], store_docs, meta_docs
2654 )),
2655 ));
2656 }
2657 }
2658
2659 log::info!(
2660 "[merge] index={} loaded {} segment readers in {:.1}s",
2661 schema.index_label(),
2662 readers.len(),
2663 load_start.elapsed().as_secs_f64()
2664 );
2665
2666 let merger = SegmentMerger::new(Arc::clone(schema))
2667 .with_posting_config(optimization, posting_codec)
2668 .with_term_dict_block_size(term_dict_block_size)
2669 .with_reorder_fields(reorder_bmp)
2670 .with_granularity(granularity)
2671 .with_bp_budget(crate::segment::BpBudget {
2672 min_partition_docs: None,
2673 time_budget: merge_bp_time_budget,
2674 })
2675 .with_cancellation(cancellation)
2676 .with_bp_memory_budget(bp_memory_budget_bytes)
2677 .with_reorder_permits(reorder_permits)
2678 .with_reorder_priority(reorder_priority)
2679 .with_background_pool(bg_cpu_pool);
2680
2681 log::info!(
2682 "[merge] index={} {} segments -> {} (trained={})",
2683 schema.index_label(),
2684 segment_ids_to_merge.len(),
2685 output_hex,
2686 trained.map_or(0, |t| t.centroids.len()),
2687 );
2688
2689 let (_merged_meta, merge_stats) = merger
2690 .merge(directory, &readers, output_segment_id, trained)
2691 .await
2692 .map_err(|error| {
2693 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2694 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2700 } else {
2701 MergeTaskError::from(error)
2702 }
2703 })?;
2704 let bp_converged = merge_stats.bp_converged;
2705 if !bp_converged {
2706 log::info!(
2707 "[merge] index={} merge-time BP hit its wall-clock budget — output marked unconverged; \
2708 the background optimizer deepens it later",
2709 schema.index_label(),
2710 );
2711 }
2712
2713 log::info!(
2714 "[merge] index={} total wall-clock: {:.1}s ({} segments, {} docs)",
2715 schema.index_label(),
2716 load_start.elapsed().as_secs_f64(),
2717 readers.len(),
2718 total_docs,
2719 );
2720
2721 Ok((output_hex, total_docs as u32, bp_converged))
2722 }
2723
2724 pub async fn abort_merges(&self) {
2734 loop {
2735 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2736 if handles.is_empty() {
2737 return;
2738 }
2739 while let Some(result) = handles.join_next().await {
2740 if let Err(error) = result
2741 && error.is_panic()
2742 {
2743 log::error!(
2744 "[merge] index={} background task panicked while draining: {}",
2745 self.schema.index_label(),
2746 error
2747 );
2748 }
2749 }
2750 }
2751 }
2752
2753 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2758 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2759 while handles.join_next().await.is_some() {}
2760 }
2761
2762 pub async fn wait_for_all_merges(self: &Arc<Self>) {
2771 loop {
2772 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2773 if handles.is_empty() {
2774 break;
2775 }
2776 while handles.join_next().await.is_some() {}
2777 }
2778 }
2779
2780 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2785 self.wait_for_all_merges().await;
2786 self.active_operations.wait_until_idle().await;
2787 loop {
2788 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2789 if handles.is_empty() {
2790 break;
2791 }
2792 for handle in handles {
2793 if let Err(error) = handle.await
2794 && error.is_panic()
2795 {
2796 log::error!(
2797 "[segment_cleanup] index={} task panicked while draining: {}",
2798 self.schema.index_label(),
2799 error
2800 );
2801 }
2802 }
2803 }
2804 }
2805
2806 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2818 self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2819 .await
2820 }
2821
2822 pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2832 self: &Arc<Self>,
2833 refresh_snapshots: F,
2834 ) -> Result<()>
2835 where
2836 F: FnMut() -> Fut,
2837 Fut: std::future::Future<Output = Result<()>>,
2838 {
2839 self.force_merge_with_compaction_and_snapshot_refresh(None, refresh_snapshots)
2840 .await
2841 }
2842
2843 pub(crate) async fn force_merge_with_compaction_and_snapshot_refresh<F, Fut>(
2844 self: &Arc<Self>,
2845 compaction_budget: Option<usize>,
2846 mut refresh_snapshots: F,
2847 ) -> Result<()>
2848 where
2849 F: FnMut() -> Fut,
2850 Fut: std::future::Future<Output = Result<()>>,
2851 {
2852 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2857 std::time::Duration::from_millis(100);
2858 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2862
2863 let (_force_merge_activity, policy_segment_docs) = {
2864 let st = self.state.lock().await;
2865 self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2869 (
2870 ForceMergeActivityGuard(&self.force_merge_active),
2871 st.merge_policy.max_segment_docs(),
2872 )
2873 };
2874
2875 let background_merges = self
2878 .merge_handles
2879 .lock()
2880 .iter()
2881 .filter(|handle| !handle.is_finished())
2882 .count();
2883 if background_merges > 0 {
2884 log::info!(
2885 "[force_merge] index={} waiting for {} in-flight background merge(s) before planning",
2886 self.schema.index_label(),
2887 background_merges,
2888 );
2889 }
2890 let drain_start = std::time::Instant::now();
2891 self.wait_for_all_merges().await;
2892 if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2893 log::info!(
2894 "[force_merge] index={} drained background merges in {:.1}s",
2895 self.schema.index_label(),
2896 drain_start.elapsed().as_secs_f64(),
2897 );
2898 }
2899
2900 let refresh_start = std::time::Instant::now();
2905 refresh_snapshots().await?;
2906 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2907 log::info!(
2908 "[force_merge] index={} initial snapshot refresh took {:.1}s",
2909 self.schema.index_label(),
2910 refresh_start.elapsed().as_secs_f64(),
2911 );
2912 }
2913
2914 let mut completed_outputs = HashSet::new();
2918 let mut logged_held_wait = false;
2920
2921 loop {
2922 if !self.active_operations.is_accepting() {
2923 return Err(Error::IndexClosed);
2924 }
2925
2926 let segments: Vec<(String, u32)> = {
2927 let st = self.state.lock().await;
2928 st.metadata
2929 .segment_metas
2930 .iter()
2931 .filter(|(id, _)| !completed_outputs.contains(*id))
2932 .map(|(id, info)| (id.clone(), info.num_docs))
2933 .collect()
2934 };
2935
2936 let active_ids = self.active_operations.snapshot();
2941 let held = segments
2942 .iter()
2943 .filter(|(id, _)| active_ids.contains(id))
2944 .count();
2945 let free_segments: Vec<_> = segments
2946 .into_iter()
2947 .filter(|(id, _)| !active_ids.contains(id))
2948 .collect();
2949 let max_docs = u64::from(u32::MAX);
2956 let planned_groups = plan_force_merge_groups(free_segments, max_docs);
2957 if let Some(cap) = policy_segment_docs {
2958 for group in planned_groups
2959 .iter()
2960 .filter(|group| group.segments.len() >= 2)
2961 .filter(|group| group.total_docs > u64::from(cap))
2962 {
2963 log::warn!(
2964 "[force_merge] index={} output of {} docs intentionally exceeds the \
2965 background merge policy cap of {} docs (force merge compacts to the \
2966 u32 format limit)",
2967 self.schema.index_label(),
2968 group.total_docs,
2969 cap,
2970 );
2971 }
2972 }
2973 let next_group = planned_groups
2974 .into_iter()
2975 .find(|group| group.segments.len() >= 2);
2976
2977 let Some(group) = next_group else {
2978 if held == 0 {
2979 if !completed_outputs.is_empty() {
2980 completed_outputs.clear();
2986 continue;
2987 }
2988 if let Some(memory_budget) = compaction_budget {
2989 let dirty: Vec<_> = {
2990 let st = self.state.lock().await;
2991 st.metadata
2992 .segment_metas
2993 .iter()
2994 .filter(|(_, meta)| meta.deletions.is_some())
2995 .map(|(id, _)| id.clone())
2996 .collect()
2997 };
2998 for id in dirty {
2999 self.compact_segment(&id, memory_budget).await?;
3000 refresh_snapshots().await?;
3001 }
3002 }
3003 let remaining = {
3008 let st = self.state.lock().await;
3009 st.metadata.segment_metas.len()
3010 };
3011 if remaining > 1 {
3012 log::warn!(
3013 "[force_merge] index={} finished with {} segments: combined \
3014 document count exceeds the u32 segment format limit, so a \
3015 single output is impossible",
3016 self.schema.index_label(),
3017 remaining,
3018 );
3019 }
3020 refresh_snapshots().await?;
3025 return Ok(());
3026 }
3027 if !logged_held_wait {
3028 log::info!(
3029 "[force_merge] index={} waiting: {} segment(s) held by active \
3030 merge/reorder operations, no free group can merge",
3031 self.schema.index_label(),
3032 held
3033 );
3034 logged_held_wait = true;
3035 } else {
3036 log::debug!(
3037 "[force_merge] index={} still waiting on {} held segment(s)",
3038 self.schema.index_label(),
3039 held
3040 );
3041 }
3042 #[cfg(test)]
3043 self.force_merge_conflict_retries
3044 .fetch_add(1, Ordering::Relaxed);
3045 tokio::select! {
3046 biased;
3047 () = self.active_operations.wait_for_shutdown() => {
3048 return Err(Error::IndexClosed);
3049 }
3050 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
3051 }
3052 continue;
3053 };
3054 logged_held_wait = false;
3055
3056 let hierarchy = plan_force_merge_hierarchy(group.segments.len());
3060 let output_ids: Vec<_> = (0..hierarchy.steps.len())
3061 .map(|_| SegmentId::new())
3062 .collect();
3063 let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
3064 let mut all_ids = source_ids.clone();
3065 all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
3066 let group_guard = {
3067 let st = self.state.lock().await;
3068 source_ids
3069 .iter()
3070 .all(|id| st.metadata.has_segment(id))
3071 .then(|| self.active_operations.try_register(all_ids))
3072 .flatten()
3073 };
3074 let _group_guard = Arc::new(match group_guard {
3075 Some(guard) => guard,
3076 None if !self.active_operations.is_accepting() => {
3077 return Err(Error::IndexClosed);
3078 }
3079 None => {
3080 #[cfg(test)]
3081 self.force_merge_conflict_retries
3082 .fetch_add(1, Ordering::Relaxed);
3083 log::debug!(
3084 "[force_merge] index={} group lost a registration race, replanning",
3085 self.schema.index_label()
3086 );
3087 let had_tracked_merges = !self.merge_handles.lock().is_empty();
3088 self.wait_for_merging_thread().await;
3089 if !had_tracked_merges {
3090 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
3091 }
3092 continue;
3093 }
3094 });
3095
3096 log::info!(
3097 "[force_merge] index={} planned final group: {} segments, {} docs, {} merge pass(es)",
3098 self.schema.index_label(),
3099 group.segments.len(),
3100 group.total_docs,
3101 output_ids.len(),
3102 );
3103
3104 let group_global_merge_permit = if self.reorder_on_merge {
3114 let capacity_start = std::time::Instant::now();
3115 let permit = tokio::select! {
3116 biased;
3117 () = self.active_operations.wait_for_shutdown() => {
3118 return Err(Error::IndexClosed);
3119 }
3120 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3121 permit.map_err(|_| {
3122 Error::Internal(
3123 "global background merge scheduler is closed".into(),
3124 )
3125 })?
3126 }
3127 };
3128 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
3129 log::info!(
3130 "[force_merge] index={} waited {:.1}s for foreground global merge capacity",
3131 self.schema.index_label(),
3132 capacity_start.elapsed().as_secs_f64(),
3133 );
3134 }
3135 Some(Arc::new(permit))
3136 } else {
3137 None
3138 };
3139 let _foreground_reorder = if self.reorder_on_merge {
3140 log::info!(
3141 "[force_merge] index={} prioritizing BP capacity ({} total pass slot(s))",
3142 self.schema.index_label(),
3143 self.reorder_permits.limit(),
3144 );
3145 let admission_start = std::time::Instant::now();
3146 let guard = Arc::clone(&self.reorder_permits)
3147 .begin_foreground()
3148 .await
3149 .map_err(|_| {
3150 Error::Internal("background reorder scheduler is closed".into())
3151 })?;
3152 if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
3153 log::info!(
3154 "[force_merge] index={} acquired foreground BP capacity in {:.1}s",
3155 self.schema.index_label(),
3156 admission_start.elapsed().as_secs_f64(),
3157 );
3158 }
3159 Some(Arc::new(guard))
3160 } else {
3161 None
3162 };
3163
3164 let source_count = group.segments.len();
3165 let mut nodes: Vec<Option<(String, u32)>> =
3166 group.segments.into_iter().map(Some).collect();
3167 nodes.resize_with(source_count + hierarchy.steps.len(), || None);
3168 for (step_index, step) in hierarchy.steps.iter().enumerate() {
3169 let final_pass = step_index + 1 == hierarchy.steps.len();
3170 let mut batch_entries = Vec::with_capacity(step.inputs.len());
3171 for &node in &step.inputs {
3172 let entry = nodes
3173 .get_mut(node)
3174 .and_then(Option::take)
3175 .expect("force-merge hierarchy must reference an available node");
3176 batch_entries.push(entry);
3177 }
3178 let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
3179 let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
3180 let output_id = output_ids[step_index];
3181
3182 let capacity_start = std::time::Instant::now();
3183 let step_global_merge_permit = if group_global_merge_permit.is_none() {
3184 Some(Arc::new(tokio::select! {
3185 biased;
3186 () = self.active_operations.wait_for_shutdown() => {
3187 return Err(Error::IndexClosed);
3188 }
3189 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3190 permit.map_err(|_| {
3191 Error::Internal(
3192 "global background merge scheduler is closed".into(),
3193 )
3194 })?
3195 }
3196 }))
3197 } else {
3198 None
3199 };
3200 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
3201 log::info!(
3202 "[force_merge] index={} waited {:.1}s for global merge capacity",
3203 self.schema.index_label(),
3204 capacity_start.elapsed().as_secs_f64(),
3205 );
3206 }
3207
3208 let reorder_bmp = final_pass && self.reorder_on_merge;
3212 log::info!(
3213 "[force_merge] index={} {} pass: {} segments ({} docs, bp={})",
3214 self.schema.index_label(),
3215 if final_pass {
3216 "final"
3217 } else {
3218 "fan-in reduction"
3219 },
3220 batch.len(),
3221 batch_docs,
3222 reorder_bmp,
3223 );
3224 let (new_segment_id, total_docs, _) = self
3225 .merge_and_replace_registered(
3226 &batch,
3227 output_id,
3228 reorder_bmp,
3229 ReorderPriority::Foreground,
3230 Arc::new((
3231 Arc::clone(&_group_guard),
3232 group_global_merge_permit.clone(),
3233 _foreground_reorder.clone(),
3234 step_global_merge_permit.clone(),
3235 )),
3236 )
3237 .await
3238 .map_err(|error| error.error)?;
3239 drop(step_global_merge_permit);
3240
3241 let refresh_start = std::time::Instant::now();
3244 refresh_snapshots().await?;
3245 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
3246 log::info!(
3247 "[force_merge] index={} post-replacement snapshot refresh took {:.1}s",
3248 self.schema.index_label(),
3249 refresh_start.elapsed().as_secs_f64(),
3250 );
3251 }
3252
3253 let output_node = source_count + step_index;
3254 debug_assert!(nodes[output_node].is_none());
3255 nodes[output_node] = Some((new_segment_id, total_docs));
3256 }
3257 let (root_id, _) = nodes[hierarchy.root]
3258 .take()
3259 .expect("force-merge hierarchy must produce its root");
3260 debug_assert!(nodes.into_iter().all(|node| node.is_none()));
3261 completed_outputs.insert(root_id);
3262 }
3263 }
3264
3265 fn segment_needs_vector_rewrite(
3266 &self,
3267 schema: &crate::dsl::Schema,
3268 reader: &SegmentReader,
3269 field_ids: &[u32],
3270 trained: &TrainedVectorStructures,
3271 rewrite_existing: bool,
3272 ) -> Result<bool> {
3273 for &field_id in field_ids {
3274 let flat = reader.flat_vectors().get(&field_id);
3275 let ann = reader.vector_indexes().get(&field_id);
3276 if ann.is_some() && flat.is_none() {
3277 return Err(Error::Corruption(format!(
3278 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
3279 reader.meta().id,
3280 )));
3281 }
3282
3283 let Some(flat) = flat else {
3284 continue;
3285 };
3286 if flat.num_vectors == 0 {
3287 continue;
3288 }
3289 if rewrite_existing {
3290 return Ok(true);
3291 }
3292 let field = crate::dsl::Field(field_id);
3293 let entry = schema.get_field_entry(field).ok_or_else(|| {
3294 Error::Corruption(format!(
3295 "segment {:032x} references unknown vector field {field_id}",
3296 reader.meta().id,
3297 ))
3298 })?;
3299 let current = match entry.field_type {
3300 crate::dsl::FieldType::DenseVector
3304 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3305 config.index_type == crate::dsl::VectorIndexType::Tq
3306 }) =>
3307 {
3308 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
3309 }
3310 crate::dsl::FieldType::DenseVector
3311 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3312 config.index_type == crate::dsl::VectorIndexType::IvfTq
3313 }) =>
3314 {
3315 let config = entry
3316 .dense_vector_config
3317 .as_ref()
3318 .expect("matched IVF-TQ configuration");
3319 match (ann, trained.centroids.get(&field_id)) {
3320 (
3321 Some(crate::segment::VectorIndex::IvfTq { index, .. }),
3322 Some(centroids),
3323 ) => {
3324 let header = index.get().header();
3325 crate::structures::is_ivf_tq_cosine_generation(centroids.version)
3326 && crate::structures::is_ivf_tq_cosine_generation(
3327 header.quantizer_version,
3328 )
3329 && header.dim == config.dim
3330 && header.num_clusters == centroids.num_clusters
3331 && header.quantizer_version == centroids.version
3332 && header.codebook_version
3333 == crate::structures::vector::quantization::tq_expected_fingerprint(
3334 config.dim,
3335 )
3336 && header.routing == config.ivf_routing
3337 }
3338 (None, None) => true,
3339 _ => false,
3340 }
3341 }
3342 crate::dsl::FieldType::DenseVector
3343 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3344 config.index_type == crate::dsl::VectorIndexType::Scann
3345 }) =>
3346 {
3347 match (ann, trained.scann_artifacts.get(&field_id)) {
3348 (Some(crate::segment::VectorIndex::ScannAh(index)), Some(artifact)) => {
3349 index
3350 .get()
3351 .validate_scann_generation(
3352 artifact.config(),
3353 artifact.generation(),
3354 artifact.artifact_id(),
3355 )
3356 .is_ok()
3357 }
3358 (None, None) => true,
3359 _ => false,
3360 }
3361 }
3362 crate::dsl::FieldType::DenseVector => false,
3365 crate::dsl::FieldType::BinaryDenseVector
3366 if entry
3367 .binary_dense_vector_config
3368 .as_ref()
3369 .is_some_and(|config| {
3370 config.index_type == crate::dsl::BinaryIndexType::Scann
3371 }) =>
3372 {
3373 match (ann, trained.scann_artifacts.get(&field_id)) {
3374 (Some(crate::segment::VectorIndex::ScannBinary(index)), Some(artifact)) => {
3375 index
3376 .get()
3377 .validate_scann_generation(
3378 artifact.config(),
3379 artifact.generation(),
3380 artifact.artifact_id(),
3381 )
3382 .is_ok()
3383 }
3384 (None, None) => true,
3385 _ => false,
3386 }
3387 }
3388 crate::dsl::FieldType::BinaryDenseVector => matches!(
3389 (ann, trained.binary_quantizers.get(&field_id)),
3390 (Some(crate::segment::VectorIndex::BinaryIvf(_)), Some(_)) | (None, None)
3391 ),
3392 _ => false,
3393 };
3394 if !current {
3395 return Ok(true);
3396 }
3397 }
3398 Ok(false)
3399 }
3400
3401 async fn acquire_maintenance_capacity(
3402 &self,
3403 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
3404 let global = tokio::select! {
3405 biased;
3406 () = self.active_operations.wait_for_shutdown() => {
3407 return Err(Error::IndexClosed);
3408 }
3409 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3410 permit.map_err(|_| Error::Internal(
3411 "global background merge scheduler is closed".into()
3412 ))?
3413 }
3414 };
3415 let local = tokio::select! {
3416 biased;
3417 () = self.active_operations.wait_for_shutdown() => {
3418 return Err(Error::IndexClosed);
3419 }
3420 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
3421 permit.map_err(|_| Error::Internal(
3422 "background merge scheduler is closed".into()
3423 ))?
3424 }
3425 };
3426 Ok((global, local))
3427 }
3428
3429 async fn build_vector_replacement(
3430 self: &Arc<Self>,
3431 schemas: (&Arc<crate::dsl::Schema>, &Arc<crate::dsl::Schema>),
3432 segment_id: &str,
3433 source_id: SegmentId,
3434 output_id: SegmentId,
3435 trained: &TrainedVectorStructures,
3436 failure_context: &'static str,
3437 ) -> Result<(String, u32, OutputCleanupGuard)> {
3438 let mut cleanup = self.output_cleanup_guard(output_id);
3439 match crate::segment::reorder::rewrite_vector_segment(
3440 self.directory.as_ref(),
3441 schemas,
3442 source_id,
3443 output_id,
3444 self.term_cache_blocks,
3445 trained,
3446 Some(self.background_cpu_pool()),
3447 )
3448 .await
3449 {
3450 Ok((new_id, doc_count)) => {
3451 self.validate_completed_segment(&new_id, doc_count).await?;
3452 Ok((new_id, doc_count, cleanup))
3453 }
3454 Err(error) => {
3455 self.delete_output_if_unregistered(output_id, failure_context)
3456 .await;
3457 cleanup.disarm();
3458 if is_deterministic_source_error(&error) {
3459 self.quarantine_segment(segment_id, &error);
3460 }
3461 Err(error)
3462 }
3463 }
3464 }
3465
3466 pub(crate) async fn stage_vector_generation(
3470 self: &Arc<Self>,
3471 _artifact_update: &VectorArtifactUpdateGuard,
3472 segment_ids: &[String],
3473 field_ids: &[u32],
3474 trained: Arc<TrainedVectorStructures>,
3475 rewrite_existing: bool,
3476 ) -> Result<Vec<StagedVectorSegment>> {
3477 let schema = self.published_generation().schema.clone();
3478 self.stage_vector_generation_with_schema(
3479 _artifact_update,
3480 segment_ids,
3481 field_ids,
3482 trained,
3483 rewrite_existing,
3484 schema,
3485 )
3486 .await
3487 }
3488
3489 pub(crate) async fn stage_vector_generation_with_schema(
3490 self: &Arc<Self>,
3491 _artifact_update: &VectorArtifactUpdateGuard,
3492 segment_ids: &[String],
3493 field_ids: &[u32],
3494 trained: Arc<TrainedVectorStructures>,
3495 rewrite_existing: bool,
3496 schema: Arc<crate::dsl::Schema>,
3497 ) -> Result<Vec<StagedVectorSegment>> {
3498 if !self.vector_artifact_update.load(Ordering::Acquire) {
3499 return Err(Error::Internal(
3500 "cannot stage a vector generation without an exclusive update lease".into(),
3501 ));
3502 }
3503
3504 let source_schema = self.published_generation().schema.clone();
3505 let mut staged = Vec::new();
3506 for segment_id in segment_ids {
3507 if self.quarantined_segments.lock().contains(segment_id) {
3508 return Err(Error::Corruption(format!(
3509 "segment {segment_id} is quarantined after a deterministic source failure"
3510 )));
3511 }
3512 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3513 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3514 })?;
3515
3516 let _capacity = self.acquire_maintenance_capacity().await?;
3519
3520 let output_id = SegmentId::new();
3521 let output_hex = output_id.to_hex();
3522 let operation = {
3523 let st = self.state.lock().await;
3524 if !st.metadata.has_segment(segment_id) {
3525 return Err(Error::Corruption(format!(
3526 "vector generation source {segment_id} disappeared while lifecycle work was paused"
3527 )));
3528 }
3529 self.active_operations
3530 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
3531 }
3532 .ok_or_else(|| {
3533 if self.active_operations.is_accepting() {
3534 Error::Internal(format!(
3535 "vector generation could not claim stable source {segment_id}"
3536 ))
3537 } else {
3538 Error::IndexClosed
3539 }
3540 })?;
3541
3542 let reader = SegmentReader::open_with_term_cache_budget(
3543 self.directory.as_ref(),
3544 source_id,
3545 Arc::clone(&source_schema),
3546 self.term_cache_blocks,
3547 self.term_cache_budget_bytes,
3548 )
3549 .await?;
3550 if !self.segment_needs_vector_rewrite(
3551 schema.as_ref(),
3552 &reader,
3553 field_ids,
3554 trained.as_ref(),
3555 rewrite_existing,
3556 )? {
3557 continue;
3558 }
3559 drop(reader);
3560
3561 let (new_id, doc_count, cleanup) = self
3562 .build_vector_replacement(
3563 (&source_schema, &schema),
3564 segment_id,
3565 source_id,
3566 output_id,
3567 trained.as_ref(),
3568 "vector generation staging failure",
3569 )
3570 .await?;
3571 debug_assert_eq!(new_id, output_hex);
3572 let output_reader = SegmentReader::open_with_term_cache_budget(
3573 self.directory.as_ref(),
3574 output_id,
3575 Arc::clone(&schema),
3576 self.term_cache_blocks,
3577 self.term_cache_budget_bytes,
3578 )
3579 .await?;
3580 if self.segment_needs_vector_rewrite(
3581 schema.as_ref(),
3582 &output_reader,
3583 field_ids,
3584 trained.as_ref(),
3585 false,
3586 )? {
3587 return Err(Error::Corruption(format!(
3588 "staged vector segment {new_id} does not match its candidate codebook generation"
3589 )));
3590 }
3591
3592 staged.push(StagedVectorSegment {
3593 source_id: segment_id.clone(),
3594 output_id,
3595 doc_count,
3596 _operation: operation,
3597 cleanup,
3598 });
3599 }
3600 Ok(staged)
3601 }
3602
3603 async fn rewrite_vector_segment_once(
3604 self: &Arc<Self>,
3605 segment_id: &str,
3606 field_ids: &[u32],
3607 ) -> Result<VectorSegmentRewriteOutcome> {
3608 if self.quarantined_segments.lock().contains(segment_id) {
3609 return Err(Error::Corruption(format!(
3610 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3611 )));
3612 }
3613 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3614 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3615 })?;
3616
3617 let _capacity = self.acquire_maintenance_capacity().await?;
3622
3623 let output_id = SegmentId::new();
3624 let output_hex = output_id.to_hex();
3625 let all_ids = vec![segment_id.to_owned(), output_hex];
3626 let operation = {
3627 let st = self.state.lock().await;
3628 if !st.metadata.has_segment(segment_id) {
3629 return Ok(VectorSegmentRewriteOutcome::SourceGone);
3630 }
3631 self.active_operations.try_register(all_ids)
3632 };
3633 let _operation = match operation {
3634 Some(operation) => operation,
3635 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3636 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3637 };
3638
3639 let Some(trained) = self.trained_for_segment_build() else {
3640 return Ok(VectorSegmentRewriteOutcome::Deferred);
3641 };
3642 let schema = self.published_generation().schema.clone();
3643
3644 let reader = SegmentReader::open_with_term_cache_budget(
3645 self.directory.as_ref(),
3646 source_id,
3647 Arc::clone(&schema),
3648 self.term_cache_blocks,
3649 self.term_cache_budget_bytes,
3650 )
3651 .await?;
3652 if !self.segment_needs_vector_rewrite(
3653 schema.as_ref(),
3654 &reader,
3655 field_ids,
3656 trained.as_ref(),
3657 false,
3658 )? {
3659 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3660 }
3661 drop(reader);
3662
3663 let (new_id, doc_count, mut output_cleanup) = self
3664 .build_vector_replacement(
3665 (&schema, &schema),
3666 segment_id,
3667 source_id,
3668 output_id,
3669 trained.as_ref(),
3670 "vector rewrite failure",
3671 )
3672 .await?;
3673
3674 if let Err(error) = self
3675 .replace_segments(
3676 &[segment_id.to_owned()],
3677 new_id,
3678 doc_count,
3679 ReplacementLayout::PreserveSingleSource,
3680 None,
3681 )
3682 .await
3683 {
3684 self.delete_output_if_unregistered(output_id, "vector replacement failure")
3685 .await;
3686 output_cleanup.disarm();
3687 return Err(error);
3688 }
3689 output_cleanup.disarm();
3690 Ok(VectorSegmentRewriteOutcome::Rewritten)
3691 }
3692
3693 pub(crate) async fn rewrite_vector_segments(
3698 self: &Arc<Self>,
3699 field_ids: &[u32],
3700 ) -> Result<usize> {
3701 if field_ids.is_empty() {
3702 return Ok(0);
3703 }
3704 let mut rewritten = 0usize;
3705 loop {
3706 let segment_ids = self.get_segment_ids().await;
3707 let mut conflicted = false;
3708 let mut changed = false;
3709 for segment_id in segment_ids {
3710 match self
3711 .rewrite_vector_segment_once(&segment_id, field_ids)
3712 .await?
3713 {
3714 VectorSegmentRewriteOutcome::Rewritten => {
3715 rewritten += 1;
3716 changed = true;
3717 }
3718 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3719 VectorSegmentRewriteOutcome::Deferred => {
3720 return Err(Error::Internal(
3721 "ANN finalization lost the published trained generation".into(),
3722 ));
3723 }
3724 VectorSegmentRewriteOutcome::AlreadyCurrent
3725 | VectorSegmentRewriteOutcome::SourceGone => {}
3726 }
3727 }
3728 if !conflicted && !changed {
3729 log::info!(
3730 "[dense_vector_rewrite] index={} ANN finalization complete ({} segment(s) rewritten)",
3731 self.schema.index_label(),
3732 rewritten,
3733 );
3734 return Ok(rewritten);
3735 }
3736 tokio::select! {
3737 biased;
3738 () = self.active_operations.wait_for_shutdown() => {
3739 return Err(Error::IndexClosed);
3740 }
3741 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3742 }
3743 }
3744 }
3745
3746 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3751 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3752 return;
3753 }
3754 let manager = Arc::clone(self);
3755 let future = async move {
3756 let field_ids = manager
3757 .read_metadata(|metadata| {
3758 metadata
3759 .vector_fields
3760 .keys()
3761 .filter(|field_id| metadata.is_field_built(**field_id))
3762 .copied()
3763 .collect::<Vec<_>>()
3764 })
3765 .await;
3766 for segment_id in segment_ids {
3767 loop {
3768 match manager
3769 .rewrite_vector_segment_once(&segment_id, &field_ids)
3770 .await
3771 {
3772 Ok(VectorSegmentRewriteOutcome::Conflict) => {
3773 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3774 }
3775 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3776 Ok(_) => break,
3777 Err(error) => {
3778 log::error!(
3779 "[dense_vector_rewrite] index={} failed to upgrade newly committed segment {}: {}",
3780 manager.schema.index_label(),
3781 segment_id,
3782 error,
3783 );
3784 break;
3785 }
3786 }
3787 }
3788 }
3789 };
3790 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3791 log::warn!(
3792 "[dense_vector_rewrite] index={} runtime unavailable; newly committed flat segment upgrade deferred",
3793 self.schema.index_label()
3794 );
3795 return;
3796 };
3797 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3798 log::warn!(
3799 "[dense_vector_rewrite] index={} runtime rejected newly committed flat segment upgrade",
3800 self.schema.index_label()
3801 );
3802 }
3803 }
3804
3805 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3811 self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3812 .await
3813 }
3814
3815 pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3818 self: &Arc<Self>,
3819 mut refresh_snapshots: F,
3820 ) -> Result<()>
3821 where
3822 F: FnMut() -> Fut,
3823 Fut: std::future::Future<Output = Result<()>>,
3824 {
3825 self.wait_for_all_merges().await;
3826 refresh_snapshots().await?;
3827 let segment_ids = self.get_segment_ids().await;
3828
3829 if segment_ids.is_empty() {
3830 log::info!(
3831 "[reorder] index={} no segments to reorder",
3832 self.schema.index_label()
3833 );
3834 return Ok(());
3835 }
3836
3837 log::info!(
3838 "[reorder] index={} reordering {} segments",
3839 self.schema.index_label(),
3840 segment_ids.len()
3841 );
3842
3843 for seg_id in segment_ids {
3844 match self
3845 .reorder_single_segment(
3846 &seg_id,
3847 Some(self.background_cpu_pool()),
3848 crate::segment::BpBudget::full(),
3849 )
3850 .await
3851 {
3852 Ok(true) => refresh_snapshots().await?,
3853 Ok(false) => log::warn!(
3854 "[reorder] index={} segment {} skipped (in merge)",
3855 self.schema.index_label(),
3856 seg_id
3857 ),
3858 Err(e) => return Err(e),
3859 }
3860 }
3861
3862 refresh_snapshots().await?;
3865 log::info!(
3866 "[reorder] index={} all segments reordered",
3867 self.schema.index_label()
3868 );
3869 Ok(())
3870 }
3871
3872 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3877 self.unreordered_segments()
3878 .await
3879 .into_iter()
3880 .map(|(id, _)| id)
3881 .collect()
3882 }
3883
3884 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3887 let quarantined = self.quarantined_segments.lock().clone();
3888 let paused = self.paused_reorder_segments();
3889 let st = self.state.lock().await;
3890 let active_ids = self.active_operations.snapshot();
3891 let has_bp_reorder = self.schema.has_reorder_fields();
3892 st.metadata
3893 .segment_metas
3894 .iter()
3895 .filter(|(id, info)| {
3896 info.bp_converged
3897 && ((has_bp_reorder && !info.reordered)
3898 || info.ann_fragmented
3899 || (info.seismic_pending_terms > 0 && info.seismic_maintenance_passes == 0))
3900 && !active_ids.contains(*id)
3901 && !quarantined.contains(*id)
3902 && !paused.contains(*id)
3903 })
3904 .map(|(id, info)| (id.clone(), info.num_docs))
3905 .collect()
3906 }
3907
3908 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3913 self.unconverged_segments_below(u32::MAX)
3914 .await
3915 .into_iter()
3916 .map(|(id, docs, _)| (id, docs))
3917 .collect()
3918 }
3919
3920 pub async fn unconverged_segments_below(
3923 &self,
3924 max_unconverged_passes: u32,
3925 ) -> Vec<(String, u32, u32)> {
3926 let quarantined = self.quarantined_segments.lock().clone();
3927 let paused = self.paused_reorder_segments();
3928 let st = self.state.lock().await;
3929 let active_ids = self.active_operations.snapshot();
3930 st.metadata
3931 .segment_metas
3932 .iter()
3933 .filter(|(id, info)| {
3934 ((!info.bp_converged && info.bp_unconverged_passes < max_unconverged_passes)
3935 || (info.seismic_pending_terms > 0
3936 && (info.seismic_maintenance_passes > 0 || !info.bp_converged)
3937 && info.seismic_no_progress_passes < max_unconverged_passes))
3938 && !active_ids.contains(*id)
3939 && !quarantined.contains(*id)
3940 && !paused.contains(*id)
3941 })
3942 .map(|(id, info)| {
3943 (
3944 id.clone(),
3945 info.num_docs,
3946 if info.seismic_pending_terms > 0
3947 && (info.seismic_maintenance_passes > 0 || !info.bp_converged)
3948 && info.seismic_no_progress_passes < max_unconverged_passes
3949 {
3950 info.seismic_no_progress_passes
3951 } else {
3952 info.bp_unconverged_passes
3953 },
3954 )
3955 })
3956 .collect()
3957 }
3958
3959 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3960 let st = self.state.lock().await;
3961 let deepening = ids.iter().any(|id| {
3962 st.metadata
3963 .segment_metas
3964 .get(id)
3965 .is_some_and(|info| !info.bp_converged)
3966 });
3967 drop(st);
3968 if deepening {
3969 log::info!(
3970 "[reorder] index={} source BP lineage unconverged — forcing record-level BP (deepening pass)",
3971 self.schema.index_label(),
3972 );
3973 crate::segment::reorder::BpGranularity::Records
3974 } else {
3975 crate::segment::reorder::BpGranularity::Auto
3976 }
3977 }
3978
3979 pub async fn reorder_single_segment(
3984 self: &Arc<Self>,
3985 seg_id: &str,
3986 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3987 bp_budget: crate::segment::BpBudget,
3988 ) -> Result<bool> {
3989 self.reorder_single_segment_with_policy(seg_id, rayon_pool, bp_budget, None)
3990 .await
3991 }
3992
3993 pub async fn optimize_single_segment(
3996 self: &Arc<Self>,
3997 seg_id: &str,
3998 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3999 bp_budget: crate::segment::BpBudget,
4000 max_bp_passes: u32,
4001 ) -> Result<bool> {
4002 self.reorder_single_segment_with_policy(seg_id, rayon_pool, bp_budget, Some(max_bp_passes))
4003 .await
4004 }
4005
4006 async fn reorder_single_segment_with_policy(
4007 self: &Arc<Self>,
4008 seg_id: &str,
4009 rayon_pool: Option<Arc<rayon::ThreadPool>>,
4010 bp_budget: crate::segment::BpBudget,
4011 automatic_bp_limit: Option<u32>,
4012 ) -> Result<bool> {
4013 let source_id = SegmentId::from_hex(seg_id)
4014 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
4015 if self.quarantined_segments.lock().contains(seg_id) {
4016 return Err(Error::Corruption(format!(
4017 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
4018 seg_id
4019 )));
4020 }
4021 if self.force_merge_active.load(Ordering::Acquire) > 0 {
4022 log::debug!(
4023 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
4024 self.schema.index_label(),
4025 seg_id,
4026 );
4027 return Ok(false);
4028 }
4029
4030 let reorder_gate = Arc::clone(&self.reorder_permits);
4035 let _reorder_permit = tokio::select! {
4036 biased;
4037 () = self.active_operations.wait_for_shutdown() => {
4038 return Err(Error::IndexClosed);
4039 }
4040 permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
4041 permit.map_err(|_| {
4042 Error::Internal("background reorder scheduler is closed".into())
4043 })?
4044 }
4045 };
4046
4047 let output_id = SegmentId::new();
4048 let output_hex = output_id.to_hex();
4049
4050 let all_ids = vec![seg_id.to_string(), output_hex];
4056 let (_guard, source_docs, generation, source_deletions, _visibility_owner, run_bp) = {
4057 let st = self.state.lock().await;
4058 if self.force_merge_active.load(Ordering::Acquire) > 0 {
4062 log::debug!(
4063 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
4064 self.schema.index_label(),
4065 seg_id,
4066 );
4067 return Ok(false);
4068 }
4069 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
4070 log::info!(
4071 "[optimizer] index={} segment {} no longer in metadata (merged away), skipping reorder",
4072 self.schema.index_label(),
4073 seg_id
4074 );
4075 self.clear_reorder_retry(seg_id);
4076 return Ok(false);
4077 };
4078
4079 let run_bp = automatic_bp_limit.is_none_or(|limit| {
4080 self.schema.has_reorder_fields()
4081 && ((!source_meta.reordered && source_meta.bp_converged)
4082 || (!source_meta.bp_converged && source_meta.bp_unconverged_passes < limit))
4083 });
4084 let generation = self.published_generation();
4085 match self.active_operations.try_register(all_ids) {
4086 Some(guard) => {
4087 let deletion_ids: Vec<_> = source_meta
4090 .deletions
4091 .iter()
4092 .map(|deletion| deletion.id.clone())
4093 .collect();
4094 let acquired = self.tracker.acquire(&deletion_ids);
4095 if acquired.len() != deletion_ids.len() {
4096 return Err(Error::Corruption(
4097 "maintenance source visibility is already retired".into(),
4098 ));
4099 }
4100 let visibility_owner = SegmentSnapshot::with_delete_fn(
4101 Arc::clone(&self.tracker),
4102 acquired,
4103 Arc::clone(&self.delete_fn),
4104 );
4105 (
4106 guard,
4107 source_meta.num_docs,
4108 generation,
4109 source_meta.deletions.clone(),
4110 visibility_owner,
4111 run_bp,
4112 )
4113 }
4114 None if !self.active_operations.is_accepting() => {
4115 return Err(Error::IndexClosed);
4116 }
4117 None => {
4118 log::debug!(
4119 "[optimizer] index={} segment {} in active merge, skipping",
4120 self.schema.index_label(),
4121 seg_id
4122 );
4123 return Ok(false);
4124 }
4125 }
4126 };
4127
4128 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
4133 if is_deterministic_source_error(&error) {
4134 self.quarantine_segment(seg_id, &error);
4135 } else if !matches!(&error, Error::IndexClosed) {
4136 self.pause_reorder_retries(seg_id, &error).await;
4137 }
4138 return Err(error);
4139 }
4140
4141 let alive = match source_deletions {
4142 Some(deletions) => match deletions.load(self.directory.as_ref(), source_docs).await {
4143 Ok(alive) => Some(alive),
4144 Err(error) => {
4145 if is_deterministic_source_error(&error) {
4146 self.quarantine_segment(seg_id, &error);
4147 } else if !matches!(&error, Error::IndexClosed) {
4148 self.pause_reorder_retries(seg_id, &error).await;
4149 }
4150 return Err(error);
4151 }
4152 },
4153 None => None,
4154 };
4155 let mut output_cleanup = self.output_cleanup_guard(output_id);
4156
4157 let granularity = if run_bp {
4158 self.merge_granularity(&[seg_id.to_owned()]).await
4159 } else {
4160 crate::segment::reorder::BpGranularity::Auto
4161 };
4162 let reorder_result = crate::segment::reorder::reorder_segment(
4163 self.directory.as_ref(),
4164 &generation.schema,
4165 source_id,
4166 output_id,
4167 self.term_cache_blocks,
4168 self.term_cache_budget_bytes,
4169 self.bp_memory_budget_bytes,
4170 bp_budget,
4171 run_bp,
4172 granularity,
4173 self.optimization,
4174 self.posting_codec,
4175 generation.trained_vectors.as_deref(),
4176 self.term_dict_block_size,
4177 rayon_pool,
4178 Some(self.active_operations.cancellation_flag()),
4179 alive,
4180 )
4181 .await;
4182 let (new_id, total_docs, bp_converged) = match reorder_result {
4183 Ok(v) => v,
4184 Err(e) => {
4185 self.delete_output_if_unregistered(output_id, "reorder failure")
4188 .await;
4189 output_cleanup.disarm();
4190 if is_deterministic_source_error(&e) {
4191 self.quarantine_segment(seg_id, &e);
4192 } else if !matches!(&e, Error::IndexClosed) {
4193 self.pause_reorder_retries(seg_id, &e).await;
4194 }
4195 return Err(e);
4196 }
4197 };
4198
4199 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
4205 if let Err(e) = self
4206 .replace_segments(
4207 &[seg_id.to_string()],
4208 new_id,
4209 total_docs,
4210 if run_bp {
4211 ReplacementLayout::BpReordered {
4212 converged: ladder_converged,
4213 }
4214 } else {
4215 ReplacementLayout::MaintenanceOnly
4216 },
4217 None,
4218 )
4219 .await
4220 {
4221 self.delete_output_if_unregistered(output_id, "replacement failure")
4222 .await;
4223 output_cleanup.disarm();
4224 if !matches!(&e, Error::IndexClosed) {
4225 self.pause_reorder_retries(seg_id, &e).await;
4226 }
4227 return Err(e);
4228 }
4229 output_cleanup.disarm();
4230 self.clear_reorder_retry(seg_id);
4231
4232 Ok(true)
4233 }
4234
4235 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
4242 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
4243
4244 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
4245 for entry in entries {
4246 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
4247 continue;
4248 };
4249 let Some(rest) = filename.strip_prefix("seg_") else {
4250 continue;
4251 };
4252 let Some(hex_id) = rest.get(..32) else {
4253 continue;
4254 };
4255 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
4256 continue;
4257 }
4258 orphan_files
4259 .entry(hex_id.to_ascii_lowercase())
4260 .or_default()
4261 .push(entry);
4262 }
4263 }
4264
4265 let mut deleted = 0;
4266 for (hex_id, paths) in &orphan_files {
4267 let deletion_guard = {
4272 let st = self.state.lock().await;
4273 if st.metadata.owns_id(hex_id) {
4274 continue;
4275 }
4276 let Some(guard) = self
4277 .active_operations
4278 .try_register(vec![hex_id.to_string()])
4279 else {
4280 continue;
4281 };
4282 if self.tracker.is_deletion_protected(hex_id) {
4283 drop(guard);
4284 continue;
4285 }
4286 guard
4287 };
4288
4289 let results =
4294 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
4295 .await;
4296 let removed = results.into_iter().all(|result| match result {
4297 Ok(()) => true,
4298 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
4299 Err(error) => {
4300 log::warn!(
4301 "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
4302 self.schema.index_label(),
4303 hex_id,
4304 error,
4305 );
4306 false
4307 }
4308 });
4309 drop(deletion_guard);
4312 if removed {
4313 deleted += 1;
4314 log::info!(
4315 "[segment_cleanup] index={} swept orphan segment {}",
4316 self.schema.index_label(),
4317 hex_id
4318 );
4319 }
4320 }
4321
4322 Ok(deleted)
4323 }
4324}
4325
4326#[cfg(test)]
4327mod tests {
4328 use super::*;
4329 use std::sync::atomic::{AtomicBool, Ordering};
4330
4331 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
4332 let schema = crate::dsl::SchemaBuilder::default().build();
4333 let metadata = IndexMetadata::new(schema.clone());
4334 Arc::new(SegmentManager::new(
4335 Arc::new(crate::directories::RamDirectory::new()),
4336 Arc::new(schema),
4337 metadata,
4338 Box::new(crate::merge::NoMergePolicy),
4339 0,
4340 1,
4341 Arc::new(Semaphore::new(1)),
4342 None,
4343 1024,
4344 Arc::new(ReorderConcurrencyGate::new(1)),
4345 None,
4346 ))
4347 }
4348
4349 #[tokio::test]
4350 async fn queued_compaction_does_not_deadlock_a_vector_rewrite_holding_global_capacity() {
4351 use std::time::Duration;
4352 let manager = lifecycle_test_manager();
4353 let global = Arc::clone(&manager.global_merge_permits)
4354 .acquire_owned()
4355 .await
4356 .unwrap();
4357 let id = SegmentId::new().to_hex();
4358 let mut compaction = Box::pin(manager.compact_segment(&id, 1024 * 1024));
4359 assert!(
4360 tokio::time::timeout(Duration::from_millis(10), &mut compaction)
4361 .await
4362 .is_err()
4363 );
4364 let local = Arc::clone(&manager.merge_permits)
4367 .try_acquire_owned()
4368 .expect("waiting compaction must not reserve local capacity first");
4369 drop(local);
4370 drop(global);
4371 assert!(
4372 !tokio::time::timeout(Duration::from_secs(1), compaction)
4373 .await
4374 .unwrap()
4375 .unwrap()
4376 );
4377 }
4378
4379 #[test]
4380 fn force_merge_planner_pairs_large_and_small_segments() {
4381 let groups = plan_force_merge_groups(
4382 vec![
4383 ("a".into(), 6),
4384 ("b".into(), 6),
4385 ("c".into(), 4),
4386 ("d".into(), 4),
4387 ],
4388 10,
4389 );
4390
4391 assert_eq!(groups.len(), 2);
4392 assert!(groups.iter().all(|group| group.total_docs == 10));
4393 assert!(groups.iter().all(|group| group.segments.len() == 2));
4394 }
4395
4396 #[test]
4397 fn force_merge_planner_leaves_oversized_segments_alone() {
4398 let groups = plan_force_merge_groups(
4399 vec![
4400 ("oversized".into(), 11),
4401 ("small-a".into(), 5),
4402 ("small-b".into(), 5),
4403 ],
4404 10,
4405 );
4406
4407 assert_eq!(groups.len(), 2);
4408 assert_eq!(groups[0].total_docs, 10);
4409 assert_eq!(groups[0].segments.len(), 2);
4410 assert_eq!(groups[1].total_docs, 11);
4411 assert_eq!(groups[1].segments.len(), 1);
4412 }
4413
4414 #[test]
4415 fn force_merge_planner_never_exceeds_segment_format_limit() {
4416 let groups = plan_force_merge_groups(
4417 vec![
4418 ("large-a".into(), 3_000_000_000),
4419 ("large-b".into(), 2_000_000_000),
4420 ],
4421 u64::from(u32::MAX),
4422 );
4423 assert_eq!(groups.len(), 2);
4424 assert!(
4425 groups
4426 .iter()
4427 .all(|group| group.total_docs <= u64::from(u32::MAX))
4428 );
4429 }
4430
4431 #[test]
4432 fn force_merge_hierarchy_has_one_final_bp_pass() {
4433 assert_eq!(force_merge_output_count(1), 0);
4434 assert_eq!(force_merge_output_count(2), 1);
4435 assert_eq!(force_merge_output_count(64), 1);
4436 assert_eq!(force_merge_output_count(65), 2);
4437 assert_eq!(force_merge_output_count(127), 2);
4438 assert_eq!(force_merge_output_count(128), 3);
4439 assert_eq!(force_merge_output_count(1_000), 16);
4440 }
4441
4442 fn expand_force_merge_node(
4443 hierarchy: &ForceMergeHierarchy,
4444 source_count: usize,
4445 node: usize,
4446 sources: &mut Vec<usize>,
4447 ) {
4448 if node < source_count {
4449 sources.push(node);
4450 return;
4451 }
4452
4453 let step_index = node - source_count;
4454 let step = hierarchy
4455 .steps
4456 .get(step_index)
4457 .expect("merge input must refer to an existing source or output");
4458 for &input in &step.inputs {
4459 assert!(
4460 input < node,
4461 "merge step {step_index} refers to a future output node {input}"
4462 );
4463 expand_force_merge_node(hierarchy, source_count, input, sources);
4464 }
4465 }
4466
4467 #[test]
4468 fn force_merge_hierarchy_has_minimal_valid_arity() {
4469 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
4470
4471 for source_count in source_counts {
4472 let hierarchy = plan_force_merge_hierarchy(source_count);
4473 let output_count = hierarchy.steps.len();
4474
4475 assert!(
4476 hierarchy
4477 .steps
4478 .iter()
4479 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
4480 "invalid merge arity for {source_count} sources"
4481 );
4482 assert!(
4483 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
4484 "{output_count} outputs cannot reduce {source_count} sources"
4485 );
4486 assert!(
4487 output_count == 1
4488 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
4489 "{output_count} outputs are not minimal for {source_count} sources"
4490 );
4491 assert_eq!(output_count, force_merge_output_count(source_count));
4492 }
4493 }
4494
4495 #[test]
4496 fn force_merge_hierarchy_preserves_exact_source_order() {
4497 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
4498 let hierarchy = plan_force_merge_hierarchy(source_count);
4499 let mut sources = Vec::with_capacity(source_count);
4500 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
4501 assert_eq!(
4502 sources,
4503 (0..source_count).collect::<Vec<_>>(),
4504 "source order changed for {source_count} sources"
4505 );
4506 }
4507 }
4508
4509 fn force_merge_rewrite_cost(source_count: usize) -> usize {
4510 let hierarchy = plan_force_merge_hierarchy(source_count);
4511 let mut node_weights = vec![1usize; source_count];
4512 let mut rewrite_cost = 0usize;
4513
4514 for (step_index, step) in hierarchy.steps.iter().enumerate() {
4515 let output = source_count + step_index;
4516 let output_weight = step
4517 .inputs
4518 .iter()
4519 .map(|&input| {
4520 assert!(
4521 input < output,
4522 "merge step {step_index} refers to future output {input}"
4523 );
4524 node_weights[input]
4525 })
4526 .sum::<usize>();
4527 rewrite_cost += output_weight;
4528 node_weights.push(output_weight);
4529 }
4530
4531 assert_eq!(node_weights[hierarchy.root], source_count);
4532 rewrite_cost
4533 }
4534
4535 #[test]
4536 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
4537 assert_eq!(force_merge_rewrite_cost(65), 67);
4538 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
4539 }
4540
4541 #[test]
4542 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
4543 assert_eq!(
4544 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
4545 (false, false, 3),
4546 );
4547 assert_eq!(
4548 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
4549 (true, false, 4),
4550 );
4551 assert_eq!(
4552 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
4553 (true, true, 0),
4554 );
4555 }
4556
4557 #[tokio::test]
4558 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
4559 let mut schema_builder = crate::dsl::SchemaBuilder::default();
4560 let field = schema_builder.add_text_field("text", true, true);
4561 let schema = schema_builder.build();
4562 let directory = crate::directories::RamDirectory::new();
4563 let config = crate::index::IndexConfig {
4564 num_indexing_threads: 1,
4565 merge_policy: Box::new(crate::merge::NoMergePolicy),
4566 ..Default::default()
4567 };
4568 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
4569 .await
4570 .unwrap();
4571 for value in ["one", "two", "three"] {
4572 let mut document = crate::dsl::Document::new();
4573 document.add_text(field, value);
4574 writer.add_document(document).unwrap();
4575 writer.commit().await.unwrap();
4576 }
4577
4578 let manager = Arc::clone(writer.segment_manager());
4579 let held_id = manager.get_segment_ids().await.pop().unwrap();
4580 let mut held = Some(
4581 manager
4582 .active_operations
4583 .try_register(vec![held_id])
4584 .unwrap(),
4585 );
4586 let batches = Arc::new(AtomicUsize::new(0));
4587 let batch_count = Arc::clone(&batches);
4588 writer
4589 .force_merge_with_snapshot_refresh(move || {
4590 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
4591 if refresh == 2 {
4594 drop(held.take());
4595 }
4596 std::future::ready(Ok(()))
4597 })
4598 .await
4599 .unwrap();
4600
4601 assert_eq!(manager.get_segment_ids().await.len(), 1);
4602 assert_eq!(
4603 batches.load(Ordering::Relaxed),
4604 4,
4605 "initial/final refreshes plus two replacements are required after the held source releases"
4606 );
4607 }
4608
4609 #[tokio::test]
4610 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
4611 let mut schema_builder = crate::dsl::SchemaBuilder::default();
4612 schema_builder.set_reorder_on_merge(true);
4613 let schema = schema_builder.build();
4614 let mut metadata = IndexMetadata::new(schema.clone());
4615 metadata.add_segment("00000000000000000000000000000001".into(), 1);
4616 metadata.add_segment("00000000000000000000000000000002".into(), 1);
4617
4618 let global_merge_permits = Arc::new(Semaphore::new(1));
4619 let manager = Arc::new(SegmentManager::new(
4620 Arc::new(crate::directories::RamDirectory::new()),
4621 Arc::new(schema),
4622 metadata,
4623 Box::new(crate::merge::NoMergePolicy),
4624 0,
4625 1,
4626 Arc::clone(&global_merge_permits),
4627 None,
4628 1024,
4629 Arc::new(ReorderConcurrencyGate::new(1)),
4630 None,
4631 ));
4632
4633 manager.active_operations.pause_non_indexing();
4638 let force_merge = {
4639 let manager = Arc::clone(&manager);
4640 tokio::spawn(async move { manager.force_merge().await })
4641 };
4642 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4643 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4644 tokio::task::yield_now().await;
4645 }
4646 })
4647 .await
4648 .expect("force merge never reached the paused group claim");
4649
4650 assert_eq!(
4651 global_merge_permits.available_permits(),
4652 1,
4653 "force merge retained global capacity while vector staging blocked group ownership"
4654 );
4655
4656 force_merge.abort();
4657 let _ = force_merge.await;
4658 manager.active_operations.resume_non_indexing();
4659 }
4660
4661 #[test]
4662 fn output_cleanup_guard_runs_during_panic_unwind() {
4663 let cleaned = Arc::new(AtomicBool::new(false));
4664 let cleaned_in_callback = Arc::clone(&cleaned);
4665 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4666 cleaned_in_callback.store(true, Ordering::SeqCst);
4667 });
4668
4669 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4670 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4671 panic!("simulated reorder panic");
4672 }));
4673
4674 assert!(result.is_err());
4675 assert!(
4676 cleaned.load(Ordering::SeqCst),
4677 "partial output cleanup must run during unwind"
4678 );
4679 }
4680
4681 #[test]
4682 fn output_cleanup_guard_disarms_after_commit() {
4683 let cleaned = Arc::new(AtomicBool::new(false));
4684 let cleaned_in_callback = Arc::clone(&cleaned);
4685 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4686 cleaned_in_callback.store(true, Ordering::SeqCst);
4687 });
4688
4689 {
4690 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4691 guard.disarm();
4692 }
4693
4694 assert!(!cleaned.load(Ordering::SeqCst));
4695 }
4696
4697 #[test]
4698 fn test_active_operation_guard_releases_ownership() {
4699 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4700 {
4701 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4702 let snap = active.snapshot();
4703 assert!(snap.contains("a"));
4704 assert!(snap.contains("b"));
4705 }
4706 assert!(active.snapshot().is_empty());
4707 }
4708
4709 #[test]
4710 fn test_non_overlapping_operations_can_run_concurrently() {
4711 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4712 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4713 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4714 let snap = active.snapshot();
4715 assert_eq!(snap.len(), 4);
4716
4717 drop(first);
4718 let snap = active.snapshot();
4719 assert_eq!(snap.len(), 2);
4720 assert!(snap.contains("c"));
4721 assert!(snap.contains("d"));
4722 }
4723
4724 #[test]
4725 fn test_overlapping_operation_is_rejected_until_release() {
4726 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4727 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4728 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4729 drop(first);
4730 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4731 }
4732
4733 #[test]
4734 fn test_active_operation_snapshot() {
4735 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4736 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4737 let snap = active.snapshot();
4738 assert!(snap.contains("x"));
4739 assert!(snap.contains("y"));
4740 assert!(!snap.contains("z"));
4741 }
4742
4743 #[tokio::test]
4744 async fn operation_barrier_ignores_producers_started_after_snapshot() {
4745 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4746 let before_gate = active.try_register(vec!["old".into()]).unwrap();
4747 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4748 assert_eq!(parked_indexing, 0);
4749 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4750
4751 let waiter = {
4752 let active = Arc::clone(&active);
4753 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4754 };
4755 tokio::task::yield_now().await;
4756 assert!(!waiter.is_finished());
4757
4758 drop(before_gate);
4759 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4760 .await
4761 .expect("pre-gate operation barrier was starved by a post-gate producer")
4762 .unwrap();
4763 assert!(active.snapshot().contains("new-flat"));
4764 drop(after_gate);
4765 }
4766
4767 #[tokio::test]
4768 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4769 let manager = lifecycle_test_manager();
4770 let current = manager.published_generation();
4771 manager
4772 .published_generation
4773 .store(Arc::new(PublishedIndexGeneration {
4774 publication_id: current.publication_id,
4775 schema: current.schema.clone(),
4776 trained_vectors: Some(Arc::new(TrainedVectorStructures {
4777 centroids: rustc_hash::FxHashMap::default(),
4778 binary_quantizers: rustc_hash::FxHashMap::default(),
4779 ..Default::default()
4780 })),
4781 }));
4782
4783 let guard = manager.begin_vector_artifact_update().await.unwrap();
4784 assert!(
4785 manager.trained().is_some(),
4786 "search readers keep the last fully validated generation"
4787 );
4788 assert!(
4789 manager.trained_for_segment_build().is_none(),
4790 "new segment producers must stay flat during an artifact update"
4791 );
4792
4793 let detached_transaction_guard = guard.clone();
4794 drop(guard);
4795 assert!(
4796 manager.trained_for_segment_build().is_none(),
4797 "a detached lifecycle transaction must retain the producer gate after request cancellation"
4798 );
4799 drop(detached_transaction_guard);
4800 assert!(manager.trained_for_segment_build().is_some());
4801 }
4802
4803 #[tokio::test]
4804 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4805 let manager = lifecycle_test_manager();
4806 let guard = manager.begin_vector_artifact_update().await.unwrap();
4807 assert!(
4808 manager
4809 .active_operations
4810 .try_register(vec!["merge".into()])
4811 .is_none(),
4812 "ordinary merge/reorder work must not change staged sources"
4813 );
4814 let indexing = manager
4815 .active_operations
4816 .try_register_indexing(vec!["fresh".into()])
4817 .expect("indexing remains available in flat mode");
4818 drop(indexing);
4819
4820 drop(guard);
4821 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4822 assert!(
4823 manager
4824 .active_operations
4825 .try_register(vec!["merge".into()])
4826 .is_some()
4827 );
4828 }
4829
4830 #[tokio::test]
4831 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4832 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4833 let guard = active.try_register(vec!["live".into()]).unwrap();
4834 let cancellation = active.cancellation_flag();
4835 active.stop_accepting();
4836 assert!(cancellation.load(Ordering::Acquire));
4837 assert!(active.try_register(vec!["new".into()]).is_none());
4838
4839 let waiter = {
4840 let active = Arc::clone(&active);
4841 tokio::spawn(async move { active.wait_until_idle().await })
4842 };
4843 tokio::task::yield_now().await;
4844 assert!(!waiter.is_finished());
4845 drop(guard);
4846 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4847 .await
4848 .expect("shutdown waiter missed the final guard notification")
4849 .unwrap();
4850 }
4851
4852 #[tokio::test]
4853 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4854 let manager = lifecycle_test_manager();
4855 let started = Arc::new(Semaphore::new(0));
4856 let release = Arc::new(Semaphore::new(0));
4857 let completed = Arc::new(AtomicBool::new(false));
4858
4859 let request = {
4860 let manager = Arc::clone(&manager);
4861 let started = Arc::clone(&started);
4862 let release = Arc::clone(&release);
4863 let completed = Arc::clone(&completed);
4864 tokio::spawn(async move {
4865 manager
4866 .run_lifecycle_transaction(async move {
4867 started.add_permits(1);
4868 let _permit = release.acquire().await.unwrap();
4869 completed.store(true, Ordering::Release);
4870 Ok(())
4871 })
4872 .await
4873 })
4874 };
4875
4876 let _started = started.acquire().await.unwrap();
4877 request.abort();
4878 assert!(request.await.unwrap_err().is_cancelled());
4879 release.add_permits(1);
4880
4881 manager.begin_shutdown();
4882 tokio::time::timeout(
4883 std::time::Duration::from_secs(1),
4884 manager.wait_for_shutdown(),
4885 )
4886 .await
4887 .expect("shutdown did not drain detached lifecycle transaction");
4888 assert!(completed.load(Ordering::Acquire));
4889 }
4890
4891 #[tokio::test]
4892 async fn productive_seismic_maintenance_remains_eligible_beyond_three_passes() {
4893 let manager = lifecycle_test_manager();
4894 {
4895 let mut state = manager.state.lock().await;
4896 for (id, pending, attempts, stalls) in [
4897 ("eligible", 7, 100, 0),
4898 ("finished", 0, 1, 0),
4899 ("limited", 9, 3, 3),
4900 ] {
4901 state.metadata.add_segment_meta(
4902 id.into(),
4903 SegmentMetaInfo {
4904 deletions: None,
4905 num_docs: 10,
4906 ancestors: Vec::new(),
4907 generation: 1,
4908 reordered: true,
4909 bp_converged: true,
4910 bp_unconverged_passes: 0,
4911 seismic_pending_terms: pending,
4912 seismic_maintenance_passes: attempts,
4913 seismic_no_progress_passes: stalls,
4914 ann_fragmented: false,
4915 },
4916 );
4917 }
4918 }
4919 assert_eq!(
4920 manager.unconverged_segments_below(3).await,
4921 vec![("eligible".to_string(), 10, 0)]
4922 );
4923 {
4924 let mut state = manager.state.lock().await;
4925 let mixed = state.metadata.segment_metas.get_mut("eligible").unwrap();
4926 mixed.seismic_maintenance_passes = 0;
4927 mixed.bp_converged = false;
4928 mixed.bp_unconverged_passes = 3;
4929 }
4930 assert_eq!(
4931 manager.unconverged_segments_below(3).await,
4932 vec![("eligible".to_string(), 10, 0)],
4933 "inherited BMP cap must not hide newly copied Seismic debt"
4934 );
4935 assert!(manager.unreordered_segments().await.is_empty());
4936 }
4937
4938 #[test]
4939 fn seismic_replacements_preserve_stalls_until_progress_or_new_merge_inputs() {
4940 let parent = SegmentMetaInfo {
4941 deletions: None,
4942 num_docs: 10,
4943 ancestors: Vec::new(),
4944 generation: 1,
4945 reordered: true,
4946 bp_converged: true,
4947 bp_unconverged_passes: 0,
4948 seismic_pending_terms: 9,
4949 seismic_maintenance_passes: 100,
4950 seismic_no_progress_passes: 3,
4951 ann_fragmented: false,
4952 };
4953 for layout in [
4954 ReplacementLayout::PreserveSingleSource,
4955 ReplacementLayout::BlockCopy,
4956 ReplacementLayout::Compacted,
4957 ] {
4958 assert_eq!(
4959 replacement_seismic_state([&parent].into_iter(), 9, layout),
4960 (100, 3)
4961 );
4962 }
4963 assert_eq!(
4964 replacement_seismic_state([&parent].into_iter(), 8, ReplacementLayout::Compacted),
4965 (100, 0)
4966 );
4967 assert_eq!(
4968 replacement_seismic_state(
4969 [&parent].into_iter(),
4970 8,
4971 ReplacementLayout::BpReordered { converged: true }
4972 ),
4973 (101, 0)
4974 );
4975 assert_eq!(
4976 replacement_seismic_state(
4977 [&parent].into_iter(),
4978 9,
4979 ReplacementLayout::BpReordered { converged: true }
4980 ),
4981 (101, 4)
4982 );
4983 assert_eq!(
4984 replacement_seismic_state(
4985 [&parent, &parent].into_iter(),
4986 18,
4987 ReplacementLayout::BlockCopy
4988 ),
4989 (100, 0)
4990 );
4991 assert_eq!(
4992 replacement_seismic_state(
4993 [&parent].into_iter(),
4994 0,
4995 ReplacementLayout::BpReordered { converged: true }
4996 ),
4997 (0, 0)
4998 );
4999 let encoded = serde_json::to_value(&parent).unwrap();
5000 assert_eq!(encoded["seismic_no_progress_passes"], 3);
5001 let mut absent = encoded;
5002 absent
5003 .as_object_mut()
5004 .unwrap()
5005 .remove("seismic_no_progress_passes");
5006 assert_eq!(
5007 serde_json::from_value::<SegmentMetaInfo>(absent)
5008 .unwrap()
5009 .seismic_no_progress_passes,
5010 0
5011 );
5012 }
5013
5014 #[tokio::test]
5015 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
5016 let manager = lifecycle_test_manager();
5017 {
5018 let mut state = manager.state.lock().await;
5019 state.metadata.add_segment_meta(
5020 "eligible".into(),
5021 SegmentMetaInfo {
5022 deletions: None,
5023 num_docs: 10,
5024 ancestors: Vec::new(),
5025 generation: 1,
5026 reordered: true,
5027 bp_converged: false,
5028 bp_unconverged_passes: 2,
5029 seismic_pending_terms: 0,
5030 seismic_maintenance_passes: 0,
5031 seismic_no_progress_passes: 0,
5032 ann_fragmented: false,
5033 },
5034 );
5035 state.metadata.add_segment_meta(
5036 "at-limit".into(),
5037 SegmentMetaInfo {
5038 deletions: None,
5039 num_docs: 20,
5040 ancestors: Vec::new(),
5041 generation: 1,
5042 reordered: true,
5043 bp_converged: false,
5044 bp_unconverged_passes: 3,
5045 seismic_pending_terms: 0,
5046 seismic_maintenance_passes: 0,
5047 seismic_no_progress_passes: 0,
5048 ann_fragmented: false,
5049 },
5050 );
5051 state.metadata.add_segment_meta(
5052 "carried-debt".into(),
5053 SegmentMetaInfo {
5054 deletions: None,
5055 num_docs: 15,
5056 ancestors: Vec::new(),
5057 generation: 2,
5058 reordered: false,
5059 bp_converged: false,
5060 bp_unconverged_passes: 2,
5061 seismic_pending_terms: 0,
5062 seismic_maintenance_passes: 0,
5063 seismic_no_progress_passes: 0,
5064 ann_fragmented: false,
5065 },
5066 );
5067 state.metadata.add_segment_meta(
5068 "carried-debt-at-limit".into(),
5069 SegmentMetaInfo {
5070 deletions: None,
5071 num_docs: 25,
5072 ancestors: Vec::new(),
5073 generation: 2,
5074 reordered: false,
5075 bp_converged: false,
5076 bp_unconverged_passes: 3,
5077 seismic_pending_terms: 0,
5078 seismic_maintenance_passes: 0,
5079 seismic_no_progress_passes: 0,
5080 ann_fragmented: false,
5081 },
5082 );
5083 state.metadata.add_segment_meta(
5084 "converged".into(),
5085 SegmentMetaInfo {
5086 deletions: None,
5087 num_docs: 30,
5088 ancestors: Vec::new(),
5089 generation: 1,
5090 reordered: true,
5091 bp_converged: true,
5092 bp_unconverged_passes: 0,
5093 seismic_pending_terms: 0,
5094 seismic_maintenance_passes: 0,
5095 seismic_no_progress_passes: 0,
5096 ann_fragmented: false,
5097 },
5098 );
5099 state.metadata.add_segment("fresh".into(), 40);
5100 state
5101 .metadata
5102 .segment_metas
5103 .get_mut("fresh")
5104 .unwrap()
5105 .ann_fragmented = true;
5106 }
5107
5108 assert_eq!(
5109 manager.unreordered_segments().await,
5110 vec![("fresh".into(), 40)],
5111 "a block-copy output with BP debt is not a fresh first-pass candidate",
5112 );
5113 let mut eligible = manager.unconverged_segments_below(3).await;
5114 eligible.sort_unstable();
5115 assert_eq!(
5116 eligible,
5117 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
5118 );
5119 assert!(manager.unconverged_segments_below(0).await.is_empty());
5120 }
5121
5122 #[test]
5123 fn merge_retry_backoff_is_exponential_and_capped() {
5124 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
5125 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
5126 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
5127 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
5128 }
5129
5130 #[test]
5131 fn only_deterministic_source_errors_are_quarantined() {
5132 assert!(is_deterministic_source_error(&Error::Corruption(
5133 "bad footer".into()
5134 )));
5135 assert!(is_deterministic_source_error(&Error::Io(
5136 std::io::Error::from(std::io::ErrorKind::NotFound)
5137 )));
5138 assert!(!is_deterministic_source_error(&Error::Io(
5139 std::io::Error::from(std::io::ErrorKind::TimedOut)
5140 )));
5141 assert!(!is_deterministic_source_error(&Error::Io(
5142 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
5143 )));
5144 }
5145
5146 #[tokio::test]
5147 async fn transient_reorder_failure_is_backed_off_until_cleared() {
5148 let manager = lifecycle_test_manager();
5149 manager
5150 .state
5151 .lock()
5152 .await
5153 .metadata
5154 .add_segment("source".into(), 1);
5155 manager
5156 .pause_reorder_retries("source", &Error::Internal("transient".into()))
5157 .await;
5158 assert!(manager.paused_reorder_segments().contains("source"));
5159 manager.clear_reorder_retry("source");
5160 assert!(!manager.paused_reorder_segments().contains("source"));
5161 }
5162
5163 #[tokio::test]
5164 async fn optimizer_backoff_does_not_retain_retired_or_unknown_segments() {
5165 let mut schema = crate::SchemaBuilder::default();
5166 let key = schema.add_text_field("id", true, false);
5167 schema.set_primary_key(key);
5168 let index = crate::Index::create(
5169 crate::RamDirectory::new(),
5170 schema.build(),
5171 crate::IndexConfig {
5172 merge_policy: Box::new(crate::NoMergePolicy),
5173 ..Default::default()
5174 },
5175 )
5176 .await
5177 .unwrap();
5178 let mut writer = index.writer();
5179 writer.init_primary_key_dedup().await.unwrap();
5180 for value in ["dead", "live"] {
5181 let mut doc = crate::Document::new();
5182 doc.add_text(key, value);
5183 writer.add_document(doc).unwrap();
5184 }
5185 writer.commit().await.unwrap();
5186 writer.delete_primary_key("dead").unwrap();
5187 writer.commit().await.unwrap();
5188 let manager = Arc::clone(writer.segment_manager());
5189 let source = manager.get_segment_ids().await[0].clone();
5190 assert!(
5191 manager
5192 .compact_segment_if_eligible(&source, 0.3, 0)
5193 .await
5194 .is_err()
5195 );
5196 assert!(manager.reorder_retries.lock().contains_key(&source));
5197 let mut doc = crate::Document::new();
5198 doc.add_text(key, "another");
5199 writer.add_document(doc).unwrap();
5200 writer.force_merge().await.unwrap();
5201 assert!(
5202 manager.reorder_retries.lock().is_empty(),
5203 "retired segment leaked retry state"
5204 );
5205 assert!(
5207 manager
5208 .compact_segment_if_eligible(&source, 0.3, 0)
5209 .await
5210 .is_err()
5211 );
5212 assert!(manager.reorder_retries.lock().is_empty());
5213 }
5214
5215 #[derive(Default)]
5218 struct FailingExistsDirectory(crate::directories::RamDirectory);
5219
5220 #[async_trait::async_trait]
5221 impl crate::directories::Directory for FailingExistsDirectory {
5222 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
5223 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
5224 }
5225
5226 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
5227 self.0.file_size(path).await
5228 }
5229
5230 async fn open_read(
5231 &self,
5232 path: &std::path::Path,
5233 ) -> std::io::Result<crate::directories::FileHandle> {
5234 self.0.open_read(path).await
5235 }
5236
5237 async fn read_range(
5238 &self,
5239 path: &std::path::Path,
5240 range: std::ops::Range<u64>,
5241 ) -> std::io::Result<crate::directories::OwnedBytes> {
5242 self.0.read_range(path, range).await
5243 }
5244
5245 async fn list_files(
5246 &self,
5247 prefix: &std::path::Path,
5248 ) -> std::io::Result<Vec<std::path::PathBuf>> {
5249 self.0.list_files(prefix).await
5250 }
5251
5252 async fn open_lazy(
5253 &self,
5254 path: &std::path::Path,
5255 ) -> std::io::Result<crate::directories::FileHandle> {
5256 self.0.open_lazy(path).await
5257 }
5258 }
5259
5260 #[async_trait::async_trait]
5261 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
5262 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
5263 self.0.write(path, data).await
5264 }
5265
5266 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
5267 self.0.delete(path).await
5268 }
5269
5270 async fn rename(
5271 &self,
5272 from: &std::path::Path,
5273 to: &std::path::Path,
5274 ) -> std::io::Result<()> {
5275 self.0.rename(from, to).await
5276 }
5277
5278 async fn sync(&self) -> std::io::Result<()> {
5279 self.0.sync().await
5280 }
5281
5282 async fn streaming_writer(
5283 &self,
5284 path: &std::path::Path,
5285 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
5286 self.0.streaming_writer(path).await
5287 }
5288 }
5289
5290 #[derive(Debug, Clone)]
5291 struct MergeEverythingPolicy;
5292
5293 impl MergePolicy for MergeEverythingPolicy {
5294 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
5295 if segments.len() < 2 {
5296 return Vec::new();
5297 }
5298 vec![crate::merge::MergeCandidate {
5299 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
5300 }]
5301 }
5302
5303 fn clone_box(&self) -> Box<dyn MergePolicy> {
5304 Box::new(self.clone())
5305 }
5306 }
5307
5308 #[tokio::test]
5309 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
5310 let manager = lifecycle_test_manager();
5311 let parked_indexing = manager
5316 .protect_new_segment("00000000000000000000000000000abc".into())
5317 .unwrap();
5318
5319 let error = tokio::time::timeout(
5320 std::time::Duration::from_secs(2),
5321 manager.begin_vector_artifact_update(),
5322 )
5323 .await
5324 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
5325 .err()
5326 .expect("an old-generation prepared segment must block artifact replacement")
5327 .to_string();
5328 assert!(error.contains("built but uncommitted"), "{error}");
5329 assert!(
5330 !manager.vector_artifact_update.load(Ordering::Acquire),
5331 "a rejected update must release the producer gate"
5332 );
5333
5334 drop(parked_indexing);
5335
5336 let guard = manager
5337 .begin_vector_artifact_update()
5338 .await
5339 .expect("artifact update should succeed after the pending generation is resolved");
5340 drop(guard);
5341 }
5342
5343 #[tokio::test]
5344 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
5345 let manager = lifecycle_test_manager();
5346 let merge_like = manager
5347 .active_operations
5348 .try_register(vec!["merge-source".into()])
5349 .unwrap();
5350
5351 let waiter = {
5352 let manager = Arc::clone(&manager);
5353 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
5354 };
5355 for _ in 0..8 {
5356 tokio::task::yield_now().await;
5357 }
5358 assert!(
5359 !waiter.is_finished(),
5360 "artifact update must drain merge/reorder producers that may hold the previous generation"
5361 );
5362
5363 drop(merge_like);
5364 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
5365 .await
5366 .expect("artifact update missed the lifecycle guard release")
5367 .unwrap()
5368 .unwrap();
5369 }
5370
5371 #[tokio::test]
5372 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
5373 let manager = lifecycle_test_manager();
5374 let release = Arc::new(Semaphore::new(0));
5375 let merge_task = {
5376 let release = Arc::clone(&release);
5377 tokio::spawn(async move {
5378 let _permit = release.acquire().await.unwrap();
5379 })
5380 };
5381 manager.merge_handles.lock().push(merge_task);
5382
5383 let waiter = {
5384 let manager = Arc::clone(&manager);
5385 tokio::spawn(async move { manager.wait_for_all_merges().await })
5386 };
5387 for _ in 0..8 {
5388 tokio::task::yield_now().await;
5389 }
5390 assert!(!waiter.is_finished());
5391 waiter.abort();
5394 let join_error = waiter.await.unwrap_err();
5395 assert!(join_error.is_cancelled());
5396
5397 assert!(
5398 !manager.merge_handles.lock().is_empty(),
5399 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
5400 );
5401
5402 release.add_permits(1);
5404 tokio::time::timeout(
5405 std::time::Duration::from_secs(1),
5406 manager.wait_for_all_merges(),
5407 )
5408 .await
5409 .expect("subsequent drain missed the reinserted merge handle");
5410 assert!(manager.merge_handles.lock().is_empty());
5411 }
5412
5413 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5414 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
5415 let manager = lifecycle_test_manager();
5416 {
5417 let mut state = manager.state.lock().await;
5418 state
5419 .metadata
5420 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
5421 state
5422 .metadata
5423 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
5424 }
5425 let reorder_like = manager
5428 .active_operations
5429 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
5430 .unwrap();
5431
5432 let force_merge = {
5433 let manager = Arc::clone(&manager);
5434 tokio::spawn(async move { manager.force_merge().await })
5435 };
5436
5437 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
5438 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
5439 assert!(
5440 retries >= 1,
5441 "force_merge never observed the conflicting owner (retries={retries})"
5442 );
5443 assert!(
5444 retries < 20,
5445 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
5446 );
5447
5448 drop(reorder_like);
5449 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
5452 .await
5453 .expect("force_merge kept spinning after the conflicting owner released")
5454 .unwrap();
5455 assert!(result.is_err());
5456 }
5457
5458 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5459 async fn force_merge_routes_around_segments_held_by_reorder() {
5460 let manager = lifecycle_test_manager();
5461 {
5462 let mut state = manager.state.lock().await;
5463 state
5464 .metadata
5465 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
5466 state
5467 .metadata
5468 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
5469 state
5470 .metadata
5471 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
5472 }
5473 let _reorder_like = manager
5476 .active_operations
5477 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
5478 .unwrap();
5479
5480 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
5487 let manager = Arc::clone(&manager);
5488 async move { manager.force_merge().await }
5489 })
5490 .await
5491 .expect("force_merge livelocked on a segment held by an active reorder");
5492 assert!(result.is_err(), "fake segment files must fail the merge");
5493
5494 assert_eq!(
5495 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
5496 0,
5497 "batch built from the ownership snapshot must not collide with the held segment"
5498 );
5499 }
5500
5501 #[tokio::test]
5502 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
5503 let schema = crate::dsl::SchemaBuilder::default().build();
5504 let mut metadata = IndexMetadata::new(schema.clone());
5505 metadata.add_segment("00000000000000000000000000000001".into(), 10);
5506 metadata.add_segment("00000000000000000000000000000002".into(), 10);
5507 let manager = Arc::new(SegmentManager::new(
5508 Arc::new(FailingExistsDirectory::default()),
5509 Arc::new(schema),
5510 metadata,
5511 Box::new(MergeEverythingPolicy),
5512 0,
5513 1,
5514 Arc::new(Semaphore::new(1)),
5515 None,
5516 1024,
5517 Arc::new(ReorderConcurrencyGate::new(1)),
5518 None,
5519 ));
5520
5521 manager.maybe_merge().await;
5524
5525 tokio::time::timeout(
5526 std::time::Duration::from_secs(5),
5527 manager.wait_for_all_merges(),
5528 )
5529 .await
5530 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
5531 assert!(
5532 manager.merge_retry_is_paused(),
5533 "the failed merge should have armed the retry backoff"
5534 );
5535
5536 manager.begin_shutdown();
5538 tokio::time::timeout(
5539 std::time::Duration::from_secs(5),
5540 manager.wait_for_shutdown(),
5541 )
5542 .await
5543 .expect("shutdown did not drain the merge retry wakeup task");
5544 }
5545}
5546#[path = "row_mutation.rs"]
5547mod row_mutation;