Skip to main content

Monitor

Struct Monitor 

Source
pub struct Monitor { /* private fields */ }
Expand description

The wired monitor: a runtime-mutable watch set + liveliness + tick task feeding a core.

Lazy by construction (issue #84): start with empty spec.selectors declares no data-plane subscribers at all — only the zero-payload liveliness watches and the tick. Data flows only for what Monitor::watch was asked to observe, and Monitor::unwatch provably undeclares (an explicit, awaited undeclaration — not a dropped handle racing the network).

Implementations§

Source§

impl Monitor

Source

pub async fn start(session: &Session, spec: MonitorSpec) -> Result<Monitor>

Declare the spec’s subscribers on session and start watching. spec.selectors are simply the initial watches — [] is the lazy start.

Source

pub async fn watch(&self, selector: &str) -> Result<WatchId>

Observe a selector: declares a callback subscriber feeding the core.

The callback runs on zenoh’s network thread and does exactly what the old per-selector task did — one stats lock, one bounded broadcast send — so a slow UI still cannot exert backpressure into the network layer beyond the channel’s bound.

Source

pub async fn watching<S: AsRef<str>>( self, selectors: impl IntoIterator<Item = S>, ) -> Result<Monitor>

Declare selectors on this monitor, tearing it down — acknowledged — if any of them fails.

This is the judge windows’ opening move (#336): start, take the event stream, then declare what the window will observe. In that order, deliberately — a sample arriving between the subscriber’s declaration and the stream’s creation would be counted and not delivered, and these windows exist to say what they saw. But the ? on the declaration used to return with the monitor’s liveliness and tick tasks running and its subscribers left to Drop: the unacknowledged teardown shutdown exists to refuse, on the one path nobody thinks about.

Consuming and returning the monitor is what lets the failing path shutdown().await before it returns. The declaration error is the one reported — a teardown failure behind a failed declaration is noise — but the teardown itself is never skipped.

Source

pub async fn watch_seeded( &self, selector: &str, policy: SeedPolicy, ) -> Result<WatchId>

Observe a selector with a correct seed phase (issue #92; the RFC 04 §3.2 discipline of crate::seed_subscribe, run through this monitor’s bounded broadcast):

  • the subscriber is declared first, then both seed paths run as bounded GETs (@adv caches for live publishers, the selector itself for router storages — the crashed-producer case);
  • one per-key LWW merge spans seed and live samples until the boundary, so a transition in the seed window lands exactly once and a stale seed cannot regress a key. Suppressions are counted in the coverage, never silently absorbed (O6) — and they are not part of Dropped(n), which counts only broadcast lag;
  • FleetEvent::WatchSeeded fires once both paths resolve, carrying this id and the crate::SeedCoverage. After it, the merge is dropped and live samples flow untouched (the merge map is a seed-phase structure, not a per-watch leak).
Source

pub async fn unwatch(&self, id: WatchId) -> Result<()>

Stop observing: undeclares the subscriber (awaited to completion — the teardown is acknowledged, not racing a drop), then retires statistics for keys no remaining watch covers. Retired keys are counted (crate::model::stats::StatsTable::unwatched): a shrinking key set must never read as a quieting bus (RFC 09 §5.1 O6).

Source

pub async fn watched(&self) -> Vec<(WatchId, String)>

The active watch set.

Source

pub fn core(&self) -> &Arc<MonitorCore>

Source

pub fn events(&self) -> EventStream

Source

pub fn tree(&self) -> Arc<KeyTreeSnapshot>

Source

pub fn stop(self)

Stop watching. Equivalent to dropping the monitor — kept as an explicit verb for call sites that want to say so.

The teardown is the Drop one: tasks aborted, subscribers left to undeclare in the background. Where the acknowledgement matters — tearing one monitor down to declare another over the same keys — use shutdown instead.

Source

pub async fn shutdown(self) -> Result<()>

Stop watching, acknowledged: every watch undeclares and is waited for before this returns.

unwatch awaits undeclare on purpose — “the teardown is acknowledged, not racing a drop” — but the whole-monitor path had no such verb: Drop can only abort the tasks and let the subscribers undeclare on their own, in the background, which is the race that doc disavows. A frontend that re-scopes by rebuilding its monitor was therefore declaring the new subscribers while the old ones were still tearing down.

Every watch is drained even if one fails to undeclare — a monitor half torn down is worse than one torn down noisily — and the failures are reported together. Drop still runs afterwards, aborting the liveliness and tick tasks, and remains the fallback for every path that does not come through here.

Statistics are not retired the way unwatch retires them: that counter answers “the key set shrank because you stopped looking” (RFC 09 §5.1 O6) for a monitor that goes on running. This one is the end of the observation; the core goes with it unless a caller kept an Arc, and a re-scope’s next monitor starts from a fresh one.

Trait Implementations§

Source§

impl Debug for Monitor

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Drop for Monitor

Dropping a monitor stops it: the ingest tasks are aborted and the subscribers undeclare.

This is not a nicety. A JoinHandle merely detaches on drop, so without this impl every monitor that goes out of scope leaks a live subscriber and its ingest task for the lifetime of the session. zenctl never noticed — it calls Monitor::stop once and exits — but a GUI re-scopes its subscription whenever the user changes what they are watching, dropping and rebuilding the monitor each time.

Every task, which for one release meant every task but the seeded watches’ (#342): those handles live in watches, and aborting only self.tasks detached them. Each holds a cloned Session and goes on calling core.ingest/core.tick, so the re-scoping GUI above left one running per drop, each holding session teardown open for up to the seed timeout. unwatch and shutdown had aborted them all along; the async mutex is get_mut here, which needs no lock because Drop holds &mut self.

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<Source> AccessAs for Source

Source§

fn ref_as<T>(&self) -> <Source as IGuardRef<T>>::Guard<'_>
where Source: IGuardRef<T>, T: ?Sized,

Provides immutable access to a type as if it were its ABI-unstable equivalent.
Source§

fn mut_as<T>(&mut self) -> <Source as IGuardMut<T>>::GuardMut<'_>
where Source: IGuardMut<T>, T: ?Sized,

Provides mutable access to a type as if it were its ABI-unstable equivalent.
Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> AsNode<T> for T

Source§

fn as_node(&self) -> &T

Source§

impl<T> AsNodeMut<T> for T

Source§

fn as_node_mut(&mut self) -> &mut T

Source§

impl<'a, T, E> AsTaggedExplicit<'a, E> for T
where T: 'a,

Source§

fn explicit(self, class: Class, tag: u32) -> TaggedParser<'a, Explicit, Self, E>

Source§

impl<'a, T, E> AsTaggedImplicit<'a, E> for T
where T: 'a,

Source§

fn implicit( self, class: Class, constructed: bool, tag: u32, ) -> TaggedParser<'a, Implicit, Self, E>

Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, As> IGuardMut<As> for T
where T: Into<As>, As: Into<T>,

Source§

type GuardMut<'a> = MutAs<'a, T, As> where T: 'a

The type of the guard which will clean up the temporary after applying its changes to the original.
Source§

fn guard_mut_inner(&mut self) -> <T as IGuardMut<As>>::GuardMut<'_>

Construct the temporary and guard it through a mutable reference.
Source§

impl<T, As> IGuardRef<As> for T
where T: Into<As>, As: Into<T>,

Source§

type Guard<'a> = RefAs<'a, T, As> where T: 'a

The type of the guard which will clean up the temporary.
Source§

fn guard_ref_inner(&self) -> <T as IGuardRef<As>>::Guard<'_>

Construct the temporary and guard it through an immutable reference.
Source§

impl<T> Includes<End> for T

Source§

type Output = End

The result
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self>

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self>

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self>

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self>
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self>
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self>

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more