use ahash::AHashMap;
use crate::common::counter::hardware_counter::HardwareCounterCell;
use crate::common::types::PointOffsetType;
use crate::common::universal_io::UniversalRead;
use rayon::ThreadPool;
use rayon::prelude::*;
use crate::segment::common::operation_error::{OperationError, OperationResult};
use crate::segment::data_types::fully_qualified_point::{FullyQualifiedPoint, StoredPoint};
use crate::segment::types::{PointIdType, SeqNumberType};
use crate::shard::operations::CollectionUpdateOperations;
use uuid::Uuid;
use crate::edge::update_only::UpdateOnlyEdgeShard;
use crate::edge::update_only::batch::UpdateBatchPlan;
use crate::edge::update_only::holder::UpdateOnlySegmentHolder;
use crate::edge::update_only::preview::{PointAction, PointPreview, resolve_batch};
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct UpdateBatchOutcome {
pub stored: usize,
pub deleted: usize,
pub skipped: usize,
pub missing: usize,
}
#[derive(Debug, Clone, Copy)]
pub(super) struct PointLocation {
pub(super) segment: Uuid,
pub(super) internal_id: PointOffsetType,
pub(super) version: SeqNumberType,
appendable: bool,
}
impl PointLocation {
fn supersedes(&self, other: &Self) -> bool {
(self.version, self.appendable) > (other.version, other.appendable)
}
}
pub(super) struct PointLocations {
pub(super) newest: PointLocation,
pub(super) slots: Vec<(Uuid, PointOffsetType)>,
}
impl<S: UniversalRead + 'static> UpdateOnlyEdgeShard<S> {
pub fn apply_batch(
&self,
operations: impl IntoIterator<Item = (SeqNumberType, CollectionUpdateOperations)>,
) -> OperationResult<UpdateBatchOutcome> {
let plan = UpdateBatchPlan::build(operations)?;
if plan.is_empty() {
return Ok(UpdateBatchOutcome::default());
}
let hw_counter = HardwareCounterCell::disposable();
let segments = self.segments.read();
let resolved = resolve_batch(&segments, plan, &self.pool)?;
let mut outcome = UpdateBatchOutcome::default();
let mut to_store: Vec<FullyQualifiedPoint> = Vec::new();
let mut to_tombstone: AHashMap<Uuid, Vec<PointOffsetType>> = AHashMap::new();
for point in resolved {
let PointPreview {
id: _,
current: _,
slots,
action,
} = point;
match action {
PointAction::Skip => {
outcome.skipped += 1;
continue;
}
PointAction::Missing => {
outcome.missing += 1;
continue;
}
PointAction::Store(point) => {
to_store.push(*point);
outcome.stored += 1;
}
PointAction::Delete => outcome.deleted += 1,
}
for (segment, internal_id) in slots {
to_tombstone.entry(segment).or_default().push(internal_id);
}
}
if !to_store.is_empty() {
let write_target = segments.write_target()?;
write_target.write().store_points(&to_store, &hw_counter)?;
write_target.read().flush()?;
}
for (uuid, internal_ids) in to_tombstone {
let segment = segments.get(&uuid).ok_or_else(|| {
OperationError::service_error(format!("Segment {uuid} disappeared mid-batch"))
})?;
segment.write().tombstone_points(&internal_ids)?;
segment.read().flush()?;
}
Ok(outcome)
}
}
pub(super) fn locate_points<S: UniversalRead + 'static>(
segments: &UpdateOnlySegmentHolder<S>,
plan: &UpdateBatchPlan,
pool: &ThreadPool,
) -> OperationResult<AHashMap<PointIdType, PointLocations>> {
let ids: Vec<PointIdType> = plan.point_ids().collect();
let per_segment: Vec<Vec<(PointIdType, PointLocation)>> = pool.install(|| {
segments
.iter()
.collect::<Vec<_>>()
.into_par_iter()
.map(|(uuid, segment)| {
let segment = segment.read();
let appendable = segment.is_appendable();
let mut found_ids = Vec::new();
let mut internal_ids = Vec::new();
segment.locate_points(ids.iter().copied(), |id, internal_id| {
found_ids.push(id);
internal_ids.push(internal_id);
})?;
let versions = segment.point_versions(&internal_ids)?;
let located = found_ids
.into_iter()
.zip(internal_ids)
.map(|(id, internal_id)| {
let location = PointLocation {
segment: uuid,
internal_id,
version: versions.get(&internal_id).copied().unwrap_or(0),
appendable,
};
(id, location)
})
.collect();
Ok(located)
})
.collect::<OperationResult<Vec<_>>>()
})?;
let mut locations: AHashMap<PointIdType, PointLocations> = AHashMap::new();
for (id, location) in per_segment.into_iter().flatten() {
let slot = (location.segment, location.internal_id);
locations
.entry(id)
.and_modify(|current| {
current.slots.push(slot);
if location.supersedes(¤t.newest) {
current.newest = location;
}
})
.or_insert_with(|| PointLocations {
newest: location,
slots: vec![slot],
});
}
Ok(locations)
}
pub(super) fn read_stored_points<S: UniversalRead + 'static>(
segments: &UpdateOnlySegmentHolder<S>,
plan: &UpdateBatchPlan,
locations: &AHashMap<PointIdType, PointLocations>,
pool: &ThreadPool,
) -> OperationResult<AHashMap<PointIdType, StoredPoint>> {
let mut by_segment: AHashMap<Uuid, Vec<(PointIdType, PointOffsetType)>> = AHashMap::new();
for id in plan.point_ids_needing_stored_point() {
if let Some(location) = locations.get(&id) {
by_segment
.entry(location.newest.segment)
.or_default()
.push((id, location.newest.internal_id));
}
}
let per_segment: Vec<Vec<(PointIdType, StoredPoint)>> = pool.install(|| {
by_segment
.into_iter()
.collect::<Vec<_>>()
.into_par_iter()
.map(|(uuid, entries)| {
let segment = segments.get(&uuid).ok_or_else(|| {
OperationError::service_error(format!("Segment {uuid} disappeared mid-batch"))
})?;
let segment = segment.read();
let internal_ids: Vec<PointOffsetType> = entries
.iter()
.map(|(_, internal_id)| *internal_id)
.collect();
let hw_counter = HardwareCounterCell::disposable();
let points = segment.read_stored_points(&internal_ids, &hw_counter)?;
Ok(entries.into_iter().map(|(id, _)| id).zip(points).collect())
})
.collect::<OperationResult<Vec<_>>>()
})?;
let mut stored = AHashMap::new();
for (id, point) in per_segment.into_iter().flatten() {
stored.insert(id, point);
}
Ok(stored)
}