pub mod publisher;
pub mod querier;
pub mod stamp;
pub mod subscriber;
use zenoh::sample::Sample;
use crate::bus::abi::{Codec, CodecId, EncodingError, EncodingMetadata, MessagePack};
use crate::bus::contract::{Endpoint, Payload};
use crate::bus::error::{BusError, MetadataProblem, Result};
use crate::bus::metadata::BusMetadata;
pub(crate) fn decode_sample<E: Endpoint>(sample: &Sample, topic: &str) -> Result<(E, BusMetadata)> {
decode_payload::<E>(sample, topic)
}
pub(crate) fn decode_payload<B: Payload>(sample: &Sample, topic: &str) -> Result<(B, BusMetadata)> {
let malformed = |problem: MetadataProblem| BusError::metadata(topic, problem);
let encoding: EncodingMetadata = sample
.encoding()
.to_string()
.parse()
.map_err(|e: EncodingError| malformed(e.into()))?;
if encoding.codec_id() != Some(CodecId::MessagePack) {
return Err(BusError::UnsupportedCodec {
codec: encoding.codec,
topic: topic.to_string(),
});
}
let attachment = sample
.attachment()
.ok_or_else(|| malformed(MetadataProblem::MissingAttachment))?;
let metadata =
BusMetadata::decode(attachment.to_bytes().as_ref()).map_err(|e| malformed(e.into()))?;
if metadata.codec != encoding.codec {
return Err(malformed(MetadataProblem::CodecMismatch {
encoding: encoding.codec,
attachment: metadata.codec,
}));
}
if metadata.codec_id() != Some(CodecId::MessagePack) {
return Err(BusError::UnsupportedCodec {
codec: metadata.codec,
topic: topic.to_string(),
});
}
let body = MessagePack::decode::<B>(sample.payload().to_bytes().as_ref())?;
Ok((body, metadata))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::bus::abi::CodecError;
use crate::bus::test_support::{
TARGET_TOPIC as TOPIC, Target, sample, sample_with, sample_with_encoding,
};
#[test]
fn decode_accepts_a_matching_sample() {
let sample = sample(CodecId::MessagePack.as_u8());
let (body, metadata) = decode_sample::<Target>(&sample, TOPIC).unwrap();
assert_eq!(body.linear_x_mps, 1.0);
assert_eq!(metadata.codec, CodecId::MessagePack.as_u8());
}
#[test]
fn decode_rejects_encoding_attachment_codec_mismatch_before_body_decode() {
let payload = rmp_serde::to_vec_named(&Target {
linear_x_mps: 1.0,
angular_z_radps: 0.5,
})
.unwrap();
let sample = sample_with_encoding(
CodecId::MessagePack.as_u8(),
"phoxal/v0;codec=99".to_string(),
payload,
);
let error = decode_sample::<Target>(&sample, TOPIC).unwrap_err();
assert!(matches!(
error,
BusError::UnsupportedCodec { codec: 99, .. }
));
}
#[test]
fn decode_rejects_unsupported_codec() {
let error = decode_sample::<Target>(&sample(99), TOPIC).unwrap_err();
assert!(matches!(
error,
BusError::UnsupportedCodec { codec: 99, .. }
));
}
#[test]
fn decode_rejects_corrupt_payload() {
let sample = sample_with(CodecId::MessagePack.as_u8(), vec![0xc1, 0xc1, 0xc1]);
let error = decode_sample::<Target>(&sample, TOPIC).unwrap_err();
assert!(matches!(error, BusError::Codec(CodecError::Decode(_))));
}
}