+ {promptsInQueue.length > 0 && (
+
+
+
+ Queued for upcoming turns
+
+
+ Steer runs first on the next turn
+
+
+
+ {promptsInQueue.map((item, index) => (
+
+
+
+ {item.steer ? "Steer" : `Queue ${index + 1}`}
+ {item.steer ? (
+
+ Next turn
+
+ ) : null}
+
+
+ {item.prompt}
+
+
+ {!item.steer ? (
+
+ ) : (
+
+ Steering
+
+ )}
+
+ ))}
+
+
+ )}
{mentionOpen && (
@@ -311,9 +382,9 @@ export function ChatInputBar({
)}
)}
-
-
+
@@ -443,13 +515,13 @@ export function ChatInputBar({
diff --git a/sdk/apps/code/components/views/chat/chat-messages.tsx b/sdk/apps/code/components/views/chat/chat-messages.tsx
index 3d8651553f..7494332d52 100644
--- a/sdk/apps/code/components/views/chat/chat-messages.tsx
+++ b/sdk/apps/code/components/views/chat/chat-messages.tsx
@@ -13,7 +13,14 @@ import {
ShieldAlert,
Terminal,
} from "lucide-react";
-import { memo, useCallback, useEffect, useRef, useState } from "react";
+import {
+ memo,
+ useCallback,
+ useEffect,
+ useLayoutEffect,
+ useRef,
+ useState,
+} from "react";
import { Button } from "@/components/ui/button";
import type { ChatMessage, ChatSessionStatus } from "@/lib/chat-schema";
import { cn } from "@/lib/utils";
@@ -50,6 +57,8 @@ type ToolApprovalRequestItem = {
};
const IS_DEBUG = process.env.NODE_ENV === "test";
+const STICKY_BOTTOM_THRESHOLD_PX = 24;
+const SCROLL_TO_BOTTOM_BUTTON_THRESHOLD_PX = 120;
function ChatMessagesImpl({
sessionId: _sessionId,
@@ -67,7 +76,7 @@ function ChatMessagesImpl({
onStartChat,
}: ChatMessagesProps) {
const scrollAreaRef = useRef
(null);
- const hasAppliedInitialScrollRef = useRef(false);
+ const shouldStickToBottomRef = useRef(true);
const hasMessages = messages.length > 0;
const lastErrorMessage = [...messages]
.reverse()
@@ -95,6 +104,7 @@ function ChatMessagesImpl({
if (!viewport) {
return;
}
+ shouldStickToBottomRef.current = true;
viewport.scrollTo({ top: viewport.scrollHeight, behavior });
setShowScrollToBottom((prev) => (prev ? false : prev));
},
@@ -123,7 +133,10 @@ function ChatMessagesImpl({
const updateScrollToBottomVisibility = () => {
const distanceFromBottom =
viewport.scrollHeight - viewport.scrollTop - viewport.clientHeight;
- const shouldShow = distanceFromBottom > 120;
+ shouldStickToBottomRef.current =
+ distanceFromBottom <= STICKY_BOTTOM_THRESHOLD_PX;
+ const shouldShow =
+ distanceFromBottom > SCROLL_TO_BOTTOM_BUTTON_THRESHOLD_PX;
setShowScrollToBottom((prev) =>
prev === shouldShow ? prev : shouldShow,
);
@@ -137,26 +150,12 @@ function ChatMessagesImpl({
};
}, [getViewport]);
- useEffect(() => {
- if (hasAppliedInitialScrollRef.current) {
+ useLayoutEffect(() => {
+ if (!shouldStickToBottomRef.current) {
return;
}
-
- const hasScrollableContent =
- messages.length > 0 || pendingToolApprovals.length > 0;
- if (!hasScrollableContent) {
- return;
- }
-
- const frame = window.requestAnimationFrame(() => {
- scrollToBottom("auto");
- hasAppliedInitialScrollRef.current = true;
- });
-
- return () => {
- window.cancelAnimationFrame(frame);
- };
- }, [messages.length, pendingToolApprovals.length, scrollToBottom]);
+ scrollToBottom("auto");
+ }, [scrollToBottom]);
useEffect(() => {
const activeRequestIds = new Set(
diff --git a/sdk/apps/code/hooks/chat-session/types.ts b/sdk/apps/code/hooks/chat-session/types.ts
index 1088fd58ad..4d20f36eb0 100644
--- a/sdk/apps/code/hooks/chat-session/types.ts
+++ b/sdk/apps/code/hooks/chat-session/types.ts
@@ -67,6 +67,7 @@ export type ChatWsResponseEvent = {
sessionId?: string;
result?: ChatApiResult;
ok?: boolean;
+ queued?: boolean;
};
error?: string;
};
@@ -98,3 +99,9 @@ export type SerializedAttachments = {
userImages: string[];
userFiles: SerializedAttachmentFile[];
};
+
+export type PromptInQueue = {
+ id: string;
+ prompt: string;
+ steer: boolean;
+};
diff --git a/sdk/apps/code/hooks/use-chat-session.ts b/sdk/apps/code/hooks/use-chat-session.ts
index c4cd8a28f1..42a8f2f2b5 100644
--- a/sdk/apps/code/hooks/use-chat-session.ts
+++ b/sdk/apps/code/hooks/use-chat-session.ts
@@ -18,6 +18,7 @@ import type {
ChatTransportState,
CoreLogChunk,
ProcessContext,
+ PromptInQueue,
ToolApprovalRequestItem,
ToolCallEndEvent,
ToolCallStartEvent,
@@ -39,6 +40,117 @@ import type { SessionHistoryItem } from "@/lib/session-history";
export { DEFAULT_CHAT_CONFIG } from "@/hooks/chat-session/constants";
+const MAX_MESSAGES = 800;
+
+const RELEVANT_STREAMS = new Set([
+ "chat_text",
+ "chat_queued_prompt_start",
+ "chat_tool_call_start",
+ "chat_tool_call_end",
+ "chat_core_log",
+]);
+
+const BUSY_STATUSES = new Set([
+ "starting",
+ "running",
+ "stopping",
+]);
+
+// ---------------------------------------------------------------------------
+// Helpers (pure, no hooks)
+// ---------------------------------------------------------------------------
+
+function errorMessage(err: unknown): string {
+ return err instanceof Error ? err.message : String(err);
+}
+
+function makeErrorChatMessage(
+ sid: string | null,
+ content: string,
+): ChatMessage {
+ return {
+ id: makeId("error"),
+ sessionId: sid,
+ role: "error",
+ content,
+ createdAt: Date.now(),
+ };
+}
+
+function validateConfig(
+ config: ChatSessionConfig,
+):
+ | { parsed: ChatSessionConfig; error: null }
+ | { parsed: null; error: string } {
+ const runtimeConfig = normalizeRuntimeConfig(config);
+ const result = ChatSessionConfigSchema.safeParse(runtimeConfig);
+ if (!result.success) {
+ return {
+ parsed: null,
+ error: result.error.issues.map((i) => i.message).join(", "),
+ };
+ }
+ const credentialError = resolveCredentialError(result.data);
+ if (credentialError) {
+ return { parsed: null, error: credentialError };
+ }
+ return { parsed: result.data, error: null };
+}
+
+function sliceMessages(msgs: ChatMessage[]): ChatMessage[] {
+ return msgs.length > MAX_MESSAGES ? msgs.slice(-MAX_MESSAGES) : msgs;
+}
+
+function sortMessagesChronologically(messages: ChatMessage[]): ChatMessage[] {
+ return [...messages].sort((left, right) => {
+ if (left.createdAt !== right.createdAt) {
+ return left.createdAt - right.createdAt;
+ }
+ return left.id.localeCompare(right.id);
+ });
+}
+
+function updateMessageById(
+ messages: ChatMessage[],
+ id: string,
+ updater: (msg: ChatMessage) => ChatMessage,
+): ChatMessage[] {
+ let changed = false;
+ const next = messages.map((msg) => {
+ if (msg.id !== id) return msg;
+ changed = true;
+ return updater(msg);
+ });
+ return changed ? next : messages;
+}
+
+// ---------------------------------------------------------------------------
+// Core log dispatcher — avoids repeated if/else chains
+// ---------------------------------------------------------------------------
+
+const LOG_DISPATCH: Record = {
+ error: console.error,
+ warn: console.warn,
+ debug: console.debug,
+};
+
+function dispatchCoreLog(chunk: string): void {
+ let parsed: CoreLogChunk | undefined;
+ try {
+ parsed = JSON.parse(chunk) as CoreLogChunk;
+ } catch {
+ console.info("[core]", chunk);
+ return;
+ }
+ const level = parsed.level?.trim().toLowerCase() || "info";
+ const message = parsed.message?.trim() || chunk;
+ (LOG_DISPATCH[level] ?? console.info)("[core]", message, parsed.metadata);
+}
+
+// ---------------------------------------------------------------------------
+// Hook
+// ---------------------------------------------------------------------------
+
export function useChatSession() {
const [sessionId, setSessionId] = useState(null);
const [status, setStatus] = useState("idle");
@@ -62,6 +174,7 @@ export function useChatSession() {
const [pendingToolApprovals, setPendingToolApprovals] = useState<
ToolApprovalRequestItem[]
>([]);
+ const [promptsInQueue, setPromptsInQueue] = useState([]);
const liveToolMessageIdsRef = useRef>({});
const liveToolInputsRef = useRef>({});
const activeSessionIdRef = useRef(null);
@@ -71,44 +184,102 @@ export function useChatSession() {
useState(desktopClient.getTransportState());
const messagesRef = useRef([]);
+ // ---- Ref syncs ----
+
useEffect(() => {
activeSessionIdRef.current = sessionId;
}, [sessionId]);
-
useEffect(() => {
activeAssistantMessageIdRef.current = activeAssistantMessageId;
}, [activeAssistantMessageId]);
-
useEffect(() => {
messagesRef.current = messages;
}, [messages]);
+ // ---- Shared state reset helpers ----
+
+ const clearLiveToolRefs = useCallback(() => {
+ liveToolMessageIdsRef.current = {};
+ liveToolInputsRef.current = {};
+ }, []);
+
+ const resetCounters = useCallback(() => {
+ setToolCalls(0);
+ setTokensIn(0);
+ setTokensOut(0);
+ setFileDiffs([]);
+ setDiffSummary(EMPTY_DIFF_SUMMARY);
+ }, []);
+
+ const setErrorState = useCallback(
+ (msg: string, sid: string | null = null) => {
+ setError(msg);
+ setStatus("error");
+ setMessages((prev) =>
+ sliceMessages([...prev, makeErrorChatMessage(sid, msg)]),
+ );
+ },
+ [],
+ );
+
+ // ---- Data fetching ----
+
+ const postSession = useCallback(async (body: Record) => {
+ return await desktopClient.invoke<{
+ sessionId?: string;
+ result?: ChatApiResult;
+ ok?: boolean;
+ queued?: boolean;
+ promptsInQueue?: PromptInQueue[];
+ }>("chat_session_command", { request: body });
+ }, []);
+
+ const refreshPromptsInQueue = useCallback(
+ async (targetSessionId: string | null) => {
+ if (!targetSessionId) {
+ setPromptsInQueue([]);
+ return;
+ }
+ try {
+ const payload = await postSession({
+ action: "pending_prompts",
+ sessionId: targetSessionId,
+ });
+ setPromptsInQueue(
+ Array.isArray(payload.promptsInQueue) ? payload.promptsInQueue : [],
+ );
+ } catch {
+ // Ignore queue refresh failures and keep the last known state.
+ }
+ },
+ [postSession],
+ );
+
+ const applyPromptsInQueue = useCallback((value: unknown) => {
+ if (!Array.isArray(value)) {
+ return;
+ }
+ setPromptsInQueue(value as PromptInQueue[]);
+ }, []);
+
const refreshSessionDiffSummary = useCallback(
async (targetSessionId: string) => {
try {
const events = await desktopClient.invoke(
"read_session_hooks",
- {
- sessionId: targetSessionId,
- limit: 800,
- },
+ { sessionId: targetSessionId, limit: MAX_MESSAGES },
);
const diffState = buildSessionDiffState(events);
setFileDiffs(diffState.fileDiffs);
setDiffSummary(diffState.summary);
setToolCalls(
events.filter(
- (event) =>
- event.hookEventName === "tool_call" ||
- event.hookName === "tool_call",
+ (e) =>
+ e.hookEventName === "tool_call" || e.hookName === "tool_call",
).length,
);
- setTokensIn(
- events.reduce((sum, event) => sum + (event.inputTokens ?? 0), 0),
- );
- setTokensOut(
- events.reduce((sum, event) => sum + (event.outputTokens ?? 0), 0),
- );
+ setTokensIn(events.reduce((sum, e) => sum + (e.inputTokens ?? 0), 0));
+ setTokensOut(events.reduce((sum, e) => sum + (e.outputTokens ?? 0), 0));
} catch {
// Ignore in non-Tauri mode.
}
@@ -116,8 +287,12 @@ export function useChatSession() {
[],
);
+ // ---- Message helpers ----
+
const addMessage = useCallback((message: ChatMessage) => {
- setMessages((prev) => [...prev, message].slice(-800));
+ setMessages((prev) =>
+ sliceMessages(sortMessagesChronologically([...prev, message])),
+ );
}, []);
const materializeToolMessagesFromResult = useCallback(
@@ -127,19 +302,15 @@ export function useChatSession() {
toolCalls: NonNullable;
}) => {
const { sessionId: targetSessionId, turnStartedAt, toolCalls } = options;
- if (toolCalls.length === 0) {
- return;
- }
+ if (toolCalls.length === 0) return;
setMessages((prev) => {
- const hasLiveToolMessagesForTurn = prev.some(
- (message) =>
- message.sessionId === targetSessionId &&
- message.role === "tool" &&
- message.createdAt >= turnStartedAt,
+ const hasLive = prev.some(
+ (m) =>
+ m.sessionId === targetSessionId &&
+ m.role === "tool" &&
+ m.createdAt >= turnStartedAt,
);
- if (hasLiveToolMessagesForTurn) {
- return prev;
- }
+ if (hasLive) return prev;
const next = [...prev];
for (const call of toolCalls) {
@@ -154,58 +325,30 @@ export function useChatSession() {
isError: Boolean(call.error),
}),
createdAt: turnStartedAt + next.length,
- meta: {
- toolName: call.name,
- hookEventName: "tool_call_end",
- },
+ meta: { toolName: call.name, hookEventName: "tool_call_end" },
});
}
- return next.slice(-800);
+ return sliceMessages(sortMessagesChronologically(next));
});
},
[],
);
- const replaceMessageContent = useCallback((id: string, content: string) => {
+ const appendMessageContent = useCallback((id: string, chunk: string) => {
+ if (!chunk) return;
setMessages((prev) =>
- prev.map((message) => {
- if (message.id !== id) {
- return message;
- }
- return {
- ...message,
- content,
- };
+ updateMessageById(prev, id, (msg) => {
+ const existing = msg.content;
+ if (existing.endsWith(chunk)) return msg;
+ const content = chunk.startsWith(existing)
+ ? chunk
+ : `${existing}${chunk}`;
+ return { ...msg, content };
}),
);
}, []);
- const appendMessageContent = useCallback((id: string, chunk: string) => {
- if (!chunk) {
- return;
- }
- setMessages((prev) =>
- prev.map((message) => {
- if (message.id !== id) {
- return message;
- }
- const existing = message.content;
- if (existing.endsWith(chunk)) {
- return message;
- }
- if (chunk.startsWith(existing)) {
- return {
- ...message,
- content: chunk,
- };
- }
- return {
- ...message,
- content: `${existing}${chunk}`,
- };
- }),
- );
- }, []);
+ // ---- Process context ----
const applyProcessContext = useCallback(async () => {
try {
@@ -226,15 +369,19 @@ export function useChatSession() {
void applyProcessContext();
}, [applyProcessContext]);
+ // ---- Diff / approval effects ----
+
useEffect(() => {
if (!sessionId) {
setFileDiffs([]);
setDiffSummary(EMPTY_DIFF_SUMMARY);
setPendingToolApprovals([]);
+ setPromptsInQueue([]);
return;
}
void refreshSessionDiffSummary(sessionId);
- }, [refreshSessionDiffSummary, sessionId]);
+ void refreshPromptsInQueue(sessionId);
+ }, [refreshPromptsInQueue, refreshSessionDiffSummary, sessionId]);
useEffect(() => {
const activeSessionId = sessionId;
@@ -248,58 +395,86 @@ export function useChatSession() {
sessionId: activeSessionId,
limit: 20,
})
- .then((pending) => {
- setPendingToolApprovals(pending);
- })
- .catch(() => {
- // Ignore initial hydration failures.
- });
+ .then((pending) => setPendingToolApprovals(pending))
+ .catch(() => {});
return desktopClient.subscribe("tool_approval_state", (payload) => {
- if (!payload || typeof payload !== "object") {
- return;
- }
+ if (!payload || typeof payload !== "object") return;
const record = payload as {
sessionId?: string;
items?: ToolApprovalRequestItem[];
};
- if (record.sessionId !== activeSessionId) {
- return;
- }
+ if (record.sessionId !== activeSessionId) return;
setPendingToolApprovals(Array.isArray(record.items) ? record.items : []);
});
}, [sessionId]);
+ useEffect(() => {
+ return desktopClient.subscribe("prompts_in_queue_state", (payload) => {
+ if (!payload || typeof payload !== "object") return;
+ const record = payload as {
+ sessionId?: string;
+ items?: PromptInQueue[];
+ };
+ if (record.sessionId !== activeSessionIdRef.current) return;
+ setPromptsInQueue(Array.isArray(record.items) ? record.items : []);
+ });
+ }, []);
+
+ // ---- Incoming chunk handler ----
+
const handleIncomingChunk = useCallback(
(payload: AgentChunkEvent) => {
- if (
- payload.stream !== "chat_text" &&
- payload.stream !== "chat_tool_call_start" &&
- payload.stream !== "chat_tool_call_end" &&
- payload.stream !== "chat_core_log"
- ) {
- return;
- }
+ if (!RELEVANT_STREAMS.has(payload.stream)) return;
+
const listeningSessionId = activeSessionIdRef.current;
if (!listeningSessionId || payload.sessionId !== listeningSessionId) {
return;
}
- let listeningAssistantId = activeAssistantMessageIdRef.current;
+
+ if (payload.stream === "chat_queued_prompt_start") {
+ let parsed: { prompt?: string; attachmentCount?: number } = {};
+ try {
+ parsed = JSON.parse(payload.chunk) as {
+ prompt?: string;
+ attachmentCount?: number;
+ };
+ } catch {
+ parsed = { prompt: payload.chunk };
+ }
+ const prompt = parsed.prompt?.trim() ?? "";
+ const attachmentCount =
+ typeof parsed.attachmentCount === "number"
+ ? parsed.attachmentCount
+ : 0;
+ const userLabel =
+ attachmentCount > 0
+ ? `${prompt}${prompt.length > 0 ? "\n\n" : ""}[attached ${attachmentCount} file${attachmentCount === 1 ? "" : "s"}]`
+ : prompt;
+ if (userLabel) {
+ addMessage({
+ id: makeId("user"),
+ sessionId: listeningSessionId,
+ role: "user",
+ content: userLabel,
+ createdAt: payload.ts || Date.now(),
+ });
+ }
+ return;
+ }
+
+ // --- Text stream ---
if (payload.stream === "chat_text") {
- if (!listeningAssistantId) {
+ let assistantId = activeAssistantMessageIdRef.current;
+ if (!assistantId) {
const sessionMessages = messagesRef.current.filter(
- (message) => message.sessionId === listeningSessionId,
+ (m) => m.sessionId === listeningSessionId,
);
- const latestSessionMessage = sessionMessages.at(-1);
- if (latestSessionMessage?.role === "assistant") {
- listeningAssistantId = latestSessionMessage.id;
- activeAssistantMessageIdRef.current = listeningAssistantId;
- setActiveAssistantMessageId(listeningAssistantId);
+ const latest = sessionMessages.at(-1);
+ if (latest?.role === "assistant") {
+ assistantId = latest.id;
} else {
- const assistantId = makeId("assistant");
- listeningAssistantId = assistantId;
- activeAssistantMessageIdRef.current = assistantId;
- setActiveAssistantMessageId(assistantId);
+ assistantId = makeId("assistant");
addMessage({
id: assistantId,
sessionId: listeningSessionId,
@@ -308,37 +483,21 @@ export function useChatSession() {
createdAt: payload.ts || Date.now(),
});
}
+ activeAssistantMessageIdRef.current = assistantId;
+ setActiveAssistantMessageId(assistantId);
}
- appendMessageContent(listeningAssistantId, payload.chunk);
+ appendMessageContent(assistantId, payload.chunk);
setRawTranscript((prev) => `${prev}${payload.chunk}`);
return;
}
+
+ // --- Core log ---
if (payload.stream === "chat_core_log") {
- let parsed: CoreLogChunk | undefined;
- try {
- parsed = JSON.parse(payload.chunk) as CoreLogChunk;
- } catch {
- console.info("[core]", payload.chunk);
- return;
- }
- const level = parsed.level?.trim().toLowerCase() || "info";
- const message = parsed.message?.trim() || payload.chunk;
- const metadata = parsed.metadata;
- if (level === "error") {
- console.error("[core]", message, metadata);
- return;
- }
- if (level === "warn") {
- console.warn("[core]", message, metadata);
- return;
- }
- if (level === "debug") {
- console.debug("[core]", message, metadata);
- return;
- }
- console.info("[core]", message, metadata);
+ dispatchCoreLog(payload.chunk);
return;
}
+
+ // --- Tool call start ---
if (payload.stream === "chat_tool_call_start") {
let parsed: ToolCallStartEvent = {};
try {
@@ -361,14 +520,13 @@ export function useChatSession() {
output: null,
}),
createdAt: Date.now(),
- meta: {
- toolName,
- hookEventName: "tool_call_start",
- },
+ meta: { toolName, hookEventName: "tool_call_start" },
});
setToolCalls((prev) => prev + 1);
return;
}
+
+ // --- Tool call end ---
let parsed: ToolCallEndEvent = {};
try {
parsed = JSON.parse(payload.chunk) as ToolCallEndEvent;
@@ -394,21 +552,13 @@ export function useChatSession() {
delete liveToolInputsRef.current[toolCallId];
}
if (messageId) {
- replaceMessageContent(messageId, toolPayload);
+ // Single setMessages call replaces content + updates meta together.
setMessages((prev) =>
- prev.map((message) => {
- if (message.id !== messageId) {
- return message;
- }
- return {
- ...message,
- meta: {
- ...message.meta,
- toolName,
- hookEventName: "tool_call_end",
- },
- };
- }),
+ updateMessageById(prev, messageId, (msg) => ({
+ ...msg,
+ content: toolPayload,
+ meta: { ...msg.meta, toolName, hookEventName: "tool_call_end" },
+ })),
);
return;
}
@@ -418,38 +568,24 @@ export function useChatSession() {
role: "tool",
content: toolPayload,
createdAt: Date.now(),
- meta: {
- toolName,
- hookEventName: "tool_call_end",
- },
+ meta: { toolName, hookEventName: "tool_call_end" },
});
},
- [addMessage, appendMessageContent, replaceMessageContent],
+ [addMessage, appendMessageContent],
);
- const postSession = useCallback(async (body: Record) => {
- return await desktopClient.invoke<{
- sessionId?: string;
- result?: ChatApiResult;
- ok?: boolean;
- }>("chat_session_command", {
- request: body,
- });
- }, []);
+ // ---- Transport / event subscriptions ----
useEffect(() => {
const unsubscribeTransport = desktopClient.subscribeTransportState(
- (next) => {
- setChatTransportState(next);
- },
+ setChatTransportState,
);
const unsubscribeEvents = desktopClient.subscribe(
"chat_event",
(payload) => {
- if (!payload || typeof payload !== "object") {
- return;
+ if (payload && typeof payload === "object") {
+ handleIncomingChunk(payload as AgentChunkEvent);
}
- handleIncomingChunk(payload as AgentChunkEvent);
},
);
return () => {
@@ -458,56 +594,48 @@ export function useChatSession() {
};
}, [handleIncomingChunk]);
+ // ---- Shared: start a new session via RPC ----
+
+ const startSession = useCallback(
+ async (validatedConfig: ChatSessionConfig): Promise => {
+ const payload = await postSession({
+ action: "start",
+ config: validatedConfig,
+ });
+ const id = payload.sessionId;
+ if (!id) throw new Error("Missing session id from server");
+ setSessionId(id);
+ setStatus("running");
+ setConfig(validatedConfig);
+ setHydratedHistorySessionId(null);
+ return id;
+ },
+ [postSession],
+ );
+
+ // ---- Actions ----
+
const start = useCallback(
async (nextConfig: ChatSessionConfig) => {
- const runtimeConfig = normalizeRuntimeConfig(nextConfig);
- const parsed = ChatSessionConfigSchema.safeParse(runtimeConfig);
- if (!parsed.success) {
- const message = parsed.error.issues
- .map((issue) => issue.message)
- .join(", ");
- setError(message);
- setStatus("error");
- return;
- }
- const credentialError = resolveCredentialError(parsed.data);
- if (credentialError) {
- setError(credentialError);
- setStatus("error");
- addMessage({
- id: makeId("error"),
- sessionId: null,
- role: "error",
- content: credentialError,
- createdAt: Date.now(),
- });
+ const validation = validateConfig(nextConfig);
+ if (!validation.parsed) {
+ setErrorState(validation.error);
return;
}
+ const parsed = validation.parsed;
setError(null);
setStatus("starting");
setIsHydratingSession(false);
setMessages([]);
setRawTranscript("");
- setToolCalls(0);
- setTokensIn(0);
- setTokensOut(0);
- setFileDiffs([]);
- setDiffSummary(EMPTY_DIFF_SUMMARY);
- setConfig(parsed.data);
+ resetCounters();
+ setConfig(parsed);
setHydratedHistorySessionId(null);
+ setPromptsInQueue([]);
try {
- const payload = await postSession({
- action: "start",
- config: parsed.data,
- });
- const id = payload.sessionId;
- if (!id) {
- throw new Error("Missing session id from server");
- }
- setSessionId(id);
- setStatus("running");
+ const id = await startSession(parsed);
addMessage({
id: makeId("status"),
sessionId: id,
@@ -516,86 +644,39 @@ export function useChatSession() {
createdAt: Date.now(),
});
} catch (err) {
- const message = err instanceof Error ? err.message : String(err);
- setError(message);
- setStatus("error");
- addMessage({
- id: makeId("error"),
- sessionId: null,
- role: "error",
- content: message,
- createdAt: Date.now(),
- });
+ setErrorState(errorMessage(err));
}
},
- [addMessage, postSession],
+ [addMessage, resetCounters, setErrorState, startSession],
);
const sendPrompt = useCallback(
async (prompt: string, attachedFiles: File[] = []) => {
const trimmed = prompt.trim();
- if (!trimmed && attachedFiles.length === 0) {
- return;
- }
+ if (!trimmed && attachedFiles.length === 0) return;
setError(null);
setIsHydratingSession(false);
let activeSessionId = sessionId;
- const runtimeConfig = normalizeRuntimeConfig(config);
- const parsed = ChatSessionConfigSchema.safeParse(runtimeConfig);
- if (!parsed.success) {
- const message = parsed.error.issues
- .map((issue) => issue.message)
- .join(", ");
- setError(message);
- setStatus("error");
- return;
- }
- const credentialError = resolveCredentialError(parsed.data);
- if (credentialError) {
- setError(credentialError);
- setStatus("error");
- addMessage({
- id: makeId("error"),
- sessionId: activeSessionId,
- role: "error",
- content: credentialError,
- createdAt: Date.now(),
- });
+
+ const validation = validateConfig(config);
+ if (!validation.parsed) {
+ setErrorState(validation.error, activeSessionId);
return;
}
+ const parsed = validation.parsed;
if (!activeSessionId) {
try {
- const payload = await postSession({
- action: "start",
- config: parsed.data,
- });
- const id = payload.sessionId;
- if (!id) {
- throw new Error("Missing session id from server");
- }
- activeSessionId = id;
- setSessionId(id);
- setStatus("running");
- setConfig(parsed.data);
- setHydratedHistorySessionId(null);
+ activeSessionId = await startSession(parsed);
} catch (err) {
- const message = err instanceof Error ? err.message : String(err);
- setError(message);
- setStatus("error");
- addMessage({
- id: makeId("error"),
- sessionId: null,
- role: "error",
- content: message,
- createdAt: Date.now(),
- });
+ setErrorState(errorMessage(err));
return;
}
}
const now = Date.now();
+ const shouldQueue = Boolean(activeSessionId) && BUSY_STATUSES.has(status);
const serializedAttachments = await serializeAttachments(attachedFiles);
const hasAttachments =
serializedAttachments.userImages.length > 0 ||
@@ -604,32 +685,49 @@ export function useChatSession() {
const userLabel = hasAttachments
? `${trimmed}${trimmed.length > 0 ? "\n\n" : ""}[attached ${attachedFiles.length} file${attachedFiles.length === 1 ? "" : "s"}]`
: trimmed;
+ const optimisticQueuedPromptId = shouldQueue
+ ? makeId("queued_prompt")
+ : null;
- addMessage({
- id: makeId("user"),
- sessionId: activeSessionId,
- role: "user",
- content: userLabel,
- createdAt: now,
- });
-
- activeSessionIdRef.current = activeSessionId;
- activeAssistantMessageIdRef.current = null;
- setActiveAssistantMessageId(null);
- liveToolMessageIdsRef.current = {};
- liveToolInputsRef.current = {};
-
- setStatus("starting");
+ if (!shouldQueue) {
+ addMessage({
+ id: makeId("user"),
+ sessionId: activeSessionId,
+ role: "user",
+ content: userLabel,
+ createdAt: now,
+ });
+ activeSessionIdRef.current = activeSessionId;
+ activeAssistantMessageIdRef.current = null;
+ setActiveAssistantMessageId(null);
+ clearLiveToolRefs();
+ setStatus("starting");
+ } else if (optimisticQueuedPromptId) {
+ setPromptsInQueue((prev) => [
+ ...prev,
+ {
+ id: optimisticQueuedPromptId,
+ prompt: userLabel,
+ steer: false,
+ },
+ ]);
+ }
try {
const payload = await postSession({
action: "send",
sessionId: activeSessionId,
prompt: trimmed,
- config: parsed.data,
+ config: parsed,
attachments: hasAttachments ? serializedAttachments : undefined,
});
+ if (payload.ok && payload.queued) {
+ applyPromptsInQueue(payload.promptsInQueue);
+ setStatus("running");
+ return;
+ }
const result = payload.result as ChatApiResult | undefined;
+ applyPromptsInQueue(payload.promptsInQueue);
const assistantText = (result?.text ?? "").trim();
const fallbackAssistantText = extractAssistantTextFromRpcMessages(
result?.messages,
@@ -641,22 +739,14 @@ export function useChatSession() {
activeAssistantMessageIdRef.current = assistantMessageId;
setActiveAssistantMessageId(assistantMessageId);
setMessages((prev) => {
- let found = false;
- const next = prev.map((message) => {
- if (message.id !== assistantMessageId) {
- return message;
- }
- found = true;
- return {
- ...message,
- content: resolvedAssistantText,
- };
- });
- if (found) {
- return next;
- }
- return [
- ...next,
+ const updated = updateMessageById(
+ prev,
+ assistantMessageId,
+ (msg) => ({ ...msg, content: resolvedAssistantText }),
+ );
+ if (updated !== prev) return updated;
+ return sliceMessages([
+ ...prev,
{
id: assistantMessageId,
sessionId: activeSessionId,
@@ -664,17 +754,14 @@ export function useChatSession() {
content: resolvedAssistantText,
createdAt: now + 1,
},
- ].slice(-800);
+ ]);
});
} else {
- // Recovery path: if transport missed result text, load canonical messages.
+ // Recovery: load canonical messages if transport missed result text.
try {
const historyMessages = await desktopClient.invoke(
"read_session_messages",
- {
- sessionId: activeSessionId,
- maxMessages: 800,
- },
+ { sessionId: activeSessionId, maxMessages: MAX_MESSAGES },
);
if (historyMessages.length > 0) {
setMessages(historyMessages);
@@ -691,6 +778,7 @@ export function useChatSession() {
});
}
+ // Token / cost bookkeeping
const inputTokens = result?.usage?.inputTokens ?? result?.inputTokens;
if (typeof inputTokens === "number") {
setTokensIn((prev) => prev + inputTokens);
@@ -712,48 +800,41 @@ export function useChatSession() {
typeof totalCost === "number")
) {
setMessages((prev) =>
- prev.map((message) => {
- if (message.id !== assistantMessageId) {
- return message;
- }
- return {
- ...message,
- meta: {
- ...(message.meta ?? {}),
- inputTokens:
- typeof inputTokens === "number"
- ? inputTokens
- : message.meta?.inputTokens,
- outputTokens:
- typeof outputTokens === "number"
- ? outputTokens
- : message.meta?.outputTokens,
- totalCost:
- typeof totalCost === "number"
- ? totalCost
- : message.meta?.totalCost,
- providerId: config.provider,
- modelId: config.model,
- },
- };
- }),
+ updateMessageById(prev, assistantMessageId, (msg) => ({
+ ...msg,
+ meta: {
+ ...(msg.meta ?? {}),
+ inputTokens:
+ typeof inputTokens === "number"
+ ? inputTokens
+ : msg.meta?.inputTokens,
+ outputTokens:
+ typeof outputTokens === "number"
+ ? outputTokens
+ : msg.meta?.outputTokens,
+ totalCost:
+ typeof totalCost === "number"
+ ? totalCost
+ : msg.meta?.totalCost,
+ providerId: config.provider,
+ modelId: config.model,
+ },
+ })),
);
}
if (result?.finishReason === "error") {
if (!resolvedAssistantText) {
const toolError = Array.isArray(result?.toolCalls)
- ? result.toolCalls.find((call) => call.error)?.error
+ ? result.toolCalls.find((c) => c.error)?.error
: undefined;
- addMessage({
- id: makeId("error"),
- sessionId: activeSessionId,
- role: "error",
- content:
+ addMessage(
+ makeErrorChatMessage(
+ activeSessionId,
toolError?.trim() ||
- "Runtime turn failed before an assistant response was produced.",
- createdAt: Date.now(),
- });
+ "Runtime turn failed before an assistant response was produced.",
+ ),
+ );
}
setStatus("failed");
} else if (result?.finishReason === "aborted") {
@@ -763,29 +844,31 @@ export function useChatSession() {
}
void refreshSessionDiffSummary(activeSessionId);
} catch (err) {
- const message = err instanceof Error ? err.message : String(err);
- setError(message);
- setStatus("error");
- addMessage({
- id: makeId("error"),
- sessionId: activeSessionId,
- role: "error",
- content: message,
- createdAt: Date.now(),
- });
+ if (optimisticQueuedPromptId) {
+ setPromptsInQueue((prev) =>
+ prev.filter((item) => item.id !== optimisticQueuedPromptId),
+ );
+ }
+ setErrorState(errorMessage(err), activeSessionId);
} finally {
- activeAssistantMessageIdRef.current = null;
- setActiveAssistantMessageId(null);
- liveToolMessageIdsRef.current = {};
- liveToolInputsRef.current = {};
+ if (!shouldQueue) {
+ activeAssistantMessageIdRef.current = null;
+ setActiveAssistantMessageId(null);
+ clearLiveToolRefs();
+ }
}
},
[
addMessage,
+ applyPromptsInQueue,
+ clearLiveToolRefs,
config,
materializeToolMessagesFromResult,
refreshSessionDiffSummary,
sessionId,
+ setErrorState,
+ startSession,
+ status,
postSession,
],
);
@@ -793,9 +876,7 @@ export function useChatSession() {
const respondToolApproval = useCallback(
async (requestId: string, approved: boolean) => {
const activeSessionId = activeSessionIdRef.current;
- if (!activeSessionId) {
- return;
- }
+ if (!activeSessionId) return;
await desktopClient.invoke("respond_tool_approval", {
sessionId: activeSessionId,
requestId,
@@ -812,23 +893,17 @@ export function useChatSession() {
);
const approveToolApproval = useCallback(
- async (requestId: string) => {
- await respondToolApproval(requestId, true);
- },
+ (requestId: string) => respondToolApproval(requestId, true),
[respondToolApproval],
);
const rejectToolApproval = useCallback(
- async (requestId: string) => {
- await respondToolApproval(requestId, false);
- },
+ (requestId: string) => respondToolApproval(requestId, false),
[respondToolApproval],
);
const abort = useCallback(async () => {
- if (!sessionId) {
- return;
- }
+ if (!sessionId) return;
try {
await postSession({ action: "abort", sessionId });
} catch {
@@ -837,10 +912,6 @@ export function useChatSession() {
setStatus("cancelled");
}, [sessionId, postSession]);
- const stop = useCallback(async () => {
- await abort();
- }, [abort]);
-
const reset = useCallback(async () => {
const activeSessionId = sessionId;
if (activeSessionId) {
@@ -856,19 +927,15 @@ export function useChatSession() {
setMessages([]);
setRawTranscript("");
setError(null);
- setToolCalls(0);
- setTokensIn(0);
- setTokensOut(0);
- setFileDiffs([]);
- setDiffSummary(EMPTY_DIFF_SUMMARY);
+ resetCounters();
activeSessionIdRef.current = null;
activeAssistantMessageIdRef.current = null;
setActiveAssistantMessageId(null);
setHydratedHistorySessionId(null);
setPendingToolApprovals([]);
- liveToolMessageIdsRef.current = {};
- liveToolInputsRef.current = {};
- }, [sessionId, postSession]);
+ setPromptsInQueue([]);
+ clearLiveToolRefs();
+ }, [sessionId, postSession, resetCounters, clearLiveToolRefs]);
const hydrateSession = useCallback(
async (session: SessionHistoryItem) => {
@@ -890,37 +957,36 @@ export function useChatSession() {
setActiveAssistantMessageId(null);
setHydratedHistorySessionId(session.sessionId);
setPendingToolApprovals([]);
- liveToolMessageIdsRef.current = {};
- liveToolInputsRef.current = {};
+ setPromptsInQueue([]);
+ clearLiveToolRefs();
+
+ const applyHydratedMessages = (
+ msgs: ChatMessage[],
+ sessionStatus: typeof session.status,
+ ) => {
+ setMessages(msgs);
+ setRawTranscript(msgs.map((m) => m.content).join("\n\n"));
+ resetCounters();
+ setStatus(inferHydratedChatStatus(sessionStatus, msgs));
+ void refreshSessionDiffSummary(session.sessionId);
+ };
try {
const historyMessages = await desktopClient.invoke(
"read_session_messages",
- {
- sessionId: session.sessionId,
- maxMessages: 800,
- },
+ { sessionId: session.sessionId, maxMessages: MAX_MESSAGES },
);
- if (hydrationRequestIdRef.current !== requestId) {
+ if (hydrationRequestIdRef.current !== requestId) return;
+
+ if (historyMessages.length > 0) {
+ void refreshPromptsInQueue(session.sessionId);
+ applyHydratedMessages(historyMessages, session.status);
return;
}
- if (historyMessages.length > 0) {
- setMessages(historyMessages);
- setRawTranscript(
- historyMessages.map((message) => message.content).join("\n\n"),
- );
- setToolCalls(0);
- setTokensIn(0);
- setTokensOut(0);
- setFileDiffs([]);
- setStatus(inferHydratedChatStatus(session.status, historyMessages));
- void refreshSessionDiffSummary(session.sessionId);
- return;
- }
- let synthesizedMessages: ChatMessage[] = [];
+ const synthesized: ChatMessage[] = [];
if (session.prompt?.trim()) {
- synthesizedMessages.push({
+ synthesized.push({
id: makeId("history_user"),
sessionId: session.sessionId,
role: "user",
@@ -931,60 +997,59 @@ export function useChatSession() {
try {
const transcript = await desktopClient.invoke(
"read_session_transcript",
- {
- sessionId: session.sessionId,
- maxChars: 20000,
- },
+ { sessionId: session.sessionId, maxChars: 20000 },
);
const text = transcript.trim();
if (text) {
- synthesizedMessages = [
- ...synthesizedMessages,
- {
- id: makeId("history_assistant"),
- sessionId: session.sessionId,
- role: "assistant",
- content: text,
- createdAt: Date.now(),
- },
- ];
+ synthesized.push({
+ id: makeId("history_assistant"),
+ sessionId: session.sessionId,
+ role: "assistant",
+ content: text,
+ createdAt: Date.now(),
+ });
}
} catch {
// Ignore transcript fallback failures.
}
- setMessages(synthesizedMessages);
- setRawTranscript(
- synthesizedMessages.map((message) => message.content).join("\n\n"),
- );
- setToolCalls(0);
- setTokensIn(0);
- setTokensOut(0);
- setFileDiffs([]);
- setStatus(inferHydratedChatStatus(session.status, synthesizedMessages));
- void refreshSessionDiffSummary(session.sessionId);
+ void refreshPromptsInQueue(session.sessionId);
+ applyHydratedMessages(synthesized, session.status);
} catch (err) {
- if (hydrationRequestIdRef.current !== requestId) {
- return;
- }
- const message = err instanceof Error ? err.message : String(err);
- setError(message);
+ if (hydrationRequestIdRef.current !== requestId) return;
+ const msg = errorMessage(err);
+ setError(msg);
setStatus("error");
- setMessages([
- {
- id: makeId("error"),
- sessionId: session.sessionId,
- role: "error",
- content: message,
- createdAt: Date.now(),
- },
- ]);
+ setMessages([makeErrorChatMessage(session.sessionId, msg)]);
} finally {
if (hydrationRequestIdRef.current === requestId) {
setIsHydratingSession(false);
}
}
},
- [refreshSessionDiffSummary],
+ [
+ clearLiveToolRefs,
+ refreshPromptsInQueue,
+ refreshSessionDiffSummary,
+ resetCounters,
+ ],
+ );
+
+ const steerPromptInQueue = useCallback(
+ async (promptId: string) => {
+ const activeSessionId = activeSessionIdRef.current;
+ if (!activeSessionId || !promptId.trim()) {
+ return;
+ }
+ const payload = await postSession({
+ action: "steer_prompt",
+ sessionId: activeSessionId,
+ promptId,
+ });
+ setPromptsInQueue(
+ Array.isArray(payload.promptsInQueue) ? payload.promptsInQueue : [],
+ );
+ },
+ [postSession],
);
const summary = useMemo(
@@ -1016,15 +1081,17 @@ export function useChatSession() {
error,
summary,
fileDiffs,
+ promptsInQueue,
pendingToolApprovals,
setConfig,
start,
hydrateSession,
sendPrompt,
+ steerPromptInQueue,
approveToolApproval,
rejectToolApproval,
abort,
- stop,
+ stop: abort,
reset,
};
}
diff --git a/sdk/apps/code/host/runtime-bridge.ts b/sdk/apps/code/host/runtime-bridge.ts
index 873d3b3315..ff0814208b 100644
--- a/sdk/apps/code/host/runtime-bridge.ts
+++ b/sdk/apps/code/host/runtime-bridge.ts
@@ -33,31 +33,37 @@ import {
DEFAULT_RPC_CLIENT_TYPE,
type HostContext,
type JsonRecord,
+ type LiveSession,
+ type PromptInQueue,
+ type QueuedChatTurn,
type ToolApprovalRequestItem,
} from "./types";
+// ---------------------------------------------------------------------------
+// Config helpers
+// ---------------------------------------------------------------------------
+
+function getNestedString(obj: unknown, ...keys: string[]): string | undefined {
+ let current: unknown = obj;
+ for (const key of keys) {
+ if (!current || typeof current !== "object") return undefined;
+ current = (current as JsonRecord)[key];
+ }
+ return typeof current === "string" ? current : undefined;
+}
+
function setRuntimeHomeDir(config: unknown) {
- if (!config || typeof config !== "object") {
+ const homeDir = getNestedString(config, "sessions", "homeDir")?.trim();
+ if (homeDir) {
+ setHomeDir(homeDir);
+ } else {
setHomeDirIfUnset(homedir());
- return;
}
- const sessions = (config as JsonRecord).sessions;
- const homeDir =
- sessions && typeof sessions === "object"
- ? ((sessions as JsonRecord).homeDir as string | undefined)
- : undefined;
- const normalized = homeDir?.trim();
- if (normalized) {
- setHomeDir(normalized);
- return;
- }
- setHomeDirIfUnset(homedir());
}
function addRuntimeLoggerContext(config: unknown) {
- if (!config || typeof config !== "object") {
- return;
- }
+ if (!config || typeof config !== "object") return;
+
const record = config as JsonRecord;
const existing =
record.logger && typeof record.logger === "object"
@@ -67,6 +73,7 @@ function addRuntimeLoggerContext(config: unknown) {
existing.bindings && typeof existing.bindings === "object"
? { ...(existing.bindings as JsonRecord) }
: {};
+
record.logger = {
...existing,
name:
@@ -81,30 +88,35 @@ function addRuntimeLoggerContext(config: unknown) {
};
}
+// ---------------------------------------------------------------------------
+// Bridge script resolution
+// ---------------------------------------------------------------------------
+
+const BRIDGE_SCRIPT = "chat-runtime-bridge.ts";
+const BRIDGE_SEARCH_DIRS = [
+ ["apps", "code", "scripts"],
+ ["packages", "app", "scripts"],
+ ["app", "scripts"],
+];
+
function resolveChatRuntimeBridgeScriptPath(ctx: HostContext): string | null {
- const candidates = [
- join(
- ctx.workspaceRoot,
- "apps",
- "code",
- "scripts",
- "chat-runtime-bridge.ts",
- ),
- join(
- ctx.workspaceRoot,
- "packages",
- "app",
- "scripts",
- "chat-runtime-bridge.ts",
- ),
- join(ctx.workspaceRoot, "app", "scripts", "chat-runtime-bridge.ts"),
- join(process.cwd(), "app", "scripts", "chat-runtime-bridge.ts"),
- join(process.cwd(), "..", "scripts", "chat-runtime-bridge.ts"),
- join(process.cwd(), "scripts", "chat-runtime-bridge.ts"),
- ];
- return candidates.find((candidate) => existsSync(candidate)) ?? null;
+ for (const segments of BRIDGE_SEARCH_DIRS) {
+ const candidate = join(ctx.workspaceRoot, ...segments, BRIDGE_SCRIPT);
+ if (existsSync(candidate)) return candidate;
+ }
+ for (const base of [process.cwd()]) {
+ for (const rel of [["app", "scripts"], ["..", "scripts"], ["scripts"]]) {
+ const candidate = join(base, ...rel, BRIDGE_SCRIPT);
+ if (existsSync(candidate)) return candidate;
+ }
+ }
+ return null;
}
+// ---------------------------------------------------------------------------
+// Child process line reader
+// ---------------------------------------------------------------------------
+
function readChildLines(
stream: NodeJS.ReadableStream,
onLine: (line: string) => void,
@@ -112,28 +124,104 @@ function readChildLines(
let buffer = "";
stream.on("data", (chunk) => {
buffer += String(chunk);
- let newlineIndex = buffer.indexOf("\n");
- while (newlineIndex >= 0) {
- const line = buffer.slice(0, newlineIndex).trim();
- buffer = buffer.slice(newlineIndex + 1);
- if (line) {
- onLine(line);
- }
- newlineIndex = buffer.indexOf("\n");
+ let idx = buffer.indexOf("\n");
+ while (idx >= 0) {
+ const line = buffer.slice(0, idx).trim();
+ buffer = buffer.slice(idx + 1);
+ if (line) onLine(line);
+ idx = buffer.indexOf("\n");
}
});
}
+// ---------------------------------------------------------------------------
+// Bridge lifecycle
+// ---------------------------------------------------------------------------
+
+function handleBridgeStdoutLine(ctx: HostContext, parsed: JsonRecord) {
+ const type = String(parsed.type ?? "");
+ const sessionId =
+ typeof parsed.sessionId === "string" ? parsed.sessionId : "";
+
+ switch (type) {
+ case "ready":
+ ctx.bridgeReady = true;
+ return;
+
+ case "response": {
+ const requestId = String(parsed.requestId ?? "");
+ const pending = ctx.pendingBridge.get(requestId);
+ if (!pending) return;
+ ctx.pendingBridge.delete(requestId);
+ if (typeof parsed.error === "string" && parsed.error.trim()) {
+ pending.reject(new Error(parsed.error));
+ } else {
+ pending.resolve(parsed.response ?? null);
+ }
+ return;
+ }
+
+ case "chat_text":
+ emitChunk(ctx, sessionId, "chat_text", String(parsed.chunk ?? ""));
+ return;
+
+ case "tool_call_start":
+ emitChunk(
+ ctx,
+ sessionId,
+ "chat_tool_call_start",
+ JSON.stringify({
+ toolCallId: parsed.toolCallId,
+ toolName: parsed.toolName,
+ input: parsed.input,
+ }),
+ );
+ return;
+
+ case "tool_call_end":
+ emitChunk(
+ ctx,
+ sessionId,
+ "chat_tool_call_end",
+ JSON.stringify({
+ toolCallId: parsed.toolCallId,
+ toolName: parsed.toolName,
+ output: parsed.output,
+ error: parsed.error,
+ durationMs: parsed.durationMs,
+ }),
+ );
+ return;
+
+ case "error": {
+ const message =
+ typeof parsed.message === "string"
+ ? parsed.message
+ : "chat runtime bridge error";
+ if (sessionId) {
+ emitChunk(
+ ctx,
+ sessionId,
+ "chat_core_log",
+ JSON.stringify({ level: "error", message }),
+ );
+ } else {
+ console.error("[chat-runtime-bridge]", message);
+ }
+ return;
+ }
+ }
+}
+
export function ensureBridgeStarted(ctx: HostContext) {
- if (ctx.bridgeChild && ctx.bridgeChild.exitCode === null && ctx.bridgeReady) {
- return;
- }
+ if (ctx.bridgeChild?.exitCode === null && ctx.bridgeReady) return;
+
const scriptPath = resolveChatRuntimeBridgeScriptPath(ctx);
- if (!scriptPath) {
- throw new Error("chat runtime bridge script not found");
- }
+ if (!scriptPath) throw new Error("chat runtime bridge script not found");
+
ctx.bridgeReady = false;
mkdirSync(toolApprovalDir(), { recursive: true });
+
ctx.bridgeChild = spawn("bun", [scriptPath], {
cwd: ctx.workspaceRoot,
env: {
@@ -146,191 +234,139 @@ export function ensureBridgeStarted(ctx: HostContext) {
},
stdio: ["pipe", "pipe", "pipe"],
});
+
readChildLines(ctx.bridgeChild.stdout, (line) => {
- const parsed = JSON.parse(line) as JsonRecord;
- const type = String(parsed.type ?? "");
- if (type === "ready") {
- ctx.bridgeReady = true;
- return;
- }
- if (type === "response") {
- const requestId = String(parsed.requestId ?? "");
- const pending = ctx.pendingBridge.get(requestId);
- if (!pending) {
- return;
- }
- ctx.pendingBridge.delete(requestId);
- if (typeof parsed.error === "string" && parsed.error.trim()) {
- pending.reject(new Error(parsed.error));
- return;
- }
- pending.resolve(parsed.response ?? null);
- return;
- }
- if (type === "chat_text") {
- emitChunk(
- ctx,
- String(parsed.sessionId ?? ""),
- "chat_text",
- String(parsed.chunk ?? ""),
- );
- return;
- }
- if (type === "tool_call_start") {
- emitChunk(
- ctx,
- String(parsed.sessionId ?? ""),
- "chat_tool_call_start",
- JSON.stringify({
- toolCallId: parsed.toolCallId,
- toolName: parsed.toolName,
- input: parsed.input,
- }),
- );
- return;
- }
- if (type === "tool_call_end") {
- emitChunk(
- ctx,
- String(parsed.sessionId ?? ""),
- "chat_tool_call_end",
- JSON.stringify({
- toolCallId: parsed.toolCallId,
- toolName: parsed.toolName,
- output: parsed.output,
- error: parsed.error,
- durationMs: parsed.durationMs,
- }),
- );
- return;
- }
- if (type === "error") {
- const sessionId =
- typeof parsed.sessionId === "string" ? parsed.sessionId : "";
- const message =
- typeof parsed.message === "string"
- ? parsed.message
- : "chat runtime bridge error";
- if (sessionId) {
- emitChunk(
- ctx,
- sessionId,
- "chat_core_log",
- JSON.stringify({
- level: "error",
- message,
- }),
- );
- return;
- }
- console.error("[chat-runtime-bridge]", message);
- }
+ handleBridgeStdoutLine(ctx, JSON.parse(line) as JsonRecord);
});
+
readChildLines(ctx.bridgeChild.stderr, (line) => {
console.error("[chat-runtime-bridge]", line);
});
+
ctx.bridgeChild.on("exit", () => {
ctx.bridgeReady = false;
ctx.bridgeChild = null;
- for (const [requestId, pending] of ctx.pendingBridge.entries()) {
- ctx.pendingBridge.delete(requestId);
+ for (const [, pending] of ctx.pendingBridge) {
pending.reject(new Error("chat runtime bridge exited"));
}
+ ctx.pendingBridge.clear();
});
}
+// ---------------------------------------------------------------------------
+// Bridge RPC
+// ---------------------------------------------------------------------------
+
export async function runBridgeCommand(
ctx: HostContext,
command: Record,
): Promise {
ensureBridgeStarted(ctx);
const child = ctx.bridgeChild;
- if (!child || !child.stdin) {
- throw new Error("chat runtime bridge unavailable");
- }
+ if (!child?.stdin) throw new Error("chat runtime bridge unavailable");
+
const requestId = `bridge_${ctx.bridgeRequestId++}`;
- const envelope = JSON.stringify({
- type: "request",
- requestId,
- command,
- });
- return await new Promise((resolve, reject) => {
+ const envelope = JSON.stringify({ type: "request", requestId, command });
+
+ return new Promise((resolve, reject) => {
ctx.pendingBridge.set(requestId, { resolve, reject });
child.stdin.write(`${envelope}\n`, (error) => {
- if (!error) {
- return;
+ if (error) {
+ ctx.pendingBridge.delete(requestId);
+ reject(error);
}
- ctx.pendingBridge.delete(requestId);
- reject(error);
});
});
}
+// ---------------------------------------------------------------------------
+// Tool approval helpers
+// ---------------------------------------------------------------------------
+
+function sendApprovalSnapshot(ctx: HostContext, sessionId: string) {
+ sendEvent(ctx, "tool_approval_state", {
+ sessionId,
+ items: listPendingToolApprovalsForSession(sessionId, 50),
+ });
+}
+
export function listPendingToolApprovalsForSession(
sessionId: string,
limit = 20,
): ToolApprovalRequestItem[] {
const dir = toolApprovalDir();
- if (!existsSync(dir)) {
- return [];
- }
- const items: ToolApprovalRequestItem[] = [];
+ if (!existsSync(dir)) return [];
+
const prefix = toolApprovalRequestPrefix(sessionId);
+ const items: ToolApprovalRequestItem[] = [];
+
for (const entry of readdirSync(dir, { withFileTypes: true })) {
- if (!entry.isFile()) {
+ if (
+ !entry.isFile() ||
+ !entry.name.startsWith(prefix) ||
+ !entry.name.endsWith(".json")
+ )
continue;
- }
- if (!entry.name.startsWith(prefix) || !entry.name.endsWith(".json")) {
- continue;
- }
try {
- const parsed = JSON.parse(
- readFileSync(join(dir, entry.name), "utf8"),
- ) as ToolApprovalRequestItem;
- items.push(parsed);
+ items.push(
+ JSON.parse(
+ readFileSync(join(dir, entry.name), "utf8"),
+ ) as ToolApprovalRequestItem,
+ );
} catch {
// Ignore malformed approval files.
}
}
- items.sort((left, right) => left.createdAt.localeCompare(right.createdAt));
+
+ items.sort((a, b) => a.createdAt.localeCompare(b.createdAt));
return items.slice(0, Math.max(1, limit));
}
export function broadcastApprovalSnapshots(ctx: HostContext) {
const dir = toolApprovalDir();
- if (!existsSync(dir)) {
- return;
- }
+ if (!existsSync(dir)) return;
+
const sessionIds = new Set();
for (const entry of readdirSync(dir, { withFileTypes: true })) {
- if (!entry.isFile() || !entry.name.includes(".request.")) {
- continue;
- }
- const [sessionId] = entry.name.split(".request.");
- if (sessionId?.trim()) {
- sessionIds.add(sessionId.trim());
- }
+ if (!entry.isFile() || !entry.name.includes(".request.")) continue;
+ const id = entry.name.split(".request.")[0]?.trim();
+ if (id) sessionIds.add(id);
}
+
for (const sessionId of sessionIds) {
- sendEvent(ctx, "tool_approval_state", {
- sessionId,
- items: listPendingToolApprovalsForSession(sessionId, 50),
- });
+ sendApprovalSnapshot(ctx, sessionId);
}
}
+function makeQueuedTurnId(): string {
+ return `queued_${Date.now()}_${Math.random().toString(36).slice(2, 8)}`;
+}
+
+function getPromptsInQueue(session: LiveSession): PromptInQueue[] {
+ return session.pendingTurns.map((turn) => ({
+ id: turn.id,
+ prompt: turn.prompt,
+ steer: turn.steer,
+ }));
+}
+
+function sendPromptsInQueueSnapshot(ctx: HostContext, sessionId: string) {
+ const session = ctx.liveSessions.get(sessionId);
+ sendEvent(ctx, "prompts_in_queue_state", {
+ sessionId,
+ items: session ? getPromptsInQueue(session) : [],
+ });
+}
+
export function ensureApprovalWatcher(ctx: HostContext) {
- if (ctx.approvalWatcher) {
- return;
- }
+ if (ctx.approvalWatcher) return;
mkdirSync(toolApprovalDir(), { recursive: true });
ctx.approvalWatcher = watch(toolApprovalDir(), () => {
- if (ctx.approvalBroadcastTimer) {
- clearTimeout(ctx.approvalBroadcastTimer);
- }
- ctx.approvalBroadcastTimer = setTimeout(() => {
- broadcastApprovalSnapshots(ctx);
- }, 50);
+ if (ctx.approvalBroadcastTimer) clearTimeout(ctx.approvalBroadcastTimer);
+ ctx.approvalBroadcastTimer = setTimeout(
+ () => broadcastApprovalSnapshots(ctx),
+ 50,
+ );
});
}
@@ -340,9 +376,9 @@ export async function respondToolApproval(
) {
const sessionId = String(args?.sessionId ?? "").trim();
const requestId = String(args?.requestId ?? "").trim();
- if (!sessionId || !requestId) {
+ if (!sessionId || !requestId)
throw new Error("sessionId and requestId are required");
- }
+
const path = toolApprovalDecisionPath(sessionId, requestId);
mkdirSync(dirname(path), { recursive: true });
writeFileSync(
@@ -353,166 +389,317 @@ export async function respondToolApproval(
ts: nowMs(),
}),
);
+
const requestPath = join(
toolApprovalDir(),
`${sessionId}.request.${requestId}.json`,
);
- if (existsSync(requestPath)) {
- unlinkSync(requestPath);
- }
- sendEvent(ctx, "tool_approval_state", {
- sessionId,
- items: listPendingToolApprovalsForSession(sessionId, 50),
- });
+ if (existsSync(requestPath)) unlinkSync(requestPath);
+
+ sendApprovalSnapshot(ctx, sessionId);
return true;
}
-export async function handleChatSessionCommand(
+// ---------------------------------------------------------------------------
+// Chat turn execution
+// ---------------------------------------------------------------------------
+
+async function executeChatTurn(
ctx: HostContext,
- request: ChatSessionCommandRequest,
-): Promise {
- if (request.action === "start") {
- if (!request.config) {
- throw new Error("missing config for start action");
- }
- setRuntimeHomeDir(request.config);
- addRuntimeLoggerContext(request.config);
- const response = (await runBridgeCommand(ctx, {
- action: "start",
- config: request.config,
- })) as { sessionId?: string };
- const sessionId = response.sessionId?.trim();
- if (!sessionId) {
- throw new Error("chat runtime bridge start response missing session id");
- }
- await runBridgeCommand(ctx, {
- action: "set_sessions",
- sessionIds: [sessionId],
- });
- ctx.liveSessions.set(sessionId, {
- config: request.config,
- messages: [],
- busy: false,
- startedAt: nowMs(),
- status: "idle",
- });
- return { sessionId };
- }
+ sessionId: string,
+ session: LiveSession,
+ turn: QueuedChatTurn,
+): Promise {
+ if (turn.config) session.config = turn.config;
+ if (turn.prompt) session.prompt = turn.prompt;
- if (request.action === "send") {
- const prompt = request.prompt?.trim() || "";
- const hasAttachments =
- (request.attachments?.userImages?.length ?? 0) > 0 ||
- (request.attachments?.userFiles?.length ?? 0) > 0;
- if (!prompt && !hasAttachments) {
- throw new Error("prompt is required for send action");
- }
- const sessionId = request.sessionId?.trim();
- if (!sessionId) {
- throw new Error("sessionId is required for send action");
- }
+ session.busy = true;
+ session.status = "running";
+ session.endedAt = undefined;
- let session = ctx.liveSessions.get(sessionId);
- if (!session) {
- if (!request.config) {
- throw new Error("session not found. start a new session.");
- }
- const messages = readPersistedChatMessages(sessionId);
- if (!messages) {
- throw new Error("session not found. start a new session.");
- }
- session = {
- config: request.config,
- messages,
- busy: false,
- startedAt: nowMs(),
- status: "idle",
- prompt: derivePromptFromMessages(messages),
- title: readSessionMetadataTitle(sessionId),
- };
- ctx.liveSessions.set(sessionId, session);
- }
- if (request.config) {
- session.config = request.config;
- }
- if (session.busy) {
- throw new Error("session is busy. wait for current response to finish.");
- }
- session.busy = true;
- session.status = "running";
- session.endedAt = undefined;
- if (prompt) {
- session.prompt = prompt;
- }
- setRuntimeHomeDir(session.config);
- addRuntimeLoggerContext(session.config);
- await runBridgeCommand(ctx, {
- action: "set_sessions",
- sessionIds: [sessionId],
- });
+ setRuntimeHomeDir(session.config);
+ addRuntimeLoggerContext(session.config);
+
+ await runBridgeCommand(ctx, {
+ action: "set_sessions",
+ sessionIds: [sessionId],
+ });
+
+ try {
const resultEnvelope = (await runBridgeCommand(ctx, {
action: "send",
sessionId,
request: {
config: session.config,
messages: session.messages,
- prompt,
- attachments: request.attachments,
+ prompt: turn.prompt,
+ attachments: turn.attachments,
},
})) as { result?: ChatTurnResult };
- const result = resultEnvelope.result;
- if (!result) {
- throw new Error("chat runtime bridge send response missing result");
- }
- const persistedMessages = persistUsageInMessages(
+ const result = resultEnvelope.result;
+ if (!result)
+ throw new Error("chat runtime bridge send response missing result");
+
+ session.messages = persistUsageInMessages(
(Array.isArray(result.messages) ? result.messages : []) as unknown[],
session.config,
result,
);
- session.messages = persistedMessages;
- session.busy = false;
session.status = normalizeChatFinishStatus(result.finishReason);
session.endedAt = nowMs();
- persistSessionMessages(sessionId, persistedMessages);
- sendEvent(ctx, "tool_approval_state", {
- sessionId,
- items: listPendingToolApprovalsForSession(sessionId, 50),
- });
- return {
- sessionId,
- result,
- };
- }
- if (request.action === "abort") {
- const sessionId = request.sessionId?.trim();
- if (sessionId) {
- await runBridgeCommand(ctx, { action: "abort", sessionId });
- const session = ctx.liveSessions.get(sessionId);
- if (session) {
- session.busy = false;
- session.status = "cancelled";
- session.endedAt = nowMs();
- }
- }
- return {
- sessionId: request.sessionId,
- ok: true,
- };
- }
+ persistSessionMessages(sessionId, session.messages);
+ sendApprovalSnapshot(ctx, sessionId);
- if (request.action === "reset") {
- const sessionId = request.sessionId?.trim();
- if (sessionId) {
- ctx.liveSessions.delete(sessionId);
- await runBridgeCommand(ctx, { action: "reset", sessionId });
- }
- return {
- sessionId: request.sessionId,
- ok: true,
- };
+ return result;
+ } finally {
+ session.busy = false;
}
+}
- throw new Error("unsupported action");
+async function drainSessionQueue(
+ ctx: HostContext,
+ sessionId: string,
+): Promise {
+ const session = ctx.liveSessions.get(sessionId);
+ if (!session || session.busy) return;
+
+ const nextTurn = session.pendingTurns.shift();
+ if (!nextTurn) return;
+ sendPromptsInQueueSnapshot(ctx, sessionId);
+ emitChunk(
+ ctx,
+ sessionId,
+ "chat_queued_prompt_start",
+ JSON.stringify({
+ prompt: nextTurn.prompt,
+ attachmentCount:
+ (nextTurn.attachments?.userImages?.length ?? 0) +
+ (nextTurn.attachments?.userFiles?.length ?? 0),
+ }),
+ );
+
+ try {
+ await executeChatTurn(ctx, sessionId, session, nextTurn);
+ } catch (error) {
+ session.status = "error";
+ session.endedAt = nowMs();
+ emitChunk(
+ ctx,
+ sessionId,
+ "chat_core_log",
+ JSON.stringify({
+ level: "error",
+ message: error instanceof Error ? error.message : String(error),
+ }),
+ );
+ }
+
+ if (session.pendingTurns.length > 0) {
+ void drainSessionQueue(ctx, sessionId);
+ }
+}
+
+// ---------------------------------------------------------------------------
+// Command handler
+// ---------------------------------------------------------------------------
+
+function createLiveSession(
+ config: JsonRecord,
+ extra?: Partial,
+): LiveSession {
+ return {
+ config,
+ messages: [],
+ pendingTurns: [],
+ busy: false,
+ startedAt: nowMs(),
+ status: "idle",
+ ...extra,
+ };
+}
+
+async function handleStart(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+) {
+ if (!request.config) throw new Error("missing config for start action");
+
+ setRuntimeHomeDir(request.config);
+ addRuntimeLoggerContext(request.config);
+
+ const response = (await runBridgeCommand(ctx, {
+ action: "start",
+ config: request.config,
+ })) as { sessionId?: string };
+
+ const sessionId = response.sessionId?.trim();
+ if (!sessionId)
+ throw new Error("chat runtime bridge start response missing session id");
+
+ await runBridgeCommand(ctx, {
+ action: "set_sessions",
+ sessionIds: [sessionId],
+ });
+ ctx.liveSessions.set(sessionId, createLiveSession(request.config));
+
+ return { sessionId };
+}
+
+async function handleSend(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+) {
+ const prompt = request.prompt?.trim() || "";
+ const hasAttachments =
+ (request.attachments?.userImages?.length ?? 0) > 0 ||
+ (request.attachments?.userFiles?.length ?? 0) > 0;
+ if (!prompt && !hasAttachments)
+ throw new Error("prompt is required for send action");
+
+ const sessionId = request.sessionId?.trim();
+ if (!sessionId) throw new Error("sessionId is required for send action");
+
+ let session = ctx.liveSessions.get(sessionId);
+ if (!session) {
+ if (!request.config)
+ throw new Error("session not found. start a new session.");
+ const messages = readPersistedChatMessages(sessionId);
+ if (!messages) throw new Error("session not found. start a new session.");
+
+ session = createLiveSession(request.config, {
+ messages,
+ prompt: derivePromptFromMessages(messages),
+ title: readSessionMetadataTitle(sessionId),
+ });
+ ctx.liveSessions.set(sessionId, session);
+ }
+
+ if (request.config) session.config = request.config;
+
+ const turn: QueuedChatTurn = {
+ id: makeQueuedTurnId(),
+ prompt,
+ steer: false,
+ config: request.config,
+ attachments: request.attachments,
+ };
+
+ if (session.busy) {
+ session.pendingTurns.push(turn);
+ sendPromptsInQueueSnapshot(ctx, sessionId);
+ return {
+ sessionId,
+ ok: true,
+ queued: true,
+ promptsInQueue: getPromptsInQueue(session),
+ };
+ }
+
+ const result = await executeChatTurn(ctx, sessionId, session, turn);
+ if (session.pendingTurns.length > 0) {
+ sendPromptsInQueueSnapshot(ctx, sessionId);
+ void drainSessionQueue(ctx, sessionId);
+ }
+
+ return {
+ sessionId,
+ result,
+ queued: false,
+ promptsInQueue: getPromptsInQueue(session),
+ };
+}
+
+async function handleAbort(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+) {
+ const sessionId = request.sessionId?.trim();
+ if (sessionId) {
+ await runBridgeCommand(ctx, { action: "abort", sessionId });
+ const session = ctx.liveSessions.get(sessionId);
+ if (session) {
+ session.busy = false;
+ session.pendingTurns = [];
+ session.status = "cancelled";
+ session.endedAt = nowMs();
+ }
+ sendPromptsInQueueSnapshot(ctx, sessionId);
+ }
+ return { sessionId: request.sessionId, ok: true };
+}
+
+async function handleReset(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+) {
+ const sessionId = request.sessionId?.trim();
+ if (sessionId) {
+ ctx.liveSessions.delete(sessionId);
+ await runBridgeCommand(ctx, { action: "reset", sessionId });
+ sendPromptsInQueueSnapshot(ctx, sessionId);
+ }
+ return { sessionId: request.sessionId, ok: true };
+}
+
+async function handlePendingPrompts(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+) {
+ const sessionId = request.sessionId?.trim();
+ if (!sessionId) throw new Error("sessionId is required");
+ const session = ctx.liveSessions.get(sessionId);
+ return {
+ sessionId,
+ promptsInQueue: session ? getPromptsInQueue(session) : [],
+ };
+}
+
+async function handleSteerPrompt(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+) {
+ const sessionId = request.sessionId?.trim();
+ const promptId = request.promptId?.trim();
+ if (!sessionId || !promptId) {
+ throw new Error("sessionId and promptId are required");
+ }
+ const session = ctx.liveSessions.get(sessionId);
+ if (!session) {
+ return { sessionId, promptsInQueue: [] };
+ }
+ const existingIndex = session.pendingTurns.findIndex(
+ (turn) => turn.id === promptId,
+ );
+ if (existingIndex >= 0) {
+ const [turn] = session.pendingTurns.splice(existingIndex, 1);
+ session.pendingTurns.unshift({ ...turn, steer: true });
+ }
+ sendPromptsInQueueSnapshot(ctx, sessionId);
+ return {
+ sessionId,
+ promptsInQueue: getPromptsInQueue(session),
+ };
+}
+
+const ACTION_HANDLERS: Record<
+ string,
+ (ctx: HostContext, req: ChatSessionCommandRequest) => Promise
+> = {
+ start: handleStart,
+ send: handleSend,
+ abort: handleAbort,
+ reset: handleReset,
+ pending_prompts: handlePendingPrompts,
+ steer_prompt: handleSteerPrompt,
+};
+
+export async function handleChatSessionCommand(
+ ctx: HostContext,
+ request: ChatSessionCommandRequest,
+): Promise {
+ const handler = ACTION_HANDLERS[request.action];
+ if (!handler) throw new Error("unsupported action");
+ return handler(ctx, request);
}
diff --git a/sdk/apps/code/host/types.ts b/sdk/apps/code/host/types.ts
index afa13e5b35..8616e6b71a 100644
--- a/sdk/apps/code/host/types.ts
+++ b/sdk/apps/code/host/types.ts
@@ -28,16 +28,38 @@ export type ChatTurnResult = RpcChatTurnResult & {
};
export type ChatSessionCommandRequest = {
- action: "start" | "send" | "abort" | "reset";
+ action:
+ | "start"
+ | "send"
+ | "abort"
+ | "reset"
+ | "pending_prompts"
+ | "steer_prompt";
sessionId?: string;
prompt?: string;
+ promptId?: string;
config?: JsonRecord;
attachments?: ChatTurnAttachments;
};
+export type QueuedChatTurn = {
+ id: string;
+ prompt: string;
+ steer: boolean;
+ config?: JsonRecord;
+ attachments?: ChatTurnAttachments;
+};
+
+export type PromptInQueue = {
+ id: string;
+ prompt: string;
+ steer: boolean;
+};
+
export type LiveSession = {
config: JsonRecord;
messages: unknown[];
+ pendingTurns: QueuedChatTurn[];
busy: boolean;
startedAt: number;
endedAt?: number;
diff --git a/sdk/apps/code/src-tauri/src/main.rs b/sdk/apps/code/src-tauri/src/main.rs
index eed17bf831..b2fc70098c 100644
--- a/sdk/apps/code/src-tauri/src/main.rs
+++ b/sdk/apps/code/src-tauri/src/main.rs
@@ -260,6 +260,7 @@ struct ChatSessionCommandResponse {
session_id: Option,
result: Option,
ok: Option,
+ queued: Option,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -349,6 +350,7 @@ impl ChatWsBridgeState {
struct ChatRuntimeSession {
config: StartSessionRequest,
messages: Vec,
+ pending_turns: Vec,
busy: bool,
started_at: u64,
ended_at: Option,
@@ -1369,9 +1371,7 @@ fn extract_message_notice_meta(message: &Value) -> Option Option
.map(|parent| parent.join("Resources").join("code-host"))
}),
];
- candidates
- .into_iter()
- .flatten()
- .find(|path| path.exists())
+ candidates.into_iter().flatten().find(|path| path.exists())
}
fn ensure_desktop_backend_started(
@@ -3012,6 +3009,158 @@ fn run_chat_turn_via_rpc_runtime(
.map_err(|e| format!("invalid chat runtime bridge result: {e}"))
}
+fn persist_chat_turn_result(
+ state: &Arc,
+ session_id: &str,
+ config: &StartSessionRequest,
+ result: &ChatTurnResult,
+) -> Result<(), String> {
+ let mut sessions = state
+ .sessions
+ .lock()
+ .map_err(|_| "failed to lock chat session store")?;
+ if let Some(session) = sessions.get_mut(session_id) {
+ let persisted_messages = persist_usage_in_messages(&result.messages, config, result);
+ session.messages = persisted_messages.clone();
+ session.status = normalize_chat_finish_status(result.finish_reason.as_deref());
+ session.ended_at = Some(now_ms());
+ if let Some(path) = shared_session_messages_write_path(session_id) {
+ if let Some(parent) = path.parent() {
+ let _ = fs::create_dir_all(parent);
+ }
+ let body = serde_json::json!({
+ "messages": persisted_messages,
+ "ts": now_ms(),
+ });
+ if let Ok(encoded) = serde_json::to_vec(&body) {
+ let _ = fs::write(path, encoded);
+ }
+ }
+ }
+ Ok(())
+}
+
+fn mark_chat_turn_failed(state: &Arc, session_id: &str) -> Result<(), String> {
+ let mut sessions = state
+ .sessions
+ .lock()
+ .map_err(|_| "failed to lock chat session store")?;
+ if let Some(session) = sessions.get_mut(session_id) {
+ session.status = "failed".to_string();
+ session.ended_at = Some(now_ms());
+ }
+ Ok(())
+}
+
+fn dequeue_next_chat_turn(
+ state: &Arc,
+ session_id: &str,
+) -> Result