#![allow(dead_code)]
#![warn(missing_docs)]
use crate::asb::AbstractServiceBus;
use crate::calconfig::CalConfig;
use crate::uci::CalMessage;
use crate::uci::base::UUID;
use crate::uci::types::{
ClassificationEnum, ID_Type as _, MessageModeEnum, OwnerProducerChoiceType_,
};
use crate::uci::{CalError, CalErrorKind, CalResult};
use chrono::Utc;
use lazy_static::lazy_static;
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::Duration;
#[cfg(feature = "zmq")]
use crate::asb::zmq::{ZMQ_ASB_ID, ZmqAsb};
#[derive(Debug, Clone)]
pub struct MessageHeaderDefaults {
pub system_id: UUID,
pub service_id: Option<UUID>,
pub mission_id: Option<UUID>,
pub schema_version: String,
pub mode: MessageModeEnum,
pub classification: ClassificationEnum,
pub owner_producer: Vec<OwnerProducerChoiceType_>,
}
pub trait MessageListener<M: CalMessage>: Send + Sync {
fn on_message(&self, message: &Arc<M>);
}
pub trait AbstractWriter<M: CalMessage>: Send + Sync {
fn topic(&self) -> &str;
fn write(&mut self, message: &M) -> CalResult<()>;
fn close(self: Box<Self>) -> CalResult<()>;
}
pub trait AbstractReader<M: CalMessage>: Send + Sync {
fn topic(&self) -> &str;
fn add_listener(&mut self, listener: Arc<dyn MessageListener<M>>) -> CalResult<()>;
fn remove_listener(&mut self, listener: &Arc<dyn MessageListener<M>>) -> CalResult<()>;
fn read(&mut self, timeout: Option<Duration>) -> CalResult<Option<Arc<M>>>;
fn read_no_wait(&mut self) -> CalResult<Option<Arc<M>>>;
fn close(self: Box<Self>) -> CalResult<()>;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum Reliability {
#[default]
BestEffort,
Reliable,
}
#[derive(Debug, Clone)]
pub struct TimeBasedFilter {
pub min_separation: Duration,
}
#[derive(Debug, Clone)]
pub struct Expiration {
pub max_age: Duration,
}
#[derive(Debug, Clone)]
pub struct MessageBuffer {
pub max_messages: usize,
}
#[derive(Debug, Clone, Default)]
pub struct TopicQos {
pub reliability: Reliability,
pub time_based_filter: Option<TimeBasedFilter>,
pub expiration: Option<Expiration>,
pub writer_buffer: Option<MessageBuffer>,
pub reader_buffer: Option<MessageBuffer>,
}
pub trait AbstractCal: AbstractServiceBus {
fn message_header_defaults(&self) -> MessageHeaderDefaults;
}
pub trait AbstractCalExt<M: CalMessage>: AbstractCal {
fn create_writer(
&mut self,
topic: &str,
qos: TopicQos,
) -> CalResult<Box<dyn AbstractWriter<M>>>;
fn create_reader(
&mut self,
topic: &str,
qos: TopicQos,
) -> CalResult<Box<dyn AbstractReader<M>>>;
}
pub trait AbstractCalCreateMessage {
fn create_message<M: CalMessage>(&self) -> CalResult<M>;
}
impl<T: AbstractCal> AbstractCalCreateMessage for T {
fn create_message<M: CalMessage>(&self) -> CalResult<M> {
let mut msg = M::cal_create();
if let Some(mt) = msg.as_message_type_mut() {
let defaults = self.message_header_defaults();
let hdr = mt.message_header_mut();
*hdr.system_id_mut().uuid_mut() = defaults.system_id;
*hdr.schema_version_mut() = defaults.schema_version;
*hdr.mode_mut() = defaults.mode;
*hdr.timestamp_mut() = Utc::now().into();
if let (Some(sid), Some(sfield)) = (defaults.service_id, hdr.service_id_mut()) {
*sfield.uuid_mut() = sid;
}
if let (Some(mid), Some(mfield)) = (defaults.mission_id, hdr.mission_id_mut()) {
*mfield.uuid_mut() = mid;
}
let sec = mt.security_information_mut();
*sec.classification_mut() = defaults.classification;
let op = sec.owner_producer_mut();
op.clear();
op.extend(defaults.owner_producer);
}
Ok(msg)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub(crate) struct AsbKey {
pub(crate) service_identifier: String,
pub(crate) asb_identifier: String,
}
type CalInstance = Arc<Mutex<dyn AbstractCal>>;
type CalFactoryMap = HashMap<AsbKey, CalInstance>;
lazy_static! {
static ref CAL_FACTORY: tokio::sync::Mutex<CalFactoryMap> =
tokio::sync::Mutex::new(CalFactoryMap::new());
}
pub async fn get_cal(
service_identifier: impl Into<String>,
asb_identifier: impl Into<String>,
config: Arc<CalConfig>,
logger: slog::Logger,
) -> CalResult<CalInstance> {
let key = AsbKey {
service_identifier: service_identifier.into(),
asb_identifier: asb_identifier.into(),
};
let mut map = CAL_FACTORY.lock().await;
if let Some(existing) = map.get(&key) {
return Ok(Arc::clone(existing));
}
let transport = config
.get_transport(&key.asb_identifier)
.or_else(|| {
config
.system
.default_transport
.as_ref()
.and_then(|def| config.get_transport(def))
})
.ok_or_else(|| {
CalError::new(
CalErrorKind::InitializationFailure,
format!(
"No transport configured for '{}' and no default_transport available.",
key.asb_identifier
),
)
})?;
let instance: CalInstance = match transport.type_.as_str() {
#[cfg(feature = "zmq")]
ZMQ_ASB_ID => Arc::new(Mutex::new(
ZmqAsb::new(
key.service_identifier.clone(),
key.asb_identifier.clone(),
logger,
Arc::clone(&config),
transport,
)
.await?,
)),
other => {
return Err(CalError::new(
CalErrorKind::InitializationFailure,
format!("Unknown ASB transport type: '{other}'."),
));
}
};
map.insert(key, Arc::clone(&instance));
Ok(instance)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::asb::NEXT_TEST_PORT;
use crate::asb::zmq::test_config_on_ports;
use rcal_macros::init_test_logger;
use std::sync::atomic::Ordering;
#[cfg(feature = "zmq")]
#[init_test_logger]
#[tokio::test]
async fn test_cal_factory_same_key_returns_same_instance() {
let p1 = NEXT_TEST_PORT.fetch_add(1, Ordering::SeqCst);
let p2 = NEXT_TEST_PORT.fetch_add(1, Ordering::SeqCst);
let config = test_config_on_ports(&[p1, p2]);
let a = get_cal("test_svc", "TestZmq", Arc::clone(&config), logger.clone())
.await
.expect("first get_cal must succeed");
let b = get_cal("test_svc", "TestZmq", Arc::clone(&config), logger.clone())
.await
.expect("second get_cal must succeed");
let c = get_cal(
"test_svc_2",
"TestZmq2",
Arc::clone(&config),
logger.clone(),
)
.await
.expect("different service must succeed");
assert!(Arc::ptr_eq(&a, &b), "same key should return the same Arc");
assert!(
!Arc::ptr_eq(&a, &c),
"different service key must be a distinct instance"
);
}
}