Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 16 additions & 6 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

38 changes: 23 additions & 15 deletions src/lib/session-sse.ts
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@ export function connectSessionSse(
async function connect() {
if (closed) return;
controller = new AbortController();
let shouldRetry = true;

try {
const baseUrl = env.API_BASE_URL.replace(/\/+$/, "");
Expand All @@ -121,32 +122,39 @@ export function connectSessionSse(
return;
}
console.warn(
`[session-sse] /auth/events returned ${response.status} — retrying`,
`[session-sse] /auth/events returned ${response.status} — ${response.status === 401 || response.status === 403 ? "stopping" : "retrying"}`,
);
return;
}

// Reset retry count on successful connection
retryCount = 0;

const reader = response.body?.getReader();
if (!reader) {
console.warn("[session-sse] Response body has no readable stream");
return; // triggers reconnect with backoff
if (response.status === 401 || response.status === 403) {
shouldRetry = false;
}
} else {
// Reset retry count on successful connection
retryCount = 0;
}

const decoder = new TextDecoder();
let buffer = "";
if (shouldRetry) {
const reader = response.body?.getReader();
if (!reader) {
console.warn("[session-sse] Response body has no readable stream");
} else {
const decoder = new TextDecoder();
let buffer = "";

while (!closed) {
const { done, value } = await reader.read();
if (done) break;
while (!closed) {
const { done, value } = await reader.read();
if (done) break;

buffer += decoder.decode(value, { stream: true });
buffer += decoder.decode(value, { stream: true });

// Parse SSE frames
const frames = buffer.split("\n\n");
buffer = frames.pop() ?? ""; // Keep incomplete frame in buffer
// Parse SSE frames
const frames = buffer.split("\n\n");
buffer = frames.pop() ?? ""; // Keep incomplete frame in buffer

for (const frame of frames) {
const event = parseSessionSseFrame(frame);
Expand All @@ -159,7 +167,7 @@ export function connectSessionSse(
}

// Reconnect with backoff if not closed
if (!closed) {
if (!closed && shouldRetry) {
retryCount++;
const delay = getReconnectDelay();
retryTimeout = setTimeout(connect, delay);
Expand Down
129 changes: 129 additions & 0 deletions tests/session-sse.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
/** ApexChain Network Operations Intelligence Platform */
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";

import { connectSessionSse, processFrame } from "@/lib/session-sse";

describe("processFrame", () => {
const onEvent = vi.fn();

beforeEach(() => {
onEvent.mockClear();
});

afterEach(() => {
vi.restoreAllMocks();
});

it("handles a single-line session_revoked frame", () => {
processFrame(
'event: session_revoked\ndata: {"reason":"admin_logout"}\n\n',
onEvent,
);

expect(onEvent).toHaveBeenCalledTimes(1);
expect(onEvent).toHaveBeenCalledWith({
type: "session_revoked",
reason: "admin_logout",
});
});

it("keeps the current multi-line data truncation behavior", () => {
processFrame(
'event: session_revoked\ndata: {"reason":"admin_logout"}\ndata: {"reason":"password_changed"}\n\n',
onEvent,
);

expect(onEvent).toHaveBeenCalledTimes(1);
expect(onEvent).toHaveBeenCalledWith({
type: "session_revoked",
reason: "password_changed",
});
});

it("ignores malformed JSON payloads", () => {
processFrame('event: session_revoked\ndata: {"reason":\n\n', onEvent);

expect(onEvent).not.toHaveBeenCalled();
});

it("ignores unknown event types", () => {
processFrame('event: heartbeat\ndata: {"timestamp":1234567890}\n\n', onEvent);

expect(onEvent).not.toHaveBeenCalled();
});

it("ignores a session_revoked payload when the event line is missing", () => {
processFrame('data: {"reason":"session_expired"}\n\n', onEvent);

expect(onEvent).not.toHaveBeenCalled();
});

it.each([
["admin_logout", "admin_logout"],
["password_changed", "password_changed"],
["session_expired", "session_expired"],
["not_a_known_reason", "unknown"],
])("maps reason %s to %s", (reason, expected) => {
processFrame(
`event: session_revoked\ndata: {"reason":"${reason}"}\n\n`,
onEvent,
);

expect(onEvent).toHaveBeenCalledTimes(1);
expect(onEvent).toHaveBeenCalledWith({
type: "session_revoked",
reason: expected,
});
});
});

describe("connectSessionSse reconnect behavior", () => {
beforeEach(() => {
vi.useFakeTimers();
});

afterEach(() => {
vi.runOnlyPendingTimers();
vi.useRealTimers();
vi.unstubAllGlobals();
});

it("schedules a retry after a non-200 response", async () => {
const fetchMock = vi.fn().mockResolvedValue({
ok: false,
status: 503,
});
vi.stubGlobal("fetch", fetchMock);
vi.spyOn(Math, "random").mockReturnValue(0.5);

const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
const connection = connectSessionSse(() => {});

await Promise.resolve();

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(setTimeoutSpy).toHaveBeenCalledTimes(1);
expect(setTimeoutSpy).toHaveBeenCalledWith(expect.any(Function), expect.any(Number));

connection.close();
});

it("stops retrying after 401 or 403 responses", async () => {
const fetchMock = vi.fn().mockResolvedValue({
ok: false,
status: 401,
});
vi.stubGlobal("fetch", fetchMock);

const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
const connection = connectSessionSse(() => {});

await Promise.resolve();
await vi.advanceTimersByTimeAsync(1_000);

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(setTimeoutSpy).not.toHaveBeenCalled();

connection.close();
});
});
1 change: 1 addition & 0 deletions vitest-session-sse.json
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
{"numTotalTestSuites":3,"numPassedTestSuites":1,"numFailedTestSuites":2,"numPendingTestSuites":0,"numTotalTests":11,"numPassedTests":10,"numFailedTests":1,"numPendingTests":0,"numTodoTests":0,"snapshot":{"added":0,"failure":false,"filesAdded":0,"filesRemoved":0,"filesRemovedList":[],"filesUnmatched":0,"filesUpdated":0,"matched":0,"total":0,"unchecked":0,"uncheckedKeysByFile":[],"unmatched":0,"updated":0,"didUpdate":false},"startTime":1788117061551,"success":false,"testResults":[{"assertionResults":[{"ancestorTitles":["processFrame"],"fullName":"processFrame handles a single-line session_revoked frame","status":"passed","title":"handles a single-line session_revoked frame","duration":3.6824000000001433,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame keeps the current multi-line data truncation behavior","status":"passed","title":"keeps the current multi-line data truncation behavior","duration":0.9655999999999949,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame ignores malformed JSON payloads","status":"passed","title":"ignores malformed JSON payloads","duration":0.6653999999998632,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame ignores unknown event types","status":"passed","title":"ignores unknown event types","duration":0.49739999999974316,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame ignores a session_revoked payload when the event line is missing","status":"passed","title":"ignores a session_revoked payload when the event line is missing","duration":0.30960000000004584,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame maps reason admin_logout to admin_logout","status":"passed","title":"maps reason admin_logout to admin_logout","duration":0.7570999999998094,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame maps reason password_changed to password_changed","status":"passed","title":"maps reason password_changed to password_changed","duration":0.6610000000000582,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame maps reason session_expired to session_expired","status":"passed","title":"maps reason session_expired to session_expired","duration":0.9315000000001419,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["processFrame"],"fullName":"processFrame maps reason not_a_known_reason to unknown","status":"passed","title":"maps reason not_a_known_reason to unknown","duration":0.7445000000002437,"failureMessages":[],"meta":{},"tags":[]},{"ancestorTitles":["connectSessionSse reconnect behavior"],"fullName":"connectSessionSse reconnect behavior schedules a retry after a non-200 response","status":"failed","title":"schedules a retry after a non-200 response","duration":846.2370000000001,"failureMessages":["Error: Aborting after running 10000 timers, assuming an infinite loop!\nTimeout - connect\n at connect (C:/Users/a-ahmedyahaya/Desktop/dev/abdulsnk/ApexChainx-Frontend/src/lib/session-sse.ts:137:22)"],"meta":{},"tags":[]},{"ancestorTitles":["connectSessionSse reconnect behavior"],"fullName":"connectSessionSse reconnect behavior stops retrying after 401 or 403 responses","status":"passed","title":"stops retrying after 401 or 403 responses","duration":2.9220000000000255,"failureMessages":[],"meta":{},"tags":[]}],"startTime":1788117063857,"endTime":1788117064716.922,"status":"failed","message":"","name":"C:/Users/a-ahmedyahaya/Desktop/dev/abdulsnk/ApexChainx-Frontend/tests/session-sse.test.ts"}]}