use std::any::Any;
use std::marker::PhantomData;
use std::sync::Arc;
use std::time::Duration;
use slog::{error, info, trace, warn};
use crate::asb::{
AbstractServiceBus, AbstractServiceBusCreateMessage, AbstractServiceBusExt, 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: AbstractServiceBus> 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: AbstractServiceBusExt<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: AbstractServiceBusExt<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 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: AbstractServiceBus> 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,
);
Err(crate::uci::CalError::new(
crate::uci::CalErrorKind::OperationNotPermitted,
format!("reset() not implemented for service '{}'", self.service_id),
))
}
}
#[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::*;
#[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());
}
}