From d4b3b72e75fb7619dc0b229948c8155b6cfc863f Mon Sep 17 00:00:00 2001 From: Harvmaster Date: Mon, 3 Aug 2026 03:29:30 +0000 Subject: [PATCH 1/3] Formatting --- source/{app.ts => index.ts} | 4 +++- source/services/storage/database.ts | 17 +++++++++++------ .../storage/migrations/001-resources.ts | 5 +++-- source/services/storage/tables.ts | 1 + tsconfig.json | 2 +- 5 files changed, 19 insertions(+), 10 deletions(-) rename source/{app.ts => index.ts} (95%) diff --git a/source/app.ts b/source/index.ts similarity index 95% rename from source/app.ts rename to source/index.ts index 08df05e..75377a7 100644 --- a/source/app.ts +++ b/source/index.ts @@ -24,7 +24,9 @@ export class App { constructor(private readonly database: Database) {} - async start(): Promise {} + async start(): Promise { + await this.database.start(); + } /** Stop transports before releasing the database they may still use. */ async stop(): Promise { diff --git a/source/services/storage/database.ts b/source/services/storage/database.ts index c3ef093..1cf194e 100644 --- a/source/services/storage/database.ts +++ b/source/services/storage/database.ts @@ -6,6 +6,7 @@ import type { Logger } from '../../utils/logger.ts'; /** Options required to open a SQLite database connection. */ export type DatabaseOptions = { + /** Filesystem path to the SQLite database file. */ path: string; @@ -39,9 +40,6 @@ export class Database { this.kysely = new Kysely({ dialect: this.dialect, }); - - // Configure the SQLite pragmas. - this.configurePragmas(); } /** @@ -53,6 +51,13 @@ export class Database { return this.kysely; } + async start(): Promise { + this.debug('starting database connection'); + + // Configure the SQLite pragmas. + await this.configurePragmas(); + } + /** * Destroys the database connection. */ @@ -66,10 +71,10 @@ export class Database { * * WAL improves write concurrency; foreign keys enforce referential integrity. */ - private configurePragmas(): void { + private async configurePragmas(): Promise { this.debug('configuring SQLite pragmas'); - this.kysely.executeQuery(CompiledQuery.raw('PRAGMA journal_mode = WAL')); - this.kysely.executeQuery(CompiledQuery.raw('PRAGMA foreign_keys = ON')); + await this.kysely.executeQuery(CompiledQuery.raw('PRAGMA journal_mode = WAL')); + await this.kysely.executeQuery(CompiledQuery.raw('PRAGMA foreign_keys = ON')); } } diff --git a/source/services/storage/migrations/001-resources.ts b/source/services/storage/migrations/001-resources.ts index 0e548c9..27dcf73 100644 --- a/source/services/storage/migrations/001-resources.ts +++ b/source/services/storage/migrations/001-resources.ts @@ -23,7 +23,7 @@ export const up = async (db: Kysely): Promise => { .addColumn('public_key', 'text', (col) => col.notNull()) .addColumn('blob', 'blob', (col) => col.notNull()) .addColumn('timestamp', 'integer', (col) => col.notNull().defaultTo(millisecondTime)) - .addPrimaryKeyConstraint('pk_resource_data', ['resource_id', 'public_key']) + .addPrimaryKeyConstraint('pk_resource_data', [ 'resource_id', 'public_key' ]) .execute(); }; @@ -33,5 +33,6 @@ export const up = async (db: Kysely): Promise => { * @param db - Kysely database to apply the rollback against. */ export const down = async (db: Kysely): Promise => { - await db.schema.dropTable('resource_data').ifExists().execute(); + await db.schema.dropTable('resource_data').ifExists() +.execute(); }; diff --git a/source/services/storage/tables.ts b/source/services/storage/tables.ts index cedf815..d937c28 100644 --- a/source/services/storage/tables.ts +++ b/source/services/storage/tables.ts @@ -10,6 +10,7 @@ export type BlobColumn = ColumnType; * One row per (resource_id, public_key). Each publicKey owns a slot within a shared resource. */ export interface ResourceDataTable { + /** Shared resource identifier grouping related instances. */ resource_id: string; diff --git a/tsconfig.json b/tsconfig.json index bae3d55..6e15aa7 100644 --- a/tsconfig.json +++ b/tsconfig.json @@ -15,5 +15,5 @@ "declarationMap": true, "types": ["node"] }, - "exclude": ["node_modules/**/*", "dist/**/*", "test"] + "exclude": ["node_modules/**/*", "dist/**/*"] } From a36e267280b58ce82e5bd62cf3e3e67a99ab11e5 Mon Sep 17 00:00:00 2001 From: Harvmaster Date: Mon, 3 Aug 2026 03:32:11 +0000 Subject: [PATCH 2/3] Formatting --- source/constants.ts | 3 - source/routes/types.ts | 3 + source/services/router.ts | 1 + source/services/stream/base-stream.ts | 2 + test/helpers/test-connection.ts | 62 ++++---- test/services/router.test.ts | 211 ++++++++++++-------------- 6 files changed, 130 insertions(+), 152 deletions(-) diff --git a/source/constants.ts b/source/constants.ts index 140fe80..ddfd93c 100644 --- a/source/constants.ts +++ b/source/constants.ts @@ -12,6 +12,3 @@ export const HTTP_STATUS_CODE_CREATED = 201; * No content response status code. */ export const HTTP_STATUS_CODE_NO_CONTENT = 204; - - - diff --git a/source/routes/types.ts b/source/routes/types.ts index 7e39906..2d2746e 100644 --- a/source/routes/types.ts +++ b/source/routes/types.ts @@ -1,6 +1,7 @@ import type { BaseStream } from '../services/stream/base-stream.js'; export type RouteSendOptions = { + /** Defaults to `response`; any other value sends an application event. */ type?: string; @@ -15,6 +16,7 @@ export type RouteSendOptions = { * connection lifetime are shared with other requests on the same connection. */ export interface RouteStream { + /** Connection shared by every request on the same transport session. */ readonly connection: BaseStream; @@ -29,6 +31,7 @@ export type RouteHandler = (stream: RouteStream) => void | Promise; /** An exact application route with no transport-specific metadata. */ export type RouteDefinition = { + /** Exact route name. Parameter and wildcard syntax are not supported. */ url: string; handler: RouteHandler; diff --git a/source/services/router.ts b/source/services/router.ts index f4ce543..71c4d50 100644 --- a/source/services/router.ts +++ b/source/services/router.ts @@ -5,6 +5,7 @@ import type { BaseStream } from './stream/base-stream.ts'; /** Canonical request produced by every transport adapter. */ export type ApplicationRequest = { + /** Exact application route name. */ path: string; diff --git a/source/services/stream/base-stream.ts b/source/services/stream/base-stream.ts index b79f824..c76ac31 100644 --- a/source/services/stream/base-stream.ts +++ b/source/services/stream/base-stream.ts @@ -1,5 +1,6 @@ /** A normal request/response result before transport encoding. */ export type StreamResponse = { + /** Optional correlation ID for multiplexed transports. */ id?: string; @@ -15,6 +16,7 @@ export type StreamResponse = { /** An application event before a transport applies its wire encoding. */ export type StreamEvent = { + /** Optional event or correlation ID. */ id?: string; diff --git a/test/helpers/test-connection.ts b/test/helpers/test-connection.ts index caa318f..538f1b5 100644 --- a/test/helpers/test-connection.ts +++ b/test/helpers/test-connection.ts @@ -1,45 +1,43 @@ -import { - BaseStream, - type StreamMessage, -} from "../../source/services/stream/base-stream.js"; +import { BaseStream, type StreamMessage } from '../../source/services/stream/base-stream.ts'; /** Minimal observable connection used by application and broadcaster tests. */ export class TestConnection extends BaseStream { - readonly messages: StreamMessage[] = []; - readonly closeCallbacks: Array<() => void> = []; - closed = false; + readonly messages: StreamMessage[] = []; + readonly closeCallbacks: Array<() => void> = []; + closed = false; - constructor( - readonly streaming: boolean, - readonly bidirectional: boolean, - ) { - super(); - } - - async send(message: StreamMessage): Promise { - if (this.closed) { - throw new Error("connection is closed"); + constructor( + readonly streaming: boolean, + readonly bidirectional: boolean, + ) { + super(); } - this.messages.push(message); - } + async send(message: StreamMessage): Promise { + if (this.closed) { + throw new Error('connection is closed'); + } - close(): void { - if (this.closed) { - return; + this.messages.push(message); } - this.closed = true; - const callbacks = this.closeCallbacks.splice(0); - callbacks.forEach((callback) => callback()); - } + close(): void { + if (this.closed) { + return; + } - onClose(callback: () => void): void { - if (this.closed) { - callback(); - return; + this.closed = true; + const callbacks = this.closeCallbacks.splice(0); + callbacks.forEach((callback) => callback()); } - this.closeCallbacks.push(callback); - } + onClose(callback: () => void): void { + if (this.closed) { + callback(); + + return; + } + + this.closeCallbacks.push(callback); + } } diff --git a/test/services/router.test.ts b/test/services/router.test.ts index 222d3e9..20c1826 100644 --- a/test/services/router.test.ts +++ b/test/services/router.test.ts @@ -1,134 +1,111 @@ -import { describe, expect, it } from "vitest"; +import { describe, expect, it } from 'vitest'; -import type { RouteDefinition, RouteModule } from "../../source/routes/types.js"; -import { ApplicationError } from "../../source/errors/index.js"; -import { ApplicationRouter } from "../../source/services/router.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 { TestConnection } from '../helpers/test-connection.ts'; -function moduleWith(routes: RouteDefinition[]): RouteModule { - return { - async getRoutes() { - return routes; - }, - }; -} +const moduleWith = (routes: RouteDefinition[]): RouteModule => { + return { + async getRoutes(): Promise { + return routes; + }, + }; +}; -describe("ApplicationRouter initialization", () => { - it("rejects duplicate exact paths during startup", async () => { - const route = { url: "/echo", handler: () => undefined }; +describe('ApplicationRouter initialization', (): void => { + it('rejects duplicate exact paths during startup', async (): Promise => { + const route = { url: '/echo', handler: (): void => undefined }; - await expect( - ApplicationRouter.create([moduleWith([route]), moduleWith([route])]), - ).rejects.toThrow("Duplicate application route: /echo"); - }); + await expect(ApplicationRouter.create([ moduleWith([ route ]), moduleWith([ route ]) ])).rejects.toThrow('Duplicate application route: /echo'); + }); - it.each(["echo", "/", "/items/:id", "/items?active=true", "/items#active"])( - "rejects the invalid route path %s", - async (url) => { - await expect( - ApplicationRouter.create([ - moduleWith([{ url, handler: () => undefined }]), - ]), - ).rejects.toThrow("Invalid application route"); - }, - ); + it.each([ 'echo', '/', '/items/:id', '/items?active=true', '/items#active' ])('rejects the invalid route path %s', async (url): Promise => { + await expect(ApplicationRouter.create([ moduleWith([{ url, handler: (): void => undefined }]) ])).rejects.toThrow('Invalid application route'); + }); }); -describe("ApplicationRouter dispatch", () => { - it("binds the connection, body, and request ID to one route stream", async () => { - const connection = new TestConnection(false, false); - const router = await ApplicationRouter.create([ - moduleWith([ - { - url: "/echo", - handler: async (stream) => { - expect(stream.connection).toBe(connection); - await stream.send(stream.body); - }, - }, - ]), - ]); +describe('ApplicationRouter dispatch', (): void => { + it('binds the connection, body, and request ID to one route stream', async (): Promise => { + const connection = new TestConnection(false, false); + const router = await ApplicationRouter.create([ + moduleWith([ + { + url: '/echo', + handler: async (stream): Promise => { + expect(stream.connection).toBe(connection); + await stream.send(stream.body); + }, + }, + ]), + ]); - await router.dispatch( - { path: "/echo", body: { value: 1 }, requestId: "request-1" }, - connection, - ); + await router.dispatch({ path: '/echo', body: { value: 1 }, requestId: 'request-1' }, connection); - expect(connection.messages).toEqual([ - { - id: "request-1", - type: "response", - statusCode: 200, - body: { value: 1 }, - }, - ]); + expect(connection.messages).toEqual([ + { + id: 'request-1', + type: 'response', + statusCode: 200, + body: { value: 1 }, + }, + ]); - await expect( - router.dispatch({ path: "/echo/other", body: {} }, connection), - ).rejects.toMatchObject({ statusCode: 404 }); - }); + await expect(router.dispatch({ path: '/echo/other', body: {} }, connection)).rejects.toMatchObject({ statusCode: 404 }); + }); - it("preserves correlation when concurrent requests finish out of order", async () => { - const completions = new Map void>(); - const router = await ApplicationRouter.create([ - moduleWith([ - { - url: "/delayed", - handler: async (stream) => { - const key = (stream.body as { key: string }).key; - await new Promise((resolve) => completions.set(key, resolve)); - await stream.send({ key }); - }, - }, - ]), - ]); - const connection = new TestConnection(true, true); + it('preserves correlation when concurrent requests finish out of order', async (): Promise => { + const completions = new Map void>(); + const router = await ApplicationRouter.create([ + moduleWith([ + { + url: '/delayed', + handler: async (stream): Promise => { + const key = (stream.body as { key: string }).key; + await new Promise((resolve) => completions.set(key, resolve)); + await stream.send({ key }); + }, + }, + ]), + ]); + const connection = new TestConnection(true, true); - const first = router.dispatch( - { path: "/delayed", body: { key: "A" }, requestId: "A" }, - connection, - ); - const second = router.dispatch( - { path: "/delayed", body: { key: "B" }, requestId: "B" }, - connection, - ); + const first = router.dispatch({ path: '/delayed', body: { key: 'A' }, requestId: 'A' }, connection); + const second = router.dispatch({ path: '/delayed', body: { key: 'B' }, requestId: 'B' }, connection); - completions.get("B")?.(); - await second; - completions.get("A")?.(); - await first; + completions.get('B')?.(); + await second; + completions.get('A')?.(); + await first; - expect(connection.messages).toEqual([ - { - id: "B", - type: "response", - statusCode: 200, - body: { key: "B" }, - }, - { - id: "A", - type: "response", - statusCode: 200, - body: { key: "A" }, - }, - ]); - }); + expect(connection.messages).toEqual([ + { + id: 'B', + type: 'response', + statusCode: 200, + body: { key: 'B' }, + }, + { + id: 'A', + type: 'response', + statusCode: 200, + body: { key: 'A' }, + }, + ]); + }); - it("propagates route failures without infrastructure-specific cleanup", async () => { - const error = new Error("route failed"); - const router = await ApplicationRouter.create([ - moduleWith([ - { - url: "/failure", - handler: () => { - throw error; - }, - }, - ]), - ]); + it('propagates route failures without infrastructure-specific cleanup', async (): Promise => { + const error = new Error('route failed'); + const router = await ApplicationRouter.create([ + moduleWith([ + { + url: '/failure', + handler: (): void => { + throw error; + }, + }, + ]), + ]); - await expect( - router.dispatch({ path: "/failure" }, new TestConnection(false, false)), - ).rejects.toBe(error); - }); + await expect(router.dispatch({ path: '/failure' }, new TestConnection(false, false))).rejects.toBe(error); + }); }); From 161b56c756bc560c118eede7ca5157f78720d475 Mon Sep 17 00:00:00 2001 From: Harvmaster Date: Mon, 3 Aug 2026 03:35:09 +0000 Subject: [PATCH 3/3] 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: {}, + }, + ]); + }); });