Class MassiveStockStream
- Namespace
- MassiveDotNet.WebSocket
- Assembly
- MassiveDotNet.WebSocket.dll
An authenticated stock stream. Each topic is its own sequence.
[SuppressMessage("Naming", "CA1711:Identifiers should not have incorrect suffix", Justification = "\"Stream\" is this SDK's own domain term for a live market data feed (MassiveStreamClient, MassiveStreamOptions), not System.IO.Stream; the name is deliberate and consistent, not an accidental collision with the BCL suffix.")]
public sealed class MassiveStockStream : IAsyncDisposable
- Inheritance
-
MassiveStockStream
- Implements
- Inherited Members
- Extension Methods
Properties
EvictionCount
How many times the server has evicted this connection to make room for another on the same key.
public int EvictionCount { get; }
Property Value
Remarks
Eviction is the only documented disconnect cause that announces itself, and without this it is the one a consumer cannot tell apart from a slow-consumer close or a network drop. Read it from a Reconnected handler: reconnect is on by default, so an eviction normally ends in a reconnect rather than a stop, and Faulted -- which carries MassiveStreamEvictedException -- fires only for a stream configured not to reconnect (D42).
A rising count while ReconnectCount rises with it means something else is using this key: on stocks the older connection is the one closed, so two processes that both reconnect will keep displacing each other indefinitely. The SDK does not intervene -- the remedy is a second connection on the plan, or one fewer process, and neither is the SDK's call to make.
LastEvictionMessage
What the server said the last time it evicted this connection, verbatim, or null if it never has.
public string? LastEvictionMessage { get; }
Property Value
Remarks
Reported rather than categorised, as ServerMessage is (D35). Unlike the pending cause behind MassiveStreamEvictedException, this is cumulative and survives every reconnect: it answers "has this stream ever been evicted", which is a different question from whether any particular drop was one.
LastReconnected
When the connection was last re-established.
public Instant? LastReconnected { get; }
Property Value
ReconnectCount
How many times the underlying connection has been re-established.
public int ReconnectCount { get; }
Property Value
Methods
DisposeAsync()
Closes the stream and ends every topic sequence.
public ValueTask DisposeAsync()
Returns
SubscribeImbalancesAsync(IReadOnlyCollection<string>, CancellationToken)
Subscribes to net order imbalance auction events.
public Task<MassiveTopicSubscription<StockImbalance>> SubscribeImbalancesAsync(IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
tickersIReadOnlyCollection<string>Symbols, or
*for every symbol.cancellationTokenCancellationTokenCancels the subscribe.
Returns
- Task<MassiveTopicSubscription<StockImbalance>>
This stream's imbalance sequence, on the same terms as the other topics.
Exceptions
- ObjectDisposedException
The stream has been disposed.
- MassiveStreamSubscriptionException
The server acknowledged fewer subscriptions than were requested -- including every pair when the account's plan does not cover this topic, which the server refuses rather than ignores (issue #60).
SubscribeLimitUpLimitDownAsync(IReadOnlyCollection<string>, CancellationToken)
Subscribes to limit up-limit down price band events.
public Task<MassiveTopicSubscription<StockLimitUpLimitDown>> SubscribeLimitUpLimitDownAsync(IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
tickersIReadOnlyCollection<string>Symbols, or
*for every symbol.cancellationTokenCancellationTokenCancels the subscribe.
Returns
- Task<MassiveTopicSubscription<StockLimitUpLimitDown>>
This stream's band sequence, on the same terms as the other topics.
Exceptions
- ObjectDisposedException
The stream has been disposed.
- MassiveStreamSubscriptionException
The server acknowledged fewer subscriptions than were requested.
SubscribeMinuteAggregatesAsync(IReadOnlyCollection<string>, CancellationToken)
Subscribes to minute-by-minute aggregate bars.
public Task<MassiveTopicSubscription<StockAggregate>> SubscribeMinuteAggregatesAsync(IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
tickersIReadOnlyCollection<string>Symbols, or
*for every symbol.cancellationTokenCancellationTokenCancels the subscribe.
Returns
- Task<MassiveTopicSubscription<StockAggregate>>
This stream's minute-bar sequence, separate from the second-bar sequence even though both carry StockAggregate: sinks are keyed by wire code, not by type.
Exceptions
- ObjectDisposedException
The stream has been disposed.
- MassiveStreamSubscriptionException
The server acknowledged fewer subscriptions than were requested.
SubscribeQuotesAsync(IReadOnlyCollection<string>, CancellationToken)
Subscribes to NBBO quotes.
public Task<MassiveTopicSubscription<StockQuote>> SubscribeQuotesAsync(IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
tickersIReadOnlyCollection<string>Symbols, or
*for every symbol.cancellationTokenCancellationTokenCancels the subscribe.
Returns
- Task<MassiveTopicSubscription<StockQuote>>
This stream's quote sequence, on the same terms as trades.
Exceptions
- ObjectDisposedException
The stream has been disposed.
- MassiveStreamSubscriptionException
The server acknowledged fewer subscriptions than were requested.
SubscribeSecondAggregatesAsync(IReadOnlyCollection<string>, CancellationToken)
Subscribes to second-by-second aggregate bars.
public Task<MassiveTopicSubscription<StockAggregate>> SubscribeSecondAggregatesAsync(IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
tickersIReadOnlyCollection<string>Symbols, or
*for every symbol.cancellationTokenCancellationTokenCancels the subscribe.
Returns
- Task<MassiveTopicSubscription<StockAggregate>>
This stream's second-bar sequence. Calling again widens the ticker set and returns the same sequence, so a topic has one buffer and one consumer however many times it is called.
Exceptions
- ObjectDisposedException
The stream has been disposed.
- MassiveStreamSubscriptionException
The server acknowledged fewer subscriptions than were requested.
SubscribeTradesAsync(IReadOnlyCollection<string>, CancellationToken)
Subscribes to tick-level trades.
public Task<MassiveTopicSubscription<StockTrade>> SubscribeTradesAsync(IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
tickersIReadOnlyCollection<string>Symbols, or
*for every symbol.cancellationTokenCancellationTokenCancels the subscribe.
Returns
- Task<MassiveTopicSubscription<StockTrade>>
This stream's trade sequence. Calling again widens the ticker set and returns the same sequence, so a topic has one buffer and one consumer however many times it is called.
Exceptions
- ObjectDisposedException
The stream has been disposed.
- MassiveStreamSubscriptionException
The server acknowledged fewer subscriptions than were requested.
UnsubscribeAsync(StockTopic, IReadOnlyCollection<string>, CancellationToken)
Stops receiving a topic for the given symbols.
public Task UnsubscribeAsync(StockTopic topic, IReadOnlyCollection<string> tickers, CancellationToken cancellationToken = default)
Parameters
topicStockTopicThe topic.
tickersIReadOnlyCollection<string>The symbols to drop.
cancellationTokenCancellationTokenCancels the unsubscribe.
Returns
- Task
A task completing once the message has been sent.
Exceptions
- ObjectDisposedException
The stream has been disposed.
Events
DropObserved
Raised when a topic buffer overflowed and dropped an event, naming the topic's wire code and that topic's own running drop count, throttled to at most once a second so a sustained overflow does not produce an unbounded stream of notifications.
public event Action<string, long>? DropObserved
Event Type
Remarks
F7 (Task 12 review round 1): the count travels IN the event rather than requiring a caller
to separately hold every subscription and read its own
DroppedCount -- the original signature took no
topic at all, which left no correct way for the DI package's LogStreamHealth to
report on more than one topic (see that method's own remarks). Every handler is invoked
with its own try/catch (F1): a notification that something was dropped must never itself
take down the whole live feed -- exactly the defect Task 11 fixed for
MassiveStreamConnection.Faulted, twenty lines from this one.
The one-a-second throttle is shared across every topic on this stream, not scoped per topic: a sustained overflow on one topic can therefore delay, or entirely suppress, another topic's own report through this event. Nothing is lost to a consumer who reads DroppedCount instead -- that counter stays exact per topic no matter what this event manages to raise.
Faulted
Raised once the stream has stopped for good and will not reconnect.
public event Action<Exception>? Faulted
Event Type
Remarks
Three terminal stops reach this: reconnect disabled, an authentication failure surfacing
during a reconnect attempt, and a parse failure that is never reconnectable. Each raises
this, with the exception that ended the stream, before completing every subscription's
sequence -- so a consumer whose await foreach ends can learn why, rather than merely
observing that it did.
MalformedObserved
Raised when the wire sent a value a converter would not accept and the event was dropped, naming the topic's wire code, that topic's own running malformed count, and the exception saying which field and which value — throttled to at most once a second so a sustained off-schema field does not produce an unbounded stream of notifications.
public event Action<string, long, JsonException>? MalformedObserved
Event Type
Remarks
Separate from DropObserved on purpose (D-W20), and not a duplicate of it. A drop is the consumer's own backpressure, fixed by raising TopicBufferCapacity or doing less work in the loop; this is Massive's wire disagreeing with the SDK's schema, which the consumer cannot fix at all. Reporting both through one signal would hand them advice that cannot work.
The exception travels with the raise because it is the only route the refused value has out of the SDK. A malformed event was terminal until issue #65, so its message reached a consumer through Faulted; now that the connection survives, nothing else carries it. Every handler is invoked with its own try/catch (F1): a notification about one dropped event must never itself end the live feed.
The one-a-second throttle is shared across every topic on this stream, not scoped per topic (D-W20): a sustained burst of malformed events on one topic can therefore delay, or entirely suppress, another topic's own report through this event. Nothing is lost to a consumer who reads MalformedCount instead -- that counter stays exact per topic no matter what this event manages to raise.
Reconnected
Raised after a reconnect, carrying the running count.
public event Action<int>? Reconnected
Event Type
Remarks
A reconnect means messages were missed; the protocol offers no way to recover them.
SubscriptionsLost
Raised when a reconnect re-sent this stream's subscriptions and the server did not acknowledge all of them, carrying the pairs that went unacknowledged.
public event Action<MassiveStreamSubscriptionException>? SubscriptionsLost
Event Type
Remarks
Degraded, not terminal, and distinct from Faulted for that reason: the socket is up and every acknowledged topic is still delivering, so no sequence ends. It exists because the server answers an unrecognised topic with silence rather than an error (D33), so a replay nobody checks can leave a topic unsubscribed on a connection Reconnected has already reported as healthy -- a stream that looks alive and delivers nothing, indistinguishable from a quiet market.
The pairs stay in the subscription registry, so the next reconnect replays them again. A consumer that wants them back sooner can re-subscribe, which is acknowledgement-counted the ordinary way and throws if the server ignores it again.