Files

367 lines
8.5 KiB
JavaScript
Raw Permalink Normal View History

2026-07-19 13:44:11 -07:00
import { Readable, Writable, Duplex } from 'stream';
import { test, describe, beforeEach, afterEach } from 'node:test';
import { source, sink, duplex } from '../index.js';
function readableFromArray(arr) {
return new Readable({
objectMode: true,
read() {
if (arr.length > 0) {
this.push(arr.shift());
} else {
this.push(null);
}
}
});
}
function writableToArray() {
const arr = [];
return new Writable({
objectMode: true,
write(chunk, encoding, callback) {
arr.push(chunk);
callback();
}
});
}
describe('stream-to-push-stream', () => {
describe('source', () => {
test('converts readable stream to push-source', () => {
const received = [];
const readable = readableFromArray([1, 2, 3, 4, 5]);
const pushSrc = source(readable);
const output = [];
pushSrc.pipe({
write(data) {
output.push(data);
},
end() {},
get paused() { return false; },
set paused(v) {}
});
return new Promise((resolve) => {
setTimeout(() => {
if (output.length === 5) {
resolve();
} else {
resolve(new Error('Expected 5 items'));
}
}, 50);
});
});
test('handles backpressure', () => {
const received = [];
const output = [];
let paused = true;
const readable = readableFromArray([1, 2, 3, 4, 5]);
const pushSrc = source(readable);
pushSrc.pipe({
write(data) {
output.push(data);
if (output.length === 1) {
paused = true;
this.paused = true;
} else if (output.length === 3) {
paused = false;
this.paused = false;
pushSrc.resume();
}
},
end() {},
get paused() { return paused; },
set paused(v) { paused = v; }
});
return new Promise((resolve) => {
setTimeout(() => {
if (output.length === 5) {
resolve();
} else {
resolve(new Error('Expected 5 items, got ' + output.length));
}
}, 50);
});
});
test('propagates errors', () => {
const readable = new Readable({
read() {
this.destroy(new Error('Test error'));
}
});
const pushSrc = source(readable);
let errorReceived = null;
pushSrc.pipe({
write(data) {},
end(err) {
errorReceived = err;
},
get paused() { return false; },
set paused(v) {}
});
return new Promise((resolve) => {
setTimeout(() => {
if (errorReceived && errorReceived.message === 'Test error') {
resolve();
} else {
resolve(new Error('Error not propagated'));
}
}, 50);
});
});
test('handles already-ended stream', () => {
const readable = new Readable({
read() {}
});
readable.push(null);
const pushSrc = source(readable);
let endCalled = false;
pushSrc.pipe({
write(data) {},
end() { endCalled = true; },
get paused() { return false; },
set paused(v) {}
});
return new Promise((resolve) => {
setTimeout(() => {
if (endCalled) {
resolve();
} else {
resolve(new Error('End not called for already-ended stream'));
}
}, 50);
});
});
});
describe('sink', () => {
test('converts writable stream to push-sink', () => {
const writable = new Writable({
objectMode: true,
write(chunk, encoding, callback) {
callback();
}
});
const received = [];
writable.on('data', (chunk) => received.push(chunk));
const pushSnk = sink(writable);
pushSnk.write(1);
pushSnk.write(2);
pushSnk.write(3);
pushSnk.end();
return new Promise((resolve) => {
setTimeout(() => {
if (received.length === 3) {
resolve();
} else {
resolve(new Error('Expected 3 items, got ' + received.length));
}
}, 50);
});
});
test('handles backpressure from writable', () => {
let canWrite = false;
const writable = new Writable({
objectMode: true,
write(chunk, encoding, callback) {
if (!canWrite) {
callback(false);
} else {
callback();
}
}
});
const pushSnk = sink(writable);
let drainCount = 0;
writable.on('drain', () => {
drainCount++;
canWrite = true;
pushSnk.resume();
});
pushSnk.write(1);
pushSnk.write(2);
return new Promise((resolve) => {
setTimeout(() => {
if (drainCount > 0) {
resolve();
} else {
resolve(new Error('Drain not emitted'));
}
}, 50);
});
});
test('propagates writable errors', () => {
const writable = new Writable({
objectMode: true,
write(chunk, encoding, callback) {
callback(new Error('Write error'));
}
});
const pushSnk = sink(writable);
let errorReceived = null;
pushSnk.write(1);
writable.on('error', (err) => {
errorReceived = err;
});
return new Promise((resolve) => {
setTimeout(() => {
if (errorReceived) {
resolve();
} else {
resolve(new Error('Error not propagated'));
}
}, 50);
});
});
});
describe('duplex', () => {
test('converts duplex stream to push-duplex', () => {
const duplexStream = new Duplex({
objectMode: true,
read() {},
write(chunk, encoding, callback) {
this.push(chunk);
callback();
}
});
const pushDup = duplex(duplexStream);
const output = [];
pushDup.pipe({
write(data) {
output.push(data);
},
end() {},
get paused() { return false; },
set paused(v) {}
});
duplexStream.write(1, () => {
duplexStream.write(2, () => {
duplexStream.end();
});
});
return new Promise((resolve) => {
setTimeout(() => {
if (output.length === 2) {
resolve();
} else {
resolve(new Error('Expected 2 items, got ' + output.length));
}
}, 50);
});
});
test('handles source and sink independently', () => {
const duplexStream = new Duplex({
objectMode: true,
read() {},
write(chunk, encoding, callback) {
this.push(chunk);
callback();
}
});
const pushDup = duplex(duplexStream);
const sourceOutput = [];
const sinkReceived = [];
pushDup.sink.pipe({
write(data) {
sinkReceived.push(data);
},
end() {},
get paused() { return false; },
set paused(v) {}
});
pushDup.source.pipe({
write(data) {
sourceOutput.push(data);
},
end() {},
get paused() { return false; },
set paused(v) {}
});
duplexStream.write('test');
duplexStream.push('response');
return new Promise((resolve) => {
setTimeout(() => {
if (sinkReceived.length === 1 && sourceOutput.length === 1) {
resolve();
} else {
resolve(new Error('Expected 1 item each'));
}
}, 50);
});
});
});
describe('error handling', () => {
test('source throws for non-Readable', () => {
try {
source({});
throw new Error('Expected error');
} catch (e) {
if (e.message !== 'Expected a Node.js Readable stream') {
throw new Error('Wrong error: ' + e.message);
}
}
});
test('sink throws for non-Writable', () => {
try {
sink({});
throw new Error('Expected error');
} catch (e) {
if (e.message !== 'Expected a Node.js Writable stream') {
throw new Error('Wrong error: ' + e.message);
}
}
});
test('duplex throws for non-Duplex', () => {
try {
duplex({});
throw new Error('Expected error');
} catch (e) {
if (e.message !== 'Expected a Node.js Duplex stream') {
throw new Error('Wrong error: ' + e.message);
}
}
});
});
});