feat(chat): queue pending user messages with sending/error/retry states (#281)

* feat(chat): queue pending user messages with sending/error/retry states

Replace the single optimisticUserMessage slot with a FIFO queue of pending
user messages so the user gets immediate feedback when they hit send and
multiple in-flight messages do not clobber each other.

- Reshape optimistic-user-message-store around enqueuePendingMessage /
  consumeOldestSendingMessage / markPendingMessageError /
  markPendingMessageSending with a 'sending' | 'error' status per entry.
- Render queued messages via a new PendingUserMessages component using
  ChatMessage's new pendingStatus + onRetry props, applying a faded
  treatment for 'sending' and an error banner + retry link for 'error'.
- Wire the WebSocket context to consume the oldest 'sending' entry when
  the server echoes a UserMessageEvent back, so the queue drains FIFO.
- Update ChatInterface, GitControlBar, TaskCard and useHandleBuildPlanClick
  to enqueue pending messages, and flip the matching entry to 'error'
  when the send call rejects.
- Add i18n strings for Sending / Send failed / Retry.

Co-authored-by: openhands <openhands@all-hands.dev>

* fix(chat): scope pending message queue per conversation and stop double 'Sending' render

Two follow-up fixes on top of the pending-message-queue refactor:

1. **Per-conversation scoping.** The optimistic queue is global but each entry
   is now tagged with the `conversationId` it was enqueued from, and
   `PendingUserMessages` filters to entries matching the active conversation
   via `useOptionalConversationId`. This means switching conversations no
   longer carries 'Sending…' bubbles over, and the WebSocket
   `UserMessageEvent` ack for one conversation can never pop a pending entry
   belonging to another. `consumeOldestSendingMessage` now takes the
   conversation id and only matches within it.

   Call sites updated: `chat-interface.tsx`, `git-control-bar.tsx`,
   `use-handle-build-plan-click.ts`. `task-card.tsx` moves the enqueue into
   the `createConversation` `onSuccess` callback so the message is tagged
   with the newly-created conversation id.

2. **Duplicate 'Sending…' bubble.** `<PendingUserMessages />` was rendered
   from two places after the previous rebase: once at the bottom of
   `<Messages>` (where the pending queue would never re-render anyway
   because `Messages` is `React.memo`'d on event ids) and once
   unconditionally from `<ChatInterface>`. Removed the render inside
   `<Messages>` so only the ChatInterface render remains.

   Also fixed a separate double-submit path in `custom-chat-input.tsx`: the
   `submittedMessage` effect listed `onSubmit` in its deps, but
   `onSubmit` (= `handleSendMessage`) is a fresh function on every parent
   render, and the parent re-renders synchronously when the new
   `enqueuePendingMessage` call mutates the store. That made the effect
   fire twice for the same `submittedMessage` value before
   `setSubmittedMessage(null)` flushed. Pinned `onSubmit` behind a ref and
   dropped it from the dep array.

Tests updated to pass `conversationId` and a new test asserts the queue is
not consumed across conversations.

Co-authored-by: openhands <openhands@all-hands.dev>

* fix(chat): address review feedback on pending-message queue

- plan-preview.test.tsx: include `useOptionalConversationId` in the
  `#/hooks/use-conversation-id` mock so the test passes now that
  `useHandleBuildPlanClick` depends on it. This was failing all 21
  PlanPreview tests in CI on the previous push.

- optimistic-user-message-store: replace the module-level counter
  (which never reset between tests) with a `Date.now()` +
  base36 random suffix. Same uniqueness guarantee, no shared state.

- optimistic-user-message-store: collapse `consumeOldestSendingMessage`
  into a single atomic `set()` so the find and the filter can't
  observe an interleaved update from another action.

- git-control-bar: capture the `pendingId` from the
  `enqueuePendingMessage` return value and, if the `send` promise
  rejects, mark that entry as error so the user gets a retry link
  instead of a stuck 'Sending…' bubble.

- chat-message: add `role="status" aria-live="polite"` to the
  'Sending…' label and `role="alert"` to the error label so screen
  readers announce send-state transitions.

Co-authored-by: openhands <openhands@all-hands.dev>

* fix(chat): content-keyed echo matching + 60s watchdog timeout for pending messages

Addresses the bot's reposted 'critical' review threads on PR #281.

**Out-of-order echo matching.**
- `PendingUserMessage` now stores both `text` (user-visible bubble) and
  `content` (the exact string sent to the server, which may include the
  appended 'Files uploaded: …' prompt). `enqueuePendingMessage` accepts
  an optional `content` and defaults it to `text` for call sites that
  don't transform the prompt.
- New `consumeMatchingPendingMessage(conversationId, content)`. It does
  an exact content match first — so an echo of 'second' arriving before
  an echo of 'hello' correctly pops 'second' instead of the oldest
  entry — and falls back to the oldest "sending" entry in the same
  conversation if no exact match exists (lets us still drain the queue
  if the server slightly munges the body).
- `conversation-websocket-context` extracts the echoed text by joining
  the `TextContent` parts of `event.llm_message.content` and passes it
  to the new matcher. Both consumption sites (main WS + planning agent
  WS) use the matcher with the main `conversationId`, so a planning
  sub-agent echo can never consume a main-conversation pending entry.
- `chat-interface` passes `content: prompt` (text + file annotations)
  to `enqueuePendingMessage`.

**60-second watchdog timeout.**
- `enqueuePendingMessage` now schedules a `setTimeout` for
  `PENDING_MESSAGE_TIMEOUT_MS` (60s, exported). If the entry is still in
  'sending' state when the timer fires, it's flipped to 'error' with
  message 'Send timed out' so the user gets a retry link instead of a
  permanently-stuck bubble. Timeout is a no-op if the echo already
  consumed the message or if it was already marked error explicitly
  (the original `errorMessage` is preserved).
- 60s is long enough to cover legitimately slow uploads / agent-server
  latency but short enough to actually rescue stuck bubbles.

**Tests.**
- `optimistic-user-message-store.test.ts` adds coverage for: storing
  separate `text`/`content`, exact-match preference, FIFO fallback,
  skipping entries already in 'error', cross-conversation isolation
  with identical content, watchdog timeout firing, and the two no-op
  cases (echo already consumed, message already failed).

Co-authored-by: openhands <openhands@all-hands.dev>

---------

Co-authored-by: openhands <openhands@all-hands.dev>
This commit is contained in:
Robert Brennan
2026-05-10 16:49:48 -07:00
committed by GitHub
co-authored by openhands
parent a7b35f7098
commit 3e57aa6a04
17 changed files with 1147 additions and 95 deletions
@@ -1,5 +1,19 @@
import { afterEach, beforeEach, describe, expect, it, test, vi } from "vitest";
import { fireEvent, render, screen, within } from "@testing-library/react";
import {
afterEach,
beforeEach,
describe,
expect,
it,
test,
vi,
} from "vitest";
import {
fireEvent,
render,
screen,
waitFor,
within,
} from "@testing-library/react";
import { MemoryRouter, Route, Routes } from "react-router";
import { QueryClient, QueryClientProvider } from "@tanstack/react-query";
import { renderWithProviders, useParamsMock } from "test-utils";
@@ -19,6 +33,13 @@ import { useEventStore } from "#/stores/use-event-store";
import { useAgentState } from "#/hooks/use-agent-state";
import { useLoadOlderEvents } from "#/hooks/use-load-older-events";
import { AgentState } from "#/types/agent-state";
import { useConversationStore } from "#/stores/conversation-store";
import { act } from "@testing-library/react";
const mockSend = vi.fn();
vi.mock("#/hooks/use-send-message", () => ({
useSendMessage: () => ({ send: mockSend }),
}));
vi.mock("#/hooks/query/use-config");
vi.mock("#/hooks/mutation/use-get-trajectory");
@@ -112,7 +133,7 @@ describe("ChatInterface - Chat Suggestions", () => {
});
useOptimisticUserMessageStore.setState({
optimisticUserMessage: null,
pendingMessages: [],
});
useErrorMessageStore.setState({
@@ -163,8 +184,9 @@ describe("ChatInterface - Chat Suggestions", () => {
});
test("should hide chat suggestions when there is an optimistic user message", () => {
useOptimisticUserMessageStore.setState({
optimisticUserMessage: "Optimistic message",
useOptimisticUserMessageStore.getState().enqueuePendingMessage({
conversationId: "test-conversation-id",
text: "Optimistic message",
});
renderWithQueryClient(<ChatInterface />, queryClient);
@@ -202,7 +224,7 @@ describe("ChatInterface - Scroll-up loads older events", () => {
defaultOptions: { queries: { retry: false } },
});
useOptimisticUserMessageStore.setState({ optimisticUserMessage: null });
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
useErrorMessageStore.setState({ errorMessage: null });
(useConfig as unknown as ReturnType<typeof vi.fn>).mockReturnValue({
@@ -494,6 +516,118 @@ describe("ChatInterface - Scroll-up loads older events", () => {
});
});
describe("ChatInterface - Pending message queue", () => {
let queryClient: QueryClient;
beforeEach(() => {
mockSend.mockReset();
queryClient = new QueryClient({
defaultOptions: { queries: { retry: false } },
});
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
useErrorMessageStore.setState({ errorMessage: null });
(useConfig as unknown as ReturnType<typeof vi.fn>).mockReturnValue({
data: {},
});
(useGetTrajectory as unknown as ReturnType<typeof vi.fn>).mockReturnValue({
mutate: vi.fn(),
mutateAsync: vi.fn(),
isLoading: false,
});
(
useUnifiedUploadFiles as unknown as ReturnType<typeof vi.fn>
).mockReturnValue({
mutateAsync: vi
.fn()
.mockResolvedValue({ skipped_files: [], uploaded_files: [] }),
isLoading: false,
});
useEventStore.setState({
events: [],
eventIds: new Set(),
uiEvents: [],
});
});
afterEach(() => {
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
});
function submitMessage(text: string) {
// The chat input is a contenteditable div; the conversation store exposes
// `submittedMessage` which CustomChatInput watches and forwards to the
// ChatInterface's `onSubmit` handler. Driving that store directly is the
// most reliable way to simulate "the user pressed send" in jsdom.
act(() => {
useConversationStore.setState({ submittedMessage: text });
});
}
function renderInterface() {
render(
<QueryClientProvider client={queryClient}>
<MemoryRouter initialEntries={["/test-conversation-id"]}>
<Routes>
<Route path=":conversationId" element={<ChatInterface />} />
</Routes>
</MemoryRouter>
</QueryClientProvider>,
);
}
it("shows the message in 'sending' state immediately when submitted", async () => {
let resolveSend: ((value: unknown) => void) | undefined;
mockSend.mockImplementation(
() =>
new Promise((resolve) => {
resolveSend = resolve;
}),
);
renderInterface();
submitMessage("hello world");
const pendingMessage = await screen.findByTestId("user-message");
expect(pendingMessage).toHaveTextContent("hello world");
expect(pendingMessage).toHaveAttribute("data-pending-status", "sending");
expect(screen.getByTestId("chat-message-sending")).toBeInTheDocument();
expect(screen.queryByTestId("chat-message-retry")).not.toBeInTheDocument();
resolveSend?.({ queued: false });
});
it("flips the message to 'error' with a retry link when send rejects", async () => {
mockSend.mockRejectedValue(new Error("network down"));
renderInterface();
submitMessage("hello");
await waitFor(() => {
expect(screen.getByTestId("user-message")).toHaveAttribute(
"data-pending-status",
"error",
);
});
expect(screen.getByTestId("chat-message-error")).toBeInTheDocument();
expect(screen.getByTestId("chat-message-retry")).toBeInTheDocument();
});
it("queues multiple submitted messages, each with its own pending entry", async () => {
mockSend.mockResolvedValue({ queued: false });
renderInterface();
submitMessage("first");
submitMessage("second");
await waitFor(() => {
expect(screen.getAllByTestId("user-message")).toHaveLength(2);
});
const messages = screen.getAllByTestId("user-message");
expect(messages[0]).toHaveTextContent("first");
expect(messages[1]).toHaveTextContent("second");
});
});
describe("ChatInterface - Status Indicator", () => {
it("should render ChatStatusIndicator when agent is not awaiting user input / conversation is NOT ready", () => {
vi.mocked(useAgentState).mockReturnValue({
@@ -0,0 +1,153 @@
import React from "react";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { render, screen, waitFor } from "@testing-library/react";
import userEvent from "@testing-library/user-event";
import { renderWithProviders } from "test-utils";
import { PendingUserMessages } from "#/components/features/chat/pending-user-messages";
import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-store";
const ACTIVE_CONVO = "conv-active";
const mockSend = vi.fn();
vi.mock("#/hooks/use-send-message", () => ({
useSendMessage: () => ({ send: mockSend }),
}));
vi.mock("#/hooks/use-conversation-id", () => ({
useOptionalConversationId: () => ({ conversationId: ACTIVE_CONVO }),
useConversationId: () => ({ conversationId: ACTIVE_CONVO }),
}));
describe("PendingUserMessages", () => {
beforeEach(() => {
mockSend.mockReset();
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
});
afterEach(() => {
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
});
it("renders nothing when the queue is empty", () => {
const { container } = render(<PendingUserMessages />);
expect(container).toBeEmptyDOMElement();
});
it("renders each queued message with the faded 'sending' treatment", () => {
useOptimisticUserMessageStore.getState().enqueuePendingMessage({
conversationId: ACTIVE_CONVO,
text: "first message",
});
useOptimisticUserMessageStore.getState().enqueuePendingMessage({
conversationId: ACTIVE_CONVO,
text: "second message",
});
renderWithProviders(<PendingUserMessages />);
const messages = screen.getAllByTestId("user-message");
expect(messages).toHaveLength(2);
expect(messages[0]).toHaveTextContent("first message");
expect(messages[1]).toHaveTextContent("second message");
messages.forEach((message) => {
expect(message).toHaveAttribute("data-pending-status", "sending");
expect(message.className).toMatch(/opacity-60/);
});
expect(screen.getAllByTestId("chat-message-sending")).toHaveLength(2);
});
it("ignores pending entries belonging to a different conversation", () => {
useOptimisticUserMessageStore.getState().enqueuePendingMessage({
conversationId: ACTIVE_CONVO,
text: "mine",
});
useOptimisticUserMessageStore.getState().enqueuePendingMessage({
conversationId: "other-convo",
text: "from another conversation",
});
renderWithProviders(<PendingUserMessages />);
const messages = screen.getAllByTestId("user-message");
expect(messages).toHaveLength(1);
expect(messages[0]).toHaveTextContent("mine");
expect(screen.queryByText("from another conversation")).toBeNull();
});
it("shows an error state with a retry link when the message is in 'error'", () => {
const id = useOptimisticUserMessageStore
.getState()
.enqueuePendingMessage({
conversationId: ACTIVE_CONVO,
text: "broken message",
});
useOptimisticUserMessageStore
.getState()
.markPendingMessageError(id, "Server unavailable");
renderWithProviders(<PendingUserMessages />);
const message = screen.getByTestId("user-message");
expect(message).toHaveAttribute("data-pending-status", "error");
expect(screen.getByTestId("chat-message-error")).toBeInTheDocument();
expect(screen.getByTestId("chat-message-retry")).toBeInTheDocument();
});
it("re-sends and flips back to 'sending' when retry is clicked", async () => {
mockSend.mockResolvedValueOnce({ queued: false });
const id = useOptimisticUserMessageStore
.getState()
.enqueuePendingMessage({
conversationId: ACTIVE_CONVO,
text: "retry me",
});
useOptimisticUserMessageStore
.getState()
.markPendingMessageError(id, "Server unavailable");
renderWithProviders(<PendingUserMessages />);
const user = userEvent.setup();
await user.click(screen.getByTestId("chat-message-retry"));
expect(mockSend).toHaveBeenCalledTimes(1);
expect(mockSend).toHaveBeenCalledWith(
expect.objectContaining({
action: "message",
args: expect.objectContaining({ content: "retry me" }),
}),
);
await waitFor(() => {
const [entry] =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(entry.status).toBe("sending");
expect(entry.errorMessage).toBeUndefined();
});
});
it("flips back to 'error' if the retry attempt also fails", async () => {
mockSend.mockRejectedValueOnce(new Error("still broken"));
const id = useOptimisticUserMessageStore
.getState()
.enqueuePendingMessage({
conversationId: ACTIVE_CONVO,
text: "retry me",
});
useOptimisticUserMessageStore
.getState()
.markPendingMessageError(id, "Server unavailable");
renderWithProviders(<PendingUserMessages />);
const user = userEvent.setup();
await user.click(screen.getByTestId("chat-message-retry"));
await waitFor(() => {
const [entry] =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(entry.status).toBe("error");
expect(entry.errorMessage).toBe("still broken");
});
});
});
@@ -39,7 +39,7 @@ vi.mock("react-i18next", async (importOriginal) => {
});
// Mock services (underlying dependencies of the hook)
const mockSend = vi.fn();
const mockSend = vi.fn().mockResolvedValue({ queued: false });
vi.mock("#/hooks/use-send-message", () => ({
useSendMessage: vi.fn(() => ({
@@ -56,15 +56,17 @@ vi.mock("#/services/chat-service", () => ({
vi.mock("#/hooks/use-conversation-id", () => ({
useConversationId: () => ({ conversationId: "test-conversation-id" }),
useOptionalConversationId: () => ({ conversationId: "test-conversation-id" }),
}));
describe("PlanPreview", () => {
beforeEach(() => {
vi.clearAllMocks();
mockSend.mockResolvedValue({ queued: false });
// Reset store states
localStorage.clear();
useOptimisticUserMessageStore.setState({
optimisticUserMessage: null,
pendingMessages: [],
});
useConversationStore.setState({
conversationMode: "plan",
@@ -81,7 +83,7 @@ describe("PlanPreview", () => {
conversationMode: "code",
});
useOptimisticUserMessageStore.setState({
optimisticUserMessage: null,
pendingMessages: [],
});
localStorage.clear();
});
@@ -193,9 +195,9 @@ describe("PlanPreview", () => {
);
});
it("should set optimistic user message when Build button is clicked", async () => {
it("should enqueue a pending user message when Build button is clicked", async () => {
// Arrange
useOptimisticUserMessageStore.setState({ optimisticUserMessage: null });
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
const user = userEvent.setup();
const expectedPrompt =
"Execute the plan based on the .agents_tmp/PLAN.md file.";
@@ -206,9 +208,11 @@ describe("PlanPreview", () => {
await user.click(buildButton);
// Assert
expect(useOptimisticUserMessageStore.getState().optimisticUserMessage).toBe(
expectedPrompt,
);
const pending =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(pending).toHaveLength(1);
expect(pending[0].text).toBe(expectedPrompt);
expect(pending[0].status).toBe("sending");
});
it("should disable Build button when isBuildDisabled is true", () => {
@@ -40,14 +40,21 @@ export function EventStoreComponent() {
}
/**
* Test component to access and display optimistic user message store values
* Test component to access and display the queue of pending user messages
* tracked locally until the WebSocket echoes them back.
*/
export function OptimisticUserMessageStoreComponent() {
const { optimisticUserMessage } = useOptimisticUserMessageStore();
const { pendingMessages } = useOptimisticUserMessageStore();
return (
<div>
<div data-testid="optimistic-user-message">
{optimisticUserMessage || "none"}
{pendingMessages[0]?.text || "none"}
</div>
<div data-testid="optimistic-user-message-count">
{pendingMessages.length}
</div>
<div data-testid="optimistic-user-message-statuses">
{pendingMessages.map((m) => m.status).join(",") || "none"}
</div>
</div>
);
@@ -15,21 +15,29 @@ vi.mock("#/services/chat-service", () => ({
createChatMessage: vi.fn(),
}));
// The hook now scopes pending messages by conversation id; stub the lookup
// so the hook always sees a stable conversation in the test environment.
vi.mock("#/hooks/use-conversation-id", () => ({
useOptionalConversationId: () => ({ conversationId: "test-conversation-id" }),
useConversationId: () => ({ conversationId: "test-conversation-id" }),
}));
// Import mocked modules
import { useSendMessage } from "#/hooks/use-send-message";
describe("useHandleBuildPlanClick", () => {
const mockSend = vi.fn();
const mockSend = vi.fn().mockResolvedValue({ queued: false });
beforeEach(() => {
vi.clearAllMocks();
mockSend.mockResolvedValue({ queued: false });
// Reset store states
useConversationStore.setState({
conversationMode: "plan",
});
useOptimisticUserMessageStore.setState({
optimisticUserMessage: null,
pendingMessages: [],
});
// Setup send message hook mock
@@ -55,7 +63,7 @@ describe("useHandleBuildPlanClick", () => {
conversationMode: "code",
});
useOptimisticUserMessageStore.setState({
optimisticUserMessage: null,
pendingMessages: [],
});
});
@@ -103,9 +111,9 @@ describe("useHandleBuildPlanClick", () => {
);
});
it("should set optimistic user message when handleBuildPlanClick is called", () => {
it("should enqueue a pending user message when handleBuildPlanClick is called", () => {
// Arrange
useOptimisticUserMessageStore.setState({ optimisticUserMessage: null });
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
const { result } = renderHook(() => useHandleBuildPlanClick());
const expectedPrompt =
"Execute the plan based on the .agents_tmp/PLAN.md file.";
@@ -116,9 +124,11 @@ describe("useHandleBuildPlanClick", () => {
});
// Assert
expect(useOptimisticUserMessageStore.getState().optimisticUserMessage).toBe(
expectedPrompt,
);
const pending =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(pending).toHaveLength(1);
expect(pending[0].text).toBe(expectedPrompt);
expect(pending[0].status).toBe("sending");
});
it("should prevent default and stop propagation when event is provided", () => {
@@ -142,7 +152,7 @@ describe("useHandleBuildPlanClick", () => {
it("should handle call without event parameter", () => {
// Arrange
useConversationStore.setState({ conversationMode: "plan" });
useOptimisticUserMessageStore.setState({ optimisticUserMessage: null });
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
const { result } = renderHook(() => useHandleBuildPlanClick());
// Act & Assert - should not throw
@@ -153,7 +163,10 @@ describe("useHandleBuildPlanClick", () => {
// Assert all expected behaviors still occur
expect(useConversationStore.getState().conversationMode).toBe("code");
expect(mockSend).toHaveBeenCalledTimes(1);
expect(useOptimisticUserMessageStore.getState().optimisticUserMessage).toBe(
const pending =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(pending).toHaveLength(1);
expect(pending[0].text).toBe(
"Execute the plan based on the .agents_tmp/PLAN.md file.",
);
});
@@ -0,0 +1,275 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import {
PENDING_MESSAGE_TIMEOUT_MS,
useOptimisticUserMessageStore,
} from "#/stores/optimistic-user-message-store";
const CONVO = "conv-a";
describe("optimistic-user-message-store", () => {
beforeEach(() => {
vi.useFakeTimers();
useOptimisticUserMessageStore.setState({ pendingMessages: [] });
});
afterEach(() => {
vi.useRealTimers();
});
it("enqueues new messages with status 'sending' and tags them with conversationId", () => {
const store = useOptimisticUserMessageStore.getState();
const id = store.enqueuePendingMessage({
conversationId: CONVO,
text: "hello",
});
const pending = useOptimisticUserMessageStore.getState().pendingMessages;
expect(pending).toHaveLength(1);
expect(pending[0].id).toBe(id);
expect(pending[0].conversationId).toBe(CONVO);
expect(pending[0].text).toBe("hello");
expect(pending[0].status).toBe("sending");
expect(pending[0].imageUrls).toEqual([]);
expect(pending[0].fileUrls).toEqual([]);
expect(typeof pending[0].timestamp).toBe("string");
});
it("preserves FIFO order across multiple enqueues", () => {
const store = useOptimisticUserMessageStore.getState();
store.enqueuePendingMessage({ conversationId: CONVO, text: "first" });
store.enqueuePendingMessage({ conversationId: CONVO, text: "second" });
store.enqueuePendingMessage({ conversationId: CONVO, text: "third" });
const pending = useOptimisticUserMessageStore.getState().pendingMessages;
expect(pending.map((m) => m.text)).toEqual(["first", "second", "third"]);
});
it("marks a pending message as 'error' with details", () => {
const store = useOptimisticUserMessageStore.getState();
const id = store.enqueuePendingMessage({
conversationId: CONVO,
text: "broken",
});
store.markPendingMessageError(id, "boom");
const [entry] = useOptimisticUserMessageStore.getState().pendingMessages;
expect(entry.status).toBe("error");
expect(entry.errorMessage).toBe("boom");
});
it("flips an errored message back to 'sending' on retry", () => {
const store = useOptimisticUserMessageStore.getState();
const id = store.enqueuePendingMessage({
conversationId: CONVO,
text: "broken",
});
store.markPendingMessageError(id, "boom");
store.markPendingMessageSending(id);
const [entry] = useOptimisticUserMessageStore.getState().pendingMessages;
expect(entry.status).toBe("sending");
expect(entry.errorMessage).toBeUndefined();
});
it("enqueue stores `content` separately from `text` and defaults it to `text`", () => {
const store = useOptimisticUserMessageStore.getState();
const idA = store.enqueuePendingMessage({
conversationId: CONVO,
text: "hello",
});
const idB = store.enqueuePendingMessage({
conversationId: CONVO,
text: "hello",
content: "hello\n\nFiles: foo.txt",
});
const pending = useOptimisticUserMessageStore.getState().pendingMessages;
const a = pending.find((m) => m.id === idA)!;
const b = pending.find((m) => m.id === idB)!;
expect(a.content).toBe("hello");
expect(b.text).toBe("hello");
expect(b.content).toBe("hello\n\nFiles: foo.txt");
});
it("consumeMatchingPendingMessage prefers an exact content match (out-of-order echo)", () => {
const store = useOptimisticUserMessageStore.getState();
const firstId = store.enqueuePendingMessage({
conversationId: CONVO,
text: "first",
});
const secondId = store.enqueuePendingMessage({
conversationId: CONVO,
text: "second",
});
// Echo for "second" arrives before "first" — must pop "second", not the
// oldest entry. This is the case the previous FIFO-only implementation
// got wrong.
const consumed = store.consumeMatchingPendingMessage(CONVO, "second");
expect(consumed?.id).toBe(secondId);
const remaining =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(remaining).toHaveLength(1);
expect(remaining[0].id).toBe(firstId);
});
it("consumeMatchingPendingMessage falls back to oldest sending entry when no exact match exists", () => {
const store = useOptimisticUserMessageStore.getState();
const firstId = store.enqueuePendingMessage({
conversationId: CONVO,
text: "hello",
});
store.enqueuePendingMessage({ conversationId: CONVO, text: "world" });
// Server munged the echo (e.g., trimmed whitespace). FIFO fallback keeps
// the bubble from getting stuck.
const consumed = store.consumeMatchingPendingMessage(
CONVO,
"something else",
);
expect(consumed?.id).toBe(firstId);
expect(
useOptimisticUserMessageStore.getState().pendingMessages,
).toHaveLength(1);
});
it("consumeMatchingPendingMessage skips entries already in 'error' state", () => {
const store = useOptimisticUserMessageStore.getState();
const firstId = store.enqueuePendingMessage({
conversationId: CONVO,
text: "first",
});
const secondId = store.enqueuePendingMessage({
conversationId: CONVO,
text: "second",
});
store.markPendingMessageError(firstId, "boom");
const consumed = store.consumeMatchingPendingMessage(CONVO, "second");
expect(consumed?.id).toBe(secondId);
const remaining =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(remaining).toHaveLength(1);
expect(remaining[0].id).toBe(firstId);
expect(remaining[0].status).toBe("error");
});
it("consumeMatchingPendingMessage is a no-op when nothing is sending", () => {
const store = useOptimisticUserMessageStore.getState();
const id = store.enqueuePendingMessage({
conversationId: CONVO,
text: "broken",
});
store.markPendingMessageError(id, "boom");
const consumed = store.consumeMatchingPendingMessage(CONVO, "broken");
expect(consumed).toBeNull();
expect(
useOptimisticUserMessageStore.getState().pendingMessages,
).toHaveLength(1);
});
it("consumeMatchingPendingMessage only consumes entries for the given conversation", () => {
const store = useOptimisticUserMessageStore.getState();
const aId = store.enqueuePendingMessage({
conversationId: "conv-a",
text: "shared",
});
const bId = store.enqueuePendingMessage({
conversationId: "conv-b",
text: "shared",
});
// A cross-conversation ack for conv-b — even with identical content,
// must not pop conv-a's pending entry.
const consumed = store.consumeMatchingPendingMessage("conv-b", "shared");
expect(consumed?.id).toBe(bId);
const remaining =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(remaining).toHaveLength(1);
expect(remaining[0].id).toBe(aId);
});
it("enqueuePendingMessage flips the entry to 'error' after the watchdog timeout", () => {
const store = useOptimisticUserMessageStore.getState();
const id = store.enqueuePendingMessage({
conversationId: CONVO,
text: "stuck",
});
// Still sending right after enqueue.
expect(
useOptimisticUserMessageStore.getState().pendingMessages[0].status,
).toBe("sending");
// Fire the watchdog.
vi.advanceTimersByTime(PENDING_MESSAGE_TIMEOUT_MS);
const [entry] = useOptimisticUserMessageStore.getState().pendingMessages;
expect(entry.id).toBe(id);
expect(entry.status).toBe("error");
expect(entry.errorMessage).toBe("Send timed out");
});
it("watchdog timeout does nothing if the echo already consumed the message", () => {
const store = useOptimisticUserMessageStore.getState();
store.enqueuePendingMessage({ conversationId: CONVO, text: "fast" });
store.consumeMatchingPendingMessage(CONVO, "fast");
vi.advanceTimersByTime(PENDING_MESSAGE_TIMEOUT_MS);
expect(
useOptimisticUserMessageStore.getState().pendingMessages,
).toHaveLength(0);
});
it("watchdog timeout does nothing if the message already failed via send error", () => {
const store = useOptimisticUserMessageStore.getState();
const id = store.enqueuePendingMessage({
conversationId: CONVO,
text: "explicit-error",
});
store.markPendingMessageError(id, "boom");
vi.advanceTimersByTime(PENDING_MESSAGE_TIMEOUT_MS);
const [entry] = useOptimisticUserMessageStore.getState().pendingMessages;
// Should keep the original error message, not get overwritten to "Send timed out".
expect(entry.errorMessage).toBe("boom");
});
it("removePendingMessage drops a specific entry by id", () => {
const store = useOptimisticUserMessageStore.getState();
const firstId = store.enqueuePendingMessage({
conversationId: CONVO,
text: "first",
});
store.enqueuePendingMessage({ conversationId: CONVO, text: "second" });
store.removePendingMessage(firstId);
const remaining =
useOptimisticUserMessageStore.getState().pendingMessages;
expect(remaining.map((m) => m.text)).toEqual(["second"]);
});
it("clearPendingMessages wipes the queue", () => {
const store = useOptimisticUserMessageStore.getState();
store.enqueuePendingMessage({ conversationId: CONVO, text: "first" });
store.enqueuePendingMessage({ conversationId: CONVO, text: "second" });
store.clearPendingMessages();
expect(
useOptimisticUserMessageStore.getState().pendingMessages,
).toHaveLength(0);
});
});
@@ -1,8 +1,6 @@
import React from "react";
import { OpenHandsEvent } from "#/types/agent-server/core";
import { EventMessage } from "./event-message";
import { ChatMessage } from "../../features/chat/chat-message";
import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-store";
import { usePlanPreviewEvents } from "./hooks/use-plan-preview-events";
import { groupEvents } from "./group-events";
import { EventGroup } from "./event-message-components/event-group";
@@ -20,10 +18,6 @@ const getLastEventId = (events: OpenHandsEvent[]) => events.at(-1)?.id;
export const Messages: React.FC<MessagesProps> = React.memo(
({ messages, allEvents }) => {
const { getOptimisticUserMessage } = useOptimisticUserMessageStore();
const optimisticUserMessage = getOptimisticUserMessage();
// Get the set of event IDs that should render PlanPreview
// This ensures only one preview per user message "phase"
const planPreviewEventIds = usePlanPreviewEvents(allEvents);
@@ -82,10 +76,6 @@ export const Messages: React.FC<MessagesProps> = React.memo(
</EventGroup>
);
})}
{optimisticUserMessage && (
<ChatMessage type="user" message={optimisticUserMessage} />
)}
</>
);
},
+52 -13
View File
@@ -25,6 +25,7 @@ import { useErrorMessageStore } from "#/stores/error-message-store";
import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-store";
import { ErrorMessageBanner } from "./error-message-banner";
import { Messages } from "#/components/conversation-events/chat/messages";
import { PendingUserMessages } from "./pending-user-messages";
import { useUnifiedUploadFiles } from "#/hooks/mutation/use-unified-upload-files";
import { validateFiles } from "#/utils/file-validation";
import { useConversationStore } from "#/stores/conversation-store";
@@ -61,8 +62,15 @@ export function ChatInterface() {
hasSubstantiveAgentActions,
userEventsExist,
} = useFilteredEvents();
const { setOptimisticUserMessage, getOptimisticUserMessage } =
useOptimisticUserMessageStore();
const enqueuePendingMessage = useOptimisticUserMessageStore(
(state) => state.enqueuePendingMessage,
);
const markPendingMessageError = useOptimisticUserMessageStore(
(state) => state.markPendingMessageError,
);
const pendingMessages = useOptimisticUserMessageStore(
(state) => state.pendingMessages,
);
const { t } = useTranslation("openhands");
const scrollRef = React.useRef<HTMLDivElement>(null);
const {
@@ -178,7 +186,7 @@ export function ChatInterface() {
[maybeLoadOlder],
);
const optimisticUserMessage = getOptimisticUserMessage();
const hasPendingUserMessages = pendingMessages.length > 0;
// Show V1 messages immediately if events exist in store (e.g., remount),
// or once loading completes. This replaces the old transition-observation
@@ -260,15 +268,36 @@ export function ChatInterface() {
const prompt =
uploadedFiles.length > 0 ? `${content}\n\n${filePrompt}` : content;
const result = await send(
createChatMessage(prompt, imageUrls, uploadedFiles, timestamp),
);
// Only show optimistic UI if message was sent immediately via WebSocket
// If queued for later delivery, the message will appear when actually delivered
if (!result.queued) {
setOptimisticUserMessage(content);
}
// Enqueue the message into the local pending queue with status "sending"
// so the user immediately sees it in the chat with a faded treatment. The
// entry is removed when the WebSocket echoes back the corresponding
// `UserMessageEvent`. If the API call to send the message fails, the entry
// is flipped to "error" with a retry link.
const pendingId = enqueuePendingMessage({
conversationId: conversationId!,
// `text` is what the user sees in the bubble; `content` is what we
// actually hand to the server (the prompt may include an appended
// "Files uploaded: …" block) and is what the echo will be matched
// against. They're different when there are file attachments.
text: content,
content: prompt,
imageUrls,
fileUrls: uploadedFiles,
timestamp,
});
setMessageToSend("");
try {
await send(
createChatMessage(prompt, imageUrls, uploadedFiles, timestamp),
);
} catch (sendError) {
const sendErrorMessage =
sendError instanceof Error
? sendError.message
: "Failed to send message";
markPendingMessageError(pendingId, sendErrorMessage);
}
};
// Auto-scroll to bottom when new messages arrive — but only if the user is
@@ -296,7 +325,7 @@ export function ChatInterface() {
// Note: We intentionally exclude autoScroll from deps because we only want
// to scroll when message content changes, not when autoScroll state changes.
// eslint-disable-next-line react-hooks/exhaustive-deps
}, [renderableEvents.length, optimisticUserMessage, scrollDomToBottom]);
}, [renderableEvents.length, hasPendingUserMessages, scrollDomToBottom]);
// Auto-load older events when the chat content doesn't overflow the
// scroll area (no scrollbar to drag, no wheel events past 0). We
@@ -358,7 +387,7 @@ export function ChatInterface() {
<ScrollProvider value={scrollProviderValue}>
<div className="h-full flex flex-col justify-between pr-0 md:pr-4 relative">
{!hasSubstantiveAgentActions &&
!optimisticUserMessage &&
!hasPendingUserMessages &&
!userEventsExist &&
!isChatLoading && (
<ChatSuggestions
@@ -412,6 +441,16 @@ export function ChatInterface() {
allEvents={allConversationEvents}
/>
)}
{/*
Render the local pending-message queue independently so messages
the user just submitted show up immediately (with a faded "sending"
treatment) even before any real conversation event has come back
from the server. Entries drain (FIFO) when the matching
UserMessageEvent echoes back over the WebSocket, so this never
double-renders alongside the real event list.
*/}
<PendingUserMessages />
</div>
<div className="flex flex-col gap-[6px]">
@@ -1,10 +1,14 @@
import React from "react";
import { useTranslation } from "react-i18next";
import { cn } from "#/utils/utils";
import { CopyToClipboardButton } from "#/components/shared/buttons/copy-to-clipboard-button";
import type { SourceType } from "#/types/agent-server/core/base/common";
import { StyledTooltip } from "#/components/shared/buttons/styled-tooltip";
import { I18nKey } from "#/i18n/declaration";
import { MarkdownRenderer } from "../markdown/markdown-renderer";
export type ChatMessagePendingStatus = "sending" | "error";
interface ChatMessageProps {
type: SourceType;
message: string;
@@ -14,6 +18,8 @@ interface ChatMessageProps {
tooltip?: string;
}>;
isFromPlanningAgent?: boolean;
pendingStatus?: ChatMessagePendingStatus;
onRetry?: () => void;
}
export function ChatMessage({
@@ -22,7 +28,10 @@ export function ChatMessage({
children,
actions,
isFromPlanningAgent = false,
pendingStatus,
onRetry,
}: React.PropsWithChildren<ChatMessageProps>) {
const { t } = useTranslation("openhands");
const [isHovering, setIsHovering] = React.useState(false);
const [isCopy, setIsCopy] = React.useState(false);
@@ -48,6 +57,7 @@ export function ChatMessage({
return (
<article
data-testid={`${type}-message`}
data-pending-status={pendingStatus}
onMouseEnter={() => setIsHovering(true)}
onMouseLeave={() => setIsHovering(false)}
className={cn(
@@ -58,6 +68,8 @@ export function ChatMessage({
isFromPlanningAgent &&
type === "agent" &&
"border border-[#597ff4] bg-tertiary p-4 mt-2",
pendingStatus === "sending" && "opacity-60",
pendingStatus === "error" && "border border-status-fail-border",
)}
>
<div
@@ -110,6 +122,37 @@ export function ChatMessage({
<MarkdownRenderer includeStandard>{message}</MarkdownRenderer>
</div>
{pendingStatus === "sending" && (
<span
role="status"
aria-live="polite"
data-testid="chat-message-sending"
className="self-end text-xs italic text-content-muted"
>
{t(I18nKey.CHAT_INTERFACE$MESSAGE_SENDING)}
</span>
)}
{pendingStatus === "error" && (
<span
role="alert"
data-testid="chat-message-error"
className="self-end text-xs text-status-fail-text"
>
{t(I18nKey.CHAT_INTERFACE$MESSAGE_SEND_FAILED)}{" "}
{onRetry && (
<button
type="button"
onClick={onRetry}
className="underline cursor-pointer"
data-testid="chat-message-retry"
>
{t(I18nKey.CHAT_INTERFACE$MESSAGE_RETRY)}
</button>
)}
</span>
)}
{children}
</article>
);
@@ -1,4 +1,4 @@
import React, { useEffect } from "react";
import React, { useEffect, useRef } from "react";
import { useChatInputLogic } from "#/hooks/chat/use-chat-input-logic";
import { useFileHandling } from "#/hooks/chat/use-file-handling";
import { useGripResize } from "#/hooks/chat/use-grip-resize";
@@ -46,14 +46,27 @@ export function CustomChatInput({
// immediately via the WebSocket if connected, or queued via REST otherwise.
const isDisabled = disabled;
// Always call the latest `onSubmit` without making the effect re-run when
// its identity changes. `onSubmit` (typically `handleSendMessage`) is a
// fresh function on every parent render, and the parent re-renders
// whenever the pending-message queue updates synchronously inside
// `onSubmit` itself. Listing it in the dep array caused the effect to
// fire twice — once for the original submit and again from the
// mid-submit re-render, before `setSubmittedMessage(null)` was applied —
// producing a duplicate "Sending…" bubble.
const onSubmitRef = useRef(onSubmit);
useEffect(() => {
onSubmitRef.current = onSubmit;
}, [onSubmit]);
// Listen to submittedMessage state changes
useEffect(() => {
if (!submittedMessage || disabled) {
return;
}
onSubmit(submittedMessage);
onSubmitRef.current(submittedMessage);
setSubmittedMessage(null);
}, [submittedMessage, disabled, onSubmit, setSubmittedMessage]);
}, [submittedMessage, disabled, setSubmittedMessage]);
// Custom hooks
const {
@@ -30,7 +30,12 @@ export function GitControlBar({ onSuggestionsClick }: GitControlBarProps) {
const { conversationId } = useConversationId();
const [isOpenRepoModalOpen, setIsOpenRepoModalOpen] = useState(false);
const { addRecentRepository } = useHomeStore();
const { setOptimisticUserMessage } = useOptimisticUserMessageStore();
const enqueuePendingMessage = useOptimisticUserMessageStore(
(state) => state.enqueuePendingMessage,
);
const markPendingMessageError = useOptimisticUserMessageStore(
(state) => state.markPendingMessageError,
);
const { data: conversation } = useActiveConversation();
const { repositoryInfo } = useTaskPolling();
@@ -112,13 +117,25 @@ export function GitControlBar({ onSuggestionsClick }: GitControlBarProps) {
repository.git_provider.charAt(0).toUpperCase() +
repository.git_provider.slice(1);
const clonePrompt = `Clone ${repository.full_name} from ${providerName} and checkout branch ${branch.name}.`;
setOptimisticUserMessage(clonePrompt);
sendRef.current({
action: "message",
args: {
content: clonePrompt,
timestamp: new Date().toISOString(),
},
const pendingId = conversationId
? enqueuePendingMessage({ conversationId, text: clonePrompt })
: null;
// `send` returns a Promise; surface a failed send by flipping the
// matching pending entry to "error" so the user gets the retry link
// rather than a perpetual "Sending…" bubble.
Promise.resolve(
sendRef.current({
action: "message",
args: {
content: clonePrompt,
timestamp: new Date().toISOString(),
},
}),
).catch((error) => {
if (!pendingId) return;
const errorMessage =
error instanceof Error ? error.message : "Failed to send message";
markPendingMessageError(pendingId, errorMessage);
});
},
},
@@ -0,0 +1,91 @@
import React from "react";
import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-store";
import { useSendMessage } from "#/hooks/use-send-message";
import { createChatMessage } from "#/services/chat-service";
import { useOptionalConversationId } from "#/hooks/use-conversation-id";
import { ChatMessage } from "./chat-message";
/**
* Renders the queue of locally-tracked user messages that have been submitted
* but not yet echoed back through the WebSocket. Each message shows a faded
* "sending" treatment until the server echoes a real `UserMessageEvent`
* (which removes it via `consumeMatchingPendingMessage`). If the API rejects the
* send, the message switches to an "error" state with a retry button.
*
* The queue is global but each entry is tagged with the conversation id it
* was enqueued from; this component filters to only entries belonging to the
* active conversation, so switching conversations never carries pending
* bubbles over.
*/
export function PendingUserMessages() {
const { conversationId } = useOptionalConversationId();
const pendingMessages = useOptimisticUserMessageStore(
(state) => state.pendingMessages,
);
const markPendingMessageError = useOptimisticUserMessageStore(
(state) => state.markPendingMessageError,
);
const markPendingMessageSending = useOptimisticUserMessageStore(
(state) => state.markPendingMessageSending,
);
const { send } = useSendMessage();
const visibleMessages = React.useMemo(
() =>
conversationId
? pendingMessages.filter(
(message) => message.conversationId === conversationId,
)
: [],
[pendingMessages, conversationId],
);
const handleRetry = React.useCallback(
async (id: string) => {
const message = useOptimisticUserMessageStore
.getState()
.pendingMessages.find((entry) => entry.id === id);
if (!message) return;
markPendingMessageSending(id);
try {
await send(
createChatMessage(
message.text,
message.imageUrls,
message.fileUrls,
message.timestamp,
),
);
} catch (error) {
const errorMessage =
error instanceof Error ? error.message : "Failed to send message";
markPendingMessageError(id, errorMessage);
}
},
[send, markPendingMessageError, markPendingMessageSending],
);
if (visibleMessages.length === 0) {
return null;
}
return (
<>
{visibleMessages.map((message) => (
<ChatMessage
key={message.id}
type="user"
message={message.text}
pendingStatus={message.status}
onRetry={
message.status === "error"
? () => handleRetry(message.id)
: undefined
}
/>
))}
</>
);
}
@@ -21,16 +21,16 @@ interface TaskCardProps {
}
export function TaskCard({ task }: TaskCardProps) {
const { setOptimisticUserMessage } = useOptimisticUserMessageStore();
const enqueuePendingMessage = useOptimisticUserMessageStore(
(state) => state.enqueuePendingMessage,
);
const { mutate: createConversation } = useCreateConversation();
const isCreatingConversation = useIsCreatingConversation();
const { t } = useTranslation("openhands");
const { navigate } = useNavigation();
const handleLaunchConversation = () => {
setOptimisticUserMessage(t("TASK$ADDRESSING_TASK"));
return createConversation(
const handleLaunchConversation = () =>
createConversation(
{
repository: {
name: task.repo,
@@ -40,11 +40,18 @@ export function TaskCard({ task }: TaskCardProps) {
},
{
onSuccess: (data) => {
// Enqueue the pending message after the new conversation exists so
// it can be tagged with the conversation id and only show up in the
// target conversation's chat (not in whatever convo the user was
// looking at when they clicked the task).
enqueuePendingMessage({
conversationId: data.conversation_id,
text: t("TASK$ADDRESSING_TASK"),
});
navigate(`/conversations/${data.conversation_id}`);
},
},
);
};
// Determine the correct URL format based on git provider
let href: string;
+41 -10
View File
@@ -73,6 +73,24 @@ const ConversationWebSocketContext = createContext<
ConversationWebSocketContextType | undefined
>(undefined);
/**
* Extract the text body of an echoed user `MessageEvent` for matching against
* the optimistic pending-message queue. The server wraps the original
* `args.content` string in one or more `TextContent` entries (alongside any
* `ImageContent` entries for inline images), so concatenating the `text`
* fields gives us back the exact prompt we sent.
*/
function extractMessageEventText(
event: import("#/types/agent-server/core/events/message-event").MessageEvent,
): string {
return event.llm_message.content
.filter(
(part): part is { type: "text"; text: string } => part.type === "text",
)
.map((part) => part.text)
.join("");
}
export function ConversationWebSocketProvider({
children,
conversationId,
@@ -104,7 +122,9 @@ export function ConversationWebSocketProvider({
const addEvent = useEventStore((state) => state.addEvent);
const addEvents = useEventStore((state) => state.addEvents);
const { setErrorMessage, removeErrorMessage } = useErrorMessageStore();
const { removeOptimisticUserMessage } = useOptimisticUserMessageStore();
const consumeMatchingPendingMessage = useOptimisticUserMessageStore(
(state) => state.consumeMatchingPendingMessage,
);
const { setExecutionStatus } = useConversationStateStore();
const { appendInput, appendOutput } = useCommandStore();
@@ -403,11 +423,18 @@ export function ConversationWebSocketProvider({
setErrorMessage(event.error);
}
// Clear optimistic user message when a user message is confirmed
// Clear optimistic user message when a user message is confirmed.
// We match by the echoed text content (with FIFO fallback inside the
// store), so an echo for "second" pops "second" — not whichever
// pending entry happens to be oldest — protecting against any
// out-of-order delivery between conversations or sub-agents.
if (isUserMessageEvent(event)) {
removeOptimisticUserMessage();
// Clear draft from localStorage - message was successfully delivered
if (conversationId) {
consumeMatchingPendingMessage(
conversationId,
extractMessageEventText(event),
);
// Clear draft from localStorage - message was successfully delivered
setConversationState(conversationId, { draftMessage: null });
}
}
@@ -476,7 +503,7 @@ export function ConversationWebSocketProvider({
[
addEvent,
setErrorMessage,
removeOptimisticUserMessage,
consumeMatchingPendingMessage,
queryClient,
conversationId,
setExecutionStatus,
@@ -550,12 +577,16 @@ export function ConversationWebSocketProvider({
setErrorMessage(event.error);
}
// Clear optimistic user message when a user message is confirmed
// Clear optimistic user message when a user message is confirmed.
// Always scope to the main `conversationId` (where the user types)
// and match on the echoed content so the planning sub-agent's own
// events can never consume a main-conversation pending entry.
if (isUserMessageEvent(event)) {
removeOptimisticUserMessage();
// Clear draft from localStorage - message was successfully delivered
// Use main conversationId since user types in main conversation input
if (conversationId) {
consumeMatchingPendingMessage(
conversationId,
extractMessageEventText(event),
);
setConversationState(conversationId, { draftMessage: null });
}
}
@@ -648,7 +679,7 @@ export function ConversationWebSocketProvider({
isLoadingHistoryPlanning,
expectedEventCountPlanning,
setErrorMessage,
removeOptimisticUserMessage,
consumeMatchingPendingMessage,
queryClient,
subConversations,
conversationId,
+31 -5
View File
@@ -3,6 +3,7 @@ import { useConversationStore } from "#/stores/conversation-store";
import { useSendMessage } from "#/hooks/use-send-message";
import { createChatMessage } from "#/services/chat-service";
import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-store";
import { useOptionalConversationId } from "#/hooks/use-conversation-id";
/**
* Custom hook that encapsulates the logic for handling the Build button click.
@@ -13,7 +14,13 @@ import { useOptimisticUserMessageStore } from "#/stores/optimistic-user-message-
export const useHandleBuildPlanClick = () => {
const { setConversationMode } = useConversationStore();
const { send } = useSendMessage();
const { setOptimisticUserMessage } = useOptimisticUserMessageStore();
const { conversationId } = useOptionalConversationId();
const enqueuePendingMessage = useOptimisticUserMessageStore(
(state) => state.enqueuePendingMessage,
);
const markPendingMessageError = useOptimisticUserMessageStore(
(state) => state.markPendingMessageError,
);
const handleBuildPlanClick = useCallback(
(event?: React.MouseEvent<HTMLButtonElement> | KeyboardEvent) => {
@@ -26,12 +33,31 @@ export const useHandleBuildPlanClick = () => {
// Create the build prompt to execute the plan
const buildPrompt = `Execute the plan based on the .agents_tmp/PLAN.md file.`;
// Send the message to the code agent
// Show the prompt as a pending message and send it to the code agent.
// Skip the pending bubble if we somehow don't have a conversation id —
// the message still gets sent, just without the optimistic queue entry.
const timestamp = new Date().toISOString();
send(createChatMessage(buildPrompt, [], [], timestamp));
setOptimisticUserMessage(buildPrompt);
const pendingId = conversationId
? enqueuePendingMessage({
conversationId,
text: buildPrompt,
timestamp,
})
: null;
send(createChatMessage(buildPrompt, [], [], timestamp)).catch((error) => {
if (!pendingId) return;
const errorMessage =
error instanceof Error ? error.message : "Failed to send message";
markPendingMessageError(pendingId, errorMessage);
});
},
[setConversationMode, send, setOptimisticUserMessage],
[
setConversationMode,
send,
conversationId,
enqueuePendingMessage,
markPendingMessageError,
],
);
return { handleBuildPlanClick };
+51
View File
@@ -8176,6 +8176,57 @@
"uk": "Налаштування оновлено",
"ca": "Configuració actualitzada"
},
"CHAT_INTERFACE$MESSAGE_SENDING": {
"en": "Sending...",
"ja": "送信中...",
"zh-CN": "发送中...",
"zh-TW": "傳送中...",
"ko-KR": "보내는 중...",
"no": "Sender...",
"it": "Invio in corso...",
"pt": "Enviando...",
"es": "Enviando...",
"ar": "جاري الإرسال...",
"fr": "Envoi en cours...",
"tr": "Gönderiliyor...",
"de": "Wird gesendet...",
"uk": "Надсилання...",
"ca": "Enviant..."
},
"CHAT_INTERFACE$MESSAGE_SEND_FAILED": {
"en": "Failed to send.",
"ja": "送信に失敗しました。",
"zh-CN": "发送失败。",
"zh-TW": "傳送失敗。",
"ko-KR": "보내기 실패.",
"no": "Sending mislyktes.",
"it": "Invio non riuscito.",
"pt": "Falha ao enviar.",
"es": "Error al enviar.",
"ar": "فشل الإرسال.",
"fr": "Échec de l'envoi.",
"tr": "Gönderilemedi.",
"de": "Senden fehlgeschlagen.",
"uk": "Не вдалося надіслати.",
"ca": "Error en enviar."
},
"CHAT_INTERFACE$MESSAGE_RETRY": {
"en": "Retry",
"ja": "再試行",
"zh-CN": "重试",
"zh-TW": "重試",
"ko-KR": "다시 시도",
"no": "Prøv igjen",
"it": "Riprova",
"pt": "Tentar novamente",
"es": "Reintentar",
"ar": "إعادة المحاولة",
"fr": "Réessayer",
"tr": "Yeniden dene",
"de": "Wiederholen",
"uk": "Повторити",
"ca": "Torna-ho a provar"
},
"CHAT_INTERFACE$AUGMENTED_PROMPT_FILES_TITLE": {
"en": "NEW FILES ADDED",
"de": "NEUE DATEIEN HINZUGEFÜGT",
+171 -13
View File
@@ -1,36 +1,194 @@
import { create } from "zustand";
export type PendingUserMessageStatus = "sending" | "error";
/**
* How long a pending message is allowed to stay in "sending" state before we
* give up and flip it to "error" with a retry link. This guards against the
* "server crashed / websocket dropped after our send resolved, echo never
* arrives" scenario where the message would otherwise hang forever.
*
* Exported so tests can override it via vi.fakeTimers without hard-coding the
* value.
*/
export const PENDING_MESSAGE_TIMEOUT_MS = 60_000;
export interface PendingUserMessage {
id: string;
/**
* The conversation this pending message belongs to. The chat UI filters the
* global queue by the active conversation id so messages enqueued in one
* conversation never leak into another when the user switches.
*/
conversationId: string;
/** User-visible bubble text (what the user typed; no file annotations). */
text: string;
/**
* The exact string sent to the server (may include the appended
* "Files uploaded: …" prompt when attachments are present). Used as the
* primary key when matching against the echoed `UserMessageEvent`.
*/
content: string;
status: PendingUserMessageStatus;
imageUrls: string[];
fileUrls: string[];
timestamp: string;
errorMessage?: string;
}
interface OptimisticUserMessageState {
optimisticUserMessage: string | null;
pendingMessages: PendingUserMessage[];
}
export interface EnqueuePendingMessagePayload {
conversationId: string;
/** User-visible text for the bubble. */
text: string;
/**
* The exact string sent to the server. Defaults to `text` for call sites
* that don't transform the content (e.g. git-control-bar, task-card).
*/
content?: string;
imageUrls?: string[];
fileUrls?: string[];
timestamp?: string;
}
interface OptimisticUserMessageActions {
setOptimisticUserMessage: (message: string) => void;
getOptimisticUserMessage: () => string | null;
removeOptimisticUserMessage: () => void;
/**
* Append a new user message to the queue with status "sending".
* Returns the locally-generated id for later updates. Schedules a
* `PENDING_MESSAGE_TIMEOUT_MS` watchdog that flips the entry to "error" if
* it's still in "sending" state when the timer fires.
*/
enqueuePendingMessage: (payload: EnqueuePendingMessagePayload) => string;
/** Mark a pending message as failed (the API rejected it). */
markPendingMessageError: (id: string, errorMessage?: string) => void;
/** Mark a pending message as sending again (used when retrying). */
markPendingMessageSending: (id: string) => void;
/** Drop a pending message from the queue (e.g., after success/cancellation). */
removePendingMessage: (id: string) => void;
/**
* Remove the pending message that matches the given echoed `content` in
* the given conversation. Matching is done by exact content equality on
* messages still in "sending" state; if no match exists we fall back to
* removing the oldest "sending" entry in that conversation so that an echo
* with a slightly munged body (e.g. trailing-whitespace stripped by the
* server) still clears its bubble. Scoping by `conversationId` ensures a
* stale ack for one conversation never pops a pending entry belonging to
* another.
*/
consumeMatchingPendingMessage: (
conversationId: string,
content: string,
) => PendingUserMessage | null;
/** Wipe all queued messages (e.g., when changing conversations). */
clearPendingMessages: () => void;
}
type OptimisticUserMessageStore = OptimisticUserMessageState &
OptimisticUserMessageActions;
const initialState: OptimisticUserMessageState = {
optimisticUserMessage: null,
pendingMessages: [],
};
// Use a timestamp + random suffix instead of a module-level counter so ids
// stay unique across test resets and don't accumulate state between runs.
// `crypto.randomUUID` would be ideal but isn't available in older test
// environments, so a base36 random suffix is a safe lowest-common-denominator.
const generatePendingId = (): string =>
`pending-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`;
export const useOptimisticUserMessageStore = create<OptimisticUserMessageStore>(
(set, get) => ({
...initialState,
setOptimisticUserMessage: (message: string) =>
set(() => ({
optimisticUserMessage: message,
enqueuePendingMessage: (payload) => {
const id = generatePendingId();
const message: PendingUserMessage = {
id,
conversationId: payload.conversationId,
text: payload.text,
content: payload.content ?? payload.text,
status: "sending",
imageUrls: payload.imageUrls ?? [],
fileUrls: payload.fileUrls ?? [],
timestamp: payload.timestamp ?? new Date().toISOString(),
};
set((state) => ({
pendingMessages: [...state.pendingMessages, message],
}));
// Watchdog: if the server echo never lands (WS dropped, server crashed,
// network partition), flip this entry to "error" so the user gets a
// retry link instead of a permanently-pinned "Sending…" bubble.
setTimeout(() => {
const current = get().pendingMessages.find((m) => m.id === id);
if (current && current.status === "sending") {
get().markPendingMessageError(id, "Send timed out");
}
}, PENDING_MESSAGE_TIMEOUT_MS);
return id;
},
markPendingMessageError: (id, errorMessage) =>
set((state) => ({
pendingMessages: state.pendingMessages.map((message) =>
message.id === id
? { ...message, status: "error", errorMessage }
: message,
),
})),
getOptimisticUserMessage: () => get().optimisticUserMessage,
removeOptimisticUserMessage: () =>
set(() => ({
optimisticUserMessage: null,
markPendingMessageSending: (id) =>
set((state) => ({
pendingMessages: state.pendingMessages.map((message) =>
message.id === id
? { ...message, status: "sending", errorMessage: undefined }
: message,
),
})),
removePendingMessage: (id) =>
set((state) => ({
pendingMessages: state.pendingMessages.filter(
(message) => message.id !== id,
),
})),
consumeMatchingPendingMessage: (conversationId, content) => {
// Single atomic `set` so the find + filter can't observe an interleaved
// mutation from another action. We prefer an exact content match (this
// is what makes out-of-order echoes safe: an echo of "world" will pop
// the "world" bubble, not the older "hello" one). If no exact match
// exists — e.g. the server slightly munged the body — fall back to the
// oldest "sending" entry in this conversation so the user doesn't end
// up with a permanently-stuck bubble in the happy-path single-message
// case.
let consumed: PendingUserMessage | null = null;
set((state) => {
const sending = state.pendingMessages
.map((m, i) => ({ m, i }))
.filter(
({ m }) =>
m.status === "sending" && m.conversationId === conversationId,
);
if (sending.length === 0) return state;
const exact = sending.find(({ m }) => m.content === content);
const target = exact ?? sending[0];
consumed = target.m;
return {
pendingMessages: [
...state.pendingMessages.slice(0, target.i),
...state.pendingMessages.slice(target.i + 1),
],
};
});
return consumed;
},
clearPendingMessages: () => set(() => ({ ...initialState })),
}),
);