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