Skip to main content

zenoh_ext/
advanced_subscriber.rs

1//
2// Copyright (c) 2022 ZettaScale Technology
3//
4// This program and the accompanying materials are made available under the
5// terms of the Eclipse Public License 2.0 which is available at
6// http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0
7// which is available at https://www.apache.org/licenses/LICENSE-2.0.
8//
9// SPDX-License-Identifier: EPL-2.0 OR Apache-2.0
10//
11// Contributors:
12//   ZettaScale Zenoh Team, <zenoh@zettascale.tech>
13//
14use std::{
15    cmp::min, collections::BTreeMap, fmt, future::IntoFuture, hash::Hash, str::FromStr, sync::Weak,
16    time::Instant,
17};
18
19use lru::LruCache;
20use tokio_util::task::AbortOnDropHandle;
21use zenoh::{
22    config::ZenohId,
23    handlers::{Callback, CallbackDrop, CallbackParameter, IntoHandler},
24    internal::bail,
25    key_expr::KeyExpr,
26    liveliness::{LivelinessSubscriberBuilder, LivelinessToken},
27    pubsub::{SubscriberBuilder, SubscriberUndeclaration},
28    query::{
29        ConsolidationMode, Parameters, Selector, TimeBound, TimeExpr, TimeRange, ZenohParameters,
30    },
31    sample::{Locality, Sample, SampleKind},
32    session::{EntityGlobalId, EntityId, WeakSession},
33    Resolvable, Session, Wait, KE_ADV_PREFIX, KE_EMPTY, KE_PUB, KE_STAR, KE_STARSTAR, KE_SUB,
34};
35#[zenoh_macros::unstable]
36use {
37    std::collections::HashMap,
38    std::convert::TryFrom,
39    std::future::Ready,
40    std::sync::{Arc, Mutex},
41    std::time::Duration,
42    uhlc::ID,
43    zenoh::handlers::{locked, DefaultHandler},
44    zenoh::internal::{runtime::ZRuntime, zlock},
45    zenoh::pubsub::Subscriber,
46    zenoh::query::{QueryTarget, Reply, ReplyKeyExpr},
47    zenoh::time::Timestamp,
48    zenoh::Result as ZResult,
49};
50
51use crate::{
52    advanced_cache::{ke_liveliness, KE_UHLC},
53    utils::WrappingSn,
54    z_deserialize,
55};
56
57#[derive(Debug, Default, Clone)]
58/// Configure query for historical data for [`history`](crate::AdvancedSubscriberBuilder::history) method.
59#[zenoh_macros::unstable]
60pub struct HistoryConfig {
61    liveliness: bool,
62    max_samples: Option<usize>,
63    max_age: Option<f64>,
64}
65
66#[zenoh_macros::unstable]
67impl HistoryConfig {
68    /// Enable detection of late joiner publishers and query for their historical data.
69    ///
70    /// Late joiner detection can only be achieved for [`AdvancedPublishers`](crate::AdvancedPublisher) that enable publisher_detection.
71    /// History can only be retransmitted by [`AdvancedPublishers`](crate::AdvancedPublisher) that enable [`cache`](crate::AdvancedPublisherBuilder::cache).
72    #[inline]
73    #[zenoh_macros::unstable]
74    pub fn detect_late_publishers(mut self) -> Self {
75        self.liveliness = true;
76        self
77    }
78
79    /// Specify how many samples to query for each resource.
80    ///
81    /// Builder will fail if `max_samples` is set to zero.
82    #[zenoh_macros::unstable]
83    pub fn max_samples(mut self, depth: usize) -> Self {
84        self.max_samples = Some(depth);
85        self
86    }
87
88    /// Specify the maximum age of samples to query.
89    ///
90    /// Builder will fail if `max_age` is set to zero.
91    #[zenoh_macros::unstable]
92    pub fn max_age(mut self, seconds: f64) -> Self {
93        self.max_age = Some(seconds);
94        self
95    }
96}
97
98#[derive(Debug, Default, Clone, Copy)]
99/// Configure retransmission.
100#[zenoh_macros::unstable]
101pub struct RecoveryConfig<const CONFIGURED: bool = true> {
102    periodic_queries: Option<Duration>,
103    heartbeat: bool,
104    retention_period: Option<Duration>,
105}
106
107#[zenoh_macros::unstable]
108impl RecoveryConfig<false> {
109    /// Enable periodic queries for not yet received Samples and specify their period.
110    ///
111    /// This allows retrieving the last Sample(s) if the last Sample(s) is/are lost.
112    /// So it is useful for sporadic publications but useless for periodic publications
113    /// with a period smaller or equal to this period.
114    /// Retransmission can only be achieved by [`AdvancedPublishers`](crate::AdvancedPublisher)
115    /// that enable [`cache`](crate::AdvancedPublisherBuilder::cache) and
116    /// [`sample_miss_detection`](crate::AdvancedPublisherBuilder::sample_miss_detection).
117    #[zenoh_macros::unstable]
118    #[inline]
119    pub fn periodic_queries(self, period: Duration) -> RecoveryConfig<true> {
120        RecoveryConfig {
121            periodic_queries: Some(period),
122            heartbeat: false,
123            retention_period: self.retention_period,
124        }
125    }
126
127    /// Subscribe to heartbeats of [`AdvancedPublishers`](crate::AdvancedPublisher).
128    ///
129    /// This allows receiving the last published Sample's sequence number and check for misses.
130    /// Heartbeat subscriber must be paired with [`AdvancedPublishers`](crate::AdvancedPublisher)
131    /// that enable [`cache`](crate::AdvancedPublisherBuilder::cache) and
132    /// [`sample_miss_detection`](crate::AdvancedPublisherBuilder::sample_miss_detection) with
133    /// [`heartbeat`](crate::advanced_publisher::MissDetectionConfig::heartbeat) or
134    /// [`sporadic_heartbeat`](crate::advanced_publisher::MissDetectionConfig::sporadic_heartbeat).
135    #[zenoh_macros::unstable]
136    #[inline]
137    pub fn heartbeat(self) -> RecoveryConfig<true> {
138        RecoveryConfig {
139            periodic_queries: None,
140            heartbeat: true,
141            retention_period: self.retention_period,
142        }
143    }
144}
145
146#[zenoh_macros::unstable]
147impl<const CONFIGURED: bool> RecoveryConfig<CONFIGURED> {
148    const RETENTION_PERIOD_DEFAULT: Duration = Duration::from_secs(3600);
149
150    /// Set the retention period of publishers last Sample state (default to 1h).
151    #[zenoh_macros::unstable]
152    #[inline]
153    pub fn retention_period(mut self, period: Duration) -> RecoveryConfig<CONFIGURED> {
154        self.retention_period = Some(period);
155        self
156    }
157}
158
159/// The builder of an [`AdvancedSubscriber`], allowing to configure it.
160#[zenoh_macros::unstable]
161pub struct AdvancedSubscriberBuilder<'a, 'b, 'c, Handler, const BACKGROUND: bool = false> {
162    pub(crate) session: &'a Session,
163    pub(crate) key_expr: ZResult<KeyExpr<'b>>,
164    pub(crate) origin: Locality,
165    pub(crate) retransmission: Option<RecoveryConfig>,
166    pub(crate) query_target: QueryTarget,
167    pub(crate) query_timeout: Duration,
168    pub(crate) history: Option<HistoryConfig>,
169    pub(crate) liveliness: bool,
170    pub(crate) meta_key_expr: Option<ZResult<KeyExpr<'c>>>,
171    pub(crate) handler: Handler,
172}
173
174#[zenoh_macros::unstable]
175impl<Handler, const BACKGROUND: bool> fmt::Debug
176    for AdvancedSubscriberBuilder<'_, '_, '_, Handler, BACKGROUND>
177{
178    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
179        f.debug_struct("AdvancedSubscriberBuilder")
180            .field("session", &"..")
181            .field("key_expr", &self.key_expr)
182            .field("origin", &self.origin)
183            .field("retransmission", &self.retransmission)
184            .field("query_target", &self.query_target)
185            .field("query_timeout", &self.query_timeout)
186            .field("history", &self.history)
187            .field("liveliness", &self.liveliness)
188            .field("meta_key_expr", &self.meta_key_expr)
189            .field("handler", &"..")
190            .field("background", &BACKGROUND)
191            .finish()
192    }
193}
194
195#[zenoh_macros::unstable]
196impl<'a, 'b, Handler> AdvancedSubscriberBuilder<'a, 'b, '_, Handler> {
197    #[zenoh_macros::unstable]
198    pub(crate) fn new(builder: SubscriberBuilder<'a, 'b, Handler>) -> Self {
199        AdvancedSubscriberBuilder {
200            session: builder.session,
201            key_expr: builder.key_expr,
202            origin: builder.origin,
203            handler: builder.handler,
204            retransmission: None,
205            query_target: QueryTarget::All,
206            query_timeout: Duration::from_secs(10),
207            history: None,
208            liveliness: false,
209            meta_key_expr: None,
210        }
211    }
212}
213
214#[zenoh_macros::unstable]
215impl<'a, 'b, 'c> AdvancedSubscriberBuilder<'a, 'b, 'c, DefaultHandler> {
216    /// Add callback to AdvancedSubscriber.
217    #[inline]
218    #[zenoh_macros::unstable]
219    pub fn callback<F>(self, callback: F) -> AdvancedSubscriberBuilder<'a, 'b, 'c, Callback<Sample>>
220    where
221        F: Fn(Sample) + Send + Sync + 'static,
222    {
223        self.with(Callback::from(callback))
224    }
225
226    /// Add callback to `AdvancedSubscriber`.
227    ///
228    /// Using this guarantees that your callback will never be called concurrently.
229    /// If your callback is also accepted by the [`callback`](AdvancedSubscriberBuilder::callback) method, we suggest you use it instead of `callback_mut`
230    #[inline]
231    #[zenoh_macros::unstable]
232    pub fn callback_mut<F>(
233        self,
234        callback: F,
235    ) -> AdvancedSubscriberBuilder<'a, 'b, 'c, Callback<Sample>>
236    where
237        F: FnMut(Sample) + Send + Sync + 'static,
238    {
239        self.callback(locked(callback))
240    }
241
242    /// Make the built AdvancedSubscriber an [`AdvancedSubscriber`](AdvancedSubscriber).
243    #[inline]
244    #[zenoh_macros::unstable]
245    pub fn with<Handler>(self, handler: Handler) -> AdvancedSubscriberBuilder<'a, 'b, 'c, Handler>
246    where
247        Handler: IntoHandler<Sample>,
248    {
249        AdvancedSubscriberBuilder {
250            session: self.session,
251            key_expr: self.key_expr,
252            origin: self.origin,
253            retransmission: self.retransmission,
254            query_target: self.query_target,
255            query_timeout: self.query_timeout,
256            history: self.history,
257            liveliness: self.liveliness,
258            meta_key_expr: self.meta_key_expr,
259            handler,
260        }
261    }
262}
263
264#[zenoh_macros::unstable]
265impl<'a, 'b, 'c> AdvancedSubscriberBuilder<'a, 'b, 'c, Callback<Sample>> {
266    /// Make the subscriber run in background until the session is closed.
267    ///
268    /// Background builder doesn't return a `AdvancedSubscriber` object anymore.
269    pub fn background(self) -> AdvancedSubscriberBuilder<'a, 'b, 'c, Callback<Sample>, true> {
270        AdvancedSubscriberBuilder {
271            session: self.session,
272            key_expr: self.key_expr,
273            origin: self.origin,
274            retransmission: self.retransmission,
275            query_target: self.query_target,
276            query_timeout: self.query_timeout,
277            history: self.history,
278            liveliness: self.liveliness,
279            meta_key_expr: self.meta_key_expr,
280            handler: self.handler,
281        }
282    }
283}
284
285#[zenoh_macros::unstable]
286impl<'a, 'c, Handler, const BACKGROUND: bool>
287    AdvancedSubscriberBuilder<'a, '_, 'c, Handler, BACKGROUND>
288{
289    /// Restrict the matching publications that will be received by this [`Subscriber`] to the ones that have the given [`Locality`](crate::prelude::Locality).
290    #[zenoh_macros::unstable]
291    #[inline]
292    pub fn allowed_origin(mut self, origin: Locality) -> Self {
293        self.origin = origin;
294        self
295    }
296
297    /// Ask for retransmission of detected lost Samples.
298    ///
299    /// Retransmission can only be achieved by [`AdvancedPublishers`](crate::AdvancedPublisher)
300    /// that enable [`cache`](crate::AdvancedPublisherBuilder::cache) and
301    /// [`sample_miss_detection`](crate::AdvancedPublisherBuilder::sample_miss_detection).
302    #[zenoh_macros::unstable]
303    #[inline]
304    pub fn recovery(mut self, conf: RecoveryConfig) -> Self {
305        self.retransmission = Some(conf);
306        self
307    }
308
309    // /// Change the target to be used for queries.
310
311    // #[inline]
312    // pub fn query_target(mut self, query_target: QueryTarget) -> Self {
313    //     self.query_target = query_target;
314    //     self
315    // }
316
317    /// Change the timeout to be used for queries (history, retransmission).
318    #[zenoh_macros::unstable]
319    #[inline]
320    pub fn query_timeout(mut self, query_timeout: Duration) -> Self {
321        self.query_timeout = query_timeout;
322        self
323    }
324
325    /// Enable query for historical data.
326    ///
327    /// History can only be retransmitted by [`AdvancedPublishers`](crate::AdvancedPublisher) that enable [`cache`](crate::AdvancedPublisherBuilder::cache).
328    #[zenoh_macros::unstable]
329    #[inline]
330    pub fn history(mut self, config: HistoryConfig) -> Self {
331        self.history = Some(config);
332        self
333    }
334
335    /// Allow this subscriber to be detected through liveliness.
336    #[zenoh_macros::unstable]
337    pub fn subscriber_detection(mut self) -> Self {
338        self.liveliness = true;
339        self
340    }
341
342    /// A key expression added to the liveliness token key expression.
343    ///
344    /// It can be used to convey metadata.
345    #[zenoh_macros::unstable]
346    pub fn subscriber_detection_metadata<TryIntoKeyExpr>(mut self, meta: TryIntoKeyExpr) -> Self
347    where
348        TryIntoKeyExpr: TryInto<KeyExpr<'c>>,
349        <TryIntoKeyExpr as TryInto<KeyExpr<'c>>>::Error: Into<zenoh::Error>,
350    {
351        self.meta_key_expr = Some(meta.try_into().map_err(Into::into));
352        self
353    }
354
355    #[zenoh_macros::unstable]
356    fn with_static_keys(self) -> AdvancedSubscriberBuilder<'a, 'static, 'static, Handler> {
357        AdvancedSubscriberBuilder {
358            session: self.session,
359            key_expr: self.key_expr.map(|s| s.into_owned()),
360            origin: self.origin,
361            retransmission: self.retransmission,
362            query_target: self.query_target,
363            query_timeout: self.query_timeout,
364            history: self.history,
365            liveliness: self.liveliness,
366            meta_key_expr: self.meta_key_expr.map(|s| s.map(|s| s.into_owned())),
367            handler: self.handler,
368        }
369    }
370}
371
372#[zenoh_macros::unstable]
373impl<Handler> Resolvable for AdvancedSubscriberBuilder<'_, '_, '_, Handler>
374where
375    Handler: IntoHandler<Sample>,
376    Handler::Handler: Send,
377{
378    type To = ZResult<AdvancedSubscriber<Handler::Handler>>;
379}
380
381#[zenoh_macros::unstable]
382impl<Handler> Wait for AdvancedSubscriberBuilder<'_, '_, '_, Handler>
383where
384    Handler: IntoHandler<Sample> + Send,
385    Handler::Handler: Send,
386{
387    #[zenoh_macros::unstable]
388    fn wait(self) -> <Self as Resolvable>::To {
389        AdvancedSubscriber::new(self.with_static_keys())
390    }
391}
392
393#[zenoh_macros::unstable]
394impl<Handler> IntoFuture for AdvancedSubscriberBuilder<'_, '_, '_, Handler>
395where
396    Handler: IntoHandler<Sample> + Send,
397    Handler::Handler: Send,
398{
399    type Output = <Self as Resolvable>::To;
400    type IntoFuture = Ready<<Self as Resolvable>::To>;
401
402    #[zenoh_macros::unstable]
403    fn into_future(self) -> Self::IntoFuture {
404        std::future::ready(self.wait())
405    }
406}
407
408#[zenoh_macros::unstable]
409impl Resolvable for AdvancedSubscriberBuilder<'_, '_, '_, Callback<Sample>, true> {
410    type To = ZResult<()>;
411}
412
413#[zenoh_macros::unstable]
414impl Wait for AdvancedSubscriberBuilder<'_, '_, '_, Callback<Sample>, true> {
415    #[zenoh_macros::unstable]
416    fn wait(self) -> <Self as Resolvable>::To {
417        let mut sub = AdvancedSubscriber::new(self.with_static_keys())?;
418        sub.set_background_impl(true);
419        Ok(())
420    }
421}
422
423#[zenoh_macros::unstable]
424impl IntoFuture for AdvancedSubscriberBuilder<'_, '_, '_, Callback<Sample>, true> {
425    type Output = <Self as Resolvable>::To;
426    type IntoFuture = Ready<<Self as Resolvable>::To>;
427
428    #[zenoh_macros::unstable]
429    fn into_future(self) -> Self::IntoFuture {
430        std::future::ready(self.wait())
431    }
432}
433
434#[zenoh_macros::unstable]
435struct State {
436    next_id: usize,
437    global_pending_queries: u64,
438    sequenced_states: LruCache<EntityGlobalId, SourceState<WrappingSn>>,
439    timestamped_states: LruCache<ID, SourceState<Timestamp>>,
440    session: WeakSession,
441    key_expr: KeyExpr<'static>,
442    retransmission: bool,
443    period: Option<Duration>,
444    max_history_depth: usize,
445    query_target: QueryTarget,
446    query_timeout: Duration,
447    // Callback must be dropped when the underlying subscriber is undeclared
448    // (for example when session is closed), in order to "close" the advanced
449    // subscriber receiver, hence the `Option`.
450    callback: Option<Callback<Sample>>,
451    miss_handlers: HashMap<usize, Callback<Miss>>,
452    token: Option<LivelinessToken>,
453    _gc_task: AbortOnDropHandle<()>,
454}
455
456#[zenoh_macros::unstable]
457impl State {
458    #[zenoh_macros::unstable]
459    fn register_miss_callback(&mut self, callback: Callback<Miss>) -> usize {
460        let id = self.next_id;
461        self.next_id += 1;
462        self.miss_handlers.insert(id, callback);
463        id
464    }
465    #[zenoh_macros::unstable]
466    fn unregister_miss_callback(&mut self, id: &usize) {
467        self.miss_handlers.remove(id);
468    }
469}
470
471#[zenoh_macros::unstable]
472struct SourceState<T> {
473    last_delivered: Option<T>,
474    pending_queries: u64,
475    pending_samples: BTreeMap<T, Sample>,
476    /// Latest access instant used for garbage collection with retention period
477    latest_access: Instant,
478    /// Periodic queries task
479    periodic_task: Option<AbortOnDropHandle<()>>,
480    /// Alive as per liveliness subscriber
481    alive: bool,
482}
483
484impl<T> Default for SourceState<T> {
485    fn default() -> Self {
486        Self {
487            last_delivered: None,
488            pending_queries: 0,
489            pending_samples: BTreeMap::new(),
490            periodic_task: None,
491            alive: false,
492            latest_access: Instant::now(),
493        }
494    }
495}
496
497/*
498use zenoh_ext::{AdvancedSubscriberBuilderExt, HistoryConfig, RecoveryConfig};
499
500let session = zenoh::open(zenoh::Config::default()).await.unwrap();
501let subscriber = session
502    .declare_subscriber("key/expression")
503    .history(HistoryConfig::default().detect_late_publishers())
504    .recovery(RecoveryConfig::default())
505    .await
506    .unwrap();
507
508let miss_listener = subscriber.sample_miss_listener().await.unwrap();
509loop {
510    tokio::select! {
511        sample = subscriber.recv_async() => {
512            if let Ok(sample) = sample {
513                // ...
514            }
515        },
516        miss = miss_listener.recv_async() => {
517            if let Ok(miss) = miss {
518                // ...
519            }
520        },
521    }
522}
523*/
524
525/// The extension to [`Subscriber`](zenoh::pubsub::Subscriber) that provides advanced functionalities
526///
527/// The `AdvancedSubscriber` is constructed over a regular [`Subscriber`](zenoh::pubsub::Subscriber)
528/// through [`advanced`](crate::AdvancedSubscriberBuilderExt::advanced) method or by using
529/// any other method of [`AdvancedSubscriberBuilder`](crate::AdvancedSubscriberBuilder).
530///
531/// The `AdvancedSubscriber` works with [`AdvancedPublisher`](crate::AdvancedPublisher) to provide additional functionalities such as:
532/// * missing samples detection using periodic queries or heartbeat subscription configurable with [`recovery`](crate::AdvancedSubscriberBuilder::recovery) method
533/// * recovering missing samples, configured with [`history`](crate::AdvancedSubscriberBuilder::history) method
534///   (max age and sample count, late joiner detection and requesting)
535/// * liveliness-based subscriber detection with [`subscriber_detection`](crate::AdvancedSubscriberBuilder::subscriber_detection) method
536///
537/// # Examples
538/// ```no_run
539/// # #[tokio::main]
540/// # async fn main() {
541/// use zenoh_ext::{AdvancedSubscriberBuilderExt, HistoryConfig, RecoveryConfig};
542/// let session = zenoh::open(zenoh::Config::default()).await.unwrap();
543/// let subscriber = session
544///     .declare_subscriber("key/expression")
545///     .history(HistoryConfig::default().detect_late_publishers())
546///     .recovery(RecoveryConfig::default().heartbeat())
547///     .subscriber_detection()
548///     .await
549///     .unwrap();
550/// let miss_listener = subscriber.sample_miss_listener().await.unwrap();
551/// loop {
552///     tokio::select! {
553///         sample = subscriber.recv_async() => {
554///             if let Ok(sample) = sample {
555///                 // ...
556///             }
557///         },
558///         miss = miss_listener.recv_async() => {
559///             if let Ok(miss) = miss {
560///                 // ...
561///             }
562///         },
563///     }
564/// }
565/// # }
566/// ```
567#[zenoh_macros::unstable]
568pub struct AdvancedSubscriber<Receiver> {
569    statesref: Arc<Mutex<State>>,
570    subscriber: Subscriber<()>,
571    receiver: Receiver,
572    liveliness_subscriber: Option<Subscriber<()>>,
573    heartbeat_subscriber: Option<Subscriber<()>>,
574}
575
576#[zenoh_macros::unstable]
577impl<Receiver> fmt::Debug for AdvancedSubscriber<Receiver> {
578    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
579        f.debug_struct("AdvancedSubscriber")
580            .field("statesref", &"..")
581            .field("subscriber", &self.subscriber)
582            .field("receiver", &"..")
583            .field("liveliness_subscriber", &self.liveliness_subscriber)
584            .field("heartbeat_subscriber", &self.heartbeat_subscriber)
585            .finish()
586    }
587}
588
589#[zenoh_macros::unstable]
590impl<Receiver> std::ops::Deref for AdvancedSubscriber<Receiver> {
591    type Target = Receiver;
592    fn deref(&self) -> &Self::Target {
593        &self.receiver
594    }
595}
596
597#[zenoh_macros::unstable]
598impl<Receiver> std::ops::DerefMut for AdvancedSubscriber<Receiver> {
599    fn deref_mut(&mut self) -> &mut Self::Target {
600        &mut self.receiver
601    }
602}
603
604#[zenoh_macros::unstable]
605fn handle_sample(states: &mut State, sample: Sample) -> bool {
606    let Some(callback) = states.callback.as_ref() else {
607        return false;
608    };
609    if let Some(source_info) = sample.source_info().cloned() {
610        #[inline]
611        fn deliver_and_flush(
612            sample: Sample,
613            source_sn: impl Into<WrappingSn>,
614            callback: &Callback<Sample>,
615            state: &mut SourceState<WrappingSn>,
616        ) {
617            let mut source_sn = source_sn.into();
618            callback.call(sample);
619            state.last_delivered = Some(source_sn);
620            while let Some(sample) = state.pending_samples.remove(&(source_sn + 1)) {
621                callback.call(sample);
622                source_sn += 1;
623                state.last_delivered = Some(source_sn);
624            }
625        }
626
627        let mut new = false;
628        let state = states
629            .sequenced_states
630            .get_or_insert_mut(*source_info.source_id(), || {
631                new = true;
632                Default::default()
633            });
634        if state.last_delivered.is_none() && states.global_pending_queries != 0 {
635            // Avoid going through the Map if history_depth == 1
636            if states.max_history_depth == 1 {
637                state.last_delivered = Some(source_info.source_sn().into());
638                callback.call(sample);
639            } else {
640                state
641                    .pending_samples
642                    .insert(source_info.source_sn().into(), sample);
643                if state.pending_samples.len() >= states.max_history_depth {
644                    if let Some((sn, sample)) = state.pending_samples.pop_first() {
645                        deliver_and_flush(sample, sn, callback, state);
646                    }
647                }
648            }
649        } else if state.last_delivered.is_some()
650            && source_info.source_sn() != state.last_delivered.unwrap() + 1
651        {
652            if source_info.source_sn() > state.last_delivered.unwrap() {
653                if states.retransmission {
654                    state
655                        .pending_samples
656                        .insert(source_info.source_sn().into(), sample);
657                } else {
658                    tracing::info!(
659                        "Sample missed: missed {} samples from {:?}.",
660                        source_info.source_sn() - state.last_delivered.unwrap() - 1,
661                        source_info.source_id(),
662                    );
663                    for miss_callback in states.miss_handlers.values() {
664                        miss_callback.call(Miss {
665                            source: *source_info.source_id(),
666                            nb: source_info.source_sn() - state.last_delivered.unwrap() - 1,
667                        });
668                    }
669                    callback.call(sample);
670                    state.last_delivered = Some(source_info.source_sn().into());
671                }
672            }
673        } else {
674            deliver_and_flush(sample, source_info.source_sn(), callback, state);
675        }
676        state.latest_access = Instant::now();
677        new
678    } else if let Some(timestamp) = sample.timestamp() {
679        let state = states
680            .timestamped_states
681            .get_or_insert_mut(*timestamp.get_id(), Default::default);
682        if state.last_delivered.map(|t| t < *timestamp).unwrap_or(true) {
683            if (states.global_pending_queries == 0 && state.pending_queries == 0)
684                || states.max_history_depth == 1
685            {
686                state.last_delivered = Some(*timestamp);
687                callback.call(sample);
688            } else {
689                state.pending_samples.entry(*timestamp).or_insert(sample);
690                if state.pending_samples.len() >= states.max_history_depth {
691                    flush_timestamped_source(state, Some(callback));
692                }
693            }
694        }
695        state.latest_access = Instant::now();
696        false
697    } else {
698        callback.call(sample);
699        false
700    }
701}
702
703#[zenoh_macros::unstable]
704fn seq_num_range(start: Option<WrappingSn>, end: Option<WrappingSn>) -> String {
705    match (start, end) {
706        (Some(start), Some(end)) => format!("_sn={start}..{end}"),
707        (Some(start), None) => format!("_sn={start}.."),
708        (None, Some(end)) => format!("_sn=..{end}"),
709        (None, None) => "_sn=..".to_string(),
710    }
711}
712
713fn spawn_periodic_queries(
714    statesref: &Arc<Mutex<State>>,
715    period: Option<Duration>,
716    source_id: EntityGlobalId,
717) -> Option<AbortOnDropHandle<()>> {
718    let period = period?;
719    let statesref = statesref.clone();
720    Some(AbortOnDropHandle::new(ZRuntime::Application.spawn(
721        async move {
722            let mut interval = tokio::time::interval(period);
723            interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
724            interval.tick().await; // interval first tick is immediate
725            loop {
726                interval.tick().await;
727                periodic_query(&statesref, source_id);
728            }
729        },
730    )))
731}
732
733fn periodic_query(statesref: &Arc<Mutex<State>>, source_id: EntityGlobalId) {
734    let mut guard = statesref.lock().unwrap();
735    let states = &mut *guard;
736    // use peek_mut so query without samples do not prevent the state to be garbage collected
737    let Some(state) = states.sequenced_states.peek_mut(&source_id) else {
738        return;
739    };
740    state.pending_queries += 1;
741    let query_expr = &states.key_expr
742        / KE_ADV_PREFIX
743        / KE_STAR
744        / &source_id.zid().into_keyexpr()
745        / &KeyExpr::try_from(source_id.eid().to_string()).unwrap()
746        / KE_STARSTAR;
747    let seq_num_range = seq_num_range(state.last_delivered.map(|s| s + 1), None);
748
749    let session = states.session.clone();
750    let key_expr = states.key_expr.clone().into_owned();
751    let query_target = states.query_target;
752    let query_timeout = states.query_timeout;
753
754    tracing::trace!(
755        "AdvancedSubscriber{{key_expr: {}}}: Querying undelivered samples {}?{}",
756        states.key_expr,
757        query_expr,
758        seq_num_range
759    );
760    drop(guard);
761
762    let handler = SequencedRepliesHandler {
763        source_id,
764        statesref: statesref.clone(),
765    };
766    let _ = session
767        .get(Selector::from((query_expr, seq_num_range)))
768        .callback({
769            move |r: Reply| {
770                if let Ok(s) = r.into_result() {
771                    if key_expr.intersects(s.key_expr()) {
772                        let states = &mut *zlock!(handler.statesref);
773                        tracing::trace!(
774                            "AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}",
775                            states.key_expr,
776                            s.source_info(),
777                            s.timestamp()
778                        );
779                        handle_sample(states, s);
780                    }
781                }
782            }
783        })
784        .consolidation(ConsolidationMode::None)
785        .accept_replies(ReplyKeyExpr::Any)
786        .target(query_target)
787        .timeout(query_timeout)
788        .wait();
789}
790
791/// Garbage collects the source states' lists.
792///
793/// Reclamation is based on `SourceState::latest_access`; alive publishers are not reclaimed.
794async fn gc_task(statesref: Weak<Mutex<State>>, retention_period: Duration) {
795    /// Garbage collect a lists and return the oldest access.
796    fn garbage_collect<K: Copy + Eq + Hash, T>(
797        states: &mut LruCache<K, SourceState<T>>,
798        retention_period: Duration,
799        now: Instant,
800    ) -> Instant {
801        while let Some((&key, state)) = states.peek_lru() {
802            if now.duration_since(state.latest_access) <= retention_period {
803                return state.latest_access;
804            // if the publisher is still marked as alive, just update its latest access
805            // (accessing the state will also move it to the back of the LRU list)
806            } else if state.alive {
807                states.get_mut(&key).unwrap().latest_access = now;
808            } else {
809                states.pop_lru();
810            }
811        }
812        now
813    }
814    // start by sleeping for the initial retention period
815    tokio::time::sleep(retention_period).await;
816    loop {
817        let oldest_access = {
818            let Some(states) = statesref.upgrade() else {
819                // either the task was scheduled concurrently to its abortion, so we don't care
820                // sleeping, or we are in the theoretically possible but zero probability case
821                // of `new_cyclic` not having returned yet, so we still don't care sleeping.
822                tokio::time::sleep(retention_period).await;
823                continue;
824            };
825            let mut states = states.lock().unwrap();
826            let now = Instant::now();
827            min(
828                garbage_collect(&mut states.sequenced_states, retention_period, now),
829                garbage_collect(&mut states.timestamped_states, retention_period, now),
830            )
831        };
832        tokio::time::sleep_until((oldest_access + retention_period).into()).await;
833    }
834}
835
836#[zenoh_macros::unstable]
837impl<Handler> AdvancedSubscriber<Handler> {
838    fn new<H>(conf: AdvancedSubscriberBuilder<'_, '_, '_, H>) -> ZResult<Self>
839    where
840        H: IntoHandler<Sample, Handler = Handler> + Send,
841    {
842        // Check config
843        if let Some(history) = conf.history.as_ref() {
844            if history.max_samples.is_some_and(|d| d == 0) {
845                bail!("max_samples must not be zero")
846            }
847            if history.max_age.is_some_and(|a| a == 0.0) {
848                bail!("max_age must not be zero")
849            }
850        }
851        let (callback, receiver) = conf.handler.into_handler();
852        let key_expr = conf.key_expr?;
853        let meta = match conf.meta_key_expr {
854            Some(meta) => Some(meta?),
855            None => None,
856        };
857        let retransmission = conf.retransmission;
858        let query_target = conf.query_target;
859        let query_timeout = conf.query_timeout;
860        let max_history_depth = conf
861            .history
862            .as_ref()
863            .and_then(|h| h.max_samples)
864            // If the query is not bounded with `_max`, then it can receive unbounded number of responses
865            .unwrap_or(usize::MAX);
866        let retention_period = retransmission
867            .as_ref()
868            .and_then(|r| r.retention_period)
869            .unwrap_or(RecoveryConfig::<true>::RETENTION_PERIOD_DEFAULT);
870        let statesref = Arc::new_cyclic(|weak| {
871            Mutex::new(State {
872                next_id: 0,
873                sequenced_states: LruCache::unbounded(),
874                timestamped_states: LruCache::unbounded(),
875                global_pending_queries: if conf.history.is_some() { 1 } else { 0 },
876                session: conf.session.downgrade(),
877                period: retransmission.as_ref().and_then(|r| r.periodic_queries),
878                key_expr: key_expr.clone().into_owned(),
879                retransmission: retransmission.is_some(),
880                max_history_depth,
881                query_target: conf.query_target,
882                query_timeout: conf.query_timeout,
883                callback: Some(callback),
884                miss_handlers: HashMap::new(),
885                token: None,
886                _gc_task: AbortOnDropHandle::new(
887                    ZRuntime::Application.spawn(gc_task(weak.clone(), retention_period)),
888                ),
889            })
890        });
891
892        let sub_callback = {
893            let statesref = statesref.clone();
894            let session = conf.session.downgrade();
895            let key_expr = key_expr.clone().into_owned();
896
897            move |s: Sample| {
898                let mut lock = zlock!(statesref);
899                let states = &mut *lock;
900                let source_id = s.source_info().map(|si| *si.source_id());
901                let new = handle_sample(states, s);
902
903                if let Some(source_id) = source_id {
904                    if let Some(state) = states.sequenced_states.get_mut(&source_id) {
905                        if new {
906                            state.periodic_task =
907                                spawn_periodic_queries(&statesref, states.period, source_id);
908                        }
909                        if retransmission.is_some()
910                            && state.pending_queries == 0
911                            && !state.pending_samples.is_empty()
912                        {
913                            state.pending_queries += 1;
914                            let query_expr = &key_expr
915                                / KE_ADV_PREFIX
916                                / KE_STAR
917                                / &source_id.zid().into_keyexpr()
918                                / &KeyExpr::try_from(source_id.eid().to_string()).unwrap()
919                                / KE_STARSTAR;
920                            let seq_num_range =
921                                seq_num_range(state.last_delivered.map(|s| s + 1), None);
922                            tracing::trace!(
923                                "AdvancedSubscriber{{key_expr: {}}}: Querying missing samples {}?{}",
924                                states.key_expr,
925                                query_expr,
926                                seq_num_range
927                            );
928                            drop(lock);
929                            let handler = SequencedRepliesHandler {
930                                source_id,
931                                statesref: statesref.clone(),
932                            };
933                            let _ = session
934                                .get(Selector::from((query_expr, seq_num_range)))
935                                .callback({
936                                    let key_expr = key_expr.clone().into_owned();
937                                    move |r: Reply| {
938                                        if let Ok(s) = r.into_result() {
939                                            if key_expr.intersects(s.key_expr()) {
940                                                let states = &mut *zlock!(handler.statesref);
941                                                tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}", states.key_expr, s.source_info(), s.timestamp());
942                                                handle_sample(states, s);
943                                            }
944                                        }
945                                    }
946                                })
947                                .consolidation(ConsolidationMode::None)
948                                .accept_replies(ReplyKeyExpr::Any)
949                                .target(query_target)
950                                .timeout(query_timeout)
951                                .wait();
952                        }
953                    }
954                }
955            }
956        };
957
958        // When the underlying subscriber is undeclared (for example when the session is closed)
959        // the advanced subscriber callback must be dropped to "close" the receiver.
960        let drop_callback = {
961            let statesref = statesref.clone();
962            move || {
963                let mut states = statesref.lock().unwrap();
964                states.callback.take();
965                states.miss_handlers.clear();
966            }
967        };
968
969        let subscriber = conf
970            .session
971            .declare_subscriber(&key_expr)
972            .with(CallbackDrop {
973                callback: sub_callback,
974                drop: drop_callback,
975            })
976            .allowed_origin(conf.origin)
977            .wait()?;
978
979        tracing::debug!("Create AdvancedSubscriber{{key_expr: {}}}", key_expr,);
980
981        if let Some(historyconf) = conf.history.as_ref() {
982            let handler = InitialRepliesHandler {
983                statesref: statesref.clone(),
984            };
985            let mut params = Parameters::empty();
986            if let Some(max) = historyconf.max_samples {
987                params.insert("_max", max.to_string());
988            }
989            if let Some(age) = historyconf.max_age {
990                params.set_time_range(TimeRange {
991                    start: TimeBound::Inclusive(TimeExpr::Now { offset_secs: -age }),
992                    end: TimeBound::Unbounded,
993                });
994            }
995            tracing::trace!(
996                "AdvancedSubscriber{{key_expr: {}}} Querying historical samples {}?{}",
997                key_expr,
998                &key_expr / KE_ADV_PREFIX / KE_STARSTAR,
999                params
1000            );
1001            let _ = conf
1002                .session
1003                .get(Selector::from((
1004                    &key_expr / KE_ADV_PREFIX / KE_STARSTAR,
1005                    params,
1006                )))
1007                .callback({
1008                    let key_expr = key_expr.clone().into_owned();
1009                    move |r: Reply| {
1010                        if let Ok(s) = r.into_result() {
1011                            if key_expr.intersects(s.key_expr()) {
1012                                let states = &mut *zlock!(handler.statesref);
1013                                tracing::trace!(
1014                                    "AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}",
1015                                    states.key_expr,
1016                                    s.source_info(),
1017                                    s.timestamp()
1018                                );
1019                                handle_sample(states, s);
1020                            }
1021                        }
1022                    }
1023                })
1024                .consolidation(ConsolidationMode::None)
1025                .accept_replies(ReplyKeyExpr::Any)
1026                .target(query_target)
1027                .timeout(query_timeout)
1028                .wait();
1029        }
1030
1031        let liveliness_subscriber = if let Some(historyconf) = conf.history.as_ref() {
1032            if historyconf.liveliness {
1033                let live_callback = {
1034                    let session = conf.session.downgrade();
1035                    let statesref = statesref.clone();
1036                    let key_expr = key_expr.clone().into_owned();
1037                    let historyconf = historyconf.clone();
1038                    move |s: Sample| {
1039                        let Ok(parsed) = ke_liveliness::parse(s.key_expr().as_keyexpr()) else {
1040                            tracing::warn!(
1041                                "AdvancedSubscriber{{}}: Received malformed liveliness token key expression: {}",
1042                                s.key_expr()
1043                            );
1044                            return;
1045                        };
1046                        if let Ok(zid) = ZenohId::from_str(parsed.zid().as_str()) {
1047                            // TODO : If we already have a state associated to this discovered source
1048                            // we should query with the appropriate range to avoid unnecessary retransmissions
1049                            if parsed.eid() == KE_UHLC {
1050                                let mut lock = zlock!(statesref);
1051                                let states = &mut *lock;
1052                                if s.kind() == SampleKind::Delete {
1053                                    tracing::trace!(
1054                                        "AdvancedSubscriber{{key_expr: {}}}: Liveliness loss for publishers with zid={}",
1055                                        states.key_expr,
1056                                        parsed.zid().as_str()
1057                                    );
1058                                    if let Some(state) =
1059                                        states.timestamped_states.peek_mut(&ID::from(zid))
1060                                    {
1061                                        state.alive = false;
1062                                    }
1063                                    return;
1064                                }
1065                                tracing::trace!(
1066                                    "AdvancedSubscriber{{key_expr: {}}}: Detect late joiner publishers with zid={}",
1067                                    states.key_expr,
1068                                    parsed.zid().as_str()
1069                                );
1070                                let state = states
1071                                    .timestamped_states
1072                                    .get_or_insert_mut(ID::from(zid), Default::default);
1073                                state.pending_queries += 1;
1074                                state.alive = true;
1075                                state.latest_access = Instant::now();
1076
1077                                let mut params = Parameters::empty();
1078                                if let Some(max) = historyconf.max_samples {
1079                                    params.insert("_max", max.to_string());
1080                                }
1081                                if let Some(age) = historyconf.max_age {
1082                                    params.set_time_range(TimeRange {
1083                                        start: TimeBound::Inclusive(TimeExpr::Now {
1084                                            offset_secs: -age,
1085                                        }),
1086                                        end: TimeBound::Unbounded,
1087                                    });
1088                                }
1089                                tracing::trace!(
1090                                    "AdvancedSubscriber{{key_expr: {}}}: Querying historical samples {}?{}",
1091                                    states.key_expr,
1092                                    s.key_expr(),
1093                                    params
1094                                );
1095                                drop(lock);
1096
1097                                let handler = TimestampedRepliesHandler {
1098                                    id: ID::from(zid),
1099                                    statesref: statesref.clone(),
1100                                };
1101                                let _ = session
1102                                    .get(Selector::from((s.key_expr(), params)))
1103                                    .callback({
1104                                        let key_expr = key_expr.clone().into_owned();
1105                                        move |r: Reply| {
1106                                            if let Ok(s) = r.into_result() {
1107                                                if key_expr.intersects(s.key_expr()) {
1108                                                    let states =
1109                                                        &mut *zlock!(handler.statesref);
1110                                                    tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}", states.key_expr, s.source_info(), s.timestamp());
1111                                                    handle_sample(states, s);
1112                                                }
1113                                            }
1114                                        }
1115                                    })
1116                                    .consolidation(ConsolidationMode::None)
1117                                    .accept_replies(ReplyKeyExpr::Any)
1118                                    .target(query_target)
1119                                    .timeout(query_timeout)
1120                                    .wait();
1121                            } else if let Ok(eid) = EntityId::from_str(parsed.eid().as_str()) {
1122                                let source_id = EntityGlobalId::new(zid, eid);
1123                                let mut lock = zlock!(statesref);
1124                                let states = &mut *lock;
1125                                if s.kind() == SampleKind::Delete {
1126                                    tracing::trace!(
1127                                        "AdvancedSubscriber{{key_expr: {}}}: Liveliness loss for publishers with zid={}",
1128                                        states.key_expr,
1129                                        parsed.zid().as_str()
1130                                    );
1131                                    if let Some(state) =
1132                                        states.sequenced_states.peek_mut(&source_id)
1133                                    {
1134                                        state.alive = false;
1135                                    }
1136                                    return;
1137                                }
1138                                tracing::trace!(
1139                                    "AdvancedSubscriber{{key_expr: {}}}: Detect late joiner publishers with zid={}",
1140                                    states.key_expr,
1141                                    parsed.zid().as_str()
1142                                );
1143                                let mut new = false;
1144                                let state =
1145                                    states.sequenced_states.get_or_insert_mut(source_id, || {
1146                                        new = true;
1147                                        Default::default()
1148                                    });
1149                                if new {
1150                                    state.periodic_task = spawn_periodic_queries(
1151                                        &statesref,
1152                                        states.period,
1153                                        source_id,
1154                                    );
1155                                }
1156                                state.pending_queries += 1;
1157                                state.alive = true;
1158                                state.latest_access = Instant::now();
1159
1160                                let mut params = Parameters::empty();
1161                                if let Some(max) = historyconf.max_samples {
1162                                    params.insert("_max", max.to_string());
1163                                }
1164                                if let Some(age) = historyconf.max_age {
1165                                    params.set_time_range(TimeRange {
1166                                        start: TimeBound::Inclusive(TimeExpr::Now {
1167                                            offset_secs: -age,
1168                                        }),
1169                                        end: TimeBound::Unbounded,
1170                                    });
1171                                }
1172                                tracing::trace!(
1173                                    "AdvancedSubscriber{{key_expr: {}}}: Querying historical samples {}?{}",
1174                                    states.key_expr,
1175                                    s.key_expr(),
1176                                    params,
1177                                );
1178                                drop(lock);
1179
1180                                let handler = SequencedRepliesHandler {
1181                                    source_id,
1182                                    statesref: statesref.clone(),
1183                                };
1184                                let _ = session
1185                                    .get(Selector::from((s.key_expr(), params)))
1186                                    .callback({
1187                                        let key_expr = key_expr.clone().into_owned();
1188                                        move |r: Reply| {
1189                                            if let Ok(s) = r.into_result() {
1190                                                if key_expr.intersects(s.key_expr()) {
1191                                                    let states = &mut *zlock!(handler.statesref);
1192                                                    tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}", states.key_expr, s.source_info(), s.timestamp());
1193                                                    handle_sample(states, s);
1194                                                }
1195                                            }
1196                                        }
1197                                    })
1198                                    .consolidation(ConsolidationMode::None)
1199                                    .accept_replies(ReplyKeyExpr::Any)
1200                                    .target(query_target)
1201                                    .timeout(query_timeout)
1202                                    .wait();
1203                            }
1204                        } else if s.kind() == SampleKind::Put {
1205                            let mut lock = zlock!(statesref);
1206                            let states = &mut *lock;
1207                            tracing::trace!(
1208                                "AdvancedSubscriber{{key_expr: {}}}: Detect late joiner publishers with zid={}",
1209                                states.key_expr,
1210                                parsed.zid().as_str()
1211                            );
1212                            states.global_pending_queries += 1;
1213
1214                            let mut params = Parameters::empty();
1215                            if let Some(max) = historyconf.max_samples {
1216                                params.insert("_max", max.to_string());
1217                            }
1218                            if let Some(age) = historyconf.max_age {
1219                                params.set_time_range(TimeRange {
1220                                    start: TimeBound::Inclusive(TimeExpr::Now {
1221                                        offset_secs: -age,
1222                                    }),
1223                                    end: TimeBound::Unbounded,
1224                                });
1225                            }
1226                            tracing::trace!(
1227                                "AdvancedSubscriber{{key_expr: {}}}: Querying historical samples {}?{}",
1228                                states.key_expr,
1229                                s.key_expr(),
1230                                params,
1231                            );
1232                            drop(lock);
1233
1234                            let handler = InitialRepliesHandler {
1235                                statesref: statesref.clone(),
1236                            };
1237                            let _ = session
1238                                .get(Selector::from((s.key_expr(), params)))
1239                                .callback({
1240                                    let key_expr = key_expr.clone().into_owned();
1241                                    move |r: Reply| {
1242                                        if let Ok(s) = r.into_result() {
1243                                            if key_expr.intersects(s.key_expr()) {
1244                                                let states = &mut *zlock!(handler.statesref);
1245                                                tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}", states.key_expr, s.source_info(), s.timestamp());
1246                                                handle_sample(states, s);
1247                                            }
1248                                        }
1249                                    }
1250                                })
1251                                .consolidation(ConsolidationMode::None)
1252                                .accept_replies(ReplyKeyExpr::Any)
1253                                .target(query_target)
1254                                .timeout(query_timeout)
1255                                .wait();
1256                        }
1257                    }
1258                };
1259
1260                tracing::debug!(
1261                    "AdvancedSubscriber{{key_expr: {}}}: Detect late joiner publishers on {}",
1262                    key_expr,
1263                    &key_expr / KE_ADV_PREFIX / KE_PUB / KE_STARSTAR
1264                );
1265                Some(
1266                    conf.session
1267                        .liveliness()
1268                        .declare_subscriber(&key_expr / KE_ADV_PREFIX / KE_PUB / KE_STARSTAR)
1269                        // .declare_subscriber(keformat!(ke_liveliness_all::formatter(), zid = 0, eid = 0, remaining = key_expr).unwrap())
1270                        .history(true)
1271                        .callback(live_callback)
1272                        .wait()?,
1273                )
1274            } else {
1275                None
1276            }
1277        } else {
1278            None
1279        };
1280
1281        let heartbeat_subscriber = if retransmission.is_some_and(|r| r.heartbeat) {
1282            let ke_heartbeat_sub = &key_expr / KE_ADV_PREFIX / KE_PUB / KE_STARSTAR;
1283            let statesref = statesref.clone();
1284            tracing::debug!(
1285                "AdvancedSubscriber{{key_expr: {}}}: Enable heartbeat subscriber on {}",
1286                key_expr,
1287                ke_heartbeat_sub
1288            );
1289            let heartbeat_sub = conf
1290                .session
1291                .declare_subscriber(ke_heartbeat_sub)
1292                .callback(move |sample_hb| {
1293                    if sample_hb.kind() != SampleKind::Put {
1294                        return;
1295                    }
1296
1297                    let heartbeat_keyexpr = sample_hb.key_expr().as_keyexpr();
1298                    let Ok(parsed_keyexpr) = ke_liveliness::parse(heartbeat_keyexpr) else {
1299                        return;
1300                    };
1301                    let source_id = {
1302                        let Ok(zid) = ZenohId::from_str(parsed_keyexpr.zid().as_str()) else {
1303                            return;
1304                        };
1305                        let Ok(eid) = EntityId::from_str(parsed_keyexpr.eid().as_str()) else {
1306                            return;
1307                        };
1308                        EntityGlobalId::new(zid, eid)
1309                    };
1310
1311                    let Ok(heartbeat_sn) = z_deserialize::<WrappingSn>(sample_hb.payload()) else {
1312                        tracing::debug!(
1313                            "AdvancedSubscriber{{}}: Skipping invalid heartbeat payload on '{}'",
1314                            heartbeat_keyexpr
1315                        );
1316                        return;
1317                    };
1318
1319                    let mut lock = zlock!(statesref);
1320                    let states = &mut *lock;
1321                    let mut new = false;
1322                    let state = states.sequenced_states.get_or_insert_mut(source_id, ||{
1323                        new = true;
1324                        Default::default()
1325                    });
1326                    state.latest_access = Instant::now();
1327                    if new {
1328                        // NOTE: API does not allow both heartbeat and periodic_queries
1329                        state.periodic_task = spawn_periodic_queries(&statesref, states.period, source_id);
1330                        if states.global_pending_queries > 0 {
1331                            tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Skipping heartbeat on '{}' from publisher that is currently being pulled by global query", states.key_expr, heartbeat_keyexpr);
1332                            return;
1333                        }
1334                    }
1335
1336                    // check that it's not an old sn, and that there are no pending queries
1337                    if (state.last_delivered.is_none()
1338                        || state.last_delivered.is_some_and(|sn| heartbeat_sn > sn))
1339                        && state.pending_queries == 0
1340                    {
1341                        let seq_num_range = seq_num_range(
1342                            state.last_delivered.map(|s| s + 1),
1343                            Some(heartbeat_sn),
1344                        );
1345
1346                        let session = states.session.clone();
1347                        let key_expr = states.key_expr.clone().into_owned();
1348                        let query_target = states.query_target;
1349                        let query_timeout = states.query_timeout;
1350                        state.pending_queries += 1;
1351
1352                        tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Querying missing samples {}?{}", states.key_expr, heartbeat_keyexpr, seq_num_range);
1353                        drop(lock);
1354
1355                        let handler = SequencedRepliesHandler {
1356                            source_id,
1357                            statesref: statesref.clone(),
1358                        };
1359                        let _ = session
1360                            .get(Selector::from((heartbeat_keyexpr, seq_num_range)))
1361                            .callback({
1362                                move |r: Reply| {
1363                                    if let Ok(s) = r.into_result() {
1364                                        if key_expr.intersects(s.key_expr()) {
1365                                            let states = &mut *zlock!(handler.statesref);
1366                                            tracing::trace!("AdvancedSubscriber{{key_expr: {}}}: Received reply with Sample{{info:{:?}, ts:{:?}}}", states.key_expr, s.source_info(), s.timestamp());
1367                                            handle_sample(states, s);
1368                                        }
1369                                    }
1370                                }
1371                            })
1372                            .consolidation(ConsolidationMode::None)
1373                            .accept_replies(ReplyKeyExpr::Any)
1374                            .target(query_target)
1375                            .timeout(query_timeout)
1376                            .wait();
1377                    }
1378                })
1379                .allowed_origin(conf.origin)
1380                .wait()?;
1381            Some(heartbeat_sub)
1382        } else {
1383            None
1384        };
1385
1386        if conf.liveliness {
1387            let suffix = KE_ADV_PREFIX
1388                / KE_SUB
1389                / &subscriber.id().zid().into_keyexpr()
1390                / &KeyExpr::try_from(subscriber.id().eid().to_string()).unwrap();
1391            let suffix = match meta {
1392                Some(meta) => suffix / &meta,
1393                // We need this empty chunk because of a routing matching bug
1394                _ => suffix / KE_EMPTY,
1395            };
1396            tracing::debug!(
1397                "AdvancedSubscriber{{key_expr: {}}}: Declare liveliness token {}",
1398                key_expr,
1399                &key_expr / &suffix,
1400            );
1401            let token = conf
1402                .session
1403                .liveliness()
1404                .declare_token(&key_expr / &suffix)
1405                .wait()?;
1406            zlock!(statesref).token = Some(token)
1407        }
1408
1409        let reliable_subscriber = AdvancedSubscriber {
1410            statesref,
1411            subscriber,
1412            receiver,
1413            liveliness_subscriber,
1414            heartbeat_subscriber,
1415        };
1416
1417        Ok(reliable_subscriber)
1418    }
1419
1420    /// Returns the [`EntityGlobalId`] of this AdvancedSubscriber.
1421    #[zenoh_macros::unstable]
1422    pub fn id(&self) -> EntityGlobalId {
1423        self.subscriber.id()
1424    }
1425
1426    /// Returns the [`KeyExpr`] this subscriber subscribes to.
1427    #[zenoh_macros::unstable]
1428    pub fn key_expr(&self) -> &KeyExpr<'static> {
1429        self.subscriber.key_expr()
1430    }
1431
1432    /// Returns a reference to this subscriber's handler.
1433    ///
1434    /// An handler is anything that implements [`zenoh::handlers::IntoHandler`].
1435    /// The default handler is [`zenoh::handlers::DefaultHandler`].
1436    #[zenoh_macros::unstable]
1437    pub fn handler(&self) -> &Handler {
1438        &self.receiver
1439    }
1440
1441    /// Returns a mutable reference to this subscriber's handler.
1442    ///
1443    /// An handler is anything that implements [`zenoh::handlers::IntoHandler`].
1444    /// The default handler is [`zenoh::handlers::DefaultHandler`].
1445    #[zenoh_macros::unstable]
1446    pub fn handler_mut(&mut self) -> &mut Handler {
1447        &mut self.receiver
1448    }
1449
1450    /// Declares a listener to detect missed samples.
1451    ///
1452    /// Missed samples can only be detected from [`AdvancedPublisher`](crate::AdvancedPublisher) that
1453    /// enable [`sample_miss_detection`](crate::AdvancedPublisherBuilder::sample_miss_detection).
1454    #[zenoh_macros::unstable]
1455    pub fn sample_miss_listener(&self) -> SampleMissListenerBuilder<'_, DefaultHandler> {
1456        SampleMissListenerBuilder {
1457            statesref: &self.statesref,
1458            handler: DefaultHandler::default(),
1459        }
1460    }
1461
1462    /// Declares a listener to detect matching publishers.
1463    ///
1464    /// Only [`AdvancedPublisher`](crate::AdvancedPublisher) that enable
1465    /// [`publisher_detection`](crate::AdvancedPublisherBuilder::publisher_detection) can be detected.
1466    #[zenoh_macros::unstable]
1467    pub fn detect_publishers(&self) -> LivelinessSubscriberBuilder<'_, '_, DefaultHandler> {
1468        self.subscriber
1469            .session()
1470            .liveliness()
1471            .declare_subscriber(self.subscriber.key_expr() / KE_ADV_PREFIX / KE_PUB / KE_STARSTAR)
1472    }
1473
1474    /// Undeclares this AdvancedSubscriber
1475    #[inline]
1476    #[zenoh_macros::unstable]
1477    pub fn undeclare(self) -> SubscriberUndeclaration<()> {
1478        tracing::debug!(
1479            "AdvancedSubscriber{{key_expr: {}}}: Undeclare",
1480            self.key_expr()
1481        );
1482        self.subscriber.undeclare()
1483    }
1484
1485    fn set_background_impl(&mut self, background: bool) {
1486        self.subscriber.set_background(background);
1487        if let Some(mut liveliness_sub) = self.liveliness_subscriber.take() {
1488            liveliness_sub.set_background(background);
1489        }
1490        if let Some(mut heartbeat_sub) = self.heartbeat_subscriber.take() {
1491            heartbeat_sub.set_background(background);
1492        }
1493    }
1494
1495    #[zenoh_macros::internal]
1496    pub fn set_background(&mut self, background: bool) {
1497        self.set_background_impl(background)
1498    }
1499}
1500
1501#[zenoh_macros::unstable]
1502#[inline]
1503fn flush_sequenced_source(
1504    state: &mut SourceState<WrappingSn>,
1505    callback: Option<&Callback<Sample>>,
1506    source_id: &EntityGlobalId,
1507    miss_handlers: &HashMap<usize, Callback<Miss>>,
1508) {
1509    let Some(callback) = callback else {
1510        return;
1511    };
1512    if state.pending_queries == 0 && !state.pending_samples.is_empty() {
1513        let mut pending_samples = BTreeMap::new();
1514        std::mem::swap(&mut state.pending_samples, &mut pending_samples);
1515        for (seq_num, sample) in pending_samples {
1516            match state.last_delivered {
1517                None => {
1518                    state.last_delivered = Some(seq_num);
1519                    callback.call(sample);
1520                }
1521                Some(last) if seq_num == last + 1 => {
1522                    state.last_delivered = Some(seq_num);
1523                    callback.call(sample);
1524                }
1525                Some(last) if seq_num > last + 1 => {
1526                    tracing::warn!(
1527                        "Sample missed: missed {} samples from {:?}.",
1528                        seq_num - last - 1,
1529                        source_id,
1530                    );
1531                    for miss_callback in miss_handlers.values() {
1532                        miss_callback.call(Miss {
1533                            source: *source_id,
1534                            nb: seq_num - last - 1,
1535                        })
1536                    }
1537                    state.last_delivered = Some(seq_num);
1538                    callback.call(sample);
1539                }
1540                _ => {
1541                    // duplicate
1542                }
1543            }
1544        }
1545    }
1546}
1547
1548#[zenoh_macros::unstable]
1549#[inline]
1550fn flush_timestamped_source(
1551    state: &mut SourceState<Timestamp>,
1552    callback: Option<&Callback<Sample>>,
1553) {
1554    let Some(callback) = callback else {
1555        return;
1556    };
1557    if state.pending_queries == 0 && !state.pending_samples.is_empty() {
1558        for (timestamp, sample) in std::mem::take(&mut state.pending_samples) {
1559            if state
1560                .last_delivered
1561                .map(|last| timestamp > last)
1562                .unwrap_or(true)
1563            {
1564                state.last_delivered = Some(timestamp);
1565                callback.call(sample);
1566            }
1567        }
1568    }
1569}
1570
1571#[zenoh_macros::unstable]
1572#[derive(Clone)]
1573struct InitialRepliesHandler {
1574    statesref: Arc<Mutex<State>>,
1575}
1576
1577#[zenoh_macros::unstable]
1578impl Drop for InitialRepliesHandler {
1579    fn drop(&mut self) {
1580        let states = &mut *zlock!(self.statesref);
1581        states.global_pending_queries = states.global_pending_queries.saturating_sub(1);
1582        tracing::trace!(
1583            "AdvancedSubscriber{{key_expr: {}}}: Flush initial replies",
1584            states.key_expr
1585        );
1586
1587        if states.global_pending_queries == 0 {
1588            for (source_id, state) in states.sequenced_states.iter_mut() {
1589                flush_sequenced_source(
1590                    state,
1591                    states.callback.as_ref(),
1592                    source_id,
1593                    &states.miss_handlers,
1594                );
1595                state.periodic_task =
1596                    spawn_periodic_queries(&self.statesref, states.period, *source_id);
1597            }
1598            for (_, state) in states.timestamped_states.iter_mut() {
1599                flush_timestamped_source(state, states.callback.as_ref());
1600            }
1601        }
1602    }
1603}
1604
1605#[zenoh_macros::unstable]
1606#[derive(Clone)]
1607struct SequencedRepliesHandler {
1608    source_id: EntityGlobalId,
1609    statesref: Arc<Mutex<State>>,
1610}
1611
1612#[zenoh_macros::unstable]
1613impl Drop for SequencedRepliesHandler {
1614    fn drop(&mut self) {
1615        let states = &mut *zlock!(self.statesref);
1616        // use peek_mut so query without samples do not prevent the state to be garbage collected
1617        if let Some(state) = states.sequenced_states.peek_mut(&self.source_id) {
1618            state.pending_queries = state.pending_queries.saturating_sub(1);
1619            if states.global_pending_queries == 0 {
1620                tracing::trace!(
1621                    "AdvancedSubscriber{{key_expr: {}}}: Flush sequenced samples",
1622                    states.key_expr
1623                );
1624                flush_sequenced_source(
1625                    state,
1626                    states.callback.as_ref(),
1627                    &self.source_id,
1628                    &states.miss_handlers,
1629                )
1630            }
1631        }
1632    }
1633}
1634
1635#[zenoh_macros::unstable]
1636#[derive(Clone)]
1637struct TimestampedRepliesHandler {
1638    id: ID,
1639    statesref: Arc<Mutex<State>>,
1640}
1641
1642#[zenoh_macros::unstable]
1643impl Drop for TimestampedRepliesHandler {
1644    fn drop(&mut self) {
1645        let states = &mut *zlock!(self.statesref);
1646        // use peek_mut so query without samples do not prevent the state to be garbage collected
1647        if let Some(state) = states.timestamped_states.peek_mut(&self.id) {
1648            state.pending_queries = state.pending_queries.saturating_sub(1);
1649            if states.global_pending_queries == 0 {
1650                tracing::trace!(
1651                    "AdvancedSubscriber{{key_expr: {}}}: Flush timestamped samples",
1652                    states.key_expr
1653                );
1654                flush_timestamped_source(state, states.callback.as_ref());
1655            }
1656        }
1657    }
1658}
1659
1660/// A struct that represent missed samples.
1661#[zenoh_macros::unstable]
1662#[derive(Debug, Clone)]
1663pub struct Miss {
1664    source: EntityGlobalId,
1665    nb: u32,
1666}
1667
1668impl Miss {
1669    /// The source of missed samples.
1670    pub fn source(&self) -> EntityGlobalId {
1671        self.source
1672    }
1673
1674    /// The number of missed samples.
1675    pub fn nb(&self) -> u32 {
1676        self.nb
1677    }
1678}
1679
1680impl CallbackParameter for Miss {
1681    type Message<'a> = Self;
1682
1683    fn from_message(msg: Self::Message<'_>) -> Self {
1684        msg
1685    }
1686}
1687
1688/// A listener to detect missed samples.
1689///
1690/// Missed samples can only be detected from [`AdvancedPublisher`](crate::AdvancedPublisher) that
1691/// enable [`sample_miss_detection`](crate::AdvancedPublisherBuilder::sample_miss_detection).
1692#[zenoh_macros::unstable]
1693pub struct SampleMissListener<Handler> {
1694    id: usize,
1695    statesref: Arc<Mutex<State>>,
1696    handler: Handler,
1697    undeclare_on_drop: bool,
1698}
1699
1700#[zenoh_macros::unstable]
1701impl<Handler> fmt::Debug for SampleMissListener<Handler> {
1702    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1703        f.debug_struct("SampleMissListener")
1704            .field("id", &self.id)
1705            .field("statesref", &"..")
1706            .field("handler", &"..")
1707            .field("undeclare_on_drop", &self.undeclare_on_drop)
1708            .finish()
1709    }
1710}
1711
1712#[zenoh_macros::unstable]
1713impl<Handler> SampleMissListener<Handler> {
1714    #[inline]
1715    pub fn undeclare(self) -> SampleMissHandlerUndeclaration<Handler>
1716    where
1717        Handler: Send,
1718    {
1719        SampleMissHandlerUndeclaration { listener: self }
1720    }
1721
1722    fn undeclare_impl(&mut self) -> ZResult<()> {
1723        // set the flag first to avoid double panic if this function panic
1724        self.undeclare_on_drop = false;
1725        zlock!(self.statesref).unregister_miss_callback(&self.id);
1726        Ok(())
1727    }
1728
1729    #[zenoh_macros::internal]
1730    pub fn set_background(&mut self, background: bool) {
1731        self.undeclare_on_drop = !background;
1732    }
1733}
1734
1735#[cfg(feature = "unstable")]
1736impl<Handler> Drop for SampleMissListener<Handler> {
1737    fn drop(&mut self) {
1738        if self.undeclare_on_drop {
1739            if let Err(error) = self.undeclare_impl() {
1740                tracing::error!(error);
1741            }
1742        }
1743    }
1744}
1745
1746// #[zenoh_macros::unstable]
1747// impl<Handler: Send> UndeclarableSealed<()> for SampleMissHandler<Handler> {
1748//     type Undeclaration = SampleMissHandlerUndeclaration<Handler>;
1749
1750//     fn undeclare_inner(self, _: ()) -> Self::Undeclaration {
1751//         SampleMissHandlerUndeclaration(self)
1752//     }
1753// }
1754
1755#[zenoh_macros::unstable]
1756impl<Handler> std::ops::Deref for SampleMissListener<Handler> {
1757    type Target = Handler;
1758
1759    fn deref(&self) -> &Self::Target {
1760        &self.handler
1761    }
1762}
1763#[zenoh_macros::unstable]
1764impl<Handler> std::ops::DerefMut for SampleMissListener<Handler> {
1765    fn deref_mut(&mut self) -> &mut Self::Target {
1766        &mut self.handler
1767    }
1768}
1769
1770/// A [`Resolvable`] returned by [`SampleMissListener::undeclare`]
1771#[zenoh_macros::unstable]
1772pub struct SampleMissHandlerUndeclaration<Handler> {
1773    listener: SampleMissListener<Handler>,
1774}
1775
1776#[zenoh_macros::unstable]
1777impl<Handler> fmt::Debug for SampleMissHandlerUndeclaration<Handler> {
1778    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1779        f.debug_tuple("SampleMissHandlerUndeclaration")
1780            .field(&self.listener)
1781            .finish()
1782    }
1783}
1784
1785impl<Handler> SampleMissHandlerUndeclaration<Handler> {
1786    /// Block in undeclare operation until all currently running instances of sample miss listener callback (if any) return.
1787    pub fn wait_callbacks(self) -> Self {
1788        // Note: no particular synchronization is required as of now since miss listener callbacks are always executed
1789        // under state lock
1790        self
1791    }
1792}
1793
1794#[zenoh_macros::unstable]
1795impl<Handler> Resolvable for SampleMissHandlerUndeclaration<Handler> {
1796    type To = ZResult<()>;
1797}
1798
1799#[zenoh_macros::unstable]
1800impl<Handler> Wait for SampleMissHandlerUndeclaration<Handler> {
1801    fn wait(mut self) -> <Self as Resolvable>::To {
1802        self.listener.undeclare_impl()
1803    }
1804}
1805
1806#[zenoh_macros::unstable]
1807impl<Handler> IntoFuture for SampleMissHandlerUndeclaration<Handler> {
1808    type Output = <Self as Resolvable>::To;
1809    type IntoFuture = Ready<<Self as Resolvable>::To>;
1810
1811    fn into_future(self) -> Self::IntoFuture {
1812        std::future::ready(self.wait())
1813    }
1814}
1815
1816/// A builder for initializing a [`SampleMissListener`].
1817#[zenoh_macros::unstable]
1818pub struct SampleMissListenerBuilder<'a, Handler, const BACKGROUND: bool = false> {
1819    statesref: &'a Arc<Mutex<State>>,
1820    handler: Handler,
1821}
1822
1823#[zenoh_macros::unstable]
1824impl<Handler, const BACKGROUND: bool> fmt::Debug
1825    for SampleMissListenerBuilder<'_, Handler, BACKGROUND>
1826{
1827    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1828        f.debug_struct("SampleMissListenerBuilder")
1829            .field("statesref", &"..")
1830            .field("handler", &"..")
1831            .field("background", &BACKGROUND)
1832            .finish()
1833    }
1834}
1835
1836#[zenoh_macros::unstable]
1837impl<'a> SampleMissListenerBuilder<'a, DefaultHandler> {
1838    /// Receive the sample miss notification with a callback.
1839    #[inline]
1840    #[zenoh_macros::unstable]
1841    pub fn callback<F>(self, callback: F) -> SampleMissListenerBuilder<'a, Callback<Miss>>
1842    where
1843        F: Fn(Miss) + Send + Sync + 'static,
1844    {
1845        self.with(Callback::from(callback))
1846    }
1847
1848    /// Receive the sample miss notification with a mutable callback.
1849    #[inline]
1850    #[zenoh_macros::unstable]
1851    pub fn callback_mut<F>(self, callback: F) -> SampleMissListenerBuilder<'a, Callback<Miss>>
1852    where
1853        F: FnMut(Miss) + Send + Sync + 'static,
1854    {
1855        self.callback(zenoh::handlers::locked(callback))
1856    }
1857
1858    /// Receive the sample miss notification with a [`Handler`](IntoHandler).
1859    #[inline]
1860    #[zenoh_macros::unstable]
1861    pub fn with<Handler>(self, handler: Handler) -> SampleMissListenerBuilder<'a, Handler>
1862    where
1863        Handler: IntoHandler<Miss>,
1864    {
1865        SampleMissListenerBuilder {
1866            statesref: self.statesref,
1867            handler,
1868        }
1869    }
1870}
1871
1872#[zenoh_macros::unstable]
1873impl<'a> SampleMissListenerBuilder<'a, Callback<Miss>> {
1874    /// Make the sample miss notification run in the background until the advanced subscriber is undeclared.
1875    ///
1876    /// Background builder doesn't return a `SampleMissHandler` object anymore.
1877    #[zenoh_macros::unstable]
1878    pub fn background(self) -> SampleMissListenerBuilder<'a, Callback<Miss>, true> {
1879        SampleMissListenerBuilder {
1880            statesref: self.statesref,
1881            handler: self.handler,
1882        }
1883    }
1884}
1885
1886#[zenoh_macros::unstable]
1887impl<Handler> Resolvable for SampleMissListenerBuilder<'_, Handler>
1888where
1889    Handler: IntoHandler<Miss> + Send,
1890    Handler::Handler: Send,
1891{
1892    type To = ZResult<SampleMissListener<Handler::Handler>>;
1893}
1894
1895#[zenoh_macros::unstable]
1896impl<Handler> Wait for SampleMissListenerBuilder<'_, Handler>
1897where
1898    Handler: IntoHandler<Miss> + Send,
1899    Handler::Handler: Send,
1900{
1901    #[zenoh_macros::unstable]
1902    fn wait(self) -> <Self as Resolvable>::To {
1903        let (callback, handler) = self.handler.into_handler();
1904        let id = zlock!(self.statesref).register_miss_callback(callback);
1905        Ok(SampleMissListener {
1906            id,
1907            statesref: self.statesref.clone(),
1908            handler,
1909            undeclare_on_drop: true,
1910        })
1911    }
1912}
1913
1914#[zenoh_macros::unstable]
1915impl<Handler> IntoFuture for SampleMissListenerBuilder<'_, Handler>
1916where
1917    Handler: IntoHandler<Miss> + Send,
1918    Handler::Handler: Send,
1919{
1920    type Output = <Self as Resolvable>::To;
1921    type IntoFuture = Ready<<Self as Resolvable>::To>;
1922
1923    #[zenoh_macros::unstable]
1924    fn into_future(self) -> Self::IntoFuture {
1925        std::future::ready(self.wait())
1926    }
1927}
1928
1929#[zenoh_macros::unstable]
1930impl Resolvable for SampleMissListenerBuilder<'_, Callback<Miss>, true> {
1931    type To = ZResult<()>;
1932}
1933
1934#[zenoh_macros::unstable]
1935impl Wait for SampleMissListenerBuilder<'_, Callback<Miss>, true> {
1936    #[zenoh_macros::unstable]
1937    fn wait(self) -> <Self as Resolvable>::To {
1938        let (callback, _) = self.handler.into_handler();
1939        zlock!(self.statesref).register_miss_callback(callback);
1940        Ok(())
1941    }
1942}
1943
1944#[zenoh_macros::unstable]
1945impl IntoFuture for SampleMissListenerBuilder<'_, Callback<Miss>, true> {
1946    type Output = <Self as Resolvable>::To;
1947    type IntoFuture = Ready<<Self as Resolvable>::To>;
1948
1949    #[zenoh_macros::unstable]
1950    fn into_future(self) -> Self::IntoFuture {
1951        std::future::ready(self.wait())
1952    }
1953}