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