From 9d405a04b7f6c4bb5d8f5ce91c999961d0c6eb9a Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 24 Aug 2026 08:46:38 -0700 Subject: [PATCH 1/2] Scope reply-in-thread to parent via threadId cache Closes CL-6660 --- packages/chat-ui/src/use-optimistic-sends.ts | 63 +++++++- packages/chat-ui/src/use-thread-navigation.ts | 23 ++- packages/chat-ui/src/use-workbench-feed.ts | 141 ++++++++++++++---- packages/chat-ui/test/stream-apply.test.ts | 127 ++++++++++++++++ packages/chat/src/routes.ts | 71 +++++---- packages/chat/src/workbench-service.ts | 9 ++ 6 files changed, 363 insertions(+), 71 deletions(-) diff --git a/packages/chat-ui/src/use-optimistic-sends.ts b/packages/chat-ui/src/use-optimistic-sends.ts index 8c6da3e2e..3633207e8 100644 --- a/packages/chat-ui/src/use-optimistic-sends.ts +++ b/packages/chat-ui/src/use-optimistic-sends.ts @@ -22,7 +22,12 @@ import type { MentionInviteIntent } from "./mentions"; import { CHAT_STRINGS } from "./strings"; import type { PendingMessageStatus, TimelineMessageItem } from "./timeline"; import { toast } from "@corbits/react-ui"; -import { chatMessagesQueryKey } from "./use-workbench-feed"; +import { + chatMessagesQueryKey, + chatThreadsQueryKey, + ensureReplyThreadRow, + type ThreadsQueryData, +} from "./use-workbench-feed"; export type PendingSend = { readonly nonce: string; @@ -203,9 +208,6 @@ export function useOptimisticSends(args: { clientId: nonce, ...inviteOption, }); - if (sent.threadId !== undefined) { - openThreadById(sent.threadId); - } } else { sent = await sendMessage(tenantId, activeWorkbenchId, parts, { clientId: nonce, @@ -221,6 +223,12 @@ export function useOptimisticSends(args: { address: pendingSenderAddress(currentUserPrincipalId), }, clientId: sent.clientId ?? nonce, + // Stamp the server-assigned thread so the confirmed row scopes into + // the reply feed immediately (CL-6660). The stream echo may still + // arrive without threadId when assignment raced publish; keeping + // this stamp means that echo's dedupe path is a no-op rather than + // leaving the message stuck on the root feed. + ...(sent.threadId !== undefined ? { threadId: sent.threadId } : {}), }; // A refresh already in flight was issued before this send and will // report a mailbox without it — letting it resolve into the cache @@ -234,19 +242,60 @@ export function useOptimisticSends(args: { // where it has vanished from both while a refetch is still in // flight. A refresh may have folded it in already (refresh-first // interleaving); appending unconditionally would render it twice - // under one key. + // under one key. When the row is already present but still lacks + // `threadId`, patch it in place so a racing stream echo without a + // thread cannot leave the confirm stranded on the root feed. queryClient.setQueryData( chatMessagesQueryKey(tenantId, activeWorkbenchId), (current: MessagesResponse | undefined) => { if (current === undefined) return current; - const alreadyPresent = current.items.some( + const matchIndex = current.items.findIndex( (item) => item.id === confirmed.id || item.clientId === confirmed.clientId, ); - if (alreadyPresent) return current; + if (matchIndex >= 0) { + const existing = current.items[matchIndex]; + if (existing === undefined) return current; + if ( + confirmed.threadId !== undefined && + existing.threadId !== confirmed.threadId + ) { + const items = current.items.slice(); + items[matchIndex] = { + ...existing, + threadId: confirmed.threadId, + }; + return { ...current, items }; + } + return current; + } return { ...current, items: [...current.items, confirmed] }; }, ); + // Seed / bump the threads cache before navigation opens the just- + // created id (CL-6660). Without the row, `useThreadNavigation`'s + // stale-id effect would drop `openThreadId` the moment it was set. + if (sent.threadId !== undefined) { + const threadId = sent.threadId; + const parentMessageId = pendingParentMessageId; + queryClient.setQueryData( + chatThreadsQueryKey(tenantId, activeWorkbenchId), + (current: ThreadsQueryData | undefined) => { + if (current === undefined) return current; + return ensureReplyThreadRow(current, { + threadId, + createdAt: sent.createdAt, + ...(parentMessageId !== null + ? { parentMessageId } + : {}), + bumpReplyCount: true, + }); + }, + ); + if (pendingParentMessageId !== null) { + openThreadById(threadId); + } + } setPendingSends((current) => current.filter((p) => p.nonce !== nonce)); // A message just landed in a workbench with an agent in it: a reply // is owed, so show the typing indicator now rather than sitting diff --git a/packages/chat-ui/src/use-thread-navigation.ts b/packages/chat-ui/src/use-thread-navigation.ts index 5af924c56..56b94f3dd 100644 --- a/packages/chat-ui/src/use-thread-navigation.ts +++ b/packages/chat-ui/src/use-thread-navigation.ts @@ -12,7 +12,7 @@ // inside a sub-thread), so nothing here reasons about nesting beyond // grouping what it is handed. -import { useCallback, useEffect, useMemo, useState } from "react"; +import { useCallback, useEffect, useMemo, useRef, useState } from "react"; import { toast } from "@corbits/react-ui"; import { forkThread } from "./api"; import type { WorkbenchThreadRow } from "./api"; @@ -62,10 +62,17 @@ export function useThreadNavigation(args: { const [pendingParentMessageId, setPendingParentMessageId] = useState< string | null >(null); + // A thread id `openThreadById` just set (first reply create, CL-6660) + // can land one render before the seeded `GET /threads` row is visible + // to this hook. Hold it so the stale-id effect below does not drop a + // brand-new open back to the root feed; clear once the row appears or + // the workbench changes. + const heldOpenThreadIdRef = useRef(null); // Switching workbenches resets the view: a thread id belongs to the // workbench it came from. useEffect(() => { + heldOpenThreadIdRef.current = null; setOpenThreadId(null); setPendingParentMessageId(null); }, [activeWorkbenchId]); @@ -74,12 +81,18 @@ export function useThreadNavigation(args: { // restart the reconnect-ownership challenge treats the run as dead and // every id under it disappears. That is a stale reference, not a // failure: drop it and fall back to the workbench's live feed rather - // than leaving the reader staring at an empty thread. + // than leaving the reader staring at an empty thread. A just-created + // id still waiting for its seeded row is held, not stale. useEffect(() => { if (openThreadId === null || !threadsLoaded) return; - if (!threads.some((thread) => thread.id === openThreadId)) { - setOpenThreadId(null); + if (threads.some((thread) => thread.id === openThreadId)) { + if (heldOpenThreadIdRef.current === openThreadId) { + heldOpenThreadIdRef.current = null; + } + return; } + if (heldOpenThreadIdRef.current === openThreadId) return; + setOpenThreadId(null); }, [openThreadId, threads, threadsLoaded]); const replyThreadFor = useCallback( @@ -135,11 +148,13 @@ export function useThreadNavigation(args: { /** Open a thread that already exists — from the threads menu, a * breadcrumb, or a send that just created one. */ const openThreadById = useCallback((threadId: string) => { + heldOpenThreadIdRef.current = threadId; setPendingParentMessageId(null); setOpenThreadId(threadId); }, []); const closeThread = useCallback(() => { + heldOpenThreadIdRef.current = null; setOpenThreadId(null); setPendingParentMessageId(null); }, []); diff --git a/packages/chat-ui/src/use-workbench-feed.ts b/packages/chat-ui/src/use-workbench-feed.ts index 26fe6b741..7c7d6a6b1 100644 --- a/packages/chat-ui/src/use-workbench-feed.ts +++ b/packages/chat-ui/src/use-workbench-feed.ts @@ -165,21 +165,100 @@ function settleWorkbenchListRow( } } +/** + * The `GET /threads` cache shape — shared by stream apply and the + * optimistic-send path that seeds a just-created reply thread before + * `openThreadById` runs (CL-6660). + */ +export type ThreadsQueryData = { + readonly rootThreadId: string; + readonly items: readonly WorkbenchThreadRow[]; +}; + +/** + * Ensures a reply-thread row exists for `threadId`, and optionally bumps + * its activity. A first in-reply-to send uses this to seed the row with + * `parentMessageId` before navigation opens the thread; a stream echo + * uses it when the row is still missing so "N replies" and open-by-id + * have something to hang onto without a refetch (CL-6660). + * + * When the row already exists, `parentMessageId` fills in only if the + * cached row still has `null` (a stream-seeded stub catching up to the + * sender's known parent) — never overwrites a real parent. + */ +export function ensureReplyThreadRow( + current: ThreadsQueryData, + args: { + readonly threadId: string; + readonly createdAt: string; + readonly parentMessageId?: string | null; + /** When false, only insert-or-fill — never increment replyCount. */ + readonly bumpReplyCount?: boolean; + }, +): ThreadsQueryData { + const bump = args.bumpReplyCount ?? true; + const existing = current.items.find((row) => row.id === args.threadId); + if (existing !== undefined) { + const filledParent = + args.parentMessageId !== undefined && + args.parentMessageId !== null && + existing.parentMessageId === null + ? args.parentMessageId + : existing.parentMessageId; + const replyCount = bump ? existing.replyCount + 1 : existing.replyCount; + const lastActivityAt = bump ? args.createdAt : existing.lastActivityAt; + if ( + filledParent === existing.parentMessageId && + replyCount === existing.replyCount && + lastActivityAt === existing.lastActivityAt + ) { + return current; + } + return { + ...current, + items: current.items.map((row) => + row.id === args.threadId + ? { + ...row, + parentMessageId: filledParent, + replyCount, + lastActivityAt, + } + : row, + ), + }; + } + const newRow: WorkbenchThreadRow = { + id: args.threadId, + kind: "reply", + parentMessageId: args.parentMessageId ?? null, + parentThreadId: current.rootThreadId === "" ? null : current.rootThreadId, + runRef: null, + title: null, + createdAt: args.createdAt, + replyCount: bump ? 1 : 0, + lastActivityAt: bump ? args.createdAt : null, + }; + return { ...current, items: [...current.items, newRow] }; +} + /** * Folds a freshly published `chat.message` row straight into the messages * cache (CL-6328) — the §6/1.2 bar is zero refetches triggered by a stream * event, and this event already carries everything a `GET .../messages` * page item would. Deduped by `id` or `clientId` so the row this reader's * own optimistic send already wrote (`use-optimistic-sends.ts`) is never - * doubled when its own echo arrives back over the stream. A message that - * belongs to a thread also bumps that thread's `replyCount`/ - * `lastActivityAt` row in the threads cache — the one piece of thread - * metadata `MessageItem` itself doesn't carry (see `./thread-feed.ts`). - * The workbench list-cache row settles its preview/`lastActivityAt` here - * too (CL-6795) so the sidebar tracks the latest person-facing text - * without waiting for a list refetch — including when the messages - * append was a no-op because this connection's optimistic send already - * wrote the row. + * doubled when its own echo arrives back over the stream. + * + * When the matching row is already present but lacks `threadId` and the + * stream carries one, the cached row is patched in place (CL-6660) — an + * optimistic confirm that raced ahead of thread assignment used to leave + * the message stuck on the root feed. A message that belongs to a thread + * also bumps (or seeds) that thread's `replyCount`/`lastActivityAt` row + * in the threads cache — only when this apply newly scopes the message + * into the thread, so an echo of an already-counted optimistic write is + * a no-op on replyCount. The workbench list-cache row settles its + * preview/`lastActivityAt` here too (CL-6795). */ export function applyStreamMessage( queryClient: QueryClient, @@ -187,43 +266,47 @@ export function applyStreamMessage( workbenchId: string, message: MessageItem, ): void { + let threadActivityDelta = 0; queryClient.setQueryData( chatMessagesQueryKey(tenantId, workbenchId), (current: MessagesResponse | undefined) => { if (current === undefined) return current; - const alreadyPresent = current.items.some( + const matchIndex = current.items.findIndex( (item) => item.id === message.id || (message.clientId !== undefined && item.clientId === message.clientId), ); - if (alreadyPresent) return current; + if (matchIndex >= 0) { + const existing = current.items[matchIndex]; + if (existing === undefined) return current; + if ( + message.threadId !== undefined && + existing.threadId !== message.threadId + ) { + threadActivityDelta = 1; + const items = current.items.slice(); + items[matchIndex] = { ...existing, threadId: message.threadId }; + return { ...current, items }; + } + return current; + } + if (message.threadId !== undefined) threadActivityDelta = 1; return { ...current, items: [...current.items, message] }; }, ); settleWorkbenchListRow(queryClient, tenantId, workbenchId, message); if (message.threadId === undefined) return; + const threadId = message.threadId; queryClient.setQueryData( chatThreadsQueryKey(tenantId, workbenchId), - ( - current: - | { - readonly rootThreadId: string; - readonly items: readonly WorkbenchThreadRow[]; - } - | undefined, - ) => { + (current: ThreadsQueryData | undefined) => { if (current === undefined) return current; - const items = current.items.map((row) => - row.id === message.threadId - ? { - ...row, - replyCount: row.replyCount + 1, - lastActivityAt: message.createdAt, - } - : row, - ); - return { ...current, items }; + return ensureReplyThreadRow(current, { + threadId, + createdAt: message.createdAt, + bumpReplyCount: threadActivityDelta > 0, + }); }, ); } diff --git a/packages/chat-ui/test/stream-apply.test.ts b/packages/chat-ui/test/stream-apply.test.ts index 97465eb24..6e0b73cb7 100644 --- a/packages/chat-ui/test/stream-apply.test.ts +++ b/packages/chat-ui/test/stream-apply.test.ts @@ -14,12 +14,14 @@ import { chatMessagesQueryKey, chatPinsQueryKey, chatThreadsQueryKey, + ensureReplyThreadRow, } from "../src/use-workbench-feed"; import type { MessageItem, MessagesResponse, PinnedMessage, Workbench, + WorkbenchThreadRow, } from "../src/api"; import { workbenchesQueryKey } from "../src/api"; @@ -224,6 +226,131 @@ describe("applyStreamMessage", () => { expect(list?.[0]?.lastActivityAt).toBe("2026-01-01T00:02:00.000Z"); expect(list?.[0]?.preview).toBe("let's pull Scout in"); }); + + test("patches threadId onto an already-cached optimistic confirm (CL-6660)", () => { + const qc = client(); + const optimistic: MessageItem = { + ...baseMessage, + id: "m_server", + clientId: "pending_1", + }; + qc.setQueryData(chatMessagesQueryKey(TENANT, WORKBENCH), { + items: [optimistic], + } satisfies MessagesResponse); + qc.setQueryData(chatThreadsQueryKey(TENANT, WORKBENCH), { + rootThreadId: "root", + items: [] as WorkbenchThreadRow[], + }); + + applyStreamMessage(qc, TENANT, WORKBENCH, { + ...optimistic, + threadId: "thr_reply", + createdAt: "2026-01-01T00:05:00.000Z", + }); + + const messages = qc.getQueryData( + chatMessagesQueryKey(TENANT, WORKBENCH), + ); + expect(messages?.items).toHaveLength(1); + expect(messages?.items[0]?.threadId).toBe("thr_reply"); + + const threads = qc.getQueryData<{ + readonly items: readonly WorkbenchThreadRow[]; + }>(chatThreadsQueryKey(TENANT, WORKBENCH)); + expect(threads?.items).toHaveLength(1); + expect(threads?.items[0]?.id).toBe("thr_reply"); + expect(threads?.items[0]?.replyCount).toBe(1); + }); + + test("an echo of an already-threaded optimistic write does not double replyCount (CL-6660)", () => { + const qc = client(); + const confirmed: MessageItem = { + ...baseMessage, + id: "m_server", + clientId: "pending_1", + threadId: "thr_reply", + }; + qc.setQueryData(chatMessagesQueryKey(TENANT, WORKBENCH), { + items: [confirmed], + } satisfies MessagesResponse); + qc.setQueryData(chatThreadsQueryKey(TENANT, WORKBENCH), { + rootThreadId: "root", + items: [ + { + id: "thr_reply", + kind: "reply" as const, + parentMessageId: "m_parent", + parentThreadId: "root", + runRef: null, + title: null, + createdAt: "2026-01-01T00:00:00.000Z", + replyCount: 1, + lastActivityAt: "2026-01-01T00:00:00.000Z", + }, + ], + }); + + applyStreamMessage(qc, TENANT, WORKBENCH, confirmed); + + const threads = qc.getQueryData<{ + readonly items: readonly { replyCount: number }[]; + }>(chatThreadsQueryKey(TENANT, WORKBENCH)); + expect(threads?.items[0]?.replyCount).toBe(1); + }); +}); + +describe("ensureReplyThreadRow (CL-6660)", () => { + const empty: { + readonly rootThreadId: string; + readonly items: readonly WorkbenchThreadRow[]; + } = { rootThreadId: "root", items: [] }; + + test("seeds a missing reply-thread row with parentMessageId and replyCount 1", () => { + const next = ensureReplyThreadRow(empty, { + threadId: "thr_new", + createdAt: "2026-01-01T00:01:00.000Z", + parentMessageId: "m_parent", + }); + expect(next.items).toEqual([ + { + id: "thr_new", + kind: "reply", + parentMessageId: "m_parent", + parentThreadId: "root", + runRef: null, + title: null, + createdAt: "2026-01-01T00:01:00.000Z", + replyCount: 1, + lastActivityAt: "2026-01-01T00:01:00.000Z", + }, + ]); + }); + + test("fills a null parentMessageId on an existing stub without overwriting a real parent", () => { + const withStub = ensureReplyThreadRow(empty, { + threadId: "thr_new", + createdAt: "2026-01-01T00:01:00.000Z", + bumpReplyCount: false, + }); + expect(withStub.items[0]?.parentMessageId).toBeNull(); + expect(withStub.items[0]?.replyCount).toBe(0); + + const filled = ensureReplyThreadRow(withStub, { + threadId: "thr_new", + createdAt: "2026-01-01T00:02:00.000Z", + parentMessageId: "m_parent", + bumpReplyCount: false, + }); + expect(filled.items[0]?.parentMessageId).toBe("m_parent"); + + const kept = ensureReplyThreadRow(filled, { + threadId: "thr_new", + createdAt: "2026-01-01T00:03:00.000Z", + parentMessageId: "m_other", + bumpReplyCount: false, + }); + expect(kept.items[0]?.parentMessageId).toBe("m_parent"); + }); }); describe("applyStreamReaction", () => { diff --git a/packages/chat/src/routes.ts b/packages/chat/src/routes.ts index 9e50681d1..559cb2007 100644 --- a/packages/chat/src/routes.ts +++ b/packages/chat/src/routes.ts @@ -2154,6 +2154,44 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { return c.json({ command: commandResult }, 201); } + // Resolve the target thread *before* publish so the `chat.message` + // SSE payload carries `threadId` (CL-6660). Assignment still runs + // after the insert — membership needs the new message id — but + // subscribers must not see a root-feed echo of a reply send. + let targetThreadId: string | undefined; + if (deps.threads !== undefined) { + const root = await deps.threads.ensureRootThread( + ownerTenantId, + workbenchId, + ); + targetThreadId = root.id; + if (parsed.threadId !== undefined) { + const existing = await deps.threads.getThread( + ownerTenantId, + parsed.threadId, + ); + if (existing === undefined || existing.workbenchId !== workbenchId) { + return c.json(ErrorEnvelope("not_found", "thread not found"), 404); + } + targetThreadId = existing.id; + } else if (parsed.inReplyToMessageId !== undefined) { + let reply; + try { + reply = await deps.threads.openReplyThread({ + tenantId: ownerTenantId, + workbenchId, + parentMessageId: parsed.inReplyToMessageId, + }); + } catch (cause) { + if (cause instanceof ThreadDepthCapError) { + return c.json(ErrorEnvelope("conflict", cause.message), 409); + } + throw cause; + } + targetThreadId = reply.id; + } + } + // No unreachable-agent branch: the message is a room write, and // reaching an agent happens off this path — an agent that cannot // be reached answers with a notice on the timeline in its own @@ -2182,6 +2220,7 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { ...(parsed.inReplyToMessageId !== undefined ? { inReplyToMessageId: parsed.inReplyToMessageId } : {}), + ...(targetThreadId !== undefined ? { threadId: targetThreadId } : {}), ...(commandDecision !== undefined && "routeToParticipant" in commandDecision ? { forcedRecipientAddress: commandDecision.routeToParticipant } @@ -2199,37 +2238,7 @@ export function createChatRoutes(deps: CreateChatRoutesDeps): Hono { }); } - if (deps.threads !== undefined) { - const root = await deps.threads.ensureRootThread( - ownerTenantId, - workbenchId, - ); - let targetThreadId = root.id; - if (parsed.threadId !== undefined) { - const existing = await deps.threads.getThread( - ownerTenantId, - parsed.threadId, - ); - if (existing === undefined || existing.workbenchId !== workbenchId) { - return c.json(ErrorEnvelope("not_found", "thread not found"), 404); - } - targetThreadId = existing.id; - } else if (parsed.inReplyToMessageId !== undefined) { - let reply; - try { - reply = await deps.threads.openReplyThread({ - tenantId: ownerTenantId, - workbenchId, - parentMessageId: parsed.inReplyToMessageId, - }); - } catch (cause) { - if (cause instanceof ThreadDepthCapError) { - return c.json(ErrorEnvelope("conflict", cause.message), 409); - } - throw cause; - } - targetThreadId = reply.id; - } + if (deps.threads !== undefined && targetThreadId !== undefined) { await deps.threads.assignMessage({ tenantId: ownerTenantId, workbenchId, diff --git a/packages/chat/src/workbench-service.ts b/packages/chat/src/workbench-service.ts index 61d1fa56d..96d756ef5 100644 --- a/packages/chat/src/workbench-service.ts +++ b/packages/chat/src/workbench-service.ts @@ -972,6 +972,14 @@ export type SendWorkbenchMessageInput = { * mentions nobody — the reply gesture is itself an address. */ readonly inReplyToMessageId?: string; + /** + * Thread membership stamped onto the published `chat.message` so + * stream subscribers see the same scope the POST response returns + * (CL-6660). Assignment into `workbench_thread_messages` still happens + * in the route after the insert — this field only fills the row and + * the SSE payload. + */ + readonly threadId?: string; /** * A participant this message must reach regardless of what its text * mentions (CL-6451): an `@name` typed as a definition's wire name @@ -1061,6 +1069,7 @@ export async function sendWorkbenchMessage( sender: { name: null, address: input.senderAddress }, senderPrincipalId: input.principalId, parts: input.messageParts, + ...(input.threadId !== undefined ? { threadId: input.threadId } : {}), }); return { From 937bbcc9225f8467bb1c606bddece898fe169678 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Mon, 24 Aug 2026 09:23:13 -0700 Subject: [PATCH 2/2] Format files changed in this PR --- packages/chat-ui/src/use-optimistic-sends.ts | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/packages/chat-ui/src/use-optimistic-sends.ts b/packages/chat-ui/src/use-optimistic-sends.ts index 3633207e8..260438ef2 100644 --- a/packages/chat-ui/src/use-optimistic-sends.ts +++ b/packages/chat-ui/src/use-optimistic-sends.ts @@ -285,9 +285,7 @@ export function useOptimisticSends(args: { return ensureReplyThreadRow(current, { threadId, createdAt: sent.createdAt, - ...(parentMessageId !== null - ? { parentMessageId } - : {}), + ...(parentMessageId !== null ? { parentMessageId } : {}), bumpReplyCount: true, }); },