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 nodeUId was 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<IMessageType>> stream) → void
Puts stream's NetworkStream.inletHealth on events as StreamReceiveHealthEvents, until the stream is disposed.

Operators

operator ==(Object other) → bool
The equality operator.
inherited