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 an isolate that waits inside lsl_pull_sample; 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 isolate starts when the stream is listened to and stops when the subscription is cancelled. wakeInterval is how long, in seconds, that can take. Destroying the inlet stops it 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.

A listener that does not keep up. The isolate 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 isolate 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 isolate 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 isolate. debugFailAfter is for tests: the isolate throws 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,
  streamInfoAddress: streamInfo.streamInfo.address,
  wakeInterval: wakeInterval,
  maxBacklog: maxBacklog,
  backlogWarnAt: backlogWarnAt ?? _defaultBacklogWarnAt,
  onBacklog: onBacklog,
  onStarted: _listeners.add,
  onEnded: _listeners.remove,
  debugFailAfter: debugFailAfter,
);