Skip to main content

summa_core/merge/
segment_manager.rs

1//! Segment manager — coordinates segment commit, background merging, and trained structures.
2//!
3//! Architecture:
4//! - **Single mutation queue**: All metadata mutations serialize through `tokio::sync::Mutex<ManagerState>`.
5//! - **Active-operation ownership**: Every segment that is being built, merged,
6//!   or reordered is registered before its first file is written and remains
7//!   registered until it is either published in metadata or abandoned.
8//! - **Concurrent merges**: Multiple non-overlapping merges can run in parallel.
9//!   New merges are rejected only if they share segments with an active operation.
10//! - **Auto-trigger**: Each completed merge re-evaluates the merge policy and spawns
11//!   new merges if eligible (cascading merges for higher tiers).
12//! - **ArcSwap for trained**: Lock-free reads of trained vector structures.
13//!
14//! # Segment lifecycle invariant
15//!
16//! Every on-disk `seg_*` ID must be protected by at least one of these owners:
17//!
18//! 1. `state.metadata` while the segment is live and searchable;
19//! 2. `active_operations` while an indexing/merge/reorder task owns it; or
20//! 3. `tracker` while a retired segment is still visible to a reader or its
21//!    filesystem deletion is scheduled.
22//!
23//! An ID with no owner is an orphan and may be swept. Transitions are ordered
24//! so the new owner is installed before the old owner is released.
25//!
26//! # Locking model (deadlock-free by construction)
27//!
28//! ```text
29//! Lock ordering (acquire in this order):
30//!   1. state               — tokio::sync::Mutex, held for mutations + disk I/O
31//!   2. active_operations   — parking_lot::Mutex (sync), sub-μs hold, RAII guard
32//!   3. tracker.inner       — parking_lot::Mutex (sync), sub-μs hold
33//!
34//! Independent bookkeeping (never held with `state`):
35//!   trained                — arc_swap::ArcSwapOption, lock-free
36//!   vector_artifact_update — atomic producer gate, held across manual training
37//!   merge_handles          — parking_lot::Mutex, synchronous short hold
38//!   lifecycle_handles      — parking_lot::Mutex, synchronous short hold
39//!   merge/reorder permits  — tokio semaphores shared by configuration
40//! ```
41//!
42//! **Rule:** Never hold a sync lock while `.await`-ing.
43
44use 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
73/// Deterministically pack segments into near-minimal final outputs.
74///
75/// Best-fit decreasing avoids the pathological smallest-first behavior where,
76/// for example, `4 + 4` is merged before two `6 + 4` pairs under a 10-doc
77/// limit. That old order both left excess segments and made its intermediate
78/// outputs eligible for another full BP pass. Bin packing is NP-hard; BFD is a
79/// bounded, deterministic approximation that produces maximal groups (no two
80/// output groups can still fit together).
81fn 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    // Preserve deterministic small-to-large source order inside each output
119    // so logical document rebasing does not depend on merge grouping.
120    for group in &mut groups {
121        group
122            .segments
123            .sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
124                left_docs
125                    .cmp(right_docs)
126                    .then_with(|| left_id.cmp(right_id))
127            });
128    }
129
130    // Smallest outputs first limit transient disk headroom. Tie-break on the
131    // first (deterministically ordered) segment ID for reproducible plans.
132    groups.sort_unstable_by(|left, right| {
133        left.total_docs
134            .cmp(&right.total_docs)
135            .then_with(|| left.segments[0].0.cmp(&right.segments[0].0))
136    });
137    groups
138}
139
140fn force_merge_output_count(source_count: usize) -> usize {
141    if source_count < 2 {
142        return 0;
143    }
144    (source_count - 1).div_ceil(FORCE_MERGE_MAX_FAN_IN - 1)
145}
146
147#[derive(Debug)]
148struct ForceMergeStep {
149    /// Source or earlier-step node IDs, in final document order.
150    inputs: Vec<usize>,
151}
152
153#[derive(Debug)]
154struct ForceMergeHierarchy {
155    steps: Vec<ForceMergeStep>,
156    root: usize,
157}
158
159/// Build a shallow ordered 64-ary reduction tree with the minimum possible
160/// number of merge outputs. All internal nodes have full fan-in except one
161/// deepest partial node; distributing internal nodes evenly keeps source
162/// rewrite depth balanced without changing concatenation order.
163fn plan_force_merge_hierarchy(source_count: usize) -> ForceMergeHierarchy {
164    debug_assert!(source_count >= 2);
165    let internal_count = force_merge_output_count(source_count);
166    let max_leaves = 1usize
167        .checked_add(internal_count.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
168        .expect("force-merge hierarchy size exceeds usize");
169    let deficit = max_leaves - source_count;
170    debug_assert!(deficit < FORCE_MERGE_MAX_FAN_IN - 1);
171
172    fn build(
173        leaf_start: usize,
174        leaf_count: usize,
175        internal_count: usize,
176        deficit: usize,
177        source_count: usize,
178        steps: &mut Vec<ForceMergeStep>,
179    ) -> usize {
180        debug_assert!(internal_count > 0);
181        if internal_count == 1 {
182            let arity = FORCE_MERGE_MAX_FAN_IN - deficit;
183            debug_assert_eq!(leaf_count, arity);
184            debug_assert!((2..=FORCE_MERGE_MAX_FAN_IN).contains(&arity));
185            let output = source_count + steps.len();
186            steps.push(ForceMergeStep {
187                inputs: (leaf_start..leaf_start + arity).collect(),
188            });
189            return output;
190        }
191
192        // Keep this node full and place the one partial arity, if any, in the
193        // deepest/largest child. This minimizes bytes rewritten versus making
194        // the root partial (65 inputs become 2 + 63 leaves, then one final
195        // merge, rather than rewriting a 64-input prefix).
196        let child_internal_total = internal_count - 1;
197        let base = child_internal_total / FORCE_MERGE_MAX_FAN_IN;
198        let extra = child_internal_total % FORCE_MERGE_MAX_FAN_IN;
199        let mut child_internal = vec![base; FORCE_MERGE_MAX_FAN_IN];
200        for count in &mut child_internal[..extra] {
201            *count += 1;
202        }
203        let partial_child = (deficit > 0).then(|| {
204            child_internal
205                .iter()
206                .position(|&count| count > 0)
207                .expect("a non-root partial node requires an internal child")
208        });
209
210        let mut cursor = leaf_start;
211        let mut inputs = Vec::with_capacity(FORCE_MERGE_MAX_FAN_IN);
212        for (child, &child_internals) in child_internal.iter().enumerate() {
213            if child_internals == 0 {
214                inputs.push(cursor);
215                cursor += 1;
216                continue;
217            }
218            let child_deficit = usize::from(partial_child == Some(child)) * deficit;
219            let child_leaves = 1usize
220                .checked_add(child_internals.saturating_mul(FORCE_MERGE_MAX_FAN_IN - 1))
221                .and_then(|maximum| maximum.checked_sub(child_deficit))
222                .expect("force-merge child size exceeds usize");
223            inputs.push(build(
224                cursor,
225                child_leaves,
226                child_internals,
227                child_deficit,
228                source_count,
229                steps,
230            ));
231            cursor += child_leaves;
232        }
233        debug_assert_eq!(cursor, leaf_start + leaf_count);
234        let output = source_count + steps.len();
235        steps.push(ForceMergeStep { inputs });
236        output
237    }
238
239    let mut steps = Vec::with_capacity(internal_count);
240    let root = build(
241        0,
242        source_count,
243        internal_count,
244        deficit,
245        source_count,
246        &mut steps,
247    );
248    debug_assert_eq!(steps.len(), internal_count);
249    ForceMergeHierarchy { steps, root }
250}
251
252// ============================================================================
253// RAII active-operation tracking
254// ============================================================================
255
256/// Tracks every segment ID owned by an in-flight lifecycle operation.
257///
258/// Merge/reorder guards include both sources and output, providing mutual
259/// exclusion as well as orphan-sweep protection. Indexing guards contain the
260/// new output only and live from before the first write through commit/abort.
261struct ActiveOperationState {
262    segment_ids: HashSet<String>,
263    operation_tokens: HashSet<u64>,
264    /// Subset of `operation_tokens` owned by indexing producers. Their guards
265    /// travel with built-but-uncommitted `PreparedSegment`s and are released
266    /// only by a later commit/abort, so drain barriers must not wait on them:
267    /// the commit that would release them can be blocked on the barrier's own
268    /// caller (writer write lock / `&mut self`).
269    indexing_tokens: HashSet<u64>,
270    next_operation_token: u64,
271    accepting: bool,
272    /// Retraining stages a complete replacement segment generation. Ordinary
273    /// merge/reorder work is paused so its source set cannot change midway;
274    /// indexing producers remain allowed and deliberately emit flat vectors.
275    non_indexing_paused: bool,
276}
277
278struct ActiveSegmentOperations {
279    inner: parking_lot::Mutex<ActiveOperationState>,
280    idle: Notify,
281    shutdown: Notify,
282    shutdown_requested: Arc<AtomicBool>,
283    /// Owning index, so lifecycle decisions are attributable when several
284    /// indexes register and defer operations concurrently.
285    index_label: Arc<str>,
286}
287
288impl ActiveSegmentOperations {
289    fn new(index_label: Arc<str>) -> Self {
290        Self {
291            inner: parking_lot::Mutex::new(ActiveOperationState {
292                segment_ids: HashSet::new(),
293                operation_tokens: HashSet::new(),
294                indexing_tokens: HashSet::new(),
295                next_operation_token: 0,
296                accepting: true,
297                non_indexing_paused: false,
298            }),
299            idle: Notify::new(),
300            shutdown: Notify::new(),
301            shutdown_requested: Arc::new(AtomicBool::new(false)),
302            index_label,
303        }
304    }
305
306    /// Try to claim IDs for a self-draining lifecycle operation (merge,
307    /// reorder, cleanup). Returns a guard on success, `None` if any requested
308    /// ID is already owned by another active operation.
309    fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
310        self.try_register_kind(segment_ids, false, false)
311    }
312
313    /// Try to claim IDs for an indexing producer whose guard is held until
314    /// metadata publication (commit) rather than task completion.
315    fn try_register_indexing(
316        self: &Arc<Self>,
317        segment_ids: Vec<String>,
318    ) -> Option<SegmentOperationGuard> {
319        self.try_register_kind(segment_ids, true, false)
320    }
321
322    /// Claim source/output IDs for the exclusive vector-generation updater,
323    /// which is the only non-indexing producer allowed through its own pause.
324    fn try_register_vector_update(
325        self: &Arc<Self>,
326        segment_ids: Vec<String>,
327    ) -> Option<SegmentOperationGuard> {
328        self.try_register_kind(segment_ids, false, true)
329    }
330
331    fn try_register_kind(
332        self: &Arc<Self>,
333        segment_ids: Vec<String>,
334        indexing: bool,
335        vector_update: bool,
336    ) -> Option<SegmentOperationGuard> {
337        let mut inner = self.inner.lock();
338        if !inner.accepting {
339            log::debug!(
340                "[segment_lifecycle] index={} rejected operation during shutdown",
341                self.index_label
342            );
343            return None;
344        }
345        if !indexing && !vector_update && inner.non_indexing_paused {
346            log::debug!(
347                "[segment_lifecycle] index={} deferred operation during dense vector retraining",
348                self.index_label
349            );
350            return None;
351        }
352        // Check for overlap with any active lifecycle operation.
353        for id in &segment_ids {
354            if inner.segment_ids.contains(id) {
355                log::debug!(
356                    "[segment_lifecycle] index={} rejected: {} overlaps with an active operation ({} active IDs)",
357                    self.index_label,
358                    id,
359                    inner.segment_ids.len()
360                );
361                return None;
362            }
363        }
364        log::debug!(
365            "[segment_lifecycle] index={} registered {} IDs (total active: {})",
366            self.index_label,
367            segment_ids.len(),
368            inner.segment_ids.len() + segment_ids.len()
369        );
370        let operation_token = inner.next_operation_token;
371        let next_operation_token = operation_token.checked_add(1)?;
372        for id in &segment_ids {
373            inner.segment_ids.insert(id.clone());
374        }
375        inner.next_operation_token = next_operation_token;
376        inner.operation_tokens.insert(operation_token);
377        if indexing {
378            inner.indexing_tokens.insert(operation_token);
379        }
380        Some(SegmentOperationGuard {
381            active_operations: Arc::clone(self),
382            segment_ids,
383            operation_token,
384        })
385    }
386
387    /// Snapshot of all IDs owned by active operations.
388    fn snapshot(&self) -> HashSet<String> {
389        self.inner.lock().segment_ids.clone()
390    }
391
392    /// Exact identities of self-draining operations (merge/reorder/cleanup)
393    /// active at one instant, plus the number of indexing tokens excluded.
394    /// Unlike segment IDs, tokens cannot be reused by a later retry, so an
395    /// artifact-update barrier can drain only pre-gate producers without being
396    /// starved by new flat producers.
397    ///
398    /// Indexing tokens are deliberately excluded: their guards are parked in
399    /// built-but-uncommitted `PreparedSegment`s and only a later commit — which
400    /// may be blocked on the barrier's caller — releases them, so waiting on
401    /// them deadlocks (see `begin_vector_artifact_update`).
402    fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
403        let inner = self.inner.lock();
404        let tokens = inner
405            .operation_tokens
406            .difference(&inner.indexing_tokens)
407            .copied()
408            .collect();
409        (tokens, inner.indexing_tokens.len())
410    }
411
412    /// Atomically prevent new lifecycle work from starting. Existing guards
413    /// remain valid and can be drained with [`Self::wait_until_idle`].
414    fn stop_accepting(&self) {
415        self.shutdown_requested.store(true, Ordering::Release);
416        let mut inner = self.inner.lock();
417        inner.accepting = false;
418        self.shutdown.notify_waiters();
419        if inner.segment_ids.is_empty() {
420            self.idle.notify_waiters();
421        }
422    }
423
424    fn pause_non_indexing(&self) {
425        self.inner.lock().non_indexing_paused = true;
426    }
427
428    fn resume_non_indexing(&self) {
429        self.inner.lock().non_indexing_paused = false;
430        self.idle.notify_waiters();
431    }
432
433    fn is_accepting(&self) -> bool {
434        self.inner.lock().accepting
435    }
436
437    fn cancellation_flag(&self) -> Arc<AtomicBool> {
438        Arc::clone(&self.shutdown_requested)
439    }
440
441    /// Wait until every operation that started before shutdown has released
442    /// its ownership. Register/check and notification are ordered to avoid a
443    /// missed wakeup between observing a non-empty set and awaiting.
444    async fn wait_until_idle(&self) {
445        loop {
446            let notified = self.idle.notified();
447            if self.inner.lock().segment_ids.is_empty() {
448                return;
449            }
450            notified.await;
451        }
452    }
453
454    async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
455        while !operations.is_empty() {
456            let notified = self.idle.notified();
457            if self.inner.lock().operation_tokens.is_disjoint(operations) {
458                return;
459            }
460            notified.await;
461        }
462    }
463
464    /// Resolve when shutdown starts, without missing a notification between
465    /// checking the state and registering the waiter.
466    async fn wait_for_shutdown(&self) {
467        loop {
468            let notified = self.shutdown.notified();
469            if !self.inner.lock().accepting {
470                return;
471            }
472            notified.await;
473        }
474    }
475}
476
477/// RAII ownership of segment IDs used by an active lifecycle operation.
478/// Dropping on success, error, cancellation, or panic makes abandoned outputs
479/// eligible for sweeping automatically.
480pub(crate) struct SegmentOperationGuard {
481    active_operations: Arc<ActiveSegmentOperations>,
482    segment_ids: Vec<String>,
483    operation_token: u64,
484}
485
486impl Drop for SegmentOperationGuard {
487    fn drop(&mut self) {
488        let mut inner = self.active_operations.inner.lock();
489        for id in &self.segment_ids {
490            inner.segment_ids.remove(id);
491        }
492        inner.operation_tokens.remove(&self.operation_token);
493        inner.indexing_tokens.remove(&self.operation_token);
494        // Token barriers need notification on every completion, not only the
495        // transition to complete global idleness.
496        self.active_operations.idle.notify_waiters();
497        if inner.segment_ids.is_empty() {
498            debug_assert!(inner.operation_tokens.is_empty());
499        }
500    }
501}
502
503/// Exclusive gate for an index-level trained-vector artifact update.
504///
505/// Segment producers consult this gate before capturing the current trained
506/// structures. Once the gate is raised, new producers deliberately emit flat
507/// vector data; waiting for already-active producers to drain then guarantees
508/// that the committed source set stays stable while replacements are staged.
509struct VectorArtifactUpdateLease {
510    updating: Arc<AtomicBool>,
511    active_operations: Arc<ActiveSegmentOperations>,
512}
513
514impl Drop for VectorArtifactUpdateLease {
515    fn drop(&mut self) {
516        self.updating.store(false, Ordering::Release);
517        self.active_operations.resume_non_indexing();
518    }
519}
520
521#[derive(Clone)]
522pub(crate) struct VectorArtifactUpdateGuard {
523    _lease: Arc<VectorArtifactUpdateLease>,
524}
525
526/// Merge-time/manual BP pools are shared by every index in this process.
527/// A pool per `SegmentManager` multiplied a 96-core host into two 48-thread
528/// merge pools plus the optimizer pool (200+ process threads in production).
529static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
530
531const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
532const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
533
534#[derive(Default)]
535struct MergeRetryState {
536    retry_after: Option<std::time::Instant>,
537    consecutive_failures: u32,
538}
539
540fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
541    let shift = consecutive_failures.saturating_sub(1).min(16);
542    MERGE_RETRY_BASE_DELAY
543        .checked_mul(1u32 << shift)
544        .unwrap_or(MERGE_RETRY_MAX_DELAY)
545        .min(MERGE_RETRY_MAX_DELAY)
546}
547
548/// Merge JoinHandles taken out of the shared list for draining.
549///
550/// Drain futures are awaited inline by RPC handlers (force_merge/reorder) and
551/// can be dropped at any await when a client disconnects. Handles are awaited
552/// through this guard and removed only after completion, so a cancelled drain
553/// returns every un-awaited (and possibly still-running) merge to the shared
554/// list instead of silently detaching it from shutdown, abort, and
555/// force-merge tracking.
556struct DrainedMergeHandles<'a> {
557    shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
558    drained: Vec<JoinHandle<()>>,
559}
560
561impl<'a> DrainedMergeHandles<'a> {
562    fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
563        let drained = std::mem::take(&mut *shared.lock());
564        Self { shared, drained }
565    }
566
567    fn is_empty(&self) -> bool {
568        self.drained.is_empty()
569    }
570
571    /// Await the next handle. It stays owned by this guard while being polled
572    /// and is discarded only once it has completed, so cancellation at the
573    /// await reinserts it via `Drop`.
574    async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
575        let handle = self.drained.last_mut()?;
576        let result = handle.await;
577        self.drained.pop();
578        Some(result)
579    }
580}
581
582impl Drop for DrainedMergeHandles<'_> {
583    fn drop(&mut self) {
584        if !self.drained.is_empty() {
585            self.shared.lock().append(&mut self.drained);
586        }
587    }
588}
589
590/// Spawn and register auxiliary lifecycle work as one synchronous operation.
591///
592/// Registering *after* `spawn` left a small deletion race: shutdown could
593/// observe an empty handle list while the newly spawned filesystem task was
594/// already running. Holding the handle-list mutex across `Handle::spawn`
595/// makes task creation visible to the drain before either side can proceed.
596fn try_spawn_lifecycle<F>(
597    handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
598    runtime: &tokio::runtime::Handle,
599    future: F,
600) -> bool
601where
602    F: std::future::Future<Output = ()> + Send + 'static,
603{
604    std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
605        let mut handles = handles.lock();
606        handles.retain(|handle| !handle.is_finished());
607        handles.push(runtime.spawn(future));
608    }))
609    .is_ok()
610}
611
612/// Deletes an uncommitted merge/reorder output if its task unwinds.
613///
614/// Normal `Result::Err` paths delete outputs synchronously so callers observe
615/// a clean directory before returning. This guard covers the path those
616/// branches cannot: a panic after output files have been created. The cleanup
617/// callback re-checks metadata before deleting, so a panic after a successful
618/// metadata commit cannot remove a live segment.
619struct OutputCleanupGuard {
620    segment_id: SegmentId,
621    cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
622}
623
624impl OutputCleanupGuard {
625    fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
626        Self {
627            segment_id,
628            cleanup: Some(cleanup),
629        }
630    }
631
632    fn disarm(&mut self) {
633        self.cleanup = None;
634    }
635}
636
637impl Drop for OutputCleanupGuard {
638    fn drop(&mut self) {
639        if let Some(cleanup) = self.cleanup.take() {
640            cleanup(self.segment_id);
641        }
642    }
643}
644
645/// All mutable state behind the single async Mutex.
646struct ManagerState {
647    metadata: IndexMetadata,
648    merge_policy: Box<dyn MergePolicy>,
649}
650
651type ReplacementRefresh = Arc<
652    dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>
653        + Send
654        + Sync,
655>;
656
657/// Reconcile a long-lived consumer after a durable segment replacement.
658/// This runs inside the same lifecycle-owned task as metadata publication, so
659/// cancelling the caller cannot leave the replacement published but unannounced.
660async fn refresh_replacement_topology(refresh: Option<ReplacementRefresh>, index_label: &str) {
661    let Some(refresh) = refresh else {
662        return;
663    };
664    let mut last_error = None;
665    for attempt in 0..3 {
666        match refresh().await {
667            Ok(()) => return,
668            Err(error) => {
669                last_error = Some(error);
670                if attempt < 2 {
671                    tokio::time::sleep(std::time::Duration::from_secs(1 << attempt)).await;
672                }
673            }
674        }
675    }
676    if let Some(error) = last_error {
677        log::warn!(
678            "[segment_lifecycle] index={index_label} replacement topology refresh failed after 3 attempts: {}",
679            error,
680        );
681    }
682}
683
684#[cfg(feature = "native")]
685struct MergeTaskError {
686    error: Error,
687    unavailable_segments: Vec<String>,
688}
689
690#[cfg(feature = "native")]
691impl MergeTaskError {
692    fn source(segment_id: String, error: Error) -> Self {
693        Self {
694            error,
695            unavailable_segments: vec![segment_id],
696        }
697    }
698
699    fn sources(segment_ids: Vec<String>, error: Error) -> Self {
700        Self {
701            error,
702            unavailable_segments: segment_ids,
703        }
704    }
705}
706
707#[cfg(feature = "native")]
708impl From<Error> for MergeTaskError {
709    fn from(error: Error) -> Self {
710        Self {
711            error,
712            unavailable_segments: Vec::new(),
713        }
714    }
715}
716
717#[cfg(feature = "native")]
718fn is_deterministic_source_error(error: &Error) -> bool {
719    matches!(error, Error::Corruption(_) | Error::Serialization(_))
720        || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
721}
722
723#[cfg(feature = "native")]
724fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
725    if is_deterministic_source_error(&error) {
726        MergeTaskError::source(segment_id, error)
727    } else {
728        // Timeouts, interrupted reads, permission changes, and other generic
729        // I/O failures may be transient. Back them off instead of quarantining
730        // a healthy metadata segment for the rest of the process lifetime.
731        MergeTaskError::from(error)
732    }
733}
734
735#[cfg(feature = "native")]
736type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
737
738#[derive(Clone, Copy)]
739enum ReplacementLayout {
740    /// No BP ran while combining these sources. The combined segment is not
741    /// globally reordered, but an interrupted BP lineage remains owed its
742    /// record-level deepening pass.
743    BlockCopy,
744    /// Stable row/record filtering; block boundaries may need new BP refinement.
745    Compacted,
746    /// A BP pass produced the replacement layout.
747    BpReordered { converged: bool },
748    /// Seismic/ANN maintenance preserves completed or capped text/BMP order.
749    MaintenanceOnly,
750    /// A vector-only rewrite leaves document order and sparse layout exactly
751    /// unchanged, so its persisted BP progress must not be reset or advanced.
752    PreserveSingleSource,
753}
754
755fn replacement_bp_state(
756    parent_has_debt: bool,
757    parent_unconverged_passes: u32,
758    layout: ReplacementLayout,
759) -> (bool, bool, u32) {
760    match layout {
761        ReplacementLayout::Compacted => {
762            unreachable!("compaction preserves its single source lineage")
763        }
764        ReplacementLayout::BlockCopy => (
765            false,
766            !parent_has_debt,
767            if parent_has_debt {
768                parent_unconverged_passes
769            } else {
770                0
771            },
772        ),
773        ReplacementLayout::BpReordered { converged } => (
774            true,
775            converged,
776            if converged {
777                0
778            } else {
779                parent_unconverged_passes.saturating_add(1)
780            },
781        ),
782        ReplacementLayout::PreserveSingleSource | ReplacementLayout::MaintenanceOnly => {
783            unreachable!("preserved layouts retain the complete source metadata")
784        }
785    }
786}
787
788/// Account only a replacement that is about to publish. The caller installs
789/// these counters with the same durable metadata transaction as the new owner.
790fn replacement_seismic_state<'a>(
791    parents: impl Iterator<Item = &'a SegmentMetaInfo>,
792    pending_terms: u32,
793    layout: ReplacementLayout,
794) -> (u32, u32) {
795    if pending_terms == 0 {
796        return (0, 0);
797    }
798    let mut parent_count = 0usize;
799    let mut parent_pending_terms = 0u64;
800    let mut passes = 0u32;
801    let mut stalls = 0u32;
802    for parent in parents {
803        parent_count += 1;
804        parent_pending_terms += u64::from(parent.seismic_pending_terms);
805        passes = passes.max(parent.seismic_maintenance_passes);
806        stalls = stalls.max(parent.seismic_no_progress_passes);
807    }
808    // New merge inputs change the work available. An unchanged one-source
809    // replacement must not bypass a previous stall limit.
810    let made_progress = u64::from(pending_terms) < parent_pending_terms;
811    if parent_count > 1 || made_progress {
812        stalls = 0;
813    }
814    if matches!(
815        layout,
816        ReplacementLayout::BpReordered { .. } | ReplacementLayout::MaintenanceOnly
817    ) {
818        passes = passes.saturating_add(1);
819        if parent_count == 1 && !made_progress {
820            stalls = stalls.saturating_add(1);
821        }
822    }
823    (passes, stalls)
824}
825
826#[derive(Clone, Copy, Debug, Eq, PartialEq)]
827enum VectorSegmentRewriteOutcome {
828    Rewritten,
829    AlreadyCurrent,
830    SourceGone,
831    Conflict,
832    Deferred,
833}
834
835/// Complete but unpublished vector-only replacement. Its lifecycle claim and
836/// cleanup guard stay armed until the whole codebook generation commits.
837pub(crate) struct StagedVectorSegment {
838    source_id: String,
839    output_id: SegmentId,
840    doc_count: u32,
841    _operation: SegmentOperationGuard,
842    cleanup: OutputCleanupGuard,
843}
844
845/// Segment manager — coordinates segment commit, background merging, and trained structures.
846///
847/// SOLE owner of `metadata.json`. All metadata mutations go through `state` Mutex.
848pub struct SegmentManager<D: DirectoryWriter + 'static> {
849    /// Serializes ALL metadata mutations.
850    state: Arc<AsyncMutex<ManagerState>>,
851
852    /// RAII ownership for every in-flight segment lifecycle operation.
853    active_operations: Arc<ActiveSegmentOperations>,
854
855    /// Metadata-live segments involved in a deterministic source/corruption
856    /// failure. They stay searchable (and operator-visible) but are excluded
857    /// from merges for this process lifetime, preventing a bad candidate from
858    /// consuming full rewrite capacity on every retry.
859    quarantined_segments: parking_lot::Mutex<HashSet<String>>,
860
861    /// Generic merge failures pause scheduling briefly. Source-specific open
862    /// failures use `quarantined_segments` instead so healthy work can continue.
863    merge_retry: parking_lot::Mutex<MergeRetryState>,
864
865    /// Per-source backoff for non-deterministic standalone reorder failures.
866    /// Optimizer scans are periodic, but a pass can outlast the scan interval;
867    /// without completion-based backoff it would restart almost immediately.
868    reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
869
870    /// In-flight merge JoinHandles — supports multiple concurrent merges.
871    merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
872
873    /// At most one task per index waits for application-wide merge capacity.
874    /// Without this wakeup, an index denied by another index can remain idle
875    /// forever when no later commit happens to re-run merge policy evaluation.
876    global_merge_wakeup_pending: AtomicBool,
877
878    /// Non-zero while an explicit force merge is draining/running. Automatic
879    /// merges do not start during this window; otherwise a merge spawned
880    /// between the initial drain and foreground BP reservation could claim
881    /// source segments and then wait behind the foreground gate.
882    force_merge_active: AtomicUsize,
883
884    /// Times `force_merge` observed a conflicting active operation and retried.
885    /// Test-only observability for the conflict-retry backoff.
886    #[cfg(test)]
887    force_merge_conflict_retries: std::sync::atomic::AtomicU64,
888
889    /// Auxiliary lifecycle tasks: metadata transactions, deferred deletes,
890    /// and capacity wakeups. Handles registered here are drained before index
891    /// removal.
892    lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
893
894    /// Trained vector structures — lock-free reads via ArcSwap.
895    /// Wrapped in `Arc` so cancellation-safe metadata transactions can publish
896    /// the matching in-memory generation after their durable commit point.
897    published_generation: Arc<ArcSwap<PublishedIndexGeneration>>,
898
899    /// Raised while index-level trained artifacts and their metadata are being
900    /// replaced. Search readers keep using the last valid generation, while
901    /// segment producers fall back to flat output until publication completes.
902    vector_artifact_update: Arc<AtomicBool>,
903
904    /// Reference counting for safe segment deletion (sync Mutex for Drop).
905    tracker: Arc<SegmentTracker>,
906
907    /// Cached deletion callback for snapshots (avoids allocation per acquire_snapshot).
908    delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
909
910    /// Directory for segment I/O
911    directory: Arc<D>,
912    /// Schema for segment operations
913    schema: Arc<crate::dsl::Schema>,
914    /// Compression mode used for term dictionaries produced by merges.
915    optimization: crate::structures::IndexOptimization,
916    /// Posting codec used by new segments and merge re-encoding.
917    posting_codec: crate::structures::PostingCodec,
918    term_dict_block_size: crate::structures::SSTableBlockSize,
919    /// Term cache blocks for segment readers during merge
920    term_cache_blocks: usize,
921    term_cache_budget_bytes: Option<usize>,
922    /// Hard concurrency limit for background merges. A semaphore permit is
923    /// acquired before lifecycle ownership, closing the old handle-count race
924    /// where concurrent schedulers could exceed the configured maximum.
925    merge_permits: Arc<Semaphore>,
926    /// Application-wide merge limit shared across index managers.
927    global_merge_permits: Arc<Semaphore>,
928    /// Shared across every index opened from the same `IndexConfig`. This
929    /// bounds whole BP rewrites (optimizer + merge-time + manual) separately
930    /// from Rayon thread width, preventing N × memory-budget amplification.
931    reorder_permits: Arc<ReorderConcurrencyGate>,
932    /// Run BP reordering of `reorder`-attributed text and BMP fields inside merges.
933    /// Persisted index configuration (schema-level `reorder_on_merge: true`
934    /// in SDL); merged segments are marked `reordered` and skipped by the
935    /// standalone optimizer pass.
936    reorder_on_merge: bool,
937    /// Wall-clock budget for merge-time BP (from `IndexConfig`); truncated
938    /// passes mark the merged segment `bp_converged = false` so the
939    /// background optimizer deepens it later (warm-started).
940    merge_bp_time_budget: Option<std::time::Duration>,
941    /// Memory budget for the BP forward index (merge-time and background
942    /// reorder). Over-budget passes drop highest-df dims, logged loudly.
943    bp_memory_budget_bytes: usize,
944    /// Application-owned shared pool, when configured. This is the server
945    /// path and ensures optimizer and merge-time work use the same threads.
946    background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
947    /// Writer-owned topology reconciler. Every durable replacement invokes
948    /// it from tracked lifecycle work so background merges cannot leave
949    /// primary-key snapshots pinning retired sources indefinitely.
950    replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
951}
952
953struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
954
955impl Drop for ForceMergeActivityGuard<'_> {
956    fn drop(&mut self) {
957        self.0.fetch_sub(1, Ordering::AcqRel);
958    }
959}
960
961impl<D: DirectoryWriter + 'static> SegmentManager<D> {
962    /// Create a new segment manager with existing metadata
963    #[allow(clippy::too_many_arguments)]
964    pub fn new(
965        directory: Arc<D>,
966        schema: Arc<crate::dsl::Schema>,
967        metadata: IndexMetadata,
968        merge_policy: Box<dyn MergePolicy>,
969        term_cache_blocks: usize,
970        max_concurrent_merges: usize,
971        global_merge_permits: Arc<Semaphore>,
972        merge_bp_time_budget: Option<std::time::Duration>,
973        bp_memory_budget_bytes: usize,
974        reorder_permits: Arc<ReorderConcurrencyGate>,
975        background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
976    ) -> Self {
977        // Persisted index option: set via `reorder_on_merge: true` in the SDL
978        // at index creation. Absent = disabled (merges block-copy).
979        let reorder_on_merge = schema.reorder_on_merge();
980        if reorder_on_merge {
981            log::info!(
982                "[merge] index={} reorder-on-merge enabled by index schema",
983                schema.index_label()
984            );
985        }
986
987        let tracker = Arc::new(SegmentTracker::new());
988        for seg_id in metadata.owned_ids() {
989            tracker.register(&seg_id);
990        }
991
992        let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
993            Arc::new(parking_lot::Mutex::new(Vec::new()));
994        let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
995            let dir = Arc::clone(&directory);
996            let tracker = Arc::clone(&tracker);
997            let lifecycle_handles = Arc::clone(&lifecycle_handles);
998            let cleanup_index_label: Arc<str> = schema.index_label().into();
999            Arc::new(move |segment_ids| {
1000                // Guard: if the tokio runtime is gone (program exit), skip async
1001                // deletion. Segment files become orphans cleaned up on next startup.
1002                let Ok(handle) = tokio::runtime::Handle::try_current() else {
1003                    // Release in-process protection as well: if the process is
1004                    // still alive, a later sweep must be able to retry.
1005                    tracker.complete_deletion(&segment_ids);
1006                    return;
1007                };
1008                let dir = Arc::clone(&dir);
1009                let task_tracker = Arc::clone(&tracker);
1010                let task_index_label = Arc::clone(&cleanup_index_label);
1011                let cleanup_ids = segment_ids.clone();
1012                let future = async move {
1013                    for &segment_id in &segment_ids {
1014                        log::info!(
1015                            "[segment_cleanup] index={} deleting deferred segment {}",
1016                            task_index_label,
1017                            segment_id.to_hex()
1018                        );
1019                        if let Err(error) =
1020                            crate::segment::delete_segment(dir.as_ref(), segment_id).await
1021                        {
1022                            log::warn!(
1023                                "[segment_cleanup] index={} deferred delete failed for {}: {}",
1024                                task_index_label,
1025                                segment_id.to_hex(),
1026                                error,
1027                            );
1028                        }
1029                    }
1030                    task_tracker.complete_deletion(&segment_ids);
1031                };
1032                if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
1033                    // Spawning can fail only during runtime teardown. Release
1034                    // the scheduled-deletion claim so an in-process sweep can
1035                    // retry; crash recovery handles a process exit.
1036                    tracker.complete_deletion(&cleanup_ids);
1037                    log::warn!(
1038                        "[segment_cleanup] index={} runtime rejected deferred deletion; files will be swept later",
1039                        cleanup_index_label
1040                    );
1041                }
1042            })
1043        };
1044
1045        let initial_generation = Arc::new(PublishedIndexGeneration {
1046            publication_id: metadata.publication_generation,
1047            schema: Arc::clone(&schema),
1048            trained_vectors: None,
1049        });
1050        Self {
1051            state: Arc::new(AsyncMutex::new(ManagerState {
1052                metadata,
1053                merge_policy,
1054            })),
1055            active_operations: Arc::new(ActiveSegmentOperations::new(schema.index_label().into())),
1056            quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
1057            merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
1058            reorder_retries: parking_lot::Mutex::new(HashMap::new()),
1059            merge_handles: parking_lot::Mutex::new(Vec::new()),
1060            global_merge_wakeup_pending: AtomicBool::new(false),
1061            force_merge_active: AtomicUsize::new(0),
1062            #[cfg(test)]
1063            force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
1064            lifecycle_handles,
1065            published_generation: Arc::new(ArcSwap::new(initial_generation)),
1066            vector_artifact_update: Arc::new(AtomicBool::new(false)),
1067            tracker,
1068            delete_fn,
1069            directory,
1070            schema,
1071            optimization: crate::structures::IndexOptimization::default(),
1072            posting_codec: crate::structures::PostingCodec::default(),
1073            term_dict_block_size: crate::structures::SSTableBlockSize::default(),
1074            term_cache_blocks,
1075            term_cache_budget_bytes: None,
1076            merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
1077            global_merge_permits,
1078            reorder_permits,
1079            reorder_on_merge,
1080            merge_bp_time_budget,
1081            bp_memory_budget_bytes,
1082            background_reorder_pool,
1083            replacement_refresh: parking_lot::RwLock::new(None),
1084        }
1085    }
1086
1087    /// Cap decompressed dictionary blocks retained by lifecycle readers.
1088    /// None retains the block-count policy; zero disables retention.
1089    pub fn with_term_cache_budget(mut self, bytes: Option<usize>) -> Self {
1090        self.term_cache_budget_bytes = bytes;
1091        self
1092    }
1093
1094    /// Set the validated flush target for newly written term dictionaries.
1095    pub fn with_term_dict_block_size(mut self, size: crate::structures::SSTableBlockSize) -> Self {
1096        self.term_dict_block_size = size;
1097        self
1098    }
1099
1100    /// Configure posting compression for newly encoded output.
1101    pub fn with_posting_config(
1102        mut self,
1103        optimization: crate::structures::IndexOptimization,
1104        posting_codec: crate::structures::PostingCodec,
1105    ) -> Self {
1106        self.optimization = optimization;
1107        self.posting_codec = posting_codec;
1108        self
1109    }
1110
1111    pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1112    where
1113        F: Fn() -> Fut + Send + Sync + 'static,
1114        Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1115    {
1116        *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1117    }
1118
1119    /// Bounded rayon pool for background CPU (merge-time BP, manual reorder,
1120    /// and global vector-codebook training).
1121    /// Query scoring uses a dedicated search pool; keeping background work on
1122    /// this separate bounded pool prevents a merge or retrain from queueing
1123    /// every search behind its CPU passes.
1124    pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1125        if let Some(pool) = &self.background_reorder_pool {
1126            return Arc::clone(pool);
1127        }
1128        Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1129            let threads = (num_cpus::get() / 2).max(1);
1130            log::info!(
1131                "[merge] process-wide background CPU pool: {} thread(s)",
1132                threads
1133            );
1134            Arc::new(
1135                rayon::ThreadPoolBuilder::new()
1136                    .num_threads(threads)
1137                    .thread_name(|i| format!("summa-bg-cpu-{}", i))
1138                    .build()
1139                    .expect("failed to build background CPU pool"),
1140            )
1141        }))
1142    }
1143
1144    /// Stop new indexing/merge/reorder operations from claiming segment IDs.
1145    /// Used as the first half of index deletion; the writer then joins its
1146    /// workers before [`Self::wait_for_shutdown`] drains remaining ownership.
1147    pub fn begin_shutdown(&self) {
1148        self.active_operations.stop_accepting();
1149    }
1150
1151    /// Run a lifecycle mutation independently of its requesting future.
1152    ///
1153    /// Metadata writes contain an atomic rename. If an RPC is cancelled while
1154    /// awaiting that I/O, dropping the request must not abandon the matching
1155    /// in-memory/tracker transition. The spawned transaction is tracked for
1156    /// index shutdown; the oneshot only reports its result to a caller that is
1157    /// still interested.
1158    async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1159    where
1160        T: Send + 'static,
1161        F: std::future::Future<Output = Result<T>> + Send + 'static,
1162    {
1163        let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1164        let future = async move {
1165            let result = transaction.await;
1166            let _ = result_tx.send(result);
1167        };
1168        let runtime = tokio::runtime::Handle::current();
1169        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1170            return Err(Error::Internal(
1171                "runtime rejected lifecycle metadata transaction".into(),
1172            ));
1173        }
1174        result_rx.await.map_err(|_| {
1175            Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1176        })?
1177    }
1178
1179    /// Arm unwind cleanup for an output that is not visible in metadata yet.
1180    fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1181        let manager = Arc::clone(self);
1182        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1183            let Ok(handle) = tokio::runtime::Handle::try_current() else {
1184                log::warn!(
1185                    "[segment_cleanup] index={} runtime unavailable; partial output {} will be swept on startup",
1186                    manager.schema.index_label(),
1187                    segment_id.to_hex(),
1188                );
1189                return;
1190            };
1191
1192            let cleanup_manager = Arc::clone(&manager);
1193            let future = async move {
1194                cleanup_manager
1195                    .delete_output_if_unregistered(segment_id, "task unwind")
1196                    .await;
1197            };
1198            if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1199                log::warn!(
1200                    "[segment_cleanup] index={} runtime rejected output cleanup; {} will be swept on startup",
1201                    manager.schema.index_label(),
1202                    segment_id.to_hex(),
1203                );
1204            }
1205        });
1206
1207        OutputCleanupGuard::new(output_id, cleanup)
1208    }
1209
1210    /// Delete an abandoned indexing output while retaining its lifecycle
1211    /// claim until the last file operation completes. The explicit runtime
1212    /// handle makes this safe from dedicated indexing OS threads, which are
1213    /// outside Tokio's entered context.
1214    pub(crate) fn schedule_unpublished_segment_cleanup(
1215        self: &Arc<Self>,
1216        output_id: SegmentId,
1217        operation: SegmentOperationGuard,
1218        runtime: tokio::runtime::Handle,
1219    ) {
1220        let manager = Arc::clone(self);
1221        let output_hex = output_id.to_hex();
1222        let future = async move {
1223            manager
1224                .delete_output_if_unregistered(output_id, "indexing abort or failure")
1225                .await;
1226            drop(operation);
1227        };
1228        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1229            // The dropped future releases operation ownership. Startup sweep
1230            // handles its output if the runtime is already tearing down.
1231            log::warn!(
1232                "[segment_cleanup] index={} runtime unavailable; indexing output {} will be swept on startup",
1233                self.schema.index_label(),
1234                output_hex,
1235            );
1236        }
1237    }
1238
1239    /// Claim a newly generated indexing segment before its first file write.
1240    ///
1241    /// The returned guard must travel with the built segment until metadata
1242    /// publication or abort. UUID collisions are treated as corruption rather
1243    /// than silently sharing lifecycle ownership.
1244    pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1245        match self
1246            .active_operations
1247            .try_register_indexing(vec![segment_id.clone()])
1248        {
1249            Some(operation) => Ok(operation),
1250            None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1251            None => Err(Error::Corruption(format!(
1252                "new segment ID {} is already owned by an active operation",
1253                segment_id
1254            ))),
1255        }
1256    }
1257
1258    /// Validate the small, mandatory core of a completed segment before it can
1259    /// become metadata-live. Optional vector/sparse/position/fast files are
1260    /// schema- and data-dependent and are validated by `SegmentReader` when used.
1261    async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1262        let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1263            Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1264        })?;
1265        let files = SegmentFiles::new(id.0);
1266
1267        for path in files.mandatory_paths() {
1268            if !self.directory.exists(path).await.map_err(Error::Io)? {
1269                return Err(Error::Corruption(format!(
1270                    "segment {} cannot be published: mandatory file {:?} is missing",
1271                    segment_id, path
1272                )));
1273            }
1274        }
1275
1276        let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1277            Error::Corruption(format!(
1278                "segment {} cannot be published: missing/unreadable {:?}: {}",
1279                segment_id, files.meta, e
1280            ))
1281        })?;
1282        let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1283            Error::Corruption(format!(
1284                "segment {} cannot be published: failed reading {:?}: {}",
1285                segment_id, files.meta, e
1286            ))
1287        })?;
1288        let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1289            Error::Corruption(format!(
1290                "segment {} cannot be published: invalid {:?}: {}",
1291                segment_id, files.meta, e
1292            ))
1293        })?;
1294
1295        if meta.id != id.0 || meta.num_docs != expected_docs {
1296            return Err(Error::Corruption(format!(
1297                "segment {} cannot be published: metadata identity/docs mismatch \
1298                 (id={:032x}, docs={}, expected_docs={})",
1299                segment_id, meta.id, meta.num_docs, expected_docs
1300            )));
1301        }
1302
1303        Ok(())
1304    }
1305
1306    fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1307        let inserted = self
1308            .quarantined_segments
1309            .lock()
1310            .insert(segment_id.to_string());
1311        if inserted {
1312            log::error!(
1313                "[merge] index={} quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1314                 It remains metadata-live for explicit repair but is excluded from merges until restart",
1315                self.schema.index_label(),
1316                segment_id,
1317                error,
1318            );
1319        }
1320    }
1321
1322    fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1323        let mut retry = self.merge_retry.lock();
1324        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1325        let delay = merge_retry_delay(retry.consecutive_failures);
1326        retry.retry_after = std::time::Instant::now().checked_add(delay);
1327        log::warn!(
1328            "[merge] index={} pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1329            self.schema.index_label(),
1330            delay.as_secs_f64(),
1331            retry.consecutive_failures,
1332            error,
1333        );
1334        delay
1335    }
1336
1337    fn clear_merge_retry_backoff(&self) {
1338        *self.merge_retry.lock() = MergeRetryState::default();
1339    }
1340
1341    fn merge_retry_is_paused(&self) -> bool {
1342        let mut retry = self.merge_retry.lock();
1343        match retry.retry_after {
1344            Some(deadline) if deadline > std::time::Instant::now() => true,
1345            Some(_) => {
1346                retry.retry_after = None;
1347                false
1348            }
1349            None => false,
1350        }
1351    }
1352
1353    async fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1354        // Serialize failure admission with replacement publication: a worker
1355        // reporting late must not recreate retry state for a retired source.
1356        let st = self.state.lock().await;
1357        if !st.metadata.has_segment(segment_id) {
1358            return;
1359        }
1360        let mut retries = self.reorder_retries.lock();
1361        let retry = retries.entry(segment_id.to_string()).or_default();
1362        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1363        let delay = merge_retry_delay(retry.consecutive_failures);
1364        retry.retry_after = std::time::Instant::now().checked_add(delay);
1365        log::warn!(
1366            "[reorder] index={} pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1367            self.schema.index_label(),
1368            segment_id,
1369            delay.as_secs_f64(),
1370            retry.consecutive_failures,
1371            error,
1372        );
1373    }
1374
1375    fn clear_reorder_retry(&self, segment_id: &str) {
1376        self.reorder_retries.lock().remove(segment_id);
1377    }
1378
1379    /// Called with publication state held, matching failure admission order.
1380    fn retire_reorder_retries(&self, retired: &[String]) {
1381        let mut retries = self.reorder_retries.lock();
1382        for id in retired {
1383            retries.remove(id);
1384        }
1385    }
1386
1387    fn paused_reorder_segments(&self) -> HashSet<String> {
1388        let now = std::time::Instant::now();
1389        let mut retries = self.reorder_retries.lock();
1390        let mut paused = HashSet::new();
1391        for (segment_id, retry) in retries.iter_mut() {
1392            match retry.retry_after {
1393                Some(deadline) if deadline > now => {
1394                    paused.insert(segment_id.clone());
1395                }
1396                Some(_) => retry.retry_after = None,
1397                None => {}
1398            }
1399        }
1400        paused
1401    }
1402
1403    /// Re-evaluate this index when another index releases application-wide
1404    /// merge capacity. The atomic flag bounds this to one waiter per index and
1405    /// the tracked handle makes index shutdown drain it deterministically.
1406    fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1407        if self
1408            .global_merge_wakeup_pending
1409            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1410            .is_err()
1411        {
1412            return;
1413        }
1414
1415        let manager = Arc::clone(self);
1416        let future = async move {
1417            let capacity = tokio::select! {
1418                biased;
1419                () = manager.active_operations.wait_for_shutdown() => None,
1420                permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1421            };
1422
1423            manager
1424                .global_merge_wakeup_pending
1425                .store(false, Ordering::Release);
1426            if let Some(permit) = capacity {
1427                // This task is only a notification. The normal scheduler must
1428                // acquire both global and per-index permits atomically enough
1429                // for its own candidate selection.
1430                drop(permit);
1431                manager.maybe_merge().await;
1432            }
1433        };
1434        let runtime = tokio::runtime::Handle::current();
1435        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1436            self.global_merge_wakeup_pending
1437                .store(false, Ordering::Release);
1438            log::warn!(
1439                "[merge] index={} runtime rejected global-capacity wakeup task",
1440                self.schema.index_label()
1441            );
1442        }
1443    }
1444
1445    #[cfg(test)]
1446    pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1447        self.quarantined_segments.lock().contains(segment_id)
1448    }
1449
1450    /// Delete a failed output only if metadata did not make it live.
1451    ///
1452    /// Rechecking under `state` also makes unwind cleanup safe if it races
1453    /// successful publication of the same output.
1454    async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1455        let output_hex = output_id.to_hex();
1456        {
1457            let st = self.state.lock().await;
1458            if st.metadata.owns_id(&output_hex) {
1459                return;
1460            }
1461        }
1462
1463        // UUIDs are generated per producer and cannot be adopted by another
1464        // publisher after this check. Never hold the metadata mutex while a
1465        // multi-GB filesystem deletion runs.
1466        log::info!(
1467            "[segment_cleanup] index={} deleting uncommitted output {} after {}",
1468            self.schema.index_label(),
1469            output_hex,
1470            reason,
1471        );
1472        if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1473        {
1474            log::warn!(
1475                "[segment_cleanup] index={} failed deleting uncommitted output {}: {}",
1476                self.schema.index_label(),
1477                output_hex,
1478                error,
1479            );
1480        }
1481    }
1482
1483    // ========================================================================
1484    // Read path (brief lock or lock-free)
1485    // ========================================================================
1486
1487    /// Get the current segment IDs
1488    pub async fn get_segment_ids(&self) -> Vec<String> {
1489        self.state.lock().await.metadata.segment_ids()
1490    }
1491
1492    /// Get trained vector structures (lock-free via ArcSwap)
1493    pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1494        self.published_generation.load().trained_vectors.clone()
1495    }
1496
1497    /// Current schema/vector publication. Cloning the Arc gives an operation a
1498    /// stable configuration even if an ALTER publishes concurrently.
1499    pub(crate) fn published_generation(&self) -> Arc<PublishedIndexGeneration> {
1500        self.published_generation.load_full()
1501    }
1502
1503    pub(crate) fn publication_id(&self) -> u64 {
1504        self.published_generation.load().publication_id
1505    }
1506
1507    /// Capture trained structures for a segment producer.
1508    ///
1509    /// The second gate check closes the race where an update begins after the
1510    /// first check but before the ArcSwap load. Producers have lifecycle guards
1511    /// before calling this method, so an updater that raised the gate waits for
1512    /// any producer that successfully captured the previous generation.
1513    pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1514        if self.vector_artifact_update.load(Ordering::Acquire) {
1515            return None;
1516        }
1517        let trained = self.published_generation.load().trained_vectors.clone();
1518        if self.vector_artifact_update.load(Ordering::Acquire) {
1519            None
1520        } else {
1521            trained
1522        }
1523    }
1524
1525    /// Start an exclusive trained-artifact update and drain merge/reorder
1526    /// producers that may already hold the previous generation.
1527    ///
1528    /// New segment operations may continue while this waits, but they observe
1529    /// the gate through `trained_for_segment_build` and therefore emit flat
1530    /// vector data until the replacement generation is published. The guard
1531    /// is cancellation-safe: dropping the requesting future reopens ANN
1532    /// production without leaving the manager wedged.
1533    ///
1534    /// Indexing tokens cannot be waited on: their guards are parked inside
1535    /// built-but-uncommitted `PreparedSegment`s and are released only by a
1536    /// later commit. That commit typically needs the writer this update's
1537    /// caller already holds (server write lock / embedded `&mut self`), so
1538    /// waiting would permanently wedge vector-index finalization. They also
1539    /// cannot be ignored: such a segment may already contain ANN data bound to
1540    /// the previous artifact generation and could otherwise be committed after
1541    /// the new generation is published. Reject promptly and let the caller
1542    /// commit or abort the pending generation before retrying.
1543    pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1544        self.vector_artifact_update
1545            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1546            .map_err(|_| {
1547                Error::Internal("a trained-vector artifact update is already in progress".into())
1548            })?;
1549        self.active_operations.pause_non_indexing();
1550        let guard = VectorArtifactUpdateGuard {
1551            _lease: Arc::new(VectorArtifactUpdateLease {
1552                updating: Arc::clone(&self.vector_artifact_update),
1553                active_operations: Arc::clone(&self.active_operations),
1554            }),
1555        };
1556        let (preexisting, parked_indexing) =
1557            self.active_operations.draining_operation_tokens_snapshot();
1558        if parked_indexing > 0 {
1559            return Err(Error::Internal(format!(
1560                "cannot update trained-vector artifacts while {parked_indexing} indexing \
1561                 segment(s) are built but uncommitted; commit or abort the pending \
1562                 generation and retry"
1563            )));
1564        }
1565        self.active_operations
1566            .wait_until_operations_finish(&preexisting)
1567            .await;
1568        Ok(guard)
1569    }
1570
1571    /// Load trained structures from disk and publish to ArcSwap.
1572    /// Copies metadata under lock, releases lock, then does disk I/O.
1573    pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1574        // Copy vector_fields under lock (cheap clone of HashMap<u32, FieldMeta>)
1575        let (vector_fields, schema, publication_id) = {
1576            let st = self.state.lock().await;
1577            (
1578                st.metadata.vector_fields.clone(),
1579                Arc::new(st.metadata.schema.clone()),
1580                st.metadata.publication_generation,
1581            )
1582        };
1583        // Disk I/O happens WITHOUT holding the state lock
1584        let trained = IndexMetadata::try_load_trained_from_fields(
1585            &vector_fields,
1586            schema.as_ref(),
1587            self.directory.as_ref(),
1588        )
1589        .await?
1590        .map(Arc::new);
1591        // Publish exactly the validated snapshot, including None. Retaining a
1592        // previous map when metadata has no Built fields would let new segments
1593        // depend on artifacts no longer referenced durably.
1594        self.published_generation
1595            .store(Arc::new(PublishedIndexGeneration {
1596                publication_id,
1597                schema,
1598                trained_vectors: trained,
1599            }));
1600        Ok(())
1601    }
1602
1603    /// Atomically publish a fully staged vector generation.
1604    ///
1605    /// Every replacement segment and its global codebook is complete before
1606    /// this transaction starts. Metadata, tracker ownership, and the lock-free
1607    /// trained pointer advance under the same state lock; snapshots therefore
1608    /// observe either the complete old generation or the complete new one.
1609    pub(crate) async fn publish_vector_generation(
1610        self: &Arc<Self>,
1611        artifact_update: &VectorArtifactUpdateGuard,
1612        vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1613        next_trained: Arc<TrainedVectorStructures>,
1614        staged: Vec<StagedVectorSegment>,
1615    ) -> Result<()> {
1616        let schema = self.published_generation().schema.clone();
1617        self.publish_vector_generation_with_schema(
1618            artifact_update,
1619            schema,
1620            vector_fields,
1621            Some(next_trained),
1622            staged,
1623        )
1624        .await
1625    }
1626
1627    pub(crate) async fn publish_vector_generation_with_schema(
1628        self: &Arc<Self>,
1629        artifact_update: &VectorArtifactUpdateGuard,
1630        schema: Arc<crate::dsl::Schema>,
1631        vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1632        next_trained: Option<Arc<TrainedVectorStructures>>,
1633        mut staged: Vec<StagedVectorSegment>,
1634    ) -> Result<()> {
1635        if !self.vector_artifact_update.load(Ordering::Acquire) {
1636            return Err(Error::Internal(
1637                "vector generation publication lost its exclusive update lease".into(),
1638            ));
1639        }
1640
1641        for replacement in &staged {
1642            self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1643                .await?;
1644        }
1645
1646        let mut st = Arc::clone(&self.state).lock_owned().await;
1647        let mut next = st.metadata.clone();
1648        next.schema = (*schema).clone();
1649        next.vector_fields = vector_fields;
1650        next.refresh_total_vectors();
1651
1652        for replacement in &staged {
1653            let source_info = next
1654                .segment_metas
1655                .remove(&replacement.source_id)
1656                .ok_or_else(|| {
1657                    Error::Corruption(format!(
1658                        "vector generation source {} disappeared before publication",
1659                        replacement.source_id,
1660                    ))
1661                })?;
1662            let output_hex = replacement.output_id.to_hex();
1663            if next.segment_metas.contains_key(&output_hex) {
1664                return Err(Error::Corruption(format!(
1665                    "vector generation output {output_hex} is already metadata-live"
1666                )));
1667            }
1668            // A vector-only rewrite changes neither document order nor merge
1669            // lineage, so preserve the complete lifecycle record verbatim.
1670            next.add_segment_meta(output_hex, source_info);
1671        }
1672
1673        let directory = Arc::clone(&self.directory);
1674        let published_generation = Arc::clone(&self.published_generation);
1675        let tracker = Arc::clone(&self.tracker);
1676        let replacement_refresh = self.replacement_refresh.read().clone();
1677        let manager = Arc::clone(self);
1678        // Keep the producer gate raised if the requesting future is cancelled
1679        // after the metadata transaction has been detached. The last guard
1680        // clone drops only after durable metadata and ArcSwap state agree.
1681        let artifact_update = artifact_update.clone();
1682        let index_label = self.schema.index_label().to_owned();
1683        next.publication_generation =
1684            next.publication_generation.checked_add(1).ok_or_else(|| {
1685                Error::Corruption("vector publication generation exhausted u64".into())
1686            })?;
1687        let next_schema = schema;
1688        let next_publication_id = next.publication_generation;
1689        self.run_lifecycle_transaction(async move {
1690            let _artifact_update = artifact_update;
1691            next.save(directory.as_ref()).await?;
1692
1693            for replacement in &staged {
1694                tracker.register(&replacement.output_id.to_hex());
1695            }
1696            st.metadata = next;
1697            published_generation.store(Arc::new(PublishedIndexGeneration {
1698                publication_id: next_publication_id,
1699                schema: next_schema,
1700                trained_vectors: next_trained,
1701            }));
1702
1703            // Outputs are now durably live. Disarm unwind cleanup before
1704            // retiring the old generation and releasing operation ownership.
1705            for replacement in &mut staged {
1706                replacement.cleanup.disarm();
1707            }
1708            let retired = staged
1709                .iter()
1710                .map(|replacement| replacement.source_id.clone())
1711                .collect::<Vec<_>>();
1712            manager.retire_reorder_retries(&retired);
1713            let ready_to_delete = tracker.mark_for_deletion(&retired);
1714            drop(st);
1715            for &segment_id in &ready_to_delete {
1716                if let Err(error) =
1717                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
1718                {
1719                    log::warn!(
1720                        "[segment_cleanup] index={index_label} immediate dense-vector generation delete failed for {}: {}",
1721                        segment_id.to_hex(),
1722                        error,
1723                    );
1724                }
1725            }
1726            tracker.complete_deletion(&ready_to_delete);
1727            refresh_replacement_topology(replacement_refresh, &index_label).await;
1728            Ok(())
1729        })
1730        .await
1731    }
1732
1733    /// Publish query-only vector parameters without rebuilding segments or
1734    /// artifacts. The publication ID still advances so cached readers reload
1735    /// even though the segment ID set is unchanged.
1736    pub(crate) async fn publish_vector_schema_only(
1737        self: &Arc<Self>,
1738        artifact_update: &VectorArtifactUpdateGuard,
1739        schema: Arc<crate::dsl::Schema>,
1740    ) -> Result<()> {
1741        let vector_fields = self
1742            .read_metadata(|metadata| metadata.vector_fields.clone())
1743            .await;
1744        let trained = self.published_generation().trained_vectors.clone();
1745        self.publish_vector_generation_with_schema(
1746            artifact_update,
1747            schema,
1748            vector_fields,
1749            trained,
1750            Vec::new(),
1751        )
1752        .await
1753    }
1754
1755    /// Read metadata with a closure (no persist)
1756    pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1757    where
1758        F: FnOnce(&IndexMetadata) -> R,
1759    {
1760        let st = self.state.lock().await;
1761        f(&st.metadata)
1762    }
1763
1764    /// Update metadata with a closure and persist atomically
1765    pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1766    where
1767        F: FnOnce(&mut IndexMetadata),
1768    {
1769        let mut st = Arc::clone(&self.state).lock_owned().await;
1770        let mut next = st.metadata.clone();
1771        f(&mut next);
1772        let directory = Arc::clone(&self.directory);
1773        self.run_lifecycle_transaction(async move {
1774            next.save(directory.as_ref()).await?;
1775            st.metadata = next;
1776            Ok(())
1777        })
1778        .await
1779    }
1780
1781    /// Acquire a snapshot of current segments for reading.
1782    /// The snapshot holds references — segments won't be deleted while snapshot exists.
1783    pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1784        let st = self.state.lock().await;
1785        let acquired = self.tracker.acquire(&st.metadata.segment_ids());
1786        SegmentSnapshot::with_generation(
1787            Arc::clone(&self.tracker),
1788            acquired,
1789            self.published_generation.load_full(),
1790            Arc::clone(&self.delete_fn),
1791        )
1792        .with_deletions(&st.metadata)
1793    }
1794
1795    /// Get the segment tracker
1796    pub fn tracker(&self) -> Arc<SegmentTracker> {
1797        Arc::clone(&self.tracker)
1798    }
1799
1800    /// Get the directory
1801    pub fn directory(&self) -> Arc<D> {
1802        Arc::clone(&self.directory)
1803    }
1804}
1805
1806// ============================================================================
1807// Native-only: commit, merging, force_merge
1808// ============================================================================
1809
1810#[cfg(feature = "native")]
1811impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1812    /// Atomic commit: register new segments + persist metadata.
1813    pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1814        // Indexing guards still own these IDs here, so the orphan sweeper
1815        // cannot remove files between validation and metadata publication.
1816        for (segment_id, num_docs) in new_segments {
1817            self.validate_completed_segment(segment_id, *num_docs)
1818                .await?;
1819        }
1820
1821        let mut st = Arc::clone(&self.state).lock_owned().await;
1822        let mut next = st.metadata.clone();
1823        let mut added = Vec::new();
1824        for (segment_id, num_docs) in new_segments {
1825            if !next.has_segment(segment_id) {
1826                next.add_segment(segment_id.clone(), *num_docs);
1827                added.push(segment_id.clone());
1828            }
1829        }
1830
1831        // Durable-before-visible: a save failure leaves both in-memory metadata
1832        // and tracker unchanged, so callers can retry the prepared commit.
1833        // The tracked transaction continues if the requesting RPC is cancelled;
1834        // unpublished cleanup waits on this owned state guard before deciding
1835        // whether the files became metadata-live.
1836        let directory = Arc::clone(&self.directory);
1837        let tracker = Arc::clone(&self.tracker);
1838        self.run_lifecycle_transaction(async move {
1839            next.save(directory.as_ref()).await?;
1840            for segment_id in &added {
1841                tracker.register(segment_id);
1842            }
1843            st.metadata = next;
1844            Ok(())
1845        })
1846        .await
1847    }
1848
1849    /// Evaluate merge policy and spawn background merges for all eligible candidates.
1850    ///
1851    /// **Atomicity**: The entire filter → find_merges → spawn_merge sequence runs
1852    /// under the `state` lock to prevent a TOCTOU race where concurrent callers
1853    /// both see segments as eligible before either claims operation ownership.
1854    /// `spawn_merge` is non-blocking (just `try_register` + `tokio::spawn`), so
1855    /// holding the state lock through it is safe and sub-microsecond.
1856    ///
1857    /// The hard merge semaphore is acquired before lifecycle ownership, so
1858    /// concurrent triggers cannot exceed configured merge capacity.
1859    pub async fn maybe_merge(self: &Arc<Self>) {
1860        if !self.active_operations.is_accepting() {
1861            log::debug!(
1862                "[maybe_merge] index={} manager is shutting down, skipping",
1863                self.schema.index_label()
1864            );
1865            return;
1866        }
1867        if self.merge_retry_is_paused() {
1868            log::debug!(
1869                "[maybe_merge] index={} retry backoff active, skipping",
1870                self.schema.index_label()
1871            );
1872            return;
1873        }
1874
1875        // Finished handles no longer need to be retained. Concurrency itself
1876        // is enforced by `merge_permits`, not this bookkeeping vector.
1877        {
1878            let mut handles = self.merge_handles.lock();
1879            handles.retain(|h| !h.is_finished());
1880        }
1881        let local_slots = self.merge_permits.available_permits();
1882        let global_slots = self.global_merge_permits.available_permits();
1883        let slots_available = local_slots.min(global_slots);
1884
1885        // Hold state lock through spawn_merge to make filter + register atomic.
1886        // This closes the TOCTOU window where concurrent maybe_merge calls could
1887        // both see the same segments as eligible before either registers them.
1888        {
1889            let st = self.state.lock().await;
1890            let quarantined = self.quarantined_segments.lock().clone();
1891            let active_ids = self.active_operations.snapshot();
1892
1893            // Backlog pressure is based on every metadata-live segment, not
1894            // only currently available inputs. Otherwise a wave of in-flight
1895            // merges hides most of the topology and makes later candidates
1896            // look healthy while the original backlog still exists.
1897            let live_segments: Vec<SegmentInfo> = st
1898                .metadata
1899                .segment_metas
1900                .iter()
1901                .filter(|(id, _)| {
1902                    !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1903                })
1904                .map(|(id, info)| SegmentInfo {
1905                    id: id.clone(),
1906                    num_docs: info.num_docs,
1907                })
1908                .collect();
1909            let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1910
1911            // Exclude segments owned by another operation. Pending retirement
1912            // and quarantined segments were already removed above.
1913            let segments: Vec<SegmentInfo> = live_segments
1914                .iter()
1915                .filter(|segment| !active_ids.contains(&segment.id))
1916                .cloned()
1917                .collect();
1918
1919            log::debug!(
1920                "[maybe_merge] index={} {} eligible segments",
1921                self.schema.index_label(),
1922                segments.len()
1923            );
1924
1925            let candidates = st.merge_policy.find_merges(&segments);
1926
1927            if candidates.is_empty() {
1928                return;
1929            }
1930
1931            // Register a capacity waiter only for an index that actually has
1932            // eligible work. Scheduling one waiter for every idle index while
1933            // the process gate was full caused an avoidable wakeup stampede.
1934            if slots_available == 0 {
1935                if local_slots > 0 && global_slots == 0 {
1936                    self.schedule_global_merge_wakeup();
1937                }
1938                log::debug!(
1939                    "[maybe_merge] index={} at max concurrent merges, skipping",
1940                    self.schema.index_label()
1941                );
1942                return;
1943            }
1944
1945            log::debug!(
1946                "[maybe_merge] index={} {} merge candidates, {} slots available",
1947                self.schema.index_label(),
1948                candidates.len(),
1949                slots_available
1950            );
1951
1952            let mut handles = Vec::new();
1953            for c in candidates {
1954                if handles.len() >= slots_available {
1955                    break;
1956                }
1957                // Under severe topology pressure, retire segments first. BP
1958                // remains optional for correctness and the block-copy output
1959                // is explicitly marked unreordered, so the optimizer performs
1960                // one pass after the compaction wave instead of every merge
1961                // task serializing behind the whole-pass gate.
1962                let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1963                if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1964                    handles.push(h);
1965                }
1966            }
1967            if !handles.is_empty() {
1968                if severe_backlog && self.reorder_on_merge {
1969                    log::info!(
1970                        "[maybe_merge] index={} severe backlog: {} live segments; started {} fast \
1971                         block-copy merge(s), deferring BP to the optimizer",
1972                        self.schema.index_label(),
1973                        live_segments.len(),
1974                        handles.len(),
1975                    );
1976                }
1977                // Publish handles before releasing `state`. A force merge
1978                // raises its admission barrier under the same lock, then
1979                // drains this list, so no spawned merge can fall into the
1980                // otherwise-deadlocking gap between spawn and bookkeeping.
1981                self.merge_handles.lock().extend(handles);
1982            }
1983        }
1984    }
1985
1986    /// Spawn a background merge task with RAII tracking.
1987    ///
1988    /// Pre-generates the output segment ID. The operation guard registers all segment IDs
1989    /// (old + output) in `active_operations`. When the task ends (success, failure, or
1990    /// panic), the guard drops and segments are automatically unregistered.
1991    ///
1992    /// On completion, the task auto-triggers `maybe_merge` to evaluate cascading merges.
1993    /// Returns the JoinHandle if the merge was spawned, None if it was skipped.
1994    fn spawn_merge(
1995        self: &Arc<Self>,
1996        segment_ids_to_merge: Vec<String>,
1997        reorder_bmp: bool,
1998    ) -> Option<JoinHandle<()>> {
1999        if self.force_merge_active.load(Ordering::Acquire) > 0 {
2000            log::debug!(
2001                "[spawn_merge] index={} skipped: explicit force merge has priority",
2002                self.schema.index_label()
2003            );
2004            return None;
2005        }
2006        let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
2007            Ok(permit) => permit,
2008            Err(_) => {
2009                log::debug!(
2010                    "[spawn_merge] index={} skipped: global merge capacity is full",
2011                    self.schema.index_label()
2012                );
2013                self.schedule_global_merge_wakeup();
2014                return None;
2015            }
2016        };
2017        let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
2018            Ok(permit) => permit,
2019            Err(_) => {
2020                log::debug!(
2021                    "[spawn_merge] index={} skipped: no merge permit available",
2022                    self.schema.index_label()
2023                );
2024                return None;
2025            }
2026        };
2027        let output_id = SegmentId::new();
2028        let output_hex = output_id.to_hex();
2029
2030        let mut all_ids = segment_ids_to_merge.clone();
2031        all_ids.push(output_hex);
2032
2033        let guard = match self.active_operations.try_register(all_ids) {
2034            Some(g) => g,
2035            None => {
2036                log::debug!(
2037                    "[spawn_merge] index={} skipped: segments overlap with an active operation",
2038                    self.schema.index_label()
2039                );
2040                return None;
2041            }
2042        };
2043
2044        let sm = Arc::clone(self);
2045        let ids = segment_ids_to_merge;
2046
2047        let index_label = self.schema.index_label().to_owned();
2048        Some(tokio::spawn(async move {
2049            let mut reevaluate = false;
2050            let mut retry_delay = None;
2051
2052            let result = sm
2053                .merge_and_replace_registered(
2054                    &ids,
2055                    output_id,
2056                    reorder_bmp,
2057                    ReorderPriority::AutomaticMerge,
2058                    Arc::new((guard, merge_permit, global_merge_permit)),
2059                )
2060                .await;
2061
2062            match result {
2063                Ok(_) => {
2064                    sm.clear_merge_retry_backoff();
2065                    reevaluate = true;
2066                }
2067                Err(MergeTaskError {
2068                    error: Error::IndexClosed,
2069                    ..
2070                }) => {
2071                    log::debug!(
2072                        "[merge] index={index_label} background merge for segments {:?} cancelled during shutdown",
2073                        ids,
2074                    );
2075                }
2076                Err(MergeTaskError {
2077                    error,
2078                    unavailable_segments,
2079                }) => {
2080                    log::error!(
2081                        "[merge] index={index_label} background merge failed for segments {:?}: {}",
2082                        ids,
2083                        error
2084                    );
2085                    if !unavailable_segments.is_empty() {
2086                        // Recompute without known-bad inputs. This is not a
2087                        // retry of the same candidate because policy filtering
2088                        // excludes every quarantined ID.
2089                        reevaluate = true;
2090                    } else {
2091                        retry_delay = Some(sm.pause_merge_retries(&error));
2092                    }
2093                }
2094            }
2095            // Release source/output ownership before re-evaluating policy, so
2096            // the completed operation cannot artificially hide candidates.
2097
2098            if reevaluate {
2099                sm.maybe_merge().await;
2100            } else if let Some(retry_delay) = retry_delay {
2101                // A backoff without a wakeup can strand eligible segments
2102                // forever when no later commit happens. The sleep runs as a
2103                // tracked *lifecycle* task, not inside this merge JoinHandle:
2104                // wait_for_all_merges/force_merge/reorder drain merge handles,
2105                // and a pure backoff timer with no work in flight must not
2106                // stall them for up to MERGE_RETRY_MAX_DELAY.
2107                sm.schedule_merge_retry_wakeup(retry_delay);
2108            }
2109        }))
2110    }
2111
2112    /// Execute the common build → validate → durable replacement transaction
2113    /// for a batch whose source/output IDs are already lifecycle-owned.
2114    ///
2115    /// Scheduling, retry policy, and post-replacement snapshot refresh remain
2116    /// with the caller. Keeping the transaction here prevents background and
2117    /// force-merge paths from drifting on cleanup or BP metadata semantics.
2118    async fn merge_and_replace_registered(
2119        self: &Arc<Self>,
2120        ids: &[String],
2121        output_id: SegmentId,
2122        reorder_bmp: bool,
2123        priority: ReorderPriority,
2124        ownership: Arc<dyn Send + Sync>,
2125    ) -> MergeTaskResult<(String, u32, bool)> {
2126        let manager = Arc::clone(self);
2127        let ids = ids.to_vec();
2128        // Keep source/output claims and capacity through started blocking
2129        // compaction even when a force-merge requester is cancelled.
2130        self.run_lifecycle_transaction(async move {
2131            let _ownership = ownership;
2132            Ok(manager
2133                .merge_and_replace_owned(&ids, output_id, reorder_bmp, priority)
2134                .await)
2135        })
2136        .await
2137        .map_err(MergeTaskError::from)?
2138    }
2139
2140    async fn merge_and_replace_owned(
2141        self: &Arc<Self>,
2142        ids: &[String],
2143        output_id: SegmentId,
2144        reorder_bmp: bool,
2145        priority: ReorderPriority,
2146    ) -> MergeTaskResult<(String, u32, bool)> {
2147        let mut output_cleanup = self.output_cleanup_guard(output_id);
2148        let generation = self.published_generation();
2149        let trained = self.trained_for_segment_build();
2150        let granularity = if reorder_bmp {
2151            self.merge_granularity(ids).await
2152        } else {
2153            crate::segment::reorder::BpGranularity::Auto
2154        };
2155        let result = Self::do_merge(
2156            self.directory.as_ref(),
2157            &generation.schema,
2158            ids,
2159            output_id,
2160            self.term_cache_blocks,
2161            self.term_cache_budget_bytes,
2162            self.optimization,
2163            self.posting_codec,
2164            self.term_dict_block_size,
2165            trained.as_deref(),
2166            reorder_bmp,
2167            granularity,
2168            self.merge_bp_time_budget,
2169            self.bp_memory_budget_bytes,
2170            Arc::clone(&self.reorder_permits),
2171            priority,
2172            self.active_operations.cancellation_flag(),
2173            Some(self.background_cpu_pool()),
2174        )
2175        .await;
2176
2177        let (new_id, doc_count, bp_converged) = match result {
2178            Ok(value) => value,
2179            Err(error) => {
2180                for segment_id in &error.unavailable_segments {
2181                    self.quarantine_segment(segment_id, &error.error);
2182                }
2183                self.delete_output_if_unregistered(output_id, "merge failure")
2184                    .await;
2185                output_cleanup.disarm();
2186                return Err(error);
2187            }
2188        };
2189
2190        let layout = if reorder_bmp {
2191            ReplacementLayout::BpReordered {
2192                converged: bp_converged,
2193            }
2194        } else {
2195            ReplacementLayout::BlockCopy
2196        };
2197        if let Err(error) = self
2198            .replace_segments(ids, new_id.clone(), doc_count, layout, None)
2199            .await
2200        {
2201            self.delete_output_if_unregistered(output_id, "replacement failure")
2202                .await;
2203            output_cleanup.disarm();
2204            return Err(MergeTaskError::from(error));
2205        }
2206        output_cleanup.disarm();
2207        Ok((new_id, doc_count, bp_converged))
2208    }
2209
2210    /// Re-evaluate merge policy after a failure backoff, outside the tracked
2211    /// merge JoinHandles that merge waiters drain. Shutdown still drains this
2212    /// task deterministically (lifecycle handles) and interrupts its sleep.
2213    fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
2214        let manager = Arc::clone(self);
2215        let future = async move {
2216            tokio::select! {
2217                () = tokio::time::sleep(retry_delay) => {
2218                    manager.maybe_merge().await;
2219                }
2220                () = manager.active_operations.wait_for_shutdown() => {}
2221            }
2222        };
2223        let runtime = tokio::runtime::Handle::current();
2224        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
2225            log::warn!(
2226                "[merge] index={} runtime rejected merge-retry wakeup task; eligible segments may stay \
2227                 unmerged until the next commit re-runs merge policy evaluation",
2228                self.schema.index_label()
2229            );
2230        }
2231    }
2232
2233    /// Atomically replace old segments with a new merged segment.
2234    /// Computes merge generation as max(parent gens) + 1 and records ancestors.
2235    /// `reordered` marks whether the new segment was BP-reordered.
2236    async fn replace_segments(
2237        self: &Arc<Self>,
2238        old_ids: &[String],
2239        new_id: String,
2240        doc_count: u32,
2241        layout: ReplacementLayout,
2242        expected_deletions: Option<&HashMap<String, Option<crate::segment::DeletionMeta>>>,
2243    ) -> Result<()> {
2244        // The operation guard owns the output during validation. Publication
2245        // below replaces that ownership with metadata + tracker atomically.
2246        self.validate_completed_segment(&new_id, doc_count).await?;
2247        let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
2248            Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
2249        })?;
2250        let output_reader = SegmentReader::open_with_term_cache_budget(
2251            self.directory.as_ref(),
2252            output_id,
2253            self.published_generation().schema.clone(),
2254            self.term_cache_blocks,
2255            self.term_cache_budget_bytes,
2256        )
2257        .await
2258        .map_err(|error| match error {
2259            // Preserve retryable storage failures as I/O. Structural failures
2260            // are deterministic for this completed output and get explicit
2261            // corruption context.
2262            Error::Io(_) | Error::IndexClosed => error,
2263            error => Error::Corruption(format!(
2264                "replacement segment {new_id} failed full reader validation: {error}"
2265            )),
2266        })?;
2267        if output_reader.num_docs() != doc_count {
2268            return Err(Error::Corruption(format!(
2269                "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2270                output_reader.num_docs(),
2271            )));
2272        }
2273        let has_surviving_bmp = !output_reader.bmp_indexes().is_empty();
2274        let ann_fragmented = output_reader
2275            .vector_indexes()
2276            .iter()
2277            .any(|(&field, index)| {
2278                matches!(
2279                    index,
2280                    crate::segment::VectorIndex::BinaryIvf(_)
2281                        | crate::segment::VectorIndex::ScannBinary(_)
2282                ) && output_reader
2283                    .ann_health(crate::dsl::Field(field))
2284                    .is_some_and(|health| health.fragmentation() > 1.0)
2285            });
2286        let seismic_pending_terms =
2287            output_reader
2288                .seismic_indexes()
2289                .values()
2290                .try_fold(0u32, |total, index| {
2291                    total.checked_add(index.pending_terms()).ok_or_else(|| {
2292                        Error::Corruption("Seismic maintenance debt exceeds u32".into())
2293                    })
2294                })?;
2295        drop(output_reader);
2296
2297        let mut st = Arc::clone(&self.state).lock_owned().await;
2298        // Every source must still be live: callers hold operation ownership,
2299        // so a missing source means a stale merge/reorder whose input was
2300        // already replaced. Adding the output would duplicate its documents.
2301        let missing: Vec<&String> = old_ids
2302            .iter()
2303            .filter(|id| !st.metadata.has_segment(id))
2304            .collect();
2305        if !missing.is_empty() {
2306            return Err(Error::Corruption(format!(
2307                "replace_segments: source segment(s) {:?} not in metadata — \
2308                 refusing to add output {} (would duplicate documents)",
2309                missing, new_id
2310            )));
2311        }
2312
2313        if let Some(expected) = expected_deletions {
2314            if old_ids
2315                .iter()
2316                .any(|id| expected.get(id) != Some(&st.metadata.segment_metas[id].deletions))
2317            {
2318                return Err(Error::Internal(
2319                    "row visibility changed during compaction; retry with a fresh snapshot".into(),
2320                ));
2321            }
2322        } else if matches!(layout, ReplacementLayout::Compacted) {
2323            return Err(Error::Internal(
2324                "compaction requires a visibility snapshot".into(),
2325            ));
2326        }
2327        let mut replacement_info = match layout {
2328            ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2329                let generation = old_ids
2330                    .iter()
2331                    .filter_map(|id| st.metadata.segment_metas.get(id))
2332                    .map(|info| info.generation)
2333                    .max()
2334                    .unwrap_or(0)
2335                    .checked_add(1)
2336                    .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2337                let parent_unconverged_passes = old_ids
2338                    .iter()
2339                    .filter_map(|id| st.metadata.segment_metas.get(id))
2340                    .map(|info| info.bp_unconverged_passes)
2341                    .max()
2342                    .unwrap_or(0);
2343                let parent_has_debt = old_ids
2344                    .iter()
2345                    .filter_map(|id| st.metadata.segment_metas.get(id))
2346                    .any(|info| !info.bp_converged);
2347                let (reordered, bp_converged, bp_unconverged_passes) =
2348                    replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2349                SegmentMetaInfo {
2350                    deletions: None,
2351                    num_docs: doc_count,
2352                    ancestors: old_ids.to_vec(),
2353                    generation,
2354                    reordered,
2355                    bp_converged,
2356                    bp_unconverged_passes,
2357                    seismic_pending_terms,
2358                    seismic_maintenance_passes: 0,
2359                    seismic_no_progress_passes: 0,
2360                    ann_fragmented,
2361                }
2362            }
2363            ReplacementLayout::PreserveSingleSource
2364            | ReplacementLayout::Compacted
2365            | ReplacementLayout::MaintenanceOnly => {
2366                let [source_id] = old_ids else {
2367                    return Err(Error::Internal(
2368                        "layout-preserving replacement requires exactly one source".into(),
2369                    ));
2370                };
2371                let mut source = st
2372                    .metadata
2373                    .segment_metas
2374                    .get(source_id)
2375                    .cloned()
2376                    .ok_or_else(|| {
2377                        Error::Corruption(format!(
2378                            "layout-preserving replacement source {source_id} disappeared"
2379                        ))
2380                    })?;
2381                source.num_docs = doc_count;
2382                if matches!(
2383                    layout,
2384                    ReplacementLayout::Compacted | ReplacementLayout::MaintenanceOnly
2385                ) {
2386                    source.ancestors = old_ids.to_vec();
2387                    source.generation = source.generation.checked_add(1).ok_or_else(|| {
2388                        Error::Corruption("maintenance generation overflow".into())
2389                    })?;
2390                }
2391                if matches!(layout, ReplacementLayout::Compacted) {
2392                    source.deletions = None;
2393                    if has_surviving_bmp && source.reordered {
2394                        source.bp_converged = false;
2395                    }
2396                }
2397                source
2398            }
2399        };
2400        replacement_info.seismic_pending_terms = seismic_pending_terms;
2401        replacement_info.ann_fragmented = ann_fragmented;
2402        let (passes, no_progress_passes) = replacement_seismic_state(
2403            old_ids.iter().map(|id| &st.metadata.segment_metas[id]),
2404            seismic_pending_terms,
2405            layout,
2406        );
2407        replacement_info.seismic_maintenance_passes = passes;
2408        replacement_info.seismic_no_progress_passes = no_progress_passes;
2409        // Take the latest source visibility under the publication lock. A
2410        // deletion may have committed while encoded payloads were being copied.
2411        // Field-only BP and ANN rewrites preserve the same physical row order.
2412        let mut retired_ids = old_ids.to_vec();
2413        let mut deletion_output = None;
2414        let has_deletions = old_ids
2415            .iter()
2416            .any(|id| st.metadata.segment_metas[id].deletions.is_some());
2417        if matches!(layout, ReplacementLayout::Compacted) {
2418            for id in old_ids {
2419                if let Some(deletion) = &st.metadata.segment_metas[id].deletions {
2420                    retired_ids.push(deletion.id.clone());
2421                }
2422            }
2423        } else if has_deletions && old_ids.len() == 1 {
2424            let source = &st.metadata.segment_metas[&old_ids[0]];
2425            if source.num_docs != doc_count {
2426                return Err(Error::Corruption(
2427                    "layout-preserving replacement changed physical row count".into(),
2428                ));
2429            }
2430            // BP/ANN rewrites retain physical addresses; share the exact
2431            // immutable mask generation rather than copying it again.
2432            replacement_info.deletions = source.deletions.clone();
2433        } else if has_deletions {
2434            let mut alive = crate::query::DocBitset::all(doc_count);
2435            let mut offset = 0u32;
2436            for id in old_ids {
2437                let info = &st.metadata.segment_metas[id];
2438                if let Some(deletion) = &info.deletions {
2439                    let source = deletion
2440                        .load(self.directory.as_ref(), info.num_docs)
2441                        .await?;
2442                    crate::segment::deletion::append_dead_rows(
2443                        &mut alive,
2444                        &source,
2445                        info.num_docs,
2446                        offset,
2447                    )?;
2448                    retired_ids.push(deletion.id.clone());
2449                }
2450                offset = offset
2451                    .checked_add(info.num_docs)
2452                    .ok_or_else(|| Error::Corruption("merge row count overflow".into()))?;
2453            }
2454            if offset != doc_count {
2455                return Err(Error::Corruption(
2456                    "layout-preserving merge changed physical row count".into(),
2457                ));
2458            }
2459            let deletion_id = SegmentId::new();
2460            let claim = self.protect_new_segment(deletion_id.to_hex())?;
2461            deletion_output = Some((claim, self.output_cleanup_guard(deletion_id)));
2462            replacement_info.deletions = Some(
2463                crate::segment::deletion::write(
2464                    self.directory.as_ref(),
2465                    deletion_id,
2466                    doc_count,
2467                    &alive,
2468                )
2469                .await?,
2470            );
2471        }
2472        let mut next = st.metadata.clone();
2473        for id in old_ids {
2474            next.remove_segment(id);
2475        }
2476        next.add_segment_meta(new_id.clone(), replacement_info);
2477        retired_ids.sort_unstable();
2478        retired_ids.dedup();
2479        retired_ids.retain(|id| !next.owns_id(id));
2480
2481        let directory = Arc::clone(&self.directory);
2482        let tracker = Arc::clone(&self.tracker);
2483        let replacement_refresh = self.replacement_refresh.read().clone();
2484        let manager = Arc::clone(self);
2485        let index_label = self.schema.index_label().to_owned();
2486        self.run_lifecycle_transaction(async move {
2487            // Durable-before-visible. If persistence fails, old metadata and
2488            // tracker ownership stay intact and source deletion is never armed.
2489            next.save(directory.as_ref()).await?;
2490            tracker.register(&new_id);
2491            if let Some(deletion) = &next.segment_metas[&new_id].deletions {
2492                tracker.register(&deletion.id);
2493            }
2494            if let Some((_, cleanup)) = deletion_output.as_mut() {
2495                cleanup.disarm();
2496            }
2497            st.metadata = next;
2498
2499            manager.retire_reorder_retries(&retired_ids);
2500
2501            // Keep state locked until retired sources enter the tracker. The
2502            // transaction itself also performs deletion, so cancellation of
2503            // the requesting merge cannot strand pending-deletion ownership.
2504            let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2505            drop(st);
2506            for &segment_id in &ready_to_delete {
2507                if let Err(error) =
2508                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
2509                {
2510                    log::warn!(
2511                        "[segment_cleanup] index={index_label} immediate delete failed for {}: {}",
2512                        segment_id.to_hex(),
2513                        error,
2514                    );
2515                }
2516            }
2517            tracker.complete_deletion(&ready_to_delete);
2518            refresh_replacement_topology(replacement_refresh, &index_label).await;
2519            Ok(())
2520        })
2521        .await
2522    }
2523
2524    /// Perform the actual merge operation (pure function — no shared state access).
2525    /// `output_segment_id` is pre-generated by the caller so active-operation ownership
2526    /// is installed before any output file is written.
2527    /// Returns (new_segment_id_hex, total_doc_count).
2528    #[allow(clippy::too_many_arguments)]
2529    async fn do_merge(
2530        directory: &D,
2531        schema: &Arc<crate::dsl::Schema>,
2532        segment_ids_to_merge: &[String],
2533        output_segment_id: SegmentId,
2534        term_cache_blocks: usize,
2535        term_cache_budget_bytes: Option<usize>,
2536        optimization: crate::structures::IndexOptimization,
2537        posting_codec: crate::structures::PostingCodec,
2538        term_dict_block_size: crate::structures::SSTableBlockSize,
2539        trained: Option<&TrainedVectorStructures>,
2540        reorder_bmp: bool,
2541        granularity: crate::segment::reorder::BpGranularity,
2542        merge_bp_time_budget: Option<std::time::Duration>,
2543        bp_memory_budget_bytes: usize,
2544        reorder_permits: Arc<ReorderConcurrencyGate>,
2545        reorder_priority: ReorderPriority,
2546        cancellation: Arc<AtomicBool>,
2547        bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2548    ) -> MergeTaskResult<(String, u32, bool)> {
2549        let output_hex = output_segment_id.to_hex();
2550        let load_start = std::time::Instant::now();
2551
2552        let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2553        for id_str in segment_ids_to_merge {
2554            let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2555                MergeTaskError::source(
2556                    id_str.clone(),
2557                    Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2558                )
2559            })?;
2560            segment_ids.push(id);
2561        }
2562
2563        // Cheap fail-fast before opening every reader. `join_all` otherwise
2564        // waits for all healthy multi-GB inputs to load even when one source's
2565        // `.meta` is already absent, turning a known-corrupt candidate into a
2566        // large CPU/IO spike before it can be quarantined.
2567        let mut unavailable_sources = Vec::new();
2568        let mut missing_files = Vec::new();
2569        for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2570            let files = SegmentFiles::new(id.0);
2571            let mut source_unavailable = false;
2572            for path in files.mandatory_paths() {
2573                let exists = directory
2574                    .exists(path)
2575                    .await
2576                    .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2577                if !exists {
2578                    source_unavailable = true;
2579                    missing_files.push(format!("{}:{:?}", id_str, path));
2580                }
2581            }
2582            if source_unavailable {
2583                unavailable_sources.push(id_str.clone());
2584            }
2585        }
2586        if !unavailable_sources.is_empty() {
2587            return Err(MergeTaskError::sources(
2588                unavailable_sources,
2589                Error::Corruption(format!(
2590                    "merge sources are missing mandatory files: {}",
2591                    missing_files.join(", ")
2592                )),
2593            ));
2594        }
2595
2596        let schema_arc = Arc::clone(schema);
2597        let futures: Vec<_> = segment_ids
2598            .iter()
2599            .map(|&sid| {
2600                let sch = Arc::clone(&schema_arc);
2601                async move {
2602                    SegmentReader::open_with_term_cache_budget(
2603                        directory,
2604                        sid,
2605                        sch,
2606                        term_cache_blocks,
2607                        term_cache_budget_bytes,
2608                    )
2609                    .await
2610                }
2611            })
2612            .collect();
2613
2614        let results = futures::future::join_all(futures).await;
2615        let mut readers = Vec::with_capacity(results.len());
2616        let mut total_docs = 0u64;
2617        for (i, result) in results.into_iter().enumerate() {
2618            match result {
2619                Ok(r) => {
2620                    total_docs += r.meta().num_docs as u64;
2621                    readers.push(r);
2622                }
2623                Err(e) => {
2624                    log::error!(
2625                        "[merge] index={} Failed to open segment {}: {:?}",
2626                        schema.index_label(),
2627                        segment_ids_to_merge[i],
2628                        e
2629                    );
2630                    return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2631                }
2632            }
2633        }
2634        if total_docs > u32::MAX as u64 {
2635            return Err(Error::Internal(format!(
2636                "Merged segment doc count ({}) exceeds u32::MAX",
2637                total_docs
2638            ))
2639            .into());
2640        }
2641
2642        // Pre-merge validation: verify each source segment's store doc count
2643        // matches its metadata. Catching mismatches early avoids building a
2644        // corrupted merged segment and leaving orphan files on disk.
2645        for (i, reader) in readers.iter().enumerate() {
2646            let meta_docs = reader.meta().num_docs;
2647            let store_docs = reader.store().num_docs();
2648            if store_docs != meta_docs {
2649                return Err(MergeTaskError::source(
2650                    segment_ids_to_merge[i].clone(),
2651                    Error::Corruption(format!(
2652                        "pre-merge validation: segment {} store has {} docs but meta says {}",
2653                        segment_ids_to_merge[i], store_docs, meta_docs
2654                    )),
2655                ));
2656            }
2657        }
2658
2659        log::info!(
2660            "[merge] index={} loaded {} segment readers in {:.1}s",
2661            schema.index_label(),
2662            readers.len(),
2663            load_start.elapsed().as_secs_f64()
2664        );
2665
2666        let merger = SegmentMerger::new(Arc::clone(schema))
2667            .with_posting_config(optimization, posting_codec)
2668            .with_term_dict_block_size(term_dict_block_size)
2669            .with_reorder_fields(reorder_bmp)
2670            .with_granularity(granularity)
2671            .with_bp_budget(crate::segment::BpBudget {
2672                min_partition_docs: None,
2673                time_budget: merge_bp_time_budget,
2674            })
2675            .with_cancellation(cancellation)
2676            .with_bp_memory_budget(bp_memory_budget_bytes)
2677            .with_reorder_permits(reorder_permits)
2678            .with_reorder_priority(reorder_priority)
2679            .with_background_pool(bg_cpu_pool);
2680
2681        log::info!(
2682            "[merge] index={} {} segments -> {} (trained={})",
2683            schema.index_label(),
2684            segment_ids_to_merge.len(),
2685            output_hex,
2686            trained.map_or(0, |t| t.centroids.len()),
2687        );
2688
2689        let (_merged_meta, merge_stats) = merger
2690            .merge(directory, &readers, output_segment_id, trained)
2691            .await
2692            .map_err(|error| {
2693                if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2694                    // The merge has already opened every input successfully;
2695                    // a structural/serialization failure is deterministic for
2696                    // this candidate. Attribute all inputs rather than running
2697                    // the same multi-GB rewrite forever. This is deliberately
2698                    // not used for I/O errors, which may be transient/output-side.
2699                    MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2700                } else {
2701                    MergeTaskError::from(error)
2702                }
2703            })?;
2704        let bp_converged = merge_stats.bp_converged;
2705        if !bp_converged {
2706            log::info!(
2707                "[merge] index={} merge-time BP hit its wall-clock budget — output marked unconverged; \
2708                 the background optimizer deepens it later",
2709                schema.index_label(),
2710            );
2711        }
2712
2713        log::info!(
2714            "[merge] index={} total wall-clock: {:.1}s ({} segments, {} docs)",
2715            schema.index_label(),
2716            load_start.elapsed().as_secs_f64(),
2717            readers.len(),
2718            total_docs,
2719        );
2720
2721        Ok((output_hex, total_docs as u32, bp_converged))
2722    }
2723
2724    /// Drain all in-flight merge tasks safely.
2725    ///
2726    /// Merge/reorder writers use synchronous block-in-place sections, so an
2727    /// abort request takes effect only after the owned writer reaches its next
2728    /// await. Draining the complete task remains necessary before index
2729    /// deletion or orphan cleanup can safely proceed.
2730    ///
2731    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2732    /// handles to `merge_handles` so later drains still see in-flight merges.
2733    pub async fn abort_merges(&self) {
2734        loop {
2735            let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2736            if handles.is_empty() {
2737                return;
2738            }
2739            while let Some(result) = handles.join_next().await {
2740                if let Err(error) = result
2741                    && error.is_panic()
2742                {
2743                    log::error!(
2744                        "[merge] index={} background task panicked while draining: {}",
2745                        self.schema.index_label(),
2746                        error
2747                    );
2748                }
2749            }
2750        }
2751    }
2752
2753    /// Wait for all current in-flight merges to complete.
2754    ///
2755    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2756    /// handles to `merge_handles` so later drains still see in-flight merges.
2757    pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2758        let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2759        while handles.join_next().await.is_some() {}
2760    }
2761
2762    /// Wait for all eligible merges to complete, including cascading merges.
2763    ///
2764    /// Drains current handles, then loops. Each completed merge auto-triggers
2765    /// `maybe_merge` (which pushes new handles) before its JoinHandle resolves,
2766    /// so by the time `join_next` returns all cascading handles are registered.
2767    ///
2768    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2769    /// handles to `merge_handles` so later drains still see in-flight merges.
2770    pub async fn wait_for_all_merges(self: &Arc<Self>) {
2771        loop {
2772            let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2773            if handles.is_empty() {
2774                break;
2775            }
2776            while handles.join_next().await.is_some() {}
2777        }
2778    }
2779
2780    /// Complete the second half of shutdown after the owning `IndexWriter`
2781    /// has been dropped. This drains tracked merges and then waits for every
2782    /// remaining guard, including optimizer reorders that are intentionally
2783    /// launched outside the writer lock.
2784    pub async fn wait_for_shutdown(self: &Arc<Self>) {
2785        self.wait_for_all_merges().await;
2786        self.active_operations.wait_until_idle().await;
2787        loop {
2788            let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2789            if handles.is_empty() {
2790                break;
2791            }
2792            for handle in handles {
2793                if let Err(error) = handle.await
2794                    && error.is_panic()
2795                {
2796                    log::error!(
2797                        "[segment_cleanup] index={} task panicked while draining: {}",
2798                        self.schema.index_label(),
2799                        error
2800                    );
2801                }
2802            }
2803        }
2804    }
2805
2806    /// Force merge all segments into one (packing up to the u32 document
2807    /// format limit per output).
2808    ///
2809    /// An explicit force merge consolidates segments while retaining tombstones.
2810    /// Unlike background merges, it ignores the policy's `max_segment_docs` cap
2811    /// (which exists to bound background BP/merge cost, not to keep an index
2812    /// permanently split). Only the u32 doc-id format limit can force more
2813    /// than one output, and that outcome is logged loudly.
2814    ///
2815    /// Each batch is registered in `active_operations` via an RAII guard to prevent
2816    /// `maybe_merge` from spawning a conflicting background merge.
2817    pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2818        self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2819            .await
2820    }
2821
2822    /// Force merge while refreshing long-lived segment snapshots after the
2823    /// initial background-merge drain and every durable replacement.
2824    ///
2825    /// Long-lived consumers such as the primary-key index and cached
2826    /// `IndexReader` hold segment snapshots. Refreshing them only after the
2827    /// complete force merge retains every retired source file for the entire
2828    /// operation, which can temporarily double a large index on disk. The
2829    /// hook lets the owning `IndexWriter`/server advance those snapshots after
2830    /// each batch while the force merge remains otherwise memory bounded.
2831    pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2832        self: &Arc<Self>,
2833        refresh_snapshots: F,
2834    ) -> Result<()>
2835    where
2836        F: FnMut() -> Fut,
2837        Fut: std::future::Future<Output = Result<()>>,
2838    {
2839        self.force_merge_with_compaction_and_snapshot_refresh(None, refresh_snapshots)
2840            .await
2841    }
2842
2843    pub(crate) async fn force_merge_with_compaction_and_snapshot_refresh<F, Fut>(
2844        self: &Arc<Self>,
2845        compaction_budget: Option<usize>,
2846        mut refresh_snapshots: F,
2847    ) -> Result<()>
2848    where
2849        F: FnMut() -> Fut,
2850        Fut: std::future::Future<Output = Result<()>>,
2851    {
2852        // Conflicting owners that never appear in `merge_handles` (background
2853        // reorders, a concurrent force-merge) can hold a batch segment for
2854        // minutes to hours; retrying without parking would busy-spin a runtime
2855        // worker and hammer the state mutex for that whole window.
2856        const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2857            std::time::Duration::from_millis(100);
2858        // When every remaining mergeable segment is owned by another
2859        // operation (e.g. background BP reorders), there is nothing to do but
2860        // wait for a release; those passes run for minutes, so poll slowly.
2861        const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2862
2863        let (_force_merge_activity, policy_segment_docs) = {
2864            let st = self.state.lock().await;
2865            // `maybe_merge` registers and publishes every task handle under
2866            // this same state lock. Raising the barrier here therefore makes
2867            // the subsequent drain race-free.
2868            self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2869            (
2870                ForceMergeActivityGuard(&self.force_merge_active),
2871                st.merge_policy.max_segment_docs(),
2872            )
2873        };
2874
2875        // Wait for all in-flight background merges (including cascading)
2876        // before starting forced merges to avoid try_register conflicts.
2877        let background_merges = self
2878            .merge_handles
2879            .lock()
2880            .iter()
2881            .filter(|handle| !handle.is_finished())
2882            .count();
2883        if background_merges > 0 {
2884            log::info!(
2885                "[force_merge] index={} waiting for {} in-flight background merge(s) before planning",
2886                self.schema.index_label(),
2887                background_merges,
2888            );
2889        }
2890        let drain_start = std::time::Instant::now();
2891        self.wait_for_all_merges().await;
2892        if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2893            log::info!(
2894                "[force_merge] index={} drained background merges in {:.1}s",
2895                self.schema.index_label(),
2896                drain_start.elapsed().as_secs_f64(),
2897            );
2898        }
2899
2900        // A background merge may have published after the writer's preceding
2901        // commit refresh but before the drain completed. Advance PK/search
2902        // snapshots now so its retired sources do not survive until the first
2903        // forced replacement (or forever when no batch is needed).
2904        let refresh_start = std::time::Instant::now();
2905        refresh_snapshots().await?;
2906        if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2907            log::info!(
2908                "[force_merge] index={} initial snapshot refresh took {:.1}s",
2909                self.schema.index_label(),
2910                refresh_start.elapsed().as_secs_f64(),
2911            );
2912        }
2913
2914        // Outputs produced by a completed bin are already maximal according
2915        // to the BFD plan. Excluding them from later replanning prevents an
2916        // expensive final BP output from being selected and reordered again.
2917        let mut completed_outputs = HashSet::new();
2918        // One INFO line per wait episode, not per 1s poll; DEBUG afterwards.
2919        let mut logged_held_wait = false;
2920
2921        loop {
2922            if !self.active_operations.is_accepting() {
2923                return Err(Error::IndexClosed);
2924            }
2925
2926            let segments: Vec<(String, u32)> = {
2927                let st = self.state.lock().await;
2928                st.metadata
2929                    .segment_metas
2930                    .iter()
2931                    .filter(|(id, _)| !completed_outputs.contains(*id))
2932                    .map(|(id, info)| (id.clone(), info.num_docs))
2933                    .collect()
2934            };
2935
2936            // Route around segments owned by active operations instead of
2937            // retrying the same conflicting batch. Group ownership is claimed
2938            // before foreground capacity, so an in-flight optimizer/retrain
2939            // can finish without a capacity/ownership lock inversion.
2940            let active_ids = self.active_operations.snapshot();
2941            let held = segments
2942                .iter()
2943                .filter(|(id, _)| active_ids.contains(id))
2944                .count();
2945            let free_segments: Vec<_> = segments
2946                .into_iter()
2947                .filter(|(id, _)| !active_ids.contains(id))
2948                .collect();
2949            // Every segment format and doc map uses u32 document IDs.
2950            // An explicit force merge packs up to that format limit: the
2951            // policy's `max_segment_docs` bounds *background* merge cost and
2952            // must not silently leave a forced compaction above one segment
2953            // (observed in prod: an 8.8M-doc index stuck at 2 segments under
2954            // a 5M-doc policy cap).
2955            let max_docs = u64::from(u32::MAX);
2956            let planned_groups = plan_force_merge_groups(free_segments, max_docs);
2957            if let Some(cap) = policy_segment_docs {
2958                for group in planned_groups
2959                    .iter()
2960                    .filter(|group| group.segments.len() >= 2)
2961                    .filter(|group| group.total_docs > u64::from(cap))
2962                {
2963                    log::warn!(
2964                        "[force_merge] index={} output of {} docs intentionally exceeds the \
2965                         background merge policy cap of {} docs (force merge compacts to the \
2966                         u32 format limit)",
2967                        self.schema.index_label(),
2968                        group.total_docs,
2969                        cap,
2970                    );
2971                }
2972            }
2973            let next_group = planned_groups
2974                .into_iter()
2975                .find(|group| group.segments.len() >= 2);
2976
2977            let Some(group) = next_group else {
2978                if held == 0 {
2979                    if !completed_outputs.is_empty() {
2980                        // A segment held while earlier groups ran may now fit
2981                        // with one of those outputs. Reconsider completed bins
2982                        // once before declaring convergence. In the normal
2983                        // no-conflict path BFD groups are pairwise maximal, so
2984                        // this extra planning round performs no merge.
2985                        completed_outputs.clear();
2986                        continue;
2987                    }
2988                    if let Some(memory_budget) = compaction_budget {
2989                        let dirty: Vec<_> = {
2990                            let st = self.state.lock().await;
2991                            st.metadata
2992                                .segment_metas
2993                                .iter()
2994                                .filter(|(_, meta)| meta.deletions.is_some())
2995                                .map(|(id, _)| id.clone())
2996                                .collect()
2997                        };
2998                        for id in dirty {
2999                            self.compact_segment(&id, memory_budget).await?;
3000                            refresh_snapshots().await?;
3001                        }
3002                    }
3003                    // Every remaining free segment is either already at the
3004                    // u32 format limit or cannot be paired without exceeding
3005                    // it. More than one leftover segment is an exceptional,
3006                    // loudly-reported outcome — never a silent policy effect.
3007                    let remaining = {
3008                        let st = self.state.lock().await;
3009                        st.metadata.segment_metas.len()
3010                    };
3011                    if remaining > 1 {
3012                        log::warn!(
3013                            "[force_merge] index={} finished with {} segments: combined \
3014                             document count exceeds the u32 segment format limit, so a \
3015                             single output is impossible",
3016                            self.schema.index_label(),
3017                            remaining,
3018                        );
3019                    }
3020                    // A reorder that was already running when foreground
3021                    // admission began may have published after the preceding
3022                    // callback. Reconcile external readers one final time so
3023                    // they cannot retain its retired source indefinitely.
3024                    refresh_snapshots().await?;
3025                    return Ok(());
3026                }
3027                if !logged_held_wait {
3028                    log::info!(
3029                        "[force_merge] index={} waiting: {} segment(s) held by active \
3030                         merge/reorder operations, no free group can merge",
3031                        self.schema.index_label(),
3032                        held
3033                    );
3034                    logged_held_wait = true;
3035                } else {
3036                    log::debug!(
3037                        "[force_merge] index={} still waiting on {} held segment(s)",
3038                        self.schema.index_label(),
3039                        held
3040                    );
3041                }
3042                #[cfg(test)]
3043                self.force_merge_conflict_retries
3044                    .fetch_add(1, Ordering::Relaxed);
3045                tokio::select! {
3046                    biased;
3047                    () = self.active_operations.wait_for_shutdown() => {
3048                        return Err(Error::IndexClosed);
3049                    }
3050                    () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
3051                }
3052                continue;
3053            };
3054            logged_held_wait = false;
3055
3056            // Claim the complete final group and all of its future output IDs
3057            // up front. This freezes the hierarchy while allowing every
3058            // durable intermediate replacement to release its source files.
3059            let hierarchy = plan_force_merge_hierarchy(group.segments.len());
3060            let output_ids: Vec<_> = (0..hierarchy.steps.len())
3061                .map(|_| SegmentId::new())
3062                .collect();
3063            let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
3064            let mut all_ids = source_ids.clone();
3065            all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
3066            let group_guard = {
3067                let st = self.state.lock().await;
3068                source_ids
3069                    .iter()
3070                    .all(|id| st.metadata.has_segment(id))
3071                    .then(|| self.active_operations.try_register(all_ids))
3072                    .flatten()
3073            };
3074            let _group_guard = Arc::new(match group_guard {
3075                Some(guard) => guard,
3076                None if !self.active_operations.is_accepting() => {
3077                    return Err(Error::IndexClosed);
3078                }
3079                None => {
3080                    #[cfg(test)]
3081                    self.force_merge_conflict_retries
3082                        .fetch_add(1, Ordering::Relaxed);
3083                    log::debug!(
3084                        "[force_merge] index={} group lost a registration race, replanning",
3085                        self.schema.index_label()
3086                    );
3087                    let had_tracked_merges = !self.merge_handles.lock().is_empty();
3088                    self.wait_for_merging_thread().await;
3089                    if !had_tracked_merges {
3090                        tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
3091                    }
3092                    continue;
3093                }
3094            });
3095
3096            log::info!(
3097                "[force_merge] index={} planned final group: {} segments, {} docs, {} merge pass(es)",
3098                self.schema.index_label(),
3099                group.segments.len(),
3100                group.total_docs,
3101                output_ids.len(),
3102            );
3103
3104            // Claiming the complete group before any capacity wait prevents a
3105            // vector-generation pause from waiting for a global slot retained
3106            // by force merge while force merge waits for that pause to end.
3107            //
3108            // Background merges acquire global merge capacity before the
3109            // shared BP gate. Preserve that order here: wait for one group
3110            // slot while background BP remains admitted, then pause new BP
3111            // passes. The retained slot guarantees this hierarchy can make
3112            // progress even if other indexes subsequently park at the gate.
3113            let group_global_merge_permit = if self.reorder_on_merge {
3114                let capacity_start = std::time::Instant::now();
3115                let permit = tokio::select! {
3116                    biased;
3117                    () = self.active_operations.wait_for_shutdown() => {
3118                        return Err(Error::IndexClosed);
3119                    }
3120                    permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3121                        permit.map_err(|_| {
3122                            Error::Internal(
3123                                "global background merge scheduler is closed".into(),
3124                            )
3125                        })?
3126                    }
3127                };
3128                if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
3129                    log::info!(
3130                        "[force_merge] index={} waited {:.1}s for foreground global merge capacity",
3131                        self.schema.index_label(),
3132                        capacity_start.elapsed().as_secs_f64(),
3133                    );
3134                }
3135                Some(Arc::new(permit))
3136            } else {
3137                None
3138            };
3139            let _foreground_reorder = if self.reorder_on_merge {
3140                log::info!(
3141                    "[force_merge] index={} prioritizing BP capacity ({} total pass slot(s))",
3142                    self.schema.index_label(),
3143                    self.reorder_permits.limit(),
3144                );
3145                let admission_start = std::time::Instant::now();
3146                let guard = Arc::clone(&self.reorder_permits)
3147                    .begin_foreground()
3148                    .await
3149                    .map_err(|_| {
3150                        Error::Internal("background reorder scheduler is closed".into())
3151                    })?;
3152                if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
3153                    log::info!(
3154                        "[force_merge] index={} acquired foreground BP capacity in {:.1}s",
3155                        self.schema.index_label(),
3156                        admission_start.elapsed().as_secs_f64(),
3157                    );
3158                }
3159                Some(Arc::new(guard))
3160            } else {
3161                None
3162            };
3163
3164            let source_count = group.segments.len();
3165            let mut nodes: Vec<Option<(String, u32)>> =
3166                group.segments.into_iter().map(Some).collect();
3167            nodes.resize_with(source_count + hierarchy.steps.len(), || None);
3168            for (step_index, step) in hierarchy.steps.iter().enumerate() {
3169                let final_pass = step_index + 1 == hierarchy.steps.len();
3170                let mut batch_entries = Vec::with_capacity(step.inputs.len());
3171                for &node in &step.inputs {
3172                    let entry = nodes
3173                        .get_mut(node)
3174                        .and_then(Option::take)
3175                        .expect("force-merge hierarchy must reference an available node");
3176                    batch_entries.push(entry);
3177                }
3178                let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
3179                let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
3180                let output_id = output_ids[step_index];
3181
3182                let capacity_start = std::time::Instant::now();
3183                let step_global_merge_permit = if group_global_merge_permit.is_none() {
3184                    Some(Arc::new(tokio::select! {
3185                        biased;
3186                        () = self.active_operations.wait_for_shutdown() => {
3187                            return Err(Error::IndexClosed);
3188                        }
3189                        permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3190                            permit.map_err(|_| {
3191                                Error::Internal(
3192                                    "global background merge scheduler is closed".into(),
3193                                )
3194                            })?
3195                        }
3196                    }))
3197                } else {
3198                    None
3199                };
3200                if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
3201                    log::info!(
3202                        "[force_merge] index={} waited {:.1}s for global merge capacity",
3203                        self.schema.index_label(),
3204                        capacity_start.elapsed().as_secs_f64(),
3205                    );
3206                }
3207
3208                // Intermediate reductions are streaming block-copy merges.
3209                // Only the final pass pays BP, so a group with hundreds of
3210                // tiny sources never reorders the same documents repeatedly.
3211                let reorder_bmp = final_pass && self.reorder_on_merge;
3212                log::info!(
3213                    "[force_merge] index={} {} pass: {} segments ({} docs, bp={})",
3214                    self.schema.index_label(),
3215                    if final_pass {
3216                        "final"
3217                    } else {
3218                        "fan-in reduction"
3219                    },
3220                    batch.len(),
3221                    batch_docs,
3222                    reorder_bmp,
3223                );
3224                let (new_segment_id, total_docs, _) = self
3225                    .merge_and_replace_registered(
3226                        &batch,
3227                        output_id,
3228                        reorder_bmp,
3229                        ReorderPriority::Foreground,
3230                        Arc::new((
3231                            Arc::clone(&_group_guard),
3232                            group_global_merge_permit.clone(),
3233                            _foreground_reorder.clone(),
3234                            step_global_merge_permit.clone(),
3235                        )),
3236                    )
3237                    .await
3238                    .map_err(|error| error.error)?;
3239                drop(step_global_merge_permit);
3240
3241                // Advance PK/read snapshots after every hierarchy level so
3242                // retired sources do not accumulate until the final output.
3243                let refresh_start = std::time::Instant::now();
3244                refresh_snapshots().await?;
3245                if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
3246                    log::info!(
3247                        "[force_merge] index={} post-replacement snapshot refresh took {:.1}s",
3248                        self.schema.index_label(),
3249                        refresh_start.elapsed().as_secs_f64(),
3250                    );
3251                }
3252
3253                let output_node = source_count + step_index;
3254                debug_assert!(nodes[output_node].is_none());
3255                nodes[output_node] = Some((new_segment_id, total_docs));
3256            }
3257            let (root_id, _) = nodes[hierarchy.root]
3258                .take()
3259                .expect("force-merge hierarchy must produce its root");
3260            debug_assert!(nodes.into_iter().all(|node| node.is_none()));
3261            completed_outputs.insert(root_id);
3262        }
3263    }
3264
3265    fn segment_needs_vector_rewrite(
3266        &self,
3267        schema: &crate::dsl::Schema,
3268        reader: &SegmentReader,
3269        field_ids: &[u32],
3270        trained: &TrainedVectorStructures,
3271        rewrite_existing: bool,
3272    ) -> Result<bool> {
3273        for &field_id in field_ids {
3274            let flat = reader.flat_vectors().get(&field_id);
3275            let ann = reader.vector_indexes().get(&field_id);
3276            if ann.is_some() && flat.is_none() {
3277                return Err(Error::Corruption(format!(
3278                    "segment {:032x} field {field_id} has ANN data without the required flat vectors",
3279                    reader.meta().id,
3280                )));
3281            }
3282
3283            let Some(flat) = flat else {
3284                continue;
3285            };
3286            if flat.num_vectors == 0 {
3287                continue;
3288            }
3289            if rewrite_existing {
3290                return Ok(true);
3291            }
3292            let field = crate::dsl::Field(field_id);
3293            let entry = schema.get_field_entry(field).ok_or_else(|| {
3294                Error::Corruption(format!(
3295                    "segment {:032x} references unknown vector field {field_id}",
3296                    reader.meta().id,
3297                ))
3298            })?;
3299            let current = match entry.field_type {
3300                // TQ payloads carry no trained generation, so an existing one
3301                // is always current; a tq field must never be staged for a
3302                // vector-generation rewrite.
3303                crate::dsl::FieldType::DenseVector
3304                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
3305                        config.index_type == crate::dsl::VectorIndexType::Tq
3306                    }) =>
3307                {
3308                    matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
3309                }
3310                crate::dsl::FieldType::DenseVector
3311                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
3312                        config.index_type == crate::dsl::VectorIndexType::IvfTq
3313                    }) =>
3314                {
3315                    let config = entry
3316                        .dense_vector_config
3317                        .as_ref()
3318                        .expect("matched IVF-TQ configuration");
3319                    match (ann, trained.centroids.get(&field_id)) {
3320                        (
3321                            Some(crate::segment::VectorIndex::IvfTq { index, .. }),
3322                            Some(centroids),
3323                        ) => {
3324                            let header = index.get().header();
3325                            crate::structures::is_ivf_tq_cosine_generation(centroids.version)
3326                                && crate::structures::is_ivf_tq_cosine_generation(
3327                                    header.quantizer_version,
3328                                )
3329                                && header.dim == config.dim
3330                                && header.num_clusters == centroids.num_clusters
3331                                && header.quantizer_version == centroids.version
3332                                && header.codebook_version
3333                                    == crate::structures::vector::quantization::tq_expected_fingerprint(
3334                                        config.dim,
3335                                    )
3336                                && header.routing == config.ivf_routing
3337                        }
3338                        (None, None) => true,
3339                        _ => false,
3340                    }
3341                }
3342                crate::dsl::FieldType::DenseVector
3343                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
3344                        config.index_type == crate::dsl::VectorIndexType::Scann
3345                    }) =>
3346                {
3347                    match (ann, trained.scann_artifacts.get(&field_id)) {
3348                        (Some(crate::segment::VectorIndex::ScannAh(index)), Some(artifact)) => {
3349                            index
3350                                .get()
3351                                .validate_scann_generation(
3352                                    artifact.config(),
3353                                    artifact.generation(),
3354                                    artifact.artifact_id(),
3355                                )
3356                                .is_ok()
3357                        }
3358                        (None, None) => true,
3359                        _ => false,
3360                    }
3361                }
3362                // Only ivf_tq dense fields are trainable; other dense
3363                // index types never reach a vector-generation rewrite.
3364                crate::dsl::FieldType::DenseVector => false,
3365                crate::dsl::FieldType::BinaryDenseVector
3366                    if entry
3367                        .binary_dense_vector_config
3368                        .as_ref()
3369                        .is_some_and(|config| {
3370                            config.index_type == crate::dsl::BinaryIndexType::Scann
3371                        }) =>
3372                {
3373                    match (ann, trained.scann_artifacts.get(&field_id)) {
3374                        (Some(crate::segment::VectorIndex::ScannBinary(index)), Some(artifact)) => {
3375                            index
3376                                .get()
3377                                .validate_scann_generation(
3378                                    artifact.config(),
3379                                    artifact.generation(),
3380                                    artifact.artifact_id(),
3381                                )
3382                                .is_ok()
3383                        }
3384                        (None, None) => true,
3385                        _ => false,
3386                    }
3387                }
3388                crate::dsl::FieldType::BinaryDenseVector => matches!(
3389                    (ann, trained.binary_quantizers.get(&field_id)),
3390                    (Some(crate::segment::VectorIndex::BinaryIvf(_)), Some(_)) | (None, None)
3391                ),
3392                _ => false,
3393            };
3394            if !current {
3395                return Ok(true);
3396            }
3397        }
3398        Ok(false)
3399    }
3400
3401    async fn acquire_maintenance_capacity(
3402        &self,
3403    ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
3404        let global = tokio::select! {
3405            biased;
3406            () = self.active_operations.wait_for_shutdown() => {
3407                return Err(Error::IndexClosed);
3408            }
3409            permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
3410                permit.map_err(|_| Error::Internal(
3411                    "global background merge scheduler is closed".into()
3412                ))?
3413            }
3414        };
3415        let local = tokio::select! {
3416            biased;
3417            () = self.active_operations.wait_for_shutdown() => {
3418                return Err(Error::IndexClosed);
3419            }
3420            permit = Arc::clone(&self.merge_permits).acquire_owned() => {
3421                permit.map_err(|_| Error::Internal(
3422                    "background merge scheduler is closed".into()
3423                ))?
3424            }
3425        };
3426        Ok((global, local))
3427    }
3428
3429    async fn build_vector_replacement(
3430        self: &Arc<Self>,
3431        schemas: (&Arc<crate::dsl::Schema>, &Arc<crate::dsl::Schema>),
3432        segment_id: &str,
3433        source_id: SegmentId,
3434        output_id: SegmentId,
3435        trained: &TrainedVectorStructures,
3436        failure_context: &'static str,
3437    ) -> Result<(String, u32, OutputCleanupGuard)> {
3438        let mut cleanup = self.output_cleanup_guard(output_id);
3439        match crate::segment::reorder::rewrite_vector_segment(
3440            self.directory.as_ref(),
3441            schemas,
3442            source_id,
3443            output_id,
3444            self.term_cache_blocks,
3445            trained,
3446            Some(self.background_cpu_pool()),
3447        )
3448        .await
3449        {
3450            Ok((new_id, doc_count)) => {
3451                self.validate_completed_segment(&new_id, doc_count).await?;
3452                Ok((new_id, doc_count, cleanup))
3453            }
3454            Err(error) => {
3455                self.delete_output_if_unregistered(output_id, failure_context)
3456                    .await;
3457                cleanup.disarm();
3458                if is_deterministic_source_error(&error) {
3459                    self.quarantine_segment(segment_id, &error);
3460                }
3461                Err(error)
3462            }
3463        }
3464    }
3465
3466    /// Build every required replacement segment without exposing any of them.
3467    /// The returned guards keep both source and output generations alive until
3468    /// [`Self::publish_vector_generation`] commits the complete set.
3469    pub(crate) async fn stage_vector_generation(
3470        self: &Arc<Self>,
3471        _artifact_update: &VectorArtifactUpdateGuard,
3472        segment_ids: &[String],
3473        field_ids: &[u32],
3474        trained: Arc<TrainedVectorStructures>,
3475        rewrite_existing: bool,
3476    ) -> Result<Vec<StagedVectorSegment>> {
3477        let schema = self.published_generation().schema.clone();
3478        self.stage_vector_generation_with_schema(
3479            _artifact_update,
3480            segment_ids,
3481            field_ids,
3482            trained,
3483            rewrite_existing,
3484            schema,
3485        )
3486        .await
3487    }
3488
3489    pub(crate) async fn stage_vector_generation_with_schema(
3490        self: &Arc<Self>,
3491        _artifact_update: &VectorArtifactUpdateGuard,
3492        segment_ids: &[String],
3493        field_ids: &[u32],
3494        trained: Arc<TrainedVectorStructures>,
3495        rewrite_existing: bool,
3496        schema: Arc<crate::dsl::Schema>,
3497    ) -> Result<Vec<StagedVectorSegment>> {
3498        if !self.vector_artifact_update.load(Ordering::Acquire) {
3499            return Err(Error::Internal(
3500                "cannot stage a vector generation without an exclusive update lease".into(),
3501            ));
3502        }
3503
3504        let source_schema = self.published_generation().schema.clone();
3505        let mut staged = Vec::new();
3506        for segment_id in segment_ids {
3507            if self.quarantined_segments.lock().contains(segment_id) {
3508                return Err(Error::Corruption(format!(
3509                    "segment {segment_id} is quarantined after a deterministic source failure"
3510                )));
3511            }
3512            let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3513                Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3514            })?;
3515
3516            // A single rewrite may hold several gigabytes while assigning all
3517            // vectors. Reuse the ordinary local and process-wide merge bounds.
3518            let _capacity = self.acquire_maintenance_capacity().await?;
3519
3520            let output_id = SegmentId::new();
3521            let output_hex = output_id.to_hex();
3522            let operation = {
3523                let st = self.state.lock().await;
3524                if !st.metadata.has_segment(segment_id) {
3525                    return Err(Error::Corruption(format!(
3526                        "vector generation source {segment_id} disappeared while lifecycle work was paused"
3527                    )));
3528                }
3529                self.active_operations
3530                    .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
3531            }
3532            .ok_or_else(|| {
3533                if self.active_operations.is_accepting() {
3534                    Error::Internal(format!(
3535                        "vector generation could not claim stable source {segment_id}"
3536                    ))
3537                } else {
3538                    Error::IndexClosed
3539                }
3540            })?;
3541
3542            let reader = SegmentReader::open_with_term_cache_budget(
3543                self.directory.as_ref(),
3544                source_id,
3545                Arc::clone(&source_schema),
3546                self.term_cache_blocks,
3547                self.term_cache_budget_bytes,
3548            )
3549            .await?;
3550            if !self.segment_needs_vector_rewrite(
3551                schema.as_ref(),
3552                &reader,
3553                field_ids,
3554                trained.as_ref(),
3555                rewrite_existing,
3556            )? {
3557                continue;
3558            }
3559            drop(reader);
3560
3561            let (new_id, doc_count, cleanup) = self
3562                .build_vector_replacement(
3563                    (&source_schema, &schema),
3564                    segment_id,
3565                    source_id,
3566                    output_id,
3567                    trained.as_ref(),
3568                    "vector generation staging failure",
3569                )
3570                .await?;
3571            debug_assert_eq!(new_id, output_hex);
3572            let output_reader = SegmentReader::open_with_term_cache_budget(
3573                self.directory.as_ref(),
3574                output_id,
3575                Arc::clone(&schema),
3576                self.term_cache_blocks,
3577                self.term_cache_budget_bytes,
3578            )
3579            .await?;
3580            if self.segment_needs_vector_rewrite(
3581                schema.as_ref(),
3582                &output_reader,
3583                field_ids,
3584                trained.as_ref(),
3585                false,
3586            )? {
3587                return Err(Error::Corruption(format!(
3588                    "staged vector segment {new_id} does not match its candidate codebook generation"
3589                )));
3590            }
3591
3592            staged.push(StagedVectorSegment {
3593                source_id: segment_id.clone(),
3594                output_id,
3595                doc_count,
3596                _operation: operation,
3597                cleanup,
3598            });
3599        }
3600        Ok(staged)
3601    }
3602
3603    async fn rewrite_vector_segment_once(
3604        self: &Arc<Self>,
3605        segment_id: &str,
3606        field_ids: &[u32],
3607    ) -> Result<VectorSegmentRewriteOutcome> {
3608        if self.quarantined_segments.lock().contains(segment_id) {
3609            return Err(Error::Corruption(format!(
3610                "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3611            )));
3612        }
3613        let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3614            Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3615        })?;
3616
3617        // Match ordinary merge lock ordering: capacity before lifecycle
3618        // ownership. A vector rewrite can hold several gigabytes while it
3619        // assigns vectors, so it participates in both local and process-wide
3620        // merge limits.
3621        let _capacity = self.acquire_maintenance_capacity().await?;
3622
3623        let output_id = SegmentId::new();
3624        let output_hex = output_id.to_hex();
3625        let all_ids = vec![segment_id.to_owned(), output_hex];
3626        let operation = {
3627            let st = self.state.lock().await;
3628            if !st.metadata.has_segment(segment_id) {
3629                return Ok(VectorSegmentRewriteOutcome::SourceGone);
3630            }
3631            self.active_operations.try_register(all_ids)
3632        };
3633        let _operation = match operation {
3634            Some(operation) => operation,
3635            None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3636            None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3637        };
3638
3639        let Some(trained) = self.trained_for_segment_build() else {
3640            return Ok(VectorSegmentRewriteOutcome::Deferred);
3641        };
3642        let schema = self.published_generation().schema.clone();
3643
3644        let reader = SegmentReader::open_with_term_cache_budget(
3645            self.directory.as_ref(),
3646            source_id,
3647            Arc::clone(&schema),
3648            self.term_cache_blocks,
3649            self.term_cache_budget_bytes,
3650        )
3651        .await?;
3652        if !self.segment_needs_vector_rewrite(
3653            schema.as_ref(),
3654            &reader,
3655            field_ids,
3656            trained.as_ref(),
3657            false,
3658        )? {
3659            return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3660        }
3661        drop(reader);
3662
3663        let (new_id, doc_count, mut output_cleanup) = self
3664            .build_vector_replacement(
3665                (&schema, &schema),
3666                segment_id,
3667                source_id,
3668                output_id,
3669                trained.as_ref(),
3670                "vector rewrite failure",
3671            )
3672            .await?;
3673
3674        if let Err(error) = self
3675            .replace_segments(
3676                &[segment_id.to_owned()],
3677                new_id,
3678                doc_count,
3679                ReplacementLayout::PreserveSingleSource,
3680                None,
3681            )
3682            .await
3683        {
3684            self.delete_output_if_unregistered(output_id, "vector replacement failure")
3685                .await;
3686            output_cleanup.disarm();
3687            return Err(error);
3688        }
3689        output_cleanup.disarm();
3690        Ok(VectorSegmentRewriteOutcome::Rewritten)
3691    }
3692
3693    /// Finalize every committed flat vector segment against the published
3694    /// global ANN generation.
3695    /// Unlike force-merge this handles one segment and segments already at the
3696    /// merge policy's maximum size.
3697    pub(crate) async fn rewrite_vector_segments(
3698        self: &Arc<Self>,
3699        field_ids: &[u32],
3700    ) -> Result<usize> {
3701        if field_ids.is_empty() {
3702            return Ok(0);
3703        }
3704        let mut rewritten = 0usize;
3705        loop {
3706            let segment_ids = self.get_segment_ids().await;
3707            let mut conflicted = false;
3708            let mut changed = false;
3709            for segment_id in segment_ids {
3710                match self
3711                    .rewrite_vector_segment_once(&segment_id, field_ids)
3712                    .await?
3713                {
3714                    VectorSegmentRewriteOutcome::Rewritten => {
3715                        rewritten += 1;
3716                        changed = true;
3717                    }
3718                    VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3719                    VectorSegmentRewriteOutcome::Deferred => {
3720                        return Err(Error::Internal(
3721                            "ANN finalization lost the published trained generation".into(),
3722                        ));
3723                    }
3724                    VectorSegmentRewriteOutcome::AlreadyCurrent
3725                    | VectorSegmentRewriteOutcome::SourceGone => {}
3726                }
3727            }
3728            if !conflicted && !changed {
3729                log::info!(
3730                    "[dense_vector_rewrite] index={} ANN finalization complete ({} segment(s) rewritten)",
3731                    self.schema.index_label(),
3732                    rewritten,
3733                );
3734                return Ok(rewritten);
3735            }
3736            tokio::select! {
3737                biased;
3738                () = self.active_operations.wait_for_shutdown() => {
3739                    return Err(Error::IndexClosed);
3740                }
3741                () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3742            }
3743        }
3744    }
3745
3746    /// A producer that started in the force-flat phase can commit after the
3747    /// main finalization snapshot. Upgrade exactly those new segments in a
3748    /// tracked background task; ordinary producers already using the current
3749    /// generation are detected and skipped without rewriting.
3750    pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3751        if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3752            return;
3753        }
3754        let manager = Arc::clone(self);
3755        let future = async move {
3756            let field_ids = manager
3757                .read_metadata(|metadata| {
3758                    metadata
3759                        .vector_fields
3760                        .keys()
3761                        .filter(|field_id| metadata.is_field_built(**field_id))
3762                        .copied()
3763                        .collect::<Vec<_>>()
3764                })
3765                .await;
3766            for segment_id in segment_ids {
3767                loop {
3768                    match manager
3769                        .rewrite_vector_segment_once(&segment_id, &field_ids)
3770                        .await
3771                    {
3772                        Ok(VectorSegmentRewriteOutcome::Conflict) => {
3773                            tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3774                        }
3775                        Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3776                        Ok(_) => break,
3777                        Err(error) => {
3778                            log::error!(
3779                                "[dense_vector_rewrite] index={} failed to upgrade newly committed segment {}: {}",
3780                                manager.schema.index_label(),
3781                                segment_id,
3782                                error,
3783                            );
3784                            break;
3785                        }
3786                    }
3787                }
3788            }
3789        };
3790        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3791            log::warn!(
3792                "[dense_vector_rewrite] index={} runtime unavailable; newly committed flat segment upgrade deferred",
3793                self.schema.index_label()
3794            );
3795            return;
3796        };
3797        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3798            log::warn!(
3799                "[dense_vector_rewrite] index={} runtime rejected newly committed flat segment upgrade",
3800                self.schema.index_label()
3801            );
3802        }
3803    }
3804
3805    /// Apply bounded maintenance to all segments: opted-in text/BMP BP, Seismic
3806    /// nomination consolidation, and binary ANN run coalescing.
3807    /// Fields without maintenance work are copied unchanged.
3808    ///
3809    /// Uses active-operation ownership to prevent concurrent work on the same segment.
3810    pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3811        self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3812            .await
3813    }
3814
3815    /// Reorder all segments while advancing long-lived snapshots after the
3816    /// background-merge drain and every durable replacement.
3817    pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3818        self: &Arc<Self>,
3819        mut refresh_snapshots: F,
3820    ) -> Result<()>
3821    where
3822        F: FnMut() -> Fut,
3823        Fut: std::future::Future<Output = Result<()>>,
3824    {
3825        self.wait_for_all_merges().await;
3826        refresh_snapshots().await?;
3827        let segment_ids = self.get_segment_ids().await;
3828
3829        if segment_ids.is_empty() {
3830            log::info!(
3831                "[reorder] index={} no segments to reorder",
3832                self.schema.index_label()
3833            );
3834            return Ok(());
3835        }
3836
3837        log::info!(
3838            "[reorder] index={} reordering {} segments",
3839            self.schema.index_label(),
3840            segment_ids.len()
3841        );
3842
3843        for seg_id in segment_ids {
3844            match self
3845                .reorder_single_segment(
3846                    &seg_id,
3847                    Some(self.background_cpu_pool()),
3848                    crate::segment::BpBudget::full(),
3849                )
3850                .await
3851            {
3852                Ok(true) => refresh_snapshots().await?,
3853                Ok(false) => log::warn!(
3854                    "[reorder] index={} segment {} skipped (in merge)",
3855                    self.schema.index_label(),
3856                    seg_id
3857                ),
3858                Err(e) => return Err(e),
3859            }
3860        }
3861
3862        // A segment skipped because another lifecycle owner held it may have
3863        // been replaced after the preceding callback.
3864        refresh_snapshots().await?;
3865        log::info!(
3866            "[reorder] index={} all segments reordered",
3867            self.schema.index_label()
3868        );
3869        Ok(())
3870    }
3871
3872    /// Get segment IDs that have not been reordered yet.
3873    ///
3874    /// Excludes segments currently involved in a merge or reorder operation
3875    /// to avoid wasted work (the optimizer would skip them anyway).
3876    pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3877        self.unreordered_segments()
3878            .await
3879            .into_iter()
3880            .map(|(id, _)| id)
3881            .collect()
3882    }
3883
3884    /// Initial maintenance candidates from persisted layout debt and opted-in
3885    /// text/BMP BP, with document counts for budget selection. No payload scan.
3886    pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3887        let quarantined = self.quarantined_segments.lock().clone();
3888        let paused = self.paused_reorder_segments();
3889        let st = self.state.lock().await;
3890        let active_ids = self.active_operations.snapshot();
3891        let has_bp_reorder = self.schema.has_reorder_fields();
3892        st.metadata
3893            .segment_metas
3894            .iter()
3895            .filter(|(id, info)| {
3896                info.bp_converged
3897                    && ((has_bp_reorder && !info.reordered)
3898                        || info.ann_fragmented
3899                        || (info.seismic_pending_terms > 0 && info.seismic_maintenance_passes == 0))
3900                    && !active_ids.contains(*id)
3901                    && !quarantined.contains(*id)
3902                    && !paused.contains(*id)
3903            })
3904            .map(|(id, info)| (id.clone(), info.num_docs))
3905            .collect()
3906    }
3907
3908    /// Segments with incomplete text/BMP BP or Seismic nomination consolidation.
3909    /// Follow-up work shares the optimizer's cooldown. Productive Seismic passes
3910    /// remain eligible until debt is zero; only consecutive stalled passes count
3911    /// toward its bound. BMP keeps its existing total-pass lineage bound.
3912    pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3913        self.unconverged_segments_below(u32::MAX)
3914            .await
3915            .into_iter()
3916            .map(|(id, docs, _)| (id, docs))
3917            .collect()
3918    }
3919
3920    /// Segments below the BMP lineage or Seismic no-progress bound. Includes
3921    /// the applicable persisted retry count for scheduler diagnostics.
3922    pub async fn unconverged_segments_below(
3923        &self,
3924        max_unconverged_passes: u32,
3925    ) -> Vec<(String, u32, u32)> {
3926        let quarantined = self.quarantined_segments.lock().clone();
3927        let paused = self.paused_reorder_segments();
3928        let st = self.state.lock().await;
3929        let active_ids = self.active_operations.snapshot();
3930        st.metadata
3931            .segment_metas
3932            .iter()
3933            .filter(|(id, info)| {
3934                ((!info.bp_converged && info.bp_unconverged_passes < max_unconverged_passes)
3935                    || (info.seismic_pending_terms > 0
3936                        && (info.seismic_maintenance_passes > 0 || !info.bp_converged)
3937                        && info.seismic_no_progress_passes < max_unconverged_passes))
3938                    && !active_ids.contains(*id)
3939                    && !quarantined.contains(*id)
3940                    && !paused.contains(*id)
3941            })
3942            .map(|(id, info)| {
3943                (
3944                    id.clone(),
3945                    info.num_docs,
3946                    if info.seismic_pending_terms > 0
3947                        && (info.seismic_maintenance_passes > 0 || !info.bp_converged)
3948                        && info.seismic_no_progress_passes < max_unconverged_passes
3949                    {
3950                        info.seismic_no_progress_passes
3951                    } else {
3952                        info.bp_unconverged_passes
3953                    },
3954                )
3955            })
3956            .collect()
3957    }
3958
3959    async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3960        let st = self.state.lock().await;
3961        let deepening = ids.iter().any(|id| {
3962            st.metadata
3963                .segment_metas
3964                .get(id)
3965                .is_some_and(|info| !info.bp_converged)
3966        });
3967        drop(st);
3968        if deepening {
3969            log::info!(
3970                "[reorder] index={} source BP lineage unconverged — forcing record-level BP (deepening pass)",
3971                self.schema.index_label(),
3972            );
3973            crate::segment::reorder::BpGranularity::Records
3974        } else {
3975            crate::segment::reorder::BpGranularity::Auto
3976        }
3977    }
3978
3979    /// Reorder a single segment via BP. Returns Ok(true) if reordered, Ok(false) if skipped.
3980    ///
3981    /// Non-blocking: operation ownership prevents conflicts with background merges.
3982    /// Copies unchanged files, reorders text/BMP, repairs Seismic debt and coalesces binary ANN runs.
3983    pub async fn reorder_single_segment(
3984        self: &Arc<Self>,
3985        seg_id: &str,
3986        rayon_pool: Option<Arc<rayon::ThreadPool>>,
3987        bp_budget: crate::segment::BpBudget,
3988    ) -> Result<bool> {
3989        self.reorder_single_segment_with_policy(seg_id, rayon_pool, bp_budget, None)
3990            .await
3991    }
3992
3993    /// Run automatically scheduled maintenance, reusing completed text/BMP
3994    /// layouts and respecting their follow-up bound independently of Seismic.
3995    pub async fn optimize_single_segment(
3996        self: &Arc<Self>,
3997        seg_id: &str,
3998        rayon_pool: Option<Arc<rayon::ThreadPool>>,
3999        bp_budget: crate::segment::BpBudget,
4000        max_bp_passes: u32,
4001    ) -> Result<bool> {
4002        self.reorder_single_segment_with_policy(seg_id, rayon_pool, bp_budget, Some(max_bp_passes))
4003            .await
4004    }
4005
4006    async fn reorder_single_segment_with_policy(
4007        self: &Arc<Self>,
4008        seg_id: &str,
4009        rayon_pool: Option<Arc<rayon::ThreadPool>>,
4010        bp_budget: crate::segment::BpBudget,
4011        automatic_bp_limit: Option<u32>,
4012    ) -> Result<bool> {
4013        let source_id = SegmentId::from_hex(seg_id)
4014            .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
4015        if self.quarantined_segments.lock().contains(seg_id) {
4016            return Err(Error::Corruption(format!(
4017                "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
4018                seg_id
4019            )));
4020        }
4021        if self.force_merge_active.load(Ordering::Acquire) > 0 {
4022            log::debug!(
4023                "[optimizer] index={} explicit force merge active, skipping reorder of {}",
4024                self.schema.index_label(),
4025                seg_id,
4026            );
4027            return Ok(false);
4028        }
4029
4030        // Whole-pass concurrency is independent from Rayon width. One pass
4031        // can already use every configured BP worker; this permit bounds the
4032        // much larger forward-index and rewrite working set across indexes,
4033        // optimizer tasks, and merge-time BP.
4034        let reorder_gate = Arc::clone(&self.reorder_permits);
4035        let _reorder_permit = tokio::select! {
4036            biased;
4037            () = self.active_operations.wait_for_shutdown() => {
4038                return Err(Error::IndexClosed);
4039            }
4040            permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
4041                permit.map_err(|_| {
4042                    Error::Internal("background reorder scheduler is closed".into())
4043                })?
4044            }
4045        };
4046
4047        let output_id = SegmentId::new();
4048        let output_hex = output_id.to_hex();
4049
4050        // Register while holding `state`, matching orphan cleanup's deletion
4051        // barrier. Candidates are scanned ahead of time and can go stale: a
4052        // merge may have consumed this segment since. Its files may even still
4053        // be on disk (deferred deletion under a searcher snapshot) — reordering
4054        // them would re-insert a duplicate copy of docs the merge output holds.
4055        let all_ids = vec![seg_id.to_string(), output_hex];
4056        let (_guard, source_docs, generation, source_deletions, _visibility_owner, run_bp) = {
4057            let st = self.state.lock().await;
4058            // Force merge raises this barrier under the same state lock, so
4059            // this second check closes the race with the cheap early check
4060            // above and prevents optimizer starvation between final groups.
4061            if self.force_merge_active.load(Ordering::Acquire) > 0 {
4062                log::debug!(
4063                    "[optimizer] index={} explicit force merge active, skipping reorder of {}",
4064                    self.schema.index_label(),
4065                    seg_id,
4066                );
4067                return Ok(false);
4068            }
4069            let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
4070                log::info!(
4071                    "[optimizer] index={} segment {} no longer in metadata (merged away), skipping reorder",
4072                    self.schema.index_label(),
4073                    seg_id
4074                );
4075                self.clear_reorder_retry(seg_id);
4076                return Ok(false);
4077            };
4078
4079            let run_bp = automatic_bp_limit.is_none_or(|limit| {
4080                self.schema.has_reorder_fields()
4081                    && ((!source_meta.reordered && source_meta.bp_converged)
4082                        || (!source_meta.bp_converged && source_meta.bp_unconverged_passes < limit))
4083            });
4084            let generation = self.published_generation();
4085            match self.active_operations.try_register(all_ids) {
4086                Some(guard) => {
4087                    // Deletions can publish while this address-preserving rewrite
4088                    // runs. Pin this exact mask before releasing the metadata lock.
4089                    let deletion_ids: Vec<_> = source_meta
4090                        .deletions
4091                        .iter()
4092                        .map(|deletion| deletion.id.clone())
4093                        .collect();
4094                    let acquired = self.tracker.acquire(&deletion_ids);
4095                    if acquired.len() != deletion_ids.len() {
4096                        return Err(Error::Corruption(
4097                            "maintenance source visibility is already retired".into(),
4098                        ));
4099                    }
4100                    let visibility_owner = SegmentSnapshot::with_delete_fn(
4101                        Arc::clone(&self.tracker),
4102                        acquired,
4103                        Arc::clone(&self.delete_fn),
4104                    );
4105                    (
4106                        guard,
4107                        source_meta.num_docs,
4108                        generation,
4109                        source_meta.deletions.clone(),
4110                        visibility_owner,
4111                        run_bp,
4112                    )
4113                }
4114                None if !self.active_operations.is_accepting() => {
4115                    return Err(Error::IndexClosed);
4116                }
4117                None => {
4118                    log::debug!(
4119                        "[optimizer] index={} segment {} in active merge, skipping",
4120                        self.schema.index_label(),
4121                        seg_id
4122                    );
4123                    return Ok(false);
4124                }
4125            }
4126        };
4127
4128        // Fail before allocating a forward index or creating output files.
4129        // Missing mandatory files are deterministic and should remove this
4130        // segment from future optimizer scans, not consume the same CPU every
4131        // interval. Other I/O failures remain retryable.
4132        if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
4133            if is_deterministic_source_error(&error) {
4134                self.quarantine_segment(seg_id, &error);
4135            } else if !matches!(&error, Error::IndexClosed) {
4136                self.pause_reorder_retries(seg_id, &error).await;
4137            }
4138            return Err(error);
4139        }
4140
4141        let alive = match source_deletions {
4142            Some(deletions) => match deletions.load(self.directory.as_ref(), source_docs).await {
4143                Ok(alive) => Some(alive),
4144                Err(error) => {
4145                    if is_deterministic_source_error(&error) {
4146                        self.quarantine_segment(seg_id, &error);
4147                    } else if !matches!(&error, Error::IndexClosed) {
4148                        self.pause_reorder_retries(seg_id, &error).await;
4149                    }
4150                    return Err(error);
4151                }
4152            },
4153            None => None,
4154        };
4155        let mut output_cleanup = self.output_cleanup_guard(output_id);
4156
4157        let granularity = if run_bp {
4158            self.merge_granularity(&[seg_id.to_owned()]).await
4159        } else {
4160            crate::segment::reorder::BpGranularity::Auto
4161        };
4162        let reorder_result = crate::segment::reorder::reorder_segment(
4163            self.directory.as_ref(),
4164            &generation.schema,
4165            source_id,
4166            output_id,
4167            self.term_cache_blocks,
4168            self.term_cache_budget_bytes,
4169            self.bp_memory_budget_bytes,
4170            bp_budget,
4171            run_bp,
4172            granularity,
4173            self.optimization,
4174            self.posting_codec,
4175            generation.trained_vectors.as_deref(),
4176            self.term_dict_block_size,
4177            rayon_pool,
4178            Some(self.active_operations.cancellation_flag()),
4179            alive,
4180        )
4181        .await;
4182        let (new_id, total_docs, bp_converged) = match reorder_result {
4183            Ok(v) => v,
4184            Err(e) => {
4185                // A failed pass may have copied tens of GB before dying;
4186                // delete the uncommitted output before propagating.
4187                self.delete_output_if_unregistered(output_id, "reorder failure")
4188                    .await;
4189                output_cleanup.disarm();
4190                if is_deterministic_source_error(&e) {
4191                    self.quarantine_segment(seg_id, &e);
4192                } else if !matches!(&e, Error::IndexClosed) {
4193                    self.pause_reorder_retries(seg_id, &e).await;
4194                }
4195                return Err(e);
4196            }
4197        };
4198
4199        // A pass with a depth floor above block granularity has, by
4200        // definition, not converged to block-level order — record it as
4201        // unconverged so the optimizer's deepening ladder revisits it with a
4202        // full-depth (warm-started) pass. Depth caps are only used by the
4203        // optimizer's first pass on large segments.
4204        let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
4205        if let Err(e) = self
4206            .replace_segments(
4207                &[seg_id.to_string()],
4208                new_id,
4209                total_docs,
4210                if run_bp {
4211                    ReplacementLayout::BpReordered {
4212                        converged: ladder_converged,
4213                    }
4214                } else {
4215                    ReplacementLayout::MaintenanceOnly
4216                },
4217                None,
4218            )
4219            .await
4220        {
4221            self.delete_output_if_unregistered(output_id, "replacement failure")
4222                .await;
4223            output_cleanup.disarm();
4224            if !matches!(&e, Error::IndexClosed) {
4225                self.pause_reorder_retries(seg_id, &e).await;
4226            }
4227            return Err(e);
4228        }
4229        output_cleanup.disarm();
4230        self.clear_reorder_retry(seg_id);
4231
4232        Ok(true)
4233    }
4234
4235    /// Clean up orphan segment files not registered in metadata.
4236    ///
4237    /// Reads metadata, active-operation ownership, and snapshot-deferred
4238    /// deletions to determine which segments are legitimate. Filesystem
4239    /// deletion is asynchronous; in-flight outputs and retired sources still
4240    /// held by readers are both protected.
4241    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
4242        let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
4243
4244        if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
4245            for entry in entries {
4246                let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
4247                    continue;
4248                };
4249                let Some(rest) = filename.strip_prefix("seg_") else {
4250                    continue;
4251                };
4252                let Some(hex_id) = rest.get(..32) else {
4253                    continue;
4254                };
4255                if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
4256                    continue;
4257                }
4258                orphan_files
4259                    .entry(hex_id.to_ascii_lowercase())
4260                    .or_default()
4261                    .push(entry);
4262            }
4263        }
4264
4265        let mut deleted = 0;
4266        for (hex_id, paths) in &orphan_files {
4267            // Revalidate and atomically claim deletion under the same
4268            // state -> active_operations -> tracker order used by publishers.
4269            // The claim lets us release `state` before filesystem I/O: deleting
4270            // a multi-GB orphan must not freeze commits and snapshot acquisition.
4271            let deletion_guard = {
4272                let st = self.state.lock().await;
4273                if st.metadata.owns_id(hex_id) {
4274                    continue;
4275                }
4276                let Some(guard) = self
4277                    .active_operations
4278                    .try_register(vec![hex_id.to_string()])
4279                else {
4280                    continue;
4281                };
4282                if self.tracker.is_deletion_protected(hex_id) {
4283                    drop(guard);
4284                    continue;
4285                }
4286                guard
4287            };
4288
4289            // Delete what was actually discovered, not only the currently
4290            // known SegmentFiles extensions. This also removes partial files
4291            // left by older formats instead of reporting the same orphan on
4292            // every startup forever.
4293            let results =
4294                futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
4295                    .await;
4296            let removed = results.into_iter().all(|result| match result {
4297                Ok(()) => true,
4298                Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
4299                Err(error) => {
4300                    log::warn!(
4301                        "[segment_cleanup] index={} failed sweeping orphan segment {}: {}",
4302                        self.schema.index_label(),
4303                        hex_id,
4304                        error,
4305                    );
4306                    false
4307                }
4308            });
4309            // Releasing this claim is the deletion barrier. No producer can
4310            // adopt the ID while its files are being removed.
4311            drop(deletion_guard);
4312            if removed {
4313                deleted += 1;
4314                log::info!(
4315                    "[segment_cleanup] index={} swept orphan segment {}",
4316                    self.schema.index_label(),
4317                    hex_id
4318                );
4319            }
4320        }
4321
4322        Ok(deleted)
4323    }
4324}
4325
4326#[cfg(test)]
4327mod tests {
4328    use super::*;
4329    use std::sync::atomic::{AtomicBool, Ordering};
4330
4331    fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
4332        let schema = crate::dsl::SchemaBuilder::default().build();
4333        let metadata = IndexMetadata::new(schema.clone());
4334        Arc::new(SegmentManager::new(
4335            Arc::new(crate::directories::RamDirectory::new()),
4336            Arc::new(schema),
4337            metadata,
4338            Box::new(crate::merge::NoMergePolicy),
4339            0,
4340            1,
4341            Arc::new(Semaphore::new(1)),
4342            None,
4343            1024,
4344            Arc::new(ReorderConcurrencyGate::new(1)),
4345            None,
4346        ))
4347    }
4348
4349    #[tokio::test]
4350    async fn queued_compaction_does_not_deadlock_a_vector_rewrite_holding_global_capacity() {
4351        use std::time::Duration;
4352        let manager = lifecycle_test_manager();
4353        let global = Arc::clone(&manager.global_merge_permits)
4354            .acquire_owned()
4355            .await
4356            .unwrap();
4357        let id = SegmentId::new().to_hex();
4358        let mut compaction = Box::pin(manager.compact_segment(&id, 1024 * 1024));
4359        assert!(
4360            tokio::time::timeout(Duration::from_millis(10), &mut compaction)
4361                .await
4362                .is_err()
4363        );
4364        // A rewrite already holding global capacity must still be able to
4365        // acquire local capacity and finish; compaction waits behind it.
4366        let local = Arc::clone(&manager.merge_permits)
4367            .try_acquire_owned()
4368            .expect("waiting compaction must not reserve local capacity first");
4369        drop(local);
4370        drop(global);
4371        assert!(
4372            !tokio::time::timeout(Duration::from_secs(1), compaction)
4373                .await
4374                .unwrap()
4375                .unwrap()
4376        );
4377    }
4378
4379    #[test]
4380    fn force_merge_planner_pairs_large_and_small_segments() {
4381        let groups = plan_force_merge_groups(
4382            vec![
4383                ("a".into(), 6),
4384                ("b".into(), 6),
4385                ("c".into(), 4),
4386                ("d".into(), 4),
4387            ],
4388            10,
4389        );
4390
4391        assert_eq!(groups.len(), 2);
4392        assert!(groups.iter().all(|group| group.total_docs == 10));
4393        assert!(groups.iter().all(|group| group.segments.len() == 2));
4394    }
4395
4396    #[test]
4397    fn force_merge_planner_leaves_oversized_segments_alone() {
4398        let groups = plan_force_merge_groups(
4399            vec![
4400                ("oversized".into(), 11),
4401                ("small-a".into(), 5),
4402                ("small-b".into(), 5),
4403            ],
4404            10,
4405        );
4406
4407        assert_eq!(groups.len(), 2);
4408        assert_eq!(groups[0].total_docs, 10);
4409        assert_eq!(groups[0].segments.len(), 2);
4410        assert_eq!(groups[1].total_docs, 11);
4411        assert_eq!(groups[1].segments.len(), 1);
4412    }
4413
4414    #[test]
4415    fn force_merge_planner_never_exceeds_segment_format_limit() {
4416        let groups = plan_force_merge_groups(
4417            vec![
4418                ("large-a".into(), 3_000_000_000),
4419                ("large-b".into(), 2_000_000_000),
4420            ],
4421            u64::from(u32::MAX),
4422        );
4423        assert_eq!(groups.len(), 2);
4424        assert!(
4425            groups
4426                .iter()
4427                .all(|group| group.total_docs <= u64::from(u32::MAX))
4428        );
4429    }
4430
4431    #[test]
4432    fn force_merge_hierarchy_has_one_final_bp_pass() {
4433        assert_eq!(force_merge_output_count(1), 0);
4434        assert_eq!(force_merge_output_count(2), 1);
4435        assert_eq!(force_merge_output_count(64), 1);
4436        assert_eq!(force_merge_output_count(65), 2);
4437        assert_eq!(force_merge_output_count(127), 2);
4438        assert_eq!(force_merge_output_count(128), 3);
4439        assert_eq!(force_merge_output_count(1_000), 16);
4440    }
4441
4442    fn expand_force_merge_node(
4443        hierarchy: &ForceMergeHierarchy,
4444        source_count: usize,
4445        node: usize,
4446        sources: &mut Vec<usize>,
4447    ) {
4448        if node < source_count {
4449            sources.push(node);
4450            return;
4451        }
4452
4453        let step_index = node - source_count;
4454        let step = hierarchy
4455            .steps
4456            .get(step_index)
4457            .expect("merge input must refer to an existing source or output");
4458        for &input in &step.inputs {
4459            assert!(
4460                input < node,
4461                "merge step {step_index} refers to a future output node {input}"
4462            );
4463            expand_force_merge_node(hierarchy, source_count, input, sources);
4464        }
4465    }
4466
4467    #[test]
4468    fn force_merge_hierarchy_has_minimal_valid_arity() {
4469        let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
4470
4471        for source_count in source_counts {
4472            let hierarchy = plan_force_merge_hierarchy(source_count);
4473            let output_count = hierarchy.steps.len();
4474
4475            assert!(
4476                hierarchy
4477                    .steps
4478                    .iter()
4479                    .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
4480                "invalid merge arity for {source_count} sources"
4481            );
4482            assert!(
4483                source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
4484                "{output_count} outputs cannot reduce {source_count} sources"
4485            );
4486            assert!(
4487                output_count == 1
4488                    || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
4489                "{output_count} outputs are not minimal for {source_count} sources"
4490            );
4491            assert_eq!(output_count, force_merge_output_count(source_count));
4492        }
4493    }
4494
4495    #[test]
4496    fn force_merge_hierarchy_preserves_exact_source_order() {
4497        for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
4498            let hierarchy = plan_force_merge_hierarchy(source_count);
4499            let mut sources = Vec::with_capacity(source_count);
4500            expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
4501            assert_eq!(
4502                sources,
4503                (0..source_count).collect::<Vec<_>>(),
4504                "source order changed for {source_count} sources"
4505            );
4506        }
4507    }
4508
4509    fn force_merge_rewrite_cost(source_count: usize) -> usize {
4510        let hierarchy = plan_force_merge_hierarchy(source_count);
4511        let mut node_weights = vec![1usize; source_count];
4512        let mut rewrite_cost = 0usize;
4513
4514        for (step_index, step) in hierarchy.steps.iter().enumerate() {
4515            let output = source_count + step_index;
4516            let output_weight = step
4517                .inputs
4518                .iter()
4519                .map(|&input| {
4520                    assert!(
4521                        input < output,
4522                        "merge step {step_index} refers to future output {input}"
4523                    );
4524                    node_weights[input]
4525                })
4526                .sum::<usize>();
4527            rewrite_cost += output_weight;
4528            node_weights.push(output_weight);
4529        }
4530
4531        assert_eq!(node_weights[hierarchy.root], source_count);
4532        rewrite_cost
4533    }
4534
4535    #[test]
4536    fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
4537        assert_eq!(force_merge_rewrite_cost(65), 67);
4538        assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
4539    }
4540
4541    #[test]
4542    fn block_copy_carries_bp_debt_without_spending_an_attempt() {
4543        assert_eq!(
4544            replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
4545            (false, false, 3),
4546        );
4547        assert_eq!(
4548            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
4549            (true, false, 4),
4550        );
4551        assert_eq!(
4552            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
4553            (true, true, 0),
4554        );
4555    }
4556
4557    #[tokio::test]
4558    async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
4559        let mut schema_builder = crate::dsl::SchemaBuilder::default();
4560        let field = schema_builder.add_text_field("text", true, true);
4561        let schema = schema_builder.build();
4562        let directory = crate::directories::RamDirectory::new();
4563        let config = crate::index::IndexConfig {
4564            num_indexing_threads: 1,
4565            merge_policy: Box::new(crate::merge::NoMergePolicy),
4566            ..Default::default()
4567        };
4568        let mut writer = crate::index::IndexWriter::create(directory, schema, config)
4569            .await
4570            .unwrap();
4571        for value in ["one", "two", "three"] {
4572            let mut document = crate::dsl::Document::new();
4573            document.add_text(field, value);
4574            writer.add_document(document).unwrap();
4575            writer.commit().await.unwrap();
4576        }
4577
4578        let manager = Arc::clone(writer.segment_manager());
4579        let held_id = manager.get_segment_ids().await.pop().unwrap();
4580        let mut held = Some(
4581            manager
4582                .active_operations
4583                .try_register(vec![held_id])
4584                .unwrap(),
4585        );
4586        let batches = Arc::new(AtomicUsize::new(0));
4587        let batch_count = Arc::clone(&batches);
4588        writer
4589            .force_merge_with_snapshot_refresh(move || {
4590                let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
4591                // Refresh #1 follows the startup drain. Keep the source held
4592                // until refresh #2, after the two free sources were replaced.
4593                if refresh == 2 {
4594                    drop(held.take());
4595                }
4596                std::future::ready(Ok(()))
4597            })
4598            .await
4599            .unwrap();
4600
4601        assert_eq!(manager.get_segment_ids().await.len(), 1);
4602        assert_eq!(
4603            batches.load(Ordering::Relaxed),
4604            4,
4605            "initial/final refreshes plus two replacements are required after the held source releases"
4606        );
4607    }
4608
4609    #[tokio::test]
4610    async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
4611        let mut schema_builder = crate::dsl::SchemaBuilder::default();
4612        schema_builder.set_reorder_on_merge(true);
4613        let schema = schema_builder.build();
4614        let mut metadata = IndexMetadata::new(schema.clone());
4615        metadata.add_segment("00000000000000000000000000000001".into(), 1);
4616        metadata.add_segment("00000000000000000000000000000002".into(), 1);
4617
4618        let global_merge_permits = Arc::new(Semaphore::new(1));
4619        let manager = Arc::new(SegmentManager::new(
4620            Arc::new(crate::directories::RamDirectory::new()),
4621            Arc::new(schema),
4622            metadata,
4623            Box::new(crate::merge::NoMergePolicy),
4624            0,
4625            1,
4626            Arc::clone(&global_merge_permits),
4627            None,
4628            1024,
4629            Arc::new(ReorderConcurrencyGate::new(1)),
4630            None,
4631        ));
4632
4633        // Vector-generation staging rejects ordinary lifecycle claims while
4634        // acquiring the shared global merge permit separately for each source.
4635        // Force merge must wait before capacity admission; retaining the only
4636        // global slot here would deadlock both operations between sources.
4637        manager.active_operations.pause_non_indexing();
4638        let force_merge = {
4639            let manager = Arc::clone(&manager);
4640            tokio::spawn(async move { manager.force_merge().await })
4641        };
4642        tokio::time::timeout(std::time::Duration::from_secs(1), async {
4643            while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
4644                tokio::task::yield_now().await;
4645            }
4646        })
4647        .await
4648        .expect("force merge never reached the paused group claim");
4649
4650        assert_eq!(
4651            global_merge_permits.available_permits(),
4652            1,
4653            "force merge retained global capacity while vector staging blocked group ownership"
4654        );
4655
4656        force_merge.abort();
4657        let _ = force_merge.await;
4658        manager.active_operations.resume_non_indexing();
4659    }
4660
4661    #[test]
4662    fn output_cleanup_guard_runs_during_panic_unwind() {
4663        let cleaned = Arc::new(AtomicBool::new(false));
4664        let cleaned_in_callback = Arc::clone(&cleaned);
4665        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4666            cleaned_in_callback.store(true, Ordering::SeqCst);
4667        });
4668
4669        let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
4670            let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4671            panic!("simulated reorder panic");
4672        }));
4673
4674        assert!(result.is_err());
4675        assert!(
4676            cleaned.load(Ordering::SeqCst),
4677            "partial output cleanup must run during unwind"
4678        );
4679    }
4680
4681    #[test]
4682    fn output_cleanup_guard_disarms_after_commit() {
4683        let cleaned = Arc::new(AtomicBool::new(false));
4684        let cleaned_in_callback = Arc::clone(&cleaned);
4685        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
4686            cleaned_in_callback.store(true, Ordering::SeqCst);
4687        });
4688
4689        {
4690            let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
4691            guard.disarm();
4692        }
4693
4694        assert!(!cleaned.load(Ordering::SeqCst));
4695    }
4696
4697    #[test]
4698    fn test_active_operation_guard_releases_ownership() {
4699        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4700        {
4701            let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4702            let snap = active.snapshot();
4703            assert!(snap.contains("a"));
4704            assert!(snap.contains("b"));
4705        }
4706        assert!(active.snapshot().is_empty());
4707    }
4708
4709    #[test]
4710    fn test_non_overlapping_operations_can_run_concurrently() {
4711        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4712        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4713        let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
4714        let snap = active.snapshot();
4715        assert_eq!(snap.len(), 4);
4716
4717        drop(first);
4718        let snap = active.snapshot();
4719        assert_eq!(snap.len(), 2);
4720        assert!(snap.contains("c"));
4721        assert!(snap.contains("d"));
4722    }
4723
4724    #[test]
4725    fn test_overlapping_operation_is_rejected_until_release() {
4726        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4727        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
4728        assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
4729        drop(first);
4730        assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
4731    }
4732
4733    #[test]
4734    fn test_active_operation_snapshot() {
4735        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4736        let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
4737        let snap = active.snapshot();
4738        assert!(snap.contains("x"));
4739        assert!(snap.contains("y"));
4740        assert!(!snap.contains("z"));
4741    }
4742
4743    #[tokio::test]
4744    async fn operation_barrier_ignores_producers_started_after_snapshot() {
4745        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4746        let before_gate = active.try_register(vec!["old".into()]).unwrap();
4747        let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
4748        assert_eq!(parked_indexing, 0);
4749        let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
4750
4751        let waiter = {
4752            let active = Arc::clone(&active);
4753            tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
4754        };
4755        tokio::task::yield_now().await;
4756        assert!(!waiter.is_finished());
4757
4758        drop(before_gate);
4759        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4760            .await
4761            .expect("pre-gate operation barrier was starved by a post-gate producer")
4762            .unwrap();
4763        assert!(active.snapshot().contains("new-flat"));
4764        drop(after_gate);
4765    }
4766
4767    #[tokio::test]
4768    async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4769        let manager = lifecycle_test_manager();
4770        let current = manager.published_generation();
4771        manager
4772            .published_generation
4773            .store(Arc::new(PublishedIndexGeneration {
4774                publication_id: current.publication_id,
4775                schema: current.schema.clone(),
4776                trained_vectors: Some(Arc::new(TrainedVectorStructures {
4777                    centroids: rustc_hash::FxHashMap::default(),
4778                    binary_quantizers: rustc_hash::FxHashMap::default(),
4779                    ..Default::default()
4780                })),
4781            }));
4782
4783        let guard = manager.begin_vector_artifact_update().await.unwrap();
4784        assert!(
4785            manager.trained().is_some(),
4786            "search readers keep the last fully validated generation"
4787        );
4788        assert!(
4789            manager.trained_for_segment_build().is_none(),
4790            "new segment producers must stay flat during an artifact update"
4791        );
4792
4793        let detached_transaction_guard = guard.clone();
4794        drop(guard);
4795        assert!(
4796            manager.trained_for_segment_build().is_none(),
4797            "a detached lifecycle transaction must retain the producer gate after request cancellation"
4798        );
4799        drop(detached_transaction_guard);
4800        assert!(manager.trained_for_segment_build().is_some());
4801    }
4802
4803    #[tokio::test]
4804    async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4805        let manager = lifecycle_test_manager();
4806        let guard = manager.begin_vector_artifact_update().await.unwrap();
4807        assert!(
4808            manager
4809                .active_operations
4810                .try_register(vec!["merge".into()])
4811                .is_none(),
4812            "ordinary merge/reorder work must not change staged sources"
4813        );
4814        let indexing = manager
4815            .active_operations
4816            .try_register_indexing(vec!["fresh".into()])
4817            .expect("indexing remains available in flat mode");
4818        drop(indexing);
4819
4820        drop(guard);
4821        assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4822        assert!(
4823            manager
4824                .active_operations
4825                .try_register(vec!["merge".into()])
4826                .is_some()
4827        );
4828    }
4829
4830    #[tokio::test]
4831    async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4832        let active = Arc::new(ActiveSegmentOperations::new("test".into()));
4833        let guard = active.try_register(vec!["live".into()]).unwrap();
4834        let cancellation = active.cancellation_flag();
4835        active.stop_accepting();
4836        assert!(cancellation.load(Ordering::Acquire));
4837        assert!(active.try_register(vec!["new".into()]).is_none());
4838
4839        let waiter = {
4840            let active = Arc::clone(&active);
4841            tokio::spawn(async move { active.wait_until_idle().await })
4842        };
4843        tokio::task::yield_now().await;
4844        assert!(!waiter.is_finished());
4845        drop(guard);
4846        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4847            .await
4848            .expect("shutdown waiter missed the final guard notification")
4849            .unwrap();
4850    }
4851
4852    #[tokio::test]
4853    async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4854        let manager = lifecycle_test_manager();
4855        let started = Arc::new(Semaphore::new(0));
4856        let release = Arc::new(Semaphore::new(0));
4857        let completed = Arc::new(AtomicBool::new(false));
4858
4859        let request = {
4860            let manager = Arc::clone(&manager);
4861            let started = Arc::clone(&started);
4862            let release = Arc::clone(&release);
4863            let completed = Arc::clone(&completed);
4864            tokio::spawn(async move {
4865                manager
4866                    .run_lifecycle_transaction(async move {
4867                        started.add_permits(1);
4868                        let _permit = release.acquire().await.unwrap();
4869                        completed.store(true, Ordering::Release);
4870                        Ok(())
4871                    })
4872                    .await
4873            })
4874        };
4875
4876        let _started = started.acquire().await.unwrap();
4877        request.abort();
4878        assert!(request.await.unwrap_err().is_cancelled());
4879        release.add_permits(1);
4880
4881        manager.begin_shutdown();
4882        tokio::time::timeout(
4883            std::time::Duration::from_secs(1),
4884            manager.wait_for_shutdown(),
4885        )
4886        .await
4887        .expect("shutdown did not drain detached lifecycle transaction");
4888        assert!(completed.load(Ordering::Acquire));
4889    }
4890
4891    #[tokio::test]
4892    async fn productive_seismic_maintenance_remains_eligible_beyond_three_passes() {
4893        let manager = lifecycle_test_manager();
4894        {
4895            let mut state = manager.state.lock().await;
4896            for (id, pending, attempts, stalls) in [
4897                ("eligible", 7, 100, 0),
4898                ("finished", 0, 1, 0),
4899                ("limited", 9, 3, 3),
4900            ] {
4901                state.metadata.add_segment_meta(
4902                    id.into(),
4903                    SegmentMetaInfo {
4904                        deletions: None,
4905                        num_docs: 10,
4906                        ancestors: Vec::new(),
4907                        generation: 1,
4908                        reordered: true,
4909                        bp_converged: true,
4910                        bp_unconverged_passes: 0,
4911                        seismic_pending_terms: pending,
4912                        seismic_maintenance_passes: attempts,
4913                        seismic_no_progress_passes: stalls,
4914                        ann_fragmented: false,
4915                    },
4916                );
4917            }
4918        }
4919        assert_eq!(
4920            manager.unconverged_segments_below(3).await,
4921            vec![("eligible".to_string(), 10, 0)]
4922        );
4923        {
4924            let mut state = manager.state.lock().await;
4925            let mixed = state.metadata.segment_metas.get_mut("eligible").unwrap();
4926            mixed.seismic_maintenance_passes = 0;
4927            mixed.bp_converged = false;
4928            mixed.bp_unconverged_passes = 3;
4929        }
4930        assert_eq!(
4931            manager.unconverged_segments_below(3).await,
4932            vec![("eligible".to_string(), 10, 0)],
4933            "inherited BMP cap must not hide newly copied Seismic debt"
4934        );
4935        assert!(manager.unreordered_segments().await.is_empty());
4936    }
4937
4938    #[test]
4939    fn seismic_replacements_preserve_stalls_until_progress_or_new_merge_inputs() {
4940        let parent = SegmentMetaInfo {
4941            deletions: None,
4942            num_docs: 10,
4943            ancestors: Vec::new(),
4944            generation: 1,
4945            reordered: true,
4946            bp_converged: true,
4947            bp_unconverged_passes: 0,
4948            seismic_pending_terms: 9,
4949            seismic_maintenance_passes: 100,
4950            seismic_no_progress_passes: 3,
4951            ann_fragmented: false,
4952        };
4953        for layout in [
4954            ReplacementLayout::PreserveSingleSource,
4955            ReplacementLayout::BlockCopy,
4956            ReplacementLayout::Compacted,
4957        ] {
4958            assert_eq!(
4959                replacement_seismic_state([&parent].into_iter(), 9, layout),
4960                (100, 3)
4961            );
4962        }
4963        assert_eq!(
4964            replacement_seismic_state([&parent].into_iter(), 8, ReplacementLayout::Compacted),
4965            (100, 0)
4966        );
4967        assert_eq!(
4968            replacement_seismic_state(
4969                [&parent].into_iter(),
4970                8,
4971                ReplacementLayout::BpReordered { converged: true }
4972            ),
4973            (101, 0)
4974        );
4975        assert_eq!(
4976            replacement_seismic_state(
4977                [&parent].into_iter(),
4978                9,
4979                ReplacementLayout::BpReordered { converged: true }
4980            ),
4981            (101, 4)
4982        );
4983        assert_eq!(
4984            replacement_seismic_state(
4985                [&parent, &parent].into_iter(),
4986                18,
4987                ReplacementLayout::BlockCopy
4988            ),
4989            (100, 0)
4990        );
4991        assert_eq!(
4992            replacement_seismic_state(
4993                [&parent].into_iter(),
4994                0,
4995                ReplacementLayout::BpReordered { converged: true }
4996            ),
4997            (0, 0)
4998        );
4999        let encoded = serde_json::to_value(&parent).unwrap();
5000        assert_eq!(encoded["seismic_no_progress_passes"], 3);
5001        let mut absent = encoded;
5002        absent
5003            .as_object_mut()
5004            .unwrap()
5005            .remove("seismic_no_progress_passes");
5006        assert_eq!(
5007            serde_json::from_value::<SegmentMetaInfo>(absent)
5008                .unwrap()
5009                .seismic_no_progress_passes,
5010            0
5011        );
5012    }
5013
5014    #[tokio::test]
5015    async fn unconverged_scheduler_stops_at_the_lineage_limit() {
5016        let manager = lifecycle_test_manager();
5017        {
5018            let mut state = manager.state.lock().await;
5019            state.metadata.add_segment_meta(
5020                "eligible".into(),
5021                SegmentMetaInfo {
5022                    deletions: None,
5023                    num_docs: 10,
5024                    ancestors: Vec::new(),
5025                    generation: 1,
5026                    reordered: true,
5027                    bp_converged: false,
5028                    bp_unconverged_passes: 2,
5029                    seismic_pending_terms: 0,
5030                    seismic_maintenance_passes: 0,
5031                    seismic_no_progress_passes: 0,
5032                    ann_fragmented: false,
5033                },
5034            );
5035            state.metadata.add_segment_meta(
5036                "at-limit".into(),
5037                SegmentMetaInfo {
5038                    deletions: None,
5039                    num_docs: 20,
5040                    ancestors: Vec::new(),
5041                    generation: 1,
5042                    reordered: true,
5043                    bp_converged: false,
5044                    bp_unconverged_passes: 3,
5045                    seismic_pending_terms: 0,
5046                    seismic_maintenance_passes: 0,
5047                    seismic_no_progress_passes: 0,
5048                    ann_fragmented: false,
5049                },
5050            );
5051            state.metadata.add_segment_meta(
5052                "carried-debt".into(),
5053                SegmentMetaInfo {
5054                    deletions: None,
5055                    num_docs: 15,
5056                    ancestors: Vec::new(),
5057                    generation: 2,
5058                    reordered: false,
5059                    bp_converged: false,
5060                    bp_unconverged_passes: 2,
5061                    seismic_pending_terms: 0,
5062                    seismic_maintenance_passes: 0,
5063                    seismic_no_progress_passes: 0,
5064                    ann_fragmented: false,
5065                },
5066            );
5067            state.metadata.add_segment_meta(
5068                "carried-debt-at-limit".into(),
5069                SegmentMetaInfo {
5070                    deletions: None,
5071                    num_docs: 25,
5072                    ancestors: Vec::new(),
5073                    generation: 2,
5074                    reordered: false,
5075                    bp_converged: false,
5076                    bp_unconverged_passes: 3,
5077                    seismic_pending_terms: 0,
5078                    seismic_maintenance_passes: 0,
5079                    seismic_no_progress_passes: 0,
5080                    ann_fragmented: false,
5081                },
5082            );
5083            state.metadata.add_segment_meta(
5084                "converged".into(),
5085                SegmentMetaInfo {
5086                    deletions: None,
5087                    num_docs: 30,
5088                    ancestors: Vec::new(),
5089                    generation: 1,
5090                    reordered: true,
5091                    bp_converged: true,
5092                    bp_unconverged_passes: 0,
5093                    seismic_pending_terms: 0,
5094                    seismic_maintenance_passes: 0,
5095                    seismic_no_progress_passes: 0,
5096                    ann_fragmented: false,
5097                },
5098            );
5099            state.metadata.add_segment("fresh".into(), 40);
5100            state
5101                .metadata
5102                .segment_metas
5103                .get_mut("fresh")
5104                .unwrap()
5105                .ann_fragmented = true;
5106        }
5107
5108        assert_eq!(
5109            manager.unreordered_segments().await,
5110            vec![("fresh".into(), 40)],
5111            "a block-copy output with BP debt is not a fresh first-pass candidate",
5112        );
5113        let mut eligible = manager.unconverged_segments_below(3).await;
5114        eligible.sort_unstable();
5115        assert_eq!(
5116            eligible,
5117            vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
5118        );
5119        assert!(manager.unconverged_segments_below(0).await.is_empty());
5120    }
5121
5122    #[test]
5123    fn merge_retry_backoff_is_exponential_and_capped() {
5124        assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
5125        assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
5126        assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
5127        assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
5128    }
5129
5130    #[test]
5131    fn only_deterministic_source_errors_are_quarantined() {
5132        assert!(is_deterministic_source_error(&Error::Corruption(
5133            "bad footer".into()
5134        )));
5135        assert!(is_deterministic_source_error(&Error::Io(
5136            std::io::Error::from(std::io::ErrorKind::NotFound)
5137        )));
5138        assert!(!is_deterministic_source_error(&Error::Io(
5139            std::io::Error::from(std::io::ErrorKind::TimedOut)
5140        )));
5141        assert!(!is_deterministic_source_error(&Error::Io(
5142            std::io::Error::from(std::io::ErrorKind::PermissionDenied)
5143        )));
5144    }
5145
5146    #[tokio::test]
5147    async fn transient_reorder_failure_is_backed_off_until_cleared() {
5148        let manager = lifecycle_test_manager();
5149        manager
5150            .state
5151            .lock()
5152            .await
5153            .metadata
5154            .add_segment("source".into(), 1);
5155        manager
5156            .pause_reorder_retries("source", &Error::Internal("transient".into()))
5157            .await;
5158        assert!(manager.paused_reorder_segments().contains("source"));
5159        manager.clear_reorder_retry("source");
5160        assert!(!manager.paused_reorder_segments().contains("source"));
5161    }
5162
5163    #[tokio::test]
5164    async fn optimizer_backoff_does_not_retain_retired_or_unknown_segments() {
5165        let mut schema = crate::SchemaBuilder::default();
5166        let key = schema.add_text_field("id", true, false);
5167        schema.set_primary_key(key);
5168        let index = crate::Index::create(
5169            crate::RamDirectory::new(),
5170            schema.build(),
5171            crate::IndexConfig {
5172                merge_policy: Box::new(crate::NoMergePolicy),
5173                ..Default::default()
5174            },
5175        )
5176        .await
5177        .unwrap();
5178        let mut writer = index.writer();
5179        writer.init_primary_key_dedup().await.unwrap();
5180        for value in ["dead", "live"] {
5181            let mut doc = crate::Document::new();
5182            doc.add_text(key, value);
5183            writer.add_document(doc).unwrap();
5184        }
5185        writer.commit().await.unwrap();
5186        writer.delete_primary_key("dead").unwrap();
5187        writer.commit().await.unwrap();
5188        let manager = Arc::clone(writer.segment_manager());
5189        let source = manager.get_segment_ids().await[0].clone();
5190        assert!(
5191            manager
5192                .compact_segment_if_eligible(&source, 0.3, 0)
5193                .await
5194                .is_err()
5195        );
5196        assert!(manager.reorder_retries.lock().contains_key(&source));
5197        let mut doc = crate::Document::new();
5198        doc.add_text(key, "another");
5199        writer.add_document(doc).unwrap();
5200        writer.force_merge().await.unwrap();
5201        assert!(
5202            manager.reorder_retries.lock().is_empty(),
5203            "retired segment leaked retry state"
5204        );
5205        // A delayed worker failure (or invalid caller ID) must not restore it.
5206        assert!(
5207            manager
5208                .compact_segment_if_eligible(&source, 0.3, 0)
5209                .await
5210                .is_err()
5211        );
5212        assert!(manager.reorder_retries.lock().is_empty());
5213    }
5214
5215    /// Fails `exists` with the transient I/O error class that sends a
5216    /// background merge into its generic retry backoff (not source quarantine).
5217    #[derive(Default)]
5218    struct FailingExistsDirectory(crate::directories::RamDirectory);
5219
5220    #[async_trait::async_trait]
5221    impl crate::directories::Directory for FailingExistsDirectory {
5222        async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
5223            Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
5224        }
5225
5226        async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
5227            self.0.file_size(path).await
5228        }
5229
5230        async fn open_read(
5231            &self,
5232            path: &std::path::Path,
5233        ) -> std::io::Result<crate::directories::FileHandle> {
5234            self.0.open_read(path).await
5235        }
5236
5237        async fn read_range(
5238            &self,
5239            path: &std::path::Path,
5240            range: std::ops::Range<u64>,
5241        ) -> std::io::Result<crate::directories::OwnedBytes> {
5242            self.0.read_range(path, range).await
5243        }
5244
5245        async fn list_files(
5246            &self,
5247            prefix: &std::path::Path,
5248        ) -> std::io::Result<Vec<std::path::PathBuf>> {
5249            self.0.list_files(prefix).await
5250        }
5251
5252        async fn open_lazy(
5253            &self,
5254            path: &std::path::Path,
5255        ) -> std::io::Result<crate::directories::FileHandle> {
5256            self.0.open_lazy(path).await
5257        }
5258    }
5259
5260    #[async_trait::async_trait]
5261    impl crate::directories::DirectoryWriter for FailingExistsDirectory {
5262        async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
5263            self.0.write(path, data).await
5264        }
5265
5266        async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
5267            self.0.delete(path).await
5268        }
5269
5270        async fn rename(
5271            &self,
5272            from: &std::path::Path,
5273            to: &std::path::Path,
5274        ) -> std::io::Result<()> {
5275            self.0.rename(from, to).await
5276        }
5277
5278        async fn sync(&self) -> std::io::Result<()> {
5279            self.0.sync().await
5280        }
5281
5282        async fn streaming_writer(
5283            &self,
5284            path: &std::path::Path,
5285        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
5286            self.0.streaming_writer(path).await
5287        }
5288    }
5289
5290    #[derive(Debug, Clone)]
5291    struct MergeEverythingPolicy;
5292
5293    impl MergePolicy for MergeEverythingPolicy {
5294        fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
5295            if segments.len() < 2 {
5296                return Vec::new();
5297            }
5298            vec![crate::merge::MergeCandidate {
5299                segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
5300            }]
5301        }
5302
5303        fn clone_box(&self) -> Box<dyn MergePolicy> {
5304            Box::new(self.clone())
5305        }
5306    }
5307
5308    #[tokio::test]
5309    async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
5310        let manager = lifecycle_test_manager();
5311        // Simulates a memory-budget mid-cycle segment build whose guard is
5312        // parked inside a PreparedSegment: only a later commit releases this
5313        // token, and that commit can be blocked on the very caller of the
5314        // artifact update (writer write lock / &mut self).
5315        let parked_indexing = manager
5316            .protect_new_segment("00000000000000000000000000000abc".into())
5317            .unwrap();
5318
5319        let error = tokio::time::timeout(
5320            std::time::Duration::from_secs(2),
5321            manager.begin_vector_artifact_update(),
5322        )
5323        .await
5324        .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
5325        .err()
5326        .expect("an old-generation prepared segment must block artifact replacement")
5327        .to_string();
5328        assert!(error.contains("built but uncommitted"), "{error}");
5329        assert!(
5330            !manager.vector_artifact_update.load(Ordering::Acquire),
5331            "a rejected update must release the producer gate"
5332        );
5333
5334        drop(parked_indexing);
5335
5336        let guard = manager
5337            .begin_vector_artifact_update()
5338            .await
5339            .expect("artifact update should succeed after the pending generation is resolved");
5340        drop(guard);
5341    }
5342
5343    #[tokio::test]
5344    async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
5345        let manager = lifecycle_test_manager();
5346        let merge_like = manager
5347            .active_operations
5348            .try_register(vec!["merge-source".into()])
5349            .unwrap();
5350
5351        let waiter = {
5352            let manager = Arc::clone(&manager);
5353            tokio::spawn(async move { manager.begin_vector_artifact_update().await })
5354        };
5355        for _ in 0..8 {
5356            tokio::task::yield_now().await;
5357        }
5358        assert!(
5359            !waiter.is_finished(),
5360            "artifact update must drain merge/reorder producers that may hold the previous generation"
5361        );
5362
5363        drop(merge_like);
5364        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
5365            .await
5366            .expect("artifact update missed the lifecycle guard release")
5367            .unwrap()
5368            .unwrap();
5369    }
5370
5371    #[tokio::test]
5372    async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
5373        let manager = lifecycle_test_manager();
5374        let release = Arc::new(Semaphore::new(0));
5375        let merge_task = {
5376            let release = Arc::clone(&release);
5377            tokio::spawn(async move {
5378                let _permit = release.acquire().await.unwrap();
5379            })
5380        };
5381        manager.merge_handles.lock().push(merge_task);
5382
5383        let waiter = {
5384            let manager = Arc::clone(&manager);
5385            tokio::spawn(async move { manager.wait_for_all_merges().await })
5386        };
5387        for _ in 0..8 {
5388            tokio::task::yield_now().await;
5389        }
5390        assert!(!waiter.is_finished());
5391        // Simulates tonic dropping a force_merge/reorder RPC future at the
5392        // JoinHandle await when the client disconnects.
5393        waiter.abort();
5394        let join_error = waiter.await.unwrap_err();
5395        assert!(join_error.is_cancelled());
5396
5397        assert!(
5398            !manager.merge_handles.lock().is_empty(),
5399            "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
5400        );
5401
5402        // A later drain must still see and await the real in-flight merge.
5403        release.add_permits(1);
5404        tokio::time::timeout(
5405            std::time::Duration::from_secs(1),
5406            manager.wait_for_all_merges(),
5407        )
5408        .await
5409        .expect("subsequent drain missed the reinserted merge handle");
5410        assert!(manager.merge_handles.lock().is_empty());
5411    }
5412
5413    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5414    async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
5415        let manager = lifecycle_test_manager();
5416        {
5417            let mut state = manager.state.lock().await;
5418            state
5419                .metadata
5420                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
5421            state
5422                .metadata
5423                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
5424        }
5425        // A background reorder (or a concurrent force-merge) owns one segment
5426        // in the batch but never appears in merge_handles.
5427        let reorder_like = manager
5428            .active_operations
5429            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
5430            .unwrap();
5431
5432        let force_merge = {
5433            let manager = Arc::clone(&manager);
5434            tokio::spawn(async move { manager.force_merge().await })
5435        };
5436
5437        tokio::time::sleep(std::time::Duration::from_millis(600)).await;
5438        let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
5439        assert!(
5440            retries >= 1,
5441            "force_merge never observed the conflicting owner (retries={retries})"
5442        );
5443        assert!(
5444            retries < 20,
5445            "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
5446        );
5447
5448        drop(reorder_like);
5449        // With the conflict gone the loop proceeds; the batch then fails fast
5450        // in do_merge (the test IDs have no files), proving the loop exited.
5451        let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
5452            .await
5453            .expect("force_merge kept spinning after the conflicting owner released")
5454            .unwrap();
5455        assert!(result.is_err());
5456    }
5457
5458    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
5459    async fn force_merge_routes_around_segments_held_by_reorder() {
5460        let manager = lifecycle_test_manager();
5461        {
5462            let mut state = manager.state.lock().await;
5463            state
5464                .metadata
5465                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
5466            state
5467                .metadata
5468                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
5469            state
5470                .metadata
5471                .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
5472        }
5473        // A background reorder owns one segment and holds it for the whole
5474        // test (in prod: a BP pass runs for minutes while force_merge spins).
5475        let _reorder_like = manager
5476            .active_operations
5477            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
5478            .unwrap();
5479
5480        // Regression: force_merge used to rebuild the identical smallest-N
5481        // batch (including the held segment) every 100ms and retry-log
5482        // forever. It must instead skip the held segment and immediately
5483        // make progress on the two free ones — reaching do_merge (which
5484        // fails fast here: the test IDs have no files) proves the batch was
5485        // built without the held segment while the reorder is STILL active.
5486        let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
5487            let manager = Arc::clone(&manager);
5488            async move { manager.force_merge().await }
5489        })
5490        .await
5491        .expect("force_merge livelocked on a segment held by an active reorder");
5492        assert!(result.is_err(), "fake segment files must fail the merge");
5493
5494        assert_eq!(
5495            manager.force_merge_conflict_retries.load(Ordering::Relaxed),
5496            0,
5497            "batch built from the ownership snapshot must not collide with the held segment"
5498        );
5499    }
5500
5501    #[tokio::test]
5502    async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
5503        let schema = crate::dsl::SchemaBuilder::default().build();
5504        let mut metadata = IndexMetadata::new(schema.clone());
5505        metadata.add_segment("00000000000000000000000000000001".into(), 10);
5506        metadata.add_segment("00000000000000000000000000000002".into(), 10);
5507        let manager = Arc::new(SegmentManager::new(
5508            Arc::new(FailingExistsDirectory::default()),
5509            Arc::new(schema),
5510            metadata,
5511            Box::new(MergeEverythingPolicy),
5512            0,
5513            1,
5514            Arc::new(Semaphore::new(1)),
5515            None,
5516            1024,
5517            Arc::new(ReorderConcurrencyGate::new(1)),
5518            None,
5519        ));
5520
5521        // Spawns a background merge that fails with a transient I/O error and
5522        // arms the 30s..30min retry backoff.
5523        manager.maybe_merge().await;
5524
5525        tokio::time::timeout(
5526            std::time::Duration::from_secs(5),
5527            manager.wait_for_all_merges(),
5528        )
5529        .await
5530        .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
5531        assert!(
5532            manager.merge_retry_is_paused(),
5533            "the failed merge should have armed the retry backoff"
5534        );
5535
5536        // Shutdown still drains the pending backoff wakeup deterministically.
5537        manager.begin_shutdown();
5538        tokio::time::timeout(
5539            std::time::Duration::from_secs(5),
5540            manager.wait_for_shutdown(),
5541        )
5542        .await
5543        .expect("shutdown did not drain the merge retry wakeup task");
5544    }
5545}
5546#[path = "row_mutation.rs"]
5547mod row_mutation;