liblsl_coordinator 0.4.1
liblsl_coordinator: ^0.4.1 copied to clipboard
A performance-focused Dart (and Flutter) LSL-based device coordination library.
0.4.1 #
-
Streams now expose
clockSyncs, a broadcast stream ofClockSyncSamplecarrying each clock-offset estimate as it is made: the offset, its uncertainty, the peer's own clock at the moment it was measured (remoteTime), and whether the peer's clock may have been reset since the previous estimate (clockReset). LSL supplies all four —remoteTimeand the reset flag come fromlsl_time_correction_exandlsl_was_clock_resetand were previously discarded inside the inlet worker.These are deliberately not fields on
MessageTiming. An estimate is refreshed at most every 5 s while data samples arrive hundreds of times a second, so carrying them per sample would repeat one value hundreds of times across the isolate port. More importantly, a per-message view can only describe estimates that happened to be attached to a message that happened to arrive: clock drift is a property of the clocks, not of the traffic, so a quiet stream left an unbridgeable gap.clockSyncsticks on the estimate's own cadence regardless of traffic.The getter is on the base
NetworkStreamand defaults to an empty stream, so a transport that estimates no offsets needs no change and a consumer does not have to know which transport it is on.Note that reading liblsl's clock-reset flag clears it. It is read exactly once per refresh inside the inlet worker and reported on that estimate; nothing else may poll it without consuming the notification.
0.4.0+1 #
- LSL messages now carry
MessageTiming.uncertainty. The inlet isolate reads the offset throughlsl_time_correction_exinstead oflsl_time_correction, which is the same native round trip, and passes liblsl's own error bound (the probe's full RTT) along with the offset. LSL was previously the only transport that reported an offset without saying how good it was; the figure now means the same thing here as the oneClockSyncServicereports for the WebSocket and WebRTC transports.
0.4.0+0 #
The coordination layer is now transport-neutral and lives in a new pure-Dart
package, peer_coordinator. liblsl_coordinator is the Lab Streaming Layer
transport for it, and re-exports the core so existing imports keep working.
Two other transports ship in peer_coordinator: an in-memory one for testing
and a WebSocket one (with a relay hub) that also runs in the browser.
Breaking changes #
- The core moved to
package:peer_coordinator. Every library previously underpackage:liblsl_coordinator/—framework.dart,config.dart,coordination.dart,data.dart,discovery.dart,interfaces.dart,network.dart,logging.dart— is now a one-line re-export of itspeer_coordinatorcounterpart. Existing imports continue to resolve unchanged; prefer importing frompeer_coordinatorin new code.liblsl_coordinatornow depends on it. LSLCoordinationSessionis a thin subclass ofPeerSession. The session flow, election, heartbeat, membership and stream lifecycle are all inPeerSessionnow. Return types are narrowed covariantly, socreateDataStream/getDataStreamstill hand backLSLDataStreamandtransportstill returnsLSLTransport. No call-site change expected.ITransportConfigrequirescreateTransport(). Transport selection is a method on the config rather than a registry, so an application never compiles in a transport it does not name. Only affects custom transports.ITransportgainedstreamFactoryandcreateDiscovery, and lostcreateStream.createStreamhad no callers.ITransportnow also implementsIResourceManager, which it did in practice already.NetworkStreamgained the lifecycle members that previously existed only onLSLStreamMixin:started,start,stop,createOutlet,recreateOutlet,addInlet,createInletsForNodes,updateNode, plus defaultpauseStream/resumeStream/flushStreams/destroyStream. Callers no longer need a concrete LSL type to drive a stream.resume()was split.IPausable.resume()takes no arguments; the parameterised form is nowresumeWith({flushBeforeResume}). The old widenedresume({flushBeforeResume})could never be called through the base type.createResolvedInletsForStreamrenamed tocreateInletsForNodesandaddInletnow takes aPeerHandle, not anLSLStreamInfo.- Discovery is typed.
LslDiscoveryimplementsIDiscovery;startDiscovery(predicate:)/stopDiscovery()are nowstart(query:)/stop(), taking aDiscoveryQueryinstead of an XPath string. The LSL transport compiles queries to XPath viaLslPredicateCompiler, which emits byte-identical predicates (golden-tested).LSLStreamInfoHelper's string builders remain for now. - Discovery events renamed.
StreamDiscoveredEventis nowPeersDiscoveredEvent, carryingList<PeerHandle>rather thanList<StreamInfoResource>. Log.sendPortis nowLog.forwardTo(callback)andLog.logIsolateMessageisLog.replayRecord(old names deprecated). Logging no longer importsdart:isolateordart:io, so the core compiles for web. The default logger name changed fromLSLCoordinatortoPeerCoordinator.- Removed as unused:
TransportStreamConfig/TransportCoordinationStreamConfig(never implemented, and unusable — the field was dropped byDataStreamConfigFactory.fromMap, so per-stream options could never reach other nodes; put transport options onITransportConfig),IStartable,NullNode, theNetworkTopology/HierarchicalTopologyruntime classes andTopologyType(the configs remain), thesrc/events.dartEventhierarchy, andIntMessageTypeMapping.minValue/maxValue(never read, and0x7FFFFFFFFFFFFFFFcannot be represented in JavaScript). CoordinationConfig.namenow defaults to'peer_coordinator'. It surfaces only as theappIdnode metadata and is not used in discovery.
Fixed #
- Duplicate
createDataStreamleaked a live outlet. Participants build their streams automatically on the coordinator'screateStreamcommand, so application code that also calledcreateDataStreamran the setup path twice. Two concurrentcreateOutlet()calls both passed its null check (it awaited before assigning), both built an outlet, and the second overwrote the first. The orphan stayed published on the network with nothing referencing it — unreachable fromdispose(), and indistinguishable from a teardown leak — while the second listen on the single-subscription outgoing controller threwBad state: Stream has already been listened to.createDataStreamis now idempotent andcreateOutlet()serialises on an in-flight future. createInletsForNodesresolved once and threw if any producer had not published yet. Producers come up independently, so inallNodesmode one node routinely looked for another that was not ready — a race the caller could not win. It now polls a continuous resolver until the deadline.createInletsForNodesleaked everyLSLStreamInfoit resolved but did not use. Unmatched infos are destroyed.- Election's self-exclusion never excluded anything. It emitted
not(starts-with(source_id, '<nodeId>')), butsource_idbegins with the stream name, so the clause was always true. Now excludes by node uId. randomRollwas written under one metadata key and read under another ('random_roll'vs'randomRoll'), so the value silently never arrived on the skip-election path. Both sides now usePeerMetadataKeys.- XPath predicate values were never quoted, so an apostrophe in a session name or node id produced a malformed predicate that silently matched nothing.
- An unpromoted node lost its identity across the wire.
Nodealways starts withrole: 'none', andNodeFactoryreturnedNullNode()for that role, discarding the supplied config and generating a freshuId. RuntimeTypeUIDcacheduIdin astatic Map<Type, String>, so every session, data stream and coordination stream in a process shared one identity. Harmless across processes, fatal for in-process multi-node use.- Unguarded
topologyConfig as HierarchicalTopologyConfigcrashed election for any other topology type. PeerHandleownership is explicit. Continuous discovery frees the previous cycle's nativelsl_streaminfopointers, so a handle still in use became a dangling pointer.addInletnow callstake()before touching a handle, which is also what unified the two call sites that previously differed.
Added #
peer_coordinator— the transport-neutral core. Pure Dart, nodart:io,dart:isolateordart:ffi; compiles to JavaScript (tool/web_safety_check.dartenforces this in CI).- In-memory transport (
package:peer_coordinator/in_memory.dart) for testing whole multi-node sessions with no sockets and no timing luck. - WebSocket transport and relay hub
(
package:peer_coordinator/websocket.dart,.../hub.dart, ordart run peer_coordinator:hub). Control traffic is JSON; data samples use a binary frame the hub relays without parsing. On loopback at 60 Hz it measured p50 ~0.7 ms against LSL's ~6.1 ms, because LSL's figure at that rate is bounded by its polling interval rather than the network. DiscoveryQuery— a typed peer filter that relay transports evaluate directly and LSL compiles to XPath.- Cross-transport conformance scenarios (
package:peer_coordinator/testing.dart) covering everyStreamParticipationMode. All five modes pass over in-memory, WebSocket and LSL. DataStreamConfig.validateSample, shared by every transport.- A test suite for this package, where there was none: unit tests for the
message/JSON contract, handlers and state; XPath goldens; and tagged LSL
integration tests (
melos run test:lsl).
Notes #
- LSL tests must run serially (
--concurrency=1): they bind real sockets and resolve peers machine-wide, so concurrent files interfere.melos run testexcludes them;melos run test:lslruns them serially. - Two behaviours are documented in tests rather than fixed, because changing
them is a deliberate decision: a rejected node cannot distinguish rejection
from a timeout (the handler's
StateErroris swallowed by the controller's log-only catch), and a full session re-offers and re-rejects a waiting node on every discovery cycle.
0.3.0+0 #
removeInletnever actually removed or destroyed the inlet (the lazywhereIndexediterable was never consumed), so removed sources kept delivering samples and leaking native inlets.- Received samples could be tagged with another inlet's LSL time correction, because the correction index only advanced when a sample was present.
- Time corrections were never refreshed on the timer-based polling path (only
in busy-wait mode), so
lsl_time_correctionstayed at its initial value. LSLStreamMixin.create()was a permanent no-op (_createdwas initialised totrue), leavingcreatedwrong for the stream's whole lifetime.- The direct-mode (non-isolate) inlet polling timer was never cancellable and kept firing after the stream was disposed.
createInletForNodecompared a node UID against full LSL source IDs, so its duplicate check never matched and duplicate inlets could be created.leave()disposed the transport before the controller, while the controller still needed transport-owned discovery to announce leaving.- A stopped, crashed, or exited outlet/inlet isolate left in-flight requests
awaiting a response that could never arrive; they now fail with a
StateErrorinstead of hanging forever. - Pause/flush work in the busy-wait inlet loop was not awaited and could interleave with polling.
StreamControllers in the isolate manager, LSL streams, and discovery are now closed on teardown (previously disabled because awaitingclose()hung).- The outlet
StreamInfohanded to the outlet isolate is now destroyed when the worker stops, not only on outlet recreation. - Worker-side
ReceivePorts are closed and isolate log forwarding is stopped on worker shutdown, so isolates can exit rather than relying onkill. - Event waits with no timeout (
waitForMinNodes,waitForUserMessage,_waitForPhase) no longer leak their subscription; all waits now share onewaitForEventhelper that always cancels. - The session's stream-lifecycle subscription is cancelled on
dispose(). - Sends no longer block on a round-trip acknowledgement from the outlet isolate. A pool of eight native sample buffers gives backpressure, and a send only waits when every buffer is still in flight.
- Inlets are now drained (up to 100 samples per inlet per tick) instead of yielding a single sample per tick, which previously caused an unbounded backlog whenever the producer outpaced the poll interval. At 500 Hz on loopback this changes 79% sample loss and multi-second latency into zero loss at ~1.1 ms p50.
- Timer-based polling derives its interval from
sampleRate(clamped to 1–10 ms) rather than being hard-coded to 10 ms. - Removed per-sample UUID generation, per-sample lock acquisition, per-sample channel type re-scanning, and per-sample ISO-8601 string formatting from the receive path.
LSLDataStream.sendDataandsendDataTypednow returnFuture<void>. Awaiting is optional (existing fire-and-forget calls keep working) and gives backpressure when the send buffer pool is saturated.- Clock-offset estimation for the WebSocket transport. Two peers in separate
processes read monotonic clocks with unrelated epochs, so
MessageTiming.clockOffsetused to benullover WebSocket and no transit time could be computed. An NTP-style estimator modelled on liblsl'stime_receivernow runs there, one per inlet (the receiver estimates its producer's clock, as liblsl does), over the existingConnectionTestround trip. Tuned byCoordinationSessionConfig.clockSyncConfig, which defaults to liblsl's[tuning]constants. Not run on LSL (nativelsl_time_correction) or in-memory (peers share a clock).ConnectionTest/ConnectionTestResponsegainedtoNodeUId,waveId,requestSenderClockandrequestReceivedClock, all optional and additive.- Both roles now answer connection tests and consume responses. The previous
split — coordinator answers, participant listens — could not express per-inlet
probing.
CoordinatorMessageHandler.canHandlegainedconnectionTestResponseandParticipantMessageHandler.canHandlegainedconnectionTest. - A
ConnectionTestResponseis now addressed to its requester, so the other participants it is broadcast to ignore it instead of logging'Received unexpected connection test response'.
MessageTiming.uncertainty: the error bound onclockOffset, as the full round-trip time of the probe it came from — the true offset lies within±uncertainty/2, exposed astransitUncertaintySeconds. Same conservative quantitylsl_time_correction_exreports. Null on LSL, which has no Dart wrapper forlsl_time_correction_exyet.- WebSocket data samples now identify their sender by node uId rather than by a hub slot number, and carry the estimated offset, so a WS data stream's transit time is measurable too.
- Message timing moved from the metadata map to a typed
MessageTiming. Thelsl_timestamp,lsl_time_correction,received_atandsource_idmetadata keys are gone; readmessage.timinginstead, whose fields aresourceClock,clockOffset,receivedClockandsourceId, withtransitSeconds/transitMicrosdoing the clock-domain arithmetic for you. Nothing in the wild could have depended on the old keys — nothing read them and the values never reached a coordination-message consumer (see below). LSLCoordinationSession.waitForUserMessagenow matches on the message type (the first argument tosendUserMessage) rather than an auto-generated message ID, which no caller could know — the method was previously unusable. It also matches participant messages and returnsUserMessageEvent.sendUserMessage's first parameter is renamedmessageId→messageTypeto match what it actually is (positional; no call-site change needed).- New
NodeJoinRejectedEventon the coordinator's event stream; the coordinator now clears its pending-join tracking when a join is rejected, so a rejected node can be re-offered a join if capacity frees up.
0.1.1+0 #
- Initial version.