diff --git a/tests/events/events_test.ts b/tests/events/events_test.ts index 4cd3aa3..25a52b3 100644 --- a/tests/events/events_test.ts +++ b/tests/events/events_test.ts @@ -913,19 +913,34 @@ test("Documentation - waitForEvent example from docs", async () => { // Runtime-specific tests if (runtime === "deno") { test("Deno-specific - Network permissions test", async () => { - // This test only runs on Deno and requires network permissions - await using server = Deno.serve( - { port: 8080, onListen: () => null }, - () => new Response(null, { status: 200 }), + // Use an ephemeral port and force the connection closed so shutdown does + // not wait on a kept-alive client socket after the assertion passes. + const server = Deno.serve( + { + hostname: "127.0.0.1", + port: 0, + onListen: () => {}, + }, + () => new Response(null, { + status: 200, + headers: { + connection: "close", + }, + }), ); - const response = await fetch( - `http://${server.addr.hostname}:${server.addr.port}`, - ); - expect(response.status).toBe(200); - - // Dispose of the response body if it exists - await response?.body?.cancel(); + try { + const response = await fetch(`http://127.0.0.1:${server.addr.port}`, { + headers: { + connection: "close", + }, + }); + + expect(response.status).toBe(200); + await response.body?.cancel(); + } finally { + await server.shutdown(); + } }, { permissions: { net: "inherit" } }); } diff --git a/tests/helpers/utils_bdd_test.ts b/tests/helpers/utils_bdd_test.ts index 08ee245..0cc1bed 100644 --- a/tests/helpers/utils_bdd_test.ts +++ b/tests/helpers/utils_bdd_test.ts @@ -13,6 +13,13 @@ import { describe, it } from "jsr:@std/testing@^1/bdd"; import { expect } from "jsr:@std/expect@^1"; +import { + type Observable as RxObservable, + from as rxFrom, + map as rxMap, + pipe as rxPipe, + take as rxTake, +} from "npm:rxjs@7.8.2"; import { applyOperator, @@ -516,6 +523,57 @@ describe("Interop Utilities", () => { expect(await collectStream(result)).toEqual([2, 4]); }); + it("should adapt standard RxJS operator functions", async () => { + const rxOperator = fromObservableOperator< + number, + number, + RxObservable + >( + rxPipe( + rxMap((value: number) => value + 1), + rxTake(2), + ), + { sourceAdapter: (source) => rxFrom(source) as RxObservable }, + ); + + const result = applyOperator(toStream([1, 2, 3]), rxOperator); + + expect(await collectStream(result)).toEqual([2, 3]); + }); + + it("should unsubscribe foreign subscribable output when the stream is cancelled early", async () => { + let unsubscribeCount = 0; + + const foreignOperator = fromObservableOperator(() => ({ + subscribe(observer) { + observer.next?.(1); + + return { + unsubscribe() { + unsubscribeCount++; + }, + }; + }, + })); + + const result = applyOperator(toStream([1, 2, 3]), foreignOperator); + const reader = result.getReader(); + + try { + const first = await reader.read(); + + expect(first.done).toBe(false); + expect(first.value).toBe(1); + + await reader.cancel(); + await Promise.resolve(); + + expect(unsubscribeCount).toBe(1); + } finally { + reader.releaseLock(); + } + }); + it("should preserve wrapped errors from the foreign output", async () => { const failingOperator = fromObservableOperator((_source) => ({ async *[Symbol.asyncIterator]() { diff --git a/tests/observable/multiple_subscribers_test.ts b/tests/observable/multiple_subscribers_test.ts index 8e2e72a..9308c06 100644 --- a/tests/observable/multiple_subscribers_test.ts +++ b/tests/observable/multiple_subscribers_test.ts @@ -661,6 +661,16 @@ test("Observable.from throws for incompatible inputs", () => { expect(() => Observable.from(true as unknown as [])).toThrow(TypeError); }); +test("Observable.from rejects direct subscribables without Symbol.observable", () => { + expect(() => + Observable.from({ + subscribe() { + return { unsubscribe() {} }; + }, + } as unknown as never) + ).toThrow(TypeError); +}); + test("Observable.from uses the this value if it's a function", () => { let usedThisValue = false; const thisObj = function () {