From 9e2f2ac5878e5220d8c8e57fe3560603165b1979 Mon Sep 17 00:00:00 2001 From: 404-Page-Found <139850808+404-Page-Found@users.noreply.github.com> Date: Wed, 23 Sep 2026 07:35:02 +1000 Subject: [PATCH] fix(providers): preserve bounded non-streaming response bodies --- src/providers/anthropic.ts | 25 ++- src/providers/cohere.ts | 34 ++- src/providers/openai-compatible.ts | 40 +++- src/providers/request.ts | 70 +++++- .../providers/provider-body-timeout.test.mjs | 210 ++++++++++++++++++ 5 files changed, 357 insertions(+), 22 deletions(-) create mode 100644 tests/providers/provider-body-timeout.test.mjs diff --git a/src/providers/anthropic.ts b/src/providers/anthropic.ts index 1a1db23..1e6d343 100644 --- a/src/providers/anthropic.ts +++ b/src/providers/anthropic.ts @@ -1,5 +1,11 @@ import type { ChatParams, ChatResult, Provider, ProviderStreamChunk } from '../types.js'; -import { DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, fetchWithTimeout } from './request.js'; +import { + DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, + ProviderResponseTimeoutError, + fetchWithTimeout, + readResponseJsonWithTimeout, + readResponseTextWithTimeout, +} from './request.js'; import { parseAnthropicSseLine, streamSseResponse } from './sse.js'; function buildAnthropicRequestBody(params: ChatParams, options: { stream?: boolean } = {}): Record { @@ -32,6 +38,7 @@ function buildAnthropicRequestBody(params: ChatParams, options: { stream?: boole export class AnthropicProvider implements Provider { async complete(params: ChatParams): Promise { const { model, apiKey, baseUrl } = params; + const controller = new AbortController(); const url = `${baseUrl.replace(/\/+$/, '')}/messages`; const body = buildAnthropicRequestBody(params); @@ -49,19 +56,27 @@ export class AnthropicProvider implements Provider { }, 'Anthropic API request', DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, - new AbortController(), + controller, params.signal, + false, ); if (!response.ok) { - const errorBody = await response.text().catch(() => ''); + let errorBody = ''; + try { + errorBody = await readResponseTextWithTimeout(response, controller, 'Anthropic API response'); + } catch (error) { + if (error instanceof ProviderResponseTimeoutError) { + throw error; + } + } throw new Error(`Anthropic API error (${response.status}): ${errorBody || response.statusText}`); } - const data = (await response.json()) as { + const data = await readResponseJsonWithTimeout<{ content?: { type: string; text: string }[]; model?: string; - }; + }>(response, controller, 'Anthropic API response'); const textContent = data.content?.find((c) => c.type === 'text'); if (!textContent?.text) { diff --git a/src/providers/cohere.ts b/src/providers/cohere.ts index da14d97..3e9ee1c 100644 --- a/src/providers/cohere.ts +++ b/src/providers/cohere.ts @@ -1,9 +1,16 @@ import type { ChatParams, ChatResult, Provider } from '../types.js'; -import { DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, fetchWithTimeout } from './request.js'; +import { + DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, + ProviderResponseTimeoutError, + fetchWithTimeout, + readResponseJsonWithTimeout, + readResponseTextWithTimeout, +} from './request.js'; export class CohereProvider implements Provider { async complete(params: ChatParams): Promise { const { model, messages, temperature = 0.7, maxTokens = 1024, apiKey, baseUrl } = params; + const controller = new AbortController(); const url = `${baseUrl.replace(/\/+$/, '')}/chat`; @@ -41,19 +48,27 @@ export class CohereProvider implements Provider { }, 'Cohere API request', DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, - new AbortController(), + controller, params.signal, + false, ); if (!response.ok) { - const errorBody = await response.text().catch(() => ''); + let errorBody = ''; + try { + errorBody = await readResponseTextWithTimeout(response, controller, 'Cohere API response'); + } catch (error) { + if (error instanceof ProviderResponseTimeoutError) { + throw error; + } + } throw new Error(`Cohere API error (${response.status}): ${errorBody || response.statusText}`); } - const data = (await response.json()) as { + const data = await readResponseJsonWithTimeout<{ text?: string; meta?: { api_version?: { version?: string } }; - }; + }>(response, controller, 'Cohere API response'); if (!data.text) { throw new Error('Cohere returned empty response.'); @@ -66,6 +81,7 @@ export class CohereProvider implements Provider { } async fetchModels(baseUrl: string, apiKey: string): Promise { + const controller = new AbortController(); const url = `${baseUrl.replace(/\/+$/, '')}/models`; const response = await fetchWithTimeout( @@ -76,15 +92,19 @@ export class CohereProvider implements Provider { }, }, 'Cohere model request', + DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, + controller, + undefined, + false, ); if (!response.ok) { throw new Error(`Failed to fetch models (${response.status}): ${response.statusText}`); } - const data = (await response.json()) as { + const data = await readResponseJsonWithTimeout<{ models?: { name?: string; id?: string }[]; - }; + }>(response, controller, 'Cohere model response'); if (data.models && Array.isArray(data.models)) { return data.models diff --git a/src/providers/openai-compatible.ts b/src/providers/openai-compatible.ts index 84e0ad2..af9fd3a 100644 --- a/src/providers/openai-compatible.ts +++ b/src/providers/openai-compatible.ts @@ -1,5 +1,11 @@ import type { ChatParams, ChatResult, Provider, ProviderStreamChunk } from '../types.js'; -import { DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, fetchWithTimeout } from './request.js'; +import { + DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, + ProviderResponseTimeoutError, + fetchWithTimeout, + readResponseJsonWithTimeout, + readResponseTextWithTimeout, +} from './request.js'; import { parseOpenAiSseLine, streamSseResponse, SSE_STREAM_END } from './sse.js'; const MAX_BUFFERED_REASONING_CHARS = 1024 * 1024; @@ -25,6 +31,7 @@ function buildOpenAiRequestBody(params: ChatParams, options: { stream?: boolean export class OpenAICompatibleProvider implements Provider { async complete(params: ChatParams): Promise { const { model, apiKey, baseUrl } = params; + const controller = new AbortController(); const url = `${baseUrl.replace(/\/+$/, '')}/chat/completions`; @@ -44,19 +51,27 @@ export class OpenAICompatibleProvider implements Provider { }, 'OpenAI-compatible API request', DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, - new AbortController(), + controller, params.signal, + false, ); if (!response.ok) { - const errorBody = await response.text().catch(() => ''); + let errorBody = ''; + try { + errorBody = await readResponseTextWithTimeout(response, controller, 'OpenAI-compatible API response'); + } catch (error) { + if (error instanceof ProviderResponseTimeoutError) { + throw error; + } + } throw new Error(`OpenAI-compatible API error (${response.status}): ${errorBody || response.statusText}`); } - const data = (await response.json()) as { + const data = await readResponseJsonWithTimeout<{ choices?: { message?: { content?: string } }[]; model?: string; - }; + }>(response, controller, 'OpenAI-compatible API response'); const content = data.choices?.[0]?.message?.content; if (!content) { @@ -140,6 +155,7 @@ export class OpenAICompatibleProvider implements Provider { } async fetchModels(baseUrl: string, apiKey: string): Promise { + const controller = new AbortController(); const url = `${baseUrl.replace(/\/+$/, '')}/models`; const headers: Record = { @@ -149,15 +165,23 @@ export class OpenAICompatibleProvider implements Provider { headers['Authorization'] = `Bearer ${apiKey}`; } - const response = await fetchWithTimeout(url, { headers }, 'OpenAI-compatible model request'); + const response = await fetchWithTimeout( + url, + { headers }, + 'OpenAI-compatible model request', + DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, + controller, + undefined, + false, + ); if (!response.ok) { throw new Error(`Failed to fetch models (${response.status}): ${response.statusText}`); } - const data = (await response.json()) as { + const data = await readResponseJsonWithTimeout<{ data?: { id: string; object?: string }[]; - }; + }>(response, controller, 'OpenAI-compatible model response'); if (!data.data || !Array.isArray(data.data)) { throw new Error('Unexpected response format when fetching models'); diff --git a/src/providers/request.ts b/src/providers/request.ts index 433c4ba..1fa9002 100644 --- a/src/providers/request.ts +++ b/src/providers/request.ts @@ -1,5 +1,11 @@ export const DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS = 30_000; +export class ProviderResponseTimeoutError extends Error { + constructor(label: string, timeoutMs: number) { + super(`${label} timed out after ${timeoutMs}ms`); + this.name = 'ProviderResponseTimeoutError'; + } +} export async function fetchWithTimeout( url: string, init: RequestInit, @@ -10,7 +16,7 @@ export async function fetchWithTimeout( keepTimeoutThroughBody = true, ): Promise { let timedOut = false; - let timeoutError: Error | undefined; + let timeoutError: ProviderResponseTimeoutError | undefined; let externalAbortListener: (() => void) | undefined; if (externalSignal) { @@ -37,7 +43,7 @@ export async function fetchWithTimeout( timeout = setTimeout(() => { timedOut = true; - timeoutError = new Error(`${label} timed out after ${timeoutMs}ms`); + timeoutError = new ProviderResponseTimeoutError(label, timeoutMs); controller.abort(timeoutError); }, timeoutMs); @@ -104,3 +110,63 @@ export async function fetchWithTimeout( throw error; } } + +export async function readResponseTextWithTimeout( + response: Response, + controller: AbortController, + label: string, + timeoutMs = DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, +): Promise { + const reader = response.body?.getReader(); + if (!reader) { + return ''; + } + + let timedOut = false; + const timeout = setTimeout(() => { + timedOut = true; + controller.abort(); + void reader.cancel().catch(() => undefined); + }, timeoutMs); + + try { + const decoder = new TextDecoder(); + let text = ''; + + while (true) { + const { done, value } = await reader.read(); + if (done) break; + if (value) { + text += decoder.decode(value, { stream: true }); + } + if (timedOut) { + throw new ProviderResponseTimeoutError(label, timeoutMs); + } + } + + if (timedOut) { + throw new ProviderResponseTimeoutError(label, timeoutMs); + } + + text += decoder.decode(); + return text; + } catch (error) { + if (timedOut) { + throw new ProviderResponseTimeoutError(label, timeoutMs); + } + throw error; + } finally { + clearTimeout(timeout); + reader.releaseLock(); + } +} + +export async function readResponseJsonWithTimeout( + response: Response, + controller: AbortController, + label: string, + timeoutMs = DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS, +): Promise { + const text = await readResponseTextWithTimeout(response, controller, label, timeoutMs); + return JSON.parse(text) as T; +} diff --git a/tests/providers/provider-body-timeout.test.mjs b/tests/providers/provider-body-timeout.test.mjs new file mode 100644 index 0000000..b08c008 --- /dev/null +++ b/tests/providers/provider-body-timeout.test.mjs @@ -0,0 +1,210 @@ +import assert from 'node:assert/strict'; +import test from 'node:test'; + +import { AnthropicProvider } from '../../dist/providers/anthropic.js'; +import { CohereProvider } from '../../dist/providers/cohere.js'; +import { OpenAICompatibleProvider } from '../../dist/providers/openai-compatible.js'; +import { DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS } from '../../dist/providers/request.js'; + +const originalFetch = globalThis.fetch; +const originalSetTimeout = globalThis.setTimeout; + +function timeoutPattern(label) { + return new RegExp(`${label} timed out after ${DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS}ms`); +} + +function setupStalledResponse(t, status = 200) { + let cancelled = false; + let aborted = false; + + globalThis.fetch = async (_url, init = {}) => + new Response( + new ReadableStream({ + start(stream) { + init.signal?.addEventListener( + 'abort', + () => { + aborted = true; + }, + { once: true }, + ); + // Simulate headers plus a partial JSON body with no completion. + stream.enqueue(new TextEncoder().encode('{"partial":')); + }, + cancel() { + cancelled = true; + }, + }), + { status, headers: { 'Content-Type': 'application/json' } }, + ); + + t.after(() => { + globalThis.fetch = originalFetch; + }); + + return () => ({ cancelled, aborted }); +} + +async function assertBodyTimeout(t, promise, pattern) { + await new Promise((resolve) => originalSetTimeout(resolve, 0)); + t.mock.timers.tick(DEFAULT_PROVIDER_REQUEST_TIMEOUT_MS); + await assert.rejects(promise, pattern); +} + +test('OpenAI-compatible complete times out when the JSON body stalls after headers arrive', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t); + const provider = new OpenAICompatibleProvider(); + + const pending = provider.complete({ + model: 'gpt-test', + baseUrl: 'https://openai.example.com/v1', + apiKey: 'test-key', + messages: [{ role: 'user', content: 'hello' }], + }); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('OpenAI-compatible API response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('OpenAI-compatible complete preserves response-body timeout errors', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t, 500); + const provider = new OpenAICompatibleProvider(); + + const pending = provider.complete({ + model: 'gpt-test', + baseUrl: 'https://openai.example.com/v1', + apiKey: 'test-key', + messages: [{ role: 'user', content: 'hello' }], + }); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('OpenAI-compatible API response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('OpenAI-compatible fetchModels times out when the JSON body stalls after headers arrive', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t); + const provider = new OpenAICompatibleProvider(); + + const pending = provider.fetchModels('https://openai.example.com/v1', 'test-key'); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('OpenAI-compatible model response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('Anthropic complete times out when the JSON body stalls after headers arrive', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t); + const provider = new AnthropicProvider(); + + const pending = provider.complete({ + model: 'claude-test', + baseUrl: 'https://anthropic.example.com', + apiKey: 'test-key', + messages: [{ role: 'user', content: 'hello' }], + }); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('Anthropic API response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('Anthropic complete preserves response-body timeout errors', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t, 500); + const provider = new AnthropicProvider(); + + const pending = provider.complete({ + model: 'claude-test', + baseUrl: 'https://anthropic.example.com', + apiKey: 'test-key', + messages: [{ role: 'user', content: 'hello' }], + }); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('Anthropic API response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('Cohere complete times out when the JSON body stalls after headers arrive', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t); + const provider = new CohereProvider(); + + const pending = provider.complete({ + model: 'command-r', + baseUrl: 'https://cohere.example.com', + apiKey: 'test-key', + messages: [{ role: 'user', content: 'hello' }], + }); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('Cohere API response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('Cohere complete preserves response-body timeout errors', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t, 500); + const provider = new CohereProvider(); + + const pending = provider.complete({ + model: 'command-r', + baseUrl: 'https://cohere.example.com', + apiKey: 'test-key', + messages: [{ role: 'user', content: 'hello' }], + }); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('Cohere API response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +}); + +test('Cohere fetchModels times out when the JSON body stalls after headers arrive', async (t) => { + t.mock.timers.enable({ apis: ['setTimeout'] }); + const state = setupStalledResponse(t); + const provider = new CohereProvider(); + + const pending = provider.fetchModels('https://cohere.example.com', 'test-key'); + + await assertBodyTimeout( + t, + pending, + timeoutPattern('Cohere model response'), + ); + + assert.deepEqual(state(), { cancelled: true, aborted: true }); +});