Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -427,8 +427,10 @@ OPENWIKI_OPENROUTER_PROVIDER_ONLY=Novita

### Provider retry attempts

OpenWiki uses LangChain's built-in retry handling for transient provider errors.
To override the number of retries after the first provider request, set `OPENWIKI_PROVIDER_RETRY_ATTEMPTS`:
OpenWiki uses provider retry handling for transient provider errors. For OpenAI
and OpenAI-compatible transports, this also covers retryable streaming HTTP
responses surfaced by the underlying SDK. To override the number of retries
after the first provider request, set `OPENWIKI_PROVIDER_RETRY_ATTEMPTS`:

```bash
OPENWIKI_PROVIDER_RETRY_ATTEMPTS=3
Expand Down
145 changes: 132 additions & 13 deletions src/agent/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import { ChatAnthropic } from "@langchain/anthropic";
import { ChatBedrockConverse } from "@langchain/aws";
import { ChatGoogle } from "@langchain/google/node";
import { SqliteSaver } from "@langchain/langgraph-checkpoint-sqlite";
import { ChatOpenAI } from "@langchain/openai";
import { ChatOpenAI, type ClientOptions } from "@langchain/openai";
import { ChatOpenRouter } from "@langchain/openrouter";
import type { Event as ProtocolEvent } from "@langchain/protocol";
import {
Expand Down Expand Up @@ -653,6 +653,7 @@ export function createModel(
providerRetryAttempts: number,
) {
const retryOptions = { maxRetries: providerRetryAttempts };
const fetchRetryOptions = { maxRetries: DISABLE_LANGCHAIN_FETCH_RETRIES };

if (provider === "gemini") {
return new ChatGoogle({
Expand Down Expand Up @@ -720,16 +721,16 @@ export function createModel(
// every generation — including the non-streaming `.invoke()` calls
// DeepAgents' agent node issues internally.
streaming: true,
...retryOptions,
configuration: {
...fetchRetryOptions,
configuration: withProviderRetryFetch(providerRetryAttempts, {
baseURL: CODEX_RESPONSES_BASE_URL,
defaultHeaders: {
"chatgpt-account-id": tokens.accountId,
originator: CODEX_ORIGINATOR,
"OpenAI-Beta": "responses=experimental",
},
fetch: createCodexFetch(modelId),
},
}),
});
}

Expand Down Expand Up @@ -758,17 +759,135 @@ export function createModel(

return new ChatOpenAI({
apiKey: getProviderApiKey(provider),
configuration: baseURL
? {
baseURL,
}
: undefined,
configuration: withProviderRetryFetch(
providerRetryAttempts,
baseURL
? {
baseURL,
}
: undefined,
),
model: modelId,
useResponsesApi: provider === "openai",
...retryOptions,
...fetchRetryOptions,
});
}

type ProviderFetch = typeof globalThis.fetch;

const HTTP_STATUS_REQUEST_TIMEOUT = 408;
const HTTP_STATUS_CONFLICT = 409;
const HTTP_STATUS_TOO_MANY_REQUESTS = 429;
const HTTP_STATUS_INTERNAL_SERVER_ERROR = 500;
const HTTP_STATUS_BAD_GATEWAY = 502;
const HTTP_STATUS_SERVICE_UNAVAILABLE = 503;
const HTTP_STATUS_GATEWAY_TIMEOUT = 504;
const RETRY_AFTER_HEADER_NAME = "retry-after";
const ABORT_ERROR_NAME = "AbortError";
const MILLISECONDS_PER_SECOND = 1000;
const DEFAULT_PROVIDER_RETRY_DELAY_MS = 1000;
const MAX_PROVIDER_RETRY_DELAY_MS = 60_000;
const DISABLE_LANGCHAIN_FETCH_RETRIES = 0;
const PROVIDER_RETRYABLE_STATUSES = new Set<number>([
HTTP_STATUS_REQUEST_TIMEOUT,
HTTP_STATUS_CONFLICT,
HTTP_STATUS_TOO_MANY_REQUESTS,
HTTP_STATUS_INTERNAL_SERVER_ERROR,
HTTP_STATUS_BAD_GATEWAY,
HTTP_STATUS_SERVICE_UNAVAILABLE,
HTTP_STATUS_GATEWAY_TIMEOUT,
]);

function withProviderRetryFetch(
providerRetryAttempts: number,
configuration: ClientOptions | undefined,
): ClientOptions {
return {
...configuration,
fetch: createProviderRetryFetch(
providerRetryAttempts,
configuration?.fetch,
),
};
}

function createProviderRetryFetch(
providerRetryAttempts: number,
fetchImpl: ProviderFetch = globalThis.fetch,
): ProviderFetch {
return async (input, init) => {
for (let attempt = 0; ; attempt += 1) {
let response: Response;

try {
response = await fetchImpl(cloneFetchInput(input), init);
} catch (error) {
if (
attempt >= providerRetryAttempts ||
isProviderRetryAbortError(error)
) {
throw error;
}

await sleep(DEFAULT_PROVIDER_RETRY_DELAY_MS);
continue;
}

if (
attempt >= providerRetryAttempts ||
!PROVIDER_RETRYABLE_STATUSES.has(response.status)
) {
return response;
}

await sleep(resolveProviderRetryDelayMs(response));
}
};
}

function cloneFetchInput(input: Parameters<ProviderFetch>[0]) {
return input instanceof Request ? input.clone() : input;
}

function isProviderRetryAbortError(error: unknown): boolean {
return (
(error instanceof DOMException || error instanceof Error) &&
error.name === ABORT_ERROR_NAME
);
}

function resolveProviderRetryDelayMs(response: Response): number {
const retryAfter = response.headers.get(RETRY_AFTER_HEADER_NAME);

if (!retryAfter) {
return DEFAULT_PROVIDER_RETRY_DELAY_MS;
}

const retryAfterSeconds = Number(retryAfter);

if (Number.isFinite(retryAfterSeconds) && retryAfterSeconds >= 0) {
return Math.min(
retryAfterSeconds * MILLISECONDS_PER_SECOND,
MAX_PROVIDER_RETRY_DELAY_MS,
);
}

const retryAtMs = Date.parse(retryAfter);

if (!Number.isFinite(retryAtMs)) {
return DEFAULT_PROVIDER_RETRY_DELAY_MS;
}

return Math.min(
Math.max(retryAtMs - Date.now(), 0),
MAX_PROVIDER_RETRY_DELAY_MS,
);
}

function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}

const CHATGPT_LOGIN_INCOMPLETE_MESSAGE =
"ChatGPT login is incomplete. Run `openwiki code --init` or `openwiki personal --init` to sign in with your ChatGPT account.";

Expand Down Expand Up @@ -871,12 +990,12 @@ function createGeminiEnterpriseModel(
// `apiKey` is a placeholder because that header is overwritten.
return new ChatOpenAI({
apiKey: VERTEX_ADC_PLACEHOLDER_KEY,
configuration: {
configuration: withProviderRetryFetch(retryOptions.maxRetries, {
baseURL: vertexOpenAIBaseUrl(projectId, location),
fetch: createVertexAuthFetch(),
},
}),
model: toVertexPublisherModel(modelId),
...retryOptions,
maxRetries: DISABLE_LANGCHAIN_FETCH_RETRIES,
});

default:
Expand Down
12 changes: 10 additions & 2 deletions test/gemini-retry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import { afterEach, beforeEach, describe, expect, test, vi } from "vitest";
const chatGoogleArgs: Array<Record<string, unknown>> = [];
const chatAnthropicArgs: Array<[string, Record<string, unknown>]> = [];
const chatOpenAIArgs: Array<Record<string, unknown>> = [];
const DISABLED_LANGCHAIN_RETRY_ATTEMPTS = 0;
const FUNCTION_TYPE = "function";

vi.mock("@langchain/google/node", () => ({
ChatGoogle: class {
Expand Down Expand Up @@ -92,15 +94,21 @@ describe("createModel wires retry attempts into the Gemini providers", () => {
expect(chatAnthropicArgs[0]?.[1].maxRetries).toBe(3);
});

test("gemini-enterprise MaaS surface passes maxRetries to ChatOpenAI", () => {
test("gemini-enterprise MaaS surface uses fetch retry without stacking LangChain retries", () => {
createModel(
"gemini-enterprise",
"publishers/meta/models/llama-3.3-70b-instruct-maas",
2,
);

expect(chatOpenAIArgs).toHaveLength(1);
expect(chatOpenAIArgs[0]?.maxRetries).toBe(2);
expect(chatOpenAIArgs[0]?.maxRetries).toBe(
DISABLED_LANGCHAIN_RETRY_ATTEMPTS,
);
expect(
(chatOpenAIArgs[0]?.configuration as { fetch?: unknown } | undefined)
?.fetch,
).toBeTypeOf(FUNCTION_TYPE);
});
});

Expand Down
110 changes: 110 additions & 0 deletions test/openai-provider-retry-fetch.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
import { afterEach, describe, expect, test, vi } from "vitest";
import { createModel } from "../src/agent/index.ts";

const OPENAI_API_KEY_ENV_KEY = "OPENAI_API_KEY";
const TEST_API_KEY = "test-openai-key";
const TEST_MODEL_ID = "gpt-5.6-terra";
const TEST_PROVIDER_RETRY_ATTEMPTS = 1;
const DISABLED_LANGCHAIN_RETRY_ATTEMPTS = 0;
const OPENAI_PROVIDER = "openai";
const RATE_LIMIT_STATUS = 429;
const SUCCESS_STATUS = 200;
const FUNCTION_TYPE = "function";
const RETRY_AFTER_HEADER_NAME = "retry-after";
const ZERO_RETRY_AFTER_SECONDS = "0";
const TEST_PROVIDER_URL = "https://api.example.test/responses";
const TRANSIENT_FETCH_ERROR_MESSAGE = "fetch failed";
const EXPECTED_PROVIDER_CALLS_AFTER_ONE_RETRY =
TEST_PROVIDER_RETRY_ATTEMPTS + 1;

describe("createModel OpenAI provider retry fetch", () => {
let savedOpenAiApiKey: string | undefined;
let savedFetch: typeof globalThis.fetch;

afterEach(() => {
restoreEnv(OPENAI_API_KEY_ENV_KEY, savedOpenAiApiKey);
globalThis.fetch = savedFetch;
vi.restoreAllMocks();
});

test("retries provider rate-limit responses through the OpenAI SDK fetch", async () => {
savedOpenAiApiKey = process.env[OPENAI_API_KEY_ENV_KEY];
savedFetch = globalThis.fetch;
process.env[OPENAI_API_KEY_ENV_KEY] = TEST_API_KEY;

const providerFetch = vi
.fn<typeof globalThis.fetch>()
.mockResolvedValueOnce(
new Response(null, {
status: RATE_LIMIT_STATUS,
headers: { [RETRY_AFTER_HEADER_NAME]: ZERO_RETRY_AFTER_SECONDS },
}),
)
.mockResolvedValueOnce(new Response(null, { status: SUCCESS_STATUS }));
globalThis.fetch = providerFetch;

const model = createModel(
OPENAI_PROVIDER,
TEST_MODEL_ID,
TEST_PROVIDER_RETRY_ATTEMPTS,
) as { clientConfig?: { fetch?: typeof globalThis.fetch } };

const retryFetch = model.clientConfig?.fetch;

expect(retryFetch).toBeTypeOf(FUNCTION_TYPE);

const response = await retryFetch?.(TEST_PROVIDER_URL);

expect(response?.status).toBe(SUCCESS_STATUS);
expect(providerFetch).toHaveBeenCalledTimes(
EXPECTED_PROVIDER_CALLS_AFTER_ONE_RETRY,
);
});

test("retries transient fetch rejections through the same retry budget", async () => {
savedOpenAiApiKey = process.env[OPENAI_API_KEY_ENV_KEY];
savedFetch = globalThis.fetch;
process.env[OPENAI_API_KEY_ENV_KEY] = TEST_API_KEY;

const providerFetch = vi
.fn<typeof globalThis.fetch>()
.mockRejectedValueOnce(new Error(TRANSIENT_FETCH_ERROR_MESSAGE))
.mockResolvedValueOnce(new Response(null, { status: SUCCESS_STATUS }));
globalThis.fetch = providerFetch;

const model = createModel(
OPENAI_PROVIDER,
TEST_MODEL_ID,
TEST_PROVIDER_RETRY_ATTEMPTS,
) as { clientConfig?: { fetch?: typeof globalThis.fetch } };

const response = await model.clientConfig?.fetch?.(TEST_PROVIDER_URL);

expect(response?.status).toBe(SUCCESS_STATUS);
expect(providerFetch).toHaveBeenCalledTimes(
EXPECTED_PROVIDER_CALLS_AFTER_ONE_RETRY,
);
});

test("uses the fetch wrapper as the single retry budget", () => {
savedOpenAiApiKey = process.env[OPENAI_API_KEY_ENV_KEY];
savedFetch = globalThis.fetch;
process.env[OPENAI_API_KEY_ENV_KEY] = TEST_API_KEY;

const model = createModel(
OPENAI_PROVIDER,
TEST_MODEL_ID,
TEST_PROVIDER_RETRY_ATTEMPTS,
) as { caller?: { maxRetries?: number } };

expect(model.caller?.maxRetries).toBe(DISABLED_LANGCHAIN_RETRY_ATTEMPTS);
});
});

function restoreEnv(key: string, value: string | undefined): void {
if (value === undefined) {
delete process.env[key];
} else {
process.env[key] = value;
}
}