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 term_cache_blocks: usize,
872 merge_permits: Arc<Semaphore>,
876 global_merge_permits: Arc<Semaphore>,
878 reorder_permits: Arc<ReorderConcurrencyGate>,
882 reorder_on_merge: bool,
887 merge_bp_time_budget: Option<std::time::Duration>,
891 bp_memory_budget_bytes: usize,
894 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
897 replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
901}
902
903struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
904
905impl Drop for ForceMergeActivityGuard<'_> {
906 fn drop(&mut self) {
907 self.0.fetch_sub(1, Ordering::AcqRel);
908 }
909}
910
911impl<D: DirectoryWriter + 'static> SegmentManager<D> {
912 #[allow(clippy::too_many_arguments)]
914 pub fn new(
915 directory: Arc<D>,
916 schema: Arc<crate::dsl::Schema>,
917 metadata: IndexMetadata,
918 merge_policy: Box<dyn MergePolicy>,
919 term_cache_blocks: usize,
920 max_concurrent_merges: usize,
921 global_merge_permits: Arc<Semaphore>,
922 merge_bp_time_budget: Option<std::time::Duration>,
923 bp_memory_budget_bytes: usize,
924 reorder_permits: Arc<ReorderConcurrencyGate>,
925 background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
926 ) -> Self {
927 let reorder_on_merge = schema.reorder_on_merge();
930 if reorder_on_merge {
931 log::info!(
932 "[merge] index={} reorder-on-merge enabled by index schema",
933 schema.index_label()
934 );
935 }
936
937 let tracker = Arc::new(SegmentTracker::new());
938 for seg_id in metadata.segment_metas.keys() {
939 tracker.register(seg_id);
940 }
941
942 let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
943 Arc::new(parking_lot::Mutex::new(Vec::new()));
944 let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
945 let dir = Arc::clone(&directory);
946 let tracker = Arc::clone(&tracker);
947 let lifecycle_handles = Arc::clone(&lifecycle_handles);
948 let cleanup_index_label: Arc<str> = schema.index_label().into();
949 Arc::new(move |segment_ids| {
950 let Ok(handle) = tokio::runtime::Handle::try_current() else {
953 tracker.complete_deletion(&segment_ids);
956 return;
957 };
958 let dir = Arc::clone(&dir);
959 let task_tracker = Arc::clone(&tracker);
960 let task_index_label = Arc::clone(&cleanup_index_label);
961 let cleanup_ids = segment_ids.clone();
962 let future = async move {
963 for &segment_id in &segment_ids {
964 log::info!(
965 "[segment_cleanup] index={} deleting deferred segment {}",
966 task_index_label,
967 segment_id.to_hex()
968 );
969 if let Err(error) =
970 crate::segment::delete_segment(dir.as_ref(), segment_id).await
971 {
972 log::warn!(
973 "[segment_cleanup] index={} deferred delete failed for {}: {}",
974 task_index_label,
975 segment_id.to_hex(),
976 error,
977 );
978 }
979 }
980 task_tracker.complete_deletion(&segment_ids);
981 };
982 if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
983 tracker.complete_deletion(&cleanup_ids);
987 log::warn!(
988 "[segment_cleanup] index={} runtime rejected deferred deletion; files will be swept later",
989 cleanup_index_label
990 );
991 }
992 })
993 };
994
995 let initial_generation = Arc::new(PublishedIndexGeneration {
996 publication_id: metadata.publication_generation,
997 schema: Arc::clone(&schema),
998 trained_vectors: None,
999 });
1000 Self {
1001 state: Arc::new(AsyncMutex::new(ManagerState {
1002 metadata,
1003 merge_policy,
1004 })),
1005 active_operations: Arc::new(ActiveSegmentOperations::new(schema.index_label().into())),
1006 quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
1007 merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
1008 reorder_retries: parking_lot::Mutex::new(HashMap::new()),
1009 merge_handles: parking_lot::Mutex::new(Vec::new()),
1010 global_merge_wakeup_pending: AtomicBool::new(false),
1011 force_merge_active: AtomicUsize::new(0),
1012 #[cfg(test)]
1013 force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
1014 lifecycle_handles,
1015 published_generation: Arc::new(ArcSwap::new(initial_generation)),
1016 vector_artifact_update: Arc::new(AtomicBool::new(false)),
1017 tracker,
1018 delete_fn,
1019 directory,
1020 schema,
1021 term_cache_blocks,
1022 merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
1023 global_merge_permits,
1024 reorder_permits,
1025 reorder_on_merge,
1026 merge_bp_time_budget,
1027 bp_memory_budget_bytes,
1028 background_reorder_pool,
1029 replacement_refresh: parking_lot::RwLock::new(None),
1030 }
1031 }
1032
1033 pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1034 where
1035 F: Fn() -> Fut + Send + Sync + 'static,
1036 Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1037 {
1038 *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1039 }
1040
1041 pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1047 if let Some(pool) = &self.background_reorder_pool {
1048 return Arc::clone(pool);
1049 }
1050 Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1051 let threads = (num_cpus::get() / 2).max(1);
1052 log::info!(
1053 "[merge] process-wide background CPU pool: {} thread(s)",
1054 threads
1055 );
1056 Arc::new(
1057 rayon::ThreadPoolBuilder::new()
1058 .num_threads(threads)
1059 .thread_name(|i| format!("hermes-bg-cpu-{}", i))
1060 .build()
1061 .expect("failed to build background CPU pool"),
1062 )
1063 }))
1064 }
1065
1066 pub fn begin_shutdown(&self) {
1070 self.active_operations.stop_accepting();
1071 }
1072
1073 async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1081 where
1082 T: Send + 'static,
1083 F: std::future::Future<Output = Result<T>> + Send + 'static,
1084 {
1085 let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1086 let future = async move {
1087 let result = transaction.await;
1088 let _ = result_tx.send(result);
1089 };
1090 let runtime = tokio::runtime::Handle::current();
1091 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1092 return Err(Error::Internal(
1093 "runtime rejected lifecycle metadata transaction".into(),
1094 ));
1095 }
1096 result_rx.await.map_err(|_| {
1097 Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1098 })?
1099 }
1100
1101 fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1103 let manager = Arc::clone(self);
1104 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1105 let Ok(handle) = tokio::runtime::Handle::try_current() else {
1106 log::warn!(
1107 "[segment_cleanup] index={} runtime unavailable; partial output {} will be swept on startup",
1108 manager.schema.index_label(),
1109 segment_id.to_hex(),
1110 );
1111 return;
1112 };
1113
1114 let cleanup_manager = Arc::clone(&manager);
1115 let future = async move {
1116 cleanup_manager
1117 .delete_output_if_unregistered(segment_id, "task unwind")
1118 .await;
1119 };
1120 if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1121 log::warn!(
1122 "[segment_cleanup] index={} runtime rejected output cleanup; {} will be swept on startup",
1123 manager.schema.index_label(),
1124 segment_id.to_hex(),
1125 );
1126 }
1127 });
1128
1129 OutputCleanupGuard::new(output_id, cleanup)
1130 }
1131
1132 pub(crate) fn schedule_unpublished_segment_cleanup(
1137 self: &Arc<Self>,
1138 output_id: SegmentId,
1139 operation: SegmentOperationGuard,
1140 runtime: tokio::runtime::Handle,
1141 ) {
1142 let manager = Arc::clone(self);
1143 let output_hex = output_id.to_hex();
1144 let future = async move {
1145 manager
1146 .delete_output_if_unregistered(output_id, "indexing abort or failure")
1147 .await;
1148 drop(operation);
1149 };
1150 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1151 log::warn!(
1154 "[segment_cleanup] index={} runtime unavailable; indexing output {} will be swept on startup",
1155 self.schema.index_label(),
1156 output_hex,
1157 );
1158 }
1159 }
1160
1161 pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1167 match self
1168 .active_operations
1169 .try_register_indexing(vec![segment_id.clone()])
1170 {
1171 Some(operation) => Ok(operation),
1172 None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1173 None => Err(Error::Corruption(format!(
1174 "new segment ID {} is already owned by an active operation",
1175 segment_id
1176 ))),
1177 }
1178 }
1179
1180 async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1184 let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1185 Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1186 })?;
1187 let files = SegmentFiles::new(id.0);
1188
1189 for path in files.mandatory_paths() {
1190 if !self.directory.exists(path).await.map_err(Error::Io)? {
1191 return Err(Error::Corruption(format!(
1192 "segment {} cannot be published: mandatory file {:?} is missing",
1193 segment_id, path
1194 )));
1195 }
1196 }
1197
1198 let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1199 Error::Corruption(format!(
1200 "segment {} cannot be published: missing/unreadable {:?}: {}",
1201 segment_id, files.meta, e
1202 ))
1203 })?;
1204 let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1205 Error::Corruption(format!(
1206 "segment {} cannot be published: failed reading {:?}: {}",
1207 segment_id, files.meta, e
1208 ))
1209 })?;
1210 let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1211 Error::Corruption(format!(
1212 "segment {} cannot be published: invalid {:?}: {}",
1213 segment_id, files.meta, e
1214 ))
1215 })?;
1216
1217 if meta.id != id.0 || meta.num_docs != expected_docs {
1218 return Err(Error::Corruption(format!(
1219 "segment {} cannot be published: metadata identity/docs mismatch \
1220 (id={:032x}, docs={}, expected_docs={})",
1221 segment_id, meta.id, meta.num_docs, expected_docs
1222 )));
1223 }
1224
1225 Ok(())
1226 }
1227
1228 fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1229 let inserted = self
1230 .quarantined_segments
1231 .lock()
1232 .insert(segment_id.to_string());
1233 if inserted {
1234 log::error!(
1235 "[merge] index={} quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1236 It remains metadata-live for explicit repair but is excluded from merges until restart",
1237 self.schema.index_label(),
1238 segment_id,
1239 error,
1240 );
1241 }
1242 }
1243
1244 fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1245 let mut retry = self.merge_retry.lock();
1246 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1247 let delay = merge_retry_delay(retry.consecutive_failures);
1248 retry.retry_after = std::time::Instant::now().checked_add(delay);
1249 log::warn!(
1250 "[merge] index={} pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1251 self.schema.index_label(),
1252 delay.as_secs_f64(),
1253 retry.consecutive_failures,
1254 error,
1255 );
1256 delay
1257 }
1258
1259 fn clear_merge_retry_backoff(&self) {
1260 *self.merge_retry.lock() = MergeRetryState::default();
1261 }
1262
1263 fn merge_retry_is_paused(&self) -> bool {
1264 let mut retry = self.merge_retry.lock();
1265 match retry.retry_after {
1266 Some(deadline) if deadline > std::time::Instant::now() => true,
1267 Some(_) => {
1268 retry.retry_after = None;
1269 false
1270 }
1271 None => false,
1272 }
1273 }
1274
1275 fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1276 let mut retries = self.reorder_retries.lock();
1277 let retry = retries.entry(segment_id.to_string()).or_default();
1278 retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1279 let delay = merge_retry_delay(retry.consecutive_failures);
1280 retry.retry_after = std::time::Instant::now().checked_add(delay);
1281 log::warn!(
1282 "[reorder] index={} pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1283 self.schema.index_label(),
1284 segment_id,
1285 delay.as_secs_f64(),
1286 retry.consecutive_failures,
1287 error,
1288 );
1289 }
1290
1291 fn clear_reorder_retry(&self, segment_id: &str) {
1292 self.reorder_retries.lock().remove(segment_id);
1293 }
1294
1295 fn paused_reorder_segments(&self) -> HashSet<String> {
1296 let now = std::time::Instant::now();
1297 let mut retries = self.reorder_retries.lock();
1298 let mut paused = HashSet::new();
1299 for (segment_id, retry) in retries.iter_mut() {
1300 match retry.retry_after {
1301 Some(deadline) if deadline > now => {
1302 paused.insert(segment_id.clone());
1303 }
1304 Some(_) => retry.retry_after = None,
1305 None => {}
1306 }
1307 }
1308 paused
1309 }
1310
1311 fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1315 if self
1316 .global_merge_wakeup_pending
1317 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1318 .is_err()
1319 {
1320 return;
1321 }
1322
1323 let manager = Arc::clone(self);
1324 let future = async move {
1325 let capacity = tokio::select! {
1326 biased;
1327 () = manager.active_operations.wait_for_shutdown() => None,
1328 permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1329 };
1330
1331 manager
1332 .global_merge_wakeup_pending
1333 .store(false, Ordering::Release);
1334 if let Some(permit) = capacity {
1335 drop(permit);
1339 manager.maybe_merge().await;
1340 }
1341 };
1342 let runtime = tokio::runtime::Handle::current();
1343 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1344 self.global_merge_wakeup_pending
1345 .store(false, Ordering::Release);
1346 log::warn!(
1347 "[merge] index={} runtime rejected global-capacity wakeup task",
1348 self.schema.index_label()
1349 );
1350 }
1351 }
1352
1353 #[cfg(test)]
1354 pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1355 self.quarantined_segments.lock().contains(segment_id)
1356 }
1357
1358 async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1363 let output_hex = output_id.to_hex();
1364 {
1365 let st = self.state.lock().await;
1366 if st.metadata.has_segment(&output_hex) {
1367 return;
1368 }
1369 }
1370
1371 log::info!(
1375 "[segment_cleanup] index={} deleting uncommitted output {} after {}",
1376 self.schema.index_label(),
1377 output_hex,
1378 reason,
1379 );
1380 if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1381 {
1382 log::warn!(
1383 "[segment_cleanup] index={} failed deleting uncommitted output {}: {}",
1384 self.schema.index_label(),
1385 output_hex,
1386 error,
1387 );
1388 }
1389 }
1390
1391 pub async fn get_segment_ids(&self) -> Vec<String> {
1397 self.state.lock().await.metadata.segment_ids()
1398 }
1399
1400 pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1402 self.published_generation.load().trained_vectors.clone()
1403 }
1404
1405 pub(crate) fn published_generation(&self) -> Arc<PublishedIndexGeneration> {
1408 self.published_generation.load_full()
1409 }
1410
1411 pub(crate) fn publication_id(&self) -> u64 {
1412 self.published_generation.load().publication_id
1413 }
1414
1415 pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1422 if self.vector_artifact_update.load(Ordering::Acquire) {
1423 return None;
1424 }
1425 let trained = self.published_generation.load().trained_vectors.clone();
1426 if self.vector_artifact_update.load(Ordering::Acquire) {
1427 None
1428 } else {
1429 trained
1430 }
1431 }
1432
1433 pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1452 self.vector_artifact_update
1453 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1454 .map_err(|_| {
1455 Error::Internal("a trained-vector artifact update is already in progress".into())
1456 })?;
1457 self.active_operations.pause_non_indexing();
1458 let guard = VectorArtifactUpdateGuard {
1459 _lease: Arc::new(VectorArtifactUpdateLease {
1460 updating: Arc::clone(&self.vector_artifact_update),
1461 active_operations: Arc::clone(&self.active_operations),
1462 }),
1463 };
1464 let (preexisting, parked_indexing) =
1465 self.active_operations.draining_operation_tokens_snapshot();
1466 if parked_indexing > 0 {
1467 return Err(Error::Internal(format!(
1468 "cannot update trained-vector artifacts while {parked_indexing} indexing \
1469 segment(s) are built but uncommitted; commit or abort the pending \
1470 generation and retry"
1471 )));
1472 }
1473 self.active_operations
1474 .wait_until_operations_finish(&preexisting)
1475 .await;
1476 Ok(guard)
1477 }
1478
1479 pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1482 let (vector_fields, schema, publication_id) = {
1484 let st = self.state.lock().await;
1485 (
1486 st.metadata.vector_fields.clone(),
1487 Arc::new(st.metadata.schema.clone()),
1488 st.metadata.publication_generation,
1489 )
1490 };
1491 let trained = IndexMetadata::try_load_trained_from_fields(
1493 &vector_fields,
1494 schema.as_ref(),
1495 self.directory.as_ref(),
1496 )
1497 .await?
1498 .map(Arc::new);
1499 self.published_generation
1503 .store(Arc::new(PublishedIndexGeneration {
1504 publication_id,
1505 schema,
1506 trained_vectors: trained,
1507 }));
1508 Ok(())
1509 }
1510
1511 pub(crate) async fn publish_vector_generation(
1518 self: &Arc<Self>,
1519 artifact_update: &VectorArtifactUpdateGuard,
1520 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1521 next_trained: Arc<TrainedVectorStructures>,
1522 staged: Vec<StagedVectorSegment>,
1523 ) -> Result<()> {
1524 let schema = self.published_generation().schema.clone();
1525 self.publish_vector_generation_with_schema(
1526 artifact_update,
1527 schema,
1528 vector_fields,
1529 Some(next_trained),
1530 staged,
1531 )
1532 .await
1533 }
1534
1535 pub(crate) async fn publish_vector_generation_with_schema(
1536 self: &Arc<Self>,
1537 artifact_update: &VectorArtifactUpdateGuard,
1538 schema: Arc<crate::dsl::Schema>,
1539 vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1540 next_trained: Option<Arc<TrainedVectorStructures>>,
1541 mut staged: Vec<StagedVectorSegment>,
1542 ) -> Result<()> {
1543 if !self.vector_artifact_update.load(Ordering::Acquire) {
1544 return Err(Error::Internal(
1545 "vector generation publication lost its exclusive update lease".into(),
1546 ));
1547 }
1548
1549 for replacement in &staged {
1550 self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1551 .await?;
1552 }
1553
1554 let mut st = Arc::clone(&self.state).lock_owned().await;
1555 let mut next = st.metadata.clone();
1556 next.schema = (*schema).clone();
1557 next.vector_fields = vector_fields;
1558 next.refresh_total_vectors();
1559
1560 for replacement in &staged {
1561 let source_info = next
1562 .segment_metas
1563 .remove(&replacement.source_id)
1564 .ok_or_else(|| {
1565 Error::Corruption(format!(
1566 "vector generation source {} disappeared before publication",
1567 replacement.source_id,
1568 ))
1569 })?;
1570 let output_hex = replacement.output_id.to_hex();
1571 if next.segment_metas.contains_key(&output_hex) {
1572 return Err(Error::Corruption(format!(
1573 "vector generation output {output_hex} is already metadata-live"
1574 )));
1575 }
1576 next.add_segment_meta(output_hex, source_info);
1579 }
1580
1581 let directory = Arc::clone(&self.directory);
1582 let published_generation = Arc::clone(&self.published_generation);
1583 let tracker = Arc::clone(&self.tracker);
1584 let replacement_refresh = self.replacement_refresh.read().clone();
1585 let artifact_update = artifact_update.clone();
1589 let index_label = self.schema.index_label().to_owned();
1590 next.publication_generation =
1591 next.publication_generation.checked_add(1).ok_or_else(|| {
1592 Error::Corruption("vector publication generation exhausted u64".into())
1593 })?;
1594 let next_schema = schema;
1595 let next_publication_id = next.publication_generation;
1596 self.run_lifecycle_transaction(async move {
1597 let _artifact_update = artifact_update;
1598 next.save(directory.as_ref()).await?;
1599
1600 for replacement in &staged {
1601 tracker.register(&replacement.output_id.to_hex());
1602 }
1603 st.metadata = next;
1604 published_generation.store(Arc::new(PublishedIndexGeneration {
1605 publication_id: next_publication_id,
1606 schema: next_schema,
1607 trained_vectors: next_trained,
1608 }));
1609
1610 for replacement in &mut staged {
1613 replacement.cleanup.disarm();
1614 }
1615 let retired = staged
1616 .iter()
1617 .map(|replacement| replacement.source_id.clone())
1618 .collect::<Vec<_>>();
1619 let ready_to_delete = tracker.mark_for_deletion(&retired);
1620 drop(st);
1621 for &segment_id in &ready_to_delete {
1622 if let Err(error) =
1623 crate::segment::delete_segment(directory.as_ref(), segment_id).await
1624 {
1625 log::warn!(
1626 "[segment_cleanup] index={index_label} immediate dense-vector generation delete failed for {}: {}",
1627 segment_id.to_hex(),
1628 error,
1629 );
1630 }
1631 }
1632 tracker.complete_deletion(&ready_to_delete);
1633 refresh_replacement_topology(replacement_refresh, &index_label).await;
1634 Ok(())
1635 })
1636 .await
1637 }
1638
1639 pub(crate) async fn publish_vector_schema_only(
1643 self: &Arc<Self>,
1644 artifact_update: &VectorArtifactUpdateGuard,
1645 schema: Arc<crate::dsl::Schema>,
1646 ) -> Result<()> {
1647 let vector_fields = self
1648 .read_metadata(|metadata| metadata.vector_fields.clone())
1649 .await;
1650 let trained = self.published_generation().trained_vectors.clone();
1651 self.publish_vector_generation_with_schema(
1652 artifact_update,
1653 schema,
1654 vector_fields,
1655 trained,
1656 Vec::new(),
1657 )
1658 .await
1659 }
1660
1661 pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1663 where
1664 F: FnOnce(&IndexMetadata) -> R,
1665 {
1666 let st = self.state.lock().await;
1667 f(&st.metadata)
1668 }
1669
1670 pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1672 where
1673 F: FnOnce(&mut IndexMetadata),
1674 {
1675 let mut st = Arc::clone(&self.state).lock_owned().await;
1676 let mut next = st.metadata.clone();
1677 f(&mut next);
1678 let directory = Arc::clone(&self.directory);
1679 self.run_lifecycle_transaction(async move {
1680 next.save(directory.as_ref()).await?;
1681 st.metadata = next;
1682 Ok(())
1683 })
1684 .await
1685 }
1686
1687 pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1690 let (acquired, generation) = {
1691 let st = self.state.lock().await;
1692 let segment_ids = st.metadata.segment_ids();
1693 (
1694 self.tracker.acquire(&segment_ids),
1695 self.published_generation.load_full(),
1696 )
1697 };
1698
1699 SegmentSnapshot::with_generation(
1700 Arc::clone(&self.tracker),
1701 acquired,
1702 generation,
1703 Arc::clone(&self.delete_fn),
1704 )
1705 }
1706
1707 pub fn tracker(&self) -> Arc<SegmentTracker> {
1709 Arc::clone(&self.tracker)
1710 }
1711
1712 pub fn directory(&self) -> Arc<D> {
1714 Arc::clone(&self.directory)
1715 }
1716}
1717
1718#[cfg(feature = "native")]
1723impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1724 pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1726 for (segment_id, num_docs) in new_segments {
1729 self.validate_completed_segment(segment_id, *num_docs)
1730 .await?;
1731 }
1732
1733 let mut st = Arc::clone(&self.state).lock_owned().await;
1734 let mut next = st.metadata.clone();
1735 let mut added = Vec::new();
1736 for (segment_id, num_docs) in new_segments {
1737 if !next.has_segment(segment_id) {
1738 next.add_segment(segment_id.clone(), *num_docs);
1739 added.push(segment_id.clone());
1740 }
1741 }
1742
1743 let directory = Arc::clone(&self.directory);
1749 let tracker = Arc::clone(&self.tracker);
1750 self.run_lifecycle_transaction(async move {
1751 next.save(directory.as_ref()).await?;
1752 for segment_id in &added {
1753 tracker.register(segment_id);
1754 }
1755 st.metadata = next;
1756 Ok(())
1757 })
1758 .await
1759 }
1760
1761 pub async fn maybe_merge(self: &Arc<Self>) {
1772 if !self.active_operations.is_accepting() {
1773 log::debug!(
1774 "[maybe_merge] index={} manager is shutting down, skipping",
1775 self.schema.index_label()
1776 );
1777 return;
1778 }
1779 if self.merge_retry_is_paused() {
1780 log::debug!(
1781 "[maybe_merge] index={} retry backoff active, skipping",
1782 self.schema.index_label()
1783 );
1784 return;
1785 }
1786
1787 {
1790 let mut handles = self.merge_handles.lock();
1791 handles.retain(|h| !h.is_finished());
1792 }
1793 let local_slots = self.merge_permits.available_permits();
1794 let global_slots = self.global_merge_permits.available_permits();
1795 let slots_available = local_slots.min(global_slots);
1796
1797 {
1801 let st = self.state.lock().await;
1802 let quarantined = self.quarantined_segments.lock().clone();
1803 let active_ids = self.active_operations.snapshot();
1804
1805 let live_segments: Vec<SegmentInfo> = st
1810 .metadata
1811 .segment_metas
1812 .iter()
1813 .filter(|(id, _)| {
1814 !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1815 })
1816 .map(|(id, info)| SegmentInfo {
1817 id: id.clone(),
1818 num_docs: info.num_docs,
1819 })
1820 .collect();
1821 let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1822
1823 let segments: Vec<SegmentInfo> = live_segments
1826 .iter()
1827 .filter(|segment| !active_ids.contains(&segment.id))
1828 .cloned()
1829 .collect();
1830
1831 log::debug!(
1832 "[maybe_merge] index={} {} eligible segments",
1833 self.schema.index_label(),
1834 segments.len()
1835 );
1836
1837 let candidates = st.merge_policy.find_merges(&segments);
1838
1839 if candidates.is_empty() {
1840 return;
1841 }
1842
1843 if slots_available == 0 {
1847 if local_slots > 0 && global_slots == 0 {
1848 self.schedule_global_merge_wakeup();
1849 }
1850 log::debug!(
1851 "[maybe_merge] index={} at max concurrent merges, skipping",
1852 self.schema.index_label()
1853 );
1854 return;
1855 }
1856
1857 log::debug!(
1858 "[maybe_merge] index={} {} merge candidates, {} slots available",
1859 self.schema.index_label(),
1860 candidates.len(),
1861 slots_available
1862 );
1863
1864 let mut handles = Vec::new();
1865 for c in candidates {
1866 if handles.len() >= slots_available {
1867 break;
1868 }
1869 let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1875 if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1876 handles.push(h);
1877 }
1878 }
1879 if !handles.is_empty() {
1880 if severe_backlog && self.reorder_on_merge {
1881 log::info!(
1882 "[maybe_merge] index={} severe backlog: {} live segments; started {} fast \
1883 block-copy merge(s), deferring BP to the optimizer",
1884 self.schema.index_label(),
1885 live_segments.len(),
1886 handles.len(),
1887 );
1888 }
1889 self.merge_handles.lock().extend(handles);
1894 }
1895 }
1896 }
1897
1898 fn spawn_merge(
1907 self: &Arc<Self>,
1908 segment_ids_to_merge: Vec<String>,
1909 reorder_bmp: bool,
1910 ) -> Option<JoinHandle<()>> {
1911 if self.force_merge_active.load(Ordering::Acquire) > 0 {
1912 log::debug!(
1913 "[spawn_merge] index={} skipped: explicit force merge has priority",
1914 self.schema.index_label()
1915 );
1916 return None;
1917 }
1918 let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1919 Ok(permit) => permit,
1920 Err(_) => {
1921 log::debug!(
1922 "[spawn_merge] index={} skipped: global merge capacity is full",
1923 self.schema.index_label()
1924 );
1925 self.schedule_global_merge_wakeup();
1926 return None;
1927 }
1928 };
1929 let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1930 Ok(permit) => permit,
1931 Err(_) => {
1932 log::debug!(
1933 "[spawn_merge] index={} skipped: no merge permit available",
1934 self.schema.index_label()
1935 );
1936 return None;
1937 }
1938 };
1939 let output_id = SegmentId::new();
1940 let output_hex = output_id.to_hex();
1941
1942 let mut all_ids = segment_ids_to_merge.clone();
1943 all_ids.push(output_hex);
1944
1945 let guard = match self.active_operations.try_register(all_ids) {
1946 Some(g) => g,
1947 None => {
1948 log::debug!(
1949 "[spawn_merge] index={} skipped: segments overlap with an active operation",
1950 self.schema.index_label()
1951 );
1952 return None;
1953 }
1954 };
1955
1956 let sm = Arc::clone(self);
1957 let ids = segment_ids_to_merge;
1958
1959 let index_label = self.schema.index_label().to_owned();
1960 Some(tokio::spawn(async move {
1961 let mut reevaluate = false;
1962 let mut retry_delay = None;
1963
1964 let result = sm
1965 .merge_and_replace_registered(
1966 &ids,
1967 output_id,
1968 reorder_bmp,
1969 ReorderPriority::AutomaticMerge,
1970 )
1971 .await;
1972
1973 match result {
1974 Ok(_) => {
1975 sm.clear_merge_retry_backoff();
1976 reevaluate = true;
1977 }
1978 Err(MergeTaskError {
1979 error: Error::IndexClosed,
1980 ..
1981 }) => {
1982 log::debug!(
1983 "[merge] index={index_label} background merge for segments {:?} cancelled during shutdown",
1984 ids,
1985 );
1986 }
1987 Err(MergeTaskError {
1988 error,
1989 unavailable_segments,
1990 }) => {
1991 log::error!(
1992 "[merge] index={index_label} background merge failed for segments {:?}: {}",
1993 ids,
1994 error
1995 );
1996 if !unavailable_segments.is_empty() {
1997 reevaluate = true;
2001 } else {
2002 retry_delay = Some(sm.pause_merge_retries(&error));
2003 }
2004 }
2005 }
2006 drop(guard);
2009 drop(merge_permit);
2011 drop(global_merge_permit);
2012
2013 if reevaluate {
2014 sm.maybe_merge().await;
2015 } else if let Some(retry_delay) = retry_delay {
2016 sm.schedule_merge_retry_wakeup(retry_delay);
2023 }
2024 }))
2025 }
2026
2027 async fn merge_and_replace_registered(
2034 self: &Arc<Self>,
2035 ids: &[String],
2036 output_id: SegmentId,
2037 reorder_bmp: bool,
2038 priority: ReorderPriority,
2039 ) -> MergeTaskResult<(String, u32, bool)> {
2040 let mut output_cleanup = self.output_cleanup_guard(output_id);
2041 let generation = self.published_generation();
2042 let trained = self.trained_for_segment_build();
2043 let granularity = if reorder_bmp {
2044 self.merge_granularity(ids).await
2045 } else {
2046 crate::segment::reorder::BpGranularity::Auto
2047 };
2048 let result = Self::do_merge(
2049 self.directory.as_ref(),
2050 &generation.schema,
2051 ids,
2052 output_id,
2053 self.term_cache_blocks,
2054 trained.as_deref(),
2055 reorder_bmp,
2056 granularity,
2057 self.merge_bp_time_budget,
2058 self.bp_memory_budget_bytes,
2059 Arc::clone(&self.reorder_permits),
2060 priority,
2061 self.active_operations.cancellation_flag(),
2062 Some(self.background_cpu_pool()),
2063 )
2064 .await;
2065
2066 let (new_id, doc_count, bp_converged) = match result {
2067 Ok(value) => value,
2068 Err(error) => {
2069 for segment_id in &error.unavailable_segments {
2070 self.quarantine_segment(segment_id, &error.error);
2071 }
2072 self.delete_output_if_unregistered(output_id, "merge failure")
2073 .await;
2074 output_cleanup.disarm();
2075 return Err(error);
2076 }
2077 };
2078
2079 let layout = if reorder_bmp {
2080 ReplacementLayout::BpReordered {
2081 converged: bp_converged,
2082 }
2083 } else {
2084 ReplacementLayout::BlockCopy
2085 };
2086 if let Err(error) = self
2087 .replace_segments(ids, new_id.clone(), doc_count, layout)
2088 .await
2089 {
2090 self.delete_output_if_unregistered(output_id, "replacement failure")
2091 .await;
2092 output_cleanup.disarm();
2093 return Err(MergeTaskError::from(error));
2094 }
2095 output_cleanup.disarm();
2096 Ok((new_id, doc_count, bp_converged))
2097 }
2098
2099 fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
2103 let manager = Arc::clone(self);
2104 let future = async move {
2105 tokio::select! {
2106 () = tokio::time::sleep(retry_delay) => {
2107 manager.maybe_merge().await;
2108 }
2109 () = manager.active_operations.wait_for_shutdown() => {}
2110 }
2111 };
2112 let runtime = tokio::runtime::Handle::current();
2113 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2114 log::warn!(
2115 "[merge] index={} runtime rejected merge-retry wakeup task; eligible segments may stay \
2116 unmerged until the next commit re-runs merge policy evaluation",
2117 self.schema.index_label()
2118 );
2119 }
2120 }
2121
2122 async fn replace_segments(
2126 self: &Arc<Self>,
2127 old_ids: &[String],
2128 new_id: String,
2129 doc_count: u32,
2130 layout: ReplacementLayout,
2131 ) -> Result<()> {
2132 self.validate_completed_segment(&new_id, doc_count).await?;
2135 let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
2136 Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
2137 })?;
2138 let output_reader = SegmentReader::open(
2139 self.directory.as_ref(),
2140 output_id,
2141 self.published_generation().schema.clone(),
2142 self.term_cache_blocks,
2143 )
2144 .await
2145 .map_err(|error| match error {
2146 Error::Io(_) | Error::IndexClosed => error,
2150 error => Error::Corruption(format!(
2151 "replacement segment {new_id} failed full reader validation: {error}"
2152 )),
2153 })?;
2154 if output_reader.num_docs() != doc_count {
2155 return Err(Error::Corruption(format!(
2156 "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2157 output_reader.num_docs(),
2158 )));
2159 }
2160 drop(output_reader);
2161
2162 let mut st = Arc::clone(&self.state).lock_owned().await;
2163 let missing: Vec<&String> = old_ids
2167 .iter()
2168 .filter(|id| !st.metadata.has_segment(id))
2169 .collect();
2170 if !missing.is_empty() {
2171 return Err(Error::Corruption(format!(
2172 "replace_segments: source segment(s) {:?} not in metadata — \
2173 refusing to add output {} (would duplicate documents)",
2174 missing, new_id
2175 )));
2176 }
2177
2178 let replacement_info = match layout {
2179 ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2180 let generation = old_ids
2181 .iter()
2182 .filter_map(|id| st.metadata.segment_metas.get(id))
2183 .map(|info| info.generation)
2184 .max()
2185 .unwrap_or(0)
2186 .checked_add(1)
2187 .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2188 let parent_unconverged_passes = old_ids
2189 .iter()
2190 .filter_map(|id| st.metadata.segment_metas.get(id))
2191 .map(|info| info.bp_unconverged_passes)
2192 .max()
2193 .unwrap_or(0);
2194 let parent_has_debt = old_ids
2195 .iter()
2196 .filter_map(|id| st.metadata.segment_metas.get(id))
2197 .any(|info| !info.bp_converged);
2198 let (reordered, bp_converged, bp_unconverged_passes) =
2199 replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2200 SegmentMetaInfo {
2201 num_docs: doc_count,
2202 ancestors: old_ids.to_vec(),
2203 generation,
2204 reordered,
2205 bp_converged,
2206 bp_unconverged_passes,
2207 }
2208 }
2209 ReplacementLayout::PreserveSingleSource => {
2210 let [source_id] = old_ids else {
2211 return Err(Error::Internal(
2212 "layout-preserving replacement requires exactly one source".into(),
2213 ));
2214 };
2215 let mut source = st
2216 .metadata
2217 .segment_metas
2218 .get(source_id)
2219 .cloned()
2220 .ok_or_else(|| {
2221 Error::Corruption(format!(
2222 "layout-preserving replacement source {source_id} disappeared"
2223 ))
2224 })?;
2225 source.num_docs = doc_count;
2226 source
2227 }
2228 };
2229 let retired_ids = old_ids.to_vec();
2230 let mut next = st.metadata.clone();
2231 for id in old_ids {
2232 next.remove_segment(id);
2233 }
2234 next.add_segment_meta(new_id.clone(), replacement_info);
2235
2236 let directory = Arc::clone(&self.directory);
2237 let tracker = Arc::clone(&self.tracker);
2238 let replacement_refresh = self.replacement_refresh.read().clone();
2239 let index_label = self.schema.index_label().to_owned();
2240 self.run_lifecycle_transaction(async move {
2241 next.save(directory.as_ref()).await?;
2244 tracker.register(&new_id);
2245 st.metadata = next;
2246
2247 let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2251 drop(st);
2252 for &segment_id in &ready_to_delete {
2253 if let Err(error) =
2254 crate::segment::delete_segment(directory.as_ref(), segment_id).await
2255 {
2256 log::warn!(
2257 "[segment_cleanup] index={index_label} immediate delete failed for {}: {}",
2258 segment_id.to_hex(),
2259 error,
2260 );
2261 }
2262 }
2263 tracker.complete_deletion(&ready_to_delete);
2264 refresh_replacement_topology(replacement_refresh, &index_label).await;
2265 Ok(())
2266 })
2267 .await
2268 }
2269
2270 #[allow(clippy::too_many_arguments)]
2275 async fn do_merge(
2276 directory: &D,
2277 schema: &Arc<crate::dsl::Schema>,
2278 segment_ids_to_merge: &[String],
2279 output_segment_id: SegmentId,
2280 term_cache_blocks: usize,
2281 trained: Option<&TrainedVectorStructures>,
2282 reorder_bmp: bool,
2283 granularity: crate::segment::reorder::BpGranularity,
2284 merge_bp_time_budget: Option<std::time::Duration>,
2285 bp_memory_budget_bytes: usize,
2286 reorder_permits: Arc<ReorderConcurrencyGate>,
2287 reorder_priority: ReorderPriority,
2288 cancellation: Arc<AtomicBool>,
2289 bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2290 ) -> MergeTaskResult<(String, u32, bool)> {
2291 let output_hex = output_segment_id.to_hex();
2292 let load_start = std::time::Instant::now();
2293
2294 let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2295 for id_str in segment_ids_to_merge {
2296 let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2297 MergeTaskError::source(
2298 id_str.clone(),
2299 Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2300 )
2301 })?;
2302 segment_ids.push(id);
2303 }
2304
2305 let mut unavailable_sources = Vec::new();
2310 let mut missing_files = Vec::new();
2311 for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2312 let files = SegmentFiles::new(id.0);
2313 let mut source_unavailable = false;
2314 for path in files.mandatory_paths() {
2315 let exists = directory
2316 .exists(path)
2317 .await
2318 .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2319 if !exists {
2320 source_unavailable = true;
2321 missing_files.push(format!("{}:{:?}", id_str, path));
2322 }
2323 }
2324 if source_unavailable {
2325 unavailable_sources.push(id_str.clone());
2326 }
2327 }
2328 if !unavailable_sources.is_empty() {
2329 return Err(MergeTaskError::sources(
2330 unavailable_sources,
2331 Error::Corruption(format!(
2332 "merge sources are missing mandatory files: {}",
2333 missing_files.join(", ")
2334 )),
2335 ));
2336 }
2337
2338 let schema_arc = Arc::clone(schema);
2339 let futures: Vec<_> = segment_ids
2340 .iter()
2341 .map(|&sid| {
2342 let sch = Arc::clone(&schema_arc);
2343 async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
2344 })
2345 .collect();
2346
2347 let results = futures::future::join_all(futures).await;
2348 let mut readers = Vec::with_capacity(results.len());
2349 let mut total_docs = 0u64;
2350 for (i, result) in results.into_iter().enumerate() {
2351 match result {
2352 Ok(r) => {
2353 total_docs += r.meta().num_docs as u64;
2354 readers.push(r);
2355 }
2356 Err(e) => {
2357 log::error!(
2358 "[merge] index={} Failed to open segment {}: {:?}",
2359 schema.index_label(),
2360 segment_ids_to_merge[i],
2361 e
2362 );
2363 return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2364 }
2365 }
2366 }
2367 if total_docs > u32::MAX as u64 {
2368 return Err(Error::Internal(format!(
2369 "Merged segment doc count ({}) exceeds u32::MAX",
2370 total_docs
2371 ))
2372 .into());
2373 }
2374
2375 for (i, reader) in readers.iter().enumerate() {
2379 let meta_docs = reader.meta().num_docs;
2380 let store_docs = reader.store().num_docs();
2381 if store_docs != meta_docs {
2382 return Err(MergeTaskError::source(
2383 segment_ids_to_merge[i].clone(),
2384 Error::Corruption(format!(
2385 "pre-merge validation: segment {} store has {} docs but meta says {}",
2386 segment_ids_to_merge[i], store_docs, meta_docs
2387 )),
2388 ));
2389 }
2390 }
2391
2392 log::info!(
2393 "[merge] index={} loaded {} segment readers in {:.1}s",
2394 schema.index_label(),
2395 readers.len(),
2396 load_start.elapsed().as_secs_f64()
2397 );
2398
2399 let merger = SegmentMerger::new(Arc::clone(schema))
2400 .with_bmp_reorder(reorder_bmp)
2401 .with_granularity(granularity)
2402 .with_bp_budget(crate::segment::BpBudget {
2403 min_partition_docs: None,
2404 time_budget: merge_bp_time_budget,
2405 })
2406 .with_cancellation(cancellation)
2407 .with_bp_memory_budget(bp_memory_budget_bytes)
2408 .with_reorder_permits(reorder_permits)
2409 .with_reorder_priority(reorder_priority)
2410 .with_background_pool(bg_cpu_pool);
2411
2412 log::info!(
2413 "[merge] index={} {} segments -> {} (trained={})",
2414 schema.index_label(),
2415 segment_ids_to_merge.len(),
2416 output_hex,
2417 trained.map_or(0, |t| t.centroids.len()),
2418 );
2419
2420 let (_merged_meta, merge_stats) = merger
2421 .merge(directory, &readers, output_segment_id, trained)
2422 .await
2423 .map_err(|error| {
2424 if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2425 MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2431 } else {
2432 MergeTaskError::from(error)
2433 }
2434 })?;
2435 let bp_converged = merge_stats.bp_converged;
2436 if !bp_converged {
2437 log::info!(
2438 "[merge] index={} merge-time BP hit its wall-clock budget — output marked unconverged; \
2439 the background optimizer deepens it later",
2440 schema.index_label(),
2441 );
2442 }
2443
2444 log::info!(
2445 "[merge] index={} total wall-clock: {:.1}s ({} segments, {} docs)",
2446 schema.index_label(),
2447 load_start.elapsed().as_secs_f64(),
2448 readers.len(),
2449 total_docs,
2450 );
2451
2452 Ok((output_hex, total_docs as u32, bp_converged))
2453 }
2454
2455 pub async fn abort_merges(&self) {
2465 loop {
2466 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2467 if handles.is_empty() {
2468 return;
2469 }
2470 while let Some(result) = handles.join_next().await {
2471 if let Err(error) = result
2472 && error.is_panic()
2473 {
2474 log::error!(
2475 "[merge] index={} background task panicked while draining: {}",
2476 self.schema.index_label(),
2477 error
2478 );
2479 }
2480 }
2481 }
2482 }
2483
2484 pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2489 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2490 while handles.join_next().await.is_some() {}
2491 }
2492
2493 pub async fn wait_for_all_merges(self: &Arc<Self>) {
2502 loop {
2503 let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2504 if handles.is_empty() {
2505 break;
2506 }
2507 while handles.join_next().await.is_some() {}
2508 }
2509 }
2510
2511 pub async fn wait_for_shutdown(self: &Arc<Self>) {
2516 self.wait_for_all_merges().await;
2517 self.active_operations.wait_until_idle().await;
2518 loop {
2519 let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2520 if handles.is_empty() {
2521 break;
2522 }
2523 for handle in handles {
2524 if let Err(error) = handle.await
2525 && error.is_panic()
2526 {
2527 log::error!(
2528 "[segment_cleanup] index={} task panicked while draining: {}",
2529 self.schema.index_label(),
2530 error
2531 );
2532 }
2533 }
2534 }
2535 }
2536
2537 pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2549 self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2550 .await
2551 }
2552
2553 pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2563 self: &Arc<Self>,
2564 mut refresh_snapshots: F,
2565 ) -> Result<()>
2566 where
2567 F: FnMut() -> Fut,
2568 Fut: std::future::Future<Output = Result<()>>,
2569 {
2570 const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2575 std::time::Duration::from_millis(100);
2576 const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2580
2581 let (_force_merge_activity, policy_segment_docs) = {
2582 let st = self.state.lock().await;
2583 self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2587 (
2588 ForceMergeActivityGuard(&self.force_merge_active),
2589 st.merge_policy.max_segment_docs(),
2590 )
2591 };
2592
2593 let background_merges = self
2596 .merge_handles
2597 .lock()
2598 .iter()
2599 .filter(|handle| !handle.is_finished())
2600 .count();
2601 if background_merges > 0 {
2602 log::info!(
2603 "[force_merge] index={} waiting for {} in-flight background merge(s) before planning",
2604 self.schema.index_label(),
2605 background_merges,
2606 );
2607 }
2608 let drain_start = std::time::Instant::now();
2609 self.wait_for_all_merges().await;
2610 if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2611 log::info!(
2612 "[force_merge] index={} drained background merges in {:.1}s",
2613 self.schema.index_label(),
2614 drain_start.elapsed().as_secs_f64(),
2615 );
2616 }
2617
2618 let refresh_start = std::time::Instant::now();
2623 refresh_snapshots().await?;
2624 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2625 log::info!(
2626 "[force_merge] index={} initial snapshot refresh took {:.1}s",
2627 self.schema.index_label(),
2628 refresh_start.elapsed().as_secs_f64(),
2629 );
2630 }
2631
2632 let mut completed_outputs = HashSet::new();
2636 let mut logged_held_wait = false;
2638
2639 loop {
2640 if !self.active_operations.is_accepting() {
2641 return Err(Error::IndexClosed);
2642 }
2643
2644 let segments: Vec<(String, u32)> = {
2645 let st = self.state.lock().await;
2646 st.metadata
2647 .segment_metas
2648 .iter()
2649 .filter(|(id, _)| !completed_outputs.contains(*id))
2650 .map(|(id, info)| (id.clone(), info.num_docs))
2651 .collect()
2652 };
2653
2654 let active_ids = self.active_operations.snapshot();
2659 let held = segments
2660 .iter()
2661 .filter(|(id, _)| active_ids.contains(id))
2662 .count();
2663 let free_segments: Vec<_> = segments
2664 .into_iter()
2665 .filter(|(id, _)| !active_ids.contains(id))
2666 .collect();
2667 let max_docs = u64::from(u32::MAX);
2674 let planned_groups = plan_force_merge_groups(free_segments, max_docs);
2675 if let Some(cap) = policy_segment_docs {
2676 for group in planned_groups
2677 .iter()
2678 .filter(|group| group.segments.len() >= 2)
2679 .filter(|group| group.total_docs > u64::from(cap))
2680 {
2681 log::warn!(
2682 "[force_merge] index={} output of {} docs intentionally exceeds the \
2683 background merge policy cap of {} docs (force merge compacts to the \
2684 u32 format limit)",
2685 self.schema.index_label(),
2686 group.total_docs,
2687 cap,
2688 );
2689 }
2690 }
2691 let next_group = planned_groups
2692 .into_iter()
2693 .find(|group| group.segments.len() >= 2);
2694
2695 let Some(group) = next_group else {
2696 if held == 0 {
2697 if !completed_outputs.is_empty() {
2698 completed_outputs.clear();
2704 continue;
2705 }
2706 let remaining = {
2711 let st = self.state.lock().await;
2712 st.metadata.segment_metas.len()
2713 };
2714 if remaining > 1 {
2715 log::warn!(
2716 "[force_merge] index={} finished with {} segments: combined \
2717 document count exceeds the u32 segment format limit, so a \
2718 single output is impossible",
2719 self.schema.index_label(),
2720 remaining,
2721 );
2722 }
2723 refresh_snapshots().await?;
2728 return Ok(());
2729 }
2730 if !logged_held_wait {
2731 log::info!(
2732 "[force_merge] index={} waiting: {} segment(s) held by active \
2733 merge/reorder operations, no free group can merge",
2734 self.schema.index_label(),
2735 held
2736 );
2737 logged_held_wait = true;
2738 } else {
2739 log::debug!(
2740 "[force_merge] index={} still waiting on {} held segment(s)",
2741 self.schema.index_label(),
2742 held
2743 );
2744 }
2745 #[cfg(test)]
2746 self.force_merge_conflict_retries
2747 .fetch_add(1, Ordering::Relaxed);
2748 tokio::select! {
2749 biased;
2750 () = self.active_operations.wait_for_shutdown() => {
2751 return Err(Error::IndexClosed);
2752 }
2753 () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2754 }
2755 continue;
2756 };
2757 logged_held_wait = false;
2758
2759 let hierarchy = plan_force_merge_hierarchy(group.segments.len());
2763 let output_ids: Vec<_> = (0..hierarchy.steps.len())
2764 .map(|_| SegmentId::new())
2765 .collect();
2766 let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
2767 let mut all_ids = source_ids.clone();
2768 all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
2769 let group_guard = {
2770 let st = self.state.lock().await;
2771 source_ids
2772 .iter()
2773 .all(|id| st.metadata.has_segment(id))
2774 .then(|| self.active_operations.try_register(all_ids))
2775 .flatten()
2776 };
2777 let _group_guard = match group_guard {
2778 Some(guard) => guard,
2779 None if !self.active_operations.is_accepting() => {
2780 return Err(Error::IndexClosed);
2781 }
2782 None => {
2783 #[cfg(test)]
2784 self.force_merge_conflict_retries
2785 .fetch_add(1, Ordering::Relaxed);
2786 log::debug!(
2787 "[force_merge] index={} group lost a registration race, replanning",
2788 self.schema.index_label()
2789 );
2790 let had_tracked_merges = !self.merge_handles.lock().is_empty();
2791 self.wait_for_merging_thread().await;
2792 if !had_tracked_merges {
2793 tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2794 }
2795 continue;
2796 }
2797 };
2798
2799 log::info!(
2800 "[force_merge] index={} planned final group: {} segments, {} docs, {} merge pass(es)",
2801 self.schema.index_label(),
2802 group.segments.len(),
2803 group.total_docs,
2804 output_ids.len(),
2805 );
2806
2807 let group_global_merge_permit = if self.reorder_on_merge {
2817 let capacity_start = std::time::Instant::now();
2818 let permit = tokio::select! {
2819 biased;
2820 () = self.active_operations.wait_for_shutdown() => {
2821 return Err(Error::IndexClosed);
2822 }
2823 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2824 permit.map_err(|_| {
2825 Error::Internal(
2826 "global background merge scheduler is closed".into(),
2827 )
2828 })?
2829 }
2830 };
2831 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2832 log::info!(
2833 "[force_merge] index={} waited {:.1}s for foreground global merge capacity",
2834 self.schema.index_label(),
2835 capacity_start.elapsed().as_secs_f64(),
2836 );
2837 }
2838 Some(permit)
2839 } else {
2840 None
2841 };
2842 let _foreground_reorder = if self.reorder_on_merge {
2843 log::info!(
2844 "[force_merge] index={} prioritizing BP capacity ({} total pass slot(s))",
2845 self.schema.index_label(),
2846 self.reorder_permits.limit(),
2847 );
2848 let admission_start = std::time::Instant::now();
2849 let guard = Arc::clone(&self.reorder_permits)
2850 .begin_foreground()
2851 .await
2852 .map_err(|_| {
2853 Error::Internal("background reorder scheduler is closed".into())
2854 })?;
2855 if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
2856 log::info!(
2857 "[force_merge] index={} acquired foreground BP capacity in {:.1}s",
2858 self.schema.index_label(),
2859 admission_start.elapsed().as_secs_f64(),
2860 );
2861 }
2862 Some(guard)
2863 } else {
2864 None
2865 };
2866
2867 let source_count = group.segments.len();
2868 let mut nodes: Vec<Option<(String, u32)>> =
2869 group.segments.into_iter().map(Some).collect();
2870 nodes.resize_with(source_count + hierarchy.steps.len(), || None);
2871 for (step_index, step) in hierarchy.steps.iter().enumerate() {
2872 let final_pass = step_index + 1 == hierarchy.steps.len();
2873 let mut batch_entries = Vec::with_capacity(step.inputs.len());
2874 for &node in &step.inputs {
2875 let entry = nodes
2876 .get_mut(node)
2877 .and_then(Option::take)
2878 .expect("force-merge hierarchy must reference an available node");
2879 batch_entries.push(entry);
2880 }
2881 let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
2882 let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
2883 let output_id = output_ids[step_index];
2884
2885 let capacity_start = std::time::Instant::now();
2886 let step_global_merge_permit = if group_global_merge_permit.is_none() {
2887 Some(tokio::select! {
2888 biased;
2889 () = self.active_operations.wait_for_shutdown() => {
2890 return Err(Error::IndexClosed);
2891 }
2892 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2893 permit.map_err(|_| {
2894 Error::Internal(
2895 "global background merge scheduler is closed".into(),
2896 )
2897 })?
2898 }
2899 })
2900 } else {
2901 None
2902 };
2903 if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2904 log::info!(
2905 "[force_merge] index={} waited {:.1}s for global merge capacity",
2906 self.schema.index_label(),
2907 capacity_start.elapsed().as_secs_f64(),
2908 );
2909 }
2910
2911 let reorder_bmp = final_pass && self.reorder_on_merge;
2915 log::info!(
2916 "[force_merge] index={} {} pass: {} segments ({} docs, bp={})",
2917 self.schema.index_label(),
2918 if final_pass {
2919 "final"
2920 } else {
2921 "fan-in reduction"
2922 },
2923 batch.len(),
2924 batch_docs,
2925 reorder_bmp,
2926 );
2927 let (new_segment_id, total_docs, _) = self
2928 .merge_and_replace_registered(
2929 &batch,
2930 output_id,
2931 reorder_bmp,
2932 ReorderPriority::Foreground,
2933 )
2934 .await
2935 .map_err(|error| error.error)?;
2936 drop(step_global_merge_permit);
2937
2938 let refresh_start = std::time::Instant::now();
2941 refresh_snapshots().await?;
2942 if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2943 log::info!(
2944 "[force_merge] index={} post-replacement snapshot refresh took {:.1}s",
2945 self.schema.index_label(),
2946 refresh_start.elapsed().as_secs_f64(),
2947 );
2948 }
2949
2950 let output_node = source_count + step_index;
2951 debug_assert!(nodes[output_node].is_none());
2952 nodes[output_node] = Some((new_segment_id, total_docs));
2953 }
2954 let (root_id, _) = nodes[hierarchy.root]
2955 .take()
2956 .expect("force-merge hierarchy must produce its root");
2957 debug_assert!(nodes.into_iter().all(|node| node.is_none()));
2958 completed_outputs.insert(root_id);
2959 }
2960 }
2961
2962 fn segment_needs_vector_rewrite(
2963 &self,
2964 schema: &crate::dsl::Schema,
2965 reader: &SegmentReader,
2966 field_ids: &[u32],
2967 trained: &TrainedVectorStructures,
2968 rewrite_existing: bool,
2969 ) -> Result<bool> {
2970 for &field_id in field_ids {
2971 let flat = reader.flat_vectors().get(&field_id);
2972 let ann = reader.vector_indexes().get(&field_id);
2973 if ann.is_some() && flat.is_none() {
2974 return Err(Error::Corruption(format!(
2975 "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2976 reader.meta().id,
2977 )));
2978 }
2979
2980 let Some(flat) = flat else {
2981 continue;
2982 };
2983 if flat.num_vectors == 0 {
2984 continue;
2985 }
2986 if rewrite_existing {
2987 return Ok(true);
2988 }
2989 let field = crate::dsl::Field(field_id);
2990 let entry = schema.get_field_entry(field).ok_or_else(|| {
2991 Error::Corruption(format!(
2992 "segment {:032x} references unknown vector field {field_id}",
2993 reader.meta().id,
2994 ))
2995 })?;
2996 let current = match entry.field_type {
2997 crate::dsl::FieldType::DenseVector
3001 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3002 config.index_type == crate::dsl::VectorIndexType::Tq
3003 }) =>
3004 {
3005 matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
3006 }
3007 crate::dsl::FieldType::DenseVector
3008 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3009 config.index_type == crate::dsl::VectorIndexType::IvfTq
3010 }) =>
3011 {
3012 let config = entry
3013 .dense_vector_config
3014 .as_ref()
3015 .expect("matched IVF-TQ configuration");
3016 match (ann, trained.centroids.get(&field_id)) {
3017 (
3018 Some(crate::segment::VectorIndex::IvfTq { index, .. }),
3019 Some(centroids),
3020 ) => {
3021 let header = index.get().header();
3022 crate::structures::is_ivf_tq_cosine_generation(centroids.version)
3023 && crate::structures::is_ivf_tq_cosine_generation(
3024 header.quantizer_version,
3025 )
3026 && header.dim == config.dim
3027 && header.num_clusters == centroids.num_clusters
3028 && header.quantizer_version == centroids.version
3029 && header.codebook_version
3030 == crate::structures::vector::quantization::tq_expected_fingerprint(
3031 config.dim,
3032 )
3033 && header.routing == config.ivf_routing
3034 }
3035 (None, None) => true,
3036 _ => false,
3037 }
3038 }
3039 crate::dsl::FieldType::DenseVector
3040 if entry.dense_vector_config.as_ref().is_some_and(|config| {
3041 config.index_type == crate::dsl::VectorIndexType::Scann
3042 }) =>
3043 {
3044 match (ann, trained.scann_artifacts.get(&field_id)) {
3045 (Some(crate::segment::VectorIndex::ScannAh(index)), Some(artifact)) => {
3046 index
3047 .get()
3048 .validate_scann_generation(
3049 artifact.config(),
3050 artifact.generation(),
3051 artifact.artifact_id(),
3052 )
3053 .is_ok()
3054 }
3055 (None, None) => true,
3056 _ => false,
3057 }
3058 }
3059 crate::dsl::FieldType::DenseVector => false,
3062 crate::dsl::FieldType::BinaryDenseVector
3063 if entry
3064 .binary_dense_vector_config
3065 .as_ref()
3066 .is_some_and(|config| {
3067 config.index_type == crate::dsl::BinaryIndexType::Scann
3068 }) =>
3069 {
3070 match (ann, trained.scann_artifacts.get(&field_id)) {
3071 (Some(crate::segment::VectorIndex::ScannBinary(index)), Some(artifact)) => {
3072 index
3073 .get()
3074 .validate_scann_generation(
3075 artifact.config(),
3076 artifact.generation(),
3077 artifact.artifact_id(),
3078 )
3079 .is_ok()
3080 }
3081 (None, None) => true,
3082 _ => false,
3083 }
3084 }
3085 crate::dsl::FieldType::BinaryDenseVector => matches!(
3086 (ann, trained.binary_quantizers.get(&field_id)),
3087 (Some(crate::segment::VectorIndex::BinaryIvf(_)), Some(_)) | (None, None)
3088 ),
3089 _ => false,
3090 };
3091 if !current {
3092 return Ok(true);
3093 }
3094 }
3095 Ok(false)
3096 }
3097
3098 async fn acquire_vector_rewrite_capacity(
3099 &self,
3100 ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
3101 let global = tokio::select! {
3102 biased;
3103 () = self.active_operations.wait_for_shutdown() => {
3104 return Err(Error::IndexClosed);
3105 }
3106 permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3107 permit.map_err(|_| Error::Internal(
3108 "global background merge scheduler is closed".into()
3109 ))?
3110 }
3111 };
3112 let local = tokio::select! {
3113 biased;
3114 () = self.active_operations.wait_for_shutdown() => {
3115 return Err(Error::IndexClosed);
3116 }
3117 permit = Arc::clone(&self.merge_permits).acquire_owned() => {
3118 permit.map_err(|_| Error::Internal(
3119 "background merge scheduler is closed".into()
3120 ))?
3121 }
3122 };
3123 Ok((global, local))
3124 }
3125
3126 async fn build_vector_replacement(
3127 self: &Arc<Self>,
3128 schemas: (&Arc<crate::dsl::Schema>, &Arc<crate::dsl::Schema>),
3129 segment_id: &str,
3130 source_id: SegmentId,
3131 output_id: SegmentId,
3132 trained: &TrainedVectorStructures,
3133 failure_context: &'static str,
3134 ) -> Result<(String, u32, OutputCleanupGuard)> {
3135 let mut cleanup = self.output_cleanup_guard(output_id);
3136 match crate::segment::reorder::rewrite_vector_segment(
3137 self.directory.as_ref(),
3138 schemas,
3139 source_id,
3140 output_id,
3141 self.term_cache_blocks,
3142 trained,
3143 Some(self.background_cpu_pool()),
3144 )
3145 .await
3146 {
3147 Ok((new_id, doc_count)) => {
3148 self.validate_completed_segment(&new_id, doc_count).await?;
3149 Ok((new_id, doc_count, cleanup))
3150 }
3151 Err(error) => {
3152 self.delete_output_if_unregistered(output_id, failure_context)
3153 .await;
3154 cleanup.disarm();
3155 if is_deterministic_source_error(&error) {
3156 self.quarantine_segment(segment_id, &error);
3157 }
3158 Err(error)
3159 }
3160 }
3161 }
3162
3163 pub(crate) async fn stage_vector_generation(
3167 self: &Arc<Self>,
3168 _artifact_update: &VectorArtifactUpdateGuard,
3169 segment_ids: &[String],
3170 field_ids: &[u32],
3171 trained: Arc<TrainedVectorStructures>,
3172 rewrite_existing: bool,
3173 ) -> Result<Vec<StagedVectorSegment>> {
3174 let schema = self.published_generation().schema.clone();
3175 self.stage_vector_generation_with_schema(
3176 _artifact_update,
3177 segment_ids,
3178 field_ids,
3179 trained,
3180 rewrite_existing,
3181 schema,
3182 )
3183 .await
3184 }
3185
3186 pub(crate) async fn stage_vector_generation_with_schema(
3187 self: &Arc<Self>,
3188 _artifact_update: &VectorArtifactUpdateGuard,
3189 segment_ids: &[String],
3190 field_ids: &[u32],
3191 trained: Arc<TrainedVectorStructures>,
3192 rewrite_existing: bool,
3193 schema: Arc<crate::dsl::Schema>,
3194 ) -> Result<Vec<StagedVectorSegment>> {
3195 if !self.vector_artifact_update.load(Ordering::Acquire) {
3196 return Err(Error::Internal(
3197 "cannot stage a vector generation without an exclusive update lease".into(),
3198 ));
3199 }
3200
3201 let source_schema = self.published_generation().schema.clone();
3202 let mut staged = Vec::new();
3203 for segment_id in segment_ids {
3204 if self.quarantined_segments.lock().contains(segment_id) {
3205 return Err(Error::Corruption(format!(
3206 "segment {segment_id} is quarantined after a deterministic source failure"
3207 )));
3208 }
3209 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3210 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3211 })?;
3212
3213 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3216
3217 let output_id = SegmentId::new();
3218 let output_hex = output_id.to_hex();
3219 let operation = {
3220 let st = self.state.lock().await;
3221 if !st.metadata.has_segment(segment_id) {
3222 return Err(Error::Corruption(format!(
3223 "vector generation source {segment_id} disappeared while lifecycle work was paused"
3224 )));
3225 }
3226 self.active_operations
3227 .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
3228 }
3229 .ok_or_else(|| {
3230 if self.active_operations.is_accepting() {
3231 Error::Internal(format!(
3232 "vector generation could not claim stable source {segment_id}"
3233 ))
3234 } else {
3235 Error::IndexClosed
3236 }
3237 })?;
3238
3239 let reader = SegmentReader::open(
3240 self.directory.as_ref(),
3241 source_id,
3242 Arc::clone(&source_schema),
3243 self.term_cache_blocks,
3244 )
3245 .await?;
3246 if !self.segment_needs_vector_rewrite(
3247 schema.as_ref(),
3248 &reader,
3249 field_ids,
3250 trained.as_ref(),
3251 rewrite_existing,
3252 )? {
3253 continue;
3254 }
3255 drop(reader);
3256
3257 let (new_id, doc_count, cleanup) = self
3258 .build_vector_replacement(
3259 (&source_schema, &schema),
3260 segment_id,
3261 source_id,
3262 output_id,
3263 trained.as_ref(),
3264 "vector generation staging failure",
3265 )
3266 .await?;
3267 debug_assert_eq!(new_id, output_hex);
3268 let output_reader = SegmentReader::open(
3269 self.directory.as_ref(),
3270 output_id,
3271 Arc::clone(&schema),
3272 self.term_cache_blocks,
3273 )
3274 .await?;
3275 if self.segment_needs_vector_rewrite(
3276 schema.as_ref(),
3277 &output_reader,
3278 field_ids,
3279 trained.as_ref(),
3280 false,
3281 )? {
3282 return Err(Error::Corruption(format!(
3283 "staged vector segment {new_id} does not match its candidate codebook generation"
3284 )));
3285 }
3286
3287 staged.push(StagedVectorSegment {
3288 source_id: segment_id.clone(),
3289 output_id,
3290 doc_count,
3291 _operation: operation,
3292 cleanup,
3293 });
3294 }
3295 Ok(staged)
3296 }
3297
3298 async fn rewrite_vector_segment_once(
3299 self: &Arc<Self>,
3300 segment_id: &str,
3301 field_ids: &[u32],
3302 ) -> Result<VectorSegmentRewriteOutcome> {
3303 if self.quarantined_segments.lock().contains(segment_id) {
3304 return Err(Error::Corruption(format!(
3305 "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3306 )));
3307 }
3308 let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3309 Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3310 })?;
3311
3312 let _capacity = self.acquire_vector_rewrite_capacity().await?;
3317
3318 let output_id = SegmentId::new();
3319 let output_hex = output_id.to_hex();
3320 let all_ids = vec![segment_id.to_owned(), output_hex];
3321 let operation = {
3322 let st = self.state.lock().await;
3323 if !st.metadata.has_segment(segment_id) {
3324 return Ok(VectorSegmentRewriteOutcome::SourceGone);
3325 }
3326 self.active_operations.try_register(all_ids)
3327 };
3328 let _operation = match operation {
3329 Some(operation) => operation,
3330 None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3331 None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3332 };
3333
3334 let Some(trained) = self.trained_for_segment_build() else {
3335 return Ok(VectorSegmentRewriteOutcome::Deferred);
3336 };
3337 let schema = self.published_generation().schema.clone();
3338
3339 let reader = SegmentReader::open(
3340 self.directory.as_ref(),
3341 source_id,
3342 Arc::clone(&schema),
3343 self.term_cache_blocks,
3344 )
3345 .await?;
3346 if !self.segment_needs_vector_rewrite(
3347 schema.as_ref(),
3348 &reader,
3349 field_ids,
3350 trained.as_ref(),
3351 false,
3352 )? {
3353 return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3354 }
3355 drop(reader);
3356
3357 let (new_id, doc_count, mut output_cleanup) = self
3358 .build_vector_replacement(
3359 (&schema, &schema),
3360 segment_id,
3361 source_id,
3362 output_id,
3363 trained.as_ref(),
3364 "vector rewrite failure",
3365 )
3366 .await?;
3367
3368 if let Err(error) = self
3369 .replace_segments(
3370 &[segment_id.to_owned()],
3371 new_id,
3372 doc_count,
3373 ReplacementLayout::PreserveSingleSource,
3374 )
3375 .await
3376 {
3377 self.delete_output_if_unregistered(output_id, "vector replacement failure")
3378 .await;
3379 output_cleanup.disarm();
3380 return Err(error);
3381 }
3382 output_cleanup.disarm();
3383 Ok(VectorSegmentRewriteOutcome::Rewritten)
3384 }
3385
3386 pub(crate) async fn rewrite_vector_segments(
3391 self: &Arc<Self>,
3392 field_ids: &[u32],
3393 ) -> Result<usize> {
3394 if field_ids.is_empty() {
3395 return Ok(0);
3396 }
3397 let mut rewritten = 0usize;
3398 loop {
3399 let segment_ids = self.get_segment_ids().await;
3400 let mut conflicted = false;
3401 let mut changed = false;
3402 for segment_id in segment_ids {
3403 match self
3404 .rewrite_vector_segment_once(&segment_id, field_ids)
3405 .await?
3406 {
3407 VectorSegmentRewriteOutcome::Rewritten => {
3408 rewritten += 1;
3409 changed = true;
3410 }
3411 VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3412 VectorSegmentRewriteOutcome::Deferred => {
3413 return Err(Error::Internal(
3414 "ANN finalization lost the published trained generation".into(),
3415 ));
3416 }
3417 VectorSegmentRewriteOutcome::AlreadyCurrent
3418 | VectorSegmentRewriteOutcome::SourceGone => {}
3419 }
3420 }
3421 if !conflicted && !changed {
3422 log::info!(
3423 "[dense_vector_rewrite] index={} ANN finalization complete ({} segment(s) rewritten)",
3424 self.schema.index_label(),
3425 rewritten,
3426 );
3427 return Ok(rewritten);
3428 }
3429 tokio::select! {
3430 biased;
3431 () = self.active_operations.wait_for_shutdown() => {
3432 return Err(Error::IndexClosed);
3433 }
3434 () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3435 }
3436 }
3437 }
3438
3439 pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3444 if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3445 return;
3446 }
3447 let manager = Arc::clone(self);
3448 let future = async move {
3449 let field_ids = manager
3450 .read_metadata(|metadata| {
3451 metadata
3452 .vector_fields
3453 .keys()
3454 .filter(|field_id| metadata.is_field_built(**field_id))
3455 .copied()
3456 .collect::<Vec<_>>()
3457 })
3458 .await;
3459 for segment_id in segment_ids {
3460 loop {
3461 match manager
3462 .rewrite_vector_segment_once(&segment_id, &field_ids)
3463 .await
3464 {
3465 Ok(VectorSegmentRewriteOutcome::Conflict) => {
3466 tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3467 }
3468 Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3469 Ok(_) => break,
3470 Err(error) => {
3471 log::error!(
3472 "[dense_vector_rewrite] index={} failed to upgrade newly committed segment {}: {}",
3473 manager.schema.index_label(),
3474 segment_id,
3475 error,
3476 );
3477 break;
3478 }
3479 }
3480 }
3481 }
3482 };
3483 let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3484 log::warn!(
3485 "[dense_vector_rewrite] index={} runtime unavailable; newly committed flat segment upgrade deferred",
3486 self.schema.index_label()
3487 );
3488 return;
3489 };
3490 if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3491 log::warn!(
3492 "[dense_vector_rewrite] index={} runtime rejected newly committed flat segment upgrade",
3493 self.schema.index_label()
3494 );
3495 }
3496 }
3497
3498 pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3505 self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3506 .await
3507 }
3508
3509 pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3512 self: &Arc<Self>,
3513 mut refresh_snapshots: F,
3514 ) -> Result<()>
3515 where
3516 F: FnMut() -> Fut,
3517 Fut: std::future::Future<Output = Result<()>>,
3518 {
3519 self.wait_for_all_merges().await;
3520 refresh_snapshots().await?;
3521 let segment_ids = self.get_segment_ids().await;
3522
3523 if segment_ids.is_empty() {
3524 log::info!(
3525 "[reorder] index={} no segments to reorder",
3526 self.schema.index_label()
3527 );
3528 return Ok(());
3529 }
3530
3531 log::info!(
3532 "[reorder] index={} reordering {} segments",
3533 self.schema.index_label(),
3534 segment_ids.len()
3535 );
3536
3537 for seg_id in segment_ids {
3538 match self
3539 .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3540 .await
3541 {
3542 Ok(true) => refresh_snapshots().await?,
3543 Ok(false) => log::warn!(
3544 "[reorder] index={} segment {} skipped (in merge)",
3545 self.schema.index_label(),
3546 seg_id
3547 ),
3548 Err(e) => return Err(e),
3549 }
3550 }
3551
3552 refresh_snapshots().await?;
3555 log::info!(
3556 "[reorder] index={} all segments reordered",
3557 self.schema.index_label()
3558 );
3559 Ok(())
3560 }
3561
3562 pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3567 self.unreordered_segments()
3568 .await
3569 .into_iter()
3570 .map(|(id, _)| id)
3571 .collect()
3572 }
3573
3574 pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3577 let quarantined = self.quarantined_segments.lock().clone();
3578 let paused = self.paused_reorder_segments();
3579 let st = self.state.lock().await;
3580 let active_ids = self.active_operations.snapshot();
3581 st.metadata
3582 .segment_metas
3583 .iter()
3584 .filter(|(id, info)| {
3585 !info.reordered
3586 && info.bp_converged
3587 && !active_ids.contains(*id)
3588 && !quarantined.contains(*id)
3589 && !paused.contains(*id)
3590 })
3591 .map(|(id, info)| (id.clone(), info.num_docs))
3592 .collect()
3593 }
3594
3595 pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3599 self.unconverged_segments_below(u32::MAX)
3600 .await
3601 .into_iter()
3602 .map(|(id, docs, _)| (id, docs))
3603 .collect()
3604 }
3605
3606 pub async fn unconverged_segments_below(
3609 &self,
3610 max_unconverged_passes: u32,
3611 ) -> Vec<(String, u32, u32)> {
3612 let quarantined = self.quarantined_segments.lock().clone();
3613 let paused = self.paused_reorder_segments();
3614 let st = self.state.lock().await;
3615 let active_ids = self.active_operations.snapshot();
3616 st.metadata
3617 .segment_metas
3618 .iter()
3619 .filter(|(id, info)| {
3620 !info.bp_converged
3621 && info.bp_unconverged_passes < max_unconverged_passes
3622 && !active_ids.contains(*id)
3623 && !quarantined.contains(*id)
3624 && !paused.contains(*id)
3625 })
3626 .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3627 .collect()
3628 }
3629
3630 async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3640 let st = self.state.lock().await;
3641 let deepening = ids.iter().any(|id| {
3642 st.metadata
3643 .segment_metas
3644 .get(id)
3645 .is_some_and(|info| !info.bp_converged)
3646 });
3647 drop(st);
3648 if deepening {
3649 log::info!(
3650 "[reorder] index={} source BP lineage unconverged — forcing record-level BP (deepening pass)",
3651 self.schema.index_label(),
3652 );
3653 crate::segment::reorder::BpGranularity::Records
3654 } else {
3655 crate::segment::reorder::BpGranularity::Auto
3656 }
3657 }
3658
3659 pub async fn reorder_single_segment(
3664 self: &Arc<Self>,
3665 seg_id: &str,
3666 rayon_pool: Option<Arc<rayon::ThreadPool>>,
3667 bp_budget: crate::segment::BpBudget,
3668 ) -> Result<bool> {
3669 let source_id = SegmentId::from_hex(seg_id)
3670 .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3671 if self.quarantined_segments.lock().contains(seg_id) {
3672 return Err(Error::Corruption(format!(
3673 "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3674 seg_id
3675 )));
3676 }
3677 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3678 log::debug!(
3679 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3680 self.schema.index_label(),
3681 seg_id,
3682 );
3683 return Ok(false);
3684 }
3685
3686 let reorder_gate = Arc::clone(&self.reorder_permits);
3691 let _reorder_permit = tokio::select! {
3692 biased;
3693 () = self.active_operations.wait_for_shutdown() => {
3694 return Err(Error::IndexClosed);
3695 }
3696 permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3697 permit.map_err(|_| {
3698 Error::Internal("background reorder scheduler is closed".into())
3699 })?
3700 }
3701 };
3702
3703 let output_id = SegmentId::new();
3704 let output_hex = output_id.to_hex();
3705 let source_ids = [seg_id.to_string()];
3706 let granularity = self.merge_granularity(&source_ids).await;
3707
3708 let all_ids = vec![seg_id.to_string(), output_hex];
3714 let (_guard, source_docs, schema) = {
3715 let st = self.state.lock().await;
3716 if self.force_merge_active.load(Ordering::Acquire) > 0 {
3720 log::debug!(
3721 "[optimizer] index={} explicit force merge active, skipping reorder of {}",
3722 self.schema.index_label(),
3723 seg_id,
3724 );
3725 return Ok(false);
3726 }
3727 let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3728 log::info!(
3729 "[optimizer] index={} segment {} no longer in metadata (merged away), skipping reorder",
3730 self.schema.index_label(),
3731 seg_id
3732 );
3733 self.clear_reorder_retry(seg_id);
3734 return Ok(false);
3735 };
3736
3737 let schema = self.published_generation().schema.clone();
3738 match self.active_operations.try_register(all_ids) {
3739 Some(guard) => (guard, source_meta.num_docs, schema),
3740 None if !self.active_operations.is_accepting() => {
3741 return Err(Error::IndexClosed);
3742 }
3743 None => {
3744 log::debug!(
3745 "[optimizer] index={} segment {} in active merge, skipping",
3746 self.schema.index_label(),
3747 seg_id
3748 );
3749 return Ok(false);
3750 }
3751 }
3752 };
3753
3754 if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3759 if is_deterministic_source_error(&error) {
3760 self.quarantine_segment(seg_id, &error);
3761 } else if !matches!(&error, Error::IndexClosed) {
3762 self.pause_reorder_retries(seg_id, &error);
3763 }
3764 return Err(error);
3765 }
3766
3767 let mut output_cleanup = self.output_cleanup_guard(output_id);
3768
3769 let reorder_result = crate::segment::reorder::reorder_segment(
3770 self.directory.as_ref(),
3771 &schema,
3772 source_id,
3773 output_id,
3774 self.term_cache_blocks,
3775 self.bp_memory_budget_bytes,
3776 bp_budget,
3777 granularity,
3778 rayon_pool,
3779 Some(self.active_operations.cancellation_flag()),
3780 )
3781 .await;
3782 let (new_id, total_docs, bp_converged) = match reorder_result {
3783 Ok(v) => v,
3784 Err(e) => {
3785 self.delete_output_if_unregistered(output_id, "reorder failure")
3788 .await;
3789 output_cleanup.disarm();
3790 if is_deterministic_source_error(&e) {
3791 self.quarantine_segment(seg_id, &e);
3792 } else if !matches!(&e, Error::IndexClosed) {
3793 self.pause_reorder_retries(seg_id, &e);
3794 }
3795 return Err(e);
3796 }
3797 };
3798
3799 let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3805 if let Err(e) = self
3806 .replace_segments(
3807 &[seg_id.to_string()],
3808 new_id,
3809 total_docs,
3810 ReplacementLayout::BpReordered {
3811 converged: ladder_converged,
3812 },
3813 )
3814 .await
3815 {
3816 self.delete_output_if_unregistered(output_id, "replacement failure")
3817 .await;
3818 output_cleanup.disarm();
3819 if !matches!(&e, Error::IndexClosed) {
3820 self.pause_reorder_retries(seg_id, &e);
3821 }
3822 return Err(e);
3823 }
3824 output_cleanup.disarm();
3825 self.clear_reorder_retry(seg_id);
3826
3827 Ok(true)
3828 }
3829
3830 pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3837 let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3838
3839 if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3840 for entry in entries {
3841 let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3842 continue;
3843 };
3844 let Some(rest) = filename.strip_prefix("seg_") else {
3845 continue;
3846 };
3847 let Some(hex_id) = rest.get(..32) else {
3848 continue;
3849 };
3850 if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3851 continue;
3852 }
3853 orphan_files
3854 .entry(hex_id.to_ascii_lowercase())
3855 .or_default()
3856 .push(entry);
3857 }
3858 }
3859
3860 let mut deleted = 0;
3861 for (hex_id, paths) in &orphan_files {
3862 let deletion_guard = {
3867 let st = self.state.lock().await;
3868 if st.metadata.has_segment(hex_id) {
3869 continue;
3870 }
3871 let Some(guard) = self
3872 .active_operations
3873 .try_register(vec![hex_id.to_string()])
3874 else {
3875 continue;
3876 };
3877 if self.tracker.is_deletion_protected(hex_id) {
3878 drop(guard);
3879 continue;
3880 }
3881 guard
3882 };
3883
3884 let results =
3889 futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3890 .await;
3891 let removed = results.into_iter().all(|result| match result {
3892 Ok(()) => true,
3893 Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3894 Err(error) => {
3895 log::warn!(
3896 "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
3897 self.schema.index_label(),
3898 hex_id,
3899 error,
3900 );
3901 false
3902 }
3903 });
3904 drop(deletion_guard);
3907 if removed {
3908 deleted += 1;
3909 log::info!(
3910 "[segment_cleanup] index={} swept orphan segment {}",
3911 self.schema.index_label(),
3912 hex_id
3913 );
3914 }
3915 }
3916
3917 Ok(deleted)
3918 }
3919}
3920
3921#[cfg(test)]
3922mod tests {
3923 use super::*;
3924 use std::sync::atomic::{AtomicBool, Ordering};
3925
3926 fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3927 let schema = crate::dsl::SchemaBuilder::default().build();
3928 let metadata = IndexMetadata::new(schema.clone());
3929 Arc::new(SegmentManager::new(
3930 Arc::new(crate::directories::RamDirectory::new()),
3931 Arc::new(schema),
3932 metadata,
3933 Box::new(crate::merge::NoMergePolicy),
3934 0,
3935 1,
3936 Arc::new(Semaphore::new(1)),
3937 None,
3938 1024,
3939 Arc::new(ReorderConcurrencyGate::new(1)),
3940 None,
3941 ))
3942 }
3943
3944 #[test]
3945 fn force_merge_planner_pairs_large_and_small_segments() {
3946 let groups = plan_force_merge_groups(
3947 vec![
3948 ("a".into(), 6),
3949 ("b".into(), 6),
3950 ("c".into(), 4),
3951 ("d".into(), 4),
3952 ],
3953 10,
3954 );
3955
3956 assert_eq!(groups.len(), 2);
3957 assert!(groups.iter().all(|group| group.total_docs == 10));
3958 assert!(groups.iter().all(|group| group.segments.len() == 2));
3959 }
3960
3961 #[test]
3962 fn force_merge_planner_leaves_oversized_segments_alone() {
3963 let groups = plan_force_merge_groups(
3964 vec![
3965 ("oversized".into(), 11),
3966 ("small-a".into(), 5),
3967 ("small-b".into(), 5),
3968 ],
3969 10,
3970 );
3971
3972 assert_eq!(groups.len(), 2);
3973 assert_eq!(groups[0].total_docs, 10);
3974 assert_eq!(groups[0].segments.len(), 2);
3975 assert_eq!(groups[1].total_docs, 11);
3976 assert_eq!(groups[1].segments.len(), 1);
3977 }
3978
3979 #[test]
3980 fn force_merge_planner_never_exceeds_segment_format_limit() {
3981 let groups = plan_force_merge_groups(
3982 vec![
3983 ("large-a".into(), 3_000_000_000),
3984 ("large-b".into(), 2_000_000_000),
3985 ],
3986 u64::from(u32::MAX),
3987 );
3988 assert_eq!(groups.len(), 2);
3989 assert!(
3990 groups
3991 .iter()
3992 .all(|group| group.total_docs <= u64::from(u32::MAX))
3993 );
3994 }
3995
3996 #[test]
3997 fn force_merge_hierarchy_has_one_final_bp_pass() {
3998 assert_eq!(force_merge_output_count(1), 0);
3999 assert_eq!(force_merge_output_count(2), 1);
4000 assert_eq!(force_merge_output_count(64), 1);
4001 assert_eq!(force_merge_output_count(65), 2);
4002 assert_eq!(force_merge_output_count(127), 2);
4003 assert_eq!(force_merge_output_count(128), 3);
4004 assert_eq!(force_merge_output_count(1_000), 16);
4005 }
4006
4007 fn expand_force_merge_node(
4008 hierarchy: &ForceMergeHierarchy,
4009 source_count: usize,
4010 node: usize,
4011 sources: &mut Vec<usize>,
4012 ) {
4013 if node < source_count {
4014 sources.push(node);
4015 return;
4016 }
4017
4018 let step_index = node - source_count;
4019 let step = hierarchy
4020 .steps
4021 .get(step_index)
4022 .expect("merge input must refer to an existing source or output");
4023 for &input in &step.inputs {
4024 assert!(
4025 input < node,
4026 "merge step {step_index} refers to a future output node {input}"
4027 );
4028 expand_force_merge_node(hierarchy, source_count, input, sources);
4029 }
4030 }
4031
4032 #[test]
4033 fn force_merge_hierarchy_has_minimal_valid_arity() {
4034 let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
4035
4036 for source_count in source_counts {
4037 let hierarchy = plan_force_merge_hierarchy(source_count);
4038 let output_count = hierarchy.steps.len();
4039
4040 assert!(
4041 hierarchy
4042 .steps
4043 .iter()
4044 .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
4045 "invalid merge arity for {source_count} sources"
4046 );
4047 assert!(
4048 source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
4049 "{output_count} outputs cannot reduce {source_count} sources"
4050 );
4051 assert!(
4052 output_count == 1
4053 || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
4054 "{output_count} outputs are not minimal for {source_count} sources"
4055 );
4056 assert_eq!(output_count, force_merge_output_count(source_count));
4057 }
4058 }
4059
4060 #[test]
4061 fn force_merge_hierarchy_preserves_exact_source_order() {
4062 for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
4063 let hierarchy = plan_force_merge_hierarchy(source_count);
4064 let mut sources = Vec::with_capacity(source_count);
4065 expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
4066 assert_eq!(
4067 sources,
4068 (0..source_count).collect::<Vec<_>>(),
4069 "source order changed for {source_count} sources"
4070 );
4071 }
4072 }
4073
4074 fn force_merge_rewrite_cost(source_count: usize) -> usize {
4075 let hierarchy = plan_force_merge_hierarchy(source_count);
4076 let mut node_weights = vec![1usize; source_count];
4077 let mut rewrite_cost = 0usize;
4078
4079 for (step_index, step) in hierarchy.steps.iter().enumerate() {
4080 let output = source_count + step_index;
4081 let output_weight = step
4082 .inputs
4083 .iter()
4084 .map(|&input| {
4085 assert!(
4086 input < output,
4087 "merge step {step_index} refers to future output {input}"
4088 );
4089 node_weights[input]
4090 })
4091 .sum::<usize>();
4092 rewrite_cost += output_weight;
4093 node_weights.push(output_weight);
4094 }
4095
4096 assert_eq!(node_weights[hierarchy.root], source_count);
4097 rewrite_cost
4098 }
4099
4100 #[test]
4101 fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
4102 assert_eq!(force_merge_rewrite_cost(65), 67);
4103 assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
4104 }
4105
4106 #[test]
4107 fn block_copy_carries_bp_debt_without_spending_an_attempt() {
4108 assert_eq!(
4109 replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
4110 (false, false, 3),
4111 );
4112 assert_eq!(
4113 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
4114 (true, false, 4),
4115 );
4116 assert_eq!(
4117 replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
4118 (true, true, 0),
4119 );
4120 }
4121
4122 #[tokio::test]
4123 async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
4124 let mut schema_builder = crate::dsl::SchemaBuilder::default();
4125 let field = schema_builder.add_text_field("text", true, true);
4126 let schema = schema_builder.build();
4127 let directory = crate::directories::RamDirectory::new();
4128 let config = crate::index::IndexConfig {
4129 num_indexing_threads: 1,
4130 merge_policy: Box::new(crate::merge::NoMergePolicy),
4131 ..Default::default()
4132 };
4133 let mut writer = crate::index::IndexWriter::create(directory, schema, config)
4134 .await
4135 .unwrap();
4136 for value in ["one", "two", "three"] {
4137 let mut document = crate::dsl::Document::new();
4138 document.add_text(field, value);
4139 writer.add_document(document).unwrap();
4140 writer.commit().await.unwrap();
4141 }
4142
4143 let manager = Arc::clone(writer.segment_manager());
4144 let held_id = manager.get_segment_ids().await.pop().unwrap();
4145 let mut held = Some(
4146 manager
4147 .active_operations
4148 .try_register(vec![held_id])
4149 .unwrap(),
4150 );
4151 let batches = Arc::new(AtomicUsize::new(0));
4152 let batch_count = Arc::clone(&batches);
4153 writer
4154 .force_merge_with_snapshot_refresh(move || {
4155 let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
4156 if refresh == 2 {
4159 drop(held.take());
4160 }
4161 std::future::ready(Ok(()))
4162 })
4163 .await
4164 .unwrap();
4165
4166 assert_eq!(manager.get_segment_ids().await.len(), 1);
4167 assert_eq!(
4168 batches.load(Ordering::Relaxed),
4169 4,
4170 "initial/final refreshes plus two replacements are required after the held source releases"
4171 );
4172 }
4173
4174 #[tokio::test]
4175 async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
4176 let mut schema_builder = crate::dsl::SchemaBuilder::default();
4177 schema_builder.set_reorder_on_merge(true);
4178 let schema = schema_builder.build();
4179 let mut metadata = IndexMetadata::new(schema.clone());
4180 metadata.add_segment("00000000000000000000000000000001".into(), 1);
4181 metadata.add_segment("00000000000000000000000000000002".into(), 1);
4182
4183 let global_merge_permits = Arc::new(Semaphore::new(1));
4184 let manager = Arc::new(SegmentManager::new(
4185 Arc::new(crate::directories::RamDirectory::new()),
4186 Arc::new(schema),
4187 metadata,
4188 Box::new(crate::merge::NoMergePolicy),
4189 0,
4190 1,
4191 Arc::clone(&global_merge_permits),
4192 None,
4193 1024,
4194 Arc::new(ReorderConcurrencyGate::new(1)),
4195 None,
4196 ));
4197
4198 manager.active_operations.pause_non_indexing();
4203 let force_merge = {
4204 let manager = Arc::clone(&manager);
4205 tokio::spawn(async move { manager.force_merge().await })
4206 };
4207 tokio::time::timeout(std::time::Duration::from_secs(1), async {
4208 while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4209 tokio::task::yield_now().await;
4210 }
4211 })
4212 .await
4213 .expect("force merge never reached the paused group claim");
4214
4215 assert_eq!(
4216 global_merge_permits.available_permits(),
4217 1,
4218 "force merge retained global capacity while vector staging blocked group ownership"
4219 );
4220
4221 force_merge.abort();
4222 let _ = force_merge.await;
4223 manager.active_operations.resume_non_indexing();
4224 }
4225
4226 #[test]
4227 fn output_cleanup_guard_runs_during_panic_unwind() {
4228 let cleaned = Arc::new(AtomicBool::new(false));
4229 let cleaned_in_callback = Arc::clone(&cleaned);
4230 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4231 cleaned_in_callback.store(true, Ordering::SeqCst);
4232 });
4233
4234 let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4235 let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4236 panic!("simulated reorder panic");
4237 }));
4238
4239 assert!(result.is_err());
4240 assert!(
4241 cleaned.load(Ordering::SeqCst),
4242 "partial output cleanup must run during unwind"
4243 );
4244 }
4245
4246 #[test]
4247 fn output_cleanup_guard_disarms_after_commit() {
4248 let cleaned = Arc::new(AtomicBool::new(false));
4249 let cleaned_in_callback = Arc::clone(&cleaned);
4250 let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4251 cleaned_in_callback.store(true, Ordering::SeqCst);
4252 });
4253
4254 {
4255 let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4256 guard.disarm();
4257 }
4258
4259 assert!(!cleaned.load(Ordering::SeqCst));
4260 }
4261
4262 #[test]
4263 fn test_active_operation_guard_releases_ownership() {
4264 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4265 {
4266 let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4267 let snap = active.snapshot();
4268 assert!(snap.contains("a"));
4269 assert!(snap.contains("b"));
4270 }
4271 assert!(active.snapshot().is_empty());
4272 }
4273
4274 #[test]
4275 fn test_non_overlapping_operations_can_run_concurrently() {
4276 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4277 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4278 let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4279 let snap = active.snapshot();
4280 assert_eq!(snap.len(), 4);
4281
4282 drop(first);
4283 let snap = active.snapshot();
4284 assert_eq!(snap.len(), 2);
4285 assert!(snap.contains("c"));
4286 assert!(snap.contains("d"));
4287 }
4288
4289 #[test]
4290 fn test_overlapping_operation_is_rejected_until_release() {
4291 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4292 let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4293 assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4294 drop(first);
4295 assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4296 }
4297
4298 #[test]
4299 fn test_active_operation_snapshot() {
4300 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4301 let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4302 let snap = active.snapshot();
4303 assert!(snap.contains("x"));
4304 assert!(snap.contains("y"));
4305 assert!(!snap.contains("z"));
4306 }
4307
4308 #[tokio::test]
4309 async fn operation_barrier_ignores_producers_started_after_snapshot() {
4310 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4311 let before_gate = active.try_register(vec!["old".into()]).unwrap();
4312 let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4313 assert_eq!(parked_indexing, 0);
4314 let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4315
4316 let waiter = {
4317 let active = Arc::clone(&active);
4318 tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4319 };
4320 tokio::task::yield_now().await;
4321 assert!(!waiter.is_finished());
4322
4323 drop(before_gate);
4324 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4325 .await
4326 .expect("pre-gate operation barrier was starved by a post-gate producer")
4327 .unwrap();
4328 assert!(active.snapshot().contains("new-flat"));
4329 drop(after_gate);
4330 }
4331
4332 #[tokio::test]
4333 async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4334 let manager = lifecycle_test_manager();
4335 let current = manager.published_generation();
4336 manager
4337 .published_generation
4338 .store(Arc::new(PublishedIndexGeneration {
4339 publication_id: current.publication_id,
4340 schema: current.schema.clone(),
4341 trained_vectors: Some(Arc::new(TrainedVectorStructures {
4342 centroids: rustc_hash::FxHashMap::default(),
4343 binary_quantizers: rustc_hash::FxHashMap::default(),
4344 ..Default::default()
4345 })),
4346 }));
4347
4348 let guard = manager.begin_vector_artifact_update().await.unwrap();
4349 assert!(
4350 manager.trained().is_some(),
4351 "search readers keep the last fully validated generation"
4352 );
4353 assert!(
4354 manager.trained_for_segment_build().is_none(),
4355 "new segment producers must stay flat during an artifact update"
4356 );
4357
4358 let detached_transaction_guard = guard.clone();
4359 drop(guard);
4360 assert!(
4361 manager.trained_for_segment_build().is_none(),
4362 "a detached lifecycle transaction must retain the producer gate after request cancellation"
4363 );
4364 drop(detached_transaction_guard);
4365 assert!(manager.trained_for_segment_build().is_some());
4366 }
4367
4368 #[tokio::test]
4369 async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4370 let manager = lifecycle_test_manager();
4371 let guard = manager.begin_vector_artifact_update().await.unwrap();
4372 assert!(
4373 manager
4374 .active_operations
4375 .try_register(vec!["merge".into()])
4376 .is_none(),
4377 "ordinary merge/reorder work must not change staged sources"
4378 );
4379 let indexing = manager
4380 .active_operations
4381 .try_register_indexing(vec!["fresh".into()])
4382 .expect("indexing remains available in flat mode");
4383 drop(indexing);
4384
4385 drop(guard);
4386 assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4387 assert!(
4388 manager
4389 .active_operations
4390 .try_register(vec!["merge".into()])
4391 .is_some()
4392 );
4393 }
4394
4395 #[tokio::test]
4396 async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4397 let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4398 let guard = active.try_register(vec!["live".into()]).unwrap();
4399 let cancellation = active.cancellation_flag();
4400 active.stop_accepting();
4401 assert!(cancellation.load(Ordering::Acquire));
4402 assert!(active.try_register(vec!["new".into()]).is_none());
4403
4404 let waiter = {
4405 let active = Arc::clone(&active);
4406 tokio::spawn(async move { active.wait_until_idle().await })
4407 };
4408 tokio::task::yield_now().await;
4409 assert!(!waiter.is_finished());
4410 drop(guard);
4411 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4412 .await
4413 .expect("shutdown waiter missed the final guard notification")
4414 .unwrap();
4415 }
4416
4417 #[tokio::test]
4418 async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4419 let manager = lifecycle_test_manager();
4420 let started = Arc::new(Semaphore::new(0));
4421 let release = Arc::new(Semaphore::new(0));
4422 let completed = Arc::new(AtomicBool::new(false));
4423
4424 let request = {
4425 let manager = Arc::clone(&manager);
4426 let started = Arc::clone(&started);
4427 let release = Arc::clone(&release);
4428 let completed = Arc::clone(&completed);
4429 tokio::spawn(async move {
4430 manager
4431 .run_lifecycle_transaction(async move {
4432 started.add_permits(1);
4433 let _permit = release.acquire().await.unwrap();
4434 completed.store(true, Ordering::Release);
4435 Ok(())
4436 })
4437 .await
4438 })
4439 };
4440
4441 let _started = started.acquire().await.unwrap();
4442 request.abort();
4443 assert!(request.await.unwrap_err().is_cancelled());
4444 release.add_permits(1);
4445
4446 manager.begin_shutdown();
4447 tokio::time::timeout(
4448 std::time::Duration::from_secs(1),
4449 manager.wait_for_shutdown(),
4450 )
4451 .await
4452 .expect("shutdown did not drain detached lifecycle transaction");
4453 assert!(completed.load(Ordering::Acquire));
4454 }
4455
4456 #[tokio::test]
4457 async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4458 let manager = lifecycle_test_manager();
4459 {
4460 let mut state = manager.state.lock().await;
4461 state.metadata.add_segment_meta(
4462 "eligible".into(),
4463 SegmentMetaInfo {
4464 num_docs: 10,
4465 ancestors: Vec::new(),
4466 generation: 1,
4467 reordered: true,
4468 bp_converged: false,
4469 bp_unconverged_passes: 2,
4470 },
4471 );
4472 state.metadata.add_segment_meta(
4473 "at-limit".into(),
4474 SegmentMetaInfo {
4475 num_docs: 20,
4476 ancestors: Vec::new(),
4477 generation: 1,
4478 reordered: true,
4479 bp_converged: false,
4480 bp_unconverged_passes: 3,
4481 },
4482 );
4483 state.metadata.add_segment_meta(
4484 "carried-debt".into(),
4485 SegmentMetaInfo {
4486 num_docs: 15,
4487 ancestors: Vec::new(),
4488 generation: 2,
4489 reordered: false,
4490 bp_converged: false,
4491 bp_unconverged_passes: 2,
4492 },
4493 );
4494 state.metadata.add_segment_meta(
4495 "carried-debt-at-limit".into(),
4496 SegmentMetaInfo {
4497 num_docs: 25,
4498 ancestors: Vec::new(),
4499 generation: 2,
4500 reordered: false,
4501 bp_converged: false,
4502 bp_unconverged_passes: 3,
4503 },
4504 );
4505 state.metadata.add_segment_meta(
4506 "converged".into(),
4507 SegmentMetaInfo {
4508 num_docs: 30,
4509 ancestors: Vec::new(),
4510 generation: 1,
4511 reordered: true,
4512 bp_converged: true,
4513 bp_unconverged_passes: 0,
4514 },
4515 );
4516 state.metadata.add_segment("fresh".into(), 40);
4517 }
4518
4519 assert_eq!(
4520 manager.unreordered_segments().await,
4521 vec![("fresh".into(), 40)],
4522 "a block-copy output with BP debt is not a fresh first-pass candidate",
4523 );
4524 let mut eligible = manager.unconverged_segments_below(3).await;
4525 eligible.sort_unstable();
4526 assert_eq!(
4527 eligible,
4528 vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4529 );
4530 assert!(manager.unconverged_segments_below(0).await.is_empty());
4531 }
4532
4533 #[test]
4534 fn merge_retry_backoff_is_exponential_and_capped() {
4535 assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4536 assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4537 assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4538 assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4539 }
4540
4541 #[test]
4542 fn only_deterministic_source_errors_are_quarantined() {
4543 assert!(is_deterministic_source_error(&Error::Corruption(
4544 "bad footer".into()
4545 )));
4546 assert!(is_deterministic_source_error(&Error::Io(
4547 std::io::Error::from(std::io::ErrorKind::NotFound)
4548 )));
4549 assert!(!is_deterministic_source_error(&Error::Io(
4550 std::io::Error::from(std::io::ErrorKind::TimedOut)
4551 )));
4552 assert!(!is_deterministic_source_error(&Error::Io(
4553 std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4554 )));
4555 }
4556
4557 #[test]
4558 fn transient_reorder_failure_is_backed_off_until_cleared() {
4559 let manager = lifecycle_test_manager();
4560 manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4561 assert!(manager.paused_reorder_segments().contains("source"));
4562 manager.clear_reorder_retry("source");
4563 assert!(!manager.paused_reorder_segments().contains("source"));
4564 }
4565
4566 #[derive(Default)]
4569 struct FailingExistsDirectory(crate::directories::RamDirectory);
4570
4571 #[async_trait::async_trait]
4572 impl crate::directories::Directory for FailingExistsDirectory {
4573 async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4574 Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4575 }
4576
4577 async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4578 self.0.file_size(path).await
4579 }
4580
4581 async fn open_read(
4582 &self,
4583 path: &std::path::Path,
4584 ) -> std::io::Result<crate::directories::FileHandle> {
4585 self.0.open_read(path).await
4586 }
4587
4588 async fn read_range(
4589 &self,
4590 path: &std::path::Path,
4591 range: std::ops::Range<u64>,
4592 ) -> std::io::Result<crate::directories::OwnedBytes> {
4593 self.0.read_range(path, range).await
4594 }
4595
4596 async fn list_files(
4597 &self,
4598 prefix: &std::path::Path,
4599 ) -> std::io::Result<Vec<std::path::PathBuf>> {
4600 self.0.list_files(prefix).await
4601 }
4602
4603 async fn open_lazy(
4604 &self,
4605 path: &std::path::Path,
4606 ) -> std::io::Result<crate::directories::FileHandle> {
4607 self.0.open_lazy(path).await
4608 }
4609 }
4610
4611 #[async_trait::async_trait]
4612 impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4613 async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4614 self.0.write(path, data).await
4615 }
4616
4617 async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4618 self.0.delete(path).await
4619 }
4620
4621 async fn rename(
4622 &self,
4623 from: &std::path::Path,
4624 to: &std::path::Path,
4625 ) -> std::io::Result<()> {
4626 self.0.rename(from, to).await
4627 }
4628
4629 async fn sync(&self) -> std::io::Result<()> {
4630 self.0.sync().await
4631 }
4632
4633 async fn streaming_writer(
4634 &self,
4635 path: &std::path::Path,
4636 ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4637 self.0.streaming_writer(path).await
4638 }
4639 }
4640
4641 #[derive(Debug, Clone)]
4642 struct MergeEverythingPolicy;
4643
4644 impl MergePolicy for MergeEverythingPolicy {
4645 fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4646 if segments.len() < 2 {
4647 return Vec::new();
4648 }
4649 vec![crate::merge::MergeCandidate {
4650 segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4651 }]
4652 }
4653
4654 fn clone_box(&self) -> Box<dyn MergePolicy> {
4655 Box::new(self.clone())
4656 }
4657 }
4658
4659 #[tokio::test]
4660 async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4661 let manager = lifecycle_test_manager();
4662 let parked_indexing = manager
4667 .protect_new_segment("00000000000000000000000000000abc".into())
4668 .unwrap();
4669
4670 let error = tokio::time::timeout(
4671 std::time::Duration::from_secs(2),
4672 manager.begin_vector_artifact_update(),
4673 )
4674 .await
4675 .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4676 .err()
4677 .expect("an old-generation prepared segment must block artifact replacement")
4678 .to_string();
4679 assert!(error.contains("built but uncommitted"), "{error}");
4680 assert!(
4681 !manager.vector_artifact_update.load(Ordering::Acquire),
4682 "a rejected update must release the producer gate"
4683 );
4684
4685 drop(parked_indexing);
4686
4687 let guard = manager
4688 .begin_vector_artifact_update()
4689 .await
4690 .expect("artifact update should succeed after the pending generation is resolved");
4691 drop(guard);
4692 }
4693
4694 #[tokio::test]
4695 async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4696 let manager = lifecycle_test_manager();
4697 let merge_like = manager
4698 .active_operations
4699 .try_register(vec!["merge-source".into()])
4700 .unwrap();
4701
4702 let waiter = {
4703 let manager = Arc::clone(&manager);
4704 tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4705 };
4706 for _ in 0..8 {
4707 tokio::task::yield_now().await;
4708 }
4709 assert!(
4710 !waiter.is_finished(),
4711 "artifact update must drain merge/reorder producers that may hold the previous generation"
4712 );
4713
4714 drop(merge_like);
4715 tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4716 .await
4717 .expect("artifact update missed the lifecycle guard release")
4718 .unwrap()
4719 .unwrap();
4720 }
4721
4722 #[tokio::test]
4723 async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4724 let manager = lifecycle_test_manager();
4725 let release = Arc::new(Semaphore::new(0));
4726 let merge_task = {
4727 let release = Arc::clone(&release);
4728 tokio::spawn(async move {
4729 let _permit = release.acquire().await.unwrap();
4730 })
4731 };
4732 manager.merge_handles.lock().push(merge_task);
4733
4734 let waiter = {
4735 let manager = Arc::clone(&manager);
4736 tokio::spawn(async move { manager.wait_for_all_merges().await })
4737 };
4738 for _ in 0..8 {
4739 tokio::task::yield_now().await;
4740 }
4741 assert!(!waiter.is_finished());
4742 waiter.abort();
4745 let join_error = waiter.await.unwrap_err();
4746 assert!(join_error.is_cancelled());
4747
4748 assert!(
4749 !manager.merge_handles.lock().is_empty(),
4750 "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4751 );
4752
4753 release.add_permits(1);
4755 tokio::time::timeout(
4756 std::time::Duration::from_secs(1),
4757 manager.wait_for_all_merges(),
4758 )
4759 .await
4760 .expect("subsequent drain missed the reinserted merge handle");
4761 assert!(manager.merge_handles.lock().is_empty());
4762 }
4763
4764 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4765 async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4766 let manager = lifecycle_test_manager();
4767 {
4768 let mut state = manager.state.lock().await;
4769 state
4770 .metadata
4771 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4772 state
4773 .metadata
4774 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4775 }
4776 let reorder_like = manager
4779 .active_operations
4780 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4781 .unwrap();
4782
4783 let force_merge = {
4784 let manager = Arc::clone(&manager);
4785 tokio::spawn(async move { manager.force_merge().await })
4786 };
4787
4788 tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4789 let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4790 assert!(
4791 retries >= 1,
4792 "force_merge never observed the conflicting owner (retries={retries})"
4793 );
4794 assert!(
4795 retries < 20,
4796 "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4797 );
4798
4799 drop(reorder_like);
4800 let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4803 .await
4804 .expect("force_merge kept spinning after the conflicting owner released")
4805 .unwrap();
4806 assert!(result.is_err());
4807 }
4808
4809 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4810 async fn force_merge_routes_around_segments_held_by_reorder() {
4811 let manager = lifecycle_test_manager();
4812 {
4813 let mut state = manager.state.lock().await;
4814 state
4815 .metadata
4816 .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4817 state
4818 .metadata
4819 .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4820 state
4821 .metadata
4822 .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4823 }
4824 let _reorder_like = manager
4827 .active_operations
4828 .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4829 .unwrap();
4830
4831 let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4838 let manager = Arc::clone(&manager);
4839 async move { manager.force_merge().await }
4840 })
4841 .await
4842 .expect("force_merge livelocked on a segment held by an active reorder");
4843 assert!(result.is_err(), "fake segment files must fail the merge");
4844
4845 assert_eq!(
4846 manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4847 0,
4848 "batch built from the ownership snapshot must not collide with the held segment"
4849 );
4850 }
4851
4852 #[tokio::test]
4853 async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4854 let schema = crate::dsl::SchemaBuilder::default().build();
4855 let mut metadata = IndexMetadata::new(schema.clone());
4856 metadata.add_segment("00000000000000000000000000000001".into(), 10);
4857 metadata.add_segment("00000000000000000000000000000002".into(), 10);
4858 let manager = Arc::new(SegmentManager::new(
4859 Arc::new(FailingExistsDirectory::default()),
4860 Arc::new(schema),
4861 metadata,
4862 Box::new(MergeEverythingPolicy),
4863 0,
4864 1,
4865 Arc::new(Semaphore::new(1)),
4866 None,
4867 1024,
4868 Arc::new(ReorderConcurrencyGate::new(1)),
4869 None,
4870 ));
4871
4872 manager.maybe_merge().await;
4875
4876 tokio::time::timeout(
4877 std::time::Duration::from_secs(5),
4878 manager.wait_for_all_merges(),
4879 )
4880 .await
4881 .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4882 assert!(
4883 manager.merge_retry_is_paused(),
4884 "the failed merge should have armed the retry backoff"
4885 );
4886
4887 manager.begin_shutdown();
4889 tokio::time::timeout(
4890 std::time::Duration::from_secs(5),
4891 manager.wait_for_shutdown(),
4892 )
4893 .await
4894 .expect("shutdown did not drain the merge retry wakeup task");
4895 }
4896}