diff --git a/CHANGELOG.md b/CHANGELOG.md index a515839..792c602 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,66 +7,115 @@ # 1.0.0 (2026-03-18) - ### Bug Fixes -* a large number of tests ([5d0a6c5](https://github.com/okikio/observables/commit/5d0a6c58cd1bcfdda07a89c2c0007691889aeabb)) -* broken tests ([2ab9031](https://github.com/okikio/observables/commit/2ab9031b95ce97935823516a46d949888af99157)) -* **ci:** format events test for actions ([885af48](https://github.com/okikio/observables/commit/885af48ab08a9483f6c9abf768f794d3738879f6)) -* **ci:** restore import-prefix lint ignores ([848d2e0](https://github.com/okikio/observables/commit/848d2e03ecf479b1434ef323ca3c7c3d2485a009)) -* deviate from spec. on error handling for observer.next/complete/error ([98f2537](https://github.com/okikio/observables/commit/98f2537a01d07a05fcd6fd3bd37284a4651e3624)) -* edge case types in timing, pipe and core operations ([7fc9324](https://github.com/okikio/observables/commit/7fc93242aca9ee003e4b4e4a09e9f58e4f7b18c6)) -* errors breaking readable streams in pull function ([44e9c0f](https://github.com/okikio/observables/commit/44e9c0f0d8b4076e8e2f0f23bc02d1fe7354e2af)) -* even more tests ([4d7c8ba](https://github.com/okikio/observables/commit/4d7c8baf5204d88d22ace51c166255bd35e60ff1)) -* events tests ([0ebac08](https://github.com/okikio/observables/commit/0ebac08902a7a0201fe14fae6741041bbe74e857)) -* **events,operators:** revert empty-subscriber fast path and use controller.error ([f6b8d69](https://github.com/okikio/observables/commit/f6b8d6996839e9ef5998978369fce8a97b730b17)) -* fix incorrect logic for replaying item ([ad8cca5](https://github.com/okikio/observables/commit/ad8cca5506066941bb5f8b2e28f90428036d3766)) -* fix issues with observables not catch errors ([66b2948](https://github.com/okikio/observables/commit/66b29482c82c550d282e37581c4eabef4c999320)) -* fix major slow api types blocking publishing ([bd84ffb](https://github.com/okikio/observables/commit/bd84ffbbb1532d8a17cb42a25ade255dbb691d44)) -* fix operator types ([27d12e8](https://github.com/okikio/observables/commit/27d12e8fcb01356a704e23a4a2a9d514c5a3f09d)) -* improve type check of assertObservableError ([f077868](https://github.com/okikio/observables/commit/f07786805a3fac1e34984f70e6ff0e5245624a0b)) -* issues w/ delay timing ([b983ede](https://github.com/okikio/observables/commit/b983ede6714df9fc9ad14fffe1f5fe01c58a31a0)) -* non-compliance with some specific tc39 observables steps ([d5c8318](https://github.com/okikio/observables/commit/d5c83184ff46d9a4ab30adeb4781a4c7f234da9e)) -* **observable:** document bound observer methods ([b7cffa5](https://github.com/okikio/observables/commit/b7cffa50b47778137123df617859fab67e4c2297)) -* **observable:** drop detached observer binding ([f700b40](https://github.com/okikio/observables/commit/f700b40e2adf904c67ad634f223d0607de56d403)) -* **operators:** restore return short-circuit after controller.error ([f5eac66](https://github.com/okikio/observables/commit/f5eac66d370abb0748a544ea28d77e6ca6522743)) -* rejection message for waitForEvent ([c559d5f](https://github.com/okikio/observables/commit/c559d5f19ad1fa7a609af1548c13b02a2cfce868)) -* **release:** tighten validation and green ci ([95e4d4d](https://github.com/okikio/observables/commit/95e4d4d0ab8ebc6e268b6500892623183246d575)) -* remove queued teardown ([8b75e82](https://github.com/okikio/observables/commit/8b75e82536eb44f4538e75c0d73e9117ca5a2ae2)) -* **repo:** align benchmark tasks and test harnesses ([c72229c](https://github.com/okikio/observables/commit/c72229c6d8630f029b1c6e4ba9426cbde47fe767)) -* **review:** address publishing follow-up comments ([7959911](https://github.com/okikio/observables/commit/79599114d59d61bf1add0b1f885ef50a35a49e44)) -* split close subscription into mark subscription as closed and perform close subscription ([3996f5f](https://github.com/okikio/observables/commit/3996f5f12636281e004d208ec44c124169eca628)) -* **tests:** align BDD coverage with observable error and stream semantics ([d81486c](https://github.com/okikio/observables/commit/d81486c2443842ae5480aa5f9093d997aec48265)) -* **tests:** polish spec coverage naming ([e46a100](https://github.com/okikio/observables/commit/e46a100e35602b2cfa7b4d67a248b44c4813ce64)) -* **tests:** preserve standalone assertion messages ([c2ce404](https://github.com/okikio/observables/commit/c2ce404206f607a0a07e9d777602afc530c6b409)) -* **tests:** resolve remaining review-thread lint and flake issues ([cab3df1](https://github.com/okikio/observables/commit/cab3df119e1ef56e3d50418c60d8b54c5d90c8c5)) -* type issues with pipe, compose and safeCompose ([add210e](https://github.com/okikio/observables/commit/add210eeca72afe59992113830452d175d499cd7)) -* update tests to take into account new fixes ([812d5e5](https://github.com/okikio/observables/commit/812d5e526dd8a4c903f056e8bd463e0cc317e895)) -* **WIP:** fix 1/2 of the typescript type issues in dom_event_tests ([1da2571](https://github.com/okikio/observables/commit/1da2571b4b0f3df854eefd6af0008b39ed2d97b3)) - +- a large number of tests + ([5d0a6c5](https://github.com/okikio/observables/commit/5d0a6c58cd1bcfdda07a89c2c0007691889aeabb)) +- broken tests + ([2ab9031](https://github.com/okikio/observables/commit/2ab9031b95ce97935823516a46d949888af99157)) +- **ci:** format events test for actions + ([885af48](https://github.com/okikio/observables/commit/885af48ab08a9483f6c9abf768f794d3738879f6)) +- **ci:** restore import-prefix lint ignores + ([848d2e0](https://github.com/okikio/observables/commit/848d2e03ecf479b1434ef323ca3c7c3d2485a009)) +- deviate from spec. on error handling for observer.next/complete/error + ([98f2537](https://github.com/okikio/observables/commit/98f2537a01d07a05fcd6fd3bd37284a4651e3624)) +- edge case types in timing, pipe and core operations + ([7fc9324](https://github.com/okikio/observables/commit/7fc93242aca9ee003e4b4e4a09e9f58e4f7b18c6)) +- errors breaking readable streams in pull function + ([44e9c0f](https://github.com/okikio/observables/commit/44e9c0f0d8b4076e8e2f0f23bc02d1fe7354e2af)) +- even more tests + ([4d7c8ba](https://github.com/okikio/observables/commit/4d7c8baf5204d88d22ace51c166255bd35e60ff1)) +- events tests + ([0ebac08](https://github.com/okikio/observables/commit/0ebac08902a7a0201fe14fae6741041bbe74e857)) +- **events,operators:** revert empty-subscriber fast path and use + controller.error + ([f6b8d69](https://github.com/okikio/observables/commit/f6b8d6996839e9ef5998978369fce8a97b730b17)) +- fix incorrect logic for replaying item + ([ad8cca5](https://github.com/okikio/observables/commit/ad8cca5506066941bb5f8b2e28f90428036d3766)) +- fix issues with observables not catch errors + ([66b2948](https://github.com/okikio/observables/commit/66b29482c82c550d282e37581c4eabef4c999320)) +- fix major slow api types blocking publishing + ([bd84ffb](https://github.com/okikio/observables/commit/bd84ffbbb1532d8a17cb42a25ade255dbb691d44)) +- fix operator types + ([27d12e8](https://github.com/okikio/observables/commit/27d12e8fcb01356a704e23a4a2a9d514c5a3f09d)) +- improve type check of assertObservableError + ([f077868](https://github.com/okikio/observables/commit/f07786805a3fac1e34984f70e6ff0e5245624a0b)) +- issues w/ delay timing + ([b983ede](https://github.com/okikio/observables/commit/b983ede6714df9fc9ad14fffe1f5fe01c58a31a0)) +- non-compliance with some specific tc39 observables steps + ([d5c8318](https://github.com/okikio/observables/commit/d5c83184ff46d9a4ab30adeb4781a4c7f234da9e)) +- **observable:** document bound observer methods + ([b7cffa5](https://github.com/okikio/observables/commit/b7cffa50b47778137123df617859fab67e4c2297)) +- **observable:** drop detached observer binding + ([f700b40](https://github.com/okikio/observables/commit/f700b40e2adf904c67ad634f223d0607de56d403)) +- **operators:** restore return short-circuit after controller.error + ([f5eac66](https://github.com/okikio/observables/commit/f5eac66d370abb0748a544ea28d77e6ca6522743)) +- rejection message for waitForEvent + ([c559d5f](https://github.com/okikio/observables/commit/c559d5f19ad1fa7a609af1548c13b02a2cfce868)) +- **release:** tighten validation and green ci + ([95e4d4d](https://github.com/okikio/observables/commit/95e4d4d0ab8ebc6e268b6500892623183246d575)) +- remove queued teardown + ([8b75e82](https://github.com/okikio/observables/commit/8b75e82536eb44f4538e75c0d73e9117ca5a2ae2)) +- **repo:** align benchmark tasks and test harnesses + ([c72229c](https://github.com/okikio/observables/commit/c72229c6d8630f029b1c6e4ba9426cbde47fe767)) +- **review:** address publishing follow-up comments + ([7959911](https://github.com/okikio/observables/commit/79599114d59d61bf1add0b1f885ef50a35a49e44)) +- split close subscription into mark subscription as closed and perform close + subscription + ([3996f5f](https://github.com/okikio/observables/commit/3996f5f12636281e004d208ec44c124169eca628)) +- **tests:** align BDD coverage with observable error and stream semantics + ([d81486c](https://github.com/okikio/observables/commit/d81486c2443842ae5480aa5f9093d997aec48265)) +- **tests:** polish spec coverage naming + ([e46a100](https://github.com/okikio/observables/commit/e46a100e35602b2cfa7b4d67a248b44c4813ce64)) +- **tests:** preserve standalone assertion messages + ([c2ce404](https://github.com/okikio/observables/commit/c2ce404206f607a0a07e9d777602afc530c6b409)) +- **tests:** resolve remaining review-thread lint and flake issues + ([cab3df1](https://github.com/okikio/observables/commit/cab3df119e1ef56e3d50418c60d8b54c5d90c8c5)) +- type issues with pipe, compose and safeCompose + ([add210e](https://github.com/okikio/observables/commit/add210eeca72afe59992113830452d175d499cd7)) +- update tests to take into account new fixes + ([812d5e5](https://github.com/okikio/observables/commit/812d5e526dd8a4c903f056e8bd463e0cc317e895)) +- **WIP:** fix 1/2 of the typescript type issues in dom_event_tests + ([1da2571](https://github.com/okikio/observables/commit/1da2571b4b0f3df854eefd6af0008b39ed2d97b3)) ### Features -* add AbortSignal support ([299b8fe](https://github.com/okikio/observables/commit/299b8fee6ef891eedefc410339a1bad8508d0cba)) -* add helpers & operators ([7fff07d](https://github.com/okikio/observables/commit/7fff07d4b83044c8f884185acda435ea59b62f2d)) -* add isObservable + isSpecObservable methods ([faefeb4](https://github.com/okikio/observables/commit/faefeb4c751de144c8dcb7cc861d5dfb698ca956)) -* add isObservableError method ([4be5fe5](https://github.com/okikio/observables/commit/4be5fe54f9f1fa17870e0e27affbcc4d82b53d6b)) -* add observable + event bus ([75fa4ed](https://github.com/okikio/observables/commit/75fa4ed98aad7106c0ccdbcb3b4b3985e995d411)) -* add operator for throwing errors ([a214329](https://github.com/okikio/observables/commit/a21432907f68653254659622240ec543d4bf208d)) -* add the values operators, that only operate directly on values while letting errors passthrough ([7b9823d](https://github.com/okikio/observables/commit/7b9823d348a6db38ecc3aff987656da4fba27f3e)) -* add tips to ObservableError ([1d25963](https://github.com/okikio/observables/commit/1d25963e80a43c7c95cd4e44ed6d89108c8e7208)) -* add withReplay utility function ([50881d7](https://github.com/okikio/observables/commit/50881d77df674d97a15dec7e82eef8734aaecfba)) -* **bench:** add comparison benchmarks and spec coverage ([bd2bde8](https://github.com/okikio/observables/commit/bd2bde8e2efde8bb60055590182015ec31ed26e8)) -* **docs:** add comprehensive instruction files and copilot-instructions.md ([3112119](https://github.com/okikio/observables/commit/3112119e301e6054515e8d806f9b9b1130c15be1)) -* implement helpers using streams internally ([e1f5bea](https://github.com/okikio/observables/commit/e1f5beab399dbff213d969faef84f92db57bb047)) -* support passing ObservableError as part of the iterator of the `pull` function ([1ab2355](https://github.com/okikio/observables/commit/1ab2355a66162a5a8435176231323a1922cf0c47)) - +- add AbortSignal support + ([299b8fe](https://github.com/okikio/observables/commit/299b8fee6ef891eedefc410339a1bad8508d0cba)) +- add helpers & operators + ([7fff07d](https://github.com/okikio/observables/commit/7fff07d4b83044c8f884185acda435ea59b62f2d)) +- add isObservable + isSpecObservable methods + ([faefeb4](https://github.com/okikio/observables/commit/faefeb4c751de144c8dcb7cc861d5dfb698ca956)) +- add isObservableError method + ([4be5fe5](https://github.com/okikio/observables/commit/4be5fe54f9f1fa17870e0e27affbcc4d82b53d6b)) +- add observable + event bus + ([75fa4ed](https://github.com/okikio/observables/commit/75fa4ed98aad7106c0ccdbcb3b4b3985e995d411)) +- add operator for throwing errors + ([a214329](https://github.com/okikio/observables/commit/a21432907f68653254659622240ec543d4bf208d)) +- add the values operators, that only operate directly on values while letting + errors passthrough + ([7b9823d](https://github.com/okikio/observables/commit/7b9823d348a6db38ecc3aff987656da4fba27f3e)) +- add tips to ObservableError + ([1d25963](https://github.com/okikio/observables/commit/1d25963e80a43c7c95cd4e44ed6d89108c8e7208)) +- add withReplay utility function + ([50881d7](https://github.com/okikio/observables/commit/50881d77df674d97a15dec7e82eef8734aaecfba)) +- **bench:** add comparison benchmarks and spec coverage + ([bd2bde8](https://github.com/okikio/observables/commit/bd2bde8e2efde8bb60055590182015ec31ed26e8)) +- **docs:** add comprehensive instruction files and copilot-instructions.md + ([3112119](https://github.com/okikio/observables/commit/3112119e301e6054515e8d806f9b9b1130c15be1)) +- implement helpers using streams internally + ([e1f5bea](https://github.com/okikio/observables/commit/e1f5beab399dbff213d969faef84f92db57bb047)) +- support passing ObservableError as part of the iterator of the `pull` function + ([1ab2355](https://github.com/okikio/observables/commit/1ab2355a66162a5a8435176231323a1922cf0c47)) ### Performance Improvements -* add cap on replay size of 1000 ([1d52e44](https://github.com/okikio/observables/commit/1d52e446aef05bbaf1f97e8a82707355de3de616)) -* **bench:** add latency benchmarks and speed up queue wraparound ([5324282](https://github.com/okikio/observables/commit/53242822f96de446c298a00198438cb4e48da895)) -* optimize hot-paths ([91a6313](https://github.com/okikio/observables/commit/91a6313402b74f56f1feaa9609790c03927f402e)) +- add cap on replay size of 1000 + ([1d52e44](https://github.com/okikio/observables/commit/1d52e446aef05bbaf1f97e8a82707355de3de616)) +- **bench:** add latency benchmarks and speed up queue wraparound + ([5324282](https://github.com/okikio/observables/commit/53242822f96de446c298a00198438cb4e48da895)) +- optimize hot-paths + ([91a6313](https://github.com/okikio/observables/commit/91a6313402b74f56f1feaa9609790c03927f402e)) # Changelog diff --git a/README.md b/README.md index 43a477a..3e5d13f 100644 --- a/README.md +++ b/README.md @@ -1,25 +1,71 @@ # @okikio/observables -[![Open Bundle](https://bundlejs.com/badge-light.svg)](https://bundlejs.com/?q=@okikio/observables&bundle "Check the total bundle size of @okikio/observables") +[![CI](https://github.com/okikio/observables/actions/workflows/ci.yml/badge.svg)](https://github.com/okikio/observables/actions/workflows/ci.yml) +[![npm version](https://img.shields.io/npm/v/%40okikio%2Fobservables?logo=npm&label=npm)](https://www.npmjs.com/package/@okikio/observables) +[![Bundle Size](https://deno.bundlejs.com/badge?q=@okikio/observables&treeshake=[{+Observable,+pipe,+map,+filter+}]&style=flat)](https://bundlejs.com/?q=@okikio/observables&treeshake=[{+Observable,+pipe,+map,+filter+}]) +[![License: MIT](https://img.shields.io/badge/license-MIT-blue.svg)](./LICENSE) -[NPM](https://www.npmjs.com/package/@okikio/observables) -| -[GitHub](https://github.com/okikio/observables#readme) -| -[JSR](https://jsr.io/@okikio/observables) -| [Licence](./LICENSE) +[Documentation](https://jsr.io/@okikio/observables) • +[npm](https://www.npmjs.com/package/@okikio/observables) • +[GitHub](https://github.com/okikio/observables#readme) • [License](./LICENSE) A **spec-faithful** yet ergonomic TC39-inspired Observable implementation that gives you one consistent way to handle all async data in JavaScript. +Built for Deno v2+, Node, Bun, and modern browsers, `@okikio/observables` keeps +the TC39 Observable proposal's mental model while adding the parts that make +day-to-day app code easier to write: + +- **Observable pipelines that feel familiar** if you already know `Array.map()` + and `Array.filter()` +- **Web Streams-powered backpressure** so fast producers do not silently bloat + memory +- **Deterministic cleanup** via `unsubscribe()`, `using`, and `Symbol.dispose` +- **Built-in event primitives** for pub/sub and type-safe event dispatch +- **Four error modes** so you can choose between recovery, filtering, and + fail-fast behavior + +**Start here:** [Installation](#installation) • [Quick Start](#quick-start) • +[API](#api) • [Advanced Usage](#advanced-usage) • [FAQ](#faq) • +[Contributing](#contributing) + +## Start Here + +Install with the package manager that matches your runtime: + +```bash +# Deno / JSR +deno add jsr:@okikio/observables + +# npm-compatible runtimes +npm install @okikio/observables +# pnpm add @okikio/observables +# yarn add @okikio/observables +# bun add @okikio/observables +``` + +Then build a small pipeline: + +```ts +import { filter, map, Observable, pipe } from "@okikio/observables"; + +const values = pipe( + Observable.of(1, 2, 3, 4), + filter((value) => value % 2 === 0), + map((value) => value * 10), +); + +values.subscribe((value) => console.log(value)); +// 20 +// 40 +``` + **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. -[![Bundle Size](https://deno.bundlejs.com/badge?q=@okikio/observables&treeshake=[{+Observable,+pipe,+map,+filter+}]&style=flat)](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 @@ -483,40 +529,48 @@ 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). +Contributions are welcome. This project targets Deno v2+ and keeps a tight +feedback loop around formatting, linting, docs, tests, and npm packaging, so a +good contribution usually starts by getting the local validation commands +working first. + +Install Deno with [mise](https://mise.jdx.dev/) or by following the +[manual installation guide](https://deno.land/manual/getting_started/installation). -Setup Mise: +### Setup with Mise ```bash curl https://mise.run | sh echo 'eval "$(~/.local/bin/mise activate bash)"' >> ~/.bashrc ``` -Install Deno: +Then install the toolchain: ```bash mise install ``` -Then run tests: +### Validate your change ```bash +deno fmt +deno lint +deno check **/*.ts +deno doc --lint mod.ts deno task test +deno task build:npm ``` -Run benchmarks: +If your change is performance-sensitive, also run: ```bash 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. +This repository uses +[Conventional Commits](https://www.conventionalcommits.org/en/v1.0.0/), so +please format commit messages accordingly. -## Licence +## License See the [LICENSE](./LICENSE) file for license rights and limitations (MIT). diff --git a/_types.ts b/_types.ts index cb091cd..2e64153 100644 --- a/_types.ts +++ b/_types.ts @@ -2,13 +2,21 @@ /** * Public observer and subscription types for the Observable entrypoints. * - * @module - * * This module collects the runtime-facing interfaces that show up whenever you * subscribe to an Observable. It defines the enhanced `Observer` and * `Subscription` shapes used by this package, then re-exports the lower-level * spec types from `./_spec.ts` for callers that need proposal-aligned building + * spec types from `./_spec.ts` for callers that need proposal-aligned building * blocks. + * + * Use this entrypoint when you are writing libraries, adapters, or tests that + * need to talk about Observable contracts without importing the full runtime + * implementation. In day-to-day app code, you will usually consume these types + * indirectly through `Observable`, `EventBus`, or operator helpers. They live + * here so the public type surface stays easy to find and stable across + * entrypoints. + * + * @module */ import type { SpecObserver, SpecSubscription } from "./_spec.ts"; import type { Symbol } from "./symbol.ts"; diff --git a/error.ts b/error.ts index 755a67b..9997c7d 100644 --- a/error.ts +++ b/error.ts @@ -1,6 +1,16 @@ // @filename: error.ts /** - * Error handling utilities for Observable operators + * Error primitives and guards for Observable pipelines. + * + * This entrypoint explains the error values that travel through this library + * when an operator uses pass-through error handling. It exports the + * `ObservableError` class plus helper functions for narrowing and asserting + * those wrapped failures without losing the original error object, stack, or + * operator context. + * + * Reach for this module when you want to inspect failures as data, recover from + * an upstream step without throwing away buffered values, or surface richer + * debugging information than a plain `Error` can carry on its own. * * @module */ diff --git a/events.ts b/events.ts index 155f1d2..712af4f 100644 --- a/events.ts +++ b/events.ts @@ -1,5 +1,17 @@ /** - * @module EventBus + * Multicast event primitives built on top of the Observable runtime. + * + * This entrypoint is for the hot side of the library: shared event streams that + * multiple consumers can listen to at the same time. It exports `EventBus` for + * one-to-many pub/sub, `createEventDispatcher` for named and type-safe events, + * and helpers such as `withReplay` and `waitForEvent` for common coordination + * patterns. + * + * Use `Observable` when each subscription should start fresh work. Use this + * module when one emission should fan out to many listeners, such as UI events, + * app-wide notifications, or workflow status updates. + * + * @module */ import type { Observer, Subscription } from "./_types.ts"; diff --git a/helpers/mod.ts b/helpers/mod.ts index 445542e..6b1030f 100644 --- a/helpers/mod.ts +++ b/helpers/mod.ts @@ -1,27 +1,24 @@ // @filename: helpers/mod.ts /** - * Observable Operators Library - * - * @module - * - * This library provides a collection of operators for working with Observables. - * It enables functional composition of Observable transformations using the `pipe` and - * `compose` functions, with each operator implemented using Web Streams for efficiency. - * - * ## Core Features - * - * - **Observable-based API**: Clean, familiar API that works with Observables - * - **Stream-based implementation**: Uses Web Streams internally for efficiency and backpressure - * - **Functional**: Pure functions for easy composition - * - **Type-safe**: Full TypeScript support with proper type inference - * - **Tree-shakable**: Import only what you need - * - * ## Architecture - * - * This library uses a hybrid approach: - * - The public API (`pipe` function, creation functions, etc.) works with Observables - * - Internally, it converts Observables to Web Streams, applies transformations, then converts back - * - This allows for the efficiency of streams while maintaining a familiar Observable API + * High-level operator entrypoint for composing Observable pipelines. + * + * This is the ergonomic "import most things from one place" surface for the + * package. It re-exports the pipe/compose helpers, operator builders, utility + * helpers, and the built-in operator families so callers can build an entire + * pipeline without remembering which category module each operator lives in. + * + * The exports here fall into a few broad groups: + * - core transformation operators such as `map`, `filter`, `take`, and `scan` + * - timing operators such as `debounce`, `delay`, `throttle`, and `timeout` + * - combination operators such as `mergeMap`, `concatMap`, and `switchMap` + * - batching and error operators for collecting values or recovering from + * failures in a stream + * + * Everything shares the same Web Streams-based runtime, so backpressure, + * teardown, and error-mode behavior stay consistent from one operator to the + * next. Import from `./operations/*` only when you want a narrower entrypoint + * for discovery or tree-shaken docs. Import from this module when you want the + * familiar "just give me the operator toolbox" experience. * * ## Basic Usage * @@ -82,6 +79,8 @@ * - Each pipeline is limited to 9 operators due to TypeScript's recursion limits * - Use `compose` to group operators when you need more than 9 * - Composed operator groups are also limited to 9 operators + * + * @module */ // Re-export all operators from their respective modules diff --git a/helpers/operations/batch.ts b/helpers/operations/batch.ts index 0c57c13..3c444cb 100644 --- a/helpers/operations/batch.ts +++ b/helpers/operations/batch.ts @@ -1,3 +1,17 @@ +/** + * Operators that collect multiple source values into grouped results. + * + * This entrypoint contains the batching side of the operator library. Use it + * when one output value should summarize several input values, such as turning + * a whole stream into a single array with `toArray()` or bundling incoming + * items into fixed-size chunks with `batch()`. + * + * These operators trade immediacy for aggregation. They usually hold on to some + * values until a batch fills or the source completes, so they are best for + * finite streams or bounded buffers where that extra memory is intentional. + * + * @module + */ import type { ExcludeError, Operator } from "../_types.ts"; import type { ObservableError } from "../../error.ts"; diff --git a/helpers/operations/combination.ts b/helpers/operations/combination.ts index 5b69f0d..06dbbc9 100644 --- a/helpers/operations/combination.ts +++ b/helpers/operations/combination.ts @@ -1,6 +1,19 @@ -// helpers/combination.ts -// Operators that combine Observables or transform to new Observables -// Reimplemented using createStatefulOperator for better compatibility with streaming pipeline +/** + * Operators that turn each source value into another stream and combine the + * results. + * + * This entrypoint covers the "flattening" family of operators such as + * `mergeMap`, `concatMap`, and `switchMap`. Use it when one input value needs + * to start follow-up async work, for example fetching related records, reading + * files, or switching to the latest search request. + * + * The operators in this module mainly differ in concurrency and cancellation + * behavior. `mergeMap` keeps multiple inner streams alive at once, `concatMap` + * preserves order by running one at a time, and `switchMap` cancels older work + * when a newer source value arrives. + * + * @module + */ import type { SpecObservable } from "../../_spec.ts"; import type { ExcludeError, Operator } from "../_types.ts"; diff --git a/helpers/operations/conditional.ts b/helpers/operations/conditional.ts index 4186aa3..eebdb76 100644 --- a/helpers/operations/conditional.ts +++ b/helpers/operations/conditional.ts @@ -1,3 +1,17 @@ +/** + * Predicate and decision-oriented operators for Observable streams. + * + * This entrypoint exports the operators that answer questions about a stream or + * gate values based on a condition. These are the Observable equivalents of + * array helpers such as `every()`, `some()`, and `find()`, plus utilities that + * stop early once a decision has been reached. + * + * Reach for this module when you care about whether a stream contains a match, + * whether every value passes a rule, or when processing should stop as soon as + * the answer is known. + * + * @module + */ import type { ExcludeError, Operator } from "../_types.ts"; import type { ObservableError } from "../../error.ts"; diff --git a/helpers/operations/core.ts b/helpers/operations/core.ts index 1d1ac84..68b0ad5 100644 --- a/helpers/operations/core.ts +++ b/helpers/operations/core.ts @@ -4,41 +4,19 @@ import type { ObservableError } from "../../error.ts"; import { createOperator, createStatefulOperator } from "../operators.ts"; /** - * @module operations/core + * Core transformation and terminal operators for everyday stream work. * - * **Core Stream Operators - Like Array Methods, But Error-Aware** + * This module is the closest match to familiar array helpers. It exports the + * operators you reach for first when you want to transform values, filter them, + * accumulate state, or stop after a condition has been met. In practice, this + * is where most pipelines start before you add timing or concurrency behavior. * - * These operators work just like Array methods you already know, but automatically handle errors: + * The important difference from arrays is error handling. These operators are + * designed to work with the library's pass-through model, so your callbacks see + * clean data values while `ObservableError` instances continue downstream + * unchanged until a dedicated error-handling step decides what to do with them. * - * ```ts - * // Array methods: - * [1, 2, 3].map(n => n * 2).filter(n => n > 3) // [4, 6] - * - * // Stream operators (same API, error-aware): - * pipe( - * [1, 2, 3], - * map(n => n * 2), // Transforms data, passes errors through - * filter(n => n > 3) // Filters data, keeps all errors - * ) // Stream: [4, 6] + any errors - * ``` - * - * ## How Error Handling Works - * - * - **Data items**: Get processed by your functions normally - * - **ObservableError items**: Skip processing and flow through unchanged - * - **Function crashes**: Get wrapped in ObservableError automatically - * - * ```ts - * // Clean code - no manual error checking needed - * const users = pipe( - * userIds, - * map(id => fetchUser(id)), // Some API calls fail → ObservableError - * filter(user => user.isActive), // Only filters real users - * take(10) // Takes 10 real users, errors flow alongside - * ); - * ``` - * - * Your functions only receive clean data, never errors. Errors are preserved for later handling. + * @module */ /** diff --git a/helpers/operations/errors.ts b/helpers/operations/errors.ts index aa6ad77..e5b8bf7 100644 --- a/helpers/operations/errors.ts +++ b/helpers/operations/errors.ts @@ -1,3 +1,17 @@ +/** + * Error-focused operators for recovering from or reshaping stream failures. + * + * This entrypoint is for pipelines that expect some work to fail and want to + * keep going. It exports helpers for dropping wrapped errors, mapping them to + * fallback values, logging them, or converting them into a shape that fits the + * rest of the pipeline. + * + * These operators are most useful with the library's pass-through error mode, + * where failures travel as `ObservableError` values instead of immediately + * terminating the whole stream. + * + * @module + */ import type { ExcludeError, Operator } from "../_types.ts"; import { isObservableError, ObservableError } from "../../error.ts"; import { createOperator, createStatefulOperator } from "../operators.ts"; diff --git a/helpers/operations/mod.ts b/helpers/operations/mod.ts index 173df5f..232eb2d 100644 --- a/helpers/operations/mod.ts +++ b/helpers/operations/mod.ts @@ -1,12 +1,20 @@ /** - * Re-exported operator groups for the `./operations` entrypoint. + * Category-level entrypoint for the built-in Observable operations. * - * @module + * This module gathers every operator category that powers the higher-level + * `./operators` entrypoint. It is useful when you want a focused import path + * for documentation and discovery, but still want access to the full built-in + * operator set from one module. + * + * The re-exports are grouped by job: + * - `./core` covers the array-like transformations and terminal operators + * - `./timing` covers time-based coordination such as debounce and timeout + * - `./combination` covers flattening and concurrency helpers such as + * `mergeMap`, `concatMap`, and `switchMap` + * - `./batch`, `./conditional`, and `./errors` cover collection, predicate, and + * recovery-focused utilities * - * This module gives callers one place to import the built-in Observable - * operations without knowing which submodule each operator lives in. It - * forwards the batch, combination, conditional, core, error, and timing - * operator groups that make up the public `./operations` surface. + * @module */ export * from "./batch.ts"; export * from "./combination.ts"; diff --git a/helpers/operations/timing.ts b/helpers/operations/timing.ts index 5968b1c..60ebd20 100644 --- a/helpers/operations/timing.ts +++ b/helpers/operations/timing.ts @@ -1,3 +1,17 @@ +/** + * Time-based operators for spacing, delaying, and expiring stream values. + * + * This entrypoint groups the operators that make time part of your pipeline's + * behavior. Use it for UI patterns such as debouncing search input, throttling + * bursty events, delaying retries, or timing out work that takes too long. + * + * Timing operators do not just change values; they change when work is allowed + * to move downstream. That makes them especially important for coordinating + * async side effects without piling up stale requests or overwhelming slower + * consumers. + * + * @module + */ import type { ExcludeError, Operator } from "../_types.ts"; import { createOperator, createStatefulOperator } from "../operators.ts"; import { isObservableError, ObservableError } from "../../error.ts"; diff --git a/queue.ts b/queue.ts index 20f16e8..1e0b5f4 100644 --- a/queue.ts +++ b/queue.ts @@ -1,15 +1,20 @@ /** - * A lightweight, high-performance queue that delivers O(1) operations - * using a circular buffer approach. Designed for readability and ease of use - * without sacrificing performance. + * A lightweight circular-buffer queue for the library's hot paths. * - * Perfect for task queues, message buffers, or any scenario where you need - * fast FIFO (First-In-First-Out) operations without the performance penalty - * of Array.shift(). + * This entrypoint exposes the small FIFO data structure that backs replay, + * buffering, and event fan-out inside the Observable runtime. It keeps + * enqueue/dequeue work O(1) by moving `head` and `tail` pointers around a fixed + * array instead of repeatedly calling `Array.shift()`, which has to move every + * remaining element one slot to the left. + * + * Use this module when you need predictable queue performance and explicit + * capacity control. It is especially useful for buffers that grow and shrink + * frequently, where repeated array reindexing would turn steady traffic into an + * unnecessary O(n) cost. * * @example * ``` - * import { createQueue, enqueue, dequeue, peek } from './simple-queue'; + * import { createQueue, enqueue, dequeue, peek } from './queue.ts'; * * const taskQueue = createQueue(100); // capacity of 100 * enqueue(taskQueue, 'process-order-123'); // add task @@ -20,6 +25,8 @@ * * clear(taskQueue); // empty the queue instantly * ``` + * + * @module */ ///////////////////////