246 lines
5.2 KiB
JavaScript
246 lines
5.2 KiB
JavaScript
|
|
import { Readable, Writable, Duplex } from 'stream';
|
||
|
|
|
||
|
|
class PushSource {
|
||
|
|
constructor(readable) {
|
||
|
|
this.readable = readable;
|
||
|
|
this.sink = null;
|
||
|
|
this.paused = true;
|
||
|
|
this.ended = false;
|
||
|
|
this._endedWithError = null;
|
||
|
|
|
||
|
|
this._onData = this._onData.bind(this);
|
||
|
|
this._onEnd = this._onEnd.bind(this);
|
||
|
|
this._onError = this._onError.bind(this);
|
||
|
|
this._onResume = this._onResume.bind(this);
|
||
|
|
|
||
|
|
readable.on('data', this._onData);
|
||
|
|
readable.on('end', this._onEnd);
|
||
|
|
readable.on('error', this._onError);
|
||
|
|
}
|
||
|
|
|
||
|
|
pipe(sink) {
|
||
|
|
this.sink = sink;
|
||
|
|
sink.source = this;
|
||
|
|
if (!sink.paused) {
|
||
|
|
this.resume();
|
||
|
|
}
|
||
|
|
return sink;
|
||
|
|
}
|
||
|
|
|
||
|
|
resume() {
|
||
|
|
if (this.ended) return;
|
||
|
|
this.paused = false;
|
||
|
|
this.readable.resume();
|
||
|
|
this._drain();
|
||
|
|
}
|
||
|
|
|
||
|
|
abort(err) {
|
||
|
|
if (this.ended) return;
|
||
|
|
this.ended = err || new Error('Aborted');
|
||
|
|
this._cleanup();
|
||
|
|
if (this.sink) {
|
||
|
|
this.sink.end(this.ended);
|
||
|
|
this.sink.source = null;
|
||
|
|
this.sink = null;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_drain() {
|
||
|
|
// _drain is no longer needed - _onData handles all data events
|
||
|
|
// This method is kept for compatibility but does nothing
|
||
|
|
// The readable stream's 'data' event handler (_onData) is already set up in constructor
|
||
|
|
}
|
||
|
|
|
||
|
|
_onData(data) {
|
||
|
|
if (this.ended) return;
|
||
|
|
if (this.sink) {
|
||
|
|
if (this.sink.paused) {
|
||
|
|
this.paused = true;
|
||
|
|
this.readable.pause();
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
this.sink.write(data);
|
||
|
|
this.paused = this.sink.paused;
|
||
|
|
} else {
|
||
|
|
this.paused = true;
|
||
|
|
this.readable.pause();
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_onEnd() {
|
||
|
|
if (this.ended) return;
|
||
|
|
this.ended = true;
|
||
|
|
this._cleanup();
|
||
|
|
if (this.sink) {
|
||
|
|
this.sink.end(this._endedWithError);
|
||
|
|
this.sink.source = null;
|
||
|
|
this.sink = null;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_onError(err) {
|
||
|
|
if (this.ended) return;
|
||
|
|
this._endedWithError = err;
|
||
|
|
this.ended = err;
|
||
|
|
this._cleanup();
|
||
|
|
if (this.sink) {
|
||
|
|
this.sink.end(err);
|
||
|
|
this.sink.source = null;
|
||
|
|
this.sink = null;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_onResume() {
|
||
|
|
if (this.paused && !this.ended) {
|
||
|
|
this.paused = false;
|
||
|
|
this.readable.resume();
|
||
|
|
this._drain();
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_cleanup() {
|
||
|
|
this.readable.removeListener('data', this._onData);
|
||
|
|
this.readable.removeListener('end', this._onEnd);
|
||
|
|
this.readable.removeListener('error', this._onError);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
class PushSink {
|
||
|
|
constructor(writable) {
|
||
|
|
this.writable = writable;
|
||
|
|
this.source = null;
|
||
|
|
this.paused = false;
|
||
|
|
this.ended = false;
|
||
|
|
this._endedWithError = null;
|
||
|
|
this._draining = false;
|
||
|
|
|
||
|
|
this._onDrain = this._onDrain.bind(this);
|
||
|
|
this._onError = this._onError.bind(this);
|
||
|
|
|
||
|
|
writable.on('drain', this._onDrain);
|
||
|
|
writable.on('error', this._onError);
|
||
|
|
}
|
||
|
|
|
||
|
|
pipe(sink) {
|
||
|
|
if (sink && sink.write) {
|
||
|
|
this.source = sink;
|
||
|
|
sink.source = this;
|
||
|
|
}
|
||
|
|
return sink;
|
||
|
|
}
|
||
|
|
|
||
|
|
write(data) {
|
||
|
|
if (this.ended) return;
|
||
|
|
const ok = this.writable.write(data);
|
||
|
|
if (!ok) {
|
||
|
|
this.paused = true;
|
||
|
|
}
|
||
|
|
return ok;
|
||
|
|
}
|
||
|
|
|
||
|
|
end(err) {
|
||
|
|
if (this.ended) return;
|
||
|
|
this.ended = err || true;
|
||
|
|
this._cleanup();
|
||
|
|
if (this.writable.destroyed || this.writable.writableEnded) {
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
if (err) {
|
||
|
|
this.writable.destroy(err);
|
||
|
|
} else {
|
||
|
|
this.writable.end();
|
||
|
|
}
|
||
|
|
if (this.source) {
|
||
|
|
this.source.sink = null;
|
||
|
|
this.source = null;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
abort(err) {
|
||
|
|
if (this.ended) return;
|
||
|
|
this.ended = err || new Error('Aborted');
|
||
|
|
this._cleanup();
|
||
|
|
if (!this.writable.destroyed) {
|
||
|
|
this.writable.destroy(this.ended);
|
||
|
|
}
|
||
|
|
if (this.source) {
|
||
|
|
this.source.abort(this.ended);
|
||
|
|
this.source = null;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_onDrain() {
|
||
|
|
if (this.paused && !this.ended) {
|
||
|
|
this.paused = false;
|
||
|
|
if (this.source) {
|
||
|
|
this.source.resume();
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_onError(err) {
|
||
|
|
if (this.ended) return;
|
||
|
|
this._endedWithError = err;
|
||
|
|
this.ended = err;
|
||
|
|
this._cleanup();
|
||
|
|
if (this.source) {
|
||
|
|
this.source.abort(err);
|
||
|
|
this.source.sink = null;
|
||
|
|
this.source = null;
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
_cleanup() {
|
||
|
|
this.writable.removeListener('drain', this._onDrain);
|
||
|
|
this.writable.removeListener('error', this._onError);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
class PushDuplex {
|
||
|
|
constructor(duplex) {
|
||
|
|
this.source = new PushSource(duplex);
|
||
|
|
this.sink = new PushSink(duplex);
|
||
|
|
this.paused = false;
|
||
|
|
this.ended = false;
|
||
|
|
|
||
|
|
this.sink.source = this.source;
|
||
|
|
this.source.sink = this.sink;
|
||
|
|
}
|
||
|
|
|
||
|
|
pipe(sink) {
|
||
|
|
return this.source.pipe(sink);
|
||
|
|
}
|
||
|
|
|
||
|
|
resume() {
|
||
|
|
this.source.resume();
|
||
|
|
}
|
||
|
|
|
||
|
|
abort(err) {
|
||
|
|
if (this.ended) return;
|
||
|
|
this.ended = err || new Error('Aborted');
|
||
|
|
this.source.abort(err);
|
||
|
|
this.sink.end(err);
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
export function source(readableStream) {
|
||
|
|
if (!(readableStream instanceof Readable)) {
|
||
|
|
throw new Error('Expected a Node.js Readable stream');
|
||
|
|
}
|
||
|
|
return new PushSource(readableStream);
|
||
|
|
}
|
||
|
|
|
||
|
|
export function sink(writableStream) {
|
||
|
|
if (!(writableStream instanceof Writable)) {
|
||
|
|
throw new Error('Expected a Node.js Writable stream');
|
||
|
|
}
|
||
|
|
return new PushSink(writableStream);
|
||
|
|
}
|
||
|
|
|
||
|
|
export function duplex(duplexStream) {
|
||
|
|
if (!(duplexStream instanceof Duplex)) {
|
||
|
|
throw new Error('Expected a Node.js Duplex stream');
|
||
|
|
}
|
||
|
|
return new PushDuplex(duplexStream);
|
||
|
|
}
|