sampleStream method

Stream<LSLTimedSample<T>> sampleStream({
  1. double wakeInterval = 0.1,
  2. int? maxBacklog,
  3. int? backlogWarnAt,
  4. void onBacklog(
    1. LSLBacklog backlog
    )?,
  5. 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,
);