WsDataStream class
Data stream over a hub.
Samples use the binary framing: data streams are the latency-critical path and their shape is fixed by the negotiated config, so JSON would be pure overhead on every sample.
- Inheritance
-
- Object
- NetworkStream<
DataStreamConfig, IMessage< IMessageType> > - DataStream<
DataStreamConfig, IMessage< IMessageType> > - WsDataStream
- Mixed-in types
Constructors
- WsDataStream({required DataStreamConfig config, required WsConnection connection, required String sessionName, required Node streamNode, PeerClockOffsets? clockOffsets})
Properties
- carriesSenderClock → bool
-
Whether this transport already puts the sender's clock on the wire.
no setterinherited
- channelCount → int
-
Number of channels in the stream.
no setterinherited
- clockOffsets → PeerClockOffsets?
-
Per-peer clock offsets, estimated on the coordination stream and shared
across every stream on this socket. Null when nothing is estimating them,
in which case samples arrive with an unknown offset.
final
-
clockSyncs
→ Stream<
ClockSyncSample> -
The estimates behind MessageTiming.clockOffset, for the peers this
stream receives from. They are made once per peer for the whole
transport, so every stream with an inlet on a peer reports the same.
no setteroverride
- config → DataStreamConfig
-
Configuration for the stream.
This is a NetworkStreamConfig object.
finalinherited
- connection → WsConnection
-
final
-
consumers
→ List<
String> -
List of consumer node IDs.
no setterinherited
- created → bool
-
Indicates whether the ILifecycle implementation has been created.
no setterinherited
- dataType → StreamDataType
-
Data type of the stream.
This is a StreamDataType enum value.
no setterinherited
- description → String?
-
Returns a description of the identity.
no setteroverride
- descriptor → PeerDescriptor
-
no setterinherited
- disposed → bool
-
Indicates whether the ILifecycle implementation has been disposed.
no setterinherited
- endpointId → String
-
no setterinherited
- hasConsumers → bool
-
Whether the stream has any current consumers.
no setterinherited
- hashCode → int
-
The hash code for this object.
no setterinherited
- hasProducers → bool
-
Whether the stream has any current producers.
no setterinherited
- id → String
-
Identifier for the stream, derived from the config hash code.
no setterinherited
-
inbox
→ Stream<
IMessage< IMessageType> > -
Incoming messages from peers.
no setterinherited
-
inletHealth
→ Stream<
InletHealth> -
Emits when this stream stops, or resumes, receiving from one peer for a
reason of its own rather than the peer's: see InletHealth.
no setterinherited
- manager → IResourceManager?
-
gets the resource manager that manages this resource
no setterinherited
- messageClass → Type
-
no setterinherited
- name → String
-
Human-readable name for the stream.
no setterinherited
-
outbox
→ StreamSink<
IMessage< IMessageType> > -
no setterinherited
-
outletConsumerPresence
→ Stream<
bool> -
Emits when this stream's publishing endpoint gains or loses every
subscriber, where the transport can tell.
no setterinherited
- paused → bool
-
no setterinherited
-
producers
→ List<
String> -
List of producer node IDs.
no setterinherited
- relaysSourceClock → bool
-
Whether this transport can pass a sample on with the timestamp its
origin gave it (sendDataAt, upstream).
no setterinherited
- runtimeType → Type
-
A representation of the runtime type of the object.
no setterinherited
- sampleRate → double
-
Sample rate of the stream.
no setterinherited
- sessionName → String
-
final
- shadowUId ↔ String?
-
getter/setter pairinherited
-
Whether both ends of this stream read the same PeerClock.
no setterinherited
- started → bool
-
Whether the stream is currently running.
no setterinherited
- streamNode → Node
-
no setteroverride
- uId → String
-
Returns a unique identifier that is guaranteed to be globally unique.
no setterinherited
- upstream ↔ ClockChain
-
How sendDataAt's source clock maps onto this node's PeerClock: the
hops the stream crossed to get here.
getter/setter pairinherited
Methods
-
addConsumer(
Node consumer) → void -
Adds a consumer node to this stream.
inherited
-
addInlet(
PeerHandle handle) → Future< void> -
Subscribes to a peer found by discovery.
inherited
-
addProducer(
Node producer) → void -
Adds a producer node to this stream.
inherited
-
create(
) → Future< void> -
Creates.
inherited
-
createInletsForNodes(
Iterable< Node> nodes, {Duration resolveTimeout = const Duration(seconds: 10)}) → Future<void> -
Subscribes to every node in
nodes, resolving their endpoints first.inherited -
createOutlet(
) → Future< void> -
Creates this node's publishing endpoint.
inherited
-
decodeControlPayload(
Object payload) → IMessage< IMessageType> ? -
Builds a message from a relayed JSON payload.
override
-
decodeSample(
Uint8List frame) → IMessage< IMessageType> ? -
Builds a message from a binary sample frame. Null for streams that do
not carry samples.
override
-
destroyStream(
) → Future< void> -
Stops and releases the stream's endpoints.
inherited
-
dispose(
) → Future< void> -
Disposes the ILifecycle implementation, releasing any resources.
inherited
-
encodeForWire(
IMessage< IMessageType> message) → Object -
JSON-encodable form of
messagefor the control path.override -
flushStreams(
) → Future< void> -
Discards anything buffered.
inherited
-
isConsumer(
Node node) → bool -
Checks if a given node is a consumer for this stream.
inherited
-
isProducer(
Node node) → bool -
Checks if a given node is a producer for this stream.
inherited
-
noSuchMethod(
Invocation invocation) → dynamic -
Invoked when a nonexistent method or property is accessed.
inherited
-
pause(
) → FutureOr< void> -
Pauses the implementation.
inherited
-
pauseStream(
) → Future< void> -
Pauses, if running and not already paused.
inherited
-
publishSample(
List< Object?> channels) → void -
Publishes a sample using the binary framing.
inherited
-
recreateOutlet(
) → Future< void> -
Recreates the publishing endpoint to reflect a changed Node.
inherited
-
removeInlet(
String nodeUId) → Future< void> -
Unsubscribes from a peer, releasing whatever that inlet holds.
inherited
-
resume(
) → FutureOr< void> -
Resumes the implementation.
inherited
-
resumeStream(
{bool flushBeforeResume = true}) → Future< void> -
Resumes, if running and paused.
inherited
-
resumeWith(
{bool flushBeforeResume = true}) → Future< void> -
Resumes, optionally discarding data buffered while paused.
inherited
-
sendData(
Iterable data) → Future< void> -
Publishes one sample: exactly NetworkStream.channelCount values whose
runtime types match NetworkStream.dataType.
override
-
sendDataAt(
double sourceClock, Iterable data) → Future< void> -
sendData for a sample that came from somewhere else:
sourceClockis the time its origin stamped it with, in seconds on the origin's clock, and travels in place of this node's send time.inherited -
sendDataTyped<
S> (Iterable< S> data) → Future<void> -
sendData for a statically known element type.
override
-
sendMessage(
IMessage< IMessageType> message) → FutureOr<void> -
inherited
-
setStreamNode(
Node node) → void -
override
-
start(
) → Future< void> -
Begins publishing and delivering. Idempotent.
inherited
-
stop(
) → Future< void> -
Stops publishing and delivering, but keeps the stream usable.
inherited
-
toString(
) → String -
A string representation of this object.
inherited
-
updateManager(
IResourceManager? newManager) → void -
inherited
-
updateNode(
Node newNode) → void -
Replaces the Node this stream publishes as.
inherited
Operators
-
operator ==(
Object other) → bool -
The equality operator.
inherited