Retry 429, 500, and connection error API requests to the trigger.dev server
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
---
|
||||
"@trigger.dev/core": patch
|
||||
---
|
||||
|
||||
Retry 429, 500, and connection error API requests to the trigger.dev server
|
||||
@@ -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
|
||||
);
|
||||
}
|
||||
|
||||
|
||||
@@ -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: {
|
||||
|
||||
Reference in New Issue
Block a user