1use std::collections::{HashMap, HashSet};
45use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwap;
49use tokio::sync::Mutex as AsyncMutex;
50use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore};
51use tokio::task::JoinHandle;
52
53use crate::directories::DirectoryWriter;
54use crate::error::{Error, Result};
55use crate::index::{IndexMetadata, ReorderConcurrencyGate, ReorderPriority, SegmentMetaInfo};
56use crate::segment::{
57 PublishedIndexGeneration, SegmentFiles, SegmentId, SegmentMeta, SegmentSnapshot,
58 SegmentTracker, TrainedVectorStructures,
59};
60#[cfg(feature = "native")]
61use crate::segment::{SegmentMerger, SegmentReader};
62
63use super::{MergePolicy, SegmentInfo};
64
65const FORCE_MERGE_MAX_FAN_IN: usize = 64;
66
67#[derive(Debug)]
68struct ForceMergeGroup {
69 segments: Vec<(String, u32)>,
70 total_docs: u64,
71}
72
73fn plan_force_merge_groups(
82 mut segments: Vec<(String, u32)>,
83 max_docs: u64,
84) -> Vec<ForceMergeGroup> {
85 segments.sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
86 right_docs
87 .cmp(left_docs)
88 .then_with(|| left_id.cmp(right_id))
89 });
90
91 let mut groups: Vec<ForceMergeGroup> = Vec::new();
92 for segment in segments {
93 let docs = u64::from(segment.1);
94 let best_group = groups
95 .iter()
96 .enumerate()
97 .filter_map(|(index, group)| {
98 group
99 .total_docs
100 .checked_add(docs)
101 .filter(|&total| total <= max_docs)
102 .map(|_| (index, group.total_docs))
103 })
104 .max_by_key(|&(index, used)| (used, std::cmp::Reverse(index)))
105 .map(|(index, _)| index);
106
107 if let Some(index) = best_group {
108 groups[index].total_docs += docs;
109 groups[index].segments.push(segment);
110 } else {
111 groups.push(ForceMergeGroup {
112 segments: vec![segment],
113 total_docs: docs,
114 });
115 }
116 }
117
118 for group in &mut groups {
122 group
123 .segments
124 .sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
125 left_docs
126 .cmp(right_docs)
127 .then_with(|| left_id.cmp(right_id))
128 });
129 }
130
131 groups.sort_unstable_by(|left, right| {
134 left.total_docs
135 .cmp(&right.total_docs)
136 .then_with(|| left.segments[0].0.cmp(&right.segments[0].0))
137 });
138 groups
139}
140
141fn force_merge_output_count(source_count: usize) -> usize {
142 if source_count < 2 {
143 return 0;
144 }
145 (source_count - 1).div_ceil(FORCE_MERGE_MAX_FAN_IN - 1)
146}
147
148#[derive(Debug)]
149struct ForceMergeStep {
150 inputs: Vec<usize>,
152}
153
154#[derive(Debug)]
155struct ForceMergeHierarchy {
156 steps: Vec<ForceMergeStep>,
157 root: usize,
158}
159
160fn plan_force_merge_hierarchy(source_count: usize) -> ForceMergeHierarchy {
165 debug_assert!(source_count >= 2);
166 let internal_count = force_merge_output_count(source_count);
167 let max_leaves = 1usize
168 .checked_add(internal_count.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
169 .expect("force-merge hierarchy size exceeds usize");
170 let deficit = max_leaves - source_count;
171 debug_assert!(deficit < FORCE_MERGE_MAX_FAN_IN - 1);
172
173 fn build(
174 leaf_start: usize,
175 leaf_count: usize,
176 internal_count: usize,
177 deficit: usize,
178 source_count: usize,
179 steps: &mut Vec<ForceMergeStep>,
180 ) -> usize {
181 debug_assert!(internal_count > 0);
182 if internal_count == 1 {
183 let arity = FORCE_MERGE_MAX_FAN_IN - deficit;
184 debug_assert_eq!(leaf_count, arity);
185 debug_assert!((2..=FORCE_MERGE_MAX_FAN_IN).contains(&arity));
186 let output = source_count + steps.len();
187 steps.push(ForceMergeStep {
188 inputs: (leaf_start..leaf_start + arity).collect(),
189 });
190 return output;
191 }
192
193 let child_internal_total = internal_count - 1;
198 let base = child_internal_total / FORCE_MERGE_MAX_FAN_IN;
199 let extra = child_internal_total % FORCE_MERGE_MAX_FAN_IN;
200 let mut child_internal = vec![base; FORCE_MERGE_MAX_FAN_IN];
201 for count in &mut child_internal[..extra] {
202 *count += 1;
203 }
204 let partial_child = (deficit > 0).then(|| {
205 child_internal
206 .iter()
207 .position(|&count| count > 0)
208 .expect("a non-root partial node requires an internal child")
209 });
210
211 let mut cursor = leaf_start;
212 let mut inputs = Vec::with_capacity(FORCE_MERGE_MAX_FAN_IN);
213 for (child, &child_internals) in child_internal.iter().enumerate() {
214 if child_internals == 0 {
215 inputs.push(cursor);
216 cursor += 1;
217 continue;
218 }
219 let child_deficit = usize::from(partial_child == Some(child)) * deficit;
220 let child_leaves = 1usize
221 .checked_add(child_internals.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
222 .and_then(|maximum| maximum.checked_sub(child_deficit))
223 .expect("force-merge child size exceeds usize");
224 inputs.push(build(
225 cursor,
226 child_leaves,
227 child_internals,
228 child_deficit,
229 source_count,
230 steps,
231 ));
232 cursor += child_leaves;
233 }
234 debug_assert_eq!(cursor, leaf_start + leaf_count);
235 let output = source_count + steps.len();
236 steps.push(ForceMergeStep { inputs });
237 output
238 }
239
240 let mut steps = Vec::with_capacity(internal_count);
241 let root = build(
242 0,
243 source_count,
244 internal_count,
245 deficit,
246 source_count,
247 &mut steps,
248 );
249 debug_assert_eq!(steps.len(), internal_count);
250 ForceMergeHierarchy { steps, root }
251}
252
253struct ActiveOperationState {
263 segment_ids: HashSet<String>,
264 operation_tokens: HashSet<u64>,
265 indexing_tokens: HashSet<u64>,
271 next_operation_token: u64,
272 accepting: bool,
273 non_indexing_paused: bool,
277}
278
279struct ActiveSegmentOperations {
280 inner: parking_lot::Mutex<ActiveOperationState>,
281 idle: Notify,
282 shutdown: Notify,
283 shutdown_requested: Arc<AtomicBool>,
284 index_label: Arc<str>,
287}
288
289impl ActiveSegmentOperations {
290 fn new(index_label: Arc<str>) -> Self {
291 Self {
292 inner: parking_lot::Mutex::new(ActiveOperationState {
293 segment_ids: HashSet::new(),
294 operation_tokens: HashSet::new(),
295 indexing_tokens: HashSet::new(),
296 next_operation_token: 0,
297 accepting: true,
298 non_indexing_paused: false,
299 }),
300 idle: Notify::new(),
301 shutdown: Notify::new(),
302 shutdown_requested: Arc::new(AtomicBool::new(false)),
303 index_label,
304 }
305 }
306
307 fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
311 self.try_register_kind(segment_ids, false, false)
312 }
313
314 fn try_register_indexing(
317 self: &Arc<Self>,
318 segment_ids: Vec<String>,
319 ) -> Option<SegmentOperationGuard> {
320 self.try_register_kind(segment_ids, true, false)
321 }
322
323 fn try_register_vector_update(
326 self: &Arc<Self>,
327 segment_ids: Vec<String>,
328 ) -> Option<SegmentOperationGuard> {
329 self.try_register_kind(segment_ids, false, true)
330 }
331
332 fn try_register_kind(
333 self: &Arc<Self>,
334 segment_ids: Vec<String>,
335 indexing: bool,
336 vector_update: bool,
337 ) -> Option<SegmentOperationGuard> {
338 let mut inner = self.inner.lock();
339 if !inner.accepting {
340 log::debug!(
341 "[segment_lifecycle] index={} rejected operation during shutdown",
342 self.index_label
343 );
344 return None;
345 }
346 if !indexing && !vector_update && inner.non_indexing_paused {
347 log::debug!(
348 "[segment_lifecycle] index={} deferred operation during dense vector retraining",
349 self.index_label
350 );
351 return None;
352 }
353 for id in &segment_ids {
355 if inner.segment_ids.contains(id) {
356 log::debug!(
357 "[segment_lifecycle] index={} rejected: {} overlaps with an active operation ({} active IDs)",
358 self.index_label,
359 id,
360 inner.segment_ids.len()
361 );
362 return None;
363 }
364 }
365 log::debug!(
366 "[segment_lifecycle] index={} registered {} IDs (total active: {})",
367 self.index_label,
368 segment_ids.len(),
369 inner.segment_ids.len() + segment_ids.len()
370 );
371 let operation_token = inner.next_operation_token;
372 let next_operation_token = operation_token.checked_add(1)?;
373 for id in &segment_ids {
374 inner.segment_ids.insert(id.clone());
375 }
376 inner.next_operation_token = next_operation_token;
377 inner.operation_tokens.insert(operation_token);
378 if indexing {
379 inner.indexing_tokens.insert(operation_token);
380 }
381 Some(SegmentOperationGuard {
382 active_operations: Arc::clone(self),
383 segment_ids,
384 operation_token,
385 })
386 }
387
388 fn snapshot(&self) -> HashSet<String> {
390 self.inner.lock().segment_ids.clone()
391 }
392
393 fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
404 let inner = self.inner.lock();
405 let tokens = inner
406 .operation_tokens
407 .difference(&inner.indexing_tokens)
408 .copied()
409 .collect();
410 (tokens, inner.indexing_tokens.len())
411 }
412
413 fn stop_accepting(&self) {
416 self.shutdown_requested.store(true, Ordering::Release);
417 let mut inner = self.inner.lock();
418 inner.accepting = false;
419 self.shutdown.notify_waiters();
420 if inner.segment_ids.is_empty() {
421 self.idle.notify_waiters();
422 }
423 }
424
425 fn pause_non_indexing(&self) {
426 self.inner.lock().non_indexing_paused = true;
427 }
428
429 fn resume_non_indexing(&self) {
430 self.inner.lock().non_indexing_paused = false;
431 self.idle.notify_waiters();
432 }
433
434 fn is_accepting(&self) -> bool {
435 self.inner.lock().accepting
436 }
437
438 fn cancellation_flag(&self) -> Arc<AtomicBool> {
439 Arc::clone(&self.shutdown_requested)
440 }
441
442 async fn wait_until_idle(&self) {
446 loop {
447 let notified = self.idle.notified();
448 if self.inner.lock().segment_ids.is_empty() {
449 return;
450 }
451 notified.await;
452 }
453 }
454
455 async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
456 while !operations.is_empty() {
457 let notified = self.idle.notified();
458 if self.inner.lock().operation_tokens.is_disjoint(operations) {
459 return;
460 }
461 notified.await;
462 }
463 }
464
465 async fn wait_for_shutdown(&self) {
468 loop {
469 let notified = self.shutdown.notified();
470 if !self.inner.lock().accepting {
471 return;
472 }
473 notified.await;
474 }
475 }
476}
477
478pub(crate) struct SegmentOperationGuard {
482 active_operations: Arc<ActiveSegmentOperations>,
483 segment_ids: Vec<String>,
484 operation_token: u64,
485}
486
487impl Drop for SegmentOperationGuard {
488 fn drop(&mut self) {
489 let mut inner = self.active_operations.inner.lock();
490 for id in &self.segment_ids {
491 inner.segment_ids.remove(id);
492 }
493 inner.operation_tokens.remove(&self.operation_token);
494 inner.indexing_tokens.remove(&self.operation_token);
495 self.active_operations.idle.notify_waiters();
498 if inner.segment_ids.is_empty() {
499 debug_assert!(inner.operation_tokens.is_empty());
500 }
501 }
502}
503
504struct VectorArtifactUpdateLease {
511 updating: Arc<AtomicBool>,
512 active_operations: Arc<ActiveSegmentOperations>,
513}
514
515impl Drop for VectorArtifactUpdateLease {
516 fn drop(&mut self) {
517 self.updating.store(false, Ordering::Release);
518 self.active_operations.resume_non_indexing();
519 }
520}
521
522#[derive(Clone)]
523pub(crate) struct VectorArtifactUpdateGuard {
524 _lease: Arc<VectorArtifactUpdateLease>,
525}
526
527static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
531
532const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
533const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
534
535#[derive(Default)]
536struct MergeRetryState {
537 retry_after: Option<std::time::Instant>,
538 consecutive_failures: u32,
539}
540
541fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
542 let shift = consecutive_failures.saturating_sub(1).min(16);
543 MERGE_RETRY_BASE_DELAY
544 .checked_mul(1u32 << shift)
545 .unwrap_or(MERGE_RETRY_MAX_DELAY)
546 .min(MERGE_RETRY_MAX_DELAY)
547}
548
549struct DrainedMergeHandles<'a> {
558 shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
559 drained: Vec<JoinHandle<()>>,
560}
561
562impl<'a> DrainedMergeHandles<'a> {
563 fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
564 let drained = std::mem::take(&mut *shared.lock());
565 Self { shared, drained }
566 }
567
568 fn is_empty(&self) -> bool {
569 self.drained.is_empty()
570 }
571
572 async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
576 let handle = self.drained.last_mut()?;
577 let result = handle.await;
578 self.drained.pop();
579 Some(result)
580 }
581}
582
583impl Drop for DrainedMergeHandles<'_> {
584 fn drop(&mut self) {
585 if !self.drained.is_empty() {
586 self.shared.lock().append(&mut self.drained);
587 }
588 }
589}
590
591fn try_spawn_lifecycle<F>(
598 handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
599 runtime: &tokio::runtime::Handle,
600 future: F,
601) -> bool
602where
603 F: std::future::Future<Output = ()> + Send + 'static,
604{
605 std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
606 let mut handles = handles.lock();
607 handles.retain(|handle| !handle.is_finished());
608 handles.push(runtime.spawn(future));
609 }))
610 .is_ok()
611}
612
613struct OutputCleanupGuard {
621 segment_id: SegmentId,
622 cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
623}
624
625impl OutputCleanupGuard {
626 fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
627 Self {
628 segment_id,
629 cleanup: Some(cleanup),
630 }
631 }
632
633 fn disarm(&mut self) {
634 self.cleanup = None;
635 }
636}
637
638impl Drop for OutputCleanupGuard {
639 fn drop(&mut self) {
640 if let Some(cleanup) = self.cleanup.take() {
641 cleanup(self.segment_id);
642 }
643 }
644}
645
646struct ManagerState {
648 metadata: IndexMetadata,
649 merge_policy: Box<dyn MergePolicy>,
650}
651
652type ReplacementRefresh = Arc<
653 dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>
654 + Send
655 + Sync,
656>;
657
658async fn refresh_replacement_topology(refresh: Option<ReplacementRefresh>, index_label: &str) {
662 let Some(refresh) = refresh else {
663 return;
664 };
665 let mut last_error = None;
666 for attempt in 0..3 {
667 match refresh().await {
668 Ok(()) => return,
669 Err(error) => {
670 last_error = Some(error);
671 if attempt < 2 {
672 tokio::time::sleep(std::time::Duration::from_secs(1 << attempt)).await;
673 }
674 }
675 }
676 }
677 if let Some(error) = last_error {
678 log::warn!(
679 "[segment_lifecycle] index={index_label} replacement topology refresh failed after 3 attempts: {}",
680 error,
681 );
682 }
683}
684
685#[cfg(feature = "native")]
686struct MergeTaskError {
687 error: Error,
688 unavailable_segments: Vec<String>,
689}
690
691#[cfg(feature = "native")]
692impl MergeTaskError {
693 fn source(segment_id: String, error: Error) -> Self {
694 Self {
695 error,
696 unavailable_segments: vec![segment_id],
697 }
698 }
699
700 fn sources(segment_ids: Vec<String>, error: Error) -> Self {
701 Self {
702 error,
703 unavailable_segments: segment_ids,
704 }
705 }
706}
707
708#[cfg(feature = "native")]
709impl From<Error> for MergeTaskError {
710 fn from(error: Error) -> Self {
711 Self {
712 error,
713 unavailable_segments: Vec::new(),
714 }
715 }
716}
717
718#[cfg(feature = "native")]
719fn is_deterministic_source_error(error: &Error) -> bool {
720 matches!(error, Error::Corruption(_) | Error::Serialization(_))
721 || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
722}
723
724#[cfg(feature = "native")]
725fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
726 if is_deterministic_source_error(&error) {
727 MergeTaskError::source(segment_id, error)
728 } else {
729 MergeTaskError::from(error)
733 }
734}
735
736#[cfg(feature = "native")]
737type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
738
739#[derive(Clone, Copy)]
740enum ReplacementLayout {
741 BlockCopy,
745 BpReordered { converged: bool },
747 PreserveSingleSource,
750}
751
752fn replacement_bp_state(
753 parent_has_debt: bool,
754 parent_unconverged_passes: u32,
755 layout: ReplacementLayout,
756) -> (bool, bool, u32) {
757 match layout {
758 ReplacementLayout::BlockCopy => (
759 false,
760 !parent_has_debt,
761 if parent_has_debt {
762 parent_unconverged_passes
763 } else {
764 0
765 },
766 ),
767 ReplacementLayout::BpReordered { converged } => (
768 true,
769 converged,
770 if converged {
771 0
772 } else {
773 parent_unconverged_passes.saturating_add(1)
774 },
775 ),
776 ReplacementLayout::PreserveSingleSource => {
777 unreachable!("preserved layouts retain the complete source metadata")
778 }
779 }
780}
781
782#[derive(Clone, Copy, Debug, Eq, PartialEq)]
783enum VectorSegmentRewriteOutcome {
784 Rewritten,
785 AlreadyCurrent,
786 SourceGone,
787 Conflict,
788 Deferred,
789}
790
791pub(crate) struct StagedVectorSegment {
794 source_id: String,
795 output_id: SegmentId,
796 doc_count: u32,
797 _operation: SegmentOperationGuard,
798 cleanup: OutputCleanupGuard,
799}
800
801pub struct SegmentManager<D: DirectoryWriter + 'static> {
805 state: Arc<AsyncMutex<ManagerState>>,
807
808 active_operations: Arc<ActiveSegmentOperations>,
810
811 quarantined_segments: parking_lot::Mutex<HashSet<String>>,
816
817 merge_retry: parking_lot::Mutex<MergeRetryState>,
820
821 reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
825
826 merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
828
829 global_merge_wakeup_pending: AtomicBool,
833
834 force_merge_active: AtomicUsize,
839
840 #[cfg(test)]
843 force_merge_conflict_retries: std::sync::atomic::AtomicU64,
844
845 lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
849
850 published_generation: Arc<ArcSwap<PublishedIndexGeneration>>,
854
855 vector_artifact_update: Arc<AtomicBool>,
859
860 tracker: Arc<SegmentTracker>,
862
863 delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
865
866 directory: Arc<D>,
868 schema: Arc<crate::dsl::Schema>,
870 optimization: crate::structures::IndexOptimization,
872 posting_codec: crate::structures::PostingCodec,
874 term_cache_blocks: usize,
876 merge_permits: Arc<Semaphore>,
880 global_merge_permits: Arc<Semaphore>,
882 reorder_permits: Arc<ReorderConcurrencyGate>,
886 reorder_on_merge: bool,
891 merge_bp_time_budget: Option<std::time::Duration>,
895 bp_memory_budget_bytes: usize,
898 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
901 replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
905}
906
907struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
908
909impl Drop for ForceMergeActivityGuard<'_> {
910 fn drop(&mut self) {
911 self.0.fetch_sub(1, Ordering::AcqRel);
912 }
913}
914
915impl<D: DirectoryWriter + 'static> SegmentManager<D> {
916 #[allow(clippy::too_many_arguments)]
918 pub fn new(
919 directory: Arc<D>,
920 schema: Arc<crate::dsl::Schema>,
921 metadata: IndexMetadata,
922 merge_policy: Box<dyn MergePolicy>,
923 term_cache_blocks: usize,
924 max_concurrent_merges: usize,
925 global_merge_permits: Arc<Semaphore>,
926 merge_bp_time_budget: Option<std::time::Duration>,
927 bp_memory_budget_bytes: usize,
928 reorder_permits: Arc<ReorderConcurrencyGate>,
929 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
930 ) -> Self {
931 let reorder_on_merge = schema.reorder_on_merge();
934 if reorder_on_merge {
935 log::info!(
936 "[merge] index={} reorder-on-merge enabled by index schema",
937 schema.index_label()
938 );
939 }
940
941 let tracker = Arc::new(SegmentTracker::new());
942 for seg_id in metadata.segment_metas.keys() {
943 tracker.register(seg_id);
944 }
945
946 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
947 Arc::new(parking_lot::Mutex::new(Vec::new()));
948 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
949 let dir = Arc::clone(&directory);
950 let tracker = Arc::clone(&tracker);
951 let lifecycle_handles = Arc::clone(&lifecycle_handles);
952 let cleanup_index_label: Arc<str> = schema.index_label().into();
953 Arc::new(move |segment_ids| {
954 let Ok(handle) = tokio::runtime::Handle::try_current() else {
957 tracker.complete_deletion(&segment_ids);
960 return;
961 };
962 let dir = Arc::clone(&dir);
963 let task_tracker = Arc::clone(&tracker);
964 let task_index_label = Arc::clone(&cleanup_index_label);
965 let cleanup_ids = segment_ids.clone();
966 let future = async move {
967 for &segment_id in &segment_ids {
968 log::info!(
969 "[segment_cleanup] index={} deleting deferred segment {}",
970 task_index_label,
971 segment_id.to_hex()
972 );
973 if let Err(error) =
974 crate::segment::delete_segment(dir.as_ref(), segment_id).await
975 {
976 log::warn!(
977 "[segment_cleanup] index={} deferred delete failed for {}: {}",
978 task_index_label,
979 segment_id.to_hex(),
980 error,
981 );
982 }
983 }
984 task_tracker.complete_deletion(&segment_ids);
985 };
986 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
987 tracker.complete_deletion(&cleanup_ids);
991 log::warn!(
992 "[segment_cleanup] index={} runtime rejected deferred deletion; files will be swept later",
993 cleanup_index_label
994 );
995 }
996 })
997 };
998
999 let initial_generation = Arc::new(PublishedIndexGeneration {
1000 publication_id: metadata.publication_generation,
1001 schema: Arc::clone(&schema),
1002 trained_vectors: None,
1003 });
1004 Self {
1005 state: Arc::new(AsyncMutex::new(ManagerState {
1006 metadata,
1007 merge_policy,
1008 })),
1009 active_operations: Arc::new(ActiveSegmentOperations::new(schema.index_label().into())),
1010 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
1011 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
1012 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
1013 merge_handles: parking_lot::Mutex::new(Vec::new()),
1014 global_merge_wakeup_pending: AtomicBool::new(false),
1015 force_merge_active: AtomicUsize::new(0),
1016 #[cfg(test)]
1017 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
1018 lifecycle_handles,
1019 published_generation: Arc::new(ArcSwap::new(initial_generation)),
1020 vector_artifact_update: Arc::new(AtomicBool::new(false)),
1021 tracker,
1022 delete_fn,
1023 directory,
1024 schema,
1025 optimization: crate::structures::IndexOptimization::default(),
1026 posting_codec: crate::structures::PostingCodec::default(),
1027 term_cache_blocks,
1028 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
1029 global_merge_permits,
1030 reorder_permits,
1031 reorder_on_merge,
1032 merge_bp_time_budget,
1033 bp_memory_budget_bytes,
1034 background_reorder_pool,
1035 replacement_refresh: parking_lot::RwLock::new(None),
1036 }
1037 }
1038
1039 pub fn with_posting_config(
1042 mut self,
1043 optimization: crate::structures::IndexOptimization,
1044 posting_codec: crate::structures::PostingCodec,
1045 ) -> Self {
1046 self.optimization = optimization;
1047 self.posting_codec = posting_codec;
1048 self
1049 }
1050
1051 pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1052 where
1053 F: Fn() -> Fut + Send + Sync + 'static,
1054 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1055 {
1056 *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1057 }
1058
1059 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1065 if let Some(pool) = &self.background_reorder_pool {
1066 return Arc::clone(pool);
1067 }
1068 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1069 let threads = (num_cpus::get() / 2).max(1);
1070 log::info!(
1071 "[merge] process-wide background CPU pool: {} thread(s)",
1072 threads
1073 );
1074 Arc::new(
1075 rayon::ThreadPoolBuilder::new()
1076 .num_threads(threads)
1077 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
1078 .build()
1079 .expect("failed to build background CPU pool"),
1080 )
1081 }))
1082 }
1083
1084 pub fn begin_shutdown(&self) {
1088 self.active_operations.stop_accepting();
1089 }
1090
1091 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1099 where
1100 T: Send + 'static,
1101 F: std::future::Future<Output = Result<T>> + Send + 'static,
1102 {
1103 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1104 let future = async move {
1105 let result = transaction.await;
1106 let _ = result_tx.send(result);
1107 };
1108 let runtime = tokio::runtime::Handle::current();
1109 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1110 return Err(Error::Internal(
1111 "runtime rejected lifecycle metadata transaction".into(),
1112 ));
1113 }
1114 result_rx.await.map_err(|_| {
1115 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1116 })?
1117 }
1118
1119 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1121 let manager = Arc::clone(self);
1122 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1123 let Ok(handle) = tokio::runtime::Handle::try_current() else {
1124 log::warn!(
1125 "[segment_cleanup] index={} runtime unavailable; partial output {} will be swept on startup",
1126 manager.schema.index_label(),
1127 segment_id.to_hex(),
1128 );
1129 return;
1130 };
1131
1132 let cleanup_manager = Arc::clone(&manager);
1133 let future = async move {
1134 cleanup_manager
1135 .delete_output_if_unregistered(segment_id, "task unwind")
1136 .await;
1137 };
1138 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1139 log::warn!(
1140 "[segment_cleanup] index={} runtime rejected output cleanup; {} will be swept on startup",
1141 manager.schema.index_label(),
1142 segment_id.to_hex(),
1143 );
1144 }
1145 });
1146
1147 OutputCleanupGuard::new(output_id, cleanup)
1148 }
1149
1150 pub(crate) fn schedule_unpublished_segment_cleanup(
1155 self: &Arc<Self>,
1156 output_id: SegmentId,
1157 operation: SegmentOperationGuard,
1158 runtime: tokio::runtime::Handle,
1159 ) {
1160 let manager = Arc::clone(self);
1161 let output_hex = output_id.to_hex();
1162 let future = async move {
1163 manager
1164 .delete_output_if_unregistered(output_id, "indexing abort or failure")
1165 .await;
1166 drop(operation);
1167 };
1168 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1169 log::warn!(
1172 "[segment_cleanup] index={} runtime unavailable; indexing output {} will be swept on startup",
1173 self.schema.index_label(),
1174 output_hex,
1175 );
1176 }
1177 }
1178
1179 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1185 match self
1186 .active_operations
1187 .try_register_indexing(vec![segment_id.clone()])
1188 {
1189 Some(operation) => Ok(operation),
1190 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1191 None => Err(Error::Corruption(format!(
1192 "new segment ID {} is already owned by an active operation",
1193 segment_id
1194 ))),
1195 }
1196 }
1197
1198 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1202 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1203 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1204 })?;
1205 let files = SegmentFiles::new(id.0);
1206
1207 for path in files.mandatory_paths() {
1208 if !self.directory.exists(path).await.map_err(Error::Io)? {
1209 return Err(Error::Corruption(format!(
1210 "segment {} cannot be published: mandatory file {:?} is missing",
1211 segment_id, path
1212 )));
1213 }
1214 }
1215
1216 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1217 Error::Corruption(format!(
1218 "segment {} cannot be published: missing/unreadable {:?}: {}",
1219 segment_id, files.meta, e
1220 ))
1221 })?;
1222 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1223 Error::Corruption(format!(
1224 "segment {} cannot be published: failed reading {:?}: {}",
1225 segment_id, files.meta, e
1226 ))
1227 })?;
1228 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1229 Error::Corruption(format!(
1230 "segment {} cannot be published: invalid {:?}: {}",
1231 segment_id, files.meta, e
1232 ))
1233 })?;
1234
1235 if meta.id != id.0 || meta.num_docs != expected_docs {
1236 return Err(Error::Corruption(format!(
1237 "segment {} cannot be published: metadata identity/docs mismatch \
1238 (id={:032x}, docs={}, expected_docs={})",
1239 segment_id, meta.id, meta.num_docs, expected_docs
1240 )));
1241 }
1242
1243 Ok(())
1244 }
1245
1246 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1247 let inserted = self
1248 .quarantined_segments
1249 .lock()
1250 .insert(segment_id.to_string());
1251 if inserted {
1252 log::error!(
1253 "[merge] index={} quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1254 It remains metadata-live for explicit repair but is excluded from merges until restart",
1255 self.schema.index_label(),
1256 segment_id,
1257 error,
1258 );
1259 }
1260 }
1261
1262 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1263 let mut retry = self.merge_retry.lock();
1264 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1265 let delay = merge_retry_delay(retry.consecutive_failures);
1266 retry.retry_after = std::time::Instant::now().checked_add(delay);
1267 log::warn!(
1268 "[merge] index={} pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1269 self.schema.index_label(),
1270 delay.as_secs_f64(),
1271 retry.consecutive_failures,
1272 error,
1273 );
1274 delay
1275 }
1276
1277 fn clear_merge_retry_backoff(&self) {
1278 *self.merge_retry.lock() = MergeRetryState::default();
1279 }
1280
1281 fn merge_retry_is_paused(&self) -> bool {
1282 let mut retry = self.merge_retry.lock();
1283 match retry.retry_after {
1284 Some(deadline) if deadline > std::time::Instant::now() => true,
1285 Some(_) => {
1286 retry.retry_after = None;
1287 false
1288 }
1289 None => false,
1290 }
1291 }
1292
1293 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1294 let mut retries = self.reorder_retries.lock();
1295 let retry = retries.entry(segment_id.to_string()).or_default();
1296 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1297 let delay = merge_retry_delay(retry.consecutive_failures);
1298 retry.retry_after = std::time::Instant::now().checked_add(delay);
1299 log::warn!(
1300 "[reorder] index={} pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1301 self.schema.index_label(),
1302 segment_id,
1303 delay.as_secs_f64(),
1304 retry.consecutive_failures,
1305 error,
1306 );
1307 }
1308
1309 fn clear_reorder_retry(&self, segment_id: &str) {
1310 self.reorder_retries.lock().remove(segment_id);
1311 }
1312
1313 fn paused_reorder_segments(&self) -> HashSet<String> {
1314 let now = std::time::Instant::now();
1315 let mut retries = self.reorder_retries.lock();
1316 let mut paused = HashSet::new();
1317 for (segment_id, retry) in retries.iter_mut() {
1318 match retry.retry_after {
1319 Some(deadline) if deadline > now => {
1320 paused.insert(segment_id.clone());
1321 }
1322 Some(_) => retry.retry_after = None,
1323 None => {}
1324 }
1325 }
1326 paused
1327 }
1328
1329 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1333 if self
1334 .global_merge_wakeup_pending
1335 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1336 .is_err()
1337 {
1338 return;
1339 }
1340
1341 let manager = Arc::clone(self);
1342 let future = async move {
1343 let capacity = tokio::select! {
1344 biased;
1345 () = manager.active_operations.wait_for_shutdown() => None,
1346 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1347 };
1348
1349 manager
1350 .global_merge_wakeup_pending
1351 .store(false, Ordering::Release);
1352 if let Some(permit) = capacity {
1353 drop(permit);
1357 manager.maybe_merge().await;
1358 }
1359 };
1360 let runtime = tokio::runtime::Handle::current();
1361 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1362 self.global_merge_wakeup_pending
1363 .store(false, Ordering::Release);
1364 log::warn!(
1365 "[merge] index={} runtime rejected global-capacity wakeup task",
1366 self.schema.index_label()
1367 );
1368 }
1369 }
1370
1371 #[cfg(test)]
1372 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1373 self.quarantined_segments.lock().contains(segment_id)
1374 }
1375
1376 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1381 let output_hex = output_id.to_hex();
1382 {
1383 let st = self.state.lock().await;
1384 if st.metadata.has_segment(&output_hex) {
1385 return;
1386 }
1387 }
1388
1389 log::info!(
1393 "[segment_cleanup] index={} deleting uncommitted output {} after {}",
1394 self.schema.index_label(),
1395 output_hex,
1396 reason,
1397 );
1398 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1399 {
1400 log::warn!(
1401 "[segment_cleanup] index={} failed deleting uncommitted output {}: {}",
1402 self.schema.index_label(),
1403 output_hex,
1404 error,
1405 );
1406 }
1407 }
1408
1409 pub async fn get_segment_ids(&self) -> Vec<String> {
1415 self.state.lock().await.metadata.segment_ids()
1416 }
1417
1418 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1420 self.published_generation.load().trained_vectors.clone()
1421 }
1422
1423 pub(crate) fn published_generation(&self) -> Arc<PublishedIndexGeneration> {
1426 self.published_generation.load_full()
1427 }
1428
1429 pub(crate) fn publication_id(&self) -> u64 {
1430 self.published_generation.load().publication_id
1431 }
1432
1433 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1440 if self.vector_artifact_update.load(Ordering::Acquire) {
1441 return None;
1442 }
1443 let trained = self.published_generation.load().trained_vectors.clone();
1444 if self.vector_artifact_update.load(Ordering::Acquire) {
1445 None
1446 } else {
1447 trained
1448 }
1449 }
1450
1451 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1470 self.vector_artifact_update
1471 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1472 .map_err(|_| {
1473 Error::Internal("a trained-vector artifact update is already in progress".into())
1474 })?;
1475 self.active_operations.pause_non_indexing();
1476 let guard = VectorArtifactUpdateGuard {
1477 _lease: Arc::new(VectorArtifactUpdateLease {
1478 updating: Arc::clone(&self.vector_artifact_update),
1479 active_operations: Arc::clone(&self.active_operations),
1480 }),
1481 };
1482 let (preexisting, parked_indexing) =
1483 self.active_operations.draining_operation_tokens_snapshot();
1484 if parked_indexing > 0 {
1485 return Err(Error::Internal(format!(
1486 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1487 segment(s) are built but uncommitted; commit or abort the pending \
1488 generation and retry"
1489 )));
1490 }
1491 self.active_operations
1492 .wait_until_operations_finish(&preexisting)
1493 .await;
1494 Ok(guard)
1495 }
1496
1497 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1500 let (vector_fields, schema, publication_id) = {
1502 let st = self.state.lock().await;
1503 (
1504 st.metadata.vector_fields.clone(),
1505 Arc::new(st.metadata.schema.clone()),
1506 st.metadata.publication_generation,
1507 )
1508 };
1509 let trained = IndexMetadata::try_load_trained_from_fields(
1511 &vector_fields,
1512 schema.as_ref(),
1513 self.directory.as_ref(),
1514 )
1515 .await?
1516 .map(Arc::new);
1517 self.published_generation
1521 .store(Arc::new(PublishedIndexGeneration {
1522 publication_id,
1523 schema,
1524 trained_vectors: trained,
1525 }));
1526 Ok(())
1527 }
1528
1529 pub(crate) async fn publish_vector_generation(
1536 self: &Arc<Self>,
1537 artifact_update: &VectorArtifactUpdateGuard,
1538 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1539 next_trained: Arc<TrainedVectorStructures>,
1540 staged: Vec<StagedVectorSegment>,
1541 ) -> Result<()> {
1542 let schema = self.published_generation().schema.clone();
1543 self.publish_vector_generation_with_schema(
1544 artifact_update,
1545 schema,
1546 vector_fields,
1547 Some(next_trained),
1548 staged,
1549 )
1550 .await
1551 }
1552
1553 pub(crate) async fn publish_vector_generation_with_schema(
1554 self: &Arc<Self>,
1555 artifact_update: &VectorArtifactUpdateGuard,
1556 schema: Arc<crate::dsl::Schema>,
1557 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1558 next_trained: Option<Arc<TrainedVectorStructures>>,
1559 mut staged: Vec<StagedVectorSegment>,
1560 ) -> Result<()> {
1561 if !self.vector_artifact_update.load(Ordering::Acquire) {
1562 return Err(Error::Internal(
1563 "vector generation publication lost its exclusive update lease".into(),
1564 ));
1565 }
1566
1567 for replacement in &staged {
1568 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1569 .await?;
1570 }
1571
1572 let mut st = Arc::clone(&self.state).lock_owned().await;
1573 let mut next = st.metadata.clone();
1574 next.schema = (*schema).clone();
1575 next.vector_fields = vector_fields;
1576 next.refresh_total_vectors();
1577
1578 for replacement in &staged {
1579 let source_info = next
1580 .segment_metas
1581 .remove(&replacement.source_id)
1582 .ok_or_else(|| {
1583 Error::Corruption(format!(
1584 "vector generation source {} disappeared before publication",
1585 replacement.source_id,
1586 ))
1587 })?;
1588 let output_hex = replacement.output_id.to_hex();
1589 if next.segment_metas.contains_key(&output_hex) {
1590 return Err(Error::Corruption(format!(
1591 "vector generation output {output_hex} is already metadata-live"
1592 )));
1593 }
1594 next.add_segment_meta(output_hex, source_info);
1597 }
1598
1599 let directory = Arc::clone(&self.directory);
1600 let published_generation = Arc::clone(&self.published_generation);
1601 let tracker = Arc::clone(&self.tracker);
1602 let replacement_refresh = self.replacement_refresh.read().clone();
1603 let artifact_update = artifact_update.clone();
1607 let index_label = self.schema.index_label().to_owned();
1608 next.publication_generation =
1609 next.publication_generation.checked_add(1).ok_or_else(|| {
1610 Error::Corruption("vector publication generation exhausted u64".into())
1611 })?;
1612 let next_schema = schema;
1613 let next_publication_id = next.publication_generation;
1614 self.run_lifecycle_transaction(async move {
1615 let _artifact_update = artifact_update;
1616 next.save(directory.as_ref()).await?;
1617
1618 for replacement in &staged {
1619 tracker.register(&replacement.output_id.to_hex());
1620 }
1621 st.metadata = next;
1622 published_generation.store(Arc::new(PublishedIndexGeneration {
1623 publication_id: next_publication_id,
1624 schema: next_schema,
1625 trained_vectors: next_trained,
1626 }));
1627
1628 for replacement in &mut staged {
1631 replacement.cleanup.disarm();
1632 }
1633 let retired = staged
1634 .iter()
1635 .map(|replacement| replacement.source_id.clone())
1636 .collect::<Vec<_>>();
1637 let ready_to_delete = tracker.mark_for_deletion(&retired);
1638 drop(st);
1639 for &segment_id in &ready_to_delete {
1640 if let Err(error) =
1641 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1642 {
1643 log::warn!(
1644 "[segment_cleanup] index={index_label} immediate dense-vector generation delete failed for {}: {}",
1645 segment_id.to_hex(),
1646 error,
1647 );
1648 }
1649 }
1650 tracker.complete_deletion(&ready_to_delete);
1651 refresh_replacement_topology(replacement_refresh, &index_label).await;
1652 Ok(())
1653 })
1654 .await
1655 }
1656
1657 pub(crate) async fn publish_vector_schema_only(
1661 self: &Arc<Self>,
1662 artifact_update: &VectorArtifactUpdateGuard,
1663 schema: Arc<crate::dsl::Schema>,
1664 ) -> Result<()> {
1665 let vector_fields = self
1666 .read_metadata(|metadata| metadata.vector_fields.clone())
1667 .await;
1668 let trained = self.published_generation().trained_vectors.clone();
1669 self.publish_vector_generation_with_schema(
1670 artifact_update,
1671 schema,
1672 vector_fields,
1673 trained,
1674 Vec::new(),
1675 )
1676 .await
1677 }
1678
1679 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1681 where
1682 F: FnOnce(&IndexMetadata) -> R,
1683 {
1684 let st = self.state.lock().await;
1685 f(&st.metadata)
1686 }
1687
1688 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1690 where
1691 F: FnOnce(&mut IndexMetadata),
1692 {
1693 let mut st = Arc::clone(&self.state).lock_owned().await;
1694 let mut next = st.metadata.clone();
1695 f(&mut next);
1696 let directory = Arc::clone(&self.directory);
1697 self.run_lifecycle_transaction(async move {
1698 next.save(directory.as_ref()).await?;
1699 st.metadata = next;
1700 Ok(())
1701 })
1702 .await
1703 }
1704
1705 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1708 let (acquired, generation) = {
1709 let st = self.state.lock().await;
1710 let segment_ids = st.metadata.segment_ids();
1711 (
1712 self.tracker.acquire(&segment_ids),
1713 self.published_generation.load_full(),
1714 )
1715 };
1716
1717 SegmentSnapshot::with_generation(
1718 Arc::clone(&self.tracker),
1719 acquired,
1720 generation,
1721 Arc::clone(&self.delete_fn),
1722 )
1723 }
1724
1725 pub fn tracker(&self) -> Arc<SegmentTracker> {
1727 Arc::clone(&self.tracker)
1728 }
1729
1730 pub fn directory(&self) -> Arc<D> {
1732 Arc::clone(&self.directory)
1733 }
1734}
1735
1736#[cfg(feature = "native")]
1741impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1742 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1744 for (segment_id, num_docs) in new_segments {
1747 self.validate_completed_segment(segment_id, *num_docs)
1748 .await?;
1749 }
1750
1751 let mut st = Arc::clone(&self.state).lock_owned().await;
1752 let mut next = st.metadata.clone();
1753 let mut added = Vec::new();
1754 for (segment_id, num_docs) in new_segments {
1755 if !next.has_segment(segment_id) {
1756 next.add_segment(segment_id.clone(), *num_docs);
1757 added.push(segment_id.clone());
1758 }
1759 }
1760
1761 let directory = Arc::clone(&self.directory);
1767 let tracker = Arc::clone(&self.tracker);
1768 self.run_lifecycle_transaction(async move {
1769 next.save(directory.as_ref()).await?;
1770 for segment_id in &added {
1771 tracker.register(segment_id);
1772 }
1773 st.metadata = next;
1774 Ok(())
1775 })
1776 .await
1777 }
1778
1779 pub async fn maybe_merge(self: &Arc<Self>) {
1790 if !self.active_operations.is_accepting() {
1791 log::debug!(
1792 "[maybe_merge] index={} manager is shutting down, skipping",
1793 self.schema.index_label()
1794 );
1795 return;
1796 }
1797 if self.merge_retry_is_paused() {
1798 log::debug!(
1799 "[maybe_merge] index={} retry backoff active, skipping",
1800 self.schema.index_label()
1801 );
1802 return;
1803 }
1804
1805 {
1808 let mut handles = self.merge_handles.lock();
1809 handles.retain(|h| !h.is_finished());
1810 }
1811 let local_slots = self.merge_permits.available_permits();
1812 let global_slots = self.global_merge_permits.available_permits();
1813 let slots_available = local_slots.min(global_slots);
1814
1815 {
1819 let st = self.state.lock().await;
1820 let quarantined = self.quarantined_segments.lock().clone();
1821 let active_ids = self.active_operations.snapshot();
1822
1823 let live_segments: Vec<SegmentInfo> = st
1828 .metadata
1829 .segment_metas
1830 .iter()
1831 .filter(|(id, _)| {
1832 !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1833 })
1834 .map(|(id, info)| SegmentInfo {
1835 id: id.clone(),
1836 num_docs: info.num_docs,
1837 })
1838 .collect();
1839 let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1840
1841 let segments: Vec<SegmentInfo> = live_segments
1844 .iter()
1845 .filter(|segment| !active_ids.contains(&segment.id))
1846 .cloned()
1847 .collect();
1848
1849 log::debug!(
1850 "[maybe_merge] index={} {} eligible segments",
1851 self.schema.index_label(),
1852 segments.len()
1853 );
1854
1855 let candidates = st.merge_policy.find_merges(&segments);
1856
1857 if candidates.is_empty() {
1858 return;
1859 }
1860
1861 if slots_available == 0 {
1865 if local_slots > 0 && global_slots == 0 {
1866 self.schedule_global_merge_wakeup();
1867 }
1868 log::debug!(
1869 "[maybe_merge] index={} at max concurrent merges, skipping",
1870 self.schema.index_label()
1871 );
1872 return;
1873 }
1874
1875 log::debug!(
1876 "[maybe_merge] index={} {} merge candidates, {} slots available",
1877 self.schema.index_label(),
1878 candidates.len(),
1879 slots_available
1880 );
1881
1882 let mut handles = Vec::new();
1883 for c in candidates {
1884 if handles.len() >= slots_available {
1885 break;
1886 }
1887 let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1893 if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1894 handles.push(h);
1895 }
1896 }
1897 if !handles.is_empty() {
1898 if severe_backlog && self.reorder_on_merge {
1899 log::info!(
1900 "[maybe_merge] index={} severe backlog: {} live segments; started {} fast \
1901 block-copy merge(s), deferring BP to the optimizer",
1902 self.schema.index_label(),
1903 live_segments.len(),
1904 handles.len(),
1905 );
1906 }
1907 self.merge_handles.lock().extend(handles);
1912 }
1913 }
1914 }
1915
1916 fn spawn_merge(
1925 self: &Arc<Self>,
1926 segment_ids_to_merge: Vec<String>,
1927 reorder_bmp: bool,
1928 ) -> Option<JoinHandle<()>> {
1929 if self.force_merge_active.load(Ordering::Acquire) > 0 {
1930 log::debug!(
1931 "[spawn_merge] index={} skipped: explicit force merge has priority",
1932 self.schema.index_label()
1933 );
1934 return None;
1935 }
1936 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1937 Ok(permit) => permit,
1938 Err(_) => {
1939 log::debug!(
1940 "[spawn_merge] index={} skipped: global merge capacity is full",
1941 self.schema.index_label()
1942 );
1943 self.schedule_global_merge_wakeup();
1944 return None;
1945 }
1946 };
1947 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1948 Ok(permit) => permit,
1949 Err(_) => {
1950 log::debug!(
1951 "[spawn_merge] index={} skipped: no merge permit available",
1952 self.schema.index_label()
1953 );
1954 return None;
1955 }
1956 };
1957 let output_id = SegmentId::new();
1958 let output_hex = output_id.to_hex();
1959
1960 let mut all_ids = segment_ids_to_merge.clone();
1961 all_ids.push(output_hex);
1962
1963 let guard = match self.active_operations.try_register(all_ids) {
1964 Some(g) => g,
1965 None => {
1966 log::debug!(
1967 "[spawn_merge] index={} skipped: segments overlap with an active operation",
1968 self.schema.index_label()
1969 );
1970 return None;
1971 }
1972 };
1973
1974 let sm = Arc::clone(self);
1975 let ids = segment_ids_to_merge;
1976
1977 let index_label = self.schema.index_label().to_owned();
1978 Some(tokio::spawn(async move {
1979 let mut reevaluate = false;
1980 let mut retry_delay = None;
1981
1982 let result = sm
1983 .merge_and_replace_registered(
1984 &ids,
1985 output_id,
1986 reorder_bmp,
1987 ReorderPriority::AutomaticMerge,
1988 )
1989 .await;
1990
1991 match result {
1992 Ok(_) => {
1993 sm.clear_merge_retry_backoff();
1994 reevaluate = true;
1995 }
1996 Err(MergeTaskError {
1997 error: Error::IndexClosed,
1998 ..
1999 }) => {
2000 log::debug!(
2001 "[merge] index={index_label} background merge for segments {:?} cancelled during shutdown",
2002 ids,
2003 );
2004 }
2005 Err(MergeTaskError {
2006 error,
2007 unavailable_segments,
2008 }) => {
2009 log::error!(
2010 "[merge] index={index_label} background merge failed for segments {:?}: {}",
2011 ids,
2012 error
2013 );
2014 if !unavailable_segments.is_empty() {
2015 reevaluate = true;
2019 } else {
2020 retry_delay = Some(sm.pause_merge_retries(&error));
2021 }
2022 }
2023 }
2024 drop(guard);
2027 drop(merge_permit);
2029 drop(global_merge_permit);
2030
2031 if reevaluate {
2032 sm.maybe_merge().await;
2033 } else if let Some(retry_delay) = retry_delay {
2034 sm.schedule_merge_retry_wakeup(retry_delay);
2041 }
2042 }))
2043 }
2044
2045 async fn merge_and_replace_registered(
2052 self: &Arc<Self>,
2053 ids: &[String],
2054 output_id: SegmentId,
2055 reorder_bmp: bool,
2056 priority: ReorderPriority,
2057 ) -> MergeTaskResult<(String, u32, bool)> {
2058 let mut output_cleanup = self.output_cleanup_guard(output_id);
2059 let generation = self.published_generation();
2060 let trained = self.trained_for_segment_build();
2061 let granularity = if reorder_bmp {
2062 self.merge_granularity(ids).await
2063 } else {
2064 crate::segment::reorder::BpGranularity::Auto
2065 };
2066 let result = Self::do_merge(
2067 self.directory.as_ref(),
2068 &generation.schema,
2069 ids,
2070 output_id,
2071 self.term_cache_blocks,
2072 self.optimization,
2073 self.posting_codec,
2074 trained.as_deref(),
2075 reorder_bmp,
2076 granularity,
2077 self.merge_bp_time_budget,
2078 self.bp_memory_budget_bytes,
2079 Arc::clone(&self.reorder_permits),
2080 priority,
2081 self.active_operations.cancellation_flag(),
2082 Some(self.background_cpu_pool()),
2083 )
2084 .await;
2085
2086 let (new_id, doc_count, bp_converged) = match result {
2087 Ok(value) => value,
2088 Err(error) => {
2089 for segment_id in &error.unavailable_segments {
2090 self.quarantine_segment(segment_id, &error.error);
2091 }
2092 self.delete_output_if_unregistered(output_id, "merge failure")
2093 .await;
2094 output_cleanup.disarm();
2095 return Err(error);
2096 }
2097 };
2098
2099 let layout = if reorder_bmp {
2100 ReplacementLayout::BpReordered {
2101 converged: bp_converged,
2102 }
2103 } else {
2104 ReplacementLayout::BlockCopy
2105 };
2106 if let Err(error) = self
2107 .replace_segments(ids, new_id.clone(), doc_count, layout)
2108 .await
2109 {
2110 self.delete_output_if_unregistered(output_id, "replacement failure")
2111 .await;
2112 output_cleanup.disarm();
2113 return Err(MergeTaskError::from(error));
2114 }
2115 output_cleanup.disarm();
2116 Ok((new_id, doc_count, bp_converged))
2117 }
2118
2119 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
2123 let manager = Arc::clone(self);
2124 let future = async move {
2125 tokio::select! {
2126 () = tokio::time::sleep(retry_delay) => {
2127 manager.maybe_merge().await;
2128 }
2129 () = manager.active_operations.wait_for_shutdown() => {}
2130 }
2131 };
2132 let runtime = tokio::runtime::Handle::current();
2133 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2134 log::warn!(
2135 "[merge] index={} runtime rejected merge-retry wakeup task; eligible segments may stay \
2136 unmerged until the next commit re-runs merge policy evaluation",
2137 self.schema.index_label()
2138 );
2139 }
2140 }
2141
2142 async fn replace_segments(
2146 self: &Arc<Self>,
2147 old_ids: &[String],
2148 new_id: String,
2149 doc_count: u32,
2150 layout: ReplacementLayout,
2151 ) -> Result<()> {
2152 self.validate_completed_segment(&new_id, doc_count).await?;
2155 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
2156 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
2157 })?;
2158 let output_reader = SegmentReader::open(
2159 self.directory.as_ref(),
2160 output_id,
2161 self.published_generation().schema.clone(),
2162 self.term_cache_blocks,
2163 )
2164 .await
2165 .map_err(|error| match error {
2166 Error::Io(_) | Error::IndexClosed => error,
2170 error => Error::Corruption(format!(
2171 "replacement segment {new_id} failed full reader validation: {error}"
2172 )),
2173 })?;
2174 if output_reader.num_docs() != doc_count {
2175 return Err(Error::Corruption(format!(
2176 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2177 output_reader.num_docs(),
2178 )));
2179 }
2180 drop(output_reader);
2181
2182 let mut st = Arc::clone(&self.state).lock_owned().await;
2183 let missing: Vec<&String> = old_ids
2187 .iter()
2188 .filter(|id| !st.metadata.has_segment(id))
2189 .collect();
2190 if !missing.is_empty() {
2191 return Err(Error::Corruption(format!(
2192 "replace_segments: source segment(s) {:?} not in metadata — \
2193 refusing to add output {} (would duplicate documents)",
2194 missing, new_id
2195 )));
2196 }
2197
2198 let replacement_info = match layout {
2199 ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2200 let generation = old_ids
2201 .iter()
2202 .filter_map(|id| st.metadata.segment_metas.get(id))
2203 .map(|info| info.generation)
2204 .max()
2205 .unwrap_or(0)
2206 .checked_add(1)
2207 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2208 let parent_unconverged_passes = old_ids
2209 .iter()
2210 .filter_map(|id| st.metadata.segment_metas.get(id))
2211 .map(|info| info.bp_unconverged_passes)
2212 .max()
2213 .unwrap_or(0);
2214 let parent_has_debt = old_ids
2215 .iter()
2216 .filter_map(|id| st.metadata.segment_metas.get(id))
2217 .any(|info| !info.bp_converged);
2218 let (reordered, bp_converged, bp_unconverged_passes) =
2219 replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2220 SegmentMetaInfo {
2221 num_docs: doc_count,
2222 ancestors: old_ids.to_vec(),
2223 generation,
2224 reordered,
2225 bp_converged,
2226 bp_unconverged_passes,
2227 }
2228 }
2229 ReplacementLayout::PreserveSingleSource => {
2230 let [source_id] = old_ids else {
2231 return Err(Error::Internal(
2232 "layout-preserving replacement requires exactly one source".into(),
2233 ));
2234 };
2235 let mut source = st
2236 .metadata
2237 .segment_metas
2238 .get(source_id)
2239 .cloned()
2240 .ok_or_else(|| {
2241 Error::Corruption(format!(
2242 "layout-preserving replacement source {source_id} disappeared"
2243 ))
2244 })?;
2245 source.num_docs = doc_count;
2246 source
2247 }
2248 };
2249 let retired_ids = old_ids.to_vec();
2250 let mut next = st.metadata.clone();
2251 for id in old_ids {
2252 next.remove_segment(id);
2253 }
2254 next.add_segment_meta(new_id.clone(), replacement_info);
2255
2256 let directory = Arc::clone(&self.directory);
2257 let tracker = Arc::clone(&self.tracker);
2258 let replacement_refresh = self.replacement_refresh.read().clone();
2259 let index_label = self.schema.index_label().to_owned();
2260 self.run_lifecycle_transaction(async move {
2261 next.save(directory.as_ref()).await?;
2264 tracker.register(&new_id);
2265 st.metadata = next;
2266
2267 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2271 drop(st);
2272 for &segment_id in &ready_to_delete {
2273 if let Err(error) =
2274 crate::segment::delete_segment(directory.as_ref(), segment_id).await
2275 {
2276 log::warn!(
2277 "[segment_cleanup] index={index_label} immediate delete failed for {}: {}",
2278 segment_id.to_hex(),
2279 error,
2280 );
2281 }
2282 }
2283 tracker.complete_deletion(&ready_to_delete);
2284 refresh_replacement_topology(replacement_refresh, &index_label).await;
2285 Ok(())
2286 })
2287 .await
2288 }
2289
2290 #[allow(clippy::too_many_arguments)]
2295 async fn do_merge(
2296 directory: &D,
2297 schema: &Arc<crate::dsl::Schema>,
2298 segment_ids_to_merge: &[String],
2299 output_segment_id: SegmentId,
2300 term_cache_blocks: usize,
2301 optimization: crate::structures::IndexOptimization,
2302 posting_codec: crate::structures::PostingCodec,
2303 trained: Option<&TrainedVectorStructures>,
2304 reorder_bmp: bool,
2305 granularity: crate::segment::reorder::BpGranularity,
2306 merge_bp_time_budget: Option<std::time::Duration>,
2307 bp_memory_budget_bytes: usize,
2308 reorder_permits: Arc<ReorderConcurrencyGate>,
2309 reorder_priority: ReorderPriority,
2310 cancellation: Arc<AtomicBool>,
2311 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2312 ) -> MergeTaskResult<(String, u32, bool)> {
2313 let output_hex = output_segment_id.to_hex();
2314 let load_start = std::time::Instant::now();
2315
2316 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2317 for id_str in segment_ids_to_merge {
2318 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2319 MergeTaskError::source(
2320 id_str.clone(),
2321 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2322 )
2323 })?;
2324 segment_ids.push(id);
2325 }
2326
2327 let mut unavailable_sources = Vec::new();
2332 let mut missing_files = Vec::new();
2333 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2334 let files = SegmentFiles::new(id.0);
2335 let mut source_unavailable = false;
2336 for path in files.mandatory_paths() {
2337 let exists = directory
2338 .exists(path)
2339 .await
2340 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2341 if !exists {
2342 source_unavailable = true;
2343 missing_files.push(format!("{}:{:?}", id_str, path));
2344 }
2345 }
2346 if source_unavailable {
2347 unavailable_sources.push(id_str.clone());
2348 }
2349 }
2350 if !unavailable_sources.is_empty() {
2351 return Err(MergeTaskError::sources(
2352 unavailable_sources,
2353 Error::Corruption(format!(
2354 "merge sources are missing mandatory files: {}",
2355 missing_files.join(", ")
2356 )),
2357 ));
2358 }
2359
2360 let schema_arc = Arc::clone(schema);
2361 let futures: Vec<_> = segment_ids
2362 .iter()
2363 .map(|&sid| {
2364 let sch = Arc::clone(&schema_arc);
2365 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
2366 })
2367 .collect();
2368
2369 let results = futures::future::join_all(futures).await;
2370 let mut readers = Vec::with_capacity(results.len());
2371 let mut total_docs = 0u64;
2372 for (i, result) in results.into_iter().enumerate() {
2373 match result {
2374 Ok(r) => {
2375 total_docs += r.meta().num_docs as u64;
2376 readers.push(r);
2377 }
2378 Err(e) => {
2379 log::error!(
2380 "[merge] index={} Failed to open segment {}: {:?}",
2381 schema.index_label(),
2382 segment_ids_to_merge[i],
2383 e
2384 );
2385 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2386 }
2387 }
2388 }
2389 if total_docs > u32::MAX as u64 {
2390 return Err(Error::Internal(format!(
2391 "Merged segment doc count ({}) exceeds u32::MAX",
2392 total_docs
2393 ))
2394 .into());
2395 }
2396
2397 for (i, reader) in readers.iter().enumerate() {
2401 let meta_docs = reader.meta().num_docs;
2402 let store_docs = reader.store().num_docs();
2403 if store_docs != meta_docs {
2404 return Err(MergeTaskError::source(
2405 segment_ids_to_merge[i].clone(),
2406 Error::Corruption(format!(
2407 "pre-merge validation: segment {} store has {} docs but meta says {}",
2408 segment_ids_to_merge[i], store_docs, meta_docs
2409 )),
2410 ));
2411 }
2412 }
2413
2414 log::info!(
2415 "[merge] index={} loaded {} segment readers in {:.1}s",
2416 schema.index_label(),
2417 readers.len(),
2418 load_start.elapsed().as_secs_f64()
2419 );
2420
2421 let merger = SegmentMerger::new(Arc::clone(schema))
2422 .with_posting_config(optimization, posting_codec)
2423 .with_bmp_reorder(reorder_bmp)
2424 .with_granularity(granularity)
2425 .with_bp_budget(crate::segment::BpBudget {
2426 min_partition_docs: None,
2427 time_budget: merge_bp_time_budget,
2428 })
2429 .with_cancellation(cancellation)
2430 .with_bp_memory_budget(bp_memory_budget_bytes)
2431 .with_reorder_permits(reorder_permits)
2432 .with_reorder_priority(reorder_priority)
2433 .with_background_pool(bg_cpu_pool);
2434
2435 log::info!(
2436 "[merge] index={} {} segments -> {} (trained={})",
2437 schema.index_label(),
2438 segment_ids_to_merge.len(),
2439 output_hex,
2440 trained.map_or(0, |t| t.centroids.len()),
2441 );
2442
2443 let (_merged_meta, merge_stats) = merger
2444 .merge(directory, &readers, output_segment_id, trained)
2445 .await
2446 .map_err(|error| {
2447 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2448 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2454 } else {
2455 MergeTaskError::from(error)
2456 }
2457 })?;
2458 let bp_converged = merge_stats.bp_converged;
2459 if !bp_converged {
2460 log::info!(
2461 "[merge] index={} merge-time BP hit its wall-clock budget — output marked unconverged; \
2462 the background optimizer deepens it later",
2463 schema.index_label(),
2464 );
2465 }
2466
2467 log::info!(
2468 "[merge] index={} total wall-clock: {:.1}s ({} segments, {} docs)",
2469 schema.index_label(),
2470 load_start.elapsed().as_secs_f64(),
2471 readers.len(),
2472 total_docs,
2473 );
2474
2475 Ok((output_hex, total_docs as u32, bp_converged))
2476 }
2477
2478 pub async fn abort_merges(&self) {
2488 loop {
2489 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2490 if handles.is_empty() {
2491 return;
2492 }
2493 while let Some(result) = handles.join_next().await {
2494 if let Err(error) = result
2495 && error.is_panic()
2496 {
2497 log::error!(
2498 "[merge] index={} background task panicked while draining: {}",
2499 self.schema.index_label(),
2500 error
2501 );
2502 }
2503 }
2504 }
2505 }
2506
2507 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2512 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2513 while handles.join_next().await.is_some() {}
2514 }
2515
2516 pub async fn wait_for_all_merges(self: &Arc<Self>) {
2525 loop {
2526 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2527 if handles.is_empty() {
2528 break;
2529 }
2530 while handles.join_next().await.is_some() {}
2531 }
2532 }
2533
2534 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2539 self.wait_for_all_merges().await;
2540 self.active_operations.wait_until_idle().await;
2541 loop {
2542 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2543 if handles.is_empty() {
2544 break;
2545 }
2546 for handle in handles {
2547 if let Err(error) = handle.await
2548 && error.is_panic()
2549 {
2550 log::error!(
2551 "[segment_cleanup] index={} task panicked while draining: {}",
2552 self.schema.index_label(),
2553 error
2554 );
2555 }
2556 }
2557 }
2558 }
2559
2560 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2572 self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2573 .await
2574 }
2575
2576 pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2586 self: &Arc<Self>,
2587 mut refresh_snapshots: F,
2588 ) -> Result<()>
2589 where
2590 F: FnMut() -> Fut,
2591 Fut: std::future::Future<Output = Result<()>>,
2592 {
2593 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2598 std::time::Duration::from_millis(100);
2599 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2603
2604 let (_force_merge_activity, policy_segment_docs) = {
2605 let st = self.state.lock().await;
2606 self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2610 (
2611 ForceMergeActivityGuard(&self.force_merge_active),
2612 st.merge_policy.max_segment_docs(),
2613 )
2614 };
2615
2616 let background_merges = self
2619 .merge_handles
2620 .lock()
2621 .iter()
2622 .filter(|handle| !handle.is_finished())
2623 .count();
2624 if background_merges > 0 {
2625 log::info!(
2626 "[force_merge] index={} waiting for {} in-flight background merge(s) before planning",
2627 self.schema.index_label(),
2628 background_merges,
2629 );
2630 }
2631 let drain_start = std::time::Instant::now();
2632 self.wait_for_all_merges().await;
2633 if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2634 log::info!(
2635 "[force_merge] index={} drained background merges in {:.1}s",
2636 self.schema.index_label(),
2637 drain_start.elapsed().as_secs_f64(),
2638 );
2639 }
2640
2641 let refresh_start = std::time::Instant::now();
2646 refresh_snapshots().await?;
2647 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2648 log::info!(
2649 "[force_merge] index={} initial snapshot refresh took {:.1}s",
2650 self.schema.index_label(),
2651 refresh_start.elapsed().as_secs_f64(),
2652 );
2653 }
2654
2655 let mut completed_outputs = HashSet::new();
2659 let mut logged_held_wait = false;
2661
2662 loop {
2663 if !self.active_operations.is_accepting() {
2664 return Err(Error::IndexClosed);
2665 }
2666
2667 let segments: Vec<(String, u32)> = {
2668 let st = self.state.lock().await;
2669 st.metadata
2670 .segment_metas
2671 .iter()
2672 .filter(|(id, _)| !completed_outputs.contains(*id))
2673 .map(|(id, info)| (id.clone(), info.num_docs))
2674 .collect()
2675 };
2676
2677 let active_ids = self.active_operations.snapshot();
2682 let held = segments
2683 .iter()
2684 .filter(|(id, _)| active_ids.contains(id))
2685 .count();
2686 let free_segments: Vec<_> = segments
2687 .into_iter()
2688 .filter(|(id, _)| !active_ids.contains(id))
2689 .collect();
2690 let max_docs = u64::from(u32::MAX);
2697 let planned_groups = plan_force_merge_groups(free_segments, max_docs);
2698 if let Some(cap) = policy_segment_docs {
2699 for group in planned_groups
2700 .iter()
2701 .filter(|group| group.segments.len() >= 2)
2702 .filter(|group| group.total_docs > u64::from(cap))
2703 {
2704 log::warn!(
2705 "[force_merge] index={} output of {} docs intentionally exceeds the \
2706 background merge policy cap of {} docs (force merge compacts to the \
2707 u32 format limit)",
2708 self.schema.index_label(),
2709 group.total_docs,
2710 cap,
2711 );
2712 }
2713 }
2714 let next_group = planned_groups
2715 .into_iter()
2716 .find(|group| group.segments.len() >= 2);
2717
2718 let Some(group) = next_group else {
2719 if held == 0 {
2720 if !completed_outputs.is_empty() {
2721 completed_outputs.clear();
2727 continue;
2728 }
2729 let remaining = {
2734 let st = self.state.lock().await;
2735 st.metadata.segment_metas.len()
2736 };
2737 if remaining > 1 {
2738 log::warn!(
2739 "[force_merge] index={} finished with {} segments: combined \
2740 document count exceeds the u32 segment format limit, so a \
2741 single output is impossible",
2742 self.schema.index_label(),
2743 remaining,
2744 );
2745 }
2746 refresh_snapshots().await?;
2751 return Ok(());
2752 }
2753 if !logged_held_wait {
2754 log::info!(
2755 "[force_merge] index={} waiting: {} segment(s) held by active \
2756 merge/reorder operations, no free group can merge",
2757 self.schema.index_label(),
2758 held
2759 );
2760 logged_held_wait = true;
2761 } else {
2762 log::debug!(
2763 "[force_merge] index={} still waiting on {} held segment(s)",
2764 self.schema.index_label(),
2765 held
2766 );
2767 }
2768 #[cfg(test)]
2769 self.force_merge_conflict_retries
2770 .fetch_add(1, Ordering::Relaxed);
2771 tokio::select! {
2772 biased;
2773 () = self.active_operations.wait_for_shutdown() => {
2774 return Err(Error::IndexClosed);
2775 }
2776 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2777 }
2778 continue;
2779 };
2780 logged_held_wait = false;
2781
2782 let hierarchy = plan_force_merge_hierarchy(group.segments.len());
2786 let output_ids: Vec<_> = (0..hierarchy.steps.len())
2787 .map(|_| SegmentId::new())
2788 .collect();
2789 let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
2790 let mut all_ids = source_ids.clone();
2791 all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
2792 let group_guard = {
2793 let st = self.state.lock().await;
2794 source_ids
2795 .iter()
2796 .all(|id| st.metadata.has_segment(id))
2797 .then(|| self.active_operations.try_register(all_ids))
2798 .flatten()
2799 };
2800 let _group_guard = match group_guard {
2801 Some(guard) => guard,
2802 None if !self.active_operations.is_accepting() => {
2803 return Err(Error::IndexClosed);
2804 }
2805 None => {
2806 #[cfg(test)]
2807 self.force_merge_conflict_retries
2808 .fetch_add(1, Ordering::Relaxed);
2809 log::debug!(
2810 "[force_merge] index={} group lost a registration race, replanning",
2811 self.schema.index_label()
2812 );
2813 let had_tracked_merges = !self.merge_handles.lock().is_empty();
2814 self.wait_for_merging_thread().await;
2815 if !had_tracked_merges {
2816 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2817 }
2818 continue;
2819 }
2820 };
2821
2822 log::info!(
2823 "[force_merge] index={} planned final group: {} segments, {} docs, {} merge pass(es)",
2824 self.schema.index_label(),
2825 group.segments.len(),
2826 group.total_docs,
2827 output_ids.len(),
2828 );
2829
2830 let group_global_merge_permit = if self.reorder_on_merge {
2840 let capacity_start = std::time::Instant::now();
2841 let permit = tokio::select! {
2842 biased;
2843 () = self.active_operations.wait_for_shutdown() => {
2844 return Err(Error::IndexClosed);
2845 }
2846 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2847 permit.map_err(|_| {
2848 Error::Internal(
2849 "global background merge scheduler is closed".into(),
2850 )
2851 })?
2852 }
2853 };
2854 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2855 log::info!(
2856 "[force_merge] index={} waited {:.1}s for foreground global merge capacity",
2857 self.schema.index_label(),
2858 capacity_start.elapsed().as_secs_f64(),
2859 );
2860 }
2861 Some(permit)
2862 } else {
2863 None
2864 };
2865 let _foreground_reorder = if self.reorder_on_merge {
2866 log::info!(
2867 "[force_merge] index={} prioritizing BP capacity ({} total pass slot(s))",
2868 self.schema.index_label(),
2869 self.reorder_permits.limit(),
2870 );
2871 let admission_start = std::time::Instant::now();
2872 let guard = Arc::clone(&self.reorder_permits)
2873 .begin_foreground()
2874 .await
2875 .map_err(|_| {
2876 Error::Internal("background reorder scheduler is closed".into())
2877 })?;
2878 if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
2879 log::info!(
2880 "[force_merge] index={} acquired foreground BP capacity in {:.1}s",
2881 self.schema.index_label(),
2882 admission_start.elapsed().as_secs_f64(),
2883 );
2884 }
2885 Some(guard)
2886 } else {
2887 None
2888 };
2889
2890 let source_count = group.segments.len();
2891 let mut nodes: Vec<Option<(String, u32)>> =
2892 group.segments.into_iter().map(Some).collect();
2893 nodes.resize_with(source_count + hierarchy.steps.len(), || None);
2894 for (step_index, step) in hierarchy.steps.iter().enumerate() {
2895 let final_pass = step_index + 1 == hierarchy.steps.len();
2896 let mut batch_entries = Vec::with_capacity(step.inputs.len());
2897 for &node in &step.inputs {
2898 let entry = nodes
2899 .get_mut(node)
2900 .and_then(Option::take)
2901 .expect("force-merge hierarchy must reference an available node");
2902 batch_entries.push(entry);
2903 }
2904 let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
2905 let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
2906 let output_id = output_ids[step_index];
2907
2908 let capacity_start = std::time::Instant::now();
2909 let step_global_merge_permit = if group_global_merge_permit.is_none() {
2910 Some(tokio::select! {
2911 biased;
2912 () = self.active_operations.wait_for_shutdown() => {
2913 return Err(Error::IndexClosed);
2914 }
2915 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2916 permit.map_err(|_| {
2917 Error::Internal(
2918 "global background merge scheduler is closed".into(),
2919 )
2920 })?
2921 }
2922 })
2923 } else {
2924 None
2925 };
2926 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2927 log::info!(
2928 "[force_merge] index={} waited {:.1}s for global merge capacity",
2929 self.schema.index_label(),
2930 capacity_start.elapsed().as_secs_f64(),
2931 );
2932 }
2933
2934 let reorder_bmp = final_pass && self.reorder_on_merge;
2938 log::info!(
2939 "[force_merge] index={} {} pass: {} segments ({} docs, bp={})",
2940 self.schema.index_label(),
2941 if final_pass {
2942 "final"
2943 } else {
2944 "fan-in reduction"
2945 },
2946 batch.len(),
2947 batch_docs,
2948 reorder_bmp,
2949 );
2950 let (new_segment_id, total_docs, _) = self
2951 .merge_and_replace_registered(
2952 &batch,
2953 output_id,
2954 reorder_bmp,
2955 ReorderPriority::Foreground,
2956 )
2957 .await
2958 .map_err(|error| error.error)?;
2959 drop(step_global_merge_permit);
2960
2961 let refresh_start = std::time::Instant::now();
2964 refresh_snapshots().await?;
2965 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2966 log::info!(
2967 "[force_merge] index={} post-replacement snapshot refresh took {:.1}s",
2968 self.schema.index_label(),
2969 refresh_start.elapsed().as_secs_f64(),
2970 );
2971 }
2972
2973 let output_node = source_count + step_index;
2974 debug_assert!(nodes[output_node].is_none());
2975 nodes[output_node] = Some((new_segment_id, total_docs));
2976 }
2977 let (root_id, _) = nodes[hierarchy.root]
2978 .take()
2979 .expect("force-merge hierarchy must produce its root");
2980 debug_assert!(nodes.into_iter().all(|node| node.is_none()));
2981 completed_outputs.insert(root_id);
2982 }
2983 }
2984
2985 fn segment_needs_vector_rewrite(
2986 &self,
2987 schema: &crate::dsl::Schema,
2988 reader: &SegmentReader,
2989 field_ids: &[u32],
2990 trained: &TrainedVectorStructures,
2991 rewrite_existing: bool,
2992 ) -> Result<bool> {
2993 for &field_id in field_ids {
2994 let flat = reader.flat_vectors().get(&field_id);
2995 let ann = reader.vector_indexes().get(&field_id);
2996 if ann.is_some() && flat.is_none() {
2997 return Err(Error::Corruption(format!(
2998 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2999 reader.meta().id,
3000 )));
3001 }
3002
3003 let Some(flat) = flat else {
3004 continue;
3005 };
3006 if flat.num_vectors == 0 {
3007 continue;
3008 }
3009 if rewrite_existing {
3010 return Ok(true);
3011 }
3012 let field = crate::dsl::Field(field_id);
3013 let entry = schema.get_field_entry(field).ok_or_else(|| {
3014 Error::Corruption(format!(
3015 "segment {:032x} references unknown vector field {field_id}",
3016 reader.meta().id,
3017 ))
3018 })?;
3019 let current = match entry.field_type {
3020 crate::dsl::FieldType::DenseVector
3024 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3025 config.index_type == crate::dsl::VectorIndexType::Tq
3026 }) =>
3027 {
3028 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
3029 }
3030 crate::dsl::FieldType::DenseVector
3031 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3032 config.index_type == crate::dsl::VectorIndexType::IvfTq
3033 }) =>
3034 {
3035 let config = entry
3036 .dense_vector_config
3037 .as_ref()
3038 .expect("matched IVF-TQ configuration");
3039 match (ann, trained.centroids.get(&field_id)) {
3040 (
3041 Some(crate::segment::VectorIndex::IvfTq { index, .. }),
3042 Some(centroids),
3043 ) => {
3044 let header = index.get().header();
3045 crate::structures::is_ivf_tq_cosine_generation(centroids.version)
3046 && crate::structures::is_ivf_tq_cosine_generation(
3047 header.quantizer_version,
3048 )
3049 && header.dim == config.dim
3050 && header.num_clusters == centroids.num_clusters
3051 && header.quantizer_version == centroids.version
3052 && header.codebook_version
3053 == crate::structures::vector::quantization::tq_expected_fingerprint(
3054 config.dim,
3055 )
3056 && header.routing == config.ivf_routing
3057 }
3058 (None, None) => true,
3059 _ => false,
3060 }
3061 }
3062 crate::dsl::FieldType::DenseVector
3063 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3064 config.index_type == crate::dsl::VectorIndexType::Scann
3065 }) =>
3066 {
3067 match (ann, trained.scann_artifacts.get(&field_id)) {
3068 (Some(crate::segment::VectorIndex::ScannAh(index)), Some(artifact)) => {
3069 index
3070 .get()
3071 .validate_scann_generation(
3072 artifact.config(),
3073 artifact.generation(),
3074 artifact.artifact_id(),
3075 )
3076 .is_ok()
3077 }
3078 (None, None) => true,
3079 _ => false,
3080 }
3081 }
3082 crate::dsl::FieldType::DenseVector => false,
3085 crate::dsl::FieldType::BinaryDenseVector
3086 if entry
3087 .binary_dense_vector_config
3088 .as_ref()
3089 .is_some_and(|config| {
3090 config.index_type == crate::dsl::BinaryIndexType::Scann
3091 }) =>
3092 {
3093 match (ann, trained.scann_artifacts.get(&field_id)) {
3094 (Some(crate::segment::VectorIndex::ScannBinary(index)), Some(artifact)) => {
3095 index
3096 .get()
3097 .validate_scann_generation(
3098 artifact.config(),
3099 artifact.generation(),
3100 artifact.artifact_id(),
3101 )
3102 .is_ok()
3103 }
3104 (None, None) => true,
3105 _ => false,
3106 }
3107 }
3108 crate::dsl::FieldType::BinaryDenseVector => matches!(
3109 (ann, trained.binary_quantizers.get(&field_id)),
3110 (Some(crate::segment::VectorIndex::BinaryIvf(_)), Some(_)) | (None, None)
3111 ),
3112 _ => false,
3113 };
3114 if !current {
3115 return Ok(true);
3116 }
3117 }
3118 Ok(false)
3119 }
3120
3121 async fn acquire_vector_rewrite_capacity(
3122 &self,
3123 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
3124 let global = tokio::select! {
3125 biased;
3126 () = self.active_operations.wait_for_shutdown() => {
3127 return Err(Error::IndexClosed);
3128 }
3129 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3130 permit.map_err(|_| Error::Internal(
3131 "global background merge scheduler is closed".into()
3132 ))?
3133 }
3134 };
3135 let local = tokio::select! {
3136 biased;
3137 () = self.active_operations.wait_for_shutdown() => {
3138 return Err(Error::IndexClosed);
3139 }
3140 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
3141 permit.map_err(|_| Error::Internal(
3142 "background merge scheduler is closed".into()
3143 ))?
3144 }
3145 };
3146 Ok((global, local))
3147 }
3148
3149 async fn build_vector_replacement(
3150 self: &Arc<Self>,
3151 schemas: (&Arc<crate::dsl::Schema>, &Arc<crate::dsl::Schema>),
3152 segment_id: &str,
3153 source_id: SegmentId,
3154 output_id: SegmentId,
3155 trained: &TrainedVectorStructures,
3156 failure_context: &'static str,
3157 ) -> Result<(String, u32, OutputCleanupGuard)> {
3158 let mut cleanup = self.output_cleanup_guard(output_id);
3159 match crate::segment::reorder::rewrite_vector_segment(
3160 self.directory.as_ref(),
3161 schemas,
3162 source_id,
3163 output_id,
3164 self.term_cache_blocks,
3165 trained,
3166 Some(self.background_cpu_pool()),
3167 )
3168 .await
3169 {
3170 Ok((new_id, doc_count)) => {
3171 self.validate_completed_segment(&new_id, doc_count).await?;
3172 Ok((new_id, doc_count, cleanup))
3173 }
3174 Err(error) => {
3175 self.delete_output_if_unregistered(output_id, failure_context)
3176 .await;
3177 cleanup.disarm();
3178 if is_deterministic_source_error(&error) {
3179 self.quarantine_segment(segment_id, &error);
3180 }
3181 Err(error)
3182 }
3183 }
3184 }
3185
3186 pub(crate) async fn stage_vector_generation(
3190 self: &Arc<Self>,
3191 _artifact_update: &VectorArtifactUpdateGuard,
3192 segment_ids: &[String],
3193 field_ids: &[u32],
3194 trained: Arc<TrainedVectorStructures>,
3195 rewrite_existing: bool,
3196 ) -> Result<Vec<StagedVectorSegment>> {
3197 let schema = self.published_generation().schema.clone();
3198 self.stage_vector_generation_with_schema(
3199 _artifact_update,
3200 segment_ids,
3201 field_ids,
3202 trained,
3203 rewrite_existing,
3204 schema,
3205 )
3206 .await
3207 }
3208
3209 pub(crate) async fn stage_vector_generation_with_schema(
3210 self: &Arc<Self>,
3211 _artifact_update: &VectorArtifactUpdateGuard,
3212 segment_ids: &[String],
3213 field_ids: &[u32],
3214 trained: Arc<TrainedVectorStructures>,
3215 rewrite_existing: bool,
3216 schema: Arc<crate::dsl::Schema>,
3217 ) -> Result<Vec<StagedVectorSegment>> {
3218 if !self.vector_artifact_update.load(Ordering::Acquire) {
3219 return Err(Error::Internal(
3220 "cannot stage a vector generation without an exclusive update lease".into(),
3221 ));
3222 }
3223
3224 let source_schema = self.published_generation().schema.clone();
3225 let mut staged = Vec::new();
3226 for segment_id in segment_ids {
3227 if self.quarantined_segments.lock().contains(segment_id) {
3228 return Err(Error::Corruption(format!(
3229 "segment {segment_id} is quarantined after a deterministic source failure"
3230 )));
3231 }
3232 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3233 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3234 })?;
3235
3236 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3239
3240 let output_id = SegmentId::new();
3241 let output_hex = output_id.to_hex();
3242 let operation = {
3243 let st = self.state.lock().await;
3244 if !st.metadata.has_segment(segment_id) {
3245 return Err(Error::Corruption(format!(
3246 "vector generation source {segment_id} disappeared while lifecycle work was paused"
3247 )));
3248 }
3249 self.active_operations
3250 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
3251 }
3252 .ok_or_else(|| {
3253 if self.active_operations.is_accepting() {
3254 Error::Internal(format!(
3255 "vector generation could not claim stable source {segment_id}"
3256 ))
3257 } else {
3258 Error::IndexClosed
3259 }
3260 })?;
3261
3262 let reader = SegmentReader::open(
3263 self.directory.as_ref(),
3264 source_id,
3265 Arc::clone(&source_schema),
3266 self.term_cache_blocks,
3267 )
3268 .await?;
3269 if !self.segment_needs_vector_rewrite(
3270 schema.as_ref(),
3271 &reader,
3272 field_ids,
3273 trained.as_ref(),
3274 rewrite_existing,
3275 )? {
3276 continue;
3277 }
3278 drop(reader);
3279
3280 let (new_id, doc_count, cleanup) = self
3281 .build_vector_replacement(
3282 (&source_schema, &schema),
3283 segment_id,
3284 source_id,
3285 output_id,
3286 trained.as_ref(),
3287 "vector generation staging failure",
3288 )
3289 .await?;
3290 debug_assert_eq!(new_id, output_hex);
3291 let output_reader = SegmentReader::open(
3292 self.directory.as_ref(),
3293 output_id,
3294 Arc::clone(&schema),
3295 self.term_cache_blocks,
3296 )
3297 .await?;
3298 if self.segment_needs_vector_rewrite(
3299 schema.as_ref(),
3300 &output_reader,
3301 field_ids,
3302 trained.as_ref(),
3303 false,
3304 )? {
3305 return Err(Error::Corruption(format!(
3306 "staged vector segment {new_id} does not match its candidate codebook generation"
3307 )));
3308 }
3309
3310 staged.push(StagedVectorSegment {
3311 source_id: segment_id.clone(),
3312 output_id,
3313 doc_count,
3314 _operation: operation,
3315 cleanup,
3316 });
3317 }
3318 Ok(staged)
3319 }
3320
3321 async fn rewrite_vector_segment_once(
3322 self: &Arc<Self>,
3323 segment_id: &str,
3324 field_ids: &[u32],
3325 ) -> Result<VectorSegmentRewriteOutcome> {
3326 if self.quarantined_segments.lock().contains(segment_id) {
3327 return Err(Error::Corruption(format!(
3328 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3329 )));
3330 }
3331 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3332 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3333 })?;
3334
3335 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3340
3341 let output_id = SegmentId::new();
3342 let output_hex = output_id.to_hex();
3343 let all_ids = vec![segment_id.to_owned(), output_hex];
3344 let operation = {
3345 let st = self.state.lock().await;
3346 if !st.metadata.has_segment(segment_id) {
3347 return Ok(VectorSegmentRewriteOutcome::SourceGone);
3348 }
3349 self.active_operations.try_register(all_ids)
3350 };
3351 let _operation = match operation {
3352 Some(operation) => operation,
3353 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3354 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3355 };
3356
3357 let Some(trained) = self.trained_for_segment_build() else {
3358 return Ok(VectorSegmentRewriteOutcome::Deferred);
3359 };
3360 let schema = self.published_generation().schema.clone();
3361
3362 let reader = SegmentReader::open(
3363 self.directory.as_ref(),
3364 source_id,
3365 Arc::clone(&schema),
3366 self.term_cache_blocks,
3367 )
3368 .await?;
3369 if !self.segment_needs_vector_rewrite(
3370 schema.as_ref(),
3371 &reader,
3372 field_ids,
3373 trained.as_ref(),
3374 false,
3375 )? {
3376 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3377 }
3378 drop(reader);
3379
3380 let (new_id, doc_count, mut output_cleanup) = self
3381 .build_vector_replacement(
3382 (&schema, &schema),
3383 segment_id,
3384 source_id,
3385 output_id,
3386 trained.as_ref(),
3387 "vector rewrite failure",
3388 )
3389 .await?;
3390
3391 if let Err(error) = self
3392 .replace_segments(
3393 &[segment_id.to_owned()],
3394 new_id,
3395 doc_count,
3396 ReplacementLayout::PreserveSingleSource,
3397 )
3398 .await
3399 {
3400 self.delete_output_if_unregistered(output_id, "vector replacement failure")
3401 .await;
3402 output_cleanup.disarm();
3403 return Err(error);
3404 }
3405 output_cleanup.disarm();
3406 Ok(VectorSegmentRewriteOutcome::Rewritten)
3407 }
3408
3409 pub(crate) async fn rewrite_vector_segments(
3414 self: &Arc<Self>,
3415 field_ids: &[u32],
3416 ) -> Result<usize> {
3417 if field_ids.is_empty() {
3418 return Ok(0);
3419 }
3420 let mut rewritten = 0usize;
3421 loop {
3422 let segment_ids = self.get_segment_ids().await;
3423 let mut conflicted = false;
3424 let mut changed = false;
3425 for segment_id in segment_ids {
3426 match self
3427 .rewrite_vector_segment_once(&segment_id, field_ids)
3428 .await?
3429 {
3430 VectorSegmentRewriteOutcome::Rewritten => {
3431 rewritten += 1;
3432 changed = true;
3433 }
3434 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3435 VectorSegmentRewriteOutcome::Deferred => {
3436 return Err(Error::Internal(
3437 "ANN finalization lost the published trained generation".into(),
3438 ));
3439 }
3440 VectorSegmentRewriteOutcome::AlreadyCurrent
3441 | VectorSegmentRewriteOutcome::SourceGone => {}
3442 }
3443 }
3444 if !conflicted && !changed {
3445 log::info!(
3446 "[dense_vector_rewrite] index={} ANN finalization complete ({} segment(s) rewritten)",
3447 self.schema.index_label(),
3448 rewritten,
3449 );
3450 return Ok(rewritten);
3451 }
3452 tokio::select! {
3453 biased;
3454 () = self.active_operations.wait_for_shutdown() => {
3455 return Err(Error::IndexClosed);
3456 }
3457 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3458 }
3459 }
3460 }
3461
3462 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3467 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3468 return;
3469 }
3470 let manager = Arc::clone(self);
3471 let future = async move {
3472 let field_ids = manager
3473 .read_metadata(|metadata| {
3474 metadata
3475 .vector_fields
3476 .keys()
3477 .filter(|field_id| metadata.is_field_built(**field_id))
3478 .copied()
3479 .collect::<Vec<_>>()
3480 })
3481 .await;
3482 for segment_id in segment_ids {
3483 loop {
3484 match manager
3485 .rewrite_vector_segment_once(&segment_id, &field_ids)
3486 .await
3487 {
3488 Ok(VectorSegmentRewriteOutcome::Conflict) => {
3489 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3490 }
3491 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3492 Ok(_) => break,
3493 Err(error) => {
3494 log::error!(
3495 "[dense_vector_rewrite] index={} failed to upgrade newly committed segment {}: {}",
3496 manager.schema.index_label(),
3497 segment_id,
3498 error,
3499 );
3500 break;
3501 }
3502 }
3503 }
3504 }
3505 };
3506 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3507 log::warn!(
3508 "[dense_vector_rewrite] index={} runtime unavailable; newly committed flat segment upgrade deferred",
3509 self.schema.index_label()
3510 );
3511 return;
3512 };
3513 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3514 log::warn!(
3515 "[dense_vector_rewrite] index={} runtime rejected newly committed flat segment upgrade",
3516 self.schema.index_label()
3517 );
3518 }
3519 }
3520
3521 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3528 self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3529 .await
3530 }
3531
3532 pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3535 self: &Arc<Self>,
3536 mut refresh_snapshots: F,
3537 ) -> Result<()>
3538 where
3539 F: FnMut() -> Fut,
3540 Fut: std::future::Future<Output = Result<()>>,
3541 {
3542 self.wait_for_all_merges().await;
3543 refresh_snapshots().await?;
3544 let segment_ids = self.get_segment_ids().await;
3545
3546 if segment_ids.is_empty() {
3547 log::info!(
3548 "[reorder] index={} no segments to reorder",
3549 self.schema.index_label()
3550 );
3551 return Ok(());
3552 }
3553
3554 log::info!(
3555 "[reorder] index={} reordering {} segments",
3556 self.schema.index_label(),
3557 segment_ids.len()
3558 );
3559
3560 for seg_id in segment_ids {
3561 match self
3562 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3563 .await
3564 {
3565 Ok(true) => refresh_snapshots().await?,
3566 Ok(false) => log::warn!(
3567 "[reorder] index={} segment {} skipped (in merge)",
3568 self.schema.index_label(),
3569 seg_id
3570 ),
3571 Err(e) => return Err(e),
3572 }
3573 }
3574
3575 refresh_snapshots().await?;
3578 log::info!(
3579 "[reorder] index={} all segments reordered",
3580 self.schema.index_label()
3581 );
3582 Ok(())
3583 }
3584
3585 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3590 self.unreordered_segments()
3591 .await
3592 .into_iter()
3593 .map(|(id, _)| id)
3594 .collect()
3595 }
3596
3597 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3600 let quarantined = self.quarantined_segments.lock().clone();
3601 let paused = self.paused_reorder_segments();
3602 let st = self.state.lock().await;
3603 let active_ids = self.active_operations.snapshot();
3604 st.metadata
3605 .segment_metas
3606 .iter()
3607 .filter(|(id, info)| {
3608 !info.reordered
3609 && info.bp_converged
3610 && !active_ids.contains(*id)
3611 && !quarantined.contains(*id)
3612 && !paused.contains(*id)
3613 })
3614 .map(|(id, info)| (id.clone(), info.num_docs))
3615 .collect()
3616 }
3617
3618 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3622 self.unconverged_segments_below(u32::MAX)
3623 .await
3624 .into_iter()
3625 .map(|(id, docs, _)| (id, docs))
3626 .collect()
3627 }
3628
3629 pub async fn unconverged_segments_below(
3632 &self,
3633 max_unconverged_passes: u32,
3634 ) -> Vec<(String, u32, u32)> {
3635 let quarantined = self.quarantined_segments.lock().clone();
3636 let paused = self.paused_reorder_segments();
3637 let st = self.state.lock().await;
3638 let active_ids = self.active_operations.snapshot();
3639 st.metadata
3640 .segment_metas
3641 .iter()
3642 .filter(|(id, info)| {
3643 !info.bp_converged
3644 && info.bp_unconverged_passes < max_unconverged_passes
3645 && !active_ids.contains(*id)
3646 && !quarantined.contains(*id)
3647 && !paused.contains(*id)
3648 })
3649 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3650 .collect()
3651 }
3652
3653 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3663 let st = self.state.lock().await;
3664 let deepening = ids.iter().any(|id| {
3665 st.metadata
3666 .segment_metas
3667 .get(id)
3668 .is_some_and(|info| !info.bp_converged)
3669 });
3670 drop(st);
3671 if deepening {
3672 log::info!(
3673 "[reorder] index={} source BP lineage unconverged — forcing record-level BP (deepening pass)",
3674 self.schema.index_label(),
3675 );
3676 crate::segment::reorder::BpGranularity::Records
3677 } else {
3678 crate::segment::reorder::BpGranularity::Auto
3679 }
3680 }
3681
3682 pub async fn reorder_single_segment(
3687 self: &Arc<Self>,
3688 seg_id: &str,
3689 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3690 bp_budget: crate::segment::BpBudget,
3691 ) -> Result<bool> {
3692 let source_id = SegmentId::from_hex(seg_id)
3693 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3694 if self.quarantined_segments.lock().contains(seg_id) {
3695 return Err(Error::Corruption(format!(
3696 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3697 seg_id
3698 )));
3699 }
3700 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3701 log::debug!(
3702 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3703 self.schema.index_label(),
3704 seg_id,
3705 );
3706 return Ok(false);
3707 }
3708
3709 let reorder_gate = Arc::clone(&self.reorder_permits);
3714 let _reorder_permit = tokio::select! {
3715 biased;
3716 () = self.active_operations.wait_for_shutdown() => {
3717 return Err(Error::IndexClosed);
3718 }
3719 permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3720 permit.map_err(|_| {
3721 Error::Internal("background reorder scheduler is closed".into())
3722 })?
3723 }
3724 };
3725
3726 let output_id = SegmentId::new();
3727 let output_hex = output_id.to_hex();
3728 let source_ids = [seg_id.to_string()];
3729 let granularity = self.merge_granularity(&source_ids).await;
3730
3731 let all_ids = vec![seg_id.to_string(), output_hex];
3737 let (_guard, source_docs, schema) = {
3738 let st = self.state.lock().await;
3739 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3743 log::debug!(
3744 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3745 self.schema.index_label(),
3746 seg_id,
3747 );
3748 return Ok(false);
3749 }
3750 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3751 log::info!(
3752 "[optimizer] index={} segment {} no longer in metadata (merged away), skipping reorder",
3753 self.schema.index_label(),
3754 seg_id
3755 );
3756 self.clear_reorder_retry(seg_id);
3757 return Ok(false);
3758 };
3759
3760 let schema = self.published_generation().schema.clone();
3761 match self.active_operations.try_register(all_ids) {
3762 Some(guard) => (guard, source_meta.num_docs, schema),
3763 None if !self.active_operations.is_accepting() => {
3764 return Err(Error::IndexClosed);
3765 }
3766 None => {
3767 log::debug!(
3768 "[optimizer] index={} segment {} in active merge, skipping",
3769 self.schema.index_label(),
3770 seg_id
3771 );
3772 return Ok(false);
3773 }
3774 }
3775 };
3776
3777 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3782 if is_deterministic_source_error(&error) {
3783 self.quarantine_segment(seg_id, &error);
3784 } else if !matches!(&error, Error::IndexClosed) {
3785 self.pause_reorder_retries(seg_id, &error);
3786 }
3787 return Err(error);
3788 }
3789
3790 let mut output_cleanup = self.output_cleanup_guard(output_id);
3791
3792 let reorder_result = crate::segment::reorder::reorder_segment(
3793 self.directory.as_ref(),
3794 &schema,
3795 source_id,
3796 output_id,
3797 self.term_cache_blocks,
3798 self.bp_memory_budget_bytes,
3799 bp_budget,
3800 granularity,
3801 self.optimization,
3802 self.posting_codec,
3803 rayon_pool,
3804 Some(self.active_operations.cancellation_flag()),
3805 )
3806 .await;
3807 let (new_id, total_docs, bp_converged) = match reorder_result {
3808 Ok(v) => v,
3809 Err(e) => {
3810 self.delete_output_if_unregistered(output_id, "reorder failure")
3813 .await;
3814 output_cleanup.disarm();
3815 if is_deterministic_source_error(&e) {
3816 self.quarantine_segment(seg_id, &e);
3817 } else if !matches!(&e, Error::IndexClosed) {
3818 self.pause_reorder_retries(seg_id, &e);
3819 }
3820 return Err(e);
3821 }
3822 };
3823
3824 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3830 if let Err(e) = self
3831 .replace_segments(
3832 &[seg_id.to_string()],
3833 new_id,
3834 total_docs,
3835 ReplacementLayout::BpReordered {
3836 converged: ladder_converged,
3837 },
3838 )
3839 .await
3840 {
3841 self.delete_output_if_unregistered(output_id, "replacement failure")
3842 .await;
3843 output_cleanup.disarm();
3844 if !matches!(&e, Error::IndexClosed) {
3845 self.pause_reorder_retries(seg_id, &e);
3846 }
3847 return Err(e);
3848 }
3849 output_cleanup.disarm();
3850 self.clear_reorder_retry(seg_id);
3851
3852 Ok(true)
3853 }
3854
3855 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3862 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3863
3864 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3865 for entry in entries {
3866 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3867 continue;
3868 };
3869 let Some(rest) = filename.strip_prefix("seg_") else {
3870 continue;
3871 };
3872 let Some(hex_id) = rest.get(..32) else {
3873 continue;
3874 };
3875 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3876 continue;
3877 }
3878 orphan_files
3879 .entry(hex_id.to_ascii_lowercase())
3880 .or_default()
3881 .push(entry);
3882 }
3883 }
3884
3885 let mut deleted = 0;
3886 for (hex_id, paths) in &orphan_files {
3887 let deletion_guard = {
3892 let st = self.state.lock().await;
3893 if st.metadata.has_segment(hex_id) {
3894 continue;
3895 }
3896 let Some(guard) = self
3897 .active_operations
3898 .try_register(vec![hex_id.to_string()])
3899 else {
3900 continue;
3901 };
3902 if self.tracker.is_deletion_protected(hex_id) {
3903 drop(guard);
3904 continue;
3905 }
3906 guard
3907 };
3908
3909 let results =
3914 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3915 .await;
3916 let removed = results.into_iter().all(|result| match result {
3917 Ok(()) => true,
3918 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3919 Err(error) => {
3920 log::warn!(
3921 "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
3922 self.schema.index_label(),
3923 hex_id,
3924 error,
3925 );
3926 false
3927 }
3928 });
3929 drop(deletion_guard);
3932 if removed {
3933 deleted += 1;
3934 log::info!(
3935 "[segment_cleanup] index={} swept orphan segment {}",
3936 self.schema.index_label(),
3937 hex_id
3938 );
3939 }
3940 }
3941
3942 Ok(deleted)
3943 }
3944}
3945
3946#[cfg(test)]
3947mod tests {
3948 use super::*;
3949 use std::sync::atomic::{AtomicBool, Ordering};
3950
3951 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3952 let schema = crate::dsl::SchemaBuilder::default().build();
3953 let metadata = IndexMetadata::new(schema.clone());
3954 Arc::new(SegmentManager::new(
3955 Arc::new(crate::directories::RamDirectory::new()),
3956 Arc::new(schema),
3957 metadata,
3958 Box::new(crate::merge::NoMergePolicy),
3959 0,
3960 1,
3961 Arc::new(Semaphore::new(1)),
3962 None,
3963 1024,
3964 Arc::new(ReorderConcurrencyGate::new(1)),
3965 None,
3966 ))
3967 }
3968
3969 #[test]
3970 fn force_merge_planner_pairs_large_and_small_segments() {
3971 let groups = plan_force_merge_groups(
3972 vec![
3973 ("a".into(), 6),
3974 ("b".into(), 6),
3975 ("c".into(), 4),
3976 ("d".into(), 4),
3977 ],
3978 10,
3979 );
3980
3981 assert_eq!(groups.len(), 2);
3982 assert!(groups.iter().all(|group| group.total_docs == 10));
3983 assert!(groups.iter().all(|group| group.segments.len() == 2));
3984 }
3985
3986 #[test]
3987 fn force_merge_planner_leaves_oversized_segments_alone() {
3988 let groups = plan_force_merge_groups(
3989 vec![
3990 ("oversized".into(), 11),
3991 ("small-a".into(), 5),
3992 ("small-b".into(), 5),
3993 ],
3994 10,
3995 );
3996
3997 assert_eq!(groups.len(), 2);
3998 assert_eq!(groups[0].total_docs, 10);
3999 assert_eq!(groups[0].segments.len(), 2);
4000 assert_eq!(groups[1].total_docs, 11);
4001 assert_eq!(groups[1].segments.len(), 1);
4002 }
4003
4004 #[test]
4005 fn force_merge_planner_never_exceeds_segment_format_limit() {
4006 let groups = plan_force_merge_groups(
4007 vec![
4008 ("large-a".into(), 3_000_000_000),
4009 ("large-b".into(), 2_000_000_000),
4010 ],
4011 u64::from(u32::MAX),
4012 );
4013 assert_eq!(groups.len(), 2);
4014 assert!(
4015 groups
4016 .iter()
4017 .all(|group| group.total_docs <= u64::from(u32::MAX))
4018 );
4019 }
4020
4021 #[test]
4022 fn force_merge_hierarchy_has_one_final_bp_pass() {
4023 assert_eq!(force_merge_output_count(1), 0);
4024 assert_eq!(force_merge_output_count(2), 1);
4025 assert_eq!(force_merge_output_count(64), 1);
4026 assert_eq!(force_merge_output_count(65), 2);
4027 assert_eq!(force_merge_output_count(127), 2);
4028 assert_eq!(force_merge_output_count(128), 3);
4029 assert_eq!(force_merge_output_count(1_000), 16);
4030 }
4031
4032 fn expand_force_merge_node(
4033 hierarchy: &ForceMergeHierarchy,
4034 source_count: usize,
4035 node: usize,
4036 sources: &mut Vec<usize>,
4037 ) {
4038 if node < source_count {
4039 sources.push(node);
4040 return;
4041 }
4042
4043 let step_index = node - source_count;
4044 let step = hierarchy
4045 .steps
4046 .get(step_index)
4047 .expect("merge input must refer to an existing source or output");
4048 for &input in &step.inputs {
4049 assert!(
4050 input < node,
4051 "merge step {step_index} refers to a future output node {input}"
4052 );
4053 expand_force_merge_node(hierarchy, source_count, input, sources);
4054 }
4055 }
4056
4057 #[test]
4058 fn force_merge_hierarchy_has_minimal_valid_arity() {
4059 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
4060
4061 for source_count in source_counts {
4062 let hierarchy = plan_force_merge_hierarchy(source_count);
4063 let output_count = hierarchy.steps.len();
4064
4065 assert!(
4066 hierarchy
4067 .steps
4068 .iter()
4069 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
4070 "invalid merge arity for {source_count} sources"
4071 );
4072 assert!(
4073 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
4074 "{output_count} outputs cannot reduce {source_count} sources"
4075 );
4076 assert!(
4077 output_count == 1
4078 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
4079 "{output_count} outputs are not minimal for {source_count} sources"
4080 );
4081 assert_eq!(output_count, force_merge_output_count(source_count));
4082 }
4083 }
4084
4085 #[test]
4086 fn force_merge_hierarchy_preserves_exact_source_order() {
4087 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
4088 let hierarchy = plan_force_merge_hierarchy(source_count);
4089 let mut sources = Vec::with_capacity(source_count);
4090 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
4091 assert_eq!(
4092 sources,
4093 (0..source_count).collect::<Vec<_>>(),
4094 "source order changed for {source_count} sources"
4095 );
4096 }
4097 }
4098
4099 fn force_merge_rewrite_cost(source_count: usize) -> usize {
4100 let hierarchy = plan_force_merge_hierarchy(source_count);
4101 let mut node_weights = vec![1usize; source_count];
4102 let mut rewrite_cost = 0usize;
4103
4104 for (step_index, step) in hierarchy.steps.iter().enumerate() {
4105 let output = source_count + step_index;
4106 let output_weight = step
4107 .inputs
4108 .iter()
4109 .map(|&input| {
4110 assert!(
4111 input < output,
4112 "merge step {step_index} refers to future output {input}"
4113 );
4114 node_weights[input]
4115 })
4116 .sum::<usize>();
4117 rewrite_cost += output_weight;
4118 node_weights.push(output_weight);
4119 }
4120
4121 assert_eq!(node_weights[hierarchy.root], source_count);
4122 rewrite_cost
4123 }
4124
4125 #[test]
4126 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
4127 assert_eq!(force_merge_rewrite_cost(65), 67);
4128 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
4129 }
4130
4131 #[test]
4132 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
4133 assert_eq!(
4134 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
4135 (false, false, 3),
4136 );
4137 assert_eq!(
4138 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
4139 (true, false, 4),
4140 );
4141 assert_eq!(
4142 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
4143 (true, true, 0),
4144 );
4145 }
4146
4147 #[tokio::test]
4148 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
4149 let mut schema_builder = crate::dsl::SchemaBuilder::default();
4150 let field = schema_builder.add_text_field("text", true, true);
4151 let schema = schema_builder.build();
4152 let directory = crate::directories::RamDirectory::new();
4153 let config = crate::index::IndexConfig {
4154 num_indexing_threads: 1,
4155 merge_policy: Box::new(crate::merge::NoMergePolicy),
4156 ..Default::default()
4157 };
4158 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
4159 .await
4160 .unwrap();
4161 for value in ["one", "two", "three"] {
4162 let mut document = crate::dsl::Document::new();
4163 document.add_text(field, value);
4164 writer.add_document(document).unwrap();
4165 writer.commit().await.unwrap();
4166 }
4167
4168 let manager = Arc::clone(writer.segment_manager());
4169 let held_id = manager.get_segment_ids().await.pop().unwrap();
4170 let mut held = Some(
4171 manager
4172 .active_operations
4173 .try_register(vec![held_id])
4174 .unwrap(),
4175 );
4176 let batches = Arc::new(AtomicUsize::new(0));
4177 let batch_count = Arc::clone(&batches);
4178 writer
4179 .force_merge_with_snapshot_refresh(move || {
4180 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
4181 if refresh == 2 {
4184 drop(held.take());
4185 }
4186 std::future::ready(Ok(()))
4187 })
4188 .await
4189 .unwrap();
4190
4191 assert_eq!(manager.get_segment_ids().await.len(), 1);
4192 assert_eq!(
4193 batches.load(Ordering::Relaxed),
4194 4,
4195 "initial/final refreshes plus two replacements are required after the held source releases"
4196 );
4197 }
4198
4199 #[tokio::test]
4200 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
4201 let mut schema_builder = crate::dsl::SchemaBuilder::default();
4202 schema_builder.set_reorder_on_merge(true);
4203 let schema = schema_builder.build();
4204 let mut metadata = IndexMetadata::new(schema.clone());
4205 metadata.add_segment("00000000000000000000000000000001".into(), 1);
4206 metadata.add_segment("00000000000000000000000000000002".into(), 1);
4207
4208 let global_merge_permits = Arc::new(Semaphore::new(1));
4209 let manager = Arc::new(SegmentManager::new(
4210 Arc::new(crate::directories::RamDirectory::new()),
4211 Arc::new(schema),
4212 metadata,
4213 Box::new(crate::merge::NoMergePolicy),
4214 0,
4215 1,
4216 Arc::clone(&global_merge_permits),
4217 None,
4218 1024,
4219 Arc::new(ReorderConcurrencyGate::new(1)),
4220 None,
4221 ));
4222
4223 manager.active_operations.pause_non_indexing();
4228 let force_merge = {
4229 let manager = Arc::clone(&manager);
4230 tokio::spawn(async move { manager.force_merge().await })
4231 };
4232 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4233 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4234 tokio::task::yield_now().await;
4235 }
4236 })
4237 .await
4238 .expect("force merge never reached the paused group claim");
4239
4240 assert_eq!(
4241 global_merge_permits.available_permits(),
4242 1,
4243 "force merge retained global capacity while vector staging blocked group ownership"
4244 );
4245
4246 force_merge.abort();
4247 let _ = force_merge.await;
4248 manager.active_operations.resume_non_indexing();
4249 }
4250
4251 #[test]
4252 fn output_cleanup_guard_runs_during_panic_unwind() {
4253 let cleaned = Arc::new(AtomicBool::new(false));
4254 let cleaned_in_callback = Arc::clone(&cleaned);
4255 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4256 cleaned_in_callback.store(true, Ordering::SeqCst);
4257 });
4258
4259 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4260 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4261 panic!("simulated reorder panic");
4262 }));
4263
4264 assert!(result.is_err());
4265 assert!(
4266 cleaned.load(Ordering::SeqCst),
4267 "partial output cleanup must run during unwind"
4268 );
4269 }
4270
4271 #[test]
4272 fn output_cleanup_guard_disarms_after_commit() {
4273 let cleaned = Arc::new(AtomicBool::new(false));
4274 let cleaned_in_callback = Arc::clone(&cleaned);
4275 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4276 cleaned_in_callback.store(true, Ordering::SeqCst);
4277 });
4278
4279 {
4280 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4281 guard.disarm();
4282 }
4283
4284 assert!(!cleaned.load(Ordering::SeqCst));
4285 }
4286
4287 #[test]
4288 fn test_active_operation_guard_releases_ownership() {
4289 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4290 {
4291 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4292 let snap = active.snapshot();
4293 assert!(snap.contains("a"));
4294 assert!(snap.contains("b"));
4295 }
4296 assert!(active.snapshot().is_empty());
4297 }
4298
4299 #[test]
4300 fn test_non_overlapping_operations_can_run_concurrently() {
4301 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4302 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4303 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4304 let snap = active.snapshot();
4305 assert_eq!(snap.len(), 4);
4306
4307 drop(first);
4308 let snap = active.snapshot();
4309 assert_eq!(snap.len(), 2);
4310 assert!(snap.contains("c"));
4311 assert!(snap.contains("d"));
4312 }
4313
4314 #[test]
4315 fn test_overlapping_operation_is_rejected_until_release() {
4316 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4317 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4318 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4319 drop(first);
4320 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4321 }
4322
4323 #[test]
4324 fn test_active_operation_snapshot() {
4325 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4326 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4327 let snap = active.snapshot();
4328 assert!(snap.contains("x"));
4329 assert!(snap.contains("y"));
4330 assert!(!snap.contains("z"));
4331 }
4332
4333 #[tokio::test]
4334 async fn operation_barrier_ignores_producers_started_after_snapshot() {
4335 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4336 let before_gate = active.try_register(vec!["old".into()]).unwrap();
4337 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4338 assert_eq!(parked_indexing, 0);
4339 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4340
4341 let waiter = {
4342 let active = Arc::clone(&active);
4343 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4344 };
4345 tokio::task::yield_now().await;
4346 assert!(!waiter.is_finished());
4347
4348 drop(before_gate);
4349 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4350 .await
4351 .expect("pre-gate operation barrier was starved by a post-gate producer")
4352 .unwrap();
4353 assert!(active.snapshot().contains("new-flat"));
4354 drop(after_gate);
4355 }
4356
4357 #[tokio::test]
4358 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4359 let manager = lifecycle_test_manager();
4360 let current = manager.published_generation();
4361 manager
4362 .published_generation
4363 .store(Arc::new(PublishedIndexGeneration {
4364 publication_id: current.publication_id,
4365 schema: current.schema.clone(),
4366 trained_vectors: Some(Arc::new(TrainedVectorStructures {
4367 centroids: rustc_hash::FxHashMap::default(),
4368 binary_quantizers: rustc_hash::FxHashMap::default(),
4369 ..Default::default()
4370 })),
4371 }));
4372
4373 let guard = manager.begin_vector_artifact_update().await.unwrap();
4374 assert!(
4375 manager.trained().is_some(),
4376 "search readers keep the last fully validated generation"
4377 );
4378 assert!(
4379 manager.trained_for_segment_build().is_none(),
4380 "new segment producers must stay flat during an artifact update"
4381 );
4382
4383 let detached_transaction_guard = guard.clone();
4384 drop(guard);
4385 assert!(
4386 manager.trained_for_segment_build().is_none(),
4387 "a detached lifecycle transaction must retain the producer gate after request cancellation"
4388 );
4389 drop(detached_transaction_guard);
4390 assert!(manager.trained_for_segment_build().is_some());
4391 }
4392
4393 #[tokio::test]
4394 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4395 let manager = lifecycle_test_manager();
4396 let guard = manager.begin_vector_artifact_update().await.unwrap();
4397 assert!(
4398 manager
4399 .active_operations
4400 .try_register(vec!["merge".into()])
4401 .is_none(),
4402 "ordinary merge/reorder work must not change staged sources"
4403 );
4404 let indexing = manager
4405 .active_operations
4406 .try_register_indexing(vec!["fresh".into()])
4407 .expect("indexing remains available in flat mode");
4408 drop(indexing);
4409
4410 drop(guard);
4411 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4412 assert!(
4413 manager
4414 .active_operations
4415 .try_register(vec!["merge".into()])
4416 .is_some()
4417 );
4418 }
4419
4420 #[tokio::test]
4421 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4422 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4423 let guard = active.try_register(vec!["live".into()]).unwrap();
4424 let cancellation = active.cancellation_flag();
4425 active.stop_accepting();
4426 assert!(cancellation.load(Ordering::Acquire));
4427 assert!(active.try_register(vec!["new".into()]).is_none());
4428
4429 let waiter = {
4430 let active = Arc::clone(&active);
4431 tokio::spawn(async move { active.wait_until_idle().await })
4432 };
4433 tokio::task::yield_now().await;
4434 assert!(!waiter.is_finished());
4435 drop(guard);
4436 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4437 .await
4438 .expect("shutdown waiter missed the final guard notification")
4439 .unwrap();
4440 }
4441
4442 #[tokio::test]
4443 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4444 let manager = lifecycle_test_manager();
4445 let started = Arc::new(Semaphore::new(0));
4446 let release = Arc::new(Semaphore::new(0));
4447 let completed = Arc::new(AtomicBool::new(false));
4448
4449 let request = {
4450 let manager = Arc::clone(&manager);
4451 let started = Arc::clone(&started);
4452 let release = Arc::clone(&release);
4453 let completed = Arc::clone(&completed);
4454 tokio::spawn(async move {
4455 manager
4456 .run_lifecycle_transaction(async move {
4457 started.add_permits(1);
4458 let _permit = release.acquire().await.unwrap();
4459 completed.store(true, Ordering::Release);
4460 Ok(())
4461 })
4462 .await
4463 })
4464 };
4465
4466 let _started = started.acquire().await.unwrap();
4467 request.abort();
4468 assert!(request.await.unwrap_err().is_cancelled());
4469 release.add_permits(1);
4470
4471 manager.begin_shutdown();
4472 tokio::time::timeout(
4473 std::time::Duration::from_secs(1),
4474 manager.wait_for_shutdown(),
4475 )
4476 .await
4477 .expect("shutdown did not drain detached lifecycle transaction");
4478 assert!(completed.load(Ordering::Acquire));
4479 }
4480
4481 #[tokio::test]
4482 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4483 let manager = lifecycle_test_manager();
4484 {
4485 let mut state = manager.state.lock().await;
4486 state.metadata.add_segment_meta(
4487 "eligible".into(),
4488 SegmentMetaInfo {
4489 num_docs: 10,
4490 ancestors: Vec::new(),
4491 generation: 1,
4492 reordered: true,
4493 bp_converged: false,
4494 bp_unconverged_passes: 2,
4495 },
4496 );
4497 state.metadata.add_segment_meta(
4498 "at-limit".into(),
4499 SegmentMetaInfo {
4500 num_docs: 20,
4501 ancestors: Vec::new(),
4502 generation: 1,
4503 reordered: true,
4504 bp_converged: false,
4505 bp_unconverged_passes: 3,
4506 },
4507 );
4508 state.metadata.add_segment_meta(
4509 "carried-debt".into(),
4510 SegmentMetaInfo {
4511 num_docs: 15,
4512 ancestors: Vec::new(),
4513 generation: 2,
4514 reordered: false,
4515 bp_converged: false,
4516 bp_unconverged_passes: 2,
4517 },
4518 );
4519 state.metadata.add_segment_meta(
4520 "carried-debt-at-limit".into(),
4521 SegmentMetaInfo {
4522 num_docs: 25,
4523 ancestors: Vec::new(),
4524 generation: 2,
4525 reordered: false,
4526 bp_converged: false,
4527 bp_unconverged_passes: 3,
4528 },
4529 );
4530 state.metadata.add_segment_meta(
4531 "converged".into(),
4532 SegmentMetaInfo {
4533 num_docs: 30,
4534 ancestors: Vec::new(),
4535 generation: 1,
4536 reordered: true,
4537 bp_converged: true,
4538 bp_unconverged_passes: 0,
4539 },
4540 );
4541 state.metadata.add_segment("fresh".into(), 40);
4542 }
4543
4544 assert_eq!(
4545 manager.unreordered_segments().await,
4546 vec![("fresh".into(), 40)],
4547 "a block-copy output with BP debt is not a fresh first-pass candidate",
4548 );
4549 let mut eligible = manager.unconverged_segments_below(3).await;
4550 eligible.sort_unstable();
4551 assert_eq!(
4552 eligible,
4553 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4554 );
4555 assert!(manager.unconverged_segments_below(0).await.is_empty());
4556 }
4557
4558 #[test]
4559 fn merge_retry_backoff_is_exponential_and_capped() {
4560 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4561 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4562 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4563 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4564 }
4565
4566 #[test]
4567 fn only_deterministic_source_errors_are_quarantined() {
4568 assert!(is_deterministic_source_error(&Error::Corruption(
4569 "bad footer".into()
4570 )));
4571 assert!(is_deterministic_source_error(&Error::Io(
4572 std::io::Error::from(std::io::ErrorKind::NotFound)
4573 )));
4574 assert!(!is_deterministic_source_error(&Error::Io(
4575 std::io::Error::from(std::io::ErrorKind::TimedOut)
4576 )));
4577 assert!(!is_deterministic_source_error(&Error::Io(
4578 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4579 )));
4580 }
4581
4582 #[test]
4583 fn transient_reorder_failure_is_backed_off_until_cleared() {
4584 let manager = lifecycle_test_manager();
4585 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4586 assert!(manager.paused_reorder_segments().contains("source"));
4587 manager.clear_reorder_retry("source");
4588 assert!(!manager.paused_reorder_segments().contains("source"));
4589 }
4590
4591 #[derive(Default)]
4594 struct FailingExistsDirectory(crate::directories::RamDirectory);
4595
4596 #[async_trait::async_trait]
4597 impl crate::directories::Directory for FailingExistsDirectory {
4598 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4599 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4600 }
4601
4602 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4603 self.0.file_size(path).await
4604 }
4605
4606 async fn open_read(
4607 &self,
4608 path: &std::path::Path,
4609 ) -> std::io::Result<crate::directories::FileHandle> {
4610 self.0.open_read(path).await
4611 }
4612
4613 async fn read_range(
4614 &self,
4615 path: &std::path::Path,
4616 range: std::ops::Range<u64>,
4617 ) -> std::io::Result<crate::directories::OwnedBytes> {
4618 self.0.read_range(path, range).await
4619 }
4620
4621 async fn list_files(
4622 &self,
4623 prefix: &std::path::Path,
4624 ) -> std::io::Result<Vec<std::path::PathBuf>> {
4625 self.0.list_files(prefix).await
4626 }
4627
4628 async fn open_lazy(
4629 &self,
4630 path: &std::path::Path,
4631 ) -> std::io::Result<crate::directories::FileHandle> {
4632 self.0.open_lazy(path).await
4633 }
4634 }
4635
4636 #[async_trait::async_trait]
4637 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4638 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4639 self.0.write(path, data).await
4640 }
4641
4642 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4643 self.0.delete(path).await
4644 }
4645
4646 async fn rename(
4647 &self,
4648 from: &std::path::Path,
4649 to: &std::path::Path,
4650 ) -> std::io::Result<()> {
4651 self.0.rename(from, to).await
4652 }
4653
4654 async fn sync(&self) -> std::io::Result<()> {
4655 self.0.sync().await
4656 }
4657
4658 async fn streaming_writer(
4659 &self,
4660 path: &std::path::Path,
4661 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4662 self.0.streaming_writer(path).await
4663 }
4664 }
4665
4666 #[derive(Debug, Clone)]
4667 struct MergeEverythingPolicy;
4668
4669 impl MergePolicy for MergeEverythingPolicy {
4670 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4671 if segments.len() < 2 {
4672 return Vec::new();
4673 }
4674 vec![crate::merge::MergeCandidate {
4675 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4676 }]
4677 }
4678
4679 fn clone_box(&self) -> Box<dyn MergePolicy> {
4680 Box::new(self.clone())
4681 }
4682 }
4683
4684 #[tokio::test]
4685 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4686 let manager = lifecycle_test_manager();
4687 let parked_indexing = manager
4692 .protect_new_segment("00000000000000000000000000000abc".into())
4693 .unwrap();
4694
4695 let error = tokio::time::timeout(
4696 std::time::Duration::from_secs(2),
4697 manager.begin_vector_artifact_update(),
4698 )
4699 .await
4700 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4701 .err()
4702 .expect("an old-generation prepared segment must block artifact replacement")
4703 .to_string();
4704 assert!(error.contains("built but uncommitted"), "{error}");
4705 assert!(
4706 !manager.vector_artifact_update.load(Ordering::Acquire),
4707 "a rejected update must release the producer gate"
4708 );
4709
4710 drop(parked_indexing);
4711
4712 let guard = manager
4713 .begin_vector_artifact_update()
4714 .await
4715 .expect("artifact update should succeed after the pending generation is resolved");
4716 drop(guard);
4717 }
4718
4719 #[tokio::test]
4720 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4721 let manager = lifecycle_test_manager();
4722 let merge_like = manager
4723 .active_operations
4724 .try_register(vec!["merge-source".into()])
4725 .unwrap();
4726
4727 let waiter = {
4728 let manager = Arc::clone(&manager);
4729 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4730 };
4731 for _ in 0..8 {
4732 tokio::task::yield_now().await;
4733 }
4734 assert!(
4735 !waiter.is_finished(),
4736 "artifact update must drain merge/reorder producers that may hold the previous generation"
4737 );
4738
4739 drop(merge_like);
4740 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4741 .await
4742 .expect("artifact update missed the lifecycle guard release")
4743 .unwrap()
4744 .unwrap();
4745 }
4746
4747 #[tokio::test]
4748 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4749 let manager = lifecycle_test_manager();
4750 let release = Arc::new(Semaphore::new(0));
4751 let merge_task = {
4752 let release = Arc::clone(&release);
4753 tokio::spawn(async move {
4754 let _permit = release.acquire().await.unwrap();
4755 })
4756 };
4757 manager.merge_handles.lock().push(merge_task);
4758
4759 let waiter = {
4760 let manager = Arc::clone(&manager);
4761 tokio::spawn(async move { manager.wait_for_all_merges().await })
4762 };
4763 for _ in 0..8 {
4764 tokio::task::yield_now().await;
4765 }
4766 assert!(!waiter.is_finished());
4767 waiter.abort();
4770 let join_error = waiter.await.unwrap_err();
4771 assert!(join_error.is_cancelled());
4772
4773 assert!(
4774 !manager.merge_handles.lock().is_empty(),
4775 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4776 );
4777
4778 release.add_permits(1);
4780 tokio::time::timeout(
4781 std::time::Duration::from_secs(1),
4782 manager.wait_for_all_merges(),
4783 )
4784 .await
4785 .expect("subsequent drain missed the reinserted merge handle");
4786 assert!(manager.merge_handles.lock().is_empty());
4787 }
4788
4789 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4790 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4791 let manager = lifecycle_test_manager();
4792 {
4793 let mut state = manager.state.lock().await;
4794 state
4795 .metadata
4796 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4797 state
4798 .metadata
4799 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4800 }
4801 let reorder_like = manager
4804 .active_operations
4805 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4806 .unwrap();
4807
4808 let force_merge = {
4809 let manager = Arc::clone(&manager);
4810 tokio::spawn(async move { manager.force_merge().await })
4811 };
4812
4813 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4814 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4815 assert!(
4816 retries >= 1,
4817 "force_merge never observed the conflicting owner (retries={retries})"
4818 );
4819 assert!(
4820 retries < 20,
4821 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4822 );
4823
4824 drop(reorder_like);
4825 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4828 .await
4829 .expect("force_merge kept spinning after the conflicting owner released")
4830 .unwrap();
4831 assert!(result.is_err());
4832 }
4833
4834 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4835 async fn force_merge_routes_around_segments_held_by_reorder() {
4836 let manager = lifecycle_test_manager();
4837 {
4838 let mut state = manager.state.lock().await;
4839 state
4840 .metadata
4841 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4842 state
4843 .metadata
4844 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4845 state
4846 .metadata
4847 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4848 }
4849 let _reorder_like = manager
4852 .active_operations
4853 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4854 .unwrap();
4855
4856 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4863 let manager = Arc::clone(&manager);
4864 async move { manager.force_merge().await }
4865 })
4866 .await
4867 .expect("force_merge livelocked on a segment held by an active reorder");
4868 assert!(result.is_err(), "fake segment files must fail the merge");
4869
4870 assert_eq!(
4871 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4872 0,
4873 "batch built from the ownership snapshot must not collide with the held segment"
4874 );
4875 }
4876
4877 #[tokio::test]
4878 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4879 let schema = crate::dsl::SchemaBuilder::default().build();
4880 let mut metadata = IndexMetadata::new(schema.clone());
4881 metadata.add_segment("00000000000000000000000000000001".into(), 10);
4882 metadata.add_segment("00000000000000000000000000000002".into(), 10);
4883 let manager = Arc::new(SegmentManager::new(
4884 Arc::new(FailingExistsDirectory::default()),
4885 Arc::new(schema),
4886 metadata,
4887 Box::new(MergeEverythingPolicy),
4888 0,
4889 1,
4890 Arc::new(Semaphore::new(1)),
4891 None,
4892 1024,
4893 Arc::new(ReorderConcurrencyGate::new(1)),
4894 None,
4895 ));
4896
4897 manager.maybe_merge().await;
4900
4901 tokio::time::timeout(
4902 std::time::Duration::from_secs(5),
4903 manager.wait_for_all_merges(),
4904 )
4905 .await
4906 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4907 assert!(
4908 manager.merge_retry_is_paused(),
4909 "the failed merge should have armed the retry backoff"
4910 );
4911
4912 manager.begin_shutdown();
4914 tokio::time::timeout(
4915 std::time::Duration::from_secs(5),
4916 manager.wait_for_shutdown(),
4917 )
4918 .await
4919 .expect("shutdown did not drain the merge retry wakeup task");
4920 }
4921}