diff --git a/bench/interop_bench.ts b/bench/interop_bench.ts new file mode 100644 index 0000000..0f5ae5c --- /dev/null +++ b/bench/interop_bench.ts @@ -0,0 +1,214 @@ +// deno-lint-ignore-file no-import-prefix +/** + * Interop-focused benchmarks. + * + * These scenarios measure the cost of the new adapter helpers against the + * closest direct baseline so interop convenience does not hide avoidable + * overhead. + */ + +import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34"; +import { + from as rxFrom, + map as rxMap, + type Observable as RxObservable, + pipe as rxPipe, + take as rxTake, +} from "npm:rxjs@7.8.2"; + +import { isObservableError } from "../error.ts"; +import { ignoreErrors } from "../helpers/operations/errors.ts"; +import { map, take } from "../helpers/operations/core.ts"; +import { pipe } from "../helpers/pipe.ts"; +import { + applyOperator, + fromObservableOperator, + fromStreamPair, + toStream, +} from "../helpers/utils.ts"; +import { Observable } from "../observable.ts"; + +const range1000 = Array.from({ length: 1000 }, (_, index) => index); + +async function collectReadable(stream: ReadableStream): Promise { + const values: T[] = []; + const reader = stream.getReader(); + + try { + while (true) { + const { done, value } = await reader.read(); + if (done) { + break; + } + + values.push(value); + } + } finally { + reader.releaseLock(); + } + + return values; +} + +async function collectObservable(observable: Observable): Promise { + const values: T[] = []; + + for await (const value of observable) { + if (!isObservableError(value)) { + values.push(value); + } + } + + return values; +} + +function collectObservableBySubscription( + observable: Observable, +): Promise { + const values: T[] = []; + + return new Promise((resolve, reject) => { + observable.subscribe({ + next(value) { + if (!isObservableError(value)) { + values.push(value); + } + }, + error(error) { + reject(error); + }, + complete() { + resolve(values); + }, + }); + }); +} + +bench("Interop: direct TransformStream pair (1000 items)", async () => { + const transform = new TransformStream({ + transform(chunk, controller) { + controller.enqueue(String(chunk)); + }, + }); + + const values = await collectReadable( + toStream(range1000).pipeThrough(transform), + ); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: fromStreamPair wrapper (1000 items)", async () => { + const stringify = fromStreamPair(() => { + const transform = new TransformStream({ + transform(chunk, controller) { + controller.enqueue(String(chunk)); + }, + }); + + return { + readable: transform.readable, + writable: transform.writable, + }; + }); + + const values = await collectReadable(applyOperator(toStream(range1000), stringify)); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: local map/take pipeline via async iteration (1000 items)", async () => { + const result = pipe( + Observable.from(range1000), + map((value: number) => value + 1), + take(100), + ); + + const values = await collectObservable(result); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: local map/take pipeline via subscribe (1000 items)", async () => { + const result = pipe( + Observable.from(range1000), + map((value: number) => value + 1), + take(100), + ); + + const values = await collectObservableBySubscription(result); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: RxJS via fromObservableOperator via async iteration (1000 items)", async () => { + const rxOperator = fromObservableOperator< + number, + number, + RxObservable + >( + rxPipe( + rxMap((value: number) => value + 1), + rxTake(100), + ), + { sourceAdapter: (source) => rxFrom(source) as RxObservable }, + ); + + const result = pipe( + Observable.from(range1000), + ignoreErrors(), + rxOperator, + ); + + const values = await collectObservable(result); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: RxJS via fromObservableOperator via subscribe (1000 items)", async () => { + const rxOperator = fromObservableOperator< + number, + number, + RxObservable + >( + rxPipe( + rxMap((value: number) => value + 1), + rxTake(100), + ), + { sourceAdapter: (source) => rxFrom(source) as RxObservable }, + ); + + const result = pipe( + Observable.from(range1000), + ignoreErrors(), + rxOperator, + ); + + const values = await collectObservableBySubscription(result); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: direct RxJS pipeline (1000 items)", async () => { + const values: number[] = []; + + await new Promise((resolve) => { + rxFrom(range1000) + .pipe( + rxMap((value: number) => value + 1), + rxTake(100), + ) + .subscribe({ + next(value) { + values.push(value); + }, + complete() { + resolve(); + }, + }); + }); + + do_not_optimize(values); +}).gc("inner"); + +await run(); \ No newline at end of file diff --git a/bench/operators_bench.ts b/bench/operators_bench.ts index 78335b9..f7e2f4c 100644 --- a/bench/operators_bench.ts +++ b/bench/operators_bench.ts @@ -9,6 +9,14 @@ import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34"; import { Observable } from "../observable.ts"; import { isObservableError } from "../error.ts"; +import { + zipWith, +} from "../helpers/operations/combination.ts"; +import { + elementAt, + findIndex, + first, +} from "../helpers/operations/conditional.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"; @@ -182,4 +190,74 @@ bench("Operators: deep chain (10 operators)", async () => { do_not_optimize(values); }).gc("inner"); +bench("Operators: findIndex early match (1000 items)", async () => { + const result = pipe( + Observable.from(numberStream(1000)), + findIndex((value: number) => value === 10), + ); + + const values: number[] = []; + for await (const val of result) { + if (!isObservableError(val)) values.push(val); + } + + do_not_optimize(values); +}).gc("inner"); + +bench("Operators: elementAt middle index (1000 items)", async () => { + const result = pipe( + Observable.from(numberStream(1000)), + elementAt(500), + ); + + const values: number[] = []; + for await (const val of result) { + if (!isObservableError(val)) values.push(val); + } + + do_not_optimize(values); +}).gc("inner"); + +bench("Operators: first predicate early termination (1000 items)", async () => { + const result = pipe( + Observable.from(numberStream(1000)), + first((value: number) => value >= 10), + ); + + const values: number[] = []; + for await (const val of result) { + if (!isObservableError(val)) values.push(val); + } + + do_not_optimize(values); +}).gc("inner"); + +bench("Operators: zipWith balanced companions (1000 items)", async () => { + const result = pipe( + Observable.from(numberStream(1000)), + zipWith(Observable.from(numberStream(1000))), + ); + + const values: Array<[number, number]> = []; + for await (const val of result) { + if (!isObservableError(val)) values.push(val); + } + + do_not_optimize(values); +}).gc("inner"); + +bench("Operators: zipWith skewed companion backlog (1000 outputs)", async () => { + const result = pipe( + Observable.from(numberStream(1000)), + zipWith(Observable.from(numberStream(5000))), + ); + + const values: Array<[number, number]> = []; + for await (const val of result) { + if (!isObservableError(val)) values.push(val); + } + + do_not_optimize(values); +}).gc("inner"); + await run(); diff --git a/bench/run.ts b/bench/run.ts index ffde39d..cdd6463 100644 --- a/bench/run.ts +++ b/bench/run.ts @@ -24,6 +24,11 @@ console.log("Running Operator Pipeline benchmarks..."); console.log("-".repeat(80)); await import("./operators_bench.ts"); +console.log(""); +console.log("Running Interop benchmarks..."); +console.log("-".repeat(80)); +await import("./interop_bench.ts"); + console.log(""); console.log("Running Memory benchmarks..."); console.log("-".repeat(80));