feat: implement bounded LLM request recovery

This commit is contained in:
Tianyi Cui
2026-07-20 03:34:19 +08:00
parent 7cf966fc0e
commit 3b0b0cefeb
115 changed files with 3311 additions and 366 deletions

View File

@@ -16,6 +16,7 @@ The package root exposes the Cordis plugin contract and `DeepSeekAdapter`; wire
baseURL: !!js process.env.DEEPSEEK_BASE_URL # default: https://api.deepseek.com
thinking: enabled # optional; provider default is enabled
reasoningEffort: high # optional; high | max — omitted ⇒ not sent
streamIdleTimeoutMs: 300000 # optional; positive finite Node timer delay; five-minute default
models: # optional; defaults to V4 Flash and V4 Pro
- id: deepseek-v4-flash
name: DeepSeek V4 Flash
@@ -29,6 +30,8 @@ The plugin registers the single provider route `deepseek`. A request selects it
`thinking`/`reasoningEffort` are adapter-level request defaults serialized as the official top-level `thinking: {type}` / `reasoning_effort` wire fields. They live in adapter config (not `GenerateOptions`) to keep the core vocabulary provider-neutral.
`streamIdleTimeoutMs` bounds each outstanding provider read, including the initial `fetch`, without counting time the consumer spends between chunks. One stable abort signal reaches the request and body reader for the whole call; expiry stops the transport and throws `LlmError('TIMEOUT')`, while an earlier caller abort throws `LlmError('ABORTED')`. The adapter makes exactly one provider request per `stream()` call; agent-level retry is a separate plugin policy.
## App attribution
Every request carries the shared attribution header from dsh-llm's `attributionHeaders()` - the mandatory `User-Agent` baseline identifying the harness (see [dsh-llm § App attribution](../llm/README.md#app-attribution-attributionts)). Direct DeepSeek requests and OpenAI-compatible gateway requests get no provider-specific app-attribution headers under this adapter contract; OpenRouter app attribution is deferred to a future explicit OpenRouter adapter or mode.
@@ -42,11 +45,11 @@ Every request carries the shared attribution header from dsh-llm's `attributionH
## Errors
Non-2xx responses throw `LlmError` with stable codes: `AUTH` (401/403), `RATE_LIMIT` (429), `CONTEXT_WINDOW_EXCEEDED` (a 400 whose provider code, type, or message identifies context overflow), `INVALID_REQUEST` (other 400s), `SERVER` (5xx), `HTTP_<status>` otherwise. Protocol violations throw `STREAM_CLOSED` (no `[DONE]`) or `MALFORMED_RESPONSE` (bad JSON payload). Unknown wire `finish_reason`s (e.g. `content_filter`, `insufficient_system_resource`) become `finish {kind: 'error', code: <REASON>}` chunks.
Non-2xx responses throw `LlmError` with stable codes: `AUTH` (401/403), `QUOTA` (a response whose provider details identify exhausted quota, balance, or credits), `RATE_LIMIT` (other 429s), `CONTEXT_WINDOW_EXCEEDED` (a 400 whose provider code, type, or message identifies context overflow), `INVALID_REQUEST` (other 400s), `SERVER` (5xx), `HTTP_<status>` otherwise. Its serializable `failure` retains the HTTP status plus a valid positive `Retry-After` seconds/date delay and `x-request-id` / `x-deepseek-request-id` when present. Connection failures are `TRANSPORT`; protocol violations throw `STREAM_CLOSED` (no `[DONE]`) or `MALFORMED_RESPONSE` (bad JSON payload). Unknown wire `finish_reason`s (e.g. `content_filter`, `insufficient_system_resource`) become `finish {kind: 'error', failure}` chunks.
## Testing
Unit suites run against a local `node:http` mock SSE server (no network). Real-API coverage lives in `tests/adapter.e2e.ts` (`pnpm run test:e2e`, key-gated): V4 Flash + V4 Pro across thinking enabled/disabled and both official effort levels, including the thinking+tools round trip with reasoning passback.
Unit suites run against a local `node:http` mock SSE server (no network), including structured HTTP facts, malformed/truncated streams, caller abort, connection failure, and proof that idle timeout aborts the actual body. Real-API coverage lives in `tests/adapter.e2e.ts` (`pnpm run test:e2e`, key-gated): V4 Flash + V4 Pro across thinking enabled/disabled and both official effort levels, including the thinking+tools round trip with reasoning passback.
## Model Experience

View File

@@ -23,6 +23,7 @@
"license": "BSD-3-Clause",
"peerDependencies": {
"@deepseek-ai/dsh-llm": "^0.0.1",
"@deepseek-ai/dsh-timeout": "^0.0.1",
"cordis": "^4.0.0-rc.7"
},
"dependencies": {
@@ -30,6 +31,7 @@
},
"devDependencies": {
"@deepseek-ai/dsh-llm": "workspace:^",
"@deepseek-ai/dsh-timeout": "workspace:^",
"cordis": "^4.0.0-rc.7"
}
}

View File

@@ -5,8 +5,9 @@
* @module dsh-llm-deepseek/adapter
*/
import { attributionHeaders, CONTEXT_WINDOW_EXCEEDED_CODE, isContextWindowExceededError, LlmAdapter, LlmError } from '@deepseek-ai/dsh-llm'
import { attributionHeaders, CONTEXT_WINDOW_EXCEEDED_CODE, isContextWindowExceededError, isQuotaExceededError, LlmAdapter, LlmError, ProviderRequestId, QUOTA_EXCEEDED_CODE } from '@deepseek-ai/dsh-llm'
import type { GenerateOptions, LlmModelInfo, LlmProviderInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
import { idleWatchdog, MAX_TIMER_DELAY_MS, timeoutOf } from '@deepseek-ai/dsh-timeout'
import { serializeRequest } from './serialize.ts'
import type { RequestDefaults } from './serialize.ts'
import { parseSse } from './sse.ts'
@@ -33,6 +34,31 @@ export interface DeepSeekAdapterOptions {
defaults?: RequestDefaults
/** Advisory models exposed to discovery consumers; requests remain unrestricted. */
models?: readonly DeepSeekCatalogModel[]
/** Maximum provider idle time while one stream read is outstanding. */
streamIdleTimeoutMs?: number
}
/** Default maximum idle interval while an adapter stream read is outstanding. */
export const DEFAULT_STREAM_IDLE_TIMEOUT_MS = 300_000
const STREAM_IDLE_TIMEOUT_CODE = 'LLM_STREAM_IDLE_TIMEOUT'
function retryAfterMs(value: string | null): number | undefined {
if (value === null) return undefined
if (/^\d+$/.test(value)) {
const delay = Number(value) * 1_000
return Number.isFinite(delay) && delay > 0 ? delay : undefined
}
const delay = Date.parse(value) - Date.now()
return Number.isFinite(delay) && delay > 0 ? delay : undefined
}
function requestId(headers: Headers): ReturnType<typeof ProviderRequestId> | undefined {
const value = headers.get('x-request-id') ?? headers.get('x-deepseek-request-id')
return value === null || value.length === 0 ? undefined : ProviderRequestId(value)
}
function errorMessage(value: unknown): string {
return value instanceof Error ? value.message : String(value)
}
/**
@@ -43,9 +69,10 @@ export interface DeepSeekAdapterOptions {
*/
export function httpErrorCode(status: number, error?: WireError['error']): string {
if (status === 401 || status === 403) return 'AUTH'
const detail = [error?.code, error?.type, error?.message].filter(Boolean).join(' ')
if (isQuotaExceededError(detail)) return QUOTA_EXCEEDED_CODE
if (status === 429) return 'RATE_LIMIT'
if (status === 400) {
const detail = [error?.code, error?.type, error?.message].filter(Boolean).join(' ')
if (isContextWindowExceededError(detail)) return CONTEXT_WINDOW_EXCEEDED_CODE
return 'INVALID_REQUEST'
}
@@ -57,13 +84,22 @@ export function httpErrorCode(status: number, error?: WireError['error']): strin
* The first real `LlmAdapter`. One instance serves every model name it was
* registered under (the harness model name IS the wire model name).
*
* Abort: `options.signal` is handed to fetch — both the initial request and
* the body stream reject on abort, which surfaces to the loop as a rejected
* step (the loop already contains step errors).
* One stable signal reaches both initial fetch and body reads. Caller aborts
* map to `ABORTED`; the configured per-read idle watchdog maps to `TIMEOUT`.
*/
export class DeepSeekAdapter extends LlmAdapter {
private readonly streamIdleTimeoutMs: number
constructor(private readonly options: DeepSeekAdapterOptions) {
super()
this.streamIdleTimeoutMs = options.streamIdleTimeoutMs ?? DEFAULT_STREAM_IDLE_TIMEOUT_MS
if (!Number.isFinite(this.streamIdleTimeoutMs)
|| this.streamIdleTimeoutMs <= 0
|| this.streamIdleTimeoutMs > MAX_TIMER_DELAY_MS) {
throw new Error(
`llm-deepseek: streamIdleTimeoutMs must be a positive finite number no greater than ${MAX_TIMER_DELAY_MS}`,
)
}
}
override providerInfo(provider: string): LlmProviderInfo {
@@ -80,6 +116,48 @@ export class DeepSeekAdapter extends LlmAdapter {
}
async * stream(options: GenerateOptions): AsyncIterable<StreamChunk> {
const consumer = new AbortController()
const upstream = options.signal === undefined
? consumer.signal
: AbortSignal.any([options.signal, consumer.signal])
using watchdog = idleWatchdog(upstream, this.streamIdleTimeoutMs, STREAM_IDLE_TIMEOUT_CODE)
const iterator = this.request(options, watchdog.signal)[Symbol.asyncIterator]()
let exhausted = false
try {
while (true) {
const result = await watchdog.next(iterator)
if (result.done) {
exhausted = true
return
}
yield result.value
}
} catch (error: unknown) {
if (timeoutOf(watchdog.signal, STREAM_IDLE_TIMEOUT_CODE) !== undefined) {
throw new LlmError(
`DeepSeek stream idle timeout after ${this.streamIdleTimeoutMs}ms`,
'TIMEOUT',
{ cause: error },
)
}
if (options.signal?.aborted) {
throw new LlmError('DeepSeek request aborted by caller', 'ABORTED', { cause: error })
}
if (error instanceof LlmError) throw error
throw new LlmError(`DeepSeek transport failed: ${errorMessage(error)}`, 'TRANSPORT', { cause: error })
} finally {
consumer.abort('DeepSeek stream consumer stopped')
if (!exhausted && iterator.return !== undefined) {
try {
await iterator.return()
} catch (_abortedTransportTeardown) {
// The consumer controller already owns termination; a return-time abort cannot add a second outcome.
}
}
}
}
private async * request(options: GenerateOptions, signal: AbortSignal): AsyncIterable<StreamChunk> {
const body = serializeRequest(options, this.options.defaults ?? {})
// TODO(http): adopt the Cordis HTTP service when shared transport configuration
@@ -96,7 +174,7 @@ export class DeepSeekAdapter extends LlmAdapter {
: {},
},
body: JSON.stringify(body),
...options.signal ? { signal: options.signal } : {},
signal,
})
if (!response.ok) {
@@ -110,7 +188,13 @@ export class DeepSeekAdapter extends LlmAdapter {
// Only swallow error-body parsing: the HTTP status still identifies the
// failure, so malformed gateway JSON must not mask it.
}
throw new LlmError(message, httpErrorCode(response.status, providerError))
const delay = retryAfterMs(response.headers.get('retry-after'))
const id = requestId(response.headers)
throw new LlmError(message, httpErrorCode(response.status, providerError), {
status: response.status,
...delay === undefined ? {} : { retryAfterMs: delay },
...id === undefined ? {} : { requestId: id },
})
}
if (!response.body) {
throw new LlmError('DeepSeek API returned no response body', 'EMPTY_RESPONSE')

View File

@@ -8,7 +8,8 @@
import type { Context } from 'cordis'
import z from 'schemastery'
import type {} from '@deepseek-ai/dsh-llm'
import { DeepSeekAdapter } from './adapter.ts'
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
import { DEFAULT_STREAM_IDLE_TIMEOUT_MS, DeepSeekAdapter } from './adapter.ts'
import type { DeepSeekCatalogModel } from './adapter.ts'
export { DeepSeekAdapter } from './adapter.ts'
@@ -41,6 +42,8 @@ export interface Config {
reasoningEffort?: 'high' | 'max'
/** Advisory models shown by discovery consumers; defaults to V4 Flash and V4 Pro. */
models?: DeepSeekCatalogModel[]
/** Maximum provider idle time while one stream read is outstanding (default five minutes). */
streamIdleTimeoutMs?: number
}
const catalogModel: z<DeepSeekCatalogModel> = z.object({
@@ -55,6 +58,7 @@ export const Config: z<Config> = z.object({
thinking: z.union(['enabled', 'disabled']),
reasoningEffort: z.union(['high', 'max']),
models: z.array(catalogModel).default(DEFAULT_MODELS),
streamIdleTimeoutMs: z.number().min(Number.MIN_VALUE).max(MAX_TIMER_DELAY_MS).default(DEFAULT_STREAM_IDLE_TIMEOUT_MS),
})
/** Public API default; the internal endpoint comes from $DEEPSEEK_BASE_URL. */
@@ -92,5 +96,6 @@ export function apply(ctx: Context, config: Config): void {
reasoningEffort: config.reasoningEffort,
},
models: resolveModels(config.models),
streamIdleTimeoutMs: config.streamIdleTimeoutMs ?? DEFAULT_STREAM_IDLE_TIMEOUT_MS,
}))
}

View File

@@ -35,7 +35,10 @@ export function mapFinishReason(reason: string): FinishReason {
case 'length': return { kind: 'max-tokens' }
default:
// content_filter, insufficient_system_resource, future additions.
return { kind: 'error', message: `model stopped: ${reason}`, code: reason.toUpperCase() }
return {
kind: 'error',
failure: { message: `model stopped: ${reason}`, code: reason.toUpperCase() },
}
}
}

View File

@@ -2,7 +2,14 @@ import { createServer } from 'node:http'
import type { IncomingMessage, Server, ServerResponse } from 'node:http'
import { afterEach, describe, expect, it, vi } from 'vitest'
import { Context } from 'cordis'
import LlmService, { CONTEXT_WINDOW_EXCEEDED_CODE, LlmError, userAgent } from '@deepseek-ai/dsh-llm'
import LlmService, {
CONTEXT_WINDOW_EXCEEDED_CODE,
LlmError,
ProviderRequestId,
QUOTA_EXCEEDED_CODE,
userAgent,
} from '@deepseek-ai/dsh-llm'
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
import { SessionId } from '@deepseek-ai/dsh-session'
import * as LlmDeepSeek from '@deepseek-ai/dsh-llm-deepseek'
import { DeepSeekAdapter } from '@deepseek-ai/dsh-llm-deepseek'
@@ -12,7 +19,7 @@ import { assemble } from './assemble.ts'
/** One scripted behavior for the next request the mock server receives. */
type Behavior =
| { kind: 'sse'; events: string[]; delayMs?: number }
| { kind: 'http-error'; status: number; body: string; contentType?: string }
| { kind: 'http-error'; status: number; body: string; contentType?: string; headers?: Record<string, string> }
| { kind: 'close-early'; events: string[] }
interface MockServer {
@@ -30,6 +37,7 @@ const servers: Server[] = []
afterEach(async () => {
await Promise.all(servers.splice(0).map(server => new Promise(resolve => server.close(resolve))))
vi.unstubAllEnvs()
vi.useRealTimers()
})
/** Local chat-completions stand-in: replays scripted behaviors per request. */
@@ -48,7 +56,10 @@ async function mockServer(script: Behavior[]): Promise<MockServer> {
return
}
if (behavior.kind === 'http-error') {
response.writeHead(behavior.status, { 'content-type': behavior.contentType ?? 'application/json' })
response.writeHead(behavior.status, {
'content-type': behavior.contentType ?? 'application/json',
...behavior.headers,
})
response.end(behavior.body)
return
}
@@ -202,6 +213,84 @@ describe('DeepSeekAdapter against a mock server', () => {
expect(code).toBe(CONTEXT_WINDOW_EXCEEDED_CODE)
})
it('retains status, Retry-After seconds, and provider request id as structured facts', async () => {
const server = await mockServer([{
kind: 'http-error',
status: 429,
body: JSON.stringify({ error: { message: 'slow down' } }),
headers: { 'retry-after': '2', 'x-request-id': 'req-429' },
}])
const ctx = await harness(server.url)
let thrown: unknown
try {
await assemble(ctx, { model: 'deepseek-v4-flash', messages: [] })
} catch (error: unknown) {
thrown = error
}
expect(thrown).toBeInstanceOf(LlmError)
expect((thrown as LlmError).failure).toEqual({
message: 'slow down',
code: 'RATE_LIMIT',
status: 429,
retryAfterMs: 2_000,
requestId: ProviderRequestId('req-429'),
})
})
it('parses a future Retry-After HTTP date and the DeepSeek request-id fallback', async () => {
const now = 1_800_000_000_000
const dateNow = vi.spyOn(Date, 'now').mockReturnValue(now)
try {
const server = await mockServer([{
kind: 'http-error',
status: 503,
body: JSON.stringify({ error: { message: 'come back later' } }),
headers: {
'retry-after': new Date(now + 3_000).toUTCString(),
'x-deepseek-request-id': 'deepseek-503',
},
}])
const ctx = await harness(server.url)
await expect(assemble(ctx, { model: 'deepseek-v4-flash', messages: [] }))
.rejects.toMatchObject({
failure: {
message: 'come back later',
code: 'SERVER',
status: 503,
retryAfterMs: 3_000,
requestId: ProviderRequestId('deepseek-503'),
},
})
} finally {
dateNow.mockRestore()
}
})
it('omits zero, non-finite, invalid, and past Retry-After values', async () => {
const values = [
'0',
'9'.repeat(400),
'not-a-date',
new Date(0).toUTCString(),
]
for (const value of values) {
const server = await mockServer([{
kind: 'http-error',
status: 429,
body: JSON.stringify({ error: { message: 'retry later' } }),
headers: { 'retry-after': value },
}])
const ctx = await harness(server.url)
let thrown: LlmError | undefined
try {
await assemble(ctx, { model: 'deepseek-v4-flash', messages: [] })
} catch (error: unknown) {
if (error instanceof LlmError) thrown = error
}
expect(thrown?.failure).toEqual({ message: 'retry later', code: 'RATE_LIMIT', status: 429 })
}
})
it('classifies only context-capacity HTTP 400 details as context overflow', () => {
expect(httpErrorCode(400, { message: 'request too large for model context' }))
.toBe(CONTEXT_WINDOW_EXCEEDED_CODE)
@@ -210,6 +299,12 @@ describe('DeepSeekAdapter against a mock server', () => {
expect(httpErrorCode(413, { code: 'context_length_exceeded' })).toBe('HTTP_413')
})
it('distinguishes terminal quota exhaustion from transient HTTP 429 throttling', () => {
expect(httpErrorCode(429, { code: 'insufficient_quota', message: 'account credits exhausted' }))
.toBe(QUOTA_EXCEEDED_CODE)
expect(httpErrorCode(429, { message: 'request rate limit exceeded' })).toBe('RATE_LIMIT')
})
it('keeps the status-line message for JSON error bodies without a message', async () => {
const server = await mockServer([{ kind: 'http-error', status: 500, body: '{"error":{"type":"x"}}' }])
const ctx = await harness(server.url)
@@ -272,7 +367,76 @@ describe('DeepSeekAdapter against a mock server', () => {
})()
setTimeout(() => { controller.abort() }, 30)
await expect(pending).rejects.toThrow()
await expect(pending).rejects.toMatchObject({ code: 'ABORTED' })
})
it('maps connection failures to TRANSPORT without losing the cause', async () => {
const cause = new TypeError('connection refused')
const fetchSpy = vi.spyOn(globalThis, 'fetch').mockRejectedValue(cause)
const adapter = new DeepSeekAdapter({ apiKey: 'k', baseURL: 'https://example.invalid' })
try {
const drain = async (): Promise<void> => {
for await (const _chunk of adapter.stream({ provider: 'deepseek', model: 'm', messages: [] })) { /* drain */ }
}
await expect(drain()).rejects.toMatchObject({ code: 'TRANSPORT', cause })
} finally {
fetchSpy.mockRestore()
}
})
it('renders a non-Error transport rejection without losing its cause', async () => {
const fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation(() => {
const failed = Promise.withResolvers<Response>()
failed.reject('offline')
return failed.promise
})
const adapter = new DeepSeekAdapter({ apiKey: 'k', baseURL: 'https://example.invalid' })
try {
const drain = async (): Promise<void> => {
for await (const _chunk of adapter.stream({ provider: 'deepseek', model: 'm', messages: [] })) { /* drain */ }
}
await expect(drain()).rejects.toMatchObject({
message: 'DeepSeek transport failed: offline',
code: 'TRANSPORT',
cause: 'offline',
})
} finally {
fetchSpy.mockRestore()
}
})
it('aborts the underlying body when the stream stays idle past its watchdog', async () => {
vi.useFakeTimers()
let stopped = false
const fetchSpy = vi.spyOn(globalThis, 'fetch').mockImplementation((_input, init) => {
const signal = init?.signal
const body = new ReadableStream<Uint8Array>({
start(controller) {
signal?.addEventListener('abort', () => {
stopped = true
controller.error(signal.reason)
}, { once: true })
},
})
return Promise.resolve(new Response(body, { status: 200 }))
})
const adapter = new DeepSeekAdapter({
apiKey: 'k',
baseURL: 'https://example.invalid',
streamIdleTimeoutMs: 100,
})
try {
const drain = (async () => {
for await (const _chunk of adapter.stream({ provider: 'deepseek', model: 'm', messages: [] })) { /* drain */ }
})()
const rejected = expect(drain).rejects.toMatchObject({ code: 'TIMEOUT' })
await vi.advanceTimersByTimeAsync(0)
await vi.advanceTimersByTimeAsync(100)
await rejected
expect(stopped).toBe(true)
} finally {
fetchSpy.mockRestore()
}
})
})
@@ -419,4 +583,30 @@ describe('plugin registration and config', () => {
expect(adapter).toBeInstanceOf(DeepSeekAdapter)
await expect(adapter.listModels('deepseek')).resolves.toEqual([])
})
it('rejects invalid idle watchdog bounds for direct and plugin composition', async () => {
expect(() => new DeepSeekAdapter({
apiKey: 'k',
baseURL: 'http://127.0.0.1:1',
streamIdleTimeoutMs: Number.POSITIVE_INFINITY,
})).toThrow(/streamIdleTimeoutMs.*positive finite/)
expect(() => new DeepSeekAdapter({
apiKey: 'k',
baseURL: 'http://127.0.0.1:1',
streamIdleTimeoutMs: MAX_TIMER_DELAY_MS + 1,
})).toThrow(/streamIdleTimeoutMs.*no greater/)
const ctx = new Context()
await ctx.plugin(LlmService)
await expect(ctx.plugin(LlmDeepSeek, {
apiKey: 'k',
baseURL: 'http://127.0.0.1:1',
streamIdleTimeoutMs: 0,
})).rejects.toThrow(/streamIdleTimeoutMs/)
await expect(ctx.plugin(LlmDeepSeek, {
apiKey: 'k',
baseURL: 'http://127.0.0.1:1',
streamIdleTimeoutMs: MAX_TIMER_DELAY_MS + 1,
})).rejects.toThrow(/streamIdleTimeoutMs/)
})
})

View File

@@ -232,8 +232,7 @@ describe('mapFinishReason', () => {
(wire) => {
expect(mapFinishReason(wire)).toEqual({
kind: 'error',
message: `model stopped: ${wire}`,
code: wire.toUpperCase(),
failure: { message: `model stopped: ${wire}`, code: wire.toUpperCase() },
})
},
)

View File

@@ -19,6 +19,9 @@
},
{
"path": "../../llm/llm"
},
{
"path": "../../util/timeout"
}
]
}