fix(worker): add jitter to default reconnect backoff to prevent thundering herd (#116279)

* fix(worker): stagger gateway reconnect attempts

Co-authored-by: 赵旺0668001248 <zhao.wang1@xydigit.com>

* test(worker): preserve gateway lifecycle during reconnect proof

* test(worker): decode websocket frames with canonical helper

---------

Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
zw-xysk
2026-08-26 14:17:13 +08:00
committed by GitHub
parent e3958f882f
commit 4e4d595476
2 changed files with 135 additions and 2 deletions
+134 -1
View File
@@ -1,6 +1,8 @@
import { once } from "node:events";
import net from "node:net";
import { rawDataToString } from "@openclaw/gateway-client/websocket-data";
import { describe, expect, it, vi } from "vitest";
import type { WebSocket } from "ws";
import { WebSocketServer, type WebSocket } from "ws";
import {
GATEWAY_CLIENT_IDS,
GATEWAY_CLIENT_MODES,
@@ -255,6 +257,137 @@ describe("worker connection endpoint failures", () => {
});
});
describe("worker connection reconnect backoff", () => {
it("staggers twenty workers recovering from transient Gateway transport loss", async () => {
const server = new WebSocketServer({ host: "127.0.0.1", port: 0 });
await once(server, "listening");
const address = server.address();
if (!address || typeof address === "string") {
throw new Error("test gateway did not allocate a TCP port");
}
let available = true;
let transportInterrupted = false;
const unavailableWorkers = new Set<string>();
const recoveredWorkers = new Set<string>();
server.on("connection", (socket) => {
socket.on("message", (data) => {
const frame = JSON.parse(rawDataToString(data)) as {
id: string;
params: { admission: WorkerConnectParams["admission"] };
};
const admission = frame.params.admission;
if (!available) {
unavailableWorkers.add(admission.environmentId);
socket.send(
JSON.stringify({
type: "res",
id: frame.id,
ok: false,
error: {
code: "INVALID_REQUEST",
message: "gateway temporarily unavailable",
details: { reason: "gateway-unavailable" },
retryable: true,
},
}),
);
return;
}
if (transportInterrupted) {
recoveredWorkers.add(admission.environmentId);
}
socket.send(
JSON.stringify({
type: "res",
id: frame.id,
ok: true,
payload: {
type: "worker-hello-ok",
environmentId: admission.environmentId,
sessionId: admission.sessionId,
ownerEpoch: admission.ownerEpoch,
rpcSetVersion: admission.rpcSetVersion,
protocolFeatures: [...admission.handshake.protocolFeatures],
credentialExpiresAtMs: Date.now() + 60_000,
policy: { heartbeatIntervalMs: 60_000, maxPayload: 25 * 1024 * 1024 },
},
}),
);
});
});
const workers = Array.from({ length: 20 }, (_, index) =>
createWorkerConnection({
endpoint: {
kind: "websocket",
url: `ws://127.0.0.1:${address.port}${WORKER_PUBLIC_INGRESS_PATH}`,
},
connectParams: {
...FRAME_CONNECT_PARAMS,
admission: {
...FRAME_CONNECT_PARAMS.admission,
environmentId: `reconnect-worker-${index}`,
},
},
admissionDeadlineMs: 10_000,
}),
);
let randomDraw = 0;
const random = vi
.spyOn(Math, "random")
.mockImplementation(() => ((randomDraw++ % 20) + 1) / 21);
const retryDelays: number[] = [];
const originalSetTimeout = globalThis.setTimeout;
const timeout = vi
.spyOn(globalThis, "setTimeout")
.mockImplementation((callback, delay, ...args) => {
if (typeof delay === "number" && delay >= 250 && delay <= 275) {
retryDelays.push(delay);
}
return originalSetTimeout(callback, delay, ...args);
});
try {
await Promise.all(workers.map((worker) => worker.start()));
const readyAgain = workers.map(
(worker) =>
new Promise<void>((resolve) => {
const unsubscribe = worker.onReady(() => {
unsubscribe();
resolve();
});
}),
);
available = false;
transportInterrupted = true;
for (const socket of server.clients) {
socket.close(1012, "gateway-unavailable");
}
await vi.waitFor(() => expect(unavailableWorkers.size).toBe(20), {
timeout: 3_000,
interval: 5,
});
available = true;
await Promise.all(readyAgain);
expect(recoveredWorkers.size).toBe(20);
expect(retryDelays).toHaveLength(20);
expect(new Set(retryDelays).size).toBeGreaterThanOrEqual(10);
expect(Math.max(...retryDelays)).toBeLessThanOrEqual(30_000);
} finally {
random.mockRestore();
timeout.mockRestore();
await Promise.all(workers.map((worker) => worker.stop()));
await new Promise<void>((resolve, reject) => {
server.close((error) => (error ? reject(error) : resolve()));
});
}
});
});
describe("worker connection error coercion", () => {
it("preserves structured non-Error causes", () => {
const cause = { code: "ECONNRESET", status: 503 };
+1 -1
View File
@@ -58,7 +58,7 @@ const DEFAULT_RECONNECT_BACKOFF: BackoffPolicy = {
initialMs: 250,
maxMs: 30_000,
factor: 2,
jitter: 0,
jitter: 0.1,
};
const DEFAULT_ADMISSION_TIMEOUT_MS = DEFAULT_PREAUTH_HANDSHAKE_TIMEOUT_MS;