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, Ordering};
46use std::sync::{Arc, OnceLock};
47
48use arc_swap::ArcSwapOption;
49use tokio::sync::Mutex as AsyncMutex;
50use tokio::sync::{Notify, Semaphore};
51use tokio::task::JoinHandle;
52
53use crate::directories::DirectoryWriter;
54use crate::error::{Error, Result};
55use crate::index::{IndexMetadata, 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    next_operation_token: u64,
77    accepting: bool,
78}
79
80struct ActiveSegmentOperations {
81    inner: parking_lot::Mutex<ActiveOperationState>,
82    idle: Notify,
83    shutdown: Notify,
84}
85
86impl ActiveSegmentOperations {
87    fn new() -> Self {
88        Self {
89            inner: parking_lot::Mutex::new(ActiveOperationState {
90                segment_ids: HashSet::new(),
91                operation_tokens: HashSet::new(),
92                next_operation_token: 0,
93                accepting: true,
94            }),
95            idle: Notify::new(),
96            shutdown: Notify::new(),
97        }
98    }
99
100    /// Try to claim IDs for an operation. Returns a guard on success, `None`
101    /// if any requested ID is already owned by another active operation.
102    fn try_register(self: &Arc<Self>, segment_ids: Vec<String>) -> Option<SegmentOperationGuard> {
103        let mut inner = self.inner.lock();
104        if !inner.accepting {
105            log::debug!("[segment_lifecycle] rejected operation during shutdown");
106            return None;
107        }
108        // Check for overlap with any active lifecycle operation.
109        for id in &segment_ids {
110            if inner.segment_ids.contains(id) {
111                log::debug!(
112                    "[segment_lifecycle] rejected: {} overlaps with an active operation ({} active IDs)",
113                    id,
114                    inner.segment_ids.len()
115                );
116                return None;
117            }
118        }
119        log::debug!(
120            "[segment_lifecycle] registered {} IDs (total active: {})",
121            segment_ids.len(),
122            inner.segment_ids.len() + segment_ids.len()
123        );
124        let operation_token = inner.next_operation_token;
125        let next_operation_token = operation_token.checked_add(1)?;
126        for id in &segment_ids {
127            inner.segment_ids.insert(id.clone());
128        }
129        inner.next_operation_token = next_operation_token;
130        inner.operation_tokens.insert(operation_token);
131        Some(SegmentOperationGuard {
132            active_operations: Arc::clone(self),
133            segment_ids,
134            operation_token,
135        })
136    }
137
138    /// Snapshot of all IDs owned by active operations.
139    fn snapshot(&self) -> HashSet<String> {
140        self.inner.lock().segment_ids.clone()
141    }
142
143    /// Exact operation identities active at one instant. Unlike segment IDs,
144    /// tokens cannot be reused by a later retry, so an artifact-update barrier
145    /// can drain only pre-gate producers without being starved by new flat
146    /// producers.
147    fn operation_tokens_snapshot(&self) -> HashSet<u64> {
148        self.inner.lock().operation_tokens.clone()
149    }
150
151    /// Atomically prevent new lifecycle work from starting. Existing guards
152    /// remain valid and can be drained with [`Self::wait_until_idle`].
153    fn stop_accepting(&self) {
154        let mut inner = self.inner.lock();
155        inner.accepting = false;
156        self.shutdown.notify_waiters();
157        if inner.segment_ids.is_empty() {
158            self.idle.notify_waiters();
159        }
160    }
161
162    fn is_accepting(&self) -> bool {
163        self.inner.lock().accepting
164    }
165
166    /// Wait until every operation that started before shutdown has released
167    /// its ownership. Register/check and notification are ordered to avoid a
168    /// missed wakeup between observing a non-empty set and awaiting.
169    async fn wait_until_idle(&self) {
170        loop {
171            let notified = self.idle.notified();
172            if self.inner.lock().segment_ids.is_empty() {
173                return;
174            }
175            notified.await;
176        }
177    }
178
179    async fn wait_until_operations_finish(&self, operations: &HashSet<u64>) {
180        while !operations.is_empty() {
181            let notified = self.idle.notified();
182            if self.inner.lock().operation_tokens.is_disjoint(operations) {
183                return;
184            }
185            notified.await;
186        }
187    }
188
189    /// Resolve when shutdown starts, without missing a notification between
190    /// checking the state and registering the waiter.
191    async fn wait_for_shutdown(&self) {
192        loop {
193            let notified = self.shutdown.notified();
194            if !self.inner.lock().accepting {
195                return;
196            }
197            notified.await;
198        }
199    }
200}
201
202/// RAII ownership of segment IDs used by an active lifecycle operation.
203/// Dropping on success, error, cancellation, or panic makes abandoned outputs
204/// eligible for sweeping automatically.
205pub(crate) struct SegmentOperationGuard {
206    active_operations: Arc<ActiveSegmentOperations>,
207    segment_ids: Vec<String>,
208    operation_token: u64,
209}
210
211impl Drop for SegmentOperationGuard {
212    fn drop(&mut self) {
213        let mut inner = self.active_operations.inner.lock();
214        for id in &self.segment_ids {
215            inner.segment_ids.remove(id);
216        }
217        inner.operation_tokens.remove(&self.operation_token);
218        // Token barriers need notification on every completion, not only the
219        // transition to complete global idleness.
220        self.active_operations.idle.notify_waiters();
221        if inner.segment_ids.is_empty() {
222            debug_assert!(inner.operation_tokens.is_empty());
223        }
224    }
225}
226
227/// Exclusive gate for an index-level trained-vector artifact update.
228///
229/// Segment producers consult this gate before capturing the current trained
230/// structures. Once the gate is raised, new producers deliberately emit flat
231/// vector data; waiting for already-active producers to drain then guarantees
232/// that no segment using the previous generation can appear after a rebuild
233/// safety check.
234struct VectorArtifactUpdateLease {
235    updating: Arc<AtomicBool>,
236}
237
238impl Drop for VectorArtifactUpdateLease {
239    fn drop(&mut self) {
240        self.updating.store(false, Ordering::Release);
241    }
242}
243
244#[derive(Clone)]
245pub(crate) struct VectorArtifactUpdateGuard {
246    _lease: Arc<VectorArtifactUpdateLease>,
247}
248
249/// Merge-time/manual BP pools are shared by every index in this process.
250/// A pool per `SegmentManager` multiplied a 96-core host into two 48-thread
251/// merge pools plus the optimizer pool (200+ process threads in production).
252static BACKGROUND_CPU_POOL: OnceLock<Arc<rayon::ThreadPool>> = OnceLock::new();
253
254const MERGE_RETRY_BASE_DELAY: std::time::Duration = std::time::Duration::from_secs(30);
255const MERGE_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(30 * 60);
256
257#[derive(Default)]
258struct MergeRetryState {
259    retry_after: Option<std::time::Instant>,
260    consecutive_failures: u32,
261}
262
263fn merge_retry_delay(consecutive_failures: u32) -> std::time::Duration {
264    let shift = consecutive_failures.saturating_sub(1).min(16);
265    MERGE_RETRY_BASE_DELAY
266        .checked_mul(1u32 << shift)
267        .unwrap_or(MERGE_RETRY_MAX_DELAY)
268        .min(MERGE_RETRY_MAX_DELAY)
269}
270
271/// Spawn and register auxiliary lifecycle work as one synchronous operation.
272///
273/// Registering *after* `spawn` left a small deletion race: shutdown could
274/// observe an empty handle list while the newly spawned filesystem task was
275/// already running. Holding the handle-list mutex across `Handle::spawn`
276/// makes task creation visible to the drain before either side can proceed.
277fn try_spawn_lifecycle<F>(
278    handles: &parking_lot::Mutex<Vec<JoinHandle<()>>>,
279    runtime: &tokio::runtime::Handle,
280    future: F,
281) -> bool
282where
283    F: std::future::Future<Output = ()> + Send + 'static,
284{
285    std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
286        let mut handles = handles.lock();
287        handles.retain(|handle| !handle.is_finished());
288        handles.push(runtime.spawn(future));
289    }))
290    .is_ok()
291}
292
293/// Deletes an uncommitted merge/reorder output if its task unwinds.
294///
295/// Normal `Result::Err` paths delete outputs synchronously so callers observe
296/// a clean directory before returning. This guard covers the path those
297/// branches cannot: a panic after output files have been created. The cleanup
298/// callback re-checks metadata before deleting, so a panic after a successful
299/// metadata commit cannot remove a live segment.
300struct OutputCleanupGuard {
301    segment_id: SegmentId,
302    cleanup: Option<Arc<dyn Fn(SegmentId) + Send + Sync>>,
303}
304
305impl OutputCleanupGuard {
306    fn new(segment_id: SegmentId, cleanup: Arc<dyn Fn(SegmentId) + Send + Sync>) -> Self {
307        Self {
308            segment_id,
309            cleanup: Some(cleanup),
310        }
311    }
312
313    fn disarm(&mut self) {
314        self.cleanup = None;
315    }
316}
317
318impl Drop for OutputCleanupGuard {
319    fn drop(&mut self) {
320        if let Some(cleanup) = self.cleanup.take() {
321            cleanup(self.segment_id);
322        }
323    }
324}
325
326/// All mutable state behind the single async Mutex.
327struct ManagerState {
328    metadata: IndexMetadata,
329    merge_policy: Box<dyn MergePolicy>,
330}
331
332#[cfg(feature = "native")]
333struct MergeTaskError {
334    error: Error,
335    unavailable_segments: Vec<String>,
336}
337
338#[cfg(feature = "native")]
339impl MergeTaskError {
340    fn source(segment_id: String, error: Error) -> Self {
341        Self {
342            error,
343            unavailable_segments: vec![segment_id],
344        }
345    }
346
347    fn sources(segment_ids: Vec<String>, error: Error) -> Self {
348        Self {
349            error,
350            unavailable_segments: segment_ids,
351        }
352    }
353}
354
355#[cfg(feature = "native")]
356impl From<Error> for MergeTaskError {
357    fn from(error: Error) -> Self {
358        Self {
359            error,
360            unavailable_segments: Vec::new(),
361        }
362    }
363}
364
365#[cfg(feature = "native")]
366fn is_deterministic_source_error(error: &Error) -> bool {
367    matches!(error, Error::Corruption(_) | Error::Serialization(_))
368        || matches!(error, Error::Io(error) if error.kind() == std::io::ErrorKind::NotFound)
369}
370
371#[cfg(feature = "native")]
372fn classify_source_error(segment_id: String, error: Error) -> MergeTaskError {
373    if is_deterministic_source_error(&error) {
374        MergeTaskError::source(segment_id, error)
375    } else {
376        // Timeouts, interrupted reads, permission changes, and other generic
377        // I/O failures may be transient. Back them off instead of quarantining
378        // a healthy metadata segment for the rest of the process lifetime.
379        MergeTaskError::from(error)
380    }
381}
382
383#[cfg(feature = "native")]
384type MergeTaskResult<T> = std::result::Result<T, MergeTaskError>;
385
386/// Segment manager — coordinates segment commit, background merging, and trained structures.
387///
388/// SOLE owner of `metadata.json`. All metadata mutations go through `state` Mutex.
389pub struct SegmentManager<D: DirectoryWriter + 'static> {
390    /// Serializes ALL metadata mutations.
391    state: Arc<AsyncMutex<ManagerState>>,
392
393    /// RAII ownership for every in-flight segment lifecycle operation.
394    active_operations: Arc<ActiveSegmentOperations>,
395
396    /// Metadata-live segments involved in a deterministic source/corruption
397    /// failure. They stay searchable (and operator-visible) but are excluded
398    /// from merges for this process lifetime, preventing a bad candidate from
399    /// consuming full rewrite capacity on every retry.
400    quarantined_segments: parking_lot::Mutex<HashSet<String>>,
401
402    /// Generic merge failures pause scheduling briefly. Source-specific open
403    /// failures use `quarantined_segments` instead so healthy work can continue.
404    merge_retry: parking_lot::Mutex<MergeRetryState>,
405
406    /// Per-source backoff for non-deterministic standalone reorder failures.
407    /// Optimizer scans are periodic, but a pass can outlast the scan interval;
408    /// without completion-based backoff it would restart almost immediately.
409    reorder_retries: parking_lot::Mutex<HashMap<String, MergeRetryState>>,
410
411    /// In-flight merge JoinHandles — supports multiple concurrent merges.
412    merge_handles: parking_lot::Mutex<Vec<JoinHandle<()>>>,
413
414    /// At most one task per index waits for application-wide merge capacity.
415    /// Without this wakeup, an index denied by another index can remain idle
416    /// forever when no later commit happens to re-run merge policy evaluation.
417    global_merge_wakeup_pending: AtomicBool,
418
419    /// Auxiliary lifecycle tasks: metadata transactions, deferred deletes,
420    /// and capacity wakeups. Handles registered here are drained before index
421    /// removal.
422    lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>>,
423
424    /// Trained vector structures — lock-free reads via ArcSwap.
425    /// Wrapped in `Arc` so cancellation-safe metadata transactions can publish
426    /// the matching in-memory generation after their durable commit point.
427    trained: Arc<ArcSwapOption<TrainedVectorStructures>>,
428
429    /// Raised while index-level trained artifacts and their metadata are being
430    /// replaced. Search readers keep using the last valid generation, while
431    /// segment producers fall back to flat output until publication completes.
432    vector_artifact_update: Arc<AtomicBool>,
433
434    /// Reference counting for safe segment deletion (sync Mutex for Drop).
435    tracker: Arc<SegmentTracker>,
436
437    /// Cached deletion callback for snapshots (avoids allocation per acquire_snapshot).
438    delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync>,
439
440    /// Directory for segment I/O
441    directory: Arc<D>,
442    /// Schema for segment operations
443    schema: Arc<crate::dsl::Schema>,
444    /// Term cache blocks for segment readers during merge
445    term_cache_blocks: usize,
446    /// Hard concurrency limit for background merges. A semaphore permit is
447    /// acquired before lifecycle ownership, closing the old handle-count race
448    /// where concurrent schedulers could exceed the configured maximum.
449    merge_permits: Arc<Semaphore>,
450    /// Application-wide merge limit shared across index managers.
451    global_merge_permits: Arc<Semaphore>,
452    /// Shared across every index opened from the same `IndexConfig`. This
453    /// bounds whole BP rewrites (optimizer + merge-time + manual) separately
454    /// from Rayon thread width, preventing N × memory-budget amplification.
455    reorder_permits: Arc<Semaphore>,
456    /// Run BP reordering of `reorder`-attributed BMP fields inside merges.
457    /// Persisted index configuration (schema-level `reorder_on_merge: true`
458    /// in SDL); merged segments are marked `reordered` and skipped by the
459    /// standalone optimizer pass.
460    reorder_on_merge: bool,
461    /// Wall-clock budget for merge-time BP (from `IndexConfig`); truncated
462    /// passes mark the merged segment `bp_converged = false` so the
463    /// background optimizer deepens it later (warm-started).
464    merge_bp_time_budget: Option<std::time::Duration>,
465    /// Memory budget for the BP forward index (merge-time and background
466    /// reorder). Over-budget passes drop highest-df dims, logged loudly.
467    bp_memory_budget_bytes: usize,
468    /// Application-owned shared pool, when configured. This is the server
469    /// path and ensures optimizer and merge-time work use the same threads.
470    background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
471}
472
473impl<D: DirectoryWriter + 'static> SegmentManager<D> {
474    /// Create a new segment manager with existing metadata
475    #[allow(clippy::too_many_arguments)]
476    pub fn new(
477        directory: Arc<D>,
478        schema: Arc<crate::dsl::Schema>,
479        metadata: IndexMetadata,
480        merge_policy: Box<dyn MergePolicy>,
481        term_cache_blocks: usize,
482        max_concurrent_merges: usize,
483        global_merge_permits: Arc<Semaphore>,
484        merge_bp_time_budget: Option<std::time::Duration>,
485        bp_memory_budget_bytes: usize,
486        reorder_permits: Arc<Semaphore>,
487        background_reorder_pool: Option<Arc<rayon::ThreadPool>>,
488    ) -> Self {
489        // Persisted index option: set via `reorder_on_merge: true` in the SDL
490        // at index creation. Absent = disabled (merges block-copy).
491        let reorder_on_merge = schema.reorder_on_merge();
492        if reorder_on_merge {
493            log::info!("[merge] reorder-on-merge enabled by index schema");
494        }
495
496        let tracker = Arc::new(SegmentTracker::new());
497        for seg_id in metadata.segment_metas.keys() {
498            tracker.register(seg_id);
499        }
500
501        let lifecycle_handles: Arc<parking_lot::Mutex<Vec<JoinHandle<()>>>> =
502            Arc::new(parking_lot::Mutex::new(Vec::new()));
503        let delete_fn: Arc<dyn Fn(Vec<SegmentId>) + Send + Sync> = {
504            let dir = Arc::clone(&directory);
505            let tracker = Arc::clone(&tracker);
506            let lifecycle_handles = Arc::clone(&lifecycle_handles);
507            Arc::new(move |segment_ids| {
508                // Guard: if the tokio runtime is gone (program exit), skip async
509                // deletion. Segment files become orphans cleaned up on next startup.
510                let Ok(handle) = tokio::runtime::Handle::try_current() else {
511                    // Release in-process protection as well: if the process is
512                    // still alive, a later sweep must be able to retry.
513                    tracker.complete_deletion(&segment_ids);
514                    return;
515                };
516                let dir = Arc::clone(&dir);
517                let task_tracker = Arc::clone(&tracker);
518                let cleanup_ids = segment_ids.clone();
519                let future = async move {
520                    for &segment_id in &segment_ids {
521                        log::info!(
522                            "[segment_cleanup] deleting deferred segment {}",
523                            segment_id.to_hex()
524                        );
525                        if let Err(error) =
526                            crate::segment::delete_segment(dir.as_ref(), segment_id).await
527                        {
528                            log::warn!(
529                                "[segment_cleanup] deferred delete failed for {}: {}",
530                                segment_id.to_hex(),
531                                error,
532                            );
533                        }
534                    }
535                    task_tracker.complete_deletion(&segment_ids);
536                };
537                if !try_spawn_lifecycle(&lifecycle_handles, &handle, future) {
538                    // Spawning can fail only during runtime teardown. Release
539                    // the scheduled-deletion claim so an in-process sweep can
540                    // retry; crash recovery handles a process exit.
541                    tracker.complete_deletion(&cleanup_ids);
542                    log::warn!(
543                        "[segment_cleanup] runtime rejected deferred deletion; files will be swept later"
544                    );
545                }
546            })
547        };
548
549        Self {
550            state: Arc::new(AsyncMutex::new(ManagerState {
551                metadata,
552                merge_policy,
553            })),
554            active_operations: Arc::new(ActiveSegmentOperations::new()),
555            quarantined_segments: parking_lot::Mutex::new(HashSet::new()),
556            merge_retry: parking_lot::Mutex::new(MergeRetryState::default()),
557            reorder_retries: parking_lot::Mutex::new(HashMap::new()),
558            merge_handles: parking_lot::Mutex::new(Vec::new()),
559            global_merge_wakeup_pending: AtomicBool::new(false),
560            lifecycle_handles,
561            trained: Arc::new(ArcSwapOption::new(None)),
562            vector_artifact_update: Arc::new(AtomicBool::new(false)),
563            tracker,
564            delete_fn,
565            directory,
566            schema,
567            term_cache_blocks,
568            merge_permits: Arc::new(Semaphore::new(max_concurrent_merges.max(1))),
569            global_merge_permits,
570            reorder_permits,
571            reorder_on_merge,
572            merge_bp_time_budget,
573            bp_memory_budget_bytes,
574            background_reorder_pool,
575        }
576    }
577
578    /// Bounded rayon pool for background CPU (merge-time BP, manual reorder).
579    /// Query scoring uses the global rayon pool; keeping background BP off it
580    /// prevents a large merge from queueing every search behind gain passes.
581    pub fn background_cpu_pool(&self) -> Arc<rayon::ThreadPool> {
582        if let Some(pool) = &self.background_reorder_pool {
583            return Arc::clone(pool);
584        }
585        Arc::clone(BACKGROUND_CPU_POOL.get_or_init(|| {
586            let threads = (num_cpus::get() / 2).max(1);
587            log::info!(
588                "[merge] process-wide background CPU pool: {} thread(s)",
589                threads
590            );
591            Arc::new(
592                rayon::ThreadPoolBuilder::new()
593                    .num_threads(threads)
594                    .thread_name(|i| format!("hermes-bg-cpu-{}", i))
595                    .build()
596                    .expect("failed to build background CPU pool"),
597            )
598        }))
599    }
600
601    /// Stop new indexing/merge/reorder operations from claiming segment IDs.
602    /// Used as the first half of index deletion; the writer then joins its
603    /// workers before [`Self::wait_for_shutdown`] drains remaining ownership.
604    pub fn begin_shutdown(&self) {
605        self.active_operations.stop_accepting();
606    }
607
608    /// Run a lifecycle mutation independently of its requesting future.
609    ///
610    /// Metadata writes contain an atomic rename. If an RPC is cancelled while
611    /// awaiting that I/O, dropping the request must not abandon the matching
612    /// in-memory/tracker transition. The spawned transaction is tracked for
613    /// index shutdown; the oneshot only reports its result to a caller that is
614    /// still interested.
615    async fn run_lifecycle_transaction<T, F>(&self, transaction: F) -> Result<T>
616    where
617        T: Send + 'static,
618        F: std::future::Future<Output = Result<T>> + Send + 'static,
619    {
620        let (result_tx, result_rx) = tokio::sync::oneshot::channel();
621        let future = async move {
622            let result = transaction.await;
623            let _ = result_tx.send(result);
624        };
625        let runtime = tokio::runtime::Handle::current();
626        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
627            return Err(Error::Internal(
628                "runtime rejected lifecycle metadata transaction".into(),
629            ));
630        }
631        result_rx.await.map_err(|_| {
632            Error::Internal("lifecycle metadata transaction terminated unexpectedly".into())
633        })?
634    }
635
636    /// Arm unwind cleanup for an output that is not visible in metadata yet.
637    fn output_cleanup_guard(self: &Arc<Self>, output_id: SegmentId) -> OutputCleanupGuard {
638        let manager = Arc::clone(self);
639        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |segment_id| {
640            let Ok(handle) = tokio::runtime::Handle::try_current() else {
641                log::warn!(
642                    "[segment_cleanup] runtime unavailable; partial output {} will be swept on startup",
643                    segment_id.to_hex(),
644                );
645                return;
646            };
647
648            let cleanup_manager = Arc::clone(&manager);
649            let future = async move {
650                cleanup_manager
651                    .delete_output_if_unregistered(segment_id, "task unwind")
652                    .await;
653            };
654            if !try_spawn_lifecycle(&manager.lifecycle_handles, &handle, future) {
655                log::warn!(
656                    "[segment_cleanup] runtime rejected output cleanup; {} will be swept on startup",
657                    segment_id.to_hex(),
658                );
659            }
660        });
661
662        OutputCleanupGuard::new(output_id, cleanup)
663    }
664
665    /// Delete an abandoned indexing output while retaining its lifecycle
666    /// claim until the last file operation completes. The explicit runtime
667    /// handle makes this safe from dedicated indexing OS threads, which are
668    /// outside Tokio's entered context.
669    pub(crate) fn schedule_unpublished_segment_cleanup(
670        self: &Arc<Self>,
671        output_id: SegmentId,
672        operation: SegmentOperationGuard,
673        runtime: tokio::runtime::Handle,
674    ) {
675        let manager = Arc::clone(self);
676        let output_hex = output_id.to_hex();
677        let future = async move {
678            manager
679                .delete_output_if_unregistered(output_id, "indexing abort or failure")
680                .await;
681            drop(operation);
682        };
683        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
684            // The dropped future releases operation ownership. Startup sweep
685            // handles its output if the runtime is already tearing down.
686            log::warn!(
687                "[segment_cleanup] runtime unavailable; indexing output {} will be swept on startup",
688                output_hex,
689            );
690        }
691    }
692
693    /// Claim a newly generated indexing segment before its first file write.
694    ///
695    /// The returned guard must travel with the built segment until metadata
696    /// publication or abort. UUID collisions are treated as corruption rather
697    /// than silently sharing lifecycle ownership.
698    pub(crate) fn protect_new_segment(&self, segment_id: String) -> Result<SegmentOperationGuard> {
699        match self
700            .active_operations
701            .try_register(vec![segment_id.clone()])
702        {
703            Some(operation) => Ok(operation),
704            None if !self.active_operations.is_accepting() => Err(Error::IndexClosed),
705            None => Err(Error::Corruption(format!(
706                "new segment ID {} is already owned by an active operation",
707                segment_id
708            ))),
709        }
710    }
711
712    /// Validate the small, mandatory core of a completed segment before it can
713    /// become metadata-live. Optional vector/sparse/position/fast files are
714    /// schema- and data-dependent and are validated by `SegmentReader` when used.
715    async fn validate_completed_segment(&self, segment_id: &str, expected_docs: u32) -> Result<()> {
716        let id = SegmentId::from_hex(segment_id).ok_or_else(|| {
717            Error::Corruption(format!("invalid completed segment ID: {}", segment_id))
718        })?;
719        let files = SegmentFiles::new(id.0);
720
721        for path in files.mandatory_paths() {
722            if !self.directory.exists(path).await.map_err(Error::Io)? {
723                return Err(Error::Corruption(format!(
724                    "segment {} cannot be published: mandatory file {:?} is missing",
725                    segment_id, path
726                )));
727            }
728        }
729
730        let meta_slice = self.directory.open_read(&files.meta).await.map_err(|e| {
731            Error::Corruption(format!(
732                "segment {} cannot be published: missing/unreadable {:?}: {}",
733                segment_id, files.meta, e
734            ))
735        })?;
736        let meta_bytes = meta_slice.read_bytes().await.map_err(|e| {
737            Error::Corruption(format!(
738                "segment {} cannot be published: failed reading {:?}: {}",
739                segment_id, files.meta, e
740            ))
741        })?;
742        let meta = SegmentMeta::deserialize(meta_bytes.as_slice()).map_err(|e| {
743            Error::Corruption(format!(
744                "segment {} cannot be published: invalid {:?}: {}",
745                segment_id, files.meta, e
746            ))
747        })?;
748
749        if meta.id != id.0 || meta.num_docs != expected_docs {
750            return Err(Error::Corruption(format!(
751                "segment {} cannot be published: metadata identity/docs mismatch \
752                 (id={:032x}, docs={}, expected_docs={})",
753                segment_id, meta.id, meta.num_docs, expected_docs
754            )));
755        }
756
757        Ok(())
758    }
759
760    fn quarantine_segment(&self, segment_id: &str, error: &Error) {
761        let inserted = self
762            .quarantined_segments
763            .lock()
764            .insert(segment_id.to_string());
765        if inserted {
766            log::error!(
767                "[merge] quarantined metadata-live segment {} after deterministic source/validation failure: {}. \
768                 It remains metadata-live for explicit repair but is excluded from merges until restart",
769                segment_id,
770                error,
771            );
772        }
773    }
774
775    fn pause_merge_retries(&self, error: &Error) -> std::time::Duration {
776        let mut retry = self.merge_retry.lock();
777        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
778        let delay = merge_retry_delay(retry.consecutive_failures);
779        retry.retry_after = std::time::Instant::now().checked_add(delay);
780        log::warn!(
781            "[merge] pausing background merge scheduling for {:.0}s after consecutive failure #{}: {}",
782            delay.as_secs_f64(),
783            retry.consecutive_failures,
784            error,
785        );
786        delay
787    }
788
789    fn clear_merge_retry_backoff(&self) {
790        *self.merge_retry.lock() = MergeRetryState::default();
791    }
792
793    fn merge_retry_is_paused(&self) -> bool {
794        let mut retry = self.merge_retry.lock();
795        match retry.retry_after {
796            Some(deadline) if deadline > std::time::Instant::now() => true,
797            Some(_) => {
798                retry.retry_after = None;
799                false
800            }
801            None => false,
802        }
803    }
804
805    fn pause_reorder_retries(&self, segment_id: &str, error: &Error) {
806        let mut retries = self.reorder_retries.lock();
807        let retry = retries.entry(segment_id.to_string()).or_default();
808        retry.consecutive_failures = retry.consecutive_failures.saturating_add(1);
809        let delay = merge_retry_delay(retry.consecutive_failures);
810        retry.retry_after = std::time::Instant::now().checked_add(delay);
811        log::warn!(
812            "[reorder] pausing optimizer retries for segment {} for {:.0}s after failure #{}: {}",
813            segment_id,
814            delay.as_secs_f64(),
815            retry.consecutive_failures,
816            error,
817        );
818    }
819
820    fn clear_reorder_retry(&self, segment_id: &str) {
821        self.reorder_retries.lock().remove(segment_id);
822    }
823
824    fn paused_reorder_segments(&self) -> HashSet<String> {
825        let now = std::time::Instant::now();
826        let mut retries = self.reorder_retries.lock();
827        let mut paused = HashSet::new();
828        for (segment_id, retry) in retries.iter_mut() {
829            match retry.retry_after {
830                Some(deadline) if deadline > now => {
831                    paused.insert(segment_id.clone());
832                }
833                Some(_) => retry.retry_after = None,
834                None => {}
835            }
836        }
837        paused
838    }
839
840    /// Re-evaluate this index when another index releases application-wide
841    /// merge capacity. The atomic flag bounds this to one waiter per index and
842    /// the tracked handle makes index shutdown drain it deterministically.
843    fn schedule_global_merge_wakeup(self: &Arc<Self>) {
844        if self
845            .global_merge_wakeup_pending
846            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
847            .is_err()
848        {
849            return;
850        }
851
852        let manager = Arc::clone(self);
853        let future = async move {
854            let capacity = tokio::select! {
855                biased;
856                () = manager.active_operations.wait_for_shutdown() => None,
857                permit = Arc::clone(&manager.global_merge_permits).acquire_owned() => permit.ok(),
858            };
859
860            manager
861                .global_merge_wakeup_pending
862                .store(false, Ordering::Release);
863            if let Some(permit) = capacity {
864                // This task is only a notification. The normal scheduler must
865                // acquire both global and per-index permits atomically enough
866                // for its own candidate selection.
867                drop(permit);
868                manager.maybe_merge().await;
869            }
870        };
871        let runtime = tokio::runtime::Handle::current();
872        if !try_spawn_lifecycle(&self.lifecycle_handles, &runtime, future) {
873            self.global_merge_wakeup_pending
874                .store(false, Ordering::Release);
875            log::warn!("[merge] runtime rejected global-capacity wakeup task");
876        }
877    }
878
879    #[cfg(test)]
880    pub(crate) fn is_segment_quarantined(&self, segment_id: &str) -> bool {
881        self.quarantined_segments.lock().contains(segment_id)
882    }
883
884    /// Delete a failed output only if metadata did not make it live.
885    ///
886    /// Rechecking under `state` also makes unwind cleanup safe if it races
887    /// successful publication of the same output.
888    async fn delete_output_if_unregistered(&self, output_id: SegmentId, reason: &str) {
889        let output_hex = output_id.to_hex();
890        {
891            let st = self.state.lock().await;
892            if st.metadata.has_segment(&output_hex) {
893                return;
894            }
895        }
896
897        // UUIDs are generated per producer and cannot be adopted by another
898        // publisher after this check. Never hold the metadata mutex while a
899        // multi-GB filesystem deletion runs.
900        log::info!(
901            "[segment_cleanup] deleting uncommitted output {} after {}",
902            output_hex,
903            reason,
904        );
905        if let Err(error) = crate::segment::delete_segment(self.directory.as_ref(), output_id).await
906        {
907            log::warn!(
908                "[segment_cleanup] failed deleting uncommitted output {}: {}",
909                output_hex,
910                error,
911            );
912        }
913    }
914
915    // ========================================================================
916    // Read path (brief lock or lock-free)
917    // ========================================================================
918
919    /// Get the current segment IDs
920    pub async fn get_segment_ids(&self) -> Vec<String> {
921        self.state.lock().await.metadata.segment_ids()
922    }
923
924    /// Get trained vector structures (lock-free via ArcSwap)
925    pub fn trained(&self) -> Option<Arc<TrainedVectorStructures>> {
926        self.trained.load_full()
927    }
928
929    /// Capture trained structures for a segment producer.
930    ///
931    /// The second gate check closes the race where an update begins after the
932    /// first check but before the ArcSwap load. Producers have lifecycle guards
933    /// before calling this method, so an updater that raised the gate waits for
934    /// any producer that successfully captured the previous generation.
935    pub(crate) fn trained_for_segment_build(&self) -> Option<Arc<TrainedVectorStructures>> {
936        if self.vector_artifact_update.load(Ordering::Acquire) {
937            return None;
938        }
939        let trained = self.trained.load_full();
940        if self.vector_artifact_update.load(Ordering::Acquire) {
941            None
942        } else {
943            trained
944        }
945    }
946
947    /// Start an exclusive trained-artifact update and drain producers that may
948    /// already hold the previous generation.
949    ///
950    /// New segment operations may continue while this waits, but they observe
951    /// the gate through `trained_for_segment_build` and therefore emit flat
952    /// vector data. The guard is cancellation-safe: dropping the requesting
953    /// future reopens ANN production without leaving the manager wedged.
954    pub(crate) async fn begin_vector_artifact_update(&self) -> Result<VectorArtifactUpdateGuard> {
955        self.vector_artifact_update
956            .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
957            .map_err(|_| {
958                Error::Internal("a trained-vector artifact update is already in progress".into())
959            })?;
960        let guard = VectorArtifactUpdateGuard {
961            _lease: Arc::new(VectorArtifactUpdateLease {
962                updating: Arc::clone(&self.vector_artifact_update),
963            }),
964        };
965        let preexisting = self.active_operations.operation_tokens_snapshot();
966        self.active_operations
967            .wait_until_operations_finish(&preexisting)
968            .await;
969        Ok(guard)
970    }
971
972    /// Compatibility entry point: load the complete trained set or clear the
973    /// published generation and log the validation failure.
974    pub async fn load_and_publish_trained(&self) {
975        if let Err(error) = self.try_load_and_publish_trained().await {
976            self.trained.store(None);
977            log::error!("[trained] refusing to publish trained artifacts: {error}");
978        }
979    }
980
981    /// Load trained structures from disk and publish to ArcSwap.
982    /// Copies metadata under lock, releases lock, then does disk I/O.
983    pub(crate) async fn try_load_and_publish_trained(&self) -> Result<()> {
984        // Copy vector_fields under lock (cheap clone of HashMap<u32, FieldMeta>)
985        let vector_fields = {
986            let st = self.state.lock().await;
987            st.metadata.vector_fields.clone()
988        };
989        // Disk I/O happens WITHOUT holding the state lock
990        let trained = IndexMetadata::try_load_trained_from_fields(
991            &vector_fields,
992            self.schema.as_ref(),
993            self.directory.as_ref(),
994        )
995        .await?
996        .map(Arc::new);
997        // Publish exactly the validated snapshot, including None. Retaining a
998        // previous map when metadata has no Built fields would let new segments
999        // depend on artifacts no longer referenced durably.
1000        self.trained.store(trained);
1001        Ok(())
1002    }
1003
1004    /// Persist vector metadata and publish its fully validated artifact set as
1005    /// one cancellation-safe lifecycle transaction.
1006    ///
1007    /// Validation happens before the durable metadata commit. Once the rename
1008    /// commits, both in-memory metadata and ArcSwap publication are completed by
1009    /// the tracked transaction even if the requesting RPC is cancelled.
1010    pub(crate) async fn update_vector_metadata_and_publish<F>(
1011        self: &Arc<Self>,
1012        artifact_update: &VectorArtifactUpdateGuard,
1013        update: F,
1014    ) -> Result<()>
1015    where
1016        F: FnOnce(&mut IndexMetadata),
1017    {
1018        let mut st = Arc::clone(&self.state).lock_owned().await;
1019        let mut next = st.metadata.clone();
1020        update(&mut next);
1021
1022        let next_trained = IndexMetadata::try_load_trained_from_fields(
1023            &next.vector_fields,
1024            self.schema.as_ref(),
1025            self.directory.as_ref(),
1026        )
1027        .await?
1028        .map(Arc::new);
1029
1030        let directory = Arc::clone(&self.directory);
1031        let trained = Arc::clone(&self.trained);
1032        // Keep the producer gate raised if the requesting future is cancelled
1033        // after the metadata transaction has been detached. The last guard
1034        // clone drops only after durable metadata and ArcSwap state agree.
1035        let artifact_update = artifact_update.clone();
1036        self.run_lifecycle_transaction(async move {
1037            let _artifact_update = artifact_update;
1038            next.save(directory.as_ref()).await?;
1039            st.metadata = next;
1040            trained.store(next_trained);
1041            Ok(())
1042        })
1043        .await
1044    }
1045
1046    /// Read metadata with a closure (no persist)
1047    pub(crate) async fn read_metadata<F, R>(&self, f: F) -> R
1048    where
1049        F: FnOnce(&IndexMetadata) -> R,
1050    {
1051        let st = self.state.lock().await;
1052        f(&st.metadata)
1053    }
1054
1055    /// Update metadata with a closure and persist atomically
1056    pub(crate) async fn update_metadata<F>(self: &Arc<Self>, f: F) -> Result<()>
1057    where
1058        F: FnOnce(&mut IndexMetadata),
1059    {
1060        let mut st = Arc::clone(&self.state).lock_owned().await;
1061        let mut next = st.metadata.clone();
1062        f(&mut next);
1063        let directory = Arc::clone(&self.directory);
1064        self.run_lifecycle_transaction(async move {
1065            next.save(directory.as_ref()).await?;
1066            st.metadata = next;
1067            Ok(())
1068        })
1069        .await
1070    }
1071
1072    /// Acquire a snapshot of current segments for reading.
1073    /// The snapshot holds references — segments won't be deleted while snapshot exists.
1074    pub async fn acquire_snapshot(&self) -> SegmentSnapshot {
1075        let acquired = {
1076            let st = self.state.lock().await;
1077            let segment_ids = st.metadata.segment_ids();
1078            self.tracker.acquire(&segment_ids)
1079        };
1080
1081        SegmentSnapshot::with_delete_fn(
1082            Arc::clone(&self.tracker),
1083            acquired,
1084            Arc::clone(&self.delete_fn),
1085        )
1086    }
1087
1088    /// Get the segment tracker
1089    pub fn tracker(&self) -> Arc<SegmentTracker> {
1090        Arc::clone(&self.tracker)
1091    }
1092
1093    /// Get the directory
1094    pub fn directory(&self) -> Arc<D> {
1095        Arc::clone(&self.directory)
1096    }
1097}
1098
1099// ============================================================================
1100// Native-only: commit, merging, force_merge
1101// ============================================================================
1102
1103#[cfg(feature = "native")]
1104impl<D: DirectoryWriter + 'static> SegmentManager<D> {
1105    /// Atomic commit: register new segments + persist metadata.
1106    pub async fn commit(self: &Arc<Self>, new_segments: &[(String, u32)]) -> Result<()> {
1107        // Indexing guards still own these IDs here, so the orphan sweeper
1108        // cannot remove files between validation and metadata publication.
1109        for (segment_id, num_docs) in new_segments {
1110            self.validate_completed_segment(segment_id, *num_docs)
1111                .await?;
1112        }
1113
1114        let mut st = Arc::clone(&self.state).lock_owned().await;
1115        let mut next = st.metadata.clone();
1116        let mut added = Vec::new();
1117        for (segment_id, num_docs) in new_segments {
1118            if !next.has_segment(segment_id) {
1119                next.add_segment(segment_id.clone(), *num_docs);
1120                added.push(segment_id.clone());
1121            }
1122        }
1123
1124        // Durable-before-visible: a save failure leaves both in-memory metadata
1125        // and tracker unchanged, so callers can retry the prepared commit.
1126        // The tracked transaction continues if the requesting RPC is cancelled;
1127        // unpublished cleanup waits on this owned state guard before deciding
1128        // whether the files became metadata-live.
1129        let directory = Arc::clone(&self.directory);
1130        let tracker = Arc::clone(&self.tracker);
1131        self.run_lifecycle_transaction(async move {
1132            next.save(directory.as_ref()).await?;
1133            for segment_id in &added {
1134                tracker.register(segment_id);
1135            }
1136            st.metadata = next;
1137            Ok(())
1138        })
1139        .await
1140    }
1141
1142    /// Evaluate merge policy and spawn background merges for all eligible candidates.
1143    ///
1144    /// **Atomicity**: The entire filter → find_merges → spawn_merge sequence runs
1145    /// under the `state` lock to prevent a TOCTOU race where concurrent callers
1146    /// both see segments as eligible before either claims operation ownership.
1147    /// `spawn_merge` is non-blocking (just `try_register` + `tokio::spawn`), so
1148    /// holding the state lock through it is safe and sub-microsecond.
1149    ///
1150    /// The hard merge semaphore is acquired before lifecycle ownership, so
1151    /// concurrent triggers cannot exceed configured merge capacity.
1152    pub async fn maybe_merge(self: &Arc<Self>) {
1153        if !self.active_operations.is_accepting() {
1154            log::debug!("[maybe_merge] manager is shutting down, skipping");
1155            return;
1156        }
1157        if self.merge_retry_is_paused() {
1158            log::debug!("[maybe_merge] retry backoff active, skipping");
1159            return;
1160        }
1161
1162        // Finished handles no longer need to be retained. Concurrency itself
1163        // is enforced by `merge_permits`, not this bookkeeping vector.
1164        {
1165            let mut handles = self.merge_handles.lock();
1166            handles.retain(|h| !h.is_finished());
1167        }
1168        let local_slots = self.merge_permits.available_permits();
1169        let global_slots = self.global_merge_permits.available_permits();
1170        let slots_available = local_slots.min(global_slots);
1171
1172        // Hold state lock through spawn_merge to make filter + register atomic.
1173        // This closes the TOCTOU window where concurrent maybe_merge calls could
1174        // both see the same segments as eligible before either registers them.
1175        let new_handles = {
1176            let st = self.state.lock().await;
1177            let quarantined = self.quarantined_segments.lock().clone();
1178            let active_ids = self.active_operations.snapshot();
1179
1180            // Exclude segments owned by another operation, pending retirement,
1181            // or quarantined after a persistent open/validation failure.
1182            let segments: Vec<SegmentInfo> = st
1183                .metadata
1184                .segment_metas
1185                .iter()
1186                .filter(|(id, _)| {
1187                    !self.tracker.is_pending_deletion(id)
1188                        && !active_ids.contains(*id)
1189                        && !quarantined.contains(*id)
1190                })
1191                .map(|(id, info)| SegmentInfo {
1192                    id: id.clone(),
1193                    num_docs: info.num_docs,
1194                })
1195                .collect();
1196
1197            log::debug!("[maybe_merge] {} eligible segments", segments.len());
1198
1199            let candidates = st.merge_policy.find_merges(&segments);
1200
1201            if candidates.is_empty() {
1202                return;
1203            }
1204
1205            // Register a capacity waiter only for an index that actually has
1206            // eligible work. Scheduling one waiter for every idle index while
1207            // the process gate was full caused an avoidable wakeup stampede.
1208            if slots_available == 0 {
1209                if local_slots > 0 && global_slots == 0 {
1210                    self.schedule_global_merge_wakeup();
1211                }
1212                log::debug!("[maybe_merge] at max concurrent merges, skipping");
1213                return;
1214            }
1215
1216            log::debug!(
1217                "[maybe_merge] {} merge candidates, {} slots available",
1218                candidates.len(),
1219                slots_available
1220            );
1221
1222            let mut handles = Vec::new();
1223            for c in candidates {
1224                if handles.len() >= slots_available {
1225                    break;
1226                }
1227                if let Some(h) = self.spawn_merge(c.segment_ids) {
1228                    handles.push(h);
1229                }
1230            }
1231            handles
1232            // State lock released after spawn_merge claimed operation ownership.
1233        };
1234
1235        if !new_handles.is_empty() {
1236            // Synchronous insertion is part of spawning: there must be no
1237            // cancellation point where a live task exists but shutdown and
1238            // force-merge draining cannot see its JoinHandle.
1239            self.merge_handles.lock().extend(new_handles);
1240        }
1241    }
1242
1243    /// Spawn a background merge task with RAII tracking.
1244    ///
1245    /// Pre-generates the output segment ID. The operation guard registers all segment IDs
1246    /// (old + output) in `active_operations`. When the task ends (success, failure, or
1247    /// panic), the guard drops and segments are automatically unregistered.
1248    ///
1249    /// On completion, the task auto-triggers `maybe_merge` to evaluate cascading merges.
1250    /// Returns the JoinHandle if the merge was spawned, None if it was skipped.
1251    fn spawn_merge(self: &Arc<Self>, segment_ids_to_merge: Vec<String>) -> Option<JoinHandle<()>> {
1252        let global_merge_permit = match Arc::clone(&self.global_merge_permits).try_acquire_owned() {
1253            Ok(permit) => permit,
1254            Err(_) => {
1255                log::debug!("[spawn_merge] skipped: global merge capacity is full");
1256                self.schedule_global_merge_wakeup();
1257                return None;
1258            }
1259        };
1260        let merge_permit = match Arc::clone(&self.merge_permits).try_acquire_owned() {
1261            Ok(permit) => permit,
1262            Err(_) => {
1263                log::debug!("[spawn_merge] skipped: no merge permit available");
1264                return None;
1265            }
1266        };
1267        let output_id = SegmentId::new();
1268        let output_hex = output_id.to_hex();
1269
1270        let mut all_ids = segment_ids_to_merge.clone();
1271        all_ids.push(output_hex);
1272
1273        let guard = match self.active_operations.try_register(all_ids) {
1274            Some(g) => g,
1275            None => {
1276                log::debug!("[spawn_merge] skipped: segments overlap with an active operation");
1277                return None;
1278            }
1279        };
1280
1281        let sm = Arc::clone(self);
1282        let ids = segment_ids_to_merge;
1283
1284        Some(tokio::spawn(async move {
1285            let mut output_cleanup = sm.output_cleanup_guard(output_id);
1286            let mut reevaluate = false;
1287            let mut retry_delay = None;
1288
1289            let trained_snap = sm.trained_for_segment_build();
1290            let granularity = sm.merge_granularity(&ids).await;
1291            let result = Self::do_merge(
1292                sm.directory.as_ref(),
1293                &sm.schema,
1294                &ids,
1295                output_id,
1296                sm.term_cache_blocks,
1297                trained_snap.as_deref(),
1298                sm.reorder_on_merge,
1299                granularity,
1300                sm.merge_bp_time_budget,
1301                sm.bp_memory_budget_bytes,
1302                Arc::clone(&sm.reorder_permits),
1303                Some(sm.background_cpu_pool()),
1304            )
1305            .await;
1306
1307            match result {
1308                Ok((new_id, doc_count, bp_converged)) => {
1309                    match sm
1310                        .replace_segments(
1311                            &ids,
1312                            new_id,
1313                            doc_count,
1314                            sm.reorder_on_merge,
1315                            bp_converged,
1316                        )
1317                        .await
1318                    {
1319                        Ok(()) => {
1320                            output_cleanup.disarm();
1321                            sm.clear_merge_retry_backoff();
1322                            reevaluate = true;
1323                        }
1324                        Err(e) => {
1325                            sm.delete_output_if_unregistered(output_id, "replacement failure")
1326                                .await;
1327                            output_cleanup.disarm();
1328                            retry_delay = Some(sm.pause_merge_retries(&e));
1329                            log::error!("[merge] failed to publish merged segment: {}", e);
1330                        }
1331                    }
1332                }
1333                Err(MergeTaskError {
1334                    error,
1335                    unavailable_segments,
1336                }) => {
1337                    log::error!(
1338                        "[merge] background merge failed for segments {:?}: {}",
1339                        ids,
1340                        error
1341                    );
1342                    if !unavailable_segments.is_empty() {
1343                        for segment_id in &unavailable_segments {
1344                            sm.quarantine_segment(segment_id, &error);
1345                        }
1346                        // Recompute without this known-bad input. This is not a
1347                        // retry of the same candidate because policy filtering
1348                        // excludes every quarantined ID.
1349                        reevaluate = true;
1350                    } else {
1351                        retry_delay = Some(sm.pause_merge_retries(&error));
1352                    }
1353                    sm.delete_output_if_unregistered(output_id, "merge failure")
1354                        .await;
1355                    output_cleanup.disarm();
1356                }
1357            }
1358            // Release source/output ownership before re-evaluating policy, so
1359            // the completed operation cannot artificially hide candidates.
1360            drop(guard);
1361            // A failed merge must not reserve capacity during its retry delay.
1362            drop(merge_permit);
1363            drop(global_merge_permit);
1364
1365            if reevaluate {
1366                sm.maybe_merge().await;
1367            } else if let Some(retry_delay) = retry_delay {
1368                // A backoff without a wakeup can strand eligible segments
1369                // forever when no later commit happens. Keep this sleep inside
1370                // the tracked merge task so shutdown can await it safely.
1371                tokio::select! {
1372                    () = tokio::time::sleep(retry_delay) => {
1373                        sm.maybe_merge().await;
1374                    }
1375                    () = sm.active_operations.wait_for_shutdown() => {}
1376                }
1377            }
1378        }))
1379    }
1380
1381    /// Atomically replace old segments with a new merged segment.
1382    /// Computes merge generation as max(parent gens) + 1 and records ancestors.
1383    /// `reordered` marks whether the new segment was BP-reordered.
1384    async fn replace_segments(
1385        self: &Arc<Self>,
1386        old_ids: &[String],
1387        new_id: String,
1388        doc_count: u32,
1389        reordered: bool,
1390        bp_converged: bool,
1391    ) -> Result<()> {
1392        // The operation guard owns the output during validation. Publication
1393        // below replaces that ownership with metadata + tracker atomically.
1394        self.validate_completed_segment(&new_id, doc_count).await?;
1395        let output_id = SegmentId::from_hex(&new_id).ok_or_else(|| {
1396            Error::Corruption(format!("invalid replacement segment ID: {new_id}"))
1397        })?;
1398        let output_reader = SegmentReader::open(
1399            self.directory.as_ref(),
1400            output_id,
1401            Arc::clone(&self.schema),
1402            self.term_cache_blocks,
1403        )
1404        .await
1405        .map_err(|error| match error {
1406            // Preserve retryable storage failures as I/O. Structural failures
1407            // are deterministic for this completed output and get explicit
1408            // corruption context.
1409            Error::Io(_) | Error::IndexClosed => error,
1410            error => Error::Corruption(format!(
1411                "replacement segment {new_id} failed full reader validation: {error}"
1412            )),
1413        })?;
1414        if output_reader.num_docs() != doc_count {
1415            return Err(Error::Corruption(format!(
1416                "replacement segment {new_id} opened with {} docs, expected {doc_count}",
1417                output_reader.num_docs(),
1418            )));
1419        }
1420        drop(output_reader);
1421
1422        let mut st = Arc::clone(&self.state).lock_owned().await;
1423        // Every source must still be live: callers hold operation ownership,
1424        // so a missing source means a stale merge/reorder whose input was
1425        // already replaced. Adding the output would duplicate its documents.
1426        let missing: Vec<&String> = old_ids
1427            .iter()
1428            .filter(|id| !st.metadata.has_segment(id))
1429            .collect();
1430        if !missing.is_empty() {
1431            return Err(Error::Corruption(format!(
1432                "replace_segments: source segment(s) {:?} not in metadata — \
1433                 refusing to add output {} (would duplicate documents)",
1434                missing, new_id
1435            )));
1436        }
1437
1438        let parent_generation = old_ids
1439            .iter()
1440            .filter_map(|id| st.metadata.segment_metas.get(id))
1441            .map(|info| info.generation)
1442            .max()
1443            .unwrap_or(0)
1444            .checked_add(1)
1445            .ok_or_else(|| Error::Corruption("merge generation exceeds u32::MAX".into()))?;
1446        let parent_unconverged_passes = old_ids
1447            .iter()
1448            .filter_map(|id| st.metadata.segment_metas.get(id))
1449            .map(|info| info.bp_unconverged_passes)
1450            .max()
1451            .unwrap_or(0);
1452        let bp_unconverged_passes = if reordered && !bp_converged {
1453            parent_unconverged_passes.saturating_add(1)
1454        } else {
1455            0
1456        };
1457        let retired_ids = old_ids.to_vec();
1458        let mut next = st.metadata.clone();
1459        for id in old_ids {
1460            next.remove_segment(id);
1461        }
1462        next.add_segment_meta(
1463            new_id.clone(),
1464            SegmentMetaInfo {
1465                num_docs: doc_count,
1466                ancestors: retired_ids.clone(),
1467                generation: parent_generation,
1468                reordered,
1469                bp_converged,
1470                bp_unconverged_passes,
1471            },
1472        );
1473
1474        let directory = Arc::clone(&self.directory);
1475        let tracker = Arc::clone(&self.tracker);
1476        self.run_lifecycle_transaction(async move {
1477            // Durable-before-visible. If persistence fails, old metadata and
1478            // tracker ownership stay intact and source deletion is never armed.
1479            next.save(directory.as_ref()).await?;
1480            tracker.register(&new_id);
1481            st.metadata = next;
1482
1483            // Keep state locked until retired sources enter the tracker. The
1484            // transaction itself also performs deletion, so cancellation of
1485            // the requesting merge cannot strand pending-deletion ownership.
1486            let ready_to_delete = tracker.mark_for_deletion(&retired_ids);
1487            drop(st);
1488            for &segment_id in &ready_to_delete {
1489                if let Err(error) =
1490                    crate::segment::delete_segment(directory.as_ref(), segment_id).await
1491                {
1492                    log::warn!(
1493                        "[segment_cleanup] immediate delete failed for {}: {}",
1494                        segment_id.to_hex(),
1495                        error,
1496                    );
1497                }
1498            }
1499            tracker.complete_deletion(&ready_to_delete);
1500            Ok(())
1501        })
1502        .await
1503    }
1504
1505    /// Perform the actual merge operation (pure function — no shared state access).
1506    /// `output_segment_id` is pre-generated by the caller so active-operation ownership
1507    /// is installed before any output file is written.
1508    /// Returns (new_segment_id_hex, total_doc_count).
1509    #[allow(clippy::too_many_arguments)]
1510    async fn do_merge(
1511        directory: &D,
1512        schema: &Arc<crate::dsl::Schema>,
1513        segment_ids_to_merge: &[String],
1514        output_segment_id: SegmentId,
1515        term_cache_blocks: usize,
1516        trained: Option<&TrainedVectorStructures>,
1517        reorder_bmp: bool,
1518        granularity: crate::segment::reorder::BpGranularity,
1519        merge_bp_time_budget: Option<std::time::Duration>,
1520        bp_memory_budget_bytes: usize,
1521        reorder_permits: Arc<Semaphore>,
1522        bg_cpu_pool: Option<Arc<rayon::ThreadPool>>,
1523    ) -> MergeTaskResult<(String, u32, bool)> {
1524        let output_hex = output_segment_id.to_hex();
1525        let load_start = std::time::Instant::now();
1526
1527        let mut segment_ids = Vec::with_capacity(segment_ids_to_merge.len());
1528        for id_str in segment_ids_to_merge {
1529            let id = SegmentId::from_hex(id_str).ok_or_else(|| {
1530                MergeTaskError::source(
1531                    id_str.clone(),
1532                    Error::Corruption(format!("Invalid segment ID: {}", id_str)),
1533                )
1534            })?;
1535            segment_ids.push(id);
1536        }
1537
1538        // Cheap fail-fast before opening every reader. `join_all` otherwise
1539        // waits for all healthy multi-GB inputs to load even when one source's
1540        // `.meta` is already absent, turning a known-corrupt candidate into a
1541        // large CPU/IO spike before it can be quarantined.
1542        let mut unavailable_sources = Vec::new();
1543        let mut missing_files = Vec::new();
1544        for (id_str, id) in segment_ids_to_merge.iter().zip(&segment_ids) {
1545            let files = SegmentFiles::new(id.0);
1546            let mut source_unavailable = false;
1547            for path in files.mandatory_paths() {
1548                let exists = directory
1549                    .exists(path)
1550                    .await
1551                    .map_err(|error| MergeTaskError::from(Error::Io(error)))?;
1552                if !exists {
1553                    source_unavailable = true;
1554                    missing_files.push(format!("{}:{:?}", id_str, path));
1555                }
1556            }
1557            if source_unavailable {
1558                unavailable_sources.push(id_str.clone());
1559            }
1560        }
1561        if !unavailable_sources.is_empty() {
1562            return Err(MergeTaskError::sources(
1563                unavailable_sources,
1564                Error::Corruption(format!(
1565                    "merge sources are missing mandatory files: {}",
1566                    missing_files.join(", ")
1567                )),
1568            ));
1569        }
1570
1571        let schema_arc = Arc::clone(schema);
1572        let futures: Vec<_> = segment_ids
1573            .iter()
1574            .map(|&sid| {
1575                let sch = Arc::clone(&schema_arc);
1576                async move { SegmentReader::open(directory, sid, sch, term_cache_blocks).await }
1577            })
1578            .collect();
1579
1580        let results = futures::future::join_all(futures).await;
1581        let mut readers = Vec::with_capacity(results.len());
1582        let mut total_docs = 0u64;
1583        for (i, result) in results.into_iter().enumerate() {
1584            match result {
1585                Ok(r) => {
1586                    total_docs += r.meta().num_docs as u64;
1587                    readers.push(r);
1588                }
1589                Err(e) => {
1590                    log::error!(
1591                        "[merge] Failed to open segment {}: {:?}",
1592                        segment_ids_to_merge[i],
1593                        e
1594                    );
1595                    return Err(classify_source_error(segment_ids_to_merge[i].clone(), e));
1596                }
1597            }
1598        }
1599        if total_docs > u32::MAX as u64 {
1600            return Err(Error::Internal(format!(
1601                "Merged segment doc count ({}) exceeds u32::MAX",
1602                total_docs
1603            ))
1604            .into());
1605        }
1606
1607        // Pre-merge validation: verify each source segment's store doc count
1608        // matches its metadata. Catching mismatches early avoids building a
1609        // corrupted merged segment and leaving orphan files on disk.
1610        for (i, reader) in readers.iter().enumerate() {
1611            let meta_docs = reader.meta().num_docs;
1612            let store_docs = reader.store().num_docs();
1613            if store_docs != meta_docs {
1614                return Err(MergeTaskError::source(
1615                    segment_ids_to_merge[i].clone(),
1616                    Error::Corruption(format!(
1617                        "pre-merge validation: segment {} store has {} docs but meta says {}",
1618                        segment_ids_to_merge[i], store_docs, meta_docs
1619                    )),
1620                ));
1621            }
1622        }
1623
1624        log::info!(
1625            "[merge] loaded {} segment readers in {:.1}s",
1626            readers.len(),
1627            load_start.elapsed().as_secs_f64()
1628        );
1629
1630        let merger = SegmentMerger::new(Arc::clone(schema))
1631            .with_bmp_reorder(reorder_bmp)
1632            .with_granularity(granularity)
1633            .with_bp_budget(crate::segment::BpBudget {
1634                min_partition_docs: None,
1635                time_budget: merge_bp_time_budget,
1636            })
1637            .with_bp_memory_budget(bp_memory_budget_bytes)
1638            .with_reorder_permits(reorder_permits)
1639            .with_background_pool(bg_cpu_pool);
1640
1641        log::info!(
1642            "[merge] {} segments -> {} (trained={})",
1643            segment_ids_to_merge.len(),
1644            output_hex,
1645            trained.map_or(0, |t| t.centroids.len()),
1646        );
1647
1648        let (_merged_meta, merge_stats) = merger
1649            .merge(directory, &readers, output_segment_id, trained)
1650            .await
1651            .map_err(|error| {
1652                if matches!(error, Error::Corruption(_) | Error::Serialization(_)) {
1653                    // The merge has already opened every input successfully;
1654                    // a structural/serialization failure is deterministic for
1655                    // this candidate. Attribute all inputs rather than running
1656                    // the same multi-GB rewrite forever. This is deliberately
1657                    // not used for I/O errors, which may be transient/output-side.
1658                    MergeTaskError::sources(segment_ids_to_merge.to_vec(), error)
1659                } else {
1660                    MergeTaskError::from(error)
1661                }
1662            })?;
1663        let bp_converged = merge_stats.bp_converged;
1664        if !bp_converged {
1665            log::info!(
1666                "[merge] merge-time BP hit its wall-clock budget — output marked unconverged; \
1667                 the background optimizer deepens it later",
1668            );
1669        }
1670
1671        log::info!(
1672            "[merge] total wall-clock: {:.1}s ({} segments, {} docs)",
1673            load_start.elapsed().as_secs_f64(),
1674            readers.len(),
1675            total_docs,
1676        );
1677
1678        Ok((output_hex, total_docs as u32, bp_converged))
1679    }
1680
1681    /// Drain all in-flight merge tasks safely.
1682    ///
1683    /// Tokio cannot abort a `spawn_blocking` closure once it has started. The
1684    /// old implementation aborted only the async wrapper and returned while
1685    /// merge-time BP still owned an `OffsetWriter`, allowing index deletion or
1686    /// orphan cleanup to race a live writer. Awaiting is the only sound generic
1687    /// behavior until every merge phase supports cooperative cancellation.
1688    pub async fn abort_merges(&self) {
1689        loop {
1690            let handles: Vec<JoinHandle<()>> = { std::mem::take(&mut *self.merge_handles.lock()) };
1691            if handles.is_empty() {
1692                return;
1693            }
1694            for handle in handles {
1695                if let Err(error) = handle.await
1696                    && error.is_panic()
1697                {
1698                    log::error!("[merge] background task panicked while draining: {}", error);
1699                }
1700            }
1701        }
1702    }
1703
1704    /// Wait for all current in-flight merges to complete.
1705    pub async fn wait_for_merging_thread(self: &Arc<Self>) {
1706        let handles: Vec<JoinHandle<()>> = { std::mem::take(&mut *self.merge_handles.lock()) };
1707        for h in handles {
1708            let _ = h.await;
1709        }
1710    }
1711
1712    /// Wait for all eligible merges to complete, including cascading merges.
1713    ///
1714    /// Drains current handles, then loops. Each completed merge auto-triggers
1715    /// `maybe_merge` (which pushes new handles) before its JoinHandle resolves,
1716    /// so by the time `h.await` returns all cascading handles are registered.
1717    pub async fn wait_for_all_merges(self: &Arc<Self>) {
1718        loop {
1719            let handles: Vec<JoinHandle<()>> = { std::mem::take(&mut *self.merge_handles.lock()) };
1720            if handles.is_empty() {
1721                break;
1722            }
1723            for h in handles {
1724                let _ = h.await;
1725            }
1726        }
1727    }
1728
1729    /// Complete the second half of shutdown after the owning `IndexWriter`
1730    /// has been dropped. This drains tracked merges and then waits for every
1731    /// remaining guard, including optimizer reorders that are intentionally
1732    /// launched outside the writer lock.
1733    pub async fn wait_for_shutdown(self: &Arc<Self>) {
1734        self.wait_for_all_merges().await;
1735        self.active_operations.wait_until_idle().await;
1736        loop {
1737            let handles = { std::mem::take(&mut *self.lifecycle_handles.lock()) };
1738            if handles.is_empty() {
1739                break;
1740            }
1741            for handle in handles {
1742                if let Err(error) = handle.await
1743                    && error.is_panic()
1744                {
1745                    log::error!("[segment_cleanup] task panicked while draining: {}", error);
1746                }
1747            }
1748        }
1749    }
1750
1751    /// Force merge segments into the fewest possible segments, respecting
1752    /// `max_segment_docs` from the merge policy.
1753    ///
1754    /// If the policy defines a max segment size, segments are merged in batches
1755    /// that stay within that limit. Otherwise, all segments are merged into one.
1756    ///
1757    /// Each batch is registered in `active_operations` via an RAII guard to prevent
1758    /// `maybe_merge` from spawning a conflicting background merge.
1759    pub async fn force_merge(self: &Arc<Self>) -> Result<()> {
1760        const FORCE_MERGE_BATCH: usize = 64;
1761
1762        let max_segment_docs = {
1763            let st = self.state.lock().await;
1764            st.merge_policy.max_segment_docs()
1765        };
1766
1767        // Wait for all in-flight background merges (including cascading)
1768        // before starting forced merges to avoid try_register conflicts.
1769        self.wait_for_all_merges().await;
1770
1771        loop {
1772            if !self.active_operations.is_accepting() {
1773                return Err(Error::IndexClosed);
1774            }
1775            // Get segment IDs with their doc counts, sorted ascending by size
1776            let mut segments: Vec<(String, u32)> = {
1777                let st = self.state.lock().await;
1778                st.metadata
1779                    .segment_metas
1780                    .iter()
1781                    .map(|(id, info)| (id.clone(), info.num_docs))
1782                    .collect()
1783            };
1784
1785            if segments.len() < 2 {
1786                return Ok(());
1787            }
1788
1789            segments.sort_by_key(|(_, docs)| *docs);
1790
1791            // Build a batch respecting max_segment_docs
1792            let max_docs = max_segment_docs.map(|m| m as u64).unwrap_or(u64::MAX);
1793            let mut batch = Vec::new();
1794            let mut batch_docs = 0u64;
1795
1796            for (id, docs) in &segments {
1797                if batch.len() >= FORCE_MERGE_BATCH {
1798                    break;
1799                }
1800                let next_total = batch_docs + *docs as u64;
1801                if next_total > max_docs && !batch.is_empty() {
1802                    break;
1803                }
1804                batch.push(id.clone());
1805                batch_docs += *docs as u64;
1806            }
1807
1808            if batch.len() < 2 {
1809                return Ok(());
1810            }
1811
1812            log::info!(
1813                "[force_merge] merging batch of {} segments ({} docs)",
1814                batch.len(),
1815                batch_docs
1816            );
1817
1818            let _global_merge_permit = tokio::select! {
1819                biased;
1820                () = self.active_operations.wait_for_shutdown() => {
1821                    return Err(Error::IndexClosed);
1822                }
1823                permit = Arc::clone(&self.global_merge_permits).acquire_owned() => {
1824                    permit.map_err(|_| {
1825                        Error::Internal("global background merge scheduler is closed".into())
1826                    })?
1827                }
1828            };
1829
1830            let output_id = SegmentId::new();
1831            let output_hex = output_id.to_hex();
1832
1833            // Register batch + output under `state`, matching orphan cleanup's
1834            // deletion barrier and preventing a stale batch from starting.
1835            let mut all_ids = batch.clone();
1836            all_ids.push(output_hex);
1837            let guard = {
1838                let st = self.state.lock().await;
1839                batch
1840                    .iter()
1841                    .all(|id| st.metadata.has_segment(id))
1842                    .then(|| self.active_operations.try_register(all_ids))
1843                    .flatten()
1844            };
1845            let _guard = match guard {
1846                Some(g) => g,
1847                None if !self.active_operations.is_accepting() => {
1848                    return Err(Error::IndexClosed);
1849                }
1850                None => {
1851                    // A background merge slipped in — wait for it, then retry the loop
1852                    self.wait_for_merging_thread().await;
1853                    continue;
1854                }
1855            };
1856            let mut output_cleanup = self.output_cleanup_guard(output_id);
1857
1858            let trained_snap = self.trained_for_segment_build();
1859            let granularity = self.merge_granularity(&batch).await;
1860            let merge_result = Self::do_merge(
1861                self.directory.as_ref(),
1862                &self.schema,
1863                &batch,
1864                output_id,
1865                self.term_cache_blocks,
1866                trained_snap.as_deref(),
1867                self.reorder_on_merge,
1868                granularity,
1869                self.merge_bp_time_budget,
1870                self.bp_memory_budget_bytes,
1871                Arc::clone(&self.reorder_permits),
1872                Some(self.background_cpu_pool()),
1873            )
1874            .await;
1875            let (new_segment_id, total_docs, bp_converged) = match merge_result {
1876                Ok(v) => v,
1877                Err(MergeTaskError {
1878                    error,
1879                    unavailable_segments,
1880                }) => {
1881                    for segment_id in &unavailable_segments {
1882                        self.quarantine_segment(segment_id, &error);
1883                    }
1884                    self.delete_output_if_unregistered(output_id, "force-merge failure")
1885                        .await;
1886                    output_cleanup.disarm();
1887                    return Err(error);
1888                }
1889            };
1890
1891            if let Err(e) = self
1892                .replace_segments(
1893                    &batch,
1894                    new_segment_id,
1895                    total_docs,
1896                    self.reorder_on_merge,
1897                    bp_converged,
1898                )
1899                .await
1900            {
1901                self.delete_output_if_unregistered(output_id, "replacement failure")
1902                    .await;
1903                output_cleanup.disarm();
1904                return Err(e);
1905            }
1906            output_cleanup.disarm();
1907
1908            // _guard drops here, releasing operation ownership.
1909        }
1910    }
1911
1912    /// Reorder all segments via Recursive Graph Bisection (BP) for better BMP pruning.
1913    ///
1914    /// Each segment is individually rebuilt with reordered BMP blocks.
1915    /// Non-BMP fields are copied unchanged via streaming file copy.
1916    ///
1917    /// Uses active-operation ownership to prevent concurrent work on the same segment.
1918    pub async fn reorder_segments(self: &Arc<Self>) -> Result<()> {
1919        self.wait_for_all_merges().await;
1920        let segment_ids = self.get_segment_ids().await;
1921
1922        if segment_ids.is_empty() {
1923            log::info!("[reorder] no segments to reorder");
1924            return Ok(());
1925        }
1926
1927        log::info!("[reorder] reordering {} segments", segment_ids.len());
1928
1929        for seg_id in segment_ids {
1930            match self
1931                .reorder_single_segment(&seg_id, None, crate::segment::BpBudget::full())
1932                .await
1933            {
1934                Ok(true) => {}
1935                Ok(false) => log::warn!("[reorder] segment {} skipped (in merge)", seg_id),
1936                Err(e) => return Err(e),
1937            }
1938        }
1939
1940        log::info!("[reorder] all segments reordered");
1941        Ok(())
1942    }
1943
1944    /// Get segment IDs that have not been reordered yet.
1945    ///
1946    /// Excludes segments currently involved in a merge or reorder operation
1947    /// to avoid wasted work (the optimizer would skip them anyway).
1948    pub async fn unreordered_segment_ids(&self) -> Vec<String> {
1949        self.unreordered_segments()
1950            .await
1951            .into_iter()
1952            .map(|(id, _)| id)
1953            .collect()
1954    }
1955
1956    /// Segments never reordered, with doc counts — for the optimizer to pick
1957    /// a size-appropriate BP budget.
1958    pub async fn unreordered_segments(&self) -> Vec<(String, u32)> {
1959        let quarantined = self.quarantined_segments.lock().clone();
1960        let paused = self.paused_reorder_segments();
1961        let st = self.state.lock().await;
1962        let active_ids = self.active_operations.snapshot();
1963        st.metadata
1964            .segment_metas
1965            .iter()
1966            .filter(|(id, info)| {
1967                !info.reordered
1968                    && !active_ids.contains(*id)
1969                    && !quarantined.contains(*id)
1970                    && !paused.contains(*id)
1971            })
1972            .map(|(id, info)| (id.clone(), info.num_docs))
1973            .collect()
1974    }
1975
1976    /// Segments whose last BP pass hit its wall-clock budget before finishing
1977    /// (`bp_converged == false`). A warm-started follow-up pass deepens the
1978    /// ordering; the optimizer revisits these at low priority.
1979    pub async fn unconverged_segments(&self) -> Vec<(String, u32)> {
1980        self.unconverged_segments_below(u32::MAX)
1981            .await
1982            .into_iter()
1983            .map(|(id, docs, _)| (id, docs))
1984            .collect()
1985    }
1986
1987    /// Unconverged segments still below a hard replacement-lineage work
1988    /// bound. Includes the persisted attempt count for scheduler diagnostics.
1989    pub async fn unconverged_segments_below(
1990        &self,
1991        max_unconverged_passes: u32,
1992    ) -> Vec<(String, u32, u32)> {
1993        let quarantined = self.quarantined_segments.lock().clone();
1994        let paused = self.paused_reorder_segments();
1995        let st = self.state.lock().await;
1996        let active_ids = self.active_operations.snapshot();
1997        st.metadata
1998            .segment_metas
1999            .iter()
2000            .filter(|(id, info)| {
2001                info.reordered
2002                    && !info.bp_converged
2003                    && info.bp_unconverged_passes < max_unconverged_passes
2004                    && !active_ids.contains(*id)
2005                    && !quarantined.contains(*id)
2006                    && !paused.contains(*id)
2007            })
2008            .map(|(id, info)| (id.clone(), info.num_docs, info.bp_unconverged_passes))
2009            .collect()
2010    }
2011
2012    /// Granularity for a BP pass whose sources are `ids`: `Records` when any
2013    /// source is an unconverged partial reorder, `Auto` otherwise.
2014    ///
2015    /// Alignment with the depth budget (docs/block-level-reorder.md): an
2016    /// unconverged segment is owed a deepening pass, and the output of this
2017    /// pass will be marked `bp_converged`. `Auto` would measure the partial
2018    /// pass's residual coherence, potentially take the blockwise path — which
2019    /// cannot deepen record clustering — and end the cascade at partial
2020    /// quality. Only record-level BP discharges the debt.
2021    async fn merge_granularity(&self, ids: &[String]) -> crate::segment::reorder::BpGranularity {
2022        let st = self.state.lock().await;
2023        let deepening = ids.iter().any(|id| {
2024            st.metadata
2025                .segment_metas
2026                .get(id)
2027                .is_some_and(|info| info.reordered && !info.bp_converged)
2028        });
2029        drop(st);
2030        if deepening {
2031            log::info!(
2032                "[reorder] source segment(s) unconverged — forcing record-level BP (deepening pass)",
2033            );
2034            crate::segment::reorder::BpGranularity::Records
2035        } else {
2036            crate::segment::reorder::BpGranularity::Auto
2037        }
2038    }
2039
2040    /// Reorder a single segment via BP. Returns Ok(true) if reordered, Ok(false) if skipped.
2041    ///
2042    /// Non-blocking: operation ownership prevents conflicts with background merges.
2043    /// Copies unchanged files and rebuilds only the sparse file with reordered BMP data.
2044    pub async fn reorder_single_segment(
2045        self: &Arc<Self>,
2046        seg_id: &str,
2047        rayon_pool: Option<Arc<rayon::ThreadPool>>,
2048        bp_budget: crate::segment::BpBudget,
2049    ) -> Result<bool> {
2050        let source_id = SegmentId::from_hex(seg_id)
2051            .ok_or_else(|| Error::Corruption(format!("Invalid segment ID: {}", seg_id)))?;
2052        if self.quarantined_segments.lock().contains(seg_id) {
2053            return Err(Error::Corruption(format!(
2054                "segment {} is quarantined after a deterministic source failure; repair it and restart before reordering",
2055                seg_id
2056            )));
2057        }
2058
2059        // Whole-pass concurrency is independent from Rayon width. One pass
2060        // can already use every configured BP worker; this permit bounds the
2061        // much larger forward-index and rewrite working set across indexes,
2062        // optimizer tasks, and merge-time BP.
2063        let _reorder_permit = tokio::select! {
2064            biased;
2065            () = self.active_operations.wait_for_shutdown() => {
2066                return Err(Error::IndexClosed);
2067            }
2068            permit = Arc::clone(&self.reorder_permits).acquire_owned() => {
2069                permit.map_err(|_| {
2070                    Error::Internal("background reorder scheduler is closed".into())
2071                })?
2072            }
2073        };
2074
2075        let output_id = SegmentId::new();
2076        let output_hex = output_id.to_hex();
2077        let source_ids = [seg_id.to_string()];
2078        let granularity = self.merge_granularity(&source_ids).await;
2079
2080        // Register while holding `state`, matching orphan cleanup's deletion
2081        // barrier. Candidates are scanned ahead of time and can go stale: a
2082        // merge may have consumed this segment since. Its files may even still
2083        // be on disk (deferred deletion under a searcher snapshot) — reordering
2084        // them would re-insert a duplicate copy of docs the merge output holds.
2085        let all_ids = vec![seg_id.to_string(), output_hex];
2086        let (_guard, source_docs) = {
2087            let st = self.state.lock().await;
2088            let Some(source_meta) = st.metadata.segment_metas.get(seg_id) else {
2089                log::info!(
2090                    "[optimizer] segment {} no longer in metadata (merged away), skipping reorder",
2091                    seg_id
2092                );
2093                self.clear_reorder_retry(seg_id);
2094                return Ok(false);
2095            };
2096
2097            match self.active_operations.try_register(all_ids) {
2098                Some(guard) => (guard, source_meta.num_docs),
2099                None if !self.active_operations.is_accepting() => {
2100                    return Err(Error::IndexClosed);
2101                }
2102                None => {
2103                    log::debug!("[optimizer] segment {} in active merge, skipping", seg_id);
2104                    return Ok(false);
2105                }
2106            }
2107        };
2108
2109        // Fail before allocating a forward index or creating output files.
2110        // Missing mandatory files are deterministic and should remove this
2111        // segment from future optimizer scans, not consume the same CPU every
2112        // interval. Other I/O failures remain retryable.
2113        if let Err(error) = self.validate_completed_segment(seg_id, source_docs).await {
2114            if is_deterministic_source_error(&error) {
2115                self.quarantine_segment(seg_id, &error);
2116            } else if !matches!(&error, Error::IndexClosed) {
2117                self.pause_reorder_retries(seg_id, &error);
2118            }
2119            return Err(error);
2120        }
2121
2122        let mut output_cleanup = self.output_cleanup_guard(output_id);
2123
2124        let reorder_result = crate::segment::reorder::reorder_segment(
2125            self.directory.as_ref(),
2126            &self.schema,
2127            source_id,
2128            output_id,
2129            self.term_cache_blocks,
2130            self.bp_memory_budget_bytes,
2131            bp_budget,
2132            granularity,
2133            rayon_pool,
2134        )
2135        .await;
2136        let (new_id, total_docs, bp_converged) = match reorder_result {
2137            Ok(v) => v,
2138            Err(e) => {
2139                // A failed pass may have copied tens of GB before dying;
2140                // delete the uncommitted output before propagating.
2141                self.delete_output_if_unregistered(output_id, "reorder failure")
2142                    .await;
2143                output_cleanup.disarm();
2144                if is_deterministic_source_error(&e) {
2145                    self.quarantine_segment(seg_id, &e);
2146                } else if !matches!(&e, Error::IndexClosed) {
2147                    self.pause_reorder_retries(seg_id, &e);
2148                }
2149                return Err(e);
2150            }
2151        };
2152
2153        // A pass with a depth floor above block granularity has, by
2154        // definition, not converged to block-level order — record it as
2155        // unconverged so the optimizer's deepening ladder revisits it with a
2156        // full-depth (warm-started) pass. Depth caps are only used by the
2157        // optimizer's first pass on large segments.
2158        let ladder_converged = bp_converged && bp_budget.min_partition_docs.is_none();
2159        if let Err(e) = self
2160            .replace_segments(
2161                &[seg_id.to_string()],
2162                new_id,
2163                total_docs,
2164                true,
2165                ladder_converged,
2166            )
2167            .await
2168        {
2169            self.delete_output_if_unregistered(output_id, "replacement failure")
2170                .await;
2171            output_cleanup.disarm();
2172            if !matches!(&e, Error::IndexClosed) {
2173                self.pause_reorder_retries(seg_id, &e);
2174            }
2175            return Err(e);
2176        }
2177        output_cleanup.disarm();
2178        self.clear_reorder_retry(seg_id);
2179
2180        Ok(true)
2181    }
2182
2183    /// Clean up orphan segment files not registered in metadata.
2184    ///
2185    /// Reads metadata, active-operation ownership, and snapshot-deferred
2186    /// deletions to determine which segments are legitimate. Filesystem
2187    /// deletion is asynchronous; in-flight outputs and retired sources still
2188    /// held by readers are both protected.
2189    pub async fn cleanup_orphan_segments(&self) -> Result<usize> {
2190        let mut orphan_files: HashMap<String, Vec<std::path::PathBuf>> = HashMap::new();
2191
2192        if let Ok(entries) = self.directory.list_files(std::path::Path::new("")).await {
2193            for entry in entries {
2194                let Some(filename) = entry.file_name().and_then(|name| name.to_str()) else {
2195                    continue;
2196                };
2197                let Some(rest) = filename.strip_prefix("seg_") else {
2198                    continue;
2199                };
2200                let Some(hex_id) = rest.get(..32) else {
2201                    continue;
2202                };
2203                if !hex_id.bytes().all(|byte| byte.is_ascii_hexdigit()) {
2204                    continue;
2205                }
2206                orphan_files
2207                    .entry(hex_id.to_ascii_lowercase())
2208                    .or_default()
2209                    .push(entry);
2210            }
2211        }
2212
2213        let mut deleted = 0;
2214        for (hex_id, paths) in &orphan_files {
2215            // Revalidate and atomically claim deletion under the same
2216            // state -> active_operations -> tracker order used by publishers.
2217            // The claim lets us release `state` before filesystem I/O: deleting
2218            // a multi-GB orphan must not freeze commits and snapshot acquisition.
2219            let deletion_guard = {
2220                let st = self.state.lock().await;
2221                if st.metadata.has_segment(hex_id) {
2222                    continue;
2223                }
2224                let Some(guard) = self
2225                    .active_operations
2226                    .try_register(vec![hex_id.to_string()])
2227                else {
2228                    continue;
2229                };
2230                if self.tracker.is_deletion_protected(hex_id) {
2231                    drop(guard);
2232                    continue;
2233                }
2234                guard
2235            };
2236
2237            // Delete what was actually discovered, not only the currently
2238            // known SegmentFiles extensions. This also removes partial files
2239            // left by older formats instead of reporting the same orphan on
2240            // every startup forever.
2241            let results =
2242                futures::future::join_all(paths.iter().map(|path| self.directory.delete(path)))
2243                    .await;
2244            let removed = results.into_iter().all(|result| match result {
2245                Ok(()) => true,
2246                Err(error) if error.kind() == std::io::ErrorKind::NotFound => true,
2247                Err(error) => {
2248                    log::warn!(
2249                        "[segment_cleanup] failed sweeping orphan segment {}: {}",
2250                        hex_id,
2251                        error,
2252                    );
2253                    false
2254                }
2255            });
2256            // Releasing this claim is the deletion barrier. No producer can
2257            // adopt the ID while its files are being removed.
2258            drop(deletion_guard);
2259            if removed {
2260                deleted += 1;
2261                log::info!("[segment_cleanup] swept orphan segment {}", hex_id);
2262            }
2263        }
2264
2265        Ok(deleted)
2266    }
2267}
2268
2269#[cfg(test)]
2270mod tests {
2271    use super::*;
2272    use std::sync::atomic::{AtomicBool, Ordering};
2273
2274    fn lifecycle_test_manager() -> Arc<SegmentManager<crate::directories::RamDirectory>> {
2275        let schema = crate::dsl::SchemaBuilder::default().build();
2276        let metadata = IndexMetadata::new(schema.clone());
2277        Arc::new(SegmentManager::new(
2278            Arc::new(crate::directories::RamDirectory::new()),
2279            Arc::new(schema),
2280            metadata,
2281            Box::new(crate::merge::NoMergePolicy),
2282            0,
2283            1,
2284            Arc::new(Semaphore::new(1)),
2285            None,
2286            1024,
2287            Arc::new(Semaphore::new(1)),
2288            None,
2289        ))
2290    }
2291
2292    #[test]
2293    fn output_cleanup_guard_runs_during_panic_unwind() {
2294        let cleaned = Arc::new(AtomicBool::new(false));
2295        let cleaned_in_callback = Arc::clone(&cleaned);
2296        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
2297            cleaned_in_callback.store(true, Ordering::SeqCst);
2298        });
2299
2300        let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
2301            let _guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
2302            panic!("simulated reorder panic");
2303        }));
2304
2305        assert!(result.is_err());
2306        assert!(
2307            cleaned.load(Ordering::SeqCst),
2308            "partial output cleanup must run during unwind"
2309        );
2310    }
2311
2312    #[test]
2313    fn output_cleanup_guard_disarms_after_commit() {
2314        let cleaned = Arc::new(AtomicBool::new(false));
2315        let cleaned_in_callback = Arc::clone(&cleaned);
2316        let cleanup: Arc<dyn Fn(SegmentId) + Send + Sync> = Arc::new(move |_| {
2317            cleaned_in_callback.store(true, Ordering::SeqCst);
2318        });
2319
2320        {
2321            let mut guard = OutputCleanupGuard::new(SegmentId::new(), cleanup);
2322            guard.disarm();
2323        }
2324
2325        assert!(!cleaned.load(Ordering::SeqCst));
2326    }
2327
2328    #[test]
2329    fn test_active_operation_guard_releases_ownership() {
2330        let active = Arc::new(ActiveSegmentOperations::new());
2331        {
2332            let _guard = active.try_register(vec!["a".into(), "b".into()]).unwrap();
2333            let snap = active.snapshot();
2334            assert!(snap.contains("a"));
2335            assert!(snap.contains("b"));
2336        }
2337        assert!(active.snapshot().is_empty());
2338    }
2339
2340    #[test]
2341    fn test_non_overlapping_operations_can_run_concurrently() {
2342        let active = Arc::new(ActiveSegmentOperations::new());
2343        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
2344        let _second = active.try_register(vec!["c".into(), "d".into()]).unwrap();
2345        let snap = active.snapshot();
2346        assert_eq!(snap.len(), 4);
2347
2348        drop(first);
2349        let snap = active.snapshot();
2350        assert_eq!(snap.len(), 2);
2351        assert!(snap.contains("c"));
2352        assert!(snap.contains("d"));
2353    }
2354
2355    #[test]
2356    fn test_overlapping_operation_is_rejected_until_release() {
2357        let active = Arc::new(ActiveSegmentOperations::new());
2358        let first = active.try_register(vec!["a".into(), "b".into()]).unwrap();
2359        assert!(active.try_register(vec!["b".into(), "c".into()]).is_none());
2360        drop(first);
2361        assert!(active.try_register(vec!["b".into(), "c".into()]).is_some());
2362    }
2363
2364    #[test]
2365    fn test_active_operation_snapshot() {
2366        let active = Arc::new(ActiveSegmentOperations::new());
2367        let _guard = active.try_register(vec!["x".into(), "y".into()]).unwrap();
2368        let snap = active.snapshot();
2369        assert!(snap.contains("x"));
2370        assert!(snap.contains("y"));
2371        assert!(!snap.contains("z"));
2372    }
2373
2374    #[tokio::test]
2375    async fn operation_barrier_ignores_producers_started_after_snapshot() {
2376        let active = Arc::new(ActiveSegmentOperations::new());
2377        let before_gate = active.try_register(vec!["old".into()]).unwrap();
2378        let barrier = active.operation_tokens_snapshot();
2379        let after_gate = active.try_register(vec!["new-flat".into()]).unwrap();
2380
2381        let waiter = {
2382            let active = Arc::clone(&active);
2383            tokio::spawn(async move { active.wait_until_operations_finish(&barrier).await })
2384        };
2385        tokio::task::yield_now().await;
2386        assert!(!waiter.is_finished());
2387
2388        drop(before_gate);
2389        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
2390            .await
2391            .expect("pre-gate operation barrier was starved by a post-gate producer")
2392            .unwrap();
2393        assert!(active.snapshot().contains("new-flat"));
2394        drop(after_gate);
2395    }
2396
2397    #[tokio::test]
2398    async fn artifact_update_gate_preserves_search_generation_but_forces_flat_producers() {
2399        let manager = lifecycle_test_manager();
2400        manager
2401            .trained
2402            .store(Some(Arc::new(TrainedVectorStructures {
2403                centroids: rustc_hash::FxHashMap::default(),
2404                codebooks: rustc_hash::FxHashMap::default(),
2405            })));
2406
2407        let guard = manager.begin_vector_artifact_update().await.unwrap();
2408        assert!(
2409            manager.trained().is_some(),
2410            "search readers keep the last fully validated generation"
2411        );
2412        assert!(
2413            manager.trained_for_segment_build().is_none(),
2414            "new segment producers must stay flat during an artifact update"
2415        );
2416
2417        let detached_transaction_guard = guard.clone();
2418        drop(guard);
2419        assert!(
2420            manager.trained_for_segment_build().is_none(),
2421            "a detached lifecycle transaction must retain the producer gate after request cancellation"
2422        );
2423        drop(detached_transaction_guard);
2424        assert!(manager.trained_for_segment_build().is_some());
2425    }
2426
2427    #[tokio::test]
2428    async fn shutdown_rejects_new_work_and_waits_for_existing_guard() {
2429        let active = Arc::new(ActiveSegmentOperations::new());
2430        let guard = active.try_register(vec!["live".into()]).unwrap();
2431        active.stop_accepting();
2432        assert!(active.try_register(vec!["new".into()]).is_none());
2433
2434        let waiter = {
2435            let active = Arc::clone(&active);
2436            tokio::spawn(async move { active.wait_until_idle().await })
2437        };
2438        tokio::task::yield_now().await;
2439        assert!(!waiter.is_finished());
2440        drop(guard);
2441        tokio::time::timeout(std::time::Duration::from_secs(1), waiter)
2442            .await
2443            .expect("shutdown waiter missed the final guard notification")
2444            .unwrap();
2445    }
2446
2447    #[tokio::test]
2448    async fn lifecycle_transaction_survives_request_cancellation_and_is_drained() {
2449        let manager = lifecycle_test_manager();
2450        let started = Arc::new(Semaphore::new(0));
2451        let release = Arc::new(Semaphore::new(0));
2452        let completed = Arc::new(AtomicBool::new(false));
2453
2454        let request = {
2455            let manager = Arc::clone(&manager);
2456            let started = Arc::clone(&started);
2457            let release = Arc::clone(&release);
2458            let completed = Arc::clone(&completed);
2459            tokio::spawn(async move {
2460                manager
2461                    .run_lifecycle_transaction(async move {
2462                        started.add_permits(1);
2463                        let _permit = release.acquire().await.unwrap();
2464                        completed.store(true, Ordering::Release);
2465                        Ok(())
2466                    })
2467                    .await
2468            })
2469        };
2470
2471        let _started = started.acquire().await.unwrap();
2472        request.abort();
2473        assert!(request.await.unwrap_err().is_cancelled());
2474        release.add_permits(1);
2475
2476        manager.begin_shutdown();
2477        tokio::time::timeout(
2478            std::time::Duration::from_secs(1),
2479            manager.wait_for_shutdown(),
2480        )
2481        .await
2482        .expect("shutdown did not drain detached lifecycle transaction");
2483        assert!(completed.load(Ordering::Acquire));
2484    }
2485
2486    #[tokio::test]
2487    async fn unconverged_scheduler_stops_at_the_lineage_limit() {
2488        let manager = lifecycle_test_manager();
2489        {
2490            let mut state = manager.state.lock().await;
2491            state.metadata.add_segment_meta(
2492                "eligible".into(),
2493                SegmentMetaInfo {
2494                    num_docs: 10,
2495                    ancestors: Vec::new(),
2496                    generation: 1,
2497                    reordered: true,
2498                    bp_converged: false,
2499                    bp_unconverged_passes: 2,
2500                },
2501            );
2502            state.metadata.add_segment_meta(
2503                "at-limit".into(),
2504                SegmentMetaInfo {
2505                    num_docs: 20,
2506                    ancestors: Vec::new(),
2507                    generation: 1,
2508                    reordered: true,
2509                    bp_converged: false,
2510                    bp_unconverged_passes: 3,
2511                },
2512            );
2513            state.metadata.add_segment_meta(
2514                "converged".into(),
2515                SegmentMetaInfo {
2516                    num_docs: 30,
2517                    ancestors: Vec::new(),
2518                    generation: 1,
2519                    reordered: true,
2520                    bp_converged: true,
2521                    bp_unconverged_passes: 0,
2522                },
2523            );
2524            state.metadata.add_segment("fresh".into(), 40);
2525        }
2526
2527        assert_eq!(
2528            manager.unconverged_segments_below(3).await,
2529            vec![("eligible".into(), 10, 2)]
2530        );
2531        assert!(manager.unconverged_segments_below(0).await.is_empty());
2532    }
2533
2534    #[test]
2535    fn merge_retry_backoff_is_exponential_and_capped() {
2536        assert_eq!(merge_retry_delay(1), std::time::Duration::from_secs(30));
2537        assert_eq!(merge_retry_delay(2), std::time::Duration::from_secs(60));
2538        assert_eq!(merge_retry_delay(3), std::time::Duration::from_secs(120));
2539        assert_eq!(merge_retry_delay(100), MERGE_RETRY_MAX_DELAY);
2540    }
2541
2542    #[test]
2543    fn only_deterministic_source_errors_are_quarantined() {
2544        assert!(is_deterministic_source_error(&Error::Corruption(
2545            "bad footer".into()
2546        )));
2547        assert!(is_deterministic_source_error(&Error::Io(
2548            std::io::Error::from(std::io::ErrorKind::NotFound)
2549        )));
2550        assert!(!is_deterministic_source_error(&Error::Io(
2551            std::io::Error::from(std::io::ErrorKind::TimedOut)
2552        )));
2553        assert!(!is_deterministic_source_error(&Error::Io(
2554            std::io::Error::from(std::io::ErrorKind::PermissionDenied)
2555        )));
2556    }
2557
2558    #[test]
2559    fn transient_reorder_failure_is_backed_off_until_cleared() {
2560        let manager = lifecycle_test_manager();
2561        manager.pause_reorder_retries("source", &Error::Internal("transient".into()));
2562        assert!(manager.paused_reorder_segments().contains("source"));
2563        manager.clear_reorder_retry("source");
2564        assert!(!manager.paused_reorder_segments().contains("source"));
2565    }
2566}