use std::sync::Arc;
use std::time::Duration;
use slog::{error, info};
use rcal::asb::zmq::ZmqAsb;
use rcal::asb::{
AbstractServiceBus, AbstractServiceBusCreateMessage, AbstractServiceBusExt, TopicQos,
};
use rcal::uci::CalMessage;
use rcal::uci::base::UUID;
use rcal::uci::types::*;
const TOPIC: &str = "SystemStatus";
const ITERATIONS: usize = 10;
#[rcal_macros::rcal_main(config = "examples/CALConfig.toml")]
async fn main() {
let config = Arc::new(rcal_config);
let transport_id = config
.system
.default_transport
.clone()
.unwrap_or_else(|| "default".to_string());
let tconfig = config
.get_transport(&transport_id)
.unwrap_or_else(|| panic!("transport '{transport_id}' not in config"))
.clone();
let mut bus = ZmqAsb::new(
"SystemStatusExample",
transport_id,
root_logger.clone(),
config,
&tconfig,
)
.await
.expect("ASB init failed");
let mut reader = <ZmqAsb as AbstractServiceBusExt<SystemStatus_>>::create_reader(
&mut bus,
TOPIC,
TopicQos::default(),
)
.expect("create_reader failed");
let mut writer = <ZmqAsb as AbstractServiceBusExt<SystemStatus_>>::create_writer(
&mut bus,
TOPIC,
TopicQos::default(),
)
.expect("create_writer failed");
let pretty = tconfig.format == rcal::uci::base::SerializationFormat::PrettyXml;
let reader_handle = tokio::task::spawn_blocking(move || {
let expected_type = SystemStatus_::message_type_name();
loop {
match reader.read(Some(Duration::from_millis(500))) {
Ok(Some(msg)) => {
assert_eq!(SystemStatus_::message_type_name(), expected_type);
let xml_result: Result<String, String> = if pretty {
let mut buf = String::new();
(|| {
let mut ser =
quick_xml::se::Serializer::with_root(&mut buf, Some(TOPIC))
.map_err(|e| e.to_string())?;
ser.indent(' ', 4);
use serde::Serialize as _;
(*msg).serialize(ser).map_err(|e| e.to_string())?;
Ok(buf)
})()
} else {
quick_xml::se::to_string_with_root(TOPIC, &*msg).map_err(|e| e.to_string())
};
match xml_result {
Ok(xml) => println!("[received]\n{xml}\n"),
Err(e) => eprintln!("[reader] serialize error: {e}"),
}
}
Ok(None) => {} Err(_) => break, }
}
});
let mut msg = bus
.create_message::<SystemStatus_>()
.expect("create_message failed");
*(*msg).security_information_mut().classification_mut() = ClassificationEnum::U;
(*msg).security_information_mut().owner_producer_mut().push(
OwnerProducerChoiceType_::GovernmentIdentifier {
inner: OwnerProducerEnum::Usa,
},
);
*(*msg).message_header_mut().system_id_mut().uuid_mut() = UUID::generate(None);
(*msg)
.message_header_mut()
.system_id_mut()
.descriptive_label_mut()
.replace(&mut "This is an example system".to_string());
(*msg)
.message_header_mut()
.schema_version_mut()
.clone_from(&bus.oms_schema_version().to_string());
*(*msg).message_header_mut().mode_mut() = MessageModeEnum::Simulation;
let sysid = (*msg).message_header().system_id().clone();
*(*msg).message_data_mut().system_id_mut() = sysid;
*(*msg).message_data_mut().system_state_mut() = SystemStateEnum::Operational;
*(*msg).message_data_mut().source_mut() = SystemSourceEnum::Actual;
tokio::time::sleep(Duration::from_millis(100)).await;
for i in 0..ITERATIONS {
*msg.message_data_mut().system_state_mut() = if i % 2 == 0 {
SystemStateEnum::Operational
} else {
SystemStateEnum::Degraded
};
if let Err(e) = writer.write(&msg) {
error!(root_logger, "Unable to write to ASB: {e}");
break;
}
info!(root_logger, "sent"; "iteration" => i + 1, "of" => ITERATIONS);
tokio::time::sleep(Duration::from_secs(1)).await;
}
bus.close().ok();
reader_handle.await.ok();
}