Skip to content

Commit 2f3769f

Browse files
committed
fix(tlon): clear SSE connect timeout after failed openStream
1 parent 242cdca commit 2f3769f

3 files changed

Lines changed: 190 additions & 32 deletions

File tree

Lines changed: 120 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,120 @@
1+
import { createServer, type Server } from "node:http";
2+
import type { AddressInfo } from "node:net";
3+
import type { LookupFn } from "openclaw/plugin-sdk/ssrf-runtime";
4+
import { afterEach, describe, expect, it, vi } from "vitest";
5+
import { UrbitSSEClient } from "./sse-client.js";
6+
7+
const CONNECT_TIMEOUT_MS = 60_000;
8+
const STORM_ATTEMPTS = 5;
9+
10+
const lookupLoopback = (async () => [{ address: "127.0.0.1", family: 4 }]) as unknown as LookupFn;
11+
12+
const runningServers: Server[] = [];
13+
14+
async function startStreamServer(
15+
handler: (
16+
req: import("node:http").IncomingMessage,
17+
res: import("node:http").ServerResponse,
18+
) => void,
19+
): Promise<{ baseUrl: string; requests: string[] }> {
20+
const requests: string[] = [];
21+
const server = createServer((req, res) => {
22+
requests.push(`${req.method ?? "GET"} ${req.url ?? "/"}`);
23+
handler(req, res);
24+
});
25+
await new Promise<void>((resolve, reject) => {
26+
server.once("error", reject);
27+
server.listen(0, "127.0.0.1", () => resolve());
28+
});
29+
runningServers.push(server);
30+
const address = server.address() as AddressInfo;
31+
return {
32+
baseUrl: `http://127.0.0.1:${address.port}`,
33+
requests,
34+
};
35+
}
36+
37+
function connectTimeoutHandles(
38+
setTimeoutSpy: ReturnType<typeof vi.spyOn>,
39+
): ReturnType<typeof setTimeout>[] {
40+
return setTimeoutSpy.mock.results
41+
.filter((_result: { value: unknown }, index: number) => {
42+
return setTimeoutSpy.mock.calls[index]?.[1] === CONNECT_TIMEOUT_MS;
43+
})
44+
.map((result: { value: unknown }) => result.value as ReturnType<typeof setTimeout>);
45+
}
46+
47+
afterEach(async () => {
48+
vi.restoreAllMocks();
49+
await Promise.all(
50+
runningServers.splice(0).map(
51+
(server) =>
52+
new Promise<void>((resolve) => {
53+
server.close(() => resolve());
54+
}),
55+
),
56+
);
57+
});
58+
59+
describe("UrbitSSEClient openStream connect-timeout proof", () => {
60+
it("clears connect timers after a real non-OK stream response storm", async () => {
61+
const { baseUrl, requests } = await startStreamServer((_req, res) => {
62+
res.writeHead(503, { "Content-Type": "text/plain" });
63+
res.end("unavailable");
64+
});
65+
66+
const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
67+
const clearTimeoutSpy = vi.spyOn(globalThis, "clearTimeout");
68+
69+
const client = new UrbitSSEClient(baseUrl, "urbauth-~zod=proof", {
70+
autoReconnect: false,
71+
ship: "zod",
72+
ssrfPolicy: { allowPrivateNetwork: true },
73+
lookupFn: lookupLoopback,
74+
});
75+
76+
for (let attempt = 0; attempt < STORM_ATTEMPTS; attempt += 1) {
77+
await expect(client.openStream()).rejects.toThrow("Stream connection failed: 503");
78+
}
79+
80+
const armed = connectTimeoutHandles(setTimeoutSpy);
81+
expect(armed).toHaveLength(STORM_ATTEMPTS);
82+
for (const handle of armed) {
83+
expect(clearTimeoutSpy).toHaveBeenCalledWith(handle);
84+
}
85+
expect(requests.filter((entry) => entry.startsWith("GET /~/channel/"))).toHaveLength(
86+
STORM_ATTEMPTS,
87+
);
88+
});
89+
90+
it("clears the connect timer when production urbitFetch rejects", async () => {
91+
const setTimeoutSpy = vi.spyOn(globalThis, "setTimeout");
92+
const clearTimeoutSpy = vi.spyOn(globalThis, "clearTimeout");
93+
94+
// Controlled failure: accept then immediately destroy the socket so the
95+
// production urbitFetch / SSRF path rejects without depending on a free port.
96+
const server = createServer();
97+
server.on("connection", (socket) => {
98+
socket.destroy();
99+
});
100+
await new Promise<void>((resolve, reject) => {
101+
server.once("error", reject);
102+
server.listen(0, "127.0.0.1", () => resolve());
103+
});
104+
runningServers.push(server);
105+
const address = server.address() as AddressInfo;
106+
107+
const client = new UrbitSSEClient(`http://127.0.0.1:${address.port}`, "urbauth-~zod=proof", {
108+
autoReconnect: false,
109+
ship: "zod",
110+
ssrfPolicy: { allowPrivateNetwork: true },
111+
lookupFn: lookupLoopback,
112+
});
113+
114+
await expect(client.openStream()).rejects.toThrow();
115+
116+
const armed = connectTimeoutHandles(setTimeoutSpy);
117+
expect(armed).toHaveLength(1);
118+
expect(clearTimeoutSpy).toHaveBeenCalledWith(armed[0]);
119+
});
120+
});

extensions/tlon/src/urbit/sse-client.test.ts

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@ describe("UrbitSSEClient", () => {
3131
});
3232

3333
afterEach(() => {
34+
vi.useRealTimers();
3435
vi.restoreAllMocks();
3536
});
3637

@@ -119,6 +120,44 @@ describe("UrbitSSEClient", () => {
119120
});
120121
});
121122

123+
describe("openStream", () => {
124+
it("clears the connect timeout when urbitFetch rejects", async () => {
125+
vi.useFakeTimers();
126+
const clearTimeoutSpy = vi.spyOn(globalThis, "clearTimeout");
127+
const mockUrbitFetch = vi.mocked(urbitFetch);
128+
mockUrbitFetch.mockRejectedValueOnce(new Error("dns failed"));
129+
130+
const client = new UrbitSSEClient("https://example.com", "urbauth-~zod=123", {
131+
autoReconnect: false,
132+
});
133+
134+
await expect(client.openStream()).rejects.toThrow("dns failed");
135+
expect(clearTimeoutSpy).toHaveBeenCalled();
136+
expect(vi.getTimerCount()).toBe(0);
137+
});
138+
139+
it("clears the connect timeout when the stream response is not ok", async () => {
140+
vi.useFakeTimers();
141+
const clearTimeoutSpy = vi.spyOn(globalThis, "clearTimeout");
142+
const release = vi.fn().mockResolvedValue(undefined);
143+
const mockUrbitFetch = vi.mocked(urbitFetch);
144+
mockUrbitFetch.mockResolvedValueOnce({
145+
response: { ok: false, status: 503 } as unknown as Response,
146+
finalUrl: "https://example.com",
147+
release,
148+
});
149+
150+
const client = new UrbitSSEClient("https://example.com", "urbauth-~zod=123", {
151+
autoReconnect: false,
152+
});
153+
154+
await expect(client.openStream()).rejects.toThrow("Stream connection failed: 503");
155+
expect(release).toHaveBeenCalledOnce();
156+
expect(clearTimeoutSpy).toHaveBeenCalled();
157+
expect(vi.getTimerCount()).toBe(0);
158+
});
159+
});
160+
122161
describe("reconnection", () => {
123162
it("has autoReconnect enabled by default", () => {
124163
const client = new UrbitSSEClient("https://example.com", "urbauth-~zod=123");

extensions/tlon/src/urbit/sse-client.ts

Lines changed: 31 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -188,44 +188,43 @@ export class UrbitSSEClient {
188188

189189
this.streamController = controller;
190190

191-
const { response, release } = await urbitFetch({
192-
baseUrl: this.url,
193-
path: `/~/channel/${this.channelId}`,
194-
init: {
195-
method: "GET",
196-
headers: {
197-
Accept: "text/event-stream",
198-
Cookie: this.cookie,
191+
try {
192+
const { response, release } = await urbitFetch({
193+
baseUrl: this.url,
194+
path: `/~/channel/${this.channelId}`,
195+
init: {
196+
method: "GET",
197+
headers: {
198+
Accept: "text/event-stream",
199+
Cookie: this.cookie,
200+
},
199201
},
200-
},
201-
ssrfPolicy: this.ssrfPolicy,
202-
lookupFn: this.lookupFn,
203-
fetchImpl: this.fetchImpl,
204-
signal: controller.signal,
205-
auditContext: "tlon-urbit-sse-stream",
206-
});
207-
208-
this.streamRelease = release;
202+
ssrfPolicy: this.ssrfPolicy,
203+
lookupFn: this.lookupFn,
204+
fetchImpl: this.fetchImpl,
205+
signal: controller.signal,
206+
auditContext: "tlon-urbit-sse-stream",
207+
});
209208

210-
// Clear timeout once connection established (headers received).
211-
clearTimeout(timeoutId);
209+
this.streamRelease = release;
212210

213-
if (!response.ok) {
214-
await release();
215-
this.streamRelease = null;
216-
throw new Error(`Stream connection failed: ${response.status}`);
217-
}
211+
if (!response.ok) {
212+
await release();
213+
this.streamRelease = null;
214+
throw new Error(`Stream connection failed: ${response.status}`);
215+
}
218216

219-
this.processStream(response.body).catch((error: unknown) => {
220-
if (!this.aborted) {
221-
this.logger.error?.(`Stream error: ${String(error)}`);
222-
for (const { err } of this.eventHandlers.values()) {
223-
if (err) {
224-
err(error);
217+
this.processStream(response.body).catch((error: unknown) => {
218+
if (!this.aborted) {
219+
this.logger.error?.(`Stream error: ${String(error)}`);
220+
for (const { err } of this.eventHandlers.values()) {
221+
err?.(error);
225222
}
226223
}
227-
}
228-
});
224+
});
225+
} finally {
226+
clearTimeout(timeoutId); // success+failure: avoid dangling reconnect timers
227+
}
229228
}
230229

231230
async processStream(body: unknown) {

0 commit comments

Comments
 (0)