use std::{collections::HashSet, future::Future, sync::Arc, time::Instant};
use arrow_array::Decimal128Array;
use futures::future::try_join_all;
use roaring::RoaringBitmap;
use tracing::trace;
use uuid::Uuid;
use super::SuperfileHit;
use crate::{
runtime_metrics::op_stats::{self, OpStatsCollector},
storage::StorageProvider,
superfile::{SuperfileReader, builder::VectorConfig, vector::layout::VectorLayout},
supertable::{
error::QueryError,
handle::SupertableReader,
manifest::SuperfileEntry,
query::{
exec::common::{stamp_stable_ids, take_rows_byte_source, take_rows_object_store},
superfile_reader::superfile_reader,
vector::row_id_from_manifest_entry,
},
reader_cache::{DiskCacheStore, SuperfileReaderCache},
tombstones::SidecarCache,
},
};
#[cfg_attr(
feature = "detailed-tracing",
tracing::instrument(skip_all, fields(uri = ?entry.uri))
)]
pub(crate) async fn open_reader(
store: &Arc<dyn SuperfileReaderCache>,
disk_cache: Option<&Arc<DiskCacheStore>>,
storage: Option<&Arc<dyn StorageProvider>>,
entry: &SuperfileEntry,
allow_background_fill: bool,
) -> Result<Arc<SuperfileReader>, QueryError> {
superfile_reader(
store,
disk_cache,
storage,
&entry.uri,
entry.subsection_offsets.as_ref(),
allow_background_fill,
)
.await
.map_err(|e| QueryError::Store(e.to_string()))
}
pub(crate) fn verify_superfile_vector_codecs(
reader: &SuperfileReader,
expected: &[VectorConfig],
) -> Result<(), QueryError> {
if expected.is_empty() {
return Ok(());
}
let vector = reader.vec().ok_or_else(|| {
QueryError::Execute("superfile is missing configured vector index".into())
})?;
for config in expected {
let mut matched = false;
for column in vector
.vector_columns_config()
.filter(|column| column.name == config.column)
{
matched = true;
let stored = column.rerank_codec;
let usable = stored.supports_metric(config.metric)
&& (!vector.is_multi_cell() || stored.is_ivf_mergeable());
if !usable {
return Err(QueryError::Execute(format!(
"vector codec {} stored for {:?} cannot serve this table (metric {:?}{})",
stored.name(),
config.column,
config.metric,
if vector.is_multi_cell() {
"; a multi-cell unit requires an IVF-mergeable codec"
} else {
""
},
)));
}
}
if !matched {
return Err(QueryError::Execute(format!(
"superfile is missing configured vector column {:?}",
config.column
)));
}
}
Ok(())
}
pub(crate) async fn open_compaction_input(
store: &Arc<dyn SuperfileReaderCache>,
disk_cache: Option<&Arc<DiskCacheStore>>,
storage: Option<&Arc<dyn StorageProvider>>,
entry: &SuperfileEntry,
) -> Result<Arc<SuperfileReader>, QueryError> {
if let Some(storage) = storage {
if let Some(cache) = disk_cache {
let reader = cache
.reader_synchronous_with_storage(&entry.uri, Arc::clone(storage))
.await
.map_err(|e| QueryError::Store(e.to_string()));
if let Ok(reader) = reader
&& reader.is_fully_resident()
{
return Ok(reader);
}
}
let path = entry.uri.storage_path();
let (bytes, _) = storage
.get(&path)
.await
.map_err(|e| QueryError::Store(e.to_string()))?;
let reader = SuperfileReader::open(bytes).map_err(|e| QueryError::Store(e.to_string()))?;
return Ok(Arc::new(reader));
}
open_reader(store, disk_cache, storage, entry, true).await
}
pub(crate) fn tag_hits(entry: &SuperfileEntry, hits: Vec<(u32, f32)>) -> Vec<SuperfileHit> {
let base = row_id_from_manifest_entry(entry, 0);
hits.into_iter()
.map(|(local_doc_id, score)| SuperfileHit {
superfile: entry.uri,
local_doc_id,
score,
stable_id: base.map(|b| b + i128::from(local_doc_id)),
})
.collect()
}
pub(crate) fn tombstone_deny_set(
cache: &SidecarCache,
superfile_id: Uuid,
now: Instant,
) -> Result<Option<Arc<RoaringBitmap>>, QueryError> {
let bitmap = cache
.bitmap_for(superfile_id, now)
.map_err(|e| QueryError::Store(format!("tombstone cache: {e}")))?;
Ok((!bitmap.is_empty()).then_some(bitmap))
}
pub(crate) fn apply_tombstone_filter(
cache: Option<&Arc<SidecarCache>>,
entry: &SuperfileEntry,
hits: &mut Vec<SuperfileHit>,
now: Instant,
) -> Result<(), QueryError> {
let Some(cache) = cache else {
return Ok(());
};
let Some(bitmap) = tombstone_deny_set(cache, entry.superfile_id, now)? else {
return Ok(());
};
hits.retain(|h| !bitmap.contains(h.local_doc_id));
Ok(())
}
pub(crate) async fn attach_stable_ids(
reader: &SuperfileReader,
entry: &SuperfileEntry,
hits: &mut [SuperfileHit],
fetch_lazy_id_page: bool,
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<(), QueryError> {
if hits.is_empty() {
return Ok(());
}
if let Some(base) = row_id_from_manifest_entry(entry, 0) {
for hit in hits.iter_mut() {
hit.stable_id = Some(base + i128::from(hit.local_doc_id));
}
return Ok(());
}
let locals: Vec<u32> = hits.iter().map(|h| h.local_doc_id).collect();
if let Some(ids) = stable_ids_for_tagged_hits(reader, &locals).await? {
if let Some(stats) = op_stats {
stats.add_planned_read_ranges(1);
}
for (hit, id) in hits.iter_mut().zip(ids) {
hit.stable_id = Some(id);
}
return Ok(());
}
if !fetch_lazy_id_page {
return Ok(());
}
let id_column = reader.id_column();
let batch = if reader.can_take_by_local_doc_ids() {
let (batch, decode_ns) = op_stats::timed_section(|| {
reader
.take_by_local_doc_ids(&locals, &[id_column])
.map_err(|error| QueryError::Execute(error.to_string()))
});
if let Some(stats) = op_stats {
stats.add_kernel_cpu_ns(decode_ns);
}
batch?
} else {
take_rows_byte_source(reader, &locals, &[id_column])
.await
.map_err(|error| QueryError::Execute(error.to_string()))?
};
let ids = batch
.column(0)
.as_any()
.downcast_ref::<Decimal128Array>()
.ok_or_else(|| QueryError::Execute("_id column missing".into()))?;
for (hit, id) in hits.iter_mut().zip(ids.values()) {
hit.stable_id = Some(*id);
}
Ok(())
}
pub(crate) async fn attach_stable_ids_to_hits(
table_reader: &SupertableReader,
hits: &mut [SuperfileHit],
) -> Result<(), QueryError> {
stamp_stable_ids(table_reader, hits)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?;
if let Some(missing) = hits.iter().find(|h| h.stable_id.is_none()) {
return Err(QueryError::Execute(format!(
"hit {:?}/{} missing stable _id after search-wave stamping",
missing.superfile, missing.local_doc_id
)));
}
Ok(())
}
pub(crate) async fn apply_resolved_tombstone_filter(
reader: &SuperfileReader,
storage: Option<&Arc<dyn StorageProvider>>,
cache: Option<&Arc<SidecarCache>>,
entry: &SuperfileEntry,
hits: &mut Vec<SuperfileHit>,
now: Instant,
op_stats: &Option<Arc<OpStatsCollector>>,
) -> Result<(), QueryError> {
if entry.vector_layout != VectorLayout::MultiCellIvf {
return apply_tombstone_filter(cache, entry, hits, now);
}
let Some(cache) = cache else {
return Ok(());
};
let bitmap = cache
.bitmap_for(entry.superfile_id, now)
.map_err(|e| QueryError::Store(format!("tombstone cache: {e}")))?;
if bitmap.is_empty() {
return Ok(());
}
let locals: Vec<u32> = bitmap.iter().collect();
let id_column = reader.id_column();
let batch = if reader.parquet_bytes().is_some() {
let (batch, decode_ns) = op_stats::timed_section(|| {
reader
.take_by_local_doc_ids(&locals, &[id_column])
.map_err(|e| QueryError::Execute(e.to_string()))
});
if let Some(stats) = op_stats {
stats.add_kernel_cpu_ns(decode_ns);
}
batch?
} else {
let storage = storage.ok_or_else(|| {
QueryError::Execute(
"MultiCell tombstone resolve needs resident bytes or storage".into(),
)
})?;
let (object_store, path) = storage
.object_store_handle(&entry.uri.storage_path())
.ok_or_else(|| QueryError::Execute("no object_store handle for superfile".into()))?;
let file_size = entry
.subsection_offsets
.as_ref()
.map(|offsets| offsets.total_size);
take_rows_object_store(
object_store,
path,
file_size,
reader.schema(),
reader.n_docs(),
&locals,
&[id_column],
)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?
};
let ids = batch
.column(0)
.as_any()
.downcast_ref::<Decimal128Array>()
.ok_or_else(|| QueryError::Execute("_id column missing".into()))?;
let deleted: HashSet<i128> = ids.values().iter().copied().collect();
hits.retain(|hit| hit.stable_id.is_none_or(|id| !deleted.contains(&id)));
Ok(())
}
async fn stable_ids_for_tagged_hits(
reader: &SuperfileReader,
locals: &[u32],
) -> Result<Option<Vec<i128>>, QueryError> {
if locals.is_empty() {
return Ok(Some(Vec::new()));
}
if let Some(v) = reader.vec()
&& let Some(ids) = v.inline_stable_ids_for_locals(locals)
{
return Ok(Some(ids));
}
if let Some(v) = reader.vec()
&& let Some(ids) = v
.inline_stable_ids_for_locals_async(locals)
.await
.map_err(|e| QueryError::Execute(e.to_string()))?
{
return Ok(Some(ids));
}
if locals
.iter()
.any(|&local| u64::from(local) >= reader.n_docs())
{
return Ok(None);
}
if reader.parquet_bytes().is_none() {
return Ok(None);
}
let id_column = reader.id_column();
let batch = reader
.take_by_local_doc_ids(locals, &[id_column])
.map_err(|e| QueryError::Execute(e.to_string()))?;
let array = batch
.column(0)
.as_any()
.downcast_ref::<Decimal128Array>()
.ok_or_else(|| QueryError::Execute("_id column missing".into()))?;
Ok(Some(array.values().to_vec()))
}
pub(crate) async fn fanout_local_hits<P, K, Fut>(
reader: &SupertableReader,
units: Vec<(Arc<SuperfileEntry>, P)>,
kernel: K,
) -> Result<Vec<Vec<SuperfileHit>>, QueryError>
where
P: Send + 'static,
K: Fn(Arc<SuperfileReader>, P) -> Fut + Clone + Send + 'static,
Fut: Future<Output = Result<Vec<(u32, f32)>, QueryError>> + Send + 'static,
{
fanout_with(
reader,
units,
true,
true, move |r, entry, tombstone_cache, now, params| {
let kernel = kernel.clone();
async move {
let hits = kernel(r, params).await?;
let mut tagged = tag_hits(&entry, hits);
apply_tombstone_filter(tombstone_cache.as_ref(), &entry, &mut tagged, now)?;
Ok::<Vec<SuperfileHit>, QueryError>(tagged)
}
},
)
.await
}
pub(crate) async fn fanout_with<P, R, B, Fut>(
reader: &SupertableReader,
units: Vec<(Arc<SuperfileEntry>, P)>,
prefetch_tombstones: bool,
allow_background_fill: bool,
body: B,
) -> Result<Vec<R>, QueryError>
where
P: Send + 'static,
R: Send + 'static,
B: Fn(Arc<SuperfileReader>, Arc<SuperfileEntry>, Option<Arc<SidecarCache>>, Instant, P) -> Fut
+ Clone
+ Send
+ 'static,
Fut: Future<Output = Result<R, QueryError>> + Send + 'static,
{
if units.is_empty() {
return Ok(Vec::new());
}
trace!(units = units.len(), "fanning query out across superfiles");
let manifest = reader.manifest();
let store = Arc::clone(&manifest.options.store);
let disk_cache = manifest.options.disk_cache.as_ref().map(Arc::clone);
let storage = manifest.options.storage.as_ref().map(Arc::clone);
let vector_columns = Arc::new(manifest.options.vector_columns.clone());
let tombstone_cache = reader.tombstone_cache.clone();
let now = Instant::now();
if prefetch_tombstones && let Some(cache) = tombstone_cache.as_ref() {
let mut ids: Vec<Uuid> = units.iter().map(|(e, _)| e.superfile_id).collect();
ids.sort_unstable();
ids.dedup();
cache.prefetch(&ids, now).await;
}
if units.len() == 1 {
let (entry, params) = units.into_iter().next().expect("len == 1");
let r = open_reader(
&store,
disk_cache.as_ref(),
storage.as_ref(),
&entry,
allow_background_fill,
)
.await?;
verify_superfile_vector_codecs(&r, &vector_columns)?;
let out = body(r, entry, tombstone_cache, now, params).await?;
return Ok(vec![out]);
}
let handles = units.into_iter().map(|(entry, params)| {
let store = Arc::clone(&store);
let disk_cache = disk_cache.clone();
let storage = storage.clone();
let tombstone_cache = tombstone_cache.clone();
let body = body.clone();
let vector_columns = Arc::clone(&vector_columns);
let handle = tokio::spawn(async move {
let r = open_reader(
&store,
disk_cache.as_ref(),
storage.as_ref(),
&entry,
allow_background_fill,
)
.await?;
verify_superfile_vector_codecs(&r, &vector_columns)?;
body(r, entry, tombstone_cache, now, params).await
});
async move {
handle
.await
.map_err(|e| QueryError::Store(format!("fan-out task join: {e}")))?
}
});
try_join_all(handles).await
}
#[cfg(test)]
mod codec_verify_tests {
use std::sync::Arc;
use arrow_array::RecordBatch;
use arrow_schema::Schema;
use bytes::Bytes;
use super::verify_superfile_vector_codecs;
use crate::{
superfile::{
SuperfileReader,
builder::{BuilderOptions, SuperfileBuilder, VectorConfig},
vector::{distance::Metric, rerank_codec::RerankCodec},
},
test_helpers::{decimal128_id_field, decimal128_ids, default_vector_config},
};
const DIM: usize = 16;
fn cfg(metric: Metric, codec: RerankCodec) -> VectorConfig {
VectorConfig {
metric,
..default_vector_config("emb", 7).with_rerank_codec(codec)
}
}
fn build_reader(metric: Metric, codec: RerankCodec) -> SuperfileReader {
let schema = Arc::new(Schema::new(vec![decimal128_id_field("doc_id")]));
let opts = BuilderOptions::new(
schema.clone(),
"doc_id",
vec![],
vec![cfg(metric, codec)],
None,
);
let mut b = SuperfileBuilder::new(opts).expect("new builder");
let batch = RecordBatch::try_new(schema, vec![Arc::new(decimal128_ids(vec![10u64, 11]))])
.expect("batch");
let mut v = vec![0.0f32; 2 * DIM]; v[0] = 1.0;
v[DIM + 1] = 1.0;
b.add_batch(&batch, &[v.as_slice()]).expect("add_batch");
SuperfileReader::open(Bytes::from(b.finish().expect("finish"))).expect("open")
}
#[test]
fn reopen_after_default_flip_accepts_stored_codec() {
let reader = build_reader(Metric::L2Sq, RerankCodec::Sq8Residual);
let expected = [cfg(Metric::L2Sq, RerankCodec::Sq16Adaptive)];
verify_superfile_vector_codecs(&reader, &expected)
.expect("stored Sq8Residual must be accepted under the Sq16Adaptive default");
}
#[test]
fn verifier_rejects_metric_incompatible_stored_codec() {
let reader = build_reader(Metric::Cosine, RerankCodec::Sq16);
let expected = [cfg(Metric::L2Sq, RerankCodec::Sq16Adaptive)];
assert!(
verify_superfile_vector_codecs(&reader, &expected).is_err(),
"a cosine-only Sq16 codec must not be accepted for an L2Sq table"
);
}
}