CoordinationController class
Controls the coordination flow with clear phases and event-driven logic.
Owns election, the heartbeat and node-timeout timers, the discovery loop, and the coordinator/participant command API. None of that is transport-specific: it talks to an ITransport, a CoordinationStream and an IDiscovery, and the transport supplies whichever implementations it has.
Emits all coordination events through a single events stream. Use the ControllerEventStreamExtensions for convenient filtering.
Constructors
-
CoordinationController({required CoordinationConfig coordinationConfig, required ITransport<
ITransportConfig> transport, required Node thisNode, required CoordinationSession session})
Properties
- clockOffsets → PeerClockOffsets?
-
Per-peer clock offsets, or null when this transport does not estimate
them. Shared with the transport so data streams read the same table.
no setter
-
connectedNodes
→ List<
Node> -
no setter
-
connectedParticipantNodes
→ List<
Node> -
no setter
-
coordinationClockSyncs
→ Stream<
ClockSyncSample> -
Clock-offset estimates for the coordination stream's peers.
no setter
- coordinationConfig → CoordinationConfig
-
final
-
coordinationOutletConsumers
→ Stream<
bool> -
Emits when this node's coordination outlet gains or loses every consumer.
no setter
- coordinationSendFailures → int
-
How many consecutive coordination sends have failed. Zero when healthy.
no setter
- coordinatorUId → String?
-
no setter
- currentPhase → CoordinationPhase
-
no setter
- endReason → SessionEndReason?
-
Why the session ended, or null while it is live.
no setter
-
events
→ Stream<
ControllerEvent> -
Single public stream for all coordination events.
no setter
- hashCode → int
-
The hash code for this object.
no setterinherited
- heartbeatFailures → int
-
How many consecutive heartbeat sends have failed. Zero when healthy.
no setter
- isAcceptingNodes → bool
-
no setter
- isCoordinator → bool
-
no setter
- runtimeType → Type
-
A representation of the runtime type of the object.
no setterinherited
- session → CoordinationSession
-
final
- sessionEnded → bool
-
Whether this session has ended and must not send anything further.
no setter
- thisNode → Node
-
no setter
-
transport
→ ITransport<
ITransportConfig> -
final
Methods
-
createStream(
String streamName, DataStreamConfig config) → Future< void> -
destroyStream(
String streamName) → Future< void> -
dispose(
) → Future< void> - Dispose and cleanup
-
flushStream(
String streamName) → Future< void> -
initialize(
) → Future< void> - Initialize the controller - creates streams and discovery
-
markStreamReady(
String streamName) → Future< void> -
noSuchMethod(
Invocation invocation) → dynamic -
Invoked when a nonexistent method or property is accessed.
inherited
-
noteNodeActivity(
String nodeUId) → void -
pauseAcceptingNodes(
) → Future< void> -
pauseStream(
String streamName) → Future< void> -
resumeAcceptingNodes(
) → Future< void> -
resumeStream(
String streamName, {bool flushBeforeResume = true}) → Future< void> -
sendUserMessage(
String messageType, String description, Map< String, dynamic> payload, {String? parentMessageId}) → Future<void> -
sinceLastHeard(
String nodeUId) → Duration? -
Records that
nodeUIdwas heard from on some stream other than the coordination one, so the liveness sweep counts it as alive. -
start(
[Duration? timeout]) → Future< void> - Start the coordination process - begins election TODO: Don't always become coordinator based on capabilities
-
startStream(
String streamName, DataStreamConfig config, {DateTime? startAt}) → Future< void> -
stopStream(
String streamName) → Future< void> -
toString(
) → String -
A string representation of this object.
inherited
-
updateConfig(
Map< String, dynamic> config) → Future<void> -
watchStreamReceiveHealth(
NetworkStream< NetworkStreamConfig, IMessage< stream) → voidIMessageType> > -
Puts
stream's NetworkStream.inletHealth on events as StreamReceiveHealthEvents, until the stream is disposed.
Operators
-
operator ==(
Object other) → bool -
The equality operator.
inherited