use std::sync::{Arc, LazyLock};
use crate::library::support::cache::OnceCache;
use vectordata::TestDataGroup;
use vectordata::TestDataView;
use vectordata::catalog::resolver::Catalog;
use vectordata::catalog::sources::CatalogSources;
use vectordata::io::{VectorReader, VvecReader};
static DATASET_CACHE: LazyLock<OnceCache<String, Arc<TestDataGroup>>> =
LazyLock::new(OnceCache::new);
type FacetCache = OnceCache<(String, String, String), Arc<dyn std::any::Any + Send + Sync>>;
static FACET_CACHE: LazyLock<FacetCache> = LazyLock::new(OnceCache::new);
static PREBUFFER_CACHE: LazyLock<OnceCache<String, Arc<DatasetHandle>>> =
LazyLock::new(OnceCache::new);
fn parse_source_specifier(source: &str) -> (&str, &str) {
if source.starts_with("http://") || source.starts_with("https://") {
return (source, "default");
}
if let Some(pos) = source.find(':') {
(&source[..pos], &source[pos + 1..])
} else {
(source, "default")
}
}
fn run_blocking_io<R>(body: impl FnOnce() -> R) -> R {
#[cfg(feature = "vectordata")]
if tokio::runtime::Handle::try_current().is_ok() {
return tokio::task::block_in_place(body);
}
body()
}
pub(crate) fn load_dataset_group(source: &str) -> Result<Arc<TestDataGroup>, String> {
let (dataset_name, _profile) = parse_source_specifier(source);
DATASET_CACHE.get_or_init(dataset_name.to_string(), || {
run_blocking_io(|| {
let catalog = Catalog::of(&CatalogSources::new().configure_default());
catalog
.open(dataset_name)
.map(Arc::new)
.map_err(|e| format!("failed to load dataset '{dataset_name}': {e}"))
})
})
}
pub(crate) struct UniformDataset<T: Send + Sync + 'static> {
reader: Arc<dyn VectorReader<T>>,
count: usize,
dim: usize,
}
fn load_uniform_facet<T: Send + Sync + 'static>(
source: &str,
profile: &str,
facet: &str,
open_fn: impl FnOnce(
&dyn TestDataView,
) -> std::result::Result<Arc<dyn VectorReader<T>>, vectordata::Error>,
) -> Result<Arc<UniformDataset<T>>, String> {
let key = (source.to_string(), profile.to_string(), facet.to_string());
let any = FACET_CACHE.get_or_init(key, || {
let group = load_dataset_group(source)?;
let view = group
.profile(profile)
.ok_or_else(|| format!("profile '{profile}' not found in '{source}'"))?;
crate::library::support::audit::record_opened(source, profile, facet, "uniform");
let reader = run_blocking_io(|| open_fn(view.as_ref()))
.map_err(|e| format!("failed to access {facet} from '{source}': {e}"))?;
let count = reader.count();
let dim = reader.dim();
let arc: Arc<UniformDataset<T>> = Arc::new(UniformDataset { reader, count, dim });
Ok(arc as Arc<dyn std::any::Any + Send + Sync>)
})?;
any.downcast::<UniformDataset<T>>().map_err(|_| {
format!(
"facet cache type mismatch for '{source}:{profile}/{facet}' — \
this should be impossible; please file a bug."
)
})
}
type F32Dataset = UniformDataset<f32>;
type I32Dataset = UniformDataset<i32>;
impl F32Dataset {
fn load(source: &str, profile: &str, facet: &str) -> Result<Arc<Self>, String> {
let facet_name = facet.to_string();
load_uniform_facet(source, profile, facet, move |view| {
match facet_name.as_str() {
"base" => view.base_vectors(),
"query" => view.query_vectors(),
"neighbor_distances" => view.neighbor_distances(),
"filtered_neighbor_distances" => view.prefiltered_neighbor_distances(),
other => Err(vectordata::Error::MissingFacet(format!(
"unknown f32 facet: '{other}'"
))),
}
})
}
}
impl I32Dataset {
fn load(source: &str, profile: &str, facet: &str) -> Result<Arc<Self>, String> {
let facet_name = facet.to_string();
load_uniform_facet(source, profile, facet, move |view| {
match facet_name.as_str() {
"neighbor_indices" => view.neighbor_indices(),
"filtered_neighbor_indices" => view.prefiltered_neighbor_indices(),
other => Err(vectordata::Error::MissingFacet(format!(
"unknown i32 facet: '{other}'"
))),
}
})
}
}
#[derive(Clone)]
pub(crate) enum DatasetHandle {
F32(Arc<F32Dataset>),
I32(Arc<I32Dataset>),
Ivvec32(Arc<Ivvec32Dataset>),
Generic(Arc<GenericFacetDataset>),
Group(Arc<TestDataGroup>),
Prebuffered {
group: Arc<TestDataGroup>,
source: String,
},
}
impl DatasetHandle {
fn open(source: &str, facet: &str) -> Result<Self, String> {
let (_, profile) = parse_source_specifier(source);
match facet {
"base" | "query" | "neighbor_distances" | "filtered_neighbor_distances" => {
F32Dataset::load(source, profile, facet).map(DatasetHandle::F32)
}
"neighbor_indices" | "filtered_neighbor_indices" => {
I32Dataset::load(source, profile, facet).map(DatasetHandle::I32)
}
"metadata_results" => Ivvec32Dataset::load(source, profile).map(DatasetHandle::Ivvec32),
_ => GenericFacetDataset::load(source, profile, facet).map(DatasetHandle::Generic),
}
}
fn open_group(source: &str) -> Result<Self, String> {
load_dataset_group(source).map(DatasetHandle::Group)
}
}
#[crate::polydat_node(category = RealData)]
fn dataset_open(source: &str, facet: &str) -> Arc<DatasetHandle> {
match DatasetHandle::open(source, facet) {
Ok(h) => Arc::new(h),
Err(e) => {
let msg = format!("dataset_open: failed to resolve '{source}' facet='{facet}': {e}");
crate::library::support::audit::error(&msg);
panic!("{msg}");
}
}
}
#[crate::polydat_node(category = RealData)]
fn dataset_group_open(source: &str) -> Arc<DatasetHandle> {
match DatasetHandle::open_group(source) {
Ok(h) => Arc::new(h),
Err(e) => {
let msg = format!("dataset_group_open: failed to resolve '{source}': {e}");
crate::library::support::audit::error(&msg);
panic!("{msg}");
}
}
}
impl DatasetHandle {
fn resolve_facet<'a>(&'a self, facet: &str) -> std::borrow::Cow<'a, DatasetHandle> {
match self {
DatasetHandle::Prebuffered { source, .. } => match DatasetHandle::open(source, facet) {
Ok(opened) => std::borrow::Cow::Owned(opened),
Err(e) => panic!(
"DatasetHandle::resolve_facet: failed to open \
'{facet}' from prebuffered '{source}': {e}"
),
},
_ => std::borrow::Cow::Borrowed(self),
}
}
fn resolve_group(&self) -> &TestDataGroup {
match self {
DatasetHandle::Group(g) | DatasetHandle::Prebuffered { group: g, .. } => g.as_ref(),
other => panic!(
"expected a group or prebuffered handle, got {}",
dataset_handle_kind(other)
),
}
}
}
macro_rules! facet_kind {
($name:ident, $facet:literal) => {
#[doc = $facet]
pub struct $name;
impl $name {
pub const FACET: &'static str = $facet;
}
impl crate::derive_support::ResolverKind for $name {
const RESOLVER: crate::dsl::registry::DefaultResolver =
crate::dsl::registry::DefaultResolver::Facet($facet);
}
};
}
facet_kind!(BaseFacet, "base");
facet_kind!(QueryFacet, "query");
facet_kind!(NeighborIndicesFacet, "neighbor_indices");
facet_kind!(NeighborDistancesFacet, "neighbor_distances");
facet_kind!(FilteredNeighborIndicesFacet, "filtered_neighbor_indices");
facet_kind!(
FilteredNeighborDistancesFacet,
"filtered_neighbor_distances"
);
facet_kind!(MetadataResultsFacet, "metadata_results");
facet_kind!(MetadataContentFacet, "metadata_content");
facet_kind!(MetadataPredicatesFacet, "metadata_predicates");
type Facet<R> = crate::derive_support::Resolved<R, DatasetHandle>;
type Group = crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>;
fn facet_record_count(h: &DatasetHandle, who: &str, facet: &str) -> u64 {
match h.resolve_facet(facet).as_ref() {
DatasetHandle::F32(d) => d.count as u64,
DatasetHandle::I32(d) => d.count as u64,
DatasetHandle::Ivvec32(d) => d.count as u64,
DatasetHandle::Generic(d) => d.count as u64,
other => panic!(
"{who}: expected a facet handle, got {}",
dataset_handle_kind(other)
),
}
}
fn f32_vec_typed(h: &DatasetHandle, index: usize) -> crate::ast::SliceArc<f32> {
match h {
DatasetHandle::F32(d) => slice_arc_from_uniform(d, index),
other => panic!(
"expected F32 dataset handle, got {}",
dataset_handle_kind(other)
),
}
}
fn i32_vec_typed(h: &DatasetHandle, index: usize) -> crate::ast::SliceArc<i32> {
match h {
DatasetHandle::I32(d) => slice_arc_from_uniform(d, index),
other => panic!(
"expected I32 dataset handle, got {}",
dataset_handle_kind(other)
),
}
}
fn ivvec32_vec_typed(h: &DatasetHandle, index: usize) -> crate::ast::SliceArc<i32> {
match h {
DatasetHandle::Ivvec32(d) => {
if d.count == 0 {
return crate::ast::SliceArc::from_vec(Vec::<i32>::new());
}
let v = d.reader.get(index % d.count).unwrap_or_default();
crate::ast::SliceArc::from_vec(v)
}
other => panic!(
"expected Ivvec32 dataset handle, got {}",
dataset_handle_kind(other)
),
}
}
fn slice_arc_from_uniform<T>(d: &Arc<UniformDataset<T>>, index: usize) -> crate::ast::SliceArc<T>
where
T: Send + Sync + Copy + 'static,
{
if d.count == 0 {
return crate::ast::SliceArc::from_vec(Vec::<T>::new());
}
let idx = index % d.count;
if let Some(slice) = d.reader.get_slice(idx) {
let ptr_len = (slice.as_ptr(), slice.len());
let owner = d.clone();
let owner_dyn: Arc<dyn std::any::Any + Send + Sync> = owner;
return unsafe {
crate::ast::SliceArc::from_borrowed(
owner_dyn,
std::slice::from_raw_parts(ptr_len.0, ptr_len.1),
)
};
}
crate::ast::SliceArc::from_vec(d.reader.get(idx).unwrap_or_default())
}
fn dataset_handle_kind(h: &DatasetHandle) -> &'static str {
match h {
DatasetHandle::F32(_) => "F32",
DatasetHandle::I32(_) => "I32",
DatasetHandle::Ivvec32(_) => "Ivvec32",
DatasetHandle::Generic(_) => "Generic",
DatasetHandle::Group(_) => "Group",
DatasetHandle::Prebuffered { .. } => "Prebuffered",
}
}
#[crate::polydat_node(category = RealData)]
fn vector_at(handle: Facet<BaseFacet>, index: u64) -> crate::ast::SliceArc<f32> {
f32_vec_typed(
handle.resolve_facet(BaseFacet::FACET).as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn query_vector_at(handle: Facet<QueryFacet>, index: u64) -> crate::ast::SliceArc<f32> {
f32_vec_typed(
handle.resolve_facet(QueryFacet::FACET).as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn neighbor_indices_at(
handle: Facet<NeighborIndicesFacet>,
index: u64,
) -> crate::ast::SliceArc<i32> {
i32_vec_typed(
handle.resolve_facet(NeighborIndicesFacet::FACET).as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn neighbor_distances_at(
handle: Facet<NeighborDistancesFacet>,
index: u64,
) -> crate::ast::SliceArc<f32> {
f32_vec_typed(
handle.resolve_facet(NeighborDistancesFacet::FACET).as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn filtered_neighbor_indices_at(
handle: Facet<FilteredNeighborIndicesFacet>,
index: u64,
) -> crate::ast::SliceArc<i32> {
i32_vec_typed(
handle
.resolve_facet(FilteredNeighborIndicesFacet::FACET)
.as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn filtered_neighbor_distances_at(
handle: Facet<FilteredNeighborDistancesFacet>,
index: u64,
) -> crate::ast::SliceArc<f32> {
f32_vec_typed(
handle
.resolve_facet(FilteredNeighborDistancesFacet::FACET)
.as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn vector_dim(handle: Facet<BaseFacet>) -> u64 {
match handle.resolve_facet(BaseFacet::FACET).as_ref() {
DatasetHandle::F32(d) => d.dim as u64,
DatasetHandle::I32(d) => d.dim as u64,
_ => 0,
}
}
#[crate::polydat_node(category = RealData)]
fn vector_count(handle: Facet<BaseFacet>) -> u64 {
facet_record_count(&handle, "vector_count", BaseFacet::FACET)
}
#[crate::polydat_node(category = RealData)]
fn query_count(handle: Facet<QueryFacet>) -> u64 {
facet_record_count(&handle, "query_count", QueryFacet::FACET)
}
#[crate::polydat_node(category = RealData)]
fn neighbor_count(handle: Facet<NeighborIndicesFacet>) -> u64 {
match handle.resolve_facet(NeighborIndicesFacet::FACET).as_ref() {
DatasetHandle::I32(d) => d.dim as u64,
_ => 0,
}
}
#[crate::polydat_node(category = RealData)]
fn dataset_distance_function(group: Group) -> String {
let group = group.resolve_group();
let raw = group
.attribute("distance_function")
.and_then(|v| v.as_str())
.unwrap_or("unknown");
match raw.to_uppercase().as_str() {
"L2" | "EUCLIDEAN" => "EUCLIDEAN",
"L1" | "MANHATTAN" => "MANHATTAN",
"COSINE" => "COSINE",
"DOT_PRODUCT" | "DOTPRODUCT" | "DOT" | "INNER_PRODUCT" | "IP" => "DOT_PRODUCT",
_ => raw,
}
.to_string()
}
pub(crate) struct Ivvec32Dataset {
reader: Arc<dyn VvecReader<i32>>,
count: usize,
}
impl Ivvec32Dataset {
fn load(source: &str, profile: &str) -> Result<Arc<Self>, String> {
let key = (
source.to_string(),
profile.to_string(),
"metadata_results".to_string(),
);
let any = FACET_CACHE.get_or_init(key, || {
let group = load_dataset_group(source)?;
let view = group
.profile(profile)
.ok_or_else(|| format!("profile '{profile}' not found in '{source}'"))?;
crate::library::support::audit::record_opened(
source,
profile,
"metadata_results",
"ivvec32",
);
let reader = run_blocking_io(|| view.metadata_results())
.map_err(|e| format!("failed to access metadata_results from '{source}': {e}"))?;
let count = reader.count();
let arc: Arc<Self> = Arc::new(Self { reader, count });
Ok(arc as Arc<dyn std::any::Any + Send + Sync>)
})?;
any.downcast::<Self>().map_err(|_| {
format!("facet cache type mismatch for '{source}:{profile}/metadata_results'")
})
}
}
#[crate::polydat_node(category = RealData)]
fn metadata_results_at(
handle: Facet<MetadataResultsFacet>,
index: u64,
) -> crate::ast::SliceArc<i32> {
ivvec32_vec_typed(
handle.resolve_facet(MetadataResultsFacet::FACET).as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn metadata_results_len_at(handle: Facet<MetadataResultsFacet>, index: u64) -> u64 {
match handle.resolve_facet(MetadataResultsFacet::FACET).as_ref() {
DatasetHandle::Ivvec32(d) if d.count > 0 => {
d.reader.dim_at(index as usize % d.count).unwrap_or(0) as u64
}
_ => 0,
}
}
#[crate::polydat_node(category = RealData)]
fn metadata_results_count(handle: Facet<MetadataResultsFacet>) -> u64 {
match handle.resolve_facet(MetadataResultsFacet::FACET).as_ref() {
DatasetHandle::Ivvec32(d) => d.count as u64,
_ => 0,
}
}
#[crate::polydat_node(category = RealData)]
fn dataset_facets(group: Group) -> String {
let group = group.resolve_group();
let names = group.profile_names();
if let Some(first) = names.first()
&& let Some(view) = group.profile(first)
{
let manifest = view.facet_manifest();
let mut names: Vec<&String> = manifest.keys().collect();
names.sort();
return names
.iter()
.map(|s| s.as_str())
.collect::<Vec<_>>()
.join(", ");
}
String::new()
}
#[crate::polydat_node(category = RealData)]
fn dataset_profile_count(group: Group) -> u64 {
group.resolve_group().profile_names().len() as u64
}
#[crate::polydat_node(category = RealData)]
fn dataset_profile_names(group: Group) -> String {
group.resolve_group().profile_names().join(", ")
}
#[crate::polydat_node(category = RealData)]
fn matching_profiles(
group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
prefix: &str,
) -> String {
let group: &TestDataGroup = group.resolve_group();
let all = group.profile_names();
let mut matched: Vec<&str> = if prefix.is_empty() {
all.iter().map(|s| s.as_str()).collect()
} else {
all.iter()
.filter(|s| s.starts_with(prefix))
.map(|s| s.as_str())
.collect()
};
matched.sort_by(|a, b| natural_cmp(a, b));
matched.join(",")
}
fn natural_cmp(a: &str, b: &str) -> std::cmp::Ordering {
let mut ai = a.chars().peekable();
let mut bi = b.chars().peekable();
loop {
match (ai.peek().copied(), bi.peek().copied()) {
(None, None) => return std::cmp::Ordering::Equal,
(None, _) => return std::cmp::Ordering::Less,
(_, None) => return std::cmp::Ordering::Greater,
(Some(ac), Some(bc)) => {
if ac.is_ascii_digit() && bc.is_ascii_digit() {
let mut na: u64 = 0;
while let Some(c) = ai.peek().copied()
&& c.is_ascii_digit()
{
na = na
.saturating_mul(10)
.saturating_add((c as u8 - b'0') as u64);
ai.next();
}
let mut nb: u64 = 0;
while let Some(c) = bi.peek().copied()
&& c.is_ascii_digit()
{
nb = nb
.saturating_mul(10)
.saturating_add((c as u8 - b'0') as u64);
bi.next();
}
match na.cmp(&nb) {
std::cmp::Ordering::Equal => continue,
non_eq => return non_eq,
}
} else {
match ac.cmp(&bc) {
std::cmp::Ordering::Equal => {
ai.next();
bi.next();
}
non_eq => return non_eq,
}
}
}
}
}
}
#[crate::polydat_node(category = RealData)]
fn dataset_profile_name_at(
group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
index: u64,
) -> String {
let group: &TestDataGroup = group.resolve_group();
let names = group.profile_names();
if names.is_empty() {
String::new()
} else {
names[(index as usize) % names.len()].clone()
}
}
#[crate::polydat_node(category = RealData)]
fn profile_base_count(
group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
index: u64,
) -> u64 {
let group: &TestDataGroup = group.resolve_group();
let names = group.profile_names();
if names.is_empty() {
0
} else {
let name = &names[(index as usize) % names.len()];
group
.profile(name)
.and_then(|view| view.base_count())
.unwrap_or(0)
}
}
#[crate::polydat_node(category = RealData)]
fn profile_facets(
group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
index: u64,
) -> String {
let group: &TestDataGroup = group.resolve_group();
let names = group.profile_names();
if names.is_empty() {
String::new()
} else {
let name = &names[(index as usize) % names.len()];
match group.profile(name) {
Some(view) => {
let manifest = view.facet_manifest();
let mut fnames: Vec<&String> = manifest.keys().collect();
fnames.sort();
fnames
.iter()
.map(|s| s.as_str())
.collect::<Vec<_>>()
.join(", ")
}
None => String::new(),
}
}
}
#[crate::polydat_node(category = RealData)]
fn profile_partitions(
group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
pattern: &str,
) -> crate::derive_support::Ext<crate::iteration::cursor_partition::PartitionList> {
let group: &TestDataGroup = group.resolve_group();
let parts = build_profile_partitions(group, pattern);
crate::derive_support::Ext(crate::iteration::cursor_partition::PartitionList::new(
parts,
))
}
fn masked_profile_tiers(group: &TestDataGroup, pattern: &str) -> Vec<(String, u64)> {
let (re, _) = crate::library::support::pattern::compile_pattern(pattern)
.unwrap_or_else(|e| panic!("profile pattern: {e}"));
group
.profile_names()
.iter()
.filter(|n| re.is_match(n))
.filter_map(|n| {
group
.profile(n)
.and_then(|v| v.base_count())
.filter(|&c| c > 0)
.map(|c| (n.clone(), c))
})
.collect()
}
pub(crate) fn build_profile_partitions(
group: &TestDataGroup,
pattern: &str,
) -> Vec<crate::iteration::cursor_partition::Partition> {
let masked: Vec<u64> = masked_profile_tiers(group, pattern)
.into_iter()
.map(|(_, c)| c)
.collect();
let base_extent = masked.last().copied().unwrap_or(0);
let count = masked.len() as u64;
let pct = |o: u64| {
if base_extent == 0 {
0.0
} else {
(o as f64 / base_extent as f64) * 100.0
}
};
let mut parts: Vec<crate::iteration::cursor_partition::Partition> =
Vec::with_capacity(masked.len());
let mut prev: u64 = 0;
for (k, &this) in masked.iter().enumerate() {
parts.push(crate::iteration::cursor_partition::Partition {
idx: k as u64,
count,
start_ord: prev,
end_ord: this,
start_pct: pct(prev),
end_pct: pct(this),
base_extent,
});
prev = this;
}
parts
}
#[crate::polydat_node(category = RealData)]
fn matching_profile_name_at(
group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
pattern: &str,
index: u64,
) -> String {
let group: &TestDataGroup = group.resolve_group();
let tiers = masked_profile_tiers(group, pattern);
if tiers.is_empty() {
String::new()
} else {
tiers[(index as usize) % tiers.len()].0.clone()
}
}
#[crate::polydat_node(category = RealData)]
fn dataset_prebuffer(source: &str) -> Option<Arc<dyn std::any::Any + Send + Sync>> {
match PREBUFFER_CACHE.get_or_init(source.to_string(), || {
crate::library::support::audit::record_prebuffer_entered(source);
do_dataset_prebuffer_inner(source)
}) {
Ok(handle) => Some(handle as Arc<dyn std::any::Any + Send + Sync>),
Err(_) => None,
}
}
fn do_dataset_prebuffer_inner(source: &str) -> Result<Arc<DatasetHandle>, String> {
let group = match load_dataset_group(source) {
Ok(g) => g,
Err(e) => {
let msg = format!("dataset_prebuffer: cannot resolve '{source}': {e}");
crate::library::support::audit::error(&msg);
return Err(msg);
}
};
let group_for_handle = group.clone();
let (_, profile) = parse_source_specifier(source);
let view = match group.profile(profile) {
Some(v) => v,
None => {
crate::library::support::audit::error(&format!(
"dataset_prebuffer: profile '{profile}' not found in '{source}'"
));
return Ok(Arc::new(DatasetHandle::Prebuffered {
group: group_for_handle,
source: source.to_string(),
}));
}
};
let mut facet_count: u64 = 0;
for (name, _descriptor) in view.facet_manifest() {
if view.facet_element_type(&name).is_err() {
continue;
}
let storage = match run_blocking_io(|| view.open_facet_storage(&name)) {
Ok(s) => s,
Err(e) => {
crate::library::support::audit::warn(&format!(
"dataset_prebuffer: open '{name}' for prebuffer failed: {e}"
));
continue;
}
};
let prebuf_source = source.to_string();
let prebuf_profile = profile.to_string();
let prebuf_facet = name.clone();
let mut last_done: u64 = 0;
let facet_start = std::time::Instant::now();
let mut last_log_at = facet_start;
let prebuf_result = run_blocking_io(|| {
storage.prebuffer_with_progress(|p| {
let now = std::time::Instant::now();
let interval =
if now.duration_since(facet_start) < std::time::Duration::from_secs(10) {
std::time::Duration::from_secs(1)
} else {
std::time::Duration::from_secs(10)
};
let since_last = now.duration_since(last_log_at);
if since_last < interval {
return;
}
last_log_at = now;
let total_b = p.total_bytes();
let done_b = p.downloaded_bytes();
let total_c = p.total_chunks();
let done_c = p.completed_chunks();
let pct = p.fraction() * 100.0;
let delta_mb = (done_b.saturating_sub(last_done)) as f64 / (1024.0 * 1024.0);
let rate_mb_s = delta_mb / since_last.as_secs_f64().max(0.001);
last_done = done_b;
crate::library::support::audit::info(&format!(
"prebuffer: progress {prebuf_source}:{prebuf_profile}/{prebuf_facet} \
{pct:5.1}% ({done_c}/{total_c} chunks, \
{done_mb:.1}/{total_mb:.1} MB, {rate_mb_s:.1} MB/s)",
done_mb = done_b as f64 / (1024.0 * 1024.0),
total_mb = total_b as f64 / (1024.0 * 1024.0),
));
})
});
if let Err(e) = prebuf_result {
crate::library::support::audit::warn(&format!(
"dataset_prebuffer: download error for '{source}' facet '{name}': {e}"
));
continue;
}
facet_count = facet_count.saturating_add(1);
crate::library::support::audit::record_prebuffered(source, profile, &name);
}
crate::library::support::audit::log_prebuffer_summary(source, profile, facet_count);
Ok(Arc::new(DatasetHandle::Prebuffered {
group: group_for_handle,
source: source.to_string(),
}))
}
pub(crate) struct GenericFacetDataset {
reader: vectordata::typed_access::TypedReader<i64>,
count: usize,
}
impl GenericFacetDataset {
fn load(source: &str, profile: &str, facet: &str) -> Result<Arc<Self>, String> {
let key = (source.to_string(), profile.to_string(), facet.to_string());
let any = FACET_CACHE.get_or_init(key, || {
let group = load_dataset_group(source)?;
let gv = group
.generic_view(profile)
.ok_or_else(|| format!("profile '{profile}' not found in '{source}'"))?;
crate::library::support::audit::record_opened(source, profile, facet, "generic-typed");
let reader = run_blocking_io(|| gv.open_facet_typed::<i64>(facet))
.map_err(|e| format!("failed to open {facet} from '{source}:{profile}': {e}"))?;
let count = reader.count();
let arc: Arc<Self> = Arc::new(Self { reader, count });
Ok(arc as Arc<dyn std::any::Any + Send + Sync>)
})?;
any.downcast::<Self>()
.map_err(|_| format!("facet cache type mismatch for '{source}:{profile}/{facet}'"))
}
fn get_scalar(&self, index: usize) -> i64 {
if self.count == 0 {
return 0;
}
self.reader.get_value(index % self.count).unwrap_or(0)
}
fn format_scalar(&self, index: usize) -> String {
self.get_scalar(index).to_string()
}
}
fn generic_str_typed(h: &DatasetHandle, idx: usize) -> String {
match h {
DatasetHandle::Generic(d) => d.format_scalar(idx),
_ => String::new(),
}
}
#[crate::polydat_node(category = RealData)]
fn metadata_value_at(handle: Facet<MetadataContentFacet>, index: u64) -> String {
generic_str_typed(
handle.resolve_facet(MetadataContentFacet::FACET).as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn predicate_value_at(handle: Facet<MetadataPredicatesFacet>, index: u64) -> String {
generic_str_typed(
handle
.resolve_facet(MetadataPredicatesFacet::FACET)
.as_ref(),
index as usize,
)
}
#[crate::polydat_node(category = RealData)]
fn metadata_content_count(handle: Facet<MetadataContentFacet>) -> u64 {
match handle.resolve_facet(MetadataContentFacet::FACET).as_ref() {
DatasetHandle::Generic(d) => d.count as u64,
_ => 0,
}
}
fn generic_count_typed(h: &DatasetHandle, value: i64) -> u64 {
match h {
DatasetHandle::Generic(d) => {
(0..d.count).filter(|i| d.get_scalar(*i) == value).count() as u64
}
_ => 0,
}
}
#[crate::polydat_node(category = RealData)]
fn metadata_count_of(handle: Facet<MetadataContentFacet>, value: i64) -> u64 {
generic_count_typed(
handle.resolve_facet(MetadataContentFacet::FACET).as_ref(),
value,
)
}
#[crate::polydat_node(category = RealData)]
fn predicate_count_of(handle: Facet<MetadataPredicatesFacet>, value: i64) -> u64 {
generic_count_typed(
handle
.resolve_facet(MetadataPredicatesFacet::FACET)
.as_ref(),
value,
)
}
#[cfg(feature = "vectordata")]
fn vectordata_sugar(
source_name: &str,
constructor: &crate::dsl::ast::Expr,
) -> Result<Option<crate::dsl::cursor_sugar::CursorSugar>, String> {
use crate::dsl::ast::{Arg, CallExpr, Expr};
use crate::dsl::compile::positional_str_lit;
use crate::dsl::cursor_sugar::{AuxBinding, CursorSugar};
let Expr::Call(call) = constructor else {
return Ok(None);
};
let (dataset, profile, facet) = match call.func.as_str() {
"vectordata_source" => {
let d = positional_str_lit(call.args.first()).ok_or_else(|| format!(
"cursor '{source_name}': vectordata_source(dataset, profile, facet) — first arg must be a string literal"
))?;
let p = positional_str_lit(call.args.get(1)).ok_or_else(|| format!(
"cursor '{source_name}': vectordata_source(dataset, profile, facet) — second arg must be a string literal"
))?;
let f = positional_str_lit(call.args.get(2)).ok_or_else(|| format!(
"cursor '{source_name}': vectordata_source(dataset, profile, facet) — third arg must be a string literal (\"base\" or \"query\")"
))?;
(d, p, f)
}
"vectordata_base" | "vectordata_query" => {
let f = call.func.strip_prefix("vectordata_").unwrap().to_string();
let d = positional_str_lit(call.args.first()).ok_or_else(|| format!(
"cursor '{source_name}': {}(dataset, profile) — first arg must be a string literal",
call.func,
))?;
let p = positional_str_lit(call.args.get(1)).ok_or_else(|| format!(
"cursor '{source_name}': {}(dataset, profile) — second arg must be a string literal",
call.func,
))?;
(d, p, f)
}
_ => return Ok(None),
};
if facet != "base" && facet != "query" {
return Err(format!(
"cursor '{source_name}': vectordata facet must be \"base\" or \"query\", got \"{facet}\""
));
}
let (count_func, vector_func) = match facet.as_str() {
"base" => ("vector_count", "vector_at"),
"query" => ("query_count", "query_vector_at"),
_ => unreachable!(),
};
let combined = format!("{dataset}:{profile}");
let span = call.span;
let lit = |s: String| Expr::StringLit(s, span);
let positional = |e: Expr| Arg::Positional(e);
let effective_constructor = Expr::Call(CallExpr {
func: "range".into(),
args: vec![
positional(Expr::IntLit(0, span)),
positional(Expr::Call(CallExpr {
func: count_func.into(),
args: vec![positional(lit(combined.clone()))],
span,
})),
],
span,
});
let prebuffer_binding = AuxBinding {
name: format!("__{source_name}_prebuffer"),
value: Expr::Call(CallExpr {
func: "dataset_prebuffer".into(),
args: vec![positional(lit(combined.clone()))],
span,
}),
projection: None,
};
let vector_binding = AuxBinding {
name: format!("{source_name}__vector"),
value: Expr::Call(CallExpr {
func: vector_func.into(),
args: vec![
positional(lit(combined)),
positional(Expr::Ident(format!("{source_name}__ordinal"), span)),
],
span,
}),
projection: Some(("vector".into(), crate::ast::PortType::VecF32)),
};
Ok(Some(CursorSugar {
effective_constructor,
aux_bindings: vec![prebuffer_binding, vector_binding],
}))
}
#[cfg(feature = "vectordata")]
inventory::submit! {
crate::dsl::cursor_sugar::CursorSugarRegistration {
handler: vectordata_sugar,
name: "vectordata",
}
}
#[cfg(test)]
mod count_of_tests {
use super::*;
use crate::ast::{PolydatNode, PortType, Slot};
use crate::dsl::registry::DefaultResolver;
#[test]
fn the_count_of_nodes_take_a_handle_and_a_signed_value() {
for node in [
Box::new(MetadataCountOf::new()) as Box<dyn PolydatNode>,
Box::new(PredicateCountOf::new()) as Box<dyn PolydatNode>,
] {
let meta = node.meta();
assert_eq!(meta.outs.len(), 1);
assert_eq!(meta.outs[0].typ, PortType::U64, "{}", meta.name);
assert_eq!(meta.ins.len(), 2, "{}", meta.name);
let Slot::Wire(handle) = &meta.ins[0] else {
panic!("{}: the first input is a wire", meta.name)
};
let Slot::Wire(value) = &meta.ins[1] else {
panic!("{}: the second input is a wire", meta.name)
};
assert_eq!(handle.typ, PortType::Handle, "{}", meta.name);
assert_eq!(value.typ, PortType::I64, "{}", meta.name);
}
}
#[test]
fn a_handle_of_another_shape_counts_nothing() {
let group_shaped = DatasetHandle::Prebuffered {
group: match load_dataset_group("nonexistent:profile") {
Ok(g) => g,
Err(_) => return,
},
source: "nonexistent:profile".into(),
};
assert_eq!(generic_count_typed(&group_shaped, 1), 0);
}
#[test]
fn both_names_are_registered_with_a_facet_resolver() {
for (name, facet) in [
("metadata_count_of", "metadata_content"),
("predicate_count_of", "metadata_predicates"),
] {
let sig = crate::dsl::registry::lookup(name)
.unwrap_or_else(|| panic!("{name} is registered"));
assert_eq!(sig.outputs, 1, "{name}");
assert_eq!(sig.params.len(), 2, "{name}");
match sig.default_resolver {
Some(DefaultResolver::Facet(f)) => assert_eq!(f, facet, "{name}"),
other => panic!("{name}: expected a facet resolver, got {other:?}"),
}
}
}
}