1use 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 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 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 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 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 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 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 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}