use std::sync::Mutex;
use std::sync::atomic::Ordering;
use std::time::Duration;
use crate::model::bounded::BoundedLru;
use crate::Result;
use zenkey::schema::decode::{DecodeError, DecodedPayload, DecoderRegistry};
use zenkey::schema::validate::{NotValidated, Verdict};
use zenkey::schema::{SchemaSet, TypeSchema, WireEncoding};
use zenoh::Session;
use crate::model::registry::SliceSet;
use crate::report::{DriftVerdict, SchemaDrift, SchemaServer, TotalityGap};
pub const DEFAULT_MAX_PRODUCERS: usize = 1_024;
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct StoreBounds {
pub max_producers: usize,
pub producers: usize,
pub sets_evicted: u64,
pub queriers_evicted: u64,
pub gates_evicted: u64,
}
#[derive(Debug)]
struct Entry<V> {
value: V,
seen: u64,
}
pub struct SchemaStore {
base: String,
timeout: Duration,
sets: Mutex<BoundedLru<String, Entry<Cached>>>,
queriers: Mutex<BoundedLru<String, Entry<std::sync::Arc<crate::bus::query::RepeatingQuery>>>>,
inflight: Mutex<BoundedLru<String, Entry<std::sync::Arc<tokio::sync::Mutex<()>>>>>,
clock: std::sync::atomic::AtomicU64,
sets_evicted: std::sync::atomic::AtomicU64,
queriers_evicted: std::sync::atomic::AtomicU64,
gates_evicted: std::sync::atomic::AtomicU64,
decoders: std::sync::RwLock<DecoderRegistry>,
sealed: std::sync::atomic::AtomicBool,
}
pub struct Sealed<'a> {
store: &'a SchemaStore,
}
impl Drop for Sealed<'_> {
fn drop(&mut self) {
self.store
.sealed
.store(false, std::sync::atomic::Ordering::Release);
}
}
const NOT_SERVED_TTL: Duration = Duration::from_secs(60);
const NO_REPLY_BACKOFF: Duration = Duration::from_millis(250);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum MissReason {
NoReplies,
AnsweredUnusable,
}
#[derive(Debug, Clone, Copy)]
struct Missing {
reason: MissReason,
asked: std::time::Instant,
attempts: u32,
}
impl Missing {
fn backoff(&self) -> Duration {
match self.reason {
MissReason::AnsweredUnusable => NOT_SERVED_TTL,
MissReason::NoReplies => NO_REPLY_BACKOFF
.saturating_mul(1u32 << self.attempts.saturating_sub(1).min(16))
.min(NOT_SERVED_TTL),
}
}
fn may_reask(&self) -> bool {
self.asked.elapsed() >= self.backoff()
}
}
enum Cached {
Served(std::sync::Arc<SchemaSet>),
Missing(Missing),
}
enum Lookup {
Answered(Option<std::sync::Arc<SchemaSet>>),
Ask(u32),
}
enum Fetched {
Served(SchemaSet),
NoReplies,
AnsweredUnusable,
}
impl SchemaStore {
pub fn new(base: impl Into<String>, timeout: Duration) -> Self {
SchemaStore::bounded(base, timeout, DEFAULT_MAX_PRODUCERS)
}
pub fn bounded(base: impl Into<String>, timeout: Duration, max_producers: usize) -> Self {
SchemaStore {
base: base.into(),
timeout,
sets: Mutex::new(BoundedLru::with_capacity(max_producers)),
queriers: Mutex::new(BoundedLru::with_capacity(max_producers)),
inflight: Mutex::new(BoundedLru::with_capacity(max_producers)),
decoders: std::sync::RwLock::new(DecoderRegistry::new()),
sealed: std::sync::atomic::AtomicBool::new(false),
clock: std::sync::atomic::AtomicU64::new(0),
sets_evicted: std::sync::atomic::AtomicU64::new(0),
queriers_evicted: std::sync::atomic::AtomicU64::new(0),
gates_evicted: std::sync::atomic::AtomicU64::new(0),
}
}
fn tick(&self) -> u64 {
self.clock.fetch_add(1, Ordering::Relaxed)
}
pub fn bounds(&self) -> StoreBounds {
let sets = self.sets.lock().expect("store lock");
StoreBounds {
max_producers: sets.max_keys(),
producers: sets.len(),
sets_evicted: self.sets_evicted.load(Ordering::Relaxed),
queriers_evicted: self.queriers_evicted.load(Ordering::Relaxed),
gates_evicted: self.gates_evicted.load(Ordering::Relaxed),
}
}
pub fn seal(&self) -> Sealed<'_> {
self.sealed
.store(true, std::sync::atomic::Ordering::Release);
Sealed { store: self }
}
pub fn register_decoder(&self, decoder: Box<dyn zenkey::schema::decode::PayloadDecoder>) {
self.decoders
.write()
.expect("decoder lock")
.register(decoder);
}
pub fn insert(&self, producer: impl Into<String>, set: SchemaSet) {
self.remember(producer.into(), Cached::Served(std::sync::Arc::new(set)));
}
fn remember(&self, producer: String, cached: Cached) {
let seen = self.tick();
let mut sets = self.sets.lock().expect("store lock");
if sets.get(producer.as_str()).is_none() {
let dropped = sets.admit(|e| e.seen) as u64;
if dropped > 0 {
self.sets_evicted.fetch_add(dropped, Ordering::Relaxed);
}
}
sets.insert(
producer,
Entry {
value: cached,
seen,
},
);
}
pub async fn schema_for(
&self,
session: &Session,
producer: &str,
type_name: &str,
) -> Option<TypeSchema> {
self.set_for(session, producer)
.await
.and_then(|set| set.get(type_name).cloned())
}
pub async fn set_for(
&self,
session: &Session,
producer: &str,
) -> Option<std::sync::Arc<SchemaSet>> {
let may_ask = !self.sealed.load(std::sync::atomic::Ordering::Acquire);
self.set_for_within(session, producer, may_ask).await
}
async fn set_for_within(
&self,
session: &Session,
producer: &str,
may_ask: bool,
) -> Option<std::sync::Arc<SchemaSet>> {
if let Lookup::Answered(hit) = self.lookup(producer) {
return hit;
}
if !may_ask {
return None;
}
let gate = self.gate_for(producer);
let _held = gate.lock().await;
let attempts = match self.lookup(producer) {
Lookup::Answered(hit) => return hit,
Lookup::Ask(attempts) => attempts,
};
let entry = match self.fetch(session, producer).await {
Fetched::Served(set) => Cached::Served(std::sync::Arc::new(set)),
Fetched::NoReplies => Cached::Missing(Missing {
reason: MissReason::NoReplies,
asked: std::time::Instant::now(),
attempts: attempts.saturating_add(1),
}),
Fetched::AnsweredUnusable => Cached::Missing(Missing {
reason: MissReason::AnsweredUnusable,
asked: std::time::Instant::now(),
attempts: 0,
}),
};
let served = match &entry {
Cached::Served(set) => Some(std::sync::Arc::clone(set)),
Cached::Missing(_) => None,
};
self.remember(producer.to_string(), entry);
served
}
fn gate_for(&self, producer: &str) -> std::sync::Arc<tokio::sync::Mutex<()>> {
let seen = self.tick();
let mut inflight = self.inflight.lock().expect("inflight lock");
if let Some(entry) = inflight.get_mut(producer) {
entry.seen = seen;
return std::sync::Arc::clone(&entry.value);
}
let dropped = inflight.admit(|e| e.seen) as u64;
if dropped > 0 {
self.gates_evicted.fetch_add(dropped, Ordering::Relaxed);
}
let gate = std::sync::Arc::new(tokio::sync::Mutex::new(()));
inflight.insert(
producer.to_string(),
Entry {
value: std::sync::Arc::clone(&gate),
seen,
},
);
gate
}
fn lookup(&self, producer: &str) -> Lookup {
let seen = self.tick();
let mut sets = self.sets.lock().expect("store lock");
let Some(entry) = sets.get_mut(producer) else {
return Lookup::Ask(0);
};
entry.seen = seen;
match &entry.value {
Cached::Served(set) => Lookup::Answered(Some(std::sync::Arc::clone(set))),
Cached::Missing(m) if !m.may_reask() => Lookup::Answered(None),
Cached::Missing(m) => Lookup::Ask(m.attempts),
}
}
pub fn forget(&self, producer: &str) {
self.sets.lock().expect("store lock").remove(producer);
}
pub fn forget_all(&self) {
self.sets.lock().expect("store lock").clear();
}
pub fn known(&self) -> Vec<(String, bool)> {
let sets = self.sets.lock().expect("store lock");
let mut out: Vec<(String, bool)> = sets
.iter()
.map(|(p, e)| (p.clone(), matches!(e.value, Cached::Served(_))))
.collect();
out.sort();
out
}
async fn fetch(&self, session: &Session, producer: &str) -> Fetched {
let cached = {
let seen = self.tick();
let mut queriers = self.queriers.lock().expect("querier lock");
queriers.get_mut(producer).map(|e| {
e.seen = seen;
std::sync::Arc::clone(&e.value)
})
};
let querier = match cached {
Some(q) => q,
None => {
let key = zenkey::grammar::with_base(
&self.base,
zenkey::selector::fleet_rpc(producer, &["describe"]),
);
let fleet = crate::Fleet::new(session, &self.base);
let declared =
match crate::bus::query::declare_repeating(&fleet, &key, self.timeout).await {
Ok(q) => std::sync::Arc::new(q),
Err(_) => return Fetched::NoReplies,
};
let seen = self.tick();
let mut queriers = self.queriers.lock().expect("querier lock");
if let Some(entry) = queriers.get_mut(producer) {
entry.seen = seen;
std::sync::Arc::clone(&entry.value)
} else {
let dropped = queriers.admit(|e| e.seen);
if dropped > 0 {
self.queriers_evicted
.fetch_add(dropped as u64, Ordering::Relaxed);
}
queriers.insert(
producer.to_string(),
Entry {
value: std::sync::Arc::clone(&declared),
seen,
},
);
declared
}
}
};
let Ok(answers) = querier.fetch().await else {
return Fetched::NoReplies;
};
if answers.is_empty() {
return Fetched::NoReplies;
}
for a in answers {
if let crate::bus::query::Answer::Value(bytes) = a.answer {
let cow = bytes.to_bytes();
if let Ok(text) = std::str::from_utf8(&cow)
&& let Ok(set) = SchemaSet::parse(text)
{
return Fetched::Served(set);
}
}
}
Fetched::AnsweredUnusable
}
pub fn decode(
&self,
schema: &TypeSchema,
encoding: &WireEncoding,
bytes: &[u8],
) -> Result<DecodedPayload, DecodeError> {
self.decoders
.read()
.expect("decoder lock")
.decode(schema, encoding, bytes)
}
pub fn encode(
&self,
schema: &TypeSchema,
value: &serde_json::Value,
target: &WireEncoding,
) -> Result<Vec<u8>, DecodeError> {
self.decoders
.read()
.expect("decoder lock")
.encode(schema, value, target)
}
}
fn referenced_types(slice: &zenkey::slice::RegistrySlice) -> Vec<String> {
let mut names: Vec<&str> = slice
.subjects
.iter()
.map(|s| s.type_name.as_str())
.filter(|t| !t.is_empty())
.collect();
for p in &slice.procedures {
names.extend(p.request.as_deref());
names.extend(p.reply.as_deref());
}
for b in &slice.blob {
names.extend(b.reference.as_deref());
}
for m in &slice.media {
names.extend(m.attachment.as_deref());
}
names.sort_unstable();
names.dedup();
names.into_iter().map(str::to_string).collect()
}
fn row(
producer: &str,
type_name: &str,
schema: &TypeSchema,
full: bool,
) -> crate::report::SchemaRow {
crate::report::SchemaRow {
producer: producer.to_string(),
type_name: type_name.to_string(),
kind: schema.kind_str().to_string(),
hash: schema.hash().unwrap_or_default().to_string(),
document: full.then(|| schema_document(schema)),
}
}
fn schema_document(schema: &TypeSchema) -> serde_json::Value {
if let Some(doc) = schema.json_document() {
return doc.clone();
}
let mut obj = serde_json::Map::new();
obj.insert(
"kind".into(),
serde_json::Value::String(schema.kind_str().to_string()),
);
if let Some(m) = schema.protobuf_message() {
obj.insert("message".into(), serde_json::Value::String(m.to_string()));
}
if let Some(bytes) = schema.protobuf_descriptor_set() {
obj.insert(
"descriptor_set_bytes".into(),
serde_json::Value::from(bytes.len()),
);
}
if let Some(fields) = schema.cdr_fields() {
obj.insert("fields".into(), fields.clone());
}
if let Some(types) = schema.cdr_types() {
obj.insert("types".into(), serde_json::Value::Object(types.clone()));
}
serde_json::Value::Object(obj)
}
pub async fn schema_dump(
store: &SchemaStore,
session: &Session,
slices: Option<&SliceSet>,
producer: &str,
type_filter: Option<&str>,
full: bool,
) -> crate::report::SchemaDump {
let set = store.set_for(session, producer).await;
let Some(set) = set else {
return crate::report::SchemaDump {
producer: producer.to_string(),
served: false,
app: None,
types: Vec::new(),
missing: crate::report::Asked::NotAsked,
};
};
let types: Vec<crate::report::SchemaRow> = set
.iter()
.filter(|(name, _)| type_filter.is_none_or(|f| f == *name))
.map(|(name, schema)| row(producer, name, schema, full || type_filter.is_some()))
.collect();
let missing = slices.map(|slices| {
slices
.get(producer)
.map(|slice| {
referenced_types(slice)
.into_iter()
.filter(|n| set.get(n).is_none())
.collect()
})
.unwrap_or_default()
});
crate::report::SchemaDump {
producer: producer.to_string(),
served: true,
app: Some(set.app().to_string()),
types,
missing: missing.into(),
}
}
pub async fn schemas_for_type(
store: &SchemaStore,
session: &Session,
producers: &[String],
type_name: &str,
full: bool,
) -> Vec<crate::report::SchemaRow> {
let mut out = Vec::new();
for producer in producers {
if let Some(schema) = store.schema_for(session, producer, type_name).await {
out.push(row(producer, type_name, &schema, full));
}
}
out
}
pub fn schema_drift(described: &[(String, SchemaSet)]) -> Vec<SchemaDrift> {
use std::collections::BTreeMap;
let mut by_name: BTreeMap<&str, Vec<SchemaServer>> = BTreeMap::new();
for (producer, set) in described {
for (name, schema) in set.iter() {
by_name.entry(name).or_default().push(SchemaServer {
producer: producer.clone(),
hash: schema.hash().map(str::to_string).into(),
});
}
}
by_name
.into_iter()
.filter(|(_, servers)| servers.len() > 1)
.filter_map(|(name, servers)| {
let claimed: Vec<&String> = servers.iter().filter_map(|s| s.hash.as_option()).collect();
let verdict = if claimed.len() < servers.len() {
DriftVerdict::Unjudgeable
} else if claimed.iter().any(|h| *h != claimed[0]) {
DriftVerdict::Disagree
} else {
return None;
};
Some(SchemaDrift {
type_name: name.to_string(),
servers,
verdict,
})
})
.collect()
}
pub fn totality_gaps(described: &[(String, SchemaSet)], slices: &SliceSet) -> Vec<TotalityGap> {
let mut gaps = Vec::new();
for (producer, set) in described {
let Some(slice) = slices.get(producer) else {
continue;
};
let mut names: Vec<&str> = Vec::new();
names.extend(
slice
.subjects
.iter()
.map(|s| s.type_name.as_str())
.filter(|t| !t.is_empty()),
);
for p in &slice.procedures {
names.extend(p.request.as_deref());
names.extend(p.reply.as_deref());
}
for b in &slice.blob {
names.extend(b.reference.as_deref());
}
names.sort();
names.dedup();
let missing: Vec<String> = names
.into_iter()
.filter(|n| set.get(n).is_none())
.map(str::to_string)
.collect();
if !missing.is_empty() {
gaps.push(TotalityGap {
producer: producer.clone(),
missing,
});
}
}
gaps
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Rendering {
Typed(DecodedPayload),
Structural(String),
}
pub fn resolve_encoding(
sample_encoding: Option<&str>,
registry_encoding: Option<&WireEncoding>,
bytes: &[u8],
) -> WireEncoding {
if let Some(e) = sample_encoding
&& e != "zenoh/bytes"
{
return WireEncoding::from_encoding_str(e);
}
if let Some(e) = registry_encoding {
return e.clone();
}
match bytes.first() {
Some(b'{' | b'[' | b'"') => WireEncoding::Json,
_ => WireEncoding::Cbor,
}
}
pub const OBSERVE_LIMIT: usize = 64 * 1024;
pub fn structural_value(bytes: &[u8]) -> Option<serde_json::Value> {
let looks_json = bytes.first().is_some_and(|b| {
matches!(
b,
b'{' | b'[' | b'"' | b'-' | b'0'..=b'9' | b't' | b'f' | b'n'
)
});
if looks_json && let Ok(v) = serde_json::from_slice::<serde_json::Value>(bytes) {
return Some(v);
}
let is_text = std::str::from_utf8(bytes).is_ok_and(|t| !t.is_empty());
if let Some(v) = cbor_whole(bytes)
&& !(is_text && is_scalar(&v))
&& let Ok(value) = serde_json::to_value(&v)
{
return Some(value);
}
None
}
pub fn structural(bytes: &[u8]) -> String {
if let Some(v) = structural_value(bytes) {
return serde_json::to_string(&v).unwrap_or_default();
}
match std::str::from_utf8(bytes).ok().filter(|t| !t.is_empty()) {
Some(text) => text.to_string(),
None => format!("<{} bytes>", bytes.len()),
}
}
fn cbor_whole(bytes: &[u8]) -> Option<ciborium::Value> {
let mut cursor = std::io::Cursor::new(bytes);
let value = ciborium::from_reader::<ciborium::Value, _>(&mut cursor).ok()?;
(cursor.position() as usize == bytes.len()).then_some(value)
}
fn is_scalar(v: &ciborium::Value) -> bool {
!matches!(v, ciborium::Value::Map(_) | ciborium::Value::Array(_))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DecodedSample {
pub type_name: Option<String>,
pub rendering: Rendering,
pub verdict: Verdict,
pub decode_error: Option<String>,
}
impl DecodedSample {
fn structural(type_name: Option<String>, reason: NotValidated, bytes: &[u8]) -> DecodedSample {
DecodedSample {
type_name,
rendering: Rendering::Structural(structural(bytes)),
verdict: Verdict::NotValidated(reason),
decode_error: None,
}
}
}
pub async fn prewarm(
fleet: &crate::Fleet<'_>,
store: &SchemaStore,
slices: Option<&SliceSet>,
) -> usize {
let Some(slices) = slices else { return 0 };
let mut served = 0;
for slice in slices.slices() {
if store
.set_for_within(fleet.session(), &slice.name, true)
.await
.is_some()
{
served += 1;
}
}
served
}
pub async fn decode_sample(
fleet: &crate::Fleet<'_>,
store: &SchemaStore,
slices: Option<&SliceSet>,
wire_key: &str,
sample_encoding: Option<&str>,
bytes: &[u8],
) -> DecodedSample {
use zenkey::grammar::ClassOrPlane;
let (session, base) = (fleet.session(), fleet.base());
let Some(slices) = slices else {
return DecodedSample::structural(None, NotValidated::NoRegistry, bytes);
};
let refined = zenkey::grammar::parse_full(base, wire_key).and_then(|parsed| {
let producer = match (parsed.producer(), &parsed.origin) {
(Some(p), _) => p.name().to_string(),
(None, zenkey::grammar::Origin::Service(s)) => {
slices.by_service_origin(s.as_str())?.name.clone()
}
_ => return None,
};
let ClassOrPlane::Class(class) = parsed.class else {
return None;
};
let (subject, _) = slices.refine(&producer, class.chunk(), &parsed.subject)?;
Some((
producer,
subject.type_name.clone(),
subject.encoding.clone(),
))
});
let Some((producer, type_name, registry_encoding)) = refined else {
return DecodedSample::structural(None, NotValidated::NoSchema, bytes);
};
let encoding = resolve_encoding(sample_encoding, registry_encoding.as_ref(), bytes);
match store.schema_for(session, &producer, &type_name).await {
Some(schema) => match store.decode(&schema, &encoding, bytes) {
Ok(decoded) => {
let verdict = decoded.verdict.clone();
DecodedSample {
type_name: Some(type_name),
rendering: Rendering::Typed(decoded),
verdict,
decode_error: None,
}
}
Err(e) => DecodedSample {
type_name: Some(type_name),
rendering: Rendering::Structural(structural(bytes)),
verdict: Verdict::NotValidated(NotValidated::Undecodable),
decode_error: Some(e.to_string()),
},
},
None => DecodedSample::structural(Some(type_name), NotValidated::NoSchema, bytes),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_store_is_bounded_and_says_what_the_bound_cost() {
const PRODUCERS: usize = 200;
let set = || {
SchemaSet::parse(
r#"{"schema_version":1,"app":"t",
"types":{"W":{"kind":"cddl","hash":"sha256:00","spec":"x = int"}}}"#,
)
.expect("fixture parses")
};
let store = SchemaStore::bounded("", Duration::from_millis(1), 16);
for i in 0..PRODUCERS {
store.insert(format!("p{i:04}"), set());
}
let bounds = store.bounds();
assert_eq!(bounds.max_producers, 16);
assert!(bounds.producers <= 16, "the bound bit: {bounds:?}");
assert_eq!(
bounds.producers as u64 + bounds.sets_evicted,
PRODUCERS as u64,
"every producer is held or counted: {bounds:?}"
);
assert_eq!(store.known().len(), bounds.producers, "known() agrees");
assert_eq!(bounds.queriers_evicted, 0);
assert_eq!(bounds.gates_evicted, 0);
}
#[test]
fn a_producer_still_being_read_survives_the_bound() {
let set = || {
SchemaSet::parse(
r#"{"schema_version":1,"app":"t",
"types":{"W":{"kind":"cddl","hash":"sha256:00","spec":"x = int"}}}"#,
)
.expect("fixture parses")
};
let store = SchemaStore::bounded("", Duration::from_millis(1), 8);
store.insert("hot", set());
for i in 0..7 {
store.insert(format!("cold{i}"), set());
}
for i in 7..64 {
assert!(
matches!(store.lookup("hot"), Lookup::Answered(Some(_))),
"the hot producer was evicted at insert {i}"
);
store.insert(format!("cold{i}"), set());
}
assert!(store.bounds().sets_evicted > 0, "the bound did bite");
assert!(
store
.known()
.iter()
.any(|(p, served)| p == "hot" && *served),
"the producer in use survived: {:?}",
store.known()
);
}
#[test]
fn the_totality_set_counts_every_referenced_type() {
let slice = zenkey::slice::parse_slice(
r#"
[registry]
version = "1.0"
app = "acme"
convention = 1
[producer]
name = "netring"
[[subject]]
path = "health"
class = "state"
type = "Health"
[[procedure]]
path = "capture/trigger"
kind = "write"
request = "CaptureSpec"
reply = "Ack"
[[blob]]
tier = "artifact"
endpoints = ["manifest"]
reference = "PcapRef"
[[media]]
path = "front/video/h264"
encoding = "video/h264"
attachment = "FrameMeta"
"#,
)
.unwrap();
assert_eq!(
referenced_types(&slice),
["Ack", "CaptureSpec", "FrameMeta", "Health", "PcapRef"]
);
}
#[test]
fn encoding_resolution_order() {
assert_eq!(
resolve_encoding(Some("application/json"), Some(&WireEncoding::Cbor), b"x"),
WireEncoding::Json
);
assert_eq!(
resolve_encoding(Some("zenoh/bytes"), Some(&WireEncoding::Cbor), b"{"),
WireEncoding::Cbor
);
assert_eq!(
resolve_encoding(None, None, b"{\"a\":1}"),
WireEncoding::Json
);
assert_eq!(resolve_encoding(None, None, &[0xa1]), WireEncoding::Cbor);
}
#[test]
fn structural_rendering_is_honest() {
assert_eq!(structural(b"{\"a\":1}"), "{\"a\":1}");
let mut cbor = Vec::new();
ciborium::into_writer(&serde_json::json!({"x": 1}), &mut cbor).unwrap();
assert!(structural(&cbor).contains("\"x\""));
assert_eq!(structural(&[0xff, 0xfe, 0x00]), "<3 bytes>");
}
#[test]
fn structural_value_yields_documents_and_nothing_else() {
assert_eq!(
structural_value(br#"{"value":42.0}"#),
Some(serde_json::json!({"value": 42.0}))
);
let mut cbor = Vec::new();
ciborium::into_writer(&serde_json::json!({"x": 1}), &mut cbor).unwrap();
assert_eq!(structural_value(&cbor), Some(serde_json::json!({"x": 1})));
assert_eq!(structural_value(b"just a plain string"), None);
assert_eq!(structural_value(&[0xff, 0xfe, 0x00]), None);
assert_eq!(structural_value(b""), None);
}
#[test]
fn the_rendering_agrees_with_the_value() {
for payload in [
&br#"{"a":1}"#[..],
&b"[1,2,3]"[..],
&b"just a plain string"[..],
&[0xff, 0xfe, 0x00][..],
] {
if let Some(v) = structural_value(payload) {
assert_eq!(structural(payload), serde_json::to_string(&v).unwrap());
}
}
}
#[test]
fn plain_text_is_not_mistaken_for_cbor() {
assert_eq!(structural(b"just a plain string"), "just a plain string");
assert_eq!(
structural(b"a v2 key: not this convention"),
"a v2 key: not this convention"
);
for first in b'a'..=b'z' {
let mut payload = vec![first];
payload.extend_from_slice(b" some trailing words here");
let text = String::from_utf8(payload.clone()).unwrap();
assert_eq!(structural(&payload), text, "mangled {text:?}");
}
}
#[test]
fn an_exact_cbor_text_string_still_reads_as_text() {
let payload = b"just a plai";
assert!(cbor_whole(payload).is_some(), "setup: this is valid CBOR");
assert_eq!(structural(payload), "just a plai");
}
#[test]
fn structured_cbor_still_wins_over_text() {
let mut cbor = Vec::new();
ciborium::into_writer(&serde_json::json!({"ok": true}), &mut cbor).unwrap();
let rendered = structural(&cbor);
assert!(rendered.contains("\"ok\""), "{rendered}");
assert!(rendered.starts_with('{'), "{rendered}");
}
#[test]
fn cbor_must_account_for_every_byte() {
let mut cbor = Vec::new();
ciborium::into_writer(&serde_json::json!({"x": 1}), &mut cbor).unwrap();
assert!(cbor_whole(&cbor).is_some());
cbor.push(0x00);
assert!(cbor_whole(&cbor).is_none(), "trailing byte must reject");
}
fn set_with(name: &str, schema: serde_json::Value) -> SchemaSet {
SchemaSet::builder("app")
.entry(name, zenkey::schema::TypeSchema::json_schema(schema))
.build()
}
#[test]
fn drift_findings_name_every_server() {
let a = SchemaSet::builder("app")
.entry(
"T",
zenkey::schema::TypeSchema::json_schema(serde_json::json!({"type":"object"})),
)
.build();
let b = SchemaSet::builder("app")
.entry(
"T",
zenkey::schema::TypeSchema::json_schema(serde_json::json!({"type":"string"})),
)
.build();
let c = SchemaSet::builder("app")
.entry(
"T",
zenkey::schema::TypeSchema::json_schema(serde_json::json!({"type":"object"})),
)
.build();
let described = vec![
("p1".to_string(), a),
("p2".to_string(), b),
("p3".to_string(), c),
];
let drift = schema_drift(&described);
assert_eq!(drift.len(), 1);
assert_eq!(drift[0].type_name, "T");
assert_eq!(drift[0].servers.len(), 3, "every server is named");
assert_eq!(drift[0].verdict, DriftVerdict::Disagree);
assert_eq!(drift[0].servers[0].hash, drift[0].servers[2].hash);
assert_ne!(drift[0].servers[0].hash, drift[0].servers[1].hash);
let unhashed = |app: &str| {
SchemaSet::parse(&format!(
r#"{{"schema_version":1,"app":"{app}","types":{{"T":{{"kind":"json-schema","hash":"","schema":{{}}}}}}}}"#
))
.unwrap()
};
let silent = vec![
("p1".to_string(), unhashed("app")),
("p2".to_string(), unhashed("app")),
];
let drift = schema_drift(&silent);
assert_eq!(drift.len(), 1, "silence is reported, not read as agreement");
assert_eq!(drift[0].verdict, DriftVerdict::Unjudgeable);
assert!(
drift[0].servers.iter().all(|s| s.hash.is_not_asked()),
"and it names who did not say"
);
let mixed = vec![
("p1".to_string(), unhashed("app")),
(
"p2".to_string(),
SchemaSet::builder("app")
.entry(
"T",
zenkey::schema::TypeSchema::json_schema(
serde_json::json!({"type":"object"}),
),
)
.build(),
),
];
assert_eq!(schema_drift(&mixed)[0].verdict, DriftVerdict::Unjudgeable);
assert!(schema_drift(&[("p1".to_string(), unhashed("app"))]).is_empty());
let described = vec![
(
"p1".to_string(),
set_with("T", serde_json::json!({"type":"object"})),
),
(
"p3".to_string(),
set_with("T", serde_json::json!({"type":"object"})),
),
];
assert!(schema_drift(&described).is_empty());
}
#[test]
fn totality_gaps_check_only_served_producers() {
use zenkey::slice::{RegistrySlice, SubjectDecl};
let mut subject = SubjectDecl::new("cpu", zenkey::Class::Telemetry);
subject.type_name = "TelemetryPoint".into();
let mut slice = RegistrySlice::new("1", "a", "sysinfo");
slice.subjects = vec![subject];
let slices = crate::model::registry::SliceSet::from_slices(vec![slice]);
let incomplete = SchemaSet::builder("a")
.entry(
"Other",
zenkey::schema::TypeSchema::json_schema(serde_json::json!({"type":"object"})),
)
.build();
let gaps = totality_gaps(&[("sysinfo".to_string(), incomplete)], &slices);
assert_eq!(gaps.len(), 1);
assert_eq!(gaps[0].missing, ["TelemetryPoint"]);
assert!(totality_gaps(&[], &slices).is_empty());
}
#[test]
fn an_untyped_subject_is_not_a_totality_gap() {
use zenkey::slice::{RegistrySlice, SubjectDecl};
let mut subject = SubjectDecl::new("raw", zenkey::Class::Telemetry);
subject.type_name = String::new();
let mut slice = RegistrySlice::new("1", "a", "sysinfo");
slice.subjects = vec![subject];
let slices = crate::model::registry::SliceSet::from_slices(vec![slice]);
let served = SchemaSet::builder("a").build();
assert!(
totality_gaps(&[("sysinfo".to_string(), served)], &slices).is_empty(),
"empty type names must be filtered, not reported as gaps"
);
}
#[test]
fn a_zero_reply_ask_backs_off_fast_and_an_answered_one_does_not() {
let now = std::time::Instant::now();
let no_reply = |attempts| Missing {
reason: MissReason::NoReplies,
asked: now,
attempts,
};
assert_eq!(no_reply(1).backoff(), NO_REPLY_BACKOFF);
assert_eq!(no_reply(2).backoff(), NO_REPLY_BACKOFF * 2);
assert_eq!(no_reply(3).backoff(), NO_REPLY_BACKOFF * 4);
assert_eq!(no_reply(30).backoff(), NOT_SERVED_TTL);
let answered = Missing {
reason: MissReason::AnsweredUnusable,
asked: now,
attempts: 0,
};
assert_eq!(
answered.backoff(),
NOT_SERVED_TTL,
"a producer that answered and served nothing is asked once per TTL"
);
}
#[test]
fn the_first_reask_is_sub_second() {
let m = Missing {
reason: MissReason::NoReplies,
asked: std::time::Instant::now(),
attempts: 1,
};
assert!(m.backoff() < Duration::from_secs(1));
assert!(!m.may_reask(), "and not before it elapses");
}
}