use super::{
Arc, Change, DocumentChanges, DocumentSelection, DocumentStore, StorageBackendResult,
StorageReadControl,
};
use uqa_core::memory::{Budgeted, BudgetedVec, MemoryReservation};
use uqa_sql::ast::ColumnDef;
use uqa_storage::{
diskann_index::{DiskANNReadChanges, DiskANNReadSnapshot},
vector_index::{
RetainedVectorIndexesBuilder, SelectedVectorRead, VectorIndexSource, VectorReadSnapshot,
},
};
struct VectorSource {
field: String,
column: Option<[u8; 16]>,
source: Option<DiskANNReadSnapshot>,
values: VectorReadSnapshot,
_memory: MemoryReservation,
}
pub(super) struct CapturedRows {
pub(super) documents: Arc<dyn DocumentStore>,
vectors: BudgetedVec<VectorSource>,
}
impl CapturedRows {
fn vector(&self, field: &str, column: Option<&ColumnDef>) -> Option<&VectorSource> {
let identity = column.and_then(|column| column.object_id);
self.vectors
.iter()
.find(|source| match (source.column, identity) {
(Some(source), Some(target)) => source == target,
_ => source.field == field,
})
}
}
impl DocumentChanges {
pub fn with_retained_vectors(
mut self,
documents: Arc<dyn DocumentStore>,
desired: DocumentSelection,
columns: &[ColumnDef],
indexes: &dyn VectorIndexSource,
control: &StorageReadControl,
) -> StorageBackendResult<Self> {
control.check()?;
let mut vectors = BudgetedVec::new(control.memory());
indexes.visit(&mut |field, index| {
let source = index.diskann_read_snapshot(control)?;
let values = if let Some(source) = &source {
Some(
Budgeted::new(source.clone(), control.memory().empty_reservation())
.into_shared()? as VectorReadSnapshot,
)
} else {
index.vector_read_snapshot(control)?
};
let Some(values) = values else {
return Ok(());
};
vectors.reserve(1)?;
let column = columns
.iter()
.find(|column| column.name == field)
.and_then(|column| column.object_id);
let (field, memory) = RetainedVectorIndexesBuilder::copy_field(field, control)?;
vectors.push(VectorSource {
field,
column,
source,
values,
_memory: memory,
})?;
Ok(())
})?;
if vectors.is_empty() {
return self.with_retained(documents, desired, control);
}
let source = Budgeted::new(
CapturedRows { documents, vectors },
control.memory().empty_reservation(),
)
.into_shared()?;
let desired = desired.finish(control)?;
let mut newer = Self::default();
for (id, present) in desired.entries() {
control.check()?;
let present = present && source.documents.contains_doc_id(id)?;
newer.insert(id, Change::Captured(source.clone(), present), control)?;
}
self.extend(newer, control)?;
Ok(self)
}
pub(in crate::query) fn diskann_read_changes(
&self,
field: &str,
column: Option<&ColumnDef>,
control: &StorageReadControl,
) -> StorageBackendResult<Option<DiskANNReadChanges>> {
for (_, change) in self.rows() {
control.check()?;
if change.present()
&& !matches!(change, Change::Captured(source, true) if source.vector(field, column).is_some_and(|source| source.source.is_some()))
{
return Ok(None);
}
}
DiskANNReadChanges::capture(
self.rows().iter().map(|(document, change)| {
if !change.present() {
return Ok((*document, None));
}
let Change::Captured(source, _) = change else {
unreachable!("validated private source");
};
Ok((
*document,
Some(
source
.vector(field, column)
.expect("validated private column")
.source
.as_ref()
.expect("validated private physical source")
.clone(),
),
))
}),
control,
)
.map(Some)
}
pub(in crate::query) fn vector_read_selection(
&self,
base: Option<VectorReadSnapshot>,
dimensions: u32,
field: &str,
column: Option<&ColumnDef>,
control: &StorageReadControl,
) -> StorageBackendResult<Option<VectorReadSnapshot>> {
for (_, change) in self.rows() {
control.check()?;
if change.present()
&& !matches!(change, Change::Captured(source, true) if source.vector(field, column).is_some())
{
return Ok(None);
}
}
SelectedVectorRead::capture(
base,
dimensions,
self.rows().iter().map(|(document, change)| {
let source = match change {
Change::Captured(source, true) => Some(
source
.vector(field, column)
.expect("validated private column")
.values
.clone(),
),
_ => None,
};
(*document, source)
}),
control,
)
.map(Some)
}
}