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 rewrite_existing: bool,
2757 ) -> Result<bool> {
2758 for &field_id in field_ids {
2759 let flat = reader.flat_vectors().get(&field_id);
2760 let ann = reader.vector_indexes().get(&field_id);
2761 if ann.is_some() && flat.is_none() {
2762 return Err(Error::Corruption(format!(
2763 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2764 reader.meta().id,
2765 )));
2766 }
2767
2768 let Some(flat) = flat else {
2769 continue;
2770 };
2771 if flat.num_vectors == 0 {
2772 continue;
2773 }
2774 if rewrite_existing {
2775 return Ok(true);
2776 }
2777 let field = crate::dsl::Field(field_id);
2778 let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2779 Error::Corruption(format!(
2780 "segment {:032x} references unknown vector field {field_id}",
2781 reader.meta().id,
2782 ))
2783 })?;
2784 let current = match entry.field_type {
2785 crate::dsl::FieldType::DenseVector
2789 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2790 config.index_type == crate::dsl::VectorIndexType::Tq
2791 }) =>
2792 {
2793 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
2794 }
2795 crate::dsl::FieldType::DenseVector
2796 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2797 config.index_type == crate::dsl::VectorIndexType::IvfTq
2798 }) =>
2799 {
2800 matches!(ann, Some(crate::segment::VectorIndex::IvfTq { .. }))
2801 }
2802 crate::dsl::FieldType::DenseVector => false,
2805 crate::dsl::FieldType::BinaryDenseVector => {
2806 matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2807 }
2808 _ => false,
2809 };
2810 if !current {
2811 return Ok(true);
2812 }
2813 }
2814 Ok(false)
2815 }
2816
2817 async fn acquire_vector_rewrite_capacity(
2818 &self,
2819 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2820 let global = tokio::select! {
2821 biased;
2822 () = self.active_operations.wait_for_shutdown() => {
2823 return Err(Error::IndexClosed);
2824 }
2825 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2826 permit.map_err(|_| Error::Internal(
2827 "global background merge scheduler is closed".into()
2828 ))?
2829 }
2830 };
2831 let local = tokio::select! {
2832 biased;
2833 () = self.active_operations.wait_for_shutdown() => {
2834 return Err(Error::IndexClosed);
2835 }
2836 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2837 permit.map_err(|_| Error::Internal(
2838 "background merge scheduler is closed".into()
2839 ))?
2840 }
2841 };
2842 Ok((global, local))
2843 }
2844
2845 async fn build_vector_replacement(
2846 self: &Arc<Self>,
2847 segment_id: &str,
2848 source_id: SegmentId,
2849 output_id: SegmentId,
2850 trained: &TrainedVectorStructures,
2851 failure_context: &'static str,
2852 ) -> Result<(String, u32, OutputCleanupGuard)> {
2853 let mut cleanup = self.output_cleanup_guard(output_id);
2854 match crate::segment::reorder::rewrite_vector_segment(
2855 self.directory.as_ref(),
2856 &self.schema,
2857 source_id,
2858 output_id,
2859 self.term_cache_blocks,
2860 trained,
2861 Some(self.background_cpu_pool()),
2862 )
2863 .await
2864 {
2865 Ok((new_id, doc_count)) => {
2866 self.validate_completed_segment(&new_id, doc_count).await?;
2867 Ok((new_id, doc_count, cleanup))
2868 }
2869 Err(error) => {
2870 self.delete_output_if_unregistered(output_id, failure_context)
2871 .await;
2872 cleanup.disarm();
2873 if is_deterministic_source_error(&error) {
2874 self.quarantine_segment(segment_id, &error);
2875 }
2876 Err(error)
2877 }
2878 }
2879 }
2880
2881 pub(crate) async fn stage_vector_generation(
2885 self: &Arc<Self>,
2886 _artifact_update: &VectorArtifactUpdateGuard,
2887 segment_ids: &[String],
2888 field_ids: &[u32],
2889 trained: Arc<TrainedVectorStructures>,
2890 rewrite_existing: bool,
2891 ) -> Result<Vec<StagedVectorSegment>> {
2892 if !self.vector_artifact_update.load(Ordering::Acquire) {
2893 return Err(Error::Internal(
2894 "cannot stage a vector generation without an exclusive update lease".into(),
2895 ));
2896 }
2897
2898 let mut staged = Vec::new();
2899 for segment_id in segment_ids {
2900 if self.quarantined_segments.lock().contains(segment_id) {
2901 return Err(Error::Corruption(format!(
2902 "segment {segment_id} is quarantined after a deterministic source failure"
2903 )));
2904 }
2905 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2906 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2907 })?;
2908
2909 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2912
2913 let output_id = SegmentId::new();
2914 let output_hex = output_id.to_hex();
2915 let operation = {
2916 let st = self.state.lock().await;
2917 if !st.metadata.has_segment(segment_id) {
2918 return Err(Error::Corruption(format!(
2919 "vector generation source {segment_id} disappeared while lifecycle work was paused"
2920 )));
2921 }
2922 self.active_operations
2923 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
2924 }
2925 .ok_or_else(|| {
2926 if self.active_operations.is_accepting() {
2927 Error::Internal(format!(
2928 "vector generation could not claim stable source {segment_id}"
2929 ))
2930 } else {
2931 Error::IndexClosed
2932 }
2933 })?;
2934
2935 let reader = SegmentReader::open(
2936 self.directory.as_ref(),
2937 source_id,
2938 Arc::clone(&self.schema),
2939 self.term_cache_blocks,
2940 )
2941 .await?;
2942 if !self.segment_needs_vector_rewrite(&reader, field_ids, rewrite_existing)? {
2943 continue;
2944 }
2945 drop(reader);
2946
2947 let (new_id, doc_count, cleanup) = self
2948 .build_vector_replacement(
2949 segment_id,
2950 source_id,
2951 output_id,
2952 trained.as_ref(),
2953 "vector generation staging failure",
2954 )
2955 .await?;
2956 debug_assert_eq!(new_id, output_hex);
2957 let output_reader = SegmentReader::open(
2958 self.directory.as_ref(),
2959 output_id,
2960 Arc::clone(&self.schema),
2961 self.term_cache_blocks,
2962 )
2963 .await?;
2964 if self.segment_needs_vector_rewrite(&output_reader, field_ids, false)? {
2965 return Err(Error::Corruption(format!(
2966 "staged vector segment {new_id} does not match its candidate codebook generation"
2967 )));
2968 }
2969
2970 staged.push(StagedVectorSegment {
2971 source_id: segment_id.clone(),
2972 output_id,
2973 doc_count,
2974 _operation: operation,
2975 cleanup,
2976 });
2977 }
2978 Ok(staged)
2979 }
2980
2981 async fn rewrite_vector_segment_once(
2982 self: &Arc<Self>,
2983 segment_id: &str,
2984 field_ids: &[u32],
2985 ) -> Result<VectorSegmentRewriteOutcome> {
2986 if self.quarantined_segments.lock().contains(segment_id) {
2987 return Err(Error::Corruption(format!(
2988 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
2989 )));
2990 }
2991 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2992 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2993 })?;
2994
2995 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3000
3001 let output_id = SegmentId::new();
3002 let output_hex = output_id.to_hex();
3003 let all_ids = vec![segment_id.to_owned(), output_hex];
3004 let operation = {
3005 let st = self.state.lock().await;
3006 if !st.metadata.has_segment(segment_id) {
3007 return Ok(VectorSegmentRewriteOutcome::SourceGone);
3008 }
3009 self.active_operations.try_register(all_ids)
3010 };
3011 let _operation = match operation {
3012 Some(operation) => operation,
3013 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3014 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3015 };
3016
3017 let Some(trained) = self.trained_for_segment_build() else {
3018 return Ok(VectorSegmentRewriteOutcome::Deferred);
3019 };
3020
3021 let reader = SegmentReader::open(
3022 self.directory.as_ref(),
3023 source_id,
3024 Arc::clone(&self.schema),
3025 self.term_cache_blocks,
3026 )
3027 .await?;
3028 if !self.segment_needs_vector_rewrite(&reader, field_ids, false)? {
3029 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3030 }
3031 drop(reader);
3032
3033 let (new_id, doc_count, mut output_cleanup) = self
3034 .build_vector_replacement(
3035 segment_id,
3036 source_id,
3037 output_id,
3038 trained.as_ref(),
3039 "vector rewrite failure",
3040 )
3041 .await?;
3042
3043 if let Err(error) = self
3044 .replace_segments(
3045 &[segment_id.to_owned()],
3046 new_id,
3047 doc_count,
3048 ReplacementLayout::PreserveSingleSource,
3049 )
3050 .await
3051 {
3052 self.delete_output_if_unregistered(output_id, "vector replacement failure")
3053 .await;
3054 output_cleanup.disarm();
3055 return Err(error);
3056 }
3057 output_cleanup.disarm();
3058 Ok(VectorSegmentRewriteOutcome::Rewritten)
3059 }
3060
3061 pub(crate) async fn rewrite_vector_segments(
3066 self: &Arc<Self>,
3067 field_ids: &[u32],
3068 ) -> Result<usize> {
3069 if field_ids.is_empty() {
3070 return Ok(0);
3071 }
3072 let mut rewritten = 0usize;
3073 loop {
3074 let segment_ids = self.get_segment_ids().await;
3075 let mut conflicted = false;
3076 let mut changed = false;
3077 for segment_id in segment_ids {
3078 match self
3079 .rewrite_vector_segment_once(&segment_id, field_ids)
3080 .await?
3081 {
3082 VectorSegmentRewriteOutcome::Rewritten => {
3083 rewritten += 1;
3084 changed = true;
3085 }
3086 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3087 VectorSegmentRewriteOutcome::Deferred => {
3088 return Err(Error::Internal(
3089 "ANN finalization lost the published trained generation".into(),
3090 ));
3091 }
3092 VectorSegmentRewriteOutcome::AlreadyCurrent
3093 | VectorSegmentRewriteOutcome::SourceGone => {}
3094 }
3095 }
3096 if !conflicted && !changed {
3097 log::info!(
3098 "[dense_vector_rewrite] ANN finalization complete ({} segment(s) rewritten)",
3099 rewritten,
3100 );
3101 return Ok(rewritten);
3102 }
3103 tokio::select! {
3104 biased;
3105 () = self.active_operations.wait_for_shutdown() => {
3106 return Err(Error::IndexClosed);
3107 }
3108 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3109 }
3110 }
3111 }
3112
3113 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3118 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3119 return;
3120 }
3121 let manager = Arc::clone(self);
3122 let future = async move {
3123 let field_ids = manager
3124 .read_metadata(|metadata| {
3125 metadata
3126 .vector_fields
3127 .keys()
3128 .filter(|field_id| metadata.is_field_built(**field_id))
3129 .copied()
3130 .collect::<Vec<_>>()
3131 })
3132 .await;
3133 for segment_id in segment_ids {
3134 loop {
3135 match manager
3136 .rewrite_vector_segment_once(&segment_id, &field_ids)
3137 .await
3138 {
3139 Ok(VectorSegmentRewriteOutcome::Conflict) => {
3140 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3141 }
3142 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3143 Ok(_) => break,
3144 Err(error) => {
3145 log::error!(
3146 "[dense_vector_rewrite] failed to upgrade newly committed segment {}: {}",
3147 segment_id,
3148 error,
3149 );
3150 break;
3151 }
3152 }
3153 }
3154 }
3155 };
3156 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3157 log::warn!(
3158 "[dense_vector_rewrite] runtime unavailable; newly committed flat segment upgrade deferred"
3159 );
3160 return;
3161 };
3162 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3163 log::warn!(
3164 "[dense_vector_rewrite] runtime rejected newly committed flat segment upgrade"
3165 );
3166 }
3167 }
3168
3169 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3176 self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3177 .await
3178 }
3179
3180 pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3183 self: &Arc<Self>,
3184 mut refresh_snapshots: F,
3185 ) -> Result<()>
3186 where
3187 F: FnMut() -> Fut,
3188 Fut: std::future::Future<Output = Result<()>>,
3189 {
3190 self.wait_for_all_merges().await;
3191 refresh_snapshots().await?;
3192 let segment_ids = self.get_segment_ids().await;
3193
3194 if segment_ids.is_empty() {
3195 log::info!("[reorder] no segments to reorder");
3196 return Ok(());
3197 }
3198
3199 log::info!("[reorder] reordering {} segments", segment_ids.len());
3200
3201 for seg_id in segment_ids {
3202 match self
3203 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3204 .await
3205 {
3206 Ok(true) => refresh_snapshots().await?,
3207 Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
3208 Err(e) => return Err(e),
3209 }
3210 }
3211
3212 refresh_snapshots().await?;
3215 log::info!("[reorder] all segments reordered");
3216 Ok(())
3217 }
3218
3219 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3224 self.unreordered_segments()
3225 .await
3226 .into_iter()
3227 .map(|(id, _)| id)
3228 .collect()
3229 }
3230
3231 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3234 let quarantined = self.quarantined_segments.lock().clone();
3235 let paused = self.paused_reorder_segments();
3236 let st = self.state.lock().await;
3237 let active_ids = self.active_operations.snapshot();
3238 st.metadata
3239 .segment_metas
3240 .iter()
3241 .filter(|(id, info)| {
3242 !info.reordered
3243 && info.bp_converged
3244 && !active_ids.contains(*id)
3245 && !quarantined.contains(*id)
3246 && !paused.contains(*id)
3247 })
3248 .map(|(id, info)| (id.clone(), info.num_docs))
3249 .collect()
3250 }
3251
3252 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3256 self.unconverged_segments_below(u32::MAX)
3257 .await
3258 .into_iter()
3259 .map(|(id, docs, _)| (id, docs))
3260 .collect()
3261 }
3262
3263 pub async fn unconverged_segments_below(
3266 &self,
3267 max_unconverged_passes: u32,
3268 ) -> Vec<(String, u32, 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.bp_converged
3278 && info.bp_unconverged_passes < max_unconverged_passes
3279 && !active_ids.contains(*id)
3280 && !quarantined.contains(*id)
3281 && !paused.contains(*id)
3282 })
3283 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3284 .collect()
3285 }
3286
3287 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3297 let st = self.state.lock().await;
3298 let deepening = ids.iter().any(|id| {
3299 st.metadata
3300 .segment_metas
3301 .get(id)
3302 .is_some_and(|info| !info.bp_converged)
3303 });
3304 drop(st);
3305 if deepening {
3306 log::info!(
3307 "[reorder] source BP lineage unconverged — forcing record-level BP (deepening pass)",
3308 );
3309 crate::segment::reorder::BpGranularity::Records
3310 } else {
3311 crate::segment::reorder::BpGranularity::Auto
3312 }
3313 }
3314
3315 pub async fn reorder_single_segment(
3320 self: &Arc<Self>,
3321 seg_id: &str,
3322 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3323 bp_budget: crate::segment::BpBudget,
3324 ) -> Result<bool> {
3325 let source_id = SegmentId::from_hex(seg_id)
3326 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3327 if self.quarantined_segments.lock().contains(seg_id) {
3328 return Err(Error::Corruption(format!(
3329 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3330 seg_id
3331 )));
3332 }
3333 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3334 log::debug!(
3335 "[optimizer] explicit force merge active, skipping reorder of {}",
3336 seg_id,
3337 );
3338 return Ok(false);
3339 }
3340
3341 let reorder_gate = Arc::clone(&self.reorder_permits);
3346 let _reorder_permit = tokio::select! {
3347 biased;
3348 () = self.active_operations.wait_for_shutdown() => {
3349 return Err(Error::IndexClosed);
3350 }
3351 permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3352 permit.map_err(|_| {
3353 Error::Internal("background reorder scheduler is closed".into())
3354 })?
3355 }
3356 };
3357
3358 let output_id = SegmentId::new();
3359 let output_hex = output_id.to_hex();
3360 let source_ids = [seg_id.to_string()];
3361 let granularity = self.merge_granularity(&source_ids).await;
3362
3363 let all_ids = vec![seg_id.to_string(), output_hex];
3369 let (_guard, source_docs) = {
3370 let st = self.state.lock().await;
3371 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3375 log::debug!(
3376 "[optimizer] explicit force merge active, skipping reorder of {}",
3377 seg_id,
3378 );
3379 return Ok(false);
3380 }
3381 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3382 log::info!(
3383 "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
3384 seg_id
3385 );
3386 self.clear_reorder_retry(seg_id);
3387 return Ok(false);
3388 };
3389
3390 match self.active_operations.try_register(all_ids) {
3391 Some(guard) => (guard, source_meta.num_docs),
3392 None if !self.active_operations.is_accepting() => {
3393 return Err(Error::IndexClosed);
3394 }
3395 None => {
3396 log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
3397 return Ok(false);
3398 }
3399 }
3400 };
3401
3402 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3407 if is_deterministic_source_error(&error) {
3408 self.quarantine_segment(seg_id, &error);
3409 } else if !matches!(&error, Error::IndexClosed) {
3410 self.pause_reorder_retries(seg_id, &error);
3411 }
3412 return Err(error);
3413 }
3414
3415 let mut output_cleanup = self.output_cleanup_guard(output_id);
3416
3417 let reorder_result = crate::segment::reorder::reorder_segment(
3418 self.directory.as_ref(),
3419 &self.schema,
3420 source_id,
3421 output_id,
3422 self.term_cache_blocks,
3423 self.bp_memory_budget_bytes,
3424 bp_budget,
3425 granularity,
3426 rayon_pool,
3427 Some(self.active_operations.cancellation_flag()),
3428 )
3429 .await;
3430 let (new_id, total_docs, bp_converged) = match reorder_result {
3431 Ok(v) => v,
3432 Err(e) => {
3433 self.delete_output_if_unregistered(output_id, "reorder failure")
3436 .await;
3437 output_cleanup.disarm();
3438 if is_deterministic_source_error(&e) {
3439 self.quarantine_segment(seg_id, &e);
3440 } else if !matches!(&e, Error::IndexClosed) {
3441 self.pause_reorder_retries(seg_id, &e);
3442 }
3443 return Err(e);
3444 }
3445 };
3446
3447 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3453 if let Err(e) = self
3454 .replace_segments(
3455 &[seg_id.to_string()],
3456 new_id,
3457 total_docs,
3458 ReplacementLayout::BpReordered {
3459 converged: ladder_converged,
3460 },
3461 )
3462 .await
3463 {
3464 self.delete_output_if_unregistered(output_id, "replacement failure")
3465 .await;
3466 output_cleanup.disarm();
3467 if !matches!(&e, Error::IndexClosed) {
3468 self.pause_reorder_retries(seg_id, &e);
3469 }
3470 return Err(e);
3471 }
3472 output_cleanup.disarm();
3473 self.clear_reorder_retry(seg_id);
3474
3475 Ok(true)
3476 }
3477
3478 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3485 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3486
3487 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3488 for entry in entries {
3489 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3490 continue;
3491 };
3492 let Some(rest) = filename.strip_prefix("seg_") else {
3493 continue;
3494 };
3495 let Some(hex_id) = rest.get(..32) else {
3496 continue;
3497 };
3498 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3499 continue;
3500 }
3501 orphan_files
3502 .entry(hex_id.to_ascii_lowercase())
3503 .or_default()
3504 .push(entry);
3505 }
3506 }
3507
3508 let mut deleted = 0;
3509 for (hex_id, paths) in &orphan_files {
3510 let deletion_guard = {
3515 let st = self.state.lock().await;
3516 if st.metadata.has_segment(hex_id) {
3517 continue;
3518 }
3519 let Some(guard) = self
3520 .active_operations
3521 .try_register(vec![hex_id.to_string()])
3522 else {
3523 continue;
3524 };
3525 if self.tracker.is_deletion_protected(hex_id) {
3526 drop(guard);
3527 continue;
3528 }
3529 guard
3530 };
3531
3532 let results =
3537 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3538 .await;
3539 let removed = results.into_iter().all(|result| match result {
3540 Ok(()) => true,
3541 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3542 Err(error) => {
3543 log::warn!(
3544 "[segment_cleanup] failed sweeping orphan segment {}: {}",
3545 hex_id,
3546 error,
3547 );
3548 false
3549 }
3550 });
3551 drop(deletion_guard);
3554 if removed {
3555 deleted += 1;
3556 log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
3557 }
3558 }
3559
3560 Ok(deleted)
3561 }
3562}
3563
3564#[cfg(test)]
3565mod tests {
3566 use super::*;
3567 use std::sync::atomic::{AtomicBool, Ordering};
3568
3569 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3570 let schema = crate::dsl::SchemaBuilder::default().build();
3571 let metadata = IndexMetadata::new(schema.clone());
3572 Arc::new(SegmentManager::new(
3573 Arc::new(crate::directories::RamDirectory::new()),
3574 Arc::new(schema),
3575 metadata,
3576 Box::new(crate::merge::NoMergePolicy),
3577 0,
3578 1,
3579 Arc::new(Semaphore::new(1)),
3580 None,
3581 1024,
3582 Arc::new(ReorderConcurrencyGate::new(1)),
3583 None,
3584 ))
3585 }
3586
3587 #[test]
3588 fn force_merge_planner_pairs_large_and_small_segments() {
3589 let groups = plan_force_merge_groups(
3590 vec![
3591 ("a".into(), 6),
3592 ("b".into(), 6),
3593 ("c".into(), 4),
3594 ("d".into(), 4),
3595 ],
3596 10,
3597 );
3598
3599 assert_eq!(groups.len(), 2);
3600 assert!(groups.iter().all(|group| group.total_docs == 10));
3601 assert!(groups.iter().all(|group| group.segments.len() == 2));
3602 }
3603
3604 #[test]
3605 fn force_merge_planner_leaves_oversized_segments_alone() {
3606 let groups = plan_force_merge_groups(
3607 vec![
3608 ("oversized".into(), 11),
3609 ("small-a".into(), 5),
3610 ("small-b".into(), 5),
3611 ],
3612 10,
3613 );
3614
3615 assert_eq!(groups.len(), 2);
3616 assert_eq!(groups[0].total_docs, 10);
3617 assert_eq!(groups[0].segments.len(), 2);
3618 assert_eq!(groups[1].total_docs, 11);
3619 assert_eq!(groups[1].segments.len(), 1);
3620 }
3621
3622 #[test]
3623 fn force_merge_planner_never_exceeds_segment_format_limit() {
3624 let groups = plan_force_merge_groups(
3625 vec![
3626 ("large-a".into(), 3_000_000_000),
3627 ("large-b".into(), 2_000_000_000),
3628 ],
3629 u64::from(u32::MAX),
3630 );
3631 assert_eq!(groups.len(), 2);
3632 assert!(
3633 groups
3634 .iter()
3635 .all(|group| group.total_docs <= u64::from(u32::MAX))
3636 );
3637 }
3638
3639 #[test]
3640 fn force_merge_hierarchy_has_one_final_bp_pass() {
3641 assert_eq!(force_merge_output_count(1), 0);
3642 assert_eq!(force_merge_output_count(2), 1);
3643 assert_eq!(force_merge_output_count(64), 1);
3644 assert_eq!(force_merge_output_count(65), 2);
3645 assert_eq!(force_merge_output_count(127), 2);
3646 assert_eq!(force_merge_output_count(128), 3);
3647 assert_eq!(force_merge_output_count(1_000), 16);
3648 }
3649
3650 fn expand_force_merge_node(
3651 hierarchy: &ForceMergeHierarchy,
3652 source_count: usize,
3653 node: usize,
3654 sources: &mut Vec<usize>,
3655 ) {
3656 if node < source_count {
3657 sources.push(node);
3658 return;
3659 }
3660
3661 let step_index = node - source_count;
3662 let step = hierarchy
3663 .steps
3664 .get(step_index)
3665 .expect("merge input must refer to an existing source or output");
3666 for &input in &step.inputs {
3667 assert!(
3668 input < node,
3669 "merge step {step_index} refers to a future output node {input}"
3670 );
3671 expand_force_merge_node(hierarchy, source_count, input, sources);
3672 }
3673 }
3674
3675 #[test]
3676 fn force_merge_hierarchy_has_minimal_valid_arity() {
3677 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
3678
3679 for source_count in source_counts {
3680 let hierarchy = plan_force_merge_hierarchy(source_count);
3681 let output_count = hierarchy.steps.len();
3682
3683 assert!(
3684 hierarchy
3685 .steps
3686 .iter()
3687 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
3688 "invalid merge arity for {source_count} sources"
3689 );
3690 assert!(
3691 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
3692 "{output_count} outputs cannot reduce {source_count} sources"
3693 );
3694 assert!(
3695 output_count == 1
3696 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
3697 "{output_count} outputs are not minimal for {source_count} sources"
3698 );
3699 assert_eq!(output_count, force_merge_output_count(source_count));
3700 }
3701 }
3702
3703 #[test]
3704 fn force_merge_hierarchy_preserves_exact_source_order() {
3705 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
3706 let hierarchy = plan_force_merge_hierarchy(source_count);
3707 let mut sources = Vec::with_capacity(source_count);
3708 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
3709 assert_eq!(
3710 sources,
3711 (0..source_count).collect::<Vec<_>>(),
3712 "source order changed for {source_count} sources"
3713 );
3714 }
3715 }
3716
3717 fn force_merge_rewrite_cost(source_count: usize) -> usize {
3718 let hierarchy = plan_force_merge_hierarchy(source_count);
3719 let mut node_weights = vec![1usize; source_count];
3720 let mut rewrite_cost = 0usize;
3721
3722 for (step_index, step) in hierarchy.steps.iter().enumerate() {
3723 let output = source_count + step_index;
3724 let output_weight = step
3725 .inputs
3726 .iter()
3727 .map(|&input| {
3728 assert!(
3729 input < output,
3730 "merge step {step_index} refers to future output {input}"
3731 );
3732 node_weights[input]
3733 })
3734 .sum::<usize>();
3735 rewrite_cost += output_weight;
3736 node_weights.push(output_weight);
3737 }
3738
3739 assert_eq!(node_weights[hierarchy.root], source_count);
3740 rewrite_cost
3741 }
3742
3743 #[test]
3744 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
3745 assert_eq!(force_merge_rewrite_cost(65), 67);
3746 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
3747 }
3748
3749 #[test]
3750 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
3751 assert_eq!(
3752 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
3753 (false, false, 3),
3754 );
3755 assert_eq!(
3756 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
3757 (true, false, 4),
3758 );
3759 assert_eq!(
3760 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
3761 (true, true, 0),
3762 );
3763 }
3764
3765 #[tokio::test]
3766 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
3767 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3768 let field = schema_builder.add_text_field("text", true, true);
3769 let schema = schema_builder.build();
3770 let directory = crate::directories::RamDirectory::new();
3771 let config = crate::index::IndexConfig {
3772 num_indexing_threads: 1,
3773 merge_policy: Box::new(crate::merge::NoMergePolicy),
3774 ..Default::default()
3775 };
3776 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
3777 .await
3778 .unwrap();
3779 for value in ["one", "two", "three"] {
3780 let mut document = crate::dsl::Document::new();
3781 document.add_text(field, value);
3782 writer.add_document(document).unwrap();
3783 writer.commit().await.unwrap();
3784 }
3785
3786 let manager = Arc::clone(writer.segment_manager());
3787 let held_id = manager.get_segment_ids().await.pop().unwrap();
3788 let mut held = Some(
3789 manager
3790 .active_operations
3791 .try_register(vec![held_id])
3792 .unwrap(),
3793 );
3794 let batches = Arc::new(AtomicUsize::new(0));
3795 let batch_count = Arc::clone(&batches);
3796 writer
3797 .force_merge_with_snapshot_refresh(move || {
3798 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
3799 if refresh == 2 {
3802 drop(held.take());
3803 }
3804 std::future::ready(Ok(()))
3805 })
3806 .await
3807 .unwrap();
3808
3809 assert_eq!(manager.get_segment_ids().await.len(), 1);
3810 assert_eq!(
3811 batches.load(Ordering::Relaxed),
3812 4,
3813 "initial/final refreshes plus two replacements are required after the held source releases"
3814 );
3815 }
3816
3817 #[tokio::test]
3818 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
3819 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3820 schema_builder.set_reorder_on_merge(true);
3821 let schema = schema_builder.build();
3822 let mut metadata = IndexMetadata::new(schema.clone());
3823 metadata.add_segment("00000000000000000000000000000001".into(), 1);
3824 metadata.add_segment("00000000000000000000000000000002".into(), 1);
3825
3826 let global_merge_permits = Arc::new(Semaphore::new(1));
3827 let manager = Arc::new(SegmentManager::new(
3828 Arc::new(crate::directories::RamDirectory::new()),
3829 Arc::new(schema),
3830 metadata,
3831 Box::new(crate::merge::NoMergePolicy),
3832 0,
3833 1,
3834 Arc::clone(&global_merge_permits),
3835 None,
3836 1024,
3837 Arc::new(ReorderConcurrencyGate::new(1)),
3838 None,
3839 ));
3840
3841 manager.active_operations.pause_non_indexing();
3846 let force_merge = {
3847 let manager = Arc::clone(&manager);
3848 tokio::spawn(async move { manager.force_merge().await })
3849 };
3850 tokio::time::timeout(std::time::Duration::from_secs(1), async {
3851 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
3852 tokio::task::yield_now().await;
3853 }
3854 })
3855 .await
3856 .expect("force merge never reached the paused group claim");
3857
3858 assert_eq!(
3859 global_merge_permits.available_permits(),
3860 1,
3861 "force merge retained global capacity while vector staging blocked group ownership"
3862 );
3863
3864 force_merge.abort();
3865 let _ = force_merge.await;
3866 manager.active_operations.resume_non_indexing();
3867 }
3868
3869 #[test]
3870 fn output_cleanup_guard_runs_during_panic_unwind() {
3871 let cleaned = Arc::new(AtomicBool::new(false));
3872 let cleaned_in_callback = Arc::clone(&cleaned);
3873 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3874 cleaned_in_callback.store(true, Ordering::SeqCst);
3875 });
3876
3877 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3878 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3879 panic!("simulated reorder panic");
3880 }));
3881
3882 assert!(result.is_err());
3883 assert!(
3884 cleaned.load(Ordering::SeqCst),
3885 "partial output cleanup must run during unwind"
3886 );
3887 }
3888
3889 #[test]
3890 fn output_cleanup_guard_disarms_after_commit() {
3891 let cleaned = Arc::new(AtomicBool::new(false));
3892 let cleaned_in_callback = Arc::clone(&cleaned);
3893 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3894 cleaned_in_callback.store(true, Ordering::SeqCst);
3895 });
3896
3897 {
3898 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3899 guard.disarm();
3900 }
3901
3902 assert!(!cleaned.load(Ordering::SeqCst));
3903 }
3904
3905 #[test]
3906 fn test_active_operation_guard_releases_ownership() {
3907 let active = Arc::new(ActiveSegmentOperations::new());
3908 {
3909 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3910 let snap = active.snapshot();
3911 assert!(snap.contains("a"));
3912 assert!(snap.contains("b"));
3913 }
3914 assert!(active.snapshot().is_empty());
3915 }
3916
3917 #[test]
3918 fn test_non_overlapping_operations_can_run_concurrently() {
3919 let active = Arc::new(ActiveSegmentOperations::new());
3920 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3921 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
3922 let snap = active.snapshot();
3923 assert_eq!(snap.len(), 4);
3924
3925 drop(first);
3926 let snap = active.snapshot();
3927 assert_eq!(snap.len(), 2);
3928 assert!(snap.contains("c"));
3929 assert!(snap.contains("d"));
3930 }
3931
3932 #[test]
3933 fn test_overlapping_operation_is_rejected_until_release() {
3934 let active = Arc::new(ActiveSegmentOperations::new());
3935 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3936 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
3937 drop(first);
3938 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
3939 }
3940
3941 #[test]
3942 fn test_active_operation_snapshot() {
3943 let active = Arc::new(ActiveSegmentOperations::new());
3944 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
3945 let snap = active.snapshot();
3946 assert!(snap.contains("x"));
3947 assert!(snap.contains("y"));
3948 assert!(!snap.contains("z"));
3949 }
3950
3951 #[tokio::test]
3952 async fn operation_barrier_ignores_producers_started_after_snapshot() {
3953 let active = Arc::new(ActiveSegmentOperations::new());
3954 let before_gate = active.try_register(vec!["old".into()]).unwrap();
3955 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
3956 assert_eq!(parked_indexing, 0);
3957 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
3958
3959 let waiter = {
3960 let active = Arc::clone(&active);
3961 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
3962 };
3963 tokio::task::yield_now().await;
3964 assert!(!waiter.is_finished());
3965
3966 drop(before_gate);
3967 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3968 .await
3969 .expect("pre-gate operation barrier was starved by a post-gate producer")
3970 .unwrap();
3971 assert!(active.snapshot().contains("new-flat"));
3972 drop(after_gate);
3973 }
3974
3975 #[tokio::test]
3976 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
3977 let manager = lifecycle_test_manager();
3978 manager
3979 .trained
3980 .store(Some(Arc::new(TrainedVectorStructures {
3981 centroids: rustc_hash::FxHashMap::default(),
3982 binary_quantizers: rustc_hash::FxHashMap::default(),
3983 ..Default::default()
3984 })));
3985
3986 let guard = manager.begin_vector_artifact_update().await.unwrap();
3987 assert!(
3988 manager.trained().is_some(),
3989 "search readers keep the last fully validated generation"
3990 );
3991 assert!(
3992 manager.trained_for_segment_build().is_none(),
3993 "new segment producers must stay flat during an artifact update"
3994 );
3995
3996 let detached_transaction_guard = guard.clone();
3997 drop(guard);
3998 assert!(
3999 manager.trained_for_segment_build().is_none(),
4000 "a detached lifecycle transaction must retain the producer gate after request cancellation"
4001 );
4002 drop(detached_transaction_guard);
4003 assert!(manager.trained_for_segment_build().is_some());
4004 }
4005
4006 #[tokio::test]
4007 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4008 let manager = lifecycle_test_manager();
4009 let guard = manager.begin_vector_artifact_update().await.unwrap();
4010 assert!(
4011 manager
4012 .active_operations
4013 .try_register(vec!["merge".into()])
4014 .is_none(),
4015 "ordinary merge/reorder work must not change staged sources"
4016 );
4017 let indexing = manager
4018 .active_operations
4019 .try_register_indexing(vec!["fresh".into()])
4020 .expect("indexing remains available in flat mode");
4021 drop(indexing);
4022
4023 drop(guard);
4024 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4025 assert!(
4026 manager
4027 .active_operations
4028 .try_register(vec!["merge".into()])
4029 .is_some()
4030 );
4031 }
4032
4033 #[tokio::test]
4034 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4035 let active = Arc::new(ActiveSegmentOperations::new());
4036 let guard = active.try_register(vec!["live".into()]).unwrap();
4037 let cancellation = active.cancellation_flag();
4038 active.stop_accepting();
4039 assert!(cancellation.load(Ordering::Acquire));
4040 assert!(active.try_register(vec!["new".into()]).is_none());
4041
4042 let waiter = {
4043 let active = Arc::clone(&active);
4044 tokio::spawn(async move { active.wait_until_idle().await })
4045 };
4046 tokio::task::yield_now().await;
4047 assert!(!waiter.is_finished());
4048 drop(guard);
4049 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4050 .await
4051 .expect("shutdown waiter missed the final guard notification")
4052 .unwrap();
4053 }
4054
4055 #[tokio::test]
4056 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4057 let manager = lifecycle_test_manager();
4058 let started = Arc::new(Semaphore::new(0));
4059 let release = Arc::new(Semaphore::new(0));
4060 let completed = Arc::new(AtomicBool::new(false));
4061
4062 let request = {
4063 let manager = Arc::clone(&manager);
4064 let started = Arc::clone(&started);
4065 let release = Arc::clone(&release);
4066 let completed = Arc::clone(&completed);
4067 tokio::spawn(async move {
4068 manager
4069 .run_lifecycle_transaction(async move {
4070 started.add_permits(1);
4071 let _permit = release.acquire().await.unwrap();
4072 completed.store(true, Ordering::Release);
4073 Ok(())
4074 })
4075 .await
4076 })
4077 };
4078
4079 let _started = started.acquire().await.unwrap();
4080 request.abort();
4081 assert!(request.await.unwrap_err().is_cancelled());
4082 release.add_permits(1);
4083
4084 manager.begin_shutdown();
4085 tokio::time::timeout(
4086 std::time::Duration::from_secs(1),
4087 manager.wait_for_shutdown(),
4088 )
4089 .await
4090 .expect("shutdown did not drain detached lifecycle transaction");
4091 assert!(completed.load(Ordering::Acquire));
4092 }
4093
4094 #[tokio::test]
4095 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4096 let manager = lifecycle_test_manager();
4097 {
4098 let mut state = manager.state.lock().await;
4099 state.metadata.add_segment_meta(
4100 "eligible".into(),
4101 SegmentMetaInfo {
4102 num_docs: 10,
4103 ancestors: Vec::new(),
4104 generation: 1,
4105 reordered: true,
4106 bp_converged: false,
4107 bp_unconverged_passes: 2,
4108 },
4109 );
4110 state.metadata.add_segment_meta(
4111 "at-limit".into(),
4112 SegmentMetaInfo {
4113 num_docs: 20,
4114 ancestors: Vec::new(),
4115 generation: 1,
4116 reordered: true,
4117 bp_converged: false,
4118 bp_unconverged_passes: 3,
4119 },
4120 );
4121 state.metadata.add_segment_meta(
4122 "carried-debt".into(),
4123 SegmentMetaInfo {
4124 num_docs: 15,
4125 ancestors: Vec::new(),
4126 generation: 2,
4127 reordered: false,
4128 bp_converged: false,
4129 bp_unconverged_passes: 2,
4130 },
4131 );
4132 state.metadata.add_segment_meta(
4133 "carried-debt-at-limit".into(),
4134 SegmentMetaInfo {
4135 num_docs: 25,
4136 ancestors: Vec::new(),
4137 generation: 2,
4138 reordered: false,
4139 bp_converged: false,
4140 bp_unconverged_passes: 3,
4141 },
4142 );
4143 state.metadata.add_segment_meta(
4144 "converged".into(),
4145 SegmentMetaInfo {
4146 num_docs: 30,
4147 ancestors: Vec::new(),
4148 generation: 1,
4149 reordered: true,
4150 bp_converged: true,
4151 bp_unconverged_passes: 0,
4152 },
4153 );
4154 state.metadata.add_segment("fresh".into(), 40);
4155 }
4156
4157 assert_eq!(
4158 manager.unreordered_segments().await,
4159 vec![("fresh".into(), 40)],
4160 "a block-copy output with BP debt is not a fresh first-pass candidate",
4161 );
4162 let mut eligible = manager.unconverged_segments_below(3).await;
4163 eligible.sort_unstable();
4164 assert_eq!(
4165 eligible,
4166 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4167 );
4168 assert!(manager.unconverged_segments_below(0).await.is_empty());
4169 }
4170
4171 #[test]
4172 fn merge_retry_backoff_is_exponential_and_capped() {
4173 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4174 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4175 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4176 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4177 }
4178
4179 #[test]
4180 fn only_deterministic_source_errors_are_quarantined() {
4181 assert!(is_deterministic_source_error(&Error::Corruption(
4182 "bad footer".into()
4183 )));
4184 assert!(is_deterministic_source_error(&Error::Io(
4185 std::io::Error::from(std::io::ErrorKind::NotFound)
4186 )));
4187 assert!(!is_deterministic_source_error(&Error::Io(
4188 std::io::Error::from(std::io::ErrorKind::TimedOut)
4189 )));
4190 assert!(!is_deterministic_source_error(&Error::Io(
4191 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4192 )));
4193 }
4194
4195 #[test]
4196 fn transient_reorder_failure_is_backed_off_until_cleared() {
4197 let manager = lifecycle_test_manager();
4198 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4199 assert!(manager.paused_reorder_segments().contains("source"));
4200 manager.clear_reorder_retry("source");
4201 assert!(!manager.paused_reorder_segments().contains("source"));
4202 }
4203
4204 #[derive(Default)]
4207 struct FailingExistsDirectory(crate::directories::RamDirectory);
4208
4209 #[async_trait::async_trait]
4210 impl crate::directories::Directory for FailingExistsDirectory {
4211 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4212 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4213 }
4214
4215 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4216 self.0.file_size(path).await
4217 }
4218
4219 async fn open_read(
4220 &self,
4221 path: &std::path::Path,
4222 ) -> std::io::Result<crate::directories::FileHandle> {
4223 self.0.open_read(path).await
4224 }
4225
4226 async fn read_range(
4227 &self,
4228 path: &std::path::Path,
4229 range: std::ops::Range<u64>,
4230 ) -> std::io::Result<crate::directories::OwnedBytes> {
4231 self.0.read_range(path, range).await
4232 }
4233
4234 async fn list_files(
4235 &self,
4236 prefix: &std::path::Path,
4237 ) -> std::io::Result<Vec<std::path::PathBuf>> {
4238 self.0.list_files(prefix).await
4239 }
4240
4241 async fn open_lazy(
4242 &self,
4243 path: &std::path::Path,
4244 ) -> std::io::Result<crate::directories::FileHandle> {
4245 self.0.open_lazy(path).await
4246 }
4247 }
4248
4249 #[async_trait::async_trait]
4250 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4251 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4252 self.0.write(path, data).await
4253 }
4254
4255 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4256 self.0.delete(path).await
4257 }
4258
4259 async fn rename(
4260 &self,
4261 from: &std::path::Path,
4262 to: &std::path::Path,
4263 ) -> std::io::Result<()> {
4264 self.0.rename(from, to).await
4265 }
4266
4267 async fn sync(&self) -> std::io::Result<()> {
4268 self.0.sync().await
4269 }
4270
4271 async fn streaming_writer(
4272 &self,
4273 path: &std::path::Path,
4274 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4275 self.0.streaming_writer(path).await
4276 }
4277 }
4278
4279 #[derive(Debug, Clone)]
4280 struct MergeEverythingPolicy;
4281
4282 impl MergePolicy for MergeEverythingPolicy {
4283 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4284 if segments.len() < 2 {
4285 return Vec::new();
4286 }
4287 vec![crate::merge::MergeCandidate {
4288 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4289 }]
4290 }
4291
4292 fn clone_box(&self) -> Box<dyn MergePolicy> {
4293 Box::new(self.clone())
4294 }
4295 }
4296
4297 #[tokio::test]
4298 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4299 let manager = lifecycle_test_manager();
4300 let parked_indexing = manager
4305 .protect_new_segment("00000000000000000000000000000abc".into())
4306 .unwrap();
4307
4308 let error = tokio::time::timeout(
4309 std::time::Duration::from_secs(2),
4310 manager.begin_vector_artifact_update(),
4311 )
4312 .await
4313 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4314 .err()
4315 .expect("an old-generation prepared segment must block artifact replacement")
4316 .to_string();
4317 assert!(error.contains("built but uncommitted"), "{error}");
4318 assert!(
4319 !manager.vector_artifact_update.load(Ordering::Acquire),
4320 "a rejected update must release the producer gate"
4321 );
4322
4323 drop(parked_indexing);
4324
4325 let guard = manager
4326 .begin_vector_artifact_update()
4327 .await
4328 .expect("artifact update should succeed after the pending generation is resolved");
4329 drop(guard);
4330 }
4331
4332 #[tokio::test]
4333 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4334 let manager = lifecycle_test_manager();
4335 let merge_like = manager
4336 .active_operations
4337 .try_register(vec!["merge-source".into()])
4338 .unwrap();
4339
4340 let waiter = {
4341 let manager = Arc::clone(&manager);
4342 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4343 };
4344 for _ in 0..8 {
4345 tokio::task::yield_now().await;
4346 }
4347 assert!(
4348 !waiter.is_finished(),
4349 "artifact update must drain merge/reorder producers that may hold the previous generation"
4350 );
4351
4352 drop(merge_like);
4353 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4354 .await
4355 .expect("artifact update missed the lifecycle guard release")
4356 .unwrap()
4357 .unwrap();
4358 }
4359
4360 #[tokio::test]
4361 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4362 let manager = lifecycle_test_manager();
4363 let release = Arc::new(Semaphore::new(0));
4364 let merge_task = {
4365 let release = Arc::clone(&release);
4366 tokio::spawn(async move {
4367 let _permit = release.acquire().await.unwrap();
4368 })
4369 };
4370 manager.merge_handles.lock().push(merge_task);
4371
4372 let waiter = {
4373 let manager = Arc::clone(&manager);
4374 tokio::spawn(async move { manager.wait_for_all_merges().await })
4375 };
4376 for _ in 0..8 {
4377 tokio::task::yield_now().await;
4378 }
4379 assert!(!waiter.is_finished());
4380 waiter.abort();
4383 let join_error = waiter.await.unwrap_err();
4384 assert!(join_error.is_cancelled());
4385
4386 assert!(
4387 !manager.merge_handles.lock().is_empty(),
4388 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4389 );
4390
4391 release.add_permits(1);
4393 tokio::time::timeout(
4394 std::time::Duration::from_secs(1),
4395 manager.wait_for_all_merges(),
4396 )
4397 .await
4398 .expect("subsequent drain missed the reinserted merge handle");
4399 assert!(manager.merge_handles.lock().is_empty());
4400 }
4401
4402 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4403 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4404 let manager = lifecycle_test_manager();
4405 {
4406 let mut state = manager.state.lock().await;
4407 state
4408 .metadata
4409 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4410 state
4411 .metadata
4412 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4413 }
4414 let reorder_like = manager
4417 .active_operations
4418 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4419 .unwrap();
4420
4421 let force_merge = {
4422 let manager = Arc::clone(&manager);
4423 tokio::spawn(async move { manager.force_merge().await })
4424 };
4425
4426 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4427 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4428 assert!(
4429 retries >= 1,
4430 "force_merge never observed the conflicting owner (retries={retries})"
4431 );
4432 assert!(
4433 retries < 20,
4434 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4435 );
4436
4437 drop(reorder_like);
4438 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4441 .await
4442 .expect("force_merge kept spinning after the conflicting owner released")
4443 .unwrap();
4444 assert!(result.is_err());
4445 }
4446
4447 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4448 async fn force_merge_routes_around_segments_held_by_reorder() {
4449 let manager = lifecycle_test_manager();
4450 {
4451 let mut state = manager.state.lock().await;
4452 state
4453 .metadata
4454 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4455 state
4456 .metadata
4457 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4458 state
4459 .metadata
4460 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4461 }
4462 let _reorder_like = manager
4465 .active_operations
4466 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4467 .unwrap();
4468
4469 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4476 let manager = Arc::clone(&manager);
4477 async move { manager.force_merge().await }
4478 })
4479 .await
4480 .expect("force_merge livelocked on a segment held by an active reorder");
4481 assert!(result.is_err(), "fake segment files must fail the merge");
4482
4483 assert_eq!(
4484 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4485 0,
4486 "batch built from the ownership snapshot must not collide with the held segment"
4487 );
4488 }
4489
4490 #[tokio::test]
4491 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4492 let schema = crate::dsl::SchemaBuilder::default().build();
4493 let mut metadata = IndexMetadata::new(schema.clone());
4494 metadata.add_segment("00000000000000000000000000000001".into(), 10);
4495 metadata.add_segment("00000000000000000000000000000002".into(), 10);
4496 let manager = Arc::new(SegmentManager::new(
4497 Arc::new(FailingExistsDirectory::default()),
4498 Arc::new(schema),
4499 metadata,
4500 Box::new(MergeEverythingPolicy),
4501 0,
4502 1,
4503 Arc::new(Semaphore::new(1)),
4504 None,
4505 1024,
4506 Arc::new(ReorderConcurrencyGate::new(1)),
4507 None,
4508 ));
4509
4510 manager.maybe_merge().await;
4513
4514 tokio::time::timeout(
4515 std::time::Duration::from_secs(5),
4516 manager.wait_for_all_merges(),
4517 )
4518 .await
4519 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4520 assert!(
4521 manager.merge_retry_is_paused(),
4522 "the failed merge should have armed the retry backoff"
4523 );
4524
4525 manager.begin_shutdown();
4527 tokio::time::timeout(
4528 std::time::Duration::from_secs(5),
4529 manager.wait_for_shutdown(),
4530 )
4531 .await
4532 .expect("shutdown did not drain the merge retry wakeup task");
4533 }
4534}