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}
283
284impl ActiveSegmentOperations {
285 fn new() -> Self {
286 Self {
287 inner: parking_lot::Mutex::new(ActiveOperationState {
288 segment_ids: HashSet::new(),
289 operation_tokens: HashSet::new(),
290 indexing_tokens: HashSet::new(),
291 next_operation_token: 0,
292 accepting: true,
293 non_indexing_paused: false,
294 }),
295 idle: Notify::new(),
296 shutdown: Notify::new(),
297 }
298 }
299
300 fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
304 self.try_register_kind(segment_ids, false, false)
305 }
306
307 fn try_register_indexing(
310 self: &Arc<Self>,
311 segment_ids: Vec<String>,
312 ) -> Option<SegmentOperationGuard> {
313 self.try_register_kind(segment_ids, true, false)
314 }
315
316 fn try_register_vector_update(
319 self: &Arc<Self>,
320 segment_ids: Vec<String>,
321 ) -> Option<SegmentOperationGuard> {
322 self.try_register_kind(segment_ids, false, true)
323 }
324
325 fn try_register_kind(
326 self: &Arc<Self>,
327 segment_ids: Vec<String>,
328 indexing: bool,
329 vector_update: bool,
330 ) -> Option<SegmentOperationGuard> {
331 let mut inner = self.inner.lock();
332 if !inner.accepting {
333 log::debug!("[segment_lifecycle] rejected operation during shutdown");
334 return None;
335 }
336 if !indexing && !vector_update && inner.non_indexing_paused {
337 log::debug!("[segment_lifecycle] deferred operation during dense vector retraining");
338 return None;
339 }
340 for id in &segment_ids {
342 if inner.segment_ids.contains(id) {
343 log::debug!(
344 "[segment_lifecycle] rejected: {} overlaps with an active operation ({} active IDs)",
345 id,
346 inner.segment_ids.len()
347 );
348 return None;
349 }
350 }
351 log::debug!(
352 "[segment_lifecycle] registered {} IDs (total active: {})",
353 segment_ids.len(),
354 inner.segment_ids.len() + segment_ids.len()
355 );
356 let operation_token = inner.next_operation_token;
357 let next_operation_token = operation_token.checked_add(1)?;
358 for id in &segment_ids {
359 inner.segment_ids.insert(id.clone());
360 }
361 inner.next_operation_token = next_operation_token;
362 inner.operation_tokens.insert(operation_token);
363 if indexing {
364 inner.indexing_tokens.insert(operation_token);
365 }
366 Some(SegmentOperationGuard {
367 active_operations: Arc::clone(self),
368 segment_ids,
369 operation_token,
370 })
371 }
372
373 fn snapshot(&self) -> HashSet<String> {
375 self.inner.lock().segment_ids.clone()
376 }
377
378 fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
389 let inner = self.inner.lock();
390 let tokens = inner
391 .operation_tokens
392 .difference(&inner.indexing_tokens)
393 .copied()
394 .collect();
395 (tokens, inner.indexing_tokens.len())
396 }
397
398 fn stop_accepting(&self) {
401 let mut inner = self.inner.lock();
402 inner.accepting = false;
403 self.shutdown.notify_waiters();
404 if inner.segment_ids.is_empty() {
405 self.idle.notify_waiters();
406 }
407 }
408
409 fn pause_non_indexing(&self) {
410 self.inner.lock().non_indexing_paused = true;
411 }
412
413 fn resume_non_indexing(&self) {
414 self.inner.lock().non_indexing_paused = false;
415 self.idle.notify_waiters();
416 }
417
418 fn is_accepting(&self) -> bool {
419 self.inner.lock().accepting
420 }
421
422 async fn wait_until_idle(&self) {
426 loop {
427 let notified = self.idle.notified();
428 if self.inner.lock().segment_ids.is_empty() {
429 return;
430 }
431 notified.await;
432 }
433 }
434
435 async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
436 while !operations.is_empty() {
437 let notified = self.idle.notified();
438 if self.inner.lock().operation_tokens.is_disjoint(operations) {
439 return;
440 }
441 notified.await;
442 }
443 }
444
445 async fn wait_for_shutdown(&self) {
448 loop {
449 let notified = self.shutdown.notified();
450 if !self.inner.lock().accepting {
451 return;
452 }
453 notified.await;
454 }
455 }
456}
457
458pub(crate) struct SegmentOperationGuard {
462 active_operations: Arc<ActiveSegmentOperations>,
463 segment_ids: Vec<String>,
464 operation_token: u64,
465}
466
467impl Drop for SegmentOperationGuard {
468 fn drop(&mut self) {
469 let mut inner = self.active_operations.inner.lock();
470 for id in &self.segment_ids {
471 inner.segment_ids.remove(id);
472 }
473 inner.operation_tokens.remove(&self.operation_token);
474 inner.indexing_tokens.remove(&self.operation_token);
475 self.active_operations.idle.notify_waiters();
478 if inner.segment_ids.is_empty() {
479 debug_assert!(inner.operation_tokens.is_empty());
480 }
481 }
482}
483
484struct VectorArtifactUpdateLease {
491 updating: Arc<AtomicBool>,
492 active_operations: Arc<ActiveSegmentOperations>,
493}
494
495impl Drop for VectorArtifactUpdateLease {
496 fn drop(&mut self) {
497 self.updating.store(false, Ordering::Release);
498 self.active_operations.resume_non_indexing();
499 }
500}
501
502#[derive(Clone)]
503pub(crate) struct VectorArtifactUpdateGuard {
504 _lease: Arc<VectorArtifactUpdateLease>,
505}
506
507static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
511
512const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
513const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
514
515#[derive(Default)]
516struct MergeRetryState {
517 retry_after: Option<std::time::Instant>,
518 consecutive_failures: u32,
519}
520
521fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
522 let shift = consecutive_failures.saturating_sub(1).min(16);
523 MERGE_RETRY_BASE_DELAY
524 .checked_mul(1u32 << shift)
525 .unwrap_or(MERGE_RETRY_MAX_DELAY)
526 .min(MERGE_RETRY_MAX_DELAY)
527}
528
529struct DrainedMergeHandles<'a> {
538 shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
539 drained: Vec<JoinHandle<()>>,
540}
541
542impl<'a> DrainedMergeHandles<'a> {
543 fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
544 let drained = std::mem::take(&mut *shared.lock());
545 Self { shared, drained }
546 }
547
548 fn is_empty(&self) -> bool {
549 self.drained.is_empty()
550 }
551
552 async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
556 let handle = self.drained.last_mut()?;
557 let result = handle.await;
558 self.drained.pop();
559 Some(result)
560 }
561}
562
563impl Drop for DrainedMergeHandles<'_> {
564 fn drop(&mut self) {
565 if !self.drained.is_empty() {
566 self.shared.lock().append(&mut self.drained);
567 }
568 }
569}
570
571fn try_spawn_lifecycle<F>(
578 handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
579 runtime: &tokio::runtime::Handle,
580 future: F,
581) -> bool
582where
583 F: std::future::Future<Output = ()> + Send + 'static,
584{
585 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
586 let mut handles = handles.lock();
587 handles.retain(|handle| !handle.is_finished());
588 handles.push(runtime.spawn(future));
589 }))
590 .is_ok()
591}
592
593struct OutputCleanupGuard {
601 segment_id: SegmentId,
602 cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
603}
604
605impl OutputCleanupGuard {
606 fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
607 Self {
608 segment_id,
609 cleanup: Some(cleanup),
610 }
611 }
612
613 fn disarm(&mut self) {
614 self.cleanup = None;
615 }
616}
617
618impl Drop for OutputCleanupGuard {
619 fn drop(&mut self) {
620 if let Some(cleanup) = self.cleanup.take() {
621 cleanup(self.segment_id);
622 }
623 }
624}
625
626struct ManagerState {
628 metadata: IndexMetadata,
629 merge_policy: Box<dyn MergePolicy>,
630}
631
632type ReplacementRefresh = Arc<
633 dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>
634 + Send
635 + Sync,
636>;
637
638async fn refresh_replacement_topology(refresh: Option<ReplacementRefresh>) {
642 let Some(refresh) = refresh else {
643 return;
644 };
645 let mut last_error = None;
646 for attempt in 0..3 {
647 match refresh().await {
648 Ok(()) => return,
649 Err(error) => {
650 last_error = Some(error);
651 if attempt < 2 {
652 tokio::time::sleep(std::time::Duration::from_secs(1 << attempt)).await;
653 }
654 }
655 }
656 }
657 if let Some(error) = last_error {
658 log::warn!(
659 "[segment_lifecycle] replacement topology refresh failed after 3 attempts: {}",
660 error,
661 );
662 }
663}
664
665#[cfg(feature = "native")]
666struct MergeTaskError {
667 error: Error,
668 unavailable_segments: Vec<String>,
669}
670
671#[cfg(feature = "native")]
672impl MergeTaskError {
673 fn source(segment_id: String, error: Error) -> Self {
674 Self {
675 error,
676 unavailable_segments: vec![segment_id],
677 }
678 }
679
680 fn sources(segment_ids: Vec<String>, error: Error) -> Self {
681 Self {
682 error,
683 unavailable_segments: segment_ids,
684 }
685 }
686}
687
688#[cfg(feature = "native")]
689impl From<Error> for MergeTaskError {
690 fn from(error: Error) -> Self {
691 Self {
692 error,
693 unavailable_segments: Vec::new(),
694 }
695 }
696}
697
698#[cfg(feature = "native")]
699fn is_deterministic_source_error(error: &Error) -> bool {
700 matches!(error, Error::Corruption(_) | Error::Serialization(_))
701 || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
702}
703
704#[cfg(feature = "native")]
705fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
706 if is_deterministic_source_error(&error) {
707 MergeTaskError::source(segment_id, error)
708 } else {
709 MergeTaskError::from(error)
713 }
714}
715
716#[cfg(feature = "native")]
717type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
718
719#[derive(Clone, Copy)]
720enum ReplacementLayout {
721 BlockCopy,
725 BpReordered { converged: bool },
727 PreserveSingleSource,
730}
731
732fn replacement_bp_state(
733 parent_has_debt: bool,
734 parent_unconverged_passes: u32,
735 layout: ReplacementLayout,
736) -> (bool, bool, u32) {
737 match layout {
738 ReplacementLayout::BlockCopy => (
739 false,
740 !parent_has_debt,
741 if parent_has_debt {
742 parent_unconverged_passes
743 } else {
744 0
745 },
746 ),
747 ReplacementLayout::BpReordered { converged } => (
748 true,
749 converged,
750 if converged {
751 0
752 } else {
753 parent_unconverged_passes.saturating_add(1)
754 },
755 ),
756 ReplacementLayout::PreserveSingleSource => {
757 unreachable!("preserved layouts retain the complete source metadata")
758 }
759 }
760}
761
762#[derive(Clone, Copy, Debug, Eq, PartialEq)]
763enum VectorSegmentRewriteOutcome {
764 Rewritten,
765 AlreadyCurrent,
766 SourceGone,
767 Conflict,
768 Deferred,
769}
770
771pub(crate) struct StagedVectorSegment {
774 source_id: String,
775 output_id: SegmentId,
776 doc_count: u32,
777 _operation: SegmentOperationGuard,
778 cleanup: OutputCleanupGuard,
779}
780
781pub struct SegmentManager<D: DirectoryWriter + 'static> {
785 state: Arc<AsyncMutex<ManagerState>>,
787
788 active_operations: Arc<ActiveSegmentOperations>,
790
791 quarantined_segments: parking_lot::Mutex<HashSet<String>>,
796
797 merge_retry: parking_lot::Mutex<MergeRetryState>,
800
801 reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
805
806 merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
808
809 global_merge_wakeup_pending: AtomicBool,
813
814 force_merge_active: AtomicUsize,
819
820 #[cfg(test)]
823 force_merge_conflict_retries: std::sync::atomic::AtomicU64,
824
825 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
829
830 trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
834
835 vector_artifact_update: Arc<AtomicBool>,
839
840 tracker: Arc<SegmentTracker>,
842
843 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
845
846 directory: Arc<D>,
848 schema: Arc<crate::dsl::Schema>,
850 term_cache_blocks: usize,
852 merge_permits: Arc<Semaphore>,
856 global_merge_permits: Arc<Semaphore>,
858 reorder_permits: Arc<ReorderConcurrencyGate>,
862 reorder_on_merge: bool,
867 merge_bp_time_budget: Option<std::time::Duration>,
871 bp_memory_budget_bytes: usize,
874 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
877 replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
881}
882
883struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
884
885impl Drop for ForceMergeActivityGuard<'_> {
886 fn drop(&mut self) {
887 self.0.fetch_sub(1, Ordering::AcqRel);
888 }
889}
890
891impl<D: DirectoryWriter + 'static> SegmentManager<D> {
892 #[allow(clippy::too_many_arguments)]
894 pub fn new(
895 directory: Arc<D>,
896 schema: Arc<crate::dsl::Schema>,
897 metadata: IndexMetadata,
898 merge_policy: Box<dyn MergePolicy>,
899 term_cache_blocks: usize,
900 max_concurrent_merges: usize,
901 global_merge_permits: Arc<Semaphore>,
902 merge_bp_time_budget: Option<std::time::Duration>,
903 bp_memory_budget_bytes: usize,
904 reorder_permits: Arc<ReorderConcurrencyGate>,
905 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
906 ) -> Self {
907 let reorder_on_merge = schema.reorder_on_merge();
910 if reorder_on_merge {
911 log::info!("[merge] reorder-on-merge enabled by index schema");
912 }
913
914 let tracker = Arc::new(SegmentTracker::new());
915 for seg_id in metadata.segment_metas.keys() {
916 tracker.register(seg_id);
917 }
918
919 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
920 Arc::new(parking_lot::Mutex::new(Vec::new()));
921 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
922 let dir = Arc::clone(&directory);
923 let tracker = Arc::clone(&tracker);
924 let lifecycle_handles = Arc::clone(&lifecycle_handles);
925 Arc::new(move |segment_ids| {
926 let Ok(handle) = tokio::runtime::Handle::try_current() else {
929 tracker.complete_deletion(&segment_ids);
932 return;
933 };
934 let dir = Arc::clone(&dir);
935 let task_tracker = Arc::clone(&tracker);
936 let cleanup_ids = segment_ids.clone();
937 let future = async move {
938 for &segment_id in &segment_ids {
939 log::info!(
940 "[segment_cleanup] deleting deferred segment {}",
941 segment_id.to_hex()
942 );
943 if let Err(error) =
944 crate::segment::delete_segment(dir.as_ref(), segment_id).await
945 {
946 log::warn!(
947 "[segment_cleanup] deferred delete failed for {}: {}",
948 segment_id.to_hex(),
949 error,
950 );
951 }
952 }
953 task_tracker.complete_deletion(&segment_ids);
954 };
955 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
956 tracker.complete_deletion(&cleanup_ids);
960 log::warn!(
961 "[segment_cleanup] runtime rejected deferred deletion; files will be swept later"
962 );
963 }
964 })
965 };
966
967 Self {
968 state: Arc::new(AsyncMutex::new(ManagerState {
969 metadata,
970 merge_policy,
971 })),
972 active_operations: Arc::new(ActiveSegmentOperations::new()),
973 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
974 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
975 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
976 merge_handles: parking_lot::Mutex::new(Vec::new()),
977 global_merge_wakeup_pending: AtomicBool::new(false),
978 force_merge_active: AtomicUsize::new(0),
979 #[cfg(test)]
980 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
981 lifecycle_handles,
982 trained: Arc::new(ArcSwapOption::new(None)),
983 vector_artifact_update: Arc::new(AtomicBool::new(false)),
984 tracker,
985 delete_fn,
986 directory,
987 schema,
988 term_cache_blocks,
989 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
990 global_merge_permits,
991 reorder_permits,
992 reorder_on_merge,
993 merge_bp_time_budget,
994 bp_memory_budget_bytes,
995 background_reorder_pool,
996 replacement_refresh: parking_lot::RwLock::new(None),
997 }
998 }
999
1000 pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1001 where
1002 F: Fn() -> Fut + Send + Sync + 'static,
1003 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1004 {
1005 *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1006 }
1007
1008 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1014 if let Some(pool) = &self.background_reorder_pool {
1015 return Arc::clone(pool);
1016 }
1017 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1018 let threads = (num_cpus::get() / 2).max(1);
1019 log::info!(
1020 "[merge] process-wide background CPU pool: {} thread(s)",
1021 threads
1022 );
1023 Arc::new(
1024 rayon::ThreadPoolBuilder::new()
1025 .num_threads(threads)
1026 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
1027 .build()
1028 .expect("failed to build background CPU pool"),
1029 )
1030 }))
1031 }
1032
1033 pub fn begin_shutdown(&self) {
1037 self.active_operations.stop_accepting();
1038 }
1039
1040 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1048 where
1049 T: Send + 'static,
1050 F: std::future::Future<Output = Result<T>> + Send + 'static,
1051 {
1052 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1053 let future = async move {
1054 let result = transaction.await;
1055 let _ = result_tx.send(result);
1056 };
1057 let runtime = tokio::runtime::Handle::current();
1058 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1059 return Err(Error::Internal(
1060 "runtime rejected lifecycle metadata transaction".into(),
1061 ));
1062 }
1063 result_rx.await.map_err(|_| {
1064 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1065 })?
1066 }
1067
1068 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1070 let manager = Arc::clone(self);
1071 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1072 let Ok(handle) = tokio::runtime::Handle::try_current() else {
1073 log::warn!(
1074 "[segment_cleanup] runtime unavailable; partial output {} will be swept on startup",
1075 segment_id.to_hex(),
1076 );
1077 return;
1078 };
1079
1080 let cleanup_manager = Arc::clone(&manager);
1081 let future = async move {
1082 cleanup_manager
1083 .delete_output_if_unregistered(segment_id, "task unwind")
1084 .await;
1085 };
1086 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1087 log::warn!(
1088 "[segment_cleanup] runtime rejected output cleanup; {} will be swept on startup",
1089 segment_id.to_hex(),
1090 );
1091 }
1092 });
1093
1094 OutputCleanupGuard::new(output_id, cleanup)
1095 }
1096
1097 pub(crate) fn schedule_unpublished_segment_cleanup(
1102 self: &Arc<Self>,
1103 output_id: SegmentId,
1104 operation: SegmentOperationGuard,
1105 runtime: tokio::runtime::Handle,
1106 ) {
1107 let manager = Arc::clone(self);
1108 let output_hex = output_id.to_hex();
1109 let future = async move {
1110 manager
1111 .delete_output_if_unregistered(output_id, "indexing abort or failure")
1112 .await;
1113 drop(operation);
1114 };
1115 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1116 log::warn!(
1119 "[segment_cleanup] runtime unavailable; indexing output {} will be swept on startup",
1120 output_hex,
1121 );
1122 }
1123 }
1124
1125 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1131 match self
1132 .active_operations
1133 .try_register_indexing(vec![segment_id.clone()])
1134 {
1135 Some(operation) => Ok(operation),
1136 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1137 None => Err(Error::Corruption(format!(
1138 "new segment ID {} is already owned by an active operation",
1139 segment_id
1140 ))),
1141 }
1142 }
1143
1144 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1148 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1149 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1150 })?;
1151 let files = SegmentFiles::new(id.0);
1152
1153 for path in files.mandatory_paths() {
1154 if !self.directory.exists(path).await.map_err(Error::Io)? {
1155 return Err(Error::Corruption(format!(
1156 "segment {} cannot be published: mandatory file {:?} is missing",
1157 segment_id, path
1158 )));
1159 }
1160 }
1161
1162 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1163 Error::Corruption(format!(
1164 "segment {} cannot be published: missing/unreadable {:?}: {}",
1165 segment_id, files.meta, e
1166 ))
1167 })?;
1168 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1169 Error::Corruption(format!(
1170 "segment {} cannot be published: failed reading {:?}: {}",
1171 segment_id, files.meta, e
1172 ))
1173 })?;
1174 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1175 Error::Corruption(format!(
1176 "segment {} cannot be published: invalid {:?}: {}",
1177 segment_id, files.meta, e
1178 ))
1179 })?;
1180
1181 if meta.id != id.0 || meta.num_docs != expected_docs {
1182 return Err(Error::Corruption(format!(
1183 "segment {} cannot be published: metadata identity/docs mismatch \
1184 (id={:032x}, docs={}, expected_docs={})",
1185 segment_id, meta.id, meta.num_docs, expected_docs
1186 )));
1187 }
1188
1189 Ok(())
1190 }
1191
1192 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1193 let inserted = self
1194 .quarantined_segments
1195 .lock()
1196 .insert(segment_id.to_string());
1197 if inserted {
1198 log::error!(
1199 "[merge] quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1200 It remains metadata-live for explicit repair but is excluded from merges until restart",
1201 segment_id,
1202 error,
1203 );
1204 }
1205 }
1206
1207 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1208 let mut retry = self.merge_retry.lock();
1209 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1210 let delay = merge_retry_delay(retry.consecutive_failures);
1211 retry.retry_after = std::time::Instant::now().checked_add(delay);
1212 log::warn!(
1213 "[merge] pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1214 delay.as_secs_f64(),
1215 retry.consecutive_failures,
1216 error,
1217 );
1218 delay
1219 }
1220
1221 fn clear_merge_retry_backoff(&self) {
1222 *self.merge_retry.lock() = MergeRetryState::default();
1223 }
1224
1225 fn merge_retry_is_paused(&self) -> bool {
1226 let mut retry = self.merge_retry.lock();
1227 match retry.retry_after {
1228 Some(deadline) if deadline > std::time::Instant::now() => true,
1229 Some(_) => {
1230 retry.retry_after = None;
1231 false
1232 }
1233 None => false,
1234 }
1235 }
1236
1237 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1238 let mut retries = self.reorder_retries.lock();
1239 let retry = retries.entry(segment_id.to_string()).or_default();
1240 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1241 let delay = merge_retry_delay(retry.consecutive_failures);
1242 retry.retry_after = std::time::Instant::now().checked_add(delay);
1243 log::warn!(
1244 "[reorder] pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1245 segment_id,
1246 delay.as_secs_f64(),
1247 retry.consecutive_failures,
1248 error,
1249 );
1250 }
1251
1252 fn clear_reorder_retry(&self, segment_id: &str) {
1253 self.reorder_retries.lock().remove(segment_id);
1254 }
1255
1256 fn paused_reorder_segments(&self) -> HashSet<String> {
1257 let now = std::time::Instant::now();
1258 let mut retries = self.reorder_retries.lock();
1259 let mut paused = HashSet::new();
1260 for (segment_id, retry) in retries.iter_mut() {
1261 match retry.retry_after {
1262 Some(deadline) if deadline > now => {
1263 paused.insert(segment_id.clone());
1264 }
1265 Some(_) => retry.retry_after = None,
1266 None => {}
1267 }
1268 }
1269 paused
1270 }
1271
1272 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1276 if self
1277 .global_merge_wakeup_pending
1278 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1279 .is_err()
1280 {
1281 return;
1282 }
1283
1284 let manager = Arc::clone(self);
1285 let future = async move {
1286 let capacity = tokio::select! {
1287 biased;
1288 () = manager.active_operations.wait_for_shutdown() => None,
1289 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1290 };
1291
1292 manager
1293 .global_merge_wakeup_pending
1294 .store(false, Ordering::Release);
1295 if let Some(permit) = capacity {
1296 drop(permit);
1300 manager.maybe_merge().await;
1301 }
1302 };
1303 let runtime = tokio::runtime::Handle::current();
1304 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1305 self.global_merge_wakeup_pending
1306 .store(false, Ordering::Release);
1307 log::warn!("[merge] runtime rejected global-capacity wakeup task");
1308 }
1309 }
1310
1311 #[cfg(test)]
1312 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1313 self.quarantined_segments.lock().contains(segment_id)
1314 }
1315
1316 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1321 let output_hex = output_id.to_hex();
1322 {
1323 let st = self.state.lock().await;
1324 if st.metadata.has_segment(&output_hex) {
1325 return;
1326 }
1327 }
1328
1329 log::info!(
1333 "[segment_cleanup] deleting uncommitted output {} after {}",
1334 output_hex,
1335 reason,
1336 );
1337 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1338 {
1339 log::warn!(
1340 "[segment_cleanup] failed deleting uncommitted output {}: {}",
1341 output_hex,
1342 error,
1343 );
1344 }
1345 }
1346
1347 pub async fn get_segment_ids(&self) -> Vec<String> {
1353 self.state.lock().await.metadata.segment_ids()
1354 }
1355
1356 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1358 self.trained.load_full()
1359 }
1360
1361 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1368 if self.vector_artifact_update.load(Ordering::Acquire) {
1369 return None;
1370 }
1371 let trained = self.trained.load_full();
1372 if self.vector_artifact_update.load(Ordering::Acquire) {
1373 None
1374 } else {
1375 trained
1376 }
1377 }
1378
1379 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1398 self.vector_artifact_update
1399 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1400 .map_err(|_| {
1401 Error::Internal("a trained-vector artifact update is already in progress".into())
1402 })?;
1403 self.active_operations.pause_non_indexing();
1404 let guard = VectorArtifactUpdateGuard {
1405 _lease: Arc::new(VectorArtifactUpdateLease {
1406 updating: Arc::clone(&self.vector_artifact_update),
1407 active_operations: Arc::clone(&self.active_operations),
1408 }),
1409 };
1410 let (preexisting, parked_indexing) =
1411 self.active_operations.draining_operation_tokens_snapshot();
1412 if parked_indexing > 0 {
1413 return Err(Error::Internal(format!(
1414 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1415 segment(s) are built but uncommitted; commit or abort the pending \
1416 generation and retry"
1417 )));
1418 }
1419 self.active_operations
1420 .wait_until_operations_finish(&preexisting)
1421 .await;
1422 Ok(guard)
1423 }
1424
1425 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1428 let vector_fields = {
1430 let st = self.state.lock().await;
1431 st.metadata.vector_fields.clone()
1432 };
1433 let trained = IndexMetadata::try_load_trained_from_fields(
1435 &vector_fields,
1436 self.schema.as_ref(),
1437 self.directory.as_ref(),
1438 )
1439 .await?
1440 .map(Arc::new);
1441 self.trained.store(trained);
1445 Ok(())
1446 }
1447
1448 pub(crate) async fn publish_vector_generation(
1455 self: &Arc<Self>,
1456 artifact_update: &VectorArtifactUpdateGuard,
1457 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1458 next_trained: Arc<TrainedVectorStructures>,
1459 mut staged: Vec<StagedVectorSegment>,
1460 ) -> Result<()> {
1461 if !self.vector_artifact_update.load(Ordering::Acquire) {
1462 return Err(Error::Internal(
1463 "vector generation publication lost its exclusive update lease".into(),
1464 ));
1465 }
1466
1467 for replacement in &staged {
1468 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1469 .await?;
1470 }
1471
1472 let mut st = Arc::clone(&self.state).lock_owned().await;
1473 let mut next = st.metadata.clone();
1474 next.vector_fields = vector_fields;
1475 next.refresh_total_vectors();
1476
1477 for replacement in &staged {
1478 let source_info = next
1479 .segment_metas
1480 .remove(&replacement.source_id)
1481 .ok_or_else(|| {
1482 Error::Corruption(format!(
1483 "vector generation source {} disappeared before publication",
1484 replacement.source_id,
1485 ))
1486 })?;
1487 let output_hex = replacement.output_id.to_hex();
1488 if next.segment_metas.contains_key(&output_hex) {
1489 return Err(Error::Corruption(format!(
1490 "vector generation output {output_hex} is already metadata-live"
1491 )));
1492 }
1493 next.add_segment_meta(output_hex, source_info);
1496 }
1497
1498 let directory = Arc::clone(&self.directory);
1499 let trained = Arc::clone(&self.trained);
1500 let tracker = Arc::clone(&self.tracker);
1501 let replacement_refresh = self.replacement_refresh.read().clone();
1502 let artifact_update = artifact_update.clone();
1506 self.run_lifecycle_transaction(async move {
1507 let _artifact_update = artifact_update;
1508 next.save(directory.as_ref()).await?;
1509
1510 for replacement in &staged {
1511 tracker.register(&replacement.output_id.to_hex());
1512 }
1513 st.metadata = next;
1514 trained.store(Some(next_trained));
1515
1516 for replacement in &mut staged {
1519 replacement.cleanup.disarm();
1520 }
1521 let retired = staged
1522 .iter()
1523 .map(|replacement| replacement.source_id.clone())
1524 .collect::<Vec<_>>();
1525 let ready_to_delete = tracker.mark_for_deletion(&retired);
1526 drop(st);
1527 for &segment_id in &ready_to_delete {
1528 if let Err(error) =
1529 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1530 {
1531 log::warn!(
1532 "[segment_cleanup] immediate dense-vector generation delete failed for {}: {}",
1533 segment_id.to_hex(),
1534 error,
1535 );
1536 }
1537 }
1538 tracker.complete_deletion(&ready_to_delete);
1539 refresh_replacement_topology(replacement_refresh).await;
1540 Ok(())
1541 })
1542 .await
1543 }
1544
1545 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1547 where
1548 F: FnOnce(&IndexMetadata) -> R,
1549 {
1550 let st = self.state.lock().await;
1551 f(&st.metadata)
1552 }
1553
1554 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1556 where
1557 F: FnOnce(&mut IndexMetadata),
1558 {
1559 let mut st = Arc::clone(&self.state).lock_owned().await;
1560 let mut next = st.metadata.clone();
1561 f(&mut next);
1562 let directory = Arc::clone(&self.directory);
1563 self.run_lifecycle_transaction(async move {
1564 next.save(directory.as_ref()).await?;
1565 st.metadata = next;
1566 Ok(())
1567 })
1568 .await
1569 }
1570
1571 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1574 let (acquired, trained) = {
1575 let st = self.state.lock().await;
1576 let segment_ids = st.metadata.segment_ids();
1577 (self.tracker.acquire(&segment_ids), self.trained.load_full())
1578 };
1579
1580 SegmentSnapshot::with_generation(
1581 Arc::clone(&self.tracker),
1582 acquired,
1583 trained,
1584 Arc::clone(&self.delete_fn),
1585 )
1586 }
1587
1588 pub fn tracker(&self) -> Arc<SegmentTracker> {
1590 Arc::clone(&self.tracker)
1591 }
1592
1593 pub fn directory(&self) -> Arc<D> {
1595 Arc::clone(&self.directory)
1596 }
1597}
1598
1599#[cfg(feature = "native")]
1604impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1605 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1607 for (segment_id, num_docs) in new_segments {
1610 self.validate_completed_segment(segment_id, *num_docs)
1611 .await?;
1612 }
1613
1614 let mut st = Arc::clone(&self.state).lock_owned().await;
1615 let mut next = st.metadata.clone();
1616 let mut added = Vec::new();
1617 for (segment_id, num_docs) in new_segments {
1618 if !next.has_segment(segment_id) {
1619 next.add_segment(segment_id.clone(), *num_docs);
1620 added.push(segment_id.clone());
1621 }
1622 }
1623
1624 let directory = Arc::clone(&self.directory);
1630 let tracker = Arc::clone(&self.tracker);
1631 self.run_lifecycle_transaction(async move {
1632 next.save(directory.as_ref()).await?;
1633 for segment_id in &added {
1634 tracker.register(segment_id);
1635 }
1636 st.metadata = next;
1637 Ok(())
1638 })
1639 .await
1640 }
1641
1642 pub async fn maybe_merge(self: &Arc<Self>) {
1653 if !self.active_operations.is_accepting() {
1654 log::debug!("[maybe_merge] manager is shutting down, skipping");
1655 return;
1656 }
1657 if self.merge_retry_is_paused() {
1658 log::debug!("[maybe_merge] retry backoff active, skipping");
1659 return;
1660 }
1661
1662 {
1665 let mut handles = self.merge_handles.lock();
1666 handles.retain(|h| !h.is_finished());
1667 }
1668 let local_slots = self.merge_permits.available_permits();
1669 let global_slots = self.global_merge_permits.available_permits();
1670 let slots_available = local_slots.min(global_slots);
1671
1672 {
1676 let st = self.state.lock().await;
1677 let quarantined = self.quarantined_segments.lock().clone();
1678 let active_ids = self.active_operations.snapshot();
1679
1680 let segments: Vec<SegmentInfo> = st
1683 .metadata
1684 .segment_metas
1685 .iter()
1686 .filter(|(id, _)| {
1687 !self.tracker.is_pending_deletion(id)
1688 && !active_ids.contains(*id)
1689 && !quarantined.contains(*id)
1690 })
1691 .map(|(id, info)| SegmentInfo {
1692 id: id.clone(),
1693 num_docs: info.num_docs,
1694 })
1695 .collect();
1696
1697 log::debug!("[maybe_merge] {} eligible segments", segments.len());
1698
1699 let candidates = st.merge_policy.find_merges(&segments);
1700
1701 if candidates.is_empty() {
1702 return;
1703 }
1704
1705 if slots_available == 0 {
1709 if local_slots > 0 && global_slots == 0 {
1710 self.schedule_global_merge_wakeup();
1711 }
1712 log::debug!("[maybe_merge] at max concurrent merges, skipping");
1713 return;
1714 }
1715
1716 log::debug!(
1717 "[maybe_merge] {} merge candidates, {} slots available",
1718 candidates.len(),
1719 slots_available
1720 );
1721
1722 let mut handles = Vec::new();
1723 for c in candidates {
1724 if handles.len() >= slots_available {
1725 break;
1726 }
1727 if let Some(h) = self.spawn_merge(c.segment_ids) {
1728 handles.push(h);
1729 }
1730 }
1731 if !handles.is_empty() {
1732 self.merge_handles.lock().extend(handles);
1737 }
1738 }
1739 }
1740
1741 fn spawn_merge(self: &Arc<Self>, segment_ids_to_merge: Vec<String>) -> Option<JoinHandle<()>> {
1750 if self.force_merge_active.load(Ordering::Acquire) > 0 {
1751 log::debug!("[spawn_merge] skipped: explicit force merge has priority");
1752 return None;
1753 }
1754 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1755 Ok(permit) => permit,
1756 Err(_) => {
1757 log::debug!("[spawn_merge] skipped: global merge capacity is full");
1758 self.schedule_global_merge_wakeup();
1759 return None;
1760 }
1761 };
1762 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1763 Ok(permit) => permit,
1764 Err(_) => {
1765 log::debug!("[spawn_merge] skipped: no merge permit available");
1766 return None;
1767 }
1768 };
1769 let output_id = SegmentId::new();
1770 let output_hex = output_id.to_hex();
1771
1772 let mut all_ids = segment_ids_to_merge.clone();
1773 all_ids.push(output_hex);
1774
1775 let guard = match self.active_operations.try_register(all_ids) {
1776 Some(g) => g,
1777 None => {
1778 log::debug!("[spawn_merge] skipped: segments overlap with an active operation");
1779 return None;
1780 }
1781 };
1782
1783 let sm = Arc::clone(self);
1784 let ids = segment_ids_to_merge;
1785
1786 Some(tokio::spawn(async move {
1787 let mut reevaluate = false;
1788 let mut retry_delay = None;
1789
1790 let result = sm
1791 .merge_and_replace_registered(
1792 &ids,
1793 output_id,
1794 sm.reorder_on_merge,
1795 ReorderPriority::AutomaticMerge,
1796 )
1797 .await;
1798
1799 match result {
1800 Ok(_) => {
1801 sm.clear_merge_retry_backoff();
1802 reevaluate = true;
1803 }
1804 Err(MergeTaskError {
1805 error,
1806 unavailable_segments,
1807 }) => {
1808 log::error!(
1809 "[merge] background merge failed for segments {:?}: {}",
1810 ids,
1811 error
1812 );
1813 if !unavailable_segments.is_empty() {
1814 reevaluate = true;
1818 } else {
1819 retry_delay = Some(sm.pause_merge_retries(&error));
1820 }
1821 }
1822 }
1823 drop(guard);
1826 drop(merge_permit);
1828 drop(global_merge_permit);
1829
1830 if reevaluate {
1831 sm.maybe_merge().await;
1832 } else if let Some(retry_delay) = retry_delay {
1833 sm.schedule_merge_retry_wakeup(retry_delay);
1840 }
1841 }))
1842 }
1843
1844 async fn merge_and_replace_registered(
1851 self: &Arc<Self>,
1852 ids: &[String],
1853 output_id: SegmentId,
1854 reorder_bmp: bool,
1855 priority: ReorderPriority,
1856 ) -> MergeTaskResult<(String, u32, bool)> {
1857 let mut output_cleanup = self.output_cleanup_guard(output_id);
1858 let trained = self.trained_for_segment_build();
1859 let granularity = if reorder_bmp {
1860 self.merge_granularity(ids).await
1861 } else {
1862 crate::segment::reorder::BpGranularity::Auto
1863 };
1864 let result = Self::do_merge(
1865 self.directory.as_ref(),
1866 &self.schema,
1867 ids,
1868 output_id,
1869 self.term_cache_blocks,
1870 trained.as_deref(),
1871 reorder_bmp,
1872 granularity,
1873 self.merge_bp_time_budget,
1874 self.bp_memory_budget_bytes,
1875 Arc::clone(&self.reorder_permits),
1876 priority,
1877 Some(self.background_cpu_pool()),
1878 )
1879 .await;
1880
1881 let (new_id, doc_count, bp_converged) = match result {
1882 Ok(value) => value,
1883 Err(error) => {
1884 for segment_id in &error.unavailable_segments {
1885 self.quarantine_segment(segment_id, &error.error);
1886 }
1887 self.delete_output_if_unregistered(output_id, "merge failure")
1888 .await;
1889 output_cleanup.disarm();
1890 return Err(error);
1891 }
1892 };
1893
1894 let layout = if reorder_bmp {
1895 ReplacementLayout::BpReordered {
1896 converged: bp_converged,
1897 }
1898 } else {
1899 ReplacementLayout::BlockCopy
1900 };
1901 if let Err(error) = self
1902 .replace_segments(ids, new_id.clone(), doc_count, layout)
1903 .await
1904 {
1905 self.delete_output_if_unregistered(output_id, "replacement failure")
1906 .await;
1907 output_cleanup.disarm();
1908 return Err(MergeTaskError::from(error));
1909 }
1910 output_cleanup.disarm();
1911 Ok((new_id, doc_count, bp_converged))
1912 }
1913
1914 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
1918 let manager = Arc::clone(self);
1919 let future = async move {
1920 tokio::select! {
1921 () = tokio::time::sleep(retry_delay) => {
1922 manager.maybe_merge().await;
1923 }
1924 () = manager.active_operations.wait_for_shutdown() => {}
1925 }
1926 };
1927 let runtime = tokio::runtime::Handle::current();
1928 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1929 log::warn!(
1930 "[merge] runtime rejected merge-retry wakeup task; eligible segments may stay \
1931 unmerged until the next commit re-runs merge policy evaluation"
1932 );
1933 }
1934 }
1935
1936 async fn replace_segments(
1940 self: &Arc<Self>,
1941 old_ids: &[String],
1942 new_id: String,
1943 doc_count: u32,
1944 layout: ReplacementLayout,
1945 ) -> Result<()> {
1946 self.validate_completed_segment(&new_id, doc_count).await?;
1949 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
1950 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
1951 })?;
1952 let output_reader = SegmentReader::open(
1953 self.directory.as_ref(),
1954 output_id,
1955 Arc::clone(&self.schema),
1956 self.term_cache_blocks,
1957 )
1958 .await
1959 .map_err(|error| match error {
1960 Error::Io(_) | Error::IndexClosed => error,
1964 error => Error::Corruption(format!(
1965 "replacement segment {new_id} failed full reader validation: {error}"
1966 )),
1967 })?;
1968 if output_reader.num_docs() != doc_count {
1969 return Err(Error::Corruption(format!(
1970 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
1971 output_reader.num_docs(),
1972 )));
1973 }
1974 drop(output_reader);
1975
1976 let mut st = Arc::clone(&self.state).lock_owned().await;
1977 let missing: Vec<&String> = old_ids
1981 .iter()
1982 .filter(|id| !st.metadata.has_segment(id))
1983 .collect();
1984 if !missing.is_empty() {
1985 return Err(Error::Corruption(format!(
1986 "replace_segments: source segment(s) {:?} not in metadata — \
1987 refusing to add output {} (would duplicate documents)",
1988 missing, new_id
1989 )));
1990 }
1991
1992 let replacement_info = match layout {
1993 ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
1994 let generation = old_ids
1995 .iter()
1996 .filter_map(|id| st.metadata.segment_metas.get(id))
1997 .map(|info| info.generation)
1998 .max()
1999 .unwrap_or(0)
2000 .checked_add(1)
2001 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2002 let parent_unconverged_passes = old_ids
2003 .iter()
2004 .filter_map(|id| st.metadata.segment_metas.get(id))
2005 .map(|info| info.bp_unconverged_passes)
2006 .max()
2007 .unwrap_or(0);
2008 let parent_has_debt = old_ids
2009 .iter()
2010 .filter_map(|id| st.metadata.segment_metas.get(id))
2011 .any(|info| !info.bp_converged);
2012 let (reordered, bp_converged, bp_unconverged_passes) =
2013 replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2014 SegmentMetaInfo {
2015 num_docs: doc_count,
2016 ancestors: old_ids.to_vec(),
2017 generation,
2018 reordered,
2019 bp_converged,
2020 bp_unconverged_passes,
2021 }
2022 }
2023 ReplacementLayout::PreserveSingleSource => {
2024 let [source_id] = old_ids else {
2025 return Err(Error::Internal(
2026 "layout-preserving replacement requires exactly one source".into(),
2027 ));
2028 };
2029 let mut source = st
2030 .metadata
2031 .segment_metas
2032 .get(source_id)
2033 .cloned()
2034 .ok_or_else(|| {
2035 Error::Corruption(format!(
2036 "layout-preserving replacement source {source_id} disappeared"
2037 ))
2038 })?;
2039 source.num_docs = doc_count;
2040 source
2041 }
2042 };
2043 let retired_ids = old_ids.to_vec();
2044 let mut next = st.metadata.clone();
2045 for id in old_ids {
2046 next.remove_segment(id);
2047 }
2048 next.add_segment_meta(new_id.clone(), replacement_info);
2049
2050 let directory = Arc::clone(&self.directory);
2051 let tracker = Arc::clone(&self.tracker);
2052 let replacement_refresh = self.replacement_refresh.read().clone();
2053 self.run_lifecycle_transaction(async move {
2054 next.save(directory.as_ref()).await?;
2057 tracker.register(&new_id);
2058 st.metadata = next;
2059
2060 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2064 drop(st);
2065 for &segment_id in &ready_to_delete {
2066 if let Err(error) =
2067 crate::segment::delete_segment(directory.as_ref(), segment_id).await
2068 {
2069 log::warn!(
2070 "[segment_cleanup] immediate delete failed for {}: {}",
2071 segment_id.to_hex(),
2072 error,
2073 );
2074 }
2075 }
2076 tracker.complete_deletion(&ready_to_delete);
2077 refresh_replacement_topology(replacement_refresh).await;
2078 Ok(())
2079 })
2080 .await
2081 }
2082
2083 #[allow(clippy::too_many_arguments)]
2088 async fn do_merge(
2089 directory: &D,
2090 schema: &Arc<crate::dsl::Schema>,
2091 segment_ids_to_merge: &[String],
2092 output_segment_id: SegmentId,
2093 term_cache_blocks: usize,
2094 trained: Option<&TrainedVectorStructures>,
2095 reorder_bmp: bool,
2096 granularity: crate::segment::reorder::BpGranularity,
2097 merge_bp_time_budget: Option<std::time::Duration>,
2098 bp_memory_budget_bytes: usize,
2099 reorder_permits: Arc<ReorderConcurrencyGate>,
2100 reorder_priority: ReorderPriority,
2101 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2102 ) -> MergeTaskResult<(String, u32, bool)> {
2103 let output_hex = output_segment_id.to_hex();
2104 let load_start = std::time::Instant::now();
2105
2106 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2107 for id_str in segment_ids_to_merge {
2108 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2109 MergeTaskError::source(
2110 id_str.clone(),
2111 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2112 )
2113 })?;
2114 segment_ids.push(id);
2115 }
2116
2117 let mut unavailable_sources = Vec::new();
2122 let mut missing_files = Vec::new();
2123 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2124 let files = SegmentFiles::new(id.0);
2125 let mut source_unavailable = false;
2126 for path in files.mandatory_paths() {
2127 let exists = directory
2128 .exists(path)
2129 .await
2130 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2131 if !exists {
2132 source_unavailable = true;
2133 missing_files.push(format!("{}:{:?}", id_str, path));
2134 }
2135 }
2136 if source_unavailable {
2137 unavailable_sources.push(id_str.clone());
2138 }
2139 }
2140 if !unavailable_sources.is_empty() {
2141 return Err(MergeTaskError::sources(
2142 unavailable_sources,
2143 Error::Corruption(format!(
2144 "merge sources are missing mandatory files: {}",
2145 missing_files.join(", ")
2146 )),
2147 ));
2148 }
2149
2150 let schema_arc = Arc::clone(schema);
2151 let futures: Vec<_> = segment_ids
2152 .iter()
2153 .map(|&sid| {
2154 let sch = Arc::clone(&schema_arc);
2155 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
2156 })
2157 .collect();
2158
2159 let results = futures::future::join_all(futures).await;
2160 let mut readers = Vec::with_capacity(results.len());
2161 let mut total_docs = 0u64;
2162 for (i, result) in results.into_iter().enumerate() {
2163 match result {
2164 Ok(r) => {
2165 total_docs += r.meta().num_docs as u64;
2166 readers.push(r);
2167 }
2168 Err(e) => {
2169 log::error!(
2170 "[merge] Failed to open segment {}: {:?}",
2171 segment_ids_to_merge[i],
2172 e
2173 );
2174 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2175 }
2176 }
2177 }
2178 if total_docs > u32::MAX as u64 {
2179 return Err(Error::Internal(format!(
2180 "Merged segment doc count ({}) exceeds u32::MAX",
2181 total_docs
2182 ))
2183 .into());
2184 }
2185
2186 for (i, reader) in readers.iter().enumerate() {
2190 let meta_docs = reader.meta().num_docs;
2191 let store_docs = reader.store().num_docs();
2192 if store_docs != meta_docs {
2193 return Err(MergeTaskError::source(
2194 segment_ids_to_merge[i].clone(),
2195 Error::Corruption(format!(
2196 "pre-merge validation: segment {} store has {} docs but meta says {}",
2197 segment_ids_to_merge[i], store_docs, meta_docs
2198 )),
2199 ));
2200 }
2201 }
2202
2203 log::info!(
2204 "[merge] loaded {} segment readers in {:.1}s",
2205 readers.len(),
2206 load_start.elapsed().as_secs_f64()
2207 );
2208
2209 let merger = SegmentMerger::new(Arc::clone(schema))
2210 .with_bmp_reorder(reorder_bmp)
2211 .with_granularity(granularity)
2212 .with_bp_budget(crate::segment::BpBudget {
2213 min_partition_docs: None,
2214 time_budget: merge_bp_time_budget,
2215 })
2216 .with_bp_memory_budget(bp_memory_budget_bytes)
2217 .with_reorder_permits(reorder_permits)
2218 .with_reorder_priority(reorder_priority)
2219 .with_background_pool(bg_cpu_pool);
2220
2221 log::info!(
2222 "[merge] {} segments -> {} (trained={})",
2223 segment_ids_to_merge.len(),
2224 output_hex,
2225 trained.map_or(0, |t| t.centroids.len()),
2226 );
2227
2228 let (_merged_meta, merge_stats) = merger
2229 .merge(directory, &readers, output_segment_id, trained)
2230 .await
2231 .map_err(|error| {
2232 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2233 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2239 } else {
2240 MergeTaskError::from(error)
2241 }
2242 })?;
2243 let bp_converged = merge_stats.bp_converged;
2244 if !bp_converged {
2245 log::info!(
2246 "[merge] merge-time BP hit its wall-clock budget — output marked unconverged; \
2247 the background optimizer deepens it later",
2248 );
2249 }
2250
2251 log::info!(
2252 "[merge] total wall-clock: {:.1}s ({} segments, {} docs)",
2253 load_start.elapsed().as_secs_f64(),
2254 readers.len(),
2255 total_docs,
2256 );
2257
2258 Ok((output_hex, total_docs as u32, bp_converged))
2259 }
2260
2261 pub async fn abort_merges(&self) {
2271 loop {
2272 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2273 if handles.is_empty() {
2274 return;
2275 }
2276 while let Some(result) = handles.join_next().await {
2277 if let Err(error) = result
2278 && error.is_panic()
2279 {
2280 log::error!("[merge] background task panicked while draining: {}", error);
2281 }
2282 }
2283 }
2284 }
2285
2286 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2291 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2292 while handles.join_next().await.is_some() {}
2293 }
2294
2295 pub async fn wait_for_all_merges(self: &Arc<Self>) {
2304 loop {
2305 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2306 if handles.is_empty() {
2307 break;
2308 }
2309 while handles.join_next().await.is_some() {}
2310 }
2311 }
2312
2313 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2318 self.wait_for_all_merges().await;
2319 self.active_operations.wait_until_idle().await;
2320 loop {
2321 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2322 if handles.is_empty() {
2323 break;
2324 }
2325 for handle in handles {
2326 if let Err(error) = handle.await
2327 && error.is_panic()
2328 {
2329 log::error!("[segment_cleanup] task panicked while draining: {}", error);
2330 }
2331 }
2332 }
2333 }
2334
2335 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2344 self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2345 .await
2346 }
2347
2348 pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2358 self: &Arc<Self>,
2359 mut refresh_snapshots: F,
2360 ) -> Result<()>
2361 where
2362 F: FnMut() -> Fut,
2363 Fut: std::future::Future<Output = Result<()>>,
2364 {
2365 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2370 std::time::Duration::from_millis(100);
2371 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2375
2376 let (_force_merge_activity, max_segment_docs) = {
2377 let st = self.state.lock().await;
2378 self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2382 (
2383 ForceMergeActivityGuard(&self.force_merge_active),
2384 st.merge_policy.max_segment_docs(),
2385 )
2386 };
2387
2388 let drain_start = std::time::Instant::now();
2391 self.wait_for_all_merges().await;
2392 if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2393 log::info!(
2394 "[force_merge] drained background merges in {:.1}s",
2395 drain_start.elapsed().as_secs_f64(),
2396 );
2397 }
2398
2399 let refresh_start = std::time::Instant::now();
2404 refresh_snapshots().await?;
2405 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2406 log::info!(
2407 "[force_merge] initial snapshot refresh took {:.1}s",
2408 refresh_start.elapsed().as_secs_f64(),
2409 );
2410 }
2411
2412 let mut completed_outputs = HashSet::new();
2416 let mut logged_held_wait = false;
2418
2419 loop {
2420 if !self.active_operations.is_accepting() {
2421 return Err(Error::IndexClosed);
2422 }
2423
2424 let segments: Vec<(String, u32)> = {
2425 let st = self.state.lock().await;
2426 st.metadata
2427 .segment_metas
2428 .iter()
2429 .filter(|(id, _)| !completed_outputs.contains(*id))
2430 .map(|(id, info)| (id.clone(), info.num_docs))
2431 .collect()
2432 };
2433
2434 let active_ids = self.active_operations.snapshot();
2439 let held = segments
2440 .iter()
2441 .filter(|(id, _)| active_ids.contains(id))
2442 .count();
2443 let free_segments: Vec<_> = segments
2444 .into_iter()
2445 .filter(|(id, _)| !active_ids.contains(id))
2446 .collect();
2447 let max_docs = max_segment_docs
2451 .map(u64::from)
2452 .unwrap_or(u64::from(u32::MAX));
2453 let next_group = plan_force_merge_groups(free_segments, max_docs)
2454 .into_iter()
2455 .find(|group| group.segments.len() >= 2);
2456
2457 let Some(group) = next_group else {
2458 if held == 0 {
2459 if !completed_outputs.is_empty() {
2460 completed_outputs.clear();
2466 continue;
2467 }
2468 refresh_snapshots().await?;
2475 return Ok(());
2476 }
2477 if !logged_held_wait {
2478 log::info!(
2479 "[force_merge] waiting: {} segment(s) held by active \
2480 merge/reorder operations, no free group can merge",
2481 held
2482 );
2483 logged_held_wait = true;
2484 } else {
2485 log::debug!("[force_merge] still waiting on {} held segment(s)", held);
2486 }
2487 #[cfg(test)]
2488 self.force_merge_conflict_retries
2489 .fetch_add(1, Ordering::Relaxed);
2490 tokio::select! {
2491 biased;
2492 () = self.active_operations.wait_for_shutdown() => {
2493 return Err(Error::IndexClosed);
2494 }
2495 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2496 }
2497 continue;
2498 };
2499 logged_held_wait = false;
2500
2501 let hierarchy = plan_force_merge_hierarchy(group.segments.len());
2505 let output_ids: Vec<_> = (0..hierarchy.steps.len())
2506 .map(|_| SegmentId::new())
2507 .collect();
2508 let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
2509 let mut all_ids = source_ids.clone();
2510 all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
2511 let group_guard = {
2512 let st = self.state.lock().await;
2513 source_ids
2514 .iter()
2515 .all(|id| st.metadata.has_segment(id))
2516 .then(|| self.active_operations.try_register(all_ids))
2517 .flatten()
2518 };
2519 let _group_guard = match group_guard {
2520 Some(guard) => guard,
2521 None if !self.active_operations.is_accepting() => {
2522 return Err(Error::IndexClosed);
2523 }
2524 None => {
2525 #[cfg(test)]
2526 self.force_merge_conflict_retries
2527 .fetch_add(1, Ordering::Relaxed);
2528 log::debug!("[force_merge] group lost a registration race, replanning");
2529 let had_tracked_merges = !self.merge_handles.lock().is_empty();
2530 self.wait_for_merging_thread().await;
2531 if !had_tracked_merges {
2532 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2533 }
2534 continue;
2535 }
2536 };
2537
2538 log::info!(
2539 "[force_merge] planned final group: {} segments, {} docs, {} merge pass(es)",
2540 group.segments.len(),
2541 group.total_docs,
2542 output_ids.len(),
2543 );
2544
2545 let group_global_merge_permit = if self.reorder_on_merge {
2555 let capacity_start = std::time::Instant::now();
2556 let permit = tokio::select! {
2557 biased;
2558 () = self.active_operations.wait_for_shutdown() => {
2559 return Err(Error::IndexClosed);
2560 }
2561 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2562 permit.map_err(|_| {
2563 Error::Internal(
2564 "global background merge scheduler is closed".into(),
2565 )
2566 })?
2567 }
2568 };
2569 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2570 log::info!(
2571 "[force_merge] waited {:.1}s for foreground global merge capacity",
2572 capacity_start.elapsed().as_secs_f64(),
2573 );
2574 }
2575 Some(permit)
2576 } else {
2577 None
2578 };
2579 let _foreground_reorder = if self.reorder_on_merge {
2580 log::info!(
2581 "[force_merge] prioritizing BP capacity ({} total pass slot(s))",
2582 self.reorder_permits.limit(),
2583 );
2584 let admission_start = std::time::Instant::now();
2585 let guard = Arc::clone(&self.reorder_permits)
2586 .begin_foreground()
2587 .await
2588 .map_err(|_| {
2589 Error::Internal("background reorder scheduler is closed".into())
2590 })?;
2591 if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
2592 log::info!(
2593 "[force_merge] acquired foreground BP capacity in {:.1}s",
2594 admission_start.elapsed().as_secs_f64(),
2595 );
2596 }
2597 Some(guard)
2598 } else {
2599 None
2600 };
2601
2602 let source_count = group.segments.len();
2603 let mut nodes: Vec<Option<(String, u32)>> =
2604 group.segments.into_iter().map(Some).collect();
2605 nodes.resize_with(source_count + hierarchy.steps.len(), || None);
2606 for (step_index, step) in hierarchy.steps.iter().enumerate() {
2607 let final_pass = step_index + 1 == hierarchy.steps.len();
2608 let mut batch_entries = Vec::with_capacity(step.inputs.len());
2609 for &node in &step.inputs {
2610 let entry = nodes
2611 .get_mut(node)
2612 .and_then(Option::take)
2613 .expect("force-merge hierarchy must reference an available node");
2614 batch_entries.push(entry);
2615 }
2616 let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
2617 let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
2618 let output_id = output_ids[step_index];
2619
2620 let capacity_start = std::time::Instant::now();
2621 let step_global_merge_permit = if group_global_merge_permit.is_none() {
2622 Some(tokio::select! {
2623 biased;
2624 () = self.active_operations.wait_for_shutdown() => {
2625 return Err(Error::IndexClosed);
2626 }
2627 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2628 permit.map_err(|_| {
2629 Error::Internal(
2630 "global background merge scheduler is closed".into(),
2631 )
2632 })?
2633 }
2634 })
2635 } else {
2636 None
2637 };
2638 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2639 log::info!(
2640 "[force_merge] waited {:.1}s for global merge capacity",
2641 capacity_start.elapsed().as_secs_f64(),
2642 );
2643 }
2644
2645 let reorder_bmp = final_pass && self.reorder_on_merge;
2649 log::info!(
2650 "[force_merge] {} pass: {} segments ({} docs, bp={})",
2651 if final_pass {
2652 "final"
2653 } else {
2654 "fan-in reduction"
2655 },
2656 batch.len(),
2657 batch_docs,
2658 reorder_bmp,
2659 );
2660 let (new_segment_id, total_docs, _) = self
2661 .merge_and_replace_registered(
2662 &batch,
2663 output_id,
2664 reorder_bmp,
2665 ReorderPriority::Foreground,
2666 )
2667 .await
2668 .map_err(|error| error.error)?;
2669 drop(step_global_merge_permit);
2670
2671 let refresh_start = std::time::Instant::now();
2674 refresh_snapshots().await?;
2675 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2676 log::info!(
2677 "[force_merge] post-replacement snapshot refresh took {:.1}s",
2678 refresh_start.elapsed().as_secs_f64(),
2679 );
2680 }
2681
2682 let output_node = source_count + step_index;
2683 debug_assert!(nodes[output_node].is_none());
2684 nodes[output_node] = Some((new_segment_id, total_docs));
2685 }
2686 let (root_id, _) = nodes[hierarchy.root]
2687 .take()
2688 .expect("force-merge hierarchy must produce its root");
2689 debug_assert!(nodes.into_iter().all(|node| node.is_none()));
2690 completed_outputs.insert(root_id);
2691 }
2692 }
2693
2694 fn segment_needs_vector_rewrite(
2695 &self,
2696 reader: &SegmentReader,
2697 field_ids: &[u32],
2698 rewrite_existing: bool,
2699 ) -> Result<bool> {
2700 for &field_id in field_ids {
2701 let flat = reader.flat_vectors().get(&field_id);
2702 let ann = reader.vector_indexes().get(&field_id);
2703 if ann.is_some() && flat.is_none() {
2704 return Err(Error::Corruption(format!(
2705 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2706 reader.meta().id,
2707 )));
2708 }
2709
2710 let Some(flat) = flat else {
2711 continue;
2712 };
2713 if flat.num_vectors == 0 {
2714 continue;
2715 }
2716 if rewrite_existing {
2717 return Ok(true);
2718 }
2719 let field = crate::dsl::Field(field_id);
2720 let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2721 Error::Corruption(format!(
2722 "segment {:032x} references unknown vector field {field_id}",
2723 reader.meta().id,
2724 ))
2725 })?;
2726 let current = match entry.field_type {
2727 crate::dsl::FieldType::DenseVector
2731 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2732 config.index_type == crate::dsl::VectorIndexType::Tq
2733 }) =>
2734 {
2735 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
2736 }
2737 crate::dsl::FieldType::DenseVector
2738 if entry.dense_vector_config.as_ref().is_some_and(|config| {
2739 config.index_type == crate::dsl::VectorIndexType::IvfTq
2740 }) =>
2741 {
2742 matches!(ann, Some(crate::segment::VectorIndex::IvfTq { .. }))
2743 }
2744 crate::dsl::FieldType::DenseVector => false,
2747 crate::dsl::FieldType::BinaryDenseVector => {
2748 matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2749 }
2750 _ => false,
2751 };
2752 if !current {
2753 return Ok(true);
2754 }
2755 }
2756 Ok(false)
2757 }
2758
2759 async fn acquire_vector_rewrite_capacity(
2760 &self,
2761 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2762 let global = tokio::select! {
2763 biased;
2764 () = self.active_operations.wait_for_shutdown() => {
2765 return Err(Error::IndexClosed);
2766 }
2767 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2768 permit.map_err(|_| Error::Internal(
2769 "global background merge scheduler is closed".into()
2770 ))?
2771 }
2772 };
2773 let local = tokio::select! {
2774 biased;
2775 () = self.active_operations.wait_for_shutdown() => {
2776 return Err(Error::IndexClosed);
2777 }
2778 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2779 permit.map_err(|_| Error::Internal(
2780 "background merge scheduler is closed".into()
2781 ))?
2782 }
2783 };
2784 Ok((global, local))
2785 }
2786
2787 async fn build_vector_replacement(
2788 self: &Arc<Self>,
2789 segment_id: &str,
2790 source_id: SegmentId,
2791 output_id: SegmentId,
2792 trained: &TrainedVectorStructures,
2793 failure_context: &'static str,
2794 ) -> Result<(String, u32, OutputCleanupGuard)> {
2795 let mut cleanup = self.output_cleanup_guard(output_id);
2796 match crate::segment::reorder::rewrite_vector_segment(
2797 self.directory.as_ref(),
2798 &self.schema,
2799 source_id,
2800 output_id,
2801 self.term_cache_blocks,
2802 trained,
2803 Some(self.background_cpu_pool()),
2804 )
2805 .await
2806 {
2807 Ok((new_id, doc_count)) => {
2808 self.validate_completed_segment(&new_id, doc_count).await?;
2809 Ok((new_id, doc_count, cleanup))
2810 }
2811 Err(error) => {
2812 self.delete_output_if_unregistered(output_id, failure_context)
2813 .await;
2814 cleanup.disarm();
2815 if is_deterministic_source_error(&error) {
2816 self.quarantine_segment(segment_id, &error);
2817 }
2818 Err(error)
2819 }
2820 }
2821 }
2822
2823 pub(crate) async fn stage_vector_generation(
2827 self: &Arc<Self>,
2828 _artifact_update: &VectorArtifactUpdateGuard,
2829 segment_ids: &[String],
2830 field_ids: &[u32],
2831 trained: Arc<TrainedVectorStructures>,
2832 rewrite_existing: bool,
2833 ) -> Result<Vec<StagedVectorSegment>> {
2834 if !self.vector_artifact_update.load(Ordering::Acquire) {
2835 return Err(Error::Internal(
2836 "cannot stage a vector generation without an exclusive update lease".into(),
2837 ));
2838 }
2839
2840 let mut staged = Vec::new();
2841 for segment_id in segment_ids {
2842 if self.quarantined_segments.lock().contains(segment_id) {
2843 return Err(Error::Corruption(format!(
2844 "segment {segment_id} is quarantined after a deterministic source failure"
2845 )));
2846 }
2847 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2848 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2849 })?;
2850
2851 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2854
2855 let output_id = SegmentId::new();
2856 let output_hex = output_id.to_hex();
2857 let operation = {
2858 let st = self.state.lock().await;
2859 if !st.metadata.has_segment(segment_id) {
2860 return Err(Error::Corruption(format!(
2861 "vector generation source {segment_id} disappeared while lifecycle work was paused"
2862 )));
2863 }
2864 self.active_operations
2865 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
2866 }
2867 .ok_or_else(|| {
2868 if self.active_operations.is_accepting() {
2869 Error::Internal(format!(
2870 "vector generation could not claim stable source {segment_id}"
2871 ))
2872 } else {
2873 Error::IndexClosed
2874 }
2875 })?;
2876
2877 let reader = SegmentReader::open(
2878 self.directory.as_ref(),
2879 source_id,
2880 Arc::clone(&self.schema),
2881 self.term_cache_blocks,
2882 )
2883 .await?;
2884 if !self.segment_needs_vector_rewrite(&reader, field_ids, rewrite_existing)? {
2885 continue;
2886 }
2887 drop(reader);
2888
2889 let (new_id, doc_count, cleanup) = self
2890 .build_vector_replacement(
2891 segment_id,
2892 source_id,
2893 output_id,
2894 trained.as_ref(),
2895 "vector generation staging failure",
2896 )
2897 .await?;
2898 debug_assert_eq!(new_id, output_hex);
2899 let output_reader = SegmentReader::open(
2900 self.directory.as_ref(),
2901 output_id,
2902 Arc::clone(&self.schema),
2903 self.term_cache_blocks,
2904 )
2905 .await?;
2906 if self.segment_needs_vector_rewrite(&output_reader, field_ids, false)? {
2907 return Err(Error::Corruption(format!(
2908 "staged vector segment {new_id} does not match its candidate codebook generation"
2909 )));
2910 }
2911
2912 staged.push(StagedVectorSegment {
2913 source_id: segment_id.clone(),
2914 output_id,
2915 doc_count,
2916 _operation: operation,
2917 cleanup,
2918 });
2919 }
2920 Ok(staged)
2921 }
2922
2923 async fn rewrite_vector_segment_once(
2924 self: &Arc<Self>,
2925 segment_id: &str,
2926 field_ids: &[u32],
2927 ) -> Result<VectorSegmentRewriteOutcome> {
2928 if self.quarantined_segments.lock().contains(segment_id) {
2929 return Err(Error::Corruption(format!(
2930 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
2931 )));
2932 }
2933 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2934 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2935 })?;
2936
2937 let _capacity = self.acquire_vector_rewrite_capacity().await?;
2942
2943 let output_id = SegmentId::new();
2944 let output_hex = output_id.to_hex();
2945 let all_ids = vec![segment_id.to_owned(), output_hex];
2946 let operation = {
2947 let st = self.state.lock().await;
2948 if !st.metadata.has_segment(segment_id) {
2949 return Ok(VectorSegmentRewriteOutcome::SourceGone);
2950 }
2951 self.active_operations.try_register(all_ids)
2952 };
2953 let _operation = match operation {
2954 Some(operation) => operation,
2955 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
2956 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
2957 };
2958
2959 let Some(trained) = self.trained_for_segment_build() else {
2960 return Ok(VectorSegmentRewriteOutcome::Deferred);
2961 };
2962
2963 let reader = SegmentReader::open(
2964 self.directory.as_ref(),
2965 source_id,
2966 Arc::clone(&self.schema),
2967 self.term_cache_blocks,
2968 )
2969 .await?;
2970 if !self.segment_needs_vector_rewrite(&reader, field_ids, false)? {
2971 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
2972 }
2973 drop(reader);
2974
2975 let (new_id, doc_count, mut output_cleanup) = self
2976 .build_vector_replacement(
2977 segment_id,
2978 source_id,
2979 output_id,
2980 trained.as_ref(),
2981 "vector rewrite failure",
2982 )
2983 .await?;
2984
2985 if let Err(error) = self
2986 .replace_segments(
2987 &[segment_id.to_owned()],
2988 new_id,
2989 doc_count,
2990 ReplacementLayout::PreserveSingleSource,
2991 )
2992 .await
2993 {
2994 self.delete_output_if_unregistered(output_id, "vector replacement failure")
2995 .await;
2996 output_cleanup.disarm();
2997 return Err(error);
2998 }
2999 output_cleanup.disarm();
3000 Ok(VectorSegmentRewriteOutcome::Rewritten)
3001 }
3002
3003 pub(crate) async fn rewrite_vector_segments(
3008 self: &Arc<Self>,
3009 field_ids: &[u32],
3010 ) -> Result<usize> {
3011 if field_ids.is_empty() {
3012 return Ok(0);
3013 }
3014 let mut rewritten = 0usize;
3015 loop {
3016 let segment_ids = self.get_segment_ids().await;
3017 let mut conflicted = false;
3018 let mut changed = false;
3019 for segment_id in segment_ids {
3020 match self
3021 .rewrite_vector_segment_once(&segment_id, field_ids)
3022 .await?
3023 {
3024 VectorSegmentRewriteOutcome::Rewritten => {
3025 rewritten += 1;
3026 changed = true;
3027 }
3028 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3029 VectorSegmentRewriteOutcome::Deferred => {
3030 return Err(Error::Internal(
3031 "ANN finalization lost the published trained generation".into(),
3032 ));
3033 }
3034 VectorSegmentRewriteOutcome::AlreadyCurrent
3035 | VectorSegmentRewriteOutcome::SourceGone => {}
3036 }
3037 }
3038 if !conflicted && !changed {
3039 log::info!(
3040 "[dense_vector_rewrite] ANN finalization complete ({} segment(s) rewritten)",
3041 rewritten,
3042 );
3043 return Ok(rewritten);
3044 }
3045 tokio::select! {
3046 biased;
3047 () = self.active_operations.wait_for_shutdown() => {
3048 return Err(Error::IndexClosed);
3049 }
3050 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3051 }
3052 }
3053 }
3054
3055 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3060 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3061 return;
3062 }
3063 let manager = Arc::clone(self);
3064 let future = async move {
3065 let field_ids = manager
3066 .read_metadata(|metadata| {
3067 metadata
3068 .vector_fields
3069 .keys()
3070 .filter(|field_id| metadata.is_field_built(**field_id))
3071 .copied()
3072 .collect::<Vec<_>>()
3073 })
3074 .await;
3075 for segment_id in segment_ids {
3076 loop {
3077 match manager
3078 .rewrite_vector_segment_once(&segment_id, &field_ids)
3079 .await
3080 {
3081 Ok(VectorSegmentRewriteOutcome::Conflict) => {
3082 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3083 }
3084 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3085 Ok(_) => break,
3086 Err(error) => {
3087 log::error!(
3088 "[dense_vector_rewrite] failed to upgrade newly committed segment {}: {}",
3089 segment_id,
3090 error,
3091 );
3092 break;
3093 }
3094 }
3095 }
3096 }
3097 };
3098 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3099 log::warn!(
3100 "[dense_vector_rewrite] runtime unavailable; newly committed flat segment upgrade deferred"
3101 );
3102 return;
3103 };
3104 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3105 log::warn!(
3106 "[dense_vector_rewrite] runtime rejected newly committed flat segment upgrade"
3107 );
3108 }
3109 }
3110
3111 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3118 self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3119 .await
3120 }
3121
3122 pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3125 self: &Arc<Self>,
3126 mut refresh_snapshots: F,
3127 ) -> Result<()>
3128 where
3129 F: FnMut() -> Fut,
3130 Fut: std::future::Future<Output = Result<()>>,
3131 {
3132 self.wait_for_all_merges().await;
3133 refresh_snapshots().await?;
3134 let segment_ids = self.get_segment_ids().await;
3135
3136 if segment_ids.is_empty() {
3137 log::info!("[reorder] no segments to reorder");
3138 return Ok(());
3139 }
3140
3141 log::info!("[reorder] reordering {} segments", segment_ids.len());
3142
3143 for seg_id in segment_ids {
3144 match self
3145 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3146 .await
3147 {
3148 Ok(true) => refresh_snapshots().await?,
3149 Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
3150 Err(e) => return Err(e),
3151 }
3152 }
3153
3154 refresh_snapshots().await?;
3157 log::info!("[reorder] all segments reordered");
3158 Ok(())
3159 }
3160
3161 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3166 self.unreordered_segments()
3167 .await
3168 .into_iter()
3169 .map(|(id, _)| id)
3170 .collect()
3171 }
3172
3173 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3176 let quarantined = self.quarantined_segments.lock().clone();
3177 let paused = self.paused_reorder_segments();
3178 let st = self.state.lock().await;
3179 let active_ids = self.active_operations.snapshot();
3180 st.metadata
3181 .segment_metas
3182 .iter()
3183 .filter(|(id, info)| {
3184 !info.reordered
3185 && info.bp_converged
3186 && !active_ids.contains(*id)
3187 && !quarantined.contains(*id)
3188 && !paused.contains(*id)
3189 })
3190 .map(|(id, info)| (id.clone(), info.num_docs))
3191 .collect()
3192 }
3193
3194 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3198 self.unconverged_segments_below(u32::MAX)
3199 .await
3200 .into_iter()
3201 .map(|(id, docs, _)| (id, docs))
3202 .collect()
3203 }
3204
3205 pub async fn unconverged_segments_below(
3208 &self,
3209 max_unconverged_passes: u32,
3210 ) -> Vec<(String, u32, u32)> {
3211 let quarantined = self.quarantined_segments.lock().clone();
3212 let paused = self.paused_reorder_segments();
3213 let st = self.state.lock().await;
3214 let active_ids = self.active_operations.snapshot();
3215 st.metadata
3216 .segment_metas
3217 .iter()
3218 .filter(|(id, info)| {
3219 !info.bp_converged
3220 && info.bp_unconverged_passes < max_unconverged_passes
3221 && !active_ids.contains(*id)
3222 && !quarantined.contains(*id)
3223 && !paused.contains(*id)
3224 })
3225 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3226 .collect()
3227 }
3228
3229 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3239 let st = self.state.lock().await;
3240 let deepening = ids.iter().any(|id| {
3241 st.metadata
3242 .segment_metas
3243 .get(id)
3244 .is_some_and(|info| !info.bp_converged)
3245 });
3246 drop(st);
3247 if deepening {
3248 log::info!(
3249 "[reorder] source BP lineage unconverged — forcing record-level BP (deepening pass)",
3250 );
3251 crate::segment::reorder::BpGranularity::Records
3252 } else {
3253 crate::segment::reorder::BpGranularity::Auto
3254 }
3255 }
3256
3257 pub async fn reorder_single_segment(
3262 self: &Arc<Self>,
3263 seg_id: &str,
3264 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3265 bp_budget: crate::segment::BpBudget,
3266 ) -> Result<bool> {
3267 let source_id = SegmentId::from_hex(seg_id)
3268 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3269 if self.quarantined_segments.lock().contains(seg_id) {
3270 return Err(Error::Corruption(format!(
3271 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3272 seg_id
3273 )));
3274 }
3275 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3276 log::debug!(
3277 "[optimizer] explicit force merge active, skipping reorder of {}",
3278 seg_id,
3279 );
3280 return Ok(false);
3281 }
3282
3283 let reorder_gate = Arc::clone(&self.reorder_permits);
3288 let _reorder_permit = tokio::select! {
3289 biased;
3290 () = self.active_operations.wait_for_shutdown() => {
3291 return Err(Error::IndexClosed);
3292 }
3293 permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3294 permit.map_err(|_| {
3295 Error::Internal("background reorder scheduler is closed".into())
3296 })?
3297 }
3298 };
3299
3300 let output_id = SegmentId::new();
3301 let output_hex = output_id.to_hex();
3302 let source_ids = [seg_id.to_string()];
3303 let granularity = self.merge_granularity(&source_ids).await;
3304
3305 let all_ids = vec![seg_id.to_string(), output_hex];
3311 let (_guard, source_docs) = {
3312 let st = self.state.lock().await;
3313 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3317 log::debug!(
3318 "[optimizer] explicit force merge active, skipping reorder of {}",
3319 seg_id,
3320 );
3321 return Ok(false);
3322 }
3323 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3324 log::info!(
3325 "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
3326 seg_id
3327 );
3328 self.clear_reorder_retry(seg_id);
3329 return Ok(false);
3330 };
3331
3332 match self.active_operations.try_register(all_ids) {
3333 Some(guard) => (guard, source_meta.num_docs),
3334 None if !self.active_operations.is_accepting() => {
3335 return Err(Error::IndexClosed);
3336 }
3337 None => {
3338 log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
3339 return Ok(false);
3340 }
3341 }
3342 };
3343
3344 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3349 if is_deterministic_source_error(&error) {
3350 self.quarantine_segment(seg_id, &error);
3351 } else if !matches!(&error, Error::IndexClosed) {
3352 self.pause_reorder_retries(seg_id, &error);
3353 }
3354 return Err(error);
3355 }
3356
3357 let mut output_cleanup = self.output_cleanup_guard(output_id);
3358
3359 let reorder_result = crate::segment::reorder::reorder_segment(
3360 self.directory.as_ref(),
3361 &self.schema,
3362 source_id,
3363 output_id,
3364 self.term_cache_blocks,
3365 self.bp_memory_budget_bytes,
3366 bp_budget,
3367 granularity,
3368 rayon_pool,
3369 )
3370 .await;
3371 let (new_id, total_docs, bp_converged) = match reorder_result {
3372 Ok(v) => v,
3373 Err(e) => {
3374 self.delete_output_if_unregistered(output_id, "reorder failure")
3377 .await;
3378 output_cleanup.disarm();
3379 if is_deterministic_source_error(&e) {
3380 self.quarantine_segment(seg_id, &e);
3381 } else if !matches!(&e, Error::IndexClosed) {
3382 self.pause_reorder_retries(seg_id, &e);
3383 }
3384 return Err(e);
3385 }
3386 };
3387
3388 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3394 if let Err(e) = self
3395 .replace_segments(
3396 &[seg_id.to_string()],
3397 new_id,
3398 total_docs,
3399 ReplacementLayout::BpReordered {
3400 converged: ladder_converged,
3401 },
3402 )
3403 .await
3404 {
3405 self.delete_output_if_unregistered(output_id, "replacement failure")
3406 .await;
3407 output_cleanup.disarm();
3408 if !matches!(&e, Error::IndexClosed) {
3409 self.pause_reorder_retries(seg_id, &e);
3410 }
3411 return Err(e);
3412 }
3413 output_cleanup.disarm();
3414 self.clear_reorder_retry(seg_id);
3415
3416 Ok(true)
3417 }
3418
3419 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3426 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3427
3428 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3429 for entry in entries {
3430 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3431 continue;
3432 };
3433 let Some(rest) = filename.strip_prefix("seg_") else {
3434 continue;
3435 };
3436 let Some(hex_id) = rest.get(..32) else {
3437 continue;
3438 };
3439 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3440 continue;
3441 }
3442 orphan_files
3443 .entry(hex_id.to_ascii_lowercase())
3444 .or_default()
3445 .push(entry);
3446 }
3447 }
3448
3449 let mut deleted = 0;
3450 for (hex_id, paths) in &orphan_files {
3451 let deletion_guard = {
3456 let st = self.state.lock().await;
3457 if st.metadata.has_segment(hex_id) {
3458 continue;
3459 }
3460 let Some(guard) = self
3461 .active_operations
3462 .try_register(vec![hex_id.to_string()])
3463 else {
3464 continue;
3465 };
3466 if self.tracker.is_deletion_protected(hex_id) {
3467 drop(guard);
3468 continue;
3469 }
3470 guard
3471 };
3472
3473 let results =
3478 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3479 .await;
3480 let removed = results.into_iter().all(|result| match result {
3481 Ok(()) => true,
3482 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3483 Err(error) => {
3484 log::warn!(
3485 "[segment_cleanup] failed sweeping orphan segment {}: {}",
3486 hex_id,
3487 error,
3488 );
3489 false
3490 }
3491 });
3492 drop(deletion_guard);
3495 if removed {
3496 deleted += 1;
3497 log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
3498 }
3499 }
3500
3501 Ok(deleted)
3502 }
3503}
3504
3505#[cfg(test)]
3506mod tests {
3507 use super::*;
3508 use std::sync::atomic::{AtomicBool, Ordering};
3509
3510 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3511 let schema = crate::dsl::SchemaBuilder::default().build();
3512 let metadata = IndexMetadata::new(schema.clone());
3513 Arc::new(SegmentManager::new(
3514 Arc::new(crate::directories::RamDirectory::new()),
3515 Arc::new(schema),
3516 metadata,
3517 Box::new(crate::merge::NoMergePolicy),
3518 0,
3519 1,
3520 Arc::new(Semaphore::new(1)),
3521 None,
3522 1024,
3523 Arc::new(ReorderConcurrencyGate::new(1)),
3524 None,
3525 ))
3526 }
3527
3528 #[test]
3529 fn force_merge_planner_pairs_large_and_small_segments() {
3530 let groups = plan_force_merge_groups(
3531 vec![
3532 ("a".into(), 6),
3533 ("b".into(), 6),
3534 ("c".into(), 4),
3535 ("d".into(), 4),
3536 ],
3537 10,
3538 );
3539
3540 assert_eq!(groups.len(), 2);
3541 assert!(groups.iter().all(|group| group.total_docs == 10));
3542 assert!(groups.iter().all(|group| group.segments.len() == 2));
3543 }
3544
3545 #[test]
3546 fn force_merge_planner_leaves_oversized_segments_alone() {
3547 let groups = plan_force_merge_groups(
3548 vec![
3549 ("oversized".into(), 11),
3550 ("small-a".into(), 5),
3551 ("small-b".into(), 5),
3552 ],
3553 10,
3554 );
3555
3556 assert_eq!(groups.len(), 2);
3557 assert_eq!(groups[0].total_docs, 10);
3558 assert_eq!(groups[0].segments.len(), 2);
3559 assert_eq!(groups[1].total_docs, 11);
3560 assert_eq!(groups[1].segments.len(), 1);
3561 }
3562
3563 #[test]
3564 fn force_merge_planner_never_exceeds_segment_format_limit() {
3565 let groups = plan_force_merge_groups(
3566 vec![
3567 ("large-a".into(), 3_000_000_000),
3568 ("large-b".into(), 2_000_000_000),
3569 ],
3570 u64::from(u32::MAX),
3571 );
3572 assert_eq!(groups.len(), 2);
3573 assert!(
3574 groups
3575 .iter()
3576 .all(|group| group.total_docs <= u64::from(u32::MAX))
3577 );
3578 }
3579
3580 #[test]
3581 fn force_merge_hierarchy_has_one_final_bp_pass() {
3582 assert_eq!(force_merge_output_count(1), 0);
3583 assert_eq!(force_merge_output_count(2), 1);
3584 assert_eq!(force_merge_output_count(64), 1);
3585 assert_eq!(force_merge_output_count(65), 2);
3586 assert_eq!(force_merge_output_count(127), 2);
3587 assert_eq!(force_merge_output_count(128), 3);
3588 assert_eq!(force_merge_output_count(1_000), 16);
3589 }
3590
3591 fn expand_force_merge_node(
3592 hierarchy: &ForceMergeHierarchy,
3593 source_count: usize,
3594 node: usize,
3595 sources: &mut Vec<usize>,
3596 ) {
3597 if node < source_count {
3598 sources.push(node);
3599 return;
3600 }
3601
3602 let step_index = node - source_count;
3603 let step = hierarchy
3604 .steps
3605 .get(step_index)
3606 .expect("merge input must refer to an existing source or output");
3607 for &input in &step.inputs {
3608 assert!(
3609 input < node,
3610 "merge step {step_index} refers to a future output node {input}"
3611 );
3612 expand_force_merge_node(hierarchy, source_count, input, sources);
3613 }
3614 }
3615
3616 #[test]
3617 fn force_merge_hierarchy_has_minimal_valid_arity() {
3618 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
3619
3620 for source_count in source_counts {
3621 let hierarchy = plan_force_merge_hierarchy(source_count);
3622 let output_count = hierarchy.steps.len();
3623
3624 assert!(
3625 hierarchy
3626 .steps
3627 .iter()
3628 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
3629 "invalid merge arity for {source_count} sources"
3630 );
3631 assert!(
3632 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
3633 "{output_count} outputs cannot reduce {source_count} sources"
3634 );
3635 assert!(
3636 output_count == 1
3637 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
3638 "{output_count} outputs are not minimal for {source_count} sources"
3639 );
3640 assert_eq!(output_count, force_merge_output_count(source_count));
3641 }
3642 }
3643
3644 #[test]
3645 fn force_merge_hierarchy_preserves_exact_source_order() {
3646 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
3647 let hierarchy = plan_force_merge_hierarchy(source_count);
3648 let mut sources = Vec::with_capacity(source_count);
3649 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
3650 assert_eq!(
3651 sources,
3652 (0..source_count).collect::<Vec<_>>(),
3653 "source order changed for {source_count} sources"
3654 );
3655 }
3656 }
3657
3658 fn force_merge_rewrite_cost(source_count: usize) -> usize {
3659 let hierarchy = plan_force_merge_hierarchy(source_count);
3660 let mut node_weights = vec![1usize; source_count];
3661 let mut rewrite_cost = 0usize;
3662
3663 for (step_index, step) in hierarchy.steps.iter().enumerate() {
3664 let output = source_count + step_index;
3665 let output_weight = step
3666 .inputs
3667 .iter()
3668 .map(|&input| {
3669 assert!(
3670 input < output,
3671 "merge step {step_index} refers to future output {input}"
3672 );
3673 node_weights[input]
3674 })
3675 .sum::<usize>();
3676 rewrite_cost += output_weight;
3677 node_weights.push(output_weight);
3678 }
3679
3680 assert_eq!(node_weights[hierarchy.root], source_count);
3681 rewrite_cost
3682 }
3683
3684 #[test]
3685 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
3686 assert_eq!(force_merge_rewrite_cost(65), 67);
3687 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
3688 }
3689
3690 #[test]
3691 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
3692 assert_eq!(
3693 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
3694 (false, false, 3),
3695 );
3696 assert_eq!(
3697 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
3698 (true, false, 4),
3699 );
3700 assert_eq!(
3701 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
3702 (true, true, 0),
3703 );
3704 }
3705
3706 #[tokio::test]
3707 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
3708 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3709 let field = schema_builder.add_text_field("text", true, true);
3710 let schema = schema_builder.build();
3711 let directory = crate::directories::RamDirectory::new();
3712 let config = crate::index::IndexConfig {
3713 num_indexing_threads: 1,
3714 merge_policy: Box::new(crate::merge::NoMergePolicy),
3715 ..Default::default()
3716 };
3717 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
3718 .await
3719 .unwrap();
3720 for value in ["one", "two", "three"] {
3721 let mut document = crate::dsl::Document::new();
3722 document.add_text(field, value);
3723 writer.add_document(document).unwrap();
3724 writer.commit().await.unwrap();
3725 }
3726
3727 let manager = Arc::clone(writer.segment_manager());
3728 let held_id = manager.get_segment_ids().await.pop().unwrap();
3729 let mut held = Some(
3730 manager
3731 .active_operations
3732 .try_register(vec![held_id])
3733 .unwrap(),
3734 );
3735 let batches = Arc::new(AtomicUsize::new(0));
3736 let batch_count = Arc::clone(&batches);
3737 writer
3738 .force_merge_with_snapshot_refresh(move || {
3739 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
3740 if refresh == 2 {
3743 drop(held.take());
3744 }
3745 std::future::ready(Ok(()))
3746 })
3747 .await
3748 .unwrap();
3749
3750 assert_eq!(manager.get_segment_ids().await.len(), 1);
3751 assert_eq!(
3752 batches.load(Ordering::Relaxed),
3753 4,
3754 "initial/final refreshes plus two replacements are required after the held source releases"
3755 );
3756 }
3757
3758 #[tokio::test]
3759 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
3760 let mut schema_builder = crate::dsl::SchemaBuilder::default();
3761 schema_builder.set_reorder_on_merge(true);
3762 let schema = schema_builder.build();
3763 let mut metadata = IndexMetadata::new(schema.clone());
3764 metadata.add_segment("00000000000000000000000000000001".into(), 1);
3765 metadata.add_segment("00000000000000000000000000000002".into(), 1);
3766
3767 let global_merge_permits = Arc::new(Semaphore::new(1));
3768 let manager = Arc::new(SegmentManager::new(
3769 Arc::new(crate::directories::RamDirectory::new()),
3770 Arc::new(schema),
3771 metadata,
3772 Box::new(crate::merge::NoMergePolicy),
3773 0,
3774 1,
3775 Arc::clone(&global_merge_permits),
3776 None,
3777 1024,
3778 Arc::new(ReorderConcurrencyGate::new(1)),
3779 None,
3780 ));
3781
3782 manager.active_operations.pause_non_indexing();
3787 let force_merge = {
3788 let manager = Arc::clone(&manager);
3789 tokio::spawn(async move { manager.force_merge().await })
3790 };
3791 tokio::time::timeout(std::time::Duration::from_secs(1), async {
3792 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
3793 tokio::task::yield_now().await;
3794 }
3795 })
3796 .await
3797 .expect("force merge never reached the paused group claim");
3798
3799 assert_eq!(
3800 global_merge_permits.available_permits(),
3801 1,
3802 "force merge retained global capacity while vector staging blocked group ownership"
3803 );
3804
3805 force_merge.abort();
3806 let _ = force_merge.await;
3807 manager.active_operations.resume_non_indexing();
3808 }
3809
3810 #[test]
3811 fn output_cleanup_guard_runs_during_panic_unwind() {
3812 let cleaned = Arc::new(AtomicBool::new(false));
3813 let cleaned_in_callback = Arc::clone(&cleaned);
3814 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3815 cleaned_in_callback.store(true, Ordering::SeqCst);
3816 });
3817
3818 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3819 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3820 panic!("simulated reorder panic");
3821 }));
3822
3823 assert!(result.is_err());
3824 assert!(
3825 cleaned.load(Ordering::SeqCst),
3826 "partial output cleanup must run during unwind"
3827 );
3828 }
3829
3830 #[test]
3831 fn output_cleanup_guard_disarms_after_commit() {
3832 let cleaned = Arc::new(AtomicBool::new(false));
3833 let cleaned_in_callback = Arc::clone(&cleaned);
3834 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3835 cleaned_in_callback.store(true, Ordering::SeqCst);
3836 });
3837
3838 {
3839 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3840 guard.disarm();
3841 }
3842
3843 assert!(!cleaned.load(Ordering::SeqCst));
3844 }
3845
3846 #[test]
3847 fn test_active_operation_guard_releases_ownership() {
3848 let active = Arc::new(ActiveSegmentOperations::new());
3849 {
3850 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3851 let snap = active.snapshot();
3852 assert!(snap.contains("a"));
3853 assert!(snap.contains("b"));
3854 }
3855 assert!(active.snapshot().is_empty());
3856 }
3857
3858 #[test]
3859 fn test_non_overlapping_operations_can_run_concurrently() {
3860 let active = Arc::new(ActiveSegmentOperations::new());
3861 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3862 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
3863 let snap = active.snapshot();
3864 assert_eq!(snap.len(), 4);
3865
3866 drop(first);
3867 let snap = active.snapshot();
3868 assert_eq!(snap.len(), 2);
3869 assert!(snap.contains("c"));
3870 assert!(snap.contains("d"));
3871 }
3872
3873 #[test]
3874 fn test_overlapping_operation_is_rejected_until_release() {
3875 let active = Arc::new(ActiveSegmentOperations::new());
3876 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3877 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
3878 drop(first);
3879 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
3880 }
3881
3882 #[test]
3883 fn test_active_operation_snapshot() {
3884 let active = Arc::new(ActiveSegmentOperations::new());
3885 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
3886 let snap = active.snapshot();
3887 assert!(snap.contains("x"));
3888 assert!(snap.contains("y"));
3889 assert!(!snap.contains("z"));
3890 }
3891
3892 #[tokio::test]
3893 async fn operation_barrier_ignores_producers_started_after_snapshot() {
3894 let active = Arc::new(ActiveSegmentOperations::new());
3895 let before_gate = active.try_register(vec!["old".into()]).unwrap();
3896 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
3897 assert_eq!(parked_indexing, 0);
3898 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
3899
3900 let waiter = {
3901 let active = Arc::clone(&active);
3902 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
3903 };
3904 tokio::task::yield_now().await;
3905 assert!(!waiter.is_finished());
3906
3907 drop(before_gate);
3908 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3909 .await
3910 .expect("pre-gate operation barrier was starved by a post-gate producer")
3911 .unwrap();
3912 assert!(active.snapshot().contains("new-flat"));
3913 drop(after_gate);
3914 }
3915
3916 #[tokio::test]
3917 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
3918 let manager = lifecycle_test_manager();
3919 manager
3920 .trained
3921 .store(Some(Arc::new(TrainedVectorStructures {
3922 centroids: rustc_hash::FxHashMap::default(),
3923 binary_quantizers: rustc_hash::FxHashMap::default(),
3924 ..Default::default()
3925 })));
3926
3927 let guard = manager.begin_vector_artifact_update().await.unwrap();
3928 assert!(
3929 manager.trained().is_some(),
3930 "search readers keep the last fully validated generation"
3931 );
3932 assert!(
3933 manager.trained_for_segment_build().is_none(),
3934 "new segment producers must stay flat during an artifact update"
3935 );
3936
3937 let detached_transaction_guard = guard.clone();
3938 drop(guard);
3939 assert!(
3940 manager.trained_for_segment_build().is_none(),
3941 "a detached lifecycle transaction must retain the producer gate after request cancellation"
3942 );
3943 drop(detached_transaction_guard);
3944 assert!(manager.trained_for_segment_build().is_some());
3945 }
3946
3947 #[tokio::test]
3948 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
3949 let manager = lifecycle_test_manager();
3950 let guard = manager.begin_vector_artifact_update().await.unwrap();
3951 assert!(
3952 manager
3953 .active_operations
3954 .try_register(vec!["merge".into()])
3955 .is_none(),
3956 "ordinary merge/reorder work must not change staged sources"
3957 );
3958 let indexing = manager
3959 .active_operations
3960 .try_register_indexing(vec!["fresh".into()])
3961 .expect("indexing remains available in flat mode");
3962 drop(indexing);
3963
3964 drop(guard);
3965 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
3966 assert!(
3967 manager
3968 .active_operations
3969 .try_register(vec!["merge".into()])
3970 .is_some()
3971 );
3972 }
3973
3974 #[tokio::test]
3975 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
3976 let active = Arc::new(ActiveSegmentOperations::new());
3977 let guard = active.try_register(vec!["live".into()]).unwrap();
3978 active.stop_accepting();
3979 assert!(active.try_register(vec!["new".into()]).is_none());
3980
3981 let waiter = {
3982 let active = Arc::clone(&active);
3983 tokio::spawn(async move { active.wait_until_idle().await })
3984 };
3985 tokio::task::yield_now().await;
3986 assert!(!waiter.is_finished());
3987 drop(guard);
3988 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3989 .await
3990 .expect("shutdown waiter missed the final guard notification")
3991 .unwrap();
3992 }
3993
3994 #[tokio::test]
3995 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
3996 let manager = lifecycle_test_manager();
3997 let started = Arc::new(Semaphore::new(0));
3998 let release = Arc::new(Semaphore::new(0));
3999 let completed = Arc::new(AtomicBool::new(false));
4000
4001 let request = {
4002 let manager = Arc::clone(&manager);
4003 let started = Arc::clone(&started);
4004 let release = Arc::clone(&release);
4005 let completed = Arc::clone(&completed);
4006 tokio::spawn(async move {
4007 manager
4008 .run_lifecycle_transaction(async move {
4009 started.add_permits(1);
4010 let _permit = release.acquire().await.unwrap();
4011 completed.store(true, Ordering::Release);
4012 Ok(())
4013 })
4014 .await
4015 })
4016 };
4017
4018 let _started = started.acquire().await.unwrap();
4019 request.abort();
4020 assert!(request.await.unwrap_err().is_cancelled());
4021 release.add_permits(1);
4022
4023 manager.begin_shutdown();
4024 tokio::time::timeout(
4025 std::time::Duration::from_secs(1),
4026 manager.wait_for_shutdown(),
4027 )
4028 .await
4029 .expect("shutdown did not drain detached lifecycle transaction");
4030 assert!(completed.load(Ordering::Acquire));
4031 }
4032
4033 #[tokio::test]
4034 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4035 let manager = lifecycle_test_manager();
4036 {
4037 let mut state = manager.state.lock().await;
4038 state.metadata.add_segment_meta(
4039 "eligible".into(),
4040 SegmentMetaInfo {
4041 num_docs: 10,
4042 ancestors: Vec::new(),
4043 generation: 1,
4044 reordered: true,
4045 bp_converged: false,
4046 bp_unconverged_passes: 2,
4047 },
4048 );
4049 state.metadata.add_segment_meta(
4050 "at-limit".into(),
4051 SegmentMetaInfo {
4052 num_docs: 20,
4053 ancestors: Vec::new(),
4054 generation: 1,
4055 reordered: true,
4056 bp_converged: false,
4057 bp_unconverged_passes: 3,
4058 },
4059 );
4060 state.metadata.add_segment_meta(
4061 "carried-debt".into(),
4062 SegmentMetaInfo {
4063 num_docs: 15,
4064 ancestors: Vec::new(),
4065 generation: 2,
4066 reordered: false,
4067 bp_converged: false,
4068 bp_unconverged_passes: 2,
4069 },
4070 );
4071 state.metadata.add_segment_meta(
4072 "carried-debt-at-limit".into(),
4073 SegmentMetaInfo {
4074 num_docs: 25,
4075 ancestors: Vec::new(),
4076 generation: 2,
4077 reordered: false,
4078 bp_converged: false,
4079 bp_unconverged_passes: 3,
4080 },
4081 );
4082 state.metadata.add_segment_meta(
4083 "converged".into(),
4084 SegmentMetaInfo {
4085 num_docs: 30,
4086 ancestors: Vec::new(),
4087 generation: 1,
4088 reordered: true,
4089 bp_converged: true,
4090 bp_unconverged_passes: 0,
4091 },
4092 );
4093 state.metadata.add_segment("fresh".into(), 40);
4094 }
4095
4096 assert_eq!(
4097 manager.unreordered_segments().await,
4098 vec![("fresh".into(), 40)],
4099 "a block-copy output with BP debt is not a fresh first-pass candidate",
4100 );
4101 let mut eligible = manager.unconverged_segments_below(3).await;
4102 eligible.sort_unstable();
4103 assert_eq!(
4104 eligible,
4105 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4106 );
4107 assert!(manager.unconverged_segments_below(0).await.is_empty());
4108 }
4109
4110 #[test]
4111 fn merge_retry_backoff_is_exponential_and_capped() {
4112 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4113 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4114 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4115 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4116 }
4117
4118 #[test]
4119 fn only_deterministic_source_errors_are_quarantined() {
4120 assert!(is_deterministic_source_error(&Error::Corruption(
4121 "bad footer".into()
4122 )));
4123 assert!(is_deterministic_source_error(&Error::Io(
4124 std::io::Error::from(std::io::ErrorKind::NotFound)
4125 )));
4126 assert!(!is_deterministic_source_error(&Error::Io(
4127 std::io::Error::from(std::io::ErrorKind::TimedOut)
4128 )));
4129 assert!(!is_deterministic_source_error(&Error::Io(
4130 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4131 )));
4132 }
4133
4134 #[test]
4135 fn transient_reorder_failure_is_backed_off_until_cleared() {
4136 let manager = lifecycle_test_manager();
4137 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4138 assert!(manager.paused_reorder_segments().contains("source"));
4139 manager.clear_reorder_retry("source");
4140 assert!(!manager.paused_reorder_segments().contains("source"));
4141 }
4142
4143 #[derive(Default)]
4146 struct FailingExistsDirectory(crate::directories::RamDirectory);
4147
4148 #[async_trait::async_trait]
4149 impl crate::directories::Directory for FailingExistsDirectory {
4150 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4151 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4152 }
4153
4154 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4155 self.0.file_size(path).await
4156 }
4157
4158 async fn open_read(
4159 &self,
4160 path: &std::path::Path,
4161 ) -> std::io::Result<crate::directories::FileHandle> {
4162 self.0.open_read(path).await
4163 }
4164
4165 async fn read_range(
4166 &self,
4167 path: &std::path::Path,
4168 range: std::ops::Range<u64>,
4169 ) -> std::io::Result<crate::directories::OwnedBytes> {
4170 self.0.read_range(path, range).await
4171 }
4172
4173 async fn list_files(
4174 &self,
4175 prefix: &std::path::Path,
4176 ) -> std::io::Result<Vec<std::path::PathBuf>> {
4177 self.0.list_files(prefix).await
4178 }
4179
4180 async fn open_lazy(
4181 &self,
4182 path: &std::path::Path,
4183 ) -> std::io::Result<crate::directories::FileHandle> {
4184 self.0.open_lazy(path).await
4185 }
4186 }
4187
4188 #[async_trait::async_trait]
4189 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4190 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4191 self.0.write(path, data).await
4192 }
4193
4194 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4195 self.0.delete(path).await
4196 }
4197
4198 async fn rename(
4199 &self,
4200 from: &std::path::Path,
4201 to: &std::path::Path,
4202 ) -> std::io::Result<()> {
4203 self.0.rename(from, to).await
4204 }
4205
4206 async fn sync(&self) -> std::io::Result<()> {
4207 self.0.sync().await
4208 }
4209
4210 async fn streaming_writer(
4211 &self,
4212 path: &std::path::Path,
4213 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4214 self.0.streaming_writer(path).await
4215 }
4216 }
4217
4218 #[derive(Debug, Clone)]
4219 struct MergeEverythingPolicy;
4220
4221 impl MergePolicy for MergeEverythingPolicy {
4222 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4223 if segments.len() < 2 {
4224 return Vec::new();
4225 }
4226 vec![crate::merge::MergeCandidate {
4227 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4228 }]
4229 }
4230
4231 fn clone_box(&self) -> Box<dyn MergePolicy> {
4232 Box::new(self.clone())
4233 }
4234 }
4235
4236 #[tokio::test]
4237 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4238 let manager = lifecycle_test_manager();
4239 let parked_indexing = manager
4244 .protect_new_segment("00000000000000000000000000000abc".into())
4245 .unwrap();
4246
4247 let error = tokio::time::timeout(
4248 std::time::Duration::from_secs(2),
4249 manager.begin_vector_artifact_update(),
4250 )
4251 .await
4252 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4253 .err()
4254 .expect("an old-generation prepared segment must block artifact replacement")
4255 .to_string();
4256 assert!(error.contains("built but uncommitted"), "{error}");
4257 assert!(
4258 !manager.vector_artifact_update.load(Ordering::Acquire),
4259 "a rejected update must release the producer gate"
4260 );
4261
4262 drop(parked_indexing);
4263
4264 let guard = manager
4265 .begin_vector_artifact_update()
4266 .await
4267 .expect("artifact update should succeed after the pending generation is resolved");
4268 drop(guard);
4269 }
4270
4271 #[tokio::test]
4272 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4273 let manager = lifecycle_test_manager();
4274 let merge_like = manager
4275 .active_operations
4276 .try_register(vec!["merge-source".into()])
4277 .unwrap();
4278
4279 let waiter = {
4280 let manager = Arc::clone(&manager);
4281 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4282 };
4283 for _ in 0..8 {
4284 tokio::task::yield_now().await;
4285 }
4286 assert!(
4287 !waiter.is_finished(),
4288 "artifact update must drain merge/reorder producers that may hold the previous generation"
4289 );
4290
4291 drop(merge_like);
4292 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4293 .await
4294 .expect("artifact update missed the lifecycle guard release")
4295 .unwrap()
4296 .unwrap();
4297 }
4298
4299 #[tokio::test]
4300 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4301 let manager = lifecycle_test_manager();
4302 let release = Arc::new(Semaphore::new(0));
4303 let merge_task = {
4304 let release = Arc::clone(&release);
4305 tokio::spawn(async move {
4306 let _permit = release.acquire().await.unwrap();
4307 })
4308 };
4309 manager.merge_handles.lock().push(merge_task);
4310
4311 let waiter = {
4312 let manager = Arc::clone(&manager);
4313 tokio::spawn(async move { manager.wait_for_all_merges().await })
4314 };
4315 for _ in 0..8 {
4316 tokio::task::yield_now().await;
4317 }
4318 assert!(!waiter.is_finished());
4319 waiter.abort();
4322 let join_error = waiter.await.unwrap_err();
4323 assert!(join_error.is_cancelled());
4324
4325 assert!(
4326 !manager.merge_handles.lock().is_empty(),
4327 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4328 );
4329
4330 release.add_permits(1);
4332 tokio::time::timeout(
4333 std::time::Duration::from_secs(1),
4334 manager.wait_for_all_merges(),
4335 )
4336 .await
4337 .expect("subsequent drain missed the reinserted merge handle");
4338 assert!(manager.merge_handles.lock().is_empty());
4339 }
4340
4341 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4342 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4343 let manager = lifecycle_test_manager();
4344 {
4345 let mut state = manager.state.lock().await;
4346 state
4347 .metadata
4348 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4349 state
4350 .metadata
4351 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4352 }
4353 let reorder_like = manager
4356 .active_operations
4357 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4358 .unwrap();
4359
4360 let force_merge = {
4361 let manager = Arc::clone(&manager);
4362 tokio::spawn(async move { manager.force_merge().await })
4363 };
4364
4365 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4366 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4367 assert!(
4368 retries >= 1,
4369 "force_merge never observed the conflicting owner (retries={retries})"
4370 );
4371 assert!(
4372 retries < 20,
4373 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4374 );
4375
4376 drop(reorder_like);
4377 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4380 .await
4381 .expect("force_merge kept spinning after the conflicting owner released")
4382 .unwrap();
4383 assert!(result.is_err());
4384 }
4385
4386 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4387 async fn force_merge_routes_around_segments_held_by_reorder() {
4388 let manager = lifecycle_test_manager();
4389 {
4390 let mut state = manager.state.lock().await;
4391 state
4392 .metadata
4393 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4394 state
4395 .metadata
4396 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4397 state
4398 .metadata
4399 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4400 }
4401 let _reorder_like = manager
4404 .active_operations
4405 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4406 .unwrap();
4407
4408 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4415 let manager = Arc::clone(&manager);
4416 async move { manager.force_merge().await }
4417 })
4418 .await
4419 .expect("force_merge livelocked on a segment held by an active reorder");
4420 assert!(result.is_err(), "fake segment files must fail the merge");
4421
4422 assert_eq!(
4423 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4424 0,
4425 "batch built from the ownership snapshot must not collide with the held segment"
4426 );
4427 }
4428
4429 #[tokio::test]
4430 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4431 let schema = crate::dsl::SchemaBuilder::default().build();
4432 let mut metadata = IndexMetadata::new(schema.clone());
4433 metadata.add_segment("00000000000000000000000000000001".into(), 10);
4434 metadata.add_segment("00000000000000000000000000000002".into(), 10);
4435 let manager = Arc::new(SegmentManager::new(
4436 Arc::new(FailingExistsDirectory::default()),
4437 Arc::new(schema),
4438 metadata,
4439 Box::new(MergeEverythingPolicy),
4440 0,
4441 1,
4442 Arc::new(Semaphore::new(1)),
4443 None,
4444 1024,
4445 Arc::new(ReorderConcurrencyGate::new(1)),
4446 None,
4447 ));
4448
4449 manager.maybe_merge().await;
4452
4453 tokio::time::timeout(
4454 std::time::Duration::from_secs(5),
4455 manager.wait_for_all_merges(),
4456 )
4457 .await
4458 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4459 assert!(
4460 manager.merge_retry_is_paused(),
4461 "the failed merge should have armed the retry backoff"
4462 );
4463
4464 manager.begin_shutdown();
4466 tokio::time::timeout(
4467 std::time::Duration::from_secs(5),
4468 manager.wait_for_shutdown(),
4469 )
4470 .await
4471 .expect("shutdown did not drain the merge retry wakeup task");
4472 }
4473}