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); } } }); }); });