peer_coordinator 0.3.1
peer_coordinator: ^0.3.1 copied to clipboard
Transport-neutral peer coordination: election, membership, heartbeat and synchronised data streams over a pluggable backend. Pure Dart, web-safe.
0.3.1 #
-
New
ClockSyncSample: one clock-offset estimate for one peer —offset,uncertainty(full probe RTT, so the ± bound isbound),remoteTime(the peer's own clock when the estimate was taken),receivedClock, andclockReset. Exposed onNetworkStream.clockSyncs, which defaults to an empty broadcast stream so transports that estimate no offsets need no change.Complements
MessageTimingrather than duplicating it:MessageTiminganswers "when did this message arrive",ClockSyncSampleanswers "how well do the two clocks agree right now, and how is that changing". The(remoteTime, offset)pairs are what make drift fittable; an offset alone is not. And because the estimates arrive on their own cadence, a stream that carries no traffic for a while no longer leaves a hole in the record.
0.3.0 #
Defines what happens when the coordinator goes away. Previously nothing did: the liveness sweep ran on the coordinator only, so a participant whose coordinator had gone kept heartbeating and publishing into a stream with no consumer, indefinitely and silently, and there was no way for a departing coordinator to say it was leaving.
Behaviour changes #
-
A session now ends when its coordinator does. New
CoordinationSessionConfig.coordinatorLossPolicy, defaulting toCoordinatorLossPolicy.endSession: the node stops its timers, drops the topology, and refuses further sends.CoordinatorLossPolicy.reelectinstead has the survivors re-elect between themselves;CoordinatorLossPolicy.remainOpenis the previous behaviour, for applications that drive their own recovery.This is the breaking change in this release. Anything relying on a session outliving its coordinator must now say
remainOpen. -
Sends after a session ends throw
StateError.sendUserMessage,createStreamandstartStreamfail loudly rather than publishing into a stream nobody is reading. -
A departing coordinator announces it. New
CoordinationMessageType.sessionEnd/SessionEndMessage, the counterpart to the participant-onlynodeLeaving. Survivors learn within a round trip instead of waiting outnodeTimeout. It is accepted in any phase — a node still handshaking is the one that most needs to hear it — and only from the node that actually holds the coordinator role. -
An evicted node is told it was evicted. The timeout sweep sends
SessionEndMessage(evicted)to the node it is dropping, best effort. -
New
SessionEndedEventandevents.sessionEnded, carrying aSessionEndReason(coordinatorLeft,coordinatorTimedOut,coordinatorTransportLost,evicted) and the policy applied. Emitted under every policy, so there is one place to listen.
WebSocket protocol #
-
wsProtocolVersionis now 2, addingWsControl.signal: an opaque payload the hub forwards verbatim to one named endpoint, checking only that the sender owns thefromendpoint. It is the hub's only unicast and the only frame whose contents it never inspects.This exists so a peer-to-peer transport can use this hub for discovery and connection setup — WebRTC offer/answer/candidate exchange — while its data never touches the hub. That exchange cannot ride the coordination stream, because the coordination stream is what is being established.
WsFrame.decoderejects unknown versions rather than guessing, so a hub and its clients must be deployed together. Client side:WsConnection.sendSignalandWsConnection.signals, plus the exportedWsSignal.
Fixes #
-
The node-timeout sweep no longer spins. Its period was
Duration(seconds: nodeTimeout.inSeconds ~/ 2), which truncates to zero for any timeout under two seconds — a periodic timer firing every event-loop turn. Computed in microseconds now. -
A failed join no longer promotes the node to coordinator. Election wrapped both discovery and role setup in one
catch, so a participant that could not reach its coordinator became one instead. In a re-election that is how two survivors both take the role and split the session. Only the discovery call is caught now. -
Departure notices reach the wire.
announceLeavingenqueued onto the handler's outgoing stream, which is drained by a listener on a microtask, and the stream was disposed on the next line —WsStreamMixin.sendMessagedrops silently once disposed. Departures are now sent directly and awaited before teardown. -
Heartbeat entries are cleared for nodes that never joined. The clear sat inside the "was it in the topology?" branch, so an entry with no matching node survived every removal and was re-reported stale on every tick, re-broadcasting the topology each time.
-
CoordinationSessionConfig.hashCodeand==no longer recurse.idis derived fromhashCode, andhashCodehashedid; either would have overflowed the stack. Nothing called them, which is why it went unnoticed. -
copyWithandfromMapstopped dropping fields.copyWithresetclockSyncConfigto the default;fromMapdroppeddiscoveryIntervalandconsumeCoordinationStreamAsCoordinator.
0.2.0 #
Fixes a stream lifecycle race in which a startStream command could overtake
the createStream it belonged to, leaving a participant holding a stream it
never started while the coordinator believed it was running.
Behaviour changes #
-
A
coordinatorOnlycreateDataStreamnow waits for its consumers. Creation previously waited only for producers, and acoordinatorOnlystream has none — so the coordinator returned immediately and callers that issuestartStreamnext (the normal pattern) broadcast a start command milliseconds later, while participants still needed the best part of a second to resolve the outlet and build an inlet. Creation now completes only once every consumer has acked, which makes create-then-start ordering a guarantee rather than a hope.This adds roughly one participant's inlet-build time to each
coordinatorOnlystream creation. Unlike the producers wait, a consumer that does not ack within the timeout is logged atsevereand creation continues: the outlet is valid regardless, and late consumers can still attach off their ownstreamReady. -
getDataStreamwaits for an in-flight create instead of throwing. The stream lock only covers registration, not the outlet and inlet wiring that follows, so a lookup landing in that window was told the stream did not exist. It now queues behind the create. It still throwsArgumentErrorfor a stream nobody is creating. -
Concurrent
createDataStreamcalls for one name share a single setup. The second call used to find nothing registered, fall through, and re-run the wiring on the stream the first had just registered — callingcreateOutlet()on a stream that already had one. -
A
startStreamfor a stream that does not exist logssevere, notwarning. Nothing retries it, so the failure surfaces here or not at all.
Internal #
_waitForParticipantStreamsReadyis replaced by_StreamReadyAcks, which begins collecting acks when it is constructed rather than when it is awaited. The previous wait was built on the forward-onlywaitForEventand was reached only after the outlet existed, so an ack from a fast participant could land in the gap and be missed — costing the full timeout for a message already sent.