fix(kiro): improve direct session cache reuse

Reshape Kiro direct requests so resumed client sessions reuse Kiro's
cache-affinity fields instead of starting unrelated CodeWhisperer
conversations.

- keep conversationState.conversationId stable when the client sends an
  explicit session id (x-session-id, session_id, conversation_id, Claude
  Code session metadata)
- add a stable conversationState.agentContinuationId per Kiro session
- send conversationState.agentTaskType: "vibe" and agentMode: "vibe",
  matching the normal Kiro CLI/KAS chat path
- move Kiro thinking instructions into Kiro-compatible systemPrompt /
  additionalModelRequestFields instead of generic top-level thinking
- keep volatile timestamp context out of the top-level systemPrompt; it
  remains only in user content fallback
- suppress additionalModelRequestFields for legacy 4.5-era Claude/Kiro
  models that reject it, while defaulting future Claude/Kiro model ids
  to supported
- preserve Kiro meteringEvent credit usage internally for accounting
  without leaking provider-specific fields into OpenAI-compatible usage
- prevent unrelated headerless Kiro requests from sharing one
  connection-wide continuation
- cap/evict continuation sessions so long-running processes do not grow
  the continuation map unbounded
- treat generated headerless Kiro sessions as one-shot so they do not
  evict real explicit-session continuations
- keep credit-only Kiro metering valid for internal persistence when
  token metrics are unavailable
This commit is contained in:
Edison42
2026-07-16 15:13:28 +07:00
committed by decolua
parent 70e8dc4974
commit 9c58ba645e
10 changed files with 786 additions and 90 deletions

View File

@@ -131,6 +131,50 @@ export function resolveKiroThinkingBudget(body, headers, model) {
return null;
}
export function extractKiroEffortLevel(body) {
const effort =
body?.output_config?.effort ??
body?.reasoning_effort ??
(typeof body?.reasoning === "object" ? body.reasoning?.effort : null);
if (typeof effort !== "string") return null;
const normalized = effort.toLowerCase();
if (normalized === "none" || normalized === "off" || normalized === "disabled") return null;
if (normalized === "xhigh" || normalized === "max") return "high";
if (["low", "medium", "high"].includes(normalized)) return normalized;
return null;
}
export function buildKiroAdditionalModelRequestFields(body) {
const effort = extractKiroEffortLevel(body);
if (!effort) return undefined;
// Mirrors Kiro CLI/KAS buildEffortRequestFields("output_config").
return {
thinking: { type: "adaptive", display: "summarized" },
output_config: { effort },
};
}
export function supportsKiroAdditionalModelRequestFields(model) {
if (typeof model !== "string") return false;
const normalized = model.toLowerCase().replace(/-/g, ".");
if (!normalized.includes("claude")) return false;
const match = normalized.match(/(?:^|[/.])claude(?:[/.][a-z]+)*[/.](\d+)(?:[/.](\d+))?(?:[/.]|$)/);
if (!match) return false;
const [, majorText, minorText] = match;
const major = Number(majorText);
const minor = minorText === undefined ? null : Number(minorText);
const dateSuffixMinor = minor !== null && minor >= 1000;
// Kiro rejected additionalModelRequestFields on legacy 4.5 models in live smoke.
// Default future Claude/Kiro models to supported so new model releases do not
// need a code allowlist update.
return !(major < 4 || (major === 4 && (minor === null || minor <= 5 || dateSuffixMinor)));
}
export function buildKiroAdditionalModelRequestFieldsForModel(body, model) {
if (!supportsKiroAdditionalModelRequestFields(model)) return undefined;
return buildKiroAdditionalModelRequestFields(body);
}
/**
* Detect whether an inbound request is asking for reasoning / thinking output.
* Thin wrapper over resolveKiroThinkingBudget (single source of truth).

View File

@@ -103,8 +103,16 @@ export function translateRequest(sourceFormat, targetFormat, model, body, stream
}
}
// Normalize thinking to the target provider-native format (config-driven, capability-aware)
applyThinking(targetFormat, model, result, provider, thinkingIntent);
// Normalize thinking to the target provider-native format (config-driven, capability-aware).
// Kiro's GenerateAssistantResponse request does not accept the generic top-level
// `thinking` field; its translators map thinking intent to KAS-compatible
// systemPrompt/additionalModelRequestFields instead.
const kiroThinkingMappedByTranslator =
targetFormat === FORMATS.KIRO &&
(sourceFormat === FORMATS.OPENAI || sourceFormat === FORMATS.CLAUDE);
if (!kiroThinkingMappedByTranslator) {
applyThinking(targetFormat, model, result, provider, thinkingIntent);
}
// Always normalize to clean OpenAI format when target is OpenAI
// This handles hybrid requests (e.g., OpenAI messages + Claude tools)

View File

@@ -24,13 +24,15 @@
*/
import { register } from "../index.js";
import { FORMATS } from "../formats.js";
import { v4 as uuidv4 } from "uuid";
import { applyKiroSessionReplay } from "../../utils/kiroSessionReplay.js";
import { resolveContinuationId, resolveSessionIdentity } from "../../utils/sessionManager.js";
import {
resolveKiroModel,
resolveKiroThinkingBudget,
buildThinkingSystemPrefix,
KIRO_AGENTIC_SYSTEM_PROMPT,
resolveDefaultProfileArn,
buildKiroAdditionalModelRequestFieldsForModel,
} from "../../config/kiroConstants.js";
import { DEFAULT_IMAGE_MIME } from "../schema/index.js";
import { ROLE, CLAUDE_BLOCK } from "../schema/index.js";
@@ -363,6 +365,18 @@ function reconcileOrphanedToolResults(history, currentMessage) {
}
}
function extractClaudeSystemText(system) {
if (!system) return "";
if (typeof system === "string") return system;
if (Array.isArray(system)) {
return system.map((s) => {
if (typeof s === "string") return s;
return s?.text || "";
}).filter(Boolean).join("\n");
}
return "";
}
/**
* Build a Kiro payload directly from a Claude Messages API request body.
*/
@@ -402,62 +416,75 @@ export function claudeToKiroRequest(model, body, stream, credentials) {
? (credentials?.providerSpecificData?.profileArn || "")
: (credentials?.providerSpecificData?.profileArn || resolveDefaultProfileArn(authMethod));
let finalContent = currentMessage?.userInputMessage?.content || "";
// System prompt: pass via native systemInstruction field (Kiro/Q API supports it)
// and also prepend as <instructions> in user content as fallback for upstreams
// that don't support the native field.
let systemInstruction = undefined;
if (body.system) {
let systemText = "";
if (typeof body.system === "string") {
systemText = body.system;
} else if (Array.isArray(body.system)) {
systemText = body.system.map((s) => s.text || "").join("\n");
}
if (systemText) {
systemInstruction = systemText;
finalContent = `<instructions>\n${systemText}\n</instructions>\n\n${finalContent}`;
}
}
// Prefix order: thinking_mode tag, timestamp marker, then agentic prompt.
// Kiro CLI/KAS sends system prompt as top-level `systemPrompt`. Keep a
// content fallback too because the CodeWhisperer surface does not always
// enforce top-level systemPrompt for direct calls.
const timestamp = new Date().toISOString();
const prefixParts = [];
if (thinkingBudget !== null) prefixParts.push(buildThinkingSystemPrefix(thinkingBudget));
prefixParts.push(`[Context: Current time is ${timestamp}]`);
if (agentic) prefixParts.push(KIRO_AGENTIC_SYSTEM_PROMPT);
finalContent = `${prefixParts.join("\n\n")}\n\n${finalContent}`;
const systemPromptParts = [];
if (thinkingBudget !== null) systemPromptParts.push(buildThinkingSystemPrefix(thinkingBudget));
if (agentic) systemPromptParts.push(KIRO_AGENTIC_SYSTEM_PROMPT);
const systemInstruction = extractClaudeSystemText(body.system);
if (systemInstruction) systemPromptParts.push(systemInstruction);
const systemPrompt = systemPromptParts.filter(Boolean).join("\n\n");
const currentTimeContext = `[Context: Current time is ${timestamp}]`;
const contentPrefix = [systemPrompt, currentTimeContext].filter(Boolean).join("\n\n");
const sessionIdentity = resolveSessionIdentity({
headers: credentials?.rawHeaders,
body,
connectionId: credentials?.connectionId,
scope: "kiro",
});
const conversationId = sessionIdentity.sessionId;
const continuationId = resolveContinuationId({
sessionId: conversationId,
connectionId: credentials?.connectionId,
scope: "kiro",
ephemeral: sessionIdentity.ephemeral,
});
const replay = applyKiroSessionReplay({
conversationId,
connectionId: credentials?.connectionId,
modelId: upstreamModel,
systemPrompt,
contentPrefix,
currentContentPrefix: currentTimeContext,
history,
currentMessage,
});
const replayCurrent = replay.currentMessage?.userInputMessage || {};
const userInputMessage = {
content: finalContent,
content: replayCurrent.content || "",
modelId: upstreamModel,
origin: "AI_EDITOR",
...(currentMessage?.userInputMessage?.userInputMessageContext && {
userInputMessageContext:
currentMessage.userInputMessage.userInputMessageContext,
...(replayCurrent.userInputMessageContext && {
userInputMessageContext: replayCurrent.userInputMessageContext,
}),
...(currentMessage?.userInputMessage?.images && {
images: currentMessage.userInputMessage.images,
...(replayCurrent.images && {
images: replayCurrent.images,
}),
};
if (systemInstruction) {
userInputMessage.systemInstruction = systemInstruction;
}
const payload = {
conversationState: {
chatTriggerType: "MANUAL",
conversationId: uuidv4(),
conversationId,
agentContinuationId: continuationId,
agentTaskType: "vibe",
currentMessage: {
userInputMessage,
},
history,
history: replay.history,
},
agentMode: "vibe",
};
if (profileArn) payload.profileArn = profileArn;
if (systemPrompt) payload.systemPrompt = systemPrompt;
const additionalModelRequestFields = buildKiroAdditionalModelRequestFieldsForModel(body, upstreamModel);
if (additionalModelRequestFields) {
payload.additionalModelRequestFields = additionalModelRequestFields;
}
if (maxTokens || temperature !== undefined || topP !== undefined) {
payload.inferenceConfig = {};

View File

@@ -5,13 +5,15 @@
import { register } from "../index.js";
import { FORMATS } from "../formats.js";
import { v4 as uuidv4 } from "uuid";
import { resolveSessionId } from "../../utils/sessionManager.js";
import { applyKiroSessionReplay } from "../../utils/kiroSessionReplay.js";
import { resolveContinuationId, resolveSessionIdentity } from "../../utils/sessionManager.js";
import {
resolveKiroModel,
resolveKiroThinkingBudget,
buildThinkingSystemPrefix,
KIRO_AGENTIC_SYSTEM_PROMPT,
resolveDefaultProfileArn
resolveDefaultProfileArn,
buildKiroAdditionalModelRequestFieldsForModel
} from "../../config/kiroConstants.js";
import { parseDataUri } from "../concerns/image.js";
import { DEFAULT_IMAGE_MIME } from "../schema/index.js";
@@ -546,47 +548,74 @@ export function openaiToKiroRequest(model, body, stream, credentials) {
? (credentials?.providerSpecificData?.profileArn || "")
: (credentials?.providerSpecificData?.profileArn || resolveDefaultProfileArn(authMethod));
let finalContent = currentMessage?.userInputMessage?.content || "";
const timestamp = new Date().toISOString();
// Build the system-prompt prefix that goes ABOVE the user message body.
// Order: thinking_mode tag first (so Kiro sees it before any user text),
// then context/timestamp marker, then optional agentic chunked-write prompt.
const prefixParts = [];
// Kiro CLI/KAS sends these as top-level systemPrompt. Keep a content fallback
// too because the CodeWhisperer surface does not always enforce top-level
// systemPrompt for direct calls.
const systemPromptParts = [];
if (thinkingBudget !== null) {
prefixParts.push(buildThinkingSystemPrefix(thinkingBudget));
systemPromptParts.push(buildThinkingSystemPrefix(thinkingBudget));
}
prefixParts.push(`[Context: Current time is ${timestamp}]`);
if (agentic) {
prefixParts.push(KIRO_AGENTIC_SYSTEM_PROMPT);
systemPromptParts.push(KIRO_AGENTIC_SYSTEM_PROMPT);
}
finalContent = `${prefixParts.join("\n\n")}\n\n${finalContent}`;
const systemPrompt = systemPromptParts.filter(Boolean).join("\n\n");
const currentTimeContext = `[Context: Current time is ${timestamp}]`;
const contentPrefix = [systemPrompt, currentTimeContext].filter(Boolean).join("\n\n");
const sessionIdentity = resolveSessionIdentity({ headers: credentials?.rawHeaders, body, connectionId: credentials?.connectionId, scope: "kiro" });
const conversationId = sessionIdentity.sessionId;
const continuationId = resolveContinuationId({
sessionId: conversationId,
connectionId: credentials?.connectionId,
scope: "kiro",
ephemeral: sessionIdentity.ephemeral,
});
const replay = applyKiroSessionReplay({
conversationId,
connectionId: credentials?.connectionId,
modelId: upstreamModel,
systemPrompt,
contentPrefix,
currentContentPrefix: currentTimeContext,
history,
currentMessage,
});
const replayCurrent = replay.currentMessage?.userInputMessage || {};
const payload = {
conversationState: {
chatTriggerType: "MANUAL",
conversationId: resolveSessionId({ headers: credentials?.rawHeaders, body, connectionId: credentials?.connectionId, scope: "kiro" }),
conversationId,
agentContinuationId: continuationId,
agentTaskType: "vibe",
currentMessage: {
userInputMessage: {
content: finalContent,
content: replayCurrent.content || "",
modelId: upstreamModel,
origin: "AI_EDITOR",
...(currentMessage?.userInputMessage?.images?.length > 0 && {
images: currentMessage.userInputMessage.images
...(replayCurrent.images?.length > 0 && {
images: replayCurrent.images
}),
...(currentMessage?.userInputMessage?.userInputMessageContext && {
userInputMessageContext: currentMessage.userInputMessage.userInputMessageContext
...(replayCurrent.userInputMessageContext && {
userInputMessageContext: replayCurrent.userInputMessageContext
})
}
},
history: history
}
history: replay.history
},
agentMode: "vibe",
};
if (profileArn) {
payload.profileArn = profileArn;
}
if (systemPrompt) payload.systemPrompt = systemPrompt;
const additionalModelRequestFields = buildKiroAdditionalModelRequestFieldsForModel(body, upstreamModel);
if (additionalModelRequestFields) {
payload.additionalModelRequestFields = additionalModelRequestFields;
}
if (maxTokens || temperature !== undefined || topP !== undefined) {
payload.inferenceConfig = {};

View File

@@ -0,0 +1,125 @@
import { MEMORY_CONFIG } from "../config/runtimeConfig.js";
const sessionStartStore = new Map();
const MAX_SESSION_STARTS = 5000;
function clone(value) {
return value == null ? value : JSON.parse(JSON.stringify(value));
}
function sessionKey(connectionId, conversationId) {
return `${connectionId || ""}:${conversationId || ""}`;
}
function ensureUserMessageModelId(message, modelId) {
if (message?.userInputMessage && !message.userInputMessage.modelId && modelId) {
message.userInputMessage.modelId = modelId;
}
return message;
}
function ensureHistoryModelIds(history, modelId) {
for (const item of history || []) {
ensureUserMessageModelId(item, modelId);
}
return history;
}
function prefixUserMessage(message, contentPrefix, modelId) {
const out = clone(message) || { userInputMessage: { content: "" } };
if (!out.userInputMessage) out.userInputMessage = { content: "" };
ensureUserMessageModelId(out, modelId);
if (contentPrefix) {
const content = out.userInputMessage.content || "";
out.userInputMessage.content = content
? `${contentPrefix}\n\n${content}`
: contentPrefix;
}
return out;
}
function findFirstUserIndex(history) {
return history.findIndex((item) => item?.userInputMessage);
}
function rememberSessionStart(key, entry) {
if (sessionStartStore.size >= MAX_SESSION_STARTS) {
sessionStartStore.delete(sessionStartStore.keys().next().value);
}
sessionStartStore.set(key, { ...entry, lastUsed: Date.now() });
}
/**
* Preserve Kiro cacheability by freezing the first user message (`msg0`) for a
* session, replaying that exact message as the first history user on later
* turns, and injecting volatile current-time context only into the current turn.
*/
export function applyKiroSessionReplay({
conversationId,
connectionId,
modelId,
systemPrompt = "",
contentPrefix = "",
currentContentPrefix = "",
history = [],
currentMessage,
} = {}) {
const key = sessionKey(connectionId, conversationId);
const existing = conversationId ? sessionStartStore.get(key) : null;
const baseHistory = clone(history) || [];
const baseCurrent = clone(currentMessage) || { userInputMessage: { content: "" } };
if (existing && existing.modelId === modelId && existing.systemPrompt === systemPrompt) {
existing.lastUsed = Date.now();
const firstUserIndex = findFirstUserIndex(baseHistory);
const sessionStart = ensureUserMessageModelId(clone(existing.sessionStart), modelId);
if (firstUserIndex >= 0) {
baseHistory[firstUserIndex] = sessionStart;
} else {
baseHistory.unshift(sessionStart);
}
return {
history: ensureHistoryModelIds(baseHistory, modelId),
currentMessage: prefixUserMessage(baseCurrent, currentContentPrefix, modelId),
replayed: true,
};
}
const firstUserIndex = findFirstUserIndex(baseHistory);
let sessionStart;
let nextCurrent = ensureUserMessageModelId(baseCurrent, modelId);
if (firstUserIndex >= 0) {
sessionStart = prefixUserMessage(baseHistory[firstUserIndex], contentPrefix, modelId);
baseHistory[firstUserIndex] = clone(sessionStart);
nextCurrent = prefixUserMessage(baseCurrent, currentContentPrefix, modelId);
} else {
sessionStart = prefixUserMessage(baseCurrent, contentPrefix, modelId);
nextCurrent = clone(sessionStart);
}
if (conversationId) {
rememberSessionStart(key, {
sessionStart: clone(sessionStart),
modelId,
systemPrompt,
});
}
return {
history: ensureHistoryModelIds(baseHistory, modelId),
currentMessage: nextCurrent,
replayed: false,
};
}
export function clearKiroSessionReplayStore() {
sessionStartStore.clear();
}
const cleanup = setInterval(() => {
const now = Date.now();
for (const [key, entry] of sessionStartStore) {
if (now - entry.lastUsed > MEMORY_CONFIG.sessionTtlMs) sessionStartStore.delete(key);
}
}, MEMORY_CONFIG.sessionCleanupIntervalMs);
if (cleanup.unref) cleanup.unref();

View File

@@ -13,6 +13,7 @@ import { MEMORY_CONFIG } from "../config/runtimeConfig.js";
// Runtime storage: Key = connectionId, Value = { sessionId, lastUsed }
const runtimeSessionStore = new Map();
const continuationStore = new Map();
// Periodically evict entries that haven't been used within TTL
const cleanupInterval = setInterval(() => {
@@ -80,6 +81,7 @@ export function generateBinaryStyleId() {
export function clearSessionStore() {
runtimeSessionStore.clear();
assistantSessionStore.clear();
continuationStore.clear();
}
// Conversation-stable session store: Key = hash(scope+assistant text), Value = { sessionId, lastUsed }
@@ -87,9 +89,10 @@ const assistantSessionStore = new Map();
const ASSISTANT_MIN_LEN = 50;
const ASSISTANT_CAP_LEN = 50;
const MAX_ASSISTANT_SESSIONS = 5000;
const MAX_CONTINUATION_SESSIONS = 5000;
// Client headers/body fields that carry an upstream session id (priority order)
const SESSION_HEADER_KEYS = ["x-session-id", "session-id", "session_id", "x-amp-thread-id", "x-client-request-id"];
const SESSION_HEADER_KEYS = ["x-session-id", "session-id", "session_id", "x-amp-thread-id"];
const CLAUDE_CODE_SESSION_RE = /_session_([a-f0-9-]+)$/;
function sha16(text) {
@@ -131,7 +134,7 @@ function extractAntigravitySession(body) {
return m ? normalizeSessionId(m[1]) : null;
}
function extractClientSessionId(headers, body) {
function extractClientSessionId(headers, body, scope = "") {
const claude = extractClaudeCodeSession(body?.metadata?.user_id);
if (claude) return `claude:${claude}`;
const antigravity = extractAntigravitySession(body);
@@ -140,18 +143,25 @@ function extractClientSessionId(headers, body) {
const v = headerValue(headers, key);
if (v) return v;
}
const requestId = scope === "kiro" ? null : headerValue(headers, "x-client-request-id");
if (requestId) return requestId;
const fromBody =
normalizeSessionId(body?.prompt_cache_key) ||
normalizeSessionId(body?.session_id) ||
normalizeSessionId(body?.conversation_id) ||
normalizeSessionId(body?.metadata?.user_id);
(scope === "kiro" ? null : normalizeSessionId(body?.metadata?.user_id));
return fromBody || null;
}
function requestMessages(body) {
if (Array.isArray(body?.messages)) return body.messages;
if (Array.isArray(body?.input)) return body.input;
return [];
}
// Accumulate assistant text from OpenAI/Responses-style input/messages (cap-limited)
function accumulateAssistantText(body) {
const items = Array.isArray(body?.input) ? body.input
: Array.isArray(body?.messages) ? body.messages : null;
const items = requestMessages(body);
if (!items) return "";
let text = "";
for (const item of items) {
@@ -193,16 +203,39 @@ function assistantTextSessionId(scope, body) {
* @param {string} [opts.connectionId] - Connection identifier (fallback scope)
* @param {string} [opts.workspaceId] - Provider workspace id (account-wide fallback)
* @param {string} [opts.scope] - Provider scope to isolate cache keys across providers
* @returns {string} A stable session id
* @returns {{sessionId: string, ephemeral: boolean}} A session id plus whether it is one-shot
*/
export function resolveSessionId({ headers, body, connectionId, workspaceId, scope = "" } = {}) {
const client = extractClientSessionId(headers, body);
if (client) return client;
const fromAssistant = assistantTextSessionId(`${scope}:${connectionId || ""}`, body);
if (fromAssistant) return fromAssistant;
export function resolveSessionIdentity({ headers, body, connectionId, workspaceId, scope = "" } = {}) {
const client = extractClientSessionId(headers, body, scope);
if (client) return { sessionId: client, ephemeral: false };
const fromAssistant = scope === "kiro" ? null : assistantTextSessionId(`${scope}:${connectionId || ""}`, body);
if (fromAssistant) return { sessionId: fromAssistant, ephemeral: false };
const ws = normalizeSessionId(workspaceId);
if (ws) return ws;
return deriveSessionId(connectionId);
if (ws) return { sessionId: ws, ephemeral: false };
if (scope === "kiro") return { sessionId: generateBinaryStyleId(), ephemeral: true };
return { sessionId: deriveSessionId(connectionId), ephemeral: false };
}
export function resolveSessionId(opts = {}) {
return resolveSessionIdentity(opts).sessionId;
}
export function resolveContinuationId({ sessionId, connectionId, scope = "", ephemeral = false } = {}) {
if (ephemeral) return crypto.randomUUID();
const key = `${scope}:${connectionId || ""}:${sessionId || ""}`;
const existing = continuationStore.get(key);
if (existing) {
existing.lastUsed = Date.now();
continuationStore.delete(key);
continuationStore.set(key, existing);
return existing.continuationId;
}
const continuationId = crypto.randomUUID();
if (continuationStore.size >= MAX_CONTINUATION_SESSIONS) {
continuationStore.delete(continuationStore.keys().next().value);
}
continuationStore.set(key, { continuationId, lastUsed: Date.now() });
return continuationId;
}
// Capture session id from request body + credentials (envelope still intact here)
@@ -227,5 +260,8 @@ const assistantCleanup = setInterval(() => {
for (const [key, entry] of assistantSessionStore) {
if (now - entry.lastUsed > MEMORY_CONFIG.sessionTtlMs) assistantSessionStore.delete(key);
}
for (const [key, entry] of continuationStore) {
if (now - entry.lastUsed > MEMORY_CONFIG.sessionTtlMs) continuationStore.delete(key);
}
}, MEMORY_CONFIG.sessionCleanupIntervalMs);
if (assistantCleanup.unref) assistantCleanup.unref();