Files
sync-server-v2/test/services/broadcaster.test.ts
T
2026-07-27 10:22:01 +00:00

200 lines
6.6 KiB
TypeScript

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";
function createBroadcaster(): Broadcaster {
return new Broadcaster(new Logger("broadcaster-test"));
}
function routeStream(connection: TestConnection): ApplicationRouteStream {
return new ApplicationRouteStream(connection, undefined);
}
async function expectPending(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 () => {
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" },
});
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();
});
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<void>((resolve) => {
releaseFirst = resolve;
});
let markSecondSent: () => void = () => undefined;
const secondSent = new Promise<void>((resolve) => {
markSecondSent = resolve;
});
first.send = async (message) => {
await firstReleased;
await originalFirstSend(message);
};
second.send = async (message) => {
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 () => {
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();
});
});