use crate::report::SampleRow;
impl SampleRow {
pub fn of_key(key: &str, base: &str) -> SampleRow {
let parsed = zenkey::grammar::parse_full(base, key);
SampleRow {
key: key.to_string(),
origin: parsed.as_ref().map(|p| p.origin.chunk().to_string()),
subject: parsed.as_ref().map(|p| p.subject.join("/")),
..SampleRow::default()
}
}
pub fn with_wire(mut self, view: &crate::SampleView) -> SampleRow {
self.delete = view.kind == zenoh::sample::SampleKind::Delete;
if !view.encoding.is_empty() {
self.encoding = Some(view.encoding.clone());
}
self.timestamp = view.timestamp.map(|t| t.to_string());
self.qos = zenkey::qos::QosProfile::ALL
.into_iter()
.find(|p| view.qos_matches(*p))
.map(|p| p.name().to_string());
self
}
pub fn with_payload_bytes(mut self, payload: &[u8]) -> SampleRow {
self.bytes = Some(b64(payload));
self
}
pub fn to_line(&self) -> String {
serde_json::to_string(self).expect("a sample row serializes")
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct IngestRow {
pub key: String,
pub payload: Vec<u8>,
pub encoding: Option<String>,
pub qos: Option<String>,
pub delete: bool,
pub attachment: Option<Vec<u8>>,
}
fn value_bytes(v: &serde_json::Value) -> Vec<u8> {
match v {
serde_json::Value::String(s) => s.clone().into_bytes(),
other => serde_json::to_vec(other).unwrap_or_default(),
}
}
pub(crate) fn b64(bytes: &[u8]) -> String {
use base64::Engine as _;
base64::engine::general_purpose::STANDARD.encode(bytes)
}
fn b64_bytes(
obj: &serde_json::Map<String, serde_json::Value>,
field: &str,
) -> Result<Option<Vec<u8>>, String> {
use base64::Engine as _;
match obj.get(field) {
None | Some(serde_json::Value::Null) => Ok(None),
Some(serde_json::Value::String(s)) => base64::engine::general_purpose::STANDARD
.decode(s)
.map(Some)
.map_err(|e| format!("\"{field}\" is not base64: {e}")),
Some(_) => Err(format!("\"{field}\" is not a base64 string")),
}
}
pub fn parse_row(line: &str) -> Result<IngestRow, String> {
let v: serde_json::Value =
serde_json::from_str(line).map_err(|e| format!("not a JSON object: {e}"))?;
let obj = v.as_object().ok_or("not a JSON object")?;
let key = obj
.get("key")
.and_then(|k| k.as_str())
.ok_or("no \"key\" field")?
.to_string();
if key.is_empty() {
return Err("empty \"key\"".into());
}
let delete = obj.get("delete").and_then(|d| d.as_bool()).unwrap_or(false);
let payload = match (b64_bytes(obj, "bytes")?, obj.get("value")) {
(Some(raw), _) => raw,
(None, Some(serde_json::Value::Null) | None) if delete => Vec::new(),
(None, Some(serde_json::Value::Null) | None) => {
return Err("no \"value\" or \"bytes\" field (and not a delete row)".into());
}
(None, Some(v)) => value_bytes(v),
};
let attachment = match b64_bytes(obj, "attachment_b64")? {
Some(raw) => Some(raw),
None => obj.get("attachment").map(value_bytes),
};
Ok(IngestRow {
key,
payload,
encoding: obj
.get("encoding")
.and_then(|e| e.as_str())
.filter(|e| !e.is_empty())
.map(str::to_string),
qos: obj.get("qos").and_then(|q| q.as_str()).map(str::to_string),
delete,
attachment,
})
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum StreamLine {
Sample(IngestRow),
Meta(String),
}
pub fn parse_stream_line(line: &str) -> Result<StreamLine, String> {
if let Ok(serde_json::Value::Object(obj)) = serde_json::from_str::<serde_json::Value>(line)
&& let Some(tag) = obj.get("row").and_then(serde_json::Value::as_str)
&& tag != "sample"
{
return Ok(StreamLine::Meta(tag.to_string()));
}
parse_row(line).map(StreamLine::Sample)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_hand_written_row_reads_back_with_its_extras_ignored() {
let row = parse_row(
r#"{"key":"v1/h-1/state/p/health","origin":"h-1","type":"Health","typed":true,
"encoding":"application/json","timestamp":null,"delete":false,
"value":{"status":"ok"}}"#,
)
.unwrap();
assert_eq!(row.key, "v1/h-1/state/p/health");
assert_eq!(row.payload, br#"{"status":"ok"}"#);
assert_eq!(row.encoding.as_deref(), Some("application/json"));
assert!(!row.delete);
let text = parse_row(r#"{"key":"k","value":"just words"}"#).unwrap();
assert_eq!(text.payload, b"just words");
}
#[test]
fn tombstones_and_malformed_rows_are_told_apart() {
let del = parse_row(r#"{"key":"k","delete":true,"value":null}"#).unwrap();
assert!(del.delete);
assert!(del.payload.is_empty());
let err = parse_row(r#"{"key":"k"}"#).unwrap_err();
assert!(err.contains("value"), "{err}");
let err = parse_row(r#"{"value":1}"#).unwrap_err();
assert!(err.contains("key"), "{err}");
let err = parse_row("not json").unwrap_err();
assert!(err.contains("JSON"), "{err}");
}
#[test]
fn lossless_bytes_win_over_the_rendering() {
let row = parse_row(
r#"{"key":"k","t":125000,"bytes":"AAEC/w==","value":"lossy render",
"attachment_b64":"3q0=","encoding":"application/octet-stream"}"#,
)
.unwrap();
assert_eq!(row.payload, vec![0x00, 0x01, 0x02, 0xff]);
assert_eq!(row.attachment.as_deref(), Some([0xde, 0xad].as_ref()));
let row = parse_row(r#"{"key":"k","bytes":"aGk="}"#).unwrap();
assert_eq!(row.payload, b"hi");
let err = parse_row(r#"{"key":"k","bytes":"not base64!"}"#).unwrap_err();
assert!(err.contains("base64"), "{err}");
let err = parse_row(r#"{"key":"k","bytes":7}"#).unwrap_err();
assert!(err.contains("base64"), "{err}");
}
fn view(payload: &[u8], profile: Option<zenkey::qos::QosProfile>) -> crate::SampleView {
use zenkey::qos::QosProfile;
let p = profile.unwrap_or(QosProfile::Sampled);
crate::SampleView {
key: "v1/h-3fa9c2d41b7e/state/sysinfo/health".into(),
payload: zenoh::bytes::ZBytes::from(payload.to_vec()),
encoding: "application/json".into(),
kind: zenoh::sample::SampleKind::Put,
timestamp: None,
stamped_by: None,
attachment: None,
priority: if profile.is_some() {
p.priority()
} else {
zenoh::qos::Priority::Background
},
congestion_control: p.congestion_control(),
reliability: p.reliability(),
express: p.express(),
source: None,
received: std::time::Instant::now(),
}
}
#[test]
fn the_qos_field_only_ever_carries_a_name_the_reader_can_resolve() {
let named = SampleRow::of_key("k", "")
.with_wire(&view(b"{}", Some(zenkey::qos::QosProfile::Alert)));
assert_eq!(named.qos.as_deref(), Some("alert"));
assert!(
zenkey::qos::QosProfile::from_name(named.qos.as_deref().unwrap()).is_some(),
"whatever lands in `qos` must resolve, or every row of the pipe is malformed"
);
let unnamed = SampleRow::of_key("k", "").with_wire(&view(b"{}", None));
assert_eq!(
unnamed.qos, None,
"axes matching no profile are omitted, never spelled into `qos`"
);
}
#[test]
fn a_row_this_crate_wrote_is_a_row_this_crate_reads() {
let v = view(
br#"{"status":"ok"}"#,
Some(zenkey::qos::QosProfile::Refreshed),
);
let mut observed = SampleRow::of_key(&v.key, "").with_wire(&v);
observed.value = Some(serde_json::json!({"status": "ok"}));
observed.qos_axes = Some("data/drop/reliable".into());
observed.payload_bytes = Some(15);
let back = parse_row(&observed.to_line()).expect("the observer's row reads back");
assert_eq!(back.payload, br#"{"status":"ok"}"#);
assert_eq!(back.qos.as_deref(), Some("refreshed"));
assert_eq!(back.encoding.as_deref(), Some("application/json"));
let captured = SampleRow {
key: v.key.clone(),
t: Some(1_250),
..SampleRow::default()
}
.with_wire(&v)
.with_payload_bytes(&v.payload.to_bytes());
let back = parse_row(&captured.to_line()).expect("the capture's row reads back");
assert_eq!(back.payload, br#"{"status":"ok"}"#);
assert_eq!(back.qos.as_deref(), Some("refreshed"));
}
#[test]
fn an_unheld_fact_is_absent_from_the_row_rather_than_null() {
let row = SampleRow {
key: "demo/foreign".into(),
delete: true,
..SampleRow::default()
};
assert_eq!(
serde_json::to_value(&row).unwrap(),
serde_json::json!({"key": "demo/foreign", "delete": true}),
"a key that did not parse carries no origin/subject, an unstamped \
sample carries no timestamp, and an undecoded one carries no type"
);
}
#[test]
fn tagged_meta_rows_are_skipped_not_malformed() {
assert_eq!(
parse_stream_line(r#"{"dropped":7,"row":"dropped"}"#).unwrap(),
StreamLine::Meta("dropped".into())
);
assert_eq!(
parse_stream_line(r#"{"row":"seed","seed_complete":{"superseded":0}}"#).unwrap(),
StreamLine::Meta("seed".into())
);
assert!(parse_stream_line(r#"{"row":7}"#).is_err());
assert!(matches!(
parse_stream_line(r#"{"key":"k","value":1}"#).unwrap(),
StreamLine::Sample(r) if r.key == "k"
));
let err = parse_stream_line(r#"{"key":"k"}"#).unwrap_err();
assert!(err.contains("value"), "{err}");
}
#[test]
fn attachments_ride_rows() {
let row =
parse_row(r#"{"key":"k","value":1,"attachment":{"who":"me"},"qos":"alert"}"#).unwrap();
assert_eq!(row.attachment.as_deref(), Some(br#"{"who":"me"}"#.as_ref()));
assert_eq!(row.qos.as_deref(), Some("alert"));
}
}