Skip to content
Open
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
395 changes: 395 additions & 0 deletions load-test/reconnection-storm.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,395 @@
/**
* CrowdStream — Reconnection-Storm (Thundering-Herd) Load Test
* ============================================================
* Measures RECOVERY (not cold connect) when a whole signaling pod's clients
* reconnect at once. socket.io auto-reconnection is ENABLED, so the client-side
* backoff + jitter is exercised exactly as it would be in production behind
* sticky sessions + the Socket.IO Redis adapter.
*
* A prior 10k/2s spike test saw 8,139 failures — but that only measured cold
* connect. This test quantifies how the surviving clients recover after a mass
* drop, i.e. the thundering-herd re-join.
*
* Two trigger modes:
* forcedrop (default) — the harness drops every connected client itself
* (io.engine.close ⇒ socket.io auto-reconnect); no infra
* access required.
* killpod — you kill a real signaling pod by hand; the harness does
* NOT drop clients and relies on socket.io's own disconnect
* detection.
*
* Recovery time per client = (successful re-join ack) − (storm trigger).
*
* Never hardcodes secrets: the JWT is passed via --token and sent as the
* `accessToken` cookie, matching CrowdStream's socket auth middleware.
*
* > Authored by Claude (Anthropic), via Claude Code — 2026-08-27.
*/

const { io } = require("socket.io-client");

// ---------- CLI args ----------
function arg(name, def) {
const i = process.argv.indexOf(`--${name}`);
if (i === -1) return def;
return process.argv[i + 1];
}

const URL = arg("url", "http://localhost:3000");
const URLS_ARG = arg("urls", null);
const URLS = URLS_ARG
? URLS_ARG.split(",").map((s) => s.trim()).filter(Boolean)
: [URL];

const ROOM_ID = arg("room", null);
const TOKEN = arg("token", null);
const NUM_CLIENTS = parseInt(arg("clients", "500"), 10);

Check warning on line 46 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.parseInt` over `parseInt`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzgw&open=AaBCJkPEgAAgQ-9SJzgw&pullRequest=85
const RAMP_MS = parseInt(arg("rampMs", "10000"), 10);

Check warning on line 47 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.parseInt` over `parseInt`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzgx&open=AaBCJkPEgAAgQ-9SJzgx&pullRequest=85
const MODE = arg("mode", "forcedrop");
const TRIGGER_AFTER_MS = parseInt(arg("triggerAfterMs", "15000"), 10);

Check warning on line 49 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.parseInt` over `parseInt`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzgy&open=AaBCJkPEgAAgQ-9SJzgy&pullRequest=85
const RECOVER_WINDOW_MS = parseInt(arg("recoverWindowMs", "30000"), 10);

Check warning on line 50 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.parseInt` over `parseInt`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzgz&open=AaBCJkPEgAAgQ-9SJzgz&pullRequest=85
const ACK_TIMEOUT_MS = parseInt(arg("timeout", "8000"), 10);

Check warning on line 51 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.parseInt` over `parseInt`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzg0&open=AaBCJkPEgAAgQ-9SJzg0&pullRequest=85

if (!ROOM_ID) {
console.error(
"Missing --room <roomId>. Start a broadcast first, then pass its id."
);
process.exit(1);
}

if (!TOKEN) {
console.error(
"Missing --token <jwt>. The socket auth middleware rejects unauthenticated sockets."
);
process.exit(1);
}

if (MODE !== "forcedrop" && MODE !== "killpod") {
console.error(`Invalid --mode "${MODE}". Use "forcedrop" or "killpod".`);
process.exit(1);
}

// ---------- percentile ---------- (identical to signaling-latency.js)
function percentile(sortedArr, p) {
if (sortedArr.length === 0) return NaN;

Check warning on line 74 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.NaN` over `NaN`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzg1&open=AaBCJkPEgAAgQ-9SJzg1&pullRequest=85

const idx = Math.ceil((p / 100) * sortedArr.length) - 1;

return sortedArr[Math.min(Math.max(idx, 0), sortedArr.length - 1)];
}

function summarize(label, samples) {
const clean = samples
.filter((n) => Number.isFinite(n))
.sort((a, b) => a - b);

if (clean.length === 0) {
console.log(`${label}: no samples`);
return;
}

const avg = clean.reduce((a, b) => a + b, 0) / clean.length;

console.log(
`${label.padEnd(32)} ` +
`n=${clean.length.toString().padEnd(5)} ` +
`avg=${avg.toFixed(1)}ms ` +
`p50=${percentile(clean, 50)}ms ` +
`p90=${percentile(clean, 90)}ms ` +
`p95=${percentile(clean, 95)}ms ` +
`p99=${percentile(clean, 99)}ms ` +
`max=${clean[clean.length - 1]}ms`
);
}

// ---------- Socket.IO ACK helper ---------- (same shape as signaling-latency.js)
function ackWithTimeout(socket, event, ...args) {
return new Promise((resolve, reject) => {
const timer = setTimeout(() => {
reject(new Error(`${event} ack timeout`));
}, ACK_TIMEOUT_MS);

const t0 = performance.now();

socket.emit(event, ...args, (response) => {
clearTimeout(timer);

resolve({
response,
latencyMs: performance.now() - t0,
});
});
});
}

function sleep(ms) {
return new Promise((resolve) => setTimeout(resolve, ms));
}

// ---------- shared state ----------
const clients = [];
let stormTriggered = false;
let triggerAt = 0;
const recoverSamples = []; // recovery ms for clients that re-joined within the window
let connectErrors = 0;
let reconnectAttempts = 0; // manager-level auto-reconnect attempts after the trigger

function urlFor(idx) {
return URLS[idx % URLS.length];
}

// ---------- per-client join ----------
async function joinRoomOnce(client) {
try {
const { response } = await ackWithTimeout(
client.socket,
"joinRoom",
ROOM_ID
);

if (response?.success) return true;

client.lastError = `joinRoom: ${response?.code || "unsuccessful ack"}`;
return false;
} catch (err) {
client.lastError = err?.message || String(err);
return false;
}
}

async function onConnect(client) {
if (!stormTriggered) {
// Baseline join: initial connect (or a benign pre-trigger reconnect).
client.baselineJoined = await joinRoomOnce(client);
return;
}

// Post-trigger (re)connect — this is the recovery path we are measuring.
if (client.recovered) return;

const ok = await joinRoomOnce(client);

if (ok && !client.recovered) {
client.recovered = true;
client.recoverMs = performance.now() - triggerAt;

if (client.recoverMs <= RECOVER_WINDOW_MS) {
recoverSamples.push(client.recoverMs);
}
}
}

function spawnClient(idx) {
const url = urlFor(idx);

// reconnection ENABLED so socket.io's own backoff + jitter drives recovery.
const socket = io(url, {
transports: ["websocket"],
reconnection: true,
reconnectionAttempts: Infinity,
forceNew: true,
extraHeaders: {
cookie: `accessToken=${TOKEN}`,
},
});

const client = {
idx,
url,
socket,
baselineJoined: false,
droppedAfterTrigger: false,
recovered: false,
recoverMs: NaN,

Check warning on line 203 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer `Number.NaN` over `NaN`.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzg2&open=AaBCJkPEgAAgQ-9SJzg2&pullRequest=85
lastError: null,
};

clients.push(client);

// socket.io re-emits "connect" on every successful (re)connection.
socket.on("connect", () => {
onConnect(client).catch((err) => {
client.lastError = err?.message || String(err);
});
});

socket.on("disconnect", () => {
if (stormTriggered) client.droppedAfterTrigger = true;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: In killpod mode, unrelated disconnects after the prompt are scored as victims, so recovery percentiles and success rate do not represent the killed pod. Track the affected pod/client set or otherwise distinguish the kill-induced disconnects before tallying.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At load-test/reconnection-storm.js, line 217:

<comment>In `killpod` mode, unrelated disconnects after the prompt are scored as victims, so recovery percentiles and success rate do not represent the killed pod. Track the affected pod/client set or otherwise distinguish the kill-induced disconnects before tallying.</comment>

<file context>
@@ -0,0 +1,395 @@
+  });
+
+  socket.on("disconnect", () => {
+    if (stormTriggered) client.droppedAfterTrigger = true;
+  });
+
</file context>

});

socket.on("connect_error", () => {
connectErrors += 1;
});

// Manager-level signal: how many auto-reconnect attempts the herd generated.
socket.io.on("reconnect_attempt", () => {
if (stormTriggered) reconnectAttempts += 1;
});

return client;
}

// ---------- storm triggers ----------
function triggerForcedrop() {
for (const client of clients) {
// Only drop clients that were actually in steady state; a client that never
// connected was not part of the herd and should not count against recovery.
if (!client.socket.connected) continue;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: In forcedrop mode, connected clients whose baseline joinRoom failed or is still pending are counted as recovered/stalled storm members. Require baselineJoined when selecting clients to drop.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At load-test/reconnection-storm.js, line 237:

<comment>In `forcedrop` mode, connected clients whose baseline `joinRoom` failed or is still pending are counted as recovered/stalled storm members. Require `baselineJoined` when selecting clients to drop.</comment>

<file context>
@@ -0,0 +1,395 @@
+  for (const client of clients) {
+    // Only drop clients that were actually in steady state; a client that never
+    // connected was not part of the herd and should not count against recovery.
+    if (!client.socket.connected) continue;
+
+    client.droppedAfterTrigger = true;
</file context>
Suggested change
if (!client.socket.connected) continue;
if (!client.socket.connected || !client.baselineJoined) continue;


client.droppedAfterTrigger = true;

const engine = client.socket.io && client.socket.io.engine;

Check warning on line 241 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Prefer using an optional chain expression instead, as it's more concise and easier to read.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzg3&open=AaBCJkPEgAAgQ-9SJzg3&pullRequest=85

if (engine) {
// Abrupt transport close ⇒ socket.io's automatic reconnection (backoff +
// jitter) re-establishes the socket. This is the path that actually
// smooths — or fails to smooth — the reconnection herd.
engine.close();
} else {
// Fallback if the engine is unavailable: manual bounce (immediate, no backoff).
client.socket.disconnect();
client.socket.connect();
}
}
}

function triggerKillpod() {
console.log(
"\n *** killpod mode: KILL ONE SIGNALING POD NOW ***\n" +
" The harness will NOT drop clients; it relies on socket.io's own\n" +
" disconnect detection. Recovery timing starts at this instant.\n"
);
}

// ---------- main ----------
async function main() {

Check failure on line 265 in load-test/reconnection-storm.js

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Refactor this function to reduce its Cognitive Complexity from 18 to the 15 allowed.

See more on https://sonarcloud.io/project/issues?id=Harxhit_CrowdStream&issues=AaBCJkPEgAAgQ-9SJzg4&open=AaBCJkPEgAAgQ-9SJzg4&pullRequest=85
console.log(
`Reconnection-storm test: ${NUM_CLIENTS} clients, mode=${MODE}, ` +
`ramp=${RAMP_MS}ms, triggerAfter=${TRIGGER_AFTER_MS}ms, ` +
`recoverWindow=${RECOVER_WINDOW_MS}ms, room=${ROOM_ID}\n` +
`endpoints: ${URLS.join(", ")}`
);

const t0 = performance.now();
const delayBetween = NUM_CLIENTS > 0 ? RAMP_MS / NUM_CLIENTS : 0;

for (let i = 0; i < NUM_CLIENTS; i++) {
spawnClient(i);

if (delayBetween > 0) {
await sleep(delayBetween);
}
}

// Hold until the scheduled trigger time (measured from launch) so the fleet
// can connect, join, and settle into steady state before the storm hits.
const elapsed = performance.now() - t0;
const waitBeforeTrigger = Math.max(0, TRIGGER_AFTER_MS - elapsed);
await sleep(waitBeforeTrigger);

const stable = clients.filter(
(c) => c.socket.connected && c.baselineJoined
).length;

console.log(
`\nPre-storm stability: ${stable}/${NUM_CLIENTS} connected+joined.`
);

if (stable < NUM_CLIENTS) {
console.log(
" (Not all clients reached steady state — raise --triggerAfterMs/--timeout " +
"or lower --clients. Recovery % is measured only over clients that dropped.)"
);
}

// ---- trigger the storm ----
stormTriggered = true;
triggerAt = performance.now();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: In killpod mode, operator delay is included in every recovery time and can consume the recovery window before the pod is killed. Wait for an explicit post-kill confirmation or set triggerAt when the actual kill is initiated.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. At load-test/reconnection-storm.js, line 307:

<comment>In `killpod` mode, operator delay is included in every recovery time and can consume the recovery window before the pod is killed. Wait for an explicit post-kill confirmation or set `triggerAt` when the actual kill is initiated.</comment>

<file context>
@@ -0,0 +1,395 @@
+
+  // ---- trigger the storm ----
+  stormTriggered = true;
+  triggerAt = performance.now();
+
+  if (MODE === "forcedrop") {
</file context>


if (MODE === "forcedrop") {
triggerForcedrop();
} else {
triggerKillpod();
}

console.log(
`Storm triggered. Watching recovery for ${RECOVER_WINDOW_MS}ms...`
);
await sleep(RECOVER_WINDOW_MS);

// ---- tally ----
const dropped = clients.filter((c) => c.droppedAfterTrigger);
const recovered = dropped.filter(
(c) => c.recovered && c.recoverMs <= RECOVER_WINDOW_MS
);
const failed = dropped.filter(
(c) => !(c.recovered && c.recoverMs <= RECOVER_WINDOW_MS)
);
const successPct = dropped.length
? (recovered.length / dropped.length) * 100
: 0;

console.log("\n=== Reconnection-storm results ===");
console.log(
`mode=${MODE} clients=${NUM_CLIENTS} dropped/affected=${dropped.length}`
);

summarize("reconnect + re-join", recoverSamples);

console.log(
`recovery success within ${RECOVER_WINDOW_MS}ms: ` +
`${recovered.length}/${dropped.length} (${successPct.toFixed(1)}%)`
);
console.log(`failed to recover in window: ${failed.length}`);
console.log(`auto-reconnect attempts (post-trigger): ${reconnectAttempts}`);
console.log(`connect_error events (whole run): ${connectErrors}`);

// Data-driven note on herd smoothing — computed from real samples, never fabricated.
const clean = recoverSamples
.filter((n) => Number.isFinite(n))
.sort((a, b) => a - b);

if (clean.length >= 2) {
const p50 = percentile(clean, 50);
const p99 = percentile(clean, 99);
const spread = p50 > 0 ? p99 / p50 : Infinity;

if (spread >= 3) {
console.log(
`\nHerd note: recoveries are spread out (p99/p50=${spread.toFixed(1)}x) — ` +
"backoff + jitter appears to be de-synchronising the herd."
);
} else {
console.log(
`\nHerd note: recoveries are tightly clustered (p99/p50=${spread.toFixed(1)}x) — ` +
"clients reconnected in a near-simultaneous burst; verify " +
"reconnectionDelayMax / randomizationFactor and LB connection limits."
);
}
}

if (failed.length > 0) {
const samples = failed
.slice(0, 10)
.map(
(c) =>
`client ${c.idx}@${c.url}: ${
c.lastError || "no successful re-join in window"
}`
);

console.log("\nsample failures:");
console.log(samples.join("\n"));

if (failed.length > 10) {
console.log(`...and ${failed.length - 10} more`);
}
}

// Teardown: reconnection is on (attempts=Infinity), so we must close every
// socket and exit explicitly or the process would never terminate.
for (const client of clients) client.socket.disconnect();
process.exit(0);
}

main();
Loading