use std::any::Any;
use std::marker::PhantomData;
use std::sync::Arc;
use std::time::Duration;
use slog::{error, info, trace, warn};
use crate::cal::{
AbstractCal, AbstractCalCreateMessage, AbstractCalExt, AbstractReader, AbstractWriter,
MessageListener, TopicQos,
};
use crate::calconfig::CalConfig;
use crate::uci::{CalMessage, CalResult};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ServiceLifecycleState {
Inactive,
Active,
}
pub trait AbstractService: Send + Sync {
fn system_id(&self) -> &str;
fn service_id(&self) -> &str;
fn subsystem_ids(&self) -> &[String];
fn lifecycle_state(&self) -> ServiceLifecycleState;
fn activate(&mut self) -> CalResult<()>;
fn deactivate(&mut self) -> CalResult<()>;
fn reset(&mut self) -> CalResult<()>;
}
struct CallbackListener<M, F> {
topic: String,
callback: F,
_p: PhantomData<M>,
}
impl<M, F> MessageListener<M> for CallbackListener<M, F>
where
M: CalMessage,
F: Fn(Arc<M>, &str) + Send + Sync,
{
fn on_message(&self, message: &Arc<M>) {
(self.callback)(Arc::clone(message), &self.topic);
}
}
#[rcal_macros::monitor]
pub struct AbstractServiceImpl<A> {
service_id: String,
system_id: String,
subsystem_ids: Vec<String>,
state: ServiceLifecycleState,
asb: A,
#[allow(dead_code)]
config: Arc<CalConfig>,
logger: slog::Logger,
_readers: std::sync::Mutex<Vec<Box<dyn Any + Send>>>,
}
impl<A: AbstractCal> AbstractServiceImpl<A> {
pub fn new(
service_id: impl Into<String>,
system_id: impl Into<String>,
subsystem_ids: Vec<String>,
asb: A,
config: Arc<CalConfig>,
logger: slog::Logger,
) -> Self {
let service_id = service_id.into();
let system_id = system_id.into();
trace!(logger, "AbstractServiceImpl::new";
"service_id" => &service_id,
"system_id" => &system_id,
);
Self {
service_id,
system_id,
subsystem_ids,
state: ServiceLifecycleState::Inactive,
asb,
config,
logger,
_readers: std::sync::Mutex::new(Vec::new()),
__monitor_mutex: ::std::sync::Mutex::new(()),
}
}
pub fn create_message<M: CalMessage>(&self) -> CalResult<M> {
trace!(self.logger, "AbstractServiceImpl::create_message";
"service_id" => &self.service_id,
);
self.asb.create_message::<M>()
}
pub fn create_writer<M>(
&mut self,
topic: &str,
qos: TopicQos,
) -> CalResult<Box<dyn AbstractWriter<M>>>
where
M: CalMessage,
A: AbstractCalExt<M>,
{
trace!(self.logger, "AbstractServiceImpl::create_writer";
"service_id" => &self.service_id,
"topic" => topic,
);
let _guard = self.__monitor_mutex.lock().unwrap();
self.asb.create_writer(topic, qos)
}
pub fn create_reader<M, F>(&mut self, topic: &str, qos: TopicQos, callback: F) -> CalResult<()>
where
M: CalMessage + 'static,
A: AbstractCalExt<M>,
F: Fn(Arc<M>, &str) + Send + Sync + 'static,
{
trace!(self.logger, "AbstractServiceImpl::create_reader";
"service_id" => &self.service_id,
"topic" => topic,
);
let _guard = self.__monitor_mutex.lock().unwrap();
let mut reader = self.asb.create_reader(topic, qos).map_err(|e| {
error!(self.logger, "create_reader failed";
"service_id" => &self.service_id,
"topic" => topic,
"error" => %e,
);
e
})?;
let listener: Arc<dyn MessageListener<M>> = Arc::new(CallbackListener {
topic: topic.to_string(),
callback,
_p: PhantomData,
});
reader.add_listener(listener).map_err(|e| {
error!(self.logger, "add_listener failed";
"service_id" => &self.service_id,
"topic" => topic,
"error" => %e,
);
e
})?;
self._readers
.lock()
.unwrap()
.push(Box::new(reader) as Box<dyn Any + Send>);
Ok(())
}
pub fn create_polling_reader<M>(
&mut self,
topic: &str,
qos: TopicQos,
) -> CalResult<Box<dyn AbstractReader<M>>>
where
M: CalMessage + 'static,
A: AbstractCalExt<M>,
{
trace!(self.logger, "AbstractServiceImpl::create_polling_reader";
"service_id" => &self.service_id,
"topic" => topic,
);
let _guard = self.__monitor_mutex.lock().unwrap();
self.asb.create_reader(topic, qos).map_err(|e| {
error!(self.logger, "create_polling_reader failed";
"service_id" => &self.service_id,
"topic" => topic,
"error" => %e,
);
e
})
}
pub fn status_delay(&self) -> Option<Duration> {
let cfg = self.config.get_service(&self.service_id)?;
let s = cfg.status_delay.as_deref()?;
parse_duration(s)
.map_err(|e| {
warn!(self.logger, "invalid status_delay, ignored";
"service_id" => &self.service_id,
"value" => s,
"error" => e,
);
})
.ok()
}
pub fn status_data_request_enabled(&self) -> bool {
self.config
.get_service(&self.service_id)
.map(|s| s.service_status_data_request_enable)
.unwrap_or(false)
}
pub fn asb_mut(&mut self) -> &mut A {
&mut self.asb
}
}
impl<A: AbstractCal> AbstractService for AbstractServiceImpl<A> {
fn system_id(&self) -> &str {
&self.system_id
}
fn service_id(&self) -> &str {
&self.service_id
}
fn subsystem_ids(&self) -> &[String] {
&self.subsystem_ids
}
fn lifecycle_state(&self) -> ServiceLifecycleState {
trace!(self.logger, "AbstractService::lifecycle_state";
"service_id" => &self.service_id,
);
let _guard = self.__monitor_mutex.lock().unwrap();
self.state
}
fn activate(&mut self) -> CalResult<()> {
trace!(self.logger, "AbstractService::activate";
"service_id" => &self.service_id,
);
let _guard = self.__monitor_mutex.lock().unwrap();
if self.state == ServiceLifecycleState::Active {
warn!(self.logger, "activate called while already active";
"service_id" => &self.service_id,
);
return Ok(());
}
self.state = ServiceLifecycleState::Active;
info!(self.logger, "service activated";
"service_id" => &self.service_id,
);
Ok(())
}
fn deactivate(&mut self) -> CalResult<()> {
trace!(self.logger, "AbstractService::deactivate";
"service_id" => &self.service_id,
);
let _guard = self.__monitor_mutex.lock().unwrap();
if self.state == ServiceLifecycleState::Inactive {
warn!(self.logger, "deactivate called while already inactive";
"service_id" => &self.service_id,
);
return Ok(());
}
self.state = ServiceLifecycleState::Inactive;
info!(self.logger, "service deactivated";
"service_id" => &self.service_id,
);
Ok(())
}
fn reset(&mut self) -> CalResult<()> {
trace!(self.logger, "AbstractService::reset";
"service_id" => &self.service_id,
);
let _guard = self.__monitor_mutex.lock().unwrap();
self.state = ServiceLifecycleState::Inactive;
self._readers.lock().unwrap().clear();
info!(self.logger, "service reset";
"service_id" => &self.service_id,
);
Ok(())
}
}
#[macro_export]
macro_rules! service_status_loop {
($interval:expr, $body:block) => {
::tokio::spawn(async move {
loop {
$body;
::tokio::time::sleep($interval).await;
}
})
};
}
fn parse_duration(s: &str) -> Result<Duration, &'static str> {
if let Some(ms) = s.strip_suffix("ms") {
ms.parse::<u64>()
.map(Duration::from_millis)
.map_err(|_| "invalid milliseconds value")
} else if let Some(secs) = s.strip_suffix('s') {
secs.parse::<u64>()
.map(Duration::from_secs)
.map_err(|_| "invalid seconds value")
} else if let Some(mins) = s.strip_suffix('m') {
mins.parse::<u64>()
.map(|m| Duration::from_secs(m * 60))
.map_err(|_| "invalid minutes value")
} else {
Err("unrecognised duration suffix (expected ms, s, or m)")
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::asb::{AbstractServiceBus, AsbStatus, AsbStatusListener};
use crate::cal::MessageHeaderDefaults;
use crate::calconfig::CalConfig;
use crate::uci::base::UUID;
struct NullAsb;
impl AbstractServiceBus for NullAsb {
fn get_logger(&self) -> &slog::Logger {
unimplemented!()
}
fn service_identifier(&self) -> &str {
"null"
}
fn asb_identifier(&self) -> &str {
"null"
}
fn get_system_uuid(&self) -> UUID {
unimplemented!()
}
fn get_service_uuid(&self) -> Option<UUID> {
None
}
fn get_subsystem_uuid(&self) -> Option<UUID> {
None
}
fn get_component_uuid(&self, _: &str) -> Option<UUID> {
None
}
fn get_capability_uuid(&self, _: &str) -> Option<UUID> {
None
}
fn oms_schema_version(&self) -> &str {
""
}
fn oms_schema_compiler_version(&self) -> &str {
""
}
fn get_system_label(&self) -> Option<&str> {
None
}
fn get_asb_connection_version(&self) -> &str {
""
}
fn get_oms_api_version(&self) -> &str {
""
}
fn connection_status(&self) -> &AsbStatus {
unimplemented!()
}
fn register_status_listener(
&mut self,
_: Arc<dyn AsbStatusListener>,
) -> crate::uci::CalResult<()> {
Ok(())
}
fn unregister_status_listener(
&mut self,
_: &Arc<dyn AsbStatusListener>,
) -> crate::uci::CalResult<()> {
Ok(())
}
fn close(&mut self) -> crate::uci::CalResult<()> {
Ok(())
}
}
impl AbstractCal for NullAsb {
fn message_header_defaults(&self) -> MessageHeaderDefaults {
unimplemented!()
}
}
fn make_svc() -> AbstractServiceImpl<NullAsb> {
let logger = slog::Logger::root(slog::Discard, slog::o!());
AbstractServiceImpl::new(
"svc",
"sys",
vec![],
NullAsb,
Arc::new(CalConfig::default()),
logger,
)
}
#[test]
fn test_reset_sets_state_inactive() {
let mut svc = make_svc();
svc.activate().unwrap();
assert_eq!(svc.lifecycle_state(), ServiceLifecycleState::Active);
svc.reset().unwrap();
assert_eq!(svc.lifecycle_state(), ServiceLifecycleState::Inactive);
}
#[test]
fn test_reset_idempotent_when_inactive() {
let mut svc = make_svc();
assert!(svc.reset().is_ok());
assert_eq!(svc.lifecycle_state(), ServiceLifecycleState::Inactive);
}
#[test]
fn test_parse_duration_seconds() {
assert_eq!(parse_duration("1s").unwrap(), Duration::from_secs(1));
assert_eq!(parse_duration("60s").unwrap(), Duration::from_secs(60));
}
#[test]
fn test_parse_duration_milliseconds() {
assert_eq!(parse_duration("500ms").unwrap(), Duration::from_millis(500));
}
#[test]
fn test_parse_duration_minutes() {
assert_eq!(parse_duration("2m").unwrap(), Duration::from_secs(120));
}
#[test]
fn test_parse_duration_invalid() {
assert!(parse_duration("bad").is_err());
assert!(parse_duration("1x").is_err());
}
}