Table of Contents

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

T

The 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

long

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

long

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

cancellationToken CancellationToken

Ends 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.