From 62019a83c78e64078113d5ea9bf62ad32f42eda1 Mon Sep 17 00:00:00 2001
From: "copilot-swe-agent[bot]" <198982749+Copilot@users.noreply.github.com>
Date: Wed, 18 Mar 2026 03:10:45 +0000
Subject: [PATCH] build(publish): address publishing review feedback
Co-authored-by: okikio <17222836+okikio@users.noreply.github.com>
---
.github/workflows/publish.yml | 4 +-
README.md | 277 ++---
bench/events_bench.ts | 46 +-
bench/latency_bench.ts | 23 +-
bench/memory_bench.ts | 145 +--
bench/observable_bench.ts | 51 +-
bench/operators_bench.ts | 120 +--
bench/queue_bench.ts | 90 +-
deno.jsonc | 8 +-
deno.lock | 109 ++
events.ts | 112 ++-
helpers/_types.ts | 120 ++-
helpers/operations/combination.ts | 56 +-
helpers/operators.ts | 501 ++++++---
helpers/pipe.ts | 341 +++++--
observable.ts | 1074 +++++++++++---------
scripts/build_npm.ts | 150 +--
tests/_utils/_assert.ts | 5 +-
tests/events/events_test.ts | 1230 ++++++++++++-----------
tests/events_bdd_test.ts | 382 ++++---
tests/helpers/operations/timing_test.ts | 31 +-
tests/helpers/operators_bdd_test.ts | 451 +++++----
tests/helpers/utils_bdd_test.ts | 387 +++----
tests/integration_bdd_test.ts | 283 +++---
tests/publishing_setup_test.ts | 56 +-
tests/queue_bdd_test.ts | 419 ++++----
26 files changed, 3641 insertions(+), 2830 deletions(-)
create mode 100644 deno.lock
diff --git a/.github/workflows/publish.yml b/.github/workflows/publish.yml
index bde4f92..1092bb6 100644
--- a/.github/workflows/publish.yml
+++ b/.github/workflows/publish.yml
@@ -40,8 +40,8 @@ jobs:
- uses: actions/setup-node@v6
with:
- node-version: '24'
- registry-url: 'https://registry.npmjs.org'
+ node-version: "24"
+ registry-url: "https://registry.npmjs.org"
- name: Publish to npm
# npm trusted publishing uses the GitHub OIDC token minted from
diff --git a/README.md b/README.md
index 3c3e583..43a477a 100644
--- a/README.md
+++ b/README.md
@@ -2,42 +2,56 @@
[](https://bundlejs.com/?q=@okikio/observables&bundle "Check the total bundle size of @okikio/observables")
-[NPM](https://www.npmjs.com/package/@okikio/observables) | [GitHub](https://github.com/okikio/observables#readme) | [JSR](https://jsr.io/@okikio/observables) | [Licence](./LICENSE)
-
-A **spec-faithful** yet ergonomic TC39-inspired Observable implementation that gives you one consistent way to handle all async data in JavaScript.
-
-**Observables** are a **push‑based stream abstraction** for events, data, and long‑running operations. Think of them as a **multi‑value Promise** that keeps sending values until you tell it to stop, where a Promise gives you one value eventually, an Observable can give you many values over time: mouse clicks, search results, chat messages, sensor readings.
+[NPM](https://www.npmjs.com/package/@okikio/observables)
+|
+[GitHub](https://github.com/okikio/observables#readme)
+|
+[JSR](https://jsr.io/@okikio/observables)
+| [Licence](./LICENSE)
+
+A **spec-faithful** yet ergonomic TC39-inspired Observable implementation that
+gives you one consistent way to handle all async data in JavaScript.
+
+**Observables** are a **push‑based stream abstraction** for events, data, and
+long‑running operations. Think of them as a **multi‑value Promise** that keeps
+sending values until you tell it to stop, where a Promise gives you one value
+eventually, an Observable can give you many values over time: mouse clicks,
+search results, chat messages, sensor readings.
[](https://bundlejs.com/?q=@okikio/observables&treeshake=[{+Observable,+pipe,+map,+filter+}])
-If you've ever built a web app, you know this all too well: user clicks, API responses, WebSocket messages, timers, file uploads, they all arrive at different times and need different handling. Before Observables, we'd all end up with a mess of callbacks, Promise chains, event listeners, and async/await scattered throughout our code.
+If you've ever built a web app, you know this all too well: user clicks, API
+responses, WebSocket messages, timers, file uploads, they all arrive at
+different times and need different handling. Before Observables, we'd all end up
+with a mess of callbacks, Promise chains, event listeners, and async/await
+scattered throughout our code.
-Let's say you're building a search box. You've probably written something like this:
+Let's say you're building a search box. You've probably written something like
+this:
```ts
// We've all been here: callbacks, timers, and manual cleanup 😫
let searchTimeout: number;
let lastRequest: Promise | null = null;
-searchInput.addEventListener('input', async (event) => {
+searchInput.addEventListener("input", async (event) => {
const query = event.target.value;
-
+
// Debounce: wait 300ms after user stops typing
clearTimeout(searchTimeout);
searchTimeout = setTimeout(async () => {
-
// Cancel previous request somehow?
if (lastRequest) {
// How do you cancel a fetch? 🤔
}
-
+
if (query.length < 3) return; // Skip short queries
-
+
try {
lastRequest = fetch(`/search?q=${query}`);
const response = await lastRequest;
const results = await response.json();
-
+
// Update UI, but what if user already typed something new?
updateSearchResults(results);
} catch (error) {
@@ -51,13 +65,18 @@ searchInput.addEventListener('input', async (event) => {
// (Spoiler: we all forget this and create memory leaks)
```
-This works but it's fragile, hard to test, and easy to mess up. Plus, you have to remember to clean up event listeners, cancel timers, and handle edge cases manually.
+This works but it's fragile, hard to test, and easy to mess up. Plus, you have
+to remember to clean up event listeners, cancel timers, and handle edge cases
+manually.
We've all felt this pain before:
-- **Memory Leaks**: Forgot to remove an event listener? Your app slowly eats memory
-- **Race Conditions**: User clicks fast, requests arrive out of order, wrong results appear
-- **Error Handling**: Network failed? Now you need custom backoff and error recovery
+- **Memory Leaks**: Forgot to remove an event listener? Your app slowly eats
+ memory
+- **Race Conditions**: User clicks fast, requests arrive out of order, wrong
+ results appear
+- **Error Handling**: Network failed? Now you need custom backoff and error
+ recovery
- **Backpressure**: Producer too fast for consumer? Memory bloats until crash
- **Testing**: Complex async flows become nearly impossible to test reliably
- **Maintenance**: Each async pattern needs its own cleanup and error handling
@@ -66,32 +85,35 @@ Here's the same search box with Observables:
```ts
// Much cleaner: composable and robust ✨
-import { pipe, debounce, filter, switchMap, map } from "@okikio/observables";
+import { debounce, filter, map, pipe, switchMap } from "@okikio/observables";
const searchResults = pipe(
- inputEvents, // Stream of input events
- debounce(300), // Wait 300ms after user stops typing
- filter(query => query.length >= 3), // Skip short queries
- switchMap(query => // Cancel previous requests automatically
+ inputEvents, // Stream of input events
+ debounce(300), // Wait 300ms after user stops typing
+ filter((query) => query.length >= 3), // Skip short queries
+ switchMap((query) =>
+ // Cancel previous requests automatically
Observable.from(fetch(`/search?q=${query}`))
),
- map(response => response.json()) // Parse response
+ map((response) => response.json()), // Parse response
);
// Subscribe to results (with automatic cleanup!)
using subscription = searchResults.subscribe({
- next: results => updateSearchResults(results),
- error: error => handleSearchError(error)
+ next: (results) => updateSearchResults(results),
+ error: (error) => handleSearchError(error),
});
// Subscription automatically cleaned up when leaving scope
```
-Notice the difference? No manual timers, no cancellation logic, no memory leaks. The operators handle all the complex async coordination for you.
+Notice the difference? No manual timers, no cancellation logic, no memory leaks.
+The operators handle all the complex async coordination for you.
-This library was built by developers who've felt these same frustrations. It focuses on:
+This library was built by developers who've felt these same frustrations. It
+focuses on:
- **Familiarity**: If you know `Array.map()`, you already understand operators
-- **Performance**: Built on Web Streams with pre-compiled error handling
+- **Performance**: Built on Web Streams with pre-compiled error handling
- **Type Safety**: Full TypeScript support with intelligent inference
- **Standards**: Follows the TC39 Observable proposal for future compatibility
- **Practicality**: <4KB but includes everything you need for real apps
@@ -102,7 +124,7 @@ This library was built by developers who've felt these same frustrations. It foc
### Deno
```ts
-import { Observable, pipe, map } from "jsr:@okikio/observables";
+import { map, Observable, pipe } from "jsr:@okikio/observables";
```
Or
@@ -133,7 +155,7 @@ npx jsr add @okikio/observables
pnpm add jsr:@okikio/observables
```
-Or
+Or
```bash
yarn add @okikio/observables@jsr:latest
@@ -151,34 +173,34 @@ bunx jsr add @okikio/observables
You can also use it via a CDN:
-```ts
-import { Observable, pipe, map } from "https://esm.sh/jsr/@okikio/observables";
+```ts ignore
+import { map, Observable, pipe } from "https://esm.sh/jsr/@okikio/observables";
```
## Quick Start
```ts
-import { Observable, pipe, map, filter, debounce } from "@okikio/observables";
+import { debounce, filter, map, Observable, pipe } from "@okikio/observables";
// Create from anything async
-const clicks = new Observable(observer => {
- const handler = e => observer.next(e);
- button.addEventListener('click', handler);
- return () => button.removeEventListener('click', handler);
+const clicks = new Observable((observer) => {
+ const handler = (e) => observer.next(e);
+ button.addEventListener("click", handler);
+ return () => button.removeEventListener("click", handler);
});
// Transform with operators (like Array.map, but for async data)
const doubleClicks = pipe(
clicks,
- debounce(300), // Wait 300ms between clicks
- filter((_, index) => index % 2), // Only odd-numbered clicks
- map(event => ({ x: event.clientX, y: event.clientY }))
+ debounce(300), // Wait 300ms between clicks
+ filter((_, index) => index % 2), // Only odd-numbered clicks
+ map((event) => ({ x: event.clientX, y: event.clientY })),
);
// Subscribe to results
using subscription = doubleClicks.subscribe({
- next: coords => console.log('Double click at:', coords),
- error: err => console.error('Error:', err)
+ next: (coords) => console.log("Double click at:", coords),
+ error: (err) => console.error("Error:", err),
});
// Automatically cleaned up when leaving scope
```
@@ -191,7 +213,8 @@ A couple sites/projects that use `@okikio/observables`:
## API
-The API of `@okikio/observables` provides everything you need for reactive programming:
+The API of `@okikio/observables` provides everything you need for reactive
+programming:
### Core Observable
@@ -199,40 +222,40 @@ The API of `@okikio/observables` provides everything you need for reactive progr
import { Observable } from "@okikio/observables";
// Create observables
-const timer = new Observable(observer => {
+const timer = new Observable((observer) => {
const id = setInterval(() => observer.next(Date.now()), 1000);
return () => clearInterval(id);
});
// Factory methods
-Observable.of(1, 2, 3); // From values
-Observable.from(fetch('/api/data')); // From promises/iterables
+Observable.of(1, 2, 3); // From values
+Observable.from(fetch("/api/data")); // From promises/iterables
```
### Operators (19+ included)
```ts
-import { pipe, map, filter, debounce, switchMap } from "@okikio/observables";
+import { debounce, filter, map, pipe, switchMap } from "@okikio/observables";
// Transform data as it flows
pipe(
source,
- map(x => x * 2), // Transform each value
- filter(x => x > 10), // Keep only values > 10
- debounce(300), // Wait for quiet periods
- switchMap(x => fetchData(x)) // Cancel previous requests
+ map((x) => x * 2), // Transform each value
+ filter((x) => x > 10), // Keep only values > 10
+ debounce(300), // Wait for quiet periods
+ switchMap((x) => fetchData(x)), // Cancel previous requests
);
```
### EventBus & EventDispatcher
```ts
-import { EventBus, createEventDispatcher } from "@okikio/observables";
+import { createEventDispatcher, EventBus } from "@okikio/observables";
// Simple pub/sub
const bus = new EventBus();
-bus.events.subscribe(msg => console.log(msg));
-bus.emit('Hello world!');
+bus.events.subscribe((msg) => console.log(msg));
+bus.emit("Hello world!");
// Type-safe events
interface AppEvents {
@@ -241,8 +264,8 @@ interface AppEvents {
}
const events = createEventDispatcher();
-events.emit('userLogin', { userId: '123' });
-events.on('cartUpdate', data => updateUI(data.items));
+events.emit("userLogin", { userId: "123" });
+events.on("cartUpdate", (data) => updateUI(data.items));
```
### Error Handling (4 modes)
@@ -252,14 +275,14 @@ import { createOperator } from "@okikio/observables";
// Choose your error handling strategy
const processor = createOperator({
- errorMode: 'pass-through', // Errors become values (default)
+ errorMode: "pass-through", // Errors become values (default)
// errorMode: 'ignore', // Skip errors silently
// errorMode: 'throw', // Fail fast
// errorMode: 'manual', // You handle everything
-
+
transform(value, controller) {
controller.enqueue(processValue(value));
- }
+ },
});
```
@@ -283,55 +306,66 @@ async function example() {
```ts
// Process large datasets with backpressure
-for await (const chunk of bigDataStream.pull({
- strategy: { highWaterMark: 8 } // Small buffer for large files
-})) {
+for await (
+ const chunk of bigDataStream.pull({
+ strategy: { highWaterMark: 8 }, // Small buffer for large files
+ })
+) {
await processChunk(chunk);
}
```
-Look through the [tests/](./tests/) and [bench/](./bench/) folders for complex examples and multiple usage patterns.
+Look through the [tests/](./tests/) and [bench/](./bench/) folders for complex
+examples and multiple usage patterns.
## Advanced Usage
### Smart Search with Cancellation
```ts
-import { pipe, debounce, filter, switchMap, map, catchErrors } from "@okikio/observables";
+import {
+ catchErrors,
+ debounce,
+ filter,
+ map,
+ pipe,
+ switchMap,
+} from "@okikio/observables";
const searchResults = pipe(
searchInput,
- debounce(300), // Wait for typing pause
- filter(query => query.length > 2), // Skip short queries
- switchMap(query => // Cancel old requests automatically
+ debounce(300), // Wait for typing pause
+ filter((query) => query.length > 2), // Skip short queries
+ switchMap((query) =>
+ // Cancel old requests automatically
pipe(
Observable.from(fetch(`/search?q=${query}`)),
- map(res => res.json()),
- catchErrors([]) // Return empty array on error
+ map((res) => res.json()),
+ catchErrors([]), // Return empty array on error
)
- )
+ ),
);
-searchResults.subscribe(results => updateUI(results));
+searchResults.subscribe((results) => updateUI(results));
```
### Real-Time Dashboard
```ts
-import { pipe, filter, scan, throttle } from "@okikio/observables";
+import { filter, pipe, scan, throttle } from "@okikio/observables";
const dashboardData = pipe(
webSocketEvents,
- filter(event => event.type === 'metric'), // Only metric events
- scan((acc, event) => ({ // Build running totals
+ filter((event) => event.type === "metric"), // Only metric events
+ scan((acc, event) => ({ // Build running totals
total: acc.total + event.value,
count: acc.count + 1,
- average: (acc.total + event.value) / (acc.count + 1)
+ average: (acc.total + event.value) / (acc.count + 1),
}), { total: 0, count: 0, average: 0 }),
- throttle(1000) // Update UI max once per second
+ throttle(1000), // Update UI max once per second
);
-dashboardData.subscribe(stats => updateDashboard(stats));
+dashboardData.subscribe((stats) => updateDashboard(stats));
```
### Custom Operators
@@ -342,102 +376,116 @@ import { createOperator, createStatefulOperator } from "@okikio/observables";
// Simple transformation
function double() {
return createOperator({
- name: 'double',
+ name: "double",
transform(value, controller) {
controller.enqueue(value * 2);
- }
+ },
});
}
// Stateful operation
function movingAverage(windowSize: number) {
return createStatefulOperator({
- name: 'movingAverage',
+ name: "movingAverage",
createState: () => [],
-
+
transform(value, arr, controller) {
arr.push(value);
if (arr.length > windowSize) arr.shift();
-
+
const avg = arr.reduce((sum, n) => sum + n, 0) / arr.length;
controller.enqueue(avg);
- }
+ },
});
}
```
## Performance
-We built this on Web Streams for good reason, native backpressure and memory efficiency come for free. Here's what you get:
+We built this on Web Streams for good reason, native backpressure and memory
+efficiency come for free. Here's what you get:
-- **Web Streams Foundation**: Handles backpressure automatically, no memory bloat
-- **Pre-compiled Error Modes**: Skip runtime checks in hot paths
+- **Web Streams Foundation**: Handles backpressure automatically, no memory
+ bloat
+- **Pre-compiled Error Modes**: Skip runtime checks in hot paths
- **Tree Shaking**: Import only what you use (most apps need <4KB)
- **TypeScript Native**: Zero runtime overhead for type safety
Performance varies by use case, but here's how different error modes stack up:
-| Error Mode | Performance | When We Use It |
-|------------|-------------|-----------------|
-| `manual` | Fastest | Hot paths, custom logic |
-| `ignore` | Very fast | Filtering bad data |
-| `pass-through` | Fast | Error recovery, debugging |
-| `throw` | Good | Fail-fast validation |
+| Error Mode | Performance | When We Use It |
+| -------------- | ----------- | ------------------------- |
+| `manual` | Fastest | Hot paths, custom logic |
+| `ignore` | Very fast | Filtering bad data |
+| `pass-through` | Fast | Error recovery, debugging |
+| `throw` | Good | Fail-fast validation |
## Comparison
-| Feature | @okikio/observables | RxJS | zen-observable |
-|---------|-------------------|------|----------------|
-| Bundle Size | <4KB | ~35KB | ~2KB |
-| Operators | 19+ | 100+ | 5 |
-| Error Modes | 4 modes | 1 mode | 1 mode |
-| EventBus | ✅ Built-in | ❌ Separate | ❌ None |
-| TC39 Compliance | ✅ Yes | ⚠️ Partial | ✅ Yes |
-| TypeScript | ✅ Native | ✅ Yes | ⚠️ Basic |
-| Tree Shaking | ✅ Perfect | ⚠️ Partial | ✅ Yes |
-| Learning Curve | 🟢 Gentle | 🔴 Steep | 🟢 Gentle |
+| Feature | @okikio/observables | RxJS | zen-observable |
+| --------------- | ------------------- | ----------- | -------------- |
+| Bundle Size | <4KB | ~35KB | ~2KB |
+| Operators | 19+ | 100+ | 5 |
+| Error Modes | 4 modes | 1 mode | 1 mode |
+| EventBus | ✅ Built-in | ❌ Separate | ❌ None |
+| TC39 Compliance | ✅ Yes | ⚠️ Partial | ✅ Yes |
+| TypeScript | ✅ Native | ✅ Yes | ⚠️ Basic |
+| Tree Shaking | ✅ Perfect | ⚠️ Partial | ✅ Yes |
+| Learning Curve | 🟢 Gentle | 🔴 Steep | 🟢 Gentle |
## Browser Support
-| Chrome | Edge | Firefox | Safari | Node | Deno | Bun |
-| ------ | ---- | ------- | ------ | ---- | ---- | --- |
+| Chrome | Edge | Firefox | Safari | Node | Deno | Bun |
+| ------ | ---- | ------- | ------ | ---- | ---- | ---- |
| 80+ | 80+ | 72+ | 13+ | 16+ | 1.0+ | 1.0+ |
-> Native support for Observables is excellent. Some advanced features like `Symbol.dispose` require newer environments or polyfills.
+> Native support for Observables is excellent. Some advanced features like
+> `Symbol.dispose` require newer environments or polyfills.
## FAQ
### What are Observables exactly?
-Think of them as Promises that can send multiple values over time. Where a Promise gives you one result eventually, an Observable can keep sending values, like a stream of search results, mouse movements, or WebSocket messages.
+Think of them as Promises that can send multiple values over time. Where a
+Promise gives you one result eventually, an Observable can keep sending values,
+like a stream of search results, mouse movements, or WebSocket messages.
### Why not just use RxJS?
-RxJS is powerful but can be overwhelming. We've all been there, 100+ operators, steep learning curve, 35KB bundle size. This library gives you the essential Observable patterns you actually use day-to-day, following the TC39 proposal so you're future-ready.
+RxJS is powerful but can be overwhelming. We've all been there, 100+ operators,
+steep learning curve, 35KB bundle size. This library gives you the essential
+Observable patterns you actually use day-to-day, following the TC39 proposal so
+you're future-ready.
### EventBus vs Observable, when do I use which?
Good question! Here's how we think about it:
-- **Observable**: When you're transforming data one-to-one (API calls, processing user input)
-- **EventBus**: When you need one-to-many communication (notifications, cross-component events)
+- **Observable**: When you're transforming data one-to-one (API calls,
+ processing user input)
+- **EventBus**: When you need one-to-many communication (notifications,
+ cross-component events)
### How should I handle errors?
Pick the mode that fits your situation:
- **`pass-through`**: Errors become values you can recover from
-- **`ignore`**: Skip errors silently (great for filtering noisy data)
+- **`ignore`**: Skip errors silently (great for filtering noisy data)
- **`throw`**: Fail fast for validation
- **`manual`**: Handle everything yourself
### Is this actually production ready?
-We use it in production. It follows the TC39 proposal, has comprehensive tests, and handles resource management properly. The Web Streams foundation is battle-tested across browsers and runtimes.
+We use it in production. It follows the TC39 proposal, has comprehensive tests,
+and handles resource management properly. The Web Streams foundation is
+battle-tested across browsers and runtimes.
## Contributing
-I encourage you to use [deno](https://deno.com/) to contribute to this repo, to setup deno you can install it via [mise](https://mise.jdx.dev/) or [manually](https://deno.land/manual/getting_started/installation).
+I encourage you to use [deno](https://deno.com/) to contribute to this repo, to
+setup deno you can install it via [mise](https://mise.jdx.dev/) or
+[manually](https://deno.land/manual/getting_started/installation).
Setup Mise:
@@ -464,7 +512,10 @@ Run benchmarks:
deno task bench
```
-> **Note**: This project uses [Conventional Commits](https://www.conventionalcommits.org/en/v1.0.0/) standard for commits, so please format your commits using the rules it sets out.
+> **Note**: This project uses
+> [Conventional Commits](https://www.conventionalcommits.org/en/v1.0.0/)
+> standard for commits, so please format your commits using the rules it sets
+> out.
## Licence
diff --git a/bench/events_bench.ts b/bench/events_bench.ts
index 6504cfd..46df732 100644
--- a/bench/events_bench.ts
+++ b/bench/events_bench.ts
@@ -1,3 +1,4 @@
+// deno-lint-ignore-file no-import-prefix
/**
* Event system benchmarks.
*
@@ -7,9 +8,9 @@
* subscription setup on every iteration.
*/
-import { bench, do_not_optimize, run } from 'npm:mitata';
+import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34";
-import { EventBus, createEventDispatcher } from '../events.ts';
+import { createEventDispatcher, EventBus } from "../events.ts";
let singleSubscriberValue = 0;
const singleSubscriberBus = new EventBus();
@@ -19,18 +20,22 @@ const singleSubscriberSub = singleSubscriberBus.subscribe((value) => {
let tenSubscriberValue = 0;
const tenSubscriberBus = new EventBus();
-const tenSubscriberSubs = Array.from({ length: 10 }, (_, index) =>
- tenSubscriberBus.subscribe((value) => {
- tenSubscriberValue = value + index;
- })
+const tenSubscriberSubs = Array.from(
+ { length: 10 },
+ (_, index) =>
+ tenSubscriberBus.subscribe((value) => {
+ tenSubscriberValue = value + index;
+ }),
);
let hundredSubscriberValue = 0;
const hundredSubscriberBus = new EventBus();
-const hundredSubscriberSubs = Array.from({ length: 100 }, (_, index) =>
- hundredSubscriberBus.subscribe((value) => {
- hundredSubscriberValue = value + index;
- })
+const hundredSubscriberSubs = Array.from(
+ { length: 100 },
+ (_, index) =>
+ hundredSubscriberBus.subscribe((value) => {
+ hundredSubscriberValue = value + index;
+ }),
);
const dispatcher = createEventDispatcher<{
@@ -39,33 +44,36 @@ const dispatcher = createEventDispatcher<{
}>();
let dispatcherValue = 0;
-const dispatcherSub = dispatcher.on('message', (payload) => {
+const dispatcherSub = dispatcher.on("message", (payload) => {
dispatcherValue = payload.id;
});
-bench('Events: EventBus emit -> 1 subscriber', () => {
+bench("Events: EventBus emit -> 1 subscriber", () => {
singleSubscriberBus.emit(1);
do_not_optimize(singleSubscriberValue);
});
-bench('Events: EventBus emit -> 10 subscribers', () => {
+bench("Events: EventBus emit -> 10 subscribers", () => {
tenSubscriberBus.emit(10);
do_not_optimize(tenSubscriberValue);
});
-bench('Events: EventBus emit -> 100 subscribers', () => {
+bench("Events: EventBus emit -> 100 subscribers", () => {
hundredSubscriberBus.emit(100);
do_not_optimize(hundredSubscriberValue);
});
-bench('Events: typed dispatcher emit -> handler', () => {
- dispatcher.emit('message', { id: 42, text: 'ok' });
+bench("Events: typed dispatcher emit -> handler", () => {
+ dispatcher.emit("message", { id: 42, text: "ok" });
do_not_optimize(dispatcherValue);
});
-bench('Events: setup-inclusive bus + 1000 subscribe/unsubscribe cycles', () => {
+bench("Events: setup-inclusive bus + 1000 subscribe/unsubscribe cycles", () => {
const bus = new EventBus();
- const subscriptions = Array.from({ length: 1000 }, () => bus.subscribe(() => {}));
+ const subscriptions = Array.from(
+ { length: 1000 },
+ () => bus.subscribe(() => {}),
+ );
for (const subscription of subscriptions) {
subscription.unsubscribe();
@@ -73,7 +81,7 @@ bench('Events: setup-inclusive bus + 1000 subscribe/unsubscribe cycles', () => {
do_not_optimize(subscriptions);
do_not_optimize(bus);
-}).gc('inner');
+}).gc("inner");
await run();
diff --git a/bench/latency_bench.ts b/bench/latency_bench.ts
index 3f65d78..d27ad68 100644
--- a/bench/latency_bench.ts
+++ b/bench/latency_bench.ts
@@ -1,3 +1,4 @@
+// deno-lint-ignore-file no-import-prefix
/**
* Low-latency benchmarks for primitive operations.
*
@@ -6,26 +7,26 @@
* single-value operator hop.
*/
-import { bench, do_not_optimize, run } from 'npm:mitata';
+import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34";
-import { EventBus } from '../events.ts';
-import { isObservableError } from '../error.ts';
-import { Observable } from '../observable.ts';
-import { pipe } from '../helpers/pipe.ts';
-import { filter, map } from '../helpers/operations/core.ts';
+import { EventBus } from "../events.ts";
+import { isObservableError } from "../error.ts";
+import { Observable } from "../observable.ts";
+import { pipe } from "../helpers/pipe.ts";
+import { filter, map } from "../helpers/operations/core.ts";
const noEmissionObservable = new Observable(() => {
return () => {};
});
-bench('Latency: subscribe + unsubscribe (no emissions)', () => {
+bench("Latency: subscribe + unsubscribe (no emissions)", () => {
const subscription = noEmissionObservable.subscribe(() => {});
subscription.unsubscribe();
do_not_optimize(subscription);
do_not_optimize(subscription.closed);
});
-bench('Latency: Observable.of(1) -> subscribe', () => {
+bench("Latency: Observable.of(1) -> subscribe", () => {
let lastValue = 0;
const subscription = Observable.of(1).subscribe((value) => {
lastValue = value;
@@ -35,7 +36,7 @@ bench('Latency: Observable.of(1) -> subscribe', () => {
subscription.unsubscribe();
});
-bench('Latency: single value through map', () => {
+bench("Latency: single value through map", () => {
let lastValue = 0;
const result = pipe(
Observable.of(1),
@@ -52,7 +53,7 @@ bench('Latency: single value through map', () => {
subscription.unsubscribe();
});
-bench('Latency: single value through map + filter', () => {
+bench("Latency: single value through map + filter", () => {
let lastValue = 0;
const result = pipe(
Observable.of(1),
@@ -80,7 +81,7 @@ const warmPathSub = warmPathBus.subscribe((value) => {
latestEventBusValue = value;
});
-bench('Latency: EventBus emit -> 1 subscriber', () => {
+bench("Latency: EventBus emit -> 1 subscriber", () => {
warmPathBus.emit(7);
do_not_optimize(latestEventBusValue);
});
diff --git a/bench/memory_bench.ts b/bench/memory_bench.ts
index facd5db..5452b61 100644
--- a/bench/memory_bench.ts
+++ b/bench/memory_bench.ts
@@ -1,21 +1,22 @@
+// deno-lint-ignore-file no-import-prefix
/**
* Memory allocation and GC pressure benchmarks.
- *
+ *
* Measures memory usage patterns, allocation rates, and GC behavior
* under various Observable usage scenarios. Critical for understanding
* real-world performance characteristics beyond simple execution time.
*/
-import { bench, run, do_not_optimize } from 'npm:mitata';
-import { Observable } from '../observable.ts';
-import { isObservableError } from '../error.ts';
-import { pipe } from '../helpers/pipe.ts';
-import { map, filter, take, scan } from '../helpers/operations/core.ts';
-import { createQueue, enqueue, dequeue } from '../queue.ts';
+import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34";
+import { Observable } from "../observable.ts";
+import { isObservableError } from "../error.ts";
+import { pipe } from "../helpers/pipe.ts";
+import { map, scan, take } from "../helpers/operations/core.ts";
+import { createQueue, dequeue, enqueue } from "../queue.ts";
// Memory tracking helper
-function getMemoryUsage(): number {
- if (typeof Deno !== 'undefined' && Deno.memoryUsage) {
+function _getMemoryUsage(): number {
+ if (typeof Deno !== "undefined" && Deno.memoryUsage) {
return Deno.memoryUsage().heapUsed;
}
return 0;
@@ -25,200 +26,202 @@ function getMemoryUsage(): number {
function* largeDataGenerator(sizeBytes: number) {
const chunkSize = 1024; // 1KB chunks
const numChunks = Math.floor(sizeBytes / chunkSize);
-
+
for (let i = 0; i < numChunks; i++) {
yield new Uint8Array(chunkSize);
}
}
-bench('Memory: Observable creation overhead (1000x)', () => {
+bench("Memory: Observable creation overhead (1000x)", () => {
const observables: Observable[] = [];
-
+
for (let i = 0; i < 1000; i++) {
observables.push(Observable.of(i));
}
-
+
do_not_optimize(observables);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Subscription tracking (1000 subs)', () => {
+bench("Memory: Subscription tracking (1000 subs)", () => {
const obs = Observable.of(1, 2, 3);
const subscriptions = [];
-
+
for (let i = 0; i < 1000; i++) {
subscriptions.push(obs.subscribe(() => {}));
}
-
+
// Cleanup
for (const sub of subscriptions) {
sub.unsubscribe();
}
-
+
do_not_optimize(subscriptions);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Queue buffer reuse (10000 cycles)', () => {
+bench("Memory: Queue buffer reuse (10000 cycles)", () => {
const queue = createQueue(1000);
-
+
// Fill queue
for (let i = 0; i < 1000; i++) {
enqueue(queue, i);
}
-
+
// Cycle through: should reuse buffer slots
for (let i = 0; i < 10000; i++) {
dequeue(queue);
enqueue(queue, i);
}
-
+
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Large data streaming (1MB)', async () => {
+bench("Memory: Large data streaming (1MB)", async () => {
const oneMB = 1024 * 1024;
const obs = Observable.from(largeDataGenerator(oneMB));
-
+
let totalBytes = 0;
for await (const chunk of obs) {
totalBytes += chunk.length;
}
-
+
do_not_optimize(totalBytes);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Large data streaming (10MB)', async () => {
+bench("Memory: Large data streaming (10MB)", async () => {
const tenMB = 10 * 1024 * 1024;
const obs = Observable.from(largeDataGenerator(tenMB));
-
+
let totalBytes = 0;
for await (const chunk of obs) {
totalBytes += chunk.length;
}
-
+
do_not_optimize(totalBytes);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Large data streaming (100MB)', async () => {
+bench("Memory: Large data streaming (100MB)", async () => {
const hundredMB = 100 * 1024 * 1024;
const obs = Observable.from(largeDataGenerator(hundredMB));
-
+
let totalBytes = 0;
for await (const chunk of obs) {
totalBytes += chunk.length;
}
-
+
do_not_optimize(totalBytes);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Operator chain with large data (10MB)', async () => {
+bench("Memory: Operator chain with large data (10MB)", async () => {
const tenMB = 10 * 1024 * 1024;
-
+
const result = pipe(
Observable.from(largeDataGenerator(tenMB)),
map((chunk: Uint8Array) => chunk.length),
scan((acc: number, len: number) => acc + len, 0),
- take(1000)
+ take(1000),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Rapid subscribe/unsubscribe cycles (10000x)', () => {
+bench("Memory: Rapid subscribe/unsubscribe cycles (10000x)", () => {
const obs = Observable.of(1, 2, 3, 4, 5);
-
+
for (let i = 0; i < 10000; i++) {
const sub = obs.subscribe(() => {});
sub.unsubscribe();
}
-
+
do_not_optimize(obs);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Concurrent subscriptions with cleanup (1000x)', () => {
+bench("Memory: Concurrent subscriptions with cleanup (1000x)", () => {
const obs = new Observable((observer) => {
const id = setInterval(() => {
observer.next(Math.random());
}, 100);
-
+
return () => clearInterval(id);
});
-
+
const subs = [];
for (let i = 0; i < 1000; i++) {
subs.push(obs.subscribe(() => {}));
}
-
+
// Cleanup all
for (const sub of subs) {
sub.unsubscribe();
}
-
+
do_not_optimize(subs);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Queue allocation patterns (100K items)', () => {
+bench("Memory: Queue allocation patterns (100K items)", () => {
const queue = createQueue(100000);
-
+
for (let i = 0; i < 100000; i++) {
enqueue(queue, i);
}
-
+
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('Memory: Observable value buffering (10000 items)', async () => {
+bench("Memory: Observable value buffering (10000 items)", async () => {
const values: number[] = [];
-
+
const obs = new Observable((observer) => {
for (let i = 0; i < 10000; i++) {
observer.next(i);
}
observer.complete();
});
-
+
for await (const val of obs) {
values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
// Stress test: Push system to limits
-bench('STRESS: 1GB data streaming', async () => {
+bench("STRESS: 1GB data streaming", async () => {
const oneGB = 1024 * 1024 * 1024;
const obs = Observable.from(largeDataGenerator(oneGB));
-
+
let totalBytes = 0;
let chunkCount = 0;
-
+
for await (const chunk of obs) {
totalBytes += chunk.length;
chunkCount++;
}
-
+
do_not_optimize({ totalBytes, chunkCount });
-}).gc('inner');
+}).gc("inner");
-bench('STRESS: 1M subscriptions lifecycle', () => {
+bench("STRESS: 1M subscriptions lifecycle", () => {
const obs = Observable.of(42);
let completedCount = 0;
-
+
// Note: This is intentionally stressful
// Real code shouldn't do this pattern
for (let i = 0; i < 1000000; i++) {
const sub = obs.subscribe({
next: () => {},
- complete: () => { completedCount++; }
+ complete: () => {
+ completedCount++;
+ },
});
sub.unsubscribe();
}
-
+
do_not_optimize(completedCount);
-}).gc('inner');
+}).gc("inner");
await run();
diff --git a/bench/observable_bench.ts b/bench/observable_bench.ts
index 1bef6ac..cc82366 100644
--- a/bench/observable_bench.ts
+++ b/bench/observable_bench.ts
@@ -1,14 +1,15 @@
+// deno-lint-ignore-file no-import-prefix
/**
* Observable creation and subscription benchmarks.
- *
+ *
* Measures overhead of Observable creation, subscription, and teardown
* across various scenarios from simple to complex patterns.
*/
-import { bench, run, do_not_optimize } from 'npm:mitata';
-import { Observable } from '../observable.ts';
-import { pipe } from '../helpers/pipe.ts';
-import { map, filter, take } from '../helpers/operations/core.ts';
+import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34";
+import { Observable } from "../observable.ts";
+import { pipe } from "../helpers/pipe.ts";
+import { filter, map, take } from "../helpers/operations/core.ts";
// Test data generation
function* numberGenerator(count: number) {
@@ -25,41 +26,41 @@ const simpleObservable = new Observable((observer) => {
const rangeObservable = Observable.from(numberGenerator(1000));
-const asyncObservable = new Observable((observer) => {
+const _asyncObservable = new Observable((observer) => {
const id = setInterval(() => {
observer.next(Math.random());
}, 10);
return () => clearInterval(id);
});
-bench('Observable.of(single value)', () => {
+bench("Observable.of(single value)", () => {
// Act & Assert
do_not_optimize(Observable.of(42));
});
-bench('Observable.of(10 values)', () => {
+bench("Observable.of(10 values)", () => {
do_not_optimize(Observable.of(1, 2, 3, 4, 5, 6, 7, 8, 9, 10));
});
-bench('Observable.from(array 100 items)', () => {
+bench("Observable.from(array 100 items)", () => {
const arr = Array.from({ length: 100 }, (_, i) => i);
do_not_optimize(Observable.from(arr));
});
-bench('Observable.from(generator 1000 items)', () => {
+bench("Observable.from(generator 1000 items)", () => {
do_not_optimize(Observable.from(numberGenerator(1000)));
-}).gc('inner');
+}).gc("inner");
-bench('new Observable (simple sync)', () => {
+bench("new Observable (simple sync)", () => {
do_not_optimize(
new Observable((observer) => {
observer.next(42);
observer.complete();
- })
+ }),
);
});
-bench('new Observable (with cleanup)', () => {
+bench("new Observable (with cleanup)", () => {
do_not_optimize(
new Observable((observer) => {
observer.next(42);
@@ -67,17 +68,17 @@ bench('new Observable (with cleanup)', () => {
return () => {
// Cleanup
};
- })
+ }),
);
});
-bench('subscribe + immediate complete', () => {
+bench("subscribe + immediate complete", () => {
const sub = simpleObservable.subscribe(() => {});
do_not_optimize(sub);
sub.unsubscribe();
});
-bench('subscribe + collect 1000 values', async () => {
+bench("subscribe + collect 1000 values", async () => {
const values: number[] = [];
await new Promise((resolve) => {
rangeObservable.subscribe({
@@ -86,35 +87,35 @@ bench('subscribe + collect 1000 values', async () => {
});
});
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('async iteration over 1000 values', async () => {
+bench("async iteration over 1000 values", async () => {
const values: number[] = [];
for await (const val of rangeObservable) {
values.push(val);
}
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('multiple subscribers (cold semantics)', () => {
+bench("multiple subscribers (cold semantics)", () => {
const obs = Observable.of(1, 2, 3);
const sub1 = obs.subscribe(() => {});
const sub2 = obs.subscribe(() => {});
const sub3 = obs.subscribe(() => {});
-
+
do_not_optimize([sub1, sub2, sub3]);
-
+
sub1.unsubscribe();
sub2.unsubscribe();
sub3.unsubscribe();
});
-bench('pipe with 3 operators', () => {
+bench("pipe with 3 operators", () => {
const result = pipe(
Observable.from(numberGenerator(100)),
map((x: number) => x * 2),
filter((x: number) => x % 4 === 0),
- take(10)
+ take(10),
);
do_not_optimize(result);
});
diff --git a/bench/operators_bench.ts b/bench/operators_bench.ts
index 88e8a07..78335b9 100644
--- a/bench/operators_bench.ts
+++ b/bench/operators_bench.ts
@@ -1,17 +1,17 @@
+// deno-lint-ignore-file no-import-prefix
/**
* Operator pipeline benchmarks.
- *
+ *
* Measures performance of various operator combinations that represent
* real-world Observable usage patterns.
*/
-import { bench, run, do_not_optimize } from 'npm:mitata';
-import { Observable } from '../observable.ts';
-import { isObservableError } from '../error.ts';
-import { pipe } from '../helpers/pipe.ts';
-import { map, filter, scan, take, tap } from '../helpers/operations/core.ts';
-import { debounce, delay, throttle } from '../helpers/operations/timing.ts';
-import { batch, toArray } from '../helpers/operations/batch.ts';
+import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34";
+import { Observable } from "../observable.ts";
+import { isObservableError } from "../error.ts";
+import { pipe } from "../helpers/pipe.ts";
+import { filter, map, scan, take, tap } from "../helpers/operations/core.ts";
+import { batch, toArray } from "../helpers/operations/batch.ts";
function* numberStream(count: number) {
for (let i = 0; i < count; i++) {
@@ -19,145 +19,147 @@ function* numberStream(count: number) {
}
}
-bench('Operators: map only (1000 items)', async () => {
+bench("Operators: map only (1000 items)", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
- map((x: number) => x * 2)
+ map((x: number) => x * 2),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: filter only (1000 items)', async () => {
+bench("Operators: filter only (1000 items)", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
- filter((x: number) => x % 2 === 0)
+ filter((x: number) => x % 2 === 0),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: map + filter chain', async () => {
+bench("Operators: map + filter chain", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
map((x: number) => x * 2),
filter((x: number) => x % 4 === 0),
- map((x: number) => x / 2)
+ map((x: number) => x / 2),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: map + filter + take', async () => {
+bench("Operators: map + filter + take", async () => {
const result = pipe(
Observable.from(numberStream(10000)),
map((x: number) => x * 2),
filter((x: number) => x % 4 === 0),
- take(100)
+ take(100),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
});
-bench('Operators: scan (running sum)', async () => {
+bench("Operators: scan (running sum)", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
- scan((acc: number, val: number) => acc + val, 0)
+ scan((acc: number, val: number) => acc + val, 0),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: complex chain (5 operators)', async () => {
+bench("Operators: complex chain (5 operators)", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
map((x: number) => x * 2),
filter((x: number) => x % 4 === 0),
scan((acc: number, val: number) => acc + val, 0),
map((x: number) => x / 100),
- take(500)
+ take(500),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: tap (side effects)', async () => {
+bench("Operators: tap (side effects)", async () => {
let sideEffectCount = 0;
-
+
const result = pipe(
Observable.from(numberStream(1000)),
- tap(() => { sideEffectCount++; }),
- map((x: number) => x * 2)
+ tap(() => {
+ sideEffectCount++;
+ }),
+ map((x: number) => x * 2),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize({ values, sideEffectCount });
-}).gc('inner');
+}).gc("inner");
-bench('Operators: batch (groups of 10)', async () => {
+bench("Operators: batch (groups of 10)", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
- batch(10)
+ batch(10),
);
-
+
const batches: number[][] = [];
for await (const batch of result) {
if (!isObservableError(batch)) batches.push(batch);
}
-
+
do_not_optimize(batches);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: toArray collector', async () => {
+bench("Operators: toArray collector", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
map((x: number) => x * 2),
- toArray()
+ toArray(),
);
-
+
const arrays: number[][] = [];
for await (const arr of result) {
if (!isObservableError(arr)) arrays.push(arr);
}
-
+
do_not_optimize(arrays);
-}).gc('inner');
+}).gc("inner");
-bench('Operators: deep chain (10 operators)', async () => {
+bench("Operators: deep chain (10 operators)", async () => {
const result = pipe(
Observable.from(numberStream(1000)),
map((x: number) => x + 1),
@@ -169,15 +171,15 @@ bench('Operators: deep chain (10 operators)', async () => {
filter((x: number) => x > 0),
map((x: number) => Math.floor(x)),
take(100),
- tap(() => {})
+ tap(() => {}),
);
-
+
const values: number[] = [];
for await (const val of result) {
if (!isObservableError(val)) values.push(val);
}
-
+
do_not_optimize(values);
-}).gc('inner');
+}).gc("inner");
await run();
diff --git a/bench/queue_bench.ts b/bench/queue_bench.ts
index dcec8fd..15d10ed 100644
--- a/bench/queue_bench.ts
+++ b/bench/queue_bench.ts
@@ -1,48 +1,48 @@
+// deno-lint-ignore-file no-import-prefix
/**
* Queue operations benchmarks focusing on O(1) performance claims.
- *
+ *
* Tests circular buffer operations at various scales to verify constant-time
* performance and compare against naive Array.shift() implementations.
*/
-import { bench, run, do_not_optimize } from 'npm:mitata';
+import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34";
import {
+ clear,
createQueue,
- enqueue,
dequeue,
- peek,
+ enqueue,
isEmpty,
- isFull,
- clear,
+ peek,
toArray,
-} from '../queue.ts';
+} from "../queue.ts";
// Baseline: Array.shift() for comparison
class ArrayQueue {
private items: T[] = [];
-
+
enqueue(item: T): void {
this.items.push(item);
}
-
+
dequeue(): T | undefined {
return this.items.shift();
}
-
+
peek(): T | undefined {
return this.items[0];
}
-
+
isEmpty(): boolean {
return this.items.length === 0;
}
}
-bench('Queue: create empty queue', () => {
+bench("Queue: create empty queue", () => {
do_not_optimize(createQueue(1000));
});
-bench('Queue: enqueue 100 items', () => {
+bench("Queue: enqueue 100 items", () => {
const queue = createQueue(100);
for (let i = 0; i < 100; i++) {
enqueue(queue, i);
@@ -50,28 +50,28 @@ bench('Queue: enqueue 100 items', () => {
do_not_optimize(queue);
});
-bench('Queue: enqueue 1000 items', () => {
+bench("Queue: enqueue 1000 items", () => {
const queue = createQueue(1000);
for (let i = 0; i < 1000; i++) {
enqueue(queue, i);
}
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('Queue: enqueue 10000 items', () => {
+bench("Queue: enqueue 10000 items", () => {
const queue = createQueue(10000);
for (let i = 0; i < 10000; i++) {
enqueue(queue, i);
}
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('Queue: dequeue 100 items', () => {
+bench("Queue: dequeue 100 items", () => {
const queue = createQueue(100);
for (let i = 0; i < 100; i++) {
enqueue(queue, i);
}
-
+
const results: number[] = [];
for (let i = 0; i < 100; i++) {
const val = dequeue(queue);
@@ -80,23 +80,23 @@ bench('Queue: dequeue 100 items', () => {
do_not_optimize(results);
});
-bench('Queue: dequeue 1000 items', () => {
+bench("Queue: dequeue 1000 items", () => {
const queue = createQueue(1000);
for (let i = 0; i < 1000; i++) {
enqueue(queue, i);
}
-
+
const results: number[] = [];
for (let i = 0; i < 1000; i++) {
const val = dequeue(queue);
if (val !== undefined) results.push(val);
}
do_not_optimize(results);
-}).gc('inner');
+}).gc("inner");
-bench('Queue: enqueue+dequeue mixed 1000 ops', () => {
+bench("Queue: enqueue+dequeue mixed 1000 ops", () => {
const queue = createQueue(500);
-
+
for (let i = 0; i < 1000; i++) {
if (i % 2 === 0) {
enqueue(queue, i);
@@ -105,31 +105,31 @@ bench('Queue: enqueue+dequeue mixed 1000 ops', () => {
}
}
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('Queue: circular wrap 1000 cycles', () => {
+bench("Queue: circular wrap 1000 cycles", () => {
const queue = createQueue(100);
-
+
// Fill queue
for (let i = 0; i < 100; i++) {
enqueue(queue, i);
}
-
+
// Cycle: dequeue one, enqueue one (causes wrapping)
for (let i = 0; i < 1000; i++) {
dequeue(queue);
enqueue(queue, i + 100);
}
-
+
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('Queue: peek 1000 times', () => {
+bench("Queue: peek 1000 times", () => {
const queue = createQueue(100);
for (let i = 0; i < 100; i++) {
enqueue(queue, i);
}
-
+
let sum = 0;
for (let i = 0; i < 1000; i++) {
const val = peek(queue);
@@ -138,51 +138,51 @@ bench('Queue: peek 1000 times', () => {
do_not_optimize(sum);
});
-bench('Queue: toArray with 1000 items', () => {
+bench("Queue: toArray with 1000 items", () => {
const queue = createQueue(1000);
for (let i = 0; i < 1000; i++) {
enqueue(queue, i);
}
-
+
do_not_optimize(toArray(queue));
-}).gc('inner');
+}).gc("inner");
-bench('Queue: clear 1000 items', () => {
+bench("Queue: clear 1000 items", () => {
const queue = createQueue(1000);
for (let i = 0; i < 1000; i++) {
enqueue(queue, i);
}
-
+
clear(queue);
do_not_optimize(queue);
});
// Comparison benchmarks: circular buffer vs Array.shift()
-bench('[Baseline] Array.shift: enqueue 1000 items', () => {
+bench("[Baseline] Array.shift: enqueue 1000 items", () => {
const queue = new ArrayQueue();
for (let i = 0; i < 1000; i++) {
queue.enqueue(i);
}
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
-bench('[Baseline] Array.shift: dequeue 1000 items', () => {
+bench("[Baseline] Array.shift: dequeue 1000 items", () => {
const queue = new ArrayQueue();
for (let i = 0; i < 1000; i++) {
queue.enqueue(i);
}
-
+
const results: number[] = [];
for (let i = 0; i < 1000; i++) {
const val = queue.dequeue();
if (val !== undefined) results.push(val);
}
do_not_optimize(results);
-}).gc('inner');
+}).gc("inner");
-bench('[Baseline] Array.shift: mixed 1000 ops', () => {
+bench("[Baseline] Array.shift: mixed 1000 ops", () => {
const queue = new ArrayQueue();
-
+
for (let i = 0; i < 1000; i++) {
if (i % 2 === 0) {
queue.enqueue(i);
@@ -191,6 +191,6 @@ bench('[Baseline] Array.shift: mixed 1000 ops', () => {
}
}
do_not_optimize(queue);
-}).gc('inner');
+}).gc("inner");
await run();
diff --git a/deno.jsonc b/deno.jsonc
index fa51e24..121e576 100644
--- a/deno.jsonc
+++ b/deno.jsonc
@@ -34,6 +34,11 @@
},
"license": "MIT",
"publish": {
+ "include": [
+ "**/*.ts",
+ "LICENSE",
+ "README.md"
+ ],
"exclude": [
".devcontainer/",
".github/",
@@ -55,6 +60,7 @@
},
"exclude": [
"coverage/",
- "bench/results/"
+ "bench/results/",
+ "npm/"
]
}
diff --git a/deno.lock b/deno.lock
new file mode 100644
index 0000000..b4a5dad
--- /dev/null
+++ b/deno.lock
@@ -0,0 +1,109 @@
+{
+ "version": "5",
+ "specifiers": {
+ "jsr:@david/code-block-writer@^13.0.3": "13.0.3",
+ "jsr:@deno/dnt@*": "0.42.3",
+ "jsr:@std/assert@^1.0.17": "1.0.19",
+ "jsr:@std/assert@^1.0.19": "1.0.19",
+ "jsr:@std/expect@1": "1.0.18",
+ "jsr:@std/fmt@1": "1.0.9",
+ "jsr:@std/fs@1": "1.0.23",
+ "jsr:@std/internal@^1.0.12": "1.0.12",
+ "jsr:@std/json@^1.0.2": "1.0.3",
+ "jsr:@std/jsonc@*": "1.0.2",
+ "jsr:@std/path@1": "1.1.4",
+ "jsr:@std/path@^1.1.4": "1.1.4",
+ "jsr:@std/testing@1": "1.0.17",
+ "jsr:@ts-morph/bootstrap@0.27": "0.27.0",
+ "jsr:@ts-morph/common@0.27": "0.27.0",
+ "npm:mitata@^1.0.34": "1.0.34"
+ },
+ "jsr": {
+ "@david/code-block-writer@13.0.3": {
+ "integrity": "f98c77d320f5957899a61bfb7a9bead7c6d83ad1515daee92dbacc861e13bb7f"
+ },
+ "@deno/dnt@0.42.3": {
+ "integrity": "62a917a0492f3c8af002dce90605bb0d41f7d29debc06aca40dba72ab65d8ae3",
+ "dependencies": [
+ "jsr:@david/code-block-writer",
+ "jsr:@std/fmt",
+ "jsr:@std/fs",
+ "jsr:@std/path@1",
+ "jsr:@ts-morph/bootstrap"
+ ]
+ },
+ "@std/assert@1.0.19": {
+ "integrity": "eaada96ee120cb980bc47e040f82814d786fe8162ecc53c91d8df60b8755991e",
+ "dependencies": [
+ "jsr:@std/internal"
+ ]
+ },
+ "@std/expect@1.0.18": {
+ "integrity": "8566eab35200466f8609eb7e7aed062ed0db314e9a258d5d201b1b8997ce801a",
+ "dependencies": [
+ "jsr:@std/assert@^1.0.19",
+ "jsr:@std/internal",
+ "jsr:@std/path@^1.1.4"
+ ]
+ },
+ "@std/fmt@1.0.9": {
+ "integrity": "2487343e8899fb2be5d0e3d35013e54477ada198854e52dd05ed0422eddcabe0"
+ },
+ "@std/fs@1.0.23": {
+ "integrity": "3ecbae4ce4fee03b180fa710caff36bb5adb66631c46a6460aaad49515565a37",
+ "dependencies": [
+ "jsr:@std/internal",
+ "jsr:@std/path@^1.1.4"
+ ]
+ },
+ "@std/internal@1.0.12": {
+ "integrity": "972a634fd5bc34b242024402972cd5143eac68d8dffaca5eaa4dba30ce17b027"
+ },
+ "@std/json@1.0.3": {
+ "integrity": "97d5710996293a027b7aa5f0d1f4fa29f246f269e6b5597e08807613f37d426c"
+ },
+ "@std/jsonc@1.0.2": {
+ "integrity": "909605dae3af22bd75b1cbda8d64a32cf1fd2cf6efa3f9e224aba6d22c0f44c7",
+ "dependencies": [
+ "jsr:@std/json"
+ ]
+ },
+ "@std/path@1.1.4": {
+ "integrity": "1d2d43f39efb1b42f0b1882a25486647cb851481862dc7313390b2bb044314b5",
+ "dependencies": [
+ "jsr:@std/internal"
+ ]
+ },
+ "@std/testing@1.0.17": {
+ "integrity": "87bdc2700fa98249d48a17cd72413352d3d3680dcfbdb64947fd0982d6bbf681",
+ "dependencies": [
+ "jsr:@std/assert@^1.0.17",
+ "jsr:@std/internal"
+ ]
+ },
+ "@ts-morph/bootstrap@0.27.0": {
+ "integrity": "b8d7bc8f7942ce853dde4161b28f9aa96769cef3d8eebafb379a81800b9e2448",
+ "dependencies": [
+ "jsr:@ts-morph/common"
+ ]
+ },
+ "@ts-morph/common@0.27.0": {
+ "integrity": "c7b73592d78ce8479b356fd4f3d6ec3c460d77753a8680ff196effea7a939052"
+ }
+ },
+ "npm": {
+ "mitata@1.0.34": {
+ "integrity": "sha512-Mc3zrtNBKIMeHSCQ0XqRLo1vbdIx1wvFV9c8NJAiyho6AjNfMY8bVhbS12bwciUdd1t4rj8099CH3N3NFahaUA=="
+ }
+ },
+ "workspace": {
+ "dependencies": [
+ "jsr:@libs/testing@5",
+ "jsr:@std/assert@^1.0.14",
+ "jsr:@std/async@^1.0.14",
+ "jsr:@std/expect@^1.0.17",
+ "jsr:@std/testing@^1.0.15",
+ "npm:mitata@^1.0.34"
+ ]
+ }
+}
diff --git a/events.ts b/events.ts
index e6f194e..62cd6f6 100644
--- a/events.ts
+++ b/events.ts
@@ -3,12 +3,19 @@
*/
import type { Observer, Subscription } from "./_types.ts";
-import type { SubscriptionObserver } from './observable.ts';
+import type { SubscriptionObserver } from "./observable.ts";
-import { Observable } from './observable.ts';
+import { Observable } from "./observable.ts";
import { Symbol } from "./symbol.ts";
-import { createQueue, enqueue, dequeue, isFull, clear, forEach } from './queue.ts'; // Assume path to your queue utils
+import {
+ clear,
+ createQueue,
+ dequeue,
+ enqueue,
+ forEach,
+ isFull,
+} from "./queue.ts"; // Assume path to your queue utils
/**
* A multicast event bus that extends {@link Observable}, allowing
@@ -17,7 +24,6 @@ import { createQueue, enqueue, dequeue, isFull, clear, forEach } from './queue.t
*
* @typeParam T - The type of values emitted by this bus.
*
- *
* - Calling {@link emit} delivers the value to all active subscribers.
* - Calling {@link close} completes all subscribers and prevents further emissions.
* - Implements both {@link Symbol.dispose} and {@link Symbol.asyncDispose}
@@ -53,12 +59,11 @@ export class EventBus extends Observable {
/**
* Construct a new EventBus instance.
*
- *
* The base {@link Observable} constructor is invoked with the subscriber
* registration logic, adding and removing subscribers to the internal set.
*/
constructor() {
- super(subscriber => {
+ super((subscriber) => {
if (this.#closed) {
subscriber.complete?.();
return;
@@ -109,7 +114,6 @@ export class EventBus extends Observable {
/**
* Synchronous disposal method (for `using` syntax).
*
- *
* Alias for {@link close}.
*/
[Symbol.dispose](): void {
@@ -119,7 +123,6 @@ export class EventBus extends Observable {
/**
* Asynchronous disposal method.
*
- *
* Alias for {@link close}.
*/
async [Symbol.asyncDispose](): Promise {
@@ -138,7 +141,7 @@ export class EventBus extends Observable {
* }
* ```
*/
-export type EventMap = {};
+export type EventMap = object;
/**
* The return type of {@link createEventDispatcher}.
@@ -163,7 +166,7 @@ export interface EventDispatcher {
*/
on(
name: Name,
- handler: (payload: E[Name]) => void
+ handler: (payload: E[Name]) => void,
): Subscription;
/**
@@ -222,7 +225,9 @@ export interface EventDispatcher {
* bus.close();
* ```
*/
-export function createEventDispatcher(): EventDispatcher {
+export function createEventDispatcher(): EventDispatcher<
+ E
+> {
// Internal bus carries a union of all event types and payloads
const bus = new EventBus<{ type: keyof E; payload: E[keyof E] }>();
@@ -245,14 +250,14 @@ export function createEventDispatcher(): EventDispatcher
*/
on(
name: Name,
- handler: (payload: E[Name]) => void
+ handler: (payload: E[Name]) => void,
) {
return bus.events.subscribe({
next(event) {
if (event.type === name) {
handler(event.payload as E[Name]);
}
- }
+ },
});
},
@@ -282,7 +287,7 @@ export function createEventDispatcher(): EventDispatcher
*/
close(): void {
bus.close();
- }
+ },
};
}
@@ -334,13 +339,15 @@ export interface WaitForEventOptions {
*/
export function waitForEvent<
E extends EventMap,
- K extends keyof E
+ K extends keyof E,
>(
bus: { events: Observable<{ type: keyof E; payload: E[keyof E] }> },
type: K,
- { signal, throwOnClose = false }: WaitForEventOptions = {}
+ { signal, throwOnClose = false }: WaitForEventOptions = {},
): Promise {
- const { resolve, reject, promise } = Promise.withResolvers();
+ const { resolve, reject, promise } = Promise.withResolvers<
+ E[K] | undefined
+ >();
// Immediate abort
if (signal?.aborted) {
@@ -348,13 +355,13 @@ export function waitForEvent<
return promise;
}
- let subscription: Subscription | undefined;
+ const subscription_ref: { current?: Subscription } = {};
- subscription = bus.events.subscribe({
+ subscription_ref.current = bus.events.subscribe({
next(event) {
if (event.type === type) {
cleanup();
-
+
// cast payload to the correct type
resolve(event.payload as E[K]);
}
@@ -375,8 +382,8 @@ export function waitForEvent<
});
function cleanup() {
- subscription?.unsubscribe?.();
- signal?.removeEventListener?.('abort', onAbort);
+ subscription_ref.current?.unsubscribe?.();
+ signal?.removeEventListener?.("abort", onAbort);
}
function onAbort() {
@@ -384,12 +391,11 @@ export function waitForEvent<
reject(signal!.reason);
}
- signal?.addEventListener?.('abort', onAbort, { once: true });
+ signal?.addEventListener?.("abort", onAbort, { once: true });
return promise;
}
-
/**
* Controls when the replay buffer connects to the source Observable.
*
@@ -401,7 +407,7 @@ export function waitForEvent<
* Choose 'eager' for system-critical events you never want to miss.
* Choose 'lazy' for expensive operations that shouldn't run without consumers.
*/
-export type ReplayMode = 'eager' | 'lazy';
+export type ReplayMode = "eager" | "lazy";
/**
* Configuration options for replay behavior.
@@ -410,14 +416,14 @@ export interface ReplayOptions {
/**
* Maximum number of values to buffer.
* When the buffer is full, the oldest value is discarded (FIFO).
- *
+ *
* @default Infinity (unlimited buffer, use with caution)
*/
count?: number;
/**
* Determines when to connect to the source Observable.
- *
+ *
* @default 'lazy' (resource-efficient, connects on-demand)
*/
mode?: ReplayMode;
@@ -426,28 +432,28 @@ export interface ReplayOptions {
/**
* Adds replay capability to an Observable, multicasting values to multiple subscribers
* while maintaining a buffer of recent emissions.
- *
+ *
* Without replay, each new subscriber triggers a fresh execution of the source Observable:
* ```ts
* const apiCall = new Observable(subscriber => {
* console.log('Making expensive API call...');
* fetch('/api/data').then(response => subscriber.next(response));
* });
- *
+ *
* apiCall.subscribe(data1 => {}); // Triggers API call #1
* apiCall.subscribe(data2 => {}); // Triggers API call #2 (duplicate!)
* ```
- *
+ *
* With replay, the source executes once and shares results:
* ```ts
* const sharedApi = withReplay(apiCall, { count: 1, mode: 'lazy' });
- *
+ *
* sharedApi.subscribe(data1 => {}); // Triggers API call
* sharedApi.subscribe(data2 => {}); // Gets cached result, no new call!
* ```
- *
+ *
* ## Memory Considerations
- *
+ *
* - Buffer size directly impacts memory usage: `count * sizeof(T)`
* - 'eager' mode holds references even with no subscribers (potential memory leak)
* - 'lazy' mode clears buffer when all subscribers disconnect (automatic cleanup)
@@ -456,49 +462,49 @@ export interface ReplayOptions {
* > **Note**: Infinite buffers are by default capped at 1000 items to prevent memory issues.
* > The primary reason for this cap is because some runtimes such as Deno and Node.js
* > litereally crash when you try to allocate Ininity-sized arrays.
- *
+ *
* ## Performance Characteristics
- *
+ *
* - Enqueue/Dequeue: O(1) constant time
* - New subscriber replay: O(n) where n = buffer size
* - Memory overhead: One queue + subscriber set + source subscription
- *
+ *
* ## Edge Cases & Gotchas
- *
+ *
* 1. **Late subscribers in eager mode**: May receive very old values if the source
* emitted long ago and no cleanup occurred.
- *
+ *
* 2. **Infinite buffers**: Without a count limit, buffers grow indefinitely.
* Always set a reasonable count for production use.
- *
+ *
* 3. **Error handling**: Errors are multicast to all subscribers but don't clear
* the buffer. New subscribers still get the replay before the error.
- *
+ *
* 4. **Completion**: The source completion is multicast, but the replay buffer
* remains accessible to new subscribers (they get replay + completion).
- *
+ *
* @param source The source Observable to add replay behavior to
* @param options Configuration for replay behavior
* @returns A new Observable with replay capability
- *
+ *
* @example
* ```ts
* // Lazy mode - only buffers when subscribers are present
- * const shared = withReplay(expensive, {
- * count: 5,
+ * const shared = withReplay(expensive, {
+ * count: 5,
* mode: 'lazy' // Only run expensive when needed
* });
- *
+ *
* // Eager mode - always buffering, like a flight recorder
- * const eventLog = withReplay(systemEvents, {
- * count: 100,
+ * const eventLog = withReplay(systemEvents, {
+ * count: 100,
* mode: 'eager' // Capture events even if no one's listening
* });
* ```
*/
export function withReplay(
source: Observable,
- { count = Infinity, mode = "lazy" }: ReplayOptions = {}
+ { count = Infinity, mode = "lazy" }: ReplayOptions = {},
): Observable {
// Validate inputs
if (count <= 0) {
@@ -514,7 +520,7 @@ export function withReplay(
next(value) {
// Manage buffer capacity
if (count !== Infinity && isFull(buffer)) {
- dequeue(buffer); // Remove oldest
+ dequeue(buffer); // Remove oldest
}
enqueue(buffer, value);
@@ -532,7 +538,7 @@ export function withReplay(
for (const sub of subscribers) {
sub.complete();
}
- }
+ },
};
const isEager = mode === "eager";
@@ -541,9 +547,9 @@ export function withReplay(
/**
* Creates the replay Observable that new subscribers will receive.
*/
- return new Observable(subscriber => {
+ return new Observable((subscriber) => {
// Step 1: Replay buffered values to the new subscriber
- forEach(buffer, item => subscriber.next(item));
+ forEach(buffer, (item) => subscriber.next(item));
// Step 2: Add to active subscribers for future emissions
subscribers.add(subscriber);
@@ -561,7 +567,7 @@ export function withReplay(
if (!isEager && subscribers.size === 0 && shared) {
shared.unsubscribe();
shared = null;
- clear(buffer); // Clear shared buffer when fully disconnected
+ clear(buffer); // Clear shared buffer when fully disconnected
}
// In eager mode, we keep the connection alive regardless
};
diff --git a/helpers/_types.ts b/helpers/_types.ts
index 219060e..6eab603 100644
--- a/helpers/_types.ts
+++ b/helpers/_types.ts
@@ -6,49 +6,69 @@ import type { Observable } from "../observable.ts";
* Type representing a stream operator function
* Transforms a ReadableStream of type In to a ReadableStream of type Out
*/
-export type Operator = (stream: ReadableStream) => ReadableStream;
+export type Operator = (
+ stream: ReadableStream,
+) => ReadableStream;
+
+/**
+ * Removes `ObservableError` from a type so operators can describe
+ * error-filtered output channels.
+ */
export type ExcludeError = Exclude;
-// Combines Operator and SafeOperator.
-// Why: Supports mixed operators. Solves pipeline flexibility.
+/**
+ * Represents an operator slot in a pipeline where the previous operator may
+ * or may not have filtered `ObservableError` values out of the stream.
+ */
export type OperatorItem = Operator | Operator>;
-// Inference Types
-// Figures out the source type (Observable or Operator).
-// Why: Ensures correct input type. Solves type safety in pipelines.
-export type InferSourceType =
- TSource extends Observable ? InferObservableType :
- TSource extends OperatorItem ? InferOperatorItemOutputType :
- TSource;
-
-// Gets the output type of an OperatorItem.
-// Why: Tracks operator output. Solves pipeline type resolution.
-export type InferOperatorItemOutputType> =
- ReturnType extends ReadableStream ? T : never;
-
-// Gets the data type an Observable emits.
-// Why: Extracts Observable data type. Solves type-safe data access.
-export type InferObservableType> =
+/**
+ * Infers the item type contributed by either an Observable source or an
+ * operator in a pipe chain.
+ */
+export type InferSourceType = TSource extends
+ Observable ? InferObservableType
+ : TSource extends OperatorItem
+ ? InferOperatorItemOutputType
+ : TSource;
+
+/**
+ * Extracts the chunk type emitted by an operator.
+ */
+export type InferOperatorItemOutputType<
+ TSource extends OperatorItem,
+> = ReturnType extends ReadableStream ? T : never;
+
+/**
+ * Extracts the value type emitted by an Observable.
+ */
+export type InferObservableType> =
TSource extends Observable ? R : any;
-// Utility Types
-// Gets the first item of a tuple.
-// Why: Accesses pipeline start. Solves type extraction.
-export type FirstTupleItem = TTuple[0];
+/**
+ * Returns the first item in a non-empty tuple.
+ */
+export type FirstTupleItem =
+ TTuple[0];
-// Gets the last item of a tuple.
-// Why: Finds final operator. Solves pipeline output typing.
-export type GenericLastTupleItem =
- T extends [...infer _, infer L] ? L : never;
+/**
+ * Returns the last item in any tuple shape.
+ */
+export type GenericLastTupleItem = T extends
+ [...infer _, infer L] ? L : never;
-// Ensures last tuple item is an OperatorItem.
-// Why: Validates pipeline end. Solves output type safety.
+/**
+ * Returns the last tuple item only when it is a valid operator item.
+ */
export type LastTupleItem =
- GenericLastTupleItem extends OperatorItem ? GenericLastTupleItem : never;
+ GenericLastTupleItem extends OperatorItem
+ ? GenericLastTupleItem
+ : never;
-// Pipeline final type
-// Defines Observable output based on last operator.
-// Why: Sets pipeline result type. Solves type-safe output.
+/**
+ * Computes the Observable type returned by a `pipe()` call from the source and
+ * its final operator.
+ */
export type ObservableWithPipe<
TPipe extends readonly [Observable, ...OperatorItem[]],
> = Observable>>;
@@ -57,7 +77,7 @@ export type ObservableWithPipe<
* Type representing how to handle errors in operators
* - "ignore" => errors wrapped in ObservableError will be used as values
* and can then be transformed as the operator sees fit
- * - "pass-through" => errors automatically pass through, meaning ObservableError
+ * - "pass-through" => errors automatically pass through, meaning ObservableError
* will not appear as a value
* - "throw" => errors will cause the stream to throw and terminate
* - "manual" => errors are passed to the transform function to handle manually
@@ -96,7 +116,7 @@ export interface TransformFunctionOptions extends BaseTransformOptions {
* How to handle errors in the stream:
* - "ignore" => errors wrapped in ObservableError will be used as values
* and can then be transformed as the operator sees fit
- * - "pass-through" => errors automatically pass through, meaning ObservableError
+ * - "pass-through" => errors automatically pass through, meaning ObservableError
* will not appear as a value
* - "throw" => errors will cause the stream to throw and terminate
* - "manual" => errors are passed to the transform function to handle manually
@@ -112,7 +132,7 @@ export interface TransformFunctionOptions extends BaseTransformOptions {
*/
transform: (
chunk: T,
- controller: TransformStreamDefaultController
+ controller: TransformStreamDefaultController,
) => R | undefined | void | null | Promise;
/**
@@ -121,7 +141,7 @@ export interface TransformFunctionOptions extends BaseTransformOptions {
* @param controller - The TransformStreamDefaultController
*/
flush?: (
- controller: TransformStreamDefaultController
+ controller: TransformStreamDefaultController,
) => void | Promise;
/**
@@ -129,7 +149,7 @@ export interface TransformFunctionOptions extends BaseTransformOptions {
* @param controller - The TransformStreamDefaultController
*/
start?: (
- controller: TransformStreamDefaultController
+ controller: TransformStreamDefaultController,
) => void | Promise;
/**
@@ -152,12 +172,13 @@ export type CreateOperatorOptions =
/**
* Options for stateful transformation logic
*/
-export interface StatefulTransformFunctionOptions extends BaseTransformOptions {
+export interface StatefulTransformFunctionOptions
+ extends BaseTransformOptions {
/**
* How to handle errors in the stream:
* - "ignore" => errors wrapped in ObservableError will be used as values
* and can then be transformed as the operator sees fit
- * - "pass-through" => errors automatically pass through, meaning ObservableError
+ * - "pass-through" => errors automatically pass through, meaning ObservableError
* will not appear as a value
* - "throw" => errors will cause the stream to throw and terminate
* - "manual" => errors are passed to the transform function to handle manually
@@ -180,7 +201,7 @@ export interface StatefulTransformFunctionOptions extends BaseTransform
transform: (
chunk: T,
state: S,
- controller: TransformStreamDefaultController
+ controller: TransformStreamDefaultController,
) => void | Promise;
/**
@@ -191,7 +212,7 @@ export interface StatefulTransformFunctionOptions extends BaseTransform
*/
flush?: (
state: S,
- controller: TransformStreamDefaultController
+ controller: TransformStreamDefaultController,
) => void | Promise;
/**
@@ -201,7 +222,7 @@ export interface StatefulTransformFunctionOptions extends BaseTransform
*/
start?: (
state: S,
- controller: TransformStreamDefaultController
+ controller: TransformStreamDefaultController,
) => void | Promise;
/**
@@ -210,7 +231,7 @@ export interface StatefulTransformFunctionOptions extends BaseTransform
*/
cancel?: (
state: S,
- reason?: unknown
+ reason?: unknown,
) => void | Promise;
}
@@ -222,7 +243,16 @@ export interface StatefulTransformFunctionOptions extends BaseTransform
* @typeParam S - State type (if applicable)
*/
export interface TransformHandlerContext {
+ /**
+ * Human-readable operator name used in wrapped error messages.
+ */
operatorName?: string;
+ /**
+ * Indicates that the lifecycle handler should pass shared state through.
+ */
isStateful?: boolean;
+ /**
+ * Shared operator state for stateful transforms.
+ */
state?: S;
-}
\ No newline at end of file
+}
diff --git a/helpers/operations/combination.ts b/helpers/operations/combination.ts
index bfe2fe1..5b69f0d 100644
--- a/helpers/operations/combination.ts
+++ b/helpers/operations/combination.ts
@@ -6,7 +6,7 @@ import type { SpecObservable } from "../../_spec.ts";
import type { ExcludeError, Operator } from "../_types.ts";
import { createStatefulOperator } from "../operators.ts";
-import { ObservableError, isObservableError } from "../../error.ts";
+import { isObservableError, ObservableError } from "../../error.ts";
import { pull } from "../../observable.ts";
/**
@@ -56,7 +56,7 @@ import { pull } from "../../observable.ts";
*/
export function mergeMap(
project: (value: ExcludeError, index: number) => SpecObservable,
- concurrent: number = Infinity
+ concurrent: number = Infinity,
): Operator {
return createStatefulOperator(
buffer: [],
sourceCompleted: false,
index: 0,
- activeCount: 0
+ activeCount: 0,
}),
// Process each incoming chunk
@@ -96,7 +96,13 @@ export function mergeMap(
innerObservable = project(value as ExcludeError, innerIndex);
} catch (err) {
// Forward any errors from the projection function
- controller.enqueue(ObservableError.from(err, "operator:stateful:mergeMap:project", value) as R);
+ controller.enqueue(
+ ObservableError.from(
+ err,
+ "operator:stateful:mergeMap:project",
+ value,
+ ) as R,
+ );
return;
}
@@ -104,11 +110,19 @@ export function mergeMap(
// Use pull to iterate asynchronously
try {
- for await (const innerValue of pull(innerObservable, { throwError: false })) {
+ for await (
+ const innerValue of pull(innerObservable, { throwError: false })
+ ) {
controller.enqueue(innerValue as R | ObservableError);
}
} catch (err) {
- controller.enqueue(ObservableError.from(err, "operator:stateful:mergeMap:innerObservable", value) as R);
+ controller.enqueue(
+ ObservableError.from(
+ err,
+ "operator:stateful:mergeMap:innerObservable",
+ value,
+ ) as R,
+ );
} finally {
// Clean up after inner Observable completes
state.activeSubscriptions.delete(innerIndex);
@@ -155,7 +169,7 @@ export function mergeMap(
state.buffer.length = 0;
state.activeSubscriptions.clear();
state.activeCount = 0;
- }
+ },
});
}
@@ -209,7 +223,7 @@ export function mergeMap(
* @returns An operator function that maps and concatenates values
*/
export function concatMap(
- project: (value: ExcludeError, index: number) => SpecObservable
+ project: (value: ExcludeError, index: number) => SpecObservable,
): Operator {
// concatMap is just mergeMap with concurrency = 1
return mergeMap(project, 1);
@@ -260,7 +274,7 @@ export function concatMap(
* @returns An operator function that maps and switches between values
*/
export function switchMap(
- project: (value: ExcludeError, index: number) => SpecObservable
+ project: (value: ExcludeError, index: number) => SpecObservable,
): Operator {
return createStatefulOperator(
currentTask: null,
currentTaskToken: null,
sourceCompleted: false,
- index: 0
+ index: 0,
}),
// Process each incoming chunk
@@ -302,7 +316,13 @@ export function switchMap(
innerObservable = project(chunk as ExcludeError, state.index++);
} catch (err) {
// Forward any errors from the projection function
- controller.enqueue(ObservableError.from(err, "operator:stateful:switchMap:project", chunk) as R);
+ controller.enqueue(
+ ObservableError.from(
+ err,
+ "operator:stateful:switchMap:project",
+ chunk,
+ ) as R,
+ );
return;
}
@@ -312,8 +332,7 @@ export function switchMap(
// Subscribe to the new inner Observable
const currentTaskToken = {};
- let currentTask: Promise;
- currentTask = (async () => {
+ const currentTask: Promise = (async () => {
const enqueueIfActive = (value: R | ObservableError): void => {
if (
abortController.signal.aborted ||
@@ -330,7 +349,8 @@ export function switchMap(
};
try {
- const iterator = pull(innerObservable, { throwError: false })[Symbol.asyncIterator]();
+ const iterator = pull(innerObservable, { throwError: false })
+ [Symbol.asyncIterator]();
while (!abortController.signal.aborted) {
const { value, done } = await iterator.next();
@@ -344,7 +364,11 @@ export function switchMap(
} catch (err) {
if (!abortController.signal.aborted) {
enqueueIfActive(
- ObservableError.from(err, "operator:stateful:switchMap:innerObservable", chunk),
+ ObservableError.from(
+ err,
+ "operator:stateful:switchMap:innerObservable",
+ chunk,
+ ),
);
}
} finally {
@@ -386,6 +410,6 @@ export function switchMap(
}
state.currentTask = null;
state.currentTaskToken = null;
- }
+ },
});
}
diff --git a/helpers/operators.ts b/helpers/operators.ts
index a8c51c3..2aab8b4 100644
--- a/helpers/operators.ts
+++ b/helpers/operators.ts
@@ -1,19 +1,19 @@
/**
* Operators are the building blocks of Observable pipelines.
- *
+ *
* If you've ever used `Array.map` or `Array.filter`, you already know the core idea:
* an **operator** takes a sequence of values and transforms, filters, or combines them
* into a new sequence. Operators let you build data pipelines, think of them as the
* Lego bricks for working with streams of data.
- *
+ *
* Think of an operator as a function that takes a stream of values and returns a new stream,
* transforming, filtering, or combining the data as it flows through.
- *
+ *
* For example, to double every number in an array:
* ```ts
* [1, 2, 3].map(x => x * 2); // [2, 4, 6]
* ```
- *
+ *
* With Observables, you want to do the same thing, but for values that arrive over time:
* ```ts
* // Double every number in a stream
@@ -22,14 +22,14 @@
* controller.enqueue(chunk * 2);
* }
* });
- *
+ *
* // Only allow even numbers through
* const evens = createOperator({
* transform(chunk, controller) {
* if (chunk % 2 === 0) controller.enqueue(chunk);
* }
* });
- *
+ *
* // Use them together in a pipeline
* pipe(
* Observable.from([1, 2, 3, 4]),
@@ -37,19 +37,19 @@
* evens
* ).subscribe(console.log); // Output: 4, 8
* ```
- *
+ *
* This module lets you build your own operators using the Web Streams API under the hood.
* Why streams? Because they're fast, memory-efficient, and let you process data as it arrives,
* not just after everything is loaded. This is especially useful for things like file processing,
* network requests, or any situation where you want to handle data piece-by-piece.
- *
+ *
* ## Why Streams? Why Not Just Arrays?
*
* Arrays are great for data you already have. But what about data that arrives slowly,
* or is too big to fit in memory? Think files, network responses, or user events.
* That's where **streams** shine: they let you process data piece-by-piece, as it arrives,
* without waiting for everything or loading it all at once.
- *
+ *
* The Web Streams API (and Node.js streams) are the standard way to do this in modern JavaScript.
* But using them directly is verbose and error-prone:
* ```ts
@@ -69,16 +69,16 @@
* }
* });
* ```
- *
+ *
* By building operators on top of streams, you get:
* - **Backpressure**: Slow consumers don't overwhelm fast producers.
* - **Low memory usage**: Process data chunk-by-chunk, not all at once.
* - **Composable pipelines**: Easily chain transformations.
- *
+ *
* ## Connecting Operators: Pipelines
*
* Operators are most powerful when you chain them together. This is called a pipeline.
- *
+ *
* It's just like chaining `map` and `filter` on arrays, but for streams:
* ```ts
* pipe(
@@ -95,7 +95,7 @@
* })
* ).subscribe(console.log); // Output: 6
* ```
- *
+ *
* Compare to arrays:
* ```ts
* [1, 2, 3, 4]
@@ -103,11 +103,11 @@
* .filter(x => x % 3 === 0)
* .forEach(console.log); // [2, 4, 8]
* ```
- *
- * Of course, no one wants to write operators from scratch every time.
+ *
+ * Of course, no one wants to write operators from scratch every time.
* So we provide some core operations via basic familiar operators,
* plus error handling utilities to make your pipelines robust.
- *
+ *
* Aka, `map`, `filter`, `reduce`, `batch`, `catchErrors`, `ignoreErrors`, and more.
* So really the example above becomes:
* ```ts
@@ -117,20 +117,20 @@
* filter(x => x % 3 === 0)
* ).subscribe(console.log); // Output: 2, 4, 8
* ```
- *
- * The example is not ideal given arrays have functions for this already,
+ *
+ * The example is not ideal given arrays have functions for this already,
* but you get the idea. It's meant more for streams of data that arrive over time.
- *
+ *
* ## Error Handling: Real-World Data is Messy
- *
+ *
* Real-world data is messy. Sometimes things go wrong aka, maybe a chunk is malformed, or a network
* request fails. Our operators let you choose how to handle errors, with four modes:
- *
+ *
* - `"pass-through"` (default): Errors become special values in the stream, so you can handle them downstream. Imagine almost like bubble wrap over error since they are dangerous allowing us to make sure we don't break the flow.
- * - `"ignore"`: Errors are silently skipped. The stream keeps going as if nothing happened. Imagine that we're basically just remove any errors from the stream while it's flowing (pretty stressful ngl).
+ * - `"ignore"`: Errors are silently skipped. The stream keeps going as if nothing happened. Imagine that we're basically just remove any errors from the stream while it's flowing (pretty stressful ngl).
* - `"throw"`: The stream stops immediately on the first error. Basically start screaming bloody murder, an error has occured so everything must stop.
* - `"manual"`: You handle all errors yourself. If you don't catch them, the stream will error. This is primarily for operators who have special error handling requirements.
- *
+ *
* Example: parsing JSON safely
* ```ts
* // Pass-through: errors become ObservableError values (
@@ -158,7 +158,7 @@
* }
* });
* ```
- *
+ *
* Compare to native TransformStream error handling:
* ```ts
* // Native: you must handle errors yourself
@@ -172,12 +172,12 @@
* }
* });
* ```
- *
+ *
* ## Stateful Operators: Remembering Across Chunks
- *
+ *
* Sometimes you need to keep track of things as data flows through, like running totals,
* buffers, or windows. Your `createStatefulOperator` lets you do this easily:
- *
+ *
* ```ts
* // Running sum
* const runningSum = createStatefulOperator({
@@ -187,16 +187,16 @@
* controller.enqueue(state.sum);
* }
* });
- *
+ *
* pipe(
* Observable.from([1, 2, 3]),
* runningSum
* ).subscribe(console.log); // Output: 1, 3, 6
* ```
- *
+ *
* Native TransformStream can't do this as cleanly, you'd have to manage state outside the stream,
* which gets messy, error-prone and annoying real quick.
- *
+ *
* ## Performance and Memory
*
* - **Hot path optimization**: The error handling logic is generated for each operator,
@@ -204,41 +204,50 @@
* - **Memory safety**: Only the functions and state you need are kept alive; everything else
* can be garbage collected.
* - **Streams scale**: You can process gigabytes of data with minimal RAM, and your operators
- * work just as well for infinite streams as for arrays (though arrays have better performance through
+ * work just as well for infinite streams as for arrays (though arrays have better performance through
* their built-in `filter`, `map`, `forEach`, etc..., methods).
- *
+ *
* ## Summary
- *
+ *
* - Operators are like `Array.map`/`filter`, but for async streams of data.
* - You can build pipelines that transform, filter, buffer, or combine data.
* - Error handling is flexible and explicit.
* - Streams make your code scalable and memory-efficient.
* - State is easy to manage for advanced use cases.
* - The helpers make working with streams as easy as working with arrays.
- *
+ *
* @module
*/
-import type { Operator, CreateOperatorOptions, StatefulTransformFunctionOptions, TransformFunctionOptions, TransformStreamOptions, ExcludeError, OperatorErrorMode, TransformHandlerContext } from "./_types.ts";
+import type {
+ CreateOperatorOptions,
+ ExcludeError,
+ Operator,
+ OperatorErrorMode,
+ StatefulTransformFunctionOptions,
+ TransformFunctionOptions,
+ TransformHandlerContext,
+ TransformStreamOptions,
+} from "./_types.ts";
import { injectError, isTransformStreamOptions } from "./utils.ts";
-import { ObservableError, isObservableError } from "../error.ts";
+import { isObservableError, ObservableError } from "../error.ts";
/**
* Creates optimized stream operators with consistent error handling
- *
+ *
* The Web Streams API's TransformStream is powerful but requires boilerplate for
* error handling, lifecycle management, and memory optimization. This function
* eliminates that complexity while providing four error handling strategies:
- *
+ *
* - **pass-through**: Errors become observable values in the stream (default)
* - **ignore**: Silently skip errors and continue processing
* - **throw**: Stop stream immediately on first error
* - **manual**: No automatic error handling - you're in full control
- *
+ *
* Performance: Pre-compiles error handling logic to avoid runtime checks on every chunk.
* Memory: Extracts only needed functions from options to enable garbage collection.
- *
+ *
* @example
* ```ts
* // Stream continues even if mapping fails for some items
@@ -249,16 +258,16 @@ import { ObservableError, isObservableError } from "../error.ts";
* controller.enqueue(fn(chunk)); // If fn() throws, error gets enqueued
* }
* });
- *
+ *
* // Stream stops immediately on any error
* const strictMap = (fn: (x: T) => R) => createOperator({
- * name: 'strictMap',
+ * name: 'strictMap',
* errorMode: 'throw', // Stream terminates on first error
* transform(chunk, controller) {
* controller.enqueue(fn(chunk));
* }
* });
- *
+ *
* // You handle all errors manually
* const customMap = (fn: (x: T) => R) => createOperator({
* name: 'customMap',
@@ -275,56 +284,107 @@ import { ObservableError, isObservableError } from "../error.ts";
* ```
*/
// For "pass-through" error mode - output includes ObservableErrors
-export function createOperator(
- options: TransformFunctionOptions & { errorMode?: "pass-through" }
+export function createOperator<
+ T,
+ R,
+ O extends R | ObservableError = R | ObservableError,
+>(
+ options: TransformFunctionOptions & { errorMode?: "pass-through" },
): Operator;
-export function createOperator(
- options: TransformStreamOptions & { errorMode?: "pass-through" }
+/**
+ * Creates a pass-through operator from a pre-built TransformStream factory.
+ */
+export function createOperator<
+ T,
+ R,
+ O extends R | ObservableError = R | ObservableError,
+>(
+ options: TransformStreamOptions & { errorMode?: "pass-through" },
): Operator;
-// For "ignore" error mode - no ObservableErrors in output
-export function createOperator = ExcludeError>(
- options: TransformFunctionOptions & { errorMode: "ignore" }
+/**
+ * Creates an ignore-mode operator from a transform callback.
+ */
+export function createOperator<
+ T,
+ R,
+ O extends ExcludeError = ExcludeError,
+>(
+ options: TransformFunctionOptions & { errorMode: "ignore" },
): Operator;
-export function createOperator = ExcludeError>(
- options: TransformStreamOptions & { errorMode: "ignore" }
+/**
+ * Creates an ignore-mode operator from a TransformStream factory.
+ */
+export function createOperator<
+ T,
+ R,
+ O extends ExcludeError = ExcludeError,
+>(
+ options: TransformStreamOptions & { errorMode: "ignore" },
): Operator;
-// For "throw" error mode - no ObservableErrors in output
-export function createOperator = ExcludeError>(
- options: TransformFunctionOptions & { errorMode: "throw" }
+/**
+ * Creates a throw-mode operator from a transform callback.
+ */
+export function createOperator<
+ T,
+ R,
+ O extends ExcludeError = ExcludeError,
+>(
+ options: TransformFunctionOptions & { errorMode: "throw" },
): Operator;
-export function createOperator = ExcludeError>(
- options: TransformStreamOptions & { errorMode: "throw" }
+/**
+ * Creates a throw-mode operator from a TransformStream factory.
+ */
+export function createOperator<
+ T,
+ R,
+ O extends ExcludeError = ExcludeError,
+>(
+ options: TransformStreamOptions & { errorMode: "throw" },
): Operator;
-// For "manual" error mode - output is entirely up to the implementation
+/**
+ * Creates a manual-mode operator from a transform callback.
+ */
export function createOperator(
- options: TransformFunctionOptions & { errorMode: "manual" }
+ options: TransformFunctionOptions & { errorMode: "manual" },
): Operator;
+/**
+ * Creates a manual-mode operator from a TransformStream factory.
+ */
export function createOperator(
- options: TransformStreamOptions & { errorMode: "manual" }
+ options: TransformStreamOptions & { errorMode: "manual" },
): Operator;
// Default case
-export function createOperator | ObservableError = R | ExcludeError | ObservableError>(
- options: CreateOperatorOptions
+export function createOperator<
+ T,
+ R,
+ O extends R | ExcludeError | ObservableError =
+ | R
+ | ExcludeError
+ | ObservableError,
+>(
+ options: CreateOperatorOptions,
): Operator {
// Extract operator name from options or the function name for better error reporting
- const operatorName = `operator:${options.name || 'unknown'}`;
- const errorMode = (options as TransformFunctionOptions)?.errorMode ?? "pass-through";
+ const operatorName = `operator:${options.name || "unknown"}`;
+ const errorMode = (options as TransformFunctionOptions)?.errorMode ??
+ "pass-through";
// Extract only what we need to avoid retaining the full options object
const transform = (options as TransformFunctionOptions)?.transform;
const start = (options as TransformFunctionOptions)?.start;
const flush = (options as TransformFunctionOptions)?.flush;
-
+
return (source) => {
try {
// Create a transform stream with the provided options
- const transformStream = isTransformStreamOptions(options) ?
- options.stream(options) :
- new TransformStream({
+ const transformStream = isTransformStreamOptions(options)
+ ? options.stream(options)
+ : new TransformStream(
+ {
// Transform function to process each chunk
transform: handleTransform(errorMode, transform, { operatorName }),
@@ -335,30 +395,32 @@ export function createOperator | ObservableE
flush: handleFlush(errorMode, flush, { operatorName }),
},
{ highWaterMark: 1 },
- { highWaterMark: 0 }
+ { highWaterMark: 0 },
);
-
+
// Pipe the source through the transform
return source.pipeThrough(transformStream);
} catch (err) {
// If setup fails, return a stream that errors immediately
- return source.pipeThrough(injectError(err, `${operatorName}:setup`, options));
+ return source.pipeThrough(
+ injectError(err, `${operatorName}:setup`, options),
+ );
}
};
}
/**
* Hot-path optimized error handling for transform functions
- *
+ *
* Problem: Transform functions are called for EVERY chunk in a stream. Doing
* error mode checks and type checks on every call kills performance.
- *
+ *
* Solution: Pre-compile the error handling logic into optimized functions.
* Each error mode gets its own specialized function with zero runtime overhead.
- *
+ *
* Memory optimization: Only references the specific transform function and state,
* not the entire options object, enabling garbage collection of unused properties.
- *
+ *
* @example
* ```ts
* // Instead of this slow approach:
@@ -370,7 +432,7 @@ export function createOperator | ObservableE
* }
* // ... more runtime checks
* }
- *
+ *
* // handleTransform pre-compiles to this:
* function fastIgnoreTransform(chunk, controller) {
* if (isObservableError(chunk)) return; // Only one check needed
@@ -379,7 +441,7 @@ export function createOperator | ObservableE
* } catch (_) { return; } // Pre-compiled error handling
* }
* ```
- *
+ *
* @typeParam T - Input chunk type
* @typeParam O - Output chunk type
* @typeParam S - State type (for stateful operators)
@@ -390,97 +452,143 @@ export function createOperator | ObservableE
*/
export function handleTransform(
errorMode: OperatorErrorMode,
- transform:
- TransformFunctionOptions['transform'] |
- StatefulTransformFunctionOptions['transform'],
- context: TransformHandlerContext = { }
-): Transformer['transform'] {
+ transform:
+ | TransformFunctionOptions["transform"]
+ | StatefulTransformFunctionOptions["transform"],
+ context: TransformHandlerContext = {},
+): Transformer["transform"] {
const operatorName = context.operatorName || `operator:unknown`;
const isStateful = context.isStateful || false;
const state = context.state;
switch (errorMode) {
case "pass-through":
- return async function (chunk: T, controller: TransformStreamDefaultController) {
+ return async function (
+ chunk: T,
+ controller: TransformStreamDefaultController,
+ ) {
if (isObservableError(chunk)) {
controller.enqueue(chunk as O);
return;
}
-
+
try {
if (isStateful) {
// If stateful, pass the state along
- return await (transform as StatefulTransformFunctionOptions['transform'])(chunk, state as S, controller);
+ return await (transform as StatefulTransformFunctionOptions<
+ T,
+ O,
+ S
+ >["transform"])(chunk, state as S, controller);
}
- await (transform as TransformFunctionOptions['transform'])(chunk, controller);
+ await (transform as TransformFunctionOptions["transform"])(
+ chunk,
+ controller,
+ );
} catch (err) {
- controller.enqueue(ObservableError.from(err, operatorName, chunk) as O);
+ controller.enqueue(
+ ObservableError.from(err, operatorName, chunk) as O,
+ );
}
};
-
+
case "ignore":
- return async function (chunk: T, controller: TransformStreamDefaultController) {
+ return async function (
+ chunk: T,
+ controller: TransformStreamDefaultController