Skip to content

Commit 7432f89

Browse files
committed
fix(core): time out API requests so a dead socket can't hang a call
API client requests had no timeout, so a request reusing a half-open keep-alive connection could hang until the runtime's socket default (minutes) rather than failing fast, and the retry logic never engaged because a hang produces no response to react to. Requests now run under a default 30s timeout (configurable per request via timeoutInMs); on timeout the request is aborted and the retry reconnects on a fresh socket. Caller cancellations are not retried.
1 parent 9f29eb2 commit 7432f89

3 files changed

Lines changed: 178 additions & 2 deletions

File tree

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,13 @@
1+
---
2+
"@trigger.dev/core": patch
3+
---
4+
5+
API requests now run under a request timeout (default 30s), so a half-open keep-alive connection can no longer hang a call until the runtime's socket default (minutes). On timeout the request aborts and retries on a fresh connection, and the retry is duplicate-safe via the request idempotency key. This mostly affects long-lived processes that reuse connections, including tasks that trigger other tasks.
6+
7+
Set the timeout per request, per client, or globally (most specific wins, `0` disables):
8+
9+
```ts
10+
tasks.trigger("my-task", payload, undefined, { timeoutInMs: 10_000 }); // per request
11+
new TriggerClient({ accessToken, requestOptions: { timeoutInMs: 10_000 } }); // per client
12+
// or globally: TRIGGER_API_REQUEST_TIMEOUT_MS=10000
13+
```
Lines changed: 131 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,131 @@
1+
import { afterEach, describe, expect, it, vi } from "vitest";
2+
import { z } from "zod";
3+
import { zodfetch } from "./core.js";
4+
5+
vi.setConfig({ testTimeout: 5_000 });
6+
7+
const schema = z.object({ ok: z.boolean() });
8+
9+
describe("zodfetch request timeout", () => {
10+
const originalFetch = globalThis.fetch;
11+
12+
afterEach(() => {
13+
globalThis.fetch = originalFetch;
14+
vi.restoreAllMocks();
15+
delete process.env.TRIGGER_API_REQUEST_TIMEOUT_MS;
16+
});
17+
18+
// Never settles unless the request is aborted (models a half-open keep-alive socket).
19+
function neverResponds(init?: RequestInit): Promise<Response> {
20+
return new Promise<Response>((_resolve, reject) => {
21+
init?.signal?.addEventListener("abort", () =>
22+
reject(init.signal!.reason ?? new DOMException("aborted", "AbortError"))
23+
);
24+
});
25+
}
26+
27+
it("aborts a hung request after the timeout instead of hanging forever", async () => {
28+
globalThis.fetch = vi.fn().mockImplementation((_url: string, init?: RequestInit) =>
29+
neverResponds(init)
30+
);
31+
32+
const start = Date.now();
33+
await expect(
34+
zodfetch(schema, "http://localhost/x", { method: "POST" }, {
35+
timeoutInMs: 150,
36+
retry: { maxAttempts: 1 },
37+
})
38+
).rejects.toThrow();
39+
expect(Date.now() - start).toBeLessThan(2_000);
40+
});
41+
42+
it("retries on a fresh connection after a timeout and recovers", async () => {
43+
let calls = 0;
44+
globalThis.fetch = vi.fn().mockImplementation((_url: string, init?: RequestInit) => {
45+
calls++;
46+
// First connection hangs; the retry's fresh connection responds.
47+
if (calls === 1) {
48+
return neverResponds(init);
49+
}
50+
return Promise.resolve(
51+
new Response(JSON.stringify({ ok: true }), {
52+
status: 200,
53+
headers: { "content-type": "application/json" },
54+
})
55+
);
56+
});
57+
58+
const result = await zodfetch(schema, "http://localhost/x", { method: "POST" }, {
59+
timeoutInMs: 150,
60+
retry: { maxAttempts: 3, minTimeoutInMs: 10, maxTimeoutInMs: 50, factor: 1, randomize: false },
61+
});
62+
63+
expect(result).toEqual({ ok: true });
64+
expect(calls).toBe(2);
65+
});
66+
67+
it("reuses the same request idempotency key across a timeout retry so the server can dedupe", async () => {
68+
const keys: Array<string | null> = [];
69+
let calls = 0;
70+
globalThis.fetch = vi.fn().mockImplementation((_url: string, init?: RequestInit) => {
71+
calls++;
72+
keys.push(new Headers(init?.headers).get("x-trigger-request-idempotency-key"));
73+
if (calls === 1) {
74+
return neverResponds(init);
75+
}
76+
return Promise.resolve(
77+
new Response(JSON.stringify({ ok: true }), {
78+
status: 200,
79+
headers: { "content-type": "application/json" },
80+
})
81+
);
82+
});
83+
84+
await zodfetch(schema, "http://localhost/x", { method: "POST" }, {
85+
timeoutInMs: 150,
86+
retry: { maxAttempts: 3, minTimeoutInMs: 10, maxTimeoutInMs: 50, factor: 1, randomize: false },
87+
});
88+
89+
expect(keys).toHaveLength(2);
90+
expect(keys[0]).toBeTruthy();
91+
expect(keys[0]).toBe(keys[1]);
92+
});
93+
94+
it("does not time out when timeoutInMs is 0 (for long-lived requests)", async () => {
95+
globalThis.fetch = vi.fn().mockImplementation(
96+
() =>
97+
new Promise<Response>((resolve) =>
98+
setTimeout(
99+
() =>
100+
resolve(
101+
new Response(JSON.stringify({ ok: true }), {
102+
status: 200,
103+
headers: { "content-type": "application/json" },
104+
})
105+
),
106+
300
107+
)
108+
)
109+
);
110+
111+
const result = await zodfetch(schema, "http://localhost/x", { method: "POST" }, {
112+
timeoutInMs: 0,
113+
retry: { maxAttempts: 1 },
114+
});
115+
116+
expect(result).toEqual({ ok: true });
117+
});
118+
119+
it("uses TRIGGER_API_REQUEST_TIMEOUT_MS as the default when no timeout is passed", async () => {
120+
process.env.TRIGGER_API_REQUEST_TIMEOUT_MS = "120";
121+
globalThis.fetch = vi.fn().mockImplementation((_url: string, init?: RequestInit) =>
122+
neverResponds(init)
123+
);
124+
125+
const start = Date.now();
126+
await expect(
127+
zodfetch(schema, "http://localhost/x", { method: "POST" }, { retry: { maxAttempts: 1 } })
128+
).rejects.toThrow();
129+
expect(Date.now() - start).toBeLessThan(2_000);
130+
});
131+
});

‎packages/core/src/v3/apiClient/core.ts‎

Lines changed: 34 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import {
1919
} from "./pagination.js";
2020
import { EventSource, type ErrorEvent } from "eventsource";
2121
import { randomUUID } from "../utils/crypto.js";
22+
import { getNumberEnvVar } from "../utils/getEnv.js";
2223

2324
export const defaultRetryOptions = {
2425
maxAttempts: 3,
@@ -28,8 +29,11 @@ export const defaultRetryOptions = {
2829
randomize: false,
2930
} satisfies RetryOptions;
3031

32+
export const defaultRequestTimeoutInMs = 30_000;
33+
3134
export type ZodFetchOptions<TData = any> = {
3235
retry?: RetryOptions;
36+
timeoutInMs?: number;
3337
tracer?: TriggerTracer;
3438
name?: string;
3539
attributes?: Attributes;
@@ -40,7 +44,7 @@ export type ZodFetchOptions<TData = any> = {
4044

4145
export type AnyZodFetchOptions = ZodFetchOptions<any>;
4246

43-
export type ApiRequestOptions = Pick<ZodFetchOptions, "retry"> & {
47+
export type ApiRequestOptions = Pick<ZodFetchOptions, "retry" | "timeoutInMs"> & {
4448
additionalHeaders?: Record<string, string>;
4549
};
4650

@@ -51,6 +55,7 @@ type KeysEnum<T> = { [P in keyof Required<T>]: true };
5155
// compiler such that any missing / extraneous keys will cause an error.
5256
const requestOptionsKeys: KeysEnum<ApiRequestOptions> = {
5357
retry: true,
58+
timeoutInMs: true,
5459
additionalHeaders: true,
5560
};
5661

@@ -230,9 +235,28 @@ async function _doZodFetchWithRetries<TResponseBodySchema extends z.ZodTypeAny>(
230235
options?: ZodFetchOptions,
231236
attempt = 1
232237
): Promise<ZodFetchResult<z.output<TResponseBodySchema>>> {
238+
// Precedence: per-request option > per-client requestOptions > env var > built-in default.
239+
const timeoutInMs =
240+
options?.timeoutInMs ?? getNumberEnvVar("TRIGGER_API_REQUEST_TIMEOUT_MS") ?? defaultRequestTimeoutInMs;
241+
const controller = new AbortController();
242+
const callerSignal = requestInit?.signal ?? undefined;
243+
const onCallerAbort = () => controller.abort(callerSignal?.reason);
244+
if (callerSignal) {
245+
if (callerSignal.aborted) controller.abort(callerSignal.reason);
246+
else callerSignal.addEventListener("abort", onCallerAbort, { once: true });
247+
}
248+
// timeoutInMs <= 0 disables the timeout, for long-lived requests that set their own bound.
249+
const timeout =
250+
timeoutInMs > 0
251+
? setTimeout(
252+
() => controller.abort(new Error(`Request to ${url} timed out after ${timeoutInMs}ms`)),
253+
timeoutInMs
254+
)
255+
: undefined;
256+
233257
try {
234258
const response = await context.with(suppressTracing(context.active()), () =>
235-
fetch(url, requestInitWithCache(requestInit))
259+
fetch(url, { ...requestInitWithCache(requestInit), signal: controller.signal })
236260
);
237261

238262
const responseHeaders = createResponseHeaders(response.headers);
@@ -277,6 +301,11 @@ async function _doZodFetchWithRetries<TResponseBodySchema extends z.ZodTypeAny>(
277301
if (error instanceof ValidationError) {
278302
}
279303

304+
// A caller-supplied signal aborting is a cancellation, not a transient failure.
305+
if (callerSignal?.aborted) {
306+
throw error;
307+
}
308+
280309
if (options?.retry) {
281310
const retry = { ...defaultRetryOptions, ...options.retry };
282311

@@ -290,6 +319,9 @@ async function _doZodFetchWithRetries<TResponseBodySchema extends z.ZodTypeAny>(
290319
}
291320

292321
throw new ApiConnectionError({ cause: castToError(error) });
322+
} finally {
323+
clearTimeout(timeout);
324+
callerSignal?.removeEventListener("abort", onCallerAbort);
293325
}
294326
}
295327

0 commit comments

Comments
 (0)