mirror of
https://github.com/alphaXiv/OpenResearch.git
synced 2026-10-02 01:34:34 +08:00
OR-166 Queue chat messages sent while a turn is running (mid-run steering) (#178)
* Queue messages sent while a turn is running
Sending a chat message while a turn is in flight previously errored with
"session is busy — interrupt it first". Instead, park the message and run
it automatically when the current turn finishes — Claude-desktop-style
mid-run steering. Lives in ChatHost above the harnesses, so Claude, Codex,
and OpenCode are all covered.
Backend: a per-session in-memory queue on ChatHost; send_message enqueues
on a busy claim and emits chat.queued; a natural turn completion drains the
next message (FIFO); user Stop and session delete clear the queue; a new
DELETE .../queue/{itemId} cancels one parked message; the messages endpoint
returns the queued snapshot for reload recovery.
Frontend: send() parks instead of bailing while busy; queued messages render
as chips above the composer with a cancel ✕; chat.queued drives the reducer;
the seed snapshot restores chips after a reload.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
* Address review: drop vs re-park on drain failure, clear queue on project delete
- drain_queue re-parks a message only on a genuine busy-race (is_busy after
the failed send); genuine setup errors and deleting sessions now drop+log
instead of stranding a chip that re-fails on every future drain.
- delete_project clears each session's queue (parity with delete_session), so
no parked entries orphan a deleted project.
- Image/file-only parked messages get an "N attachment(s)" chip label instead
of a blank one.
- Import ordering + comment trims from review.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
* Coalesce queued messages into one turn
When several messages are parked behind a running turn, drain them all into a
single combined turn (texts joined by blank lines, attachments concatenated,
most-recent composer overrides win) instead of running a full turn per message
— so successive steering prompts run together as soon as the turn ends.
Re-park on a busy-race restores the individual messages (front, arrival order)
so they re-coalesce on the next drain.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
---------
Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.8
parent
82ea072291
commit
a94c399458
+19
-2
@@ -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<AppState>, Path(id): Path<String>) -
|
||||
);
|
||||
}
|
||||
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<String>) -> ApiResult {
|
||||
async fn chat_messages(State(state): State<AppState>, Path(id): Path<String>) -> 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<AppState>, Path(id): Path<String>) -
|
||||
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<AppState>,
|
||||
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 {
|
||||
|
||||
+188
-6
@@ -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<HashSet<String>>,
|
||||
/// The port `orx up` bound, for the bridge env contract.
|
||||
up_port: std::sync::OnceLock<u16>,
|
||||
/// 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<HashMap<String, VecDeque<QueuedMessage>>>,
|
||||
}
|
||||
|
||||
/// 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<ImageAttachment>,
|
||||
}
|
||||
|
||||
/// 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<Value> {
|
||||
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<Self>,
|
||||
session_id: &'a str,
|
||||
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
|
||||
Box::pin(async move {
|
||||
// Scope the guard: a std mutex must never be held across an await.
|
||||
let items: Vec<QueuedMessage> = {
|
||||
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::<Vec<_>>()
|
||||
.join("\n\n");
|
||||
let images: Vec<ImageAttachment> = 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<Self>,
|
||||
@@ -1059,7 +1207,7 @@ impl ChatHost {
|
||||
overrides: TurnOverrides,
|
||||
images: Vec<ImageAttachment>,
|
||||
) -> 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<String>,
|
||||
overrides: TurnOverrides,
|
||||
images: Vec<ImageAttachment>,
|
||||
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<String>,
|
||||
pub permission_mode: Option<String>,
|
||||
|
||||
+15
-3
@@ -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 {
|
||||
|
||||
@@ -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<string, ChatMessage[]>;
|
||||
busySessions: Set<string>;
|
||||
// Messages parked behind a running turn, per session, oldest first.
|
||||
queuedBySession: Record<string, QueuedMessage[]>;
|
||||
}
|
||||
|
||||
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<string>(),
|
||||
queuedBySession: {},
|
||||
});
|
||||
const [harnesses, setHarnesses] = useState<Harness[]>([]);
|
||||
const [selection, setSelection] = useState<ModelSelection | null>(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 && (
|
||||
<div className="composer-queued flex flex-col gap-1 mb-1.5">
|
||||
{queued.map((q) => (
|
||||
<div
|
||||
key={q.id}
|
||||
className="queued-chip flex items-center gap-2 py-1.5 px-2.5 text-sm text-subtext bg-surface border border-border rounded-sm"
|
||||
title={q.text}
|
||||
>
|
||||
<Clock size={13} className="shrink-0 text-muted" />
|
||||
<span className="flex-1 overflow-hidden text-ellipsis whitespace-nowrap">
|
||||
{q.text}
|
||||
</span>
|
||||
<span className="shrink-0 text-xs text-muted uppercase tracking-wide">Queued</span>
|
||||
<button
|
||||
title="Cancel queued message"
|
||||
aria-label="Cancel queued message"
|
||||
onClick={() => cancelQueued(q.id)}
|
||||
className="shrink-0 inline-flex items-center justify-center w-4 h-4 p-0 border-0 rounded-full text-muted cursor-pointer [&:hover]:bg-text [&:hover]:text-background"
|
||||
>
|
||||
<X size={11} />
|
||||
</button>
|
||||
</div>
|
||||
))}
|
||||
</div>
|
||||
)}
|
||||
<div className="composer-box relative flex flex-col border border-border rounded-md bg-background" data-onboarding="composer">
|
||||
{activeHarness && !activeHarness.agentReady && (
|
||||
<div className="composer-harness-warning py-2 px-3 text-subtext text-xs leading-normal border-b border-b-border-variant [&_strong]:text-accent-amber [&_strong]:font-medium [&_code]:font-mono [&_code]:text-text">
|
||||
|
||||
@@ -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).
|
||||
|
||||
+14
-1
@@ -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<HarnessAuthEvent>(e as MessageEvent);
|
||||
if (d?.harness && d.authState) emitHarnessAuth(d);
|
||||
|
||||
Reference in New Issue
Block a user