1use 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#[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 #[inline]
73 #[zenoh_macros::unstable]
74 pub fn detect_late_publishers(mut self) -> Self {
75 self.liveliness = true;
76 self
77 }
78
79 #[zenoh_macros::unstable]
83 pub fn max_samples(mut self, depth: usize) -> Self {
84 self.max_samples = Some(depth);
85 self
86 }
87
88 #[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#[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 #[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 #[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 #[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#[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 #[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 #[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 #[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 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 #[zenoh_macros::unstable]
291 #[inline]
292 pub fn allowed_origin(mut self, origin: Locality) -> Self {
293 self.origin = origin;
294 self
295 }
296
297 #[zenoh_macros::unstable]
303 #[inline]
304 pub fn recovery(mut self, conf: RecoveryConfig) -> Self {
305 self.retransmission = Some(conf);
306 self
307 }
308
309 #[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 #[zenoh_macros::unstable]
329 #[inline]
330 pub fn history(mut self, config: HistoryConfig) -> Self {
331 self.history = Some(config);
332 self
333 }
334
335 #[zenoh_macros::unstable]
337 pub fn subscriber_detection(mut self) -> Self {
338 self.liveliness = true;
339 self
340 }
341
342 #[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: 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,
478 periodic_task: Option<AbortOnDropHandle<()>>,
480 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#[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 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; 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 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
791async fn gc_task(statesref: Weak<Mutex<State>>, retention_period: Duration) {
795 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 } 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 tokio::time::sleep(retention_period).await;
816 loop {
817 let oldest_access = {
818 let Some(states) = statesref.upgrade() else {
819 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 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 .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 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 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 .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 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 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 _ => 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 #[zenoh_macros::unstable]
1422 pub fn id(&self) -> EntityGlobalId {
1423 self.subscriber.id()
1424 }
1425
1426 #[zenoh_macros::unstable]
1428 pub fn key_expr(&self) -> &KeyExpr<'static> {
1429 self.subscriber.key_expr()
1430 }
1431
1432 #[zenoh_macros::unstable]
1437 pub fn handler(&self) -> &Handler {
1438 &self.receiver
1439 }
1440
1441 #[zenoh_macros::unstable]
1446 pub fn handler_mut(&mut self) -> &mut Handler {
1447 &mut self.receiver
1448 }
1449
1450 #[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 #[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 #[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 }
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 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 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#[zenoh_macros::unstable]
1662#[derive(Debug, Clone)]
1663pub struct Miss {
1664 source: EntityGlobalId,
1665 nb: u32,
1666}
1667
1668impl Miss {
1669 pub fn source(&self) -> EntityGlobalId {
1671 self.source
1672 }
1673
1674 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#[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 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]
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#[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 pub fn wait_callbacks(self) -> Self {
1788 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#[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 #[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 #[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 #[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 #[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}