diff --git a/src/commands/up.rs b/src/commands/up.rs index abd10554..ec5d22cf 100644 --- a/src/commands/up.rs +++ b/src/commands/up.rs @@ -392,6 +392,10 @@ fn router(state: AppState) -> Router { .route("/api/chat/sessions/{id}/worktree", get(session_worktree)) .route("/api/chat/sessions/{id}/message", post(send_chat_message)) .route("/api/chat/sessions/{id}/interrupt", post(interrupt_chat)) + .route( + "/api/chat/sessions/{id}/queue/{itemId}", + axum::routing::delete(cancel_queued_chat), + ) .route("/api/chat/sessions/{id}/respond", post(respond_chat)) // Internal: the `orx mcp-gate` permission bridge's long-poll (plan // mode). Token-authenticated in the handler; blocks until the surfaced @@ -1427,6 +1431,7 @@ async fn delete_project(State(state): State, Path(id): Path) - ); } for session in &sessions { + state.chat.clear_queue(&session.id); let _ = state.chat.interrupt(&session.id).await; state.chat.opencode.kill_session(&session.id).await; state.chat.codex.kill_session(&session.id).await; @@ -3876,12 +3881,15 @@ async fn update_chat_session( )) } -async fn chat_messages(Path(id): Path) -> ApiResult { +async fn chat_messages(State(state): State, Path(id): Path) -> ApiResult { Store::open()? .get_chat_session(&id)? .ok_or_else(|| not_found("chat session"))?; let messages = local::chat::list_messages(&id)?; - Ok(Json(json!({ "messages": messages }))) + // Parked messages are in-memory, so a reload mid-turn recovers them here + // rather than from the store. + let queued = state.chat.queued_items(&id); + Ok(Json(json!({ "messages": messages, "queued": queued }))) } #[derive(Deserialize)] @@ -3953,6 +3961,15 @@ async fn interrupt_chat(State(state): State, Path(id): Path) - Ok(Json(json!({ "ok": true }))) } +/// Cancel one message parked behind a running turn (the ✕ on a queued chip). +async fn cancel_queued_chat( + State(state): State, + Path((id, item_id)): Path<(String, String)>, +) -> ApiResult { + let removed = state.chat.cancel_queued(&id, &item_id); + Ok(Json(json!({ "ok": true, "removed": removed }))) +} + #[derive(Deserialize)] #[serde(rename_all = "camelCase")] struct RespondReq { diff --git a/src/local/chat/mod.rs b/src/local/chat/mod.rs index 3ba769fd..e7fcee29 100644 --- a/src/local/chat/mod.rs +++ b/src/local/chat/mod.rs @@ -8,7 +8,7 @@ //! normalized parts into the per-turn assistant message; every flush persists //! the message and broadcasts it as a `chat.message` SSE event. -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::sync::Arc; use std::time::{Duration, Instant}; @@ -296,7 +296,7 @@ impl WirePart { // --- image attachments --------------------------------------------------------- /// A pasted image or uploaded file riding the send-message request. -#[derive(Debug, Deserialize)] +#[derive(Debug, Clone, Deserialize)] #[serde(rename_all = "camelCase")] pub struct ImageAttachment { pub media_type: String, @@ -677,6 +677,31 @@ pub struct ChatHost { bridge_prompted: std::sync::Mutex>, /// The port `orx up` bound, for the bridge env contract. up_port: std::sync::OnceLock, + /// Messages the user sent while the session's turn was in flight, oldest + /// first. `drain_queue` runs the front one when a turn finishes naturally; + /// a user Stop clears the whole queue. In-memory and uncommitted — a queued + /// message only becomes a transcript bubble once it actually runs. + queued: std::sync::Mutex>>, +} + +/// A user message parked while the session was busy, replayed verbatim through +/// the normal send path once the running turn ends. +#[derive(Clone)] +struct QueuedMessage { + id: String, + text: String, + overrides: TurnOverrides, + images: Vec, +} + +/// Chip label for a parked message: its text, or an attachment count for an +/// image/file-only send (which carries no text to show). +fn queued_label(m: &QueuedMessage) -> String { + if !m.text.trim().is_empty() || m.images.is_empty() { + return m.text.clone(); + } + let n = m.images.len(); + format!("{n} attachment{}", if n == 1 { "" } else { "s" }) } /// Reserves a session's turn slot for the duration of `send_message`'s setup. @@ -771,6 +796,7 @@ impl ChatHost { gate_tokens: std::sync::Mutex::new(HashMap::new()), bridge_prompted: std::sync::Mutex::new(HashSet::new()), up_port: std::sync::OnceLock::new(), + queued: std::sync::Mutex::new(HashMap::new()), } } @@ -1051,6 +1077,128 @@ impl ChatHost { }) } + /// The session's parked messages, oldest first — for the reload snapshot. + pub fn queued_items(&self, session_id: &str) -> Vec { + self.queued + .lock() + .unwrap() + .get(session_id) + .map(|q| { + q.iter() + .map(|m| json!({ "id": m.id, "text": queued_label(m) })) + .collect() + }) + .unwrap_or_default() + } + + fn queued_json(&self, session_id: &str) -> Value { + json!({ "sessionId": session_id, "items": self.queued_items(session_id) }) + } + + fn emit_queued(&self, session_id: &str) { + self.emit("chat.queued", self.queued_json(session_id)); + } + + /// Drop every parked message for a session (user Stop / delete). Emits an + /// empty `chat.queued` only if there was something to clear. + pub fn clear_queue(&self, session_id: &str) { + let had = self + .queued + .lock() + .unwrap() + .remove(session_id) + .is_some_and(|q| !q.is_empty()); + if had { + self.emit_queued(session_id); + } + } + + /// Remove one parked message by id (the ✕ on a queued chip). + pub fn cancel_queued(&self, session_id: &str, item_id: &str) -> bool { + let removed = { + let mut map = self.queued.lock().unwrap(); + let Some(q) = map.get_mut(session_id) else { + return false; + }; + let before = q.len(); + q.retain(|m| m.id != item_id); + let removed = q.len() != before; + if q.is_empty() { + map.remove(session_id); + } + removed + }; + if removed { + self.emit_queued(session_id); + } + removed + } + + /// Drain the whole parked queue into a single turn once the current turn + /// finishes naturally — successive steering messages run together in one + /// turn, not a full turn each. Boxed return: `drain_queue` → + /// `send_message_showing` → (spawned) `drain_queue` is an async recursion + /// cycle the auto-`Send` solver can't close on its own, so we assert the + /// boxed future is `Send` to break it. + fn drain_queue<'a>( + self: &'a Arc, + session_id: &'a str, + ) -> std::pin::Pin + Send + 'a>> { + Box::pin(async move { + // Scope the guard: a std mutex must never be held across an await. + let items: Vec = { + let mut map = self.queued.lock().unwrap(); + match map.remove(session_id) { + Some(q) => q.into(), + None => return, + } + }; + if items.is_empty() { + return; + } + self.emit_queued(session_id); + // Coalesce every parked message into one turn: join their texts and + // concatenate their attachments; the most recent composer overrides win. + let text = items + .iter() + .map(|m| m.text.trim()) + .filter(|t| !t.is_empty()) + .collect::>() + .join("\n\n"); + let images: Vec = items + .iter() + .flat_map(|m| m.images.iter().cloned()) + .collect(); + let overrides = items + .last() + .map(|m| m.overrides.clone()) + .unwrap_or_default(); + if let Err(err) = self + .send_message_showing(session_id, text, None, overrides, images, false) + .await + { + // Re-park only for the genuine race: a fresh send claimed the + // slot in the gap after `finish_turn` freed it (session busy + // again), so restore the messages up front and let that turn + // drain them. Any other failure (a real setup error, or a session + // being deleted) has no turn to retry against — drop them rather + // than strand chips that re-fail on every future drain. + if self.is_busy(session_id).await { + { + let mut map = self.queued.lock().unwrap(); + let q = map.entry(session_id.to_string()).or_default(); + for item in items.into_iter().rev() { + q.push_front(item); + } + } + self.emit_queued(session_id); + } else { + eprintln!("orx up: dropped queued messages after send failure: {err}"); + } + } + }) + } + /// Persist the user message and run one harness turn in the background. pub async fn send_message( self: &Arc, @@ -1059,7 +1207,7 @@ impl ChatHost { overrides: TurnOverrides, images: Vec, ) -> Result<()> { - self.send_message_showing(session_id, text, None, overrides, images) + self.send_message_showing(session_id, text, None, overrides, images, true) .await } @@ -1078,6 +1226,7 @@ impl ChatHost { transcript_text: Option, overrides: TurnOverrides, images: Vec, + queue_if_busy: bool, ) -> Result<()> { // Atomically claim the session's turn slot: the busy-check and the // reservation happen under one lock so two concurrent sends (or a @@ -1085,6 +1234,24 @@ impl ChatHost { // same session. `_guard` releases the reservation on any early error. let _guard = match TurnGuard::claim(self, session_id).await { Some(guard) => guard, + // Busy: park a genuine user send (Claude-desktop steering) so it + // runs when the turn ends, instead of rejecting it. System/resume + // sends pass `queue_if_busy = false` and keep the old rejection. + None if queue_if_busy && !(text.trim().is_empty() && images.is_empty()) => { + self.queued + .lock() + .unwrap() + .entry(session_id.to_string()) + .or_default() + .push_back(QueuedMessage { + id: format!("q_{}", uuid::Uuid::new_v4()), + text, + overrides, + images, + }); + self.emit_queued(session_id); + return Ok(()); + } None => return Err(anyhow!("session is busy — interrupt it first")), }; let store = Store::open()?; @@ -1297,6 +1464,10 @@ impl ChatHost { } } ctx.host.finish_turn(&ctx.session_id).await; + // Natural completion only: a user Stop aborts this task before + // it reaches here (and clears the queue itself), so an + // interrupted turn never drains. + ctx.host.drain_queue(&ctx.session_id).await; }); turns.insert(sid, Some(task.abort_handle())); if let Some(seed) = title_seed { @@ -1387,6 +1558,9 @@ impl ChatHost { // still paint them in arrival order for a few ms; a reload converges // on the stored order.) let created_at = now_ms(); + // Stop means stop everything: drop any messages parked behind this turn + // so they don't fire the moment it aborts. + self.clear_queue(session_id); if !self.interrupt(session_id).await? { return Ok(()); } @@ -1504,8 +1678,15 @@ impl ChatHost { "plan" | "permission" => Some(req.note.clone().unwrap_or_default()), _ => None, }; - self.send_message_showing(&req.session_id, text, transcript, overrides, Vec::new()) - .await?; + self.send_message_showing( + &req.session_id, + text, + transcript, + overrides, + Vec::new(), + false, + ) + .await?; self.resolve_prompt_card(&req); Ok(()) } @@ -1646,6 +1827,7 @@ impl ChatHost { let _deleting = self .begin_session_delete(session_id) .ok_or_else(|| anyhow!("session deletion is already in progress"))?; + self.clear_queue(session_id); let _ = self.interrupt(session_id).await; // A live opencode serve child would keep running in (and lock) the // session's worktree; the resident claude child's cwd is that worktree @@ -1895,7 +2077,7 @@ impl ChatHost { /// Composer selections a single message can override, mirroring the sticky /// per-session settings. Empty/None fields leave the stored value in place. -#[derive(Debug, Default)] +#[derive(Debug, Default, Clone)] pub struct TurnOverrides { pub model: Option, pub permission_mode: Option, diff --git a/ui/src/api.ts b/ui/src/api.ts index ae7e271f..ce65695c 100644 --- a/ui/src/api.ts +++ b/ui/src/api.ts @@ -1162,10 +1162,22 @@ export const renameChatSession = (sessionId: string, title: string) => (r) => r.session, ); +/** A message the user sent while a turn was running, parked to run next. */ +export interface QueuedMessage { + id: string; + text: string; +} + export const getChatMessages = (sessionId: string) => - get<{ messages: ChatMessage[] }>(`/api/chat/sessions/${sessionId}/messages`).then( - (r) => r.messages, - ); + get<{ messages: ChatMessage[]; queued?: QueuedMessage[] }>( + `/api/chat/sessions/${sessionId}/messages`, + ).then((r) => ({ messages: r.messages, queued: r.queued ?? [] })); + +/** Cancel a still-parked message (the ✕ on a queued chip). */ +export const cancelQueuedMessage = (sessionId: string, itemId: string) => + fetch(`/api/chat/sessions/${sessionId}/queue/${encodeURIComponent(itemId)}`, { + method: "DELETE", + }).then((r) => json<{ ok: boolean; removed: boolean }>(r)); /** A pasted image or uploaded file riding a chat message. */ export interface ChatImageAttachment { diff --git a/ui/src/components/ChatPanel.tsx b/ui/src/components/ChatPanel.tsx index beff36e9..82dce9b2 100644 --- a/ui/src/components/ChatPanel.tsx +++ b/ui/src/components/ChatPanel.tsx @@ -5,6 +5,7 @@ import { ChartSpline, Check, ChevronRight, + Clock, CornerDownLeft, FileText, FlaskConical, @@ -33,6 +34,7 @@ import { } from "react"; import { BrandMark } from "./Wordmark"; import { + cancelQueuedMessage, chatAttachmentUrl, createChatSession, deleteChatSession, @@ -57,6 +59,7 @@ import { type ChatSession, type Harness, type PromptAnswer, + type QueuedMessage, type SkillInfo, } from "../api"; import { onChatEvent } from "../events"; @@ -128,11 +131,19 @@ const PROMPT_ACTIONS_CLASS_NAME = [ interface ChatState { messagesBySession: Record; busySessions: Set; + // Messages parked behind a running turn, per session, oldest first. + queuedBySession: Record; } type Action = | { type: "reset" } - | { type: "seed"; sessionId: string; messages: ChatMessage[]; onlyIfAbsent?: boolean } + | { + type: "seed"; + sessionId: string; + messages: ChatMessage[]; + queued?: QueuedMessage[]; + onlyIfAbsent?: boolean; + } | { type: "upsertMessage"; sessionId: string; message: ChatMessage } | { type: "optimisticUser"; @@ -144,6 +155,7 @@ type Action = // `known` scopes the reseed: flags for sessions outside it (other projects — // busy events aren't project-filtered) are carried forward, not wiped. | { type: "seedBusy"; sessions: string[]; known: string[] } + | { type: "setQueued"; sessionId: string; items: QueuedMessage[] } | { type: "forget"; sessionId: string }; const LOCAL_PREFIX = "local-"; @@ -164,7 +176,7 @@ function upsertMessage(list: ChatMessage[], message: ChatMessage): ChatMessage[] function reducer(state: ChatState, action: Action): ChatState { switch (action.type) { case "reset": - return { messagesBySession: {}, busySessions: new Set() }; + return { messagesBySession: {}, busySessions: new Set(), queuedBySession: {} }; case "seed": // onlyIfAbsent: recover a failed fetch without clobbering messages that // streamed in via SSE during it (a `message` event already created the key). @@ -172,6 +184,12 @@ function reducer(state: ChatState, action: Action): ChatState { return { ...state, messagesBySession: { ...state.messagesBySession, [action.sessionId]: action.messages }, + // A seed is the authoritative snapshot, so it also (re)sets the parked + // queue — recovering it after a reload or an SSE gap. + queuedBySession: { + ...state.queuedBySession, + [action.sessionId]: action.queued ?? [], + }, }; case "upsertMessage": { const list = state.messagesBySession[action.sessionId] ?? []; @@ -215,6 +233,12 @@ function reducer(state: ChatState, action: Action): ChatState { for (const id of state.busySessions) if (!known.has(id)) busySessions.add(id); return { ...state, busySessions }; } + case "setQueued": { + return { + ...state, + queuedBySession: { ...state.queuedBySession, [action.sessionId]: action.items }, + }; + } case "forget": { // Deleted session: drop its transcript and busy flag so a same-id event // arriving late can't render stale state. @@ -222,7 +246,9 @@ function reducer(state: ChatState, action: Action): ChatState { delete messagesBySession[action.sessionId]; const busySessions = new Set(state.busySessions); busySessions.delete(action.sessionId); - return { messagesBySession, busySessions }; + const queuedBySession = { ...state.queuedBySession }; + delete queuedBySession[action.sessionId]; + return { messagesBySession, busySessions, queuedBySession }; } } } @@ -1345,6 +1371,7 @@ export function ChatPanel({ const [state, dispatch] = useReducer(reducer, { messagesBySession: {}, busySessions: new Set(), + queuedBySession: {}, }); const [harnesses, setHarnesses] = useState([]); const [selection, setSelection] = useState(preferredAgent); @@ -1641,7 +1668,9 @@ export function ChatPanel({ if (!activeId || loadedSessions.current.has(activeId)) return; loadedSessions.current.add(activeId); getChatMessages(activeId) - .then((messages) => dispatch({ type: "seed", sessionId: activeId, messages })) + .then(({ messages, queued }) => + dispatch({ type: "seed", sessionId: activeId, messages, queued }), + ) .catch(() => { // Recover from a failed fetch to a usable state rather than a stuck // "Loading conversation…" spinner: seed an empty transcript (clears @@ -1707,6 +1736,9 @@ export function ChatPanel({ case "busy": dispatch({ type: "busy", sessionId: ev.sessionId, busy: ev.busy }); break; + case "queued": + dispatch({ type: "setQueued", sessionId: ev.sessionId, items: ev.items }); + break; case "usage": setSessions((cur) => cur.map((s) => (s.id === ev.sessionId ? { ...s, contextUsage: ev.usage } : s)), @@ -1736,8 +1768,8 @@ export function ChatPanel({ const reseed = (allowRetry: boolean) => { const gen = msgGen.current; getChatMessages(activeId) - .then((messages) => { - dispatch({ type: "seed", sessionId: activeId, messages }); + .then(({ messages, queued }) => { + dispatch({ type: "seed", sessionId: activeId, messages, queued }); if (allowRetry && msgGen.current !== gen) reseed(false); }) .catch(() => {}); @@ -1748,6 +1780,9 @@ export function ChatPanel({ const messages = activeId ? (state.messagesBySession[activeId] ?? []) : []; const busy = activeId ? state.busySessions.has(activeId) : false; + // Messages the user parked behind the running turn (oldest first). Populated + // by chat.queued events and the seed snapshot; each runs when its turn ends. + const queued = activeId ? (state.queuedBySession[activeId] ?? []) : []; // A session whose transcript hasn't been seeded yet: its key is absent from // messagesBySession (vs. present-but-empty for a genuinely empty session). // Switching to an existing session leaves this true for the getChatMessages @@ -1929,7 +1964,39 @@ export function ChatPanel({ }); return; } - if (busy) return; + if (busy) { + // A turn is already running: park this message (Claude-desktop steering) + // so it runs when the turn ends, instead of dropping it. The server + // enqueues it and echoes chat.queued to render the chip — no optimistic + // transcript bubble, since it hasn't run yet. + if (!activeId || !activeHarness?.agentReady) return; + const sid = activeId; + setDraft(""); + setPickedSkill(null); + setAttachments([]); + setAttachError(null); + const turnOpts = composerSelection + ? { + model: composerSelection.model, + permissionMode: composerSelection.permissionMode, + reasoningLevel: composerSelection.reasoningLevel, + } + : {}; + setSessionOverride({}); + const images: ChatImageAttachment[] = pending.map((a) => ({ + mediaType: a.mediaType, + dataBase64: a.dataUrl.slice(a.dataUrl.indexOf(",") + 1), + name: a.name, + })); + try { + await sendChatMessage(sid, text, turnOpts, images.length ? images : undefined); + } catch { + // Never reached the queue — restore the composer so a retry is one keypress. + setDraft((cur) => cur || text); + setAttachments((cur) => (cur.length ? cur : pending)); + } + return; + } if (!activeHarness?.agentReady) return; // `composerSelection` already resolves to the open session's settings (+ any // unsent tweak) or, for a new session, the global preference. @@ -2038,6 +2105,19 @@ export function ChatPanel({ if (activeId) void interruptChat(activeId); } + // Optimistic: drop locally now; the server's chat.queued echo reconciles. A + // message that already started running server-side simply isn't found. + function cancelQueued(itemId: string) { + if (!activeId) return; + const sid = activeId; + dispatch({ + type: "setQueued", + sessionId: sid, + items: queued.filter((q) => q.id !== itemId), + }); + void cancelQueuedMessage(sid, itemId).catch(() => {}); + } + // Escape stops the streaming turn and drops focus back into the composer, // mirroring the Claude Code desktop app. Harness-agnostic — `stop()` → // `interruptChat` interrupts whichever harness (Claude, Codex, OpenCode, …) @@ -2145,7 +2225,7 @@ export function ChatPanel({ // just-started optimistic flag), so the optimistic dispatch above // can't wedge true after a no-op or failure. getChatMessages(sid) - .then((messages) => dispatch({ type: "seed", sessionId: sid, messages })) + .then(({ messages, queued }) => dispatch({ type: "seed", sessionId: sid, messages, queued })) .catch(() => {}); listChatSessions(projectId) .then((list) => @@ -2495,6 +2575,31 @@ export function ChatPanel({ }} /> )} + {queued.length > 0 && ( +
+ {queued.map((q) => ( +
+ + + {q.text} + + Queued + +
+ ))} +
+ )}
{activeHarness && !activeHarness.agentReady && (
diff --git a/ui/src/components/SubagentTab.tsx b/ui/src/components/SubagentTab.tsx index 76f2a716..60f82112 100644 --- a/ui/src/components/SubagentTab.tsx +++ b/ui/src/components/SubagentTab.tsx @@ -32,7 +32,7 @@ export function SubagentTab({ useEffect(() => { let live = true; getChatMessages(sessionId) - .then((m) => live && setMessages(m)) + .then(({ messages }) => live && setMessages(messages)) .catch(() => live && setMessages([])); // Live updates: replace the message the event carries (assistant turns // re-broadcast the whole message on every flush). diff --git a/ui/src/events.ts b/ui/src/events.ts index fc3ee698..34aa9eec 100644 --- a/ui/src/events.ts +++ b/ui/src/events.ts @@ -3,7 +3,15 @@ // terminals can subscribe without threading props everywhere. import { useEffect, useRef } from "react"; -import type { ChatMessage, ChatSession, ContextUsage, Experiment, Project, Run } from "./api"; +import type { + ChatMessage, + ChatSession, + ContextUsage, + Experiment, + Project, + QueuedMessage, + Run, +} from "./api"; export interface RunLogEvent { runId: string; @@ -38,6 +46,7 @@ export type ChatEvent = | { type: "message"; sessionId: string; message: ChatMessage } | { type: "busy"; sessionId: string; busy: boolean } | { type: "usage"; sessionId: string; usage: ContextUsage } + | { type: "queued"; sessionId: string; items: QueuedMessage[] } /** The EventSource re-connected after a drop. Chat events are edge-only (no * snapshot on connect), so frames emitted during the gap are lost for good — * subscribers must refetch whatever they render from chat events. */ @@ -176,6 +185,10 @@ export function useOrxEvents(handlers: OrxEventHandlers) { const d = parse<{ sessionId: string; usage: ContextUsage }>(e as MessageEvent); if (d?.sessionId && d.usage) emitChat({ type: "usage", sessionId: d.sessionId, usage: d.usage }); }); + es.addEventListener("chat.queued", (e) => { + const d = parse<{ sessionId: string; items: QueuedMessage[] }>(e as MessageEvent); + if (d?.sessionId) emitChat({ type: "queued", sessionId: d.sessionId, items: d.items ?? [] }); + }); es.addEventListener("harness.auth", (e) => { const d = parse(e as MessageEvent); if (d?.harness && d.authState) emitHarnessAuth(d);