From ab5d74787ac42af503b3ed8ec022d0f5c25b2f6b Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 30 Sep 2026 11:48:52 +0900 Subject: [PATCH 1/3] fix(kiro): settle held text across bounded completion retry --- .../src/content/docs/guides/providers.md | 5 + scripts/test-layout/layout.json | 1 + src/adapters/kiro/stream.ts | 75 ++++++---- structure/providers/kiro.md | 12 ++ tests/fixtures/test-layout-expected.json | 1 + .../providers/kiro/kiro-single-final.test.ts | 138 ++++++++++++++++++ tests/providers/kiro/kiro-stream.test.ts | 4 - .../server/server-kiro-completion-e2e.test.ts | 7 +- 8 files changed, 207 insertions(+), 36 deletions(-) create mode 100644 tests/providers/kiro/kiro-single-final.test.ts diff --git a/docs-site/src/content/docs/guides/providers.md b/docs-site/src/content/docs/guides/providers.md index ee2191f61b5..14f4e0accba 100644 --- a/docs-site/src/content/docs/guides/providers.md +++ b/docs-site/src/content/docs/guides/providers.md @@ -406,6 +406,11 @@ the dashboard Codex account pool also performs. See ### Kiro request credits +On tool-enabled turns, opencodex holds Kiro's ordinary text until completion is validated. +If Kiro ends with plain text instead of its private final-answer tool, one bounded retry +still runs, and only the resulting final answer is displayed. Progress accompanying a real +tool call remains visible. A normal private final answer needs no completion retry. + When Kiro emits credit metering, request logs preserve the reported spend as `usage.providerCredits`, including in the persisted usage ledger. These are Kiro credits; token counts may still be estimated, and the credit value does not replace USD cost estimates. diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 9ffc50ee6d1..08f69e7d5b4 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -964,6 +964,7 @@ "kiro-review-regressions.test.ts": "providers/kiro", "kiro-metering-events.test.ts": "providers/kiro", "kiro-metering-usage.test.ts": "providers/kiro", + "kiro-single-final.test.ts": "providers/kiro", "kiro-stream.test.ts": "providers/kiro", "kiro-transport-parity.test.ts": "providers/kiro", "kiro-usage-quota.test.ts": "providers/kiro", diff --git a/src/adapters/kiro/stream.ts b/src/adapters/kiro/stream.ts index f10cc6fce99..b0490d7870f 100644 --- a/src/adapters/kiro/stream.ts +++ b/src/adapters/kiro/stream.ts @@ -38,6 +38,8 @@ interface KiroAttemptParseResult { } interface KiroAttemptResult extends KiroAttemptParseResult { + drainDeferred(supersededByCompletion?: boolean): AsyncGenerator; + releaseCollectors(): void; releaseRetained(): void; } @@ -45,6 +47,7 @@ interface KiroAttemptRetention { trackReplacement(previousBytes: number, nextBytes: number): void; retainEvent(event: AdapterEvent, bytes: number): void; releaseEvent(event: AdapterEvent): void; + releaseCollectors(): void; releaseAll(): void; } @@ -66,6 +69,11 @@ function createKiroAttemptRetention(budget: TranslatorBudget): KiroAttemptRetent retainedBytes = Math.max(0, retainedBytes - bytes); budget.releaseRetained(bytes, { kind: "retained_collectors" }); }, + releaseCollectors() { + const pendingBytes = [...eventBytes.values()].reduce((sum, bytes) => sum + bytes, 0); + budget.releaseRetained(retainedBytes - pendingBytes, { kind: "retained_collectors" }); + retainedBytes = pendingBytes; + }, releaseAll() { if (retainedBytes > 0) budget.releaseRetained(retainedBytes, { kind: "retained_collectors" }); retainedBytes = 0; @@ -228,13 +236,20 @@ async function* parseKiroAttempt( nameMap: Map | undefined, conversationId: string | undefined, contextInputEstimate?: number, - /** True when an earlier attempt already flushed visible content to the client (#520). */ + /** True when an earlier attempt has output that must survive a failed retry. */ priorEmittedOutput = false, + priorAttempt?: KiroAttemptResult, ): AsyncGenerator { - // `required` mode holds staged commentary until a real tool call or terminal metadata identifies - // the attempt boundary. Anything the inner parser leaves behind is flushed before the terminal. + // Hold commentary through completion validation; tools and failures still release progress. const deferred: AdapterEvent[] = []; const retention = createKiroAttemptRetention(budget); + const drainDeferred = async function* (supersededByCompletion = false): AsyncGenerator { + for (const event of deferred.splice(0)) { + try { + if (!supersededByCompletion || event.type !== "text_delta") yield event; + } finally { retention.releaseEvent(event); } + } + }; // Shared box: the inner parser stages its calibration observation here on the completion path, // and this wrapper decides whether the attempt was terminal enough to commit it. A box rather // than a return field because the completion path has a dozen terminal returns and threading a @@ -254,6 +269,7 @@ async function* parseKiroAttempt( attemptCalibration, contextInputEstimate, priorEmittedOutput, + priorAttempt, ); let handedOff = false; try { @@ -267,11 +283,14 @@ async function* parseKiroAttempt( if (staged && !result.needsFallback) { recordKiroCalibration(staged.conversationId, staged.estimated, staged.charged); } - for (const event of deferred.splice(0)) { - try { yield event; } finally { retention.releaseEvent(event); } - } + if (priorAttempt) yield* priorAttempt.drainDeferred(); + if (!result.needsFallback) yield* drainDeferred(); handedOff = true; - return { ...result, releaseRetained: () => retention.releaseAll() }; + return { + ...result, drainDeferred, + releaseCollectors: () => retention.releaseCollectors(), + releaseRetained: () => retention.releaseAll(), + }; } finally { if (!handedOff) retention.releaseAll(); } @@ -291,6 +310,7 @@ async function* parseKiroAttemptEvents( attemptCalibration: { value?: { conversationId: string; estimated: number; charged: number } }, contextInputEstimate?: number, priorEmittedOutput = false, + priorAttempt?: KiroAttemptResult, ): AsyncGenerator { const emptyResult = (): KiroAttemptParseResult => ({ assistantText: "", sawReasoning: false }); // Every early return below is a failure path that stages nothing; only the completion path @@ -343,10 +363,8 @@ async function* parseKiroAttemptEvents( // (#2819 follow-up). Consume the collection instead — drop the redundant text, keep every // non-text event, and release retention either way. // - // This is deliberately the ONLY suppression site. The outer drain in `parseKiroAttempt` is also - // the leftover flush for early terminal returns (stream, protocol, and provider failures), so - // teaching it to discard text would hide the only commentary a failed turn ever produced. - // Splicing here leaves that drain empty on the completion path and untouched everywhere else. + // The preceding attempt uses the same rule when bounded validation succeeds. Failures retain + // the ordinary leftover flush so a failed turn's only progress is still delivered. const consumeSupersededByCompletion = async function* ( events: AdapterEvent[], ): AsyncGenerator { @@ -421,7 +439,7 @@ async function* parseKiroAttemptEvents( message, usage(), providerState(), - // First-attempt progress was already flushed before this bounded fallback (#520). + // Failed validation releases first-attempt progress before the terminal. !priorEmittedOutput, ); } @@ -469,7 +487,7 @@ async function* parseKiroAttemptEvents( // In `required` mode Kiro's stop reason only arrives on the terminal metadata event, so staged // commentary is held until either a real tool call proves the turn continues (flush as - // commentary) or the stream ends (relabel as the final answer when END_TURN says so). A heartbeat + // commentary) or bounded validation settles the held text. A heartbeat // stands in for each held event so the bridge's stall watchdog stays armed. const defer = (event: AdapterEvent): AdapterEvent[] => { if (sawRealTool) return [...deferred.splice(0), event]; @@ -802,14 +820,14 @@ async function* parseKiroAttemptEvents( assistantChars: assistantText.length, }); - if (mode === "required") { - // A valid completion answer makes this inference's staged prose redundant; anything else - // still flushes exactly as before (bounded fallback, explicit stops, real tool calls). - if (completionAnswer !== undefined) yield* consumeSupersededByCompletion(deferred); - else yield* emitRetained(deferred.splice(0)); + if (mode === "required" && completionAnswer !== undefined) { + yield* consumeSupersededByCompletion(deferred); } if (mode === "text_fallback") { + if (priorAttempt) { + yield* priorAttempt.drainDeferred(completionAnswer !== undefined || (sawText && !sawRealTool)); + } if (completionAnswer !== undefined) { yield* consumeSupersededByCompletion(fallbackEvents); yield { type: "text_delta", text: completionAnswer, phase: "final_answer" }; @@ -853,7 +871,7 @@ async function* parseKiroAttemptEvents( : "Kiro produced no final answer on its bounded completion retry", finalUsage, finalProviderState, - // First-attempt progress was already flushed before this bounded fallback (#520). + // Failed validation releases first-attempt progress before the terminal. !priorEmittedOutput, ), }; @@ -1047,6 +1065,7 @@ export async function* parseKiroStream( return; } if (!fallbackFactory) { + yield* firstResult.drainDeferred(); yield retryableKiroIncomplete( "uncompleted_kiro_response", "Kiro produced progress without an explicit final answer and no bounded retry transport was available", @@ -1057,9 +1076,8 @@ export async function* parseKiroStream( } yield { type: "heartbeat" }; - // First attempt already flushed deferred progress before this point. Gate fallback - // setup/HTTP failures the same way as the second-stream catch so a replay cannot - // duplicate visible commentary (#520). + // Failed validation releases held progress. Keep those failures non-retryable so a later + // replay cannot duplicate it; successful validation instead discards the superseded text. const priorEmittedOutput = Boolean(firstResult.assistantText.trim()) || firstResult.sawReasoning; let firstAssistantText = firstResult.assistantText; const firstHadAssistantText = firstAssistantText.length > 0; @@ -1072,6 +1090,7 @@ export async function* parseKiroStream( budget, ); } catch (err) { + yield* firstResult.drainDeferred(); firstAssistantText = ""; firstResult.assistantText = ""; firstResult.releaseRetained(); @@ -1096,14 +1115,14 @@ export async function* parseKiroStream( }; return; } - // The factory has finished using the live first-attempt alias and has retained its own retry - // serialization through the fetch boundary. The discarded parser collectors can now release - // before the second attempt begins on the same turn budget. + // The factory has retained its retry serialization. First-attempt progress remains charged + // until the second attempt decides whether it is superseded or must be released. firstAssistantText = ""; firstResult.assistantText = ""; - firstResult.releaseRetained(); + firstResult.releaseCollectors(); fallback.releaseRequestBody?.(); if (!fallback.response.ok) { + yield* firstResult.drainDeferred(); const payload = await readDisplaySafeErrorPayloadText(fallback.response, fallback.abortSignal); const failure = classifyKiroHttpError(fallback.response.status, fallback.response.headers, payload); yield { @@ -1128,9 +1147,9 @@ export async function* parseKiroStream( fallback.nameMap, fallback.conversationId, fallback.contextInputEstimate, - // First attempt already flushed deferred progress to the client before this fallback. - // A zero-output transport failure here must stay non-retryable to avoid duplicating that text. + // Failed validation will release the held first-attempt progress. priorEmittedOutput, + firstResult, ); try { if (!secondResult.terminal) { diff --git a/structure/providers/kiro.md b/structure/providers/kiro.md index afb64a91b07..47981213ec1 100644 --- a/structure/providers/kiro.md +++ b/structure/providers/kiro.md @@ -122,6 +122,18 @@ raw body. ## Bounded fallback HTTP errors +Tool-enabled turns in `src/adapters/kiro/stream.ts` hold ordinary text through the one +bounded completion retry. A valid private final answer or accepted retry text supersedes +first-attempt prose, so the client receives one final answer. A real tool call releases +held progress as commentary before the tool; failed validation also releases progress +and preserves the non-retryable boundary. Held events stay charged to the translator +budget until emitted, discarded, or cancelled; replay collectors are released after +retry construction. Native `END_TURN` and `STOP_SEQUENCE` alone do not distinguish +progress from an answer and therefore still require validation. Normal private completion +and real tool calls need no completion retry. +Coverage: `tests/providers/kiro/kiro-single-final.test.ts` and +`tests/server/server-kiro-completion-e2e.test.ts`. + `src/adapters/kiro-retry.ts` uses the configured executor for every generation send and may try the existing `q.{region}.amazonaws.com` host once after a canonical-host HTTP 502/503/504 before output, subject to the same send budget. Reset, 429, alternate, and completion-fallback sends wait for a pacing slot; only the first send is pre-paid. Kiro web-search turns are paced as well. A Kiro-local wrapper maps its header deadline to HTTP 504 without changing shared or Google fetch behavior; caller cancellation remains an abort. Final HTTP 5xx text is fixed for clients, and opt-in provider diagnostics carry only closed-set status and classification codes. When a first Kiro stream needs a completion fallback, the fallback response's non-success diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 607da16a816..0617e37d721 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -972,6 +972,7 @@ "kiro-review-regressions.test.ts": "providers/kiro", "kiro-metering-events.test.ts": "providers/kiro", "kiro-metering-usage.test.ts": "providers/kiro", + "kiro-single-final.test.ts": "providers/kiro", "kiro-stream.test.ts": "providers/kiro", "kiro-transport-parity.test.ts": "providers/kiro", "kiro-usage-quota.test.ts": "providers/kiro", diff --git a/tests/providers/kiro/kiro-single-final.test.ts b/tests/providers/kiro/kiro-single-final.test.ts new file mode 100644 index 00000000000..717f9bce29d --- /dev/null +++ b/tests/providers/kiro/kiro-single-final.test.ts @@ -0,0 +1,138 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { createKiroAdapter } from "../../../src/adapters/kiro"; +import { KIRO_COMPLETION_TOOL_NAME } from "../../../src/adapters/kiro-constants"; +import { resetKiroThrottleStateForTests } from "../../../src/adapters/kiro-retry"; +import { encodeMessage } from "../../../src/lib/eventstream-decoder"; +import { createTranslatorBudget } from "../../../src/lib/translator-budget"; +import type { AdapterEvent, OcxParsedRequest, OcxProviderConfig } from "../../../src/types"; + +const provider: OcxProviderConfig = { + adapter: "kiro", baseUrl: "https://runtime.us-east-1.kiro.dev", apiKey: "ksk_test", +}; +const parsed: OcxParsedRequest = { + modelId: "claude-opus-5.5", stream: true, options: {}, + context: { + messages: [{ role: "user", content: "Inspect the workspace." }], + tools: [{ name: "bash", description: "Run a command", parameters: { type: "object" } }], + }, +}; +const frame = (type: string, payload: object) => encodeMessage( + { ":message-type": "event", ":event-type": type }, + new TextEncoder().encode(JSON.stringify(payload)), +); +const text = (content: string) => frame("assistantResponseEvent", { content }); +function tool(name: string, input: object): Uint8Array[] { + return [ + frame("toolUseEvent", { name, toolUseId: "call-1", input: JSON.stringify(input) }), + frame("toolUseEvent", { name, toolUseId: "call-1", stop: true }), + ]; +} +const completion = (answer: string) => tool(KIRO_COMPLETION_TOOL_NAME, { answer }); +function response(frames: Uint8Array[]): Response { + return new Response(new ReadableStream({ + start(controller) { + for (const value of frames) controller.enqueue(value); + controller.close(); + }, + })); +} +afterEach(resetKiroThrottleStateForTests); + +async function run(first: Uint8Array[], retry: Uint8Array[], buffered = false) { + const adapter = createKiroAdapter(provider); + const budget = createTranslatorBudget(); + const sends: number[] = []; + const events: AdapterEvent[] = []; + let visibleAtRetry: AdapterEvent[] = []; + let physicalRequests = 0; + try { + const request = await adapter.buildRequest(structuredClone(parsed)); + const upstream = await adapter.fetchResponse!(request, { + executor: (async () => { + if (++physicalRequests === 1) return response(first); + visibleAtRetry = events.filter(event => event.type === "text_delta"); + return response(retry); + }) as typeof fetch, + onPhysicalSend: send => { sends.push(send.ordinal); }, + }); + if (buffered) events.push(...await adapter.parseResponse!(upstream, budget)); + else for await (const event of adapter.parseStream(upstream, budget)) events.push(event); + if (!buffered) expect(budget.snapshot().currentBytes).toBe(0); + return { events, sends, physicalRequests, visibleAtRetry }; + } finally { + budget.dispose(); + } +} + +describe("Kiro single final answer (#6270)", () => { + for (const buffered of [false, true]) { + test.each(["END_TURN", "STOP_SEQUENCE", undefined])( + `plain text ending is held through validation (stop=%s, buffered=${buffered})`, + async stopReason => { + const answer = "The workspace is ready."; + const { events, sends, physicalRequests, visibleAtRetry } = await run( + [text("The workspace "), text("is ready."), + ...(stopReason ? [frame("metadataEvent", { stopReason })] : [])], + [text(answer), ...completion(answer)], + buffered, + ); + expect(events.filter(event => event.type === "text_delta")).toEqual([ + { type: "text_delta", text: answer, phase: "final_answer" }, + ]); + expect(events.at(-1)).toMatchObject({ type: "done", endTurn: true }); + expect(physicalRequests).toBe(2); + expect(sends).toEqual([1, 2]); + expect(visibleAtRetry).toEqual([]); + }, + ); + } + + test("a real tool ending releases genuine progress without a completion retry", async () => { + const { events, sends, physicalRequests } = await run( + [text("Checking the workspace."), ...tool("bash", { command: "pwd" })], [], + ); + expect(events.filter(event => event.type !== "heartbeat").map(event => event.type)) + .toEqual(["text_delta", "tool_call_start", "tool_call_delta", "tool_call_end", "done"]); + expect(events.find(event => event.type === "text_delta")).toEqual({ + type: "text_delta", text: "Checking the workspace.", phase: "commentary", + }); + expect(events.at(-1)).toMatchObject({ type: "done", endTurn: false }); + expect(physicalRequests).toBe(1); + expect(sends).toEqual([1]); + }); + + test("normal private final_answer supersedes prose without a completion retry", async () => { + const answer = "The workspace is ready."; + const { events, sends, physicalRequests } = await run([text(answer), ...completion(answer)], []); + expect(events.filter(event => event.type === "text_delta")).toEqual([ + { type: "text_delta", text: answer, phase: "final_answer" }, + ]); + expect(events.at(-1)).toMatchObject({ type: "done", endTurn: true }); + expect(physicalRequests).toBe(1); + expect(sends).toEqual([1]); + }); + + test("a retry tool call releases first-attempt progress before the tool", async () => { + const { events } = await run([text("Checking the workspace.")], tool("bash", { command: "pwd" })); + expect(events.filter(event => event.type !== "heartbeat").map(event => event.type)) + .toEqual(["text_delta", "tool_call_start", "tool_call_delta", "tool_call_end", "done"]); + expect(events.at(-1)).toMatchObject({ type: "done", endTurn: false }); + }); + + test("an accepted plain-text retry also replaces first-attempt text", async () => { + const { events } = await run([text("The workspace is ready.")], [text("The workspace is ready.")]); + expect(events.filter(event => event.type === "text_delta")).toEqual([ + { type: "text_delta", text: "The workspace is ready.", phase: "final_answer" }, + ]); + expect(events.at(-1)).toMatchObject({ type: "done", endTurn: true }); + }); + + test("an empty retry preserves held progress and stays non-retryable", async () => { + const { events, physicalRequests } = await run([text("Checking the workspace.")], []); + expect(events.filter(event => event.type === "text_delta")).toEqual([ + { type: "text_delta", text: "Checking the workspace.", phase: "commentary" }, + ]); + expect(events.at(-1)).toMatchObject({ type: "incomplete", retryable: false, endTurn: false }); + expect(physicalRequests).toBe(2); + }); +}); diff --git a/tests/providers/kiro/kiro-stream.test.ts b/tests/providers/kiro/kiro-stream.test.ts index 5ad687ef6f7..c09fac50430 100644 --- a/tests/providers/kiro/kiro-stream.test.ts +++ b/tests/providers/kiro/kiro-stream.test.ts @@ -359,7 +359,6 @@ describe("kiro adapter — parseStream", () => { (tool: { toolSpecification: { name: string } }) => tool.toolSpecification.name, )).toEqual(["bash", KIRO_COMPLETION_TOOL_NAME]); expect(events.filter(event => event.type === "text_delta")).toEqual([ - { type: "text_delta", text: "I am checking.", phase: "commentary" }, { type: "text_delta", text: "Final from fallback.", phase: "final_answer" }, ]); expect(events.at(-1)).toMatchObject({ @@ -617,8 +616,6 @@ describe("kiro adapter — parseStream", () => { expect(fetches).toBe(1); expect(events.filter(event => event.type === "text_delta")).toEqual([ - { type: "text_delta", text: "The file has ", phase: "commentary" }, - { type: "text_delta", text: "three lines.", phase: "commentary" }, { type: "text_delta", text: "The file has three lines.", phase: "final_answer" }, ]); expect(events.at(-1)).toMatchObject({ type: "done", endTurn: true }); @@ -689,7 +686,6 @@ describe("kiro adapter — parseStream", () => { )))); expect(events.filter(event => event.type === "text_delta")).toEqual([ - { type: "text_delta", text: "Done.", phase: "commentary" }, { type: "text_delta", text: "Done.", phase: "final_answer" }, ]); expect(fetches).toBe(1); diff --git a/tests/server/server-kiro-completion-e2e.test.ts b/tests/server/server-kiro-completion-e2e.test.ts index 4d35b2100d3..12f0566ddeb 100644 --- a/tests/server/server-kiro-completion-e2e.test.ts +++ b/tests/server/server-kiro-completion-e2e.test.ts @@ -131,7 +131,7 @@ function anthropicEvents(sse: string): Array<{ name: string; data: Record { - test("/v1/responses keeps progress nonterminal and lets only the bounded fallback complete", async () => { + test("/v1/responses releases only the final answer after bounded validation", async () => { const upstream = scriptedKiroUpstream([ [textFrame("Checking the workspace."), eventFrame("meteringEvent", { unit: "credit", usage: 0.04582331509121062 })], [...completionFrames("The workspace is ready."), eventFrame("meteringEvent", { unit: "credit", amount: 0.01 })], @@ -155,14 +155,13 @@ describe("Kiro completion through public server endpoints", () => { const events = responseEvents(wire); const text = events.filter(event => event.name === "response.output_text.delta"); expect(text.map(event => [event.data.delta, event.data.phase])).toEqual([ - ["Checking the workspace.", undefined], ["The workspace is ready.", undefined], ]); const completed = events.filter(event => event.name === "response.completed"); expect(completed).toHaveLength(1); expect(events.at(-1)?.name).toBe("response.completed"); const messages = completed[0].data.response.output.filter((item: { type: string }) => item.type === "message"); - expect(messages.map((item: { phase?: string }) => item.phase)).toEqual(["commentary", "final_answer"]); + expect(messages.map((item: { phase?: string }) => item.phase)).toEqual(["final_answer"]); expect(wire).not.toContain(KIRO_COMPLETION_TOOL_NAME); const expectedCredits = 0.04582331509121062 + 0.01; @@ -212,7 +211,7 @@ describe("Kiro completion through public server endpoints", () => { const deltas = events .filter(event => event.name === "content_block_delta" && event.data.delta?.type === "text_delta") .map(event => event.data.delta.text); - expect(deltas).toEqual(["I am checking the Claude task.", "The Claude task is complete."]); + expect(deltas).toEqual(["The Claude task is complete."]); expect(events.filter(event => event.name === "message_delta")).toHaveLength(1); expect(events.find(event => event.name === "message_delta")?.data.delta.stop_reason).toBe("end_turn"); expect(events.at(-1)?.name).toBe("message_stop"); From 81f2ed057bd27ae7aa9c6106b0b5366d7cd9210c Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 30 Sep 2026 12:10:32 +0900 Subject: [PATCH 2/3] fix(kiro): release retry tool progress before EOF --- src/adapters/kiro/stream.ts | 2 + .../providers/kiro/kiro-single-final.test.ts | 42 +++++++++++++++++++ 2 files changed, 44 insertions(+) diff --git a/src/adapters/kiro/stream.ts b/src/adapters/kiro/stream.ts index b0490d7870f..dbd39e83764 100644 --- a/src/adapters/kiro/stream.ts +++ b/src/adapters/kiro/stream.ts @@ -726,6 +726,7 @@ async function* parseKiroAttemptEvents( if (ev.stop === true) { const flushed = flushOpen(); if (flushed.terminal) return { assistantText, sawReasoning, terminal: flushed.terminal }; + if (priorAttempt && flushed.events.length) yield* priorAttempt.drainDeferred(); for (const event of flushed.events) { yield* emitRetained(stage(event)); } @@ -762,6 +763,7 @@ async function* parseKiroAttemptEvents( } const flushed = flushOpen(); if (flushed.terminal) return { assistantText, sawReasoning, terminal: flushed.terminal }; + if (priorAttempt && flushed.events.length) yield* priorAttempt.drainDeferred(); for (const event of flushed.events) { yield* emitRetained(stage(event)); } diff --git a/tests/providers/kiro/kiro-single-final.test.ts b/tests/providers/kiro/kiro-single-final.test.ts index 717f9bce29d..2e7fd2d3143 100644 --- a/tests/providers/kiro/kiro-single-final.test.ts +++ b/tests/providers/kiro/kiro-single-final.test.ts @@ -119,6 +119,48 @@ describe("Kiro single final answer (#6270)", () => { expect(events.at(-1)).toMatchObject({ type: "done", endTurn: false }); }); + test("a complete retry tool releases held progress while its stream is still open", async () => { + let releaseEOF!: () => void; + const eof = new Promise(resolve => { releaseEOF = resolve; }); + let reachedOpenStream!: () => void; + const openStream = new Promise(resolve => { reachedOpenStream = resolve; }); + const frames = tool("bash", { command: "pwd" }); + const retry = new Response(new ReadableStream({ + async pull(controller) { + const next = frames.shift(); + if (next) { controller.enqueue(next); return; } + reachedOpenStream(); + await eof; + controller.close(); + }, + }, { highWaterMark: 0 })); + const adapter = createKiroAdapter(provider); + const budget = createTranslatorBudget(); + const events: AdapterEvent[] = []; + let physicalRequests = 0; + try { + const request = await adapter.buildRequest(structuredClone(parsed)); + const first = await adapter.fetchResponse!(request, { + executor: (async () => ++physicalRequests === 1 + ? response([text("Checking the workspace.")]) : retry) as typeof fetch, + }); + const draining = (async () => { + for await (const event of adapter.parseStream(first, budget)) { + events.push(event); + } + })(); + try { + await openStream; + expect(events.filter(event => event.type === "text_delta")).toEqual([ + { type: "text_delta", text: "Checking the workspace.", phase: "commentary" }, + ]); + } finally { releaseEOF(); await draining; } + expect(events.at(-1)).toMatchObject({ type: "done", endTurn: false }); + expect(physicalRequests).toBe(2); + expect(budget.snapshot().currentBytes).toBe(0); + } finally { budget.dispose(); } + }); + test("an accepted plain-text retry also replaces first-attempt text", async () => { const { events } = await run([text("The workspace is ready.")], [text("The workspace is ready.")]); expect(events.filter(event => event.type === "text_delta")).toEqual([ From c3b7fc5c09b0be10d312e538c79edc8ed8cfd0ec Mon Sep 17 00:00:00 2001 From: JUN Date: Wed, 30 Sep 2026 12:18:16 +0900 Subject: [PATCH 3/3] test(kiro): assert buffered retention after event release --- tests/providers/kiro/kiro-single-final.test.ts | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/tests/providers/kiro/kiro-single-final.test.ts b/tests/providers/kiro/kiro-single-final.test.ts index 2e7fd2d3143..a4bac94faed 100644 --- a/tests/providers/kiro/kiro-single-final.test.ts +++ b/tests/providers/kiro/kiro-single-final.test.ts @@ -3,7 +3,7 @@ import { createKiroAdapter } from "../../../src/adapters/kiro"; import { KIRO_COMPLETION_TOOL_NAME } from "../../../src/adapters/kiro-constants"; import { resetKiroThrottleStateForTests } from "../../../src/adapters/kiro-retry"; import { encodeMessage } from "../../../src/lib/eventstream-decoder"; -import { createTranslatorBudget } from "../../../src/lib/translator-budget"; +import { createTranslatorBudget, releaseTranslatedEvent } from "../../../src/lib/translator-budget"; import type { AdapterEvent, OcxParsedRequest, OcxProviderConfig } from "../../../src/types"; const provider: OcxProviderConfig = { @@ -57,7 +57,8 @@ async function run(first: Uint8Array[], retry: Uint8Array[], buffered = false) { }); if (buffered) events.push(...await adapter.parseResponse!(upstream, budget)); else for await (const event of adapter.parseStream(upstream, budget)) events.push(event); - if (!buffered) expect(budget.snapshot().currentBytes).toBe(0); + if (buffered) for (const event of events) releaseTranslatedEvent(event, budget); + expect(budget.snapshot().currentBytes).toBe(0); return { events, sends, physicalRequests, visibleAtRetry }; } finally { budget.dispose();