FN-5790: harden room chat send handling and optimistic reconciliation

Prevent duplicate room sends and avoid restoring composer text when ambiguous failures still delivered the message.

- guard room composer sends with an in-flight ref to prevent double-dispatch
- treat ambiguous post-send refresh failures as delivered when recovery confirms the user message
- centralize optimistic/SSE user-message reconciliation to replace temp messages instead of duplicating them
- add and update chat room hook/component tests that cover double-send prevention and delivered-on-ambiguous-failure behavior
- add a patch changeset for @runfusion/fusion documenting the room chat reliability fix

Files changed:
 .changeset/few-dingos-smile.md                     |  5 ++
 packages/dashboard/app/components/ChatView.tsx     |  8 +++
 .../components/__tests__/ChatView.rooms.test.tsx   | 27 ++++++++++
 .../app/hooks/__tests__/useChatRooms.test.ts       | 48 +++++++++++++++--
 packages/dashboard/app/hooks/useChatRooms.ts       | 62 +++++++++++++++-------
 5 files changed, 126 insertions(+), 24 deletions(-)

Fusion-Task-Id: FN-5790

Fusion-Task-Lineage: a1ac5e4f-4cc2-4cf9-8bc1-42587d0d1d1b
This commit is contained in:
gsxdsm
2026-05-31 16:57:42 -07:00
parent 2140ab2dcf
commit 716f3964c5
5 changed files with 126 additions and 24 deletions

View File

@@ -0,0 +1,5 @@
---
"@runfusion/fusion": patch
---
Fix room chat send reliability by preventing concurrent in-flight room dispatches, classifying ambiguous delivered sends as delivered (so composer text is not restored), and hardening optimistic/SSE reconciliation to avoid duplicate user message rendering.

View File

@@ -1070,6 +1070,7 @@ export function ChatView({ projectId, addToast, experimentalFeatures }: ChatView
const pendingAttachmentsRef = useRef<PendingAttachment[]>([]); const pendingAttachmentsRef = useRef<PendingAttachment[]>([]);
const mentionCursorPosRef = useRef(0); const mentionCursorPosRef = useRef(0);
const copyFeedbackTimeoutsRef = useRef<Map<string, number>>(new Map()); const copyFeedbackTimeoutsRef = useRef<Map<string, number>>(new Map());
const roomSendInFlightRef = useRef(false);
const mode = useViewportMode(); const mode = useViewportMode();
const isMobile = mode === "mobile"; const isMobile = mode === "mobile";
@@ -1901,6 +1902,11 @@ export function ChatView({ projectId, addToast, experimentalFeatures }: ChatView
return; return;
} }
if (roomSendInFlightRef.current) {
return;
}
roomSendInFlightRef.current = true;
const previousInput = messageInput; const previousInput = messageInput;
clearComposerState(); clearComposerState();
@@ -1920,6 +1926,8 @@ export function ChatView({ projectId, addToast, experimentalFeatures }: ChatView
? error.message ? error.message
: "Failed to send room message"; : "Failed to send room message";
addToast(message, "error"); addToast(message, "error");
} finally {
roomSendInFlightRef.current = false;
} }
return; return;
} }

View File

@@ -370,6 +370,33 @@ describe("ChatView — rooms (FN-3805..FN-3811 contract)", () => {
expect(addToast).not.toHaveBeenCalledWith(expect.stringMatching(/attach/i), "warning"); expect(addToast).not.toHaveBeenCalledWith(expect.stringMatching(/attach/i), "warning");
}); });
it("blocks concurrent room send dispatches while send is in flight", async () => {
let resolveSend: () => void;
const sendPromise = new Promise<void>((resolve) => {
resolveSend = resolve;
});
const sendRoomMessage = vi.fn().mockReturnValue(sendPromise);
setup({}, { sendRoomMessage, activeRoom: roomA });
render(<ChatView projectId="proj-123" addToast={vi.fn()} experimentalFeatures={{ chatRooms: true }} />);
const textarea = screen.getByTestId("chat-input") as HTMLTextAreaElement;
await userEvent.type(textarea, "single send");
fireEvent.keyDown(textarea, { key: "Enter" });
fireEvent.keyDown(textarea, { key: "Enter" });
await waitFor(() => {
expect(sendRoomMessage).toHaveBeenCalledTimes(1);
expect(sendRoomMessage).toHaveBeenCalledWith("single send", { files: [] });
});
resolveSend!();
await act(async () => {
await sendPromise;
});
});
it("FN-5360 keeps room composer cleared when delivery succeeded but reply generation failed", async () => { it("FN-5360 keeps room composer cleared when delivery succeeded but reply generation failed", async () => {
const addToast = vi.fn(); const addToast = vi.fn();
const sendRoomMessage = vi const sendRoomMessage = vi

View File

@@ -295,6 +295,44 @@ describe("useChatRooms", () => {
expect(result.current.messages.map((message) => message.id)).toEqual(["msg-user", "msg-assistant"]); expect(result.current.messages.map((message) => message.id)).toEqual(["msg-user", "msg-assistant"]);
}); });
it("deduplicates user message across optimistic add, SSE echo, and post resolution interleaving", async () => {
const active = room("room-1", "one", "2026-05-09T01:00:00.000Z");
mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] });
const { result } = renderHook(() => useChatRooms("proj-1"));
await waitFor(() => expect(result.current.rooms.length).toBe(1));
mockFetchChatRoomMembers.mockResolvedValueOnce({ members: [] });
mockFetchChatRoomMessages.mockResolvedValueOnce({ messages: [] });
act(() => result.current.selectRoom("room-1"));
await waitFor(() => expect(result.current.activeRoom?.id).toBe("room-1"));
let resolvePost: ((value: { message: ChatRoomMessage }) => void) | undefined;
const postPromise = new Promise<{ message: ChatRoomMessage }>((resolve) => {
resolvePost = resolve;
});
mockPostChatRoomMessage.mockReturnValueOnce(postPromise);
mockFetchChatRoomMessages.mockResolvedValueOnce({ messages: [roomMessage("msg-user", "room-1", "hello")] });
let sendPromise!: Promise<void>;
await act(async () => {
sendPromise = result.current.sendRoomMessage("hello");
});
act(() => {
capturedEvents["chat:room:message:added"]?.({ data: JSON.stringify(roomMessage("msg-user", "room-1", "hello")) } as MessageEvent);
});
resolvePost?.({ message: roomMessage("msg-user", "room-1", "hello") });
await act(async () => {
await sendPromise;
});
const matchingMessages = result.current.messages.filter((message) => message.role === "user" && message.content === "hello");
expect(matchingMessages).toHaveLength(1);
expect(matchingMessages[0]?.id).toBe("msg-user");
});
it("uploads files before posting room message", async () => { it("uploads files before posting room message", async () => {
const active = room("room-1", "one", "2026-05-09T01:00:00.000Z"); const active = room("room-1", "one", "2026-05-09T01:00:00.000Z");
mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] }); mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] });
@@ -351,7 +389,7 @@ describe("useChatRooms", () => {
expect(mockPostChatRoomMessage).not.toHaveBeenCalledWith("room-1", expect.objectContaining({ content: "hello" }), "proj-1"); expect(mockPostChatRoomMessage).not.toHaveBeenCalledWith("room-1", expect.objectContaining({ content: "hello" }), "proj-1");
}); });
it("rejects with original error when post fails before delivery", async () => { it("rejects with original error when post fails and recovery transcript has no persisted user message", async () => {
const active = room("room-1", "one", "2026-05-09T01:00:00.000Z"); const active = room("room-1", "one", "2026-05-09T01:00:00.000Z");
mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] }); mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] });
const { result } = renderHook(() => useChatRooms("proj-1")); const { result } = renderHook(() => useChatRooms("proj-1"));
@@ -363,7 +401,9 @@ describe("useChatRooms", () => {
await waitFor(() => expect(result.current.activeRoom?.id).toBe("room-1")); await waitFor(() => expect(result.current.activeRoom?.id).toBe("room-1"));
mockPostChatRoomMessage.mockRejectedValueOnce(new Error("POST failed")); mockPostChatRoomMessage.mockRejectedValueOnce(new Error("POST failed"));
mockFetchChatRoomMessages.mockRejectedValueOnce(new Error("refresh failed")); mockFetchChatRoomMessages.mockResolvedValueOnce({
messages: [{ ...roomMessage("msg-assistant", "room-1", "Room reply"), role: "assistant", senderAgentId: "agent-1" }],
});
let postError: unknown; let postError: unknown;
await act(async () => { await act(async () => {
@@ -437,7 +477,7 @@ describe("useChatRooms", () => {
); );
}); });
it("refreshes persisted room messages even when room reply generation fails", async () => { const active = room("room-1", "one", "2026-05-09T01:00:00.000Z"); it("classifies post rejection as delivered when recovery transcript includes persisted user message", async () => { const active = room("room-1", "one", "2026-05-09T01:00:00.000Z");
mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] }); mockFetchChatRooms.mockResolvedValueOnce({ rooms: [active] });
const { result } = renderHook(() => useChatRooms("proj-1")); const { result } = renderHook(() => useChatRooms("proj-1"));
await waitFor(() => expect(result.current.rooms.length).toBe(1)); await waitFor(() => expect(result.current.rooms.length).toBe(1));
@@ -452,7 +492,7 @@ describe("useChatRooms", () => {
mockFetchChatRoomMessages.mockResolvedValueOnce({ messages: [persistedUserMessage] }); mockFetchChatRoomMessages.mockResolvedValueOnce({ messages: [persistedUserMessage] });
const sendPromise = act(async () => { const sendPromise = act(async () => {
await expect(result.current.sendRoomMessage("hello")).rejects.toThrow("No active room responders available for room room-1"); await expect(result.current.sendRoomMessage("hello")).rejects.toBeInstanceOf(RoomMessageDeliveredButReplyFailedError);
}); });
await sendPromise; await sendPromise;
await waitFor(() => { await waitFor(() => {

View File

@@ -86,6 +86,39 @@ function createOptimisticRoomMessage(roomId: string, content: string, attachment
}; };
} }
function reconcileOptimisticUserMessage(
previous: ChatRoomMessage[],
delivered: ChatRoomMessage,
optimisticId?: string,
): ChatRoomMessage[] {
if (previous.some((candidate) => candidate.id === delivered.id)) {
return previous;
}
const optimisticIndexById = optimisticId
? previous.findIndex((candidate) => candidate.id === optimisticId)
: -1;
if (optimisticIndexById >= 0) {
const next = [...previous];
next[optimisticIndexById] = delivered;
return next;
}
if (delivered.role === "user") {
const optimisticIndexByContent = previous.findIndex((candidate) =>
candidate.role === "user"
&& candidate.id.startsWith("temp-")
&& candidate.content.trim() === delivered.content.trim());
if (optimisticIndexByContent >= 0) {
const next = [...previous];
next[optimisticIndexByContent] = delivered;
return next;
}
}
return [...previous, delivered];
}
export function useChatRooms( export function useChatRooms(
projectId?: string, projectId?: string,
addToast?: (msg: string, type?: "success" | "error" | "warning") => void, addToast?: (msg: string, type?: "success" | "error" | "warning") => void,
@@ -332,8 +365,7 @@ export function useChatRooms(
if (activeRoomRef.current?.id === roomId) { if (activeRoomRef.current?.id === roomId) {
setMessages((previous) => { setMessages((previous) => {
const next = previous.map((message) => const next = reconcileOptimisticUserMessage(previous, postResult.message, optimisticMessage.id);
message.id === optimisticMessage.id ? postResult.message : message);
// Snapshot mirrors server `order: desc` shape. // Snapshot mirrors server `order: desc` shape.
writeCache(messagesCacheKey(roomId), next, { maxBytes: 500_000 }); writeCache(messagesCacheKey(roomId), next, { maxBytes: 500_000 });
return next; return next;
@@ -349,8 +381,10 @@ export function useChatRooms(
setMessages(latestMessages.messages); setMessages(latestMessages.messages);
timer.mark("hydrate"); timer.mark("hydrate");
} catch (error) { } catch (error) {
let recoveredMessages: ChatRoomMessage[] | null = null;
try { try {
const latestMessages = await fetchChatRoomMessages(roomId, { limit: 100, order: "desc" }, projectId); const latestMessages = await fetchChatRoomMessages(roomId, { limit: 100, order: "desc" }, projectId);
recoveredMessages = latestMessages.messages;
// Snapshot mirrors server `order: desc` shape. // Snapshot mirrors server `order: desc` shape.
writeCache(messagesCacheKey(roomId), latestMessages.messages, { maxBytes: 500_000 }); writeCache(messagesCacheKey(roomId), latestMessages.messages, { maxBytes: 500_000 });
if (activeRoomRef.current?.id === roomId) { if (activeRoomRef.current?.id === roomId) {
@@ -368,7 +402,10 @@ export function useChatRooms(
} }
} }
if (userMessageDelivered) { const messageDeliveredAfterRecovery = recoveredMessages?.some((candidate) =>
candidate.role === "user" && candidate.content.trim() === optimisticMessage.content.trim());
if (userMessageDelivered || messageDeliveredAfterRecovery) {
const message = error instanceof Error && error.message.trim() const message = error instanceof Error && error.message.trim()
? error.message ? error.message
: "Message delivered, but failed to refresh room replies"; : "Message delivered, but failed to refresh room replies";
@@ -488,25 +525,10 @@ export function useChatRooms(
if (activeRoomRef.current?.id !== message.roomId) return; if (activeRoomRef.current?.id !== message.roomId) return;
setMessages((previous) => { setMessages((previous) => {
if (previous.some((candidate) => candidate.id === message.id)) { const next = reconcileOptimisticUserMessage(previous, message);
if (next === previous) {
return previous; return previous;
} }
if (message.role === "user") {
const optimisticIndex = previous.findIndex((candidate) =>
candidate.role === "user"
&& candidate.id.startsWith("temp-")
&& candidate.content.trim() === message.content.trim());
if (optimisticIndex >= 0) {
const next = [...previous];
next[optimisticIndex] = message;
// Snapshot mirrors server `order: desc` shape.
writeCache(messagesCacheKey(message.roomId), next, { maxBytes: 500_000 });
return next;
}
}
const next = [...previous, message];
// Snapshot mirrors server `order: desc` shape. // Snapshot mirrors server `order: desc` shape.
writeCache(messagesCacheKey(message.roomId), next, { maxBytes: 500_000 }); writeCache(messagesCacheKey(message.roomId), next, { maxBytes: 500_000 });
return next; return next;