From 161b56c756bc560c118eede7ca5157f78720d475 Mon Sep 17 00:00:00 2001 From: Harvmaster Date: Mon, 3 Aug 2026 03:35:09 +0000 Subject: [PATCH] Formatting --- source/services/broadcaster.ts | 2 + test/services/broadcaster.test.ts | 335 ++++++++++++++---------------- test/subscription-flow.test.ts | 127 ++++++----- 3 files changed, 219 insertions(+), 245 deletions(-) diff --git a/source/services/broadcaster.ts b/source/services/broadcaster.ts index a59e8a4..463c26c 100644 --- a/source/services/broadcaster.ts +++ b/source/services/broadcaster.ts @@ -6,6 +6,7 @@ import { ApplicationError } from '../errors/index.ts'; /** Request-scoped view from which the broadcaster obtains a stable connection. */ export interface BroadcastStream { + /** Connection identity shared by every request on the same transport session. */ readonly connection: BaseStream; @@ -15,6 +16,7 @@ export interface BroadcastStream { /** One pending subscribe call and the topics whose removal will resolve it. */ interface SubscriptionWaiter { + /** Only topics newly introduced by this particular subscribe call. */ readonly remainingTopics: Set; diff --git a/test/services/broadcaster.test.ts b/test/services/broadcaster.test.ts index 0fc6f3e..3acc018 100644 --- a/test/services/broadcaster.test.ts +++ b/test/services/broadcaster.test.ts @@ -1,199 +1,180 @@ -import { describe, expect, it, vi } from "vitest"; +import { describe, expect, it, vi } from 'vitest'; -import { ApplicationError } from "../../source/errors/index.js"; -import { Broadcaster } from "../../source/services/broadcaster.js"; -import { ApplicationRouteStream } from "../../source/services/route-stream.js"; -import { Logger } from "../../source/utils/logger.js"; -import { TestConnection } from "../helpers/test-connection.js"; +import { Broadcaster } from '../../source/services/broadcaster.ts'; +import { ApplicationRouteStream } from '../../source/services/route-stream.ts'; +import { Logger } from '../../source/utils/logger.ts'; +import { TestConnection } from '../helpers/test-connection.ts'; -function createBroadcaster(): Broadcaster { - return new Broadcaster(new Logger("broadcaster-test")); -} +const createBroadcaster = (): Broadcaster => { + return new Broadcaster(new Logger('broadcaster-test')); +}; -function routeStream(connection: TestConnection): ApplicationRouteStream { - return new ApplicationRouteStream(connection, undefined); -} +const routeStream = (connection: TestConnection): ApplicationRouteStream => { + return new ApplicationRouteStream(connection, undefined); +}; -async function expectPending(promise: Promise): Promise { - const settled = vi.fn(); - void promise.then(settled); - await Promise.resolve(); - expect(settled).not.toHaveBeenCalled(); -} +const expectPending = async (promise: Promise): Promise => { + const settled = vi.fn(); + void promise.then(settled); + await Promise.resolve(); + expect(settled).not.toHaveBeenCalled(); +}; -describe("Broadcaster subscriptions", () => { - it("delivers events and resolves after a later request removes the topic", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, true); - const subscribed = broadcaster.subscribe(routeStream(connection), [ - "items", - ]); +describe('Broadcaster subscriptions', () => { + it('delivers events and resolves after a later request removes the topic', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, true); + const subscribed = broadcaster.subscribe(routeStream(connection), [ 'items' ]); - await expectPending(subscribed); - await broadcaster.publish("items", { - type: "item-changed", - data: { id: "a" }, + await expectPending(subscribed); + await broadcaster.publish('items', { + type: 'item-changed', + data: { id: 'a' }, + }); + + expect(connection.messages).toEqual([ + expect.objectContaining({ + type: 'item-changed', + data: { id: 'a' }, + }), + ]); + + // A different request-scoped facade still resolves the connection's + // original subscription. + await broadcaster.unsubscribe(routeStream(connection), [ 'items' ]); + await expect(subscribed).resolves.toBeUndefined(); }); - expect(connection.messages).toEqual([ - expect.objectContaining({ - type: "item-changed", - data: { id: "a" }, - }), - ]); + it('resolves fully duplicate subscriptions immediately', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, true); + const first = broadcaster.subscribe(routeStream(connection), [ 'items', 'items' ]); + const duplicate = broadcaster.subscribe(routeStream(connection), [ 'items' ]); - // A different request-scoped facade still resolves the connection's - // original subscription. - await broadcaster.unsubscribe(routeStream(connection), ["items"]); - await expect(subscribed).resolves.toBeUndefined(); - }); + await expect(duplicate).resolves.toBeUndefined(); + await expectPending(first); + expect(connection.closeCallbacks).toHaveLength(1); - it("resolves fully duplicate subscriptions immediately", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, true); - const first = broadcaster.subscribe(routeStream(connection), [ - "items", - "items", - ]); - const duplicate = broadcaster.subscribe(routeStream(connection), ["items"]); - - await expect(duplicate).resolves.toBeUndefined(); - await expectPending(first); - expect(connection.closeCallbacks).toHaveLength(1); - - await broadcaster.unsubscribe(routeStream(connection), ["items"]); - await first; - }); - - it("waits only for topics newly added by a partially overlapping call", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, true); - const first = broadcaster.subscribe(routeStream(connection), ["a"]); - const second = broadcaster.subscribe(routeStream(connection), ["a", "b"]); - - await broadcaster.unsubscribe(routeStream(connection), ["a"]); - await expect(first).resolves.toBeUndefined(); - await expectPending(second); - - await broadcaster.unsubscribe(routeStream(connection), ["b"]); - await expect(second).resolves.toBeUndefined(); - }); - - it("resolves every pending subscription and removes topics on close", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, false); - const first = broadcaster.subscribe(routeStream(connection), ["a"]); - const second = broadcaster.subscribe(routeStream(connection), ["b"]); - - connection.close(); - await Promise.all([first, second]); - await broadcaster.publish("a", { type: "changed", data: null }); - await broadcaster.publish("b", { type: "changed", data: null }); - - expect(connection.messages).toEqual([]); - }); - - it("immediately resolves registration against an already-closed connection", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, false); - connection.close(); - - await expect( - broadcaster.subscribe(routeStream(connection), ["items"]), - ).resolves.toBeUndefined(); - }); - - it("rejects subscriptions on a non-streaming connection", () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(false, false); - - expect(() => - broadcaster.subscribe(routeStream(connection), ["items"]), - ).toThrowError( - expect.objectContaining({ statusCode: 406 }), - ); - }); - - it("treats an empty subscription and repeated unsubscribe as no-ops", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, true); - - await expect( - broadcaster.subscribe(routeStream(connection), []), - ).resolves.toBeUndefined(); - await broadcaster.unsubscribe(routeStream(connection), ["missing"]); - await broadcaster.unsubscribe(routeStream(connection)); - - expect(connection.closeCallbacks).toHaveLength(0); - }); - - it("fans out concurrently to independent connections", async () => { - const broadcaster = createBroadcaster(); - const first = new TestConnection(true, false); - const second = new TestConnection(true, false); - const originalFirstSend = first.send.bind(first); - let releaseFirst: () => void = () => undefined; - const firstReleased = new Promise((resolve) => { - releaseFirst = resolve; - }); - let markSecondSent: () => void = () => undefined; - const secondSent = new Promise((resolve) => { - markSecondSent = resolve; + await broadcaster.unsubscribe(routeStream(connection), [ 'items' ]); + await first; }); - first.send = async (message) => { - await firstReleased; - await originalFirstSend(message); - }; - second.send = async (message) => { - await TestConnection.prototype.send.call(second, message); - markSecondSent(); - }; + it('waits only for topics newly added by a partially overlapping call', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, true); + const first = broadcaster.subscribe(routeStream(connection), [ 'a' ]); + const second = broadcaster.subscribe(routeStream(connection), [ 'a', 'b' ]); - const firstSubscription = broadcaster.subscribe(routeStream(first), [ - "items", - ]); - const secondSubscription = broadcaster.subscribe(routeStream(second), [ - "items", - ]); - const publication = broadcaster.publish("items", { - type: "item-changed", - data: {}, + await broadcaster.unsubscribe(routeStream(connection), [ 'a' ]); + await expect(first).resolves.toBeUndefined(); + await expectPending(second); + + await broadcaster.unsubscribe(routeStream(connection), [ 'b' ]); + await expect(second).resolves.toBeUndefined(); }); - await secondSent; - releaseFirst(); - await publication; + it('resolves every pending subscription and removes topics on close', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, false); + const first = broadcaster.subscribe(routeStream(connection), [ 'a' ]); + const second = broadcaster.subscribe(routeStream(connection), [ 'b' ]); - expect(first.messages).toHaveLength(1); - expect(second.messages).toHaveLength(1); + connection.close(); + await Promise.all([ first, second ]); + await broadcaster.publish('a', { type: 'changed', data: null }); + await broadcaster.publish('b', { type: 'changed', data: null }); - first.close(); - second.close(); - await Promise.all([firstSubscription, secondSubscription]); - }); - - it("closes and removes a connection whose event delivery fails", async () => { - const broadcaster = createBroadcaster(); - const connection = new TestConnection(true, true); - connection.send = vi.fn().mockRejectedValue(new Error("socket failed")); - const subscribed = broadcaster.subscribe(routeStream(connection), [ - "items", - ]); - - await broadcaster.publish("items", { - type: "item-changed", - data: {}, + expect(connection.messages).toEqual([]); }); - await subscribed; - expect(connection.closed).toBe(true); - expect(connection.send).toHaveBeenCalledOnce(); + it('immediately resolves registration against an already-closed connection', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, false); + connection.close(); - await broadcaster.publish("items", { - type: "item-changed", - data: {}, + await expect(broadcaster.subscribe(routeStream(connection), [ 'items' ])).resolves.toBeUndefined(); + }); + + it('rejects subscriptions on a non-streaming connection', (): void => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(false, false); + + expect(() => broadcaster.subscribe(routeStream(connection), [ 'items' ])).toThrowError(expect.objectContaining({ statusCode: 406 })); + }); + + it('treats an empty subscription and repeated unsubscribe as no-ops', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, true); + + await expect(broadcaster.subscribe(routeStream(connection), [])).resolves.toBeUndefined(); + await broadcaster.unsubscribe(routeStream(connection), [ 'missing' ]); + await broadcaster.unsubscribe(routeStream(connection)); + + expect(connection.closeCallbacks).toHaveLength(0); + }); + + it('fans out concurrently to independent connections', async (): Promise => { + const broadcaster = createBroadcaster(); + const first = new TestConnection(true, false); + const second = new TestConnection(true, false); + const originalFirstSend = first.send.bind(first); + let releaseFirst: () => void = () => undefined; + const firstReleased = new Promise((resolve) => { + releaseFirst = resolve; + }); + let markSecondSent: () => void = () => undefined; + const secondSent = new Promise((resolve) => { + markSecondSent = resolve; + }); + + first.send = async (message): Promise => { + await firstReleased; + await originalFirstSend(message); + }; + + second.send = async (message): Promise => { + await TestConnection.prototype.send.call(second, message); + markSecondSent(); + }; + + const firstSubscription = broadcaster.subscribe(routeStream(first), [ 'items' ]); + const secondSubscription = broadcaster.subscribe(routeStream(second), [ 'items' ]); + const publication = broadcaster.publish('items', { + type: 'item-changed', + data: {}, + }); + + await secondSent; + releaseFirst(); + await publication; + + expect(first.messages).toHaveLength(1); + expect(second.messages).toHaveLength(1); + + first.close(); + second.close(); + await Promise.all([ firstSubscription, secondSubscription ]); + }); + + it('closes and removes a connection whose event delivery fails', async (): Promise => { + const broadcaster = createBroadcaster(); + const connection = new TestConnection(true, true); + connection.send = vi.fn().mockRejectedValue(new Error('socket failed')); + const subscribed = broadcaster.subscribe(routeStream(connection), [ 'items' ]); + + await broadcaster.publish('items', { + type: 'item-changed', + data: {}, + }); + await subscribed; + + expect(connection.closed).toBe(true); + expect(connection.send).toHaveBeenCalledOnce(); + + await broadcaster.publish('items', { + type: 'item-changed', + data: {}, + }); + expect(connection.send).toHaveBeenCalledOnce(); }); - expect(connection.send).toHaveBeenCalledOnce(); - }); }); diff --git a/test/subscription-flow.test.ts b/test/subscription-flow.test.ts index 79c8400..d585d11 100644 --- a/test/subscription-flow.test.ts +++ b/test/subscription-flow.test.ts @@ -1,76 +1,67 @@ -import { describe, expect, it, vi } from "vitest"; +import { describe, expect, it, vi } from 'vitest'; -import type { RouteDefinition, RouteModule } from "../source/routes/types.js"; -import { ApplicationRouter } from "../source/services/router.js"; -import { Broadcaster } from "../source/services/broadcaster.js"; -import { Logger } from "../source/utils/logger.js"; -import { TestConnection } from "./helpers/test-connection.js"; +import type { RouteDefinition, RouteModule } from '../source/routes/types.ts'; +import { ApplicationRouter } from '../source/services/router.ts'; +import { Broadcaster } from '../source/services/broadcaster.ts'; +import { Logger } from '../source/utils/logger.ts'; +import { TestConnection } from './helpers/test-connection.ts'; -function moduleWith(routes: RouteDefinition[]): RouteModule { - return { - async getRoutes() { - return routes; - }, - }; -} - -async function expectPending(promise: Promise): Promise { - const settled = vi.fn(); - void promise.then(settled); - await Promise.resolve(); - expect(settled).not.toHaveBeenCalled(); -} - -describe("long-lived subscription dispatch", () => { - it("allows duplicate subscribe and unsubscribe requests while the original request waits", async () => { - const broadcaster = new Broadcaster(new Logger("subscription-flow-test")); - const router = await ApplicationRouter.create([ - moduleWith([ - { - url: "/items/subscribe", - handler: async (stream) => { - await broadcaster.subscribe(stream, ["items"]); - }, +const moduleWith = (routes: RouteDefinition[]): RouteModule => { + return { + async getRoutes(): Promise { + return routes; }, - { - url: "/items/unsubscribe", - handler: async (stream) => { - await broadcaster.unsubscribe(stream, ["items"]); - await stream.send({}); - }, - }, - ]), - ]); - const connection = new TestConnection(true, true); + }; +}; - const original = router.dispatch( - { path: "/items/subscribe", requestId: "subscribe-1" }, - connection, - ); - await vi.waitFor(() => expect(connection.closeCallbacks).toHaveLength(1)); - await expectPending(original); +const expectPending = async (promise: Promise): Promise => { + const settled = vi.fn(); + void promise.then(settled); + await Promise.resolve(); + expect(settled).not.toHaveBeenCalled(); +}; - // This request uses a different ApplicationRouteStream over the same - // connection. Since the topic already exists, its dispatch completes. - await router.dispatch( - { path: "/items/subscribe", requestId: "subscribe-2" }, - connection, - ); - await expectPending(original); +describe('long-lived subscription dispatch', (): void => { + it('allows duplicate subscribe and unsubscribe requests while the original request waits', async (): Promise => { + const broadcaster = new Broadcaster(new Logger('subscription-flow-test')); + const router = await ApplicationRouter.create([ + moduleWith([ + { + url: '/items/subscribe', + handler: async (stream): Promise => { + await broadcaster.subscribe(stream, [ 'items' ]); + }, + }, + { + url: '/items/unsubscribe', + handler: async (stream): Promise => { + await broadcaster.unsubscribe(stream, [ 'items' ]); + await stream.send({}); + }, + }, + ]), + ]); + const connection = new TestConnection(true, true); - await router.dispatch( - { path: "/items/unsubscribe", requestId: "unsubscribe-1" }, - connection, - ); - await original; + const original = router.dispatch({ path: '/items/subscribe', requestId: 'subscribe-1' }, connection); + await vi.waitFor(() => expect(connection.closeCallbacks).toHaveLength(1)); + await expectPending(original); - expect(connection.messages).toEqual([ - { - id: "unsubscribe-1", - type: "response", - statusCode: 200, - body: {}, - }, - ]); - }); + // This request uses a different ApplicationRouteStream over the same + // connection. Since the topic already exists, its dispatch completes. + await router.dispatch({ path: '/items/subscribe', requestId: 'subscribe-2' }, connection); + await expectPending(original); + + await router.dispatch({ path: '/items/unsubscribe', requestId: 'unsubscribe-1' }, connection); + await original; + + expect(connection.messages).toEqual([ + { + id: 'unsubscribe-1', + type: 'response', + statusCode: 200, + body: {}, + }, + ]); + }); });