213 lines
5.3 KiB
JavaScript
213 lines
5.3 KiB
JavaScript
|
|
import { fileURLToPath } from 'node:url';
|
||
|
|
import { dirname } from 'node:path';
|
||
|
|
import { test } from 'node:test';
|
||
|
|
|
||
|
|
const __filename = fileURLToPath(import.meta.url);
|
||
|
|
const __dirname = dirname(__filename);
|
||
|
|
const indexPath = __filename.replace('/test/index.js', '/index.js');
|
||
|
|
const indexUrl = `file://${indexPath}`;
|
||
|
|
|
||
|
|
let takeUntil;
|
||
|
|
|
||
|
|
test('setup', async () => {
|
||
|
|
const mod = await import(indexUrl);
|
||
|
|
takeUntil = mod.takeUntil || mod.default || mod;
|
||
|
|
});
|
||
|
|
|
||
|
|
function createAsyncSource(values, delay = 5) {
|
||
|
|
let index = 0;
|
||
|
|
return {
|
||
|
|
pipe(s) {
|
||
|
|
this.sink = s;
|
||
|
|
try { s.source = this; } catch (e) {}
|
||
|
|
this.resume();
|
||
|
|
return s;
|
||
|
|
},
|
||
|
|
resume() { this.writeNext(); },
|
||
|
|
writeNext() {
|
||
|
|
if (index >= values.length) {
|
||
|
|
if (this.sink && this.sink.end) {
|
||
|
|
this.sink.end();
|
||
|
|
}
|
||
|
|
return;
|
||
|
|
}
|
||
|
|
if (!this.sink) return;
|
||
|
|
this.sink.write(values[index++]);
|
||
|
|
setTimeout(() => this.writeNext(), delay);
|
||
|
|
},
|
||
|
|
write() {},
|
||
|
|
end() {},
|
||
|
|
abort() {},
|
||
|
|
get sink() { return this._sink; },
|
||
|
|
set sink(s) { this._sink = s; },
|
||
|
|
get paused() { return this._sink ? this._sink.paused : true; },
|
||
|
|
get ended() { return false; }
|
||
|
|
};
|
||
|
|
}
|
||
|
|
|
||
|
|
function createSignalStream(delay) {
|
||
|
|
return {
|
||
|
|
pipe(s) {
|
||
|
|
this.sink = s;
|
||
|
|
try { s.source = this; } catch (e) {}
|
||
|
|
if (delay === 0) {
|
||
|
|
if (this.sink && this.sink.end) {
|
||
|
|
this.sink.end();
|
||
|
|
}
|
||
|
|
} else {
|
||
|
|
setTimeout(() => {
|
||
|
|
if (this.sink && this.sink.end) {
|
||
|
|
this.sink.end();
|
||
|
|
}
|
||
|
|
}, delay);
|
||
|
|
}
|
||
|
|
return s;
|
||
|
|
},
|
||
|
|
resume() {},
|
||
|
|
write() {},
|
||
|
|
end() {},
|
||
|
|
abort() {},
|
||
|
|
get sink() { return this._sink; },
|
||
|
|
set sink(s) { this._sink = s; },
|
||
|
|
get paused() { return this._sink ? this._sink.paused : true; },
|
||
|
|
get ended() { return false; }
|
||
|
|
};
|
||
|
|
}
|
||
|
|
|
||
|
|
test('ends when signal emits', async () => {
|
||
|
|
const results = [];
|
||
|
|
let streamEnded = false;
|
||
|
|
let endTime = 0;
|
||
|
|
const startTime = Date.now();
|
||
|
|
|
||
|
|
const signalStream = createSignalStream(30);
|
||
|
|
const { sink, pipe } = takeUntil(signalStream);
|
||
|
|
|
||
|
|
const finalSink = {
|
||
|
|
write: (v) => results.push(v),
|
||
|
|
end: (err) => { streamEnded = true; endTime = Date.now() - startTime; },
|
||
|
|
get paused() { return false; }
|
||
|
|
};
|
||
|
|
|
||
|
|
pipe(finalSink);
|
||
|
|
createAsyncSource([1, 2, 3, 4, 5, 6, 7, 8, 9, 10], 5).pipe(sink);
|
||
|
|
|
||
|
|
let waited = 0;
|
||
|
|
while (!streamEnded && waited < 200) {
|
||
|
|
await new Promise(resolve => setTimeout(resolve, 10));
|
||
|
|
waited++;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (streamEnded === false) {
|
||
|
|
throw new Error('Stream should have ended');
|
||
|
|
}
|
||
|
|
if (endTime < 20) {
|
||
|
|
throw new Error(`Stream ended too quickly (${endTime}ms) - signal may not be working`);
|
||
|
|
}
|
||
|
|
if (endTime > 100) {
|
||
|
|
throw new Error(`Stream took too long to end (${endTime}ms)`);
|
||
|
|
}
|
||
|
|
if (results.length >= 10) {
|
||
|
|
throw new Error(`Expected fewer than 10 values (signal should have stopped it), got ${results.length}`);
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
test('signal immediately ends', async () => {
|
||
|
|
const results = [];
|
||
|
|
let streamEnded = false;
|
||
|
|
|
||
|
|
const signalStream = createSignalStream(0);
|
||
|
|
const { sink, pipe } = takeUntil(signalStream);
|
||
|
|
|
||
|
|
const finalSink = {
|
||
|
|
write: (v) => results.push(v),
|
||
|
|
end: () => { streamEnded = true; },
|
||
|
|
get paused() { return false; }
|
||
|
|
};
|
||
|
|
|
||
|
|
pipe(finalSink);
|
||
|
|
createAsyncSource([1, 2, 3, 4, 5], 5).pipe(sink);
|
||
|
|
|
||
|
|
let waited = 0;
|
||
|
|
while (!streamEnded && waited < 100) {
|
||
|
|
await new Promise(resolve => setTimeout(resolve, 10));
|
||
|
|
waited++;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (results.length !== 0) {
|
||
|
|
throw new Error(`Expected 0 values (immediate signal), got ${results.length}`);
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
test('signal never emits', async () => {
|
||
|
|
const results = [];
|
||
|
|
let streamEnded = false;
|
||
|
|
|
||
|
|
const signalStream = {
|
||
|
|
pipe(s) {
|
||
|
|
this.sink = s;
|
||
|
|
try { s.source = this; } catch (e) {}
|
||
|
|
return s;
|
||
|
|
},
|
||
|
|
resume() {},
|
||
|
|
write() {},
|
||
|
|
end() {},
|
||
|
|
abort() {},
|
||
|
|
get sink() { return this._sink; },
|
||
|
|
set sink(s) { this._sink = s; },
|
||
|
|
get paused() { return this._sink ? this._sink.paused : true; },
|
||
|
|
get ended() { return false; }
|
||
|
|
};
|
||
|
|
|
||
|
|
const takeUntilStream = takeUntil(signalStream);
|
||
|
|
|
||
|
|
const finalSink = {
|
||
|
|
write: (v) => results.push(v),
|
||
|
|
end: () => { streamEnded = true; },
|
||
|
|
get paused() { return false; }
|
||
|
|
};
|
||
|
|
|
||
|
|
takeUntilStream.pipe(finalSink);
|
||
|
|
createAsyncSource([1, 2, 3], 5).pipe(takeUntilStream);
|
||
|
|
|
||
|
|
let waited = 0;
|
||
|
|
while (!streamEnded && waited < 200) {
|
||
|
|
await new Promise(resolve => setTimeout(resolve, 10));
|
||
|
|
waited++;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (results.length !== 3) {
|
||
|
|
throw new Error(`Expected 3 values, got ${results.length}`);
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
test('passes through all values before signal', async () => {
|
||
|
|
const results = [];
|
||
|
|
let streamEnded = false;
|
||
|
|
|
||
|
|
const signalStream = createSignalStream(100);
|
||
|
|
const takeUntilStream = takeUntil(signalStream);
|
||
|
|
|
||
|
|
const finalSink = {
|
||
|
|
write: (v) => results.push(v),
|
||
|
|
end: () => { streamEnded = true; },
|
||
|
|
get paused() { return false; }
|
||
|
|
};
|
||
|
|
|
||
|
|
takeUntilStream.pipe(finalSink);
|
||
|
|
createAsyncSource([1, 2], 5).pipe(takeUntilStream);
|
||
|
|
|
||
|
|
let waited = 0;
|
||
|
|
while (!streamEnded && waited < 200) {
|
||
|
|
await new Promise(resolve => setTimeout(resolve, 10));
|
||
|
|
waited++;
|
||
|
|
}
|
||
|
|
|
||
|
|
if (results.length !== 2) {
|
||
|
|
throw new Error(`Expected 2 values before signal, got ${results.length}`);
|
||
|
|
}
|
||
|
|
if (results[0] !== 1 || results[1] !== 2) {
|
||
|
|
throw new Error(`Values wrong: ${JSON.stringify(results)}`);
|
||
|
|
}
|
||
|
|
});
|