Skip to content

Commit 65211b9

Browse files
authored
Merge pull request #125 from native-federation/fix/sse-bounded-pool
fix(sse): bound the build-notification connection pool
2 parents 3d6dac6 + d2c5405 commit 65211b9

3 files changed

Lines changed: 229 additions & 1 deletion

File tree

Lines changed: 197 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,197 @@
1+
import type { IncomingMessage, ServerResponse } from "http";
2+
import { afterEach, describe, expect, it, vi } from "vitest";
3+
4+
vi.mock("@softarc/native-federation/internal", () => ({
5+
logger: { info: vi.fn(), verbose: vi.fn(), error: vi.fn(), warn: vi.fn() },
6+
}));
7+
8+
import { BuildNotificationType } from "@softarc/native-federation";
9+
10+
import { federationBuildNotifier } from "./federation-build-notifier.js";
11+
12+
const ENDPOINT = "/@angular-architects/native-federation:build-notifications";
13+
14+
// Minimal stand-ins for the node req/res pair: the notifier only writes to the response
15+
// and subscribes to the request's close/error events.
16+
function createClient() {
17+
const listeners = new Map<string, () => void>();
18+
const written: string[] = [];
19+
20+
const req = {
21+
destroyed: false,
22+
on(event: string, callback: () => void) {
23+
listeners.set(event, callback);
24+
return req;
25+
},
26+
} as unknown as IncomingMessage;
27+
28+
const res = {
29+
destroyed: false,
30+
writable: true,
31+
writeHead: vi.fn(),
32+
write: vi.fn((chunk: string) => {
33+
written.push(chunk);
34+
return true;
35+
}),
36+
end: vi.fn(() => {
37+
res.destroyed = true;
38+
res.writable = false;
39+
}),
40+
};
41+
42+
return {
43+
req,
44+
res: res as unknown as ServerResponse,
45+
written,
46+
ended: () => res.end.mock.calls.length > 0,
47+
fire: (event: string) => listeners.get(event)?.(),
48+
};
49+
}
50+
51+
function connect(clients: number) {
52+
const middleware = federationBuildNotifier.createEventMiddleware(() => ENDPOINT);
53+
const next = vi.fn();
54+
55+
return Array.from({ length: clients }, () => {
56+
const client = createClient();
57+
middleware(client.req, client.res, next);
58+
return client;
59+
});
60+
}
61+
62+
// The notifier is a module-level singleton, so each test has to hand back a clean pool.
63+
afterEach(() => {
64+
federationBuildNotifier.stopEventServer();
65+
vi.clearAllMocks();
66+
});
67+
68+
describe("createEventMiddleware", () => {
69+
it("passes the request along when the notifier is inactive", () => {
70+
const middleware = federationBuildNotifier.createEventMiddleware(() => ENDPOINT);
71+
const client = createClient();
72+
const next = vi.fn();
73+
74+
middleware(client.req, client.res, next);
75+
76+
expect(next).toHaveBeenCalled();
77+
expect(client.res.writeHead).not.toHaveBeenCalled();
78+
});
79+
80+
it("passes the request along when the url is not the endpoint", () => {
81+
federationBuildNotifier.initialize(ENDPOINT);
82+
const middleware = federationBuildNotifier.createEventMiddleware(() => "/main.js");
83+
const client = createClient();
84+
const next = vi.fn();
85+
86+
middleware(client.req, client.res, next);
87+
88+
expect(next).toHaveBeenCalled();
89+
expect(client.res.writeHead).not.toHaveBeenCalled();
90+
});
91+
92+
it("opens an event stream on the endpoint", () => {
93+
federationBuildNotifier.initialize(ENDPOINT);
94+
95+
const [client] = connect(1);
96+
97+
expect(client!.res.writeHead).toHaveBeenCalledWith(
98+
200,
99+
expect.objectContaining({ "Content-Type": "text/event-stream" }),
100+
);
101+
expect(federationBuildNotifier.activeConnections).toBe(1);
102+
});
103+
104+
it("declares the reconnect delay before the first event", () => {
105+
federationBuildNotifier.initialize(ENDPOINT);
106+
107+
const [client] = connect(1);
108+
109+
expect(client!.written[0]).toBe("retry: 5000\n");
110+
expect(client!.written[1]).toContain('"type":"connected"');
111+
});
112+
113+
it("drops a connection from the pool when the request closes", () => {
114+
federationBuildNotifier.initialize(ENDPOINT);
115+
const [client] = connect(1);
116+
117+
client!.fire("close");
118+
119+
expect(federationBuildNotifier.activeConnections).toBe(0);
120+
});
121+
});
122+
123+
describe("connection limit", () => {
124+
it("holds at most 16 connections", () => {
125+
federationBuildNotifier.initialize(ENDPOINT);
126+
127+
connect(20);
128+
129+
expect(federationBuildNotifier.activeConnections).toBe(16);
130+
});
131+
132+
it("evicts the oldest connection rather than refusing the newest", () => {
133+
federationBuildNotifier.initialize(ENDPOINT);
134+
135+
const clients = connect(17);
136+
137+
expect(clients[0]!.ended()).toBe(true);
138+
expect(clients[16]!.ended()).toBe(false);
139+
});
140+
141+
it("stops broadcasting to an evicted connection", () => {
142+
federationBuildNotifier.initialize(ENDPOINT);
143+
const clients = connect(17);
144+
const evicted = clients[0]!;
145+
const writesBefore = evicted.written.length;
146+
147+
federationBuildNotifier.broadcastBuildCompletion();
148+
149+
expect(evicted.written.length).toBe(writesBefore);
150+
expect(clients[16]!.written.at(-1)).toContain(BuildNotificationType.COMPLETED);
151+
});
152+
});
153+
154+
describe("broadcasts", () => {
155+
it("sends the completion event to every connection", () => {
156+
federationBuildNotifier.initialize(ENDPOINT);
157+
const clients = connect(3);
158+
159+
federationBuildNotifier.broadcastBuildCompletion();
160+
161+
for (const client of clients) {
162+
expect(client.written.at(-1)).toContain(BuildNotificationType.COMPLETED);
163+
}
164+
});
165+
166+
it("sends the error message with the error event", () => {
167+
federationBuildNotifier.initialize(ENDPOINT);
168+
const [client] = connect(1);
169+
170+
federationBuildNotifier.broadcastBuildError(new Error("boom"));
171+
172+
expect(client!.written.at(-1)).toContain(BuildNotificationType.ERROR);
173+
expect(client!.written.at(-1)).toContain("boom");
174+
});
175+
176+
it("does nothing when the notifier is inactive", () => {
177+
const [client] = connect(1);
178+
const writesBefore = client!.written.length;
179+
180+
federationBuildNotifier.broadcastBuildCancellation();
181+
182+
expect(client!.written.length).toBe(writesBefore);
183+
});
184+
});
185+
186+
describe("stopEventServer", () => {
187+
it("closes every connection and empties the pool", () => {
188+
federationBuildNotifier.initialize(ENDPOINT);
189+
const clients = connect(3);
190+
191+
federationBuildNotifier.stopEventServer();
192+
193+
expect(clients.every((client) => client.ended())).toBe(true);
194+
expect(federationBuildNotifier.activeConnections).toBe(0);
195+
expect(federationBuildNotifier.isRunning).toBe(false);
196+
});
197+
});

‎src/builders/build/federation-build-notifier.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,9 @@ interface FederationEvent {
2121
type NextFunction = (error?: Error) => void;
2222
type MiddlewareFunction = (req: IncomingMessage, res: ServerResponse, next: NextFunction) => void;
2323

24+
const MAX_CONNECTIONS = 16;
25+
const RECONNECT_DELAY_MS = 5000;
26+
2427
/**
2528
* Manages Server-Sent Events for federation hot reload in local development
2629
* Only active when running in development mode with dev server
@@ -71,6 +74,8 @@ class FederationBuildNotifier {
7174
* Sets up a new SSE connection
7275
*/
7376
private _setupSSEConnection(req: IncomingMessage, res: ServerResponse): void {
77+
this._evictOverflow();
78+
7479
res.writeHead(200, {
7580
'Content-Type': 'text/event-stream',
7681
'Cache-Control': 'no-cache',
@@ -79,6 +84,9 @@ class FederationBuildNotifier {
7984
'Access-Control-Allow-Headers': 'Cache-Control',
8085
});
8186

87+
// Pins the reconnect backoff instead of leaving it to the browser default.
88+
res.write(`retry: ${RECONNECT_DELAY_MS}\n`);
89+
8290
// Send initial connection event
8391
this._sendEvent(res, {
8492
type: 'connected',
@@ -98,6 +106,28 @@ class FederationBuildNotifier {
98106
);
99107
}
100108

109+
/**
110+
* Drops the oldest connections once the pool is full
111+
*
112+
* Behind a reverse proxy the peer is the proxy rather than the browser, so a stream the
113+
* client already abandoned still looks writable here and cannot be detected as stale.
114+
* Bounding the pool caps how many such connections can pile up.
115+
*/
116+
private _evictOverflow(): void {
117+
while (this.connections.length >= MAX_CONNECTIONS) {
118+
const oldest = this.connections.shift();
119+
if (!oldest) return;
120+
121+
try {
122+
oldest.response.end();
123+
} catch {
124+
// Connection might already be closed
125+
}
126+
127+
logger.info('[Federation SSE] Connection limit reached, dropped the oldest connection');
128+
}
129+
}
130+
101131
/**
102132
* Removes a connection from the pool
103133
*/

‎src/builders/build/schema.json‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,11 +71,12 @@
7171
},
7272
"buildNotifications": {
7373
"type": "object",
74+
"description": "Opt in to build completion notifications. Omitting this option leaves them off; the defaults below only apply once it is present.",
7475
"properties": {
7576
"enable": {
7677
"type": "boolean",
7778
"default": true,
78-
"description": "Enable build completion notifications for local development. It will send events to notify when the federation build is complete."
79+
"description": "Enable build completion notifications for local development. It will send events to notify when the federation build is complete. Each connected browser tab holds an event stream open, so the runtime also has to opt in with 'sse: true'."
7980
},
8081
"endpoint": {
8182
"type": "string",

0 commit comments

Comments
 (0)