fix(llm): isolate retry policy histories
This commit is contained in:
@@ -11,6 +11,7 @@ import type { Agent, RequestError, RequestErrorDecision } from '@deepseek-ai/dsh
|
||||
import type { LlmFailure, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'
|
||||
import type { SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import { providerForClosedStep } from './history.ts'
|
||||
import { retryPolicyKey } from './policy-key.ts'
|
||||
|
||||
declare module '@deepseek-ai/dsh-session' {
|
||||
interface SessionEventMap {
|
||||
@@ -20,6 +21,7 @@ declare module '@deepseek-ai/dsh-session' {
|
||||
step: number
|
||||
provider: string
|
||||
mode: 'normal'
|
||||
policyKey: string
|
||||
retry: number
|
||||
maxRetries: number
|
||||
delayMs: number
|
||||
@@ -29,6 +31,7 @@ declare module '@deepseek-ai/dsh-session' {
|
||||
step: number
|
||||
provider: string
|
||||
mode: 'always'
|
||||
policyKey: string
|
||||
retry: number
|
||||
delayMs: number
|
||||
failure: LlmFailure
|
||||
@@ -121,6 +124,7 @@ export function apply(ctx: Context, config: Config = {}, internals: RetryInterna
|
||||
failure: LlmFailure,
|
||||
provider: string,
|
||||
policy: ResolvedRetryPolicy,
|
||||
policyKey: string,
|
||||
retry: number,
|
||||
delayMs: number,
|
||||
signal: AbortSignal,
|
||||
@@ -133,6 +137,7 @@ export function apply(ctx: Context, config: Config = {}, internals: RetryInterna
|
||||
step,
|
||||
provider,
|
||||
mode: policy.mode,
|
||||
policyKey,
|
||||
retry,
|
||||
maxRetries: policy.maxRetries,
|
||||
delayMs,
|
||||
@@ -143,6 +148,7 @@ export function apply(ctx: Context, config: Config = {}, internals: RetryInterna
|
||||
step,
|
||||
provider,
|
||||
mode: policy.mode,
|
||||
policyKey,
|
||||
retry,
|
||||
delayMs,
|
||||
failure,
|
||||
@@ -192,6 +198,7 @@ export function apply(ctx: Context, config: Config = {}, internals: RetryInterna
|
||||
return next()
|
||||
}
|
||||
|
||||
const policyKey = retryPolicyKey(policy)
|
||||
const firstPriorStep = step - priorFailures.length
|
||||
const priorPolicyRetry = agent.session.events.findLast((event): event is SessionEvent<'llm/retry'> =>
|
||||
event.type === 'llm/retry'
|
||||
@@ -199,7 +206,7 @@ export function apply(ctx: Context, config: Config = {}, internals: RetryInterna
|
||||
&& event.data.step >= firstPriorStep
|
||||
&& event.data.step < step
|
||||
&& event.data.provider === provider
|
||||
&& event.data.mode === policy.mode,
|
||||
&& event.data.policyKey === policyKey,
|
||||
)
|
||||
const previousRetry = priorPolicyRetry?.data.retry ?? 0
|
||||
if (policy.mode === 'normal' && previousRetry >= policy.maxRetries) return next()
|
||||
@@ -218,7 +225,7 @@ export function apply(ctx: Context, config: Config = {}, internals: RetryInterna
|
||||
delayMs = localDelay(policy, retry, random)
|
||||
}
|
||||
|
||||
return backoff(agent, turn, step, failure, provider, policy, retry, delayMs, signal)
|
||||
return backoff(agent, turn, step, failure, provider, policy, policyKey, retry, delayMs, signal)
|
||||
}
|
||||
|
||||
const disposeListener = ctx.on('agent/request-error', (
|
||||
|
||||
@@ -2,9 +2,9 @@
|
||||
|
||||
import type { Context } from 'cordis'
|
||||
import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
|
||||
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
|
||||
import type { InvariantFailure, InvariantInstaller } from '@deepseek-ai/dsh-invariants'
|
||||
import { providerForClosedStep } from './history.ts'
|
||||
import { parseRetryPolicyKey } from './policy-key.ts'
|
||||
import type {} from './index.ts'
|
||||
|
||||
const PACKAGE_NAME = '@deepseek-ai/dsh-llm-retry'
|
||||
@@ -20,30 +20,46 @@ function validateRetry(
|
||||
event: SessionEvent<'llm/retry'>,
|
||||
fail: InvariantFailure,
|
||||
): void {
|
||||
const { turn, step, provider, mode, retry, delayMs } = event.data
|
||||
const { turn, step, provider, mode, policyKey, retry, delayMs } = event.data
|
||||
if (!Number.isSafeInteger(retry) || retry < 1) {
|
||||
fail('llm/retry retry must be a positive safe integer')
|
||||
}
|
||||
if (typeof provider !== 'string' || provider.length === 0) {
|
||||
fail('llm/retry provider must be non-empty string')
|
||||
}
|
||||
const keyedPolicy = parseRetryPolicyKey(policyKey)
|
||||
if (keyedPolicy === undefined) {
|
||||
fail('llm/retry policyKey must encode a canonical resolved policy')
|
||||
}
|
||||
switch (mode) {
|
||||
case 'normal': {
|
||||
const { maxRetries } = event.data
|
||||
if (!Number.isSafeInteger(maxRetries) || maxRetries < 1 || retry > maxRetries) {
|
||||
fail(`llm/retry retry ${retry} must not exceed a positive safe maxRetries ${maxRetries}`)
|
||||
}
|
||||
if (keyedPolicy.mode !== 'normal') {
|
||||
fail(`llm/retry mode normal must match policyKey mode ${keyedPolicy.mode}`)
|
||||
}
|
||||
if (keyedPolicy.maxRetries !== maxRetries) {
|
||||
fail(`llm/retry maxRetries ${maxRetries} must match policyKey`)
|
||||
}
|
||||
if (!keyedPolicy.retryableCodes.includes(event.data.failure.code)) {
|
||||
fail(`llm/retry failure code ${event.data.failure.code} must be eligible under policyKey`)
|
||||
}
|
||||
break
|
||||
}
|
||||
case 'always':
|
||||
if (keyedPolicy.mode !== 'always') {
|
||||
fail(`llm/retry mode always must match policyKey mode ${keyedPolicy.mode}`)
|
||||
}
|
||||
if ('maxRetries' in event.data) fail('llm/retry always mode must omit maxRetries')
|
||||
break
|
||||
default:
|
||||
fail(`llm/retry mode must be normal or always, got ${String(mode)}`)
|
||||
}
|
||||
if (typeof delayMs !== 'number' || !Number.isFinite(delayMs)
|
||||
|| delayMs < 0 || delayMs > MAX_TIMER_DELAY_MS) {
|
||||
fail(`llm/retry delayMs must be a finite number within 0..${MAX_TIMER_DELAY_MS}`)
|
||||
|| delayMs < 0 || delayMs > keyedPolicy.maxDelayMs) {
|
||||
fail(`llm/retry delayMs must be a finite number within policyKey range 0..${keyedPolicy.maxDelayMs}`)
|
||||
}
|
||||
|
||||
const turnStartIndex = history.findLastIndex(prior =>
|
||||
@@ -86,7 +102,7 @@ function validateRetry(
|
||||
index > lastSuccessIndex
|
||||
&& prior.type === 'llm/retry'
|
||||
&& prior.data.provider === provider
|
||||
&& prior.data.mode === mode
|
||||
&& prior.data.policyKey === policyKey
|
||||
))
|
||||
const expectedRetry = (priorPolicyRetry?.data.retry ?? 0) + 1
|
||||
if (retry !== expectedRetry) {
|
||||
|
||||
100
packages/llm/llm-retry/src/policy-key.ts
Normal file
100
packages/llm/llm-retry/src/policy-key.ts
Normal file
@@ -0,0 +1,100 @@
|
||||
/** Canonical durable identity for resolved retry policies. @module @deepseek-ai/dsh-llm-retry/policy-key */
|
||||
|
||||
import type { ResolvedRetryBackoff, ResolvedRetryPolicy } from '@deepseek-ai/dsh-llm'
|
||||
import { MAX_TIMER_DELAY_MS } from '@deepseek-ai/dsh-timeout'
|
||||
|
||||
function parseBackoff(
|
||||
tuple: readonly unknown[],
|
||||
offset: number,
|
||||
): ResolvedRetryBackoff | undefined {
|
||||
const initialDelayMs = tuple[offset]
|
||||
const maxDelayMs = tuple[offset + 1]
|
||||
const jitterRatio = tuple[offset + 2]
|
||||
if (typeof initialDelayMs !== 'number' || !Number.isFinite(initialDelayMs)
|
||||
|| initialDelayMs <= 0 || initialDelayMs > MAX_TIMER_DELAY_MS
|
||||
|| typeof maxDelayMs !== 'number' || !Number.isFinite(maxDelayMs)
|
||||
|| maxDelayMs <= 0 || maxDelayMs > MAX_TIMER_DELAY_MS
|
||||
|| initialDelayMs > maxDelayMs
|
||||
|| typeof jitterRatio !== 'number' || !Number.isFinite(jitterRatio)
|
||||
|| jitterRatio < 0 || jitterRatio > 1) {
|
||||
return undefined
|
||||
}
|
||||
return { initialDelayMs, maxDelayMs, jitterRatio }
|
||||
}
|
||||
|
||||
/**
|
||||
* Derive the canonical durable key for one fully resolved provider policy.
|
||||
* Retryable-code order is normalized because eligibility uses set membership.
|
||||
* @param policy - immutable policy captured from the serving registration.
|
||||
* @returns canonical JSON tuple containing every behavior-affecting field.
|
||||
*/
|
||||
export function retryPolicyKey(policy: ResolvedRetryPolicy): string {
|
||||
if (policy.mode === 'always') {
|
||||
return JSON.stringify([
|
||||
policy.mode,
|
||||
policy.initialDelayMs,
|
||||
policy.maxDelayMs,
|
||||
policy.jitterRatio,
|
||||
])
|
||||
}
|
||||
return JSON.stringify([
|
||||
policy.mode,
|
||||
policy.maxRetries,
|
||||
[...policy.retryableCodes].sort(),
|
||||
policy.initialDelayMs,
|
||||
policy.maxDelayMs,
|
||||
policy.jitterRatio,
|
||||
])
|
||||
}
|
||||
|
||||
/**
|
||||
* Parse a producer-canonical policy key from durable input.
|
||||
* @param value - untrusted persisted event field.
|
||||
* @returns the resolved policy encoded by the key, or `undefined` for any non-canonical value.
|
||||
*/
|
||||
export function parseRetryPolicyKey(value: unknown): ResolvedRetryPolicy | undefined {
|
||||
if (typeof value !== 'string' || value.length === 0) return undefined
|
||||
let tuple: unknown
|
||||
try {
|
||||
tuple = JSON.parse(value) as unknown
|
||||
} catch (_invalidPolicyKeyJson) {
|
||||
return undefined
|
||||
}
|
||||
if (!Array.isArray(tuple)) return undefined
|
||||
const items = tuple as readonly unknown[]
|
||||
const mode = items[0]
|
||||
let policy: ResolvedRetryPolicy
|
||||
switch (mode) {
|
||||
case 'always': {
|
||||
if (items.length !== 4) return undefined
|
||||
const backoff = parseBackoff(items, 1)
|
||||
if (backoff === undefined) return undefined
|
||||
policy = Object.freeze({ mode, ...backoff })
|
||||
break
|
||||
}
|
||||
case 'normal': {
|
||||
if (items.length !== 6) return undefined
|
||||
const maxRetries = items[1]
|
||||
const retryableCodes = items[2]
|
||||
const backoff = parseBackoff(items, 3)
|
||||
if (!Number.isSafeInteger(maxRetries) || (maxRetries as number) < 0
|
||||
|| !Array.isArray(retryableCodes) || retryableCodes.length === 0
|
||||
|| (retryableCodes as readonly unknown[])
|
||||
.some(code => typeof code !== 'string' || code.length === 0)
|
||||
|| new Set(retryableCodes).size !== retryableCodes.length
|
||||
|| backoff === undefined) {
|
||||
return undefined
|
||||
}
|
||||
policy = Object.freeze({
|
||||
mode,
|
||||
maxRetries: maxRetries as number,
|
||||
retryableCodes: Object.freeze(retryableCodes as string[]),
|
||||
...backoff,
|
||||
})
|
||||
break
|
||||
}
|
||||
default:
|
||||
return undefined
|
||||
}
|
||||
return retryPolicyKey(policy) === value ? policy : undefined
|
||||
}
|
||||
Reference in New Issue
Block a user