Files
qdrant/lib/edge/src/read_view.rs
xzfc db2c68e004 Replace DynConditionChecker with ConditionCheckerEnum (#9560)
* Add UniversalReadExt

Need this trait for the upcoming `ConditionCheckerEnum`.

* Replace DynConditionChecker with ConditionCheckerEnum
2026-06-26 12:29:01 +00:00

184 lines
6.8 KiB
Rust

use std::path::Path;
use std::sync::Arc;
use common::counter::hardware_accumulator::HwMeasurementAcc;
use common::types::ScoreType;
use parking_lot::{RwLock, RwLockReadGuard};
use segment::common::operation_error::OperationResult;
use segment::data_types::facets::FacetResponse;
use segment::entry::ReadSegmentEntry;
use segment::index::UniversalReadExt;
use segment::index::query_optimization::rescore_formula::parsed_formula::ParsedFormula;
use segment::segment::read_only::ReadOnlySegment;
use segment::types::{ExtendedPointId, PointIdType, ScoredPoint, WithPayloadInterface, WithVector};
use shard::count::CountRequestInternal;
use shard::facet::FacetRequestInternal;
use shard::locked_segment::LockedSegment;
use shard::query::ShardQueryRequest;
use shard::query::scroll::QueryScrollRequestInternal;
use shard::retrieve::record_internal::RecordInternal;
use shard::scroll::ScrollRequestInternal;
use shard::search::CoreSearchRequest;
use crate::{EdgeConfig, ShardInfo};
/// A handle to a single segment that can be read-locked to yield a [`ReadSegmentEntry`].
///
/// Abstracting over the handle (rather than the segment type) keeps the read path monomorphic for
/// homogeneous callers — a read-only follower's handle is the concrete
/// `Arc<RwLock<ReadOnlySegment<S>>>` — while the read-write shard, whose holder is intrinsically
/// heterogeneous (`Segment` and `ProxySegment` coexist during optimization), uses the
/// [`LockedSegment`] enum. Dynamic dispatch is confined to the `LockedSegment` impl, mirroring the
/// pre-existing [`LockedSegment::get_read`].
pub trait ReadSegmentHandle {
type Segment: ReadSegmentEntry + ?Sized;
/// Acquire a read guard. One guard is held for the whole per-segment operation, so reads that
/// call several methods on the same segment observe a consistent state.
fn read_segment(&self) -> RwLockReadGuard<'_, Self::Segment>;
/// Owned handle for the retrieval / version-dedup path ([`retrieve_over`]).
///
/// [`retrieve_over`]: shard::retrieve::retrieve_blocking::retrieve_over
fn segment_arc(&self) -> Arc<RwLock<Self::Segment>>;
}
impl<S: UniversalReadExt + 'static> ReadSegmentHandle for Arc<RwLock<ReadOnlySegment<S>>> {
type Segment = ReadOnlySegment<S>;
fn read_segment(&self) -> RwLockReadGuard<'_, ReadOnlySegment<S>> {
self.read()
}
fn segment_arc(&self) -> Arc<RwLock<ReadOnlySegment<S>>> {
self.clone()
}
}
impl ReadSegmentHandle for LockedSegment {
type Segment = dyn ReadSegmentEntry;
fn read_segment(&self) -> RwLockReadGuard<'_, dyn ReadSegmentEntry> {
self.get_read().read()
}
fn segment_arc(&self) -> Arc<RwLock<dyn ReadSegmentEntry>> {
self.get_read_arc()
}
}
/// A consistent read snapshot of an edge shard: owned segment handles (collected in retrieval order,
/// non-appendable first then appendable) plus an immutable config snapshot.
///
/// All edge read logic is implemented exactly once, here, generic over the segment handle `H`.
/// Callers do not use this type directly — they go through [`EdgeShardRead`], which builds a snapshot
/// and delegates. Because the holder lock is released once the handles are collected, a single
/// top-level read runs over one immutable snapshot; sub-reads (e.g. the searches inside a `query`)
/// share that snapshot instead of re-locking the holder.
pub struct EdgeReadView<H: ReadSegmentHandle> {
pub(crate) segments: Vec<H>,
pub(crate) config: Arc<EdgeConfig>,
}
impl<H: ReadSegmentHandle> EdgeReadView<H> {
pub(crate) fn new(segments: Vec<H>, config: Arc<EdgeConfig>) -> Self {
Self { segments, config }
}
/// Owned read handles for the retrieval / version-dedup path.
pub(crate) fn segment_arcs(&self) -> Vec<Arc<RwLock<H::Segment>>> {
self.segments
.iter()
.map(ReadSegmentHandle::segment_arc)
.collect()
}
}
/// Read API shared by the read-write [`EdgeShard`](crate::EdgeShard) and the read-only follower
/// shard.
///
/// An implementer only provides how to snapshot its segments ([`read_segments`](Self::read_segments))
/// and config ([`config_snapshot`](Self::config_snapshot)); every read operation is a default method
/// that builds an [`EdgeReadView`] from that snapshot and runs the shared logic, so the read code is
/// never duplicated.
pub trait EdgeShardRead {
/// Concrete segment handle backing this shard. A follower uses the monomorphic
/// `Arc<RwLock<ReadOnlySegment<S>>>`; the read-write shard uses `LockedSegment`.
type Handle: ReadSegmentHandle;
/// Snapshot the current segments in retrieval order (non-appendable first, then appendable).
fn read_segments(&self) -> Vec<Self::Handle>;
/// Snapshot the current config.
fn config_snapshot(&self) -> Arc<EdgeConfig>;
fn path(&self) -> &Path;
/// This method is DEPRECATED and should be replaced with query.
fn search(&self, search: CoreSearchRequest) -> OperationResult<Vec<ScoredPoint>> {
view(self).search(search)
}
fn query(&self, request: ShardQueryRequest) -> OperationResult<Vec<ScoredPoint>> {
view(self).query(request)
}
fn query_scroll(
&self,
request: &QueryScrollRequestInternal,
) -> OperationResult<Vec<ScoredPoint>> {
view(self).query_scroll(request)
}
fn scroll(
&self,
request: ScrollRequestInternal,
) -> OperationResult<(Vec<RecordInternal>, Option<PointIdType>)> {
view(self).scroll(request)
}
fn retrieve(
&self,
point_ids: &[ExtendedPointId],
with_payload: Option<WithPayloadInterface>,
with_vector: Option<WithVector>,
) -> OperationResult<Vec<RecordInternal>> {
view(self).retrieve(point_ids, with_payload, with_vector)
}
fn count(&self, request: CountRequestInternal) -> OperationResult<usize> {
view(self).count(request)
}
fn facet(&self, request: FacetRequestInternal) -> OperationResult<FacetResponse> {
view(self).facet(request)
}
fn info(&self) -> ShardInfo {
view(self).info()
}
fn rescore_with_formula(
&self,
formula: ParsedFormula,
prefetches_results: Vec<Vec<ScoredPoint>>,
limit: usize,
score_threshold: Option<ScoreType>,
hw_measurement_acc: HwMeasurementAcc,
) -> OperationResult<Vec<ScoredPoint>> {
view(self).rescore_with_formula(
formula,
prefetches_results,
limit,
score_threshold,
hw_measurement_acc,
)
}
}
/// Build a one-shot read snapshot for a shard. Private so it is not part of the trait's surface —
/// the snapshot is an implementation detail of the default read methods.
fn view<T: EdgeShardRead + ?Sized>(shard: &T) -> EdgeReadView<T::Handle> {
EdgeReadView::new(shard.read_segments(), shard.config_snapshot())
}