diff --git a/sdk/apps/cli/src/commands/rpc-runtime.ts b/sdk/apps/cli/src/commands/rpc-runtime.ts index 873841ff26..f0b61c4951 100644 --- a/sdk/apps/cli/src/commands/rpc-runtime.ts +++ b/sdk/apps/cli/src/commands/rpc-runtime.ts @@ -417,9 +417,10 @@ export function createRpcRuntimeHandlers(): RpcRuntimeHandlers { prompt: input, userImages, userFiles: fileMaterialized.paths, + delivery: request.delivery, }); if (!result) { - throw new Error("runtime send returned no result"); + return { queued: true }; } runtimeLogger.info?.("RPC runtime turn send completed", { sessionId, @@ -468,6 +469,7 @@ export function createRpcRuntimeHandlers(): RpcRuntimeHandlers { prompt: input, userImages, userFiles: fileMaterialized.paths, + delivery: request.delivery, }); } catch (restoredError) { runtimeLogger.error?.( diff --git a/sdk/apps/cli/src/commands/rpc-runtime/event-bridge.ts b/sdk/apps/cli/src/commands/rpc-runtime/event-bridge.ts index 1ad59ae791..c76be15bc0 100644 --- a/sdk/apps/cli/src/commands/rpc-runtime/event-bridge.ts +++ b/sdk/apps/cli/src/commands/rpc-runtime/event-bridge.ts @@ -147,6 +147,17 @@ export function subscribeRuntimeEventBridge(input: { }); return; } + if (coreEvent.type === "pending_prompts") { + publishRuntimeEvent({ + eventClient: input.eventClient, + sessionId: coreEvent.payload.sessionId, + eventType: "runtime.chat.pending_prompts", + payload: { + prompts: coreEvent.payload.prompts, + }, + }); + return; + } if (coreEvent.type !== "team_progress") { return; } diff --git a/sdk/apps/cli/src/connectors/runtime-turn.ts b/sdk/apps/cli/src/connectors/runtime-turn.ts index 02253a2de0..67db7b9928 100644 --- a/sdk/apps/cli/src/connectors/runtime-turn.ts +++ b/sdk/apps/cli/src/connectors/runtime-turn.ts @@ -252,6 +252,9 @@ export function createConnectorRuntimeTurnStream(input: { const runTurn = input.client .sendRuntimeSession(input.sessionId, input.request) .then(async (response) => { + if (!response.result) { + throw new Error("connector runtime turn unexpectedly queued"); + } const finalText = response.result.text ?? ""; await input.onCompleted?.({ text: finalText, diff --git a/sdk/apps/cli/src/runtime/run-interactive.ts b/sdk/apps/cli/src/runtime/run-interactive.ts index 3410f1ba14..ad3f691432 100644 --- a/sdk/apps/cli/src/runtime/run-interactive.ts +++ b/sdk/apps/cli/src/runtime/run-interactive.ts @@ -30,7 +30,11 @@ import { resolveClineWelcomeLine, } from "./interactive-welcome"; import { buildUserInputMessage } from "./prompt"; -import { getUIEventEmitter, subscribeToAgentEvents } from "./session-events"; +import { + getUIEventEmitter, + subscribeToAgentEvents, + subscribeToPendingPromptEvents, +} from "./session-events"; export async function runInteractive( config: Config, @@ -91,7 +95,18 @@ export async function runInteractive( const onAgentEvent = (event: AgentEvent) => { uiEvents.emit("agent", event); }; - const unsubscribe = subscribeToAgentEvents(sessionManager, onAgentEvent); + const unsubscribeAgent = subscribeToAgentEvents(sessionManager, onAgentEvent); + const unsubscribePendingPrompts = subscribeToPendingPromptEvents( + sessionManager, + { + onPendingPrompts: (event) => { + uiEvents.emit("pending-prompts", event); + }, + onPendingPromptSubmitted: (event) => { + uiEvents.emit("pending-prompt-submitted", event); + }, + }, + ); const initialMessages = await loadInteractiveResumeMessages( sessionManager, @@ -303,7 +318,8 @@ export async function runInteractive( requestExit(); process.off("SIGINT", handleSigint); process.off("SIGTERM", handleSigterm); - unsubscribe(); + unsubscribeAgent(); + unsubscribePendingPrompts(); try { await sessionManager.stop(activeSessionId); } finally { @@ -334,17 +350,28 @@ export async function runInteractive( cwd: config.cwd, workspaceRoot: config.workspaceRoot?.trim() || config.cwd, }), - subscribeToEvents: ({ onAgentEvent: onAgent, onTeamEvent: onTeam }) => { + subscribeToEvents: ({ + onAgentEvent: onAgent, + onTeamEvent: onTeam, + onPendingPrompts, + onPendingPromptSubmitted, + }) => { uiEvents.on("agent", onAgent); uiEvents.on("team", onTeam); + uiEvents.on("pending-prompts", onPendingPrompts); + uiEvents.on("pending-prompt-submitted", onPendingPromptSubmitted); return () => { uiEvents.off("agent", onAgent); uiEvents.off("team", onTeam); + uiEvents.off("pending-prompts", onPendingPrompts); + uiEvents.off("pending-prompt-submitted", onPendingPromptSubmitted); }; }, - onSubmit: async (input, _mode) => { + onSubmit: async (input, _mode, delivery) => { abortRequested = false; - isRunning = true; + if (!delivery) { + isRunning = true; + } try { let commandOutput: string | undefined; if ( @@ -399,9 +426,17 @@ export async function runInteractive( const result = await sessionManager.send({ sessionId: activeSessionId, prompt: userInput, + delivery, }); if (!result) { - throw new Error("session manager did not return a result"); + return { + usage: { + inputTokens: 0, + outputTokens: 0, + }, + iterations: 0, + queued: delivery === "queue" || delivery === "steer", + }; } const usage = (await sessionManager.getAccumulatedUsage(activeSessionId)) ?? @@ -411,7 +446,9 @@ export async function runInteractive( iterations: result.iterations, }; } finally { - isRunning = false; + if (!delivery) { + isRunning = false; + } } }, onAbort: () => { diff --git a/sdk/apps/cli/src/runtime/session-events.ts b/sdk/apps/cli/src/runtime/session-events.ts index 9bbba4faaf..6474e85d66 100644 --- a/sdk/apps/cli/src/runtime/session-events.ts +++ b/sdk/apps/cli/src/runtime/session-events.ts @@ -1,16 +1,56 @@ import { EventEmitter } from "node:events"; import type { AgentEvent, TeamEvent } from "@clinebot/agents"; +import type { CoreSessionEvent } from "@clinebot/core"; export const getUIEventEmitter = () => new EventEmitter() as InteractiveEventBridge; +export interface PendingPromptSnapshot { + sessionId: string; + prompts: Array<{ + id: string; + prompt: string; + delivery: "queue" | "steer"; + attachmentCount: number; + }>; +} + +export interface PendingPromptSubmittedEvent { + sessionId: string; + id: string; + prompt: string; + delivery: "queue" | "steer"; + attachmentCount: number; +} + interface InteractiveEventBridge { on(event: "agent", listener: (event: AgentEvent) => void): this; on(event: "team", listener: (event: TeamEvent) => void): this; + on( + event: "pending-prompts", + listener: (event: PendingPromptSnapshot) => void, + ): this; + on( + event: "pending-prompt-submitted", + listener: (event: PendingPromptSubmittedEvent) => void, + ): this; off(event: "agent", listener: (event: AgentEvent) => void): this; off(event: "team", listener: (event: TeamEvent) => void): this; + off( + event: "pending-prompts", + listener: (event: PendingPromptSnapshot) => void, + ): this; + off( + event: "pending-prompt-submitted", + listener: (event: PendingPromptSubmittedEvent) => void, + ): this; emit(event: "agent", payload: AgentEvent): boolean; emit(event: "team", payload: TeamEvent): boolean; + emit(event: "pending-prompts", payload: PendingPromptSnapshot): boolean; + emit( + event: "pending-prompt-submitted", + payload: PendingPromptSubmittedEvent, + ): boolean; } type SessionManagerSubscriber = { @@ -60,3 +100,22 @@ export function subscribeToAgentEvents( } }); } + +export function subscribeToPendingPromptEvents( + sessionManager: SessionManagerSubscriber, + handlers: { + onPendingPrompts: (event: PendingPromptSnapshot) => void; + onPendingPromptSubmitted: (event: PendingPromptSubmittedEvent) => void; + }, +): () => void { + return sessionManager.subscribe((event: unknown) => { + const typedEvent = event as CoreSessionEvent; + if (typedEvent.type === "pending_prompts") { + handlers.onPendingPrompts(typedEvent.payload); + return; + } + if (typedEvent.type === "pending_prompt_submitted") { + handlers.onPendingPromptSubmitted(typedEvent.payload); + } + }); +} diff --git a/sdk/apps/cli/src/tui/interactive-tui.ts b/sdk/apps/cli/src/tui/interactive-tui.ts index 8e5dbbb708..4e1a7280ed 100644 --- a/sdk/apps/cli/src/tui/interactive-tui.ts +++ b/sdk/apps/cli/src/tui/interactive-tui.ts @@ -16,6 +16,10 @@ import { type InteractiveSlashCommand, searchWorkspaceFilesForMention, } from "../runtime/interactive-welcome"; +import type { + PendingPromptSnapshot, + PendingPromptSubmittedEvent, +} from "../runtime/session-events"; import { formatToolInput, formatToolOutput, truncate } from "../utils/helpers"; import { c, formatUsd } from "../utils/output"; import { type RepoStatus, readRepoStatus } from "../utils/repo-status"; @@ -30,6 +34,13 @@ interface InteractiveTurnResult { }; iterations: number; commandOutput?: string; + queued?: boolean; +} + +interface QueuedPromptItem { + id: string; + prompt: string; + steer: boolean; } interface InteractiveTuiProps { @@ -42,10 +53,13 @@ interface InteractiveTuiProps { subscribeToEvents: (handlers: { onAgentEvent: (event: AgentEvent) => void; onTeamEvent: (event: TeamEvent) => void; + onPendingPrompts: (event: PendingPromptSnapshot) => void; + onPendingPromptSubmitted: (event: PendingPromptSubmittedEvent) => void; }) => () => void; onSubmit: ( input: string, mode: "act" | "plan", + delivery?: "queue" | "steer", ) => Promise; onAbort: () => void; onExit: () => void; @@ -262,6 +276,7 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { const [autoApproveAll, setAutoApproveAll] = useState( config.toolPolicies["*"]?.autoApprove !== false, ); + const [queuedPrompts, setQueuedPrompts] = useState([]); // File mention completion const [fileMentionResults, setFileMentionResults] = useState([]); @@ -307,6 +322,7 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { const mentionSearchCounterRef = useRef(0); const turnErrorReportedRef = useRef(false); const configLoadCounterRef = useRef(0); + const knownPendingPromptIdsRef = useRef(new Set()); const workspaceName = useMemo( () => basename(config.cwd) || config.cwd, @@ -528,6 +544,7 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { (event: AgentEvent) => { switch (event.type) { case "iteration_start": + setIsRunning(true); closeInlineStream(); break; case "iteration_end": @@ -593,9 +610,11 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { break; } case "done": + setIsRunning(false); closeInlineStream(); break; case "error": + setIsRunning(false); closeInlineStream(); turnErrorReportedRef.current = true; onTurnErrorReported(true); @@ -700,13 +719,44 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { [appendLine], ); + const handlePendingPrompts = useCallback((event: PendingPromptSnapshot) => { + const nextIds = new Set(); + for (const entry of event.prompts) { + nextIds.add(entry.id); + } + knownPendingPromptIdsRef.current = nextIds; + setQueuedPrompts( + event.prompts.map((entry, index) => ({ + id: entry.id || `${entry.delivery}:${index}:${entry.prompt}`, + prompt: entry.prompt, + steer: entry.delivery === "steer", + })), + ); + }, []); + + const handlePendingPromptSubmitted = useCallback( + (event: PendingPromptSubmittedEvent) => { + knownPendingPromptIdsRef.current.delete(event.id); + appendLine(`${c.green}>${c.reset} ${event.prompt}`); + }, + [appendLine], + ); + useEffect( () => subscribeToEvents({ onAgentEvent: handleAgentEvent, onTeamEvent: handleTeamEvent, + onPendingPrompts: handlePendingPrompts, + onPendingPromptSubmitted: handlePendingPromptSubmitted, }), - [handleAgentEvent, handleTeamEvent, subscribeToEvents], + [ + handleAgentEvent, + handlePendingPrompts, + handlePendingPromptSubmitted, + handleTeamEvent, + subscribeToEvents, + ], ); useEffect(() => { @@ -714,21 +764,32 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { }, [isRunning, onRunningChange]); const submitPrompt = useCallback( - async (prompt: string) => { + async (prompt: string, delivery?: "queue" | "steer") => { setHasSubmitted(true); - setIsRunning(true); - setAbortRequested(false); - turnErrorReportedRef.current = false; - onTurnErrorReported(false); - appendLine(`${c.green}>${c.reset} ${prompt}`); + if (!delivery) { + setIsRunning(true); + setAbortRequested(false); + turnErrorReportedRef.current = false; + onTurnErrorReported(false); + } + const prefix = + delivery === "steer" + ? `${c.yellow}[steer]${c.reset} ` + : delivery === "queue" + ? `${c.dim}[queued]${c.reset} ` + : ""; + appendLine(`${c.green}>${c.reset} ${prefix}${prompt}`); setInput(""); const startedAt = performance.now(); try { - const result = await onSubmit(prompt, uiMode); + const result = await onSubmit(prompt, uiMode, delivery); if (result.commandOutput) { appendLine(result.commandOutput); } + if (result.queued) { + return; + } const tokens = result.usage.inputTokens + result.usage.outputTokens; setLastTotalTokens(tokens); if (typeof result.usage.totalCost === "number") { @@ -758,8 +819,10 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { ); } } finally { - setIsRunning(false); - refreshRepoStatus(); + if (!delivery) { + setIsRunning(false); + refreshRepoStatus(); + } } }, [ @@ -896,6 +959,15 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { } } + if (key.ctrl && value === "s") { + const trimmed = input.trim(); + if (!trimmed || !isRunning) { + return; + } + void submitPrompt(trimmed, "steer"); + return; + } + if (key.return) { if (hasMentionMenu) { const selectedPath = mentionResults[fileMentionSelectedIndex]; @@ -927,15 +999,15 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { return; } const trimmed = input.trim(); - if (!trimmed || isRunning) { + if (!trimmed) { return; } - if (isConfigCommand(trimmed)) { + if (!isRunning && isConfigCommand(trimmed)) { setInput(""); openConfigView(); return; } - void submitPrompt(trimmed); + void submitPrompt(trimmed, isRunning ? "queue" : undefined); return; } @@ -1054,8 +1126,46 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { const renderInputBox = !isConfigViewOpen ? React.createElement( Box, - { borderStyle: "round", paddingX: 1 }, - React.createElement(Text, null, `${c.green}>${c.reset} ${input}`), + { flexDirection: "column" }, + queuedPrompts.length > 0 + ? React.createElement( + Box, + { + borderStyle: "round", + paddingX: 1, + paddingY: 0, + marginBottom: 1, + flexDirection: "column", + }, + React.createElement( + Text, + { color: "gray" }, + "Queued for upcoming turns", + ), + React.createElement( + Text, + { color: "gray" }, + "Enter queues while running. Ctrl+S steers the next turn.", + ), + ...queuedPrompts.map((item, index) => + React.createElement( + Text, + { + key: item.id, + color: item.steer ? "yellow" : undefined, + }, + item.steer + ? `Steer: ${truncate(item.prompt, 100)}` + : `Queue ${index + 1}: ${truncate(item.prompt, 100)}`, + ), + ), + ) + : null, + React.createElement( + Box, + { borderStyle: "round", paddingX: 1 }, + React.createElement(Text, null, `${c.green}>${c.reset} ${input}`), + ), ) : null; @@ -1260,6 +1370,16 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { "Auto-approve all disabled (Shift+Tab)", ); + const renderQueueHint = !isConfigViewOpen + ? React.createElement( + Text, + { color: "gray" }, + isRunning + ? "Enter queues while running · Ctrl+S steers the next turn" + : "Enter submits · / for commands · @ for files", + ) + : null; + const renderStatusBar = React.createElement( Box, { flexDirection: "column", marginTop: 1 }, @@ -1271,12 +1391,13 @@ export function InteractiveTui(props: InteractiveTuiProps): React.ReactElement { { color: "gray" }, isConfigViewOpen ? "Config mode: \u2190/\u2192 tabs \u00b7 \u2191/\u2193 navigate \u00b7 Esc close" - : "/ for commands · @ for files", + : undefined, ), isConfigViewOpen ? React.createElement(Text, { color: "gray" }, "(Esc)") : renderModeSelector, ), + renderQueueHint, React.createElement( Box, null, diff --git a/sdk/apps/cli/src/utils/session.ts b/sdk/apps/cli/src/utils/session.ts index fc74a6cceb..4137385703 100644 --- a/sdk/apps/cli/src/utils/session.ts +++ b/sdk/apps/cli/src/utils/session.ts @@ -82,6 +82,7 @@ export interface CliSessionManager { prompt: string; userImages?: string[]; userFiles?: string[]; + delivery?: "queue" | "steer"; }): Promise; getAccumulatedUsage( sessionId: string, diff --git a/sdk/apps/code/hooks/chat-session/types.ts b/sdk/apps/code/hooks/chat-session/types.ts index 4d20f36eb0..7b2c012c89 100644 --- a/sdk/apps/code/hooks/chat-session/types.ts +++ b/sdk/apps/code/hooks/chat-session/types.ts @@ -104,4 +104,5 @@ export type PromptInQueue = { id: string; prompt: string; steer: boolean; + attachmentCount?: number; }; diff --git a/sdk/apps/code/hooks/use-chat-session.ts b/sdk/apps/code/hooks/use-chat-session.ts index 42a8f2f2b5..88910685b1 100644 --- a/sdk/apps/code/hooks/use-chat-session.ts +++ b/sdk/apps/code/hooks/use-chat-session.ts @@ -452,6 +452,10 @@ export function useChatSession() { ? `${prompt}${prompt.length > 0 ? "\n\n" : ""}[attached ${attachmentCount} file${attachmentCount === 1 ? "" : "s"}]` : prompt; if (userLabel) { + activeAssistantMessageIdRef.current = null; + setActiveAssistantMessageId(null); + clearLiveToolRefs(); + setStatus("running"); addMessage({ id: makeId("user"), sessionId: listeningSessionId, @@ -562,16 +566,10 @@ export function useChatSession() { ); return; } - addMessage({ - id: makeId("tool"), - sessionId: listeningSessionId, - role: "tool", - content: toolPayload, - createdAt: Date.now(), - meta: { toolName, hookEventName: "tool_call_end" }, - }); + // Ignore unmatched tool end events so stale completions from the + // previous turn do not get rendered under a newer streaming turn. }, - [addMessage, appendMessageContent], + [addMessage, appendMessageContent, clearLiveToolRefs], ); // ---- Transport / event subscriptions ---- @@ -717,6 +715,7 @@ export function useChatSession() { action: "send", sessionId: activeSessionId, prompt: trimmed, + delivery: shouldQueue ? "queue" : undefined, config: parsed, attachments: hasAttachments ? serializedAttachments : undefined, }); @@ -823,6 +822,9 @@ export function useChatSession() { ); } + const hasQueuedFollowUps = + Array.isArray(payload.promptsInQueue) && + payload.promptsInQueue.length > 0; if (result?.finishReason === "error") { if (!resolvedAssistantText) { const toolError = Array.isArray(result?.toolCalls) @@ -839,6 +841,8 @@ export function useChatSession() { setStatus("failed"); } else if (result?.finishReason === "aborted") { setStatus("cancelled"); + } else if (hasQueuedFollowUps) { + setStatus("running"); } else { setStatus("completed"); } diff --git a/sdk/apps/code/host/runtime-bridge.ts b/sdk/apps/code/host/runtime-bridge.ts index ff0814208b..c86fc2c5f4 100644 --- a/sdk/apps/code/host/runtime-bridge.ts +++ b/sdk/apps/code/host/runtime-bridge.ts @@ -35,7 +35,6 @@ import { type JsonRecord, type LiveSession, type PromptInQueue, - type QueuedChatTurn, type ToolApprovalRequestItem, } from "./types"; @@ -193,6 +192,53 @@ function handleBridgeStdoutLine(ctx: HostContext, parsed: JsonRecord) { ); return; + case "pending_prompts": { + const prompts = Array.isArray(parsed.prompts) + ? ( + parsed.prompts as Array<{ + id?: unknown; + prompt?: unknown; + delivery?: unknown; + attachmentCount?: unknown; + }> + ) + .map((item) => ({ + id: typeof item.id === "string" ? item.id : "", + prompt: typeof item.prompt === "string" ? item.prompt : "", + steer: item.delivery === "steer", + attachmentCount: + typeof item.attachmentCount === "number" + ? item.attachmentCount + : 0, + })) + .filter((item) => item.id && item.prompt) + : []; + if (sessionId) { + const session = ctx.liveSessions.get(sessionId); + const previous = session?.promptsInQueue ?? []; + if (session) { + session.promptsInQueue = prompts; + } + if ( + previous.length > prompts.length && + previous[0] && + previous[0].id !== prompts[0]?.id + ) { + emitChunk( + ctx, + sessionId, + "chat_queued_prompt_start", + JSON.stringify({ + prompt: previous[0].prompt, + attachmentCount: previous[0].attachmentCount ?? 0, + }), + ); + } + sendPromptsInQueueSnapshot(ctx, sessionId); + } + return; + } + case "error": { const message = typeof parsed.message === "string" @@ -338,16 +384,8 @@ export function broadcastApprovalSnapshots(ctx: HostContext) { } } -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, - })); + return session.promptsInQueue; } function sendPromptsInQueueSnapshot(ctx: HostContext, sessionId: string) { @@ -408,10 +446,18 @@ async function executeChatTurn( ctx: HostContext, sessionId: string, session: LiveSession, - turn: QueuedChatTurn, -): Promise { - if (turn.config) session.config = turn.config; - if (turn.prompt) session.prompt = turn.prompt; + input: { + prompt: string; + config?: JsonRecord; + attachments?: { + userImages?: string[]; + userFiles?: Array<{ name: string; content: string }>; + }; + delivery?: "queue" | "steer"; + }, +): Promise<{ result?: ChatTurnResult; queued?: boolean }> { + if (input.config) session.config = input.config; + if (input.prompt) session.prompt = input.prompt; session.busy = true; session.status = "running"; @@ -432,14 +478,19 @@ async function executeChatTurn( request: { config: session.config, messages: session.messages, - prompt: turn.prompt, - attachments: turn.attachments, + prompt: input.prompt, + attachments: input.attachments, + delivery: input.delivery, }, - })) as { result?: ChatTurnResult }; + })) as { result?: ChatTurnResult; queued?: boolean }; + if (resultEnvelope.queued) { + return { queued: true }; + } const result = resultEnvelope.result; - if (!result) + if (!result) { throw new Error("chat runtime bridge send response missing result"); + } session.messages = persistUsageInMessages( (Array.isArray(result.messages) ? result.messages : []) as unknown[], @@ -452,55 +503,12 @@ async function executeChatTurn( persistSessionMessages(sessionId, session.messages); sendApprovalSnapshot(ctx, sessionId); - return result; + return { result, queued: false }; } finally { session.busy = false; } } -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 // --------------------------------------------------------------------------- @@ -512,7 +520,7 @@ function createLiveSession( return { config, messages: [], - pendingTurns: [], + promptsInQueue: [], busy: false, startedAt: nowMs(), status: "idle", @@ -577,36 +585,23 @@ async function handleSend( } if (request.config) session.config = request.config; - - const turn: QueuedChatTurn = { - id: makeQueuedTurnId(), + const delivery = + request.delivery === "queue" || request.delivery === "steer" + ? request.delivery + : session.busy + ? "queue" + : undefined; + const { result, queued } = await executeChatTurn(ctx, sessionId, session, { 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); - } + delivery, + }); return { sessionId, result, - queued: false, + queued: queued === true, promptsInQueue: getPromptsInQueue(session), }; } @@ -621,7 +616,7 @@ async function handleAbort( const session = ctx.liveSessions.get(sessionId); if (session) { session.busy = false; - session.pendingTurns = []; + session.promptsInQueue = []; session.status = "cancelled"; session.endedAt = nowMs(); } @@ -669,14 +664,16 @@ async function handleSteerPrompt( if (!session) { return { sessionId, promptsInQueue: [] }; } - const existingIndex = session.pendingTurns.findIndex( + const prompt = session.promptsInQueue.find( (turn) => turn.id === promptId, - ); - if (existingIndex >= 0) { - const [turn] = session.pendingTurns.splice(existingIndex, 1); - session.pendingTurns.unshift({ ...turn, steer: true }); + )?.prompt; + if (prompt) { + await executeChatTurn(ctx, sessionId, session, { + prompt, + config: session.config, + delivery: "steer", + }); } - sendPromptsInQueueSnapshot(ctx, sessionId); return { sessionId, promptsInQueue: getPromptsInQueue(session), diff --git a/sdk/apps/code/host/types.ts b/sdk/apps/code/host/types.ts index 8616e6b71a..7b6ff20283 100644 --- a/sdk/apps/code/host/types.ts +++ b/sdk/apps/code/host/types.ts @@ -38,14 +38,7 @@ export type ChatSessionCommandRequest = { sessionId?: string; prompt?: string; promptId?: string; - config?: JsonRecord; - attachments?: ChatTurnAttachments; -}; - -export type QueuedChatTurn = { - id: string; - prompt: string; - steer: boolean; + delivery?: "queue" | "steer"; config?: JsonRecord; attachments?: ChatTurnAttachments; }; @@ -54,12 +47,13 @@ export type PromptInQueue = { id: string; prompt: string; steer: boolean; + attachmentCount?: number; }; export type LiveSession = { config: JsonRecord; messages: unknown[]; - pendingTurns: QueuedChatTurn[]; + promptsInQueue: PromptInQueue[]; busy: boolean; startedAt: number; endedAt?: number; diff --git a/sdk/apps/code/scripts/chat-runtime-bridge.ts b/sdk/apps/code/scripts/chat-runtime-bridge.ts index 742eb719e6..a814dc3863 100644 --- a/sdk/apps/code/scripts/chat-runtime-bridge.ts +++ b/sdk/apps/code/scripts/chat-runtime-bridge.ts @@ -1,9 +1,5 @@ import { homedir } from "node:os"; -import { - type RpcChatTurnResult, - setHomeDir, - setHomeDirIfUnset, -} from "@clinebot/core"; +import { setHomeDir, setHomeDirIfUnset } from "@clinebot/core"; import { type RpcRuntimeBridgeCommandOutputLine, runRpcRuntimeCommandBridge, @@ -82,7 +78,7 @@ async function main() { setRuntimeHomeDir(config); addRuntimeLoggerContext(config); }, - parseSendResult: (resultRaw) => resultRaw as RpcChatTurnResult, + parseSendResult: (resultRaw) => resultRaw, }); } diff --git a/sdk/apps/code/src-tauri/src/main.rs b/sdk/apps/code/src-tauri/src/main.rs index b2fc70098c..ffee9e3745 100644 --- a/sdk/apps/code/src-tauri/src/main.rs +++ b/sdk/apps/code/src-tauri/src/main.rs @@ -112,6 +112,7 @@ struct ChatRunTurnRequest { prompt: String, #[serde(default)] attachments: Option, + delivery: Option, } fn default_agent_mode() -> String { @@ -225,6 +226,7 @@ struct ChatRuntimeBridgeLine { error: Option, duration_ms: Option, message: Option, + prompts: Option>, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -249,6 +251,8 @@ struct ChatSessionCommandRequest { action: String, session_id: Option, prompt: Option, + prompt_id: Option, + delivery: Option, config: Option, #[serde(default)] attachments: Option, @@ -261,6 +265,26 @@ struct ChatSessionCommandResponse { result: Option, ok: Option, queued: Option, + #[serde(default)] + prompts_in_queue: Vec, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +struct PendingPromptSnapshot { + id: String, + prompt: String, + delivery: String, + attachment_count: Option, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +struct PromptInQueue { + id: String, + prompt: String, + steer: bool, + attachment_count: Option, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -350,7 +374,7 @@ impl ChatWsBridgeState { struct ChatRuntimeSession { config: StartSessionRequest, messages: Vec, - pending_turns: Vec, + prompts_in_queue: Vec, busy: bool, started_at: u64, ended_at: Option, @@ -2684,6 +2708,7 @@ fn ensure_chat_runtime_bridge_started( std::sync::mpsc::Sender>, >::new())); let stdout_app = app.clone(); + let stdout_state = state.clone(); let stdout_pending = pending.clone(); thread::spawn(move || { let mut reader = BufReader::new(stdout); @@ -2765,6 +2790,40 @@ fn ensure_chat_runtime_bridge_started( payload.to_string(), ); } + "pending_prompts" => { + let Some(session_id) = parsed.session_id.as_deref() else { + continue; + }; + let prompts = map_pending_prompts(parsed.prompts.unwrap_or_default()); + let previous = { + let mut sessions = match stdout_state.sessions.lock() { + Ok(value) => value, + Err(_) => continue, + }; + let Some(session) = sessions.get_mut(session_id) else { + continue; + }; + let previous = session.prompts_in_queue.clone(); + session.prompts_in_queue = prompts.clone(); + previous + }; + if previous.len() > prompts.len() + && !previous.is_empty() + && previous[0].id != prompts.first().map(|item| item.id.as_str()).unwrap_or("") + { + let payload = serde_json::json!({ + "prompt": previous[0].prompt, + "attachmentCount": previous[0].attachment_count.unwrap_or(0), + }); + emit_chunk( + &stdout_app, + session_id, + "chat_queued_prompt_start", + payload.to_string(), + ); + } + send_prompts_in_queue_snapshot(&stdout_app, &stdout_state, session_id); + } "error" => { if let Some(session_id) = parsed.session_id.as_deref() { let payload = serde_json::json!({ @@ -2990,7 +3049,7 @@ fn run_chat_turn_via_rpc_runtime( context: &AppContext, session_id: &str, request: &ChatRunTurnRequest, -) -> Result { +) -> Result { let response = run_chat_runtime_bridge_command( app, state, @@ -3001,12 +3060,20 @@ fn run_chat_turn_via_rpc_runtime( "request": request, }), )?; - let result_value = response - .get("result") - .cloned() - .ok_or_else(|| "chat runtime bridge send response missing result".to_string())?; - serde_json::from_value::(result_value) - .map_err(|e| format!("invalid chat runtime bridge result: {e}")) + serde_json::from_value::(response) + .map_err(|e| format!("invalid chat runtime bridge response: {e}")) +} + +fn map_pending_prompts(prompts: Vec) -> Vec { + prompts + .into_iter() + .map(|item| PromptInQueue { + id: item.id, + prompt: item.prompt, + steer: item.delivery == "steer", + attachment_count: item.attachment_count, + }) + .collect() } fn persist_chat_turn_result( @@ -3046,119 +3113,32 @@ fn mark_chat_turn_failed(state: &Arc, session_id: &str) -> Res .lock() .map_err(|_| "failed to lock chat session store")?; if let Some(session) = sessions.get_mut(session_id) { + session.busy = false; session.status = "failed".to_string(); session.ended_at = Some(now_ms()); } Ok(()) } -fn dequeue_next_chat_turn( - state: &Arc, - session_id: &str, -) -> Result, String> { - let mut sessions = state - .sessions - .lock() - .map_err(|_| "failed to lock chat session store")?; - let Some(session) = sessions.get_mut(session_id) else { - return Ok(None); - }; - if session.pending_turns.is_empty() { - session.busy = false; - return Ok(None); - } - let mut next = session.pending_turns.remove(0); - session.busy = true; - session.status = "running".to_string(); - session.ended_at = None; - if !next.prompt.trim().is_empty() { - session.prompt = Some(next.prompt.clone()); - } - next.config = session.config.clone(); - next.messages = session.messages.clone(); - Ok(Some((session.config.clone(), next))) -} - -fn queue_chat_turn( - state: &Arc, - session_id: &str, - turn_request: ChatRunTurnRequest, -) -> Result<(), String> { - let mut sessions = state - .sessions - .lock() - .map_err(|_| "failed to lock chat session store")?; - let session = sessions - .get_mut(session_id) - .ok_or_else(|| "session not found. start a new session.".to_string())?; - session.pending_turns.push(turn_request); - Ok(()) -} - -fn spawn_next_queued_chat_turn( +fn send_prompts_in_queue_snapshot( app: &AppHandle, state: &Arc, - context: &AppContext, - session_id: String, + session_id: &str, ) { - let next = match dequeue_next_chat_turn(state, &session_id) { - Ok(value) => value, - Err(error) => { - eprintln!("[chat-queue] failed to dequeue next turn for {session_id}: {error}"); - return; - } - }; - let Some((config, turn_request)) = next else { - return; - }; - - let app_for_turn = app.clone(); - let state_for_turn = state.clone(); - let context_for_turn = context.clone(); - tauri::async_runtime::spawn(async move { - let turn_result = tauri::async_runtime::spawn_blocking({ - let app_for_run = app_for_turn.clone(); - let state_for_run = state_for_turn.clone(); - let context_for_run = context_for_turn.clone(); - let session_id_for_run = session_id.clone(); - let request_for_run = turn_request.clone(); - move || { - run_chat_turn_via_rpc_runtime( - &app_for_run, - &state_for_run, - &context_for_run, - &session_id_for_run, - &request_for_run, - ) - } - }) - .await - .map_err(|e| format!("chat turn task failed: {e}")); - - match turn_result { - Ok(Ok(result)) => { - if let Err(error) = - persist_chat_turn_result(&state_for_turn, &session_id, &config, &result) - { - eprintln!( - "[chat-queue] failed to persist turn result for {session_id}: {error}" - ); - } - } - _ => { - if let Err(error) = mark_chat_turn_failed(&state_for_turn, &session_id) { - eprintln!("[chat-queue] failed to mark turn failed for {session_id}: {error}"); - } - } - } - - spawn_next_queued_chat_turn( - &app_for_turn, - &state_for_turn, - &context_for_turn, - session_id, - ); - }); + let items = state + .sessions + .lock() + .ok() + .and_then(|sessions| sessions.get(session_id).cloned()) + .map(|session| session.prompts_in_queue) + .unwrap_or_default(); + let payload = serde_json::json!({ "items": items }); + emit_chunk( + app, + session_id, + "prompts_in_queue_state", + payload.to_string(), + ); } fn abort_chat_session_via_rpc_runtime( @@ -4265,7 +4245,7 @@ async fn handle_chat_session_command( ChatRuntimeSession { config, messages: Vec::new(), - pending_turns: Vec::new(), + prompts_in_queue: Vec::new(), busy: false, started_at: now_ms(), ended_at: None, @@ -4279,6 +4259,7 @@ async fn handle_chat_session_command( result: None, ok: None, queued: None, + prompts_in_queue: Vec::new(), }) } "send" => { @@ -4322,7 +4303,7 @@ async fn handle_chat_session_command( prompt: derive_prompt_from_messages(&messages), title: read_session_metadata_title(&session_id), messages, - pending_turns: Vec::new(), + prompts_in_queue: Vec::new(), busy: false, started_at: now_ms(), ended_at: None, @@ -4330,7 +4311,7 @@ async fn handle_chat_session_command( }); } - let (config, messages) = { + let (config, messages, delivery) = { let mut sessions = state .sessions .lock() @@ -4343,30 +4324,23 @@ async fn handle_chat_session_command( ensure_start_session_home_dir(&mut next_config); session.config = next_config; } - if session.busy { - let queued_request = ChatRunTurnRequest { - config: session.config.clone(), - messages: Vec::new(), - prompt: prompt.clone(), - attachments: attachments.clone(), - }; - drop(sessions); - queue_chat_turn(state, &session_id, queued_request)?; - return Ok(ChatSessionCommandResponse { - session_id: Some(session_id), - result: None, - ok: Some(true), - queued: Some(true), - }); - } + let delivery = match request.delivery.as_deref() { + Some("queue") => Some("queue".to_string()), + Some("steer") => Some("steer".to_string()), + _ if session.busy => Some("queue".to_string()), + _ => None, + }; session.busy = true; session.status = "running".to_string(); session.ended_at = None; if !prompt.is_empty() { - // Keep sidebar/session discovery title in sync while turn is in-flight. session.prompt = Some(prompt.clone()); } - (session.config.clone(), session.messages.clone()) + ( + session.config.clone(), + session.messages.clone(), + delivery, + ) }; let session_id_for_turn = session_id.clone(); @@ -4379,6 +4353,7 @@ async fn handle_chat_session_command( messages, prompt: prompt.clone(), attachments, + delivery, }; let turn_result = tauri::async_runtime::spawn_blocking(move || { @@ -4393,20 +4368,62 @@ async fn handle_chat_session_command( .await .map_err(|e| format!("chat turn task failed: {e}")); - if let Ok(Ok(result)) = &turn_result { + if let Ok(Ok(response)) = &turn_result { + if response.queued == Some(true) { + let prompts_in_queue = { + let sessions = state + .sessions + .lock() + .map_err(|_| "failed to lock chat session store")?; + sessions + .get(&session_id) + .cloned() + .map(|session| session.prompts_in_queue) + .unwrap_or_default() + }; + return Ok(ChatSessionCommandResponse { + session_id: Some(session_id), + result: None, + ok: Some(true), + queued: Some(true), + prompts_in_queue, + }); + } + let result = response + .result + .as_ref() + .ok_or_else(|| "chat runtime bridge send response missing result".to_string())?; persist_chat_turn_result(state, &session_id, &config, result)?; + let prompts_in_queue = { + let mut sessions = state + .sessions + .lock() + .map_err(|_| "failed to lock chat session store")?; + let session = sessions + .get_mut(&session_id) + .ok_or_else(|| "session not found. start a new session.".to_string())?; + session.busy = !session.prompts_in_queue.is_empty(); + session.prompts_in_queue.clone() + }; + return Ok(ChatSessionCommandResponse { + session_id: Some(session_id), + result: Some(result.clone()), + ok: None, + queued: Some(false), + prompts_in_queue, + }); } else { mark_chat_turn_failed(state, &session_id)?; } - spawn_next_queued_chat_turn(app, state, context, session_id.clone()); let turn_result = turn_result?; - let result = turn_result?; + let response = turn_result?; Ok(ChatSessionCommandResponse { session_id: Some(session_id), - result: Some(result), - ok: None, - queued: Some(false), + result: response.result, + ok: response.ok, + queued: response.queued, + prompts_in_queue: Vec::new(), }) } "abort" => Ok(ChatSessionCommandResponse { @@ -4419,16 +4436,19 @@ async fn handle_chat_session_command( .map_err(|_| "failed to lock chat session store")?; if let Some(session) = sessions.get_mut(&session_id) { session.busy = false; - session.pending_turns.clear(); + session.prompts_in_queue.clear(); session.status = "cancelled".to_string(); session.ended_at = Some(now_ms()); } + drop(sessions); + send_prompts_in_queue_snapshot(app, state, &session_id); } request.session_id }, result: None, ok: Some(true), queued: None, + prompts_in_queue: Vec::new(), }), "reset" => { if let Some(session_id) = request.session_id.clone() { @@ -4444,12 +4464,94 @@ async fn handle_chat_session_command( Some(session_id.as_str()), ); let _ = remove_chat_stream_subscription(app, state, context, &session_id); + send_prompts_in_queue_snapshot(app, state, &session_id); } Ok(ChatSessionCommandResponse { session_id: request.session_id, result: None, ok: Some(true), queued: None, + prompts_in_queue: Vec::new(), + }) + } + "pending_prompts" => { + let Some(session_id) = request.session_id else { + return Err("sessionId is required for pending_prompts action".to_string()); + }; + let prompts_in_queue = state + .sessions + .lock() + .ok() + .and_then(|sessions| sessions.get(&session_id).cloned()) + .map(|session| session.prompts_in_queue) + .unwrap_or_default(); + Ok(ChatSessionCommandResponse { + session_id: Some(session_id), + result: None, + ok: Some(true), + queued: None, + prompts_in_queue, + }) + } + "steer_prompt" => { + let Some(session_id) = request.session_id.clone() else { + return Err("sessionId is required for steer_prompt action".to_string()); + }; + let Some(prompt_id) = request.prompt_id.clone() else { + return Err("promptId is required for steer_prompt action".to_string()); + }; + let session = state + .sessions + .lock() + .ok() + .and_then(|sessions| sessions.get(&session_id).cloned()); + if let Some(session) = session { + let prompt = session + .prompts_in_queue + .iter() + .find(|item| item.id == prompt_id) + .map(|item| item.prompt.clone()); + if let Some(prompt) = prompt { + ensure_chat_stream_subscription(app, state, context, &session_id)?; + let _ = tauri::async_runtime::spawn_blocking({ + let app_for_turn = app.clone(); + let state_for_turn = state.clone(); + let context_for_turn = context.clone(); + let session_id_for_turn = session_id.clone(); + let turn_request = ChatRunTurnRequest { + config: session.config.clone(), + messages: session.messages.clone(), + prompt, + attachments: None, + delivery: Some("steer".to_string()), + }; + move || { + run_chat_turn_via_rpc_runtime( + &app_for_turn, + &state_for_turn, + &context_for_turn, + &session_id_for_turn, + &turn_request, + ) + } + }) + .await + .map_err(|e| format!("chat turn task failed: {e}"))??; + } + } + let prompts_in_queue = state + .sessions + .lock() + .ok() + .and_then(|sessions| sessions.get(&session_id).cloned()) + .map(|session| session.prompts_in_queue) + .unwrap_or_default(); + Ok(ChatSessionCommandResponse { + session_id: Some(session_id), + result: None, + ok: Some(true), + queued: Some(true), + prompts_in_queue, }) } _ => Err("unsupported action".to_string()), diff --git a/sdk/packages/core/src/session/default-session-manager.test.ts b/sdk/packages/core/src/session/default-session-manager.test.ts index c78312c413..fe49a6a1a7 100644 --- a/sdk/packages/core/src/session/default-session-manager.test.ts +++ b/sdk/packages/core/src/session/default-session-manager.test.ts @@ -55,6 +55,55 @@ function createManifest(sessionId: string): SessionManifest { }; } +type PluginEventTestHarness = { + handlePluginEvent: ( + rootSessionId: string, + event: { name: string; payload?: unknown }, + ) => Promise; + getPendingPrompts: ( + sessionId: string, + ) => Array<{ prompt: string; delivery: "queue" | "steer" }>; +}; + +function createPluginEventHarness( + manager: DefaultSessionManager, +): PluginEventTestHarness { + const target = manager as object; + return { + handlePluginEvent: async (rootSessionId, event) => { + const handler = Reflect.get(target, "handlePluginEvent"); + if (typeof handler !== "function") { + throw new Error("handlePluginEvent test hook unavailable"); + } + await Reflect.apply( + handler as ( + rootSessionId: string, + event: { name: string; payload?: unknown }, + ) => Promise, + target, + [rootSessionId, event], + ); + }, + getPendingPrompts: (sessionId) => { + const getter = Reflect.get(target, "getSessionOrThrow"); + if (typeof getter !== "function") { + throw new Error("getSessionOrThrow test hook unavailable"); + } + const session = Reflect.apply( + getter as (sessionId: string) => { + pendingPrompts: Array<{ + prompt: string; + delivery: "queue" | "steer"; + }>; + }, + target, + [sessionId], + ); + return session.pendingPrompts; + }, + }; +} + function createConfig( overrides: Partial = {}, ): CoreSessionConfig { @@ -381,14 +430,8 @@ describe("DefaultSessionManager", () => { interactive: true, }); - await ( - manager as DefaultSessionManager & { - handlePluginEvent: ( - rootSessionId: string, - event: { name: string; payload?: unknown }, - ) => Promise; - } - ).handlePluginEvent(sessionId, { + const harness = createPluginEventHarness(manager); + await harness.handlePluginEvent(sessionId, { name: "steer_message", payload: { prompt: "async result" }, }); @@ -451,30 +494,22 @@ describe("DefaultSessionManager", () => { interactive: true, }); - const internal = manager as DefaultSessionManager & { - handlePluginEvent: ( - rootSessionId: string, - event: { name: string; payload?: unknown }, - ) => Promise; - getSessionOrThrow: (id: string) => { - pendingPrompts: Array<{ prompt: string; delivery: "queue" | "steer" }>; - }; - }; + const harness = createPluginEventHarness(manager); - await internal.handlePluginEvent(sessionId, { + await harness.handlePluginEvent(sessionId, { name: "queue_message", payload: { prompt: "queued first" }, }); - await internal.handlePluginEvent(sessionId, { + await harness.handlePluginEvent(sessionId, { name: "queue_message", payload: { prompt: "queued second" }, }); - await internal.handlePluginEvent(sessionId, { + await harness.handlePluginEvent(sessionId, { name: "steer_message", payload: { prompt: "queued first" }, }); - expect(internal.getSessionOrThrow(sessionId).pendingPrompts).toEqual([ + expect(harness.getPendingPrompts(sessionId)).toEqual([ { prompt: "queued first", delivery: "steer" }, { prompt: "queued second", delivery: "queue" }, ]); @@ -824,6 +859,112 @@ describe("DefaultSessionManager", () => { }); }); + it("queues sends with explicit queue or steer delivery and emits snapshots", async () => { + const sessionId = "sess-delivery-queue"; + const manifest = createManifest(sessionId); + const sessionService = { + ensureSessionsDir: vi.fn().mockReturnValue("/tmp/sessions"), + createRootSessionWithArtifacts: vi.fn().mockResolvedValue({ + manifestPath: "/tmp/manifest-queue.json", + transcriptPath: "/tmp/transcript-queue.log", + hookPath: "/tmp/hook-queue.log", + messagesPath: "/tmp/messages-queue.json", + manifest, + }), + persistSessionMessages: vi.fn(), + updateSessionStatus: vi.fn().mockResolvedValue({ updated: true }), + writeSessionManifest: vi.fn(), + listSessions: vi.fn().mockResolvedValue([]), + deleteSession: vi.fn().mockResolvedValue({ deleted: true }), + }; + const runtimeBuilder = { + build: vi.fn().mockReturnValue({ + tools: [], + shutdown: vi.fn(), + }), + }; + let canStartRun = false; + const run = vi.fn().mockResolvedValue(createResult({ text: "first" })); + const continueFn = vi + .fn() + .mockResolvedValue(createResult({ text: "next" })); + const manager = new DefaultSessionManager({ + distinctId, + sessionService: sessionService as never, + runtimeBuilder, + createAgent: () => + ({ + run, + continue: continueFn, + canStartRun: vi.fn(() => canStartRun), + abort: vi.fn(), + shutdown: vi.fn().mockResolvedValue(undefined), + getMessages: vi.fn().mockReturnValue([]), + messages: [], + }) as never, + }); + const events: Array = []; + manager.subscribe((event) => { + events.push(event); + }); + + await manager.start({ + config: createConfig({ sessionId }), + interactive: true, + }); + + await expect( + manager.send({ sessionId, prompt: "queued first", delivery: "queue" }), + ).resolves.toBeUndefined(); + await expect( + manager.send({ sessionId, prompt: "queued second", delivery: "steer" }), + ).resolves.toBeUndefined(); + + expect(run).not.toHaveBeenCalled(); + expect(continueFn).not.toHaveBeenCalled(); + const promptSnapshots = events + .filter((event) => { + return ( + typeof event === "object" && + event !== null && + "type" in event && + event.type === "pending_prompts" + ); + }) + .map((event) => (event as { payload: { prompts: unknown[] } }).payload); + expect(promptSnapshots.at(-1)).toEqual({ + prompts: [ + expect.objectContaining({ + prompt: "queued second", + delivery: "steer", + attachmentCount: 0, + }), + expect.objectContaining({ + prompt: "queued first", + delivery: "queue", + attachmentCount: 0, + }), + ], + sessionId, + }); + + canStartRun = true; + await manager.send({ sessionId, prompt: "run now" }); + expect(run).toHaveBeenCalledTimes(1); + expect( + events.some((event) => { + return ( + typeof event === "object" && + event !== null && + "type" in event && + event.type === "pending_prompt_submitted" && + "payload" in event && + (event.payload as { prompt?: string }).prompt === "queued second" + ); + }), + ).toBe(true); + }); + it("returns undefined accumulated usage for unknown sessions", async () => { const manager = new DefaultSessionManager({ distinctId, diff --git a/sdk/packages/core/src/session/default-session-manager.ts b/sdk/packages/core/src/session/default-session-manager.ts index 8aa7610fbe..dc92b7984f 100644 --- a/sdk/packages/core/src/session/default-session-manager.ts +++ b/sdk/packages/core/src/session/default-session-manager.ts @@ -288,9 +288,6 @@ export class DefaultSessionManager implements SessionManager { onConsecutiveMistakeLimitReached: configWithProvider.onConsecutiveMistakeLimitReached, completionGuard: runtime.completionGuard, - toolContextMetadata: { - sessionId, - }, logger: runtime.logger ?? configWithProvider.logger, onEvent: (event: AgentEvent) => this.onAgentEvent(sessionId, configWithProvider, event), @@ -355,8 +352,18 @@ export class DefaultSessionManager implements SessionManager { promptLength: input.prompt.length, userImageCount: input.userImages?.length ?? 0, userFileCount: input.userFiles?.length ?? 0, + delivery: input.delivery ?? "immediate", }, }); + if (input.delivery === "queue" || input.delivery === "steer") { + this.enqueuePendingPrompt(input.sessionId, { + prompt: input.prompt, + delivery: input.delivery, + userImages: input.userImages, + userFiles: input.userFiles, + }); + return undefined; + } try { const result = await this.runTurn(session, { prompt: input.prompt, @@ -808,36 +815,64 @@ export class DefaultSessionManager implements SessionManager { : payload?.delivery === "steer" ? "steer" : "queue"; - this.enqueuePendingPrompt(targetSessionId, prompt, delivery); + this.enqueuePendingPrompt(targetSessionId, { + prompt, + delivery, + }); } private enqueuePendingPrompt( sessionId: string, - prompt: string, - delivery: "queue" | "steer", + entry: { + prompt: string; + delivery: "queue" | "steer"; + userImages?: string[]; + userFiles?: string[]; + }, ): void { const session = this.sessions.get(sessionId); if (!session) { return; } + const { prompt, delivery, userImages, userFiles } = entry; const existingIndex = session.pendingPrompts.findIndex( - (entry) => entry.prompt === prompt, + (queued) => queued.prompt === prompt, ); if (existingIndex >= 0) { const [existing] = session.pendingPrompts.splice(existingIndex, 1); if (delivery === "steer" || existing.delivery === "steer") { session.pendingPrompts.unshift({ + id: existing.id, prompt, delivery: "steer", + userImages: userImages ?? existing.userImages, + userFiles: userFiles ?? existing.userFiles, }); } else { - session.pendingPrompts.push(existing); + session.pendingPrompts.push({ + ...existing, + userImages: userImages ?? existing.userImages, + userFiles: userFiles ?? existing.userFiles, + }); } } else if (delivery === "steer") { - session.pendingPrompts.unshift({ prompt, delivery }); + session.pendingPrompts.unshift({ + id: `pending_${Date.now()}_${nanoid(5)}`, + prompt, + delivery, + userImages, + userFiles, + }); } else { - session.pendingPrompts.push({ prompt, delivery }); + session.pendingPrompts.push({ + id: `pending_${Date.now()}_${nanoid(5)}`, + prompt, + delivery, + userImages, + userFiles, + }); } + this.emitPendingPrompts(session); queueMicrotask(() => { void this.drainPendingPrompts(sessionId); }); @@ -864,13 +899,21 @@ export class DefaultSessionManager implements SessionManager { if (!next) { return; } + this.emitPendingPrompts(session); + this.emitPendingPromptSubmitted(session, next); session.drainingPendingPrompts = true; try { - await this.send({ sessionId, prompt: next.prompt }); + await this.send({ + sessionId, + prompt: next.prompt, + userImages: next.userImages, + userFiles: next.userFiles, + }); } catch (error) { const message = error instanceof Error ? error.message : String(error); if (message.includes("already in progress")) { session.pendingPrompts.unshift(next); + this.emitPendingPrompts(session); } else { throw error; } @@ -909,6 +952,45 @@ export class DefaultSessionManager implements SessionManager { handleAgentEvent(ctx, event); } + private emitPendingPrompts(session: ActiveSession): void { + this.emit({ + type: "pending_prompts", + payload: { + sessionId: session.sessionId, + prompts: session.pendingPrompts.map((entry) => ({ + id: entry.id, + prompt: entry.prompt, + delivery: entry.delivery, + attachmentCount: + (entry.userImages?.length ?? 0) + (entry.userFiles?.length ?? 0), + })), + }, + }); + } + + private emitPendingPromptSubmitted( + session: ActiveSession, + entry: { + id: string; + prompt: string; + delivery: "queue" | "steer"; + userImages?: string[]; + userFiles?: string[]; + }, + ): void { + this.emit({ + type: "pending_prompt_submitted", + payload: { + sessionId: session.sessionId, + id: entry.id, + prompt: entry.prompt, + delivery: entry.delivery, + attachmentCount: + (entry.userImages?.length ?? 0) + (entry.userFiles?.length ?? 0), + }, + }); + } + // ── Spawn / sub-agents ────────────────────────────────────────────── private createSpawnTool( diff --git a/sdk/packages/core/src/session/session-manager.ts b/sdk/packages/core/src/session/session-manager.ts index b1f4785588..2d81d5486f 100644 --- a/sdk/packages/core/src/session/session-manager.ts +++ b/sdk/packages/core/src/session/session-manager.ts @@ -38,6 +38,7 @@ export interface SendSessionInput { prompt: string; userImages?: string[]; userFiles?: string[]; + delivery?: "queue" | "steer"; } export interface SessionAccumulatedUsage { diff --git a/sdk/packages/core/src/session/utils/types.ts b/sdk/packages/core/src/session/utils/types.ts index 50a028b2fa..9f7335dec8 100644 --- a/sdk/packages/core/src/session/utils/types.ts +++ b/sdk/packages/core/src/session/utils/types.ts @@ -29,8 +29,11 @@ export type ActiveSession = { }; export type PendingPrompt = { + id: string; prompt: string; delivery: "queue" | "steer"; + userImages?: string[]; + userFiles?: string[]; }; export type TeamRunUpdate = { diff --git a/sdk/packages/core/src/types/events.ts b/sdk/packages/core/src/types/events.ts index 991c386163..dbef48d68e 100644 --- a/sdk/packages/core/src/types/events.ts +++ b/sdk/packages/core/src/types/events.ts @@ -36,6 +36,24 @@ export interface SessionTeamProgressEvent { summary: import("@clinebot/shared").TeamProgressSummary; } +export interface SessionPendingPromptsEvent { + sessionId: string; + prompts: Array<{ + id: string; + prompt: string; + delivery: "queue" | "steer"; + attachmentCount: number; + }>; +} + +export interface SessionPendingPromptSubmittedEvent { + sessionId: string; + id: string; + prompt: string; + delivery: "queue" | "steer"; + attachmentCount: number; +} + export type CoreSessionEvent = | { type: "chunk"; payload: SessionChunkEvent } | { @@ -46,6 +64,11 @@ export type CoreSessionEvent = }; } | { type: "team_progress"; payload: SessionTeamProgressEvent } + | { type: "pending_prompts"; payload: SessionPendingPromptsEvent } + | { + type: "pending_prompt_submitted"; + payload: SessionPendingPromptSubmittedEvent; + } | { type: "ended"; payload: SessionEndedEvent } | { type: "hook"; payload: SessionToolEvent } | { type: "status"; payload: { sessionId: string; status: string } }; diff --git a/sdk/packages/rpc/src/client.ts b/sdk/packages/rpc/src/client.ts index 062eb1ff0c..f2533f2097 100644 --- a/sdk/packages/rpc/src/client.ts +++ b/sdk/packages/rpc/src/client.ts @@ -449,7 +449,7 @@ export class RpcSessionClient { public async sendRuntimeSession( sessionId: string, request: RpcChatRunTurnRequest, - ): Promise<{ result: RpcChatTurnResult }> { + ): Promise<{ result?: RpcChatTurnResult; queued?: boolean }> { const runtimeRequest = { config: { workspaceRoot: request.config.workspaceRoot, @@ -506,6 +506,7 @@ export class RpcSessionClient { content: toProtoValue(message.content), })), prompt: request.prompt, + delivery: request.delivery, attachments: request.attachments ? { userImages: request.attachments.userImages ?? [], @@ -524,31 +525,34 @@ export class RpcSessionClient { ); }, ); + if (!response.result) { + return { queued: true }; + } return { result: { - text: response.result?.text ?? "", + text: response.result.text ?? "", usage: { - inputTokens: Number(response.result?.usage?.inputTokens ?? 0), - outputTokens: Number(response.result?.usage?.outputTokens ?? 0), - cacheReadTokens: response.result?.usage?.hasCacheReadTokens - ? Number(response.result?.usage?.cacheReadTokens ?? 0) + inputTokens: Number(response.result.usage?.inputTokens ?? 0), + outputTokens: Number(response.result.usage?.outputTokens ?? 0), + cacheReadTokens: response.result.usage?.hasCacheReadTokens + ? Number(response.result.usage?.cacheReadTokens ?? 0) : undefined, - cacheWriteTokens: response.result?.usage?.hasCacheWriteTokens - ? Number(response.result?.usage?.cacheWriteTokens ?? 0) + cacheWriteTokens: response.result.usage?.hasCacheWriteTokens + ? Number(response.result.usage?.cacheWriteTokens ?? 0) : undefined, - totalCost: response.result?.usage?.hasTotalCost - ? Number(response.result?.usage?.totalCost ?? 0) + totalCost: response.result.usage?.hasTotalCost + ? Number(response.result.usage?.totalCost ?? 0) : undefined, }, - inputTokens: Number(response.result?.inputTokens ?? 0), - outputTokens: Number(response.result?.outputTokens ?? 0), - iterations: Number(response.result?.iterations ?? 0), - finishReason: response.result?.finishReason ?? "", - messages: (response.result?.messages ?? []).map((message) => ({ + inputTokens: Number(response.result.inputTokens ?? 0), + outputTokens: Number(response.result.outputTokens ?? 0), + iterations: Number(response.result.iterations ?? 0), + finishReason: response.result.finishReason ?? "", + messages: (response.result.messages ?? []).map((message) => ({ role: message.role ?? "", content: fromProtoValue(message.content), })), - toolCalls: (response.result?.toolCalls ?? []).map((call) => ({ + toolCalls: (response.result.toolCalls ?? []).map((call) => ({ name: call.name ?? "", input: call.hasInput ? fromProtoValue(call.input) : undefined, output: call.hasOutput ? fromProtoValue(call.output) : undefined, @@ -558,6 +562,7 @@ export class RpcSessionClient { : undefined, })), }, + queued: false, }; } diff --git a/sdk/packages/rpc/src/proto/rpc.proto b/sdk/packages/rpc/src/proto/rpc.proto index e9e68fa6d9..0a62ce856e 100644 --- a/sdk/packages/rpc/src/proto/rpc.proto +++ b/sdk/packages/rpc/src/proto/rpc.proto @@ -264,6 +264,7 @@ message RuntimeTurnRequest { repeated RuntimeChatMessage messages = 2; string prompt = 3; RuntimeAttachments attachments = 4; + string delivery = 5; } message RuntimeUsage { diff --git a/sdk/packages/rpc/src/runtime-chat-client.ts b/sdk/packages/rpc/src/runtime-chat-client.ts index 082f8885a4..86cc9f8c85 100644 --- a/sdk/packages/rpc/src/runtime-chat-client.ts +++ b/sdk/packages/rpc/src/runtime-chat-client.ts @@ -36,9 +36,8 @@ export class RpcRuntimeChatClient { async sendSession( sessionId: string, request: RpcChatRunTurnRequest, - ): Promise { - const response = await this.client.sendRuntimeSession(sessionId, request); - return response.result; + ): Promise<{ result?: RpcChatTurnResult; queued?: boolean }> { + return await this.client.sendRuntimeSession(sessionId, request); } async abortSession(sessionId: string): Promise { diff --git a/sdk/packages/rpc/src/runtime-chat-command-bridge.ts b/sdk/packages/rpc/src/runtime-chat-command-bridge.ts index 9045d90104..915e3018df 100644 --- a/sdk/packages/rpc/src/runtime-chat-command-bridge.ts +++ b/sdk/packages/rpc/src/runtime-chat-command-bridge.ts @@ -130,10 +130,16 @@ export async function runRpcRuntimeCommandBridge(options: { const parsedResult = options.parseSendResult ? options.parseSendResult(resultRaw) : resultRaw; + const responsePayload = + parsedResult && + typeof parsedResult === "object" && + ("result" in parsedResult || "queued" in parsedResult) + ? (parsedResult as Record) + : { result: parsedResult }; respond({ type: "response", requestId, - response: { result: parsedResult }, + response: responsePayload, }); return false; } diff --git a/sdk/packages/rpc/src/runtime-chat-stream-relay.ts b/sdk/packages/rpc/src/runtime-chat-stream-relay.ts index 78ace3b8a4..cee1d38487 100644 --- a/sdk/packages/rpc/src/runtime-chat-stream-relay.ts +++ b/sdk/packages/rpc/src/runtime-chat-stream-relay.ts @@ -26,6 +26,11 @@ export type RpcRuntimeBridgeStreamLine = error?: string; durationMs?: number; } + | { + type: "pending_prompts"; + sessionId: string; + prompts: unknown[]; + } | { type: "error"; message: string; @@ -168,6 +173,14 @@ export function createRpcRuntimeStreamRelay(options: { ? payload.durationMs : undefined, }); + return; + } + if (event.eventType === "runtime.chat.pending_prompts") { + options.writeLine({ + type: "pending_prompts", + sessionId: event.sessionId, + prompts: Array.isArray(payload.prompts) ? payload.prompts : [], + }); } }, onError: (error: Error) => { diff --git a/sdk/packages/rpc/src/server/runtime.ts b/sdk/packages/rpc/src/server/runtime.ts index 9695b346ed..de7116c5f2 100644 --- a/sdk/packages/rpc/src/server/runtime.ts +++ b/sdk/packages/rpc/src/server/runtime.ts @@ -1,6 +1,9 @@ import { randomUUID } from "node:crypto"; import type { SchedulerService } from "@clinebot/scheduler"; -import type { RpcProviderActionRequest } from "@clinebot/shared"; +import type { + RpcChatStartSessionRequest, + RpcProviderActionRequest, +} from "@clinebot/shared"; import type * as grpc from "@grpc/grpc-js"; import { fromProtoStruct, @@ -353,7 +356,7 @@ export class ClineGatewayRuntime { if (!handler) { throw new Error("runtime start handler is not configured"); } - const payload = request.request + const payload: RpcChatStartSessionRequest | undefined = request.request ? { sessionId: safeString(request.request.sessionId), workspaceRoot: safeString(request.request.workspaceRoot), @@ -446,6 +449,11 @@ export class ClineGatewayRuntime { if (!sessionId) { throw new Error("sessionId is required"); } + const delivery: "queue" | "steer" | undefined = + request.request?.delivery === "queue" || + request.request?.delivery === "steer" + ? request.request.delivery + : undefined; const payload = request.request ? { config: { @@ -513,6 +521,7 @@ export class ClineGatewayRuntime { content: fromProtoValue(message.content), })), prompt: safeString(request.request.prompt), + delivery, attachments: request.request.attachments ? { userImages: request.request.attachments.userImages ?? [], @@ -530,6 +539,9 @@ export class ClineGatewayRuntime { throw new Error("runtime send request is required"); } const result = await handler(sessionId, payload); + if (!result.result) { + return {}; + } return { result: { text: safeString(result.result.text), diff --git a/sdk/packages/rpc/src/server/server-start.ts b/sdk/packages/rpc/src/server/server-start.ts index 8fe6178934..5c230af7a9 100644 --- a/sdk/packages/rpc/src/server/server-start.ts +++ b/sdk/packages/rpc/src/server/server-start.ts @@ -124,7 +124,18 @@ export async function startRpcServer( ? new SchedulerService({ runtimeHandlers: { startSession: options.runtimeHandlers.startSession, - sendSession: options.runtimeHandlers.sendSession, + sendSession: async (sessionId, request) => { + const result = await options.runtimeHandlers?.sendSession?.( + sessionId, + request, + ); + if (!result?.result) { + throw new Error( + "scheduler runtime send unexpectedly queued a turn", + ); + } + return { result: result.result }; + }, abortSession: options.runtimeHandlers.abortSession, stopSession: options.runtimeHandlers.stopSession, }, diff --git a/sdk/packages/rpc/src/types.ts b/sdk/packages/rpc/src/types.ts index a241e48ba5..d60afe7b2f 100644 --- a/sdk/packages/rpc/src/types.ts +++ b/sdk/packages/rpc/src/types.ts @@ -37,7 +37,7 @@ export interface RpcRuntimeHandlers { sendSession?: ( sessionId: string, request: RpcChatRunTurnRequest, - ) => Promise<{ result: RpcChatTurnResult }>; + ) => Promise<{ result?: RpcChatTurnResult; queued?: boolean }>; stopSession?: (sessionId: string) => Promise<{ applied: boolean }>; abortSession?: (sessionId: string) => Promise<{ applied: boolean }>; runProviderAction?: ( diff --git a/sdk/packages/shared/src/rpc/runtime.ts b/sdk/packages/shared/src/rpc/runtime.ts index 2e48f4d846..7b72ef56e7 100644 --- a/sdk/packages/shared/src/rpc/runtime.ts +++ b/sdk/packages/shared/src/rpc/runtime.ts @@ -75,6 +75,7 @@ export interface RpcChatRunTurnRequest { messages?: RpcChatMessage[]; prompt: string; attachments?: RpcChatAttachments; + delivery?: "queue" | "steer"; } export interface RpcChatToolCallResult {