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 86 | 4x 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;
}
}
|