Skip to main content

summa_core/merge/
row_mutation.rs

1//! Row visibility publication uses the segment manager's existing transaction.
2use super::*;
3use crate::query::DocBitset;
4
5impl<D: DirectoryWriter + 'static> SegmentManager<D> {
6    pub(crate) async fn commit_with_deletes(
7        self: &Arc<Self>,
8        new_segments: &[(String, u32)],
9        keys: Vec<String>,
10        staged_deletions: Vec<(String, Arc<crate::index::staged_row::StagedSegment>)>,
11    ) -> Result<()> {
12        if keys.is_empty() && staged_deletions.is_empty() {
13            return self.commit(new_segments).await;
14        }
15        let field = self
16            .schema
17            .primary_field()
18            .ok_or_else(|| Error::Schema("row deletion requires a primary key".into()))?;
19        for (id, count) in new_segments {
20            self.validate_completed_segment(id, *count).await?;
21        }
22        let mut st = Arc::clone(&self.state).lock_owned().await;
23        let manager = Arc::clone(self);
24        let new_segments = new_segments.to_vec();
25        self.run_lifecycle_transaction(async move {
26            let keys: Arc<HashSet<String>> = Arc::new(keys.into_iter().collect());
27            let mut next = st.metadata.clone();
28            let mut retired = Vec::new();
29            let mut outputs = Vec::new();
30            // Only committed rows are targets. The same transaction's replacement
31            // insertions must survive the deletion of their old primary keys.
32            for (segment_id, info) in st
33                .metadata
34                .segment_metas
35                .iter()
36                .filter(|_| !keys.is_empty())
37            {
38                let sid = SegmentId::from_hex(segment_id)
39                    .ok_or_else(|| Error::Corruption("invalid deletion source".into()))?;
40                let mut fields = crate::segment::reader::loader::load_fast_fields_file(
41                    manager.directory.as_ref(),
42                    &SegmentFiles::new(sid.0),
43                    &manager.schema,
44                )
45                .await?;
46                let ff = fields
47                    .remove(&field.0)
48                    .ok_or_else(|| Error::Corruption("primary-key column is missing".into()))?;
49                // Bounded by the admitted deletion-key set, independent of
50                // the source dictionary/row count. Decode global ordinals in
51                // batches instead of looking up and hashing each row's text.
52                let cancellation = manager.active_operations.cancellation_flag();
53                let lookup_cancellation = Arc::clone(&cancellation);
54                let keys = Arc::clone(&keys);
55                let num_docs = info.num_docs;
56                // Publication holds the state lock. Run its bounded serial scan
57                // directly on Tokio's blocking executor, never behind bulk BP.
58                let (ff, ordinals) = tokio::task::spawn_blocking(move || {
59                    let ordinals = crate::segment::deletion::target_ordinals(
60                        &ff,
61                        num_docs,
62                        keys.iter().map(String::as_str),
63                        || {
64                            if lookup_cancellation.load(Ordering::Acquire) {
65                                Err(Error::IndexClosed)
66                            } else {
67                                Ok(())
68                            }
69                        },
70                    )?;
71                    Ok::<_, Error>((ff, ordinals))
72                })
73                .await
74                .map_err(|error| {
75                    Error::Internal(format!("deletion key lookup failed: {error}"))
76                })??;
77                if ordinals.is_empty() {
78                    continue;
79                }
80                let mut alive = match &info.deletions {
81                    Some(meta) => {
82                        (*meta.load(manager.directory.as_ref(), info.num_docs).await?).clone()
83                    }
84                    None => DocBitset::all(info.num_docs),
85                };
86                let (alive, changed) = tokio::task::spawn_blocking(move || {
87                    let changed = crate::segment::deletion::clear_target_rows(
88                        &ff,
89                        &ordinals,
90                        &mut alive,
91                        || {
92                            if cancellation.load(Ordering::Acquire) {
93                                Err(Error::IndexClosed)
94                            } else {
95                                Ok(())
96                            }
97                        },
98                    )?;
99                    Ok::<_, Error>((alive, changed))
100                })
101                .await
102                .map_err(|error| Error::Internal(format!("deletion worker failed: {error}")))??;
103                if !changed {
104                    continue;
105                }
106                let id = SegmentId::new();
107                let claim = manager.protect_new_segment(id.to_hex())?;
108                let cleanup = manager.output_cleanup_guard(id);
109                // Install the unwind owner before writing the first byte.
110                outputs.push((claim, cleanup));
111                let deletion = crate::segment::deletion::write(
112                    manager.directory.as_ref(),
113                    id,
114                    info.num_docs,
115                    &alive,
116                )
117                .await?;
118                if let Some(old) = &info.deletions {
119                    retired.push(old.id.clone());
120                }
121                next.segment_metas.get_mut(segment_id).unwrap().deletions = Some(deletion);
122            }
123            for (id, count) in &new_segments {
124                if !next.has_segment(id) {
125                    next.add_segment(id.clone(), *count);
126                }
127            }
128            for (segment_id, staged) in &staged_deletions {
129                if st.metadata.has_segment(segment_id) {
130                    return Err(Error::Corruption(
131                        "staged visibility must refer to a new segment".into(),
132                    ));
133                }
134                let info = next
135                    .segment_metas
136                    .get_mut(segment_id)
137                    .ok_or_else(|| Error::Corruption("staged segment is missing".into()))?;
138                let alive = staged.live_rows(info.num_docs).ok_or_else(|| {
139                    Error::Internal("staged visibility has no cancelled rows".into())
140                })?;
141                let id = SegmentId::new();
142                let claim = manager.protect_new_segment(id.to_hex())?;
143                let cleanup = manager.output_cleanup_guard(id);
144                outputs.push((claim, cleanup));
145                info.deletions = Some(
146                    crate::segment::deletion::write(
147                        manager.directory.as_ref(),
148                        id,
149                        info.num_docs,
150                        &alive,
151                    )
152                    .await?,
153                );
154            }
155            next.publication_generation = next
156                .publication_generation
157                .checked_add(1)
158                .ok_or_else(|| Error::Corruption("publication generation overflow".into()))?;
159            next.save(manager.directory.as_ref()).await?;
160            for id in next.owned_ids() {
161                manager.tracker.register(&id);
162            }
163            let previous = manager.published_generation();
164            manager
165                .published_generation
166                .store(Arc::new(PublishedIndexGeneration {
167                    publication_id: next.publication_generation,
168                    schema: Arc::clone(&previous.schema),
169                    trained_vectors: previous.trained_vectors.clone(),
170                }));
171            retired.retain(|id| !next.owns_id(id));
172            st.metadata = next;
173            for (_, cleanup) in &mut outputs {
174                cleanup.disarm();
175            }
176            let ready = manager.tracker.mark_for_deletion(&retired);
177            drop(st);
178            (manager.delete_fn)(ready);
179            Ok(())
180        })
181        .await
182    }
183}
184
185impl<D: DirectoryWriter + 'static> SegmentManager<D> {
186    /// Rewrite one segment with a dense live-row space, including single-segment
187    /// indexes. Returns false if it has no tombstones or is no longer current.
188    pub async fn compact_segment(
189        self: &Arc<Self>,
190        segment_id: &str,
191        memory_budget: usize,
192    ) -> Result<bool> {
193        self.compact_segment_admitted(segment_id, memory_budget, None)
194            .await
195    }
196
197    /// Metadata-only candidates for the existing periodic optimizer. Zero
198    /// threshold disables automatic work; invalid ratios fail before I/O.
199    pub async fn compaction_candidates(&self, threshold: f64) -> Result<Vec<(String, u32, u32)>> {
200        validate_threshold(threshold)?;
201        if threshold == 0.0 || self.force_merge_active.load(Ordering::Acquire) > 0 {
202            return Ok(Vec::new());
203        }
204        let active = self.active_operations.snapshot();
205        let paused = self.paused_reorder_segments();
206        let quarantined = self.quarantined_segments.lock().clone();
207        let st = self.state.lock().await;
208        let mut candidates: Vec<_> = st
209            .metadata
210            .segment_metas
211            .iter()
212            .filter(|(id, meta)| {
213                meta.num_deleted_docs() > 0
214                    && meta.deleted_ratio() >= threshold
215                    && !active.contains(*id)
216                    && !paused.contains(*id)
217                    && !quarantined.contains(*id)
218            })
219            .map(|(id, meta)| (id.clone(), meta.num_docs, meta.num_deleted_docs()))
220            .collect();
221        candidates.sort_unstable_by(|a, b| {
222            (f64::from(b.2) / f64::from(b.1))
223                .total_cmp(&(f64::from(a.2) / f64::from(a.1)))
224                .then_with(|| a.1.cmp(&b.1))
225                .then_with(|| a.0.cmp(&b.0))
226        });
227        Ok(candidates)
228    }
229
230    /// Nonblocking maintenance admission, using the same ownership, pool,
231    /// capacity and bounded retry state as existing optimizer work.
232    pub async fn compact_segment_if_eligible(
233        self: &Arc<Self>,
234        segment_id: &str,
235        threshold: f64,
236        memory_budget: usize,
237    ) -> Result<bool> {
238        validate_threshold(threshold)?;
239        if threshold == 0.0 {
240            return Ok(false);
241        }
242        let result = self
243            .compact_segment_admitted(segment_id, memory_budget, Some(threshold))
244            .await;
245        match &result {
246            Ok(true) => self.clear_reorder_retry(segment_id),
247            Err(error) if !matches!(error, Error::IndexClosed) => {
248                self.pause_reorder_retries(segment_id, error).await
249            }
250            _ => {}
251        }
252        result
253    }
254
255    async fn compact_segment_admitted(
256        self: &Arc<Self>,
257        segment_id: &str,
258        memory_budget: usize,
259        threshold: Option<f64>,
260    ) -> Result<bool> {
261        if memory_budget < 1024 * 1024 {
262            return Err(Error::Schema(
263                "compaction memory budget must be at least 1 MiB".into(),
264            ));
265        }
266        let sid = SegmentId::from_hex(segment_id)
267            .ok_or_else(|| Error::Document("invalid compaction segment ID".into()))?;
268        if !self.active_operations.is_accepting() {
269            return Err(Error::IndexClosed);
270        }
271        if self.quarantined_segments.lock().contains(segment_id) {
272            return Err(Error::Corruption("compaction source is quarantined".into()));
273        }
274        if threshold.is_some()
275            && (self.force_merge_active.load(Ordering::Acquire) > 0
276                || self.paused_reorder_segments().contains(segment_id))
277        {
278            return Ok(false);
279        }
280        let (global_permit, permit, reorder_permit) = if threshold.is_some() {
281            let Ok(global) = Arc::clone(&self.global_merge_permits).try_acquire_owned() else {
282                return Ok(false);
283            };
284            let Ok(local) = Arc::clone(&self.merge_permits).try_acquire_owned() else {
285                return Ok(false);
286            };
287            let Ok(reorder) = self.reorder_permits.try_acquire_optimizer() else {
288                return Ok(false);
289            };
290            (global, local, reorder)
291        } else {
292            let (global, local) = self.acquire_maintenance_capacity().await?;
293            let reorder = tokio::select! {
294                biased;
295                () = self.active_operations.wait_for_shutdown() => return Err(Error::IndexClosed),
296                permit = self.reorder_permits.acquire(ReorderPriority::Optimizer) => permit.map_err(|_| Error::IndexClosed)?,
297            };
298            (global, local, reorder)
299        };
300        let manager = Arc::clone(self);
301        let source_id = segment_id.to_owned();
302        self.run_lifecycle_transaction(async move {
303            let _permit = permit;
304            let _global_permit = global_permit;
305            let _reorder_permit = reorder_permit;
306            let start = std::time::Instant::now();
307            let output = SegmentId::new();
308            let (claim, snapshot, visibility, schema) = {
309                let st = manager.state.lock().await;
310                let Some(info) = st.metadata.segment_metas.get(&source_id) else {
311                    return Ok(false);
312                };
313                if let Some(threshold) = threshold
314                    && (manager.force_merge_active.load(Ordering::Acquire) > 0 || info.deleted_ratio() < threshold) {
315                    return Ok(false);
316                }
317                let Some(visibility) = info.deletions.clone() else {
318                    return Ok(false);
319                };
320                let Some(claim) = manager.active_operations.try_register(vec![source_id.clone(), output.to_hex()]) else {
321                    return if threshold.is_some() { Ok(false) } else {
322                        Err(Error::Internal("compaction source is busy; retry after maintenance finishes".into()))
323                    };
324                };
325                let acquired = manager
326                    .tracker
327                    .acquire(&[source_id.clone(), visibility.id.clone()]);
328                let snapshot = SegmentSnapshot::with_delete_fn(
329                    Arc::clone(&manager.tracker),
330                    acquired,
331                    Arc::clone(&manager.delete_fn),
332                );
333                (
334                    claim,
335                    snapshot,
336                    visibility,
337                    manager.published_generation().schema.clone(),
338                )
339            };
340            let _claim = claim;
341            let _snapshot = snapshot;
342            let mut cleanup = manager.output_cleanup_guard(output);
343            let mut reader = SegmentReader::open_with_term_cache_budget(
344                manager.directory.as_ref(),
345                sid,
346                schema.clone(),
347                manager.term_cache_blocks,
348                manager.term_cache_budget_bytes,
349            )
350            .await?;
351            reader
352                .load_deletions(manager.directory.as_ref(), visibility.clone())
353                .await?;
354            let physical = reader.num_docs();
355            let meta = manager
356                .build_compacted(reader, output, memory_budget)
357                .await?;
358            let removed = visibility.num_deleted;
359            if meta.num_docs.checked_add(removed) != Some(physical) {
360                return Err(Error::Corruption("compaction row count mismatch".into()));
361            }
362            let expected = HashMap::from([(source_id.clone(), Some(visibility))]);
363            manager
364                .replace_segments(
365                    std::slice::from_ref(&source_id),
366                    output.to_hex(),
367                    meta.num_docs,
368                    ReplacementLayout::Compacted,
369                    Some(&expected),
370                )
371                .await?;
372            cleanup.disarm();
373            log::info!("[compaction] index={} source={} physical_rows={} removed_rows={} removed_ratio={:.6} live_rows={} elapsed_secs={:.3}",
374                manager.schema.index_label(), source_id, physical, removed,
375                f64::from(removed) / f64::from(physical), meta.num_docs, start.elapsed().as_secs_f64());
376            Ok(true)
377        })
378        .await
379    }
380
381    pub(super) async fn build_compacted(
382        &self,
383        reader: SegmentReader,
384        output: SegmentId,
385        memory_budget: usize,
386    ) -> Result<crate::segment::SegmentMeta> {
387        let directory = Arc::clone(&self.directory);
388        let merger = crate::segment::SegmentMerger::new(self.published_generation().schema.clone())
389            .with_posting_config(self.optimization, self.posting_codec)
390            .with_term_dict_block_size(self.term_dict_block_size)
391            .with_cancellation(self.active_operations.cancellation_flag())
392            .with_background_pool(Some(self.background_cpu_pool()));
393        let runtime = tokio::runtime::Handle::current();
394        let (meta, _) = tokio::task::spawn_blocking(move || {
395            runtime.block_on(merger.compact(directory.as_ref(), &reader, output, memory_budget))
396        })
397        .await
398        .map_err(|error| Error::Internal(format!("compaction worker failed: {error}")))??;
399        Ok(meta)
400    }
401}
402
403fn validate_threshold(threshold: f64) -> Result<()> {
404    if !threshold.is_finite() || !(0.0..=1.0).contains(&threshold) {
405        return Err(Error::Schema("compaction deleted ratio must be finite and in 0..=1 (0 disables automatic compaction)".into()));
406    }
407    Ok(())
408}