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