e0988abc74
# Conflicts: # .agents/notes/implemented/architecture/2026-06-21-bounded-llm-request-recovery.i18n.yaml # .agents/notes/implemented/feature/2026-06-18-compaction-capability-seam.i18n.yaml # .agents/notes/implemented/feature/2026-06-18-compaction-capability-seam.md # .agents/notes/implemented/feature/2026-06-18-compaction-capability-seam.zh.md # .agents/notes/implemented/feature/2026-07-06-sandbox.i18n.yaml # .agents/notes/implemented/feature/2026-07-06-sandbox.md # .agents/notes/implemented/feature/2026-07-06-sandbox.zh.md # .agents/notes/implemented/feature/2026-07-27-tmux-location-context.i18n.yaml # .agents/notes/implemented/simplification/2026-06-20-public-agent-stop-surface.i18n.yaml # .agents/notes/implemented/simplification/2026-07-28-remove-synthetic-log-only-turns.i18n.yaml # .agents/notes/implemented/simplification/2026-07-30-private-agent-send.i18n.yaml # docs/architecture.i18n.yaml # docs/architecture.md # docs/architecture.zh.md # docs/config-catalog.md # docs/cordis-catalog/events.md # docs/cordis-catalog/services.md # docs/core-data-structures/compaction.i18n.yaml # docs/core-data-structures/core.i18n.yaml # docs/core-data-structures/core.md # docs/core-data-structures/core.zh.md # docs/core-data-structures/llm-streaming.i18n.yaml # docs/core-data-structures/llm-streaming.md # docs/core-data-structures/llm-streaming.zh.md # docs/core-data-structures/session.i18n.yaml # docs/event-producer-consumer.md # docs/module-graph.md # docs/persistence-catalog.md # examples/acp-agent/tests/goal-snapshots/goal-session/session.expected.jsonl # examples/acp-agent/tests/snapshots/advanced-toolchain/session.1.jsonl # examples/acp-agent/tests/snapshots/advanced-toolchain/session.2.jsonl # examples/acp-agent/tests/snapshots/advanced-toolchain/session.jsonl # examples/acp-agent/tests/snapshots/bash-spill/session.jsonl # examples/acp-agent/tests/snapshots/bash-tool-turn/session.jsonl # examples/acp-agent/tests/snapshots/both-mode-turn/session.jsonl # examples/acp-agent/tests/snapshots/cancel-tool-calls/session.jsonl # examples/acp-agent/tests/snapshots/cancel/session.jsonl # examples/acp-agent/tests/snapshots/code-mode-turn/session.jsonl # examples/acp-agent/tests/snapshots/code-mode-workspace-context/session.jsonl # examples/acp-agent/tests/snapshots/cordis-inspect-jsdoc/session.jsonl # examples/acp-agent/tests/snapshots/empty-response-retry/session.jsonl # examples/acp-agent/tests/snapshots/error-finish/session.jsonl # examples/acp-agent/tests/snapshots/escalation-approved/session.jsonl # examples/acp-agent/tests/snapshots/escalation-rejected/session.jsonl # examples/acp-agent/tests/snapshots/fs-edit/session.jsonl # examples/acp-agent/tests/snapshots/fs-escalation-approved/session.jsonl # examples/acp-agent/tests/snapshots/fs-glob-sampling/session.jsonl # examples/acp-agent/tests/snapshots/fs-policy-reject/session.jsonl # examples/acp-agent/tests/snapshots/fs-read-window/session.jsonl # examples/acp-agent/tests/snapshots/fs-read/session.jsonl # examples/acp-agent/tests/snapshots/fs-write-overwrite/session.jsonl # examples/acp-agent/tests/snapshots/fs-write/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-invalid-matcher/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-posttool-block/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-posttool-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-pretool-ask/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-pretool-deny/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-promptsubmit-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-cc-stop-continue/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-invalid-matcher/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-posttool-block/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-posttool-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-pretool-block/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-promptsubmit-context/session.jsonl # examples/acp-agent/tests/snapshots/hook-codex-stop-continue/session.jsonl # examples/acp-agent/tests/snapshots/lsp-definition/session.jsonl # examples/acp-agent/tests/snapshots/multi-turn/session.jsonl # examples/acp-agent/tests/snapshots/packed-chunks/session.jsonl # examples/acp-agent/tests/snapshots/parallel-tool-calls/session.jsonl # examples/acp-agent/tests/snapshots/pty-tools/session.jsonl # examples/acp-agent/tests/snapshots/repeat-tool-guard/session.jsonl # examples/acp-agent/tests/snapshots/session-query-spill/session.jsonl # examples/acp-agent/tests/snapshots/session-sandbox-root/session.jsonl # examples/acp-agent/tests/snapshots/session-title-after-turn/session.jsonl # examples/acp-agent/tests/snapshots/skill-load/session.jsonl # examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.2.jsonl # examples/acp-agent/tests/snapshots/subagent-depth-two-rejection/session.jsonl # examples/acp-agent/tests/snapshots/subagent-fork/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-fork/session.jsonl # examples/acp-agent/tests/snapshots/subagent-mixed/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-mixed/session.2.jsonl # examples/acp-agent/tests/snapshots/subagent-mixed/session.jsonl # examples/acp-agent/tests/snapshots/subagent-multi/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-multi/session.2.jsonl # examples/acp-agent/tests/snapshots/subagent-multi/session.jsonl # examples/acp-agent/tests/snapshots/subagent-spawn/session.1.jsonl # examples/acp-agent/tests/snapshots/subagent-spawn/session.jsonl # examples/acp-agent/tests/snapshots/text-turn/session.jsonl # examples/acp-agent/tests/snapshots/todo-write/session.jsonl # examples/acp-agent/tests/snapshots/tool-call-turn/session.jsonl # examples/acp-agent/tests/snapshots/web-fetch/session.jsonl # examples/acp-agent/tests/snapshots/workflow-run/session.1.jsonl # examples/acp-agent/tests/snapshots/workflow-run/session.jsonl # examples/acp-agent/tests/snapshots/workspace-context/session.jsonl # examples/acp-agent/tests/snapshots/workspace-edit/session.jsonl # examples/headless-agent/tests/semantic-checkpoint-snapshots/tool-outcome-unknown/session.expected.jsonl # examples/headless-agent/tests/snapshots/advanced-toolchain/session.1.jsonl # examples/headless-agent/tests/snapshots/advanced-toolchain/session.2.jsonl # examples/headless-agent/tests/snapshots/advanced-toolchain/session.jsonl # examples/headless-agent/tests/snapshots/advanced-toolchain/stream-json.expected.jsonl # examples/headless-agent/tests/snapshots/goal-tools/stream-json.expected.jsonl # examples/headless-agent/tests/snapshots/missing-credential/stream-json.expected.jsonl # examples/headless-agent/tests/snapshots/provider-retry/stream-json.expected.jsonl # examples/headless-agent/tests/snapshots/pty-tools/session.jsonl # examples/headless-agent/tests/snapshots/pty-tools/stream-json.expected.jsonl # examples/headless-agent/tests/snapshots/ralph-loop/stream-json.expected.jsonl # examples/headless-agent/tests/subagent-inheritance-snapshots/parent-override/child.expected.jsonl # examples/headless-agent/tests/subagent-inheritance-snapshots/parent-override/parent.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/bash-tool/notifications.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/bash-tool/session.jsonl # examples/jsonrpc-agent/tests/snapshots/persistent-tools/notifications.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/persistent-tools/session.jsonl # examples/jsonrpc-agent/tests/snapshots/subagent-spawn/notifications.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/subagent-spawn/session.1.jsonl # examples/jsonrpc-agent/tests/snapshots/subagent-spawn/session.jsonl # examples/jsonrpc-agent/tests/snapshots/text-turn/notifications.expected.jsonl # examples/jsonrpc-agent/tests/snapshots/text-turn/session.jsonl # packages/client/runtime/README.i18n.yaml # packages/client/runtime/src/client/sessions/request-inspection.ts # packages/compact/compact-basic/README.i18n.yaml # packages/compact/compact-basic/README.md # packages/compact/compact-basic/README.zh.md # packages/compact/compact-basic/src/index.ts # packages/context/time-context/tests/time-context.spec.ts # packages/context/tmux-context/README.i18n.yaml # packages/context/tmux-context/tests/tmux-context.spec.ts # packages/context/workspace-context/tests/workspace-context.spec.ts # packages/cordis/tool-cordis/src/api-catalog.ts # packages/core/agent-loop/README.i18n.yaml # packages/core/agent-loop/README.md # packages/core/agent-loop/README.zh.md # packages/core/agent-loop/src/agent.ts # packages/core/agent/README.i18n.yaml # packages/core/agent/README.md # packages/core/agent/README.zh.md # packages/core/agent/src/types.ts # packages/core/session/README.i18n.yaml # packages/core/session/README.md # packages/core/session/README.zh.md # packages/fs/tool-str-replace-editor/tests/tools.spec.ts # packages/goal/command-goal/tests/command-goal.spec.ts # packages/goal/goal/tests/goal.spec.ts # packages/host/apiproxy/README.i18n.yaml # packages/host/apiproxy/README.md # packages/host/apiproxy/README.zh.md # packages/host/apiproxy/src/api/index.ts # packages/host/apiproxy/tests/api-proxy-workspace.spec.ts # packages/llm/llm/README.i18n.yaml # packages/llm/llm/README.md # packages/llm/llm/README.zh.md # packages/llm/llm/src/index.ts # packages/pty/pty-local/tests/index.spec.ts # packages/pty/pty-local/tests/local.spec.ts # packages/pty/pty/tests/service.spec.ts # packages/pty/tool-bash-persistent/tests/loader-composition.spec.ts # packages/pty/tool-bash-persistent/tests/tools.spec.ts # packages/pty/tool-pty/tests/loader-composition.spec.ts # packages/pty/tool-pty/tests/tools.spec.ts # packages/session-persistence/session-checkpoint-policy/tests/crash-recovery.e2e.ts # packages/skill/tool-skill/tests/tool-skill.spec.ts # packages/tasks/tasks-local/tests/tasks.spec.ts # packages/ui/tui/README.i18n.yaml # packages/ui/tui/tests/tui.spec.ts # packages/ui/user-approval/src/index.ts # packages/ui/user-approval/tests/approval.spec.ts
249 lines
8.1 KiB
TypeScript
249 lines
8.1 KiB
TypeScript
/**
|
|
* Provider-routed model-request retry policy on the agent loop's request
|
|
* recovery seam. Each scheduled retry is durable before its cancellable wait.
|
|
*
|
|
* @module @deepseek-ai/dsh-llm-retry
|
|
*/
|
|
|
|
import type { Context } from 'cordis'
|
|
import z from 'schemastery'
|
|
import type { Agent, RequestErrorAction, RequestFailureContext } from '@deepseek-ai/dsh-agent'
|
|
import type { LlmFailure, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'
|
|
import type { SessionEvent } from '@deepseek-ai/dsh-session'
|
|
|
|
declare module '@deepseek-ai/dsh-session' {
|
|
interface SessionEventMap {
|
|
/** Durable, non-surface record of one provider-routed retry scheduled after a failed request attempt. */
|
|
'llm/retry': {
|
|
turn: number
|
|
step: number
|
|
provider: string
|
|
mode: 'normal'
|
|
policyKey: string
|
|
retry: number
|
|
maxRetries: number
|
|
delayMs: number
|
|
failure: LlmFailure
|
|
} | {
|
|
turn: number
|
|
step: number
|
|
provider: string
|
|
mode: 'always'
|
|
policyKey: string
|
|
retry: number
|
|
delayMs: number
|
|
failure: LlmFailure
|
|
}
|
|
}
|
|
}
|
|
|
|
export type { LlmRetryEventData } from './types.ts'
|
|
|
|
export const name = 'llm-retry'
|
|
export const inject = ['agents']
|
|
|
|
/** This policy executor has no config; providers own `retryPolicy`. */
|
|
export type Config = Readonly<Record<string, never>>
|
|
|
|
/** Runtime schema for {@link Config}. */
|
|
export const Config = z.object({}) as unknown as z<Config>
|
|
|
|
function validateConfig(config: Config): void {
|
|
const [key] = Object.keys(config)
|
|
if (key === undefined) return
|
|
if (key === 'retryPolicy') {
|
|
throw new Error('llm-retry: retryPolicy belongs under each provider configuration')
|
|
}
|
|
throw new Error(`llm-retry: unknown key "${key}"`)
|
|
}
|
|
|
|
/** Non-serializable seams used to make timing policy deterministic in tests. */
|
|
export interface RetryInternals {
|
|
/** Random sample in the inclusive zero-to-one range used for jitter. */
|
|
random?: () => number
|
|
}
|
|
|
|
type DownstreamOutcome =
|
|
| { readonly type: 'decision'; readonly decision: RequestErrorAction }
|
|
| { readonly type: 'error'; readonly error: unknown }
|
|
|
|
async function settleDownstream(
|
|
next: () => Promise<RequestErrorAction>,
|
|
): Promise<DownstreamOutcome> {
|
|
try {
|
|
return { type: 'decision', decision: await next() }
|
|
} catch (error: unknown) {
|
|
return { type: 'error', error }
|
|
}
|
|
}
|
|
|
|
function localDelay(config: ResolvedRetryPolicy, retry: number, random: () => number): number {
|
|
const exponent = Math.min(retry - 1, 1024)
|
|
const exponential = Math.min(config.initialDelayMs * 2 ** exponent, config.maxDelayMs)
|
|
const jitter = 1 - config.jitterRatio + 2 * config.jitterRatio * random()
|
|
return Math.min(exponential * jitter, config.maxDelayMs)
|
|
}
|
|
|
|
function retryPolicyKey(policy: ResolvedRetryPolicy): string {
|
|
return policy.mode === 'always'
|
|
? JSON.stringify([policy.mode, policy.initialDelayMs, policy.maxDelayMs, policy.jitterRatio])
|
|
: JSON.stringify([
|
|
policy.mode,
|
|
policy.maxRetries,
|
|
[...policy.retryableCodes].sort(),
|
|
policy.initialDelayMs,
|
|
policy.maxDelayMs,
|
|
policy.jitterRatio,
|
|
])
|
|
}
|
|
|
|
function cancellableDelay(delayMs: number, signal: AbortSignal): Promise<boolean> {
|
|
if (signal.aborted) return Promise.resolve(false)
|
|
return new Promise((resolve) => {
|
|
const timer = setTimeout(() => {
|
|
signal.removeEventListener('abort', onAbort)
|
|
resolve(true)
|
|
}, delayMs)
|
|
function onAbort(): void {
|
|
clearTimeout(timer)
|
|
resolve(false)
|
|
}
|
|
signal.addEventListener('abort', onAbort, { once: true })
|
|
})
|
|
}
|
|
|
|
/**
|
|
* Install provider-routed normal or unbounded request recovery.
|
|
* @param ctx - plugin context that owns the listener and active waits.
|
|
* @param config - empty executor config; provider registrations own policy.
|
|
* @param internals - non-serializable deterministic seams for tests.
|
|
*/
|
|
export function apply(ctx: Context, config: Config = {}, internals: RetryInternals = {}): void {
|
|
validateConfig(config)
|
|
const random = internals.random ?? Math.random
|
|
const lifetime = new AbortController()
|
|
const active = new Set<Promise<RequestErrorAction>>()
|
|
|
|
function track(operation: Promise<RequestErrorAction>): Promise<RequestErrorAction> {
|
|
const tracked = operation.finally(() => active.delete(tracked))
|
|
active.add(tracked)
|
|
return tracked
|
|
}
|
|
|
|
async function backoff(
|
|
agent: Agent,
|
|
turn: number,
|
|
step: number,
|
|
failure: LlmFailure,
|
|
provider: string,
|
|
policy: ResolvedRetryPolicy,
|
|
policyKey: string,
|
|
retry: number,
|
|
delayMs: number,
|
|
signal: AbortSignal,
|
|
): Promise<RequestErrorAction> {
|
|
const fusedSignal = AbortSignal.any([signal, lifetime.signal])
|
|
if (fusedSignal.aborted) return
|
|
const eventData = policy.mode === 'normal'
|
|
? {
|
|
turn,
|
|
step,
|
|
provider,
|
|
mode: policy.mode,
|
|
policyKey,
|
|
retry,
|
|
maxRetries: policy.maxRetries,
|
|
delayMs,
|
|
failure,
|
|
}
|
|
: {
|
|
turn,
|
|
step,
|
|
provider,
|
|
mode: policy.mode,
|
|
policyKey,
|
|
retry,
|
|
delayMs,
|
|
failure,
|
|
}
|
|
agent.session.append('llm/retry', eventData)
|
|
if (!await cancellableDelay(delayMs, fusedSignal)) return
|
|
return { kind: 'retry' }
|
|
}
|
|
|
|
async function recover(
|
|
agent: Agent,
|
|
context: RequestFailureContext,
|
|
signal: AbortSignal,
|
|
next: () => Promise<RequestErrorAction>,
|
|
): Promise<RequestErrorAction> {
|
|
const { turn, step, provider, failure, retryPolicy: policy } = context
|
|
if (policy === undefined) return next()
|
|
if (policy.mode === 'always') {
|
|
if (signal.aborted || lifetime.signal.aborted) return
|
|
const fusedSignal = AbortSignal.any([signal, lifetime.signal])
|
|
// The loop and plugin lifetime stay open until delegated recovery settles.
|
|
// An abort then wins before the decision or fallback can mutate later state.
|
|
const downstream = await settleDownstream(next)
|
|
if (fusedSignal.aborted) return
|
|
if (downstream.type === 'error') {
|
|
ctx.logger.warn(
|
|
`llm-retry: provider "${provider}" always policy ignored a downstream recovery failure: %o`,
|
|
downstream.error,
|
|
)
|
|
}
|
|
if (downstream.type === 'decision' && downstream.decision?.kind === 'retry') {
|
|
return downstream.decision
|
|
}
|
|
} else if (!policy.retryableCodes.includes(failure.code)) {
|
|
return next()
|
|
}
|
|
|
|
const policyKey = retryPolicyKey(policy)
|
|
const priorPolicyRetry = agent.session.events.findLast((event): event is SessionEvent<'llm/retry'> =>
|
|
event.type === 'llm/retry'
|
|
&& event.data.turn === turn
|
|
&& event.data.step === step
|
|
&& event.data.provider === provider
|
|
&& event.data.policyKey === policyKey,
|
|
)
|
|
const previousRetry = priorPolicyRetry?.data.retry ?? 0
|
|
if (policy.mode === 'normal' && previousRetry >= policy.maxRetries) return next()
|
|
const retry = previousRetry + 1
|
|
let delayMs: number
|
|
if (failure.providerRetryAfterMs !== undefined
|
|
&& Number.isFinite(failure.providerRetryAfterMs)
|
|
&& failure.providerRetryAfterMs > 0) {
|
|
if (failure.providerRetryAfterMs > policy.maxDelayMs) {
|
|
if (policy.mode === 'normal') return next()
|
|
delayMs = localDelay(policy, retry, random)
|
|
} else {
|
|
delayMs = failure.providerRetryAfterMs
|
|
}
|
|
} else {
|
|
delayMs = localDelay(policy, retry, random)
|
|
}
|
|
|
|
return backoff(agent, turn, step, failure, provider, policy, policyKey, retry, delayMs, signal)
|
|
}
|
|
|
|
const disposeListener = ctx.on('agent/request-error', (
|
|
agent: Agent,
|
|
context: RequestFailureContext,
|
|
signal: AbortSignal,
|
|
next: () => Promise<RequestErrorAction>,
|
|
) => {
|
|
// A waterfall may have captured this callback before its registration was
|
|
// removed. Lifetime cancellation must prevent that stale callback from
|
|
// entering a downstream policy after disposal.
|
|
if (lifetime.signal.aborted) return Promise.resolve<RequestErrorAction>(undefined)
|
|
return track(recover(agent, context, signal, next))
|
|
})
|
|
|
|
ctx.effect(() => async () => {
|
|
disposeListener()
|
|
lifetime.abort(new Error('llm-retry plugin disposed'))
|
|
await Promise.allSettled([...active])
|
|
}, 'llm-retry: abort and drain active recovery')
|
|
}
|