1
/**2
* Provider-routed model-request retry policy on the agent loop's request3
* recovery extension point. Each scheduled retry is durable before its cancellable wait.4
*5
* @module @deepseek-ai/dsh-llm-retry6
*/8
import { randomUUID } from 'node:crypto'9
import type { Context, Events } from '@deepseek-ai/cordis'10
import z from '@deepseek-ai/schemastery'11
import { z as zod } from 'zod'12
import type { Agent, RequestErrorAction } from '@deepseek-ai/dsh-agent'13
import type { LlmFailure, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'14
import type {} from '@deepseek-ai/dsh-session-projection'15
import { RetryId } from './brand.ts'16
import type { LlmRetryEventData } from './types.ts'18
export type { LlmRetryEventData, LlmRetryStartedEventData } from './types.ts'19
export { RetryId } from './brand.ts'21
export const name = 'llm-retry'22
export const inject = ['agents', 'sessionProjections']24
/** This policy executor has no config; providers own `retryPolicy`. */25
export type Config = Readonly<Record<string, never>>27
/** Runtime schema for {@link Config}. */28
export const Config = z.object({}) as unknown as z<Config>30
function validateConfig(config: Config): void {31
const [key] = Object.keys(config)32
if (key === undefined) return33
if (key === 'retryPolicy') {34
throw new Error('llm-retry: retryPolicy belongs under each provider configuration')35
}36
throw new Error(`llm-retry: unknown key "${key}"`)37
}39
/** Non-serializable hooks used to make timing policy deterministic in tests. */40
export interface RetryInternals {41
/** Random sample in the inclusive zero-to-one range used for jitter. */42
random?: () => number43
}45
type DownstreamOutcome =46
| { readonly type: 'decision'; readonly decision: RequestErrorAction }47
| { readonly type: 'error'; readonly error: unknown }49
async function settleDownstream(50
next: () => Promise<RequestErrorAction>,51
): Promise<DownstreamOutcome> {52
try {53
return { type: 'decision', decision: await next() }54
} catch (error: unknown) {55
return { type: 'error', error }56
}57
}59
function localDelay(config: ResolvedRetryPolicy, retry: number, random: () => number): number {60
const exponent = Math.min(retry - 1, 1024)61
const exponential = Math.min(config.initialDelayMs * 2 ** exponent, config.maxDelayMs)62
const jitter = 1 - config.jitterRatio + 2 * config.jitterRatio * random()63
return Math.min(exponential * jitter, config.maxDelayMs)64
}66
function retryPolicyKey(policy: ResolvedRetryPolicy): string {67
return policy.mode === 'always'68
? JSON.stringify([policy.mode, policy.initialDelayMs, policy.maxDelayMs, policy.jitterRatio])69
: JSON.stringify([70
policy.mode,71
policy.maxRetries,72
[...policy.retryableCodes].sort(),73
policy.initialDelayMs,74
policy.maxDelayMs,75
policy.jitterRatio,76
])77
}79
function retryStateKey(provider: string, policyKey: string): string {80
return JSON.stringify([provider, policyKey])81
}83
function cancellableDelay(delayMs: number, signal: AbortSignal): Promise<boolean> {84
if (signal.aborted) return Promise.resolve(false)85
return new Promise((resolve) => {86
const timer = setTimeout(() => {87
signal.removeEventListener('abort', onAbort)88
resolve(true)89
}, delayMs)90
function onAbort(): void {91
clearTimeout(timer)92
resolve(false)93
}94
signal.addEventListener('abort', onAbort, { once: true })95
})96
}98
/**99
* Install provider-routed normal or unbounded request recovery.100
* @param ctx - plugin context that owns the listener and active waits.101
* @param config - empty executor config; provider registrations own policy.102
* @param internals - non-serializable deterministic hooks for tests.103
*/104
interface RetryStateEntry {105
retry: number106
retryId: RetryId107
}109
type LlmRetryState = Record<string, RetryStateEntry>111
// The cast bridges the branded retry id, which Zod cannot express directly.112
const llmRetryStateSchema: zod.ZodType<LlmRetryState> = zod.record(zod.string(), zod.object({113
retry: zod.number().int().nonnegative(),114
retryId: zod.string(),115
})) as unknown as zod.ZodType<LlmRetryState>116
declare module '@deepseek-ai/dsh-session-projection/types' {117
interface SessionProjectionStateMap {118
/** Retry state for the current step by provider and policy. */119
llmRetry: LlmRetryState120
}121
}123
export function apply(ctx: Context, config: Config = {}, internals: RetryInternals = {}): void {124
validateConfig(config)125
ctx.sessionProjections.register({126
key: 'llmRetry',127
stateVersion: 1,128
stateSchema: llmRetryStateSchema,129
init: () => ({}),130
apply: (state, event) => {131
if (event.type === 'step/start' || event.type === 'turn/end') return {}132
if (event.type !== 'llm/retry') return state133
const key = retryStateKey(event.data.provider, event.data.policyKey)134
const entry = state[key]135
if (entry?.retry === event.data.retry && entry.retryId === event.data.retryId) return state136
return { ...state, [key]: { retry: event.data.retry, retryId: event.data.retryId } }137
},138
})139
const random = internals.random ?? Math.random140
const lifetime = new AbortController()141
const active = new Set<Promise<RequestErrorAction>>()143
function track(operation: Promise<RequestErrorAction>): Promise<RequestErrorAction> {144
const tracked = operation.finally(() => active.delete(tracked))145
active.add(tracked)146
return tracked147
}149
async function backoff(150
agent: Agent,151
turn: number,152
step: number,153
failure: LlmFailure,154
provider: string,155
policy: ResolvedRetryPolicy,156
policyKey: string,157
retry: number,158
retryId: RetryId,159
delayMs: number,160
signal: AbortSignal,161
): Promise<RequestErrorAction> {162
const fusedSignal = AbortSignal.any([signal, lifetime.signal])163
if (fusedSignal.aborted) return164
const eventData: LlmRetryEventData = policy.mode === 'normal'165
? {166
retryId,167
turn,168
step,169
provider,170
mode: policy.mode,171
policyKey,172
retry,173
maxRetries: policy.maxRetries,174
delayMs,175
failure,176
}177
: {178
retryId,179
turn,180
step,181
provider,182
mode: policy.mode,183
policyKey,184
retry,185
delayMs,186
failure,187
}188
agent.session.append('llm/retry', eventData)189
if (!await cancellableDelay(delayMs, fusedSignal)) return190
agent.session.append('llm/retry-started', { retryId, turn, step, retry })191
return { kind: 'retry' }192
}194
async function recover(195
{ agent, turn, step, provider, failure, retryPolicy: policy, signal }: Parameters<Events['agent/request-error']>[0],196
next: () => Promise<RequestErrorAction>,197
): Promise<RequestErrorAction> {198
if (policy === undefined) return next()199
if (policy.mode === 'always') {200
if (signal.aborted || lifetime.signal.aborted) return201
const fusedSignal = AbortSignal.any([signal, lifetime.signal])202
// The loop and plugin lifetime stay open until delegated recovery settles.203
// An abort then wins before the decision or fallback can mutate later state.204
const downstream = await settleDownstream(next)205
if (fusedSignal.aborted) return206
if (downstream.type === 'error') {207
ctx.logger.warn(208
`llm-retry: provider "${provider}" always policy ignored a downstream recovery failure: %o`,209
downstream.error,210
)211
}212
if (downstream.type === 'decision' && downstream.decision?.kind === 'retry') {213
return downstream.decision214
}215
} else if (!policy.retryableCodes.includes(failure.code)) {216
return next()217
}219
const policyKey = retryPolicyKey(policy)220
const retryState = ctx.sessionProjections.stateOf(agent.session, 'llmRetry') as LlmRetryState221
const previous = retryState[retryStateKey(provider, policyKey)]222
const previousRetry = previous?.retry ?? 0223
if (policy.mode === 'normal' && previousRetry >= policy.maxRetries) return next()224
const retry = previousRetry + 1225
const retryId = previous?.retryId ?? RetryId(randomUUID())226
let delayMs: number227
if (failure.providerRetryAfterMs !== undefined228
&& Number.isFinite(failure.providerRetryAfterMs)229
&& failure.providerRetryAfterMs > 0) {230
if (failure.providerRetryAfterMs > policy.maxDelayMs) {231
if (policy.mode === 'normal') return next()232
delayMs = localDelay(policy, retry, random)233
} else {234
delayMs = failure.providerRetryAfterMs235
}236
} else {237
delayMs = localDelay(policy, retry, random)238
}240
return backoff(agent, turn, step, failure, provider, policy, policyKey, retry, retryId, delayMs, signal)241
}243
const disposeListener = ctx.on('agent/request-error', (244
payload,245
next: () => Promise<RequestErrorAction>,246
) => {247
// A waterfall may have captured this callback before its registration was248
// removed. Lifetime cancellation must prevent that stale callback from249
// entering a downstream policy after disposal.250
if (lifetime.signal.aborted) return Promise.resolve<RequestErrorAction>(undefined)251
return track(recover(payload, next))252
})254
ctx.effect(() => async () => {255
disposeListener()256
lifetime.abort(new Error('llm-retry plugin disposed'))257
await Promise.allSettled([...active])258
}, 'llm-retry: abort and drain active recovery')259
}