use std::{
convert::TryInto,
future::{IntoFuture, Ready},
time::Duration,
};
use tracing::error;
use zenoh_core::{Resolvable, Resolve, Result as ZResult, Wait};
use crate::{
api::{
handlers::{locked, DefaultHandler, IntoHandler},
key_expr::KeyExpr,
query::Reply,
sample::{Locality, Sample},
session::{Session, UndeclarableSealed, WeakSession},
subscriber::{Subscriber, SubscriberInner},
Id,
},
handlers::Callback,
};
pub struct Liveliness<'a> {
pub(crate) session: &'a Session,
}
impl<'a> Liveliness<'a> {
pub fn declare_token<'b, TryIntoKeyExpr>(
&self,
key_expr: TryIntoKeyExpr,
) -> LivelinessTokenBuilder<'a, 'b>
where
TryIntoKeyExpr: TryInto<KeyExpr<'b>>,
<TryIntoKeyExpr as TryInto<KeyExpr<'b>>>::Error: Into<zenoh_core::Error>,
{
LivelinessTokenBuilder {
session: self.session,
key_expr: TryIntoKeyExpr::try_into(key_expr).map_err(Into::into),
}
}
pub fn declare_subscriber<'b, TryIntoKeyExpr>(
&self,
key_expr: TryIntoKeyExpr,
) -> LivelinessSubscriberBuilder<'a, 'b, DefaultHandler>
where
TryIntoKeyExpr: TryInto<KeyExpr<'b>>,
<TryIntoKeyExpr as TryInto<KeyExpr<'b>>>::Error: Into<zenoh_result::Error>,
{
LivelinessSubscriberBuilder {
session: self.session,
key_expr: TryIntoKeyExpr::try_into(key_expr).map_err(Into::into),
handler: DefaultHandler::default(),
history: false,
}
}
pub fn get<'b, TryIntoKeyExpr>(
&self,
key_expr: TryIntoKeyExpr,
) -> LivelinessGetBuilder<'a, 'b, DefaultHandler>
where
TryIntoKeyExpr: TryInto<KeyExpr<'b>>,
<TryIntoKeyExpr as TryInto<KeyExpr<'b>>>::Error: Into<zenoh_result::Error>,
{
let key_expr = key_expr.try_into().map_err(Into::into);
let timeout = {
Duration::from_millis(
self.session
.0
.runtime
.get_config()
.queries_default_timeout_ms(),
)
};
LivelinessGetBuilder {
session: self.session,
key_expr,
timeout,
handler: DefaultHandler::default(),
}
}
}
#[must_use = "Resolvables do nothing unless you resolve them using `.await` or `zenoh::Wait::wait`"]
#[derive(Debug)]
pub struct LivelinessTokenBuilder<'a, 'b> {
pub(crate) session: &'a Session,
pub(crate) key_expr: ZResult<KeyExpr<'b>>,
}
impl Resolvable for LivelinessTokenBuilder<'_, '_> {
type To = ZResult<LivelinessToken>;
}
impl Wait for LivelinessTokenBuilder<'_, '_> {
#[inline]
fn wait(self) -> <Self as Resolvable>::To {
let session = self.session;
let key_expr = self.key_expr?.into_owned();
session
.0
.declare_liveliness_inner(&key_expr)
.map(|id| LivelinessToken {
session: self.session.downgrade(),
id,
undeclare_on_drop: true,
})
}
}
impl IntoFuture for LivelinessTokenBuilder<'_, '_> {
type Output = <Self as Resolvable>::To;
type IntoFuture = Ready<<Self as Resolvable>::To>;
fn into_future(self) -> Self::IntoFuture {
std::future::ready(self.wait())
}
}
#[must_use = "Liveliness tokens will be immediately dropped and undeclared if not bound to a variable"]
#[derive(Debug)]
pub struct LivelinessToken {
session: WeakSession,
id: Id,
undeclare_on_drop: bool,
}
#[must_use = "Resolvables do nothing unless you resolve them using `.await` or `zenoh::Wait::wait`"]
pub struct LivelinessTokenUndeclaration(LivelinessToken);
impl Resolvable for LivelinessTokenUndeclaration {
type To = ZResult<()>;
}
impl Wait for LivelinessTokenUndeclaration {
fn wait(mut self) -> <Self as Resolvable>::To {
self.0.undeclare_impl()
}
}
impl IntoFuture for LivelinessTokenUndeclaration {
type Output = <Self as Resolvable>::To;
type IntoFuture = Ready<<Self as Resolvable>::To>;
fn into_future(self) -> Self::IntoFuture {
std::future::ready(self.wait())
}
}
impl LivelinessToken {
#[inline]
pub fn undeclare(self) -> impl Resolve<ZResult<()>> {
UndeclarableSealed::undeclare_inner(self, ())
}
fn undeclare_impl(&mut self) -> ZResult<()> {
self.undeclare_on_drop = false;
self.session.undeclare_liveliness(self.id)
}
}
impl UndeclarableSealed<()> for LivelinessToken {
type Undeclaration = LivelinessTokenUndeclaration;
fn undeclare_inner(self, _: ()) -> Self::Undeclaration {
LivelinessTokenUndeclaration(self)
}
}
impl Drop for LivelinessToken {
fn drop(&mut self) {
if self.undeclare_on_drop {
if let Err(error) = self.undeclare_impl() {
error!(error);
}
}
}
}
#[must_use = "Resolvables do nothing unless you resolve them using `.await` or `zenoh::Wait::wait`"]
#[derive(Debug)]
pub struct LivelinessSubscriberBuilder<'a, 'b, Handler, const BACKGROUND: bool = false> {
pub session: &'a Session,
pub key_expr: ZResult<KeyExpr<'b>>,
pub handler: Handler,
pub history: bool,
}
impl<'a, 'b> LivelinessSubscriberBuilder<'a, 'b, DefaultHandler> {
#[inline]
pub fn callback<F>(self, callback: F) -> LivelinessSubscriberBuilder<'a, 'b, Callback<Sample>>
where
F: Fn(Sample) + Send + Sync + 'static,
{
self.with(Callback::from(callback))
}
#[inline]
pub fn callback_mut<F>(
self,
callback: F,
) -> LivelinessSubscriberBuilder<'a, 'b, Callback<Sample>>
where
F: FnMut(Sample) + Send + Sync + 'static,
{
self.callback(locked(callback))
}
#[inline]
pub fn with<Handler>(self, handler: Handler) -> LivelinessSubscriberBuilder<'a, 'b, Handler>
where
Handler: IntoHandler<Sample>,
{
let LivelinessSubscriberBuilder {
session,
key_expr,
handler: _,
history,
} = self;
LivelinessSubscriberBuilder {
session,
key_expr,
handler,
history,
}
}
}
impl<'a, 'b> LivelinessSubscriberBuilder<'a, 'b, Callback<Sample>> {
pub fn background(self) -> LivelinessSubscriberBuilder<'a, 'b, Callback<Sample>, true> {
LivelinessSubscriberBuilder {
session: self.session,
key_expr: self.key_expr,
handler: self.handler,
history: self.history,
}
}
}
impl<Handler, const BACKGROUND: bool> LivelinessSubscriberBuilder<'_, '_, Handler, BACKGROUND> {
#[inline]
pub fn history(mut self, history: bool) -> Self {
self.history = history;
self
}
}
impl<Handler> Resolvable for LivelinessSubscriberBuilder<'_, '_, Handler>
where
Handler: IntoHandler<Sample> + Send,
Handler::Handler: Send,
{
type To = ZResult<Subscriber<Handler::Handler>>;
}
impl<Handler> Wait for LivelinessSubscriberBuilder<'_, '_, Handler>
where
Handler: IntoHandler<Sample> + Send,
Handler::Handler: Send,
{
fn wait(self) -> <Self as Resolvable>::To {
use super::subscriber::SubscriberKind;
let key_expr = self.key_expr?;
let session = self.session;
let (callback, handler) = self.handler.into_handler();
session
.0
.declare_liveliness_subscriber_inner(
&key_expr,
Locality::default(),
self.history,
callback,
)
.map(|sub_state| Subscriber {
inner: SubscriberInner {
session: self.session.downgrade(),
id: sub_state.id,
key_expr: sub_state.key_expr.clone(),
kind: SubscriberKind::LivelinessSubscriber,
undeclare_on_drop: true,
},
handler,
})
}
}
impl<Handler> IntoFuture for LivelinessSubscriberBuilder<'_, '_, Handler>
where
Handler: IntoHandler<Sample> + Send,
Handler::Handler: Send,
{
type Output = <Self as Resolvable>::To;
type IntoFuture = Ready<<Self as Resolvable>::To>;
fn into_future(self) -> Self::IntoFuture {
std::future::ready(self.wait())
}
}
impl Resolvable for LivelinessSubscriberBuilder<'_, '_, Callback<Sample>, true> {
type To = ZResult<()>;
}
impl Wait for LivelinessSubscriberBuilder<'_, '_, Callback<Sample>, true> {
fn wait(self) -> <Self as Resolvable>::To {
self.session.0.declare_liveliness_subscriber_inner(
&self.key_expr?,
Locality::default(),
self.history,
self.handler,
)?;
Ok(())
}
}
impl IntoFuture for LivelinessSubscriberBuilder<'_, '_, Callback<Sample>, true> {
type Output = <Self as Resolvable>::To;
type IntoFuture = Ready<<Self as Resolvable>::To>;
fn into_future(self) -> Self::IntoFuture {
std::future::ready(self.wait())
}
}
#[must_use = "Resolvables do nothing unless you resolve them using `.await` or `zenoh::Wait::wait`"]
#[derive(Debug)]
pub struct LivelinessGetBuilder<'a, 'b, Handler> {
pub(crate) session: &'a Session,
pub(crate) key_expr: ZResult<KeyExpr<'b>>,
pub(crate) timeout: Duration,
pub(crate) handler: Handler,
}
impl<'a, 'b> LivelinessGetBuilder<'a, 'b, DefaultHandler> {
#[inline]
pub fn callback<F>(self, callback: F) -> LivelinessGetBuilder<'a, 'b, Callback<Reply>>
where
F: Fn(Reply) + Send + Sync + 'static,
{
self.with(Callback::from(callback))
}
#[inline]
pub fn callback_mut<F>(self, callback: F) -> LivelinessGetBuilder<'a, 'b, Callback<Reply>>
where
F: FnMut(Reply) + Send + Sync + 'static,
{
self.callback(locked(callback))
}
#[inline]
pub fn with<Handler>(self, handler: Handler) -> LivelinessGetBuilder<'a, 'b, Handler>
where
Handler: IntoHandler<Reply>,
{
let LivelinessGetBuilder {
session,
key_expr,
timeout,
handler: _,
} = self;
LivelinessGetBuilder {
session,
key_expr,
timeout,
handler,
}
}
}
impl<Handler> LivelinessGetBuilder<'_, '_, Handler> {
#[inline]
pub fn timeout(mut self, timeout: Duration) -> Self {
self.timeout = timeout;
self
}
}
impl<Handler> Resolvable for LivelinessGetBuilder<'_, '_, Handler>
where
Handler: IntoHandler<Reply> + Send,
Handler::Handler: Send,
{
type To = ZResult<Handler::Handler>;
}
impl<Handler> Wait for LivelinessGetBuilder<'_, '_, Handler>
where
Handler: IntoHandler<Reply> + Send,
Handler::Handler: Send,
{
fn wait(self) -> <Self as Resolvable>::To {
let (callback, receiver) = self.handler.into_handler();
self.session
.0
.liveliness_query(&self.key_expr?, self.timeout, callback)
.map(|_| receiver)
}
}
impl<Handler> IntoFuture for LivelinessGetBuilder<'_, '_, Handler>
where
Handler: IntoHandler<Reply> + Send,
Handler::Handler: Send,
{
type Output = <Self as Resolvable>::To;
type IntoFuture = Ready<<Self as Resolvable>::To>;
fn into_future(self) -> Self::IntoFuture {
std::future::ready(self.wait())
}
}