fix: memory leak is sse events cleanup
This commit is contained in:
parent
a1ac62bce2
commit
b5c8f71d88
2 changed files with 59 additions and 10 deletions
|
|
@ -1,5 +1,6 @@
|
||||||
import { test, describe, expect } from "bun:test";
|
import { test, describe, expect } from "bun:test";
|
||||||
import { createApp } from "~/server/app";
|
import { createApp } from "~/server/app";
|
||||||
|
import { serverEvents } from "~/server/core/events";
|
||||||
import { createTestSession, getAuthHeaders } from "~/test/helpers/auth";
|
import { createTestSession, getAuthHeaders } from "~/test/helpers/auth";
|
||||||
|
|
||||||
const app = createApp();
|
const app = createApp();
|
||||||
|
|
@ -30,6 +31,32 @@ describe("events security", () => {
|
||||||
|
|
||||||
expect(res.status).toBe(200);
|
expect(res.status).toBe(200);
|
||||||
expect(res.headers.get("Content-Type")).toBe("text/event-stream");
|
expect(res.headers.get("Content-Type")).toBe("text/event-stream");
|
||||||
|
await res.body?.cancel();
|
||||||
|
});
|
||||||
|
|
||||||
|
test("should cleanup SSE listeners when client disconnects", async () => {
|
||||||
|
const { token } = await createTestSession();
|
||||||
|
const initialCount = serverEvents.listenerCount("doctor:cancelled");
|
||||||
|
|
||||||
|
const res = await app.request("/api/v1/events", {
|
||||||
|
headers: getAuthHeaders(token),
|
||||||
|
});
|
||||||
|
|
||||||
|
expect(res.status).toBe(200);
|
||||||
|
|
||||||
|
for (let i = 0; i < 20 && serverEvents.listenerCount("doctor:cancelled") < initialCount + 1; i++) {
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||||
|
}
|
||||||
|
|
||||||
|
expect(serverEvents.listenerCount("doctor:cancelled")).toBe(initialCount + 1);
|
||||||
|
|
||||||
|
await res.body?.cancel();
|
||||||
|
|
||||||
|
for (let i = 0; i < 20 && serverEvents.listenerCount("doctor:cancelled") > initialCount; i++) {
|
||||||
|
await new Promise((resolve) => setTimeout(resolve, 10));
|
||||||
|
}
|
||||||
|
|
||||||
|
expect(serverEvents.listenerCount("doctor:cancelled")).toBe(initialCount);
|
||||||
});
|
});
|
||||||
|
|
||||||
describe("unauthenticated access", () => {
|
describe("unauthenticated access", () => {
|
||||||
|
|
|
||||||
|
|
@ -162,10 +162,13 @@ export const eventsController = new Hono().use(requireAuth).get("/", (c) => {
|
||||||
serverEvents.on("doctor:cancelled", onDoctorCancelled);
|
serverEvents.on("doctor:cancelled", onDoctorCancelled);
|
||||||
|
|
||||||
let keepAlive = true;
|
let keepAlive = true;
|
||||||
|
let cleanedUp = false;
|
||||||
|
|
||||||
stream.onAbort(() => {
|
function cleanup() {
|
||||||
logger.info("Client disconnected from SSE endpoint");
|
if (cleanedUp) return;
|
||||||
keepAlive = false;
|
cleanedUp = true;
|
||||||
|
|
||||||
|
c.req.raw.signal.removeEventListener("abort", onRequestAbort);
|
||||||
serverEvents.off("backup:started", onBackupStarted);
|
serverEvents.off("backup:started", onBackupStarted);
|
||||||
serverEvents.off("backup:progress", onBackupProgress);
|
serverEvents.off("backup:progress", onBackupProgress);
|
||||||
serverEvents.off("backup:completed", onBackupCompleted);
|
serverEvents.off("backup:completed", onBackupCompleted);
|
||||||
|
|
@ -177,14 +180,33 @@ export const eventsController = new Hono().use(requireAuth).get("/", (c) => {
|
||||||
serverEvents.off("doctor:started", onDoctorStarted);
|
serverEvents.off("doctor:started", onDoctorStarted);
|
||||||
serverEvents.off("doctor:completed", onDoctorCompleted);
|
serverEvents.off("doctor:completed", onDoctorCompleted);
|
||||||
serverEvents.off("doctor:cancelled", onDoctorCancelled);
|
serverEvents.off("doctor:cancelled", onDoctorCancelled);
|
||||||
});
|
}
|
||||||
|
|
||||||
while (keepAlive) {
|
function handleDisconnect() {
|
||||||
await stream.writeSSE({
|
if (!keepAlive) return;
|
||||||
data: JSON.stringify({ timestamp: Date.now() }),
|
logger.info("Client disconnected from SSE endpoint");
|
||||||
event: "heartbeat",
|
keepAlive = false;
|
||||||
});
|
cleanup();
|
||||||
await stream.sleep(5000);
|
}
|
||||||
|
|
||||||
|
function onRequestAbort() {
|
||||||
|
handleDisconnect();
|
||||||
|
stream.abort();
|
||||||
|
}
|
||||||
|
|
||||||
|
stream.onAbort(handleDisconnect);
|
||||||
|
c.req.raw.signal.addEventListener("abort", onRequestAbort, { once: true });
|
||||||
|
|
||||||
|
try {
|
||||||
|
while (keepAlive && !c.req.raw.signal.aborted && !stream.aborted) {
|
||||||
|
await stream.writeSSE({
|
||||||
|
data: JSON.stringify({ timestamp: Date.now() }),
|
||||||
|
event: "heartbeat",
|
||||||
|
});
|
||||||
|
await stream.sleep(5000);
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
cleanup();
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue