commit 10f2d98ada05ffb4626d79a7e8982e8be7ce0924 Author: Tenere CI Date: Sun Jul 19 13:44:10 2026 -0700 initial import diff --git a/.gitea/workflows/ci.yaml b/.gitea/workflows/ci.yaml new file mode 100644 index 0000000..0ca1503 --- /dev/null +++ b/.gitea/workflows/ci.yaml @@ -0,0 +1,28 @@ +name: CI + +on: + push: + branches: [master] + tags: ['v*'] + pull_request: + branches: [master] + +jobs: + test: + runs-on: ubuntu-latest + strategy: + matrix: + node: [18, 20, 22] + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-node@v4 + with: + node-version: ${{ matrix.node }} + - name: Auth for Gitea npm registry + run: | + echo "@${{ github.repository_owner }}:registry=https://hub.kl1.tenere.ai/api/packages/${{ github.repository_owner }}/npm/" > .npmrc + echo "//hub.kl1.tenere.ai/api/packages/${{ github.repository_owner }}/npm/:_authToken=${{ secrets.PACKAGE_TOKEN }}" >> .npmrc + - run: npm install + - run: npm test + - run: npm run bench + if: matrix.node == 22 diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..d570088 --- /dev/null +++ b/.gitignore @@ -0,0 +1,2 @@ +node_modules/ + diff --git a/.npmrc b/.npmrc new file mode 100644 index 0000000..a23de45 --- /dev/null +++ b/.npmrc @@ -0,0 +1,2 @@ +@push-stream-std:registry=https://hub.kl1.tenere.ai/api/packages/push-stream-std/npm/ +//hub.kl1.tenere.ai/api/packages/push-stream-std/npm/:_authToken=${PACKAGE_TOKEN} diff --git a/LICENSE b/LICENSE new file mode 100644 index 0000000..e5bd26f --- /dev/null +++ b/LICENSE @@ -0,0 +1,4 @@ +Tenere Labs Inc. — All Rights Reserved. + +PROPRIETARY AND CONFIDENTIAL. +Source-available for auditing and inspection only. diff --git a/README.md b/README.md new file mode 100644 index 0000000..c053753 --- /dev/null +++ b/README.md @@ -0,0 +1,58 @@ +# push-take-until + +Take values from a source stream until a signal stream emits or ends. + +Published as `@push-stream-std/push-take-until`. + +## Synopsis + +```js +import takeUntil from '@push-stream-std/push-take-until'; + +const signal = pushable(); +const untilSignal = takeUntil(signal); + +values([1, 2, 3, 4, 5]).pipe(untilSignal).pipe(collect(function (err, result) { + console.log(result); // [1, 2] — signal ended after 2 values +})); + +// End the signal stream to stop the outer stream +signal.end(); +``` + +## Install + +```bash +npm config set @push-stream-std:registry https://hub.kl1.tenere.ai/api/packages/push-stream-std/npm/ +npm config set -- //hub.kl1.tenere.ai/api/packages/push-stream-std/npm/:_authToken "$PACKAGE_TOKEN" +npm install @push-stream-std/push-take-until +``` + +## Description + +A push-stream through that passes through values from the source until the provided signal stream ends. When the signal stream ends (either normally or via abort), the take-until stream immediately ends the upstream source and downstream sink, closing the pipeline. This is useful for implementing cancellation patterns — for example, taking values from an event stream until a "stop" button is clicked. + +## API + +### `takeUntil(signalStream)` + +Creates a take-until through stream. + +- **signalStream** `object` – A push-stream source. The take-until stream pipes the signal stream to an internal sink that tracks when the signal ends. When the signal's `end()` or `abort()` fires, the take-until stream ends. +- **Returns** A through object with `pipe()`, `write()`, `end()`, `abort()`, `resume()`, `source`, `paused`, `ended`. + +## Tests + +```bash +npm test +``` + +Tests cover passthrough before signal ends, stopping after signal end, normal end without signal, signal abort propagation, and idempotent signal end handling. + +## Requirements + +Node.js 18 or later. ESM only. + +## License + +MIT diff --git a/bench/benchlib.js b/bench/benchlib.js new file mode 100644 index 0000000..8cd139e --- /dev/null +++ b/bench/benchlib.js @@ -0,0 +1,99 @@ +import fs from 'fs'; + +const results = []; + +async function bench(name, fn, options = {}) { + const iterations = options.iterations || 1; + const warmup = options.warmup || 0; + const unit = options.unit || 'ops/s'; + const itemsPerIter = options.itemsPerIter || null; + + for (let i = 0; i < warmup; i++) { + await fn(); + } + + let totalTime = 0; + for (let i = 0; i < iterations; i++) { + const start = process.hrtime.bigint(); + await fn(); + const end = process.hrtime.bigint(); + totalTime += Number(end - start) / 1e6; + } + + const avgTime = totalTime / iterations; + const iterPerSec = 1000 / avgTime; + const itemsPerSec = itemsPerIter ? Math.round(iterPerSec * itemsPerIter) : null; + + const resultEntry = { + name, + avgTimeMs: parseFloat(avgTime.toFixed(3)), + iterPerSec: Math.round(iterPerSec), + itemsPerSec, + unit, + itemsPerIter, + samples: iterations, + timestamp: new Date().toISOString() + }; + + results.push(resultEntry); + + return resultEntry; +} + +async function run(options = {}) { + const outputFile = options.output || './benchmark-results.json'; + const threshold = options.regressionThreshold || 0.9; + + const output = { + timestamp: new Date().toISOString(), + package: process.env.npm_package_name || 'unknown', + version: process.env.npm_package_version || '0.0.0', + node: process.version, + benchmarks: results + }; + + if (fs.existsSync(outputFile)) { + const previous = JSON.parse(fs.readFileSync(outputFile, 'utf8')); + if (previous.benchmarks) { + for (const newBench of results) { + const prev = previous.benchmarks.find(b => b.name === newBench.name); + if (prev) { + const ratio = newBench.iterPerSec / prev.iterPerSec; + newBench.regression = ratio < threshold; + newBench.comparedToPrevious = { + iterPerSec: prev.iterPerSec, + ratio: parseFloat(ratio.toFixed(4)) + }; + } + } + } + } + + const regression = results.some(b => b.regression); + + console.log('\nBenchmark Results:'); + console.log('-'.repeat(90)); + console.log(`${'Benchmark'.padEnd(50)} ${'Iter/s'.padStart(10)} ${'Items/s'.padStart(12)} Unit`); + console.log('-'.repeat(90)); + for (const r of results) { + const reg = r.regression ? ' [REGRESSION]' : ''; + const prev = r.comparedToPrevious ? ` (prev: ${r.comparedToPrevious.iterPerSec}, ratio: ${r.comparedToPrevious.ratio})` : ''; + const itemsRate = r.itemsPerSec ? r.itemsPerSec.toString().padStart(12) : ' -'; + console.log(`${r.name.padEnd(50)} ${r.iterPerSec.toString().padStart(10)} ${itemsRate} ${r.unit}${reg}${prev}`); + } + console.log('-'.repeat(90)); + + if (outputFile) { + fs.writeFileSync(outputFile, JSON.stringify(output, null, 2)); + console.log(`\nResults saved to: ${outputFile}`); + } + + if (regression) { + console.error('\nWARNING: Performance regression detected!'); + process.exit(1); + } + + return results; +} + +export { bench, run }; diff --git a/bench/index.js b/bench/index.js new file mode 100644 index 0000000..f0a2021 --- /dev/null +++ b/bench/index.js @@ -0,0 +1,84 @@ +import { bench, run } from './benchlib.js'; +import takeUntil from '../index.js'; + +function createSource() { + let index = 0; + let sink = null; + return { + pipe(s) { sink = s; return s; }, + resume() { + while (sink && !sink.paused && index < 10000) { + sink.write(index++); + } + }, + write() {}, + end() {}, + abort() {}, + get sink() { return sink; }, + get paused() { return sink ? sink.paused : true; } + }; +} + +function createCollectSink() { + const items = []; + return { + items, + write(data) { items.push(data); }, + end() {}, + abort() {}, + paused: false, + ended: false + }; +} + +function createSignalStream() { + let sink = null; + return { + pipe(s) { sink = s; return s; }, + write() {}, + end() { if (sink) sink.end(); }, + abort() {}, + get paused() { return false; }, + trigger() { if (sink) sink.end(); } + }; +} + +async function benchmark() { + await bench('takeUntil with early signal (10000 items)', () => { + const sink = createCollectSink(); + const signal = createSignalStream(); + const stream = takeUntil(signal); + stream.pipe(sink); + const source = createSource(); + source.pipe(stream); + source.resume(); + signal.trigger(); + return sink.items.length; + }, { iterations: 20, warmup: 2, itemsPerIter: 10000, unit: 'items/s' }); + + await bench('takeUntil with late signal (10000 items)', () => { + const sink = createCollectSink(); + const signal = createSignalStream(); + const stream = takeUntil(signal); + stream.pipe(sink); + const source = createSource(); + source.pipe(stream); + source.resume(); + return sink.items.length; + }, { iterations: 20, warmup: 2, itemsPerIter: 10000, unit: 'items/s' }); + + await bench('takeUntil immediate signal (0 items)', () => { + const sink = createCollectSink(); + const signal = createSignalStream(); + const stream = takeUntil(signal); + stream.pipe(sink); + const source = createSource(); + source.pipe(stream); + signal.trigger(); + return sink.items.length; + }, { iterations: 100, warmup: 5 }); + + await run({ output: './benchmark-results.json' }); +} + +benchmark().catch(console.error); diff --git a/index.js b/index.js new file mode 100644 index 0000000..7673ca6 --- /dev/null +++ b/index.js @@ -0,0 +1,94 @@ +/** + * push-take-until - Take elements until signal stream emits + * + * Uses canonical push-stream pattern: + * - This is a sink that wraps a signal stream + * - When signal ends, the take-until ends + * - Pass through all values until signal fires + * + * API: + * import takeUntil from 'push-take-until' + * + * const untilClick = takeUntil(clickStream) + * source.pipe(untilClick).pipe(sink) + */ + +export function takeUntil(signalStream) { + let source = null; + let sink = null; + let isEnded = false; + let signalEnded = false; + + const stream = { + get paused() { return false; }, + get ended() { return isEnded; }, + + pipe(s) { + if (isEnded) { + s.end(); + return s; + } + sink = s; + try { s.source = stream; } catch (e) {} + return s; + }, + + write(data) { + if (isEnded || signalEnded) return; + if (sink && !sink.paused) { + sink.write(data); + } + }, + + end(err) { + if (isEnded) return; + isEnded = true; + if (source && source.end) source.end(err); + if (sink) sink.end(err); + }, + + abort(err) { + if (isEnded) return; + isEnded = true; + if (source && source.abort) source.abort(err); + if (sink) sink.end(err); + }, + + resume() { + if (source && source.resume) source.resume(); + }, + + get source() { return source; }, + set source(s) { source = s; } + }; + + signalStream.pipe({ + write() {}, + end() { + if (signalEnded) return; + signalEnded = true; + + if (!isEnded) { + isEnded = true; + if (source && source.end) source.end(); + if (sink) sink.end(); + } + }, + abort(err) { + if (signalEnded) return; + signalEnded = true; + + if (!isEnded) { + isEnded = true; + if (source && source.abort) source.abort(err); + if (sink) sink.end(err); + } + }, + get paused() { return false; }, + get ended() { return signalEnded; } + }); + + return stream; +} + +export default takeUntil; \ No newline at end of file diff --git a/package.json b/package.json new file mode 100644 index 0000000..9bab939 --- /dev/null +++ b/package.json @@ -0,0 +1,36 @@ +{ + "name": "@push-stream-std/push-take-until", + "version": "1.0.0", + "description": "Take elements until signal stream emits", + "type": "module", + "main": "index.js", + "exports": { + ".": "./index.js" + }, + "scripts": { + "test": "node --test test/", + "bench": "node bench/index.js", + "lint": "node --check index.js test/*.js" + }, + "keywords": [ + "push-stream", + "stream", + "take", + "until", + "signal", + "cancel" + ], + "license": "SEE LICENSE IN LICENSE", + "engines": { + "node": ">=22.0.0" + }, + "files": [ + "index.js", + "bench/", + "README.md", + "LICENSE" + ], + "publishConfig": { + "registry": "https://hub.kl1.tenere.ai/api/packages/push-stream-std/npm/" + } +} diff --git a/test/index.js b/test/index.js new file mode 100644 index 0000000..c4839bf --- /dev/null +++ b/test/index.js @@ -0,0 +1,212 @@ +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)}`); + } +});