polydat_core/library/vectors.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Vector dataset access nodes via the `vectordata` crate.
5//!
6//! Each node takes a dataset source string (URL, local path, or
7//! catalog name) as a const parameter and loads the dataset handle at
8//! construction time.
9//!
10//! Source specifier formats:
11//! - `"dataset"` — catalog lookup, uses default profile
12//! - `"dataset:profile"` — catalog lookup with explicit profile
13//! - `"https://..."` — direct URL
14//! - `"/path/to/dir"` — local filesystem path
15//!
16//! ## Prebuffering
17//!
18//! Use `dataset_prebuffer("source")` to eagerly download all facets
19//! for a dataset before workload execution. After prebuffering, data
20//! access uses local mmap readers (zero HTTP overhead).
21//!
22//! ## Cache-aware loading
23//!
24//! For catalog-resolved datasets, the loader checks the local cache
25//! at `~/.cache/vectordata/<dataset>/` before issuing HTTP requests.
26//! Prebuffering populates this cache; subsequent loads are local.
27//!
28//! Feature-gated behind `vectordata`.
29//!
30//! ## Upstream API reference
31//!
32//! See the vectordata consumer API docs for the full dataset access,
33//! catalog, caching, and prebuffer model:
34//! <https://github.com/nosqlbench/vectordata-rs/blob/main/docs/sysref/02-api.md>
35
36use std::sync::{Arc, LazyLock};
37
38use crate::library::support::cache::OnceCache;
39
40use vectordata::TestDataGroup;
41use vectordata::TestDataView;
42use vectordata::catalog::resolver::Catalog;
43use vectordata::catalog::sources::CatalogSources;
44use vectordata::io::{VectorReader, VvecReader};
45
46/// Global cache for loaded dataset groups keyed by source string.
47/// Ensures each dataset is loaded exactly once regardless of how many
48/// node functions reference it. The race-free init pattern lives in
49/// [`crate::library::support::cache::OnceCache`].
50static DATASET_CACHE: LazyLock<OnceCache<String, Arc<TestDataGroup>>> =
51 LazyLock::new(OnceCache::new);
52
53/// Type-erased facet cache: (source, profile, facet) → Arc<dyn Any + Send + Sync>.
54/// Ensures each reader (of any element type) is opened exactly once
55/// and shared across all node instances that reference the same data.
56/// The concrete type inside is `Arc<UniformDataset<T>>` or `Arc<Ivvec32Dataset>`
57/// or `Arc<GenericFacetDataset>`. The race-free init pattern lives
58/// in [`OnceCache`].
59type FacetCache = OnceCache<(String, String, String), Arc<dyn std::any::Any + Send + Sync>>;
60static FACET_CACHE: LazyLock<FacetCache> = LazyLock::new(OnceCache::new);
61
62/// Whole-dataset prebuffer cache keyed by source string. Wraps
63/// [`do_dataset_prebuffer_inner`] so N fibers all calling
64/// `dataset_prebuffer("ds:profile")` at init time serialize on
65/// one OnceLock and share the resulting handle. Without this,
66/// each fiber's own kernel-init pass invokes the real prebuffer
67/// body and they all race into the manifest walk + per-facet
68/// download — exactly the thundering herd we hit. The inner
69/// caches (DATASET_CACHE, FACET_CACHE) only protect the
70/// group-resolve and per-facet-reader steps, not the outer
71/// "walk every facet and pull all chunks" work.
72static PREBUFFER_CACHE: LazyLock<OnceCache<String, Arc<DatasetHandle>>> =
73 LazyLock::new(OnceCache::new);
74
75// =================================================================
76// Dataset resolution — catalog-aware, cache-aware
77// =================================================================
78
79/// Parse a source specifier into (dataset_name, profile_name).
80///
81/// Supports `"dataset:profile"` syntax. If no colon is present,
82/// the profile defaults to `"default"`.
83fn parse_source_specifier(source: &str) -> (&str, &str) {
84 // Don't split on colon in URLs
85 if source.starts_with("http://") || source.starts_with("https://") {
86 return (source, "default");
87 }
88 if let Some(pos) = source.find(':') {
89 (&source[..pos], &source[pos + 1..])
90 } else {
91 (source, "default")
92 }
93}
94
95/// Run a synchronous body that may internally drive
96/// `reqwest::blocking` (and thus spin up a private tokio
97/// runtime per HTTP request). If we're sitting on an outer
98/// async runtime — we are, in every per-cycle and per-phase
99/// path — the inner runtime panics on drop with "Cannot drop a
100/// runtime in a context where blocking is not allowed".
101/// `block_in_place` parks the outer multi-thread worker for
102/// the duration of the call so the inner runtime sees a
103/// non-async drop context. Falls back to a direct call when no
104/// runtime is current (e.g. unit tests). Single helper used by
105/// all three vectordata-facing entry points
106/// ([`load_dataset_group`], [`load_uniform_facet`],
107/// [`GenericFacetDataset::load`]).
108fn run_blocking_io<R>(body: impl FnOnce() -> R) -> R {
109 // tokio rides the `vectordata` feature: without the crate there is
110 // no HTTP client to park a worker for.
111 #[cfg(feature = "vectordata")]
112 if tokio::runtime::Handle::try_current().is_ok() {
113 return tokio::task::block_in_place(body);
114 }
115 body()
116}
117
118/// Load a dataset group by name.
119///
120/// Uses the vectordata catalog API: `catalog.open(name)` handles
121/// catalog discovery, cache resolution, and download transparently.
122pub(crate) fn load_dataset_group(source: &str) -> Result<Arc<TestDataGroup>, String> {
123 let (dataset_name, _profile) = parse_source_specifier(source);
124 // Both keys (dataset name and full source spec) point at the
125 // same `Arc<TestDataGroup>`; the `OnceCache` slot for one is
126 // primed by the other on first hit.
127 DATASET_CACHE.get_or_init(dataset_name.to_string(), || {
128 run_blocking_io(|| {
129 let catalog = Catalog::of(&CatalogSources::new().configure_default());
130 catalog
131 .open(dataset_name)
132 .map(Arc::new)
133 .map_err(|e| format!("failed to load dataset '{dataset_name}': {e}"))
134 })
135 })
136}
137
138// =================================================================
139// Dataset handles — loaded once at node construction, shared via Arc
140// =================================================================
141
142/// Generic handle to a loaded uniform vector facet. Thread-safe, random-access.
143/// Supports any element type provided by the vectordata API (f32, f64,
144/// i32, i16, u8, i8, u16, u32, u64, i64, f16).
145pub(crate) struct UniformDataset<T: Send + Sync + 'static> {
146 reader: Arc<dyn VectorReader<T>>,
147 count: usize,
148 dim: usize,
149}
150
151/// Cache-aware loader for uniform vector facets.
152/// Returns a shared `Arc<UniformDataset<T>>`, creating and caching it
153/// on first access. Subsequent loads for the same (source, profile, facet)
154/// return the cached instance.
155fn load_uniform_facet<T: Send + Sync + 'static>(
156 source: &str,
157 profile: &str,
158 facet: &str,
159 open_fn: impl FnOnce(
160 &dyn TestDataView,
161 ) -> std::result::Result<Arc<dyn VectorReader<T>>, vectordata::Error>,
162) -> Result<Arc<UniformDataset<T>>, String> {
163 let key = (source.to_string(), profile.to_string(), facet.to_string());
164 let any = FACET_CACHE.get_or_init(key, || {
165 let group = load_dataset_group(source)?;
166 let view = group
167 .profile(profile)
168 .ok_or_else(|| format!("profile '{profile}' not found in '{source}'"))?;
169 // Audit: log the open *before* `open_fn` runs so the
170 // line appears even if the open errors. Inside
171 // `get_or_init`'s closure, this fires exactly once per
172 // (source, profile, facet) — the prior shape logged
173 // once per concurrent miss (storms of N for N fibers).
174 crate::library::support::audit::record_opened(source, profile, facet, "uniform");
175 let reader = run_blocking_io(|| open_fn(view.as_ref()))
176 .map_err(|e| format!("failed to access {facet} from '{source}': {e}"))?;
177 let count = reader.count();
178 let dim = reader.dim();
179 let arc: Arc<UniformDataset<T>> = Arc::new(UniformDataset { reader, count, dim });
180 Ok(arc as Arc<dyn std::any::Any + Send + Sync>)
181 })?;
182 any.downcast::<UniformDataset<T>>().map_err(|_| {
183 format!(
184 "facet cache type mismatch for '{source}:{profile}/{facet}' — \
185 this should be impossible; please file a bug."
186 )
187 })
188}
189
190// Type aliases for backward compatibility
191type F32Dataset = UniformDataset<f32>;
192type I32Dataset = UniformDataset<i32>;
193
194impl F32Dataset {
195 fn load(source: &str, profile: &str, facet: &str) -> Result<Arc<Self>, String> {
196 let facet_name = facet.to_string();
197 load_uniform_facet(source, profile, facet, move |view| {
198 match facet_name.as_str() {
199 "base" => view.base_vectors(),
200 "query" => view.query_vectors(),
201 "neighbor_distances" => view.neighbor_distances(),
202 "filtered_neighbor_distances" => view.prefiltered_neighbor_distances(),
203 other => Err(vectordata::Error::MissingFacet(format!(
204 "unknown f32 facet: '{other}'"
205 ))),
206 }
207 })
208 }
209}
210
211impl I32Dataset {
212 fn load(source: &str, profile: &str, facet: &str) -> Result<Arc<Self>, String> {
213 let facet_name = facet.to_string();
214 load_uniform_facet(source, profile, facet, move |view| {
215 match facet_name.as_str() {
216 "neighbor_indices" => view.neighbor_indices(),
217 "filtered_neighbor_indices" => view.prefiltered_neighbor_indices(),
218 other => Err(vectordata::Error::MissingFacet(format!(
219 "unknown i32 facet: '{other}'"
220 ))),
221 }
222 })
223 }
224}
225
226// =================================================================
227// Dataset handle — typed enum wrapping all per-facet dataset shapes
228// =================================================================
229//
230// Per type_system.md §1.8: `dataset_open(source, facet)` is
231// the resolver node; per-cycle accessors take a handle wire and
232// downcast to the concrete variant they expect. A single Value
233// variant `Value::Handle(Arc<dyn Any>)` carries any of the
234// concrete shapes — the enum below is what's actually inside the
235// Arc, and accessors `match` on the variant.
236//
237// One handle may be opened against multiple facets at the
238// `dataset_open` boundary (the user calls `dataset_open(spec,
239// "base")` vs. `dataset_open(spec, "query")` and gets two
240// distinct handles); the downstream accessor matches on whatever
241// variant came back.
242
243/// Typed wrapper around a resolved dataset facet or group. Held
244/// inside `Value::Handle` as `Arc<DatasetHandle>` and downcast by
245/// accessor nodes via [`Value::as_handle::<DatasetHandle>()`].
246/// Cloning a `Value::Handle` is one `Arc::clone` — one atomic
247/// increment, no allocation — which is the design contract this
248/// enum exists to satisfy.
249#[derive(Clone)]
250pub(crate) enum DatasetHandle {
251 /// Uniform `f32` vector facet (base, query, neighbor_distances,
252 /// filtered_neighbor_distances).
253 F32(Arc<F32Dataset>),
254 /// Uniform `i32` vector facet (neighbor_indices, filtered_neighbor_indices).
255 I32(Arc<I32Dataset>),
256 /// Variable-length `i32` facet (metadata_results).
257 Ivvec32(Arc<Ivvec32Dataset>),
258 /// Type-erased generic-typed scalar facet (metadata_content,
259 /// metadata_predicates, ...).
260 Generic(Arc<GenericFacetDataset>),
261 /// Dataset-group handle — used by group-level metadata
262 /// accessors (`dataset_profile_count`, `dataset_facets`,
263 /// `dataset_distance_function`, ...) that operate on the
264 /// `TestDataGroup` before any profile/facet is chosen.
265 Group(Arc<TestDataGroup>),
266 /// Prebuffered-and-resident dataset, returned by
267 /// `dataset_prebuffer(source)`. Carries both the group and
268 /// the source spec: a facet accessor opens its facet from
269 /// `<source>:<profile>` (`query_vector_at(prebuffered, q)`
270 /// reads the `query` facet), and a group accessor reads the
271 /// group. [`DatasetHandle::resolve_facet`] and
272 /// [`DatasetHandle::resolve_group`] are those two readings.
273 ///
274 /// The `group` field also keeps the prebuffered `TestDataGroup`
275 /// alive for the duration of the handle — vectordata's
276 /// internal storage cache is keyed off the group instance,
277 /// so dropping the group prematurely would force per-facet
278 /// readers to re-open against transport.
279 Prebuffered {
280 group: Arc<TestDataGroup>,
281 source: String,
282 },
283}
284
285impl DatasetHandle {
286 fn open(source: &str, facet: &str) -> Result<Self, String> {
287 let (_, profile) = parse_source_specifier(source);
288 match facet {
289 "base" | "query" | "neighbor_distances" | "filtered_neighbor_distances" => {
290 F32Dataset::load(source, profile, facet).map(DatasetHandle::F32)
291 }
292 "neighbor_indices" | "filtered_neighbor_indices" => {
293 I32Dataset::load(source, profile, facet).map(DatasetHandle::I32)
294 }
295 "metadata_results" => Ivvec32Dataset::load(source, profile).map(DatasetHandle::Ivvec32),
296 // Anything else routes through GenericFacetDataset (typed
297 // scalar reader), which covers metadata_content,
298 // metadata_predicates, and any future scalar facet.
299 _ => GenericFacetDataset::load(source, profile, facet).map(DatasetHandle::Generic),
300 }
301 }
302
303 fn open_group(source: &str) -> Result<Self, String> {
304 load_dataset_group(source).map(DatasetHandle::Group)
305 }
306}
307
308/// `dataset_open(source: str, facet: str) -> Handle`
309///
310/// The single resolver node. Provenance follows its two wire
311/// inputs; when both are scope-extern constants (the iter-var
312/// case), this evaluates exactly once at iteration entry and
313/// stays cached for every cycle in that iteration. When source
314/// or facet is cycle-time, it re-evaluates accordingly.
315///
316/// All per-cycle accessors take the resulting handle on a wire
317/// — the catalog/HTTP/mmap path is never on the cycle hot path.
318///
319/// A resolve failure panics with the underlying error rather than
320/// yielding `Value::None`, so a catalog miss cannot degrade into a
321/// garbage value downstream: the engine's `enrich_eval_panic` adds
322/// node provenance, `eval_const_expr` traps the panic into
323/// `Result::Err`, and `evaluate_spec` propagates it as a
324/// clause-level diagnostic.
325#[crate::polydat_node(category = RealData)]
326fn dataset_open(source: &str, facet: &str) -> Arc<DatasetHandle> {
327 match DatasetHandle::open(source, facet) {
328 Ok(h) => Arc::new(h),
329 Err(e) => {
330 let msg = format!("dataset_open: failed to resolve '{source}' facet='{facet}': {e}");
331 crate::library::support::audit::error(&msg);
332 panic!("{msg}");
333 }
334 }
335}
336
337/// `dataset_group_open(source: str) -> Handle`
338///
339/// Group-level resolver. Returns a handle wrapping `Arc<TestDataGroup>`
340/// — used by group-level metadata accessors (`dataset_profile_count`,
341/// `dataset_facets`, ...) that operate on the dataset as a whole
342/// before any profile/facet is selected.
343///
344/// A resolve failure panics, as in `dataset_open` and for the
345/// same reason.
346#[crate::polydat_node(category = RealData)]
347fn dataset_group_open(source: &str) -> Arc<DatasetHandle> {
348 match DatasetHandle::open_group(source) {
349 Ok(h) => Arc::new(h),
350 Err(e) => {
351 let msg = format!("dataset_group_open: failed to resolve '{source}': {e}");
352 crate::library::support::audit::error(&msg);
353 panic!("{msg}");
354 }
355 }
356}
357
358// =================================================================
359// Base vector nodes
360// =================================================================
361
362// All indexed-accessor nodes share the same shape: a handle wire
363// and an `index` wire (u64). A source string in the handle
364// position is promoted to `dataset_open(source, "<facet>")` for
365// the facet the accessor names, and the facet reader is opened
366// once per (source, profile, facet) through `FACET_CACHE`, so
367// the per-cycle work is the facet read alone.
368
369// =================================================================
370// Per-cycle indexed accessors
371// =================================================================
372
373/// Resolve a Prebuffered handle to a specific-facet handle by
374/// opening the named facet via the existing FACET_CACHE path
375/// (which is `OnceCache`-backed — concurrent callers serialize
376/// per (source, profile, facet)). For non-Prebuffered handles
377/// this is a borrow-through; the macro callers use the result as
378/// `&DatasetHandle` regardless of which arm fires.
379///
380/// Lives on `DatasetHandle` so the resolution logic — and the
381/// "what does Prebuffered mean to a per-facet accessor" question
382/// — sits next to the variant declaration.
383impl DatasetHandle {
384 fn resolve_facet<'a>(&'a self, facet: &str) -> std::borrow::Cow<'a, DatasetHandle> {
385 match self {
386 DatasetHandle::Prebuffered { source, .. } => match DatasetHandle::open(source, facet) {
387 Ok(opened) => std::borrow::Cow::Owned(opened),
388 Err(e) => panic!(
389 "DatasetHandle::resolve_facet: failed to open \
390 '{facet}' from prebuffered '{source}': {e}"
391 ),
392 },
393 _ => std::borrow::Cow::Borrowed(self),
394 }
395 }
396
397 /// The dataset group a group accessor reads: the group of a
398 /// `dataset_group_open` handle, or the group a prebuffered handle
399 /// carries. A facet handle has no group, and asking one for it
400 /// panics with the handle's kind.
401 fn resolve_group(&self) -> &TestDataGroup {
402 match self {
403 DatasetHandle::Group(g) | DatasetHandle::Prebuffered { group: g, .. } => g.as_ref(),
404 other => panic!(
405 "expected a group or prebuffered handle, got {}",
406 dataset_handle_kind(other)
407 ),
408 }
409 }
410}
411
412/// The facet each accessor opens, as a type.
413///
414/// A node states its facet twice — in the auto-resolver that splices
415/// `dataset_open(_, "<facet>")` when the source is a string, and in
416/// the `resolve_facet` call that follows a prebuffered handle — and
417/// both readings come from here, so the two cannot drift apart.
418macro_rules! facet_kind {
419 ($name:ident, $facet:literal) => {
420 /// The `
421 #[doc = $facet]
422 /// ` facet of a dataset.
423 pub struct $name;
424 impl $name {
425 /// The facet's name, as `dataset_open` takes it.
426 pub const FACET: &'static str = $facet;
427 }
428 impl crate::derive_support::ResolverKind for $name {
429 const RESOLVER: crate::dsl::registry::DefaultResolver =
430 crate::dsl::registry::DefaultResolver::Facet($facet);
431 }
432 };
433}
434
435facet_kind!(BaseFacet, "base");
436facet_kind!(QueryFacet, "query");
437facet_kind!(NeighborIndicesFacet, "neighbor_indices");
438facet_kind!(NeighborDistancesFacet, "neighbor_distances");
439facet_kind!(FilteredNeighborIndicesFacet, "filtered_neighbor_indices");
440facet_kind!(
441 FilteredNeighborDistancesFacet,
442 "filtered_neighbor_distances"
443);
444facet_kind!(MetadataResultsFacet, "metadata_results");
445facet_kind!(MetadataContentFacet, "metadata_content");
446facet_kind!(MetadataPredicatesFacet, "metadata_predicates");
447
448/// A handle wire that carries its facet in its type.
449type Facet<R> = crate::derive_support::Resolved<R, DatasetHandle>;
450
451/// A handle wire on the dataset group rather than one of its facets.
452type Group = crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>;
453
454/// The record count of a facet handle, opening `facet` first when the
455/// handle is prebuffered. `who` names the caller in the diagnostics,
456/// which is what the reader needs to see when a workload asks a
457/// facet for a count it cannot give.
458fn facet_record_count(h: &DatasetHandle, who: &str, facet: &str) -> u64 {
459 match h.resolve_facet(facet).as_ref() {
460 DatasetHandle::F32(d) => d.count as u64,
461 DatasetHandle::I32(d) => d.count as u64,
462 DatasetHandle::Ivvec32(d) => d.count as u64,
463 DatasetHandle::Generic(d) => d.count as u64,
464 other => panic!(
465 "{who}: expected a facet handle, got {}",
466 dataset_handle_kind(other)
467 ),
468 }
469}
470
471/// The `f32` record at `index`, as the typed vector it is.
472fn f32_vec_typed(h: &DatasetHandle, index: usize) -> crate::ast::SliceArc<f32> {
473 match h {
474 DatasetHandle::F32(d) => slice_arc_from_uniform(d, index),
475 other => panic!(
476 "expected F32 dataset handle, got {}",
477 dataset_handle_kind(other)
478 ),
479 }
480}
481
482/// The `i32` record at `index`, as the typed vector it is.
483fn i32_vec_typed(h: &DatasetHandle, index: usize) -> crate::ast::SliceArc<i32> {
484 match h {
485 DatasetHandle::I32(d) => slice_arc_from_uniform(d, index),
486 other => panic!(
487 "expected I32 dataset handle, got {}",
488 dataset_handle_kind(other)
489 ),
490 }
491}
492
493fn ivvec32_vec_typed(h: &DatasetHandle, index: usize) -> crate::ast::SliceArc<i32> {
494 match h {
495 // Variable-length records — vectordata's trait doesn't expose
496 // a zero-copy slice path for these (per-record dim is read
497 // from the file), so we always allocate. One Vec<i32> per
498 // cycle; same bound as the upstream trait API.
499 DatasetHandle::Ivvec32(d) => {
500 if d.count == 0 {
501 return crate::ast::SliceArc::from_vec(Vec::<i32>::new());
502 }
503 let v = d.reader.get(index % d.count).unwrap_or_default();
504 crate::ast::SliceArc::from_vec(v)
505 }
506 other => panic!(
507 "expected Ivvec32 dataset handle, got {}",
508 dataset_handle_kind(other)
509 ),
510 }
511}
512
513/// Build a [`SliceArc<T>`] from any uniform-stride dataset at
514/// `index`. One generic helper serves every element type — element type is
515/// erased by the upstream `VectorReader<T>` trait, and `SliceArc`
516/// is element-type-generic, so one body suffices for `f32`,
517/// `i32`, and any other element type the trait supports.
518///
519/// Tries the zero-copy `VectorReader::get_slice` path first —
520/// that returns a borrow into the mmap'd file pages, kept alive
521/// by the `Arc<UniformDataset<T>>` that we move into the
522/// `SliceArc` as owner. No allocation, no element decoding, no
523/// copy.
524///
525/// Falls back to `VectorReader::get(idx) -> Vec<T>` for readers
526/// that don't support zero-copy (HTTP-backed, or merkle-cached
527/// storage that hasn't been promoted to mmap yet); that path
528/// allocates one `Vec<T>` per cycle.
529fn slice_arc_from_uniform<T>(d: &Arc<UniformDataset<T>>, index: usize) -> crate::ast::SliceArc<T>
530where
531 T: Send + Sync + Copy + 'static,
532{
533 if d.count == 0 {
534 return crate::ast::SliceArc::from_vec(Vec::<T>::new());
535 }
536 let idx = index % d.count;
537 if let Some(slice) = d.reader.get_slice(idx) {
538 // SAFETY: the slice points into the mmap pages owned by
539 // `d`'s reader; cloning `d` into the `SliceArc`'s owner
540 // keeps the mmap alive for as long as the slice is held.
541 let ptr_len = (slice.as_ptr(), slice.len());
542 let owner = d.clone();
543 let owner_dyn: Arc<dyn std::any::Any + Send + Sync> = owner;
544 return unsafe {
545 crate::ast::SliceArc::from_borrowed(
546 owner_dyn,
547 std::slice::from_raw_parts(ptr_len.0, ptr_len.1),
548 )
549 };
550 }
551 crate::ast::SliceArc::from_vec(d.reader.get(idx).unwrap_or_default())
552}
553
554fn dataset_handle_kind(h: &DatasetHandle) -> &'static str {
555 match h {
556 DatasetHandle::F32(_) => "F32",
557 DatasetHandle::I32(_) => "I32",
558 DatasetHandle::Ivvec32(_) => "Ivvec32",
559 DatasetHandle::Generic(_) => "Generic",
560 DatasetHandle::Group(_) => "Group",
561 DatasetHandle::Prebuffered { .. } => "Prebuffered",
562 }
563}
564
565/// Access an `f32` vector by index, returning a typed `VecF32`.
566/// Reads the base facet of a source string or prebuffered handle,
567/// and whatever facet an opened F32 handle names.
568///
569/// Signature: `vector_at(handle, index: u64) -> VecF32`
570#[crate::polydat_node(category = RealData)]
571fn vector_at(handle: Facet<BaseFacet>, index: u64) -> crate::ast::SliceArc<f32> {
572 f32_vec_typed(
573 handle.resolve_facet(BaseFacet::FACET).as_ref(),
574 index as usize,
575 )
576}
577
578/// Access a query vector by index. Reads the query facet of a
579/// source string or prebuffered handle, as `vector_at` reads the
580/// base facet.
581///
582/// Signature: `query_vector_at(handle, index: u64) -> VecF32`
583#[crate::polydat_node(category = RealData)]
584fn query_vector_at(handle: Facet<QueryFacet>, index: u64) -> crate::ast::SliceArc<f32> {
585 f32_vec_typed(
586 handle.resolve_facet(QueryFacet::FACET).as_ref(),
587 index as usize,
588 )
589}
590
591/// Access ground-truth neighbor indices for a query. Expects an
592/// I32 handle.
593///
594/// Signature: `neighbor_indices_at(handle, index: u64) -> VecI32`
595#[crate::polydat_node(category = RealData)]
596fn neighbor_indices_at(
597 handle: Facet<NeighborIndicesFacet>,
598 index: u64,
599) -> crate::ast::SliceArc<i32> {
600 i32_vec_typed(
601 handle.resolve_facet(NeighborIndicesFacet::FACET).as_ref(),
602 index as usize,
603 )
604}
605
606/// Access ground-truth neighbor distances for a query. Expects an
607/// F32 handle.
608///
609/// Signature: `neighbor_distances_at(handle, index: u64) -> VecF32`
610#[crate::polydat_node(category = RealData)]
611fn neighbor_distances_at(
612 handle: Facet<NeighborDistancesFacet>,
613 index: u64,
614) -> crate::ast::SliceArc<f32> {
615 f32_vec_typed(
616 handle.resolve_facet(NeighborDistancesFacet::FACET).as_ref(),
617 index as usize,
618 )
619}
620
621/// Access filtered ground-truth neighbor indices. Expects an I32 handle.
622#[crate::polydat_node(category = RealData)]
623fn filtered_neighbor_indices_at(
624 handle: Facet<FilteredNeighborIndicesFacet>,
625 index: u64,
626) -> crate::ast::SliceArc<i32> {
627 i32_vec_typed(
628 handle
629 .resolve_facet(FilteredNeighborIndicesFacet::FACET)
630 .as_ref(),
631 index as usize,
632 )
633}
634
635/// Access filtered ground-truth neighbor distances. Expects an F32 handle.
636#[crate::polydat_node(category = RealData)]
637fn filtered_neighbor_distances_at(
638 handle: Facet<FilteredNeighborDistancesFacet>,
639 index: u64,
640) -> crate::ast::SliceArc<f32> {
641 f32_vec_typed(
642 handle
643 .resolve_facet(FilteredNeighborDistancesFacet::FACET)
644 .as_ref(),
645 index as usize,
646 )
647}
648
649// =================================================================
650// Metadata nodes (constant per dataset)
651// =================================================================
652
653// =================================================================
654// Per-handle metadata nodes
655// =================================================================
656//
657// Take `(handle: Handle)` and read shape/metadata directly from the
658// resolved dataset. Provenance bounded by the handle, which itself
659// is bounded by the resolver's externs — so these collapse to a
660// single per-iteration eval chain via the standard provenance
661// caching. No per-cycle work.
662
663/// Return the dimensionality (`f32` count per record) of a
664/// vector facet handle.
665///
666/// Signature: `vector_dim(handle) -> (u64)`
667#[crate::polydat_node(category = RealData)]
668fn vector_dim(handle: Facet<BaseFacet>) -> u64 {
669 match handle.resolve_facet(BaseFacet::FACET).as_ref() {
670 DatasetHandle::F32(d) => d.dim as u64,
671 DatasetHandle::I32(d) => d.dim as u64,
672 _ => 0,
673 }
674}
675
676/// Return the count of records in the facet a handle was opened
677/// against. This is the canonical "how many vectors / queries /
678/// neighbor-rows" accessor — `vector_count(base_handle)` for
679/// base vectors, `vector_count(query_handle)` for query
680/// vectors, etc. A source string or prebuffered handle counts
681/// the base facet.
682///
683/// Signature: `vector_count(handle) -> (u64)`
684#[crate::polydat_node(category = RealData)]
685fn vector_count(handle: Facet<BaseFacet>) -> u64 {
686 facet_record_count(&handle, "vector_count", BaseFacet::FACET)
687}
688
689/// Return the count of query vectors. Counts the query facet of a
690/// source string or prebuffered handle, as `vector_count` counts
691/// the base facet.
692///
693/// Signature: `query_count(handle) -> (u64)`
694#[crate::polydat_node(category = RealData)]
695fn query_count(handle: Facet<QueryFacet>) -> u64 {
696 // `query_count(prebuffered)` is a legitimate idiom in the
697 // pvs_query workload (`cursor q = range(0, query_count(...))`),
698 // so the prebuffered case resolves the query facet and reads
699 // its count, through the same OnceCache-backed open the
700 // per-cycle accessors use.
701 facet_record_count(&handle, "query_count", QueryFacet::FACET)
702}
703
704/// Return the per-record neighbor count (k) for an I32
705/// neighbor-indices handle.
706///
707/// Signature: `neighbor_count(handle) -> (u64)`
708#[crate::polydat_node(category = RealData)]
709fn neighbor_count(handle: Facet<NeighborIndicesFacet>) -> u64 {
710 match handle.resolve_facet(NeighborIndicesFacet::FACET).as_ref() {
711 DatasetHandle::I32(d) => d.dim as u64,
712 _ => 0,
713 }
714}
715
716/// Return the dataset's distance function (e.g., "COSINE",
717/// "EUCLIDEAN"). Operates on the dataset group; takes a group or
718/// prebuffered handle.
719///
720/// Signature: `dataset_distance_function(group) -> (String)`
721#[crate::polydat_node(category = RealData)]
722fn dataset_distance_function(group: Group) -> String {
723 let group = group.resolve_group();
724 let raw = group
725 .attribute("distance_function")
726 .and_then(|v| v.as_str())
727 .unwrap_or("unknown");
728 match raw.to_uppercase().as_str() {
729 "L2" | "EUCLIDEAN" => "EUCLIDEAN",
730 "L1" | "MANHATTAN" => "MANHATTAN",
731 "COSINE" => "COSINE",
732 "DOT_PRODUCT" | "DOTPRODUCT" | "DOT" | "INNER_PRODUCT" | "IP" => "DOT_PRODUCT",
733 _ => raw,
734 }
735 .to_string()
736}
737
738// =================================================================
739// Metadata facet nodes
740// =================================================================
741
742/// Handle to a loaded variable-length i32 facet (metadata_results).
743pub(crate) struct Ivvec32Dataset {
744 reader: Arc<dyn VvecReader<i32>>,
745 count: usize,
746}
747
748impl Ivvec32Dataset {
749 fn load(source: &str, profile: &str) -> Result<Arc<Self>, String> {
750 let key = (
751 source.to_string(),
752 profile.to_string(),
753 "metadata_results".to_string(),
754 );
755 let any = FACET_CACHE.get_or_init(key, || {
756 let group = load_dataset_group(source)?;
757 let view = group
758 .profile(profile)
759 .ok_or_else(|| format!("profile '{profile}' not found in '{source}'"))?;
760 crate::library::support::audit::record_opened(
761 source,
762 profile,
763 "metadata_results",
764 "ivvec32",
765 );
766 let reader = run_blocking_io(|| view.metadata_results())
767 .map_err(|e| format!("failed to access metadata_results from '{source}': {e}"))?;
768 let count = reader.count();
769 let arc: Arc<Self> = Arc::new(Self { reader, count });
770 Ok(arc as Arc<dyn std::any::Any + Send + Sync>)
771 })?;
772 any.downcast::<Self>().map_err(|_| {
773 format!("facet cache type mismatch for '{source}:{profile}/metadata_results'")
774 })
775 }
776}
777
778/// Access metadata indices (variable-length matching base ordinals
779/// per query) at an index. Expects an Ivvec32 handle.
780///
781/// Signature: `metadata_results_at(handle, index: u64) -> VecI32`
782#[crate::polydat_node(category = RealData)]
783fn metadata_results_at(
784 handle: Facet<MetadataResultsFacet>,
785 index: u64,
786) -> crate::ast::SliceArc<i32> {
787 ivvec32_vec_typed(
788 handle.resolve_facet(MetadataResultsFacet::FACET).as_ref(),
789 index as usize,
790 )
791}
792
793/// Return the length of one metadata_results record without
794/// loading the data (reads only the 4-byte header). Expects an
795/// Ivvec32 handle.
796///
797/// Signature: `metadata_results_len_at(handle, index: u64) -> (u64)`
798#[crate::polydat_node(category = RealData)]
799fn metadata_results_len_at(handle: Facet<MetadataResultsFacet>, index: u64) -> u64 {
800 match handle.resolve_facet(MetadataResultsFacet::FACET).as_ref() {
801 DatasetHandle::Ivvec32(d) if d.count > 0 => {
802 d.reader.dim_at(index as usize % d.count).unwrap_or(0) as u64
803 }
804 _ => 0,
805 }
806}
807
808/// Return the metadata indices count (number of predicate result
809/// sets). Expects an Ivvec32 handle.
810#[crate::polydat_node(category = RealData)]
811fn metadata_results_count(handle: Facet<MetadataResultsFacet>) -> u64 {
812 match handle.resolve_facet(MetadataResultsFacet::FACET).as_ref() {
813 DatasetHandle::Ivvec32(d) => d.count as u64,
814 _ => 0,
815 }
816}
817
818/// Report which facets are available for the default profile of
819/// a dataset group as a comma-separated list. Expects a group or
820/// prebuffered handle.
821///
822/// Signature: `dataset_facets(group) -> (String)`
823#[crate::polydat_node(category = RealData)]
824fn dataset_facets(group: Group) -> String {
825 let group = group.resolve_group();
826 // Use the default profile — group-level callers want a
827 // dataset-wide manifest; specific profile facets come via
828 // `profile_facets(group, idx)`.
829 let names = group.profile_names();
830 if let Some(first) = names.first()
831 && let Some(view) = group.profile(first)
832 {
833 let manifest = view.facet_manifest();
834 let mut names: Vec<&String> = manifest.keys().collect();
835 names.sort();
836 return names
837 .iter()
838 .map(|s| s.as_str())
839 .collect::<Vec<_>>()
840 .join(", ");
841 }
842 String::new()
843}
844
845// =================================================================
846// Profile enumeration — discover and iterate over dataset profiles
847// =================================================================
848//
849// Group-level: take a Group handle, read the dataset's sorted profile
850// list. Sort order is canonical (by base_count via `profile_sort_by_size`)
851// and computed once per group via the underlying TestDataGroup.
852
853/// Total number of profiles in a dataset group.
854///
855/// Signature: `dataset_profile_count(group) -> (u64)`
856#[crate::polydat_node(category = RealData)]
857fn dataset_profile_count(group: Group) -> u64 {
858 group.resolve_group().profile_names().len() as u64
859}
860
861/// Comma-separated list of all profile names in canonical sort
862/// order (by base_count).
863///
864/// Signature: `dataset_profile_names(group) -> (String)`
865#[crate::polydat_node(category = RealData)]
866fn dataset_profile_names(group: Group) -> String {
867 group.resolve_group().profile_names().join(", ")
868}
869
870/// Return profile names matching a prefix, comma-separated.
871///
872/// Signature: `matching_profiles(group, prefix: str) -> (String)`
873///
874/// If prefix is empty, returns all profiles. Used by `for_each:`
875/// phase templates to discover profiles dynamically.
876///
877/// `group` declares its source-string auto-resolver via the
878/// `Resolved<GroupResolver, _>` marker wrapper — the macro reads
879/// `<Resolved<GroupResolver, DatasetHandle> as Wire>::RESOLVER` at
880/// codegen and emits the matching `FuncSig.default_resolver`
881/// (`DefaultResolver::Group`). The spliced `dataset_group_open` yields
882/// the canonical `Value::Handle(Arc<DatasetHandle>)`, so the wire is
883/// resolved as `DatasetHandle` and the group is taken via
884/// [`DatasetHandle::resolve_group`] (downcasting straight to
885/// `TestDataGroup` would fail — the handle is always the unified
886/// `DatasetHandle` enum).
887#[crate::polydat_node(category = RealData)]
888fn matching_profiles(
889 group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
890 prefix: &str,
891) -> String {
892 let group: &TestDataGroup = group.resolve_group();
893 let all = group.profile_names();
894 let mut matched: Vec<&str> = if prefix.is_empty() {
895 all.iter().map(|s| s.as_str()).collect()
896 } else {
897 all.iter()
898 .filter(|s| s.starts_with(prefix))
899 .map(|s| s.as_str())
900 .collect()
901 };
902 // Natural-order sort: alphabetic with numeric runs
903 // compared as numbers so `label_03` sorts before
904 // `label_10`. The upstream `group.profile_names()`
905 // orders by `base_count` (vectordata's choice — useful
906 // for index-based lookups), but `for_each` iteration
907 // wants stable, human-natural order so users see
908 // label_01, label_02, label_03 instead of whatever the
909 // size-sort happens to produce.
910 matched.sort_by(|a, b| natural_cmp(a, b));
911 matched.join(",")
912}
913
914/// Natural ordering: split each string into alternating text
915/// and numeric runs and compare run-by-run, comparing numeric
916/// runs as integers. Beats lexicographic on `label_03` vs
917/// `label_10` (lex: "10" < "3"; natural: 3 < 10).
918fn natural_cmp(a: &str, b: &str) -> std::cmp::Ordering {
919 let mut ai = a.chars().peekable();
920 let mut bi = b.chars().peekable();
921 loop {
922 match (ai.peek().copied(), bi.peek().copied()) {
923 (None, None) => return std::cmp::Ordering::Equal,
924 (None, _) => return std::cmp::Ordering::Less,
925 (_, None) => return std::cmp::Ordering::Greater,
926 (Some(ac), Some(bc)) => {
927 if ac.is_ascii_digit() && bc.is_ascii_digit() {
928 let mut na: u64 = 0;
929 while let Some(c) = ai.peek().copied()
930 && c.is_ascii_digit()
931 {
932 na = na
933 .saturating_mul(10)
934 .saturating_add((c as u8 - b'0') as u64);
935 ai.next();
936 }
937 let mut nb: u64 = 0;
938 while let Some(c) = bi.peek().copied()
939 && c.is_ascii_digit()
940 {
941 nb = nb
942 .saturating_mul(10)
943 .saturating_add((c as u8 - b'0') as u64);
944 bi.next();
945 }
946 match na.cmp(&nb) {
947 std::cmp::Ordering::Equal => continue,
948 non_eq => return non_eq,
949 }
950 } else {
951 match ac.cmp(&bc) {
952 std::cmp::Ordering::Equal => {
953 ai.next();
954 bi.next();
955 }
956 non_eq => return non_eq,
957 }
958 }
959 }
960 }
961 }
962}
963
964/// Look up a profile name by index from the canonical sorted list.
965///
966/// Signature: `dataset_profile_name_at(group, index: u64) -> (String)`
967///
968/// Index wraps modulo the number of profiles.
969#[crate::polydat_node(category = RealData)]
970fn dataset_profile_name_at(
971 group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
972 index: u64,
973) -> String {
974 let group: &TestDataGroup = group.resolve_group();
975 let names = group.profile_names();
976 if names.is_empty() {
977 String::new()
978 } else {
979 names[(index as usize) % names.len()].clone()
980 }
981}
982
983/// Return the base vector count for the profile at a given index.
984///
985/// Signature: `profile_base_count(group, index: u64) -> (u64)`
986#[crate::polydat_node(category = RealData)]
987fn profile_base_count(
988 group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
989 index: u64,
990) -> u64 {
991 let group: &TestDataGroup = group.resolve_group();
992 let names = group.profile_names();
993 if names.is_empty() {
994 0
995 } else {
996 let name = &names[(index as usize) % names.len()];
997 group
998 .profile(name)
999 .and_then(|view| view.base_count())
1000 .unwrap_or(0)
1001 }
1002}
1003
1004/// Return the comma-separated facet list for the profile at a given index.
1005///
1006/// Signature: `profile_facets(group, index: u64) -> (String)`
1007#[crate::polydat_node(category = RealData)]
1008fn profile_facets(
1009 group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
1010 index: u64,
1011) -> String {
1012 let group: &TestDataGroup = group.resolve_group();
1013 let names = group.profile_names();
1014 if names.is_empty() {
1015 String::new()
1016 } else {
1017 let name = &names[(index as usize) % names.len()];
1018 match group.profile(name) {
1019 Some(view) => {
1020 let manifest = view.facet_manifest();
1021 let mut fnames: Vec<&String> = manifest.keys().collect();
1022 fnames.sort();
1023 fnames
1024 .iter()
1025 .map(|s| s.as_str())
1026 .collect::<Vec<_>>()
1027 .join(", ")
1028 }
1029 None => String::new(),
1030 }
1031 }
1032}
1033
1034/// Partition a dataset's vector space by its **profiles matching a
1035/// pattern**, treated as cumulative size tiers (a partition source,
1036/// cursor_partitions.md §7.1). One partition per masked profile, in canonical
1037/// (base-count-ascending) order: partition `k` spans
1038/// `[prev_masked_base_count, this_masked_base_count)` — exactly the
1039/// vectors added at that tier — so a sweep's "load only the increment
1040/// since the previously loaded set" is just the partition's
1041/// `[start_of(p), end_of(p))`, and a partition inherently knows its
1042/// start (no cross-iteration carry needed). `idx_of(p)` is the 0-based
1043/// masked position (pairs with `matching_profile_name_at` to address
1044/// the tier's own ground-truth facets); `count_of(p)` is the number of
1045/// masked tiers; `base_extent` is the largest masked tier's size.
1046///
1047/// `pattern` follows the literal / glob / regex promotion of
1048/// [`crate::library::support::pattern::compile_pattern`]; `*` (or any pattern that
1049/// matches all names) selects every profile.
1050///
1051/// Signature: `profile_partitions(group, pattern: str) -> (PartitionList)`
1052#[crate::polydat_node(category = RealData)]
1053fn profile_partitions(
1054 group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
1055 pattern: &str,
1056) -> crate::derive_support::Ext<crate::iteration::cursor_partition::PartitionList> {
1057 let group: &TestDataGroup = group.resolve_group();
1058 let parts = build_profile_partitions(group, pattern);
1059 crate::derive_support::Ext(crate::iteration::cursor_partition::PartitionList::new(
1060 parts,
1061 ))
1062}
1063
1064/// Build the cumulative size-tier partitions for the profiles of `group`
1065/// matching `pattern` (literal/glob/regex promotion). Shared by the
1066/// [`profile_partitions`] node and the comprehension-source desugaring
1067/// in `iteration::comprehension::eval` so a `for: "p in
1068/// profile_partitions(...)"` sweep produces identical partitions either
1069/// way. Partition `k` spans `[prev masked base_count, this base_count)`.
1070/// The masked profile size-tiers for `group` matching `pattern`, in
1071/// canonical (base-count-ascending) order: `(name, cumulative_base_count)`
1072/// for each matching profile WITH a non-zero base count. Profiles with no
1073/// base facet, or a zero count, are skipped — they would otherwise emit a
1074/// degenerate empty `[prev, prev)` partition (and a query against an empty
1075/// tier). Shared by [`build_profile_partitions`] and
1076/// [`matching_profile_name_at`] so a partition's `idx_of(p)` and the tier
1077/// name resolve against the SAME masked sequence.
1078fn masked_profile_tiers(group: &TestDataGroup, pattern: &str) -> Vec<(String, u64)> {
1079 let (re, _) = crate::library::support::pattern::compile_pattern(pattern)
1080 .unwrap_or_else(|e| panic!("profile pattern: {e}"));
1081 group
1082 .profile_names()
1083 .iter()
1084 .filter(|n| re.is_match(n))
1085 .filter_map(|n| {
1086 group
1087 .profile(n)
1088 .and_then(|v| v.base_count())
1089 .filter(|&c| c > 0)
1090 .map(|c| (n.clone(), c))
1091 })
1092 .collect()
1093}
1094
1095pub(crate) fn build_profile_partitions(
1096 group: &TestDataGroup,
1097 pattern: &str,
1098) -> Vec<crate::iteration::cursor_partition::Partition> {
1099 // Masked size-tiers (matching profiles with a non-zero base count),
1100 // canonical (ascending) order — skipping zero-count profiles avoids a
1101 // degenerate empty first partition.
1102 let masked: Vec<u64> = masked_profile_tiers(group, pattern)
1103 .into_iter()
1104 .map(|(_, c)| c)
1105 .collect();
1106 let base_extent = masked.last().copied().unwrap_or(0);
1107 let count = masked.len() as u64;
1108 let pct = |o: u64| {
1109 if base_extent == 0 {
1110 0.0
1111 } else {
1112 (o as f64 / base_extent as f64) * 100.0
1113 }
1114 };
1115 let mut parts: Vec<crate::iteration::cursor_partition::Partition> =
1116 Vec::with_capacity(masked.len());
1117 let mut prev: u64 = 0;
1118 for (k, &this) in masked.iter().enumerate() {
1119 parts.push(crate::iteration::cursor_partition::Partition {
1120 idx: k as u64,
1121 count,
1122 start_ord: prev,
1123 end_ord: this,
1124 start_pct: pct(prev),
1125 end_pct: pct(this),
1126 base_extent,
1127 });
1128 prev = this;
1129 }
1130 parts
1131}
1132
1133/// Name of the `index`-th profile **matching `pattern`**, in canonical
1134/// (base-count-ascending) order. Pairs with `profile_partitions`'s
1135/// masked-position `idx_of` so a sweep can prebuffer the active tier's
1136/// own ground-truth facets (`dataset_prebuffer(str_concat("ds:", name))`).
1137/// `pattern` follows the literal / glob / regex promotion; the index
1138/// wraps modulo the number of matching profiles; empty string if none
1139/// match.
1140///
1141/// Signature: `matching_profile_name_at(group, pattern: str, index: u64) -> (String)`
1142#[crate::polydat_node(category = RealData)]
1143fn matching_profile_name_at(
1144 group: crate::derive_support::Resolved<crate::derive_support::GroupResolver, DatasetHandle>,
1145 pattern: &str,
1146 index: u64,
1147) -> String {
1148 let group: &TestDataGroup = group.resolve_group();
1149 // Same masked sequence `profile_partitions` uses (matching profiles
1150 // with a non-zero base count, canonical order), so `idx_of(p)` from a
1151 // partition resolves to the right tier name.
1152 let tiers = masked_profile_tiers(group, pattern);
1153 if tiers.is_empty() {
1154 String::new()
1155 } else {
1156 tiers[(index as usize) % tiers.len()].0.clone()
1157 }
1158}
1159
1160// =================================================================
1161// Prebuffering — eagerly download dataset facets before workload run
1162// =================================================================
1163
1164/// Eagerly download all facets for a dataset profile into the
1165/// local cache, returning a dataset handle that every facet
1166/// accessor and every group accessor takes as its first
1167/// argument. After this returns, every subsequent facet read
1168/// served by [`vectordata::TestDataView`] hits the merkle-
1169/// verified mmap fast path with no further network traffic.
1170///
1171/// Signature: `dataset_prebuffer(source) -> Handle`
1172///
1173/// **Why a handle, not a count.** Returning a value the
1174/// downstream accessors *consume* makes prebuffer part of
1175/// the dataflow graph: DCE keeps the chain alive because
1176/// `vector_at(prebuffered, q)` needs `prebuffered` to be
1177/// resolvable, which forces evaluation. A binding such as
1178/// `init prebuffered = dataset_prebuffer(...)` whose result
1179/// nothing reads is pruned, so a workload threads the handle
1180/// through its accessors.
1181///
1182/// `source` is the canonical `dataset:profile` string. The
1183/// returned handle names the dataset rather than one facet: a
1184/// facet accessor opens its own facet of `source` from it
1185/// ([`DatasetHandle::resolve_facet`]), and a group accessor
1186/// (`dataset_facets`, `dataset_distance_function`, ...) reads
1187/// the dataset's group ([`DatasetHandle::resolve_group`]).
1188///
1189/// A dataset that does not resolve yields no handle, and the
1190/// resolve error goes to the audit log. A profile the dataset
1191/// lacks, or a facet whose download fails, is logged there too
1192/// and still yields the handle: prebuffering is a best-effort
1193/// warm-up, and an accessor that then reads a missing facet
1194/// fails with that facet's own error.
1195#[crate::polydat_node(category = RealData)]
1196fn dataset_prebuffer(source: &str) -> Option<Arc<dyn std::any::Any + Send + Sync>> {
1197 // The const binding contract (evaluation_model.md §"Const
1198 // Binding Contract") means this runs exactly once per scope
1199 // activation: the activation kernel pulls the binding and the host propagates the result to every
1200 // fiber. PREBUFFER_CACHE is belt-and-suspenders: it serializes
1201 // anyone who reaches here by another path, so a per-source
1202 // thundering herd is structurally impossible whatever the caller
1203 // does.
1204 match PREBUFFER_CACHE.get_or_init(source.to_string(), || {
1205 // Only the first concurrent caller for this source runs the
1206 // inner body — the audit "entered" event reflects that one
1207 // download, not per-caller noise.
1208 crate::library::support::audit::record_prebuffer_entered(source);
1209 do_dataset_prebuffer_inner(source)
1210 }) {
1211 Ok(handle) => Some(handle as Arc<dyn std::any::Any + Send + Sync>),
1212 Err(_) => None,
1213 }
1214}
1215
1216/// Inner body of [`dataset_prebuffer`]. Runs **at most once
1217/// per (source, process)** under the [`PREBUFFER_CACHE`] OnceLock.
1218fn do_dataset_prebuffer_inner(source: &str) -> Result<Arc<DatasetHandle>, String> {
1219 // Every exit after the group resolves returns a prebuffered
1220 // handle, so the downstream `*_at(prebuffered, q)` accessors can
1221 // resolve. A failed download still hands back the handle — the
1222 // accessors then read through the transport, and the operator
1223 // sees the prebuffer error in the audit log.
1224 let group = match load_dataset_group(source) {
1225 Ok(g) => g,
1226 Err(e) => {
1227 let msg = format!("dataset_prebuffer: cannot resolve '{source}': {e}");
1228 crate::library::support::audit::error(&msg);
1229 // No group → no handle to hand back. Sticky error in
1230 // the cache; downstream accessors will produce a
1231 // diagnostic when they fail to downcast.
1232 return Err(msg);
1233 }
1234 };
1235 let group_for_handle = group.clone();
1236 let (_, profile) = parse_source_specifier(source);
1237 let view = match group.profile(profile) {
1238 Some(v) => v,
1239 None => {
1240 crate::library::support::audit::error(&format!(
1241 "dataset_prebuffer: profile '{profile}' not found in '{source}'"
1242 ));
1243 return Ok(Arc::new(DatasetHandle::Prebuffered {
1244 group: group_for_handle,
1245 source: source.to_string(),
1246 }));
1247 }
1248 };
1249 // vectordata's default `prebuffer_all_with_progress`
1250 // walks the manifest and calls `FacetStorage::prebuffer`
1251 // per facet — works for both record-shaped (xvec) and
1252 // scalar (typed) facets after the storage-transport
1253 // refactor in vectordata 1.0.0.
1254 //
1255 // Audit instrumentation: log every facet the prebuffer
1256 // *covers* with its key (`source:profile/facet`). Compared
1257 // against the `vectordata: opened …` lines emitted by
1258 // [`load_uniform_facet`] / [`GenericFacetDataset::load`],
1259 // this surfaces any facet the workload reads at cycle time
1260 // that prebuffer did NOT pull — the typical cause of a
1261 // "still hitting HTTP after prebuffer" symptom.
1262 // Manual facet walk so we can hook the per-chunk
1263 // `DownloadProgress` callback (vectordata's
1264 // `view.prebuffer_all_with_progress` only fires its outer cb
1265 // once per facet *completion*; the per-chunk cb is exposed
1266 // via `FacetStorage::prebuffer_with_progress`). Every facet's
1267 // chunk-level progress is throttled per facet: ~1 Hz for the first
1268 // 10s of its download (fine detail while it ramps), then one line
1269 // per 10s for the long tail so a multi-gigabyte facet doesn't spew
1270 // thousands of lines into session.log; a final per-facet `covered`
1271 // line lands when the download finishes.
1272 let mut facet_count: u64 = 0;
1273 for (name, _descriptor) in view.facet_manifest() {
1274 // Skip facets with unrecognised element types (vectordata's
1275 // own default impl skips these — they're not data facets the
1276 // typed reader would touch).
1277 if view.facet_element_type(&name).is_err() {
1278 continue;
1279 }
1280
1281 let storage = match run_blocking_io(|| view.open_facet_storage(&name)) {
1282 Ok(s) => s,
1283 Err(e) => {
1284 crate::library::support::audit::warn(&format!(
1285 "dataset_prebuffer: open '{name}' for prebuffer failed: {e}"
1286 ));
1287 continue;
1288 }
1289 };
1290
1291 // Inline the closure at the call site so Rust's type
1292 // inference can read `&DownloadProgress` straight from
1293 // the trait bound on `prebuffer_with_progress`.
1294 // The `DownloadProgress` type is `pub(crate)` upstream
1295 // (vectordata 1.0.2) so we can't name it ourselves.
1296 let prebuf_source = source.to_string();
1297 let prebuf_profile = profile.to_string();
1298 let prebuf_facet = name.clone();
1299 let mut last_done: u64 = 0;
1300 let facet_start = std::time::Instant::now();
1301 let mut last_log_at = facet_start;
1302 // `prebuffer_with_progress` drives `reqwest::blocking`,
1303 // which spins up a private tokio runtime per request.
1304 // Without parking the outer worker via `block_in_place`,
1305 // dropping that inner runtime inside an async context
1306 // panics with "Cannot drop a runtime in a context where
1307 // blocking is not allowed". Same pattern as
1308 // `load_dataset_group` and `load_uniform_facet`.
1309 let prebuf_result = run_blocking_io(|| {
1310 storage.prebuffer_with_progress(|p| {
1311 // Per-facet throttle: ~1 Hz for the first 10s of this
1312 // facet's download, then once per 10s, so a gigabyte-sized
1313 // facet doesn't spew thousands of progress lines.
1314 let now = std::time::Instant::now();
1315 let interval =
1316 if now.duration_since(facet_start) < std::time::Duration::from_secs(10) {
1317 std::time::Duration::from_secs(1)
1318 } else {
1319 std::time::Duration::from_secs(10)
1320 };
1321 let since_last = now.duration_since(last_log_at);
1322 if since_last < interval {
1323 return;
1324 }
1325 last_log_at = now;
1326 let total_b = p.total_bytes();
1327 let done_b = p.downloaded_bytes();
1328 let total_c = p.total_chunks();
1329 let done_c = p.completed_chunks();
1330 let pct = p.fraction() * 100.0;
1331 // Throughput over the actual gap between log lines, so the
1332 // rate reads correctly whichever interval is in force.
1333 let delta_mb = (done_b.saturating_sub(last_done)) as f64 / (1024.0 * 1024.0);
1334 let rate_mb_s = delta_mb / since_last.as_secs_f64().max(0.001);
1335 last_done = done_b;
1336 crate::library::support::audit::info(&format!(
1337 "prebuffer: progress {prebuf_source}:{prebuf_profile}/{prebuf_facet} \
1338 {pct:5.1}% ({done_c}/{total_c} chunks, \
1339 {done_mb:.1}/{total_mb:.1} MB, {rate_mb_s:.1} MB/s)",
1340 done_mb = done_b as f64 / (1024.0 * 1024.0),
1341 total_mb = total_b as f64 / (1024.0 * 1024.0),
1342 ));
1343 })
1344 });
1345 if let Err(e) = prebuf_result {
1346 crate::library::support::audit::warn(&format!(
1347 "dataset_prebuffer: download error for '{source}' facet '{name}': {e}"
1348 ));
1349 continue;
1350 }
1351 facet_count = facet_count.saturating_add(1);
1352 crate::library::support::audit::record_prebuffered(source, profile, &name);
1353 }
1354 crate::library::support::audit::log_prebuffer_summary(source, profile, facet_count);
1355 Ok(Arc::new(DatasetHandle::Prebuffered {
1356 group: group_for_handle,
1357 source: source.to_string(),
1358 }))
1359}
1360
1361// =================================================================
1362// Generic facet access — type-aware scalar/vector readers
1363// =================================================================
1364
1365/// A type-erased facet reader that stores values as i64.
1366/// Uses `generic_view().open_facet_typed::<i64>()` which handles all
1367/// element types (u8, i32, etc.), caching, and local/remote access.
1368pub(crate) struct GenericFacetDataset {
1369 reader: vectordata::typed_access::TypedReader<i64>,
1370 count: usize,
1371}
1372
1373impl GenericFacetDataset {
1374 fn load(source: &str, profile: &str, facet: &str) -> Result<Arc<Self>, String> {
1375 let key = (source.to_string(), profile.to_string(), facet.to_string());
1376 let any = FACET_CACHE.get_or_init(key, || {
1377 let group = load_dataset_group(source)?;
1378 let gv = group
1379 .generic_view(profile)
1380 .ok_or_else(|| format!("profile '{profile}' not found in '{source}'"))?;
1381 crate::library::support::audit::record_opened(source, profile, facet, "generic-typed");
1382 let reader = run_blocking_io(|| gv.open_facet_typed::<i64>(facet))
1383 .map_err(|e| format!("failed to open {facet} from '{source}:{profile}': {e}"))?;
1384 let count = reader.count();
1385 let arc: Arc<Self> = Arc::new(Self { reader, count });
1386 Ok(arc as Arc<dyn std::any::Any + Send + Sync>)
1387 })?;
1388 any.downcast::<Self>()
1389 .map_err(|_| format!("facet cache type mismatch for '{source}:{profile}/{facet}'"))
1390 }
1391
1392 fn get_scalar(&self, index: usize) -> i64 {
1393 if self.count == 0 {
1394 return 0;
1395 }
1396 self.reader.get_value(index % self.count).unwrap_or(0)
1397 }
1398
1399 fn format_scalar(&self, index: usize) -> String {
1400 self.get_scalar(index).to_string()
1401 }
1402}
1403
1404// Generic-facet readers (typed scalar). Take a Generic handle and
1405// read the i64-cast value at the requested index.
1406/// The scalar of a Generic facet at `idx`, rendered.
1407fn generic_str_typed(h: &DatasetHandle, idx: usize) -> String {
1408 match h {
1409 DatasetHandle::Generic(d) => d.format_scalar(idx),
1410 _ => String::new(),
1411 }
1412}
1413
1414/// Access a metadata value per base ordinal. Expects a Generic
1415/// handle opened against the `metadata_content` facet.
1416///
1417/// Signature: `metadata_value_at(handle, index: u64) -> (String)`
1418#[crate::polydat_node(category = RealData)]
1419fn metadata_value_at(handle: Facet<MetadataContentFacet>, index: u64) -> String {
1420 generic_str_typed(
1421 handle.resolve_facet(MetadataContentFacet::FACET).as_ref(),
1422 index as usize,
1423 )
1424}
1425
1426/// Access a predicate value per query ordinal. Expects a Generic
1427/// handle opened against the `metadata_predicates` facet.
1428///
1429/// Signature: `predicate_value_at(handle, index: u64) -> (String)`
1430#[crate::polydat_node(category = RealData)]
1431fn predicate_value_at(handle: Facet<MetadataPredicatesFacet>, index: u64) -> String {
1432 generic_str_typed(
1433 handle
1434 .resolve_facet(MetadataPredicatesFacet::FACET)
1435 .as_ref(),
1436 index as usize,
1437 )
1438}
1439
1440/// Count of metadata content records. Expects a Generic handle.
1441///
1442/// Signature: `metadata_content_count(handle) -> (u64)`
1443#[crate::polydat_node(category = RealData)]
1444fn metadata_content_count(handle: Facet<MetadataContentFacet>) -> u64 {
1445 match handle.resolve_facet(MetadataContentFacet::FACET).as_ref() {
1446 DatasetHandle::Generic(d) => d.count as u64,
1447 _ => 0,
1448 }
1449}
1450
1451// ── Counting a facet's matching records ─────────────────────────────
1452
1453/// How many scalars of a Generic facet equal `value`.
1454///
1455/// A linear read of the facet, which is what the question is: the
1456/// reader holds the values and nothing indexes them by content. Both
1457/// callers below pass a handle and a value that are fixed for a scope,
1458/// so the count is computed once at scope init and reused, the way the
1459/// facet itself is loaded once.
1460/// How many scalars of a Generic facet equal `value`.
1461fn generic_count_typed(h: &DatasetHandle, value: i64) -> u64 {
1462 match h {
1463 DatasetHandle::Generic(d) => {
1464 (0..d.count).filter(|i| d.get_scalar(*i) == value).count() as u64
1465 }
1466 _ => 0,
1467 }
1468}
1469
1470/// How many records carry the metadata value `value`.
1471///
1472/// The per-label cardinality the recall audit compares against a
1473/// profile's declared count: not a counter advancing as records are
1474/// visited, but a property of the facet and the value, so it
1475/// replays and agrees on every engine and every fiber
1476/// (docs/design/runtime_model.md §9.1).
1477///
1478/// `metadata_value_at(handle, i)` formats the same scalar as text,
1479/// so a host comparing labels as strings and one counting them
1480/// here are reading one facet.
1481#[crate::polydat_node(category = RealData)]
1482fn metadata_count_of(handle: Facet<MetadataContentFacet>, value: i64) -> u64 {
1483 generic_count_typed(
1484 handle.resolve_facet(MetadataContentFacet::FACET).as_ref(),
1485 value,
1486 )
1487}
1488
1489/// How many queries carry the predicate value `value`.
1490///
1491/// The denominator of a per-predicate recall figure, and the bound
1492/// on a rank within one predicate's queries.
1493///
1494/// Signature: `predicate_count_of(handle, value: i64) -> (u64)`
1495#[crate::polydat_node(category = RealData)]
1496fn predicate_count_of(handle: Facet<MetadataPredicatesFacet>, value: i64) -> u64 {
1497 generic_count_typed(
1498 handle
1499 .resolve_facet(MetadataPredicatesFacet::FACET)
1500 .as_ref(),
1501 value,
1502 )
1503}
1504
1505// =========================================================================
1506// Cursor-sugar handlers
1507// =========================================================================
1508//
1509// Three sugar forms recognized by this module — none of them
1510// known to the core compiler, which dispatches generically
1511// through `dsl::cursor_sugar`. Each desugars to a synthetic
1512// `range(0, vector_count|query_count(...))` constructor plus
1513// auxiliary bindings:
1514//
1515// - `__<cursor>_prebuffer := dataset_prebuffer("ds:profile")`
1516// loaded once at init time so cycle-time accessor calls hit
1517// prefetched memory.
1518// - `<cursor>__vector := vector_at("ds:profile", <cursor>__ordinal)`
1519// (or `query_vector_at` for the query
1520// facet) — published as the cursor's `vector` projection so
1521// workloads can reference `<cursor>.vector`.
1522//
1523// Forms:
1524// - `vectordata_source(dataset, profile, facet)` — explicit facet
1525// - `vectordata_base(dataset, profile)` — facet = "base"
1526// - `vectordata_query(dataset, profile)` — facet = "query"
1527//
1528// Facet-specific projections like `metadata` / `ground_truth` /
1529// `predicate` stay explicit (the user writes
1530// `meta := metadata_value_at(<cursor>.ordinal, "ds:profile")` by
1531// hand) because their existence is dataset-conditional — not
1532// every dataset declares a metadata column or a predicate facet.
1533
1534#[cfg(feature = "vectordata")]
1535fn vectordata_sugar(
1536 source_name: &str,
1537 constructor: &crate::dsl::ast::Expr,
1538) -> Result<Option<crate::dsl::cursor_sugar::CursorSugar>, String> {
1539 use crate::dsl::ast::{Arg, CallExpr, Expr};
1540 use crate::dsl::compile::positional_str_lit;
1541 use crate::dsl::cursor_sugar::{AuxBinding, CursorSugar};
1542
1543 let Expr::Call(call) = constructor else {
1544 return Ok(None);
1545 };
1546
1547 let (dataset, profile, facet) = match call.func.as_str() {
1548 "vectordata_source" => {
1549 let d = positional_str_lit(call.args.first()).ok_or_else(|| format!(
1550 "cursor '{source_name}': vectordata_source(dataset, profile, facet) — first arg must be a string literal"
1551 ))?;
1552 let p = positional_str_lit(call.args.get(1)).ok_or_else(|| format!(
1553 "cursor '{source_name}': vectordata_source(dataset, profile, facet) — second arg must be a string literal"
1554 ))?;
1555 let f = positional_str_lit(call.args.get(2)).ok_or_else(|| format!(
1556 "cursor '{source_name}': vectordata_source(dataset, profile, facet) — third arg must be a string literal (\"base\" or \"query\")"
1557 ))?;
1558 (d, p, f)
1559 }
1560 "vectordata_base" | "vectordata_query" => {
1561 let f = call.func.strip_prefix("vectordata_").unwrap().to_string();
1562 let d = positional_str_lit(call.args.first()).ok_or_else(|| format!(
1563 "cursor '{source_name}': {}(dataset, profile) — first arg must be a string literal",
1564 call.func,
1565 ))?;
1566 let p = positional_str_lit(call.args.get(1)).ok_or_else(|| format!(
1567 "cursor '{source_name}': {}(dataset, profile) — second arg must be a string literal",
1568 call.func,
1569 ))?;
1570 (d, p, f)
1571 }
1572 _ => return Ok(None),
1573 };
1574
1575 if facet != "base" && facet != "query" {
1576 return Err(format!(
1577 "cursor '{source_name}': vectordata facet must be \"base\" or \"query\", got \"{facet}\""
1578 ));
1579 }
1580 let (count_func, vector_func) = match facet.as_str() {
1581 "base" => ("vector_count", "vector_at"),
1582 "query" => ("query_count", "query_vector_at"),
1583 _ => unreachable!(),
1584 };
1585
1586 let combined = format!("{dataset}:{profile}");
1587 let span = call.span;
1588 let lit = |s: String| Expr::StringLit(s, span);
1589 let positional = |e: Expr| Arg::Positional(e);
1590
1591 let effective_constructor = Expr::Call(CallExpr {
1592 func: "range".into(),
1593 args: vec![
1594 positional(Expr::IntLit(0, span)),
1595 positional(Expr::Call(CallExpr {
1596 func: count_func.into(),
1597 args: vec![positional(lit(combined.clone()))],
1598 span,
1599 })),
1600 ],
1601 span,
1602 });
1603
1604 let prebuffer_binding = AuxBinding {
1605 name: format!("__{source_name}_prebuffer"),
1606 value: Expr::Call(CallExpr {
1607 func: "dataset_prebuffer".into(),
1608 args: vec![positional(lit(combined.clone()))],
1609 span,
1610 }),
1611 projection: None,
1612 };
1613
1614 // Cursor sugar emits the call in (handle, index) order. The
1615 // combined source string auto-promotes to the right facet handle
1616 // through the wire's default resolver (library_catalog.md
1617 // "Shapes", `Resolved<R, T>`).
1618 let vector_binding = AuxBinding {
1619 name: format!("{source_name}__vector"),
1620 value: Expr::Call(CallExpr {
1621 func: vector_func.into(),
1622 args: vec![
1623 positional(lit(combined)),
1624 positional(Expr::Ident(format!("{source_name}__ordinal"), span)),
1625 ],
1626 span,
1627 }),
1628 projection: Some(("vector".into(), crate::ast::PortType::VecF32)),
1629 };
1630
1631 Ok(Some(CursorSugar {
1632 effective_constructor,
1633 aux_bindings: vec![prebuffer_binding, vector_binding],
1634 }))
1635}
1636
1637#[cfg(feature = "vectordata")]
1638inventory::submit! {
1639 crate::dsl::cursor_sugar::CursorSugarRegistration {
1640 handler: vectordata_sugar,
1641 name: "vectordata",
1642 }
1643}
1644
1645#[cfg(test)]
1646mod count_of_tests {
1647 use super::*;
1648 use crate::ast::{PolydatNode, PortType, Slot};
1649 use crate::dsl::registry::DefaultResolver;
1650
1651 /// The ports the compiler types a call against: a handle and a
1652 /// signed value in, a count out. The value is `i64` because that is
1653 /// what the facet stores; a caller with a literal reaches it
1654 /// through `to_i64`, since an integer literal is a `u64` and the
1655 /// adapter catalog will not heal that pair.
1656 #[test]
1657 fn the_count_of_nodes_take_a_handle_and_a_signed_value() {
1658 for node in [
1659 Box::new(MetadataCountOf::new()) as Box<dyn PolydatNode>,
1660 Box::new(PredicateCountOf::new()) as Box<dyn PolydatNode>,
1661 ] {
1662 let meta = node.meta();
1663 assert_eq!(meta.outs.len(), 1);
1664 assert_eq!(meta.outs[0].typ, PortType::U64, "{}", meta.name);
1665 assert_eq!(meta.ins.len(), 2, "{}", meta.name);
1666 let Slot::Wire(handle) = &meta.ins[0] else {
1667 panic!("{}: the first input is a wire", meta.name)
1668 };
1669 let Slot::Wire(value) = &meta.ins[1] else {
1670 panic!("{}: the second input is a wire", meta.name)
1671 };
1672 assert_eq!(handle.typ, PortType::Handle, "{}", meta.name);
1673 assert_eq!(value.typ, PortType::I64, "{}", meta.name);
1674 }
1675 }
1676
1677 /// A handle of the wrong shape answers zero rather than failing.
1678 /// The facets these read are scalar ones; a caller that opened a
1679 /// vector facet by mistake gets a count of nothing, which is the
1680 /// same answer the facet's own readers give for a shape they do not
1681 /// hold.
1682 #[test]
1683 fn a_handle_of_another_shape_counts_nothing() {
1684 let group_shaped = DatasetHandle::Prebuffered {
1685 group: match load_dataset_group("nonexistent:profile") {
1686 Ok(g) => g,
1687 // No dataset available in this environment: the
1688 // fallback is still exercised through the match arm
1689 // below, which is what this test is about.
1690 Err(_) => return,
1691 },
1692 source: "nonexistent:profile".into(),
1693 };
1694 assert_eq!(generic_count_typed(&group_shaped, 1), 0);
1695 }
1696
1697 /// Both names resolve through the registry and carry a default
1698 /// resolver, so `metadata_count_of("ds:profile", v)` promotes the
1699 /// source string to the right facet the way its neighbours do.
1700 #[test]
1701 fn both_names_are_registered_with_a_facet_resolver() {
1702 for (name, facet) in [
1703 ("metadata_count_of", "metadata_content"),
1704 ("predicate_count_of", "metadata_predicates"),
1705 ] {
1706 let sig = crate::dsl::registry::lookup(name)
1707 .unwrap_or_else(|| panic!("{name} is registered"));
1708 assert_eq!(sig.outputs, 1, "{name}");
1709 assert_eq!(sig.params.len(), 2, "{name}");
1710 // The facet comes from the handle argument's type, so
1711 // this is the resolver the macro read off `Facet<_>`.
1712 match sig.default_resolver {
1713 Some(DefaultResolver::Facet(f)) => assert_eq!(f, facet, "{name}"),
1714 other => panic!("{name}: expected a facet resolver, got {other:?}"),
1715 }
1716 }
1717 }
1718}