From 6faf39c52914eef92f0ee79dd860a279de9bf481 Mon Sep 17 00:00:00 2001 From: elkaix Date: Fri, 21 Aug 2026 19:22:06 -0400 Subject: [PATCH 1/6] feat(agent-core-v2): port turn resilience mechanisms from reference design Add two flag-gated Agent-scope recovery domains and a token-budget continuation service: - outputTokenRecovery: when a response ends truncated at the output token limit with no tool calls, inject a resume nudge (origin kind 'retry') and continue the same turn, capped at three attempts. - modelFallback: when step retries are exhausted on persistent retryable provider errors, switch the profile to the configured loop_control.fallback_model once per turn and retry the failed step; emits ModelFallbackSwitched plus telemetry. - turnBudget: with loop_control.turn_budget_tokens set, keep a naturally-stopping turn working toward the output-token target with continuation nudges until threshold or diminishing returns. New loop_control fields (fallback_model, turn_budget_tokens) with env binding, three experimental flags, registered telemetry events, and regenerated config/state manifests. --- .../agent-core-v2/docs/config-manifest.toml | 3 + .../agent-core-v2/docs/state-manifest.d.ts | 15 +- .../src/agent/loop/configSection.ts | 4 + .../src/agent/stepRetry/stepRetryService.ts | 8 +- .../src/agent/turnBudget/flag.ts | 17 ++ .../src/agent/turnBudget/turnBudget.ts | 20 ++ .../src/agent/turnBudget/turnBudgetService.ts | 147 +++++++++++++++ .../src/agent/turnRecovery/flag.ts | 30 +++ .../src/agent/turnRecovery/modelFallback.ts | 37 ++++ .../turnRecovery/modelFallbackService.ts | 96 ++++++++++ .../agent/turnRecovery/outputTokenRecovery.ts | 19 ++ .../outputTokenRecoveryService.ts | 99 ++++++++++ .../agent-core-v2/src/app/telemetry/events.ts | 50 +++++ packages/agent-core-v2/src/index.ts | 8 + .../test/agent/turnBudget/turnBudget.test.ts | 155 ++++++++++++++++ .../agent/turnRecovery/modelFallback.test.ts | 174 ++++++++++++++++++ .../turnRecovery/outputTokenRecovery.test.ts | 147 +++++++++++++++ 17 files changed, 1027 insertions(+), 2 deletions(-) create mode 100644 packages/agent-core-v2/src/agent/turnBudget/flag.ts create mode 100644 packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts create mode 100644 packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts create mode 100644 packages/agent-core-v2/src/agent/turnRecovery/flag.ts create mode 100644 packages/agent-core-v2/src/agent/turnRecovery/modelFallback.ts create mode 100644 packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts create mode 100644 packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts create mode 100644 packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecoveryService.ts create mode 100644 packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts create mode 100644 packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts create mode 100644 packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts diff --git a/packages/agent-core-v2/docs/config-manifest.toml b/packages/agent-core-v2/docs/config-manifest.toml index d1202a318..af0fe9e4b 100644 --- a/packages/agent-core-v2/docs/config-manifest.toml +++ b/packages/agent-core-v2/docs/config-manifest.toml @@ -198,6 +198,7 @@ extra_skill_dirs = [] # env: # max_steps_per_turn <- PYTHINKER_LOOP_MAX_STEPS_PER_TURN (custom parse) # max_attempts_per_step <- PYTHINKER_LOOP_MAX_ATTEMPTS_PER_STEP (custom parse; deprecated fallback PYTHINKER_LOOP_MAX_RETRIES_PER_STEP) +# turn_budget_tokens <- PYTHINKER_LOOP_TURN_BUDGET_TOKENS (custom parse) # ########################################################################## [loop_control] @@ -206,6 +207,8 @@ extra_skill_dirs = [] # max_ralph_iterations: integer # reserved_context_size: integer # compaction_trigger_ratio: number +# fallback_model: string +# turn_budget_tokens: integer # ########################################################################## # mcp diff --git a/packages/agent-core-v2/docs/state-manifest.d.ts b/packages/agent-core-v2/docs/state-manifest.d.ts index d8749cc8c..bc9e11bd9 100644 --- a/packages/agent-core-v2/docs/state-manifest.d.ts +++ b/packages/agent-core-v2/docs/state-manifest.d.ts @@ -27,7 +27,7 @@ // references become '(circular)', and class instances collapse to a '(ClassName)' // marker — the wire shape of an entry is the JSON projection of the type here. // -// Index (App: 0 keys · Workspace: 6 keys · Session: 18 keys · Agent: 98 keys) +// Index (App: 0 keys · Workspace: 6 keys · Session: 18 keys · Agent: 103 keys) // App // Workspace // workspaceDirs.ephemeralDirs src/workspace/workspaceDirs/workspaceDirsService.ts @@ -150,6 +150,11 @@ // toolSelect.pendingLoaded src/agent/toolSelect/toolSelectService.ts // tower src/features/tower/towerOps.ts // turn src/agent/loop/turnOps.ts +// turnBudget.continuations src/agent/turnBudget/turnBudgetService.ts +// turnBudget.lastDeltaTokens src/agent/turnBudget/turnBudgetService.ts +// turnBudget.tokensUsed src/agent/turnBudget/turnBudgetService.ts +// turnRecovery.modelFallbackUsed src/agent/turnRecovery/modelFallbackService.ts +// turnRecovery.outputTokenAttempts src/agent/turnRecovery/outputTokenRecoveryService.ts // usage src/agent/usage/usageOps.ts // usage.currentTurn src/agent/usage/usageService.ts // usage.currentTurnId src/agent/usage/usageService.ts @@ -1527,6 +1532,14 @@ export interface AgentStateSnapshot { 'toolExecutor.toolCallDupTypes': Map; // src/agent/toolSelect/toolSelectService.ts 'toolSelect.pendingLoaded': Set; + // src/agent/turnBudget/turnBudgetService.ts + 'turnBudget.continuations': number; + 'turnBudget.lastDeltaTokens': number; + 'turnBudget.tokensUsed': number; + // src/agent/turnRecovery/modelFallbackService.ts + 'turnRecovery.modelFallbackUsed': boolean; + // src/agent/turnRecovery/outputTokenRecoveryService.ts + 'turnRecovery.outputTokenAttempts': number; // src/agent/usage/usageOps.ts // replayable · durable — folds: UsageRecord 'usage': /* UsageModelState — packages/agent-core-v2/src/agent/usage/usageOps.ts */ { diff --git a/packages/agent-core-v2/src/agent/loop/configSection.ts b/packages/agent-core-v2/src/agent/loop/configSection.ts index d75d2a9b5..df91b8b38 100644 --- a/packages/agent-core-v2/src/agent/loop/configSection.ts +++ b/packages/agent-core-v2/src/agent/loop/configSection.ts @@ -8,6 +8,7 @@ export const LOOP_CONTROL_SECTION = 'loopControl'; export const LOOP_MAX_STEPS_PER_TURN_ENV = 'PYTHINKER_LOOP_MAX_STEPS_PER_TURN'; export const LOOP_MAX_ATTEMPTS_PER_STEP_ENV = 'PYTHINKER_LOOP_MAX_ATTEMPTS_PER_STEP'; +export const LOOP_TURN_BUDGET_TOKENS_ENV = 'PYTHINKER_LOOP_TURN_BUDGET_TOKENS'; /** Deprecated former name of {@link LOOP_MAX_ATTEMPTS_PER_STEP_ENV}. */ export const LOOP_MAX_RETRIES_PER_STEP_ENV = 'PYTHINKER_LOOP_MAX_RETRIES_PER_STEP'; @@ -17,6 +18,8 @@ export const LoopControlSchema = z.object({ maxRalphIterations: z.number().int().min(-1).optional(), reservedContextSize: z.number().int().min(0).optional(), compactionTriggerRatio: z.number().min(0.5).max(0.99).optional(), + fallbackModel: z.string().min(1).optional(), + turnBudgetTokens: z.number().int().min(0).optional(), }); export type LoopControl = z.infer; @@ -35,6 +38,7 @@ export const loopControlEnvBindings: EnvBindings = envBindings(Loop deprecatedEnv: LOOP_MAX_RETRIES_PER_STEP_ENV, parse: parseNonNegativeInt, }, + turnBudgetTokens: { env: LOOP_TURN_BUDGET_TOKENS_ENV, parse: parseNonNegativeInt }, }); export const stripLoopControlEnv = stripEnvBoundFields(loopControlEnvBindings); diff --git a/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts b/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts index 8ab985ddf..b6537b040 100644 --- a/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts +++ b/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts @@ -23,6 +23,7 @@ import { LOOP_CONTROL_SECTION, type LoopControl } from '#/agent/loop/configSecti import { TurnStarted } from '#/agent/loop/turnEvents'; import { IAgentStateService } from '#/agent/state/agentState'; import { IEventDispatcher } from '#/state/eventDispatcher'; +import { IAgentModelFallbackService } from '#/agent/turnRecovery/modelFallback'; import { IAgentStepRetryService } from './stepRetry'; @@ -63,6 +64,7 @@ export class AgentStepRetryService extends Disposable implements IAgentStepRetry @IEventBus private readonly eventBus: IEventBus, @IEventDispatcher private readonly dispatcher: IEventDispatcher, @IAgentStateService private readonly states: IAgentStateService, + @IAgentModelFallbackService private readonly modelFallback: IAgentModelFallbackService, ) { super(); this.states.contributeState(stepRetryLastFailedDriverIdKey); @@ -121,7 +123,11 @@ export class AgentStepRetryService extends Disposable implements IAgentStepRetry ); if (this.failedAttempts >= maxAttempts) { this.resetAttempts(); - return false; + const switched = await this.modelFallback.tryFallbackSwitch(context); + if (!switched) return false; + if (context.currentStep?.signal.aborted === true) return false; + context.retry(driver, { at: 'head' }); + return true; } const error = unwrapErrorCause(context.error); diff --git a/packages/agent-core-v2/src/agent/turnBudget/flag.ts b/packages/agent-core-v2/src/agent/turnBudget/flag.ts new file mode 100644 index 000000000..80aa8600d --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnBudget/flag.ts @@ -0,0 +1,17 @@ +import { type FlagDefinitionInput, registerFlagDefinition } from '#/app/flag/flagRegistry'; + +export const TURN_BUDGET_CONTINUATION_FLAG_ID = 'turn-budget-continuation'; +export const TURN_BUDGET_CONTINUATION_FLAG_ENV = + 'PYTHINKER_CODE_EXPERIMENTAL_TURN_BUDGET_CONTINUATION'; + +export const turnBudgetContinuationFlag: FlagDefinitionInput = { + id: TURN_BUDGET_CONTINUATION_FLAG_ID, + title: 'Turn budget continuation', + description: + 'When loopControl.turn_budget_tokens is set, keep a naturally-stopping turn working toward the output-token target with continuation nudges until the threshold is reached or progress diminishes.', + env: TURN_BUDGET_CONTINUATION_FLAG_ENV, + default: false, + surface: 'core', +}; + +registerFlagDefinition(turnBudgetContinuationFlag); diff --git a/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts b/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts new file mode 100644 index 000000000..7a3be43bb --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts @@ -0,0 +1,20 @@ +import { createDecorator } from '#/_base/di/instantiation'; + +/** + * Continues a turn toward a configured output-token target by injecting + * continuation nudges while progress holds, stopping on diminishing returns. + */ +export interface IAgentTurnBudgetService { + readonly _serviceBrand: undefined; +} + +export const IAgentTurnBudgetService = + createDecorator('agentTurnBudgetService'); + +export const TURN_BUDGET_COMPLETION_THRESHOLD = 0.9; +export const TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS = 500; +export const TURN_BUDGET_MAX_DIMINISHING_CONTINUATIONS = 3; + +export function turnBudgetNudgeText(pct: number, used: number, budget: number): string { + return `Stopped at ${pct}% of token target (${used} / ${budget}). Keep working - do not summarize.`; +} diff --git a/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts b/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts new file mode 100644 index 000000000..0e60f9125 --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts @@ -0,0 +1,147 @@ +import { Disposable } from '#/_base/di/lifecycle'; +import { LifecycleScope } from '#/app/scopes'; +import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; +import { defineState } from '#/state/state'; +import { createUserMessage } from '#/kosong/contract/message'; +import type { ContextMessage } from '#/agent/contextMemory/types'; +import type { AfterStepContext } from '#/agent/loop/loop'; +import { IAgentLoopService } from '#/agent/loop/loop'; +import { TurnStarted } from '#/agent/loop/turnEvents'; +import { StepRequest } from '#/agent/loop/stepRequest'; +import { LOOP_CONTROL_SECTION, type LoopControl } from '#/agent/loop/configSection'; +import { IConfigService } from '#/app/config/config'; +import { IEventBus } from '#/app/event/eventBus'; +import { IFlagService } from '#/app/flag/flag'; +import { ITelemetryService } from '#/app/telemetry/telemetry'; +import { IAgentStateService } from '#/agent/state/agentState'; +import { TURN_BUDGET_CONTINUATION_FLAG_ID } from './flag'; +import { + IAgentTurnBudgetService, + TURN_BUDGET_COMPLETION_THRESHOLD, + TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS, + TURN_BUDGET_MAX_DIMINISHING_CONTINUATIONS, + turnBudgetNudgeText, +} from './turnBudget'; + +export const turnBudgetTokensUsedKey = defineState('turnBudget.tokensUsed', () => 0); +export const turnBudgetContinuationsKey = defineState('turnBudget.continuations', () => 0); +export const turnBudgetLastDeltaTokensKey = defineState( + 'turnBudget.lastDeltaTokens', + () => 0, +); + +class TurnBudgetContinuationRequest extends StepRequest { + readonly kind = 'turn-budget-continuation'; + + constructor(private readonly message: ContextMessage) { + super(); + } + + override resolveContextMessages(): readonly ContextMessage[] { + return [this.message]; + } +} + +export class AgentTurnBudgetService extends Disposable implements IAgentTurnBudgetService { + declare readonly _serviceBrand: undefined; + + constructor( + @IAgentLoopService private readonly loopService: IAgentLoopService, + @IFlagService private readonly flags: IFlagService, + @IConfigService private readonly config: IConfigService, + @IEventBus private readonly eventBus: IEventBus, + @ITelemetryService private readonly telemetry: ITelemetryService, + @IAgentStateService private readonly states: IAgentStateService, + ) { + super(); + this.states.contributeState(turnBudgetTokensUsedKey); + this.states.contributeState(turnBudgetContinuationsKey); + this.states.contributeState(turnBudgetLastDeltaTokensKey); + this._register(this.eventBus.subscribe(TurnStarted, () => this.reset())); + this._register( + this.loopService.hooks.onDidFinishStep.register('turn-budget', async (context, next) => { + await next(); + this.maybeContinue(context); + }), + ); + } + + private get tokensUsed(): number { + return this.states.get(turnBudgetTokensUsedKey); + } + + private set tokensUsed(value: number) { + this.states.set(turnBudgetTokensUsedKey, value); + } + + private get continuations(): number { + return this.states.get(turnBudgetContinuationsKey); + } + + private set continuations(value: number) { + this.states.set(turnBudgetContinuationsKey, value); + } + + private get lastDeltaTokens(): number { + return this.states.get(turnBudgetLastDeltaTokensKey); + } + + private set lastDeltaTokens(value: number) { + this.states.set(turnBudgetLastDeltaTokensKey, value); + } + + private reset(): void { + this.tokensUsed = 0; + this.continuations = 0; + this.lastDeltaTokens = 0; + } + + private budgetTokens(): number { + return this.config.get(LOOP_CONTROL_SECTION)?.turnBudgetTokens ?? 0; + } + + private maybeContinue(context: AfterStepContext): void { + if (!this.flags.enabled(TURN_BUDGET_CONTINUATION_FLAG_ID)) return; + if (context.stopTurn || context.signal.aborted) return; + const budget = this.budgetTokens(); + if (budget <= 0) return; + + const delta = context.usage.output; + const used = this.tokensUsed + delta; + this.lastDeltaTokens = delta; + this.tokensUsed = used; + + if (context.finishReason === 'tool_calls') return; + if (context.finishReason !== 'completed') return; + if (this.loopService.status().hasPendingRequests) return; + + const diminishing = + this.continuations >= TURN_BUDGET_MAX_DIMINISHING_CONTINUATIONS && + delta < TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS && + this.lastDeltaTokens < TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS; + if (diminishing) return; + if (used >= budget * TURN_BUDGET_COMPLETION_THRESHOLD) return; + + const pct = Math.round((used / budget) * 100); + this.continuations += 1; + this.telemetry.track2('budget_continuation', { + turn_id: context.turnId, + continuation_count: this.continuations, + tokens_used: used, + budget_tokens: budget, + }); + const message: ContextMessage = { + ...createUserMessage(turnBudgetNudgeText(pct, used, budget)), + origin: { kind: 'retry', trigger: 'token_budget' }, + }; + this.loopService.enqueue(new TurnBudgetContinuationRequest(message)); + } +} + +registerScopedService( + LifecycleScope.Agent, + IAgentTurnBudgetService, + AgentTurnBudgetService, + ScopeActivation.OnScopeCreated, + 'turnBudget', +); diff --git a/packages/agent-core-v2/src/agent/turnRecovery/flag.ts b/packages/agent-core-v2/src/agent/turnRecovery/flag.ts new file mode 100644 index 000000000..51f61f306 --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnRecovery/flag.ts @@ -0,0 +1,30 @@ +import { type FlagDefinitionInput, registerFlagDefinition } from '#/app/flag/flagRegistry'; + +export const OUTPUT_TOKEN_RECOVERY_FLAG_ID = 'output-token-recovery'; +export const OUTPUT_TOKEN_RECOVERY_FLAG_ENV = 'PYTHINKER_CODE_EXPERIMENTAL_OUTPUT_TOKEN_RECOVERY'; + +export const MODEL_FALLBACK_FLAG_ID = 'model-fallback'; +export const MODEL_FALLBACK_FLAG_ENV = 'PYTHINKER_CODE_EXPERIMENTAL_MODEL_FALLBACK'; + +export const outputTokenRecoveryFlag: FlagDefinitionInput = { + id: OUTPUT_TOKEN_RECOVERY_FLAG_ID, + title: 'Output token recovery', + description: + 'When a model response ends truncated at the output token limit with no tool calls, inject a resume nudge and continue the same turn instead of ending it truncated. Caps recoveries per turn.', + env: OUTPUT_TOKEN_RECOVERY_FLAG_ENV, + default: false, + surface: 'core', +}; + +export const modelFallbackFlag: FlagDefinitionInput = { + id: MODEL_FALLBACK_FLAG_ID, + title: 'Model fallback', + description: + 'When step retries are exhausted on persistent retryable provider errors, switch the agent to the configured loopControl.fallback_model once per turn and retry the failed step there.', + env: MODEL_FALLBACK_FLAG_ENV, + default: false, + surface: 'core', +}; + +registerFlagDefinition(outputTokenRecoveryFlag); +registerFlagDefinition(modelFallbackFlag); diff --git a/packages/agent-core-v2/src/agent/turnRecovery/modelFallback.ts b/packages/agent-core-v2/src/agent/turnRecovery/modelFallback.ts new file mode 100644 index 000000000..20b7d627b --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnRecovery/modelFallback.ts @@ -0,0 +1,37 @@ +/* oxlint-disable typescript-eslint/no-unsafe-declaration-merging, eslint-plugin-import/namespace -- Event2 class+payload-interface declaration merging is the sanctioned event-declaration idiom. */ +import { createDecorator } from '#/_base/di/instantiation'; +import { Event2 } from '#/app/event/event2'; +import type { LoopErrorContext } from '#/agent/loop/loop'; + +/** + * Switches the agent to the configured fallback model when step retries are + * exhausted on persistent retryable provider errors, so the retrying layer can + * resend the failed step on the fallback. + */ +export interface IAgentModelFallbackService { + readonly _serviceBrand: undefined; + + /** + * Switches the agent profile to `loopControl.fallback_model` when allowed + * (flag on, model configured and different from the current one, not yet + * used this turn). Returns true when the switch happened and the caller + * should retry the failed driver. + */ + tryFallbackSwitch(context: LoopErrorContext): Promise; +} + +export const IAgentModelFallbackService = + createDecorator('agentModelFallbackService'); + +export interface ModelFallbackSwitchedPayload { + readonly turnId: number; + readonly step?: number; + readonly fromModel: string; + readonly toModel: string; +} + +export class ModelFallbackSwitched extends Event2 { + static override readonly type = 'turn.model_fallback.switched'; + static override readonly observable = true; +} +export interface ModelFallbackSwitched extends ModelFallbackSwitchedPayload {} diff --git a/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts b/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts new file mode 100644 index 000000000..088cebe31 --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts @@ -0,0 +1,96 @@ +import { Disposable } from '#/_base/di/lifecycle'; +import { LifecycleScope } from '#/app/scopes'; +import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; +import { defineState } from '#/state/state'; +import { isRetryableGenerateError } from '#/kosong/contract/errors'; +import { unwrapErrorCause } from '#/errors'; +import type { LoopErrorContext } from '#/agent/loop/loop'; +import { TurnStarted } from '#/agent/loop/turnEvents'; +import { LOOP_CONTROL_SECTION, type LoopControl } from '#/agent/loop/configSection'; +import { IConfigService } from '#/app/config/config'; +import { IEventBus } from '#/app/event/eventBus'; +import { IFlagService } from '#/app/flag/flag'; +import { ITelemetryService } from '#/app/telemetry/telemetry'; +import { IAgentProfileService } from '#/agent/profile/profile'; +import { IEventDispatcher } from '#/state/eventDispatcher'; +import { IAgentStateService } from '#/agent/state/agentState'; +import { MODEL_FALLBACK_FLAG_ID } from './flag'; +import { IAgentModelFallbackService, ModelFallbackSwitched } from './modelFallback'; + +export const modelFallbackUsedKey = defineState( + 'turnRecovery.modelFallbackUsed', + () => false, +); + +export class AgentModelFallbackService extends Disposable implements IAgentModelFallbackService { + declare readonly _serviceBrand: undefined; + + constructor( + @IFlagService private readonly flags: IFlagService, + @IConfigService private readonly config: IConfigService, + @IAgentProfileService private readonly profile: IAgentProfileService, + @IEventBus private readonly eventBus: IEventBus, + @IEventDispatcher private readonly dispatcher: IEventDispatcher, + @ITelemetryService private readonly telemetry: ITelemetryService, + @IAgentStateService private readonly states: IAgentStateService, + ) { + super(); + this.states.contributeState(modelFallbackUsedKey); + this._register(this.eventBus.subscribe(TurnStarted, () => this.reset())); + } + + private get used(): boolean { + return this.states.get(modelFallbackUsedKey); + } + + private set used(value: boolean) { + this.states.set(modelFallbackUsedKey, value); + } + + private reset(): void { + this.used = false; + } + + private fallbackModel(): string | undefined { + return this.config.get(LOOP_CONTROL_SECTION)?.fallbackModel; + } + + async tryFallbackSwitch(context: LoopErrorContext): Promise { + if (!this.flags.enabled(MODEL_FALLBACK_FLAG_ID)) return false; + const target = this.fallbackModel(); + if (target === undefined || target.length === 0) return false; + if (this.used || target === this.profile.getModel()) return false; + if (!isRetryableGenerateError(unwrapErrorCause(context.error))) return false; + + const fromModel = this.profile.getModel(); + this.used = true; + try { + await this.profile.setModel(target); + } catch { + this.used = false; + return false; + } + void this.dispatcher.dispatch( + new ModelFallbackSwitched({ + turnId: context.turnId, + step: context.step, + fromModel, + toModel: target, + }), + ); + this.telemetry.track2('model_fallback_triggered', { + turn_id: context.turnId, + from_model: fromModel, + to_model: target, + }); + return true; + } +} + +registerScopedService( + LifecycleScope.Agent, + IAgentModelFallbackService, + AgentModelFallbackService, + ScopeActivation.OnScopeCreated, + 'modelFallback', +); diff --git a/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts b/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts new file mode 100644 index 000000000..9ee043c62 --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts @@ -0,0 +1,19 @@ +import { createDecorator } from '#/_base/di/instantiation'; + +/** + * Recovers turns whose model response ended truncated at the output token + * limit without tool calls, by injecting a resume nudge and continuing the + * turn. Ported from the reference design's max_output_tokens recovery ladder. + */ +export interface IAgentOutputTokenRecoveryService { + readonly _serviceBrand: undefined; +} + +export const IAgentOutputTokenRecoveryService = + createDecorator('agentOutputTokenRecoveryService'); + +export const MAX_OUTPUT_TOKEN_RECOVERY_ATTEMPTS = 3; + +export const OUTPUT_TOKEN_RECOVERY_NUDGE = + 'Output token limit hit. Resume directly - no apology, no recap of what you were doing. ' + + 'Pick up mid-thought if that is where the cut happened. Break remaining work into smaller pieces.'; diff --git a/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecoveryService.ts b/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecoveryService.ts new file mode 100644 index 000000000..38aa9bebb --- /dev/null +++ b/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecoveryService.ts @@ -0,0 +1,99 @@ +import { Disposable } from '#/_base/di/lifecycle'; +import { LifecycleScope } from '#/app/scopes'; +import { ScopeActivation, registerScopedService } from '#/_base/di/scope'; +import { defineState } from '#/state/state'; +import { createUserMessage } from '#/kosong/contract/message'; +import type { ContextMessage } from '#/agent/contextMemory/types'; +import type { AfterStepContext } from '#/agent/loop/loop'; +import { IAgentLoopService } from '#/agent/loop/loop'; +import { TurnStarted } from '#/agent/loop/turnEvents'; +import { StepRequest } from '#/agent/loop/stepRequest'; +import { IEventBus } from '#/app/event/eventBus'; +import { IFlagService } from '#/app/flag/flag'; +import { ITelemetryService } from '#/app/telemetry/telemetry'; +import { IAgentStateService } from '#/agent/state/agentState'; +import { OUTPUT_TOKEN_RECOVERY_FLAG_ID } from './flag'; +import { + IAgentOutputTokenRecoveryService, + MAX_OUTPUT_TOKEN_RECOVERY_ATTEMPTS, + OUTPUT_TOKEN_RECOVERY_NUDGE, +} from './outputTokenRecovery'; + +export const outputTokenRecoveryAttemptsKey = defineState( + 'turnRecovery.outputTokenAttempts', + () => 0, +); + +class OutputTokenRecoveryRequest extends StepRequest { + readonly kind = 'output-token-recovery'; + + constructor(private readonly message: ContextMessage) { + super(); + } + + override resolveContextMessages(): readonly ContextMessage[] { + return [this.message]; + } +} + +export class AgentOutputTokenRecoveryService extends Disposable implements IAgentOutputTokenRecoveryService { + declare readonly _serviceBrand: undefined; + + constructor( + @IAgentLoopService private readonly loopService: IAgentLoopService, + @IFlagService private readonly flags: IFlagService, + @IEventBus private readonly eventBus: IEventBus, + @ITelemetryService private readonly telemetry: ITelemetryService, + @IAgentStateService private readonly states: IAgentStateService, + ) { + super(); + this.states.contributeState(outputTokenRecoveryAttemptsKey); + this._register(this.eventBus.subscribe(TurnStarted, () => this.reset())); + this._register( + this.loopService.hooks.onDidFinishStep.register('output-token-recovery', async (context, next) => { + await next(); + this.maybeContinue(context); + }), + ); + } + + private get attempts(): number { + return this.states.get(outputTokenRecoveryAttemptsKey); + } + + private set attempts(value: number) { + this.states.set(outputTokenRecoveryAttemptsKey, value); + } + + private reset(): void { + this.attempts = 0; + } + + private maybeContinue(context: AfterStepContext): void { + if (!this.flags.enabled(OUTPUT_TOKEN_RECOVERY_FLAG_ID)) return; + if (context.stopTurn || context.signal.aborted) return; + if (context.finishReason !== 'truncated') return; + if (this.loopService.status().hasPendingRequests) return; + const attempt = this.attempts + 1; + if (attempt > MAX_OUTPUT_TOKEN_RECOVERY_ATTEMPTS) return; + this.attempts = attempt; + this.telemetry.track2('output_token_recovery', { + turn_id: context.turnId, + attempt, + max_attempts: MAX_OUTPUT_TOKEN_RECOVERY_ATTEMPTS, + }); + const message: ContextMessage = { + ...createUserMessage(OUTPUT_TOKEN_RECOVERY_NUDGE), + origin: { kind: 'retry', trigger: 'max_output_tokens' }, + }; + this.loopService.enqueue(new OutputTokenRecoveryRequest(message)); + } +} + +registerScopedService( + LifecycleScope.Agent, + IAgentOutputTokenRecoveryService, + AgentOutputTokenRecoveryService, + ScopeActivation.OnScopeCreated, + 'outputTokenRecovery', +); diff --git a/packages/agent-core-v2/src/app/telemetry/events.ts b/packages/agent-core-v2/src/app/telemetry/events.ts index f9a4b5f74..fbe972bfc 100644 --- a/packages/agent-core-v2/src/app/telemetry/events.ts +++ b/packages/agent-core-v2/src/app/telemetry/events.ts @@ -76,6 +76,25 @@ export interface TurnEndedEvent { trace_id?: string; } +export interface OutputTokenRecoveryEvent { + turn_id: number; + attempt: number; + max_attempts: number; +} + +export interface ModelFallbackEvent { + turn_id: number; + from_model: string; + to_model: string; +} + +export interface BudgetContinuationEvent { + turn_id: number; + continuation_count: number; + tokens_used: number; + budget_tokens: number; +} + export type ToolCallOutcome = 'success' | 'error' | 'cancelled'; export interface ToolCallEvent { @@ -473,6 +492,37 @@ export const telemetryEventDefinitions = { 'Trace id of the most recent LLM request in this turn; absent for non-Pythinker protocols', }, }), + output_token_recovery: defineAgentTelemetryEvent({ + owner: 'pythinker-code', + comment: + 'A truncated response (output token limit, no tool calls) is recovered with a resume nudge and the turn continues.', + properties: { + turn_id: 'Per-agent turn index (main or subagent); pair with agent_id to locate a turn within a session', + attempt: 'Recovery attempt index within the turn, starting at 1', + max_attempts: 'Configured cap on recovery attempts per turn', + }, + }), + model_fallback_triggered: defineAgentTelemetryEvent({ + owner: 'pythinker-code', + comment: + 'Step retries were exhausted on a persistent retryable provider error and the agent switched to the configured fallback model.', + properties: { + turn_id: 'Per-agent turn index (main or subagent); pair with agent_id to locate a turn within a session', + from_model: 'Model alias that kept failing before the switch', + to_model: 'Fallback model alias the agent switched to', + }, + }), + budget_continuation: defineAgentTelemetryEvent({ + owner: 'pythinker-code', + comment: + 'A naturally-stopping turn is continued toward the configured output-token target with a continuation nudge.', + properties: { + turn_id: 'Per-agent turn index (main or subagent); pair with agent_id to locate a turn within a session', + continuation_count: 'Continuation index within the turn, starting at 1', + tokens_used: 'Cumulative output tokens consumed by the turn so far', + budget_tokens: 'Configured output-token target for the turn', + }, + }), tool_call: defineAgentTelemetryEvent({ owner: 'pythinker-code', comment: 'A tool call finishes execution.', diff --git a/packages/agent-core-v2/src/index.ts b/packages/agent-core-v2/src/index.ts index 31c0f40ac..cb685c7ad 100644 --- a/packages/agent-core-v2/src/index.ts +++ b/packages/agent-core-v2/src/index.ts @@ -690,6 +690,14 @@ export * from '#/agent/shellCommand/shellCommandService'; export * from '#/agent/scopeContext/scopeContext'; export * from '#/agent/stepRetry/stepRetry'; export * from '#/agent/stepRetry/stepRetryService'; +export * from '#/agent/turnRecovery/flag'; +export * from '#/agent/turnRecovery/outputTokenRecovery'; +export * from '#/agent/turnRecovery/outputTokenRecoveryService'; +export * from '#/agent/turnRecovery/modelFallback'; +export * from '#/agent/turnRecovery/modelFallbackService'; +export * from '#/agent/turnBudget/flag'; +export * from '#/agent/turnBudget/turnBudget'; +export * from '#/agent/turnBudget/turnBudgetService'; export * from '#/features/sessionInit/sessionInit'; export * from '#/features/sessionInit/sessionInitService'; export * from '#/features/sessionInit/profile/init'; diff --git a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts new file mode 100644 index 000000000..31dfb9703 --- /dev/null +++ b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts @@ -0,0 +1,155 @@ +import { afterEach, describe, expect, it } from 'vitest'; + +import { type TokenUsage } from '#/kosong/contract/usage'; +import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; +import { IAgentLoopService } from '#/agent/loop/loop'; +import { ContinuationStepRequest } from '#/agent/loop/stepRequest'; +import { TurnStarted } from '#/agent/loop/turnEvents'; +import { IFlagService } from '#/app/flag/flag'; +import { TURN_BUDGET_CONTINUATION_FLAG_ID } from '#/agent/turnBudget/flag'; + +import { stubFlag } from '../../app/flag/stubs'; +import { appService, createTestAgent, llmGenerateServices, type TestAgentContext } from '../../harness'; + +function budgetFlags(enabled = true): ReturnType { + return appService(IFlagService, stubFlag((id) => enabled && id === TURN_BUDGET_CONTINUATION_FLAG_ID)); +} + +function outputUsage(output: number): TokenUsage { + return { inputOther: 0, inputCacheRead: 0, inputCacheCreation: 0, output }; +} + +function completedResponse(text: string, output: number) { + return { + id: `tb-${text}`, + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text }], + toolCalls: [], + }, + usage: outputUsage(output), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + }; +} + +describe('turnBudget plugin', () => { + let ctx: TestAgentContext; + + afterEach(async () => { + try { + await ctx.expectResumeMatches(); + } finally { + await ctx.dispose(); + } + }); + + async function runTurn(turnId: number): Promise>> { + void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); + const loop = ctx.get(IAgentLoopService); + loop.enqueue(new ContinuationStepRequest()); + return loop.run({ turnId }); + } + + function retryOriginTriggers(): readonly (string | undefined)[] { + return ctx + .get(IAgentContextMemoryService) + .get() + .filter((message) => message.role === 'user' && message.origin?.kind === 'retry') + .map((message) => + message.origin?.kind === 'retry' ? message.origin.trigger : undefined, + ); + } + + it('continues a naturally-stopping turn until the output-token threshold', async () => { + const outputs = [400, 400, 400]; + let calls = 0; + ctx = createTestAgent( + budgetFlags(), + llmGenerateServices(async () => { + const response = completedResponse(`step-${String(calls)}`, outputs[calls] ?? 0); + calls += 1; + return response; + }), + { + initialConfig: { + loopControl: { maxStepsPerTurn: 20, turnBudgetTokens: 1000 }, + }, + }, + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed' }); + expect(calls).toBe(3); + const triggers = retryOriginTriggers(); + expect(triggers).toHaveLength(2); + expect(triggers.every((trigger) => trigger === 'token_budget')).toBe(true); + }); + + it('stops on diminishing returns after the continuation cap', async () => { + let calls = 0; + ctx = createTestAgent( + budgetFlags(), + llmGenerateServices(async () => { + calls += 1; + return completedResponse(`tiny-${String(calls)}`, 1); + }), + { + initialConfig: { + loopControl: { maxStepsPerTurn: 20, turnBudgetTokens: 1_000_000 }, + }, + }, + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed' }); + expect(calls).toBe(4); + expect(retryOriginTriggers()).toHaveLength(3); + }); + + it('does not continue when the flag is off even with a budget configured', async () => { + let calls = 0; + ctx = createTestAgent( + budgetFlags(false), + llmGenerateServices(async () => { + calls += 1; + return completedResponse(`off-${String(calls)}`, 400); + }), + { + initialConfig: { + loopControl: { maxStepsPerTurn: 20, turnBudgetTokens: 1000 }, + }, + }, + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed' }); + expect(calls).toBe(1); + expect(retryOriginTriggers()).toHaveLength(0); + }); + + it('resets accumulation between turns', async () => { + let calls = 0; + ctx = createTestAgent( + budgetFlags(), + llmGenerateServices(async () => { + calls += 1; + return completedResponse(`reset-${String(calls)}`, 400); + }), + { + initialConfig: { + loopControl: { maxStepsPerTurn: 30, turnBudgetTokens: 1000 }, + }, + }, + ); + + await runTurn(1); + await runTurn(2); + + expect(calls).toBe(6); + expect(retryOriginTriggers()).toHaveLength(4); + }); +}); diff --git a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts new file mode 100644 index 000000000..979dab229 --- /dev/null +++ b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts @@ -0,0 +1,174 @@ +import { afterEach, describe, expect, it } from 'vitest'; + +import { APIConnectionError } from '#/kosong/contract/errors'; +import { emptyUsage } from '#/kosong/contract/usage'; +import { IFlagService } from '#/app/flag/flag'; +import { IEventBus } from '#/app/event/eventBus'; +import { IAgentLoopService } from '#/agent/loop/loop'; +import { ContinuationStepRequest } from '#/agent/loop/stepRequest'; +import { TurnStarted } from '#/agent/loop/turnEvents'; +import { IAgentProfileService } from '#/agent/profile/profile'; +import { MODEL_FALLBACK_FLAG_ID } from '#/agent/turnRecovery/flag'; +import { ModelFallbackSwitched } from '#/agent/turnRecovery/modelFallback'; + +import { stubFlag } from '../../app/flag/stubs'; +import { appService, createTestAgent, llmGenerateServices, type TestAgentContext } from '../../harness'; + +const FALLBACK_MODEL = 'fallback-model'; + +function fallbackFlags(enabled = true): ReturnType { + return appService(IFlagService, stubFlag((id) => enabled && id === MODEL_FALLBACK_FLAG_ID)); +} + +function fallbackTestConfig() { + return { + loopControl: { maxAttemptsPerStep: 1, maxStepsPerTurn: 20, fallbackModel: FALLBACK_MODEL }, + models: { + 'mock-model': { provider: 'test-provider', model: 'mock-model', maxContextSize: 1_000_000 }, + [FALLBACK_MODEL]: { + provider: 'test-provider', + model: 'fallback-model', + maxContextSize: 1_000_000, + }, + }, + }; +} + +describe('modelFallback plugin', () => { + let ctx: TestAgentContext; + + afterEach(async () => { + try { + await ctx.expectResumeMatches(); + } finally { + await ctx.dispose(); + } + }); + + async function runTurn(turnId: number): Promise>> { + void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); + const loop = ctx.get(IAgentLoopService); + loop.enqueue(new ContinuationStepRequest()); + return loop.run({ turnId }); + } + + it('switches to the configured fallback after retries are exhausted and completes', async () => { + let calls = 0; + const switched: ModelFallbackSwitched[] = []; + ctx = createTestAgent( + fallbackFlags(), + llmGenerateServices(async () => { + calls += 1; + if (calls === 1) throw new APIConnectionError('terminated'); + return { + id: `mf-${String(calls)}`, + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'recovered on fallback' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + }; + }), + { initialConfig: fallbackTestConfig() }, + ); + ctx.get(IEventBus).subscribe(ModelFallbackSwitched, (event) => switched.push(event)); + + const result = await runTurn(1); + + expect(result.type).toBe('completed'); + expect(calls).toBe(2); + expect(ctx.get(IAgentProfileService).getModel()).toBe(FALLBACK_MODEL); + expect(switched).toHaveLength(1); + expect(switched[0]?.fromModel).toBe('mock-model'); + expect(switched[0]?.toModel).toBe(FALLBACK_MODEL); + expect(switched[0]?.turnId).toBe(1); + }); + + it('falls back at most once per turn and fails when the fallback also errors', async () => { + let calls = 0; + ctx = createTestAgent( + fallbackFlags(), + llmGenerateServices(async () => { + calls += 1; + throw new APIConnectionError('still down'); + }), + { initialConfig: fallbackTestConfig() }, + ); + + const result = await runTurn(1); + + expect(result.type).toBe('failed'); + expect(calls).toBe(2); + expect(ctx.get(IAgentProfileService).getModel()).toBe(FALLBACK_MODEL); + }); + + it('stays on the primary model when no fallback is configured', async () => { + let calls = 0; + ctx = createTestAgent( + fallbackFlags(), + llmGenerateServices(async () => { + calls += 1; + if (calls === 1) throw new APIConnectionError('terminated'); + return { + id: `nofb-${String(calls)}`, + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'ok' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + }; + }), + { + initialConfig: { + loopControl: { maxAttemptsPerStep: 2, maxStepsPerTurn: 20 }, + }, + }, + ); + + const result = await runTurn(1); + + expect(result.type).toBe('completed'); + expect(calls).toBe(2); + expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); + }); + + it('does nothing when the flag is off even with a fallback configured', async () => { + let calls = 0; + ctx = createTestAgent( + fallbackFlags(false), + llmGenerateServices(async () => { + calls += 1; + if (calls === 1) throw new APIConnectionError('terminated'); + return { + id: `off-${String(calls)}`, + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'ok' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + }; + }), + { + initialConfig: { + loopControl: { maxAttemptsPerStep: 2, maxStepsPerTurn: 20, fallbackModel: FALLBACK_MODEL }, + models: fallbackTestConfig().models, + }, + }, + ); + + const result = await runTurn(1); + + expect(result.type).toBe('completed'); + expect(calls).toBeGreaterThanOrEqual(2); + expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); + }); +}); diff --git a/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts new file mode 100644 index 000000000..53444a8b2 --- /dev/null +++ b/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts @@ -0,0 +1,147 @@ +import { afterEach, describe, expect, it } from 'vitest'; + +import { emptyUsage } from '#/kosong/contract/usage'; +import { IFlagService } from '#/app/flag/flag'; +import { IAgentContextMemoryService } from '#/agent/contextMemory/contextMemory'; +import { IAgentLoopService } from '#/agent/loop/loop'; +import { ContinuationStepRequest } from '#/agent/loop/stepRequest'; +import { TurnStarted } from '#/agent/loop/turnEvents'; + +import { stubFlag } from '../../app/flag/stubs'; +import { appService, createTestAgent, llmGenerateServices, type TestAgentContext } from '../../harness'; +import { OUTPUT_TOKEN_RECOVERY_FLAG_ID } from '#/agent/turnRecovery/flag'; + +function recoveryFlags(): ReturnType { + return appService(IFlagService, stubFlag((id) => id === OUTPUT_TOKEN_RECOVERY_FLAG_ID)); +} + +describe('outputTokenRecovery plugin', () => { + let ctx: TestAgentContext; + + afterEach(async () => { + try { + await ctx.expectResumeMatches(); + } finally { + await ctx.dispose(); + } + }); + + async function runTurn(turnId: number): Promise>> { + void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); + const loop = ctx.get(IAgentLoopService); + loop.enqueue(new ContinuationStepRequest()); + return loop.run({ turnId }); + } + + function retryOriginTriggers(): readonly (string | undefined)[] { + return ctx + .get(IAgentContextMemoryService) + .get() + .filter((message) => message.role === 'user' && message.origin?.kind === 'retry') + .map((message) => + message.origin?.kind === 'retry' ? message.origin.trigger : undefined, + ); + } + + it('continues a truncated response with a resume nudge and completes', async () => { + let calls = 0; + ctx = createTestAgent( + recoveryFlags(), + llmGenerateServices(async () => { + calls += 1; + const truncated = calls === 1; + return { + id: `otr-${String(calls)}`, + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: truncated ? 'partial' : 'done' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: truncated ? ('truncated' as const) : ('completed' as const), + rawFinishReason: truncated ? 'max_tokens' : 'stop', + }; + }), + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed', truncated: false }); + expect(calls).toBe(2); + const triggers = retryOriginTriggers(); + expect(triggers).toHaveLength(1); + expect(triggers[0]).toBe('max_output_tokens'); + }); + + it('caps recoveries per turn and ends the turn truncated', async () => { + ctx = createTestAgent( + recoveryFlags(), + llmGenerateServices(async () => ({ + id: 'always-truncated', + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'cut' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'truncated' as const, + rawFinishReason: 'max_tokens', + })), + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed', truncated: true }); + expect(retryOriginTriggers()).toHaveLength(3); + }); + + it('leaves truncated turns alone when the flag is off', async () => { + ctx = createTestAgent( + appService(IFlagService, stubFlag(false)), + llmGenerateServices(async () => ({ + id: 'flag-off', + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'cut' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'truncated' as const, + rawFinishReason: 'max_tokens', + })), + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed', truncated: true }); + expect(retryOriginTriggers()).toHaveLength(0); + }); + + it('resets the attempt budget on each new turn', async () => { + let calls = 0; + ctx = createTestAgent( + recoveryFlags(), + llmGenerateServices(async () => { + calls += 1; + const truncated = calls % 2 === 1; + return { + id: `reset-${String(calls)}`, + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'x' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: truncated ? ('truncated' as const) : ('completed' as const), + rawFinishReason: truncated ? 'max_tokens' : 'stop', + }; + }), + ); + + await runTurn(1); + await runTurn(2); + + expect(calls).toBe(4); + expect(retryOriginTriggers()).toHaveLength(2); + }); +}); From a4d1ec1674c7da61e628db6ac1a6330f4d9f91e8 Mon Sep 17 00:00:00 2001 From: elkaix Date: Fri, 21 Aug 2026 19:57:49 -0400 Subject: [PATCH 2/6] fix(agent-core-v2): harden turn recovery per review findings - turnBudget: compare the preceding step delta (not the just-stored one) when evaluating diminishing returns; covers large-response-then-small sequences after the continuation cap. - modelFallback: reject an already-aborted step before switching and roll the profile back if cancellation lands after setModel; once-per-turn latch is released on rollback. - tests: exact call count in flag-off fallback case; aborted-switch regression test. --- .../src/agent/stepRetry/stepRetryService.ts | 7 ++- .../src/agent/turnBudget/turnBudget.ts | 4 ++ .../src/agent/turnBudget/turnBudgetService.ts | 3 +- .../turnRecovery/modelFallbackService.ts | 14 ++++++ .../agent/turnRecovery/outputTokenRecovery.ts | 2 + .../test/agent/turnBudget/turnBudget.test.ts | 24 +++++++++ .../agent/turnRecovery/modelFallback.test.ts | 49 +++++++++++++++++-- 7 files changed, 98 insertions(+), 5 deletions(-) diff --git a/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts b/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts index b6537b040..a22ab8216 100644 --- a/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts +++ b/packages/agent-core-v2/src/agent/stepRetry/stepRetryService.ts @@ -58,6 +58,10 @@ export const stepRetryFailedAttemptsKey = defineState( export class AgentStepRetryService extends Disposable implements IAgentStepRetryService { declare readonly _serviceBrand: undefined; + private static stepAborted(context: LoopErrorContext): boolean { + return context.signal.aborted || context.currentStep?.signal.aborted === true; + } + constructor( @IAgentLoopService private readonly loopService: IAgentLoopService, @IConfigService private readonly config: IConfigService, @@ -123,9 +127,10 @@ export class AgentStepRetryService extends Disposable implements IAgentStepRetry ); if (this.failedAttempts >= maxAttempts) { this.resetAttempts(); + if (AgentStepRetryService.stepAborted(context)) return false; const switched = await this.modelFallback.tryFallbackSwitch(context); if (!switched) return false; - if (context.currentStep?.signal.aborted === true) return false; + if (AgentStepRetryService.stepAborted(context)) return false; context.retry(driver, { at: 'head' }); return true; } diff --git a/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts b/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts index 7a3be43bb..77065d010 100644 --- a/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts +++ b/packages/agent-core-v2/src/agent/turnBudget/turnBudget.ts @@ -11,10 +11,14 @@ export interface IAgentTurnBudgetService { export const IAgentTurnBudgetService = createDecorator('agentTurnBudgetService'); +/** Fraction of the configured token target a turn must reach before stopping naturally. */ export const TURN_BUDGET_COMPLETION_THRESHOLD = 0.9; +/** Per-step output-token delta below which a step counts as low-progress. */ export const TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS = 500; +/** Continuations after which consecutive low-progress deltas stop the turn. */ export const TURN_BUDGET_MAX_DIMINISHING_CONTINUATIONS = 3; +/** Builds the meta nudge injected before each budget continuation. */ export function turnBudgetNudgeText(pct: number, used: number, budget: number): string { return `Stopped at ${pct}% of token target (${used} / ${budget}). Keep working - do not summarize.`; } diff --git a/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts b/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts index 0e60f9125..66257535e 100644 --- a/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts +++ b/packages/agent-core-v2/src/agent/turnBudget/turnBudgetService.ts @@ -108,6 +108,7 @@ export class AgentTurnBudgetService extends Disposable implements IAgentTurnBudg const delta = context.usage.output; const used = this.tokensUsed + delta; + const previousDelta = this.lastDeltaTokens; this.lastDeltaTokens = delta; this.tokensUsed = used; @@ -118,7 +119,7 @@ export class AgentTurnBudgetService extends Disposable implements IAgentTurnBudg const diminishing = this.continuations >= TURN_BUDGET_MAX_DIMINISHING_CONTINUATIONS && delta < TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS && - this.lastDeltaTokens < TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS; + previousDelta < TURN_BUDGET_DIMINISHING_MIN_DELTA_TOKENS; if (diminishing) return; if (used >= budget * TURN_BUDGET_COMPLETION_THRESHOLD) return; diff --git a/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts b/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts index 088cebe31..4eb3c6704 100644 --- a/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts +++ b/packages/agent-core-v2/src/agent/turnRecovery/modelFallbackService.ts @@ -25,6 +25,10 @@ export const modelFallbackUsedKey = defineState( export class AgentModelFallbackService extends Disposable implements IAgentModelFallbackService { declare readonly _serviceBrand: undefined; + private static stepAborted(context: LoopErrorContext): boolean { + return context.signal.aborted || context.currentStep?.signal.aborted === true; + } + constructor( @IFlagService private readonly flags: IFlagService, @IConfigService private readonly config: IConfigService, @@ -57,6 +61,7 @@ export class AgentModelFallbackService extends Disposable implements IAgentModel async tryFallbackSwitch(context: LoopErrorContext): Promise { if (!this.flags.enabled(MODEL_FALLBACK_FLAG_ID)) return false; + if (AgentModelFallbackService.stepAborted(context)) return false; const target = this.fallbackModel(); if (target === undefined || target.length === 0) return false; if (this.used || target === this.profile.getModel()) return false; @@ -70,6 +75,15 @@ export class AgentModelFallbackService extends Disposable implements IAgentModel this.used = false; return false; } + if (AgentModelFallbackService.stepAborted(context)) { + this.used = false; + try { + await this.profile.setModel(fromModel); + } catch { + return false; + } + return false; + } void this.dispatcher.dispatch( new ModelFallbackSwitched({ turnId: context.turnId, diff --git a/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts b/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts index 9ee043c62..b8d890f03 100644 --- a/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts +++ b/packages/agent-core-v2/src/agent/turnRecovery/outputTokenRecovery.ts @@ -12,8 +12,10 @@ export interface IAgentOutputTokenRecoveryService { export const IAgentOutputTokenRecoveryService = createDecorator('agentOutputTokenRecoveryService'); +/** Maximum resume-nudge continuations injected per turn for truncated output. */ export const MAX_OUTPUT_TOKEN_RECOVERY_ATTEMPTS = 3; +/** Meta user message appended before each output-token recovery continuation. */ export const OUTPUT_TOKEN_RECOVERY_NUDGE = 'Output token limit hit. Resume directly - no apology, no recap of what you were doing. ' + 'Pick up mid-thought if that is where the cut happened. Break remaining work into smaller pieces.'; diff --git a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts index 31dfb9703..020f93c52 100644 --- a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts +++ b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts @@ -109,6 +109,30 @@ describe('turnBudget plugin', () => { expect(retryOriginTriggers()).toHaveLength(3); }); + it('keeps continuing when only the current delta is small after the cap', async () => { + const outputs = [1000, 1000, 1000, 1, 1]; + let calls = 0; + ctx = createTestAgent( + budgetFlags(), + llmGenerateServices(async () => { + const response = completedResponse(`mixed-${String(calls)}`, outputs[calls] ?? 0); + calls += 1; + return response; + }), + { + initialConfig: { + loopControl: { maxStepsPerTurn: 30, turnBudgetTokens: 1_000_000 }, + }, + }, + ); + + const result = await runTurn(1); + + expect(result).toMatchObject({ type: 'completed' }); + expect(calls).toBe(5); + expect(retryOriginTriggers()).toHaveLength(4); + }); + it('does not continue when the flag is off even with a budget configured', async () => { let calls = 0; ctx = createTestAgent( diff --git a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts index 979dab229..4aae7a753 100644 --- a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts +++ b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts @@ -4,12 +4,12 @@ import { APIConnectionError } from '#/kosong/contract/errors'; import { emptyUsage } from '#/kosong/contract/usage'; import { IFlagService } from '#/app/flag/flag'; import { IEventBus } from '#/app/event/eventBus'; -import { IAgentLoopService } from '#/agent/loop/loop'; +import { IAgentLoopService, type LoopErrorContext, type Step } from '#/agent/loop/loop'; import { ContinuationStepRequest } from '#/agent/loop/stepRequest'; import { TurnStarted } from '#/agent/loop/turnEvents'; import { IAgentProfileService } from '#/agent/profile/profile'; import { MODEL_FALLBACK_FLAG_ID } from '#/agent/turnRecovery/flag'; -import { ModelFallbackSwitched } from '#/agent/turnRecovery/modelFallback'; +import { IAgentModelFallbackService, ModelFallbackSwitched } from '#/agent/turnRecovery/modelFallback'; import { stubFlag } from '../../app/flag/stubs'; import { appService, createTestAgent, llmGenerateServices, type TestAgentContext } from '../../harness'; @@ -168,7 +168,50 @@ describe('modelFallback plugin', () => { const result = await runTurn(1); expect(result.type).toBe('completed'); - expect(calls).toBeGreaterThanOrEqual(2); + expect(calls).toBe(2); + expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); + }); + + it('refuses to switch when the step is already aborted', async () => { + ctx = createTestAgent( + fallbackFlags(), + llmGenerateServices(async () => ({ + id: 'aborted', + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'ok' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + })), + { initialConfig: fallbackTestConfig() }, + ); + + const signal = AbortSignal.abort(); + const stepStub: Step = { + id: 'step-1', + turnId: 1, + state: 'running', + signal, + result: Promise.resolve({ type: 'cancelled', reason: new Error('cancelled') }), + cancel: () => false, + }; + const context: LoopErrorContext = { + turnId: 1, + step: 1, + signal, + error: new APIConnectionError('terminated'), + failedDriver: new ContinuationStepRequest(), + retry: () => stepStub, + }; + + const switched = await ctx + .get(IAgentModelFallbackService) + .tryFallbackSwitch(context); + + expect(switched).toBe(false); expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); }); }); From 1bde4b74fba4827f9866b30cb951fb997f237e76 Mon Sep 17 00:00:00 2001 From: elkaix Date: Fri, 21 Aug 2026 20:04:28 -0400 Subject: [PATCH 3/6] test(agent-core-v2): strengthen recovery tests per review findings Assert exact token_budget trigger values in the mixed-delta continuation test, and cover the current-step abort guard separately from the loop-signal abort case. --- .../test/agent/turnBudget/turnBudget.test.ts | 7 ++- .../agent/turnRecovery/modelFallback.test.ts | 47 ++++++++++++++++++- 2 files changed, 52 insertions(+), 2 deletions(-) diff --git a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts index 020f93c52..0fb2cee88 100644 --- a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts +++ b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts @@ -130,7 +130,12 @@ describe('turnBudget plugin', () => { expect(result).toMatchObject({ type: 'completed' }); expect(calls).toBe(5); - expect(retryOriginTriggers()).toHaveLength(4); + expect(retryOriginTriggers()).toEqual([ + 'token_budget', + 'token_budget', + 'token_budget', + 'token_budget', + ]); }); it('does not continue when the flag is off even with a budget configured', async () => { diff --git a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts index 4aae7a753..de7bcb57a 100644 --- a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts +++ b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts @@ -172,7 +172,7 @@ describe('modelFallback plugin', () => { expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); }); - it('refuses to switch when the step is already aborted', async () => { + it('refuses to switch when the turn signal is already aborted', async () => { ctx = createTestAgent( fallbackFlags(), llmGenerateServices(async () => ({ @@ -214,4 +214,49 @@ describe('modelFallback plugin', () => { expect(switched).toBe(false); expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); }); + + it('refuses to switch when only the current step is aborted', async () => { + ctx = createTestAgent( + fallbackFlags(), + llmGenerateServices(async () => ({ + id: 'step-aborted', + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'ok' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + })), + { initialConfig: fallbackTestConfig() }, + ); + + const stepController = new AbortController(); + stepController.abort(); + const stepStub: Step = { + id: 'step-1', + turnId: 1, + state: 'running', + signal: stepController.signal, + result: Promise.resolve({ type: 'cancelled', reason: new Error('cancelled') }), + cancel: () => false, + }; + const context: LoopErrorContext = { + turnId: 1, + step: 1, + signal: new AbortController().signal, + currentStep: stepStub, + error: new APIConnectionError('terminated'), + failedDriver: new ContinuationStepRequest(), + retry: () => stepStub, + }; + + const switched = await ctx + .get(IAgentModelFallbackService) + .tryFallbackSwitch(context); + + expect(switched).toBe(false); + expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); + }); }); From 1da28be1a50d170b962be194b9b60e5ffece8329 Mon Sep 17 00:00:00 2001 From: elkaix Date: Fri, 21 Aug 2026 20:28:58 -0400 Subject: [PATCH 4/6] test(agent-core-v2): cover cancellation during fallback setModel Add a deferred-setModel profile subclass test asserting the switch returns false and rolls back to the previous model when the step aborts mid-switch; assert exact token_budget trigger values. --- .../agent/turnRecovery/modelFallback.test.ts | 90 ++++++++++++++++++- 1 file changed, 87 insertions(+), 3 deletions(-) diff --git a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts index de7bcb57a..416364b1b 100644 --- a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts +++ b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts @@ -1,4 +1,4 @@ -import { afterEach, describe, expect, it } from 'vitest'; +import { afterEach, describe, expect, it, vi } from 'vitest'; import { APIConnectionError } from '#/kosong/contract/errors'; import { emptyUsage } from '#/kosong/contract/usage'; @@ -7,12 +7,23 @@ import { IEventBus } from '#/app/event/eventBus'; import { IAgentLoopService, type LoopErrorContext, type Step } from '#/agent/loop/loop'; import { ContinuationStepRequest } from '#/agent/loop/stepRequest'; import { TurnStarted } from '#/agent/loop/turnEvents'; -import { IAgentProfileService } from '#/agent/profile/profile'; +import { + type ProfileSetModelResult, + IAgentProfileService, +} from '#/agent/profile/profile'; +import { AgentProfileService } from '#/agent/profile/profileService'; +import { SyncDescriptor } from '#/_base/di/descriptors'; import { MODEL_FALLBACK_FLAG_ID } from '#/agent/turnRecovery/flag'; import { IAgentModelFallbackService, ModelFallbackSwitched } from '#/agent/turnRecovery/modelFallback'; import { stubFlag } from '../../app/flag/stubs'; -import { appService, createTestAgent, llmGenerateServices, type TestAgentContext } from '../../harness'; +import { + agentService, + appService, + createTestAgent, + llmGenerateServices, + type TestAgentContext, +} from '../../harness'; const FALLBACK_MODEL = 'fallback-model'; @@ -20,6 +31,24 @@ function fallbackFlags(enabled = true): ReturnType { return appService(IFlagService, stubFlag((id) => enabled && id === MODEL_FALLBACK_FLAG_ID)); } +class DeferredSetModelProfile extends AgentProfileService { + override async setModel(model: string): Promise { + setModelCalls.push(model); + if (deferFirstSetModel) { + deferFirstSetModel = false; + await new Promise((resolve) => { + resolveFirstSetModel = resolve; + }); + return { model }; + } + return super.setModel(model); + } +} + +let deferFirstSetModel = false; +let resolveFirstSetModel: (() => void) | undefined; +const setModelCalls: string[] = []; + function fallbackTestConfig() { return { loopControl: { maxAttemptsPerStep: 1, maxStepsPerTurn: 20, fallbackModel: FALLBACK_MODEL }, @@ -259,4 +288,59 @@ describe('modelFallback plugin', () => { expect(switched).toBe(false); expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); }); + + it('rolls back the model when the step aborts during setModel', async () => { + const stepController = new AbortController(); + deferFirstSetModel = true; + resolveFirstSetModel = undefined; + setModelCalls.length = 0; + ctx = createTestAgent( + fallbackFlags(), + llmGenerateServices(async () => ({ + id: 'mid-abort', + message: { + role: 'assistant' as const, + content: [{ type: 'text' as const, text: 'ok' }], + toolCalls: [], + }, + usage: emptyUsage(), + finishReason: 'completed' as const, + rawFinishReason: 'stop', + })), + { initialConfig: fallbackTestConfig() }, + agentService( + IAgentProfileService, + new SyncDescriptor(DeferredSetModelProfile), + ), + ); + + const stepStub: Step = { + id: 'step-1', + turnId: 1, + state: 'running', + signal: stepController.signal, + result: Promise.resolve({ type: 'cancelled', reason: new Error('cancelled') }), + cancel: () => false, + }; + const context: LoopErrorContext = { + turnId: 1, + step: 1, + signal: new AbortController().signal, + currentStep: stepStub, + error: new APIConnectionError('terminated'), + failedDriver: new ContinuationStepRequest(), + retry: () => stepStub, + }; + + const pending = ctx + .get(IAgentModelFallbackService) + .tryFallbackSwitch(context); + await vi.waitFor(() => expect(resolveFirstSetModel).toBeDefined()); + stepController.abort(); + resolveFirstSetModel!(); + + await expect(pending).resolves.toBe(false); + expect(setModelCalls).toEqual(['fallback-model', 'mock-model']); + expect(ctx.get(IAgentProfileService).getModel()).toBe('mock-model'); + }); }); From 42cffb0eabf6c1b584243d37e3332e3477cfab0a Mon Sep 17 00:00:00 2001 From: elkaix Date: Fri, 21 Aug 2026 20:40:49 -0400 Subject: [PATCH 5/6] test(agent-core-v2): make rollback assertion prove model restoration The deferred setModel now applies the switch after its await resolves, so the post-abort rollback assertion verifies a real state change instead of passing vacuously. --- .../agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts index 416364b1b..48845500a 100644 --- a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts +++ b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts @@ -39,7 +39,7 @@ class DeferredSetModelProfile extends AgentProfileService { await new Promise((resolve) => { resolveFirstSetModel = resolve; }); - return { model }; + return super.setModel(model); } return super.setModel(model); } From 1e135732f0f76f8589bb1190000d5493215a1b65 Mon Sep 17 00:00:00 2001 From: elkaix Date: Fri, 21 Aug 2026 23:10:09 -0400 Subject: [PATCH 6/6] test(agent-core-v2): carry agentId in turn resilience TurnStarted dispatches --- packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts | 2 +- .../agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts | 2 +- .../test/agent/turnRecovery/outputTokenRecovery.test.ts | 2 +- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts index 0fb2cee88..75b4281d0 100644 --- a/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts +++ b/packages/agent-core-v2/test/agent/turnBudget/turnBudget.test.ts @@ -45,7 +45,7 @@ describe('turnBudget plugin', () => { }); async function runTurn(turnId: number): Promise>> { - void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); + void ctx.dispatcher.dispatch(new TurnStarted({ agentId: 'main', turnId, origin: { kind: 'user' } })); const loop = ctx.get(IAgentLoopService); loop.enqueue(new ContinuationStepRequest()); return loop.run({ turnId }); diff --git a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts index 48845500a..4a1c4e8ea 100644 --- a/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts +++ b/packages/agent-core-v2/test/agent/turnRecovery/modelFallback.test.ts @@ -75,7 +75,7 @@ describe('modelFallback plugin', () => { }); async function runTurn(turnId: number): Promise>> { - void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); + void ctx.dispatcher.dispatch(new TurnStarted({ agentId: 'main', turnId, origin: { kind: 'user' } })); const loop = ctx.get(IAgentLoopService); loop.enqueue(new ContinuationStepRequest()); return loop.run({ turnId }); diff --git a/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts b/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts index 53444a8b2..56980a6fe 100644 --- a/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts +++ b/packages/agent-core-v2/test/agent/turnRecovery/outputTokenRecovery.test.ts @@ -27,7 +27,7 @@ describe('outputTokenRecovery plugin', () => { }); async function runTurn(turnId: number): Promise>> { - void ctx.dispatcher.dispatch(new TurnStarted({ turnId, origin: { kind: 'user' } })); + void ctx.dispatcher.dispatch(new TurnStarted({ agentId: 'main', turnId, origin: { kind: 'user' } })); const loop = ctx.get(IAgentLoopService); loop.enqueue(new ContinuationStepRequest()); return loop.run({ turnId });