Formatting
This commit is contained in:
@@ -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<string>;
|
||||
|
||||
|
||||
@@ -1,133 +1,119 @@
|
||||
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 {
|
||||
const routeStream = (connection: TestConnection): ApplicationRouteStream => {
|
||||
return new ApplicationRouteStream(connection, undefined);
|
||||
}
|
||||
};
|
||||
|
||||
async function expectPending(promise: Promise<void>): Promise<void> {
|
||||
const expectPending = async (promise: Promise<void>): Promise<void> => {
|
||||
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 () => {
|
||||
describe('Broadcaster subscriptions', () => {
|
||||
it('delivers events and resolves after a later request removes the topic', async (): Promise<void> => {
|
||||
const broadcaster = createBroadcaster();
|
||||
const connection = new TestConnection(true, true);
|
||||
const subscribed = broadcaster.subscribe(routeStream(connection), [
|
||||
"items",
|
||||
]);
|
||||
const subscribed = broadcaster.subscribe(routeStream(connection), [ 'items' ]);
|
||||
|
||||
await expectPending(subscribed);
|
||||
await broadcaster.publish("items", {
|
||||
type: "item-changed",
|
||||
data: { id: "a" },
|
||||
await broadcaster.publish('items', {
|
||||
type: 'item-changed',
|
||||
data: { id: 'a' },
|
||||
});
|
||||
|
||||
expect(connection.messages).toEqual([
|
||||
expect.objectContaining({
|
||||
type: "item-changed",
|
||||
data: { id: "a" },
|
||||
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 broadcaster.unsubscribe(routeStream(connection), [ 'items' ]);
|
||||
await expect(subscribed).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it("resolves fully duplicate subscriptions immediately", async () => {
|
||||
it('resolves fully duplicate subscriptions immediately', async (): Promise<void> => {
|
||||
const broadcaster = createBroadcaster();
|
||||
const connection = new TestConnection(true, true);
|
||||
const first = broadcaster.subscribe(routeStream(connection), [
|
||||
"items",
|
||||
"items",
|
||||
]);
|
||||
const duplicate = broadcaster.subscribe(routeStream(connection), ["items"]);
|
||||
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 broadcaster.unsubscribe(routeStream(connection), [ 'items' ]);
|
||||
await first;
|
||||
});
|
||||
|
||||
it("waits only for topics newly added by a partially overlapping call", async () => {
|
||||
it('waits only for topics newly added by a partially overlapping call', async (): Promise<void> => {
|
||||
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 first = broadcaster.subscribe(routeStream(connection), [ 'a' ]);
|
||||
const second = broadcaster.subscribe(routeStream(connection), [ 'a', 'b' ]);
|
||||
|
||||
await broadcaster.unsubscribe(routeStream(connection), ["a"]);
|
||||
await broadcaster.unsubscribe(routeStream(connection), [ 'a' ]);
|
||||
await expect(first).resolves.toBeUndefined();
|
||||
await expectPending(second);
|
||||
|
||||
await broadcaster.unsubscribe(routeStream(connection), ["b"]);
|
||||
await broadcaster.unsubscribe(routeStream(connection), [ 'b' ]);
|
||||
await expect(second).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it("resolves every pending subscription and removes topics on close", async () => {
|
||||
it('resolves every pending subscription and removes topics on close', async (): Promise<void> => {
|
||||
const broadcaster = createBroadcaster();
|
||||
const connection = new TestConnection(true, false);
|
||||
const first = broadcaster.subscribe(routeStream(connection), ["a"]);
|
||||
const second = broadcaster.subscribe(routeStream(connection), ["b"]);
|
||||
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 });
|
||||
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 () => {
|
||||
it('immediately resolves registration against an already-closed connection', async (): Promise<void> => {
|
||||
const broadcaster = createBroadcaster();
|
||||
const connection = new TestConnection(true, false);
|
||||
connection.close();
|
||||
|
||||
await expect(
|
||||
broadcaster.subscribe(routeStream(connection), ["items"]),
|
||||
).resolves.toBeUndefined();
|
||||
await expect(broadcaster.subscribe(routeStream(connection), [ 'items' ])).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
it("rejects subscriptions on a non-streaming connection", () => {
|
||||
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 }),
|
||||
);
|
||||
expect(() => broadcaster.subscribe(routeStream(connection), [ 'items' ])).toThrowError(expect.objectContaining({ statusCode: 406 }));
|
||||
});
|
||||
|
||||
it("treats an empty subscription and repeated unsubscribe as no-ops", async () => {
|
||||
it('treats an empty subscription and repeated unsubscribe as no-ops', async (): Promise<void> => {
|
||||
const broadcaster = createBroadcaster();
|
||||
const connection = new TestConnection(true, true);
|
||||
|
||||
await expect(
|
||||
broadcaster.subscribe(routeStream(connection), []),
|
||||
).resolves.toBeUndefined();
|
||||
await broadcaster.unsubscribe(routeStream(connection), ["missing"]);
|
||||
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 () => {
|
||||
it('fans out concurrently to independent connections', async (): Promise<void> => {
|
||||
const broadcaster = createBroadcaster();
|
||||
const first = new TestConnection(true, false);
|
||||
const second = new TestConnection(true, false);
|
||||
@@ -141,23 +127,20 @@ describe("Broadcaster subscriptions", () => {
|
||||
markSecondSent = resolve;
|
||||
});
|
||||
|
||||
first.send = async (message) => {
|
||||
first.send = async (message): Promise<void> => {
|
||||
await firstReleased;
|
||||
await originalFirstSend(message);
|
||||
};
|
||||
second.send = async (message) => {
|
||||
|
||||
second.send = async (message): Promise<void> => {
|
||||
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",
|
||||
const firstSubscription = broadcaster.subscribe(routeStream(first), [ 'items' ]);
|
||||
const secondSubscription = broadcaster.subscribe(routeStream(second), [ 'items' ]);
|
||||
const publication = broadcaster.publish('items', {
|
||||
type: 'item-changed',
|
||||
data: {},
|
||||
});
|
||||
|
||||
@@ -170,19 +153,17 @@ describe("Broadcaster subscriptions", () => {
|
||||
|
||||
first.close();
|
||||
second.close();
|
||||
await Promise.all([firstSubscription, secondSubscription]);
|
||||
await Promise.all([ firstSubscription, secondSubscription ]);
|
||||
});
|
||||
|
||||
it("closes and removes a connection whose event delivery fails", async () => {
|
||||
it('closes and removes a connection whose event delivery fails', async (): Promise<void> => {
|
||||
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",
|
||||
]);
|
||||
connection.send = vi.fn().mockRejectedValue(new Error('socket failed'));
|
||||
const subscribed = broadcaster.subscribe(routeStream(connection), [ 'items' ]);
|
||||
|
||||
await broadcaster.publish("items", {
|
||||
type: "item-changed",
|
||||
await broadcaster.publish('items', {
|
||||
type: 'item-changed',
|
||||
data: {},
|
||||
});
|
||||
await subscribed;
|
||||
@@ -190,8 +171,8 @@ describe("Broadcaster subscriptions", () => {
|
||||
expect(connection.closed).toBe(true);
|
||||
expect(connection.send).toHaveBeenCalledOnce();
|
||||
|
||||
await broadcaster.publish("items", {
|
||||
type: "item-changed",
|
||||
await broadcaster.publish('items', {
|
||||
type: 'item-changed',
|
||||
data: {},
|
||||
});
|
||||
expect(connection.send).toHaveBeenCalledOnce();
|
||||
|
||||
@@ -1,41 +1,41 @@
|
||||
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 {
|
||||
const moduleWith = (routes: RouteDefinition[]): RouteModule => {
|
||||
return {
|
||||
async getRoutes() {
|
||||
async getRoutes(): Promise<RouteDefinition[]> {
|
||||
return routes;
|
||||
},
|
||||
};
|
||||
}
|
||||
};
|
||||
|
||||
async function expectPending(promise: Promise<void>): Promise<void> {
|
||||
const expectPending = async (promise: Promise<void>): Promise<void> => {
|
||||
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"));
|
||||
describe('long-lived subscription dispatch', (): void => {
|
||||
it('allows duplicate subscribe and unsubscribe requests while the original request waits', async (): Promise<void> => {
|
||||
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"]);
|
||||
url: '/items/subscribe',
|
||||
handler: async (stream): Promise<void> => {
|
||||
await broadcaster.subscribe(stream, [ 'items' ]);
|
||||
},
|
||||
},
|
||||
{
|
||||
url: "/items/unsubscribe",
|
||||
handler: async (stream) => {
|
||||
await broadcaster.unsubscribe(stream, ["items"]);
|
||||
url: '/items/unsubscribe',
|
||||
handler: async (stream): Promise<void> => {
|
||||
await broadcaster.unsubscribe(stream, [ 'items' ]);
|
||||
await stream.send({});
|
||||
},
|
||||
},
|
||||
@@ -43,31 +43,22 @@ describe("long-lived subscription dispatch", () => {
|
||||
]);
|
||||
const connection = new TestConnection(true, true);
|
||||
|
||||
const original = router.dispatch(
|
||||
{ path: "/items/subscribe", requestId: "subscribe-1" },
|
||||
connection,
|
||||
);
|
||||
const original = router.dispatch({ path: '/items/subscribe', requestId: 'subscribe-1' }, connection);
|
||||
await vi.waitFor(() => expect(connection.closeCallbacks).toHaveLength(1));
|
||||
await expectPending(original);
|
||||
|
||||
// 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 router.dispatch({ path: '/items/subscribe', requestId: 'subscribe-2' }, connection);
|
||||
await expectPending(original);
|
||||
|
||||
await router.dispatch(
|
||||
{ path: "/items/unsubscribe", requestId: "unsubscribe-1" },
|
||||
connection,
|
||||
);
|
||||
await router.dispatch({ path: '/items/unsubscribe', requestId: 'unsubscribe-1' }, connection);
|
||||
await original;
|
||||
|
||||
expect(connection.messages).toEqual([
|
||||
{
|
||||
id: "unsubscribe-1",
|
||||
type: "response",
|
||||
id: 'unsubscribe-1',
|
||||
type: 'response',
|
||||
statusCode: 200,
|
||||
body: {},
|
||||
},
|
||||
|
||||
Reference in New Issue
Block a user