From ec865ed89c456dfd7f2efe1a64ab6f8704becffc Mon Sep 17 00:00:00 2001 From: _Kerman Date: Thu, 2 Jul 2026 19:20:41 +0800 Subject: [PATCH] refactor(agent-core-v2): simplify loopService --- .../src/agent/loop/loopService.ts | 166 ++++++++---------- 1 file changed, 74 insertions(+), 92 deletions(-) diff --git a/packages/agent-core-v2/src/agent/loop/loopService.ts b/packages/agent-core-v2/src/agent/loop/loopService.ts index cfc901e12..7a587f850 100644 --- a/packages/agent-core-v2/src/agent/loop/loopService.ts +++ b/packages/agent-core-v2/src/agent/loop/loopService.ts @@ -1,35 +1,29 @@ +import { randomUUID } from 'node:crypto'; + +import type { AgentEvent } from '@moonshot-ai/protocol'; + import { InstantiationType } from '#/_base/di/extensions'; -import { LifecycleScope, registerScopedService } from '#/_base/di/scope'; - -import { - APIContextOverflowError, - createToolMessage, - isToolCall, - isToolCallPart, - type ContentPart, - type StreamedMessagePart, - type TokenUsage -} from '#/app/llmProtocol'; - +import { LifecycleScope, registerScopedService } from '#/_base/di/scope'; import { IAgentContextSizeService } from '#/agent/contextSize'; import { IAgentExternalHooksService } from '#/agent/externalHooks'; -import { - IAgentLLMRequesterService, - type LLMRequestFinish, -} from '#/agent/llmRequester'; +import { IAgentLLMRequesterService, type LLMRequestFinish } from '#/agent/llmRequester'; import { IAgentProfileService } from '#/agent/profile'; import { IAgentRecordService } from '#/agent/record'; import type { ToolResult } from '#/agent/tool'; import { IAgentToolExecutorService } from '#/agent/toolExecutor'; import { IConfigService } from '#/app/config'; -import { ILogService } from '#/app/log'; import { - ErrorCodes, - isKimiError -} from "#/errors"; + APIContextOverflowError, + createToolMessage, + type ContentPart, + type StreamedMessagePart, + type TokenUsage, +} from '#/app/llmProtocol'; +import { ILogService } from '#/app/log'; +import { ErrorCodes, isKimiError } from '#/errors'; import { OrderedHookSlot } from '#/hooks'; -import type { AgentEvent } from '@moonshot-ai/protocol'; -import { randomUUID } from 'node:crypto'; + +import { IAgentContextMemoryService, newMessageId, type ContextMessage } from '../contextMemory'; import { LOOP_CONTROL_SECTION, type LoopControl } from './configSection'; import { createMaxStepsExceededError, @@ -44,7 +38,6 @@ import type { LoopTurnStopReason, TurnResult, } from './types'; -import { IAgentContextMemoryService, newMessageId, type ContextMessage } from '../contextMemory'; const TOOL_ERROR_STATUS = 'ERROR: Tool execution failed.'; const TOOL_EMPTY_STATUS = 'Tool output is empty.'; @@ -132,9 +125,7 @@ export class AgentLoopService implements IAgentLoopService { return { stopReason: 'aborted', steps }; } - const reason: LoopInterruptReason = isMaxStepsExceededError(error) - ? 'max_steps' - : 'error'; + const reason: LoopInterruptReason = isMaxStepsExceededError(error) ? 'max_steps' : 'error'; this.emitStepInterrupted(turnId, activeStep, reason, errorMessage(error)); if (isContextOverflowError(error)) { @@ -168,7 +159,7 @@ export class AgentLoopService implements IAgentLoopService { emit({ type: 'turn.step.started', turnId, step: currentStep, stepId: stepUuid }); - const emitToolCallDelta = createToolCallDeltaHandler(emit, turnId); + const emitStreamPart = this.createStreamPartHandler(turnId); const response = await this.llmRequester.request( { requestLogFields: { turnStep }, @@ -192,9 +183,7 @@ export class AgentLoopService implements IAgentLoopService { }, usageContext: { type: 'turn', turnId }, }, - (part) => { - this.emitStreamPart(turnId, emitToolCallDelta, part); - }, + emitStreamPart, signal, ); @@ -283,33 +272,6 @@ export class AgentLoopService implements IAgentLoopService { this.contextSize.measured(this.context.get().length, tokens); } - private emitStreamPart( - turnId: number, - emitToolCallDelta: (part: StreamedMessagePart) => void, - part: StreamedMessagePart, - ): void { - switch (part.type) { - case 'text': - this.record.signal({ type: 'assistant.delta', turnId, delta: part.text }); - return; - case 'think': - this.record.signal({ type: 'thinking.delta', turnId, delta: part.think }); - return; - case 'image_url': - case 'audio_url': - case 'video_url': - return; - case 'function': - case 'tool_call_part': - emitToolCallDelta(part); - return; - default: { - const _exhaustive: never = part; - return _exhaustive; - } - } - } - private emitStepCompleted( turnId: number, step: number, @@ -356,6 +318,61 @@ export class AgentLoopService implements IAgentLoopService { message, }); } + + private createStreamPartHandler(turnId: number): (part: StreamedMessagePart) => void { + // Maps a tool call's streaming index to its identity so that interleaved + // argument deltas from parallel tool calls can be routed to the right call. + // Each provider emits a `function` header before any of its `tool_call_part` + // deltas, and a delta's `index` always matches a previously-seen header's + // `_streamIndex`. The `undefined` key doubles as the single-call fallback + // for providers that stream without indices: those streams never mix indexed + // and unindexed parts, so the most recent unindexed header is always the + // target. + const callsByIndex = new Map(); + + return (part) => { + switch (part.type) { + case 'text': + this.record.signal({ type: 'assistant.delta', turnId, delta: part.text }); + return; + case 'think': + this.record.signal({ type: 'thinking.delta', turnId, delta: part.think }); + return; + case 'image_url': + case 'audio_url': + case 'video_url': + return; + case 'function': { + callsByIndex.set(part._streamIndex, { id: part.id, name: part.name }); + this.record.signal({ + type: 'tool.call.delta', + turnId, + toolCallId: part.id, + name: part.name, + argumentsPart: part.arguments ?? undefined, + }); + return; + } + case 'tool_call_part': { + if (part.argumentsPart === null) return; + const toolCall = callsByIndex.get(part.index); + if (toolCall === undefined) return; + this.record.signal({ + type: 'tool.call.delta', + turnId, + toolCallId: toolCall.id, + name: toolCall.name, + argumentsPart: part.argumentsPart, + }); + return; + } + default: { + const _exhaustive: never = part; + return _exhaustive; + } + } + }; + } } function isContextOverflowError(error: unknown): boolean { @@ -419,41 +436,6 @@ function deriveStepStopReason(response: LLMRequestFinish): LoopStepStopReason { } } -function createToolCallDeltaHandler( - emit: (event: AgentEvent) => void, - turnId: number, -): (part: StreamedMessagePart) => void { - const callsByIndex = new Map(); - let lastToolCall: { id: string; name: string } | undefined; - - return (part) => { - if (isToolCall(part)) { - lastToolCall = { id: part.id, name: part.name }; - if (part._streamIndex !== undefined) { - callsByIndex.set(part._streamIndex, lastToolCall); - } - emit({ - type: 'tool.call.delta', - turnId, - toolCallId: part.id, - name: part.name, - argumentsPart: part.arguments ?? undefined, - }); - return; - } - if (!isToolCallPart(part) || part.argumentsPart === null) return; - const toolCall = part.index !== undefined ? callsByIndex.get(part.index) : lastToolCall; - if (toolCall === undefined) return; - emit({ - type: 'tool.call.delta', - turnId, - toolCallId: toolCall.id, - name: toolCall.name, - argumentsPart: part.argumentsPart, - }); - }; -} - registerScopedService( LifecycleScope.Agent, IAgentLoopService,