From 9620a752c8cdb6e962caaa0808be3745bb6a00ea Mon Sep 17 00:00:00 2001 From: creatixchu Date: Wed, 19 Aug 2026 14:55:26 +0800 Subject: [PATCH] =?UTF-8?q?refactor(agent-loop):=20=E7=BC=A9=E5=B0=8F?= =?UTF-8?q?=E5=8F=96=E6=B6=88=E5=89=8D=E7=BC=80=E6=94=B6=E5=B0=BE=E8=8C=83?= =?UTF-8?q?=E5=9B=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...cancelled-stream-prefix-finalize.i18n.yaml | 4 +- ...-08-10-cancelled-stream-prefix-finalize.md | 4 +- ...-10-cancelled-stream-prefix-finalize.zh.md | 4 +- packages/core/agent-loop/src/agent.ts | 158 +++++++----------- 4 files changed, 71 insertions(+), 99 deletions(-) diff --git a/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.i18n.yaml b/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.i18n.yaml index 2991c6b1c8..50cdedbe5b 100644 --- a/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.i18n.yaml +++ b/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.i18n.yaml @@ -2,5 +2,5 @@ # side as of the last confirmed-consistent state. Both languages carry equal authority; # after editing either side, bring the other along and re-record with: # pnpm run verify-translation-pairing --write .agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.md -2026-08-10-cancelled-stream-prefix-finalize.md: 722d5531820efe5640068f11582aef6485f1f9e1 -2026-08-10-cancelled-stream-prefix-finalize.zh.md: a48790fe5883e341148cd82053d4ca9008e4fe81 +2026-08-10-cancelled-stream-prefix-finalize.md: 0cae25b786922fba8204d68ca9c0a669e43d76a0 +2026-08-10-cancelled-stream-prefix-finalize.zh.md: e961ea6a51f74dcc244e4ad8970eae4cbe4c9a6c diff --git a/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.md b/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.md index 722d553182..0cae25b786 100644 --- a/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.md +++ b/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.md @@ -12,9 +12,9 @@ The model history must contain assistant content that remains visible to the use ## Decision -`ReactLoopAgent.step()` retains the active `BlockAssembler`, logged chunk seqs, and provider route until the attempt commits or fails. Cancellation of an uncommitted attempt appends its delivered prefix as the step's `assistant/message` with `interrupted: true`, `surfaceOp: 'append'`, and `sourceEventSeqs` containing exactly the logged chunks. The append precedes `step/end` and the aborted `turn/end`. +`ReactLoopAgent.step()` catches cancellation while consuming a model stream, when its `BlockAssembler`, logged chunk seqs, and provider route identify the delivered prefix. It appends that prefix as the step's `assistant/message` with `interrupted: true`, `surfaceOp: 'append'`, and `sourceEventSeqs` containing exactly the logged chunks. The append precedes `step/end` and the aborted `turn/end`. -`BlockAssembler.interruptedBlocks()` returns closed and open `text` and `reasoning` blocks with non-whitespace content in stream order. It omits tool calls because interruption precedes dispatch and no real result exists; it also omits empty blocks and open unknown block types. An empty result appends no assistant message. An attempt ending with an `error` or `aborted` finish is cleared before `agent/request-error`, so provider failures and cancellation during recovery commit no content from the failed attempt. +`BlockAssembler.interruptedBlocks()` returns closed and open `text` and `reasoning` blocks with non-whitespace content in stream order. It omits tool calls because interruption precedes dispatch and no real result exists; it also omits empty blocks and open unknown block types. An empty result appends no assistant message. Provider `error` and `aborted` finishes leave the stream-consumption scope before `agent/request-error`, so provider failures and cancellation during recovery commit no content from the failed request. Chat and Trajectory Conversation Definitions read `interrupted` from the durable message. Chat renders the Stopped marker, while Trajectory keeps the provider request in the error lifecycle after `step/end` and retains the durable result seq and provenance. Cancellation during tool execution follows the tool scheduler contract because the assistant message has already committed: started calls produce real results, and undispatched calls receive `ABORTED_BEFORE_DISPATCH` results. diff --git a/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.zh.md b/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.zh.md index a48790fe58..e961ea6a51 100644 --- a/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.zh.md +++ b/.agents/notes/implemented/architecture/2026-08-10-cancelled-stream-prefix-finalize.zh.md @@ -12,9 +12,9 @@ Status: implemented ## Decision -`ReactLoopAgent.step()` 会保留活跃的 `BlockAssembler`、已记录的分片 seq 和提供方路由,直到尝试提交或失败。取消未提交的尝试时,循环把已送达前缀追加为该 step 的 `assistant/message`,并设置 `interrupted: true`、`surfaceOp: 'append'` 以及恰好包含已记录分片的 `sourceEventSeqs`。该追加先于 `step/end` 和记录 aborted 的 `turn/end`。 +`ReactLoopAgent.step()` 在消费模型流期间捕捉取消,此时 `BlockAssembler`、已记录的分片 seq 和提供方路由可以确定已送达前缀。循环把该前缀追加为 step 的 `assistant/message`,并设置 `interrupted: true`、`surfaceOp: 'append'` 以及恰好包含已记录分片的 `sourceEventSeqs`。该追加先于 `step/end` 和记录 aborted 的 `turn/end`。 -`BlockAssembler.interruptedBlocks()` 按流顺序返回内容非空白的已闭合和未闭合 `text` 与 `reasoning` 块。打断先于分派,没有真实工具结果,因此它会省略工具调用,也会省略空块和未闭合的未知块类型。返回结果为空时不追加 assistant 消息。以 `error` 或 `aborted` finish 结束的尝试会在 `agent/request-error` 前清空,因此提供方故障和恢复期间的取消都不会提交失败尝试的内容。 +`BlockAssembler.interruptedBlocks()` 按流顺序返回内容非空白的已闭合和未闭合 `text` 与 `reasoning` 块。打断先于分派,没有真实工具结果,因此它会省略工具调用,也会省略空块和未闭合的未知块类型。返回结果为空时不追加 assistant 消息。提供方的 `error` 和 `aborted` finish 会在 `agent/request-error` 前离开流消费范围,因此提供方故障和恢复期间的取消都不会提交失败请求的内容。 Chat 和 Trajectory Conversation Definition 从持久消息读取 `interrupted`。Chat 渲染 Stopped 标记,Trajectory 则在 `step/end` 后把提供方请求保持在 error 生命周期,并保留持久结果 seq 和提供方信息。工具执行期间的取消遵循工具调度器约定,因为 assistant 消息已提交:已启动的调用生成真实结果,未分派的调用获得 `ABORTED_BEFORE_DISPATCH` 结果。 diff --git a/packages/core/agent-loop/src/agent.ts b/packages/core/agent-loop/src/agent.ts index 1ba09684c3..3ef1ec7aa4 100644 --- a/packages/core/agent-loop/src/agent.ts +++ b/packages/core/agent-loop/src/agent.ts @@ -51,14 +51,6 @@ type PreparedStep = | { kind: 'reject' } | { kind: 'enter'; messages: UserMessage[]; assembly: PromptAssembly } -/** One live streaming attempt whose logged chunk prefix an abort can still finalize. */ -interface InterruptedAttempt { - readonly assembler: BlockAssembler - readonly chunkSeqs: number[] - readonly provider: string - readonly model: string -} - /** Remove adapter-derived values before plugins propose the next request config. */ function requestProposal(header: EpochHeader): LlmCallConfig { if (header.adapterDefaults === undefined) return header.config @@ -290,6 +282,8 @@ export class ReactLoopAgent implements Agent { for (const message of decision.messages) { this.session.append('user/message', message, { surfaceOp: 'append' }) } + // max-tokens is sticky: once any step hits the ceiling, later steps + // that complete normally must not downgrade the turn outcome. const stepEnd = await this.step(decision.assembly) // max-tokens stays sticky: a later completed step must not // downgrade the turn outcome. @@ -342,17 +336,13 @@ export class ReactLoopAgent implements Agent { signal.throwIfAborted() const system = renderPrompt(assembly) - // Keep the active attempt until it commits or fails so cancellation can - // preserve the same streamed prefix in durable message history. - let attempt: InterruptedAttempt | undefined - try { - while (true) { - const { request, preparedCall } = await this.buildRequest( - turn, step, assembly.tools, system, this.session.deriveMessages(), signal, - ) - const assembler = new BlockAssembler() - const chunkSeqs: number[] = [] - attempt = { assembler, chunkSeqs, provider: request.provider, model: request.model } + while (true) { + const { request, preparedCall } = await this.buildRequest( + turn, step, assembly.tools, system, this.session.deriveMessages(), signal, + ) + const assembler = new BlockAssembler() + const chunkSeqs: number[] = [] + try { const stream = preparedCall?.stream(request) ?? this.loopCtx.llm.stream(request) signal.throwIfAborted() for await (const chunk of stream) { @@ -361,92 +351,74 @@ export class ReactLoopAgent implements Agent { assembler.push(chunk) } signal.throwIfAborted() - const finish = assembler.finish - if (finish.kind === 'error' || finish.kind === 'aborted') { - // Provider failures commit no assistant content. Clearing before the - // recovery waterfall also prevents a cancellation during retry delay - // from restoring the failed attempt after clients reset its stream. - attempt = undefined - const action = await this.dispatch.waterfall( - 'agent/request-error', { + } catch (error: unknown) { + if (signal.aborted) { + const content = assembler.interruptedBlocks() + if (content.length > 0) { + this.session.append('assistant/message', { turn, step, - provider: request.provider, - failure: finish.failure, - retryPolicy: preparedCall?.retryPolicy, - signal, - }, - () => Promise.resolve(undefined), - ) - signal.throwIfAborted() - if (action?.kind !== 'retry') { - throw new LlmError(finish.failure.message, finish.failure.code, finish.failure) + message: createAssistantMessage({ + content, + source: { provider: request.provider, model: request.model }, + }), + interrupted: true, + ...assembler.usage === undefined ? {} : { usage: assembler.usage }, + }, { surfaceOp: 'append', sourceEventSeqs: chunkSeqs }) } - continue } - - const message = createAssistantMessage({ - content: assembler.blocks(), - source: { - provider: request.provider, - model: request.model, - ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {}, - }, - }) - attempt = undefined - this.session.append( - 'assistant/message', - { + throw error + } + const finish = assembler.finish + if (finish.kind === 'error' || finish.kind === 'aborted') { + const action = await this.dispatch.waterfall( + 'agent/request-error', { turn, step, - message, - ...assembler.usage === undefined ? {} : { usage: assembler.usage }, + provider: request.provider, + failure: finish.failure, + retryPolicy: preparedCall?.retryPolicy, + signal, }, - { surfaceOp: 'append', sourceEventSeqs: chunkSeqs }, + () => Promise.resolve(undefined), ) - if (finish.kind === 'max-tokens') return { kind: 'max-tokens' } + signal.throwIfAborted() + if (action?.kind !== 'retry') { + throw new LlmError(finish.failure.message, finish.failure.code, finish.failure) + } + continue + } - const toolCalls = message.content.filter(block => block.type === 'tool-call') - if (toolCalls.length === 0) return { kind: 'completed' } - const { concluded } = await executeToolCalls( - this.loopCtx, turn, step, toolCalls, signal, - context => this.inbox.splice('next-step', this.inbox.nextStep.length, 0, [context]), - ) - return concluded ? { kind: 'completed' } : null - } - } catch (error: unknown) { - if (signal.aborted && attempt !== undefined) { - this.appendInterruptedAssistant(turn, step, attempt) - } - throw error + const message = createAssistantMessage({ + content: assembler.blocks(), + source: { + provider: request.provider, + model: request.model, + ...assembler.replayState !== undefined ? { replayState: assembler.replayState } : {}, + }, + }) + this.session.append( + 'assistant/message', + { + turn, + step, + message, + ...assembler.usage === undefined ? {} : { usage: assembler.usage }, + }, + { surfaceOp: 'append', sourceEventSeqs: chunkSeqs }, + ) + if (finish.kind === 'max-tokens') return { kind: 'max-tokens' } + + const toolCalls = message.content.filter(block => block.type === 'tool-call') + if (toolCalls.length === 0) return { kind: 'completed' } + const { concluded } = await executeToolCalls( + this.loopCtx, turn, step, toolCalls, signal, + context => this.inbox.splice('next-step', this.inbox.nextStep.length, 0, [context]), + ) + return concluded ? { kind: 'completed' } : null } } - /** - * Append a cancelled attempt's delivered text and reasoning as an interrupted - * assistant message. Undispatched tool calls and empty content are omitted; - * the resulting durable history matches the prefix clients rendered. - */ - private appendInterruptedAssistant(turn: number, step: number, attempt: InterruptedAttempt): void { - const content = attempt.assembler.interruptedBlocks() - if (content.length === 0) return - const message = createAssistantMessage({ - content, - source: { provider: attempt.provider, model: attempt.model }, - }) - this.session.append( - 'assistant/message', - { - turn, - step, - message, - interrupted: true, - ...attempt.assembler.usage === undefined ? {} : { usage: attempt.assembler.usage }, - }, - { surfaceOp: 'append', sourceEventSeqs: attempt.chunkSeqs }, - ) - } - /** * Compose one frozen request and bind it to the adapter registration that * resolved its exact-model defaults.