Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 54 additions & 7 deletions packages/chat-ui/src/use-optimistic-sends.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -234,19 +242,58 @@ 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
Expand Down
23 changes: 19 additions & 4 deletions packages/chat-ui/src/use-thread-navigation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<string | null>(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]);
Expand All @@ -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(
Expand Down Expand Up @@ -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);
}, []);
Expand Down
141 changes: 112 additions & 29 deletions packages/chat-ui/src/use-workbench-feed.ts
Original file line number Diff line number Diff line change
Expand Up @@ -165,65 +165,148 @@ 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,
tenantId: string,
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,
});
},
);
}
Expand Down
Loading
Loading