use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use nmbrs_metrics::labels::Labels;
use crate::adapter::{ExecutionError, OpDispenser, OpResult, WrappingDispenser};
use crate::relevancy::{self, RelevancyFn};
use crate::wires::WireSource;
use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
pub const WRAPPER_NAME: WrapperName = WrapperName::new("validate");
fn wrapper_triggers(s: WrapperSubject) -> bool {
let Some(template) = s.op() else {
return false;
};
template.params.contains_key("verify") || template.params.contains_key("relevancy")
}
fn wrapper_describe_assignment(s: WrapperSubject) -> Option<String> {
let template = s.op()?;
let strict = template
.params
.get("strict")
.and_then(|v| v.as_bool().or_else(|| v.as_str().map(|s| s == "true")))
.unwrap_or(false);
let mut parts: Vec<String> = Vec::new();
if let Some(v) = template.params.get("verify") {
parts.push(format!(
"verify={}",
crate::wrapper_registrations::short_value(v)
));
}
if let Some(v) = template.params.get("relevancy") {
parts.push(format!(
"relevancy={}",
crate::wrapper_registrations::short_value(v)
));
}
if parts.is_empty() {
return None;
}
let body = parts.join(", ");
Some(if strict {
format!("validate: {body} (strict)")
} else {
format!("validate: {body}")
})
}
inventory::submit! {
WrapperRegistration {
name: WRAPPER_NAME,
owned_fields: &["verify", "relevancy", "strict"],
triggers: wrapper_triggers,
requires_inner: &[crate::wrappers::traverse::NAME],
forbids_outer: &[],
mutually_exclusive_with: &[],
describe_assignment: wrapper_describe_assignment,
levels: &[crate::wrapper_registry::WrapperLevel::Op],
}
}
pub const CORE_OP_PARAMS: &[&str] = &[
"batch",
"batchtype",
"max_batch_size",
"ratio",
"adapter",
"driver",
"daemon",
"daemon_cancel_grace_ms",
];
#[derive(Debug, Clone)]
pub struct RelevancyConfig {
pub actual_field: String,
pub expected_binding: String,
pub k: usize,
pub r: Option<usize>,
pub functions: Vec<RelevancyFn>,
}
#[derive(Debug, Clone)]
pub struct AssertionSpec {
pub field: String,
pub predicate: AssertionPredicate,
}
#[derive(Debug, Clone)]
pub enum AssertionPredicate {
Eq(String),
NotNull,
IsNull,
Gte(f64),
Lte(f64),
Contains(String),
MalformedBound { key: String, raw: String },
MinRows(u64),
}
impl AssertionSpec {
pub fn check(&self, result: &OpResult) -> bool {
if let AssertionPredicate::MinRows(n) = &self.predicate {
let row_count = result.body.as_ref().map(|b| b.element_count()).unwrap_or(0);
return row_count >= *n;
}
let json = match &result.body {
Some(body) => body.to_json(),
None => {
return !matches!(
self.predicate,
AssertionPredicate::NotNull | AssertionPredicate::MalformedBound { .. }
);
}
};
let field_val = extract_field_from_json(&json, &self.field);
match &self.predicate {
AssertionPredicate::NotNull => field_val.is_some(),
AssertionPredicate::IsNull => field_val.is_none(),
AssertionPredicate::Eq(expected) => {
field_val.is_some_and(|v| json_value_as_string(v) == *expected)
}
AssertionPredicate::Gte(threshold) => field_val
.and_then(|v| v.as_f64())
.is_some_and(|v| v >= *threshold),
AssertionPredicate::Lte(threshold) => field_val
.and_then(|v| v.as_f64())
.is_some_and(|v| v <= *threshold),
AssertionPredicate::Contains(substr) => {
field_val.is_some_and(|v| json_value_as_string(v).contains(substr.as_str()))
}
AssertionPredicate::MalformedBound { .. } => false,
AssertionPredicate::MinRows(_) => unreachable!("MinRows handled in early-return above"),
}
}
}
pub struct RunningAgg {
pub total_sum: f64,
pub total_count: u64,
pub window: std::collections::VecDeque<f64>,
pub window_size: usize,
}
impl RunningAgg {
pub fn new(window_size: usize) -> Self {
Self {
total_sum: 0.0,
total_count: 0,
window: std::collections::VecDeque::with_capacity(window_size),
window_size,
}
}
pub fn record(&mut self, score: f64) {
self.total_sum += score;
self.total_count += 1;
if self.window.len() == self.window_size {
self.window.pop_front();
}
self.window.push_back(score);
}
pub fn window_mean(&self) -> f64 {
if self.window.is_empty() {
0.0
} else {
self.window.iter().sum::<f64>() / self.window.len() as f64
}
}
pub fn total_mean(&self) -> f64 {
if self.total_count == 0 {
0.0
} else {
self.total_sum / self.total_count as f64
}
}
}
pub const DEFAULT_RECALL_WINDOW: usize = 10;
#[derive(Debug, Clone)]
pub struct RelevancyLive {
pub name: String,
pub window_mean: f64,
pub total_mean: f64,
pub total_count: u64,
pub window_len: usize,
}
pub struct ValidationMetrics {
pub validations_passed: AtomicU64,
pub validations_failed: AtomicU64,
pub relevancy_stats: HashMap<String, nmbrs_metrics::summaries::f64stats::F64Stats>,
pub running_aggregates: HashMap<String, std::sync::Mutex<RunningAgg>>,
}
impl ValidationMetrics {
pub fn new(labels: &Labels, functions: &[RelevancyFn], k: usize, r: Option<usize>) -> Self {
let r_value = r.unwrap_or(k);
let stats_labels = labels
.with("k", k.to_string())
.with("r", r_value.to_string());
let mut stats = HashMap::new();
let mut running = HashMap::new();
for func in functions {
let metric_name = func.metric_name().to_string();
stats.insert(
metric_name.clone(),
nmbrs_metrics::summaries::f64stats::F64Stats::new(
stats_labels.with("name", &metric_name),
),
);
running.insert(
metric_name.clone(),
std::sync::Mutex::new(RunningAgg::new(DEFAULT_RECALL_WINDOW)),
);
}
Self {
validations_passed: AtomicU64::new(0),
validations_failed: AtomicU64::new(0),
relevancy_stats: stats,
running_aggregates: running,
}
}
pub fn assertions_only() -> Self {
Self {
validations_passed: AtomicU64::new(0),
validations_failed: AtomicU64::new(0),
relevancy_stats: HashMap::new(),
running_aggregates: HashMap::new(),
}
}
pub fn record_relevancy(&self, metric_name: &str, score: f64) {
if let Some(stats) = self.relevancy_stats.get(metric_name) {
stats.record(score);
}
if let Some(agg) = self.running_aggregates.get(metric_name) {
let mut a = agg.lock().unwrap_or_else(|e| e.into_inner());
a.record(score);
}
}
pub fn live_snapshot(&self) -> Vec<RelevancyLive> {
let mut out = Vec::with_capacity(self.running_aggregates.len());
for (name, agg) in &self.running_aggregates {
let a = agg.lock().unwrap_or_else(|e| e.into_inner());
out.push(RelevancyLive {
name: name.clone(),
window_mean: a.window_mean(),
total_mean: a.total_mean(),
total_count: a.total_count,
window_len: a.window.len(),
});
}
out.sort_by(|x, y| x.name.cmp(&y.name));
out
}
pub fn passed(&self) -> u64 {
self.validations_passed.load(Ordering::Relaxed)
}
pub fn failed(&self) -> u64 {
self.validations_failed.load(Ordering::Relaxed)
}
}
pub struct ValidatingDispenser {
inner: Arc<dyn OpDispenser>,
assertions: Vec<AssertionSpec>,
relevancy: Option<RelevancyConfig>,
expected_wire_name: Option<String>,
actual_projection: Option<ActualProjection>,
metrics: Arc<ValidationMetrics>,
strict: bool,
}
struct ActualProjection {
segs: Vec<crate::wrappers::result::PathSeg>,
target: Option<polydat::ast::PortType>,
}
type WrappedDispenser = (Arc<dyn OpDispenser>, Option<Arc<ValidationMetrics>>);
impl ValidatingDispenser {
pub fn wrap(
inner: Arc<dyn OpDispenser>,
template: &nmbrs_workload::model::ParsedOp,
labels: &Labels,
program: Option<&polydat::kernel::PolydatProgram>,
fx: &mut crate::fixture::ScopeFixture,
) -> Result<WrappedDispenser, String> {
let template_owned: nmbrs_workload::model::ParsedOp;
let canonical_kernel = inner.canonical_kernel();
let template = if let Some(canonical) = &canonical_kernel {
let mut t = template.clone();
crate::scope::resolve_placeholders_in_op_params(&mut t, canonical.as_ref())?;
template_owned = t;
&template_owned
} else {
template
};
let assertions = parse_assertions(template);
let canonical_wires = canonical_kernel
.as_ref()
.map(|k| crate::wires::KernelWires(k.as_ref()));
let wires_for_parse: Option<&dyn WireSource> =
canonical_wires.as_ref().map(|w| w as &dyn WireSource);
let relevancy = parse_relevancy(template, program, wires_for_parse)?;
let strict = template
.params
.get("strict")
.and_then(|v| v.as_bool())
.unwrap_or(false);
if assertions.is_empty() && relevancy.is_none() {
return Ok((inner, None));
}
let expected_wire_name = match &relevancy {
Some(cfg) => {
let name = cfg
.expected_binding
.trim_matches(|c| c == '{' || c == '}')
.to_string();
fx.register_pull(&name).map_err(|e| {
format!("op '{op}' relevancy.expected: {e}", op = template.name,)
})?;
Some(name)
}
None => None,
};
let metrics = Arc::new(match &relevancy {
Some(config) => ValidationMetrics::new(labels, &config.functions, config.k, config.r),
None => ValidationMetrics::assertions_only(),
});
let actual_projection = match &relevancy {
Some(cfg) => {
let mut raw: Option<String> = None;
if let Some(spec) = template.result.as_ref() {
spec.walk_fragments(|frag| {
if let nmbrs_workload::model::ResultFragment::Named { name, source } = frag
&& name == cfg.actual_field
{
let s = source.trim();
if s != "count" && s != "ok" && !s.contains('(') {
raw = Some(s.to_string());
}
}
});
}
match raw {
Some(path) => {
let segs =
crate::wrappers::result::parse_path_expr(&path).map_err(|e| {
format!(
"op '{op}' relevancy.actual '{field}': result \
binding path: {e}",
op = template.name,
field = cfg.actual_field
)
})?;
let target = template
.abstract_interface
.as_ref()
.and_then(|i| i.results.get(&cfg.actual_field))
.and_then(|kw| polydat::ast::PortType::from_keyword(kw));
Some(ActualProjection { segs, target })
}
None => None,
}
}
None => None,
};
let wrapper = Arc::new(Self {
inner,
assertions,
relevancy,
expected_wire_name,
actual_projection,
metrics: metrics.clone(),
strict,
});
Ok((wrapper, Some(metrics)))
}
}
impl WrappingDispenser for ValidatingDispenser {}
impl OpDispenser for ValidatingDispenser {
fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
Some(self.inner.as_ref())
}
fn execute<'a>(
&'a self,
cycle: u64,
ctx: &'a crate::fixture::ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
let result = self.inner.execute(cycle, ctx).await?;
let mut failed_assertions: Vec<String> = Vec::new();
for assertion in &self.assertions {
if !assertion.check(&result) {
failed_assertions.push(describe_assertion_failure(assertion, &result));
}
}
let all_pass = failed_assertions.is_empty();
if let Some(config) = &self.relevancy {
let actual_ordered = if let Some(proj) = &self.actual_projection {
let projected = result.body.as_ref().and_then(|b| {
crate::wrappers::result::evaluate_path_value(
&b.to_json(),
&proj.segs,
proj.target,
)
});
projected
.as_ref()
.map(resolve_expected_from_value)
.unwrap_or_default()
} else {
match ctx.wires.get(&config.actual_field) {
Some(v) if !matches!(v, polydat::ast::Value::None) => {
resolve_expected_from_value(&v)
}
_ => extract_actual_indices(&result, &config.actual_field),
}
};
let name = self.expected_wire_name.as_deref().expect(
"ValidatingDispenser invariant violated: relevancy is \
configured but expected_wire_name was not stored. \
Construct via ValidatingDispenser::wrap.",
);
let raw_value = ctx.wires.get(name);
let expected_raw = raw_value
.as_ref()
.map(resolve_expected_from_value)
.unwrap_or_default();
if expected_raw.is_empty() {
let available: Vec<String> = ctx.wires.names().collect();
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "relevancy_error".into(),
message: format!(
"relevancy: no ground truth for '{name}'. \
Available wires: {available:?}. \
Ensure the binding exists in the Polydat program.",
),
retryable: false,
}));
}
if actual_ordered.is_empty() && result.body.is_some() {
if self.metrics.passed() + self.metrics.failed() == 0 {
let indent = crate::scene_tree::running_phase_indent();
crate::observer::log(
crate::observer::LogLevel::Warn,
&format!(
"{indent}relevancy: no values extracted for field '{}' from result",
config.actual_field
),
);
if let Some(body) = &result.body {
let preview =
serde_json::to_string(&body.to_json()).unwrap_or_default();
crate::observer::log(
crate::observer::LogLevel::Warn,
&format!(
"{indent} result preview: {}",
&preview[..preview.len().min(300)]
),
);
}
}
}
if let Some(r) = config.r
&& actual_ordered.len() != r
{
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "relevancy_error".into(),
message: format!(
"relevancy: expected exactly r={r} results from \
retrieval, got {} (k-recall@r contract). Either \
size the query LIMIT to {r}, or remove the `r:` \
declaration to fall back to first-k semantics.",
actual_ordered.len(),
),
retryable: false,
}));
}
let r_window = config.r.unwrap_or(config.k);
let expected_sorted = relevancy::truncate_and_sort(&expected_raw, config.k);
let actual_sorted = relevancy::truncate_and_sort(&actual_ordered, r_window);
for func in &config.functions {
let score =
func.compute(&expected_sorted, &actual_sorted, &actual_ordered, config.k);
self.metrics.record_relevancy(func.metric_name(), score);
if crate::observer::trace_enabled() {
let intersect =
crate::relevancy::intersection_count(&expected_sorted, &actual_sorted);
let stats_labels = self
.metrics
.relevancy_stats
.get(func.metric_name())
.map(|s| s.labels().clone())
.unwrap_or_default();
crate::observer::trace(
&stats_labels,
&format!(
"event=relevancy.score cycle={cycle} \
fn={func} k={k} r={r} \
gt_card={gt} actual_card={ac} \
intersect={inter} score={score:.6}",
func = func.metric_name(),
k = config.k,
r = config.r.unwrap_or(config.k),
gt = expected_sorted.len(),
ac = actual_sorted.len(),
inter = intersect,
),
);
}
}
}
if all_pass {
self.metrics
.validations_passed
.fetch_add(1, Ordering::Relaxed);
} else {
self.metrics
.validations_failed
.fetch_add(1, Ordering::Relaxed);
if self.strict {
return Err(ExecutionError::Op(crate::adapter::AdapterError {
error_name: "validation_failed".into(),
message: format!(
"result validation failed (strict mode): {}",
failed_assertions.join("; "),
),
retryable: false,
}));
}
}
Ok(result)
})
}
}
fn describe_assertion_failure(assertion: &AssertionSpec, result: &OpResult) -> String {
let body_tail = match &result.body {
Some(_) => format!("; body: {}", body_excerpt(result)),
None => String::new(),
};
let observed_repr = observed_field_repr(&assertion.field, result);
match &assertion.predicate {
AssertionPredicate::MinRows(n) => {
let got = result.body.as_ref().map(|b| b.element_count()).unwrap_or(0);
format!("min_rows: expected ≥{n}, got {got}{body_tail}")
}
AssertionPredicate::Eq(expected) => format!(
"field '{}' eq '{}' failed (observed: {observed_repr}){body_tail}",
assertion.field, expected
),
AssertionPredicate::NotNull => format!(
"field '{}' must not be null (observed: {observed_repr}){body_tail}",
assertion.field
),
AssertionPredicate::IsNull => format!(
"field '{}' must be null (observed: {observed_repr}){body_tail}",
assertion.field
),
AssertionPredicate::Gte(t) => format!(
"field '{}' >= {t} failed (observed: {observed_repr}){body_tail}",
assertion.field
),
AssertionPredicate::MalformedBound { key, raw } => format!(
"field '{}': `{key}: {raw}` is not a number — a numeric bound \
must be a number (`{key}: 5`) or a string holding one \
(`{key}: \"5\"`). If that is a `{{placeholder}}`, it did not \
resolve.",
assertion.field
),
AssertionPredicate::Lte(t) => format!(
"field '{}' <= {t} failed (observed: {observed_repr}){body_tail}",
assertion.field
),
AssertionPredicate::Contains(sub) => format!(
"field '{}' contains '{}' failed (observed: {observed_repr}){body_tail}",
assertion.field, sub
),
}
}
fn observed_field_repr(field: &str, result: &OpResult) -> String {
let Some(body) = &result.body else {
return "<no body returned by op — verify clause cannot read fields>".to_string();
};
let json = body.to_json();
if let Some(v) = extract_field_from_json(&json, field) {
let repr = match v {
serde_json::Value::String(s) => format!("\"{s}\""),
other => other.to_string(),
};
return truncate_for_message(&repr, 160);
}
match &json {
serde_json::Value::Object(map) => {
let keys: Vec<&str> = map.keys().map(|s| s.as_str()).collect();
if keys.is_empty() {
"<field absent; body is an empty JSON object>".to_string()
} else {
format!("<field absent; body keys: {keys:?}>")
}
}
serde_json::Value::Array(arr) => {
format!(
"<field absent; body is a JSON array of {} element{}>",
arr.len(),
if arr.len() == 1 { "" } else { "s" },
)
}
serde_json::Value::Null => "<field absent; body is JSON null>".to_string(),
other => {
let short = truncate_for_message(&other.to_string(), 80);
format!("<not-json: {short}>")
}
}
}
fn body_excerpt(result: &OpResult) -> String {
let Some(body) = &result.body else {
return "<no body>".to_string();
};
truncate_for_message(&body.to_text(), 512)
}
fn truncate_for_message(s: &str, max: usize) -> String {
if s.chars().count() <= max {
return s.to_string();
}
let cut: String = s.chars().take(max.saturating_sub(1)).collect();
format!("{cut}…")
}
fn numeric_bound(v: &serde_json::Value) -> Option<f64> {
v.as_f64()
.or_else(|| v.as_str().and_then(|s| s.trim().parse::<f64>().ok()))
}
fn parse_assertions(template: &nmbrs_workload::model::ParsedOp) -> Vec<AssertionSpec> {
let Some(verify) = template.params.get("verify") else {
return Vec::new();
};
let Some(items) = verify.as_array() else {
return Vec::new();
};
let mut assertions = Vec::new();
for item in items {
let Some(obj) = item.as_object() else {
continue;
};
if let Some(v) = obj.get("min_rows") {
let n = v.as_u64().unwrap_or(0);
assertions.push(AssertionSpec {
field: String::new(), predicate: AssertionPredicate::MinRows(n),
});
continue;
}
let Some(field) = obj.get("field").and_then(|v| v.as_str()) else {
continue;
};
let predicate = if let Some(v) = obj.get("eq") {
AssertionPredicate::Eq(json_value_as_string(v))
} else if let Some(v) = obj.get("gte") {
match numeric_bound(v) {
Some(t) => AssertionPredicate::Gte(t),
None => AssertionPredicate::MalformedBound {
key: "gte".into(),
raw: json_value_as_string(v),
},
}
} else if let Some(v) = obj.get("lte") {
match numeric_bound(v) {
Some(t) => AssertionPredicate::Lte(t),
None => AssertionPredicate::MalformedBound {
key: "lte".into(),
raw: json_value_as_string(v),
},
}
} else if let Some(v) = obj.get("contains") {
AssertionPredicate::Contains(json_value_as_string(v))
} else if let Some(v) = obj.get("is") {
match v.as_str().unwrap_or("").to_lowercase().as_str() {
"not_null" | "notnull" => AssertionPredicate::NotNull,
"null" => AssertionPredicate::IsNull,
_ => continue,
}
} else {
continue;
};
assertions.push(AssertionSpec {
field: field.to_string(),
predicate,
});
}
assertions
}
const RELEVANCY_VOCAB: &[&str] = &["actual", "expected", "k", "r", "functions"];
fn parse_relevancy(
template: &nmbrs_workload::model::ParsedOp,
_program: Option<&polydat::kernel::PolydatProgram>,
wires: Option<&dyn WireSource>,
) -> Result<Option<RelevancyConfig>, String> {
let Some(rel) = template.params.get("relevancy") else {
return Ok(None);
};
let obj = rel.as_object().ok_or_else(|| {
format!(
"op '{}': relevancy: expected a mapping, got {kind}",
template.name,
kind = match rel {
serde_json::Value::Null => "null",
serde_json::Value::Bool(_) => "boolean",
serde_json::Value::Number(_) => "number",
serde_json::Value::String(_) => "string",
serde_json::Value::Array(_) => "array",
_ => "unknown",
},
)
})?;
for k in obj.keys() {
if !RELEVANCY_VOCAB.contains(&k.as_str()) {
return Err(format!(
"op '{op}' relevancy: unknown key '{k}'. Allowed: [{vocab}]",
op = template.name,
vocab = RELEVANCY_VOCAB.join(", "),
));
}
}
let actual_field = obj
.get("actual")
.and_then(|v| v.as_str())
.ok_or_else(|| {
format!(
"op '{}' relevancy: missing required field 'actual' (string column name)",
template.name,
)
})?
.to_string();
let expected_binding = obj
.get("expected")
.and_then(|v| v.as_str())
.ok_or_else(|| {
format!(
"op '{}' relevancy: missing required field 'expected' (binding reference)",
template.name,
)
})?
.to_string();
let k_label = format!("op '{}' relevancy.k", template.name);
let k = parse_count_param(obj.get("k"), &k_label, wires)?.ok_or_else(|| {
format!(
"op '{}' relevancy: missing required field 'k' (integer)",
template.name,
)
})? as usize;
let r_label = format!("op '{}' relevancy.r", template.name);
let r: Option<usize> = parse_count_param(obj.get("r"), &r_label, wires)?.map(|n| n as usize);
if let Some(rv) = r
&& rv < k
{
return Err(format!(
"op '{op}' relevancy: r={rv} is smaller than k={k}; \
the k-recall@r contract requires r >= k",
op = template.name,
));
}
let functions: Vec<RelevancyFn> = match obj.get("functions") {
None => vec![RelevancyFn::Recall],
Some(serde_json::Value::Array(arr)) => {
let mut out: Vec<RelevancyFn> = Vec::new();
for (i, v) in arr.iter().enumerate() {
let name = v.as_str().ok_or_else(|| {
format!(
"op '{op}' relevancy.functions[{i}]: expected a string, got {kind}",
op = template.name,
kind = match v {
serde_json::Value::Null => "null",
serde_json::Value::Bool(_) => "boolean",
serde_json::Value::Number(_) => "number",
serde_json::Value::Array(_) => "array",
serde_json::Value::Object(_) => "object",
_ => "unknown",
},
)
})?;
let func = RelevancyFn::parse(name).ok_or_else(|| {
format!(
"op '{op}' relevancy.functions[{i}]: unknown function '{name}'",
op = template.name,
)
})?;
out.push(func);
}
if out.is_empty() {
return Err(format!(
"op '{op}' relevancy.functions: empty list — declare at least one \
function or remove the field to default to [recall]",
op = template.name,
));
}
out
}
Some(_) => {
return Err(format!(
"op '{op}' relevancy.functions: expected an array of strings",
op = template.name,
));
}
};
Ok(Some(RelevancyConfig {
actual_field,
expected_binding,
k,
r,
functions,
}))
}
fn parse_count_param(
val: Option<&serde_json::Value>,
field_label: &str,
wires: Option<&dyn WireSource>,
) -> Result<Option<u64>, String> {
let Some(v) = val else {
return Ok(None);
};
if let Some(n) = v.as_u64() {
return Ok(Some(n));
}
let s = match v.as_str() {
Some(s) => s,
None => {
return Err(format!(
"{field_label}: expected an integer or numeric string, got {kind}",
kind = match v {
serde_json::Value::Null => "null",
serde_json::Value::Bool(_) => "boolean",
serde_json::Value::Array(_) => "array",
serde_json::Value::Object(_) => "object",
_ => "unsupported value",
},
));
}
};
let trimmed = s.trim();
if trimmed.starts_with('{') && trimmed.ends_with('}') {
return Err(format!(
"{field_label}: '{trimmed}' was not resolved before parameter parsing — \
this is a placeholder-resolution bug, not a config-time issue. The \
single-read-path resolver should have substituted it from the kernel."
));
}
if let Some(wires) = wires
&& is_bare_ident(trimmed)
&& let Some(value) = wires.get(trimmed)
{
return value_to_u64_for_count(value)
.ok_or_else(|| {
format!(
"{field_label}: wire '{trimmed}' resolved but its value is not \
coercible to a non-negative integer"
)
})
.map(Some);
}
match trimmed.parse::<u64>() {
Ok(n) => Ok(Some(n)),
Err(_) => Err(format!(
"{field_label}: '{trimmed}' is not a valid non-negative integer \
(and not declared as a wire name on the op-template kernel)"
)),
}
}
fn is_bare_ident(s: &str) -> bool {
let mut chars = s.chars();
match chars.next() {
Some(c) if c.is_ascii_alphabetic() || c == '_' => {}
_ => return false,
}
chars.all(|c| c.is_ascii_alphanumeric() || c == '_')
}
pub(crate) fn value_to_u64_for_count(value: polydat::ast::Value) -> Option<u64> {
use polydat::ast::Value;
match value {
Value::U64(n) => Some(n),
Value::I64(n) => u64::try_from(n).ok(),
Value::Str(s) => crate::runner::parse_count(&s),
Value::F64(f) if f.is_finite() && f >= 0.0 => Some(f as u64),
Value::Bool(true) => Some(1),
Value::Bool(false) => Some(0),
_ => None,
}
}
fn extract_actual_indices(result: &OpResult, field: &str) -> Vec<i64> {
let Some(body) = &result.body else {
return Vec::new();
};
extract_indices_from_json(&body.to_json(), field)
}
fn extract_indices_from_json(json: &serde_json::Value, field: &str) -> Vec<i64> {
match json {
serde_json::Value::Array(rows) => rows
.iter()
.filter_map(|row| json_field_as_i64(row.get(field)?))
.collect(),
serde_json::Value::Object(obj) => {
if let Some(rows) = obj.get("rows") {
return extract_indices_from_json(rows, field);
}
obj.get(field)
.and_then(json_field_as_i64)
.into_iter()
.collect()
}
_ => Vec::new(),
}
}
fn json_field_as_i64(v: &serde_json::Value) -> Option<i64> {
v.as_i64().or_else(|| v.as_str()?.parse().ok())
}
fn resolve_expected_from_value(value: &polydat::ast::Value) -> Vec<i64> {
match value {
polydat::ast::Value::VecI32(slice) => slice.as_slice().iter().map(|&x| x as i64).collect(),
polydat::ast::Value::VecI64(slice) => slice.as_slice().to_vec(),
polydat::ast::Value::VecF32(slice) => {
slice.as_slice().iter().map(|&x| x as i64).collect()
}
polydat::ast::Value::VecF64(slice) => slice.as_slice().iter().map(|&x| x as i64).collect(),
polydat::ast::Value::Json(j) => match &**j {
serde_json::Value::Array(elems) => elems.iter().filter_map(json_field_as_i64).collect(),
other => json_field_as_i64(other).into_iter().collect(),
},
polydat::ast::Value::Str(s) => parse_int_array(s),
polydat::ast::Value::U64(v) => vec![*v as i64],
_ => {
let s = value.to_display_string();
parse_int_array(&s)
}
}
}
fn parse_int_array(s: &str) -> Vec<i64> {
let trimmed = s.trim().trim_start_matches('[').trim_end_matches(']');
trimmed
.split(|c: char| c == ',' || c.is_whitespace())
.filter(|s| !s.is_empty())
.filter_map(|s| s.trim().parse::<i64>().ok())
.collect()
}
fn extract_field_from_json<'a>(
json: &'a serde_json::Value,
field: &str,
) -> Option<&'a serde_json::Value> {
match json {
serde_json::Value::Object(obj) => obj.get(field).or_else(|| {
obj.get("rows")
.and_then(|r| r.as_array())
.and_then(|rows| rows.first())
.and_then(|row| row.get(field))
}),
serde_json::Value::Array(rows) => rows.first().and_then(|row| row.get(field)),
_ => None,
}
}
fn json_value_as_string(v: &serde_json::Value) -> String {
match v {
serde_json::Value::String(s) => s.clone(),
other => other.to_string(),
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::ResultBody;
use std::any::Any;
#[derive(Debug)]
struct JsonBody(serde_json::Value);
impl ResultBody for JsonBody {
fn to_json(&self) -> serde_json::Value {
self.0.clone()
}
fn as_any(&self) -> &dyn Any {
self
}
}
#[test]
fn parse_int_array_bracket_format() {
assert_eq!(parse_int_array("[1, 5, 12, 23]"), vec![1, 5, 12, 23]);
}
#[test]
fn core_op_params_disjoint_from_owned_fields() {
let registry = crate::wrapper_registry::WrapperRegistry::from_inventory();
let owned = registry.all_owned_fields();
let dupes: Vec<&str> = CORE_OP_PARAMS
.iter()
.copied()
.filter(|p| owned.contains(p))
.collect();
assert!(
dupes.is_empty(),
"these CORE_OP_PARAMS are already wrapper-owned (remove them — \
the guard accepts them via WrapperRegistry::owns_field): {dupes:?}",
);
}
#[test]
fn wrapper_field_accepted_without_cli_or_core_membership() {
let registry = crate::wrapper_registry::WrapperRegistry::from_inventory();
assert!(
registry.owns_field("readout"),
"readout must be registry-owned"
);
assert!(
registry.owns_field("errors"),
"errors must be registry-owned (was riding the CLI hatch)"
);
assert!(
registry.owns_field("tries"),
"tries must be registry-owned (was riding the CLI hatch)"
);
assert!(
!CORE_OP_PARAMS.contains(&"readout"),
"readout should NOT be in CORE_OP_PARAMS — it's wrapper-owned"
);
}
#[test]
fn parse_int_array_comma_format() {
assert_eq!(parse_int_array("1,5,12,23"), vec![1, 5, 12, 23]);
}
#[test]
fn parse_int_array_space_format() {
assert_eq!(parse_int_array("1 5 12 23"), vec![1, 5, 12, 23]);
}
#[test]
fn parse_int_array_empty() {
assert_eq!(parse_int_array("[]"), Vec::<i64>::new());
assert_eq!(parse_int_array(""), Vec::<i64>::new());
}
#[test]
fn extract_indices_from_json_array() {
let json = serde_json::json!([
{"key": 5, "distance": 0.1},
{"key": 12, "distance": 0.2},
{"key": 3, "distance": 0.3},
]);
assert_eq!(extract_indices_from_json(&json, "key"), vec![5, 12, 3]);
}
#[test]
fn extract_indices_from_json_rows_wrapper() {
let json = serde_json::json!({
"rows": [
{"key": 5},
{"key": 12},
]
});
assert_eq!(extract_indices_from_json(&json, "key"), vec![5, 12]);
}
#[test]
fn assertion_not_null() {
let result = OpResult {
body: Some(Box::new(JsonBody(serde_json::json!({"name": "alice"})))),
skipped: false,
};
let spec = AssertionSpec {
field: "name".into(),
predicate: AssertionPredicate::NotNull,
};
assert!(spec.check(&result));
let spec_missing = AssertionSpec {
field: "age".into(),
predicate: AssertionPredicate::NotNull,
};
assert!(!spec_missing.check(&result));
}
#[test]
fn assertion_eq() {
let result = OpResult {
body: Some(Box::new(JsonBody(serde_json::json!({"status": "ok"})))),
skipped: false,
};
let spec = AssertionSpec {
field: "status".into(),
predicate: AssertionPredicate::Eq("ok".into()),
};
assert!(spec.check(&result));
let spec_fail = AssertionSpec {
field: "status".into(),
predicate: AssertionPredicate::Eq("error".into()),
};
assert!(!spec_fail.check(&result));
}
#[test]
fn assertion_gte() {
let result = OpResult {
body: Some(Box::new(JsonBody(serde_json::json!({"balance": 42.5})))),
skipped: false,
};
let spec = AssertionSpec {
field: "balance".into(),
predicate: AssertionPredicate::Gte(0.0),
};
assert!(spec.check(&result));
let spec_fail = AssertionSpec {
field: "balance".into(),
predicate: AssertionPredicate::Gte(100.0),
};
assert!(!spec_fail.check(&result));
}
#[test]
fn assertion_no_body() {
let result = OpResult {
body: None,
skipped: false,
};
let spec = AssertionSpec {
field: "anything".into(),
predicate: AssertionPredicate::IsNull,
};
assert!(spec.check(&result));
let spec_not_null = AssertionSpec {
field: "anything".into(),
predicate: AssertionPredicate::NotNull,
};
assert!(!spec_not_null.check(&result));
}
#[derive(Debug)]
struct CountedBody {
rows: Vec<serde_json::Value>,
}
impl ResultBody for CountedBody {
fn to_json(&self) -> serde_json::Value {
serde_json::Value::Array(self.rows.clone())
}
fn as_any(&self) -> &dyn Any {
self
}
fn element_count(&self) -> u64 {
self.rows.len() as u64
}
}
#[test]
fn assertion_min_rows_passes_when_threshold_met() {
let result = OpResult {
body: Some(Box::new(CountedBody {
rows: vec![
serde_json::json!({"index_name": "vec_idx"}),
serde_json::json!({"index_name": "meta_idx"}),
],
})),
skipped: false,
};
let spec = AssertionSpec {
field: String::new(),
predicate: AssertionPredicate::MinRows(1),
};
assert!(spec.check(&result));
let spec_two = AssertionSpec {
field: String::new(),
predicate: AssertionPredicate::MinRows(2),
};
assert!(spec_two.check(&result));
}
#[test]
fn assertion_min_rows_fails_when_below_threshold() {
let result = OpResult {
body: Some(Box::new(CountedBody { rows: Vec::new() })),
skipped: false,
};
let spec = AssertionSpec {
field: String::new(),
predicate: AssertionPredicate::MinRows(1),
};
assert!(!spec.check(&result));
let result_none = OpResult {
body: None,
skipped: false,
};
assert!(!spec.check(&result_none));
}
#[test]
fn parse_assertions_min_rows_from_yaml() {
let mut template = nmbrs_workload::model::ParsedOp::simple("await", "test");
template.params.insert(
"verify".into(),
serde_json::json!([
{"min_rows": 1},
]),
);
let assertions = parse_assertions(&template);
assert_eq!(assertions.len(), 1);
match &assertions[0].predicate {
AssertionPredicate::MinRows(n) => assert_eq!(*n, 1),
other => panic!("expected MinRows(1), got {other:?}"),
}
}
#[test]
fn eq_failure_includes_body_and_distinguishes_absent_vs_not_json() {
let result = OpResult {
body: Some(Box::new(JsonBody(serde_json::json!({
"value": null, "request": {"type": "exec"}
})))),
skipped: false,
};
let spec = AssertionSpec {
field: "status".into(),
predicate: AssertionPredicate::Eq("200".into()),
};
let msg = describe_assertion_failure(&spec, &result);
assert!(
msg.contains("field absent"),
"json-without-field should mark observed as absent, got: {msg}"
);
assert!(
msg.contains("body keys"),
"absent message should enumerate present keys, got: {msg}"
);
assert!(
msg.contains("\"value\"") && msg.contains("\"request\""),
"key list should include both present keys, got: {msg}"
);
assert!(msg.contains("body: "), "body excerpt missing: {msg}");
assert!(
msg.contains("\"request\""),
"body excerpt should echo the actual JSON: {msg}"
);
#[derive(Debug)]
struct PlainBody(String);
impl ResultBody for PlainBody {
fn to_json(&self) -> serde_json::Value {
serde_json::Value::String(self.0.clone())
}
fn as_any(&self) -> &dyn Any {
self
}
fn to_text(&self) -> String {
self.0.clone()
}
}
let text_result = OpResult {
body: Some(Box::new(PlainBody(
"<html><body>404 Not Found</body></html>".into(),
))),
skipped: false,
};
let msg2 = describe_assertion_failure(&spec, &text_result);
assert!(
msg2.contains("not-json"),
"text body should mark observed as not-json, got: {msg2}"
);
assert!(
msg2.contains("404 Not Found"),
"body excerpt should include the text: {msg2}"
);
}
#[test]
fn eq_failure_with_no_body_explains_situation() {
let result = OpResult {
body: None,
skipped: false,
};
let spec = AssertionSpec {
field: "status".into(),
predicate: AssertionPredicate::Eq("200".into()),
};
let msg = describe_assertion_failure(&spec, &result);
assert!(
msg.contains("no body returned by op"),
"no-body case should explain why the field can't be read, got: {msg}"
);
assert!(
!msg.contains("body: "),
"no-body case should suppress the redundant body excerpt, got: {msg}"
);
}
#[test]
fn min_rows_failure_describes_actual_vs_expected() {
let result = OpResult {
body: Some(Box::new(CountedBody { rows: Vec::new() })),
skipped: false,
};
let spec = AssertionSpec {
field: String::new(),
predicate: AssertionPredicate::MinRows(1),
};
let msg = describe_assertion_failure(&spec, &result);
assert!(msg.contains("min_rows"), "got: {msg}");
assert!(msg.contains("≥1"), "got: {msg}");
assert!(msg.contains("got 0"), "got: {msg}");
}
#[test]
fn validation_metrics_record_relevancy() {
let labels = Labels::of("activity", "test");
let metrics = ValidationMetrics::new(
&labels,
&[RelevancyFn::Recall, RelevancyFn::Precision],
10,
Some(20),
);
assert!(metrics.relevancy_stats.contains_key("recall"));
assert!(metrics.relevancy_stats.contains_key("precision"));
assert!(!metrics.relevancy_stats.contains_key("f1"));
let recall_labels = metrics.relevancy_stats["recall"].labels();
assert_eq!(recall_labels.get("k"), Some("10"));
assert_eq!(recall_labels.get("r"), Some("20"));
metrics.record_relevancy("recall", 0.85);
metrics.record_relevancy("recall", 0.90);
let snap = metrics.relevancy_stats["recall"].snapshot();
assert_eq!(snap.len(), 2);
}
#[test]
fn validation_metrics_r_defaults_to_k() {
let labels = Labels::of("activity", "test");
let metrics = ValidationMetrics::new(&labels, &[RelevancyFn::Recall], 100, None);
let l = metrics.relevancy_stats["recall"].labels();
assert_eq!(l.get("k"), Some("100"));
assert_eq!(l.get("r"), Some("100"));
}
#[test]
fn resolve_expected_string_array_form() {
let v = polydat::ast::Value::Str("[1, 5, 12, 23]".into());
assert_eq!(resolve_expected_from_value(&v), vec![1, 5, 12, 23]);
}
#[test]
fn resolve_expected_string_csv_form() {
let v = polydat::ast::Value::Str("1,5,12".into());
assert_eq!(resolve_expected_from_value(&v), vec![1, 5, 12]);
}
#[test]
fn resolve_expected_native_veci32_fast_path() {
use polydat::ast::{SliceArc, Value};
let slice = SliceArc::<i32>::from_vec(vec![1, 5, 12, 23, 100]);
let v = Value::VecI32(slice);
assert_eq!(resolve_expected_from_value(&v), vec![1, 5, 12, 23, 100]);
}
#[test]
fn resolve_expected_native_vecf32_fast_path() {
use polydat::ast::{SliceArc, Value};
let slice = SliceArc::<f32>::from_vec(vec![1.0, 2.0, 3.0]);
let v = Value::VecF32(slice);
assert_eq!(resolve_expected_from_value(&v), vec![1, 2, 3]);
}
#[test]
fn parse_relevancy_from_params() {
let mut template = nmbrs_workload::model::ParsedOp::simple("test", "SELECT key FROM t");
template.params.insert(
"relevancy".into(),
serde_json::json!({
"actual": "key",
"expected": "{ground_truth}",
"k": 10,
"functions": ["recall", "precision", "f1"]
}),
);
let config = parse_relevancy(&template, None, None).unwrap().unwrap();
assert_eq!(config.actual_field, "key");
assert_eq!(config.expected_binding, "{ground_truth}");
assert_eq!(config.k, 10);
assert!(config.r.is_none(), "r should default to None when absent");
assert_eq!(config.functions.len(), 3);
assert_eq!(config.functions[0], RelevancyFn::Recall);
assert_eq!(config.functions[1], RelevancyFn::Precision);
assert_eq!(config.functions[2], RelevancyFn::F1);
}
#[test]
fn parse_relevancy_with_r_for_k_recall_at_r() {
let mut template = nmbrs_workload::model::ParsedOp::simple("test", "SELECT key FROM t");
template.params.insert(
"relevancy".into(),
serde_json::json!({
"actual": "key",
"expected": "{ground_truth}",
"k": 10,
"r": 100,
"functions": ["recall"],
}),
);
let config = parse_relevancy(&template, None, None).unwrap().unwrap();
assert_eq!(config.k, 10);
assert_eq!(config.r, Some(100));
}
#[test]
fn parse_relevancy_r_accepts_string_form() {
let mut template = nmbrs_workload::model::ParsedOp::simple("test", "SELECT key FROM t");
template.params.insert(
"relevancy".into(),
serde_json::json!({
"actual": "key",
"expected": "{ground_truth}",
"k": "10",
"r": "100",
}),
);
let config = parse_relevancy(&template, None, None).unwrap().unwrap();
assert_eq!(config.k, 10);
assert_eq!(config.r, Some(100));
}
#[test]
fn parse_relevancy_missing() {
let template = nmbrs_workload::model::ParsedOp::simple("test", "INSERT");
assert!(parse_relevancy(&template, None, None).unwrap().is_none());
}
#[test]
fn parse_relevancy_k_and_r_accept_bare_wire_names() {
let kernel = crate::scope_kernel::ScopeKernel::compile(
"input cycle: u64\n\
const k := 10\n\
const limit := 100\n",
)
.expect("compile scope kernel wires");
let wires: &dyn WireSource = &kernel;
let mut template = nmbrs_workload::model::ParsedOp::simple("test", "SELECT key FROM t");
template.params.insert(
"relevancy".into(),
serde_json::json!({
"actual": "key",
"expected": "{ground_truth}",
"k": "k", "r": "limit", "functions": ["recall"],
}),
);
let config = parse_relevancy(&template, None, Some(wires))
.unwrap()
.unwrap();
assert_eq!(config.k, 10, "bare `k:` resolved through wires");
assert_eq!(config.r, Some(100), "bare `r:` resolved through wires");
}
#[test]
fn parse_relevancy_bare_name_falls_through_to_int_parse_when_no_kernel() {
let mut template = nmbrs_workload::model::ParsedOp::simple("test", "SELECT key FROM t");
template.params.insert(
"relevancy".into(),
serde_json::json!({
"actual": "key",
"expected": "{ground_truth}",
"k": "k",
"functions": ["recall"],
}),
);
let err = parse_relevancy(&template, None, None).unwrap_err();
assert!(
err.contains("'k' is not a valid non-negative integer"),
"diagnostic should describe the parse failure: {err}"
);
}
#[test]
fn value_predicates_are_vacuous_without_a_body() {
let result = OpResult {
body: None,
..Default::default()
};
for predicate in [
AssertionPredicate::Eq("200".into()),
AssertionPredicate::Lte(5.0),
AssertionPredicate::Gte(1.0),
AssertionPredicate::Contains("ok".into()),
AssertionPredicate::IsNull,
] {
let spec = AssertionSpec {
field: "status".into(),
predicate,
};
assert!(
spec.check(&result),
"a value predicate has nothing to contradict it: {:?}",
spec.predicate
);
}
}
#[test]
fn presence_predicates_still_fail_without_a_body() {
let result = OpResult {
body: None,
..Default::default()
};
let not_null = AssertionSpec {
field: "status".into(),
predicate: AssertionPredicate::NotNull,
};
assert!(
!not_null.check(&result),
"`is: not_null` is how an author demands the field exist"
);
let min_rows = AssertionSpec {
field: String::new(),
predicate: AssertionPredicate::MinRows(1),
};
assert!(
!min_rows.check(&result),
"`min_rows: 1` is how an author demands a non-empty result"
);
}
#[test]
fn malformed_bounds_fail_even_without_a_body() {
let result = OpResult {
body: None,
..Default::default()
};
let spec = AssertionSpec {
field: "value".into(),
predicate: AssertionPredicate::MalformedBound {
key: "lte".into(),
raw: "{unresolved}".into(),
},
};
assert!(!spec.check(&result));
}
#[test]
fn numeric_bounds_accept_quoted_numbers() {
let mut template = nmbrs_workload::model::ParsedOp::simple("t", "noop");
template.params.insert(
"verify".into(),
serde_json::json!([
{"field": "value", "lte": "5"},
{"field": "value", "gte": " 2 "},
{"field": "other", "lte": 7},
]),
);
let a = parse_assertions(&template);
assert!(
matches!(a[0].predicate, AssertionPredicate::Lte(t) if t == 5.0),
"quoted lte must parse: {:?}",
a[0].predicate
);
assert!(
matches!(a[1].predicate, AssertionPredicate::Gte(t) if t == 2.0),
"surrounding whitespace is not a malformed bound: {:?}",
a[1].predicate
);
assert!(
matches!(a[2].predicate, AssertionPredicate::Lte(t) if t == 7.0),
"the unquoted form is unchanged: {:?}",
a[2].predicate
);
}
#[test]
fn non_numeric_bounds_are_malformed_not_zero() {
let mut template = nmbrs_workload::model::ParsedOp::simple("t", "noop");
template.params.insert(
"verify".into(),
serde_json::json!([
{"field": "value", "lte": "{unresolved}"},
{"field": "value", "gte": "abc"},
]),
);
let a = parse_assertions(&template);
for spec in &a {
match &spec.predicate {
AssertionPredicate::MalformedBound { raw, .. } => {
assert!(
!raw.is_empty(),
"the offending text is carried for the message"
);
}
other => panic!("expected MalformedBound, got {other:?}"),
}
}
let msg = format!("{:?}", a[0].predicate);
assert!(
msg.contains("MalformedBound"),
"the predicate stays malformed all the way to reporting: {msg}"
);
}
#[test]
fn parse_assertions_from_params() {
let mut template = nmbrs_workload::model::ParsedOp::simple("test", "SELECT");
template.params.insert(
"verify".into(),
serde_json::json!([
{"field": "name", "is": "not_null"},
{"field": "balance", "gte": 0},
{"field": "status", "eq": "active"},
]),
);
let assertions = parse_assertions(&template);
assert_eq!(assertions.len(), 3);
assert_eq!(assertions[0].field, "name");
assert!(matches!(
assertions[0].predicate,
AssertionPredicate::NotNull
));
assert_eq!(assertions[1].field, "balance");
assert!(matches!(assertions[1].predicate, AssertionPredicate::Gte(v) if v == 0.0));
assert_eq!(assertions[2].field, "status");
assert!(matches!(&assertions[2].predicate, AssertionPredicate::Eq(s) if s == "active"));
}
#[test]
fn parse_assertions_missing() {
let template = nmbrs_workload::model::ParsedOp::simple("test", "INSERT");
assert!(parse_assertions(&template).is_empty());
}
}