Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
97 changes: 85 additions & 12 deletions src/adapters/coding-agent/protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<number, OpenToolBlock>;
/** 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. */
Expand Down Expand Up @@ -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[] = [];
Expand All @@ -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;
}
Expand All @@ -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;
}

Expand Down
14 changes: 7 additions & 7 deletions src/adapters/coding-agent/turn.ts
Original file line number Diff line number Diff line change
Expand Up @@ -380,7 +380,7 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise<v
sawPartialText: false,
sawPartialThinking: false,
sawTerminalResult: false,
openToolCallId: undefined,
openToolBlocks: new Map(),
partialToolCallIds: toolBridge ? new Set<string>() : undefined,
};

Expand Down Expand Up @@ -527,8 +527,8 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise<v
toolBridge
&& !terminalEmitted
&& event.type === "done"
&& toolCallStarts > 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
Expand All @@ -552,8 +552,8 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise<v
toolBridge
&& !terminalEmitted
&& event.type === "done"
&& toolCallStarts > 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
Expand All @@ -573,8 +573,8 @@ export async function runCodingAgentTurn(input: CodingAgentTurnInput): Promise<v
toolBridge
&& !terminalEmitted
&& state.sawMessageStop
&& toolCallStarts > 0
&& (state.completedToolCalls ?? 0) !== toolCallStarts
&& (state.toolBlockStarts ?? 0) > 0
&& (state.completedToolCalls ?? 0) !== (state.toolBlockStarts ?? 0)
) {
emitOnce({
type: "error",
Expand Down
5 changes: 5 additions & 0 deletions structure/providers-and-adapters.md
Original file line number Diff line number Diff line change
@@ -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.
Expand Down
73 changes: 68 additions & 5 deletions tests/providers/codebuddy-protocol.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> }, 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<string, unknown> }, 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", () => {
Expand Down
Loading
Loading