All files stream_subscription.ts

100% Statements 37/37
100% Branches 20/20
100% Functions 8/8
100% Lines 36/36

Press n or j to go to the next uncovered block, b, p or k for the previous block.

1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 864x                             4x   136x 136x     136x 136x 136x       297x 294x 291x   166x 166x   35x 34x   35x 2x   35x   90x 47x   90x     294x       292x 292x       12x 11x 11x 2x   11x         6x 5x 5x 1x   5x 5x         48x       1x      
import {
  StreamListener,
  StreamListenOptions,
  StreamMessage,
  StreamMessageType,
} from "./types";
 
import Stream from "./stream";
 
export interface StreamSubscriptionActions {
  pause: () => void;
  resume: () => void;
  cancel: () => void;
}
 
export default class StreamSubscription<T>
  implements StreamSubscriptionActions {
  private _buffer: Array<StreamMessage<T>> = [];
  private _isPaused: boolean = false;
 
  constructor(
    private stream: Stream<T>,
    private onData: StreamListener<T>,
    private listenOptions: StreamListenOptions = {}
  ) {}
 
  private _emit() {
    if (this._isPaused) return;
    this._buffer.forEach((m) => {
      switch (m.type) {
        case StreamMessageType.Data:
          this.onData(m.data);
          break;
        case StreamMessageType.Error:
          if (this.listenOptions.onError) {
            this.listenOptions.onError(m.data);
          }
          if (this.listenOptions.cancelOnError) {
            this.cancel();
          }
          break;
        case StreamMessageType.Done:
          if (this.listenOptions.onDone) {
            this.listenOptions.onDone();
          }
          break;
      }
    });
    this._buffer = [];
  }
 
  messageHandler(message: StreamMessage<T>) {
    this._buffer.push(message);
    this._emit();
  }
 
  pause() {
    if (!this._isPaused) {
      this._isPaused = true;
      if (this.listenOptions.onPause) {
        this.listenOptions.onPause();
      }
      this.stream.pause();
    }
  }
 
  resume() {
    if (this._isPaused) {
      this._isPaused = false;
      if (this.listenOptions.onResume) {
        this.listenOptions.onResume();
      }
      this.stream.resume();
      this._emit();
    }
  }
 
  cancel() {
    this.stream.cancel(this);
  }
 
  get isPaused(): boolean {
    return this._isPaused;
  }
}