返回源码地图

packages/llm/llm-retry/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 09 / 09 / 09;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Provider-routed model-request retry policy on the agent loop's request
3 * recovery extension point. Each scheduled retry is durable before its cancellable wait.
4 *
5 * @module @deepseek-ai/dsh-llm-retry
6 */
7
8import { randomUUID } from 'node:crypto'
9import type { Context, Events } from '@deepseek-ai/cordis'
10import z from '@deepseek-ai/schemastery'
11import { z as zod } from 'zod'
12import type { Agent, RequestErrorAction } from '@deepseek-ai/dsh-agent'
13import type { LlmFailure, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'
14import type {} from '@deepseek-ai/dsh-session-projection'
15import { RetryId } from './brand.ts'
16import type { LlmRetryEventData } from './types.ts'
17
18export type { LlmRetryEventData, LlmRetryStartedEventData } from './types.ts'
19export { RetryId } from './brand.ts'
20
21export const name = 'llm-retry'
22export const inject = ['agents', 'sessionProjections']
23
24/** This policy executor has no config; providers own `retryPolicy`. */
25export type Config = Readonly<Record<string, never>>
26
27/** Runtime schema for {@link Config}. */
28export const Config = z.object({}) as unknown as z<Config>
29
30function validateConfig(config: Config): void {
31 const [key] = Object.keys(config)
32 if (key === undefined) return
33 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}
38
39/** Non-serializable hooks used to make timing policy deterministic in tests. */
40export interface RetryInternals {
41 /** Random sample in the inclusive zero-to-one range used for jitter. */
42 random?: () => number
43}
44
45type DownstreamOutcome =
46 | { readonly type: 'decision'; readonly decision: RequestErrorAction }
47 | { readonly type: 'error'; readonly error: unknown }
48
49async 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}
58
59function 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}
65
66function 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}
78
79function retryStateKey(provider: string, policyKey: string): string {
80 return JSON.stringify([provider, policyKey])
81}
82
83function 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}
97
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 */
104interface RetryStateEntry {
105 retry: number
106 retryId: RetryId
107}
108
109type LlmRetryState = Record<string, RetryStateEntry>
110
111// The cast bridges the branded retry id, which Zod cannot express directly.
112const 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>
116declare module '@deepseek-ai/dsh-session-projection/types' {
117 interface SessionProjectionStateMap {
118 /** Retry state for the current step by provider and policy. */
119 llmRetry: LlmRetryState
120 }
121}
122
123export 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 state
133 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 state
136 return { ...state, [key]: { retry: event.data.retry, retryId: event.data.retryId } }
137 },
138 })
139 const random = internals.random ?? Math.random
140 const lifetime = new AbortController()
141 const active = new Set<Promise<RequestErrorAction>>()
142
143 function track(operation: Promise<RequestErrorAction>): Promise<RequestErrorAction> {
144 const tracked = operation.finally(() => active.delete(tracked))
145 active.add(tracked)
146 return tracked
147 }
148
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) return
164 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)) return
190 agent.session.append('llm/retry-started', { retryId, turn, step, retry })
191 return { kind: 'retry' }
192 }
193
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) return
201 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) return
206 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.decision
214 }
215 } else if (!policy.retryableCodes.includes(failure.code)) {
216 return next()
217 }
218
219 const policyKey = retryPolicyKey(policy)
220 const retryState = ctx.sessionProjections.stateOf(agent.session, 'llmRetry') as LlmRetryState
221 const previous = retryState[retryStateKey(provider, policyKey)]
222 const previousRetry = previous?.retry ?? 0
223 if (policy.mode === 'normal' && previousRetry >= policy.maxRetries) return next()
224 const retry = previousRetry + 1
225 const retryId = previous?.retryId ?? RetryId(randomUUID())
226 let delayMs: number
227 if (failure.providerRetryAfterMs !== undefined
228 && 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.providerRetryAfterMs
235 }
236 } else {
237 delayMs = localDelay(policy, retry, random)
238 }
239
240 return backoff(agent, turn, step, failure, provider, policy, policyKey, retry, retryId, delayMs, signal)
241 }
242
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 was
248 // removed. Lifetime cancellation must prevent that stale callback from
249 // entering a downstream policy after disposal.
250 if (lifetime.signal.aborted) return Promise.resolve<RequestErrorAction>(undefined)
251 return track(recover(payload, next))
252 })
253
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}