Class MassiveTopicSubscription<T>
- Namespace
- MassiveDotNet.WebSocket
- Assembly
- MassiveDotNet.WebSocket.dll
A live subscription to one topic: the events, and how many were dropped.
public sealed class MassiveTopicSubscription<T> : IAsyncEnumerable<T>
Type Parameters
TThe event type.
- Inheritance
-
MassiveTopicSubscription<T>
- Implements
- Inherited Members
Remarks
One sequence per topic, and one consumer AT A TIME. Enumerating while another loop is already running throws rather than letting two loops silently steal events from each other (D-W3) -- but the claim is about concurrency, not about the sequence's lifetime, so an enumeration that ENDS releases the sequence for the next one.
Properties
DroppedCount
How many events were discarded because this topic's buffer was full when they arrived.
public long DroppedCount { get; }
Property Value
Remarks
Monotonic, and exact rather than estimated. A non-zero value means the consumer is slower than the feed: raise TopicBufferCapacity, or do less work in the loop.
MalformedCount
How many events were discarded because the wire sent a value this topic's converter would not accept.
public long MalformedCount { get; }
Property Value
Remarks
Monotonic, and exact rather than estimated. Deliberately separate from DroppedCount, because the two have opposite remedies: a drop is the consumer's own backpressure and is fixed by raising TopicBufferCapacity or doing less work in the loop, while this is Massive's wire disagreeing with the SDK's schema and there is nothing the consumer can do about it at all. One counter for both would hand a consumer advice that cannot work (D-W20). Subscribe to MalformedObserved for the exception naming which field and which value.
Methods
GetAsyncEnumerator(CancellationToken)
Enumerates this topic's events. One consumer at a time.
public IAsyncEnumerator<T> GetAsyncEnumerator(CancellationToken cancellationToken = default)
Parameters
cancellationTokenCancellationTokenEnds the enumeration.
Returns
- IAsyncEnumerator<T>
The enumerator.
Remarks
The guard is released when an enumeration ends -- by cancellation, by break, by an
exception, or by the sequence completing -- so the topic can be consumed again afterwards.
It used to be a one-way latch, which made a cancelled await foreach permanently burn
the topic's only sequence: re-subscribing hands back this same object (D-W3 gives a topic
one buffer however many times it is subscribed), so the next consumer got an
InvalidOperationException saying the sequence was "already being enumerated"
when nobody was, with no recovery short of tearing down the whole stream. Cancelling a
consumer loop is ordinary -- a timeout, a shutdown, a caller taking a break -- and D-W3's
argument is about two loops stealing from each other concurrently, which this still refuses.
Exceptions
- InvalidOperationException
Another loop is enumerating this sequence right now.