From c799682484bd65b9c9f617cdd2bea3c201d8b0fe Mon Sep 17 00:00:00 2001 From: mdwsk88 <924038395@qq.com> Date: Sat, 26 Sep 2026 09:40:02 +0800 Subject: [PATCH 1/2] fix(codebuddy): serialize parallel tool_use blocks and close same-index reuse MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CodeBuddy CLI parallel tool calls arrive as several tool_use blocks on ONE shared content-block index: intermediate blocks never receive a content_block_stop and only the final block does, while thinking deltas on other indices interleave freely. The single open-call slot mis-attributed argument fragments and miscounted completions (two starts, one counted stop), so the turn failed closed with "Coding-agent CLI ended with an incomplete tool call" 502s at message_stop — visible in Codex Desktop as reconnect banners on CodeBuddy turns that make parallel tool calls. Buffer each block under its own content-block index and emit it atomically (start, deltas, end) when it closes: at its own stop, or when a new tool_use start reuses its index — parallel argument streams are sequential, so the block open on a reused index is complete. Index-less start frames fall back to synthetic keys; index-less delta/stop frames resolve only when exactly one block is open, otherwise the turn still fails closed. --- src/adapters/coding-agent/protocol.ts | 97 ++++++++++++++-- src/adapters/coding-agent/turn.ts | 14 +-- tests/providers/codebuddy-protocol.test.ts | 73 +++++++++++- .../codebuddy-tool-bridge-turn.test.ts | 105 ++++++++++++++++++ 4 files changed, 265 insertions(+), 24 deletions(-) diff --git a/src/adapters/coding-agent/protocol.ts b/src/adapters/coding-agent/protocol.ts index cccbb05db58..09a8e38ea5c 100644 --- a/src/adapters/coding-agent/protocol.ts +++ b/src/adapters/coding-agent/protocol.ts @@ -222,9 +222,23 @@ export interface StreamParseState { sawPartialText: boolean; sawPartialThinking: boolean; sawTerminalResult: boolean; - openToolCallId?: string; /** A `message_stop` stream event arrived: the assistant message is complete. */ sawMessageStop?: boolean; + /** + * Open tool_use blocks keyed by content-block index. CodeBuddy parallel tool calls arrive + * as several tool_use blocks on ONE shared content-block index — intermediate blocks never + * receive a stop and only the final block does — while deltas for different indices (for + * example a long thinking block) interleave freely (observed 2026-09-25/26: four parallel + * calls, one counted stop, "incomplete tool call" 502 at message_stop). A single open-call + * slot both mis-attributes argument fragments and miscounts completions. Downstream + * assembly keeps a single open call, so each block is buffered and emitted atomically: + * when its own stop arrives, or when a new tool_use start reuses its index. + */ + openToolBlocks?: Map; + /** Synthetic decreasing keys for tool_use start frames that omit the block index. */ + nextSyntheticToolBlockKey?: number; + /** Tool_use blocks opened in this stream, whether or not they have closed yet. */ + toolBlockStarts?: number; /** Completed tool_use content blocks observed in this stream. */ completedToolCalls?: number; /** Tool IDs already captured through partial events, for complete-assistant deduplication. */ @@ -332,6 +346,53 @@ export function mapStreamMessageToEvents(message: StreamMessage, state: StreamPa return events; } +/** One in-flight tool_use block: identity plus its buffered argument fragments. */ +export interface OpenToolBlock { + id: string; + name: string; + argParts: string[]; +} + +/** Key a tool_use start frame by content-block index, falling back to a synthetic key. */ +function toolBlockKey(state: StreamParseState, event: StreamMessage): number { + const index = event.index; + if (typeof index === "number" && Number.isInteger(index)) return index; + const key = state.nextSyntheticToolBlockKey ?? -1; + state.nextSyntheticToolBlockKey = key - 1; + return key; +} + +/** + * Resolve a delta/stop frame to an open tool block. An indexed frame only matches a block + * opened under the same index — CodeBuddy skips stop frames for thinking blocks, and such a + * stop must not close a tool block that happens to be open. An index-less frame resolves + * only when exactly one block is open; with several open blocks attribution is unknowable, + * so the frame is dropped and the turn fails closed at the terminal accounting check. + */ +function resolveToolBlockKey(state: StreamParseState, event: StreamMessage): number | undefined { + const index = event.index; + if (typeof index === "number" && Number.isInteger(index)) { + return state.openToolBlocks?.has(index) ? index : undefined; + } + const blocks = state.openToolBlocks; + if (!blocks || blocks.size !== 1) return undefined; + return blocks.keys().next().value; +} + +/** + * Emit a closed block atomically — start, the buffered fragments in arrival order, end — + * so the strictly sequential downstream bridge never sees two calls open at once. + */ +function closeToolBlock(state: StreamParseState, key: number, events: AdapterEvent[]): void { + const block = state.openToolBlocks?.get(key); + if (!block || !state.openToolBlocks) return; + state.openToolBlocks.delete(key); + events.push({ type: "tool_call_start", id: block.id, name: block.name }); + for (const part of block.argParts) events.push({ type: "tool_call_delta", arguments: part }); + events.push({ type: "tool_call_end" }); + state.completedToolCalls = (state.completedToolCalls ?? 0) + 1; +} + /** Map a raw Anthropic SSE event (carried inside a `stream_event` frame) to AdapterEvents. */ function mapRawStreamEvent(event: StreamMessage, state: StreamParseState): AdapterEvent[] { const events: AdapterEvent[] = []; @@ -353,11 +414,16 @@ function mapRawStreamEvent(event: StreamMessage, state: StreamParseState): Adapt events.push({ type: "thinking_delta", thinking }); } } else if (deltaType === "input_json_delta") { - // Tool-input streaming. Live for capture-only bridge turns, where the advertised MCP - // catalog makes the CLI emit real tool_use blocks; parsed unconditionally so a stray - // frame on a tools-disabled turn is ignored rather than crashing. + // Tool-input streaming. Fragments are buffered under their own block index because + // CodeBuddy alternates deltas across interleaved parallel blocks; parsed + // unconditionally so a stray frame on a tools-disabled turn is ignored rather than + // crashing. const partial = asString(delta?.partial_json); - if (partial && state.openToolCallId) events.push({ type: "tool_call_delta", arguments: partial }); + if (partial) { + const key = resolveToolBlockKey(state, event); + const block = key === undefined ? undefined : state.openToolBlocks?.get(key); + if (block) block.argParts.push(partial); + } } return events; } @@ -368,20 +434,27 @@ function mapRawStreamEvent(event: StreamMessage, state: StreamParseState): Adapt const id = asString(block?.id) ?? ""; const name = asString(block?.name) ?? "tool"; if (id) { - state.openToolCallId = id; + const key = toolBlockKey(state, event); + if (state.openToolBlocks?.has(key)) { + // CodeBuddy reuses one content-block index for a parallel batch: every call in the + // batch starts on the same index, intermediate blocks never receive a stop, and only + // the final block does (observed 2026-09-26: START 2 alpha, A's complete args, START 2 + // beta, B's complete args, one STOP 2). Parallel calls stream their arguments + // sequentially — never interleaved — so the block already open on this index is + // complete, and the new start implicitly closes it. + closeToolBlock(state, key, events); + } + (state.openToolBlocks ??= new Map()).set(key, { id, name, argParts: [] }); + state.toolBlockStarts = (state.toolBlockStarts ?? 0) + 1; state.partialToolCallIds?.add(id); - events.push({ type: "tool_call_start", id, name }); } } return events; } if (eventType === "content_block_stop") { - if (state.openToolCallId) { - state.openToolCallId = undefined; - state.completedToolCalls = (state.completedToolCalls ?? 0) + 1; - events.push({ type: "tool_call_end" }); - } + const key = resolveToolBlockKey(state, event); + if (key !== undefined) closeToolBlock(state, key, events); return events; } diff --git a/src/adapters/coding-agent/turn.ts b/src/adapters/coding-agent/turn.ts index 67014960621..f66dfb43836 100644 --- a/src/adapters/coding-agent/turn.ts +++ b/src/adapters/coding-agent/turn.ts @@ -380,7 +380,7 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise() : undefined, }; @@ -527,8 +527,8 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise 0 - && (state.completedToolCalls ?? 0) !== toolCallStarts + && (state.toolBlockStarts ?? 0) > 0 + && (state.completedToolCalls ?? 0) !== (state.toolBlockStarts ?? 0) ) { // A terminal result that arrives while a captured tool call is still open must not // become a successful completion the client can accept. The message_stop check after @@ -552,8 +552,8 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise 0 - && (state.completedToolCalls ?? 0) === toolCallStarts + && (state.toolBlockStarts ?? 0) > 0 + && (state.completedToolCalls ?? 0) === (state.toolBlockStarts ?? 0) ) { // Every captured call completed and the CLI settled with a successful result before // message_stop (instead of parking on the never-answering capture server). Emitting @@ -573,8 +573,8 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise 0 - && (state.completedToolCalls ?? 0) !== toolCallStarts + && (state.toolBlockStarts ?? 0) > 0 + && (state.completedToolCalls ?? 0) !== (state.toolBlockStarts ?? 0) ) { emitOnce({ type: "error", diff --git a/tests/providers/codebuddy-protocol.test.ts b/tests/providers/codebuddy-protocol.test.ts index 3cf7690156b..ffafe2966c4 100644 --- a/tests/providers/codebuddy-protocol.test.ts +++ b/tests/providers/codebuddy-protocol.test.ts @@ -232,20 +232,83 @@ describe("codebuddy stream-json event mapping", () => { }); test("parses tool_use blocks defensively even though v1 disables tools", () => { - const state = { sawPartialText: false, sawPartialThinking: false, sawTerminalResult: false, openToolCallId: undefined as string | undefined }; + const state = { sawPartialText: false, sawPartialThinking: false, sawTerminalResult: false }; const start = mapStreamMessageToEvents( { type: "stream_event", event: { type: "content_block_start", content_block: { type: "tool_use", id: "t1", name: "exec" } } }, state, ); - expect(start).toEqual([{ type: "tool_call_start", id: "t1", name: "exec" }]); + // Tool blocks are buffered and emitted atomically at their own stop: downstream keeps a + // single open call, and CodeBuddy interleaves parallel blocks. + expect(start).toEqual([]); const delta = mapStreamMessageToEvents( { type: "stream_event", event: { type: "content_block_delta", delta: { type: "input_json_delta", partial_json: "{\"a\":1}" } } }, state, ); - expect(delta).toEqual([{ type: "tool_call_delta", arguments: "{\"a\":1}" }]); + expect(delta).toEqual([]); const stop = mapStreamMessageToEvents({ type: "stream_event", event: { type: "content_block_stop" } }, state); - expect(stop).toEqual([{ type: "tool_call_end" }]); - expect(state.openToolCallId).toBeUndefined(); + expect(stop).toEqual([ + { type: "tool_call_start", id: "t1", name: "exec" }, + { type: "tool_call_delta", arguments: "{\"a\":1}" }, + { type: "tool_call_end" }, + ]); + expect(state.completedToolCalls).toBe(1); + expect(state.openToolBlocks?.size ?? 0).toBe(0); + }); + + test("interleaved parallel tool_use blocks are serialized per block index", () => { + const state = { sawPartialText: false, sawPartialThinking: false, sawTerminalResult: false }; + const feed = (event: unknown) => mapStreamMessageToEvents({ type: "stream_event", event: event as Record }, state); + const startAt = (index: number, id: string) => + feed({ type: "content_block_start", index, content_block: { type: "tool_use", id, name: "exec" } }); + const deltaAt = (index: number, part: string) => + feed({ type: "content_block_delta", index, delta: { type: "input_json_delta", partial_json: part } }); + + expect(startAt(1, "tu_a")).toEqual([]); + expect(startAt(2, "tu_b")).toEqual([]); + expect(deltaAt(1, "{\"cmd\":\"a")).toEqual([]); + expect(deltaAt(2, "{\"cmd\":\"b")).toEqual([]); + expect(deltaAt(1, "\"}")).toEqual([]); + // A stop for a non-tool block must not close an open tool block. + expect(feed({ type: "content_block_stop", index: 0 })).toEqual([]); + expect(feed({ type: "content_block_stop", index: 2 })).toEqual([ + { type: "tool_call_start", id: "tu_b", name: "exec" }, + { type: "tool_call_delta", arguments: "{\"cmd\":\"b" }, + { type: "tool_call_end" }, + ]); + expect(feed({ type: "content_block_stop", index: 1 })).toEqual([ + { type: "tool_call_start", id: "tu_a", name: "exec" }, + { type: "tool_call_delta", arguments: "{\"cmd\":\"a" }, + { type: "tool_call_delta", arguments: "\"}" }, + { type: "tool_call_end" }, + ]); + expect(state.toolBlockStarts).toBe(2); + expect(state.completedToolCalls).toBe(2); + }); + + test("a parallel batch reuses one block index; a new start implicitly closes the open block", () => { + // Live capture 2026-09-26 (CodeBuddy 2.158.0, kimi-k3-1, two parallel calls): START idx=2 + // alpha, alpha's complete args, START idx=2 beta (alpha never stopped), beta's complete + // args, one STOP idx=2, message_stop. A start that reuses an open block's index closes + // that block — parallel argument streams are sequential, so the open block is complete. + const state = { sawPartialText: false, sawPartialThinking: false, sawTerminalResult: false }; + const feed = (event: unknown) => mapStreamMessageToEvents({ type: "stream_event", event: event as Record }, state); + + expect(feed({ type: "content_block_start", index: 2, content_block: { type: "tool_use", id: "tu_a", name: "alpha" } })).toEqual([]); + expect(feed({ type: "content_block_delta", index: 2, delta: { type: "input_json_delta", partial_json: "{\"value\":\"A\"}" } })).toEqual([]); + expect(feed({ type: "content_block_start", index: 2, content_block: { type: "tool_use", id: "tu_b", name: "beta" } })).toEqual([ + { type: "tool_call_start", id: "tu_a", name: "alpha" }, + { type: "tool_call_delta", arguments: "{\"value\":\"A\"}" }, + { type: "tool_call_end" }, + ]); + expect(feed({ type: "content_block_delta", index: 2, delta: { type: "input_json_delta", partial_json: "{\"value\":\"B\"}" } })).toEqual([]); + expect(feed({ type: "content_block_stop", index: 2 })).toEqual([ + { type: "tool_call_start", id: "tu_b", name: "beta" }, + { type: "tool_call_delta", arguments: "{\"value\":\"B\"}" }, + { type: "tool_call_end" }, + ]); + expect(state.toolBlockStarts).toBe(2); + expect(state.completedToolCalls).toBe(2); + expect(state.openToolBlocks?.size ?? 0).toBe(0); }); test("usageFromResult returns undefined when no usage is present", () => { diff --git a/tests/providers/codebuddy-tool-bridge-turn.test.ts b/tests/providers/codebuddy-tool-bridge-turn.test.ts index 99e2a241c45..6791e4aeeef 100644 --- a/tests/providers/codebuddy-tool-bridge-turn.test.ts +++ b/tests/providers/codebuddy-tool-bridge-turn.test.ts @@ -333,6 +333,111 @@ describe("CodeBuddy capture-only tool bridge turn", () => { expect(events.some(e => e.type === "done")).toBe(false); }); + test("interleaved parallel tool calls are serialized per block and complete the leg", async () => { + const p = parsed([tool("exec")]); + const bridge = buildCodeBuddyToolBridge(p); + const cliName = [...bridge.emittedNameMap.keys()][0]!; + const wireName = bridge.emittedNameMap.get(cliName)!; + const startAt = (index: number, id: string) => ({ + type: "stream_event", + event: { type: "content_block_start", index, content_block: { type: "tool_use", id, name: cliName } }, + }); + const deltaAt = (index: number, part: string) => ({ + type: "stream_event", + event: { type: "content_block_delta", index, delta: { type: "input_json_delta", partial_json: part } }, + }); + const stopAt = (index: number) => ({ type: "stream_event", event: { type: "content_block_stop", index } }); + let child: FakeChild | undefined; + const spawn: SpawnFn = (_cmd, _args) => { + child = fakeChild(frameLines([ + INIT_OK, + // CodeBuddy emits genuinely interleaved blocks for parallel calls: starts arrive + // before earlier blocks stop and argument deltas alternate across indices (captured + // from the live CLI on 2026-09-25). Single-slot accounting led to a spurious + // "incomplete tool call" 502 at message_stop. + startAt(1, "tu_a"), + startAt(2, "tu_b"), + deltaAt(1, "{\"command\":[\"ec"), + deltaAt(2, "{\"command\":[\"ls"), + deltaAt(1, "ho\"]}"), + deltaAt(2, "\"]}"), + stopAt(2), + stopAt(1), + MESSAGE_STOP, + ])); + return child as unknown as ChildProcess; + }; + const adapter = createCodeBuddyAdapter(provider(), { spawn, which: () => "/usr/bin/codebuddy" }); + const events = await run(adapter, p); + expect(events.map(e => e.type)).toEqual([ + "tool_call_start", + "tool_call_delta", + "tool_call_delta", + "tool_call_end", + "tool_call_start", + "tool_call_delta", + "tool_call_delta", + "tool_call_end", + "done", + ]); + // Blocks are emitted atomically in stop order: tu_b closed first. + expect(events[0]).toMatchObject({ type: "tool_call_start", id: "tu_b", name: wireName }); + expect(events[1]).toMatchObject({ arguments: "{\"command\":[\"ls" }); + expect(events[4]).toMatchObject({ type: "tool_call_start", id: "tu_a", name: wireName }); + expect(events[5]).toMatchObject({ arguments: "{\"command\":[\"ec" }); + expect(events[6]).toMatchObject({ arguments: "ho\"]}" }); + expect(events[8]).toMatchObject({ type: "done", stopReason: "tool_use", endTurn: false }); + expect(child?.killed).toBe(true); + }); + + test("a parallel batch on one shared block index completes every call in the leg", async () => { + const p = parsed([tool("exec")]); + const bridge = buildCodeBuddyToolBridge(p); + const cliName = [...bridge.emittedNameMap.keys()][0]!; + const wireName = bridge.emittedNameMap.get(cliName)!; + const start = (id: string) => ({ + type: "stream_event", + event: { type: "content_block_start", index: 2, content_block: { type: "tool_use", id, name: cliName } }, + }); + const delta = (part: string) => ({ + type: "stream_event", + event: { type: "content_block_delta", index: 2, delta: { type: "input_json_delta", partial_json: part } }, + }); + let child: FakeChild | undefined; + const spawn: SpawnFn = (_cmd, _args) => { + child = fakeChild(frameLines([ + INIT_OK, + // Live capture 2026-09-26 (CodeBuddy 2.158.0, kimi-k3-1): a parallel batch reuses one + // content-block index — alpha starts, streams its complete arguments, then beta starts + // on the same index with no stop for alpha; only the final block receives a stop. + start("tu_a"), + delta("{\"command\":[\"echo\"]}"), + start("tu_b"), + delta("{\"command\":[\"ls\"]}"), + { type: "stream_event", event: { type: "content_block_stop", index: 2 } }, + MESSAGE_STOP, + ])); + return child as unknown as ChildProcess; + }; + const adapter = createCodeBuddyAdapter(provider(), { spawn, which: () => "/usr/bin/codebuddy" }); + const events = await run(adapter, p); + expect(events.map(e => e.type)).toEqual([ + "tool_call_start", + "tool_call_delta", + "tool_call_end", + "tool_call_start", + "tool_call_delta", + "tool_call_end", + "done", + ]); + expect(events[0]).toMatchObject({ type: "tool_call_start", id: "tu_a", name: wireName }); + expect(events[1]).toMatchObject({ arguments: "{\"command\":[\"echo\"]}" }); + expect(events[3]).toMatchObject({ type: "tool_call_start", id: "tu_b", name: wireName }); + expect(events[4]).toMatchObject({ arguments: "{\"command\":[\"ls\"]}" }); + expect(events[6]).toMatchObject({ type: "done", stopReason: "tool_use", endTurn: false }); + expect(child?.killed).toBe(true); + }); + test("a tool call before the init frame fails closed with tool_bridge_init_missing", async () => { const p = parsed([tool("exec")]); const bridge = buildCodeBuddyToolBridge(p); From 684d7aa24e4e74b5198feadf7c211780e7d8420e Mon Sep 17 00:00:00 2001 From: mdwsk88 <924038395@qq.com> Date: Sat, 26 Sep 2026 19:45:40 +0800 Subject: [PATCH 2/2] docs(codebuddy): describe parallel tool block serialization --- structure/providers-and-adapters.md | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index 9871c973ec4..1d6d5911cd0 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -1,5 +1,10 @@ # Providers And Adapters +The coding-agent stream parser buffers each tool-use block by its content-block index +and emits a complete start/delta/end sequence on closure. A new start on an occupied +index closes the previous block; distinct indices can interleave. Turn completion +requires every opened block to close, preserving the downstream single-open-call contract. + RunTurn hosted search uses `src/web-search/run-turn-loop.ts`: synthetic calls remain private, progress reaches the bridge during collection, and a validated terminal precedes search execution. Complete search calls remain actionable at a truncated `done`; cancellation prevents subsequent queries and calls. OAuth preflight replay in `src/server/responses/run-turn-execution.ts` retains the synthetic tool while refreshing credential-scoped route state. In `src/server/responses/sidecar-execution.ts`, a search plan takes priority over image/video bridge execution for both transports; only fetch-capable adapters enter the fetch search loop. Combo preflight allows the private search tool only while a search plan is active; client tool declaration checks and replay-unsafe heartbeat protection remain enforced.