use std::sync::Arc;
use super::traverse::json_to_value;
use crate::adapter::WrappingDispenser;
use crate::adapter::{ExecutionError, OpDispenser, OpResult};
use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
pub const NAME: WrapperName = WrapperName::new("result");
fn triggers(s: WrapperSubject) -> bool {
s.op().is_some()
}
fn describe_assignment(s: WrapperSubject) -> Option<String> {
let template = s.op()?;
let spec = template.result.as_ref()?;
if spec.is_empty() {
return None;
}
let mut names: Vec<String> = Vec::new();
spec.walk_fragments(|frag| match frag {
nmbrs_workload::model::ResultFragment::Named { name, .. } => {
names.push(name.to_string());
}
nmbrs_workload::model::ResultFragment::Source(source) => {
for line in source.lines() {
let line = line.trim();
if let Some((name, _)) = line.split_once(":=") {
names.push(name.trim().to_string());
}
}
}
});
if names.is_empty() {
return None;
}
names.sort();
names.dedup();
Some(format!("result: captures {}", names.join(", ")))
}
inventory::submit! {
WrapperRegistration {
name: NAME,
owned_fields: &[],
triggers,
requires_inner: &[super::traverse::NAME],
forbids_outer: &[],
mutually_exclusive_with: &[],
describe_assignment,
levels: &[crate::wrapper_registry::WrapperLevel::Op],
}
}
pub struct ResultDispenser {
inner: Arc<dyn OpDispenser>,
specs: Vec<ResultSlot>,
populate_kernel_inputs: bool,
}
struct ResultSlot {
wire: String,
source: ResultSource,
default: Option<polydat::ast::Value>,
target: Option<polydat::ast::PortType>,
}
enum ResultSource {
Count,
Ok,
Path(Vec<PathSeg>),
#[allow(dead_code)]
PolydatCall(String),
}
#[derive(Debug, Clone)]
pub(crate) enum PathSeg {
Field(String),
Index(usize),
Wildcard,
}
pub(crate) fn parse_path_expr(src: &str) -> Result<Vec<PathSeg>, String> {
let trimmed = src.trim();
if trimmed.is_empty() {
return Err("empty path".into());
}
let mut segs = Vec::new();
let mut cur = String::new();
let mut iter = trimmed.chars().peekable();
let push_field = |segs: &mut Vec<PathSeg>, cur: &mut String| {
if !cur.is_empty() {
if cur == "*" {
segs.push(PathSeg::Wildcard);
} else if let Ok(n) = cur.parse::<usize>() {
segs.push(PathSeg::Index(n));
} else {
segs.push(PathSeg::Field(std::mem::take(cur)));
}
cur.clear();
}
};
while let Some(&c) = iter.peek() {
match c {
'.' => {
push_field(&mut segs, &mut cur);
iter.next();
}
'[' => {
push_field(&mut segs, &mut cur);
iter.next();
let mut idx = String::new();
for c2 in iter.by_ref() {
if c2 == ']' {
break;
}
idx.push(c2);
}
if idx.trim() == "*" {
segs.push(PathSeg::Wildcard);
} else {
let n: usize = idx
.trim()
.parse()
.map_err(|_| format!("path '{src}': invalid index '[{idx}]'"))?;
segs.push(PathSeg::Index(n));
}
}
_ => {
cur.push(c);
iter.next();
}
}
}
push_field(&mut segs, &mut cur);
if segs.is_empty() {
return Err(format!("path '{src}': no segments"));
}
let wildcards = segs
.iter()
.filter(|s| matches!(s, PathSeg::Wildcard))
.count();
if wildcards > 1 {
return Err(format!(
"path '{src}': at most one [*] wildcard per path \
(nested projection is not supported)"
));
}
Ok(segs)
}
fn wildcard_split(segs: &[PathSeg]) -> Option<(&[PathSeg], &[PathSeg])> {
let pos = segs.iter().position(|s| matches!(s, PathSeg::Wildcard))?;
Some((&segs[..pos], &segs[pos + 1..]))
}
fn resolve_path_projection<'a>(
json: &'a serde_json::Value,
prefix: &[PathSeg],
suffix: &[PathSeg],
) -> Option<Vec<&'a serde_json::Value>> {
let at = resolve_path(json, prefix)?;
let arr = at.as_array()?;
Some(
arr.iter()
.filter_map(|el| resolve_path(el, suffix))
.collect(),
)
}
fn leaf_as_i64(v: &serde_json::Value) -> Option<i64> {
v.as_i64().or_else(|| v.as_str()?.trim().parse().ok())
}
fn leaf_as_f64(v: &serde_json::Value) -> Option<f64> {
v.as_f64().or_else(|| v.as_str()?.trim().parse().ok())
}
pub(crate) fn evaluate_path_value(
json: &serde_json::Value,
segs: &[PathSeg],
target: Option<polydat::ast::PortType>,
) -> Option<polydat::ast::Value> {
if let Some((prefix, suffix)) = wildcard_split(segs) {
let collected = resolve_path_projection(json, prefix, suffix).unwrap_or_default();
return Some(coerce_projection(&collected, target));
}
resolve_path(json, segs).map(json_to_value)
}
fn coerce_projection(
collected: &[&serde_json::Value],
target: Option<polydat::ast::PortType>,
) -> polydat::ast::Value {
use polydat::ast::{PortType, SliceArc, Value};
let json_array = || {
Value::Json(std::sync::Arc::new(serde_json::Value::Array(
collected.iter().map(|v| (*v).clone()).collect(),
)))
};
match target {
Some(PortType::VecI64) => Value::VecI64(SliceArc::from_vec(
collected.iter().filter_map(|v| leaf_as_i64(v)).collect(),
)),
Some(PortType::VecI32) => Value::VecI32(SliceArc::from_vec(
collected
.iter()
.filter_map(|v| leaf_as_i64(v).map(|n| n as i32))
.collect(),
)),
Some(PortType::VecF64) => Value::VecF64(SliceArc::from_vec(
collected.iter().filter_map(|v| leaf_as_f64(v)).collect(),
)),
Some(PortType::VecF32) => Value::VecF32(SliceArc::from_vec(
collected
.iter()
.filter_map(|v| leaf_as_f64(v).map(|n| n as f32))
.collect(),
)),
Some(PortType::Json) => json_array(),
_ => {
let ints: Vec<i64> = collected.iter().filter_map(|v| leaf_as_i64(v)).collect();
if !collected.is_empty() && ints.len() == collected.len() {
return Value::VecI64(SliceArc::from_vec(ints));
}
let floats: Vec<f64> = collected.iter().filter_map(|v| leaf_as_f64(v)).collect();
if !collected.is_empty() && floats.len() == collected.len() {
return Value::VecF64(SliceArc::from_vec(floats));
}
json_array()
}
}
}
fn resolve_path<'a>(
json: &'a serde_json::Value,
segs: &[PathSeg],
) -> Option<&'a serde_json::Value> {
let mut cur = json;
for seg in segs {
cur = match (cur, seg) {
(serde_json::Value::Object(m), PathSeg::Field(k)) => m.get(k)?,
(serde_json::Value::Array(a), PathSeg::Index(i)) => a.get(*i)?,
_ => return None,
};
}
Some(cur)
}
fn decode_slot(
name: &str,
raw_source: &str,
target: Option<polydat::ast::PortType>,
) -> Option<ResultSlot> {
let raw = raw_source.trim();
let source = if raw == "count" {
ResultSource::Count
} else if raw == "ok" {
ResultSource::Ok
} else if raw.contains('(') {
crate::diag!(
crate::observer::LogLevel::Warn,
"result wire '{name}': GK-expression source '{raw}' is not yet \
evaluated end-to-end — slot will resolve to its default. \
SRD-66 Push 2 follow-up wires the kernel-driven path.",
);
ResultSource::PolydatCall(raw.to_string())
} else {
match parse_path_expr(raw) {
Ok(segs) => ResultSource::Path(segs),
Err(e) => {
crate::diag!(
crate::observer::LogLevel::Warn,
"result wire '{name}': source '{raw}' is not parseable as a \
path expression ({e}) — slot will be skipped.",
);
return None;
}
}
};
Some(ResultSlot {
wire: name.to_string(),
source,
default: None,
target,
})
}
impl ResultDispenser {
pub fn wrap(
inner: Arc<dyn OpDispenser>,
result_spec: Option<&nmbrs_workload::model::ResultSpec>,
interface_results: Option<&std::collections::BTreeMap<String, String>>,
) -> Arc<dyn OpDispenser> {
let mut specs: Vec<ResultSlot> = Vec::new();
let target_for = |name: &str| -> Option<polydat::ast::PortType> {
interface_results?
.get(name)
.and_then(|kw| polydat::ast::PortType::from_keyword(kw))
};
if let Some(spec) = result_spec {
spec.walk_fragments(|frag| match frag {
nmbrs_workload::model::ResultFragment::Named { name, source } => {
let raw = source.trim();
if raw == "count" || raw == "ok" {
if let Some(slot) = decode_slot(name, source, target_for(name)) {
specs.push(slot);
}
} else if !raw.contains('(') {
if let Some(slot) = decode_slot(name, source, target_for(name)) {
specs.push(slot);
}
}
}
nmbrs_workload::model::ResultFragment::Source(_source) => {
}
});
}
specs.sort_by(|a, b| a.wire.cmp(&b.wire));
Arc::new(Self {
inner,
specs,
populate_kernel_inputs: true,
})
}
fn evaluate(slot: &ResultSlot, result: &OpResult) -> Option<polydat::ast::Value> {
match &slot.source {
ResultSource::Count => {
let n = result.body.as_ref().map(|b| b.element_count()).unwrap_or(0);
Some(polydat::ast::Value::U64(n))
}
ResultSource::Ok => {
Some(polydat::ast::Value::Bool(true))
}
ResultSource::Path(segs) => {
let body = result.body.as_ref()?;
let json = body.to_json();
evaluate_path_value(&json, segs, slot.target).or_else(|| slot.default.clone())
}
ResultSource::PolydatCall(_) => slot.default.clone(),
}
}
}
impl WrappingDispenser for ResultDispenser {}
impl OpDispenser for ResultDispenser {
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?;
if result.skipped {
return Ok(result);
}
for slot in &self.specs {
if let Some(v) = Self::evaluate(slot, &result) {
let _ = ctx.wires.write(&slot.wire, v);
}
}
if self.populate_kernel_inputs {
let count = result.body.as_ref().map(|b| b.element_count()).unwrap_or(0);
let body_json = result
.body
.as_ref()
.map(|b| b.to_json())
.unwrap_or(serde_json::Value::Null);
let _ = ctx.wires.write(
"body",
polydat::ast::Value::Json(std::sync::Arc::new(body_json)),
);
let _ = ctx.wires.write("count", polydat::ast::Value::U64(count));
let _ = ctx.wires.write("ok", polydat::ast::Value::Bool(true));
}
let _ = cycle;
Ok(result)
})
}
fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
Some(self.inner.as_ref())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::{AdapterError, ExecutionError, OpResult, ResultBody};
use crate::fixture::{ExecCtx, ResolvedPulls};
use std::sync::Arc;
#[derive(Debug)]
struct ResultDispBody {
value: serde_json::Value,
count: u64,
}
impl ResultBody for ResultDispBody {
fn to_json(&self) -> serde_json::Value {
self.value.clone()
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn element_count(&self) -> u64 {
self.count
}
}
struct FakeInner {
body: Option<ResultDispBody>,
error: Option<&'static str>,
}
impl OpDispenser for FakeInner {
fn execute<'a>(
&'a self,
_cycle: u64,
_ctx: &'a ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
if let Some(msg) = self.error {
return Err(ExecutionError::Op(AdapterError {
error_name: "test".into(),
message: msg.into(),
retryable: false,
}));
}
Ok(OpResult {
body: self.body.as_ref().map(|b| {
Box::new(ResultDispBody {
value: b.value.clone(),
count: b.count,
}) as Box<dyn ResultBody>
}),
skipped: false,
})
})
}
}
fn empty_ctx() -> (crate::adapter::ResolvedFields, ResolvedPulls) {
let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
let pulls = ResolvedPulls::empty();
(fields, pulls)
}
fn kernel_with_extern_inputs(names: &[(&str, &str)]) -> crate::scope_kernel::ScopeKernel {
let mut src = String::from("input cycle: u64\n");
for (n, ty) in names {
src.push_str(&format!("extern {n}: {ty}\n"));
}
let mut k = crate::bindings::compile_scope_kernel(&src, &Default::default())
.expect("kernel_with_extern_inputs compile");
k.set_inputs(&[0]);
k
}
fn run_with_wires(
dispenser: Arc<dyn OpDispenser>,
kernel: &mut crate::scope_kernel::ScopeKernel,
) -> Result<OpResult, ExecutionError> {
let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
let pulls = ResolvedPulls::empty();
let cw = crate::wires::CycleWires::new(kernel);
let ctx = ExecCtx::with_wires(&fields, &pulls, &cw);
let rt = tokio::runtime::Builder::new_current_thread()
.build()
.unwrap();
rt.block_on(dispenser.execute(0, &ctx))
}
#[test]
fn parse_path_dotted_and_bracketed_equivalent() {
let a = parse_path_expr("rows[0].value").unwrap();
let b = parse_path_expr("rows.0.value").unwrap();
assert_eq!(a.len(), 3);
assert_eq!(b.len(), 3);
match (&a[0], &b[0]) {
(PathSeg::Field(f1), PathSeg::Field(f2)) => assert_eq!(f1, f2),
_ => panic!("expected leading field"),
}
match (&a[1], &b[1]) {
(PathSeg::Index(0), PathSeg::Index(0)) => {}
_ => panic!("expected index 0"),
}
}
fn map_spec(entries: &[(&str, &str)]) -> nmbrs_workload::model::ResultSpec {
let mut m = std::collections::BTreeMap::new();
for (k, v) in entries {
m.insert((*k).to_string(), (*v).to_string());
}
nmbrs_workload::model::ResultSpec::Map(m)
}
#[test]
fn result_dispenser_count_and_path() {
let inner = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!({"rows": [{"value": 42}]}),
count: 1,
}),
error: None,
});
let decl = map_spec(&[("row_count", "count"), ("first_value", "rows[0].value")]);
let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
let mut kernel = kernel_with_extern_inputs(&[("row_count", "u64"), ("first_value", "u64")]);
let _ = run_with_wires(wrapped, &mut kernel).unwrap();
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
assert_eq!(w.get("row_count").map(|v| v.as_u64()), Some(1));
assert_eq!(w.get("first_value").map(|v| v.as_u64()), Some(42));
}
#[test]
fn result_dispenser_ok_builtin_on_success() {
let inner = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!({}),
count: 0,
}),
error: None,
});
let decl = map_spec(&[("succeeded", "ok")]);
let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
let mut kernel = kernel_with_extern_inputs(&[("succeeded", "bool")]);
let _ = run_with_wires(wrapped, &mut kernel).unwrap();
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
match w.get("succeeded") {
Some(polydat::ast::Value::Bool(b)) => assert!(b),
other => panic!("expected Bool(true), got {other:?}"),
}
}
#[test]
fn result_dispenser_error_propagates_no_capture_write() {
let inner = Arc::new(FakeInner {
body: None,
error: Some("boom"),
});
let decl = map_spec(&[("succeeded", "ok")]);
let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
let (fields, pulls) = empty_ctx();
let ctx = ExecCtx::new(&fields, &pulls);
let rt = tokio::runtime::Builder::new_current_thread()
.build()
.unwrap();
let err = rt.block_on(wrapped.execute(0, &ctx)).unwrap_err();
assert!(format!("{err}").contains("boom"));
}
#[test]
fn result_dispenser_unresolved_path_skips_silently() {
let inner = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!({"rows": []}),
count: 0,
}),
error: None,
});
let decl = map_spec(&[("missing", "rows[0].value")]);
let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
let mut kernel = kernel_with_extern_inputs(&[("missing", "u64")]);
let _ = run_with_wires(wrapped, &mut kernel).unwrap();
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
assert!(matches!(
w.get("missing"),
Some(polydat::ast::Value::None) | None
));
}
#[test]
fn result_dispenser_always_wraps_per_srd_40b() {
let inner: Arc<dyn OpDispenser> = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!({}),
count: 0,
}),
error: None,
});
let inner_ptr = Arc::as_ptr(&inner);
let wrapped = ResultDispenser::wrap(inner.clone(), None, None);
assert_ne!(
Arc::as_ptr(&wrapped),
inner_ptr,
"ResultDispenser must always wrap so magic-extern population fires"
);
}
#[test]
fn result_dispenser_skipped_op_writes_no_captures() {
struct SkipInner;
impl OpDispenser for SkipInner {
fn execute<'a>(
&'a self,
_cycle: u64,
_ctx: &'a ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move { Ok(OpResult::skipped()) })
}
}
let decl = map_spec(&[("c", "count")]);
let wrapped = ResultDispenser::wrap(Arc::new(SkipInner), Some(&decl), None);
let mut kernel = kernel_with_extern_inputs(&[("c", "u64")]);
let result = run_with_wires(wrapped, &mut kernel).unwrap();
assert!(result.skipped);
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
assert!(matches!(w.get("c"), Some(polydat::ast::Value::None) | None));
}
#[test]
fn decode_slot_reachable() {
let slot = decode_slot("c", "count", None).expect("count slot decodes");
assert!(matches!(slot.source, ResultSource::Count));
}
#[test]
fn parse_wildcard_forms_and_reject_nested() {
for src in [
"rows[*].key",
"rows.*.key",
"[*].key",
"[*]",
"result[*].id",
] {
let segs = parse_path_expr(src).unwrap_or_else(|e| panic!("'{src}' should parse: {e}"));
assert_eq!(
segs.iter()
.filter(|s| matches!(s, PathSeg::Wildcard))
.count(),
1,
"'{src}' carries one wildcard"
);
}
let err = parse_path_expr("rows[*].items[*].id").unwrap_err();
assert!(err.contains("at most one"), "err: {err}");
}
#[test]
fn projection_collects_column_as_target_type() {
let inner = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!([
{"key": "000000000007"}, {"key": "000000000003"}]),
count: 2,
}),
error: None,
});
let decl = map_spec(&[("keys", "[*].key")]);
let mut iface = std::collections::BTreeMap::new();
iface.insert("keys".to_string(), "vec_i64".to_string());
let wrapped = ResultDispenser::wrap(inner, Some(&decl), Some(&iface));
let mut kernel = kernel_with_extern_inputs(&[("keys", "vec_i64")]);
let _ = run_with_wires(wrapped, &mut kernel).unwrap();
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
match w.get("keys") {
Some(polydat::ast::Value::VecI64(slice)) => assert_eq!(slice.as_slice(), &[7, 3]),
other => panic!("expected VecI64([7,3]), got {other:?}"),
}
}
#[test]
fn projection_nested_prefix_and_inference_ladder() {
let inner = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!({"result": [
{"id": 12, "score": 0.9},
{"id": 5, "score": 0.7}]}),
count: 2,
}),
error: None,
});
let decl = map_spec(&[("ids", "result[*].id"), ("scores", "result[*].score")]);
let wrapped = ResultDispenser::wrap(inner, Some(&decl), None);
let mut kernel = kernel_with_extern_inputs(&[("ids", "vec_i64"), ("scores", "vec_f64")]);
let _ = run_with_wires(wrapped, &mut kernel).unwrap();
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
match w.get("ids") {
Some(polydat::ast::Value::VecI64(slice)) => assert_eq!(slice.as_slice(), &[12, 5]),
other => panic!("expected VecI64, got {other:?}"),
}
match w.get("scores") {
Some(polydat::ast::Value::VecF64(slice)) => assert_eq!(slice.as_slice(), &[0.9, 0.7]),
other => panic!("expected VecF64, got {other:?}"),
}
}
#[test]
fn projection_on_missing_prefix_lands_empty_not_stale() {
let inner = Arc::new(FakeInner {
body: Some(ResultDispBody {
value: serde_json::json!({"error": "no result key"}),
count: 0,
}),
error: None,
});
let decl = map_spec(&[("keys", "result[*].id")]);
let mut iface = std::collections::BTreeMap::new();
iface.insert("keys".to_string(), "vec_i64".to_string());
let wrapped = ResultDispenser::wrap(inner, Some(&decl), Some(&iface));
let mut kernel = kernel_with_extern_inputs(&[("keys", "vec_i64")]);
let _ = run_with_wires(wrapped, &mut kernel).unwrap();
let cw = crate::wires::CycleWires::new(&mut kernel);
let w: &dyn crate::wires::WireSource = &cw;
match w.get("keys") {
Some(polydat::ast::Value::VecI64(slice)) => assert!(
slice.as_slice().is_empty(),
"missing projection target must land the EMPTY typed form"
),
other => panic!("expected empty VecI64, got {other:?}"),
}
}
}