use std::cmp;
use std::collections::{BTreeSet, HashMap};
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
use ahash::AHashMap;
use crate::common::counter::hardware_counter::HardwareCounterCell;
use crate::common::types::{DeferredBehavior, TelemetryDetail};
use crate::segment::common::Flusher;
use crate::segment::common::operation_error::{OperationError, OperationResult, SegmentFailedState};
use crate::segment::data_types::build_index_result::BuildFieldIndexResult;
use crate::segment::data_types::facets::{FacetParams, FacetValue};
use crate::segment::data_types::named_vectors::NamedVectors;
use crate::segment::data_types::order_by::OrderValue;
use crate::segment::data_types::query_context::{FormulaContext, QueryContext, SegmentQueryContext};
use crate::segment::data_types::segment_record::{SegmentRecord, SegmentRecordRaw};
use crate::segment::data_types::vector_name_config::VectorNameConfig;
use crate::segment::data_types::vectors::{QueryVector, VectorInternal};
use crate::segment::entry::StorageSegmentEntry;
use crate::segment::entry::entry_point::{NonAppendableSegmentEntry, ReadSegmentEntry, SegmentEntry};
use crate::segment::index::field_index::{CardinalityEstimation, FieldIndex};
use crate::segment::json_path::JsonPath;
use crate::segment::telemetry::SegmentTelemetry;
use crate::segment::types::*;
use uuid::Uuid;
use super::{ProxyDeletedPoint, ProxyIndexChange, ProxySegment};
use crate::shard::locked_segment::LockedSegment;
impl ProxySegment {
fn redact_and_filter_for_retrieve<'a>(
&self,
with_vector: &'a WithVector,
point_ids: &[PointIdType],
) -> (std::borrow::Cow<'a, WithVector>, Vec<PointIdType>) {
let with_vector = self
.changed_vector_names
.redact_with_vector(with_vector, &self.wrapped_config);
let filtered_point_ids = point_ids
.iter()
.copied()
.filter(|id| !self.deleted_points.contains_key(id))
.collect();
(with_vector, filtered_point_ids)
}
fn adjusted_info(&self, wrapped_info: SegmentInfo) -> SegmentInfo {
let mut vector_data = wrapped_info.vector_data;
let mut removed_num_vectors = 0usize;
let mut removed_num_indexed = 0usize;
let mut removed_num_deleted = 0usize;
vector_data.retain(|name, info| {
if self.changed_vector_names.is_wrapped_data_stale(name) {
removed_num_vectors += info.num_vectors;
removed_num_indexed += info.num_indexed_vectors;
removed_num_deleted += info.num_deleted_vectors;
false
} else {
true
}
});
let vector_name_count = vector_data.len();
let deleted_points_count = self.deleted_points.len();
let num_vectors = wrapped_info
.num_vectors
.saturating_sub(removed_num_vectors)
.saturating_sub(deleted_points_count * vector_name_count);
let num_indexed_vectors = if wrapped_info.segment_type == SegmentType::Indexed {
wrapped_info
.num_indexed_vectors
.saturating_sub(removed_num_indexed)
.saturating_sub(deleted_points_count * vector_name_count)
} else {
0
};
let num_deleted_vectors = wrapped_info
.num_deleted_vectors
.saturating_sub(removed_num_deleted)
+ deleted_points_count * vector_name_count;
SegmentInfo {
uuid: wrapped_info.uuid,
segment_type: SegmentType::Special,
num_vectors,
num_indexed_vectors,
num_points: self.available_point_count(),
num_deferred_points: Some(self.deferred_point_count()),
num_deleted_deferred_points: wrapped_info.num_deleted_deferred_points.map(
|num_deleted_deferred_points| {
num_deleted_deferred_points.saturating_add(self.deleted_deferred_count)
},
),
num_deleted_vectors,
vectors_size_bytes: wrapped_info.vectors_size_bytes,
payloads_size_bytes: wrapped_info.payloads_size_bytes,
ram_usage_bytes: wrapped_info.ram_usage_bytes,
disk_usage_bytes: wrapped_info.disk_usage_bytes,
is_appendable: false,
index_schema: wrapped_info.index_schema,
vector_data,
payload_storage_io_backend: wrapped_info.payload_storage_io_backend,
deferred_internal_id: wrapped_info.deferred_internal_id,
}
}
}
impl ReadSegmentEntry for ProxySegment {
fn is_proxy(&self) -> bool {
true
}
fn point_version(&self, point_id: PointIdType) -> Option<SeqNumberType> {
let wrapped_version = self.wrapped_segment.get().read().point_version(point_id)?;
if self
.deleted_points
.get(&point_id)
.is_some_and(|delete| wrapped_version <= delete.local_version)
{
return None;
}
Some(wrapped_version)
}
fn search_batch(
&self,
vector_name: &VectorName,
vectors: &[&QueryVector],
with_payload: &WithPayload,
with_vector: &WithVector,
filter: Option<&Filter>,
top: usize,
params: Option<&SearchParams>,
query_context: &SegmentQueryContext,
) -> OperationResult<Vec<Vec<ScoredPoint>>> {
if self.changed_vector_names.is_wrapped_data_stale(vector_name) {
return Ok(vec![Vec::new(); vectors.len()]);
}
let with_vector = self
.changed_vector_names
.redact_with_vector(with_vector, &self.wrapped_config);
let with_vector = with_vector.as_ref();
let filter = filter.map(|f| self.changed_vector_names.redact_filter(f));
let do_update_filter = !self.deleted_points.is_empty();
let wrapped_results = if do_update_filter {
if let Some(deleted_points) = self.deleted_mask.as_ref() {
let query_context_with_deleted =
query_context.fork().with_deleted_points(deleted_points);
self.wrapped_segment.get().read().search_batch(
vector_name,
vectors,
with_payload,
with_vector,
filter.as_deref(),
top,
params,
&query_context_with_deleted,
)?
} else {
let wrapped_filter = Self::add_deleted_points_condition_to_filter(
filter,
self.deleted_points.keys().copied(),
);
self.wrapped_segment.get().read().search_batch(
vector_name,
vectors,
with_payload,
with_vector,
Some(&wrapped_filter),
top,
params,
query_context,
)?
}
} else {
self.wrapped_segment.get().read().search_batch(
vector_name,
vectors,
with_payload,
with_vector,
filter.as_deref(),
top,
params,
query_context,
)?
};
Ok(wrapped_results)
}
fn rescore_with_formula(
&self,
formula_ctx: Arc<FormulaContext>,
hw_counter: &HardwareCounterCell,
) -> OperationResult<Vec<ScoredPoint>> {
let wrapped_results = self
.wrapped_segment
.get()
.read()
.rescore_with_formula(formula_ctx, hw_counter)?;
let result = {
if self.deleted_points.is_empty() {
wrapped_results
} else {
wrapped_results
.into_iter()
.filter(|point| !self.deleted_points.contains_key(&point.id))
.collect()
}
};
Ok(result)
}
fn vector(
&self,
vector_name: &VectorName,
point_id: PointIdType,
hw_counter: &HardwareCounterCell,
) -> OperationResult<Option<VectorInternal>> {
self.vector_with_behavior(
vector_name,
point_id,
DeferredBehavior::VisibleOnly,
hw_counter,
)
}
fn vector_with_behavior(
&self,
vector_name: &VectorName,
point_id: PointIdType,
deferred_behavior: DeferredBehavior,
hw_counter: &HardwareCounterCell,
) -> OperationResult<Option<VectorInternal>> {
if self.changed_vector_names.is_wrapped_data_stale(vector_name) {
return Ok(None);
}
if self.deleted_points.contains_key(&point_id) {
Ok(None)
} else {
self.wrapped_segment.get().read().vector_with_behavior(
vector_name,
point_id,
deferred_behavior,
hw_counter,
)
}
}
fn all_vectors(
&self,
point_id: PointIdType,
hw_counter: &HardwareCounterCell,
) -> OperationResult<NamedVectors<'_>> {
let mut result = NamedVectors::default();
let wrapped = self.wrapped_segment.get();
let wrapped_guard = wrapped.read();
let config = wrapped_guard.config();
let vector_names: Vec<_> = config
.vector_data
.keys()
.chain(config.sparse_vector_data.keys())
.cloned()
.collect();
drop(wrapped_guard);
for vector_name in vector_names {
if let Some(vector) = self.vector_with_behavior(
&vector_name,
point_id,
DeferredBehavior::VisibleOnly,
hw_counter,
)? {
result.insert(vector_name, vector);
}
}
Ok(result)
}
fn payload(
&self,
point_id: PointIdType,
hw_counter: &HardwareCounterCell,
) -> OperationResult<Payload> {
if self.deleted_points.contains_key(&point_id) {
Ok(Payload::default())
} else {
self.wrapped_segment
.get()
.read()
.payload(point_id, hw_counter)
}
}
fn retrieve(
&self,
point_ids: &[PointIdType],
with_payload: &WithPayload,
with_vector: &WithVector,
hw_counter: &HardwareCounterCell,
is_stopped: &AtomicBool,
deferred_behavior: DeferredBehavior,
) -> OperationResult<AHashMap<ExtendedPointId, SegmentRecord>> {
let (with_vector, filtered_point_ids) =
self.redact_and_filter_for_retrieve(with_vector, point_ids);
self.wrapped_segment.get().read().retrieve(
&filtered_point_ids,
with_payload,
with_vector.as_ref(),
hw_counter,
is_stopped,
deferred_behavior,
)
}
fn retrieve_raw(
&self,
point_ids: &[PointIdType],
with_payload: &WithPayload,
with_vector: &WithVector,
hw_counter: &HardwareCounterCell,
is_stopped: &AtomicBool,
deferred_behavior: DeferredBehavior,
) -> OperationResult<AHashMap<ExtendedPointId, SegmentRecordRaw>> {
let (with_vector, filtered_point_ids) =
self.redact_and_filter_for_retrieve(with_vector, point_ids);
self.wrapped_segment.get().read().retrieve_raw(
&filtered_point_ids,
with_payload,
with_vector.as_ref(),
hw_counter,
is_stopped,
deferred_behavior,
)
}
fn read_filtered<'a>(
&'a self,
offset: Option<PointIdType>,
limit: Option<usize>,
filter: Option<&'a Filter>,
is_stopped: &AtomicBool,
hw_counter: &HardwareCounterCell,
deferred_behavior: DeferredBehavior,
) -> OperationResult<Vec<PointIdType>> {
let filter = filter.map(|f| self.changed_vector_names.redact_filter(f));
if self.deleted_points.is_empty() {
self.wrapped_segment.get().read().read_filtered(
offset,
limit,
filter.as_deref(),
is_stopped,
hw_counter,
deferred_behavior,
)
} else {
let wrapped_filter = Self::add_deleted_points_condition_to_filter(
filter,
self.deleted_points.keys().copied(),
);
self.wrapped_segment.get().read().read_filtered(
offset,
limit,
Some(&wrapped_filter),
is_stopped,
hw_counter,
deferred_behavior,
)
}
}
fn read_ordered_filtered<'a>(
&'a self,
limit: Option<usize>,
filter: Option<&'a Filter>,
order_by: &'a crate::segment::data_types::order_by::OrderBy,
is_stopped: &AtomicBool,
hw_counter: &HardwareCounterCell,
deferred_behavior: DeferredBehavior,
) -> OperationResult<Vec<(OrderValue, PointIdType)>> {
let filter = filter.map(|f| self.changed_vector_names.redact_filter(f));
let read_points = if self.deleted_points.is_empty() {
self.wrapped_segment.get().read().read_ordered_filtered(
limit,
filter.as_deref(),
order_by,
is_stopped,
hw_counter,
deferred_behavior,
)?
} else {
let wrapped_filter = Self::add_deleted_points_condition_to_filter(
filter,
self.deleted_points.keys().copied(),
);
self.wrapped_segment.get().read().read_ordered_filtered(
limit,
Some(&wrapped_filter),
order_by,
is_stopped,
hw_counter,
deferred_behavior,
)?
};
Ok(read_points)
}
fn read_random_filtered<'a>(
&'a self,
limit: usize,
filter: Option<&'a Filter>,
is_stopped: &AtomicBool,
hw_counter: &HardwareCounterCell,
) -> OperationResult<Vec<PointIdType>> {
let filter = filter.map(|f| self.changed_vector_names.redact_filter(f));
if self.deleted_points.is_empty() {
self.wrapped_segment.get().read().read_random_filtered(
limit,
filter.as_deref(),
is_stopped,
hw_counter,
)
} else {
let wrapped_filter = Self::add_deleted_points_condition_to_filter(
filter,
self.deleted_points.keys().copied(),
);
self.wrapped_segment.get().read().read_random_filtered(
limit,
Some(&wrapped_filter),
is_stopped,
hw_counter,
)
}
}
fn read_range(&self, from: Option<PointIdType>, to: Option<PointIdType>) -> Vec<PointIdType> {
let read_points = self.wrapped_segment.get().read().read_range(from, to);
if self.deleted_points.is_empty() {
read_points
} else {
read_points
.into_iter()
.filter(|idx| !self.deleted_points.contains_key(idx))
.collect()
}
}
fn unique_values(
&self,
key: &JsonPath,
filter: Option<&Filter>,
is_stopped: &AtomicBool,
hw_counter: &HardwareCounterCell,
) -> OperationResult<BTreeSet<FacetValue>> {
let filter = filter.map(|f| self.changed_vector_names.redact_filter(f));
let values = self.wrapped_segment.get().read().unique_values(
key,
filter.as_deref(),
is_stopped,
hw_counter,
)?;
Ok(values)
}
fn facet(
&self,
request: &FacetParams,
is_stopped: &AtomicBool,
hw_counter: &HardwareCounterCell,
) -> OperationResult<HashMap<FacetValue, usize>> {
let filter = request
.filter
.as_ref()
.map(|f| self.changed_vector_names.redact_filter(f));
let hits = if self.deleted_points.is_empty() {
match filter {
None | Some(std::borrow::Cow::Borrowed(_)) => self
.wrapped_segment
.get()
.read()
.facet(request, is_stopped, hw_counter)?,
Some(std::borrow::Cow::Owned(f)) => {
let new_request = FacetParams {
filter: Some(f),
..request.clone()
};
self.wrapped_segment
.get()
.read()
.facet(&new_request, is_stopped, hw_counter)?
}
}
} else {
let wrapped_filter = Self::add_deleted_points_condition_to_filter(
filter,
self.deleted_points.keys().copied(),
);
let new_request = FacetParams {
filter: Some(wrapped_filter),
..request.clone()
};
self.wrapped_segment
.get()
.read()
.facet(&new_request, is_stopped, hw_counter)?
};
Ok(hits)
}
fn has_point(&self, point_id: PointIdType, deferred_behavior: DeferredBehavior) -> bool {
!self.deleted_points.contains_key(&point_id)
&& self
.wrapped_segment
.get()
.read()
.has_point(point_id, deferred_behavior)
}
fn is_empty(&self) -> bool {
self.wrapped_segment.get().read().is_empty()
}
fn available_point_count(&self) -> usize {
let deleted_points_count = self.deleted_points.len();
let wrapped_segment_count = self.wrapped_segment.get().read().available_point_count();
wrapped_segment_count.saturating_sub(deleted_points_count)
}
fn available_point_count_without_deferred(&self) -> usize {
let wrapped_segment_visible_count = self
.wrapped_segment
.get()
.read()
.available_point_count_without_deferred();
let deleted_visible = self
.deleted_points
.len()
.saturating_sub(self.deleted_deferred_count);
wrapped_segment_visible_count.saturating_sub(deleted_visible)
}
fn deleted_point_count(&self) -> usize {
self.wrapped_segment.get().read().deleted_point_count() + self.deleted_points.len()
}
fn available_vectors_size_in_bytes(&self, vector_name: &VectorName) -> OperationResult<usize> {
if self.changed_vector_names.is_wrapped_data_stale(vector_name) {
return Ok(0);
}
let wrapped_segment = self.wrapped_segment.get();
let wrapped_segment_guard = wrapped_segment.read();
let wrapped_size = wrapped_segment_guard.available_vectors_size_in_bytes(vector_name)?;
let wrapped_count = wrapped_segment_guard.available_point_count();
drop(wrapped_segment_guard);
let stored_points = wrapped_count;
if stored_points > 0 {
let deleted_points_count = self.deleted_points.len();
let available_points = stored_points.saturating_sub(deleted_points_count);
Ok(
((wrapped_size as u128) * available_points as u128 / stored_points as u128)
as usize,
)
} else {
Ok(0)
}
}
fn estimate_point_count<'a>(
&'a self,
filter: Option<&'a Filter>,
hw_counter: &HardwareCounterCell,
) -> OperationResult<CardinalityEstimation> {
let filter = filter.map(|f| self.changed_vector_names.redact_filter(f));
let deleted_point_count = self
.deleted_points
.len()
.saturating_sub(self.deleted_deferred_count);
let (wrapped_segment_est, total_wrapped_size) = {
let wrapped_segment = self.wrapped_segment.get();
let wrapped_segment_guard = wrapped_segment.read();
(
wrapped_segment_guard.estimate_point_count(filter.as_deref(), hw_counter)?,
wrapped_segment_guard.available_point_count_without_deferred(),
)
};
let expected_deleted_count = if total_wrapped_size > 0 {
(wrapped_segment_est.exp as f64
* (deleted_point_count as f64 / total_wrapped_size as f64)) as usize
} else {
0
};
let CardinalityEstimation {
primary_clauses,
min,
exp,
max,
} = wrapped_segment_est;
Ok(CardinalityEstimation {
primary_clauses,
min: min.saturating_sub(deleted_point_count),
exp: exp.saturating_sub(expected_deleted_count),
max,
})
}
fn segment_uuid(&self) -> Uuid {
self.wrapped_segment.get().read().segment_uuid()
}
fn segment_type(&self) -> SegmentType {
SegmentType::Special
}
fn size_info(&self) -> SegmentInfo {
self.adjusted_info(self.wrapped_segment.get().read().size_info())
}
fn info(&self) -> OperationResult<SegmentInfo> {
let wrapped_info = self.wrapped_segment.get().read().info()?;
Ok(self.adjusted_info(wrapped_info))
}
fn config(&self) -> &SegmentConfig {
&self.wrapped_config
}
fn is_appendable(&self) -> bool {
false
}
fn get_indexed_fields(&self) -> HashMap<PayloadKeyType, PayloadFieldSchema> {
let mut indexed_fields = self.wrapped_segment.get().read().get_indexed_fields();
for (field_name, change) in self.changed_indexes.iter_unordered() {
match change {
ProxyIndexChange::Create(schema, _) => {
indexed_fields.insert(field_name.to_owned(), schema.to_owned());
}
ProxyIndexChange::Delete(_) => {
indexed_fields.remove(field_name);
}
ProxyIndexChange::DeleteIfIncompatible(_, schema) => {
if let Some(existing_schema) = indexed_fields.get(field_name)
&& existing_schema != schema
{
indexed_fields.remove(field_name);
}
}
}
}
indexed_fields
}
fn vector_names(&self) -> Vec<VectorNameBuf> {
self.wrapped_segment.get().read().vector_names()
}
fn get_telemetry_data(&self, detail: TelemetryDetail) -> OperationResult<SegmentTelemetry> {
self.wrapped_segment.get().read().get_telemetry_data(detail)
}
fn fill_query_context(&self, query_context: &mut QueryContext) -> OperationResult<()> {
self.wrapped_segment
.get()
.read()
.fill_query_context(query_context)
}
fn point_is_deferred(&self, point_id: PointIdType) -> bool {
!self.deleted_points.contains_key(&point_id)
&& self
.wrapped_segment
.get()
.read()
.point_is_deferred(point_id)
}
fn deferred_point_ids(&self) -> Vec<PointIdType> {
let mut ids = self.wrapped_segment.get().read().deferred_point_ids();
if self.deleted_deferred_count > 0 {
ids.retain(|point_id| !self.deleted_points.contains_key(point_id));
}
ids
}
fn has_deferred_points(&self) -> bool {
self.wrapped_segment.get().read().has_deferred_points()
}
fn deferred_point_count(&self) -> usize {
self.wrapped_segment
.get()
.read()
.deferred_point_count()
.saturating_sub(self.deleted_deferred_count)
}
}
impl StorageSegmentEntry for ProxySegment {
fn version(&self) -> SeqNumberType {
cmp::max(self.wrapped_segment.get().read().version(), self.version)
}
fn check_error(&self) -> Option<SegmentFailedState> {
self.wrapped_segment.get().read().check_error()
}
fn persistent_version(&self) -> SeqNumberType {
self.wrapped_segment.get().read().persistent_version()
}
fn flusher(&self, force: bool) -> Option<Flusher> {
let wrapped_segment = self.wrapped_segment.get();
let wrapped_segment_guard = wrapped_segment.read();
wrapped_segment_guard.flusher(force)
}
fn drop_data(self) -> OperationResult<()> {
self.wrapped_segment.drop_data()
}
fn data_path(&self) -> PathBuf {
self.wrapped_segment.get().read().data_path()
}
}
impl SegmentEntry for ProxySegment {
fn upsert_point(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_vectors: NamedVectors,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Upsert is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn upsert_point_raw(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_vectors: &[(VectorNameBuf, Vec<u8>)],
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Upsert is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn upsert_moved_point(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_raw_vectors: &[(VectorNameBuf, Vec<u8>)],
_updated_vectors: NamedVectors,
_payload: &Payload,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Upsert is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn update_vectors(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_vectors: NamedVectors,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Update vectors is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn delete_vector(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_vector_name: &VectorName,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Delete vector is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn set_full_payload(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_full_payload: &Payload,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Set full payload is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn set_payload(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_payload: &Payload,
_key: &Option<JsonPath>,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Set payload is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn delete_payload(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_key: PayloadKeyTypeRef,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Delete payload is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
fn clear_payload(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
Err(OperationError::service_error(format!(
"Clear payload is disabled for proxy segments: operation {op_num} on point {point_id}",
)))
}
}
impl NonAppendableSegmentEntry for ProxySegment {
fn delete_point(
&mut self,
op_num: SeqNumberType,
point_id: PointIdType,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<bool> {
let mut was_deleted = false;
let was_deferred_point;
self.version = cmp::max(self.version, op_num);
let point_offset = match &self.wrapped_segment {
LockedSegment::Original(raw_segment) => {
let (point_offset, is_deferred) = {
let read_segment = raw_segment.read();
(
read_segment.get_internal_id(point_id),
read_segment.point_is_deferred(point_id),
)
};
was_deferred_point = is_deferred;
if point_offset.is_some() {
let prev = self.deleted_points.insert(
point_id,
ProxyDeletedPoint {
local_version: op_num,
operation_version: op_num,
},
);
was_deleted = prev.is_none();
if let Some(prev) = prev {
debug_assert!(
prev.operation_version < op_num,
"Overriding deleted flag {prev:?} with older op_num:{op_num}",
)
}
}
point_offset
}
LockedSegment::Proxy(proxy) => {
let (has_point, is_deferred) = {
let read_proxy = proxy.read();
(
read_proxy.has_point(point_id, DeferredBehavior::WithDeferred),
read_proxy.point_is_deferred(point_id),
)
};
was_deferred_point = is_deferred;
if has_point {
let prev = self.deleted_points.insert(
point_id,
ProxyDeletedPoint {
local_version: op_num,
operation_version: op_num,
},
);
was_deleted = prev.is_none();
if let Some(prev) = prev {
debug_assert!(
prev.operation_version < op_num,
"Overriding deleted flag {prev:?} with older op_num:{op_num}",
)
}
}
None
}
};
self.set_deleted_offset(point_offset);
if was_deleted && was_deferred_point {
self.deleted_deferred_count += 1;
}
Ok(was_deleted)
}
fn delete_field_index(&mut self, op_num: u64, key: PayloadKeyTypeRef) -> OperationResult<bool> {
if self.version() > op_num {
return Ok(false);
}
self.version = cmp::max(self.version, op_num);
self.changed_indexes
.insert(key.clone(), ProxyIndexChange::Delete(op_num));
Ok(true)
}
fn delete_field_index_if_incompatible(
&mut self,
op_num: SeqNumberType,
key: PayloadKeyTypeRef,
field_schema: &PayloadFieldSchema,
) -> OperationResult<bool> {
if self.version() > op_num {
return Ok(false);
}
self.version = cmp::max(self.version, op_num);
self.changed_indexes.insert(
key.clone(),
ProxyIndexChange::DeleteIfIncompatible(op_num, field_schema.clone()),
);
Ok(true)
}
fn build_field_index(
&self,
op_num: SeqNumberType,
_key: PayloadKeyTypeRef,
field_type: &PayloadFieldSchema,
_hw_counter: &HardwareCounterCell,
) -> OperationResult<BuildFieldIndexResult> {
if self.version() > op_num {
return Ok(BuildFieldIndexResult::SkippedByVersion);
}
Ok(BuildFieldIndexResult::Built {
indexes: vec![], schema: field_type.clone(),
})
}
fn apply_field_index(
&mut self,
op_num: SeqNumberType,
key: PayloadKeyType,
field_schema: PayloadFieldSchema,
_field_index: Vec<FieldIndex>,
) -> OperationResult<bool> {
if self.version() > op_num {
return Ok(false);
}
self.version = cmp::max(self.version, op_num);
self.changed_indexes
.insert(key, ProxyIndexChange::Create(field_schema, op_num));
Ok(true)
}
fn create_vector_name(
&mut self,
op_num: SeqNumberType,
vector_name: &VectorName,
vector_config: &VectorNameConfig,
) -> OperationResult<bool> {
if self.version() > op_num {
return Ok(false);
}
self.version = cmp::max(self.version, op_num);
self.changed_vector_names.record_create(
vector_name.to_owned(),
vector_config.clone(),
op_num,
&self.wrapped_config,
);
Ok(true)
}
fn delete_vector_name(
&mut self,
op_num: SeqNumberType,
vector_name: &VectorName,
) -> OperationResult<bool> {
if self.version() > op_num {
return Ok(false);
}
self.version = cmp::max(self.version, op_num);
self.changed_vector_names
.record_delete(vector_name.to_owned(), op_num);
Ok(true)
}
}