use std::time::Duration;
use nt_client::{Client, data::DataType, schema::{ParseFromSchemaError, PublishSchemaError, SchemaManager}, r#struct::{StructData, StructSchema, byte::{ByteBuffer, ByteReader}}, subscribe::ReceivedMessage, topic::Properties};
use protobuf::reflect::ReflectFieldRef;
#[derive(Debug, Clone, Copy, PartialEq)]
struct Position {
x: f64,
y: f64,
}
impl StructData for Position {
fn struct_type_name() -> String {
"Position".to_string()
}
fn schema() -> StructSchema {
StructSchema("double x;double y".to_string())
}
async fn publish_dependencies(_manager: &mut SchemaManager) -> Result<(), PublishSchemaError> {
Ok(())
}
fn pack(self, buf: &mut ByteBuffer) {
buf.write_f64(self.x);
buf.write_f64(self.y);
}
fn unpack(read: &mut ByteReader) -> Option<Self> {
Some(Self {
x: read.read_f64()?,
y: read.read_f64()?,
})
}
}
#[tokio::main]
async fn main() {
let client = Client::new(Default::default());
client.connect_setup(setup).await.unwrap()
}
fn setup(client: &Client) {
let mut manager = client.schema_manager();
let watch_manager = manager.clone();
tokio::spawn(async move {
watch_manager.watch().await.unwrap()
});
let mut pub_manager = manager.clone();
let struct_topic = client.topic("/position");
tokio::spawn(async move {
pub_manager.publish_struct::<Position>().await.unwrap();
let publisher = struct_topic.publish(Properties { retained: Some(true), ..Default::default() }).await.unwrap();
let position = Position {
x: 15.2,
y: 8.91,
};
publisher.set(position.into_struct_data()).await.unwrap();
});
let mut sub_struct_manager = manager.clone();
let sub_struct_topic = client.topic("/struct");
tokio::spawn(async move {
let mut sub = sub_struct_topic.subscribe(Default::default()).await.unwrap();
while let Ok(message) = sub.recv().await {
if let ReceivedMessage::Updated((topic, value)) = message
&& let DataType::Struct(type_name) = topic.r#type() {
match sub_struct_manager.parse_struct(type_name, value).await {
Ok(fields) => {
println!("{type_name} {{");
for (name, value) in fields {
println!(" {name}: {value:?}");
}
println!("}}");
},
Err(ParseFromSchemaError::SchemaNotFound) => eprintln!("the schema for struct:{type_name} was not found"),
Err(ParseFromSchemaError::InvalidData) => eprintln!("invalid struct"),
}
}
}
});
let sub_proto_topic = client.topic("/protobuf");
tokio::spawn(async move {
let mut sub = sub_proto_topic.subscribe(Default::default()).await.unwrap();
while let Ok(message) = sub.recv().await {
if let ReceivedMessage::Updated((topic, value)) = message
&& let DataType::Protobuf(type_name) = topic.r#type() {
tokio::time::sleep(Duration::from_millis(100)).await;
match manager.parse_proto(type_name, value).await {
Ok(message) => {
let descriptor = message.descriptor_dyn();
println!("{} {{", descriptor.name());
for field in descriptor.fields() {
let value = field.get_reflect(&*message);
match value {
ReflectFieldRef::Map(map) => println!(" {}: {map:#?}", field.name()),
ReflectFieldRef::Optional(optional) => println!(" {}: {:#?}", field.name(), optional.value()),
ReflectFieldRef::Repeated(repeated) => println!(" {}: {repeated:#?}", field.name()),
}
}
println!("}}");
},
Err(ParseFromSchemaError::SchemaNotFound) => eprintln!("the schema for proto:{type_name} was not found"),
Err(ParseFromSchemaError::InvalidData) => eprintln!("invalid protobuf"),
}
}
}
});
}