diff --git a/bench/interop_bench.ts b/bench/interop_bench.ts index 0f5ae5c..0930350 100644 --- a/bench/interop_bench.ts +++ b/bench/interop_bench.ts @@ -3,8 +3,14 @@ * 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 + * closest practical baselines so interop convenience does not hide avoidable * overhead. + * + * The RxJS-adapter cases include a leading `ignoreErrors()` stage because the + * foreign operator currently plugs into the pipeline after local error values + * have been filtered out. Matching local baselines include the same + * compatibility stage so the comparison separates adapter cost from that extra + * source-filtering step. */ import { bench, do_not_optimize, run } from "npm:mitata@^1.0.34"; @@ -141,6 +147,32 @@ bench("Interop: local map/take pipeline via subscribe (1000 items)", async () => do_not_optimize(values); }).gc("inner"); +bench("Interop: local map/take with ignoreErrors via async iteration (1000 items)", async () => { + const result = pipe( + Observable.from(range1000), + ignoreErrors(), + map((value: number) => value + 1), + take(100), + ); + + const values = await collectObservable(result); + + do_not_optimize(values); +}).gc("inner"); + +bench("Interop: local map/take with ignoreErrors via subscribe (1000 items)", async () => { + const result = pipe( + Observable.from(range1000), + ignoreErrors(), + 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, diff --git a/tests/helpers/utils_bdd_test.ts b/tests/helpers/utils_bdd_test.ts index 0cc1bed..53b5716 100644 --- a/tests/helpers/utils_bdd_test.ts +++ b/tests/helpers/utils_bdd_test.ts @@ -574,6 +574,24 @@ describe("Interop Utilities", () => { } }); + it("should wrap synchronous foreign subscribe failures as ObservableError values", async () => { + const foreignOperator = fromObservableOperator(() => ({ + subscribe() { + throw new Error("subscribe failed immediately"); + }, + })); + + const values = await collectStream(applyOperator(toStream([1, 2, 3]), foreignOperator)); + + expect(values).toHaveLength(1); + expect(isObservableError(values[0])).toBe(true); + + if (isObservableError(values[0])) { + expect(values[0].message).toContain("subscribe failed immediately"); + expect(values[0].operator).toBe("operator:fromObservableOperator:output"); + } + }); + it("should preserve wrapped errors from the foreign output", async () => { const failingOperator = fromObservableOperator((_source) => ({ async *[Symbol.asyncIterator]() { diff --git a/tests/observable/compat_comparison_test.ts b/tests/observable/compat_comparison_test.ts new file mode 100644 index 0000000..d08af12 --- /dev/null +++ b/tests/observable/compat_comparison_test.ts @@ -0,0 +1,362 @@ +// deno-lint-ignore-file no-import-prefix +import { expect, test } from 'jsr:@libs/testing@^5'; + +import type { + ObservableProtocol, + SpecObservable, + SpecObserver, + SpecSubscription, +} from '../../_spec.ts'; + +import { + forEach, + Observable, + SubscriptionStateMap, +} from '../../observable.ts'; +import { Symbol } from '../../symbol.ts'; + +test('subscribe preserves this for anonymous next observers', () => { + const events: Array = []; + + Observable.of(1).subscribe({ + label: 'next-context', + next(value) { + events.push(this.label); + events.push(value); + }, + } as { + label: string; + next(value: number): void; + }); + + expect(events).toEqual(['next-context', 1]); +}); + +test('subscribe preserves this for anonymous error and complete observers', () => { + const events: string[] = []; + + new Observable((observer) => { + observer.error(new Error('boom')); + }).subscribe({ + label: 'error-context', + error(error) { + events.push(`${this.label}:${(error as Error).message}`); + }, + } as { + label: string; + error(error: unknown): void; + }); + + Observable.of().subscribe({ + label: 'complete-context', + complete() { + events.push(this.label); + }, + } as { + label: string; + complete(): void; + }); + + expect(events).toEqual(['error-context:boom', 'complete-context']); +}); + +test('subscribe tolerates being called with no arguments', () => { + const source = new Observable((observer) => { + observer.next('foo'); + observer.complete(); + }); + + expect(() => { + (source as unknown as { subscribe(): unknown }).subscribe(); + }).not.toThrow(); +}); + +test('subscribe closes immediately for an already-aborted AbortSignal without running the subscriber', () => { + const controller = new AbortController(); + controller.abort(); + + let subscriberCalled = false; + let start_closed: boolean | undefined; + + const subscription = new Observable((observer) => { + subscriberCalled = true; + observer.next(1); + return () => {}; + }).subscribe({ + start(sub) { + start_closed = sub.closed; + }, + next() {}, + }, { signal: controller.signal }); + + expect(start_closed).toBe(true); + expect(subscriberCalled).toBe(false); + expect(subscription.closed).toBe(true); + expect(SubscriptionStateMap.get(subscription)).toBe(undefined); +}); + +test('unsubscribing in start clears internal subscription state before the subscriber runs', () => { + let subscriberCalled = false; + let state_after_unsubscribe: + | ReturnType + | undefined; + + const subscription = new Observable((observer) => { + subscriberCalled = true; + observer.next(1); + return () => {}; + }).subscribe({ + start(sub) { + sub.unsubscribe(); + state_after_unsubscribe = SubscriptionStateMap.get(sub); + }, + next() {}, + }); + + expect(subscriberCalled).toBe(false); + expect(subscription.closed).toBe(true); + expect(state_after_unsubscribe).toBe(undefined); +}); + +test('subscribe tears down when AbortSignal aborts after subscription starts', async () => { + const controller = new AbortController(); + const values: number[] = []; + let cleanupCount = 0; + + const subscription = new Observable((observer) => { + let value = 0; + const id = setInterval(() => observer.next(value++), 5); + + return () => { + cleanupCount++; + clearInterval(id); + }; + }).subscribe({ + next(value) { + values.push(value); + if (value === 1) { + controller.abort(); + } + }, + }, { signal: controller.signal }); + + await new Promise((resolve) => setTimeout(resolve, 30)); + + expect(values).toEqual([0, 1]); + expect(cleanupCount).toBe(1); + expect(subscription.closed).toBe(true); +}); + +test('Observable.from stops a synchronous iterable when unsubscribed mid-drain', () => { + const seen: number[] = []; + const sideEffects: number[] = []; + let iteratorClosed = false; + let subscription: + | { + unsubscribe(): void; + } + | undefined; + + const iterable = { + [Symbol.iterator]() { + let index = 0; + return { + next() { + sideEffects.push(index); + return index < 10 + ? { value: index++, done: false } + : { value: undefined, done: true }; + }, + return() { + iteratorClosed = true; + return { value: undefined, done: true }; + }, + }; + }, + }; + + Observable.from(iterable).subscribe({ + start(sub) { + subscription = sub; + }, + next(value) { + seen.push(value!); + if (value === 2) { + subscription?.unsubscribe(); + } + }, + }); + + expect(seen).toEqual([0, 1, 2]); + expect(sideEffects).toEqual([0, 1, 2]); + expect(iteratorClosed).toBe(true); +}); + +test('Symbol.asyncIterator.throw unsubscribes the source', async () => { + let state = 'idle'; + + const source = new Observable((observer) => { + state = 'subscribed'; + observer.next(0); + return () => { + state = 'unsubscribed'; + }; + }); + + const iterator = source[Symbol.asyncIterator](); + + expect(state).toBe('idle'); + await iterator.next(); + expect(state).toBe('subscribed'); + + await expect(iterator.throw?.(new Error('wee!'))).rejects.toThrow('wee!'); + expect(state).toBe('unsubscribed'); +}); + +test('forEach rejects if the callback is not a function', async () => { + await expect( + forEach( + Observable.of(1, 2, 3), + undefined as unknown as (value: number, index: number) => void, + ), + ).rejects.toThrow(TypeError); +}); + +test('forEach resolves with undefined when the source completes', async () => { + const values: number[] = []; + + await expect( + forEach(Observable.of(1, 2, 3), (value, index) => { + values.push(value + index); + }), + ).resolves.toBe(undefined); + + expect(values).toEqual([1, 3, 5]); +}); + +test('forEach rejects when the source errors', async () => { + const error = new Error('bad'); + + await expect( + forEach(new Observable((observer) => { + observer.error(error); + }), () => {}), + ).rejects.toBe(error); +}); + +test('forEach rejects if the callback throws and stops a synchronous source', async () => { + const expected = new Error('NO THREES'); + const values: number[] = []; + + await expect( + forEach(Observable.of(1, 2, 3, 4), (value) => { + values.push(value); + if (value === 3) { + throw expected; + } + }), + ).rejects.toBe(expected); + + expect(values).toEqual([1, 2, 3]); +}); + +test('forEach unsubscribes a synchronous foreign source when the callback throws before subscribe returns', async () => { + const expected = new Error('stop now'); + let unsubscribe_count = 0; + const protocol: ObservableProtocol = { + subscribe( + observer_or_next: SpecObserver | ((value: number) => void), + ): SpecSubscription { + const observer = typeof observer_or_next === 'function' + ? { next: observer_or_next } + : observer_or_next; + + observer.next?.(1); + + return { + unsubscribe() { + unsubscribe_count++; + }, + }; + }, + }; + + const foreign: SpecObservable = { + [Symbol.observable]() { + return protocol; + }, + }; + + await expect( + forEach(foreign, () => { + throw expected; + }), + ).rejects.toBe(expected); + + expect(unsubscribe_count).toBe(1); +}); + +test('forEach rejects if the callback throws and tears down an async source', async () => { + const expected = new Error('NO TWOS'); + const values: number[] = []; + let cleanupCount = 0; + + await expect( + forEach(new Observable((observer) => { + let value = 1; + const id = setInterval(() => observer.next(value++), 1); + + return () => { + cleanupCount++; + clearInterval(id); + }; + }), (value) => { + values.push(value); + if (value === 2) { + throw expected; + } + }), + ).rejects.toBe(expected); + + expect(values).toEqual([1, 2]); + expect(cleanupCount).toBe(1); +}); + +test('Observable.prototype.forEach delegates to the exported helper', async () => { + const values: string[] = []; + + await expect( + Observable.of('a', 'b').forEach((value, index) => { + values.push(`${index}:${value}`); + }), + ).resolves.toBe(undefined); + + expect(values).toEqual(['0:a', '1:b']); +}); + +test('forEach rejects when its AbortSignal aborts', async () => { + const controller = new AbortController(); + const reason = new Error('aborted'); + const values: number[] = []; + let cleanupCount = 0; + + await expect( + forEach(new Observable((observer) => { + let value = 0; + const id = setInterval(() => observer.next(value++), 5); + + return () => { + cleanupCount++; + clearInterval(id); + }; + }), (value) => { + values.push(value); + if (value === 1) { + controller.abort(reason); + } + }, { signal: controller.signal }), + ).rejects.toBe(reason); + + expect(values).toEqual([0, 1]); + expect(cleanupCount).toBe(1); +});