pub struct Updates { /* private fields */ }Expand description
The stream of one subscription’s events.
Returned by Client::subscribe. It implements
futures_util::Stream, so the usual combinators work and the plain shape
is a while let loop:
use futures_util::StreamExt;
use lightstreamer_rs::SubscriptionEvent;
while let Some(event) = updates.next().await {
if let SubscriptionEvent::Update(update) = event {
println!("{}: {:?}", update.item_name(), update.changed_fields());
}
}§Dropping this stream unsubscribes
There is no deregistration dance: drop the Updates and the client sends
the unsubscription for you. The request goes out asynchronously — the drop
itself does not block and cannot fail — so a UNSUB may still be in flight
when the value is gone. If you need to observe the unsubscription
completing, call Client::unsubscribe
instead and keep reading until SubscriptionEvent::Unsubscribed.
§Backpressure: bounded, and nothing is dropped
The stream is fed by a bounded channel whose capacity is
ConnectionOptions::with_update_capacity
— 1024 events by default. When it fills, the client blocks rather than
discarding anything, and that backpressure propagates through the session
down to the socket. Concretely, a consumer that stops polling this stream
will, once the buffer fills, stall the whole client: other
subscriptions and the session’s own liveness included.
That is a deliberate choice, and the alternatives are worse. Dropping
updates locally would corrupt the notification count that makes recovery
after a broken connection correct
[docs/spec/02-session-lifecycle.md §5.2], and it would hide loss the
protocol goes out of its way to report — TLCP has its own overflow
notification precisely so that a client is told when data was discarded
[docs/spec/04-notifications.md §3.7].
§What a long stall actually does
“Stalls the client” is literal, and it has a consequence worth planning
for. While the client is blocked on a full stream it is not reading from
the socket either, so its liveness timers do not advance. If you stop
consuming for longer than the keepalive interval the server negotiated —
reported to you as Connected::keepalive,
plus ConnectionOptions::with_keepalive_slack
— then the moment you start reading again the client will conclude the
connection went silent, tear it down, and recover
[docs/spec/02-session-lifecycle.md §8.1]. You will see a
SessionEvent::Disconnected carrying
DisconnectReason::Stalled, followed
by a reconnection.
Nothing is lost when that happens — recovery is what the notification count exists for — but it is real work, and it is avoidable. A consumer that may pause for seconds at a time should either buffer on its own side or, better, ask the server to send less.
If you cannot keep up, say so to the server, which is where the protocol
puts that decision: cap the rate with
Subscription::with_max_frequency,
or bound what the server buffers with
Subscription::with_buffer_size.
The server will then merge or drop on its side, and tell you it did with a
SubscriptionEvent::Overflow.
§The one exception: shutdown
“Nothing is dropped” holds while the client is running. It does not
survive an ordered shutdown, and it deliberately does not:
Client::disconnect and dropping the
Client are signalled out of band, every blocked delivery
gives way to them, and whatever had not yet been delivered is discarded
along with this stream.
Without that exception a consumer that walked away could make the client
impossible to stop — no disconnect, no drop, and a task and a socket alive
until the process ended. Nothing observable is lost by it: a caller still
reading frees room and the delivery completes, and a caller that has asked
to disconnect has said it wants no more. See
docs/adr/0003-typed-event-stream-as-delivery-surface.md.
Implementations§
Source§impl Updates
impl Updates
Sourcepub const fn id(&self) -> SubscriptionId
pub const fn id(&self) -> SubscriptionId
The handle identifying this subscription.
Stable for the life of the stream, including across a session that was
lost and replaced. Use it to correlate with
SessionEvent::Resubscribed.
Trait Implementations§
Source§impl Drop for Updates
impl Drop for Updates
Source§impl Stream for Updates
impl Stream for Updates
Source§type Item = SubscriptionEvent
type Item = SubscriptionEvent
Auto Trait Implementations§
impl Freeze for Updates
impl RefUnwindSafe for Updates
impl Send for Updates
impl Sync for Updates
impl Unpin for Updates
impl UnsafeUnpin for Updates
impl UnwindSafe for Updates
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> StreamExt for T
impl<T> StreamExt for T
Source§fn next(&mut self) -> Next<'_, Self>where
Self: Unpin,
fn next(&mut self) -> Next<'_, Self>where
Self: Unpin,
Source§fn into_future(self) -> StreamFuture<Self>
fn into_future(self) -> StreamFuture<Self>
Source§fn map<T, F>(self, f: F) -> Map<Self, F>
fn map<T, F>(self, f: F) -> Map<Self, F>
Source§fn enumerate(self) -> Enumerate<Self>where
Self: Sized,
fn enumerate(self) -> Enumerate<Self>where
Self: Sized,
Source§fn filter<Fut, F>(self, f: F) -> Filter<Self, Fut, F>
fn filter<Fut, F>(self, f: F) -> Filter<Self, Fut, F>
Source§fn filter_map<Fut, T, F>(self, f: F) -> FilterMap<Self, Fut, F>
fn filter_map<Fut, T, F>(self, f: F) -> FilterMap<Self, Fut, F>
Source§fn then<Fut, F>(self, f: F) -> Then<Self, Fut, F>
fn then<Fut, F>(self, f: F) -> Then<Self, Fut, F>
Source§fn collect<C>(self) -> Collect<Self, C>
fn collect<C>(self) -> Collect<Self, C>
Source§fn unzip<A, B, FromA, FromB>(self) -> Unzip<Self, FromA, FromB>
fn unzip<A, B, FromA, FromB>(self) -> Unzip<Self, FromA, FromB>
Source§fn concat(self) -> Concat<Self>
fn concat(self) -> Concat<Self>
Source§fn count(self) -> Count<Self>where
Self: Sized,
fn count(self) -> Count<Self>where
Self: Sized,
Source§fn fold<T, Fut, F>(self, init: T, f: F) -> Fold<Self, Fut, T, F>
fn fold<T, Fut, F>(self, init: T, f: F) -> Fold<Self, Fut, T, F>
Source§fn any<Fut, F>(self, f: F) -> Any<Self, Fut, F>
fn any<Fut, F>(self, f: F) -> Any<Self, Fut, F>
true if any element in stream satisfied a predicate. Read moreSource§fn all<Fut, F>(self, f: F) -> All<Self, Fut, F>
fn all<Fut, F>(self, f: F) -> All<Self, Fut, F>
true if all element in stream satisfied a predicate. Read moreSource§fn flatten(self) -> Flatten<Self>
fn flatten(self) -> Flatten<Self>
Source§fn flatten_unordered(
self,
limit: impl Into<Option<usize>>,
) -> FlattenUnorderedWithFlowController<Self, ()>
fn flatten_unordered( self, limit: impl Into<Option<usize>>, ) -> FlattenUnorderedWithFlowController<Self, ()>
Source§fn flat_map_unordered<U, F>(
self,
limit: impl Into<Option<usize>>,
f: F,
) -> FlatMapUnordered<Self, U, F>
fn flat_map_unordered<U, F>( self, limit: impl Into<Option<usize>>, f: F, ) -> FlatMapUnordered<Self, U, F>
StreamExt::map but flattens nested Streams
and polls them concurrently, yielding items in any order, as they made
available. Read moreSource§fn scan<S, B, Fut, F>(self, initial_state: S, f: F) -> Scan<Self, S, Fut, F>
fn scan<S, B, Fut, F>(self, initial_state: S, f: F) -> Scan<Self, S, Fut, F>
StreamExt::fold that holds internal state
and produces a new stream. Read moreSource§fn skip_while<Fut, F>(self, f: F) -> SkipWhile<Self, Fut, F>
fn skip_while<Fut, F>(self, f: F) -> SkipWhile<Self, Fut, F>
true. Read moreSource§fn take_while<Fut, F>(self, f: F) -> TakeWhile<Self, Fut, F>
fn take_while<Fut, F>(self, f: F) -> TakeWhile<Self, Fut, F>
true. Read moreSource§fn take_until<Fut>(self, fut: Fut) -> TakeUntil<Self, Fut>
fn take_until<Fut>(self, fut: Fut) -> TakeUntil<Self, Fut>
Source§fn for_each<Fut, F>(self, f: F) -> ForEach<Self, Fut, F>
fn for_each<Fut, F>(self, f: F) -> ForEach<Self, Fut, F>
Source§fn for_each_concurrent<Fut, F>(
self,
limit: impl Into<Option<usize>>,
f: F,
) -> ForEachConcurrent<Self, Fut, F>
fn for_each_concurrent<Fut, F>( self, limit: impl Into<Option<usize>>, f: F, ) -> ForEachConcurrent<Self, Fut, F>
Source§fn take(self, n: usize) -> Take<Self>where
Self: Sized,
fn take(self, n: usize) -> Take<Self>where
Self: Sized,
n items of the underlying stream. Read moreSource§fn skip(self, n: usize) -> Skip<Self>where
Self: Sized,
fn skip(self, n: usize) -> Skip<Self>where
Self: Sized,
n items of the underlying stream. Read moreSource§fn catch_unwind(self) -> CatchUnwind<Self>where
Self: Sized + UnwindSafe,
fn catch_unwind(self) -> CatchUnwind<Self>where
Self: Sized + UnwindSafe,
Source§fn boxed<'a>(self) -> Pin<Box<dyn Stream<Item = Self::Item> + Send + 'a>>
fn boxed<'a>(self) -> Pin<Box<dyn Stream<Item = Self::Item> + Send + 'a>>
Source§fn boxed_local<'a>(self) -> Pin<Box<dyn Stream<Item = Self::Item> + 'a>>where
Self: Sized + 'a,
fn boxed_local<'a>(self) -> Pin<Box<dyn Stream<Item = Self::Item> + 'a>>where
Self: Sized + 'a,
Source§fn buffered(self, n: usize) -> Buffered<Self>
fn buffered(self, n: usize) -> Buffered<Self>
Source§fn buffer_unordered(self, n: usize) -> BufferUnordered<Self>
fn buffer_unordered(self, n: usize) -> BufferUnordered<Self>
Source§fn zip<St>(self, other: St) -> Zip<Self, St>
fn zip<St>(self, other: St) -> Zip<Self, St>
Source§fn peekable(self) -> Peekable<Self>where
Self: Sized,
fn peekable(self) -> Peekable<Self>where
Self: Sized,
peek method. Read moreSource§fn chunks(self, capacity: usize) -> Chunks<Self>where
Self: Sized,
fn chunks(self, capacity: usize) -> Chunks<Self>where
Self: Sized,
Source§fn ready_chunks(self, capacity: usize) -> ReadyChunks<Self>where
Self: Sized,
fn ready_chunks(self, capacity: usize) -> ReadyChunks<Self>where
Self: Sized,
Source§fn forward<S>(self, sink: S) -> Forward<Self, S>
fn forward<S>(self, sink: S) -> Forward<Self, S>
Source§fn split<Item>(self) -> (SplitSink<Self, Item>, SplitStream<Self>)
fn split<Item>(self) -> (SplitSink<Self, Item>, SplitStream<Self>)
Source§fn inspect<F>(self, f: F) -> Inspect<Self, F>
fn inspect<F>(self, f: F) -> Inspect<Self, F>
Source§fn left_stream<B>(self) -> Either<Self, B>
fn left_stream<B>(self) -> Either<Self, B>
Source§fn right_stream<B>(self) -> Either<B, Self>
fn right_stream<B>(self) -> Either<B, Self>
Source§fn poll_next_unpin(&mut self, cx: &mut Context<'_>) -> Poll<Option<Self::Item>>where
Self: Unpin,
fn poll_next_unpin(&mut self, cx: &mut Context<'_>) -> Poll<Option<Self::Item>>where
Self: Unpin,
Stream::poll_next on Unpin
stream types.