#![cfg(feature = "thread-safe")]
use std::collections::HashSet;
use std::fs;
use lazily::{ThreadSafeComputedMap, ThreadSafeContext};
use serde_json::Value;
const SPEC_DIR: &str = "../lazily-spec/conformance/materialization";
type V = i64;
fn present() -> bool {
std::path::Path::new(SPEC_DIR).exists()
}
fn load(name: &str) -> Value {
let path = format!("{SPEC_DIR}/{name}");
serde_json::from_str(&fs::read_to_string(&path).unwrap_or_else(|e| panic!("read {path}: {e}")))
.unwrap_or_else(|e| panic!("parse {path}: {e}"))
}
fn val_entries(fixture: &Value) -> Vec<(String, V)> {
fixture
.get("spec")
.and_then(|s| s.get("val"))
.and_then(|v| v.as_object())
.expect("spec.val")
.iter()
.map(|(k, v)| (k.clone(), v.as_i64().expect("int val")))
.collect()
}
fn str_array(v: &Value, path: &str) -> Vec<String> {
v.get(path)
.and_then(|v| v.as_array())
.unwrap_or_else(|| panic!("missing array {path}"))
.iter()
.map(|k| k.as_str().expect("string").to_string())
.collect()
}
fn as_set(keys: &[String]) -> HashSet<String> {
keys.iter().cloned().collect()
}
fn lookup_fn(entries: Vec<(String, V)>) -> impl Fn(&String) -> V + Clone + Send + Sync + 'static {
move |k: &String| -> V {
entries
.iter()
.find(|(key, _)| key == k)
.map(|(_, v)| *v)
.unwrap_or_else(|| panic!("no val for {k}"))
}
}
fn eager_computed_map(
ctx: &ThreadSafeContext,
keys: Vec<String>,
entries: Vec<(String, V)>,
) -> ThreadSafeComputedMap<String, V> {
let map: ThreadSafeComputedMap<String, V> = ThreadSafeComputedMap::new(ctx);
map.materialize_all(ctx, keys, lookup_fn(entries));
map
}
fn check_val_fixture(name: &str) -> Value {
let fixture = load(name);
let entries = val_entries(&fixture);
let keys: Vec<String> = entries.iter().map(|(k, _)| k.clone()).collect();
let expected = fixture.get("expected").expect("expected");
assert_eq!(
expected.get("default_mode").and_then(|v| v.as_str()),
Some("eager")
);
let ctx = ThreadSafeContext::new();
let eager = eager_computed_map(&ctx, keys.clone(), entries.clone());
let lazy: ThreadSafeComputedMap<String, V> = ThreadSafeComputedMap::new(&ctx);
let lookup = lookup_fn(entries);
assert_eq!(eager.present_count(), keys.len());
assert_eq!(
as_set(&eager.present_keys()),
as_set(&str_array(expected, "eager_present"))
);
assert_eq!(lazy.present_count(), 0);
for (k, want) in expected.get("observe").and_then(|v| v.as_object()).unwrap() {
let want = want.as_i64().unwrap();
assert_eq!(eager.observe(&ctx, k).unwrap(), want, "eager observe {k}");
assert_eq!(
lazy.get_or_insert_with(&ctx, k.clone(), lookup.clone()),
want,
"lazy observe {k}"
);
}
fixture
}
#[test]
fn observational_transparency_thread_safe() {
if !present() {
eprintln!("skipping: {SPEC_DIR} absent");
return;
}
let fixture = check_val_fixture("observational_transparency.json");
let expected = fixture.get("expected").unwrap();
let entries = val_entries(&fixture);
let ctx = ThreadSafeContext::new();
let lazy: ThreadSafeComputedMap<String, V> = ThreadSafeComputedMap::new(&ctx);
let lookup = lookup_fn(entries);
for k in str_array(&fixture, "reads") {
lazy.get_or_insert_with(&ctx, k, lookup.clone());
}
assert_eq!(
as_set(&lazy.present_keys()),
as_set(&str_array(expected, "lazy_present_after_reads"))
);
}
#[test]
fn deferral_not_deallocation_thread_safe() {
if !present() {
eprintln!("skipping: {SPEC_DIR} absent");
return;
}
let fixture = check_val_fixture("deferral_not_deallocation.json");
let expected = fixture.get("expected").unwrap();
let entries = val_entries(&fixture);
let ctx = ThreadSafeContext::new();
let lazy: ThreadSafeComputedMap<String, V> = ThreadSafeComputedMap::new(&ctx);
let lookup = lookup_fn(entries);
let want_sizes: Vec<usize> = expected
.get("present_after_each_read")
.and_then(|v| v.as_array())
.expect("present_after_each_read")
.iter()
.map(|n| n.as_u64().unwrap() as usize)
.collect();
let mut got = Vec::new();
for k in str_array(&fixture, "reads") {
lazy.get_or_insert_with(&ctx, k, lookup.clone());
got.push(lazy.present_count());
}
assert_eq!(got, want_sizes);
let lazy_present = as_set(&lazy.present_keys());
assert!(lazy_present.is_subset(&as_set(&str_array(expected, "eager_present"))));
}
#[test]
fn materialization_confluent_under_reordering() {
if !present() {
eprintln!("skipping: {SPEC_DIR} absent");
return;
}
let fixture = load("observational_transparency.json");
let entries = val_entries(&fixture);
let keys: Vec<String> = entries.iter().map(|(k, _)| k.clone()).collect();
let ctx_fwd = ThreadSafeContext::new();
let fwd: ThreadSafeComputedMap<String, V> = ThreadSafeComputedMap::new(&ctx_fwd);
let ctx_rev = ThreadSafeContext::new();
let rev: ThreadSafeComputedMap<String, V> = ThreadSafeComputedMap::new(&ctx_rev);
let lookup = lookup_fn(entries);
let mut fwd_vals = Vec::new();
for k in keys.iter() {
fwd_vals.push((
k.clone(),
fwd.get_or_insert_with(&ctx_fwd, k.clone(), lookup.clone()),
));
}
let mut rev_vals = Vec::new();
for k in keys.iter().rev() {
rev_vals.push((
k.clone(),
rev.get_or_insert_with(&ctx_rev, k.clone(), lookup.clone()),
));
}
for (k, v) in &fwd_vals {
let rv = rev_vals
.iter()
.find(|(rk, _)| rk == k)
.map(|(_, v)| *v)
.unwrap();
assert_eq!(*v, rv, "observe {k} order-independent");
}
assert_eq!(as_set(&fwd.present_keys()), as_set(&rev.present_keys()));
}