Skip to content
Merged
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
5 changes: 5 additions & 0 deletions docs-site/src/content/docs/guides/providers.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
1 change: 1 addition & 0 deletions scripts/test-layout/layout.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
77 changes: 49 additions & 28 deletions src/adapters/kiro/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,13 +38,16 @@ interface KiroAttemptParseResult {
}

interface KiroAttemptResult extends KiroAttemptParseResult {
drainDeferred(supersededByCompletion?: boolean): AsyncGenerator<AdapterEvent>;
releaseCollectors(): void;
releaseRetained(): void;
}

interface KiroAttemptRetention {
trackReplacement(previousBytes: number, nextBytes: number): void;
retainEvent(event: AdapterEvent, bytes: number): void;
releaseEvent(event: AdapterEvent): void;
releaseCollectors(): void;
releaseAll(): void;
}

Expand All @@ -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;
Expand Down Expand Up @@ -228,13 +236,20 @@ async function* parseKiroAttempt(
nameMap: Map<string, string> | 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<AdapterEvent, KiroAttemptResult> {
// `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<AdapterEvent> {
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
Expand All @@ -254,6 +269,7 @@ async function* parseKiroAttempt(
attemptCalibration,
contextInputEstimate,
priorEmittedOutput,
priorAttempt,
);
let handedOff = false;
try {
Expand All @@ -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();
}
Expand All @@ -291,6 +310,7 @@ async function* parseKiroAttemptEvents(
attemptCalibration: { value?: { conversationId: string; estimated: number; charged: number } },
contextInputEstimate?: number,
priorEmittedOutput = false,
priorAttempt?: KiroAttemptResult,
): AsyncGenerator<AdapterEvent, KiroAttemptParseResult> {
const emptyResult = (): KiroAttemptParseResult => ({ assistantText: "", sawReasoning: false });
// Every early return below is a failure path that stages nothing; only the completion path
Expand Down Expand Up @@ -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<AdapterEvent> {
Expand Down Expand Up @@ -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,
);
}
Expand Down Expand Up @@ -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];
Expand Down Expand Up @@ -708,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));
}
Expand Down Expand Up @@ -744,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));
}
Expand Down Expand Up @@ -802,14 +822,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));
}
Comment thread
lidge-jun marked this conversation as resolved.
if (completionAnswer !== undefined) {
yield* consumeSupersededByCompletion(fallbackEvents);
yield { type: "text_delta", text: completionAnswer, phase: "final_answer" };
Expand Down Expand Up @@ -853,7 +873,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,
),
};
Expand Down Expand Up @@ -1047,6 +1067,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",
Expand All @@ -1057,9 +1078,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;
Expand All @@ -1072,6 +1092,7 @@ export async function* parseKiroStream(
budget,
);
} catch (err) {
yield* firstResult.drainDeferred();
firstAssistantText = "";
firstResult.assistantText = "";
firstResult.releaseRetained();
Expand All @@ -1096,14 +1117,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 {
Expand All @@ -1128,9 +1149,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) {
Expand Down
12 changes: 12 additions & 0 deletions structure/providers/kiro.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions tests/fixtures/test-layout-expected.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading
Loading