use zenoh::handlers::FifoChannelHandler;
use zenoh::key_expr::OwnedKeyExpr;
use zenoh::query::{Query as ZenohQuery, Queryable};
use crate::abi::{CodecId, encoding_string, parse_encoding_string};
use crate::error::{BusError, Result};
use crate::metadata::BusMetadata;
use crate::query::QueryFailure;
use crate::session::Bus;
pub struct ServerQueryable {
inner: Queryable<FifoChannelHandler<ZenohQuery>>,
topic_key: String,
}
impl ServerQueryable {
pub async fn recv(&self) -> Result<IncomingQuery> {
let query = self
.inner
.recv_async()
.await
.map_err(|_| BusError::Closed)?;
Ok(IncomingQuery {
query,
topic_key: self.topic_key.clone(),
})
}
pub fn topic_key(&self) -> &str {
&self.topic_key
}
}
pub struct IncomingQuery {
query: ZenohQuery,
topic_key: String,
}
impl IncomingQuery {
pub fn topic_key(&self) -> &str {
&self.topic_key
}
pub fn request_bytes(&self) -> Result<Vec<u8>> {
let payload = self.query.payload().ok_or_else(|| BusError::Metadata {
topic: self.topic_key.clone(),
detail: "query is missing a request payload".to_string(),
})?;
Ok(payload.to_bytes().to_vec())
}
pub fn request_metadata(&self) -> Result<BusMetadata> {
let encoding = self.query.encoding().ok_or_else(|| BusError::Metadata {
topic: self.topic_key.clone(),
detail: "query is missing a request encoding".to_string(),
})?;
let encoding =
parse_encoding_string(&encoding.to_string()).map_err(|e| BusError::Metadata {
topic: self.topic_key.clone(),
detail: format!("malformed request encoding string: {e}"),
})?;
match encoding.codec_id() {
Some(CodecId::MessagePack) => {}
None => {
return Err(BusError::UnsupportedCodec(
encoding.codec,
self.topic_key.clone(),
));
}
}
let attachment = self.query.attachment().ok_or_else(|| BusError::Metadata {
topic: self.topic_key.clone(),
detail: "query is missing a BusMetadata attachment".to_string(),
})?;
let metadata = BusMetadata::decode(attachment.to_bytes().as_ref()).map_err(|e| {
BusError::Metadata {
topic: self.topic_key.clone(),
detail: format!("malformed request BusMetadata: {e}"),
}
})?;
if metadata.codec != encoding.codec {
return Err(BusError::Metadata {
topic: self.topic_key.clone(),
detail: format!(
"request encoding/BusMetadata codec mismatch: encoding codec={}, metadata codec={}",
encoding.codec, metadata.codec
),
});
}
Ok(metadata)
}
pub async fn reply(&self, bus: &Bus, payload: Vec<u8>) -> Result<()> {
let metadata = bus.metadata(None);
self.query
.reply(self.query.key_expr(), payload)
.encoding(encoding_string(CodecId::MessagePack))
.attachment(metadata.encode())
.await
.map_err(|e| BusError::Transport(e.to_string()))
}
pub async fn reply_err(&self, failure: &QueryFailure) -> Result<()> {
self.query
.reply_err(failure.encode())
.await
.map_err(|e| BusError::Transport(e.to_string()))
}
}
impl Bus {
pub async fn declare_server(&self, topic_key: &str) -> Result<ServerQueryable> {
let full_key = self.full_key(topic_key);
let key = OwnedKeyExpr::new(full_key.clone())
.map_err(|e| BusError::Namespace(format!("invalid server key '{full_key}': {e}")))?;
let inner = self
.session()
.declare_queryable(key)
.complete(true)
.await
.map_err(|e| BusError::Transport(e.to_string()))?;
Ok(ServerQueryable {
inner,
topic_key: topic_key.to_string(),
})
}
}