sampleStream method
- double wakeInterval = 0.1,
- int? maxBacklog,
- int? backlogWarnAt,
- void onBacklog(
- LSLBacklog backlog
- int? debugFailAfter,
Samples as they arrive, each with the local clock at the moment it became available, with no polling in between. See chunkStream for streams too fast or too wide for a message per sample.
Available in both modes. With useIsolates: true the inlet's own
isolate goes on serving time correction and stream info, but must not
be asked to pull or flush while this is listened to.
Pulling is how an inlet learns of new samples, so a loop that pulls
with a zero timeout sees each one up to a poll interval late, and that
delay is in any latency it measures. This instead gives the inlet a
native thread that waits inside liblsl's pull; liblsl wakes it when a
sample is queued, and it reads lsl_local_clock() as the call returns
(LSLTimedSample.receivedClock) before handing the sample over. How
soon a listener then runs is up to its own isolate's event loop, but
the receive time is already taken.
The thread is not an isolate and costs the Dart VM nothing while its stream is quiet, so hundreds of inlets can be listened to at once. (In 1.1.0 it was an isolate, and sixteen of them were enough to starve every other isolate of the process.)
The thread starts when the stream is listened to and stops when the
subscription is cancelled. wakeInterval is how long, in seconds, that
can take, and so how long cancelling or destroy can take. It is not a
polling interval: a sample wakes the thread at once whatever its value,
and it only bounds how long the thread waits on a quiet stream before it
looks for a request to stop. It has to be above zero (an
ArgumentError otherwise), since liblsl does not wait at all on a zero
timeout. Destroying the inlet stops the thread too, and the stream then
ends with an LSLSampleListenerException. While listening, pulling
from or flushing this inlet anywhere else throws an LSLException: an
inlet's samples have one reader. Time correction calls are unaffected.
If the listening isolate exits or is killed without cancelling, or without destroying the inlet, finalizers stop the thread and then destroy the inlet, which is never destroyed while the thread is still in a pull. Nothing else should destroy it then: in particular not another isolate, by address.
A listener that does not keep up. The thread hands samples over
as fast as they arrive, whatever the listener does with them, so by
default nothing is ever held back or lost and
LSLTimedSample.receivedClock is always the arrival; a listener that
is slower than the stream then has a queue that grows without limit.
That is reported: when more than backlogWarnAt samples are queued
(by default one second of the stream, and at least 1000), onBacklog
is called, at most once a second, and once more when the queue is back
under half of that. Without onBacklog it is logged as a warning by
the liblsl logger of package:logging. The count is of samples on
their way to the stream; what a paused subscription buffers, or a
listener's own unfinished futures, is not in it.
With maxBacklog the thread stops pulling while that many samples
are queued, and while the subscription is paused. What arrives then
waits in the inlet's buffer, which is bounded (maxBuffer, where
liblsl drops the oldest when it is full), so memory is too. The price
is in the receive clock, which for a sample that waited there is when
it was taken out rather than when it arrived: leave maxBacklog unset
when measuring latency.
The stream closes without an error only when it was cancelled. If the
thread ends for any other reason, the stream delivers an
LSLSampleListenerException and then closes, so listen with onError
(and onDone): after that the inlet is no longer being read. Listening
again starts a new thread. debugFailAfter is for tests: the thread
fails after that many samples.
final inlet = await LSL.createInlet<double>(streamInfo: info, useIsolates: false);
final subscription = inlet.sampleStream().listen((sample) {
final latency = sample.receivedClock - (sample.timestamp + offset);
});
// ...
await subscription.cancel();
await inlet.destroy();
Implementation
Stream<LSLTimedSample<T>> sampleStream({
double wakeInterval = 0.1,
int? maxBacklog,
int? backlogWarnAt,
void Function(LSLBacklog backlog)? onBacklog,
int? debugFailAfter,
}) => listenToInlet<T>(
inletAddress: () => _listenAddress,
format: streamInfo.channelFormat,
channels: streamInfo.channelCount,
wakeInterval: wakeInterval,
maxBacklog: maxBacklog,
backlogWarnAt: backlogWarnAt ?? _defaultBacklogWarnAt,
onBacklog: onBacklog,
onStarted: _listeners.add,
onEnded: _listeners.remove,
debugFailAfter: debugFailAfter,
);