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        trained: &TrainedVectorStructures,
2757        rewrite_existing: bool,
2758    ) -> Result<bool> {
2759        for &field_id in field_ids {
2760            let flat = reader.flat_vectors().get(&field_id);
2761            let ann = reader.vector_indexes().get(&field_id);
2762            if ann.is_some() && flat.is_none() {
2763                return Err(Error::Corruption(format!(
2764                    "segment {:032x} field {field_id} has ANN data without the required flat vectors",
2765                    reader.meta().id,
2766                )));
2767            }
2768
2769            let Some(flat) = flat else {
2770                continue;
2771            };
2772            if flat.num_vectors == 0 {
2773                continue;
2774            }
2775            if rewrite_existing {
2776                return Ok(true);
2777            }
2778            let field = crate::dsl::Field(field_id);
2779            let entry = self.schema.get_field_entry(field).ok_or_else(|| {
2780                Error::Corruption(format!(
2781                    "segment {:032x} references unknown vector field {field_id}",
2782                    reader.meta().id,
2783                ))
2784            })?;
2785            let current = match entry.field_type {
2786                // TQ payloads carry no trained generation, so an existing one
2787                // is always current; a tq field must never be staged for a
2788                // vector-generation rewrite.
2789                crate::dsl::FieldType::DenseVector
2790                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
2791                        config.index_type == crate::dsl::VectorIndexType::Tq
2792                    }) =>
2793                {
2794                    matches!(ann, Some(crate::segment::VectorIndex::Tq { .. }))
2795                }
2796                crate::dsl::FieldType::DenseVector
2797                    if entry.dense_vector_config.as_ref().is_some_and(|config| {
2798                        config.index_type == crate::dsl::VectorIndexType::IvfTq
2799                    }) =>
2800                {
2801                    let config = entry
2802                        .dense_vector_config
2803                        .as_ref()
2804                        .expect("matched IVF-TQ configuration");
2805                    match (ann, trained.centroids.get(&field_id)) {
2806                        (
2807                            Some(crate::segment::VectorIndex::IvfTq { index, .. }),
2808                            Some(centroids),
2809                        ) => {
2810                            let header = index.get().header();
2811                            crate::structures::is_ivf_tq_cosine_generation(centroids.version)
2812                                && crate::structures::is_ivf_tq_cosine_generation(
2813                                    header.quantizer_version,
2814                                )
2815                                && header.dim == config.dim
2816                                && header.num_clusters == centroids.num_clusters
2817                                && header.quantizer_version == centroids.version
2818                                && header.codebook_version
2819                                    == crate::structures::vector::quantization::tq_expected_fingerprint(
2820                                        config.dim,
2821                                    )
2822                                && header.routing == config.ivf_routing
2823                        }
2824                        _ => false,
2825                    }
2826                }
2827                // Only ivf_tq dense fields are trainable; other dense
2828                // index types never reach a vector-generation rewrite.
2829                crate::dsl::FieldType::DenseVector => false,
2830                crate::dsl::FieldType::BinaryDenseVector => {
2831                    matches!(ann, Some(crate::segment::VectorIndex::BinaryIvf(_)))
2832                }
2833                _ => false,
2834            };
2835            if !current {
2836                return Ok(true);
2837            }
2838        }
2839        Ok(false)
2840    }
2841
2842    async fn acquire_vector_rewrite_capacity(
2843        &self,
2844    ) -> Result<(OwnedSemaphorePermit, OwnedSemaphorePermit)> {
2845        let global = tokio::select! {
2846            biased;
2847            () = self.active_operations.wait_for_shutdown() => {
2848                return Err(Error::IndexClosed);
2849            }
2850            permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
2851                permit.map_err(|_| Error::Internal(
2852                    "global background merge scheduler is closed".into()
2853                ))?
2854            }
2855        };
2856        let local = tokio::select! {
2857            biased;
2858            () = self.active_operations.wait_for_shutdown() => {
2859                return Err(Error::IndexClosed);
2860            }
2861            permit = Arc::clone(&self.merge_permits).acquire_owned() => {
2862                permit.map_err(|_| Error::Internal(
2863                    "background merge scheduler is closed".into()
2864                ))?
2865            }
2866        };
2867        Ok((global, local))
2868    }
2869
2870    async fn build_vector_replacement(
2871        self: &Arc<Self>,
2872        segment_id: &str,
2873        source_id: SegmentId,
2874        output_id: SegmentId,
2875        trained: &TrainedVectorStructures,
2876        failure_context: &'static str,
2877    ) -> Result<(String, u32, OutputCleanupGuard)> {
2878        let mut cleanup = self.output_cleanup_guard(output_id);
2879        match crate::segment::reorder::rewrite_vector_segment(
2880            self.directory.as_ref(),
2881            &self.schema,
2882            source_id,
2883            output_id,
2884            self.term_cache_blocks,
2885            trained,
2886            Some(self.background_cpu_pool()),
2887        )
2888        .await
2889        {
2890            Ok((new_id, doc_count)) => {
2891                self.validate_completed_segment(&new_id, doc_count).await?;
2892                Ok((new_id, doc_count, cleanup))
2893            }
2894            Err(error) => {
2895                self.delete_output_if_unregistered(output_id, failure_context)
2896                    .await;
2897                cleanup.disarm();
2898                if is_deterministic_source_error(&error) {
2899                    self.quarantine_segment(segment_id, &error);
2900                }
2901                Err(error)
2902            }
2903        }
2904    }
2905
2906    /// Build every required replacement segment without exposing any of them.
2907    /// The returned guards keep both source and output generations alive until
2908    /// [`Self::publish_vector_generation`] commits the complete set.
2909    pub(crate) async fn stage_vector_generation(
2910        self: &Arc<Self>,
2911        _artifact_update: &VectorArtifactUpdateGuard,
2912        segment_ids: &[String],
2913        field_ids: &[u32],
2914        trained: Arc<TrainedVectorStructures>,
2915        rewrite_existing: bool,
2916    ) -> Result<Vec<StagedVectorSegment>> {
2917        if !self.vector_artifact_update.load(Ordering::Acquire) {
2918            return Err(Error::Internal(
2919                "cannot stage a vector generation without an exclusive update lease".into(),
2920            ));
2921        }
2922
2923        let mut staged = Vec::new();
2924        for segment_id in segment_ids {
2925            if self.quarantined_segments.lock().contains(segment_id) {
2926                return Err(Error::Corruption(format!(
2927                    "segment {segment_id} is quarantined after a deterministic source failure"
2928                )));
2929            }
2930            let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
2931                Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
2932            })?;
2933
2934            // A single rewrite may hold several gigabytes while assigning all
2935            // vectors. Reuse the ordinary local and process-wide merge bounds.
2936            let _capacity = self.acquire_vector_rewrite_capacity().await?;
2937
2938            let output_id = SegmentId::new();
2939            let output_hex = output_id.to_hex();
2940            let operation = {
2941                let st = self.state.lock().await;
2942                if !st.metadata.has_segment(segment_id) {
2943                    return Err(Error::Corruption(format!(
2944                        "vector generation source {segment_id} disappeared while lifecycle work was paused"
2945                    )));
2946                }
2947                self.active_operations
2948                    .try_register_vector_update(vec![segment_id.clone(), output_hex.clone()])
2949            }
2950            .ok_or_else(|| {
2951                if self.active_operations.is_accepting() {
2952                    Error::Internal(format!(
2953                        "vector generation could not claim stable source {segment_id}"
2954                    ))
2955                } else {
2956                    Error::IndexClosed
2957                }
2958            })?;
2959
2960            let reader = SegmentReader::open(
2961                self.directory.as_ref(),
2962                source_id,
2963                Arc::clone(&self.schema),
2964                self.term_cache_blocks,
2965            )
2966            .await?;
2967            if !self.segment_needs_vector_rewrite(
2968                &reader,
2969                field_ids,
2970                trained.as_ref(),
2971                rewrite_existing,
2972            )? {
2973                continue;
2974            }
2975            drop(reader);
2976
2977            let (new_id, doc_count, cleanup) = self
2978                .build_vector_replacement(
2979                    segment_id,
2980                    source_id,
2981                    output_id,
2982                    trained.as_ref(),
2983                    "vector generation staging failure",
2984                )
2985                .await?;
2986            debug_assert_eq!(new_id, output_hex);
2987            let output_reader = SegmentReader::open(
2988                self.directory.as_ref(),
2989                output_id,
2990                Arc::clone(&self.schema),
2991                self.term_cache_blocks,
2992            )
2993            .await?;
2994            if self.segment_needs_vector_rewrite(
2995                &output_reader,
2996                field_ids,
2997                trained.as_ref(),
2998                false,
2999            )? {
3000                return Err(Error::Corruption(format!(
3001                    "staged vector segment {new_id} does not match its candidate codebook generation"
3002                )));
3003            }
3004
3005            staged.push(StagedVectorSegment {
3006                source_id: segment_id.clone(),
3007                output_id,
3008                doc_count,
3009                _operation: operation,
3010                cleanup,
3011            });
3012        }
3013        Ok(staged)
3014    }
3015
3016    async fn rewrite_vector_segment_once(
3017        self: &Arc<Self>,
3018        segment_id: &str,
3019        field_ids: &[u32],
3020    ) -> Result<VectorSegmentRewriteOutcome> {
3021        if self.quarantined_segments.lock().contains(segment_id) {
3022            return Err(Error::Corruption(format!(
3023                "segment {segment_id} is quarantined after a deterministic source failure; repair it before ANN finalization"
3024            )));
3025        }
3026        let source_id = SegmentId::from_hex(segment_id).ok_or_else(|| {
3027            Error::Corruption(format!("invalid vector rewrite segment ID: {segment_id}"))
3028        })?;
3029
3030        // Match ordinary merge lock ordering: capacity before lifecycle
3031        // ownership. A vector rewrite can hold several gigabytes while it
3032        // assigns vectors, so it participates in both local and process-wide
3033        // merge limits.
3034        let _capacity = self.acquire_vector_rewrite_capacity().await?;
3035
3036        let output_id = SegmentId::new();
3037        let output_hex = output_id.to_hex();
3038        let all_ids = vec![segment_id.to_owned(), output_hex];
3039        let operation = {
3040            let st = self.state.lock().await;
3041            if !st.metadata.has_segment(segment_id) {
3042                return Ok(VectorSegmentRewriteOutcome::SourceGone);
3043            }
3044            self.active_operations.try_register(all_ids)
3045        };
3046        let _operation = match operation {
3047            Some(operation) => operation,
3048            None if !self.active_operations.is_accepting() => return Err(Error::IndexClosed),
3049            None => return Ok(VectorSegmentRewriteOutcome::Conflict),
3050        };
3051
3052        let Some(trained) = self.trained_for_segment_build() else {
3053            return Ok(VectorSegmentRewriteOutcome::Deferred);
3054        };
3055
3056        let reader = SegmentReader::open(
3057            self.directory.as_ref(),
3058            source_id,
3059            Arc::clone(&self.schema),
3060            self.term_cache_blocks,
3061        )
3062        .await?;
3063        if !self.segment_needs_vector_rewrite(&reader, field_ids, trained.as_ref(), false)? {
3064            return Ok(VectorSegmentRewriteOutcome::AlreadyCurrent);
3065        }
3066        drop(reader);
3067
3068        let (new_id, doc_count, mut output_cleanup) = self
3069            .build_vector_replacement(
3070                segment_id,
3071                source_id,
3072                output_id,
3073                trained.as_ref(),
3074                "vector rewrite failure",
3075            )
3076            .await?;
3077
3078        if let Err(error) = self
3079            .replace_segments(
3080                &[segment_id.to_owned()],
3081                new_id,
3082                doc_count,
3083                ReplacementLayout::PreserveSingleSource,
3084            )
3085            .await
3086        {
3087            self.delete_output_if_unregistered(output_id, "vector replacement failure")
3088                .await;
3089            output_cleanup.disarm();
3090            return Err(error);
3091        }
3092        output_cleanup.disarm();
3093        Ok(VectorSegmentRewriteOutcome::Rewritten)
3094    }
3095
3096    /// Finalize every committed flat vector segment against the published
3097    /// global ANN generation.
3098    /// Unlike force-merge this handles one segment and segments already at the
3099    /// merge policy's maximum size.
3100    pub(crate) async fn rewrite_vector_segments(
3101        self: &Arc<Self>,
3102        field_ids: &[u32],
3103    ) -> Result<usize> {
3104        if field_ids.is_empty() {
3105            return Ok(0);
3106        }
3107        let mut rewritten = 0usize;
3108        loop {
3109            let segment_ids = self.get_segment_ids().await;
3110            let mut conflicted = false;
3111            let mut changed = false;
3112            for segment_id in segment_ids {
3113                match self
3114                    .rewrite_vector_segment_once(&segment_id, field_ids)
3115                    .await?
3116                {
3117                    VectorSegmentRewriteOutcome::Rewritten => {
3118                        rewritten += 1;
3119                        changed = true;
3120                    }
3121                    VectorSegmentRewriteOutcome::Conflict => conflicted = true,
3122                    VectorSegmentRewriteOutcome::Deferred => {
3123                        return Err(Error::Internal(
3124                            "ANN finalization lost the published trained generation".into(),
3125                        ));
3126                    }
3127                    VectorSegmentRewriteOutcome::AlreadyCurrent
3128                    | VectorSegmentRewriteOutcome::SourceGone => {}
3129                }
3130            }
3131            if !conflicted && !changed {
3132                log::info!(
3133                    "[dense_vector_rewrite] ANN finalization complete ({} segment(s) rewritten)",
3134                    rewritten,
3135                );
3136                return Ok(rewritten);
3137            }
3138            tokio::select! {
3139                biased;
3140                () = self.active_operations.wait_for_shutdown() => {
3141                    return Err(Error::IndexClosed);
3142                }
3143                () = tokio::time::sleep(std::time::Duration::from_millis(100)) => {}
3144            }
3145        }
3146    }
3147
3148    /// A producer that started in the force-flat phase can commit after the
3149    /// main finalization snapshot. Upgrade exactly those new segments in a
3150    /// tracked background task; ordinary producers already using the current
3151    /// generation are detected and skipped without rewriting.
3152    pub(crate) fn schedule_vector_segment_upgrades(self: &Arc<Self>, segment_ids: Vec<String>) {
3153        if segment_ids.is_empty() || self.trained_for_segment_build().is_none() {
3154            return;
3155        }
3156        let manager = Arc::clone(self);
3157        let future = async move {
3158            let field_ids = manager
3159                .read_metadata(|metadata| {
3160                    metadata
3161                        .vector_fields
3162                        .keys()
3163                        .filter(|field_id| metadata.is_field_built(**field_id))
3164                        .copied()
3165                        .collect::<Vec<_>>()
3166                })
3167                .await;
3168            for segment_id in segment_ids {
3169                loop {
3170                    match manager
3171                        .rewrite_vector_segment_once(&segment_id, &field_ids)
3172                        .await
3173                    {
3174                        Ok(VectorSegmentRewriteOutcome::Conflict) => {
3175                            tokio::time::sleep(std::time::Duration::from_millis(100)).await;
3176                        }
3177                        Ok(VectorSegmentRewriteOutcome::Deferred) => break,
3178                        Ok(_) => break,
3179                        Err(error) => {
3180                            log::error!(
3181                                "[dense_vector_rewrite] failed to upgrade newly committed segment {}: {}",
3182                                segment_id,
3183                                error,
3184                            );
3185                            break;
3186                        }
3187                    }
3188                }
3189            }
3190        };
3191        let Ok(runtime) = tokio::runtime::Handle::try_current() else {
3192            log::warn!(
3193                "[dense_vector_rewrite] runtime unavailable; newly committed flat segment upgrade deferred"
3194            );
3195            return;
3196        };
3197        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
3198            log::warn!(
3199                "[dense_vector_rewrite] runtime rejected newly committed flat segment upgrade"
3200            );
3201        }
3202    }
3203
3204    /// Reorder all segments via Recursive Graph Bisection (BP) for better BMP pruning.
3205    ///
3206    /// Each segment is individually rebuilt with reordered BMP blocks.
3207    /// Non-BMP fields are copied unchanged via streaming file copy.
3208    ///
3209    /// Uses active-operation ownership to prevent concurrent work on the same segment.
3210    pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
3211        self.reorder_segments_with_snapshot_refresh(|| std::future::ready(Ok(())))
3212            .await
3213    }
3214
3215    /// Reorder all segments while advancing long-lived snapshots after the
3216    /// background-merge drain and every durable replacement.
3217    pub(crate) async fn reorder_segments_with_snapshot_refresh<F, Fut>(
3218        self: &Arc<Self>,
3219        mut refresh_snapshots: F,
3220    ) -> Result<()>
3221    where
3222        F: FnMut() -> Fut,
3223        Fut: std::future::Future<Output = Result<()>>,
3224    {
3225        self.wait_for_all_merges().await;
3226        refresh_snapshots().await?;
3227        let segment_ids = self.get_segment_ids().await;
3228
3229        if segment_ids.is_empty() {
3230            log::info!("[reorder] no segments to reorder");
3231            return Ok(());
3232        }
3233
3234        log::info!("[reorder] reordering {} segments", segment_ids.len());
3235
3236        for seg_id in segment_ids {
3237            match self
3238                .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
3239                .await
3240            {
3241                Ok(true) => refresh_snapshots().await?,
3242                Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
3243                Err(e) => return Err(e),
3244            }
3245        }
3246
3247        // A segment skipped because another lifecycle owner held it may have
3248        // been replaced after the preceding callback.
3249        refresh_snapshots().await?;
3250        log::info!("[reorder] all segments reordered");
3251        Ok(())
3252    }
3253
3254    /// Get segment IDs that have not been reordered yet.
3255    ///
3256    /// Excludes segments currently involved in a merge or reorder operation
3257    /// to avoid wasted work (the optimizer would skip them anyway).
3258    pub async fn unreordered_segment_ids(&self) -> Vec<String> {
3259        self.unreordered_segments()
3260            .await
3261            .into_iter()
3262            .map(|(id, _)| id)
3263            .collect()
3264    }
3265
3266    /// Segments never reordered, with doc counts — for the optimizer to pick
3267    /// a size-appropriate BP budget.
3268    pub async fn unreordered_segments(&self) -> Vec<(String, 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.reordered
3278                    && info.bp_converged
3279                    && !active_ids.contains(*id)
3280                    && !quarantined.contains(*id)
3281                    && !paused.contains(*id)
3282            })
3283            .map(|(id, info)| (id.clone(), info.num_docs))
3284            .collect()
3285    }
3286
3287    /// Segments whose last BP pass hit its wall-clock budget before finishing
3288    /// (`bp_converged == false`). A warm-started follow-up pass deepens the
3289    /// ordering; the optimizer revisits these at low priority.
3290    pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
3291        self.unconverged_segments_below(u32::MAX)
3292            .await
3293            .into_iter()
3294            .map(|(id, docs, _)| (id, docs))
3295            .collect()
3296    }
3297
3298    /// Unconverged segments still below a hard replacement-lineage work
3299    /// bound. Includes the persisted attempt count for scheduler diagnostics.
3300    pub async fn unconverged_segments_below(
3301        &self,
3302        max_unconverged_passes: u32,
3303    ) -> Vec<(String, u32, u32)> {
3304        let quarantined = self.quarantined_segments.lock().clone();
3305        let paused = self.paused_reorder_segments();
3306        let st = self.state.lock().await;
3307        let active_ids = self.active_operations.snapshot();
3308        st.metadata
3309            .segment_metas
3310            .iter()
3311            .filter(|(id, info)| {
3312                !info.bp_converged
3313                    && info.bp_unconverged_passes < max_unconverged_passes
3314                    && !active_ids.contains(*id)
3315                    && !quarantined.contains(*id)
3316                    && !paused.contains(*id)
3317            })
3318            .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
3319            .collect()
3320    }
3321
3322    /// Granularity for a BP pass whose sources are `ids`: `Records` when any
3323    /// source carries unconverged BP debt, `Auto` otherwise.
3324    ///
3325    /// Alignment with the depth budget (docs/block-level-reorder.md): an
3326    /// unconverged segment is owed a deepening pass, and the output of this
3327    /// pass will be marked `bp_converged`. `Auto` would measure the partial
3328    /// pass's residual coherence, potentially take the blockwise path — which
3329    /// cannot deepen record clustering — and end the cascade at partial
3330    /// quality. Only record-level BP discharges the debt.
3331    async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
3332        let st = self.state.lock().await;
3333        let deepening = ids.iter().any(|id| {
3334            st.metadata
3335                .segment_metas
3336                .get(id)
3337                .is_some_and(|info| !info.bp_converged)
3338        });
3339        drop(st);
3340        if deepening {
3341            log::info!(
3342                "[reorder] source BP lineage unconverged — forcing record-level BP (deepening pass)",
3343            );
3344            crate::segment::reorder::BpGranularity::Records
3345        } else {
3346            crate::segment::reorder::BpGranularity::Auto
3347        }
3348    }
3349
3350    /// Reorder a single segment via BP. Returns Ok(true) if reordered, Ok(false) if skipped.
3351    ///
3352    /// Non-blocking: operation ownership prevents conflicts with background merges.
3353    /// Copies unchanged files and rebuilds only the sparse file with reordered BMP data.
3354    pub async fn reorder_single_segment(
3355        self: &Arc<Self>,
3356        seg_id: &str,
3357        rayon_pool: Option<Arc<rayon::ThreadPool>>,
3358        bp_budget: crate::segment::BpBudget,
3359    ) -> Result<bool> {
3360        let source_id = SegmentId::from_hex(seg_id)
3361            .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
3362        if self.quarantined_segments.lock().contains(seg_id) {
3363            return Err(Error::Corruption(format!(
3364                "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
3365                seg_id
3366            )));
3367        }
3368        if self.force_merge_active.load(Ordering::Acquire) > 0 {
3369            log::debug!(
3370                "[optimizer] explicit force merge active, skipping reorder of {}",
3371                seg_id,
3372            );
3373            return Ok(false);
3374        }
3375
3376        // Whole-pass concurrency is independent from Rayon width. One pass
3377        // can already use every configured BP worker; this permit bounds the
3378        // much larger forward-index and rewrite working set across indexes,
3379        // optimizer tasks, and merge-time BP.
3380        let reorder_gate = Arc::clone(&self.reorder_permits);
3381        let _reorder_permit = tokio::select! {
3382            biased;
3383            () = self.active_operations.wait_for_shutdown() => {
3384                return Err(Error::IndexClosed);
3385            }
3386            permit = reorder_gate.acquire(ReorderPriority::Optimizer) => {
3387                permit.map_err(|_| {
3388                    Error::Internal("background reorder scheduler is closed".into())
3389                })?
3390            }
3391        };
3392
3393        let output_id = SegmentId::new();
3394        let output_hex = output_id.to_hex();
3395        let source_ids = [seg_id.to_string()];
3396        let granularity = self.merge_granularity(&source_ids).await;
3397
3398        // Register while holding `state`, matching orphan cleanup's deletion
3399        // barrier. Candidates are scanned ahead of time and can go stale: a
3400        // merge may have consumed this segment since. Its files may even still
3401        // be on disk (deferred deletion under a searcher snapshot) — reordering
3402        // them would re-insert a duplicate copy of docs the merge output holds.
3403        let all_ids = vec![seg_id.to_string(), output_hex];
3404        let (_guard, source_docs) = {
3405            let st = self.state.lock().await;
3406            // Force merge raises this barrier under the same state lock, so
3407            // this second check closes the race with the cheap early check
3408            // above and prevents optimizer starvation between final groups.
3409            if self.force_merge_active.load(Ordering::Acquire) > 0 {
3410                log::debug!(
3411                    "[optimizer] explicit force merge active, skipping reorder of {}",
3412                    seg_id,
3413                );
3414                return Ok(false);
3415            }
3416            let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
3417                log::info!(
3418                    "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
3419                    seg_id
3420                );
3421                self.clear_reorder_retry(seg_id);
3422                return Ok(false);
3423            };
3424
3425            match self.active_operations.try_register(all_ids) {
3426                Some(guard) => (guard, source_meta.num_docs),
3427                None if !self.active_operations.is_accepting() => {
3428                    return Err(Error::IndexClosed);
3429                }
3430                None => {
3431                    log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
3432                    return Ok(false);
3433                }
3434            }
3435        };
3436
3437        // Fail before allocating a forward index or creating output files.
3438        // Missing mandatory files are deterministic and should remove this
3439        // segment from future optimizer scans, not consume the same CPU every
3440        // interval. Other I/O failures remain retryable.
3441        if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
3442            if is_deterministic_source_error(&error) {
3443                self.quarantine_segment(seg_id, &error);
3444            } else if !matches!(&error, Error::IndexClosed) {
3445                self.pause_reorder_retries(seg_id, &error);
3446            }
3447            return Err(error);
3448        }
3449
3450        let mut output_cleanup = self.output_cleanup_guard(output_id);
3451
3452        let reorder_result = crate::segment::reorder::reorder_segment(
3453            self.directory.as_ref(),
3454            &self.schema,
3455            source_id,
3456            output_id,
3457            self.term_cache_blocks,
3458            self.bp_memory_budget_bytes,
3459            bp_budget,
3460            granularity,
3461            rayon_pool,
3462            Some(self.active_operations.cancellation_flag()),
3463        )
3464        .await;
3465        let (new_id, total_docs, bp_converged) = match reorder_result {
3466            Ok(v) => v,
3467            Err(e) => {
3468                // A failed pass may have copied tens of GB before dying;
3469                // delete the uncommitted output before propagating.
3470                self.delete_output_if_unregistered(output_id, "reorder failure")
3471                    .await;
3472                output_cleanup.disarm();
3473                if is_deterministic_source_error(&e) {
3474                    self.quarantine_segment(seg_id, &e);
3475                } else if !matches!(&e, Error::IndexClosed) {
3476                    self.pause_reorder_retries(seg_id, &e);
3477                }
3478                return Err(e);
3479            }
3480        };
3481
3482        // A pass with a depth floor above block granularity has, by
3483        // definition, not converged to block-level order — record it as
3484        // unconverged so the optimizer's deepening ladder revisits it with a
3485        // full-depth (warm-started) pass. Depth caps are only used by the
3486        // optimizer's first pass on large segments.
3487        let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
3488        if let Err(e) = self
3489            .replace_segments(
3490                &[seg_id.to_string()],
3491                new_id,
3492                total_docs,
3493                ReplacementLayout::BpReordered {
3494                    converged: ladder_converged,
3495                },
3496            )
3497            .await
3498        {
3499            self.delete_output_if_unregistered(output_id, "replacement failure")
3500                .await;
3501            output_cleanup.disarm();
3502            if !matches!(&e, Error::IndexClosed) {
3503                self.pause_reorder_retries(seg_id, &e);
3504            }
3505            return Err(e);
3506        }
3507        output_cleanup.disarm();
3508        self.clear_reorder_retry(seg_id);
3509
3510        Ok(true)
3511    }
3512
3513    /// Clean up orphan segment files not registered in metadata.
3514    ///
3515    /// Reads metadata, active-operation ownership, and snapshot-deferred
3516    /// deletions to determine which segments are legitimate. Filesystem
3517    /// deletion is asynchronous; in-flight outputs and retired sources still
3518    /// held by readers are both protected.
3519    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
3520        let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
3521
3522        if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
3523            for entry in entries {
3524                let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
3525                    continue;
3526                };
3527                let Some(rest) = filename.strip_prefix("seg_") else {
3528                    continue;
3529                };
3530                let Some(hex_id) = rest.get(..32) else {
3531                    continue;
3532                };
3533                if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
3534                    continue;
3535                }
3536                orphan_files
3537                    .entry(hex_id.to_ascii_lowercase())
3538                    .or_default()
3539                    .push(entry);
3540            }
3541        }
3542
3543        let mut deleted = 0;
3544        for (hex_id, paths) in &orphan_files {
3545            // Revalidate and atomically claim deletion under the same
3546            // state -> active_operations -> tracker order used by publishers.
3547            // The claim lets us release `state` before filesystem I/O: deleting
3548            // a multi-GB orphan must not freeze commits and snapshot acquisition.
3549            let deletion_guard = {
3550                let st = self.state.lock().await;
3551                if st.metadata.has_segment(hex_id) {
3552                    continue;
3553                }
3554                let Some(guard) = self
3555                    .active_operations
3556                    .try_register(vec![hex_id.to_string()])
3557                else {
3558                    continue;
3559                };
3560                if self.tracker.is_deletion_protected(hex_id) {
3561                    drop(guard);
3562                    continue;
3563                }
3564                guard
3565            };
3566
3567            // Delete what was actually discovered, not only the currently
3568            // known SegmentFiles extensions. This also removes partial files
3569            // left by older formats instead of reporting the same orphan on
3570            // every startup forever.
3571            let results =
3572                futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
3573                    .await;
3574            let removed = results.into_iter().all(|result| match result {
3575                Ok(()) => true,
3576                Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
3577                Err(error) => {
3578                    log::warn!(
3579                        "[segment_cleanup] failed sweeping orphan segment {}: {}",
3580                        hex_id,
3581                        error,
3582                    );
3583                    false
3584                }
3585            });
3586            // Releasing this claim is the deletion barrier. No producer can
3587            // adopt the ID while its files are being removed.
3588            drop(deletion_guard);
3589            if removed {
3590                deleted += 1;
3591                log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
3592            }
3593        }
3594
3595        Ok(deleted)
3596    }
3597}
3598
3599#[cfg(test)]
3600mod tests {
3601    use super::*;
3602    use std::sync::atomic::{AtomicBool, Ordering};
3603
3604    fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
3605        let schema = crate::dsl::SchemaBuilder::default().build();
3606        let metadata = IndexMetadata::new(schema.clone());
3607        Arc::new(SegmentManager::new(
3608            Arc::new(crate::directories::RamDirectory::new()),
3609            Arc::new(schema),
3610            metadata,
3611            Box::new(crate::merge::NoMergePolicy),
3612            0,
3613            1,
3614            Arc::new(Semaphore::new(1)),
3615            None,
3616            1024,
3617            Arc::new(ReorderConcurrencyGate::new(1)),
3618            None,
3619        ))
3620    }
3621
3622    #[test]
3623    fn force_merge_planner_pairs_large_and_small_segments() {
3624        let groups = plan_force_merge_groups(
3625            vec![
3626                ("a".into(), 6),
3627                ("b".into(), 6),
3628                ("c".into(), 4),
3629                ("d".into(), 4),
3630            ],
3631            10,
3632        );
3633
3634        assert_eq!(groups.len(), 2);
3635        assert!(groups.iter().all(|group| group.total_docs == 10));
3636        assert!(groups.iter().all(|group| group.segments.len() == 2));
3637    }
3638
3639    #[test]
3640    fn force_merge_planner_leaves_oversized_segments_alone() {
3641        let groups = plan_force_merge_groups(
3642            vec![
3643                ("oversized".into(), 11),
3644                ("small-a".into(), 5),
3645                ("small-b".into(), 5),
3646            ],
3647            10,
3648        );
3649
3650        assert_eq!(groups.len(), 2);
3651        assert_eq!(groups[0].total_docs, 10);
3652        assert_eq!(groups[0].segments.len(), 2);
3653        assert_eq!(groups[1].total_docs, 11);
3654        assert_eq!(groups[1].segments.len(), 1);
3655    }
3656
3657    #[test]
3658    fn force_merge_planner_never_exceeds_segment_format_limit() {
3659        let groups = plan_force_merge_groups(
3660            vec![
3661                ("large-a".into(), 3_000_000_000),
3662                ("large-b".into(), 2_000_000_000),
3663            ],
3664            u64::from(u32::MAX),
3665        );
3666        assert_eq!(groups.len(), 2);
3667        assert!(
3668            groups
3669                .iter()
3670                .all(|group| group.total_docs <= u64::from(u32::MAX))
3671        );
3672    }
3673
3674    #[test]
3675    fn force_merge_hierarchy_has_one_final_bp_pass() {
3676        assert_eq!(force_merge_output_count(1), 0);
3677        assert_eq!(force_merge_output_count(2), 1);
3678        assert_eq!(force_merge_output_count(64), 1);
3679        assert_eq!(force_merge_output_count(65), 2);
3680        assert_eq!(force_merge_output_count(127), 2);
3681        assert_eq!(force_merge_output_count(128), 3);
3682        assert_eq!(force_merge_output_count(1_000), 16);
3683    }
3684
3685    fn expand_force_merge_node(
3686        hierarchy: &ForceMergeHierarchy,
3687        source_count: usize,
3688        node: usize,
3689        sources: &mut Vec<usize>,
3690    ) {
3691        if node < source_count {
3692            sources.push(node);
3693            return;
3694        }
3695
3696        let step_index = node - source_count;
3697        let step = hierarchy
3698            .steps
3699            .get(step_index)
3700            .expect("merge input must refer to an existing source or output");
3701        for &input in &step.inputs {
3702            assert!(
3703                input < node,
3704                "merge step {step_index} refers to a future output node {input}"
3705            );
3706            expand_force_merge_node(hierarchy, source_count, input, sources);
3707        }
3708    }
3709
3710    #[test]
3711    fn force_merge_hierarchy_has_minimal_valid_arity() {
3712        let source_counts = (2..=1_024).chain([4_095, 4_096, 4_097, 10_000]);
3713
3714        for source_count in source_counts {
3715            let hierarchy = plan_force_merge_hierarchy(source_count);
3716            let output_count = hierarchy.steps.len();
3717
3718            assert!(
3719                hierarchy
3720                    .steps
3721                    .iter()
3722                    .all(|step| (2..=FORCE_MERGE_MAX_FAN_IN).contains(&step.inputs.len())),
3723                "invalid merge arity for {source_count} sources"
3724            );
3725            assert!(
3726                source_count <= 1 + output_count * (FORCE_MERGE_MAX_FAN_IN - 1),
3727                "{output_count} outputs cannot reduce {source_count} sources"
3728            );
3729            assert!(
3730                output_count == 1
3731                    || source_count > 1 + (output_count - 1) * (FORCE_MERGE_MAX_FAN_IN - 1),
3732                "{output_count} outputs are not minimal for {source_count} sources"
3733            );
3734            assert_eq!(output_count, force_merge_output_count(source_count));
3735        }
3736    }
3737
3738    #[test]
3739    fn force_merge_hierarchy_preserves_exact_source_order() {
3740        for source_count in [2, 3, 63, 64, 65, 66, 126, 127, 128, 129, 1_000, 4_097] {
3741            let hierarchy = plan_force_merge_hierarchy(source_count);
3742            let mut sources = Vec::with_capacity(source_count);
3743            expand_force_merge_node(&hierarchy, source_count, hierarchy.root, &mut sources);
3744            assert_eq!(
3745                sources,
3746                (0..source_count).collect::<Vec<_>>(),
3747                "source order changed for {source_count} sources"
3748            );
3749        }
3750    }
3751
3752    fn force_merge_rewrite_cost(source_count: usize) -> usize {
3753        let hierarchy = plan_force_merge_hierarchy(source_count);
3754        let mut node_weights = vec![1usize; source_count];
3755        let mut rewrite_cost = 0usize;
3756
3757        for (step_index, step) in hierarchy.steps.iter().enumerate() {
3758            let output = source_count + step_index;
3759            let output_weight = step
3760                .inputs
3761                .iter()
3762                .map(|&input| {
3763                    assert!(
3764                        input < output,
3765                        "merge step {step_index} refers to future output {input}"
3766                    );
3767                    node_weights[input]
3768                })
3769                .sum::<usize>();
3770            rewrite_cost += output_weight;
3771            node_weights.push(output_weight);
3772        }
3773
3774        assert_eq!(node_weights[hierarchy.root], source_count);
3775        rewrite_cost
3776    }
3777
3778    #[test]
3779    fn force_merge_hierarchy_avoids_growing_prefix_rewrites() {
3780        assert_eq!(force_merge_rewrite_cost(65), 67);
3781        assert_eq!(force_merge_rewrite_cost(1_000), 1_951);
3782    }
3783
3784    #[test]
3785    fn block_copy_carries_bp_debt_without_spending_an_attempt() {
3786        assert_eq!(
3787            replacement_bp_state(true, 3, ReplacementLayout::BlockCopy),
3788            (false, false, 3),
3789        );
3790        assert_eq!(
3791            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: false },),
3792            (true, false, 4),
3793        );
3794        assert_eq!(
3795            replacement_bp_state(true, 3, ReplacementLayout::BpReordered { converged: true },),
3796            (true, true, 0),
3797        );
3798    }
3799
3800    #[tokio::test]
3801    async fn force_merge_reconsiders_outputs_after_a_held_source_releases() {
3802        let mut schema_builder = crate::dsl::SchemaBuilder::default();
3803        let field = schema_builder.add_text_field("text", true, true);
3804        let schema = schema_builder.build();
3805        let directory = crate::directories::RamDirectory::new();
3806        let config = crate::index::IndexConfig {
3807            num_indexing_threads: 1,
3808            merge_policy: Box::new(crate::merge::NoMergePolicy),
3809            ..Default::default()
3810        };
3811        let mut writer = crate::index::IndexWriter::create(directory, schema, config)
3812            .await
3813            .unwrap();
3814        for value in ["one", "two", "three"] {
3815            let mut document = crate::dsl::Document::new();
3816            document.add_text(field, value);
3817            writer.add_document(document).unwrap();
3818            writer.commit().await.unwrap();
3819        }
3820
3821        let manager = Arc::clone(writer.segment_manager());
3822        let held_id = manager.get_segment_ids().await.pop().unwrap();
3823        let mut held = Some(
3824            manager
3825                .active_operations
3826                .try_register(vec![held_id])
3827                .unwrap(),
3828        );
3829        let batches = Arc::new(AtomicUsize::new(0));
3830        let batch_count = Arc::clone(&batches);
3831        writer
3832            .force_merge_with_snapshot_refresh(move || {
3833                let refresh = batch_count.fetch_add(1, Ordering::Relaxed) + 1;
3834                // Refresh #1 follows the startup drain. Keep the source held
3835                // until refresh #2, after the two free sources were replaced.
3836                if refresh == 2 {
3837                    drop(held.take());
3838                }
3839                std::future::ready(Ok(()))
3840            })
3841            .await
3842            .unwrap();
3843
3844        assert_eq!(manager.get_segment_ids().await.len(), 1);
3845        assert_eq!(
3846            batches.load(Ordering::Relaxed),
3847            4,
3848            "initial/final refreshes plus two replacements are required after the held source releases"
3849        );
3850    }
3851
3852    #[tokio::test]
3853    async fn force_merge_does_not_hold_global_capacity_while_vector_update_pauses_claims() {
3854        let mut schema_builder = crate::dsl::SchemaBuilder::default();
3855        schema_builder.set_reorder_on_merge(true);
3856        let schema = schema_builder.build();
3857        let mut metadata = IndexMetadata::new(schema.clone());
3858        metadata.add_segment("00000000000000000000000000000001".into(), 1);
3859        metadata.add_segment("00000000000000000000000000000002".into(), 1);
3860
3861        let global_merge_permits = Arc::new(Semaphore::new(1));
3862        let manager = Arc::new(SegmentManager::new(
3863            Arc::new(crate::directories::RamDirectory::new()),
3864            Arc::new(schema),
3865            metadata,
3866            Box::new(crate::merge::NoMergePolicy),
3867            0,
3868            1,
3869            Arc::clone(&global_merge_permits),
3870            None,
3871            1024,
3872            Arc::new(ReorderConcurrencyGate::new(1)),
3873            None,
3874        ));
3875
3876        // Vector-generation staging rejects ordinary lifecycle claims while
3877        // acquiring the shared global merge permit separately for each source.
3878        // Force merge must wait before capacity admission; retaining the only
3879        // global slot here would deadlock both operations between sources.
3880        manager.active_operations.pause_non_indexing();
3881        let force_merge = {
3882            let manager = Arc::clone(&manager);
3883            tokio::spawn(async move { manager.force_merge().await })
3884        };
3885        tokio::time::timeout(std::time::Duration::from_secs(1), async {
3886            while manager.force_merge_conflict_retries.load(Ordering::Relaxed) == 0 {
3887                tokio::task::yield_now().await;
3888            }
3889        })
3890        .await
3891        .expect("force merge never reached the paused group claim");
3892
3893        assert_eq!(
3894            global_merge_permits.available_permits(),
3895            1,
3896            "force merge retained global capacity while vector staging blocked group ownership"
3897        );
3898
3899        force_merge.abort();
3900        let _ = force_merge.await;
3901        manager.active_operations.resume_non_indexing();
3902    }
3903
3904    #[test]
3905    fn output_cleanup_guard_runs_during_panic_unwind() {
3906        let cleaned = Arc::new(AtomicBool::new(false));
3907        let cleaned_in_callback = Arc::clone(&cleaned);
3908        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3909            cleaned_in_callback.store(true, Ordering::SeqCst);
3910        });
3911
3912        let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
3913            let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3914            panic!("simulated reorder panic");
3915        }));
3916
3917        assert!(result.is_err());
3918        assert!(
3919            cleaned.load(Ordering::SeqCst),
3920            "partial output cleanup must run during unwind"
3921        );
3922    }
3923
3924    #[test]
3925    fn output_cleanup_guard_disarms_after_commit() {
3926        let cleaned = Arc::new(AtomicBool::new(false));
3927        let cleaned_in_callback = Arc::clone(&cleaned);
3928        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
3929            cleaned_in_callback.store(true, Ordering::SeqCst);
3930        });
3931
3932        {
3933            let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
3934            guard.disarm();
3935        }
3936
3937        assert!(!cleaned.load(Ordering::SeqCst));
3938    }
3939
3940    #[test]
3941    fn test_active_operation_guard_releases_ownership() {
3942        let active = Arc::new(ActiveSegmentOperations::new());
3943        {
3944            let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3945            let snap = active.snapshot();
3946            assert!(snap.contains("a"));
3947            assert!(snap.contains("b"));
3948        }
3949        assert!(active.snapshot().is_empty());
3950    }
3951
3952    #[test]
3953    fn test_non_overlapping_operations_can_run_concurrently() {
3954        let active = Arc::new(ActiveSegmentOperations::new());
3955        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3956        let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
3957        let snap = active.snapshot();
3958        assert_eq!(snap.len(), 4);
3959
3960        drop(first);
3961        let snap = active.snapshot();
3962        assert_eq!(snap.len(), 2);
3963        assert!(snap.contains("c"));
3964        assert!(snap.contains("d"));
3965    }
3966
3967    #[test]
3968    fn test_overlapping_operation_is_rejected_until_release() {
3969        let active = Arc::new(ActiveSegmentOperations::new());
3970        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
3971        assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
3972        drop(first);
3973        assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
3974    }
3975
3976    #[test]
3977    fn test_active_operation_snapshot() {
3978        let active = Arc::new(ActiveSegmentOperations::new());
3979        let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
3980        let snap = active.snapshot();
3981        assert!(snap.contains("x"));
3982        assert!(snap.contains("y"));
3983        assert!(!snap.contains("z"));
3984    }
3985
3986    #[tokio::test]
3987    async fn operation_barrier_ignores_producers_started_after_snapshot() {
3988        let active = Arc::new(ActiveSegmentOperations::new());
3989        let before_gate = active.try_register(vec!["old".into()]).unwrap();
3990        let (barrier, parked_indexing) = active.draining_operation_tokens_snapshot();
3991        assert_eq!(parked_indexing, 0);
3992        let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
3993
3994        let waiter = {
3995            let active = Arc::clone(&active);
3996            tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
3997        };
3998        tokio::task::yield_now().await;
3999        assert!(!waiter.is_finished());
4000
4001        drop(before_gate);
4002        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4003            .await
4004            .expect("pre-gate operation barrier was starved by a post-gate producer")
4005            .unwrap();
4006        assert!(active.snapshot().contains("new-flat"));
4007        drop(after_gate);
4008    }
4009
4010    #[tokio::test]
4011    async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
4012        let manager = lifecycle_test_manager();
4013        manager
4014            .trained
4015            .store(Some(Arc::new(TrainedVectorStructures {
4016                centroids: rustc_hash::FxHashMap::default(),
4017                binary_quantizers: rustc_hash::FxHashMap::default(),
4018                ..Default::default()
4019            })));
4020
4021        let guard = manager.begin_vector_artifact_update().await.unwrap();
4022        assert!(
4023            manager.trained().is_some(),
4024            "search readers keep the last fully validated generation"
4025        );
4026        assert!(
4027            manager.trained_for_segment_build().is_none(),
4028            "new segment producers must stay flat during an artifact update"
4029        );
4030
4031        let detached_transaction_guard = guard.clone();
4032        drop(guard);
4033        assert!(
4034            manager.trained_for_segment_build().is_none(),
4035            "a detached lifecycle transaction must retain the producer gate after request cancellation"
4036        );
4037        drop(detached_transaction_guard);
4038        assert!(manager.trained_for_segment_build().is_some());
4039    }
4040
4041    #[tokio::test]
4042    async fn artifact_update_pauses_lifecycle_rewrites_but_allows_flat_indexing() {
4043        let manager = lifecycle_test_manager();
4044        let guard = manager.begin_vector_artifact_update().await.unwrap();
4045        assert!(
4046            manager
4047                .active_operations
4048                .try_register(vec!["merge".into()])
4049                .is_none(),
4050            "ordinary merge/reorder work must not change staged sources"
4051        );
4052        let indexing = manager
4053            .active_operations
4054            .try_register_indexing(vec!["fresh".into()])
4055            .expect("indexing remains available in flat mode");
4056        drop(indexing);
4057
4058        drop(guard);
4059        assert!(!manager.vector_artifact_update.load(Ordering::Acquire));
4060        assert!(
4061            manager
4062                .active_operations
4063                .try_register(vec!["merge".into()])
4064                .is_some()
4065        );
4066    }
4067
4068    #[tokio::test]
4069    async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
4070        let active = Arc::new(ActiveSegmentOperations::new());
4071        let guard = active.try_register(vec!["live".into()]).unwrap();
4072        let cancellation = active.cancellation_flag();
4073        active.stop_accepting();
4074        assert!(cancellation.load(Ordering::Acquire));
4075        assert!(active.try_register(vec!["new".into()]).is_none());
4076
4077        let waiter = {
4078            let active = Arc::clone(&active);
4079            tokio::spawn(async move { active.wait_until_idle().await })
4080        };
4081        tokio::task::yield_now().await;
4082        assert!(!waiter.is_finished());
4083        drop(guard);
4084        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4085            .await
4086            .expect("shutdown waiter missed the final guard notification")
4087            .unwrap();
4088    }
4089
4090    #[tokio::test]
4091    async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
4092        let manager = lifecycle_test_manager();
4093        let started = Arc::new(Semaphore::new(0));
4094        let release = Arc::new(Semaphore::new(0));
4095        let completed = Arc::new(AtomicBool::new(false));
4096
4097        let request = {
4098            let manager = Arc::clone(&manager);
4099            let started = Arc::clone(&started);
4100            let release = Arc::clone(&release);
4101            let completed = Arc::clone(&completed);
4102            tokio::spawn(async move {
4103                manager
4104                    .run_lifecycle_transaction(async move {
4105                        started.add_permits(1);
4106                        let _permit = release.acquire().await.unwrap();
4107                        completed.store(true, Ordering::Release);
4108                        Ok(())
4109                    })
4110                    .await
4111            })
4112        };
4113
4114        let _started = started.acquire().await.unwrap();
4115        request.abort();
4116        assert!(request.await.unwrap_err().is_cancelled());
4117        release.add_permits(1);
4118
4119        manager.begin_shutdown();
4120        tokio::time::timeout(
4121            std::time::Duration::from_secs(1),
4122            manager.wait_for_shutdown(),
4123        )
4124        .await
4125        .expect("shutdown did not drain detached lifecycle transaction");
4126        assert!(completed.load(Ordering::Acquire));
4127    }
4128
4129    #[tokio::test]
4130    async fn unconverged_scheduler_stops_at_the_lineage_limit() {
4131        let manager = lifecycle_test_manager();
4132        {
4133            let mut state = manager.state.lock().await;
4134            state.metadata.add_segment_meta(
4135                "eligible".into(),
4136                SegmentMetaInfo {
4137                    num_docs: 10,
4138                    ancestors: Vec::new(),
4139                    generation: 1,
4140                    reordered: true,
4141                    bp_converged: false,
4142                    bp_unconverged_passes: 2,
4143                },
4144            );
4145            state.metadata.add_segment_meta(
4146                "at-limit".into(),
4147                SegmentMetaInfo {
4148                    num_docs: 20,
4149                    ancestors: Vec::new(),
4150                    generation: 1,
4151                    reordered: true,
4152                    bp_converged: false,
4153                    bp_unconverged_passes: 3,
4154                },
4155            );
4156            state.metadata.add_segment_meta(
4157                "carried-debt".into(),
4158                SegmentMetaInfo {
4159                    num_docs: 15,
4160                    ancestors: Vec::new(),
4161                    generation: 2,
4162                    reordered: false,
4163                    bp_converged: false,
4164                    bp_unconverged_passes: 2,
4165                },
4166            );
4167            state.metadata.add_segment_meta(
4168                "carried-debt-at-limit".into(),
4169                SegmentMetaInfo {
4170                    num_docs: 25,
4171                    ancestors: Vec::new(),
4172                    generation: 2,
4173                    reordered: false,
4174                    bp_converged: false,
4175                    bp_unconverged_passes: 3,
4176                },
4177            );
4178            state.metadata.add_segment_meta(
4179                "converged".into(),
4180                SegmentMetaInfo {
4181                    num_docs: 30,
4182                    ancestors: Vec::new(),
4183                    generation: 1,
4184                    reordered: true,
4185                    bp_converged: true,
4186                    bp_unconverged_passes: 0,
4187                },
4188            );
4189            state.metadata.add_segment("fresh".into(), 40);
4190        }
4191
4192        assert_eq!(
4193            manager.unreordered_segments().await,
4194            vec![("fresh".into(), 40)],
4195            "a block-copy output with BP debt is not a fresh first-pass candidate",
4196        );
4197        let mut eligible = manager.unconverged_segments_below(3).await;
4198        eligible.sort_unstable();
4199        assert_eq!(
4200            eligible,
4201            vec![("carried-debt".into(), 15, 2), ("eligible".into(), 10, 2),],
4202        );
4203        assert!(manager.unconverged_segments_below(0).await.is_empty());
4204    }
4205
4206    #[test]
4207    fn merge_retry_backoff_is_exponential_and_capped() {
4208        assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
4209        assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
4210        assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
4211        assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
4212    }
4213
4214    #[test]
4215    fn only_deterministic_source_errors_are_quarantined() {
4216        assert!(is_deterministic_source_error(&Error::Corruption(
4217            "bad footer".into()
4218        )));
4219        assert!(is_deterministic_source_error(&Error::Io(
4220            std::io::Error::from(std::io::ErrorKind::NotFound)
4221        )));
4222        assert!(!is_deterministic_source_error(&Error::Io(
4223            std::io::Error::from(std::io::ErrorKind::TimedOut)
4224        )));
4225        assert!(!is_deterministic_source_error(&Error::Io(
4226            std::io::Error::from(std::io::ErrorKind::PermissionDenied)
4227        )));
4228    }
4229
4230    #[test]
4231    fn transient_reorder_failure_is_backed_off_until_cleared() {
4232        let manager = lifecycle_test_manager();
4233        manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
4234        assert!(manager.paused_reorder_segments().contains("source"));
4235        manager.clear_reorder_retry("source");
4236        assert!(!manager.paused_reorder_segments().contains("source"));
4237    }
4238
4239    /// Fails `exists` with the transient I/O error class that sends a
4240    /// background merge into its generic retry backoff (not source quarantine).
4241    #[derive(Default)]
4242    struct FailingExistsDirectory(crate::directories::RamDirectory);
4243
4244    #[async_trait::async_trait]
4245    impl crate::directories::Directory for FailingExistsDirectory {
4246        async fn exists(&self, _path: &std::path::Path) -> std::io::Result<bool> {
4247            Err(std::io::Error::from(std::io::ErrorKind::TimedOut))
4248        }
4249
4250        async fn file_size(&self, path: &std::path::Path) -> std::io::Result<u64> {
4251            self.0.file_size(path).await
4252        }
4253
4254        async fn open_read(
4255            &self,
4256            path: &std::path::Path,
4257        ) -> std::io::Result<crate::directories::FileHandle> {
4258            self.0.open_read(path).await
4259        }
4260
4261        async fn read_range(
4262            &self,
4263            path: &std::path::Path,
4264            range: std::ops::Range<u64>,
4265        ) -> std::io::Result<crate::directories::OwnedBytes> {
4266            self.0.read_range(path, range).await
4267        }
4268
4269        async fn list_files(
4270            &self,
4271            prefix: &std::path::Path,
4272        ) -> std::io::Result<Vec<std::path::PathBuf>> {
4273            self.0.list_files(prefix).await
4274        }
4275
4276        async fn open_lazy(
4277            &self,
4278            path: &std::path::Path,
4279        ) -> std::io::Result<crate::directories::FileHandle> {
4280            self.0.open_lazy(path).await
4281        }
4282    }
4283
4284    #[async_trait::async_trait]
4285    impl crate::directories::DirectoryWriter for FailingExistsDirectory {
4286        async fn write(&self, path: &std::path::Path, data: &[u8]) -> std::io::Result<()> {
4287            self.0.write(path, data).await
4288        }
4289
4290        async fn delete(&self, path: &std::path::Path) -> std::io::Result<()> {
4291            self.0.delete(path).await
4292        }
4293
4294        async fn rename(
4295            &self,
4296            from: &std::path::Path,
4297            to: &std::path::Path,
4298        ) -> std::io::Result<()> {
4299            self.0.rename(from, to).await
4300        }
4301
4302        async fn sync(&self) -> std::io::Result<()> {
4303            self.0.sync().await
4304        }
4305
4306        async fn streaming_writer(
4307            &self,
4308            path: &std::path::Path,
4309        ) -> std::io::Result<Box<dyn crate::directories::StreamingWriter>> {
4310            self.0.streaming_writer(path).await
4311        }
4312    }
4313
4314    #[derive(Debug, Clone)]
4315    struct MergeEverythingPolicy;
4316
4317    impl MergePolicy for MergeEverythingPolicy {
4318        fn find_merges(&self, segments: &[SegmentInfo]) -> Vec<crate::merge::MergeCandidate> {
4319            if segments.len() < 2 {
4320                return Vec::new();
4321            }
4322            vec![crate::merge::MergeCandidate {
4323                segment_ids: segments.iter().map(|s| s.id.clone()).collect(),
4324            }]
4325        }
4326
4327        fn clone_box(&self) -> Box<dyn MergePolicy> {
4328            Box::new(self.clone())
4329        }
4330    }
4331
4332    #[tokio::test]
4333    async fn artifact_update_rejects_built_uncommitted_indexing_segments() {
4334        let manager = lifecycle_test_manager();
4335        // Simulates a memory-budget mid-cycle segment build whose guard is
4336        // parked inside a PreparedSegment: only a later commit releases this
4337        // token, and that commit can be blocked on the very caller of the
4338        // artifact update (writer write lock / &mut self).
4339        let parked_indexing = manager
4340            .protect_new_segment("00000000000000000000000000000abc".into())
4341            .unwrap();
4342
4343        let error = tokio::time::timeout(
4344            std::time::Duration::from_secs(2),
4345            manager.begin_vector_artifact_update(),
4346        )
4347        .await
4348        .expect("begin_vector_artifact_update deadlocked on a built-but-uncommitted segment")
4349        .err()
4350        .expect("an old-generation prepared segment must block artifact replacement")
4351        .to_string();
4352        assert!(error.contains("built but uncommitted"), "{error}");
4353        assert!(
4354            !manager.vector_artifact_update.load(Ordering::Acquire),
4355            "a rejected update must release the producer gate"
4356        );
4357
4358        drop(parked_indexing);
4359
4360        let guard = manager
4361            .begin_vector_artifact_update()
4362            .await
4363            .expect("artifact update should succeed after the pending generation is resolved");
4364        drop(guard);
4365    }
4366
4367    #[tokio::test]
4368    async fn artifact_update_still_drains_preexisting_lifecycle_operations() {
4369        let manager = lifecycle_test_manager();
4370        let merge_like = manager
4371            .active_operations
4372            .try_register(vec!["merge-source".into()])
4373            .unwrap();
4374
4375        let waiter = {
4376            let manager = Arc::clone(&manager);
4377            tokio::spawn(async move { manager.begin_vector_artifact_update().await })
4378        };
4379        for _ in 0..8 {
4380            tokio::task::yield_now().await;
4381        }
4382        assert!(
4383            !waiter.is_finished(),
4384            "artifact update must drain merge/reorder producers that may hold the previous generation"
4385        );
4386
4387        drop(merge_like);
4388        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
4389            .await
4390            .expect("artifact update missed the lifecycle guard release")
4391            .unwrap()
4392            .unwrap();
4393    }
4394
4395    #[tokio::test]
4396    async fn cancelled_merge_drain_returns_unawaited_handles_to_shared_state() {
4397        let manager = lifecycle_test_manager();
4398        let release = Arc::new(Semaphore::new(0));
4399        let merge_task = {
4400            let release = Arc::clone(&release);
4401            tokio::spawn(async move {
4402                let _permit = release.acquire().await.unwrap();
4403            })
4404        };
4405        manager.merge_handles.lock().push(merge_task);
4406
4407        let waiter = {
4408            let manager = Arc::clone(&manager);
4409            tokio::spawn(async move { manager.wait_for_all_merges().await })
4410        };
4411        for _ in 0..8 {
4412            tokio::task::yield_now().await;
4413        }
4414        assert!(!waiter.is_finished());
4415        // Simulates tonic dropping a force_merge/reorder RPC future at the
4416        // JoinHandle await when the client disconnects.
4417        waiter.abort();
4418        let join_error = waiter.await.unwrap_err();
4419        assert!(join_error.is_cancelled());
4420
4421        assert!(
4422            !manager.merge_handles.lock().is_empty(),
4423            "cancelled drain detached an in-flight merge from shutdown/force-merge tracking"
4424        );
4425
4426        // A later drain must still see and await the real in-flight merge.
4427        release.add_permits(1);
4428        tokio::time::timeout(
4429            std::time::Duration::from_secs(1),
4430            manager.wait_for_all_merges(),
4431        )
4432        .await
4433        .expect("subsequent drain missed the reinserted merge handle");
4434        assert!(manager.merge_handles.lock().is_empty());
4435    }
4436
4437    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4438    async fn force_merge_conflict_retry_backs_off_instead_of_busy_spinning() {
4439        let manager = lifecycle_test_manager();
4440        {
4441            let mut state = manager.state.lock().await;
4442            state
4443                .metadata
4444                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4445            state
4446                .metadata
4447                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4448        }
4449        // A background reorder (or a concurrent force-merge) owns one segment
4450        // in the batch but never appears in merge_handles.
4451        let reorder_like = manager
4452            .active_operations
4453            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4454            .unwrap();
4455
4456        let force_merge = {
4457            let manager = Arc::clone(&manager);
4458            tokio::spawn(async move { manager.force_merge().await })
4459        };
4460
4461        tokio::time::sleep(std::time::Duration::from_millis(600)).await;
4462        let retries = manager.force_merge_conflict_retries.load(Ordering::Relaxed);
4463        assert!(
4464            retries >= 1,
4465            "force_merge never observed the conflicting owner (retries={retries})"
4466        );
4467        assert!(
4468            retries < 20,
4469            "force_merge busy-spun on a conflict that is not a tracked merge (retries={retries})"
4470        );
4471
4472        drop(reorder_like);
4473        // With the conflict gone the loop proceeds; the batch then fails fast
4474        // in do_merge (the test IDs have no files), proving the loop exited.
4475        let result = tokio::time::timeout(std::time::Duration::from_secs(10), force_merge)
4476            .await
4477            .expect("force_merge kept spinning after the conflicting owner released")
4478            .unwrap();
4479        assert!(result.is_err());
4480    }
4481
4482    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4483    async fn force_merge_routes_around_segments_held_by_reorder() {
4484        let manager = lifecycle_test_manager();
4485        {
4486            let mut state = manager.state.lock().await;
4487            state
4488                .metadata
4489                .add_segment("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into(), 10);
4490            state
4491                .metadata
4492                .add_segment("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb".into(), 10);
4493            state
4494                .metadata
4495                .add_segment("cccccccccccccccccccccccccccccccc".into(), 10);
4496        }
4497        // A background reorder owns one segment and holds it for the whole
4498        // test (in prod: a BP pass runs for minutes while force_merge spins).
4499        let _reorder_like = manager
4500            .active_operations
4501            .try_register(vec!["aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa".into()])
4502            .unwrap();
4503
4504        // Regression: force_merge used to rebuild the identical smallest-N
4505        // batch (including the held segment) every 100ms and retry-log
4506        // forever. It must instead skip the held segment and immediately
4507        // make progress on the two free ones — reaching do_merge (which
4508        // fails fast here: the test IDs have no files) proves the batch was
4509        // built without the held segment while the reorder is STILL active.
4510        let result = tokio::time::timeout(std::time::Duration::from_secs(5), {
4511            let manager = Arc::clone(&manager);
4512            async move { manager.force_merge().await }
4513        })
4514        .await
4515        .expect("force_merge livelocked on a segment held by an active reorder");
4516        assert!(result.is_err(), "fake segment files must fail the merge");
4517
4518        assert_eq!(
4519            manager.force_merge_conflict_retries.load(Ordering::Relaxed),
4520            0,
4521            "batch built from the ownership snapshot must not collide with the held segment"
4522        );
4523    }
4524
4525    #[tokio::test]
4526    async fn merge_failure_retry_backoff_does_not_stall_merge_waiters() {
4527        let schema = crate::dsl::SchemaBuilder::default().build();
4528        let mut metadata = IndexMetadata::new(schema.clone());
4529        metadata.add_segment("00000000000000000000000000000001".into(), 10);
4530        metadata.add_segment("00000000000000000000000000000002".into(), 10);
4531        let manager = Arc::new(SegmentManager::new(
4532            Arc::new(FailingExistsDirectory::default()),
4533            Arc::new(schema),
4534            metadata,
4535            Box::new(MergeEverythingPolicy),
4536            0,
4537            1,
4538            Arc::new(Semaphore::new(1)),
4539            None,
4540            1024,
4541            Arc::new(ReorderConcurrencyGate::new(1)),
4542            None,
4543        ));
4544
4545        // Spawns a background merge that fails with a transient I/O error and
4546        // arms the 30s..30min retry backoff.
4547        manager.maybe_merge().await;
4548
4549        tokio::time::timeout(
4550            std::time::Duration::from_secs(5),
4551            manager.wait_for_all_merges(),
4552        )
4553        .await
4554        .expect("wait_for_all_merges stalled behind a pure retry-backoff timer");
4555        assert!(
4556            manager.merge_retry_is_paused(),
4557            "the failed merge should have armed the retry backoff"
4558        );
4559
4560        // Shutdown still drains the pending backoff wakeup deterministically.
4561        manager.begin_shutdown();
4562        tokio::time::timeout(
4563            std::time::Duration::from_secs(5),
4564            manager.wait_for_shutdown(),
4565        )
4566        .await
4567        .expect("shutdown did not drain the merge retry wakeup task");
4568    }
4569}