Retry 429, 500, and connection error API requests to the trigger.dev server

This commit is contained in:
Eric Allam
2024-03-31 16:06:12 +01:00
parent a946797d95
commit 9af2570da6
4 changed files with 103 additions and 17 deletions
+5
View File
@@ -0,0 +1,5 @@
---
"@trigger.dev/core": patch
---
Retry 429, 500, and connection error API requests to the trigger.dev server
+35 -13
View File
@@ -1,5 +1,5 @@
import { context, propagation } from "@opentelemetry/api";
import { zodfetch } from "../../zodfetch";
import { ZodFetchOptions, zodfetch } from "../../zodfetch";
import { taskContextManager } from "../tasks/taskContextManager";
import { SafeAsyncLocalStorage } from "../utils/safeAsyncLocalStorage";
import { getEnvVar } from "../utils/getEnv";
@@ -15,6 +15,16 @@ export type TriggerOptions = {
spanParentAsLink?: boolean;
};
const zodFetchOptions: ZodFetchOptions = {
retry: {
maxAttempts: 5,
minTimeoutInMs: 1000,
maxTimeoutInMs: 30_000,
factor: 2,
randomize: false,
},
};
/**
* Trigger.dev v3 API client
*/
@@ -29,19 +39,29 @@ export class ApiClient {
}
triggerTask(taskId: string, body: TriggerTaskRequestBody, options?: TriggerOptions) {
return zodfetch(TriggerTaskResponse, `${this.baseUrl}/api/v1/tasks/${taskId}/trigger`, {
method: "POST",
headers: this.#getHeaders(options?.spanParentAsLink ?? false),
body: JSON.stringify(body),
});
return zodfetch(
TriggerTaskResponse,
`${this.baseUrl}/api/v1/tasks/${taskId}/trigger`,
{
method: "POST",
headers: this.#getHeaders(options?.spanParentAsLink ?? false),
body: JSON.stringify(body),
},
zodFetchOptions
);
}
batchTriggerTask(taskId: string, body: BatchTriggerTaskRequestBody, options?: TriggerOptions) {
return zodfetch(BatchTriggerTaskResponse, `${this.baseUrl}/api/v1/tasks/${taskId}/batch`, {
method: "POST",
headers: this.#getHeaders(options?.spanParentAsLink ?? false),
body: JSON.stringify(body),
});
return zodfetch(
BatchTriggerTaskResponse,
`${this.baseUrl}/api/v1/tasks/${taskId}/batch`,
{
method: "POST",
headers: this.#getHeaders(options?.spanParentAsLink ?? false),
body: JSON.stringify(body),
},
zodFetchOptions
);
}
createUploadPayloadUrl(filename: string) {
@@ -51,7 +71,8 @@ export class ApiClient {
{
method: "PUT",
headers: this.#getHeaders(false),
}
},
zodFetchOptions
);
}
@@ -62,7 +83,8 @@ export class ApiClient {
{
method: "GET",
headers: this.#getHeaders(false),
}
},
zodFetchOptions
);
}
+61 -4
View File
@@ -1,17 +1,32 @@
import { z } from "zod";
import { context, propagation } from "@opentelemetry/api";
import { RetryOptions, calculateNextRetryDelay, defaultRetryOptions } from "./v3";
type ApiResult<TSuccessResult> =
export type ApiResult<TSuccessResult> =
| { ok: true; data: TSuccessResult }
| {
ok: false;
error: string;
};
export type ZodFetchOptions = {
retry?: RetryOptions;
};
export async function zodfetch<TResponseBody extends any>(
schema: z.Schema<TResponseBody>,
url: string,
requestInit?: RequestInit
requestInit?: RequestInit,
options?: ZodFetchOptions
): Promise<ApiResult<TResponseBody>> {
return await _doZodFetch(schema, url, requestInit, options);
}
async function _doZodFetch<TResponseBody extends any>(
schema: z.Schema<TResponseBody>,
url: string,
requestInit?: RequestInit,
options?: ZodFetchOptions,
attempt = 1
): Promise<ApiResult<TResponseBody>> {
try {
const response = await fetch(url, requestInit);
@@ -23,7 +38,7 @@ export async function zodfetch<TResponseBody extends any>(
};
}
if (response.status >= 400 && response.status < 500) {
if (response.status >= 400 && response.status < 500 && response.status !== 429) {
const body = await response.json();
if (!body.error) {
return { ok: false, error: "Something went wrong" };
@@ -32,6 +47,31 @@ export async function zodfetch<TResponseBody extends any>(
return { ok: false, error: body.error };
}
// Retryable errors
if (response.status === 429 || response.status >= 500) {
if (!options?.retry) {
return {
ok: false,
error: `Failed to fetch ${url}, got status code ${response.status}`,
};
}
const retry = { ...defaultRetryOptions, ...options.retry };
if (attempt > retry.maxAttempts) {
return {
ok: false,
error: `Failed to fetch ${url}, got status code ${response.status}`,
};
}
const delay = calculateNextRetryDelay(retry, attempt);
await new Promise((resolve) => setTimeout(resolve, delay));
return await _doZodFetch(schema, url, requestInit, options, attempt + 1);
}
if (response.status !== 200) {
return {
ok: false,
@@ -55,6 +95,23 @@ export async function zodfetch<TResponseBody extends any>(
return { ok: false, error: parsedResult.error.message };
} catch (error) {
if (options?.retry) {
const retry = { ...defaultRetryOptions, ...options.retry };
if (attempt > retry.maxAttempts) {
return {
ok: false,
error: error instanceof Error ? error.message : JSON.stringify(error),
};
}
const delay = calculateNextRetryDelay(retry, attempt);
await new Promise((resolve) => setTimeout(resolve, delay));
return await _doZodFetch(schema, url, requestInit, options, attempt + 1);
}
return {
ok: false,
error: error instanceof Error ? error.message : JSON.stringify(error),
@@ -21,6 +21,8 @@ export const testConcurrency = task({
run: async ({ count = 10, delay = 5000 }: { count: number; delay: number }) => {
logger.info(`Running ${count} tasks`);
await new Promise((resolve) => setTimeout(resolve, 3000));
await testConcurrencyChild.batchTrigger({
items: Array.from({ length: count }).map((_, index) => ({
payload: {