Skip to main content

hermes_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::ArcSwapOption;
49use tokio::sync::Mutex as AsyncMutex;
50use tokio::sync::{Notify, OwnedSemaphorePermit, Semaphore};
51use tokio::task::JoinHandle;
52
53use crate::directories::DirectoryWriter;
54use crate::error::{Error, Result};
55use crate::index::{IndexMetadata, ReorderConcurrencyGate, ReorderPriority, SegmentMetaInfo};
56use crate::segment::{
57    SegmentFiles, SegmentId, SegmentMeta, SegmentSnapshot, SegmentTracker, TrainedVectorStructures,
58};
59#[cfg(feature = "native")]
60use crate::segment::{SegmentMerger, SegmentReader};
61
62use super::{MergePolicy, SegmentInfo};
63
64const FORCE_MERGE_MAX_FAN_IN: usize = 64;
65
66#[derive(Debug)]
67struct ForceMergeGroup {
68    segments: Vec<(String, u32)>,
69    total_docs: u64,
70}
71
72/// Deterministically pack segments into near-minimal final outputs.
73///
74/// Best-fit decreasing avoids the pathological smallest-first behavior where,
75/// for example, `4 + 4` is merged before two `6 + 4` pairs under a 10-doc
76/// limit. That old order both left excess segments and made its intermediate
77/// outputs eligible for another full BP pass. Bin packing is NP-hard; BFD is a
78/// bounded, deterministic approximation that produces maximal groups (no two
79/// output groups can still fit together).
80fn plan_force_merge_groups(
81    mut segments: Vec<(String, u32)>,
82    max_docs: u64,
83) -> Vec<ForceMergeGroup> {
84    segments.sort_unstable_by(|(left_id, left_docs), (right_id, right_docs)| {
85        right_docs
86            .cmp(left_docs)
87            .then_with(|| left_id.cmp(right_id))
88    });
89
90    let mut groups: Vec<ForceMergeGroup> = Vec::new();
91    for segment in segments {
92        let docs = u64::from(segment.1);
93        let best_group = groups
94            .iter()
95            .enumerate()
96            .filter_map(|(index, group)| {
97                group
98                    .total_docs
99                    .checked_add(docs)
100                    .filter(|&total| total <= max_docs)
101                    .map(|_| (index, group.total_docs))
102            })
103            .max_by_key(|&(index, used)| (used, std::cmp::Reverse(index)))
104            .map(|(index, _)| index);
105
106        if let Some(index) = best_group {
107            groups[index].total_docs += docs;
108            groups[index].segments.push(segment);
109        } else {
110            groups.push(ForceMergeGroup {
111                segments: vec![segment],
112                total_docs: docs,
113            });
114        }
115    }
116
117    // Preserve the old small-to-large source order inside each output. Besides
118    // stable document ordering, BMP block-copy compression depends on source
119    // alignment at group boundaries.
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}
284
285impl ActiveSegmentOperations {
286    fn new() -> Self {
287        Self {
288            inner: parking_lot::Mutex::new(ActiveOperationState {
289                segment_ids: HashSet::new(),
290                operation_tokens: HashSet::new(),
291                indexing_tokens: HashSet::new(),
292                next_operation_token: 0,
293                accepting: true,
294                non_indexing_paused: false,
295            }),
296            idle: Notify::new(),
297            shutdown: Notify::new(),
298            shutdown_requested: Arc::new(AtomicBool::new(false)),
299        }
300    }
301
302    /// Try to claim IDs for a self-draining lifecycle operation (merge,
303    /// reorder, cleanup). Returns a guard on success, `None` if any requested
304    /// ID is already owned by another active operation.
305    fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
306        self.try_register_kind(segment_ids, false, false)
307    }
308
309    /// Try to claim IDs for an indexing producer whose guard is held until
310    /// metadata publication (commit) rather than task completion.
311    fn try_register_indexing(
312        self: &Arc<Self>,
313        segment_ids: Vec<String>,
314    ) -> Option<SegmentOperationGuard> {
315        self.try_register_kind(segment_ids, true, false)
316    }
317
318    /// Claim source/output IDs for the exclusive vector-generation updater,
319    /// which is the only non-indexing producer allowed through its own pause.
320    fn try_register_vector_update(
321        self: &Arc<Self>,
322        segment_ids: Vec<String>,
323    ) -> Option<SegmentOperationGuard> {
324        self.try_register_kind(segment_ids, false, true)
325    }
326
327    fn try_register_kind(
328        self: &Arc<Self>,
329        segment_ids: Vec<String>,
330        indexing: bool,
331        vector_update: bool,
332    ) -> Option<SegmentOperationGuard> {
333        let mut inner = self.inner.lock();
334        if !inner.accepting {
335            log::debug!("[segment_lifecycle] rejected operation during shutdown");
336            return None;
337        }
338        if !indexing && !vector_update && inner.non_indexing_paused {
339            log::debug!("[segment_lifecycle] deferred operation during dense vector retraining");
340            return None;
341        }
342        // Check for overlap with any active lifecycle operation.
343        for id in &segment_ids {
344            if inner.segment_ids.contains(id) {
345                log::debug!(
346                    "[segment_lifecycle] rejected: {} overlaps with an active operation ({} active IDs)",
347                    id,
348                    inner.segment_ids.len()
349                );
350                return None;
351            }
352        }
353        log::debug!(
354            "[segment_lifecycle] registered {} IDs (total active: {})",
355            segment_ids.len(),
356            inner.segment_ids.len() + segment_ids.len()
357        );
358        let operation_token = inner.next_operation_token;
359        let next_operation_token = operation_token.checked_add(1)?;
360        for id in &segment_ids {
361            inner.segment_ids.insert(id.clone());
362        }
363        inner.next_operation_token = next_operation_token;
364        inner.operation_tokens.insert(operation_token);
365        if indexing {
366            inner.indexing_tokens.insert(operation_token);
367        }
368        Some(SegmentOperationGuard {
369            active_operations: Arc::clone(self),
370            segment_ids,
371            operation_token,
372        })
373    }
374
375    /// Snapshot of all IDs owned by active operations.
376    fn snapshot(&self) -> HashSet<String> {
377        self.inner.lock().segment_ids.clone()
378    }
379
380    /// Exact identities of self-draining operations (merge/reorder/cleanup)
381    /// active at one instant, plus the number of indexing tokens excluded.
382    /// Unlike segment IDs, tokens cannot be reused by a later retry, so an
383    /// artifact-update barrier can drain only pre-gate producers without being
384    /// starved by new flat producers.
385    ///
386    /// Indexing tokens are deliberately excluded: their guards are parked in
387    /// built-but-uncommitted `PreparedSegment`s and only a later commit — which
388    /// may be blocked on the barrier's caller — releases them, so waiting on
389    /// them deadlocks (see `begin_vector_artifact_update`).
390    fn draining_operation_tokens_snapshot(&self) -> (HashSet<u64>, usize) {
391        let inner = self.inner.lock();
392        let tokens = inner
393            .operation_tokens
394            .difference(&inner.indexing_tokens)
395            .copied()
396            .collect();
397        (tokens, inner.indexing_tokens.len())
398    }
399
400    /// Atomically prevent new lifecycle work from starting. Existing guards
401    /// remain valid and can be drained with [`Self::wait_until_idle`].
402    fn stop_accepting(&self) {
403        self.shutdown_requested.store(true, Ordering::Release);
404        let mut inner = self.inner.lock();
405        inner.accepting = false;
406        self.shutdown.notify_waiters();
407        if inner.segment_ids.is_empty() {
408            self.idle.notify_waiters();
409        }
410    }
411
412    fn pause_non_indexing(&self) {
413        self.inner.lock().non_indexing_paused = true;
414    }
415
416    fn resume_non_indexing(&self) {
417        self.inner.lock().non_indexing_paused = false;
418        self.idle.notify_waiters();
419    }
420
421    fn is_accepting(&self) -> bool {
422        self.inner.lock().accepting
423    }
424
425    fn cancellation_flag(&self) -> Arc<AtomicBool> {
426        Arc::clone(&self.shutdown_requested)
427    }
428
429    /// Wait until every operation that started before shutdown has released
430    /// its ownership. Register/check and notification are ordered to avoid a
431    /// missed wakeup between observing a non-empty set and awaiting.
432    async fn wait_until_idle(&self) {
433        loop {
434            let notified = self.idle.notified();
435            if self.inner.lock().segment_ids.is_empty() {
436                return;
437            }
438            notified.await;
439        }
440    }
441
442    async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
443        while !operations.is_empty() {
444            let notified = self.idle.notified();
445            if self.inner.lock().operation_tokens.is_disjoint(operations) {
446                return;
447            }
448            notified.await;
449        }
450    }
451
452    /// Resolve when shutdown starts, without missing a notification between
453    /// checking the state and registering the waiter.
454    async fn wait_for_shutdown(&self) {
455        loop {
456            let notified = self.shutdown.notified();
457            if !self.inner.lock().accepting {
458                return;
459            }
460            notified.await;
461        }
462    }
463}
464
465/// RAII ownership of segment IDs used by an active lifecycle operation.
466/// Dropping on success, error, cancellation, or panic makes abandoned outputs
467/// eligible for sweeping automatically.
468pub(crate) struct SegmentOperationGuard {
469    active_operations: Arc<ActiveSegmentOperations>,
470    segment_ids: Vec<String>,
471    operation_token: u64,
472}
473
474impl Drop for SegmentOperationGuard {
475    fn drop(&mut self) {
476        let mut inner = self.active_operations.inner.lock();
477        for id in &self.segment_ids {
478            inner.segment_ids.remove(id);
479        }
480        inner.operation_tokens.remove(&self.operation_token);
481        inner.indexing_tokens.remove(&self.operation_token);
482        // Token barriers need notification on every completion, not only the
483        // transition to complete global idleness.
484        self.active_operations.idle.notify_waiters();
485        if inner.segment_ids.is_empty() {
486            debug_assert!(inner.operation_tokens.is_empty());
487        }
488    }
489}
490
491/// Exclusive gate for an index-level trained-vector artifact update.
492///
493/// Segment producers consult this gate before capturing the current trained
494/// structures. Once the gate is raised, new producers deliberately emit flat
495/// vector data; waiting for already-active producers to drain then guarantees
496/// that the committed source set stays stable while replacements are staged.
497struct VectorArtifactUpdateLease {
498    updating: Arc<AtomicBool>,
499    active_operations: Arc<ActiveSegmentOperations>,
500}
501
502impl Drop for VectorArtifactUpdateLease {
503    fn drop(&mut self) {
504        self.updating.store(false, Ordering::Release);
505        self.active_operations.resume_non_indexing();
506    }
507}
508
509#[derive(Clone)]
510pub(crate) struct VectorArtifactUpdateGuard {
511    _lease: Arc<VectorArtifactUpdateLease>,
512}
513
514/// Merge-time/manual BP pools are shared by every index in this process.
515/// A pool per `SegmentManager` multiplied a 96-core host into two 48-thread
516/// merge pools plus the optimizer pool (200+ process threads in production).
517static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
518
519const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
520const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
521
522#[derive(Default)]
523struct MergeRetryState {
524    retry_after: Option<std::time::Instant>,
525    consecutive_failures: u32,
526}
527
528fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
529    let shift = consecutive_failures.saturating_sub(1).min(16);
530    MERGE_RETRY_BASE_DELAY
531        .checked_mul(1u32 << shift)
532        .unwrap_or(MERGE_RETRY_MAX_DELAY)
533        .min(MERGE_RETRY_MAX_DELAY)
534}
535
536/// Merge JoinHandles taken out of the shared list for draining.
537///
538/// Drain futures are awaited inline by RPC handlers (force_merge/reorder) and
539/// can be dropped at any await when a client disconnects. Handles are awaited
540/// through this guard and removed only after completion, so a cancelled drain
541/// returns every un-awaited (and possibly still-running) merge to the shared
542/// list instead of silently detaching it from shutdown, abort, and
543/// force-merge tracking.
544struct DrainedMergeHandles<'a> {
545    shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>,
546    drained: Vec<JoinHandle<()>>,
547}
548
549impl<'a> DrainedMergeHandles<'a> {
550    fn take(shared: &'a parking_lot::Mutex<Vec<JoinHandle<()>>>) -> Self {
551        let drained = std::mem::take(&mut *shared.lock());
552        Self { shared, drained }
553    }
554
555    fn is_empty(&self) -> bool {
556        self.drained.is_empty()
557    }
558
559    /// Await the next handle. It stays owned by this guard while being polled
560    /// and is discarded only once it has completed, so cancellation at the
561    /// await reinserts it via `Drop`.
562    async fn join_next(&mut self) -> Option<std::result::Result<(), tokio::task::JoinError>> {
563        let handle = self.drained.last_mut()?;
564        let result = handle.await;
565        self.drained.pop();
566        Some(result)
567    }
568}
569
570impl Drop for DrainedMergeHandles<'_> {
571    fn drop(&mut self) {
572        if !self.drained.is_empty() {
573            self.shared.lock().append(&mut self.drained);
574        }
575    }
576}
577
578/// Spawn and register auxiliary lifecycle work as one synchronous operation.
579///
580/// Registering *after* `spawn` left a small deletion race: shutdown could
581/// observe an empty handle list while the newly spawned filesystem task was
582/// already running. Holding the handle-list mutex across `Handle::spawn`
583/// makes task creation visible to the drain before either side can proceed.
584fn try_spawn_lifecycle<F>(
585    handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
586    runtime: &tokio::runtime::Handle,
587    future: F,
588) -> bool
589where
590    F: std::future::Future<Output = ()> + Send + 'static,
591{
592    std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
593        let mut handles = handles.lock();
594        handles.retain(|handle| !handle.is_finished());
595        handles.push(runtime.spawn(future));
596    }))
597    .is_ok()
598}
599
600/// Deletes an uncommitted merge/reorder output if its task unwinds.
601///
602/// Normal `Result::Err` paths delete outputs synchronously so callers observe
603/// a clean directory before returning. This guard covers the path those
604/// branches cannot: a panic after output files have been created. The cleanup
605/// callback re-checks metadata before deleting, so a panic after a successful
606/// metadata commit cannot remove a live segment.
607struct OutputCleanupGuard {
608    segment_id: SegmentId,
609    cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
610}
611
612impl OutputCleanupGuard {
613    fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
614        Self {
615            segment_id,
616            cleanup: Some(cleanup),
617        }
618    }
619
620    fn disarm(&mut self) {
621        self.cleanup = None;
622    }
623}
624
625impl Drop for OutputCleanupGuard {
626    fn drop(&mut self) {
627        if let Some(cleanup) = self.cleanup.take() {
628            cleanup(self.segment_id);
629        }
630    }
631}
632
633/// All mutable state behind the single async Mutex.
634struct ManagerState {
635    metadata: IndexMetadata,
636    merge_policy: Box<dyn MergePolicy>,
637}
638
639type ReplacementRefresh = Arc<
640    dyn Fn() -> std::pin::Pin<Box<dyn std::future::Future<Output = Result<()>> + Send>>
641        + Send
642        + Sync,
643>;
644
645/// Reconcile a long-lived consumer after a durable segment replacement.
646/// This runs inside the same lifecycle-owned task as metadata publication, so
647/// cancelling the caller cannot leave the replacement published but unannounced.
648async fn refresh_replacement_topology(refresh: Option<ReplacementRefresh>) {
649    let Some(refresh) = refresh else {
650        return;
651    };
652    let mut last_error = None;
653    for attempt in 0..3 {
654        match refresh().await {
655            Ok(()) => return,
656            Err(error) => {
657                last_error = Some(error);
658                if attempt < 2 {
659                    tokio::time::sleep(std::time::Duration::from_secs(1 << attempt)).await;
660                }
661            }
662        }
663    }
664    if let Some(error) = last_error {
665        log::warn!(
666            "[segment_lifecycle] replacement topology refresh failed after 3 attempts: {}",
667            error,
668        );
669    }
670}
671
672#[cfg(feature = "native")]
673struct MergeTaskError {
674    error: Error,
675    unavailable_segments: Vec<String>,
676}
677
678#[cfg(feature = "native")]
679impl MergeTaskError {
680    fn source(segment_id: String, error: Error) -> Self {
681        Self {
682            error,
683            unavailable_segments: vec![segment_id],
684        }
685    }
686
687    fn sources(segment_ids: Vec<String>, error: Error) -> Self {
688        Self {
689            error,
690            unavailable_segments: segment_ids,
691        }
692    }
693}
694
695#[cfg(feature = "native")]
696impl From<Error> for MergeTaskError {
697    fn from(error: Error) -> Self {
698        Self {
699            error,
700            unavailable_segments: Vec::new(),
701        }
702    }
703}
704
705#[cfg(feature = "native")]
706fn is_deterministic_source_error(error: &Error) -> bool {
707    matches!(error, Error::Corruption(_) | Error::Serialization(_))
708        || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
709}
710
711#[cfg(feature = "native")]
712fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
713    if is_deterministic_source_error(&error) {
714        MergeTaskError::source(segment_id, error)
715    } else {
716        // Timeouts, interrupted reads, permission changes, and other generic
717        // I/O failures may be transient. Back them off instead of quarantining
718        // a healthy metadata segment for the rest of the process lifetime.
719        MergeTaskError::from(error)
720    }
721}
722
723#[cfg(feature = "native")]
724type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
725
726#[derive(Clone, Copy)]
727enum ReplacementLayout {
728    /// No BP ran while combining these sources. The combined segment is not
729    /// globally reordered, but an interrupted BP lineage remains owed its
730    /// record-level deepening pass.
731    BlockCopy,
732    /// A BP pass produced the replacement layout.
733    BpReordered { converged: bool },
734    /// A vector-only rewrite leaves document order and sparse layout exactly
735    /// unchanged, so its persisted BP progress must not be reset or advanced.
736    PreserveSingleSource,
737}
738
739fn replacement_bp_state(
740    parent_has_debt: bool,
741    parent_unconverged_passes: u32,
742    layout: ReplacementLayout,
743) -> (bool, bool, u32) {
744    match layout {
745        ReplacementLayout::BlockCopy => (
746            false,
747            !parent_has_debt,
748            if parent_has_debt {
749                parent_unconverged_passes
750            } else {
751                0
752            },
753        ),
754        ReplacementLayout::BpReordered { converged } => (
755            true,
756            converged,
757            if converged {
758                0
759            } else {
760                parent_unconverged_passes.saturating_add(1)
761            },
762        ),
763        ReplacementLayout::PreserveSingleSource => {
764            unreachable!("preserved layouts retain the complete source metadata")
765        }
766    }
767}
768
769#[derive(Clone, Copy, Debug, Eq, PartialEq)]
770enum VectorSegmentRewriteOutcome {
771    Rewritten,
772    AlreadyCurrent,
773    SourceGone,
774    Conflict,
775    Deferred,
776}
777
778/// Complete but unpublished vector-only replacement. Its lifecycle claim and
779/// cleanup guard stay armed until the whole codebook generation commits.
780pub(crate) struct StagedVectorSegment {
781    source_id: String,
782    output_id: SegmentId,
783    doc_count: u32,
784    _operation: SegmentOperationGuard,
785    cleanup: OutputCleanupGuard,
786}
787
788/// Segment manager — coordinates segment commit, background merging, and trained structures.
789///
790/// SOLE owner of `metadata.json`. All metadata mutations go through `state` Mutex.
791pub struct SegmentManager<D: DirectoryWriter + 'static> {
792    /// Serializes ALL metadata mutations.
793    state: Arc<AsyncMutex<ManagerState>>,
794
795    /// RAII ownership for every in-flight segment lifecycle operation.
796    active_operations: Arc<ActiveSegmentOperations>,
797
798    /// Metadata-live segments involved in a deterministic source/corruption
799    /// failure. They stay searchable (and operator-visible) but are excluded
800    /// from merges for this process lifetime, preventing a bad candidate from
801    /// consuming full rewrite capacity on every retry.
802    quarantined_segments: parking_lot::Mutex<HashSet<String>>,
803
804    /// Generic merge failures pause scheduling briefly. Source-specific open
805    /// failures use `quarantined_segments` instead so healthy work can continue.
806    merge_retry: parking_lot::Mutex<MergeRetryState>,
807
808    /// Per-source backoff for non-deterministic standalone reorder failures.
809    /// Optimizer scans are periodic, but a pass can outlast the scan interval;
810    /// without completion-based backoff it would restart almost immediately.
811    reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
812
813    /// In-flight merge JoinHandles — supports multiple concurrent merges.
814    merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
815
816    /// At most one task per index waits for application-wide merge capacity.
817    /// Without this wakeup, an index denied by another index can remain idle
818    /// forever when no later commit happens to re-run merge policy evaluation.
819    global_merge_wakeup_pending: AtomicBool,
820
821    /// Non-zero while an explicit force merge is draining/running. Automatic
822    /// merges do not start during this window; otherwise a merge spawned
823    /// between the initial drain and foreground BP reservation could claim
824    /// source segments and then wait behind the foreground gate.
825    force_merge_active: AtomicUsize,
826
827    /// Times `force_merge` observed a conflicting active operation and retried.
828    /// Test-only observability for the conflict-retry backoff.
829    #[cfg(test)]
830    force_merge_conflict_retries: std::sync::atomic::AtomicU64,
831
832    /// Auxiliary lifecycle tasks: metadata transactions, deferred deletes,
833    /// and capacity wakeups. Handles registered here are drained before index
834    /// removal.
835    lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
836
837    /// Trained vector structures — lock-free reads via ArcSwap.
838    /// Wrapped in `Arc` so cancellation-safe metadata transactions can publish
839    /// the matching in-memory generation after their durable commit point.
840    trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
841
842    /// Raised while index-level trained artifacts and their metadata are being
843    /// replaced. Search readers keep using the last valid generation, while
844    /// segment producers fall back to flat output until publication completes.
845    vector_artifact_update: Arc<AtomicBool>,
846
847    /// Reference counting for safe segment deletion (sync Mutex for Drop).
848    tracker: Arc<SegmentTracker>,
849
850    /// Cached deletion callback for snapshots (avoids allocation per acquire_snapshot).
851    delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
852
853    /// Directory for segment I/O
854    directory: Arc<D>,
855    /// Schema for segment operations
856    schema: Arc<crate::dsl::Schema>,
857    /// Term cache blocks for segment readers during merge
858    term_cache_blocks: usize,
859    /// Hard concurrency limit for background merges. A semaphore permit is
860    /// acquired before lifecycle ownership, closing the old handle-count race
861    /// where concurrent schedulers could exceed the configured maximum.
862    merge_permits: Arc<Semaphore>,
863    /// Application-wide merge limit shared across index managers.
864    global_merge_permits: Arc<Semaphore>,
865    /// Shared across every index opened from the same `IndexConfig`. This
866    /// bounds whole BP rewrites (optimizer + merge-time + manual) separately
867    /// from Rayon thread width, preventing N × memory-budget amplification.
868    reorder_permits: Arc<ReorderConcurrencyGate>,
869    /// Run BP reordering of `reorder`-attributed BMP fields inside merges.
870    /// Persisted index configuration (schema-level `reorder_on_merge: true`
871    /// in SDL); merged segments are marked `reordered` and skipped by the
872    /// standalone optimizer pass.
873    reorder_on_merge: bool,
874    /// Wall-clock budget for merge-time BP (from `IndexConfig`); truncated
875    /// passes mark the merged segment `bp_converged = false` so the
876    /// background optimizer deepens it later (warm-started).
877    merge_bp_time_budget: Option<std::time::Duration>,
878    /// Memory budget for the BP forward index (merge-time and background
879    /// reorder). Over-budget passes drop highest-df dims, logged loudly.
880    bp_memory_budget_bytes: usize,
881    /// Application-owned shared pool, when configured. This is the server
882    /// path and ensures optimizer and merge-time work use the same threads.
883    background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
884    /// Writer-owned topology reconciler. Every durable replacement invokes
885    /// it from tracked lifecycle work so background merges cannot leave
886    /// primary-key snapshots pinning retired sources indefinitely.
887    replacement_refresh: parking_lot::RwLock<Option<ReplacementRefresh>>,
888}
889
890struct ForceMergeActivityGuard<'a>(&'a AtomicUsize);
891
892impl Drop for ForceMergeActivityGuard<'_> {
893    fn drop(&mut self) {
894        self.0.fetch_sub(1, Ordering::AcqRel);
895    }
896}
897
898impl<D: DirectoryWriter + 'static> SegmentManager<D> {
899    /// Create a new segment manager with existing metadata
900    #[allow(clippy::too_many_arguments)]
901    pub fn new(
902        directory: Arc<D>,
903        schema: Arc<crate::dsl::Schema>,
904        metadata: IndexMetadata,
905        merge_policy: Box<dyn MergePolicy>,
906        term_cache_blocks: usize,
907        max_concurrent_merges: usize,
908        global_merge_permits: Arc<Semaphore>,
909        merge_bp_time_budget: Option<std::time::Duration>,
910        bp_memory_budget_bytes: usize,
911        reorder_permits: Arc<ReorderConcurrencyGate>,
912        background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
913    ) -> Self {
914        // Persisted index option: set via `reorder_on_merge: true` in the SDL
915        // at index creation. Absent = disabled (merges block-copy).
916        let reorder_on_merge = schema.reorder_on_merge();
917        if reorder_on_merge {
918            log::info!("[merge] reorder-on-merge enabled by index schema");
919        }
920
921        let tracker = Arc::new(SegmentTracker::new());
922        for seg_id in metadata.segment_metas.keys() {
923            tracker.register(seg_id);
924        }
925
926        let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
927            Arc::new(parking_lot::Mutex::new(Vec::new()));
928        let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
929            let dir = Arc::clone(&directory);
930            let tracker = Arc::clone(&tracker);
931            let lifecycle_handles = Arc::clone(&lifecycle_handles);
932            Arc::new(move |segment_ids| {
933                // Guard: if the tokio runtime is gone (program exit), skip async
934                // deletion. Segment files become orphans cleaned up on next startup.
935                let Ok(handle) = tokio::runtime::Handle::try_current() else {
936                    // Release in-process protection as well: if the process is
937                    // still alive, a later sweep must be able to retry.
938                    tracker.complete_deletion(&segment_ids);
939                    return;
940                };
941                let dir = Arc::clone(&dir);
942                let task_tracker = Arc::clone(&tracker);
943                let cleanup_ids = segment_ids.clone();
944                let future = async move {
945                    for &segment_id in &segment_ids {
946                        log::info!(
947                            "[segment_cleanup] deleting deferred segment {}",
948                            segment_id.to_hex()
949                        );
950                        if let Err(error) =
951                            crate::segment::delete_segment(dir.as_ref(), segment_id).await
952                        {
953                            log::warn!(
954                                "[segment_cleanup] deferred delete failed for {}: {}",
955                                segment_id.to_hex(),
956                                error,
957                            );
958                        }
959                    }
960                    task_tracker.complete_deletion(&segment_ids);
961                };
962                if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
963                    // Spawning can fail only during runtime teardown. Release
964                    // the scheduled-deletion claim so an in-process sweep can
965                    // retry; crash recovery handles a process exit.
966                    tracker.complete_deletion(&cleanup_ids);
967                    log::warn!(
968                        "[segment_cleanup] runtime rejected deferred deletion; files will be swept later"
969                    );
970                }
971            })
972        };
973
974        Self {
975            state: Arc::new(AsyncMutex::new(ManagerState {
976                metadata,
977                merge_policy,
978            })),
979            active_operations: Arc::new(ActiveSegmentOperations::new()),
980            quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
981            merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
982            reorder_retries: parking_lot::Mutex::new(HashMap::new()),
983            merge_handles: parking_lot::Mutex::new(Vec::new()),
984            global_merge_wakeup_pending: AtomicBool::new(false),
985            force_merge_active: AtomicUsize::new(0),
986            #[cfg(test)]
987            force_merge_conflict_retries: std::sync::atomic::AtomicU64::new(0),
988            lifecycle_handles,
989            trained: Arc::new(ArcSwapOption::new(None)),
990            vector_artifact_update: Arc::new(AtomicBool::new(false)),
991            tracker,
992            delete_fn,
993            directory,
994            schema,
995            term_cache_blocks,
996            merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
997            global_merge_permits,
998            reorder_permits,
999            reorder_on_merge,
1000            merge_bp_time_budget,
1001            bp_memory_budget_bytes,
1002            background_reorder_pool,
1003            replacement_refresh: parking_lot::RwLock::new(None),
1004        }
1005    }
1006
1007    pub(crate) fn set_replacement_refresh<F, Fut>(&self, refresh: F)
1008    where
1009        F: Fn() -> Fut + Send + Sync + 'static,
1010        Fut: std::future::Future<Output = Result<()>> + Send + 'static,
1011    {
1012        *self.replacement_refresh.write() = Some(Arc::new(move || Box::pin(refresh())));
1013    }
1014
1015    /// Bounded rayon pool for background CPU (merge-time BP, manual reorder,
1016    /// and global vector-codebook training).
1017    /// Query scoring uses a dedicated search pool; keeping background work on
1018    /// this separate bounded pool prevents a merge or retrain from queueing
1019    /// every search behind its CPU passes.
1020    pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
1021        if let Some(pool) = &self.background_reorder_pool {
1022            return Arc::clone(pool);
1023        }
1024        Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
1025            let threads = (num_cpus::get() / 2).max(1);
1026            log::info!(
1027                "[merge] process-wide background CPU pool: {} thread(s)",
1028                threads
1029            );
1030            Arc::new(
1031                rayon::ThreadPoolBuilder::new()
1032                    .num_threads(threads)
1033                    .thread_name(|i| format!("hermes-bg-cpu-{}", i))
1034                    .build()
1035                    .expect("failed to build background CPU pool"),
1036            )
1037        }))
1038    }
1039
1040    /// Stop new indexing/merge/reorder operations from claiming segment IDs.
1041    /// Used as the first half of index deletion; the writer then joins its
1042    /// workers before [`Self::wait_for_shutdown`] drains remaining ownership.
1043    pub fn begin_shutdown(&self) {
1044        self.active_operations.stop_accepting();
1045    }
1046
1047    /// Run a lifecycle mutation independently of its requesting future.
1048    ///
1049    /// Metadata writes contain an atomic rename. If an RPC is cancelled while
1050    /// awaiting that I/O, dropping the request must not abandon the matching
1051    /// in-memory/tracker transition. The spawned transaction is tracked for
1052    /// index shutdown; the oneshot only reports its result to a caller that is
1053    /// still interested.
1054    async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
1055    where
1056        T: Send + 'static,
1057        F: std::future::Future<Output = Result<T>> + Send + 'static,
1058    {
1059        let (result_tx, result_rx) = tokio::sync::oneshot::channel();
1060        let future = async move {
1061            let result = transaction.await;
1062            let _ = result_tx.send(result);
1063        };
1064        let runtime = tokio::runtime::Handle::current();
1065        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1066            return Err(Error::Internal(
1067                "runtime rejected lifecycle metadata transaction".into(),
1068            ));
1069        }
1070        result_rx.await.map_err(|_| {
1071            Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
1072        })?
1073    }
1074
1075    /// Arm unwind cleanup for an output that is not visible in metadata yet.
1076    fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
1077        let manager = Arc::clone(self);
1078        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
1079            let Ok(handle) = tokio::runtime::Handle::try_current() else {
1080                log::warn!(
1081                    "[segment_cleanup] runtime unavailable; partial output {} will be swept on startup",
1082                    segment_id.to_hex(),
1083                );
1084                return;
1085            };
1086
1087            let cleanup_manager = Arc::clone(&manager);
1088            let future = async move {
1089                cleanup_manager
1090                    .delete_output_if_unregistered(segment_id, "task unwind")
1091                    .await;
1092            };
1093            if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
1094                log::warn!(
1095                    "[segment_cleanup] runtime rejected output cleanup; {} will be swept on startup",
1096                    segment_id.to_hex(),
1097                );
1098            }
1099        });
1100
1101        OutputCleanupGuard::new(output_id, cleanup)
1102    }
1103
1104    /// Delete an abandoned indexing output while retaining its lifecycle
1105    /// claim until the last file operation completes. The explicit runtime
1106    /// handle makes this safe from dedicated indexing OS threads, which are
1107    /// outside Tokio's entered context.
1108    pub(crate) fn schedule_unpublished_segment_cleanup(
1109        self: &Arc<Self>,
1110        output_id: SegmentId,
1111        operation: SegmentOperationGuard,
1112        runtime: tokio::runtime::Handle,
1113    ) {
1114        let manager = Arc::clone(self);
1115        let output_hex = output_id.to_hex();
1116        let future = async move {
1117            manager
1118                .delete_output_if_unregistered(output_id, "indexing abort or failure")
1119                .await;
1120            drop(operation);
1121        };
1122        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1123            // The dropped future releases operation ownership. Startup sweep
1124            // handles its output if the runtime is already tearing down.
1125            log::warn!(
1126                "[segment_cleanup] runtime unavailable; indexing output {} will be swept on startup",
1127                output_hex,
1128            );
1129        }
1130    }
1131
1132    /// Claim a newly generated indexing segment before its first file write.
1133    ///
1134    /// The returned guard must travel with the built segment until metadata
1135    /// publication or abort. UUID collisions are treated as corruption rather
1136    /// than silently sharing lifecycle ownership.
1137    pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
1138        match self
1139            .active_operations
1140            .try_register_indexing(vec![segment_id.clone()])
1141        {
1142            Some(operation) => Ok(operation),
1143            None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
1144            None => Err(Error::Corruption(format!(
1145                "new segment ID {} is already owned by an active operation",
1146                segment_id
1147            ))),
1148        }
1149    }
1150
1151    /// Validate the small, mandatory core of a completed segment before it can
1152    /// become metadata-live. Optional vector/sparse/position/fast files are
1153    /// schema- and data-dependent and are validated by `SegmentReader` when used.
1154    async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
1155        let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
1156            Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
1157        })?;
1158        let files = SegmentFiles::new(id.0);
1159
1160        for path in files.mandatory_paths() {
1161            if !self.directory.exists(path).await.map_err(Error::Io)? {
1162                return Err(Error::Corruption(format!(
1163                    "segment {} cannot be published: mandatory file {:?} is missing",
1164                    segment_id, path
1165                )));
1166            }
1167        }
1168
1169        let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
1170            Error::Corruption(format!(
1171                "segment {} cannot be published: missing/unreadable {:?}: {}",
1172                segment_id, files.meta, e
1173            ))
1174        })?;
1175        let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
1176            Error::Corruption(format!(
1177                "segment {} cannot be published: failed reading {:?}: {}",
1178                segment_id, files.meta, e
1179            ))
1180        })?;
1181        let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
1182            Error::Corruption(format!(
1183                "segment {} cannot be published: invalid {:?}: {}",
1184                segment_id, files.meta, e
1185            ))
1186        })?;
1187
1188        if meta.id != id.0 || meta.num_docs != expected_docs {
1189            return Err(Error::Corruption(format!(
1190                "segment {} cannot be published: metadata identity/docs mismatch \
1191                 (id={:032x}, docs={}, expected_docs={})",
1192                segment_id, meta.id, meta.num_docs, expected_docs
1193            )));
1194        }
1195
1196        Ok(())
1197    }
1198
1199    fn quarantine_segment(&self, segment_id: &str, error: &Error) {
1200        let inserted = self
1201            .quarantined_segments
1202            .lock()
1203            .insert(segment_id.to_string());
1204        if inserted {
1205            log::error!(
1206                "[merge] quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
1207                 It remains metadata-live for explicit repair but is excluded from merges until restart",
1208                segment_id,
1209                error,
1210            );
1211        }
1212    }
1213
1214    fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
1215        let mut retry = self.merge_retry.lock();
1216        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1217        let delay = merge_retry_delay(retry.consecutive_failures);
1218        retry.retry_after = std::time::Instant::now().checked_add(delay);
1219        log::warn!(
1220            "[merge] pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
1221            delay.as_secs_f64(),
1222            retry.consecutive_failures,
1223            error,
1224        );
1225        delay
1226    }
1227
1228    fn clear_merge_retry_backoff(&self) {
1229        *self.merge_retry.lock() = MergeRetryState::default();
1230    }
1231
1232    fn merge_retry_is_paused(&self) -> bool {
1233        let mut retry = self.merge_retry.lock();
1234        match retry.retry_after {
1235            Some(deadline) if deadline > std::time::Instant::now() => true,
1236            Some(_) => {
1237                retry.retry_after = None;
1238                false
1239            }
1240            None => false,
1241        }
1242    }
1243
1244    fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
1245        let mut retries = self.reorder_retries.lock();
1246        let retry = retries.entry(segment_id.to_string()).or_default();
1247        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
1248        let delay = merge_retry_delay(retry.consecutive_failures);
1249        retry.retry_after = std::time::Instant::now().checked_add(delay);
1250        log::warn!(
1251            "[reorder] pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
1252            segment_id,
1253            delay.as_secs_f64(),
1254            retry.consecutive_failures,
1255            error,
1256        );
1257    }
1258
1259    fn clear_reorder_retry(&self, segment_id: &str) {
1260        self.reorder_retries.lock().remove(segment_id);
1261    }
1262
1263    fn paused_reorder_segments(&self) -> HashSet<String> {
1264        let now = std::time::Instant::now();
1265        let mut retries = self.reorder_retries.lock();
1266        let mut paused = HashSet::new();
1267        for (segment_id, retry) in retries.iter_mut() {
1268            match retry.retry_after {
1269                Some(deadline) if deadline > now => {
1270                    paused.insert(segment_id.clone());
1271                }
1272                Some(_) => retry.retry_after = None,
1273                None => {}
1274            }
1275        }
1276        paused
1277    }
1278
1279    /// Re-evaluate this index when another index releases application-wide
1280    /// merge capacity. The atomic flag bounds this to one waiter per index and
1281    /// the tracked handle makes index shutdown drain it deterministically.
1282    fn schedule_global_merge_wakeup(self: &Arc<Self>) {
1283        if self
1284            .global_merge_wakeup_pending
1285            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1286            .is_err()
1287        {
1288            return;
1289        }
1290
1291        let manager = Arc::clone(self);
1292        let future = async move {
1293            let capacity = tokio::select! {
1294                biased;
1295                () = manager.active_operations.wait_for_shutdown() => None,
1296                permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
1297            };
1298
1299            manager
1300                .global_merge_wakeup_pending
1301                .store(false, Ordering::Release);
1302            if let Some(permit) = capacity {
1303                // This task is only a notification. The normal scheduler must
1304                // acquire both global and per-index permits atomically enough
1305                // for its own candidate selection.
1306                drop(permit);
1307                manager.maybe_merge().await;
1308            }
1309        };
1310        let runtime = tokio::runtime::Handle::current();
1311        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1312            self.global_merge_wakeup_pending
1313                .store(false, Ordering::Release);
1314            log::warn!("[merge] runtime rejected global-capacity wakeup task");
1315        }
1316    }
1317
1318    #[cfg(test)]
1319    pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
1320        self.quarantined_segments.lock().contains(segment_id)
1321    }
1322
1323    /// Delete a failed output only if metadata did not make it live.
1324    ///
1325    /// Rechecking under `state` also makes unwind cleanup safe if it races
1326    /// successful publication of the same output.
1327    async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
1328        let output_hex = output_id.to_hex();
1329        {
1330            let st = self.state.lock().await;
1331            if st.metadata.has_segment(&output_hex) {
1332                return;
1333            }
1334        }
1335
1336        // UUIDs are generated per producer and cannot be adopted by another
1337        // publisher after this check. Never hold the metadata mutex while a
1338        // multi-GB filesystem deletion runs.
1339        log::info!(
1340            "[segment_cleanup] deleting uncommitted output {} after {}",
1341            output_hex,
1342            reason,
1343        );
1344        if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
1345        {
1346            log::warn!(
1347                "[segment_cleanup] failed deleting uncommitted output {}: {}",
1348                output_hex,
1349                error,
1350            );
1351        }
1352    }
1353
1354    // ========================================================================
1355    // Read path (brief lock or lock-free)
1356    // ========================================================================
1357
1358    /// Get the current segment IDs
1359    pub async fn get_segment_ids(&self) -> Vec<String> {
1360        self.state.lock().await.metadata.segment_ids()
1361    }
1362
1363    /// Get trained vector structures (lock-free via ArcSwap)
1364    pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
1365        self.trained.load_full()
1366    }
1367
1368    /// Capture trained structures for a segment producer.
1369    ///
1370    /// The second gate check closes the race where an update begins after the
1371    /// first check but before the ArcSwap load. Producers have lifecycle guards
1372    /// before calling this method, so an updater that raised the gate waits for
1373    /// any producer that successfully captured the previous generation.
1374    pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
1375        if self.vector_artifact_update.load(Ordering::Acquire) {
1376            return None;
1377        }
1378        let trained = self.trained.load_full();
1379        if self.vector_artifact_update.load(Ordering::Acquire) {
1380            None
1381        } else {
1382            trained
1383        }
1384    }
1385
1386    /// Start an exclusive trained-artifact update and drain merge/reorder
1387    /// producers that may already hold the previous generation.
1388    ///
1389    /// New segment operations may continue while this waits, but they observe
1390    /// the gate through `trained_for_segment_build` and therefore emit flat
1391    /// vector data until the replacement generation is published. The guard
1392    /// is cancellation-safe: dropping the requesting future reopens ANN
1393    /// production without leaving the manager wedged.
1394    ///
1395    /// Indexing tokens cannot be waited on: their guards are parked inside
1396    /// built-but-uncommitted `PreparedSegment`s and are released only by a
1397    /// later commit. That commit typically needs the writer this update's
1398    /// caller already holds (server write lock / embedded `&mut self`), so
1399    /// waiting would permanently wedge vector-index finalization. They also
1400    /// cannot be ignored: such a segment may already contain ANN data bound to
1401    /// the previous artifact generation and could otherwise be committed after
1402    /// the new generation is published. Reject promptly and let the caller
1403    /// commit or abort the pending generation before retrying.
1404    pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
1405        self.vector_artifact_update
1406            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
1407            .map_err(|_| {
1408                Error::Internal("a trained-vector artifact update is already in progress".into())
1409            })?;
1410        self.active_operations.pause_non_indexing();
1411        let guard = VectorArtifactUpdateGuard {
1412            _lease: Arc::new(VectorArtifactUpdateLease {
1413                updating: Arc::clone(&self.vector_artifact_update),
1414                active_operations: Arc::clone(&self.active_operations),
1415            }),
1416        };
1417        let (preexisting, parked_indexing) =
1418            self.active_operations.draining_operation_tokens_snapshot();
1419        if parked_indexing > 0 {
1420            return Err(Error::Internal(format!(
1421                "cannot update trained-vector artifacts while {parked_indexing} indexing \
1422                 segment(s) are built but uncommitted; commit or abort the pending \
1423                 generation and retry"
1424            )));
1425        }
1426        self.active_operations
1427            .wait_until_operations_finish(&preexisting)
1428            .await;
1429        Ok(guard)
1430    }
1431
1432    /// Load trained structures from disk and publish to ArcSwap.
1433    /// Copies metadata under lock, releases lock, then does disk I/O.
1434    pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
1435        // Copy vector_fields under lock (cheap clone of HashMap<u32, FieldMeta>)
1436        let vector_fields = {
1437            let st = self.state.lock().await;
1438            st.metadata.vector_fields.clone()
1439        };
1440        // Disk I/O happens WITHOUT holding the state lock
1441        let trained = IndexMetadata::try_load_trained_from_fields(
1442            &vector_fields,
1443            self.schema.as_ref(),
1444            self.directory.as_ref(),
1445        )
1446        .await?
1447        .map(Arc::new);
1448        // Publish exactly the validated snapshot, including None. Retaining a
1449        // previous map when metadata has no Built fields would let new segments
1450        // depend on artifacts no longer referenced durably.
1451        self.trained.store(trained);
1452        Ok(())
1453    }
1454
1455    /// Atomically publish a fully staged vector generation.
1456    ///
1457    /// Every replacement segment and its global codebook is complete before
1458    /// this transaction starts. Metadata, tracker ownership, and the lock-free
1459    /// trained pointer advance under the same state lock; snapshots therefore
1460    /// observe either the complete old generation or the complete new one.
1461    pub(crate) async fn publish_vector_generation(
1462        self: &Arc<Self>,
1463        artifact_update: &VectorArtifactUpdateGuard,
1464        vector_fields: HashMap<u32, crate::index::FieldVectorMeta>,
1465        next_trained: Arc<TrainedVectorStructures>,
1466        mut staged: Vec<StagedVectorSegment>,
1467    ) -> Result<()> {
1468        if !self.vector_artifact_update.load(Ordering::Acquire) {
1469            return Err(Error::Internal(
1470                "vector generation publication lost its exclusive update lease".into(),
1471            ));
1472        }
1473
1474        for replacement in &staged {
1475            self.validate_completed_segment(&replacement.output_id.to_hex(), replacement.doc_count)
1476                .await?;
1477        }
1478
1479        let mut st = Arc::clone(&self.state).lock_owned().await;
1480        let mut next = st.metadata.clone();
1481        next.vector_fields = vector_fields;
1482        next.refresh_total_vectors();
1483
1484        for replacement in &staged {
1485            let source_info = next
1486                .segment_metas
1487                .remove(&replacement.source_id)
1488                .ok_or_else(|| {
1489                    Error::Corruption(format!(
1490                        "vector generation source {} disappeared before publication",
1491                        replacement.source_id,
1492                    ))
1493                })?;
1494            let output_hex = replacement.output_id.to_hex();
1495            if next.segment_metas.contains_key(&output_hex) {
1496                return Err(Error::Corruption(format!(
1497                    "vector generation output {output_hex} is already metadata-live"
1498                )));
1499            }
1500            // A vector-only rewrite changes neither document order nor merge
1501            // lineage, so preserve the complete lifecycle record verbatim.
1502            next.add_segment_meta(output_hex, source_info);
1503        }
1504
1505        let directory = Arc::clone(&self.directory);
1506        let trained = Arc::clone(&self.trained);
1507        let tracker = Arc::clone(&self.tracker);
1508        let replacement_refresh = self.replacement_refresh.read().clone();
1509        // Keep the producer gate raised if the requesting future is cancelled
1510        // after the metadata transaction has been detached. The last guard
1511        // clone drops only after durable metadata and ArcSwap state agree.
1512        let artifact_update = artifact_update.clone();
1513        self.run_lifecycle_transaction(async move {
1514            let _artifact_update = artifact_update;
1515            next.save(directory.as_ref()).await?;
1516
1517            for replacement in &staged {
1518                tracker.register(&replacement.output_id.to_hex());
1519            }
1520            st.metadata = next;
1521            trained.store(Some(next_trained));
1522
1523            // Outputs are now durably live. Disarm unwind cleanup before
1524            // retiring the old generation and releasing operation ownership.
1525            for replacement in &mut staged {
1526                replacement.cleanup.disarm();
1527            }
1528            let retired = staged
1529                .iter()
1530                .map(|replacement| replacement.source_id.clone())
1531                .collect::<Vec<_>>();
1532            let ready_to_delete = tracker.mark_for_deletion(&retired);
1533            drop(st);
1534            for &segment_id in &ready_to_delete {
1535                if let Err(error) =
1536                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
1537                {
1538                    log::warn!(
1539                        "[segment_cleanup] immediate dense-vector generation delete failed for {}: {}",
1540                        segment_id.to_hex(),
1541                        error,
1542                    );
1543                }
1544            }
1545            tracker.complete_deletion(&ready_to_delete);
1546            refresh_replacement_topology(replacement_refresh).await;
1547            Ok(())
1548        })
1549        .await
1550    }
1551
1552    /// Read metadata with a closure (no persist)
1553    pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1554    where
1555        F: FnOnce(&IndexMetadata) -> R,
1556    {
1557        let st = self.state.lock().await;
1558        f(&st.metadata)
1559    }
1560
1561    /// Update metadata with a closure and persist atomically
1562    pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1563    where
1564        F: FnOnce(&mut IndexMetadata),
1565    {
1566        let mut st = Arc::clone(&self.state).lock_owned().await;
1567        let mut next = st.metadata.clone();
1568        f(&mut next);
1569        let directory = Arc::clone(&self.directory);
1570        self.run_lifecycle_transaction(async move {
1571            next.save(directory.as_ref()).await?;
1572            st.metadata = next;
1573            Ok(())
1574        })
1575        .await
1576    }
1577
1578    /// Acquire a snapshot of current segments for reading.
1579    /// The snapshot holds references — segments won't be deleted while snapshot exists.
1580    pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1581        let (acquired, trained) = {
1582            let st = self.state.lock().await;
1583            let segment_ids = st.metadata.segment_ids();
1584            (self.tracker.acquire(&segment_ids), self.trained.load_full())
1585        };
1586
1587        SegmentSnapshot::with_generation(
1588            Arc::clone(&self.tracker),
1589            acquired,
1590            trained,
1591            Arc::clone(&self.delete_fn),
1592        )
1593    }
1594
1595    /// Get the segment tracker
1596    pub fn tracker(&self) -> Arc<SegmentTracker> {
1597        Arc::clone(&self.tracker)
1598    }
1599
1600    /// Get the directory
1601    pub fn directory(&self) -> Arc<D> {
1602        Arc::clone(&self.directory)
1603    }
1604}
1605
1606// ============================================================================
1607// Native-only: commit, merging, force_merge
1608// ============================================================================
1609
1610#[cfg(feature = "native")]
1611impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1612    /// Atomic commit: register new segments + persist metadata.
1613    pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1614        // Indexing guards still own these IDs here, so the orphan sweeper
1615        // cannot remove files between validation and metadata publication.
1616        for (segment_id, num_docs) in new_segments {
1617            self.validate_completed_segment(segment_id, *num_docs)
1618                .await?;
1619        }
1620
1621        let mut st = Arc::clone(&self.state).lock_owned().await;
1622        let mut next = st.metadata.clone();
1623        let mut added = Vec::new();
1624        for (segment_id, num_docs) in new_segments {
1625            if !next.has_segment(segment_id) {
1626                next.add_segment(segment_id.clone(), *num_docs);
1627                added.push(segment_id.clone());
1628            }
1629        }
1630
1631        // Durable-before-visible: a save failure leaves both in-memory metadata
1632        // and tracker unchanged, so callers can retry the prepared commit.
1633        // The tracked transaction continues if the requesting RPC is cancelled;
1634        // unpublished cleanup waits on this owned state guard before deciding
1635        // whether the files became metadata-live.
1636        let directory = Arc::clone(&self.directory);
1637        let tracker = Arc::clone(&self.tracker);
1638        self.run_lifecycle_transaction(async move {
1639            next.save(directory.as_ref()).await?;
1640            for segment_id in &added {
1641                tracker.register(segment_id);
1642            }
1643            st.metadata = next;
1644            Ok(())
1645        })
1646        .await
1647    }
1648
1649    /// Evaluate merge policy and spawn background merges for all eligible candidates.
1650    ///
1651    /// **Atomicity**: The entire filter → find_merges → spawn_merge sequence runs
1652    /// under the `state` lock to prevent a TOCTOU race where concurrent callers
1653    /// both see segments as eligible before either claims operation ownership.
1654    /// `spawn_merge` is non-blocking (just `try_register` + `tokio::spawn`), so
1655    /// holding the state lock through it is safe and sub-microsecond.
1656    ///
1657    /// The hard merge semaphore is acquired before lifecycle ownership, so
1658    /// concurrent triggers cannot exceed configured merge capacity.
1659    pub async fn maybe_merge(self: &Arc<Self>) {
1660        if !self.active_operations.is_accepting() {
1661            log::debug!("[maybe_merge] manager is shutting down, skipping");
1662            return;
1663        }
1664        if self.merge_retry_is_paused() {
1665            log::debug!("[maybe_merge] retry backoff active, skipping");
1666            return;
1667        }
1668
1669        // Finished handles no longer need to be retained. Concurrency itself
1670        // is enforced by `merge_permits`, not this bookkeeping vector.
1671        {
1672            let mut handles = self.merge_handles.lock();
1673            handles.retain(|h| !h.is_finished());
1674        }
1675        let local_slots = self.merge_permits.available_permits();
1676        let global_slots = self.global_merge_permits.available_permits();
1677        let slots_available = local_slots.min(global_slots);
1678
1679        // Hold state lock through spawn_merge to make filter + register atomic.
1680        // This closes the TOCTOU window where concurrent maybe_merge calls could
1681        // both see the same segments as eligible before either registers them.
1682        {
1683            let st = self.state.lock().await;
1684            let quarantined = self.quarantined_segments.lock().clone();
1685            let active_ids = self.active_operations.snapshot();
1686
1687            // Backlog pressure is based on every metadata-live segment, not
1688            // only currently available inputs. Otherwise a wave of in-flight
1689            // merges hides most of the topology and makes later candidates
1690            // look healthy while the original backlog still exists.
1691            let live_segments: Vec<SegmentInfo> = st
1692                .metadata
1693                .segment_metas
1694                .iter()
1695                .filter(|(id, _)| {
1696                    !self.tracker.is_pending_deletion(id) && !quarantined.contains(*id)
1697                })
1698                .map(|(id, info)| SegmentInfo {
1699                    id: id.clone(),
1700                    num_docs: info.num_docs,
1701                })
1702                .collect();
1703            let severe_backlog = st.merge_policy.has_severe_backlog(&live_segments);
1704
1705            // Exclude segments owned by another operation. Pending retirement
1706            // and quarantined segments were already removed above.
1707            let segments: Vec<SegmentInfo> = live_segments
1708                .iter()
1709                .filter(|segment| !active_ids.contains(&segment.id))
1710                .cloned()
1711                .collect();
1712
1713            log::debug!("[maybe_merge] {} eligible segments", segments.len());
1714
1715            let candidates = st.merge_policy.find_merges(&segments);
1716
1717            if candidates.is_empty() {
1718                return;
1719            }
1720
1721            // Register a capacity waiter only for an index that actually has
1722            // eligible work. Scheduling one waiter for every idle index while
1723            // the process gate was full caused an avoidable wakeup stampede.
1724            if slots_available == 0 {
1725                if local_slots > 0 && global_slots == 0 {
1726                    self.schedule_global_merge_wakeup();
1727                }
1728                log::debug!("[maybe_merge] at max concurrent merges, skipping");
1729                return;
1730            }
1731
1732            log::debug!(
1733                "[maybe_merge] {} merge candidates, {} slots available",
1734                candidates.len(),
1735                slots_available
1736            );
1737
1738            let mut handles = Vec::new();
1739            for c in candidates {
1740                if handles.len() >= slots_available {
1741                    break;
1742                }
1743                // Under severe topology pressure, retire segments first. BP
1744                // remains optional for correctness and the block-copy output
1745                // is explicitly marked unreordered, so the optimizer performs
1746                // one pass after the compaction wave instead of every merge
1747                // task serializing behind the whole-pass gate.
1748                let reorder_bmp = self.reorder_on_merge && !severe_backlog;
1749                if let Some(h) = self.spawn_merge(c.segment_ids, reorder_bmp) {
1750                    handles.push(h);
1751                }
1752            }
1753            if !handles.is_empty() {
1754                if severe_backlog && self.reorder_on_merge {
1755                    log::info!(
1756                        "[maybe_merge] severe backlog: {} live segments; started {} fast \
1757                         block-copy merge(s), deferring BP to the optimizer",
1758                        live_segments.len(),
1759                        handles.len(),
1760                    );
1761                }
1762                // Publish handles before releasing `state`. A force merge
1763                // raises its admission barrier under the same lock, then
1764                // drains this list, so no spawned merge can fall into the
1765                // otherwise-deadlocking gap between spawn and bookkeeping.
1766                self.merge_handles.lock().extend(handles);
1767            }
1768        }
1769    }
1770
1771    /// Spawn a background merge task with RAII tracking.
1772    ///
1773    /// Pre-generates the output segment ID. The operation guard registers all segment IDs
1774    /// (old + output) in `active_operations`. When the task ends (success, failure, or
1775    /// panic), the guard drops and segments are automatically unregistered.
1776    ///
1777    /// On completion, the task auto-triggers `maybe_merge` to evaluate cascading merges.
1778    /// Returns the JoinHandle if the merge was spawned, None if it was skipped.
1779    fn spawn_merge(
1780        self: &Arc<Self>,
1781        segment_ids_to_merge: Vec<String>,
1782        reorder_bmp: bool,
1783    ) -> Option<JoinHandle<()>> {
1784        if self.force_merge_active.load(Ordering::Acquire) > 0 {
1785            log::debug!("[spawn_merge] skipped: explicit force merge has priority");
1786            return None;
1787        }
1788        let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1789            Ok(permit) => permit,
1790            Err(_) => {
1791                log::debug!("[spawn_merge] skipped: global merge capacity is full");
1792                self.schedule_global_merge_wakeup();
1793                return None;
1794            }
1795        };
1796        let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1797            Ok(permit) => permit,
1798            Err(_) => {
1799                log::debug!("[spawn_merge] skipped: no merge permit available");
1800                return None;
1801            }
1802        };
1803        let output_id = SegmentId::new();
1804        let output_hex = output_id.to_hex();
1805
1806        let mut all_ids = segment_ids_to_merge.clone();
1807        all_ids.push(output_hex);
1808
1809        let guard = match self.active_operations.try_register(all_ids) {
1810            Some(g) => g,
1811            None => {
1812                log::debug!("[spawn_merge] skipped: segments overlap with an active operation");
1813                return None;
1814            }
1815        };
1816
1817        let sm = Arc::clone(self);
1818        let ids = segment_ids_to_merge;
1819
1820        Some(tokio::spawn(async move {
1821            let mut reevaluate = false;
1822            let mut retry_delay = None;
1823
1824            let result = sm
1825                .merge_and_replace_registered(
1826                    &ids,
1827                    output_id,
1828                    reorder_bmp,
1829                    ReorderPriority::AutomaticMerge,
1830                )
1831                .await;
1832
1833            match result {
1834                Ok(_) => {
1835                    sm.clear_merge_retry_backoff();
1836                    reevaluate = true;
1837                }
1838                Err(MergeTaskError {
1839                    error: Error::IndexClosed,
1840                    ..
1841                }) => {
1842                    log::debug!(
1843                        "[merge] background merge for segments {:?} cancelled during shutdown",
1844                        ids,
1845                    );
1846                }
1847                Err(MergeTaskError {
1848                    error,
1849                    unavailable_segments,
1850                }) => {
1851                    log::error!(
1852                        "[merge] background merge failed for segments {:?}: {}",
1853                        ids,
1854                        error
1855                    );
1856                    if !unavailable_segments.is_empty() {
1857                        // Recompute without known-bad inputs. This is not a
1858                        // retry of the same candidate because policy filtering
1859                        // excludes every quarantined ID.
1860                        reevaluate = true;
1861                    } else {
1862                        retry_delay = Some(sm.pause_merge_retries(&error));
1863                    }
1864                }
1865            }
1866            // Release source/output ownership before re-evaluating policy, so
1867            // the completed operation cannot artificially hide candidates.
1868            drop(guard);
1869            // A failed merge must not reserve capacity during its retry delay.
1870            drop(merge_permit);
1871            drop(global_merge_permit);
1872
1873            if reevaluate {
1874                sm.maybe_merge().await;
1875            } else if let Some(retry_delay) = retry_delay {
1876                // A backoff without a wakeup can strand eligible segments
1877                // forever when no later commit happens. The sleep runs as a
1878                // tracked *lifecycle* task, not inside this merge JoinHandle:
1879                // wait_for_all_merges/force_merge/reorder drain merge handles,
1880                // and a pure backoff timer with no work in flight must not
1881                // stall them for up to MERGE_RETRY_MAX_DELAY.
1882                sm.schedule_merge_retry_wakeup(retry_delay);
1883            }
1884        }))
1885    }
1886
1887    /// Execute the common build → validate → durable replacement transaction
1888    /// for a batch whose source/output IDs are already lifecycle-owned.
1889    ///
1890    /// Scheduling, retry policy, and post-replacement snapshot refresh remain
1891    /// with the caller. Keeping the transaction here prevents background and
1892    /// force-merge paths from drifting on cleanup or BP metadata semantics.
1893    async fn merge_and_replace_registered(
1894        self: &Arc<Self>,
1895        ids: &[String],
1896        output_id: SegmentId,
1897        reorder_bmp: bool,
1898        priority: ReorderPriority,
1899    ) -> MergeTaskResult<(String, u32, bool)> {
1900        let mut output_cleanup = self.output_cleanup_guard(output_id);
1901        let trained = self.trained_for_segment_build();
1902        let granularity = if reorder_bmp {
1903            self.merge_granularity(ids).await
1904        } else {
1905            crate::segment::reorder::BpGranularity::Auto
1906        };
1907        let result = Self::do_merge(
1908            self.directory.as_ref(),
1909            &self.schema,
1910            ids,
1911            output_id,
1912            self.term_cache_blocks,
1913            trained.as_deref(),
1914            reorder_bmp,
1915            granularity,
1916            self.merge_bp_time_budget,
1917            self.bp_memory_budget_bytes,
1918            Arc::clone(&self.reorder_permits),
1919            priority,
1920            self.active_operations.cancellation_flag(),
1921            Some(self.background_cpu_pool()),
1922        )
1923        .await;
1924
1925        let (new_id, doc_count, bp_converged) = match result {
1926            Ok(value) => value,
1927            Err(error) => {
1928                for segment_id in &error.unavailable_segments {
1929                    self.quarantine_segment(segment_id, &error.error);
1930                }
1931                self.delete_output_if_unregistered(output_id, "merge failure")
1932                    .await;
1933                output_cleanup.disarm();
1934                return Err(error);
1935            }
1936        };
1937
1938        let layout = if reorder_bmp {
1939            ReplacementLayout::BpReordered {
1940                converged: bp_converged,
1941            }
1942        } else {
1943            ReplacementLayout::BlockCopy
1944        };
1945        if let Err(error) = self
1946            .replace_segments(ids, new_id.clone(), doc_count, layout)
1947            .await
1948        {
1949            self.delete_output_if_unregistered(output_id, "replacement failure")
1950                .await;
1951            output_cleanup.disarm();
1952            return Err(MergeTaskError::from(error));
1953        }
1954        output_cleanup.disarm();
1955        Ok((new_id, doc_count, bp_converged))
1956    }
1957
1958    /// Re-evaluate merge policy after a failure backoff, outside the tracked
1959    /// merge JoinHandles that merge waiters drain. Shutdown still drains this
1960    /// task deterministically (lifecycle handles) and interrupts its sleep.
1961    fn schedule_merge_retry_wakeup(self: &Arc<Self>, retry_delay: std::time::Duration) {
1962        let manager = Arc::clone(self);
1963        let future = async move {
1964            tokio::select! {
1965                () = tokio::time::sleep(retry_delay) => {
1966                    manager.maybe_merge().await;
1967                }
1968                () = manager.active_operations.wait_for_shutdown() => {}
1969            }
1970        };
1971        let runtime = tokio::runtime::Handle::current();
1972        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
1973            log::warn!(
1974                "[merge] runtime rejected merge-retry wakeup task; eligible segments may stay \
1975                 unmerged until the next commit re-runs merge policy evaluation"
1976            );
1977        }
1978    }
1979
1980    /// Atomically replace old segments with a new merged segment.
1981    /// Computes merge generation as max(parent gens) + 1 and records ancestors.
1982    /// `reordered` marks whether the new segment was BP-reordered.
1983    async fn replace_segments(
1984        self: &Arc<Self>,
1985        old_ids: &[String],
1986        new_id: String,
1987        doc_count: u32,
1988        layout: ReplacementLayout,
1989    ) -> Result<()> {
1990        // The operation guard owns the output during validation. Publication
1991        // below replaces that ownership with metadata + tracker atomically.
1992        self.validate_completed_segment(&new_id, doc_count).await?;
1993        let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
1994            Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
1995        })?;
1996        let output_reader = SegmentReader::open(
1997            self.directory.as_ref(),
1998            output_id,
1999            Arc::clone(&self.schema),
2000            self.term_cache_blocks,
2001        )
2002        .await
2003        .map_err(|error| match error {
2004            // Preserve retryable storage failures as I/O. Structural failures
2005            // are deterministic for this completed output and get explicit
2006            // corruption context.
2007            Error::Io(_) | Error::IndexClosed => error,
2008            error => Error::Corruption(format!(
2009                "replacement segment {new_id} failed full reader validation: {error}"
2010            )),
2011        })?;
2012        if output_reader.num_docs() != doc_count {
2013            return Err(Error::Corruption(format!(
2014                "replacement segment {new_id} opened with {} docs, expected {doc_count}",
2015                output_reader.num_docs(),
2016            )));
2017        }
2018        drop(output_reader);
2019
2020        let mut st = Arc::clone(&self.state).lock_owned().await;
2021        // Every source must still be live: callers hold operation ownership,
2022        // so a missing source means a stale merge/reorder whose input was
2023        // already replaced. Adding the output would duplicate its documents.
2024        let missing: Vec<&String> = old_ids
2025            .iter()
2026            .filter(|id| !st.metadata.has_segment(id))
2027            .collect();
2028        if !missing.is_empty() {
2029            return Err(Error::Corruption(format!(
2030                "replace_segments: source segment(s) {:?} not in metadata — \
2031                 refusing to add output {} (would duplicate documents)",
2032                missing, new_id
2033            )));
2034        }
2035
2036        let replacement_info = match layout {
2037            ReplacementLayout::BlockCopy | ReplacementLayout::BpReordered { .. } => {
2038                let generation = old_ids
2039                    .iter()
2040                    .filter_map(|id| st.metadata.segment_metas.get(id))
2041                    .map(|info| info.generation)
2042                    .max()
2043                    .unwrap_or(0)
2044                    .checked_add(1)
2045                    .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
2046                let parent_unconverged_passes = old_ids
2047                    .iter()
2048                    .filter_map(|id| st.metadata.segment_metas.get(id))
2049                    .map(|info| info.bp_unconverged_passes)
2050                    .max()
2051                    .unwrap_or(0);
2052                let parent_has_debt = old_ids
2053                    .iter()
2054                    .filter_map(|id| st.metadata.segment_metas.get(id))
2055                    .any(|info| !info.bp_converged);
2056                let (reordered, bp_converged, bp_unconverged_passes) =
2057                    replacement_bp_state(parent_has_debt, parent_unconverged_passes, layout);
2058                SegmentMetaInfo {
2059                    num_docs: doc_count,
2060                    ancestors: old_ids.to_vec(),
2061                    generation,
2062                    reordered,
2063                    bp_converged,
2064                    bp_unconverged_passes,
2065                }
2066            }
2067            ReplacementLayout::PreserveSingleSource => {
2068                let [source_id] = old_ids else {
2069                    return Err(Error::Internal(
2070                        "layout-preserving replacement requires exactly one source".into(),
2071                    ));
2072                };
2073                let mut source = st
2074                    .metadata
2075                    .segment_metas
2076                    .get(source_id)
2077                    .cloned()
2078                    .ok_or_else(|| {
2079                        Error::Corruption(format!(
2080                            "layout-preserving replacement source {source_id} disappeared"
2081                        ))
2082                    })?;
2083                source.num_docs = doc_count;
2084                source
2085            }
2086        };
2087        let retired_ids = old_ids.to_vec();
2088        let mut next = st.metadata.clone();
2089        for id in old_ids {
2090            next.remove_segment(id);
2091        }
2092        next.add_segment_meta(new_id.clone(), replacement_info);
2093
2094        let directory = Arc::clone(&self.directory);
2095        let tracker = Arc::clone(&self.tracker);
2096        let replacement_refresh = self.replacement_refresh.read().clone();
2097        self.run_lifecycle_transaction(async move {
2098            // Durable-before-visible. If persistence fails, old metadata and
2099            // tracker ownership stay intact and source deletion is never armed.
2100            next.save(directory.as_ref()).await?;
2101            tracker.register(&new_id);
2102            st.metadata = next;
2103
2104            // Keep state locked until retired sources enter the tracker. The
2105            // transaction itself also performs deletion, so cancellation of
2106            // the requesting merge cannot strand pending-deletion ownership.
2107            let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
2108            drop(st);
2109            for &segment_id in &ready_to_delete {
2110                if let Err(error) =
2111                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
2112                {
2113                    log::warn!(
2114                        "[segment_cleanup] immediate delete failed for {}: {}",
2115                        segment_id.to_hex(),
2116                        error,
2117                    );
2118                }
2119            }
2120            tracker.complete_deletion(&ready_to_delete);
2121            refresh_replacement_topology(replacement_refresh).await;
2122            Ok(())
2123        })
2124        .await
2125    }
2126
2127    /// Perform the actual merge operation (pure function — no shared state access).
2128    /// `output_segment_id` is pre-generated by the caller so active-operation ownership
2129    /// is installed before any output file is written.
2130    /// Returns (new_segment_id_hex, total_doc_count).
2131    #[allow(clippy::too_many_arguments)]
2132    async fn do_merge(
2133        directory: &D,
2134        schema: &Arc<crate::dsl::Schema>,
2135        segment_ids_to_merge: &[String],
2136        output_segment_id: SegmentId,
2137        term_cache_blocks: usize,
2138        trained: Option<&TrainedVectorStructures>,
2139        reorder_bmp: bool,
2140        granularity: crate::segment::reorder::BpGranularity,
2141        merge_bp_time_budget: Option<std::time::Duration>,
2142        bp_memory_budget_bytes: usize,
2143        reorder_permits: Arc<ReorderConcurrencyGate>,
2144        reorder_priority: ReorderPriority,
2145        cancellation: Arc<AtomicBool>,
2146        bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
2147    ) -> MergeTaskResult<(String, u32, bool)> {
2148        let output_hex = output_segment_id.to_hex();
2149        let load_start = std::time::Instant::now();
2150
2151        let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
2152        for id_str in segment_ids_to_merge {
2153            let id = SegmentId::from_hex(id_str).ok_or_else(|| {
2154                MergeTaskError::source(
2155                    id_str.clone(),
2156                    Error::Corruption(format!("Invalid segment ID: {}", id_str)),
2157                )
2158            })?;
2159            segment_ids.push(id);
2160        }
2161
2162        // Cheap fail-fast before opening every reader. `join_all` otherwise
2163        // waits for all healthy multi-GB inputs to load even when one source's
2164        // `.meta` is already absent, turning a known-corrupt candidate into a
2165        // large CPU/IO spike before it can be quarantined.
2166        let mut unavailable_sources = Vec::new();
2167        let mut missing_files = Vec::new();
2168        for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
2169            let files = SegmentFiles::new(id.0);
2170            let mut source_unavailable = false;
2171            for path in files.mandatory_paths() {
2172                let exists = directory
2173                    .exists(path)
2174                    .await
2175                    .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
2176                if !exists {
2177                    source_unavailable = true;
2178                    missing_files.push(format!("{}:{:?}", id_str, path));
2179                }
2180            }
2181            if source_unavailable {
2182                unavailable_sources.push(id_str.clone());
2183            }
2184        }
2185        if !unavailable_sources.is_empty() {
2186            return Err(MergeTaskError::sources(
2187                unavailable_sources,
2188                Error::Corruption(format!(
2189                    "merge sources are missing mandatory files: {}",
2190                    missing_files.join(", ")
2191                )),
2192            ));
2193        }
2194
2195        let schema_arc = Arc::clone(schema);
2196        let futures: Vec<_> = segment_ids
2197            .iter()
2198            .map(|&sid| {
2199                let sch = Arc::clone(&schema_arc);
2200                async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
2201            })
2202            .collect();
2203
2204        let results = futures::future::join_all(futures).await;
2205        let mut readers = Vec::with_capacity(results.len());
2206        let mut total_docs = 0u64;
2207        for (i, result) in results.into_iter().enumerate() {
2208            match result {
2209                Ok(r) => {
2210                    total_docs += r.meta().num_docs as u64;
2211                    readers.push(r);
2212                }
2213                Err(e) => {
2214                    log::error!(
2215                        "[merge] Failed to open segment {}: {:?}",
2216                        segment_ids_to_merge[i],
2217                        e
2218                    );
2219                    return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
2220                }
2221            }
2222        }
2223        if total_docs > u32::MAX as u64 {
2224            return Err(Error::Internal(format!(
2225                "Merged segment doc count ({}) exceeds u32::MAX",
2226                total_docs
2227            ))
2228            .into());
2229        }
2230
2231        // Pre-merge validation: verify each source segment's store doc count
2232        // matches its metadata. Catching mismatches early avoids building a
2233        // corrupted merged segment and leaving orphan files on disk.
2234        for (i, reader) in readers.iter().enumerate() {
2235            let meta_docs = reader.meta().num_docs;
2236            let store_docs = reader.store().num_docs();
2237            if store_docs != meta_docs {
2238                return Err(MergeTaskError::source(
2239                    segment_ids_to_merge[i].clone(),
2240                    Error::Corruption(format!(
2241                        "pre-merge validation: segment {} store has {} docs but meta says {}",
2242                        segment_ids_to_merge[i], store_docs, meta_docs
2243                    )),
2244                ));
2245            }
2246        }
2247
2248        log::info!(
2249            "[merge] loaded {} segment readers in {:.1}s",
2250            readers.len(),
2251            load_start.elapsed().as_secs_f64()
2252        );
2253
2254        let merger = SegmentMerger::new(Arc::clone(schema))
2255            .with_bmp_reorder(reorder_bmp)
2256            .with_granularity(granularity)
2257            .with_bp_budget(crate::segment::BpBudget {
2258                min_partition_docs: None,
2259                time_budget: merge_bp_time_budget,
2260            })
2261            .with_cancellation(cancellation)
2262            .with_bp_memory_budget(bp_memory_budget_bytes)
2263            .with_reorder_permits(reorder_permits)
2264            .with_reorder_priority(reorder_priority)
2265            .with_background_pool(bg_cpu_pool);
2266
2267        log::info!(
2268            "[merge] {} segments -> {} (trained={})",
2269            segment_ids_to_merge.len(),
2270            output_hex,
2271            trained.map_or(0, |t| t.centroids.len()),
2272        );
2273
2274        let (_merged_meta, merge_stats) = merger
2275            .merge(directory, &readers, output_segment_id, trained)
2276            .await
2277            .map_err(|error| {
2278                if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
2279                    // The merge has already opened every input successfully;
2280                    // a structural/serialization failure is deterministic for
2281                    // this candidate. Attribute all inputs rather than running
2282                    // the same multi-GB rewrite forever. This is deliberately
2283                    // not used for I/O errors, which may be transient/output-side.
2284                    MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
2285                } else {
2286                    MergeTaskError::from(error)
2287                }
2288            })?;
2289        let bp_converged = merge_stats.bp_converged;
2290        if !bp_converged {
2291            log::info!(
2292                "[merge] merge-time BP hit its wall-clock budget — output marked unconverged; \
2293                 the background optimizer deepens it later",
2294            );
2295        }
2296
2297        log::info!(
2298            "[merge] total wall-clock: {:.1}s ({} segments, {} docs)",
2299            load_start.elapsed().as_secs_f64(),
2300            readers.len(),
2301            total_docs,
2302        );
2303
2304        Ok((output_hex, total_docs as u32, bp_converged))
2305    }
2306
2307    /// Drain all in-flight merge tasks safely.
2308    ///
2309    /// Merge/reorder writers use synchronous block-in-place sections, so an
2310    /// abort request takes effect only after the owned writer reaches its next
2311    /// await. Draining the complete task remains necessary before index
2312    /// deletion or orphan cleanup can safely proceed.
2313    ///
2314    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2315    /// handles to `merge_handles` so later drains still see in-flight merges.
2316    pub async fn abort_merges(&self) {
2317        loop {
2318            let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2319            if handles.is_empty() {
2320                return;
2321            }
2322            while let Some(result) = handles.join_next().await {
2323                if let Err(error) = result
2324                    && error.is_panic()
2325                {
2326                    log::error!("[merge] background task panicked while draining: {}", error);
2327                }
2328            }
2329        }
2330    }
2331
2332    /// Wait for all current in-flight merges to complete.
2333    ///
2334    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2335    /// handles to `merge_handles` so later drains still see in-flight merges.
2336    pub async fn wait_for_merging_thread(self: &Arc<Self>) {
2337        let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2338        while handles.join_next().await.is_some() {}
2339    }
2340
2341    /// Wait for all eligible merges to complete, including cascading merges.
2342    ///
2343    /// Drains current handles, then loops. Each completed merge auto-triggers
2344    /// `maybe_merge` (which pushes new handles) before its JoinHandle resolves,
2345    /// so by the time `join_next` returns all cascading handles are registered.
2346    ///
2347    /// Cancellation-safe: dropping this future mid-drain returns un-awaited
2348    /// handles to `merge_handles` so later drains still see in-flight merges.
2349    pub async fn wait_for_all_merges(self: &Arc<Self>) {
2350        loop {
2351            let mut handles = DrainedMergeHandles::take(&self.merge_handles);
2352            if handles.is_empty() {
2353                break;
2354            }
2355            while handles.join_next().await.is_some() {}
2356        }
2357    }
2358
2359    /// Complete the second half of shutdown after the owning `IndexWriter`
2360    /// has been dropped. This drains tracked merges and then waits for every
2361    /// remaining guard, including optimizer reorders that are intentionally
2362    /// launched outside the writer lock.
2363    pub async fn wait_for_shutdown(self: &Arc<Self>) {
2364        self.wait_for_all_merges().await;
2365        self.active_operations.wait_until_idle().await;
2366        loop {
2367            let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
2368            if handles.is_empty() {
2369                break;
2370            }
2371            for handle in handles {
2372                if let Err(error) = handle.await
2373                    && error.is_panic()
2374                {
2375                    log::error!("[segment_cleanup] task panicked while draining: {}", error);
2376                }
2377            }
2378        }
2379    }
2380
2381    /// Force merge segments into compact, size-balanced groups, respecting
2382    /// `max_segment_docs` from the merge policy.
2383    ///
2384    /// If the policy defines a max segment size, segments are merged in batches
2385    /// that stay within that limit. Otherwise, all segments are merged into one.
2386    ///
2387    /// Each batch is registered in `active_operations` via an RAII guard to prevent
2388    /// `maybe_merge` from spawning a conflicting background merge.
2389    pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
2390        self.force_merge_with_snapshot_refresh(|| std::future::ready(Ok(())))
2391            .await
2392    }
2393
2394    /// Force merge while refreshing long-lived segment snapshots after the
2395    /// initial background-merge drain and every durable replacement.
2396    ///
2397    /// Long-lived consumers such as the primary-key index and cached
2398    /// `IndexReader` hold segment snapshots. Refreshing them only after the
2399    /// complete force merge retains every retired source file for the entire
2400    /// operation, which can temporarily double a large index on disk. The
2401    /// hook lets the owning `IndexWriter`/server advance those snapshots after
2402    /// each batch while the force merge remains otherwise memory bounded.
2403    pub(crate) async fn force_merge_with_snapshot_refresh<F, Fut>(
2404        self: &Arc<Self>,
2405        mut refresh_snapshots: F,
2406    ) -> Result<()>
2407    where
2408        F: FnMut() -> Fut,
2409        Fut: std::future::Future<Output = Result<()>>,
2410    {
2411        // Conflicting owners that never appear in `merge_handles` (background
2412        // reorders, a concurrent force-merge) can hold a batch segment for
2413        // minutes to hours; retrying without parking would busy-spin a runtime
2414        // worker and hammer the state mutex for that whole window.
2415        const FORCE_MERGE_CONFLICT_BACKOFF: std::time::Duration =
2416            std::time::Duration::from_millis(100);
2417        // When every remaining mergeable segment is owned by another
2418        // operation (e.g. background BP reorders), there is nothing to do but
2419        // wait for a release; those passes run for minutes, so poll slowly.
2420        const FORCE_MERGE_HELD_BACKOFF: std::time::Duration = std::time::Duration::from_secs(1);
2421
2422        let (_force_merge_activity, max_segment_docs) = {
2423            let st = self.state.lock().await;
2424            // `maybe_merge` registers and publishes every task handle under
2425            // this same state lock. Raising the barrier here therefore makes
2426            // the subsequent drain race-free.
2427            self.force_merge_active.fetch_add(1, Ordering::AcqRel);
2428            (
2429                ForceMergeActivityGuard(&self.force_merge_active),
2430                st.merge_policy.max_segment_docs(),
2431            )
2432        };
2433
2434        // Wait for all in-flight background merges (including cascading)
2435        // before starting forced merges to avoid try_register conflicts.
2436        let background_merges = self
2437            .merge_handles
2438            .lock()
2439            .iter()
2440            .filter(|handle| !handle.is_finished())
2441            .count();
2442        if background_merges > 0 {
2443            log::info!(
2444                "[force_merge] waiting for {} in-flight background merge(s) before planning",
2445                background_merges,
2446            );
2447        }
2448        let drain_start = std::time::Instant::now();
2449        self.wait_for_all_merges().await;
2450        if drain_start.elapsed() >= std::time::Duration::from_secs(1) {
2451            log::info!(
2452                "[force_merge] drained background merges in {:.1}s",
2453                drain_start.elapsed().as_secs_f64(),
2454            );
2455        }
2456
2457        // A background merge may have published after the writer's preceding
2458        // commit refresh but before the drain completed. Advance PK/search
2459        // snapshots now so its retired sources do not survive until the first
2460        // forced replacement (or forever when no batch is needed).
2461        let refresh_start = std::time::Instant::now();
2462        refresh_snapshots().await?;
2463        if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2464            log::info!(
2465                "[force_merge] initial snapshot refresh took {:.1}s",
2466                refresh_start.elapsed().as_secs_f64(),
2467            );
2468        }
2469
2470        // Outputs produced by a completed bin are already maximal according
2471        // to the BFD plan. Excluding them from later replanning prevents an
2472        // expensive final BP output from being selected and reordered again.
2473        let mut completed_outputs = HashSet::new();
2474        // One INFO line per wait episode, not per 1s poll; DEBUG afterwards.
2475        let mut logged_held_wait = false;
2476
2477        loop {
2478            if !self.active_operations.is_accepting() {
2479                return Err(Error::IndexClosed);
2480            }
2481
2482            let segments: Vec<(String, u32)> = {
2483                let st = self.state.lock().await;
2484                st.metadata
2485                    .segment_metas
2486                    .iter()
2487                    .filter(|(id, _)| !completed_outputs.contains(*id))
2488                    .map(|(id, info)| (id.clone(), info.num_docs))
2489                    .collect()
2490            };
2491
2492            // Route around segments owned by active operations instead of
2493            // retrying the same conflicting batch. Group ownership is claimed
2494            // before foreground capacity, so an in-flight optimizer/retrain
2495            // can finish without a capacity/ownership lock inversion.
2496            let active_ids = self.active_operations.snapshot();
2497            let held = segments
2498                .iter()
2499                .filter(|(id, _)| active_ids.contains(id))
2500                .count();
2501            let free_segments: Vec<_> = segments
2502                .into_iter()
2503                .filter(|(id, _)| !active_ids.contains(id))
2504                .collect();
2505            // Every segment format and doc map uses u32 document IDs. A policy
2506            // with no explicit cap therefore means "up to the format limit",
2507            // not an unrepresentable u64-sized output.
2508            let max_docs = max_segment_docs
2509                .map(u64::from)
2510                .unwrap_or(u64::from(u32::MAX));
2511            let next_group = plan_force_merge_groups(free_segments, max_docs)
2512                .into_iter()
2513                .find(|group| group.segments.len() >= 2);
2514
2515            let Some(group) = next_group else {
2516                if held == 0 {
2517                    if !completed_outputs.is_empty() {
2518                        // A segment held while earlier groups ran may now fit
2519                        // with one of those outputs. Reconsider completed bins
2520                        // once before declaring convergence. In the normal
2521                        // no-conflict path BFD groups are pairwise maximal, so
2522                        // this extra planning round performs no merge.
2523                        completed_outputs.clear();
2524                        continue;
2525                    }
2526                    // Every remaining free segment is either already at the
2527                    // size cap or cannot be paired without exceeding it.
2528                    // A reorder that was already running when foreground
2529                    // admission began may have published after the preceding
2530                    // callback. Reconcile external readers one final time so
2531                    // they cannot retain its retired source indefinitely.
2532                    refresh_snapshots().await?;
2533                    return Ok(());
2534                }
2535                if !logged_held_wait {
2536                    log::info!(
2537                        "[force_merge] waiting: {} segment(s) held by active \
2538                         merge/reorder operations, no free group can merge",
2539                        held
2540                    );
2541                    logged_held_wait = true;
2542                } else {
2543                    log::debug!("[force_merge] still waiting on {} held segment(s)", held);
2544                }
2545                #[cfg(test)]
2546                self.force_merge_conflict_retries
2547                    .fetch_add(1, Ordering::Relaxed);
2548                tokio::select! {
2549                    biased;
2550                    () = self.active_operations.wait_for_shutdown() => {
2551                        return Err(Error::IndexClosed);
2552                    }
2553                    () = tokio::time::sleep(FORCE_MERGE_HELD_BACKOFF) => {}
2554                }
2555                continue;
2556            };
2557            logged_held_wait = false;
2558
2559            // Claim the complete final group and all of its future output IDs
2560            // up front. This freezes the hierarchy while allowing every
2561            // durable intermediate replacement to release its source files.
2562            let hierarchy = plan_force_merge_hierarchy(group.segments.len());
2563            let output_ids: Vec<_> = (0..hierarchy.steps.len())
2564                .map(|_| SegmentId::new())
2565                .collect();
2566            let source_ids: Vec<_> = group.segments.iter().map(|(id, _)| id.clone()).collect();
2567            let mut all_ids = source_ids.clone();
2568            all_ids.extend(output_ids.iter().map(|id| id.to_hex()));
2569            let group_guard = {
2570                let st = self.state.lock().await;
2571                source_ids
2572                    .iter()
2573                    .all(|id| st.metadata.has_segment(id))
2574                    .then(|| self.active_operations.try_register(all_ids))
2575                    .flatten()
2576            };
2577            let _group_guard = match group_guard {
2578                Some(guard) => guard,
2579                None if !self.active_operations.is_accepting() => {
2580                    return Err(Error::IndexClosed);
2581                }
2582                None => {
2583                    #[cfg(test)]
2584                    self.force_merge_conflict_retries
2585                        .fetch_add(1, Ordering::Relaxed);
2586                    log::debug!("[force_merge] group lost a registration race, replanning");
2587                    let had_tracked_merges = !self.merge_handles.lock().is_empty();
2588                    self.wait_for_merging_thread().await;
2589                    if !had_tracked_merges {
2590                        tokio::time::sleep(FORCE_MERGE_CONFLICT_BACKOFF).await;
2591                    }
2592                    continue;
2593                }
2594            };
2595
2596            log::info!(
2597                "[force_merge] planned final group: {} segments, {} docs, {} merge pass(es)",
2598                group.segments.len(),
2599                group.total_docs,
2600                output_ids.len(),
2601            );
2602
2603            // Claiming the complete group before any capacity wait prevents a
2604            // vector-generation pause from waiting for a global slot retained
2605            // by force merge while force merge waits for that pause to end.
2606            //
2607            // Background merges acquire global merge capacity before the
2608            // shared BP gate. Preserve that order here: wait for one group
2609            // slot while background BP remains admitted, then pause new BP
2610            // passes. The retained slot guarantees this hierarchy can make
2611            // progress even if other indexes subsequently park at the gate.
2612            let group_global_merge_permit = if self.reorder_on_merge {
2613                let capacity_start = std::time::Instant::now();
2614                let permit = tokio::select! {
2615                    biased;
2616                    () = self.active_operations.wait_for_shutdown() => {
2617                        return Err(Error::IndexClosed);
2618                    }
2619                    permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2620                        permit.map_err(|_| {
2621                            Error::Internal(
2622                                "global background merge scheduler is closed".into(),
2623                            )
2624                        })?
2625                    }
2626                };
2627                if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2628                    log::info!(
2629                        "[force_merge] waited {:.1}s for foreground global merge capacity",
2630                        capacity_start.elapsed().as_secs_f64(),
2631                    );
2632                }
2633                Some(permit)
2634            } else {
2635                None
2636            };
2637            let _foreground_reorder = if self.reorder_on_merge {
2638                log::info!(
2639                    "[force_merge] prioritizing BP capacity ({} total pass slot(s))",
2640                    self.reorder_permits.limit(),
2641                );
2642                let admission_start = std::time::Instant::now();
2643                let guard = Arc::clone(&self.reorder_permits)
2644                    .begin_foreground()
2645                    .await
2646                    .map_err(|_| {
2647                        Error::Internal("background reorder scheduler is closed".into())
2648                    })?;
2649                if admission_start.elapsed() >= std::time::Duration::from_secs(1) {
2650                    log::info!(
2651                        "[force_merge] acquired foreground BP capacity in {:.1}s",
2652                        admission_start.elapsed().as_secs_f64(),
2653                    );
2654                }
2655                Some(guard)
2656            } else {
2657                None
2658            };
2659
2660            let source_count = group.segments.len();
2661            let mut nodes: Vec<Option<(String, u32)>> =
2662                group.segments.into_iter().map(Some).collect();
2663            nodes.resize_with(source_count + hierarchy.steps.len(), || None);
2664            for (step_index, step) in hierarchy.steps.iter().enumerate() {
2665                let final_pass = step_index + 1 == hierarchy.steps.len();
2666                let mut batch_entries = Vec::with_capacity(step.inputs.len());
2667                for &node in &step.inputs {
2668                    let entry = nodes
2669                        .get_mut(node)
2670                        .and_then(Option::take)
2671                        .expect("force-merge hierarchy must reference an available node");
2672                    batch_entries.push(entry);
2673                }
2674                let batch: Vec<_> = batch_entries.iter().map(|(id, _)| id.clone()).collect();
2675                let batch_docs: u64 = batch_entries.iter().map(|(_, docs)| u64::from(*docs)).sum();
2676                let output_id = output_ids[step_index];
2677
2678                let capacity_start = std::time::Instant::now();
2679                let step_global_merge_permit = if group_global_merge_permit.is_none() {
2680                    Some(tokio::select! {
2681                        biased;
2682                        () = self.active_operations.wait_for_shutdown() => {
2683                            return Err(Error::IndexClosed);
2684                        }
2685                        permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2686                            permit.map_err(|_| {
2687                                Error::Internal(
2688                                    "global background merge scheduler is closed".into(),
2689                                )
2690                            })?
2691                        }
2692                    })
2693                } else {
2694                    None
2695                };
2696                if capacity_start.elapsed() >= std::time::Duration::from_secs(1) {
2697                    log::info!(
2698                        "[force_merge] waited {:.1}s for global merge capacity",
2699                        capacity_start.elapsed().as_secs_f64(),
2700                    );
2701                }
2702
2703                // Intermediate reductions are streaming block-copy merges.
2704                // Only the final pass pays BP, so a group with hundreds of
2705                // tiny sources never reorders the same documents repeatedly.
2706                let reorder_bmp = final_pass && self.reorder_on_merge;
2707                log::info!(
2708                    "[force_merge] {} pass: {} segments ({} docs, bp={})",
2709                    if final_pass {
2710                        "final"
2711                    } else {
2712                        "fan-in reduction"
2713                    },
2714                    batch.len(),
2715                    batch_docs,
2716                    reorder_bmp,
2717                );
2718                let (new_segment_id, total_docs, _) = self
2719                    .merge_and_replace_registered(
2720                        &batch,
2721                        output_id,
2722                        reorder_bmp,
2723                        ReorderPriority::Foreground,
2724                    )
2725                    .await
2726                    .map_err(|error| error.error)?;
2727                drop(step_global_merge_permit);
2728
2729                // Advance PK/read snapshots after every hierarchy level so
2730                // retired sources do not accumulate until the final output.
2731                let refresh_start = std::time::Instant::now();
2732                refresh_snapshots().await?;
2733                if refresh_start.elapsed() >= std::time::Duration::from_secs(1) {
2734                    log::info!(
2735                        "[force_merge] post-replacement snapshot refresh took {:.1}s",
2736                        refresh_start.elapsed().as_secs_f64(),
2737                    );
2738                }
2739
2740                let output_node = source_count + step_index;
2741                debug_assert!(nodes[output_node].is_none());
2742                nodes[output_node] = Some((new_segment_id, total_docs));
2743            }
2744            let (root_id, _) = nodes[hierarchy.root]
2745                .take()
2746                .expect("force-merge hierarchy must produce its root");
2747            debug_assert!(nodes.into_iter().all(|node| node.is_none()));
2748            completed_outputs.insert(root_id);
2749        }
2750    }
2751
2752    fn segment_needs_vector_rewrite(
2753        &self,
2754        reader: &SegmentReader,
2755        field_ids: &[u32],
2756        rewrite_existing: bool,
2757    ) -> Result<bool> {
2758        for &field_id in field_ids {
2759            let flat = reader.flat_vectors().get(&field_id);
2760            let ann = reader.vector_indexes().get(&field_id);
2761            if ann.is_some() && flat.is_none() {
2762                return Err(Error::Corruption(format!(
2763                    "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2764                    reader.meta().id,
2765                )));
2766            }
2767
2768            let Some(flat) = flat else {
2769                continue;
2770            };
2771            if flat.num_vectors == 0 {
2772                continue;
2773            }
2774            if rewrite_existing {
2775                return Ok(true);
2776            }
2777            let field = crate::dsl::Field(field_id);
2778            let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2779                Error::Corruption(format!(
2780                    "segment {:032x} references unknown vector field {field_id}",
2781                    reader.meta().id,
2782                ))
2783            })?;
2784            let current = match entry.field_type {
2785                // TQ payloads carry no trained generation, so an existing one
2786                // is always current; a tq field must never be staged for a
2787                // vector-generation rewrite.
2788                crate::dsl::FieldType::DenseVector
2789                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
2790                        config.index_type == crate::dsl::VectorIndexType::Tq
2791                    }) =>
2792                {
2793                    matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
2794                }
2795                crate::dsl::FieldType::DenseVector
2796                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
2797                        config.index_type == crate::dsl::VectorIndexType::IvfTq
2798                    }) =>
2799                {
2800                    matches!(ann, Some(crate::segment::VectorIndex::IvfTq { .. }))
2801                }
2802                // Only ivf_tq dense fields are trainable; other dense
2803                // index types never reach a vector-generation rewrite.
2804                crate::dsl::FieldType::DenseVector => false,
2805                crate::dsl::FieldType::BinaryDenseVector => {
2806                    matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2807                }
2808                _ => false,
2809            };
2810            if !current {
2811                return Ok(true);
2812            }
2813        }
2814        Ok(false)
2815    }
2816
2817    async fn acquire_vector_rewrite_capacity(
2818        &self,
2819    ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2820        let global = tokio::select! {
2821            biased;
2822            () = self.active_operations.wait_for_shutdown() => {
2823                return Err(Error::IndexClosed);
2824            }
2825            permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2826                permit.map_err(|_| Error::Internal(
2827                    "global background merge scheduler is closed".into()
2828                ))?
2829            }
2830        };
2831        let local = tokio::select! {
2832            biased;
2833            () = self.active_operations.wait_for_shutdown() => {
2834                return Err(Error::IndexClosed);
2835            }
2836            permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2837                permit.map_err(|_| Error::Internal(
2838                    "background merge scheduler is closed".into()
2839                ))?
2840            }
2841        };
2842        Ok((global, local))
2843    }
2844
2845    async fn build_vector_replacement(
2846        self: &Arc<Self>,
2847        segment_id: &str,
2848        source_id: SegmentId,
2849        output_id: SegmentId,
2850        trained: &TrainedVectorStructures,
2851        failure_context: &'static str,
2852    ) -> Result<(String, u32, OutputCleanupGuard)> {
2853        let mut cleanup = self.output_cleanup_guard(output_id);
2854        match crate::segment::reorder::rewrite_vector_segment(
2855            self.directory.as_ref(),
2856            &self.schema,
2857            source_id,
2858            output_id,
2859            self.term_cache_blocks,
2860            trained,
2861            Some(self.background_cpu_pool()),
2862        )
2863        .await
2864        {
2865            Ok((new_id, doc_count)) => {
2866                self.validate_completed_segment(&new_id, doc_count).await?;
2867                Ok((new_id, doc_count, cleanup))
2868            }
2869            Err(error) => {
2870                self.delete_output_if_unregistered(output_id, failure_context)
2871                    .await;
2872                cleanup.disarm();
2873                if is_deterministic_source_error(&error) {
2874                    self.quarantine_segment(segment_id, &error);
2875                }
2876                Err(error)
2877            }
2878        }
2879    }
2880
2881    /// Build every required replacement segment without exposing any of them.
2882    /// The returned guards keep both source and output generations alive until
2883    /// [`Self::publish_vector_generation`] commits the complete set.
2884    pub(crate) async fn stage_vector_generation(
2885        self: &Arc<Self>,
2886        _artifact_update: &VectorArtifactUpdateGuard,
2887        segment_ids: &[String],
2888        field_ids: &[u32],
2889        trained: Arc<TrainedVectorStructures>,
2890        rewrite_existing: bool,
2891    ) -> Result<Vec<StagedVectorSegment>> {
2892        if !self.vector_artifact_update.load(Ordering::Acquire) {
2893            return Err(Error::Internal(
2894                "cannot stage a vector generation without an exclusive update lease".into(),
2895            ));
2896        }
2897
2898        let mut staged = Vec::new();
2899        for segment_id in segment_ids {
2900            if self.quarantined_segments.lock().contains(segment_id) {
2901                return Err(Error::Corruption(format!(
2902                    "segment {segment_id} is quarantined after a deterministic source failure"
2903                )));
2904            }
2905            let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2906                Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2907            })?;
2908
2909            // A single rewrite may hold several gigabytes while assigning all
2910            // vectors. Reuse the ordinary local and process-wide merge bounds.
2911            let _capacity = self.acquire_vector_rewrite_capacity().await?;
2912
2913            let output_id = SegmentId::new();
2914            let output_hex = output_id.to_hex();
2915            let operation = {
2916                let st = self.state.lock().await;
2917                if !st.metadata.has_segment(segment_id) {
2918                    return Err(Error::Corruption(format!(
2919                        "vector generation source {segment_id} disappeared while lifecycle work was paused"
2920                    )));
2921                }
2922                self.active_operations
2923                    .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
2924            }
2925            .ok_or_else(|| {
2926                if self.active_operations.is_accepting() {
2927                    Error::Internal(format!(
2928                        "vector generation could not claim stable source {segment_id}"
2929                    ))
2930                } else {
2931                    Error::IndexClosed
2932                }
2933            })?;
2934
2935            let reader = SegmentReader::open(
2936                self.directory.as_ref(),
2937                source_id,
2938                Arc::clone(&self.schema),
2939                self.term_cache_blocks,
2940            )
2941            .await?;
2942            if !self.segment_needs_vector_rewrite(&reader, field_ids, rewrite_existing)? {
2943                continue;
2944            }
2945            drop(reader);
2946
2947            let (new_id, doc_count, cleanup) = self
2948                .build_vector_replacement(
2949                    segment_id,
2950                    source_id,
2951                    output_id,
2952                    trained.as_ref(),
2953                    "vector generation staging failure",
2954                )
2955                .await?;
2956            debug_assert_eq!(new_id, output_hex);
2957            let output_reader = SegmentReader::open(
2958                self.directory.as_ref(),
2959                output_id,
2960                Arc::clone(&self.schema),
2961                self.term_cache_blocks,
2962            )
2963            .await?;
2964            if self.segment_needs_vector_rewrite(&output_reader, field_ids, false)? {
2965                return Err(Error::Corruption(format!(
2966                    "staged vector segment {new_id} does not match its candidate codebook generation"
2967                )));
2968            }
2969
2970            staged.push(StagedVectorSegment {
2971                source_id: segment_id.clone(),
2972                output_id,
2973                doc_count,
2974                _operation: operation,
2975                cleanup,
2976            });
2977        }
2978        Ok(staged)
2979    }
2980
2981    async fn rewrite_vector_segment_once(
2982        self: &Arc<Self>,
2983        segment_id: &str,
2984        field_ids: &[u32],
2985    ) -> Result<VectorSegmentRewriteOutcome> {
2986        if self.quarantined_segments.lock().contains(segment_id) {
2987            return Err(Error::Corruption(format!(
2988                "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
2989            )));
2990        }
2991        let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2992            Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2993        })?;
2994
2995        // Match ordinary merge lock ordering: capacity before lifecycle
2996        // ownership. A vector rewrite can hold several gigabytes while it
2997        // assigns vectors, so it participates in both local and process-wide
2998        // merge limits.
2999        let _capacity = self.acquire_vector_rewrite_capacity().await?;
3000
3001        let output_id = SegmentId::new();
3002        let output_hex = output_id.to_hex();
3003        let all_ids = vec![segment_id.to_owned(), output_hex];
3004        let operation = {
3005            let st = self.state.lock().await;
3006            if !st.metadata.has_segment(segment_id) {
3007                return Ok(VectorSegmentRewriteOutcome::SourceGone);
3008            }
3009            self.active_operations.try_register(all_ids)
3010        };
3011        let _operation = match operation {
3012            Some(operation) => operation,
3013            None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3014            None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3015        };
3016
3017        let Some(trained) = self.trained_for_segment_build() else {
3018            return Ok(VectorSegmentRewriteOutcome::Deferred);
3019        };
3020
3021        let reader = SegmentReader::open(
3022            self.directory.as_ref(),
3023            source_id,
3024            Arc::clone(&self.schema),
3025            self.term_cache_blocks,
3026        )
3027        .await?;
3028        if !self.segment_needs_vector_rewrite(&reader, field_ids, false)? {
3029            return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3030        }
3031        drop(reader);
3032
3033        let (new_id, doc_count, mut output_cleanup) = self
3034            .build_vector_replacement(
3035                segment_id,
3036                source_id,
3037                output_id,
3038                trained.as_ref(),
3039                "vector rewrite failure",
3040            )
3041            .await?;
3042
3043        if let Err(error) = self
3044            .replace_segments(
3045                &[segment_id.to_owned()],
3046                new_id,
3047                doc_count,
3048                ReplacementLayout::PreserveSingleSource,
3049            )
3050            .await
3051        {
3052            self.delete_output_if_unregistered(output_id, "vector replacement failure")
3053                .await;
3054            output_cleanup.disarm();
3055            return Err(error);
3056        }
3057        output_cleanup.disarm();
3058        Ok(VectorSegmentRewriteOutcome::Rewritten)
3059    }
3060
3061    /// Finalize every committed flat vector segment against the published
3062    /// global ANN generation.
3063    /// Unlike force-merge this handles one segment and segments already at the
3064    /// merge policy's maximum size.
3065    pub(crate) async fn rewrite_vector_segments(
3066        self: &Arc<Self>,
3067        field_ids: &[u32],
3068    ) -> Result<usize> {
3069        if field_ids.is_empty() {
3070            return Ok(0);
3071        }
3072        let mut rewritten = 0usize;
3073        loop {
3074            let segment_ids = self.get_segment_ids().await;
3075            let mut conflicted = false;
3076            let mut changed = false;
3077            for segment_id in segment_ids {
3078                match self
3079                    .rewrite_vector_segment_once(&segment_id, field_ids)
3080                    .await?
3081                {
3082                    VectorSegmentRewriteOutcome::Rewritten => {
3083                        rewritten += 1;
3084                        changed = true;
3085                    }
3086                    VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3087                    VectorSegmentRewriteOutcome::Deferred => {
3088                        return Err(Error::Internal(
3089                            "ANN finalization lost the published trained generation".into(),
3090                        ));
3091                    }
3092                    VectorSegmentRewriteOutcome::AlreadyCurrent
3093                    | VectorSegmentRewriteOutcome::SourceGone => {}
3094                }
3095            }
3096            if !conflicted && !changed {
3097                log::info!(
3098                    "[dense_vector_rewrite] ANN finalization complete ({} segment(s) rewritten)",
3099                    rewritten,
3100                );
3101                return Ok(rewritten);
3102            }
3103            tokio::select! {
3104                biased;
3105                () = self.active_operations.wait_for_shutdown() => {
3106                    return Err(Error::IndexClosed);
3107                }
3108                () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3109            }
3110        }
3111    }
3112
3113    /// A producer that started in the force-flat phase can commit after the
3114    /// main finalization snapshot. Upgrade exactly those new segments in a
3115    /// tracked background task; ordinary producers already using the current
3116    /// generation are detected and skipped without rewriting.
3117    pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3118        if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3119            return;
3120        }
3121        let manager = Arc::clone(self);
3122        let future = async move {
3123            let field_ids = manager
3124                .read_metadata(|metadata| {
3125                    metadata
3126                        .vector_fields
3127                        .keys()
3128                        .filter(|field_id| metadata.is_field_built(**field_id))
3129                        .copied()
3130                        .collect::<Vec<_>>()
3131                })
3132                .await;
3133            for segment_id in segment_ids {
3134                loop {
3135                    match manager
3136                        .rewrite_vector_segment_once(&segment_id, &field_ids)
3137                        .await
3138                    {
3139                        Ok(VectorSegmentRewriteOutcome::Conflict) => {
3140                            tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3141                        }
3142                        Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3143                        Ok(_) => break,
3144                        Err(error) => {
3145                            log::error!(
3146                                "[dense_vector_rewrite] failed to upgrade newly committed segment {}: {}",
3147                                segment_id,
3148                                error,
3149                            );
3150                            break;
3151                        }
3152                    }
3153                }
3154            }
3155        };
3156        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3157            log::warn!(
3158                "[dense_vector_rewrite] runtime unavailable; newly committed flat segment upgrade deferred"
3159            );
3160            return;
3161        };
3162        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3163            log::warn!(
3164                "[dense_vector_rewrite] runtime rejected newly committed flat segment upgrade"
3165            );
3166        }
3167    }
3168
3169    /// Reorder all segments via Recursive Graph Bisection (BP) for better BMP pruning.
3170    ///
3171    /// Each segment is individually rebuilt with reordered BMP blocks.
3172    /// Non-BMP fields are copied unchanged via streaming file copy.
3173    ///
3174    /// Uses active-operation ownership to prevent concurrent work on the same segment.
3175    pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3176        self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3177            .await
3178    }
3179
3180    /// Reorder all segments while advancing long-lived snapshots after the
3181    /// background-merge drain and every durable replacement.
3182    pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3183        self: &Arc<Self>,
3184        mut refresh_snapshots: F,
3185    ) -> Result<()>
3186    where
3187        F: FnMut() -> Fut,
3188        Fut: std::future::Future<Output = Result<()>>,
3189    {
3190        self.wait_for_all_merges().await;
3191        refresh_snapshots().await?;
3192        let segment_ids = self.get_segment_ids().await;
3193
3194        if segment_ids.is_empty() {
3195            log::info!("[reorder] no segments to reorder");
3196            return Ok(());
3197        }
3198
3199        log::info!("[reorder] reordering {} segments", segment_ids.len());
3200
3201        for seg_id in segment_ids {
3202            match self
3203                .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3204                .await
3205            {
3206                Ok(true) => refresh_snapshots().await?,
3207                Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
3208                Err(e) => return Err(e),
3209            }
3210        }
3211
3212        // A segment skipped because another lifecycle owner held it may have
3213        // been replaced after the preceding callback.
3214        refresh_snapshots().await?;
3215        log::info!("[reorder] all segments reordered");
3216        Ok(())
3217    }
3218
3219    /// Get segment IDs that have not been reordered yet.
3220    ///
3221    /// Excludes segments currently involved in a merge or reorder operation
3222    /// to avoid wasted work (the optimizer would skip them anyway).
3223    pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3224        self.unreordered_segments()
3225            .await
3226            .into_iter()
3227            .map(|(id, _)| id)
3228            .collect()
3229    }
3230
3231    /// Segments never reordered, with doc counts — for the optimizer to pick
3232    /// a size-appropriate BP budget.
3233    pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
3234        let quarantined = self.quarantined_segments.lock().clone();
3235        let paused = self.paused_reorder_segments();
3236        let st = self.state.lock().await;
3237        let active_ids = self.active_operations.snapshot();
3238        st.metadata
3239            .segment_metas
3240            .iter()
3241            .filter(|(id, info)| {
3242                !info.reordered
3243                    && info.bp_converged
3244                    && !active_ids.contains(*id)
3245                    && !quarantined.contains(*id)
3246                    && !paused.contains(*id)
3247            })
3248            .map(|(id, info)| (id.clone(), info.num_docs))
3249            .collect()
3250    }
3251
3252    /// Segments whose last BP pass hit its wall-clock budget before finishing
3253    /// (`bp_converged == false`). A warm-started follow-up pass deepens the
3254    /// ordering; the optimizer revisits these at low priority.
3255    pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3256        self.unconverged_segments_below(u32::MAX)
3257            .await
3258            .into_iter()
3259            .map(|(id, docs, _)| (id, docs))
3260            .collect()
3261    }
3262
3263    /// Unconverged segments still below a hard replacement-lineage work
3264    /// bound. Includes the persisted attempt count for scheduler diagnostics.
3265    pub async fn unconverged_segments_below(
3266        &self,
3267        max_unconverged_passes: u32,
3268    ) -> Vec<(String, u32, u32)> {
3269        let quarantined = self.quarantined_segments.lock().clone();
3270        let paused = self.paused_reorder_segments();
3271        let st = self.state.lock().await;
3272        let active_ids = self.active_operations.snapshot();
3273        st.metadata
3274            .segment_metas
3275            .iter()
3276            .filter(|(id, info)| {
3277                !info.bp_converged
3278                    && info.bp_unconverged_passes < max_unconverged_passes
3279                    && !active_ids.contains(*id)
3280                    && !quarantined.contains(*id)
3281                    && !paused.contains(*id)
3282            })
3283            .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3284            .collect()
3285    }
3286
3287    /// Granularity for a BP pass whose sources are `ids`: `Records` when any
3288    /// source carries unconverged BP debt, `Auto` otherwise.
3289    ///
3290    /// Alignment with the depth budget (docs/block-level-reorder.md): an
3291    /// unconverged segment is owed a deepening pass, and the output of this
3292    /// pass will be marked `bp_converged`. `Auto` would measure the partial
3293    /// pass's residual coherence, potentially take the blockwise path — which
3294    /// cannot deepen record clustering — and end the cascade at partial
3295    /// quality. Only record-level BP discharges the debt.
3296    async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3297        let st = self.state.lock().await;
3298        let deepening = ids.iter().any(|id| {
3299            st.metadata
3300                .segment_metas
3301                .get(id)
3302                .is_some_and(|info| !info.bp_converged)
3303        });
3304        drop(st);
3305        if deepening {
3306            log::info!(
3307                "[reorder] source BP lineage unconverged — forcing record-level BP (deepening pass)",
3308            );
3309            crate::segment::reorder::BpGranularity::Records
3310        } else {
3311            crate::segment::reorder::BpGranularity::Auto
3312        }
3313    }
3314
3315    /// Reorder a single segment via BP. Returns Ok(true) if reordered, Ok(false) if skipped.
3316    ///
3317    /// Non-blocking: operation ownership prevents conflicts with background merges.
3318    /// Copies unchanged files and rebuilds only the sparse file with reordered BMP data.
3319    pub async fn reorder_single_segment(
3320        self: &Arc<Self>,
3321        seg_id: &str,
3322        rayon_pool: Option<Arc<rayon::ThreadPool>>,
3323        bp_budget: crate::segment::BpBudget,
3324    ) -> Result<bool> {
3325        let source_id = SegmentId::from_hex(seg_id)
3326            .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3327        if self.quarantined_segments.lock().contains(seg_id) {
3328            return Err(Error::Corruption(format!(
3329                "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3330                seg_id
3331            )));
3332        }
3333        if self.force_merge_active.load(Ordering::Acquire) > 0 {
3334            log::debug!(
3335                "[optimizer] explicit force merge active, skipping reorder of {}",
3336                seg_id,
3337            );
3338            return Ok(false);
3339        }
3340
3341        // Whole-pass concurrency is independent from Rayon width. One pass
3342        // can already use every configured BP worker; this permit bounds the
3343        // much larger forward-index and rewrite working set across indexes,
3344        // optimizer tasks, and merge-time BP.
3345        let reorder_gate = Arc::clone(&self.reorder_permits);
3346        let _reorder_permit = tokio::select! {
3347            biased;
3348            () = self.active_operations.wait_for_shutdown() => {
3349                return Err(Error::IndexClosed);
3350            }
3351            permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3352                permit.map_err(|_| {
3353                    Error::Internal("background reorder scheduler is closed".into())
3354                })?
3355            }
3356        };
3357
3358        let output_id = SegmentId::new();
3359        let output_hex = output_id.to_hex();
3360        let source_ids = [seg_id.to_string()];
3361        let granularity = self.merge_granularity(&source_ids).await;
3362
3363        // Register while holding `state`, matching orphan cleanup's deletion
3364        // barrier. Candidates are scanned ahead of time and can go stale: a
3365        // merge may have consumed this segment since. Its files may even still
3366        // be on disk (deferred deletion under a searcher snapshot) — reordering
3367        // them would re-insert a duplicate copy of docs the merge output holds.
3368        let all_ids = vec![seg_id.to_string(), output_hex];
3369        let (_guard, source_docs) = {
3370            let st = self.state.lock().await;
3371            // Force merge raises this barrier under the same state lock, so
3372            // this second check closes the race with the cheap early check
3373            // above and prevents optimizer starvation between final groups.
3374            if self.force_merge_active.load(Ordering::Acquire) > 0 {
3375                log::debug!(
3376                    "[optimizer] explicit force merge active, skipping reorder of {}",
3377                    seg_id,
3378                );
3379                return Ok(false);
3380            }
3381            let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3382                log::info!(
3383                    "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
3384                    seg_id
3385                );
3386                self.clear_reorder_retry(seg_id);
3387                return Ok(false);
3388            };
3389
3390            match self.active_operations.try_register(all_ids) {
3391                Some(guard) => (guard, source_meta.num_docs),
3392                None if !self.active_operations.is_accepting() => {
3393                    return Err(Error::IndexClosed);
3394                }
3395                None => {
3396                    log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
3397                    return Ok(false);
3398                }
3399            }
3400        };
3401
3402        // Fail before allocating a forward index or creating output files.
3403        // Missing mandatory files are deterministic and should remove this
3404        // segment from future optimizer scans, not consume the same CPU every
3405        // interval. Other I/O failures remain retryable.
3406        if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3407            if is_deterministic_source_error(&error) {
3408                self.quarantine_segment(seg_id, &error);
3409            } else if !matches!(&error, Error::IndexClosed) {
3410                self.pause_reorder_retries(seg_id, &error);
3411            }
3412            return Err(error);
3413        }
3414
3415        let mut output_cleanup = self.output_cleanup_guard(output_id);
3416
3417        let reorder_result = crate::segment::reorder::reorder_segment(
3418            self.directory.as_ref(),
3419            &self.schema,
3420            source_id,
3421            output_id,
3422            self.term_cache_blocks,
3423            self.bp_memory_budget_bytes,
3424            bp_budget,
3425            granularity,
3426            rayon_pool,
3427            Some(self.active_operations.cancellation_flag()),
3428        )
3429        .await;
3430        let (new_id, total_docs, bp_converged) = match reorder_result {
3431            Ok(v) => v,
3432            Err(e) => {
3433                // A failed pass may have copied tens of GB before dying;
3434                // delete the uncommitted output before propagating.
3435                self.delete_output_if_unregistered(output_id, "reorder failure")
3436                    .await;
3437                output_cleanup.disarm();
3438                if is_deterministic_source_error(&e) {
3439                    self.quarantine_segment(seg_id, &e);
3440                } else if !matches!(&e, Error::IndexClosed) {
3441                    self.pause_reorder_retries(seg_id, &e);
3442                }
3443                return Err(e);
3444            }
3445        };
3446
3447        // A pass with a depth floor above block granularity has, by
3448        // definition, not converged to block-level order — record it as
3449        // unconverged so the optimizer's deepening ladder revisits it with a
3450        // full-depth (warm-started) pass. Depth caps are only used by the
3451        // optimizer's first pass on large segments.
3452        let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3453        if let Err(e) = self
3454            .replace_segments(
3455                &[seg_id.to_string()],
3456                new_id,
3457                total_docs,
3458                ReplacementLayout::BpReordered {
3459                    converged: ladder_converged,
3460                },
3461            )
3462            .await
3463        {
3464            self.delete_output_if_unregistered(output_id, "replacement failure")
3465                .await;
3466            output_cleanup.disarm();
3467            if !matches!(&e, Error::IndexClosed) {
3468                self.pause_reorder_retries(seg_id, &e);
3469            }
3470            return Err(e);
3471        }
3472        output_cleanup.disarm();
3473        self.clear_reorder_retry(seg_id);
3474
3475        Ok(true)
3476    }
3477
3478    /// Clean up orphan segment files not registered in metadata.
3479    ///
3480    /// Reads metadata, active-operation ownership, and snapshot-deferred
3481    /// deletions to determine which segments are legitimate. Filesystem
3482    /// deletion is asynchronous; in-flight outputs and retired sources still
3483    /// held by readers are both protected.
3484    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3485        let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3486
3487        if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3488            for entry in entries {
3489                let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3490                    continue;
3491                };
3492                let Some(rest) = filename.strip_prefix("seg_") else {
3493                    continue;
3494                };
3495                let Some(hex_id) = rest.get(..32) else {
3496                    continue;
3497                };
3498                if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3499                    continue;
3500                }
3501                orphan_files
3502                    .entry(hex_id.to_ascii_lowercase())
3503                    .or_default()
3504                    .push(entry);
3505            }
3506        }
3507
3508        let mut deleted = 0;
3509        for (hex_id, paths) in &orphan_files {
3510            // Revalidate and atomically claim deletion under the same
3511            // state -> active_operations -> tracker order used by publishers.
3512            // The claim lets us release `state` before filesystem I/O: deleting
3513            // a multi-GB orphan must not freeze commits and snapshot acquisition.
3514            let deletion_guard = {
3515                let st = self.state.lock().await;
3516                if st.metadata.has_segment(hex_id) {
3517                    continue;
3518                }
3519                let Some(guard) = self
3520                    .active_operations
3521                    .try_register(vec![hex_id.to_string()])
3522                else {
3523                    continue;
3524                };
3525                if self.tracker.is_deletion_protected(hex_id) {
3526                    drop(guard);
3527                    continue;
3528                }
3529                guard
3530            };
3531
3532            // Delete what was actually discovered, not only the currently
3533            // known SegmentFiles extensions. This also removes partial files
3534            // left by older formats instead of reporting the same orphan on
3535            // every startup forever.
3536            let results =
3537                futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3538                    .await;
3539            let removed = results.into_iter().all(|result| match result {
3540                Ok(()) => true,
3541                Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3542                Err(error) => {
3543                    log::warn!(
3544                        "[segment_cleanup] failed sweeping orphan segment {}: {}",
3545                        hex_id,
3546                        error,
3547                    );
3548                    false
3549                }
3550            });
3551            // Releasing this claim is the deletion barrier. No producer can
3552            // adopt the ID while its files are being removed.
3553            drop(deletion_guard);
3554            if removed {
3555                deleted += 1;
3556                log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
3557            }
3558        }
3559
3560        Ok(deleted)
3561    }
3562}
3563
3564#[cfg(test)]
3565mod tests {
3566    use super::*;
3567    use std::sync::atomic::{AtomicBool, Ordering};
3568
3569    fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3570        let schema = crate::dsl::SchemaBuilder::default().build();
3571        let metadata = IndexMetadata::new(schema.clone());
3572        Arc::new(SegmentManager::new(
3573            Arc::new(crate::directories::RamDirectory::new()),
3574            Arc::new(schema),
3575            metadata,
3576            Box::new(crate::merge::NoMergePolicy),
3577            0,
3578            1,
3579            Arc::new(Semaphore::new(1)),
3580            None,
3581            1024,
3582            Arc::new(ReorderConcurrencyGate::new(1)),
3583            None,
3584        ))
3585    }
3586
3587    #[test]
3588    fn force_merge_planner_pairs_large_and_small_segments() {
3589        let groups = plan_force_merge_groups(
3590            vec![
3591                ("a".into(), 6),
3592                ("b".into(), 6),
3593                ("c".into(), 4),
3594                ("d".into(), 4),
3595            ],
3596            10,
3597        );
3598
3599        assert_eq!(groups.len(), 2);
3600        assert!(groups.iter().all(|group| group.total_docs == 10));
3601        assert!(groups.iter().all(|group| group.segments.len() == 2));
3602    }
3603
3604    #[test]
3605    fn force_merge_planner_leaves_oversized_segments_alone() {
3606        let groups = plan_force_merge_groups(
3607            vec![
3608                ("oversized".into(), 11),
3609                ("small-a".into(), 5),
3610                ("small-b".into(), 5),
3611            ],
3612            10,
3613        );
3614
3615        assert_eq!(groups.len(), 2);
3616        assert_eq!(groups[0].total_docs, 10);
3617        assert_eq!(groups[0].segments.len(), 2);
3618        assert_eq!(groups[1].total_docs, 11);
3619        assert_eq!(groups[1].segments.len(), 1);
3620    }
3621
3622    #[test]
3623    fn force_merge_planner_never_exceeds_segment_format_limit() {
3624        let groups = plan_force_merge_groups(
3625            vec![
3626                ("large-a".into(), 3_000_000_000),
3627                ("large-b".into(), 2_000_000_000),
3628            ],
3629            u64::from(u32::MAX),
3630        );
3631        assert_eq!(groups.len(), 2);
3632        assert!(
3633            groups
3634                .iter()
3635                .all(|group| group.total_docs <= u64::from(u32::MAX))
3636        );
3637    }
3638
3639    #[test]
3640    fn force_merge_hierarchy_has_one_final_bp_pass() {
3641        assert_eq!(force_merge_output_count(1), 0);
3642        assert_eq!(force_merge_output_count(2), 1);
3643        assert_eq!(force_merge_output_count(64), 1);
3644        assert_eq!(force_merge_output_count(65), 2);
3645        assert_eq!(force_merge_output_count(127), 2);
3646        assert_eq!(force_merge_output_count(128), 3);
3647        assert_eq!(force_merge_output_count(1_000), 16);
3648    }
3649
3650    fn expand_force_merge_node(
3651        hierarchy: &ForceMergeHierarchy,
3652        source_count: usize,
3653        node: usize,
3654        sources: &mut Vec<usize>,
3655    ) {
3656        if node < source_count {
3657            sources.push(node);
3658            return;
3659        }
3660
3661        let step_index = node - source_count;
3662        let step = hierarchy
3663            .steps
3664            .get(step_index)
3665            .expect("merge input must refer to an existing source or output");
3666        for &input in &step.inputs {
3667            assert!(
3668                input < node,
3669                "merge step {step_index} refers to a future output node {input}"
3670            );
3671            expand_force_merge_node(hierarchy, source_count, input, sources);
3672        }
3673    }
3674
3675    #[test]
3676    fn force_merge_hierarchy_has_minimal_valid_arity() {
3677        let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
3678
3679        for source_count in source_counts {
3680            let hierarchy = plan_force_merge_hierarchy(source_count);
3681            let output_count = hierarchy.steps.len();
3682
3683            assert!(
3684                hierarchy
3685                    .steps
3686                    .iter()
3687                    .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
3688                "invalid merge arity for {source_count} sources"
3689            );
3690            assert!(
3691                source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
3692                "{output_count} outputs cannot reduce {source_count} sources"
3693            );
3694            assert!(
3695                output_count == 1
3696                    || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
3697                "{output_count} outputs are not minimal for {source_count} sources"
3698            );
3699            assert_eq!(output_count, force_merge_output_count(source_count));
3700        }
3701    }
3702
3703    #[test]
3704    fn force_merge_hierarchy_preserves_exact_source_order() {
3705        for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
3706            let hierarchy = plan_force_merge_hierarchy(source_count);
3707            let mut sources = Vec::with_capacity(source_count);
3708            expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
3709            assert_eq!(
3710                sources,
3711                (0..source_count).collect::<Vec<_>>(),
3712                "source order changed for {source_count} sources"
3713            );
3714        }
3715    }
3716
3717    fn force_merge_rewrite_cost(source_count: usize) -> usize {
3718        let hierarchy = plan_force_merge_hierarchy(source_count);
3719        let mut node_weights = vec![1usize; source_count];
3720        let mut rewrite_cost = 0usize;
3721
3722        for (step_index, step) in hierarchy.steps.iter().enumerate() {
3723            let output = source_count + step_index;
3724            let output_weight = step
3725                .inputs
3726                .iter()
3727                .map(|&input| {
3728                    assert!(
3729                        input < output,
3730                        "merge step {step_index} refers to future output {input}"
3731                    );
3732                    node_weights[input]
3733                })
3734                .sum::<usize>();
3735            rewrite_cost += output_weight;
3736            node_weights.push(output_weight);
3737        }
3738
3739        assert_eq!(node_weights[hierarchy.root], source_count);
3740        rewrite_cost
3741    }
3742
3743    #[test]
3744    fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
3745        assert_eq!(force_merge_rewrite_cost(65), 67);
3746        assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
3747    }
3748
3749    #[test]
3750    fn block_copy_carries_bp_debt_without_spending_an_attempt() {
3751        assert_eq!(
3752            replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
3753            (false, false, 3),
3754        );
3755        assert_eq!(
3756            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
3757            (true, false, 4),
3758        );
3759        assert_eq!(
3760            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
3761            (true, true, 0),
3762        );
3763    }
3764
3765    #[tokio::test]
3766    async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
3767        let mut schema_builder = crate::dsl::SchemaBuilder::default();
3768        let field = schema_builder.add_text_field("text", true, true);
3769        let schema = schema_builder.build();
3770        let directory = crate::directories::RamDirectory::new();
3771        let config = crate::index::IndexConfig {
3772            num_indexing_threads: 1,
3773            merge_policy: Box::new(crate::merge::NoMergePolicy),
3774            ..Default::default()
3775        };
3776        let mut writer = crate::index::IndexWriter::create(directory, schema, config)
3777            .await
3778            .unwrap();
3779        for value in ["one", "two", "three"] {
3780            let mut document = crate::dsl::Document::new();
3781            document.add_text(field, value);
3782            writer.add_document(document).unwrap();
3783            writer.commit().await.unwrap();
3784        }
3785
3786        let manager = Arc::clone(writer.segment_manager());
3787        let held_id = manager.get_segment_ids().await.pop().unwrap();
3788        let mut held = Some(
3789            manager
3790                .active_operations
3791                .try_register(vec![held_id])
3792                .unwrap(),
3793        );
3794        let batches = Arc::new(AtomicUsize::new(0));
3795        let batch_count = Arc::clone(&batches);
3796        writer
3797            .force_merge_with_snapshot_refresh(move || {
3798                let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
3799                // Refresh #1 follows the startup drain. Keep the source held
3800                // until refresh #2, after the two free sources were replaced.
3801                if refresh == 2 {
3802                    drop(held.take());
3803                }
3804                std::future::ready(Ok(()))
3805            })
3806            .await
3807            .unwrap();
3808
3809        assert_eq!(manager.get_segment_ids().await.len(), 1);
3810        assert_eq!(
3811            batches.load(Ordering::Relaxed),
3812            4,
3813            "initial/final refreshes plus two replacements are required after the held source releases"
3814        );
3815    }
3816
3817    #[tokio::test]
3818    async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
3819        let mut schema_builder = crate::dsl::SchemaBuilder::default();
3820        schema_builder.set_reorder_on_merge(true);
3821        let schema = schema_builder.build();
3822        let mut metadata = IndexMetadata::new(schema.clone());
3823        metadata.add_segment("00000000000000000000000000000001".into(), 1);
3824        metadata.add_segment("00000000000000000000000000000002".into(), 1);
3825
3826        let global_merge_permits = Arc::new(Semaphore::new(1));
3827        let manager = Arc::new(SegmentManager::new(
3828            Arc::new(crate::directories::RamDirectory::new()),
3829            Arc::new(schema),
3830            metadata,
3831            Box::new(crate::merge::NoMergePolicy),
3832            0,
3833            1,
3834            Arc::clone(&global_merge_permits),
3835            None,
3836            1024,
3837            Arc::new(ReorderConcurrencyGate::new(1)),
3838            None,
3839        ));
3840
3841        // Vector-generation staging rejects ordinary lifecycle claims while
3842        // acquiring the shared global merge permit separately for each source.
3843        // Force merge must wait before capacity admission; retaining the only
3844        // global slot here would deadlock both operations between sources.
3845        manager.active_operations.pause_non_indexing();
3846        let force_merge = {
3847            let manager = Arc::clone(&manager);
3848            tokio::spawn(async move { manager.force_merge().await })
3849        };
3850        tokio::time::timeout(std::time::Duration::from_secs(1), async {
3851            while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
3852                tokio::task::yield_now().await;
3853            }
3854        })
3855        .await
3856        .expect("force merge never reached the paused group claim");
3857
3858        assert_eq!(
3859            global_merge_permits.available_permits(),
3860            1,
3861            "force merge retained global capacity while vector staging blocked group ownership"
3862        );
3863
3864        force_merge.abort();
3865        let _ = force_merge.await;
3866        manager.active_operations.resume_non_indexing();
3867    }
3868
3869    #[test]
3870    fn output_cleanup_guard_runs_during_panic_unwind() {
3871        let cleaned = Arc::new(AtomicBool::new(false));
3872        let cleaned_in_callback = Arc::clone(&cleaned);
3873        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3874            cleaned_in_callback.store(true, Ordering::SeqCst);
3875        });
3876
3877        let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3878            let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3879            panic!("simulated reorder panic");
3880        }));
3881
3882        assert!(result.is_err());
3883        assert!(
3884            cleaned.load(Ordering::SeqCst),
3885            "partial output cleanup must run during unwind"
3886        );
3887    }
3888
3889    #[test]
3890    fn output_cleanup_guard_disarms_after_commit() {
3891        let cleaned = Arc::new(AtomicBool::new(false));
3892        let cleaned_in_callback = Arc::clone(&cleaned);
3893        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3894            cleaned_in_callback.store(true, Ordering::SeqCst);
3895        });
3896
3897        {
3898            let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3899            guard.disarm();
3900        }
3901
3902        assert!(!cleaned.load(Ordering::SeqCst));
3903    }
3904
3905    #[test]
3906    fn test_active_operation_guard_releases_ownership() {
3907        let active = Arc::new(ActiveSegmentOperations::new());
3908        {
3909            let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3910            let snap = active.snapshot();
3911            assert!(snap.contains("a"));
3912            assert!(snap.contains("b"));
3913        }
3914        assert!(active.snapshot().is_empty());
3915    }
3916
3917    #[test]
3918    fn test_non_overlapping_operations_can_run_concurrently() {
3919        let active = Arc::new(ActiveSegmentOperations::new());
3920        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3921        let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
3922        let snap = active.snapshot();
3923        assert_eq!(snap.len(), 4);
3924
3925        drop(first);
3926        let snap = active.snapshot();
3927        assert_eq!(snap.len(), 2);
3928        assert!(snap.contains("c"));
3929        assert!(snap.contains("d"));
3930    }
3931
3932    #[test]
3933    fn test_overlapping_operation_is_rejected_until_release() {
3934        let active = Arc::new(ActiveSegmentOperations::new());
3935        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3936        assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
3937        drop(first);
3938        assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
3939    }
3940
3941    #[test]
3942    fn test_active_operation_snapshot() {
3943        let active = Arc::new(ActiveSegmentOperations::new());
3944        let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
3945        let snap = active.snapshot();
3946        assert!(snap.contains("x"));
3947        assert!(snap.contains("y"));
3948        assert!(!snap.contains("z"));
3949    }
3950
3951    #[tokio::test]
3952    async fn operation_barrier_ignores_producers_started_after_snapshot() {
3953        let active = Arc::new(ActiveSegmentOperations::new());
3954        let before_gate = active.try_register(vec!["old".into()]).unwrap();
3955        let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
3956        assert_eq!(parked_indexing, 0);
3957        let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
3958
3959        let waiter = {
3960            let active = Arc::clone(&active);
3961            tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
3962        };
3963        tokio::task::yield_now().await;
3964        assert!(!waiter.is_finished());
3965
3966        drop(before_gate);
3967        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
3968            .await
3969            .expect("pre-gate operation barrier was starved by a post-gate producer")
3970            .unwrap();
3971        assert!(active.snapshot().contains("new-flat"));
3972        drop(after_gate);
3973    }
3974
3975    #[tokio::test]
3976    async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
3977        let manager = lifecycle_test_manager();
3978        manager
3979            .trained
3980            .store(Some(Arc::new(TrainedVectorStructures {
3981                centroids: rustc_hash::FxHashMap::default(),
3982                binary_quantizers: rustc_hash::FxHashMap::default(),
3983                ..Default::default()
3984            })));
3985
3986        let guard = manager.begin_vector_artifact_update().await.unwrap();
3987        assert!(
3988            manager.trained().is_some(),
3989            "search readers keep the last fully validated generation"
3990        );
3991        assert!(
3992            manager.trained_for_segment_build().is_none(),
3993            "new segment producers must stay flat during an artifact update"
3994        );
3995
3996        let detached_transaction_guard = guard.clone();
3997        drop(guard);
3998        assert!(
3999            manager.trained_for_segment_build().is_none(),
4000            "a detached lifecycle transaction must retain the producer gate after request cancellation"
4001        );
4002        drop(detached_transaction_guard);
4003        assert!(manager.trained_for_segment_build().is_some());
4004    }
4005
4006    #[tokio::test]
4007    async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4008        let manager = lifecycle_test_manager();
4009        let guard = manager.begin_vector_artifact_update().await.unwrap();
4010        assert!(
4011            manager
4012                .active_operations
4013                .try_register(vec!["merge".into()])
4014                .is_none(),
4015            "ordinary merge/reorder work must not change staged sources"
4016        );
4017        let indexing = manager
4018            .active_operations
4019            .try_register_indexing(vec!["fresh".into()])
4020            .expect("indexing remains available in flat mode");
4021        drop(indexing);
4022
4023        drop(guard);
4024        assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4025        assert!(
4026            manager
4027                .active_operations
4028                .try_register(vec!["merge".into()])
4029                .is_some()
4030        );
4031    }
4032
4033    #[tokio::test]
4034    async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4035        let active = Arc::new(ActiveSegmentOperations::new());
4036        let guard = active.try_register(vec!["live".into()]).unwrap();
4037        let cancellation = active.cancellation_flag();
4038        active.stop_accepting();
4039        assert!(cancellation.load(Ordering::Acquire));
4040        assert!(active.try_register(vec!["new".into()]).is_none());
4041
4042        let waiter = {
4043            let active = Arc::clone(&active);
4044            tokio::spawn(async move { active.wait_until_idle().await })
4045        };
4046        tokio::task::yield_now().await;
4047        assert!(!waiter.is_finished());
4048        drop(guard);
4049        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4050            .await
4051            .expect("shutdown waiter missed the final guard notification")
4052            .unwrap();
4053    }
4054
4055    #[tokio::test]
4056    async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4057        let manager = lifecycle_test_manager();
4058        let started = Arc::new(Semaphore::new(0));
4059        let release = Arc::new(Semaphore::new(0));
4060        let completed = Arc::new(AtomicBool::new(false));
4061
4062        let request = {
4063            let manager = Arc::clone(&manager);
4064            let started = Arc::clone(&started);
4065            let release = Arc::clone(&release);
4066            let completed = Arc::clone(&completed);
4067            tokio::spawn(async move {
4068                manager
4069                    .run_lifecycle_transaction(async move {
4070                        started.add_permits(1);
4071                        let _permit = release.acquire().await.unwrap();
4072                        completed.store(true, Ordering::Release);
4073                        Ok(())
4074                    })
4075                    .await
4076            })
4077        };
4078
4079        let _started = started.acquire().await.unwrap();
4080        request.abort();
4081        assert!(request.await.unwrap_err().is_cancelled());
4082        release.add_permits(1);
4083
4084        manager.begin_shutdown();
4085        tokio::time::timeout(
4086            std::time::Duration::from_secs(1),
4087            manager.wait_for_shutdown(),
4088        )
4089        .await
4090        .expect("shutdown did not drain detached lifecycle transaction");
4091        assert!(completed.load(Ordering::Acquire));
4092    }
4093
4094    #[tokio::test]
4095    async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4096        let manager = lifecycle_test_manager();
4097        {
4098            let mut state = manager.state.lock().await;
4099            state.metadata.add_segment_meta(
4100                "eligible".into(),
4101                SegmentMetaInfo {
4102                    num_docs: 10,
4103                    ancestors: Vec::new(),
4104                    generation: 1,
4105                    reordered: true,
4106                    bp_converged: false,
4107                    bp_unconverged_passes: 2,
4108                },
4109            );
4110            state.metadata.add_segment_meta(
4111                "at-limit".into(),
4112                SegmentMetaInfo {
4113                    num_docs: 20,
4114                    ancestors: Vec::new(),
4115                    generation: 1,
4116                    reordered: true,
4117                    bp_converged: false,
4118                    bp_unconverged_passes: 3,
4119                },
4120            );
4121            state.metadata.add_segment_meta(
4122                "carried-debt".into(),
4123                SegmentMetaInfo {
4124                    num_docs: 15,
4125                    ancestors: Vec::new(),
4126                    generation: 2,
4127                    reordered: false,
4128                    bp_converged: false,
4129                    bp_unconverged_passes: 2,
4130                },
4131            );
4132            state.metadata.add_segment_meta(
4133                "carried-debt-at-limit".into(),
4134                SegmentMetaInfo {
4135                    num_docs: 25,
4136                    ancestors: Vec::new(),
4137                    generation: 2,
4138                    reordered: false,
4139                    bp_converged: false,
4140                    bp_unconverged_passes: 3,
4141                },
4142            );
4143            state.metadata.add_segment_meta(
4144                "converged".into(),
4145                SegmentMetaInfo {
4146                    num_docs: 30,
4147                    ancestors: Vec::new(),
4148                    generation: 1,
4149                    reordered: true,
4150                    bp_converged: true,
4151                    bp_unconverged_passes: 0,
4152                },
4153            );
4154            state.metadata.add_segment("fresh".into(), 40);
4155        }
4156
4157        assert_eq!(
4158            manager.unreordered_segments().await,
4159            vec![("fresh".into(), 40)],
4160            "a block-copy output with BP debt is not a fresh first-pass candidate",
4161        );
4162        let mut eligible = manager.unconverged_segments_below(3).await;
4163        eligible.sort_unstable();
4164        assert_eq!(
4165            eligible,
4166            vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4167        );
4168        assert!(manager.unconverged_segments_below(0).await.is_empty());
4169    }
4170
4171    #[test]
4172    fn merge_retry_backoff_is_exponential_and_capped() {
4173        assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4174        assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4175        assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4176        assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4177    }
4178
4179    #[test]
4180    fn only_deterministic_source_errors_are_quarantined() {
4181        assert!(is_deterministic_source_error(&Error::Corruption(
4182            "bad footer".into()
4183        )));
4184        assert!(is_deterministic_source_error(&Error::Io(
4185            std::io::Error::from(std::io::ErrorKind::NotFound)
4186        )));
4187        assert!(!is_deterministic_source_error(&Error::Io(
4188            std::io::Error::from(std::io::ErrorKind::TimedOut)
4189        )));
4190        assert!(!is_deterministic_source_error(&Error::Io(
4191            std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4192        )));
4193    }
4194
4195    #[test]
4196    fn transient_reorder_failure_is_backed_off_until_cleared() {
4197        let manager = lifecycle_test_manager();
4198        manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4199        assert!(manager.paused_reorder_segments().contains("source"));
4200        manager.clear_reorder_retry("source");
4201        assert!(!manager.paused_reorder_segments().contains("source"));
4202    }
4203
4204    /// Fails `exists` with the transient I/O error class that sends a
4205    /// background merge into its generic retry backoff (not source quarantine).
4206    #[derive(Default)]
4207    struct FailingExistsDirectory(crate::directories::RamDirectory);
4208
4209    #[async_trait::async_trait]
4210    impl crate::directories::Directory for FailingExistsDirectory {
4211        async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4212            Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4213        }
4214
4215        async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4216            self.0.file_size(path).await
4217        }
4218
4219        async fn open_read(
4220            &self,
4221            path: &std::path::Path,
4222        ) -> std::io::Result<crate::directories::FileHandle> {
4223            self.0.open_read(path).await
4224        }
4225
4226        async fn read_range(
4227            &self,
4228            path: &std::path::Path,
4229            range: std::ops::Range<u64>,
4230        ) -> std::io::Result<crate::directories::OwnedBytes> {
4231            self.0.read_range(path, range).await
4232        }
4233
4234        async fn list_files(
4235            &self,
4236            prefix: &std::path::Path,
4237        ) -> std::io::Result<Vec<std::path::PathBuf>> {
4238            self.0.list_files(prefix).await
4239        }
4240
4241        async fn open_lazy(
4242            &self,
4243            path: &std::path::Path,
4244        ) -> std::io::Result<crate::directories::FileHandle> {
4245            self.0.open_lazy(path).await
4246        }
4247    }
4248
4249    #[async_trait::async_trait]
4250    impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4251        async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4252            self.0.write(path, data).await
4253        }
4254
4255        async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4256            self.0.delete(path).await
4257        }
4258
4259        async fn rename(
4260            &self,
4261            from: &std::path::Path,
4262            to: &std::path::Path,
4263        ) -> std::io::Result<()> {
4264            self.0.rename(from, to).await
4265        }
4266
4267        async fn sync(&self) -> std::io::Result<()> {
4268            self.0.sync().await
4269        }
4270
4271        async fn streaming_writer(
4272            &self,
4273            path: &std::path::Path,
4274        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4275            self.0.streaming_writer(path).await
4276        }
4277    }
4278
4279    #[derive(Debug, Clone)]
4280    struct MergeEverythingPolicy;
4281
4282    impl MergePolicy for MergeEverythingPolicy {
4283        fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4284            if segments.len() < 2 {
4285                return Vec::new();
4286            }
4287            vec![crate::merge::MergeCandidate {
4288                segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4289            }]
4290        }
4291
4292        fn clone_box(&self) -> Box<dyn MergePolicy> {
4293            Box::new(self.clone())
4294        }
4295    }
4296
4297    #[tokio::test]
4298    async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4299        let manager = lifecycle_test_manager();
4300        // Simulates a memory-budget mid-cycle segment build whose guard is
4301        // parked inside a PreparedSegment: only a later commit releases this
4302        // token, and that commit can be blocked on the very caller of the
4303        // artifact update (writer write lock / &mut self).
4304        let parked_indexing = manager
4305            .protect_new_segment("00000000000000000000000000000abc".into())
4306            .unwrap();
4307
4308        let error = tokio::time::timeout(
4309            std::time::Duration::from_secs(2),
4310            manager.begin_vector_artifact_update(),
4311        )
4312        .await
4313        .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4314        .err()
4315        .expect("an old-generation prepared segment must block artifact replacement")
4316        .to_string();
4317        assert!(error.contains("built but uncommitted"), "{error}");
4318        assert!(
4319            !manager.vector_artifact_update.load(Ordering::Acquire),
4320            "a rejected update must release the producer gate"
4321        );
4322
4323        drop(parked_indexing);
4324
4325        let guard = manager
4326            .begin_vector_artifact_update()
4327            .await
4328            .expect("artifact update should succeed after the pending generation is resolved");
4329        drop(guard);
4330    }
4331
4332    #[tokio::test]
4333    async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4334        let manager = lifecycle_test_manager();
4335        let merge_like = manager
4336            .active_operations
4337            .try_register(vec!["merge-source".into()])
4338            .unwrap();
4339
4340        let waiter = {
4341            let manager = Arc::clone(&manager);
4342            tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4343        };
4344        for _ in 0..8 {
4345            tokio::task::yield_now().await;
4346        }
4347        assert!(
4348            !waiter.is_finished(),
4349            "artifact update must drain merge/reorder producers that may hold the previous generation"
4350        );
4351
4352        drop(merge_like);
4353        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4354            .await
4355            .expect("artifact update missed the lifecycle guard release")
4356            .unwrap()
4357            .unwrap();
4358    }
4359
4360    #[tokio::test]
4361    async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4362        let manager = lifecycle_test_manager();
4363        let release = Arc::new(Semaphore::new(0));
4364        let merge_task = {
4365            let release = Arc::clone(&release);
4366            tokio::spawn(async move {
4367                let _permit = release.acquire().await.unwrap();
4368            })
4369        };
4370        manager.merge_handles.lock().push(merge_task);
4371
4372        let waiter = {
4373            let manager = Arc::clone(&manager);
4374            tokio::spawn(async move { manager.wait_for_all_merges().await })
4375        };
4376        for _ in 0..8 {
4377            tokio::task::yield_now().await;
4378        }
4379        assert!(!waiter.is_finished());
4380        // Simulates tonic dropping a force_merge/reorder RPC future at the
4381        // JoinHandle await when the client disconnects.
4382        waiter.abort();
4383        let join_error = waiter.await.unwrap_err();
4384        assert!(join_error.is_cancelled());
4385
4386        assert!(
4387            !manager.merge_handles.lock().is_empty(),
4388            "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4389        );
4390
4391        // A later drain must still see and await the real in-flight merge.
4392        release.add_permits(1);
4393        tokio::time::timeout(
4394            std::time::Duration::from_secs(1),
4395            manager.wait_for_all_merges(),
4396        )
4397        .await
4398        .expect("subsequent drain missed the reinserted merge handle");
4399        assert!(manager.merge_handles.lock().is_empty());
4400    }
4401
4402    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4403    async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4404        let manager = lifecycle_test_manager();
4405        {
4406            let mut state = manager.state.lock().await;
4407            state
4408                .metadata
4409                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4410            state
4411                .metadata
4412                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4413        }
4414        // A background reorder (or a concurrent force-merge) owns one segment
4415        // in the batch but never appears in merge_handles.
4416        let reorder_like = manager
4417            .active_operations
4418            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4419            .unwrap();
4420
4421        let force_merge = {
4422            let manager = Arc::clone(&manager);
4423            tokio::spawn(async move { manager.force_merge().await })
4424        };
4425
4426        tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4427        let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4428        assert!(
4429            retries >= 1,
4430            "force_merge never observed the conflicting owner (retries={retries})"
4431        );
4432        assert!(
4433            retries < 20,
4434            "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4435        );
4436
4437        drop(reorder_like);
4438        // With the conflict gone the loop proceeds; the batch then fails fast
4439        // in do_merge (the test IDs have no files), proving the loop exited.
4440        let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4441            .await
4442            .expect("force_merge kept spinning after the conflicting owner released")
4443            .unwrap();
4444        assert!(result.is_err());
4445    }
4446
4447    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4448    async fn force_merge_routes_around_segments_held_by_reorder() {
4449        let manager = lifecycle_test_manager();
4450        {
4451            let mut state = manager.state.lock().await;
4452            state
4453                .metadata
4454                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4455            state
4456                .metadata
4457                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4458            state
4459                .metadata
4460                .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4461        }
4462        // A background reorder owns one segment and holds it for the whole
4463        // test (in prod: a BP pass runs for minutes while force_merge spins).
4464        let _reorder_like = manager
4465            .active_operations
4466            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4467            .unwrap();
4468
4469        // Regression: force_merge used to rebuild the identical smallest-N
4470        // batch (including the held segment) every 100ms and retry-log
4471        // forever. It must instead skip the held segment and immediately
4472        // make progress on the two free ones — reaching do_merge (which
4473        // fails fast here: the test IDs have no files) proves the batch was
4474        // built without the held segment while the reorder is STILL active.
4475        let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4476            let manager = Arc::clone(&manager);
4477            async move { manager.force_merge().await }
4478        })
4479        .await
4480        .expect("force_merge livelocked on a segment held by an active reorder");
4481        assert!(result.is_err(), "fake segment files must fail the merge");
4482
4483        assert_eq!(
4484            manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4485            0,
4486            "batch built from the ownership snapshot must not collide with the held segment"
4487        );
4488    }
4489
4490    #[tokio::test]
4491    async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4492        let schema = crate::dsl::SchemaBuilder::default().build();
4493        let mut metadata = IndexMetadata::new(schema.clone());
4494        metadata.add_segment("00000000000000000000000000000001".into(), 10);
4495        metadata.add_segment("00000000000000000000000000000002".into(), 10);
4496        let manager = Arc::new(SegmentManager::new(
4497            Arc::new(FailingExistsDirectory::default()),
4498            Arc::new(schema),
4499            metadata,
4500            Box::new(MergeEverythingPolicy),
4501            0,
4502            1,
4503            Arc::new(Semaphore::new(1)),
4504            None,
4505            1024,
4506            Arc::new(ReorderConcurrencyGate::new(1)),
4507            None,
4508        ));
4509
4510        // Spawns a background merge that fails with a transient I/O error and
4511        // arms the 30s..30min retry backoff.
4512        manager.maybe_merge().await;
4513
4514        tokio::time::timeout(
4515            std::time::Duration::from_secs(5),
4516            manager.wait_for_all_merges(),
4517        )
4518        .await
4519        .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4520        assert!(
4521            manager.merge_retry_is_paused(),
4522            "the failed merge should have armed the retry backoff"
4523        );
4524
4525        // Shutdown still drains the pending backoff wakeup deterministically.
4526        manager.begin_shutdown();
4527        tokio::time::timeout(
4528            std::time::Duration::from_secs(5),
4529            manager.wait_for_shutdown(),
4530        )
4531        .await
4532        .expect("shutdown did not drain the merge retry wakeup task");
4533    }
4534}