use std::sync::atomic::AtomicBool;
use ahash::{AHashMap, AHashSet};
use crate::common::counter::hardware_counter::HardwareCounterCell;
use crate::common::types::DeferredBehavior;
use crate::segment::common::operation_error::OperationResult;
use crate::segment::types::{Filter, PointIdType, SeqNumberType};
use crate::shard::segment_holder::SegmentHolder;
use crate::shard::update::helpers::deferred_points_to_exclude_by_filter;
const DELETION_BATCH_SIZE: usize = 512;
pub fn delete_points(
segments: &SegmentHolder,
op_num: SeqNumberType,
ids: &[PointIdType],
hw_counter: &HardwareCounterCell,
) -> OperationResult<usize> {
let mut total_deleted_points = 0;
for batch in ids.chunks(DELETION_BATCH_SIZE) {
for (_segment_id, segment) in segments.iter() {
let segment_arc = segment.get();
let mut write_segment = segment_arc.write();
for &id in batch {
if write_segment.delete_point(op_num, id, hw_counter)? {
total_deleted_points += 1;
}
}
}
}
if total_deleted_points == 0 {
segments.bump_max_segment_version_overwrite(op_num);
}
Ok(total_deleted_points)
}
pub fn delete_points_by_filter(
segments: &SegmentHolder,
op_num: SeqNumberType,
filter: &Filter,
hw_counter: &HardwareCounterCell,
) -> OperationResult<usize> {
let mut total_deleted = 0;
let is_stopped = AtomicBool::new(false);
let mut has_deferred = false;
let mut points_to_delete: AHashMap<_, _> = segments
.iter()
.map(|(segment_id, segment)| {
let segment = segment.get().read();
let point_ids = segment.read_filtered(
None,
None,
Some(filter),
&is_stopped,
hw_counter,
DeferredBehavior::WithDeferred,
)?;
has_deferred |= segment.has_deferred_points();
Ok((segment_id, point_ids))
})
.collect::<OperationResult<_>>()?;
if has_deferred {
let points_to_keep = deferred_points_to_exclude_by_filter(segments, &points_to_delete);
let all_matched_points: AHashSet<PointIdType> = points_to_delete
.values()
.flat_map(|v| v.iter().copied())
.collect();
for (segment_id, segment) in segments.iter() {
let segment = segment.get().read();
let present: Vec<_> = all_matched_points
.iter()
.copied()
.filter(|point_id| {
segment.has_point(*point_id, DeferredBehavior::WithDeferred)
&& !points_to_keep.contains(point_id)
})
.collect();
points_to_delete.insert(segment_id, present);
}
}
segments.apply_segments_batched(|s, segment_id| {
let Some(curr_points) = points_to_delete.get_mut(&segment_id) else {
return Ok(false);
};
if curr_points.is_empty() {
return Ok(false);
}
let mut deleted_in_batch = 0;
while let Some(point_id) = curr_points.pop() {
if s.delete_point(op_num, point_id, hw_counter)? {
total_deleted += 1;
deleted_in_batch += 1;
}
if deleted_in_batch >= DELETION_BATCH_SIZE {
break;
}
}
Ok(true)
})?;
if total_deleted == 0 {
segments.bump_max_segment_version_overwrite(op_num);
}
Ok(total_deleted)
}