use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
use tracing::warn;
use zenoh::Result;
use crate::{
msg::{ZMessage, ZSerializer},
pubsub::ZPub,
};
pub trait ManagedEntity: Send + Sync {
fn on_activate(&self);
fn on_deactivate(&self);
}
pub struct ZLifecyclePublisher<T: ZMessage, S: ZSerializer = <T as ZMessage>::Serdes> {
inner: ZPub<T, S>,
activated: Arc<AtomicBool>,
should_warn: AtomicBool,
}
impl<T: ZMessage, S: ZSerializer> ZLifecyclePublisher<T, S> {
pub(super) fn new(inner: ZPub<T, S>) -> Arc<Self> {
Arc::new(Self {
inner,
activated: Arc::new(AtomicBool::new(false)),
should_warn: AtomicBool::new(true),
})
}
pub fn publish(&self, msg: &T) -> Result<()>
where
T: 'static,
S: for<'a> crate::msg::ZSerializer<Input<'a> = &'a T> + 'static,
{
if !self.activated.load(Ordering::Relaxed) {
if self
.should_warn
.compare_exchange(true, false, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
warn!(
topic = %self.inner.entity.topic,
"publish() called while lifecycle publisher is deactivated — message dropped"
);
}
return Ok(());
}
self.inner.publish(msg)
}
pub fn is_activated(&self) -> bool {
self.activated.load(Ordering::Relaxed)
}
pub fn topic_name(&self) -> &str {
&self.inner.entity.topic
}
}
impl<T: ZMessage, S: ZSerializer + Send + Sync> ManagedEntity for ZLifecyclePublisher<T, S> {
fn on_activate(&self) {
self.activated.store(true, Ordering::Relaxed);
self.should_warn.store(true, Ordering::Relaxed);
}
fn on_deactivate(&self) {
self.activated.store(false, Ordering::Relaxed);
}
}
impl<T: ZMessage + std::fmt::Debug, S: ZSerializer> std::fmt::Debug for ZLifecyclePublisher<T, S> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ZLifecyclePublisher")
.field("topic", &self.inner.entity.topic)
.field("activated", &self.activated.load(Ordering::Relaxed))
.finish()
}
}