From 524e83dbf991bd956b0ce16ab4a90a3e8bb214df Mon Sep 17 00:00:00 2001 From: Soorya U Date: Wed, 5 Aug 2026 17:38:51 +0530 Subject: [PATCH 1/2] Refactor conversation folding/thread-feed and thread event bus, tidy auth UI Simplifies fold/thread-feed logic with test coverage, adds ephemeral vs persisted chunk stamping to the CLI event bus, and reorganizes auth component internals (moves show.tsx into components/helpers). Co-Authored-By: Claude Sonnet 5 --- apps/cli/package.json | 4 +- apps/cli/src/queue/bus.test.ts | 97 ++++++++++ apps/cli/src/queue/bus.ts | 52 ++++-- apps/web/src/components/auth/login-form.tsx | 127 ++++++++------ .../src/components/auth/magic-link-sent.tsx | 9 +- .../src/components/auth/provider-button.tsx | 13 +- .../auth/security/active-session.tsx | 36 ++-- .../auth/security/active-sessions.tsx | 28 +-- .../auth/security/linked-accounts.tsx | 40 +++-- .../chat/composer/agent-model-picker.tsx | 22 +-- .../composer/compact-composer-controls.tsx | 93 +++++----- .../chat/composer/footer-controls.tsx | 56 +++--- .../components/chat/feed/feed-entry-view.tsx | 3 - .../src/components/chat/work-log/diff-row.tsx | 11 +- .../src/components/chat/work-log/tool-row.tsx | 36 ++-- apps/web/src/components/helpers/show.tsx | 9 + bun.lock | 11 +- shared/hooks/package.json | 1 + .../src/conversation/conversation-cache.ts | 68 +++---- .../conversation/use-interactive-respond.ts | 7 +- .../conversation/use-thread-conversation.ts | 17 +- shared/schemas/src/rtc/chat.ts | 53 +++--- shared/schemas/src/rtc/threads.ts | 1 + shared/schemas/src/view/index.ts | 34 ++-- .../utils/src/conversations/entries.test.ts | 128 ++++++++++++++ shared/utils/src/conversations/entries.ts | 42 +++-- shared/utils/src/conversations/fold.test.ts | 123 ++++++++++++- shared/utils/src/conversations/fold.ts | 166 ++++++------------ .../src/conversations/thread-feed.test.ts | 104 +++++++++-- shared/utils/src/conversations/thread-feed.ts | 106 +++-------- 30 files changed, 953 insertions(+), 544 deletions(-) create mode 100644 apps/cli/src/queue/bus.test.ts create mode 100644 apps/web/src/components/helpers/show.tsx create mode 100644 shared/utils/src/conversations/entries.test.ts diff --git a/apps/cli/package.json b/apps/cli/package.json index 00dead1c..e1d26e86 100644 --- a/apps/cli/package.json +++ b/apps/cli/package.json @@ -24,17 +24,17 @@ "@cyrus/utils": "workspace:*", "@ff-labs/fff-node": "^0.9.6", "@orpc/server": "catalog:rpc", + "@soorya-u/better-auth-ws-ticket": "catalog:auth", "@t3-oss/env-core": "catalog:env", "@tursodatabase/database": "^0.6.1", - "@soorya-u/better-auth-ws-ticket": "catalog:auth", "async-mutex": "^0.5.0", "better-auth": "catalog:auth", "better-result": "catalog:core", "commander": "^15.0.0", "diff": "^7.0.0", "es-git": "^0.7.0", - "extract-zip": "^2.0.1", "evlog": "catalog:observability", + "extract-zip": "^2.0.1", "node-datachannel": "^0.32.3", "parse-diff": "^0.11.1", "proper-lockfile": "^4.1.2", diff --git a/apps/cli/src/queue/bus.test.ts b/apps/cli/src/queue/bus.test.ts new file mode 100644 index 00000000..5404c8ea --- /dev/null +++ b/apps/cli/src/queue/bus.test.ts @@ -0,0 +1,97 @@ +import { describe, expect, test } from "bun:test"; +import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; +import { createThreadEventBus } from "./bus"; + +function persistedChunk(seq: number, turnId = "turn-1"): ChatChunk { + return { + threadId: "thread-1", + turnId, + seq, + event: { type: "message_completed", text: "done", messageId: "m1" }, + }; +} + +function tokenChunk(turnId = "turn-1"): ChatChunk { + return { + threadId: "thread-1", + turnId, + seq: 0, + event: { type: "token", text: "hi", messageId: "m1" }, + }; +} + +function createReader(iterator: AsyncGenerator) { + let pending = iterator.next(); + return { + async next(): Promise { + const result = await pending; + pending = iterator.next(); + if (result.done) throw new Error("iterator ended unexpectedly"); + return result.value; + }, + }; +} + +describe("thread event bus", () => { + test("stamps ephemeral chunks with the last known persisted anchor and an increasing sub", async () => { + const bus = createThreadEventBus(); + bus.watch("peer-1", "thread-1"); + const reader = createReader(bus.subscribe("peer-1")); + + bus.publish(persistedChunk(5)); + bus.publish(tokenChunk()); + bus.publish(tokenChunk()); + + const first = await reader.next(); + expect(first.seq).toBe(5); + expect(first.sub).toBeUndefined(); + + expect(await reader.next()).toMatchObject({ seq: 5, sub: 1 }); + expect(await reader.next()).toMatchObject({ seq: 5, sub: 2 }); + }); + + test("does not replay ephemeral chunks on reconnect, but does replay persisted ones", async () => { + const bus = createThreadEventBus(); + bus.watch("peer-1", "thread-1"); + const first = createReader(bus.subscribe("peer-1")); + + bus.publish(persistedChunk(5)); + bus.publish(tokenChunk()); + await first.next(); + await first.next(); + + // Simulate a reconnect: same peerId subscribes again. + const second = createReader(bus.subscribe("peer-1")); + const replayed = await second.next(); + expect(replayed.seq).toBe(5); + expect(replayed.sub).toBeUndefined(); + + bus.publish(tokenChunk()); + const live = await second.next(); + // The sub counter keeps advancing across a reconnect (it's per-anchor, + // not per-connection) — the first ephemeral chunk before the reconnect + // already claimed sub 1. + expect(live).toMatchObject({ sub: 2 }); + }); + + test("advances the anchor and resets sub once new content persists", async () => { + const bus = createThreadEventBus(); + bus.watch("peer-1", "thread-1"); + const reader = createReader(bus.subscribe("peer-1")); + + bus.publish(persistedChunk(5)); + bus.publish(tokenChunk()); + bus.publish(persistedChunk(6)); + bus.publish(tokenChunk()); + + const a = await reader.next(); + expect(a.seq).toBe(5); + expect(a.sub).toBeUndefined(); + expect(await reader.next()).toMatchObject({ seq: 5, sub: 1 }); + + const b = await reader.next(); + expect(b.seq).toBe(6); + expect(b.sub).toBeUndefined(); + expect(await reader.next()).toMatchObject({ seq: 6, sub: 1 }); + }); +}); diff --git a/apps/cli/src/queue/bus.ts b/apps/cli/src/queue/bus.ts index 22fc8d28..27b0fa3f 100644 --- a/apps/cli/src/queue/bus.ts +++ b/apps/cli/src/queue/bus.ts @@ -17,6 +17,15 @@ function isTerminalEvent(event: ChatChunk["event"]): boolean { return event.type === "turn_completed" || event.type === "turn_interrupted"; } +type ThreadCursor = { + anchor: number; + nextSub: number; +}; + +function isUnstampedPlaceholder(chunk: ChatChunk): boolean { + return chunk.seq === 0; +} + export function createThreadEventBus( options: CreateThreadEventBusOptions = {} ): ThreadEventBus { @@ -26,8 +35,33 @@ export function createThreadEventBus( const watchedThreads = new Map>(); const activeTurnLogs = new Map(); const turnThreads = new Map(); + const threadCursors = new Map(); let closed = false; + function stampChunk(chunk: ChatChunk): ChatChunk { + let cursor = threadCursors.get(chunk.threadId); + if (!cursor) { + cursor = { anchor: 0, nextSub: 1 }; + threadCursors.set(chunk.threadId, cursor); + } + + if (!isUnstampedPlaceholder(chunk)) { + if (chunk.seq > cursor.anchor) { + cursor.anchor = chunk.seq; + cursor.nextSub = 1; + } + return chunk; + } + + const stamped: ChatChunk = { + ...chunk, + seq: cursor.anchor, + sub: cursor.nextSub, + }; + cursor.nextSub += 1; + return stamped; + } + function getWatchedThreads(peerId: string): Set { let set = watchedThreads.get(peerId); if (!set) { @@ -62,11 +96,6 @@ export function createThreadEventBus( } function trimTurnLog(log: ChatChunk[]): void { - while (log.length > maxChunksPerTurn) { - const deltaIndex = log.findIndex((chunk) => chunk.seq === 0); - if (deltaIndex === -1) break; - log.splice(deltaIndex, 1); - } while (log.length > maxChunksPerTurn) log.shift(); } @@ -105,13 +134,15 @@ export function createThreadEventBus( publish(chunk) { if (closed) return; - const terminal = isTerminalEvent(chunk.event); - if (!terminal) { - appendToTurnLog(chunk); + const stamped = stampChunk(chunk); + const terminal = isTerminalEvent(stamped.event); + + if (!terminal && stamped.sub === undefined) { + appendToTurnLog(stamped); } - fanOut(chunk); + fanOut(stamped); if (terminal) { - evictTurnLog(chunk.turnId); + evictTurnLog(stamped.turnId); } }, @@ -186,6 +217,7 @@ export function createThreadEventBus( watchedThreads.clear(); activeTurnLogs.clear(); turnThreads.clear(); + threadCursors.clear(); }, }; } diff --git a/apps/web/src/components/auth/login-form.tsx b/apps/web/src/components/auth/login-form.tsx index 68de1f2a..b3478c81 100644 --- a/apps/web/src/components/auth/login-form.tsx +++ b/apps/web/src/components/auth/login-form.tsx @@ -23,6 +23,7 @@ import { ProviderButtons, type SocialLayout, } from "@/components/auth/provider-buttons"; +import { Show } from "@/components/helpers/show"; import { Button } from "@/components/ui/button"; import { Card, CardContent } from "@/components/ui/card"; import { Checkbox } from "@/components/ui/checkbox"; @@ -228,14 +229,16 @@ function PasswordAuthFields({ : labels.showPassword } > - {isPasswordVisible ? : } + } when={isPasswordVisible}> + + {passwordError} - {rememberMe && ( +
- )} +
{captcha} @@ -260,7 +263,9 @@ function PasswordAuthFields({ disabled={isPending} type="submit" > - {signInEmailPending && } + + + {labels.signIn} @@ -337,7 +342,9 @@ function MagicLinkAuthFields({ disabled={isPending} type="submit" > - {signInMagicLinkPending ? : } + } when={signInMagicLinkPending}> + + {submitLabel} {authButtons} @@ -438,7 +445,7 @@ export function LoginForm({ const authButtons = ( <> - {canUseEmailAndPassword && ( + - )} + {pluginAuthButtons} ); @@ -486,31 +495,57 @@ export function LoginForm({ }); }; - const socialBlock = socialProviders && socialProviders.length > 0 && ( - + const socialBlock = ( + 0)}> + + ); - const separator = showSeparator && ( - - {localization.auth.or} - + const separator = ( + + + {localization.auth.or} + + ); - const authFields = - mode === "signIn" ? ( + const authFields = ( + + } + when={mode === "signIn"} + > {Captcha} : null + +
{Captcha}
+
} emailError={fieldErrors.email} isPasswordVisible={isPasswordVisible} @@ -542,24 +577,8 @@ export function LoginForm({ setPassword={setPassword} signInEmailPending={signInEmailPending} /> - ) : ( - - ); +
+ ); return (
@@ -573,21 +592,17 @@ export function LoginForm({ {localization.auth.signIn} - {socialPosition === "top" && ( - <> - {socialBlock} - {separator} - - )} + + {socialBlock} + {separator} + {authFields} - {socialPosition === "bottom" && ( - <> - {separator} - {socialBlock} - - )} + + {separator} + {socialBlock} +
diff --git a/apps/web/src/components/auth/magic-link-sent.tsx b/apps/web/src/components/auth/magic-link-sent.tsx index 7c2b133a..b88315e9 100644 --- a/apps/web/src/components/auth/magic-link-sent.tsx +++ b/apps/web/src/components/auth/magic-link-sent.tsx @@ -1,6 +1,7 @@ import { useAuth, useAuthPlugin } from "@better-auth-ui/react"; import { cn } from "cnfast"; import { useState } from "react"; +import { Show } from "@/components/helpers/show"; import { Card, CardContent, CardHeader, CardTitle } from "@/components/ui/card"; import { FieldDescription } from "@/components/ui/field"; import { MAGIC_LINK_SENT } from "@/constants/storage-keys"; @@ -43,10 +44,12 @@ export function MagicLinkSent({ className }: MagicLinkSentProps) { : magicLinkLocalization.magicLinkSent} - {email && } + + +
- {emailAndPassword?.enabled && ( +
- )} +
diff --git a/apps/web/src/components/auth/provider-button.tsx b/apps/web/src/components/auth/provider-button.tsx index ea624bc9..c8a9fc5f 100644 --- a/apps/web/src/components/auth/provider-button.tsx +++ b/apps/web/src/components/auth/provider-button.tsx @@ -8,6 +8,7 @@ import { useIsMutating } from "@tanstack/react-query"; import type { SocialProvider } from "better-auth/social-providers"; import { cn } from "cnfast"; import type { ComponentProps, ReactNode } from "react"; +import { Show } from "@/components/helpers/show"; import { Button } from "@/components/ui/button"; import { Spinner } from "@/components/ui/spinner"; import { LastUsedBadge } from "./last-login-method/last-used-badge"; @@ -75,16 +76,20 @@ export function ProviderButton({ variant={variant} {...props} > - {signInSocialPending && } + + + {!signInSocialPending && ProviderIcon && } {label} - {display === "icon" && ( + {getProviderName(provider)} - )} + - {view !== "signUp" && } + + + ); } diff --git a/apps/web/src/components/auth/security/active-session.tsx b/apps/web/src/components/auth/security/active-session.tsx index 0dbe6f01..cdb4b8ee 100644 --- a/apps/web/src/components/auth/security/active-session.tsx +++ b/apps/web/src/components/auth/security/active-session.tsx @@ -6,6 +6,7 @@ import Bowser from "bowser"; import { LogOut, Monitor, Smartphone, X } from "lucide-react"; import { toast } from "sonner"; +import { Show } from "@/components/helpers/show"; import { Badge } from "@/components/ui/badge"; import { Button } from "@/components/ui/button"; import { @@ -73,24 +74,29 @@ export function ActiveSession({ activeSession }: ActiveSessionProps) { return ( - {isMobile ? : } + } when={isMobile}> + + {ua.browser.name || "Unknown Browser"} {ua.os.name ? `, ${ua.os.name}` : ""} - {isCurrentSession ? ( + + + {activeSession.createdAt && timeAgo(activeSession.createdAt)} + + + } + when={isCurrentSession} + > {localization.settings.currentSession} - ) : ( - activeSession.createdAt && ( - - {timeAgo(activeSession.createdAt)} - - ) - )} + - )) - )} + ))} + when={modelsLoading} + > +
+ +
+ diff --git a/apps/web/src/components/chat/composer/compact-composer-controls.tsx b/apps/web/src/components/chat/composer/compact-composer-controls.tsx index 0b943708..fa85cac3 100644 --- a/apps/web/src/components/chat/composer/compact-composer-controls.tsx +++ b/apps/web/src/components/chat/composer/compact-composer-controls.tsx @@ -1,6 +1,7 @@ import { useAgentCatalog } from "@cyrus/hooks/agent-catalog/use-agent-catalog"; import type { RegisteredAgent } from "@cyrus/schemas/rtc/agents"; import { EllipsisIcon } from "lucide-react"; +import { Show } from "@/components/helpers/show"; import { Button } from "@/components/ui/button"; import { DropdownMenu, @@ -52,53 +53,51 @@ export function CompactComposerControls({ - {modes.length > 0 && ( - <> - Mode - - {modes.map((mode) => ( - - {mode.name} - - ))} - - - )} - {modes.length > 0 && efforts.length > 0 && } - {efforts.length > 0 && ( - <> - Effort - - {efforts.map((effort) => ( - - {effort.name} - - ))} - - - )} - {efforts.length > 0 && personas.length > 0 && } - {personas.length > 0 && ( - <> - Persona - - {personas.map((persona) => ( - - {persona.name} - - ))} - - - )} + 0}> + Mode + + {modes.map((mode) => ( + + {mode.name} + + ))} + + + 0 && efforts.length > 0}> + + + 0}> + Effort + + {efforts.map((effort) => ( + + {effort.name} + + ))} + + + 0 && personas.length > 0}> + + + 0}> + Persona + + {personas.map((persona) => ( + + {persona.name} + + ))} + + ); diff --git a/apps/web/src/components/chat/composer/footer-controls.tsx b/apps/web/src/components/chat/composer/footer-controls.tsx index 9ee3d549..0cdcd24c 100644 --- a/apps/web/src/components/chat/composer/footer-controls.tsx +++ b/apps/web/src/components/chat/composer/footer-controls.tsx @@ -4,6 +4,7 @@ import { useMediaQuery } from "@mantine/hooks"; import { AgentModelPicker } from "@/components/chat/composer/agent-model-picker"; import { CompactComposerControls } from "@/components/chat/composer/compact-composer-controls"; import { ComposerContextUsage } from "@/components/chat/composer/composer-context-usage"; +import { Show } from "@/components/helpers/show"; import { Button } from "@/components/ui/button"; import { Select, @@ -39,13 +40,13 @@ export function ComposerFooterColumn({ return (
- {catalogError ? ( +

- {catalogError.message || "Could not load agent catalog."} Select the - agent again to retry. + {catalogError?.message || "Could not load agent catalog."} Select + the agent again to retry.

- {retry ? ( + - ) : null} +
- ) : null} +
- {isCompact ? ( - - ) : ( - <> - {modes.length > 0 && ( - <> + + 0}> - - )} - {efforts.length > 0 && ( - <> + + 0}> - - )} - {personas.length > 0 && ( - <> + + 0}> - - )} - - )} + + + } + when={isCompact} + > + +
diff --git a/apps/web/src/components/chat/feed/feed-entry-view.tsx b/apps/web/src/components/chat/feed/feed-entry-view.tsx index 8095b066..f765fe87 100644 --- a/apps/web/src/components/chat/feed/feed-entry-view.tsx +++ b/apps/web/src/components/chat/feed/feed-entry-view.tsx @@ -3,7 +3,6 @@ import { ErrorRow } from "@/components/chat/feed/error-row"; import { AssistantMessage } from "@/components/chat/messages/assistant-message"; import { AssistantThinking } from "@/components/chat/messages/assistant-thinking"; import { UserMessage } from "@/components/chat/messages/user-message"; -import { DiffRow } from "@/components/chat/work-log/diff-row"; import { ToolRow } from "@/components/chat/work-log/tool-row"; export function FeedEntryView({ entry }: { entry: FeedEntry }) { @@ -18,8 +17,6 @@ export function FeedEntryView({ entry }: { entry: FeedEntry }) { ); case "tool": return ; - case "diff": - return ; case "error": return ; case "approval": diff --git a/apps/web/src/components/chat/work-log/diff-row.tsx b/apps/web/src/components/chat/work-log/diff-row.tsx index 1c8386a1..21dc60ef 100644 --- a/apps/web/src/components/chat/work-log/diff-row.tsx +++ b/apps/web/src/components/chat/work-log/diff-row.tsx @@ -3,6 +3,7 @@ import { PatchDiff } from "@pierre/diffs/react"; import { ChevronDownIcon, ChevronRightIcon, GitBranchIcon } from "lucide-react"; import { useState } from "react"; import { PATCH_DIFF_OPTIONS } from "@/components/chat/diff/patch-diff-options"; +import { Show } from "@/components/helpers/show"; export function DiffRow({ diff }: { diff: DiffView }) { const [open, setOpen] = useState(false); @@ -13,11 +14,9 @@ export function DiffRow({ diff }: { diff: DiffView }) { onClick={() => setOpen((v) => !v)} type="button" > - {open ? ( + } when={open}> - ) : ( - - )} + {diff.path} @@ -29,7 +28,7 @@ export function DiffRow({ diff }: { diff: DiffView }) { - {open && ( +
- )} +
); } diff --git a/apps/web/src/components/chat/work-log/tool-row.tsx b/apps/web/src/components/chat/work-log/tool-row.tsx index b3648947..76344714 100644 --- a/apps/web/src/components/chat/work-log/tool-row.tsx +++ b/apps/web/src/components/chat/work-log/tool-row.tsx @@ -9,6 +9,8 @@ import { XIcon, } from "lucide-react"; import { type KeyboardEvent, useState } from "react"; +import { DiffRow } from "@/components/chat/work-log/diff-row"; +import { Show } from "@/components/helpers/show"; import { KIND_PRESENTATIONS, type ToolPresentation, @@ -51,7 +53,8 @@ export function ToolRow({ tool }: { tool: ToolCallView }) { const [open, setOpen] = useState(false); const presentation = deriveToolPresentation(tool); const Icon = presentation.icon; - const canExpand = Boolean(presentation.detail?.trim()); + const hasDiffs = Boolean(tool.diffs?.length); + const canExpand = hasDiffs || Boolean(presentation.detail?.trim()); const showSuccess = tool.status === "completed"; const showFailed = tool.status === "failed"; const showPending = @@ -87,22 +90,22 @@ export function ToolRow({ tool }: { tool: ToolCallView }) { {presentation.heading} - {presentation.preview && ( + {presentation.preview} - )} +

- {canExpand && ( + - )} +
- {open && canExpand && presentation.detail && ( -
-
-						{presentation.detail}
-					
+ +
+ +
+									{presentation.detail}
+								
+
+ } + when={hasDiffs} + > + {tool.diffs?.map((diff) => ( + + ))} +
- )} +
); } diff --git a/apps/web/src/components/helpers/show.tsx b/apps/web/src/components/helpers/show.tsx new file mode 100644 index 00000000..0d18adce --- /dev/null +++ b/apps/web/src/components/helpers/show.tsx @@ -0,0 +1,9 @@ +import type { PropsWithChildren } from "react"; + +export type ShowProps = { + when: boolean; + fallback?: React.ReactNode; +} & PropsWithChildren; + +export const Show = ({ when, children, fallback }: ShowProps) => + when ? children : fallback; diff --git a/bun.lock b/bun.lock index 8f0a2cfe..26f3b542 100644 --- a/bun.lock +++ b/bun.lock @@ -325,6 +325,7 @@ "@cyrus/providers": "workspace:*", "@cyrus/schemas": "workspace:*", "@cyrus/utils": "workspace:*", + "@tanstack/pacer": "^0.21.1", "@tanstack/react-query": "catalog:rpc", "better-result": "catalog:core", "evlog": "catalog:observability", @@ -1578,7 +1579,7 @@ "@tanstack/devtools-event-bus": ["@tanstack/devtools-event-bus@0.4.2", "", { "dependencies": { "ws": "^8.18.3" } }, "sha512-2LHzhwBFlKHCcklsQrGe8TeyjHd4XAF8nuCO6wHmva5fePUkJUULbu6CsCNAlGlCi0KkEsMXZSvRdR4HgMq4yA=="], - "@tanstack/devtools-event-client": ["@tanstack/devtools-event-client@0.5.0", "", { "bin": { "intent": "./bin/intent.js" } }, "sha512-H+OH3zC6Vhu/K0NaVfQKknEKawc/+2PT+D3SB3Ox0V8SiMlTo0abbmH2rH0721R2aNYbjdMXA1oENOd8E2UVoA=="], + "@tanstack/devtools-event-client": ["@tanstack/devtools-event-client@0.4.4", "", { "bin": { "intent": "./bin/intent.js" } }, "sha512-6T5Yop/793YI+H+5J8Hsyj4kCih9sl4t3ElLgKioW5hk3ocn+ZdSJ94tT7vL7uabxSugWYBZlOTMPzEw2puvQw=="], "@tanstack/devtools-ui": ["@tanstack/devtools-ui@0.6.0", "", { "dependencies": { "clsx": "^2.1.1", "dayjs": "^1.11.19", "goober": "^2.1.16", "solid-js": "^1.9.9" } }, "sha512-CVaM6rT6Nl5ijo83vJYFa2SjofvpuOl/uOvbYGhBrRgUhhelNHhx8zZX+hnZCHmIr0/lzM65hsocnZ72592Rvg=="], @@ -1586,6 +1587,8 @@ "@tanstack/history": ["@tanstack/history@1.162.0", "", {}, "sha512-79pf/RkhteYZTRgcR4F9kbk84P2N8rugQJswxfIqovlbRiT3yI7eBE+5QorIrZaOKktsgzRlXh1l/du/xpl4iA=="], + "@tanstack/pacer": ["@tanstack/pacer@0.21.1", "", { "dependencies": { "@tanstack/devtools-event-client": "^0.4.3", "@tanstack/store": "^0.11.0" } }, "sha512-hB01dd4rlsYcTCNP7wK186jgAe6K5qimgM1Y5Jtvz+9PUaILvpmeLLjmQNUNSO1l23lIt+CeQR6mO1mjlPvRtQ=="], + "@tanstack/query-core": ["@tanstack/query-core@5.101.1", "", {}, "sha512-Y6Y92dkXtNqx67m2pMSxUsA3zOCwv862JexZRP8/EPwvKXMPu9m8rv43spiXWzOUIggQ3SQApttALStzhA8B4g=="], "@tanstack/query-devtools": ["@tanstack/query-devtools@5.101.1", "", {}, "sha512-37RQ9U2PxlXQiv1era2t+uHgVhmiyvxqTMu30+KoVf0rufiucu6rpGRKFJk61Wh5OAZFKqCQd6lxTzFWfLZiuQ=="], @@ -1612,7 +1615,7 @@ "@tanstack/router-utils": ["@tanstack/router-utils@1.162.2", "", { "dependencies": { "@babel/generator": "^7.28.5", "@babel/parser": "^7.28.5", "@babel/types": "^7.28.5", "ansis": "^4.1.0", "babel-dead-code-elimination": "^1.0.12", "diff": "^8.0.2", "pathe": "^2.0.3", "tinyglobby": "^0.2.15" } }, "sha512-hTWqJtqIFFdvuCl8WXNyrodp2L9zo2G37xKRrcVmVRWpAB2h+U1LuRAfS4tsFTiWOIoE/B+WDVFB8JpoEdw6jQ=="], - "@tanstack/store": ["@tanstack/store@0.9.3", "", {}, "sha512-8reSzl/qGWGGVKhBoxXPMWzATSbZLZFWhwBAFO9NAyp0TxzfBP0mIrGb8CP8KrQTmvzXlR/vFPPUrHTLBGyFyw=="], + "@tanstack/store": ["@tanstack/store@0.11.0", "", {}, "sha512-WlzzCt3xi0G6pCAJu1U+2jiECwabETDpQDi3hfkFZvJii9AuZqEKbOiVarX1/bWhTNjU486yQtJCCasi/0q+Cw=="], "@tanstack/virtual-file-routes": ["@tanstack/virtual-file-routes@1.162.0", "", {}, "sha512-uhOeFyxLcU41HzvrxsGpiWdcMbScY1EDgbZ5K7DVRMYInbLYWAC0EA/kx9wXAoSM8q82bUG2hRl8+EAjE6XAbA=="], @@ -3738,12 +3741,16 @@ "@tailwindcss/vite/@tailwindcss/oxide": ["@tailwindcss/oxide@4.3.2", "", { "optionalDependencies": { "@tailwindcss/oxide-android-arm64": "4.3.2", "@tailwindcss/oxide-darwin-arm64": "4.3.2", "@tailwindcss/oxide-darwin-x64": "4.3.2", "@tailwindcss/oxide-freebsd-x64": "4.3.2", "@tailwindcss/oxide-linux-arm-gnueabihf": "4.3.2", "@tailwindcss/oxide-linux-arm64-gnu": "4.3.2", "@tailwindcss/oxide-linux-arm64-musl": "4.3.2", "@tailwindcss/oxide-linux-x64-gnu": "4.3.2", "@tailwindcss/oxide-linux-x64-musl": "4.3.2", "@tailwindcss/oxide-wasm32-wasi": "4.3.2", "@tailwindcss/oxide-win32-arm64-msvc": "4.3.2", "@tailwindcss/oxide-win32-x64-msvc": "4.3.2" } }, "sha512-z8ZgnzX8gdNoWLBLqBPoh/sjnxkwvf9ZuWjnO0l0yIzbLa5/9S+eC5QxGZKRobVHIC3/1BoMWjHblqWjcgFgag=="], + "@tanstack/devtools-client/@tanstack/devtools-event-client": ["@tanstack/devtools-event-client@0.5.0", "", { "bin": { "intent": "./bin/intent.js" } }, "sha512-H+OH3zC6Vhu/K0NaVfQKknEKawc/+2PT+D3SB3Ox0V8SiMlTo0abbmH2rH0721R2aNYbjdMXA1oENOd8E2UVoA=="], + "@tanstack/devtools-event-bus/ws": ["ws@8.21.0", "", { "peerDependencies": { "bufferutil": "^4.0.1", "utf-8-validate": ">=5.0.2" }, "optionalPeers": ["bufferutil", "utf-8-validate"] }, "sha512-Vsp28b7DRcimFQvrqu2Wek3z1iYxDCWqHYB8Qsnk/S4RfaCQzPGPyBNuVjJV3cd6UiKtUtp6sNM77gWvzcCH+g=="], "@tanstack/devtools-vite/chalk": ["chalk@5.6.2", "", {}, "sha512-7NzBL0rN6fMUW+f7A6Io4h40qQlG+xGmtMxfbnH/K7TAtt8JQWVQK+6g0UXKMeVJoyV5EkkNsErQ8pVD3bLHbA=="], "@tanstack/devtools-vite/oxc-parser": ["oxc-parser@0.120.0", "", { "dependencies": { "@oxc-project/types": "^0.120.0" }, "optionalDependencies": { "@oxc-parser/binding-android-arm-eabi": "0.120.0", "@oxc-parser/binding-android-arm64": "0.120.0", "@oxc-parser/binding-darwin-arm64": "0.120.0", "@oxc-parser/binding-darwin-x64": "0.120.0", "@oxc-parser/binding-freebsd-x64": "0.120.0", "@oxc-parser/binding-linux-arm-gnueabihf": "0.120.0", "@oxc-parser/binding-linux-arm-musleabihf": "0.120.0", "@oxc-parser/binding-linux-arm64-gnu": "0.120.0", "@oxc-parser/binding-linux-arm64-musl": "0.120.0", "@oxc-parser/binding-linux-ppc64-gnu": "0.120.0", "@oxc-parser/binding-linux-riscv64-gnu": "0.120.0", "@oxc-parser/binding-linux-riscv64-musl": "0.120.0", "@oxc-parser/binding-linux-s390x-gnu": "0.120.0", "@oxc-parser/binding-linux-x64-gnu": "0.120.0", "@oxc-parser/binding-linux-x64-musl": "0.120.0", "@oxc-parser/binding-openharmony-arm64": "0.120.0", "@oxc-parser/binding-wasm32-wasi": "0.120.0", "@oxc-parser/binding-win32-arm64-msvc": "0.120.0", "@oxc-parser/binding-win32-ia32-msvc": "0.120.0", "@oxc-parser/binding-win32-x64-msvc": "0.120.0" } }, "sha512-WyPWZlcIm+Fkte63FGfgFB8mAAk33aH9h5N9lphXVOHSXEBFFsmYdOBedVKly363aWABjZdaj/m9lBfEY4wt+w=="], + "@tanstack/react-store/@tanstack/store": ["@tanstack/store@0.9.3", "", {}, "sha512-8reSzl/qGWGGVKhBoxXPMWzATSbZLZFWhwBAFO9NAyp0TxzfBP0mIrGb8CP8KrQTmvzXlR/vFPPUrHTLBGyFyw=="], + "@tanstack/router-plugin/chokidar": ["chokidar@5.0.0", "", { "dependencies": { "readdirp": "^5.0.0" } }, "sha512-TQMmc3w+5AxjpL8iIiwebF73dRDF4fBIieAqGn9RGCWaEVwQ6Fb2cGe31Yns0RRIzii5goJ1Y7xbMwo1TxMplw=="], "@tanstack/router-utils/@babel/parser": ["@babel/parser@7.29.7", "", { "dependencies": { "@babel/types": "^7.29.7" }, "bin": "./bin/babel-parser.js" }, "sha512-hnORnjP/1P/zFEndoeX+n+t1RwWRJiJpM/jO7FW32Kn9r5+sJB2JWOdYo4L6k78j15eCwY3Gm/7364B1EMwtNg=="], diff --git a/shared/hooks/package.json b/shared/hooks/package.json index e5bdb1d0..6ce785da 100644 --- a/shared/hooks/package.json +++ b/shared/hooks/package.json @@ -15,6 +15,7 @@ "@cyrus/providers": "workspace:*", "@cyrus/schemas": "workspace:*", "@cyrus/utils": "workspace:*", + "@tanstack/pacer": "^0.21.1", "@tanstack/react-query": "catalog:rpc", "better-result": "catalog:core", "evlog": "catalog:observability", diff --git a/shared/hooks/src/conversation/conversation-cache.ts b/shared/hooks/src/conversation/conversation-cache.ts index 21f0e288..abe53a12 100644 --- a/shared/hooks/src/conversation/conversation-cache.ts +++ b/shared/hooks/src/conversation/conversation-cache.ts @@ -4,7 +4,11 @@ import type { ConversationEntry, GetConversationsOutput, } from "@cyrus/schemas/rtc/threads"; -import { normalizeConversationEntries } from "@cyrus/utils/conversations/entries"; +import { + currentMaxPersistedSeq, + sortConversationEntries, +} from "@cyrus/utils/conversations/entries"; +import { Throttler } from "@tanstack/pacer"; import type { QueryClient } from "@tanstack/react-query"; /** Minimum ms between streaming delta commits — keeps token rendering readable. */ @@ -13,8 +17,10 @@ const STREAM_DELTA_MIN_MS = 120; let syntheticEntrySeq = 0; const pendingDeltas = new Map(); const completedTurnKeys = new Set(); -let deltaFlushHandle: ReturnType | null = null; -let lastDeltaFlushAt = 0; +const deltaThrottler = new Throttler( + (queryClient: QueryClient) => flushPendingDeltas(queryClient), + { wait: STREAM_DELTA_MIN_MS } +); function turnKey(threadId: string, turnId: string): string { return `${threadId}:${turnId}`; @@ -26,7 +32,7 @@ function isTerminalEvent(event: ChatChunk["event"]): boolean { function isStreamingDeltaChunk(chunk: ChatChunk): boolean { return ( - chunk.seq === 0 && + chunk.sub !== undefined && (chunk.event.type === "token" || chunk.event.type === "thought") ); } @@ -73,11 +79,23 @@ function chunkToEntry(chunk: ChatChunk, id?: string): ConversationEntry { id: id ?? `${entryIdForChunk(chunk)}-${++syntheticEntrySeq}`, threadId: chunk.threadId, seq: chunk.seq, + sub: chunk.sub, chunk, createdAt: new Date().toISOString(), }; } +function currentEntries( + queryClient: QueryClient, + threadId: string +): ConversationEntry[] { + return ( + queryClient.getQueryData( + RTC_OPERATION_KEYS.getConversations(threadId) + )?.conversations ?? [] + ); +} + function updateCache( queryClient: QueryClient, threadId: string, @@ -86,9 +104,7 @@ function updateCache( queryClient.setQueryData( RTC_OPERATION_KEYS.getConversations(threadId), (old) => ({ - conversations: normalizeConversationEntries( - updater(old?.conversations ?? []) - ), + conversations: sortConversationEntries(updater(old?.conversations ?? [])), }) ); } @@ -97,8 +113,10 @@ function shouldSkipChunk( entries: ConversationEntry[], chunk: ChatChunk ): boolean { - if (chunk.seq <= 0) return false; - return entries.some((entry) => entry.seq === chunk.seq); + if (chunk.sub !== undefined || chunk.seq <= 0) return false; + return entries.some( + (entry) => entry.sub === undefined && entry.seq === chunk.seq + ); } function applyChunkToEntries( @@ -114,12 +132,12 @@ function applyChunkToEntries( let next = [...entries]; - if (chunk.event.type === "user_message" && chunk.seq > 0) { + if (chunk.event.type === "user_message" && chunk.sub === undefined) { next = next.filter( (entry) => !( entry.chunk.turnId === chunk.turnId && - entry.seq === 0 && + entry.sub !== undefined && entry.chunk.event.type === "user_message" ) ); @@ -168,10 +186,7 @@ function commitChunk(queryClient: QueryClient, chunk: ChatChunk): void { } function flushPendingDeltas(queryClient: QueryClient): void { - if (deltaFlushHandle !== null) { - clearTimeout(deltaFlushHandle); - deltaFlushHandle = null; - } + deltaThrottler.cancel(); for (const chunk of pendingDeltas.values()) commitChunk(queryClient, chunk); pendingDeltas.clear(); @@ -195,7 +210,7 @@ function canPruneEphemeralTurn( return entries.some( (entry) => entry.chunk.turnId === turnId && - entry.seq > 0 && + entry.sub === undefined && (entry.chunk.event.type === "message_completed" || entry.chunk.event.type === "reasoning_completed" || entry.chunk.event.type === "turn_completed") @@ -211,22 +226,11 @@ export function pruneEphemeralTurnEntries( if (!canPruneEphemeralTurn(entries, turnId)) return entries; completedTurnKeys.delete(turnKey(threadId, turnId)); return entries.filter( - (entry) => !(entry.chunk.turnId === turnId && entry.seq === 0) + (entry) => !(entry.chunk.turnId === turnId && entry.sub !== undefined) ); }); } -function scheduleStreamingDeltaFlush(queryClient: QueryClient): void { - if (deltaFlushHandle !== null) return; - const elapsed = Date.now() - lastDeltaFlushAt; - const delay = Math.max(0, STREAM_DELTA_MIN_MS - elapsed); - deltaFlushHandle = setTimeout(() => { - deltaFlushHandle = null; - lastDeltaFlushAt = Date.now(); - flushPendingDeltas(queryClient); - }, delay); -} - function queueStreamingDelta(queryClient: QueryClient, chunk: ChatChunk): void { const key = streamingDeltaKey(chunk); const existing = pendingDeltas.get(key); @@ -234,7 +238,7 @@ function queueStreamingDelta(queryClient: QueryClient, chunk: ChatChunk): void { key, existing ? mergeStreamingDeltaChunks(existing, chunk) : chunk ); - scheduleStreamingDeltaFlush(queryClient); + deltaThrottler.maybeExecute(queryClient); } export function applyChunkToCache( @@ -275,7 +279,8 @@ export function appendOptimisticUserMessage( commitChunk(queryClient, { threadId, turnId, - seq: 0, + seq: currentMaxPersistedSeq(currentEntries(queryClient, threadId)), + sub: 0, event: { type: "user_message", content: message, blocks }, }); } @@ -291,7 +296,8 @@ export function appendTurnTerminal( commitChunk(queryClient, { threadId, turnId, - seq: 0, + seq: currentMaxPersistedSeq(currentEntries(queryClient, threadId)), + sub: 0, event: { type }, }); } diff --git a/shared/hooks/src/conversation/use-interactive-respond.ts b/shared/hooks/src/conversation/use-interactive-respond.ts index a130c5e5..ea4c8367 100644 --- a/shared/hooks/src/conversation/use-interactive-respond.ts +++ b/shared/hooks/src/conversation/use-interactive-respond.ts @@ -1,5 +1,6 @@ import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; import type { GetConversationsOutput } from "@cyrus/schemas/rtc/threads"; +import { currentMaxPersistedSeq } from "@cyrus/utils/conversations/entries"; import { useMutation, useQueryClient } from "@tanstack/react-query"; import { useRtc } from "../contexts/rtc"; import { applyChunkToCache } from "./conversation-cache"; @@ -23,7 +24,8 @@ export function useRespondApproval() { applyChunkToCache(queryClient, { threadId: input.threadId, turnId: "local-interactive", - seq: 0, + seq: currentMaxPersistedSeq(previous?.conversations ?? []), + sub: 0, event: { type: "approval_resolved", toolCallId: input.toolCallId, @@ -56,7 +58,8 @@ export function useRespondElicitation() { applyChunkToCache(queryClient, { threadId: input.threadId, turnId: "local-interactive", - seq: 0, + seq: currentMaxPersistedSeq(previous?.conversations ?? []), + sub: 0, event: { type: "elicitation_resolved", elicitationId: input.elicitationId, diff --git a/shared/hooks/src/conversation/use-thread-conversation.ts b/shared/hooks/src/conversation/use-thread-conversation.ts index 886e7e54..5e323258 100644 --- a/shared/hooks/src/conversation/use-thread-conversation.ts +++ b/shared/hooks/src/conversation/use-thread-conversation.ts @@ -1,6 +1,10 @@ import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; +import type { GetConversationsOutput } from "@cyrus/schemas/rtc/threads"; import type { ThreadConversation } from "@cyrus/schemas/view"; -import { mergeConversationEntries } from "@cyrus/utils/conversations/entries"; +import { + currentMaxPersistedSeq, + mergeConversationEntries, +} from "@cyrus/utils/conversations/entries"; import { fold } from "@cyrus/utils/conversations/fold"; import { keepPreviousData, @@ -14,7 +18,6 @@ import { useRtc } from "../contexts/rtc"; const EMPTY: ThreadConversation = { approvals: [], - diffs: [], elicitations: [], errors: [], messages: [], @@ -42,9 +45,13 @@ export function useThreadConversation( staleTime: Number.POSITIVE_INFINITY, refetchOnWindowFocus: false, refetchOnReconnect: false, - queryFn: async (context) => { - const fetched = await baseQueryOptions.queryFn(context); - const cached = queryClient.getQueryData(queryKey); + queryFn: async () => { + const cached = queryClient.getQueryData(queryKey); + const afterSeq = currentMaxPersistedSeq(cached?.conversations ?? []); + const fetched = await workerConnection.client.getConversations({ + threadId: threadId ?? "", + afterSeq: afterSeq || undefined, + }); return { conversations: mergeConversationEntries( cached?.conversations ?? [], diff --git a/shared/schemas/src/rtc/chat.ts b/shared/schemas/src/rtc/chat.ts index 6bce71ba..54caec5e 100644 --- a/shared/schemas/src/rtc/chat.ts +++ b/shared/schemas/src/rtc/chat.ts @@ -268,12 +268,8 @@ export const PromptInputBlockSchema = z.discriminatedUnion("type", [ PromptResourceBlockSchema, ]); -export type PromptInputBlock = z.infer; - export const ChatMessageSchema = z.array(PromptInputBlockSchema).min(1); -export type ChatMessage = z.infer; - export function formatPromptBlocks(blocks: PromptInputBlock[]): string { return blocks .map((block) => @@ -325,28 +321,6 @@ export const AgentEventSchema = z.discriminatedUnion("type", [ SessionUpdateEventSchema, ]); -export type { PermissionOptionKind } from "@cyrus/schemas/enums/permissions"; -export type { - PlanEntryPriority, - PlanEntryStatus, -} from "@cyrus/schemas/enums/plan"; -export type { ToolCallStatus, ToolKind } from "@cyrus/schemas/enums/tools"; -export type PlanEntry = z.infer; -export type ToolCallContent = z.infer; -export type ToolCallLocation = z.infer; -export type ToolCallFields = z.infer; -export type ToolCallUpdateFields = z.infer; -export type ApprovalRequest = z.infer; -export type ElicitationRequestPayload = z.infer< - typeof ElicitationRequestPayloadSchema ->; -export type RespondApprovalInput = z.infer; -export type RespondElicitationInput = z.infer< - typeof RespondElicitationInputSchema ->; -export type AgentEvent = z.infer; -export type Diff = z.infer; - export const ChatInputSchema = z.object({ agentName: z.string(), message: ChatMessageSchema, @@ -355,27 +329,42 @@ export const ChatInputSchema = z.object({ projectId: z.string(), }); -export type ChatInput = z.infer; - export const ChatOutputSchema = z.object({ threadId: z.string(), turnId: z.string(), }); -export type ChatOutput = z.infer; - export const ChatChunkSchema = z.object({ threadId: z.string(), turnId: z.string(), seq: z.number(), + sub: z.number().optional(), event: AgentEventSchema, }); -export type ChatChunk = z.infer; - export const CancelInputSchema = z.object({ agentName: z.string(), threadId: z.string(), }); +export type PromptInputBlock = z.infer; +export type ChatMessage = z.infer; +export type PlanEntry = z.infer; +export type ToolCallContent = z.infer; +export type ToolCallLocation = z.infer; +export type ToolCallFields = z.infer; +export type ToolCallUpdateFields = z.infer; +export type ApprovalRequest = z.infer; +export type ElicitationRequestPayload = z.infer< + typeof ElicitationRequestPayloadSchema +>; +export type RespondApprovalInput = z.infer; +export type RespondElicitationInput = z.infer< + typeof RespondElicitationInputSchema +>; +export type AgentEvent = z.infer; +export type Diff = z.infer; +export type ChatInput = z.infer; +export type ChatOutput = z.infer; +export type ChatChunk = z.infer; export type CancelInput = z.infer; diff --git a/shared/schemas/src/rtc/threads.ts b/shared/schemas/src/rtc/threads.ts index b16fb221..3230e287 100644 --- a/shared/schemas/src/rtc/threads.ts +++ b/shared/schemas/src/rtc/threads.ts @@ -27,6 +27,7 @@ export const ConversationEntrySchema = z.object({ id: z.string(), threadId: z.string(), seq: z.number(), + sub: z.number().optional(), chunk: ChatChunkSchema, createdAt: z.string(), }); diff --git a/shared/schemas/src/view/index.ts b/shared/schemas/src/view/index.ts index 0a943aea..0a2b63c0 100644 --- a/shared/schemas/src/view/index.ts +++ b/shared/schemas/src/view/index.ts @@ -10,6 +10,8 @@ export const MessageViewSchema = z.object({ createdAt: z.string(), turnId: z.string(), streaming: z.boolean().optional(), + seq: z.number(), + sub: z.number().optional(), }); export const ThoughtViewSchema = z.object({ @@ -18,6 +20,18 @@ export const ThoughtViewSchema = z.object({ createdAt: z.string(), turnId: z.string(), streaming: z.boolean().optional(), + seq: z.number(), + sub: z.number().optional(), +}); + +export const DiffViewSchema = z.object({ + id: z.string(), + path: z.string(), + patch: z.string(), + additions: z.number(), + deletions: z.number(), + turnId: z.string(), + toolCallId: z.string().optional(), }); export const ToolCallViewSchema = z.object({ @@ -29,16 +43,9 @@ export const ToolCallViewSchema = z.object({ rawOutput: z.unknown().optional(), createdAt: z.string(), turnId: z.string(), -}); - -export const DiffViewSchema = z.object({ - id: z.string(), - path: z.string(), - patch: z.string(), - additions: z.number(), - deletions: z.number(), - turnId: z.string(), - toolCallId: z.string().optional(), + diffs: z.array(DiffViewSchema).default([]), + seq: z.number(), + sub: z.number().optional(), }); export const TurnViewSchema = z.object({ @@ -55,6 +62,8 @@ export const ErrorViewSchema = z.object({ code: z.string().optional(), createdAt: z.string(), turnId: z.string(), + seq: z.number(), + sub: z.number().optional(), }); export const ApprovalOptionViewSchema = z.object({ @@ -73,6 +82,8 @@ export const ApprovalViewSchema = z.object({ createdAt: z.string(), turnId: z.string(), resolved: z.boolean().optional(), + seq: z.number(), + sub: z.number().optional(), }); export const ElicitationViewSchema = z.object({ @@ -87,13 +98,14 @@ export const ElicitationViewSchema = z.object({ createdAt: z.string(), turnId: z.string(), resolved: z.boolean().optional(), + seq: z.number(), + sub: z.number().optional(), }); export const ThreadConversationSchema = z.object({ messages: z.array(MessageViewSchema), thoughts: z.array(ThoughtViewSchema), toolCalls: z.array(ToolCallViewSchema), - diffs: z.array(DiffViewSchema), errors: z.array(ErrorViewSchema), approvals: z.array(ApprovalViewSchema), elicitations: z.array(ElicitationViewSchema), diff --git a/shared/utils/src/conversations/entries.test.ts b/shared/utils/src/conversations/entries.test.ts new file mode 100644 index 00000000..09d3f53a --- /dev/null +++ b/shared/utils/src/conversations/entries.test.ts @@ -0,0 +1,128 @@ +import type { ConversationEntry } from "@cyrus/schemas/rtc/threads"; +import { describe, expect, test } from "vitest"; +import { + currentMaxPersistedSeq, + mergeConversationEntries, + sortConversationEntries, +} from "./entries"; + +function entry( + id: string, + seq: number, + turnId: string, + event: ConversationEntry["chunk"]["event"], + options: { sub?: number; createdAt?: string } = {} +): ConversationEntry { + const { sub, createdAt = "2026-01-01T00:00:00.000Z" } = options; + return { + id, + threadId: "thread-1", + seq, + sub, + chunk: { threadId: "thread-1", turnId, seq, sub, event }, + createdAt, + }; +} + +describe("sortConversationEntries", () => { + test("orders persisted entries by seq, with ephemeral deltas slotted after their anchor by sub", () => { + const entries = [ + entry("delta", 1, "turn-1", { type: "token", text: "hi" }, { sub: 1 }), + entry("completed", 2, "turn-1", { type: "turn_completed" }), + entry("user", 1, "turn-1", { type: "user_message", content: "go" }), + ]; + + const sorted = sortConversationEntries(entries); + + expect(sorted.map((e) => e.id)).toEqual(["user", "delta", "completed"]); + }); + + test("does not drop an ephemeral user message even once a persisted twin exists", () => { + const entries = [ + entry( + "ephemeral", + 0, + "turn-1", + { type: "user_message", content: "hi" }, + { sub: 0 } + ), + entry("persisted", 1, "turn-1", { + type: "user_message", + content: "hi", + }), + ]; + + const sorted = sortConversationEntries(entries); + + expect(sorted.map((e) => e.id)).toEqual(["ephemeral", "persisted"]); + }); +}); + +describe("currentMaxPersistedSeq", () => { + test("ignores ephemeral entries and returns the highest persisted seq", () => { + const entries = [ + entry("persisted-1", 3, "turn-1", { type: "turn_completed" }), + entry("delta", 3, "turn-1", { type: "token", text: "hi" }, { sub: 5 }), + entry("persisted-2", 1, "turn-1", { + type: "user_message", + content: "go", + }), + ]; + + expect(currentMaxPersistedSeq(entries)).toBe(3); + }); + + test("returns 0 when nothing is persisted yet", () => { + expect(currentMaxPersistedSeq([])).toBe(0); + }); +}); + +describe("mergeConversationEntries", () => { + test("drops the ephemeral user message once a persisted twin arrives from the merge", () => { + const cached: ConversationEntry[] = [ + entry( + "ephemeral", + 0, + "turn-1", + { type: "user_message", content: "hi" }, + { sub: 0 } + ), + ]; + const fetched: ConversationEntry[] = [ + entry("persisted", 1, "turn-1", { + type: "user_message", + content: "hi", + }), + ]; + + const merged = mergeConversationEntries(cached, fetched); + + expect(merged.map((e) => e.id)).toEqual(["persisted"]); + }); + + test("keeps persisted entries from fetched over stale cached duplicates at the same seq", () => { + const cached: ConversationEntry[] = [ + entry("stale", 1, "turn-1", { type: "user_message", content: "old" }), + ]; + const fetched: ConversationEntry[] = [ + entry("fresh", 1, "turn-1", { type: "user_message", content: "new" }), + ]; + + const merged = mergeConversationEntries(cached, fetched); + + expect(merged.map((e) => e.id)).toEqual(["fresh"]); + }); + + test("never dedupes a real ephemeral delta by its shared anchor seq", () => { + const cached: ConversationEntry[] = [ + entry("delta-1", 2, "turn-1", { type: "token", text: "a" }, { sub: 1 }), + ]; + const fetched: ConversationEntry[] = [ + entry("persisted", 2, "turn-1", { type: "turn_completed" }), + ]; + + const merged = mergeConversationEntries(cached, fetched); + + expect(merged.map((e) => e.id).sort()).toEqual(["delta-1", "persisted"]); + }); +}); diff --git a/shared/utils/src/conversations/entries.ts b/shared/utils/src/conversations/entries.ts index 0435ae03..ebaf55e4 100644 --- a/shared/utils/src/conversations/entries.ts +++ b/shared/utils/src/conversations/entries.ts @@ -1,12 +1,17 @@ import type { ConversationEntry } from "@cyrus/schemas/rtc/threads"; +function isPersisted(entry: ConversationEntry): boolean { + return entry.sub === undefined; +} + function dropRedundantEphemeralUserMessages( entries: ConversationEntry[] ): ConversationEntry[] { const persistedUserTurns = new Set( entries .filter( - (entry) => entry.seq > 0 && entry.chunk.event.type === "user_message" + (entry) => + isPersisted(entry) && entry.chunk.event.type === "user_message" ) .map((entry) => entry.chunk.turnId) ); @@ -16,7 +21,7 @@ function dropRedundantEphemeralUserMessages( return entries.filter( (entry) => !( - entry.seq === 0 && + !isPersisted(entry) && entry.chunk.event.type === "user_message" && persistedUserTurns.has(entry.chunk.turnId) ) @@ -27,37 +32,38 @@ export function sortConversationEntries( entries: ConversationEntry[] ): ConversationEntry[] { return [...entries].sort((left, right) => { - if (left.seq !== right.seq) { - if (left.seq === 0) return 1; - if (right.seq === 0) return -1; - return left.seq - right.seq; - } + if (left.seq !== right.seq) return left.seq - right.seq; + const leftSub = left.sub ?? 0; + const rightSub = right.sub ?? 0; + if (leftSub !== rightSub) return leftSub - rightSub; return left.createdAt.localeCompare(right.createdAt); }); } -/** Sort and drop ephemeral user messages superseded by persisted ones. */ -export function normalizeConversationEntries( - entries: ConversationEntry[] -): ConversationEntry[] { - return dropRedundantEphemeralUserMessages(sortConversationEntries(entries)); +export function currentMaxPersistedSeq(entries: ConversationEntry[]): number { + return entries.reduce( + (max, entry) => (isPersisted(entry) && entry.seq > max ? entry.seq : max), + 0 + ); } export function mergeConversationEntries( cached: ConversationEntry[], fetched: ConversationEntry[] ): ConversationEntry[] { - if (cached.length === 0) return normalizeConversationEntries(fetched); - if (fetched.length === 0) return normalizeConversationEntries(cached); + if (cached.length === 0) + return sortConversationEntries(dropRedundantEphemeralUserMessages(fetched)); + if (fetched.length === 0) + return sortConversationEntries(dropRedundantEphemeralUserMessages(cached)); const merged = new Map(); for (const entry of fetched) { - if (entry.seq > 0) merged.set(`seq-${entry.seq}`, entry); + if (isPersisted(entry)) merged.set(`seq-${entry.seq}`, entry); } for (const entry of cached) { - if (entry.seq > 0) { + if (isPersisted(entry)) { if (!merged.has(`seq-${entry.seq}`)) merged.set(`seq-${entry.seq}`, entry); continue; @@ -65,5 +71,7 @@ export function mergeConversationEntries( merged.set(entry.id, entry); } - return normalizeConversationEntries([...merged.values()]); + return sortConversationEntries( + dropRedundantEphemeralUserMessages([...merged.values()]) + ); } diff --git a/shared/utils/src/conversations/fold.test.ts b/shared/utils/src/conversations/fold.test.ts index 14853d50..5bd3fea1 100644 --- a/shared/utils/src/conversations/fold.test.ts +++ b/shared/utils/src/conversations/fold.test.ts @@ -45,6 +45,7 @@ describe("fold", () => { createdAt: "2026-07-11T00:00:01.000Z", id: "user-turn-1", role: "user", + seq: 1, streaming: false, turnId: "turn-1", }, @@ -53,6 +54,7 @@ describe("fold", () => { createdAt: "2026-07-11T00:00:02.000Z", id: "turn-1", role: "assistant", + seq: 2, streaming: false, turnId: "turn-1", }, @@ -88,6 +90,47 @@ describe("fold", () => { }); }); + test("orders turns and messages by entry position, ignoring misleading createdAt strings", () => { + // entries always arrive pre-sorted by (seq, sub) — see + // sortConversationEntries — so fold trusts array position for order and + // no longer needs createdAt to reorder anything. + const conversation = folded([ + entry( + 1, + "turn-1", + { type: "user_message", content: "First" }, + "2026-07-11T00:00:09.000Z" + ), + entry( + 2, + "turn-1", + { type: "turn_completed" }, + "2026-07-11T00:00:08.000Z" + ), + entry( + 3, + "turn-2", + { type: "user_message", content: "Second" }, + "2026-07-11T00:00:07.000Z" + ), + entry( + 4, + "turn-2", + { type: "turn_completed" }, + "2026-07-11T00:00:06.000Z" + ), + ]); + + expect(conversation.turns.map((turn) => turn.id)).toEqual([ + "turn-1", + "turn-2", + ]); + expect(conversation.messages.map((message) => message.content)).toEqual([ + "First", + "Second", + ]); + }); + test("folds thoughts, tool calls, and diffs", () => { const conversation = folded([ entry(1, "turn-1", { type: "user_message", content: "Change it" }), @@ -121,6 +164,7 @@ describe("fold", () => { content: "Inspecting", createdAt: "2026-07-11T00:00:02.000Z", id: "turn-1:thought:thought-1", + seq: 2, streaming: false, turnId: "turn-1", }, @@ -131,19 +175,79 @@ describe("fold", () => { title: "Edit README", toolCallId: "tool-1", turnId: "turn-1", + diffs: [ + { + additions: 1, + deletions: 1, + id: "tool-1:README.md", + patch: "@@ -1 +1 @@", + path: "README.md", + toolCallId: "tool-1", + turnId: "turn-1", + }, + ], }), ]); - expect(conversation.diffs).toEqual([ - { - additions: 1, - deletions: 1, - id: "turn-1:README.md", - patch: "@@ -1 +1 @@", - path: "README.md", + }); + + test("replaces a tool call's diffs wholesale on update, doesn't merge", () => { + const conversation = folded([ + entry(1, "turn-1", { type: "user_message", content: "Change it" }), + entry(2, "turn-1", { + content: [ + { + additions: 1, + deletions: 0, + newText: "a", + oldText: "", + patch: "@@ -0,0 +1 @@", + path: "a.ts", + type: "diff", + }, + { + additions: 1, + deletions: 0, + newText: "b", + oldText: "", + patch: "@@ -0,0 +1 @@", + path: "b.ts", + type: "diff", + }, + ], + status: "in_progress", + title: "Edit files", toolCallId: "tool-1", - turnId: "turn-1", - }, + type: "tool_call", + }), + entry(3, "turn-1", { + content: [ + { + additions: 1, + deletions: 0, + newText: "a", + oldText: "", + patch: "@@ -0,0 +1 @@", + path: "a.ts", + type: "diff", + }, + ], + status: "completed", + toolCallId: "tool-1", + type: "tool_call_update", + }), + entry(4, "turn-1", { + status: "completed", + toolCallId: "tool-1", + type: "tool_call_update", + }), + entry(5, "turn-1", { type: "turn_completed" }), ]); + + const toolCall = conversation.toolCalls[0]; + // The 3rd entry replaced the initial two-diff set with just a.ts; the 4th + // entry (a status-only update with no content) must not clear that. + expect(toolCall?.diffs.map((diff) => diff.path)).toEqual(["a.ts"]); + expect(toolCall?.status).toBe("completed"); }); test("folds approval and elicitation requests", () => { @@ -252,6 +356,7 @@ describe("fold", () => { createdAt: "2026-07-11T00:00:02.000Z", id: "error-entry-2", message: "Agent crashed", + seq: 2, turnId: "turn-1", }, ]); diff --git a/shared/utils/src/conversations/fold.ts b/shared/utils/src/conversations/fold.ts index 6f1250c0..fb435fe0 100644 --- a/shared/utils/src/conversations/fold.ts +++ b/shared/utils/src/conversations/fold.ts @@ -19,7 +19,6 @@ type MutableState = { messages: Map; thoughts: Map; toolCalls: Map; - diffs: Map; errors: Map; approvals: Map; elicitations: Map; @@ -45,25 +44,23 @@ function touchTurn( }); } -function upsertDiffs( - diffs: Map, +function extractDiffs( content: ToolCallContent[] | null | undefined, turnId: string, - toolCallId?: string -): void { - if (!content) return; - for (const item of content) { - if (item.type !== "diff") continue; - diffs.set(`${turnId}:${item.path}`, { + toolCallId: string +): DiffView[] { + if (!content) return []; + return content + .filter((item) => item.type === "diff") + .map((item) => ({ additions: item.additions, deletions: item.deletions, - id: `${turnId}:${item.path}`, + id: `${toolCallId}:${item.path}`, patch: item.patch, path: item.path, toolCallId, turnId, - }); - } + })); } function applyUserMessage( @@ -77,9 +74,6 @@ function applyUserMessage( if (existing) { existing.content = event.content; existing.blocks = event.blocks; - if (entry.createdAt < existing.createdAt) { - existing.createdAt = entry.createdAt; - } return; } state.messages.set(key, { @@ -89,6 +83,8 @@ function applyUserMessage( id: key, role: "user", turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -122,6 +118,8 @@ function applyThought( createdAt: entry.createdAt, id: key, turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -144,6 +142,8 @@ function applyReasoningCompleted( createdAt: entry.createdAt, id: key, turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -167,6 +167,8 @@ function applyToken( id: key, role: "assistant", turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -178,6 +180,7 @@ function applyToolCall( ): void { state.toolCalls.set(event.toolCallId, { createdAt: entry.createdAt, + diffs: extractDiffs(event.content, turnId, event.toolCallId), kind: event.kind, rawInput: event.rawInput, rawOutput: event.rawOutput, @@ -185,8 +188,9 @@ function applyToolCall( title: event.title, toolCallId: event.toolCallId, turnId, + seq: entry.seq, + sub: entry.sub, }); - upsertDiffs(state.diffs, event.content, turnId, event.toolCallId); } function applyToolCallUpdate( @@ -202,9 +206,13 @@ function applyToolCallUpdate( if (event.kind) existing.kind = event.kind; if (event.rawInput !== undefined) existing.rawInput = event.rawInput; if (event.rawOutput !== undefined) existing.rawOutput = event.rawOutput; + if (event.content !== undefined && event.content !== null) { + existing.diffs = extractDiffs(event.content, turnId, event.toolCallId); + } } else { state.toolCalls.set(event.toolCallId, { createdAt: entry.createdAt, + diffs: extractDiffs(event.content, turnId, event.toolCallId), kind: event.kind ?? undefined, rawInput: event.rawInput, rawOutput: event.rawOutput, @@ -212,9 +220,10 @@ function applyToolCallUpdate( title: event.title ?? event.toolCallId, toolCallId: event.toolCallId, turnId, + seq: entry.seq, + sub: entry.sub, }); } - upsertDiffs(state.diffs, event.content, turnId, event.toolCallId); } function applyMessageCompleted( @@ -235,6 +244,8 @@ function applyMessageCompleted( id: key, role: "assistant", turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -251,6 +262,8 @@ function applyThreadError( id: key, message: event.message, turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -275,6 +288,8 @@ function applyApprovalRequest( title: event.request.toolCall.title ?? undefined, toolCallId, turnId, + seq: entry.seq, + sub: entry.sub, }); } @@ -297,6 +312,8 @@ function applyElicitationRequest( threadId: entry.threadId, turnId, url: event.request.mode === "url" ? event.request.url : undefined, + seq: entry.seq, + sub: entry.sub, }); } @@ -367,47 +384,6 @@ function latestTurnIdFromEntries( return latestTurnId; } -function turnStartedAt(entries: ConversationEntry[], turnId: string): string { - const userMessages = entries.filter( - (entry) => - entry.chunk.turnId === turnId && entry.chunk.event.type === "user_message" - ); - if (userMessages.length > 0) { - return userMessages.reduce( - (earliest, entry) => - entry.createdAt < earliest ? entry.createdAt : earliest, - userMessages[0]?.createdAt ?? "\uffff" - ); - } - - const turnEntries = entries.filter((entry) => entry.chunk.turnId === turnId); - return turnEntries[0]?.createdAt ?? "\uffff"; -} - -function turnMinPersistedSeq( - entries: ConversationEntry[], - turnId: string -): number { - const seqs = entries - .filter((entry) => entry.chunk.turnId === turnId && entry.seq > 0) - .map((entry) => entry.seq); - if (seqs.length === 0) return Number.POSITIVE_INFINITY; - return Math.min(...seqs); -} - -function compareTurnOrder( - entries: ConversationEntry[], - leftId: string, - rightId: string -): number { - const leftSeq = turnMinPersistedSeq(entries, leftId); - const rightSeq = turnMinPersistedSeq(entries, rightId); - if (leftSeq !== rightSeq) return leftSeq - rightSeq; - return turnStartedAt(entries, leftId).localeCompare( - turnStartedAt(entries, rightId) - ); -} - function applyEvent( state: MutableState, entry: ConversationEntry, @@ -465,7 +441,6 @@ export function fold( ): Result { const state: MutableState = { approvals: new Map(), - diffs: new Map(), elicitations: new Map(), errors: new Map(), messages: new Map(), @@ -485,9 +460,11 @@ export function fold( entriesByTurn.set(entry.chunk.turnId, turnEntries); } - const orderedTurns = [...turns.values()].sort((left, right) => - compareTurnOrder(entries, left.id, right.id) - ); + // `entries` arrives pre-sorted by (seq, sub) — see sortConversationEntries + // in @cyrus/utils/conversations/entries — so each Map's insertion order + // (first-touch order, per key) already is the correct display order. No + // resort needed here, by turn or otherwise. + const orderedTurns = [...turns.values()]; orderedTurns.forEach((turn, index) => { turn.index = index; }); @@ -502,42 +479,12 @@ export function fold( ? orderedTurns.find((turn) => turn.id === latestTurnId) : orderedTurns.at(-1); - const turnOrder = new Map( - orderedTurns.map((turn, index) => [turn.id, index] as const) - ); - const parsed = ThreadConversationSchema.safeParse({ - approvals: [...state.approvals.values()].sort((left, right) => { - const leftTurn = turnOrder.get(left.turnId) ?? Number.MAX_SAFE_INTEGER; - const rightTurn = turnOrder.get(right.turnId) ?? Number.MAX_SAFE_INTEGER; - if (leftTurn !== rightTurn) return leftTurn - rightTurn; - return left.createdAt.localeCompare(right.createdAt); - }), - diffs: [...state.diffs.values()], - elicitations: [...state.elicitations.values()].sort((left, right) => { - const leftTurn = turnOrder.get(left.turnId) ?? Number.MAX_SAFE_INTEGER; - const rightTurn = turnOrder.get(right.turnId) ?? Number.MAX_SAFE_INTEGER; - if (leftTurn !== rightTurn) return leftTurn - rightTurn; - return left.createdAt.localeCompare(right.createdAt); - }), - errors: [...state.errors.values()].sort((left, right) => { - const leftTurn = turnOrder.get(left.turnId) ?? Number.MAX_SAFE_INTEGER; - const rightTurn = turnOrder.get(right.turnId) ?? Number.MAX_SAFE_INTEGER; - if (leftTurn !== rightTurn) return leftTurn - rightTurn; - return left.createdAt.localeCompare(right.createdAt); - }), + approvals: [...state.approvals.values()], + elicitations: [...state.elicitations.values()], + errors: [...state.errors.values()], thoughts: [...state.thoughts.values()] .filter((thought) => thought.content.trim().length > 0) - .sort((left, right) => { - const leftTurn = left.turnId - ? (turnOrder.get(left.turnId) ?? Number.MAX_SAFE_INTEGER) - : Number.MAX_SAFE_INTEGER; - const rightTurn = right.turnId - ? (turnOrder.get(right.turnId) ?? Number.MAX_SAFE_INTEGER) - : Number.MAX_SAFE_INTEGER; - if (leftTurn !== rightTurn) return leftTurn - rightTurn; - return left.createdAt.localeCompare(right.createdAt); - }) .map((thought) => ({ ...thought, streaming: @@ -545,28 +492,13 @@ export function fold( latestTurn?.state === "running" && !completedThoughtIds.has(thought.id), })), - messages: [...state.messages.values()] - .sort((left, right) => { - const leftTurn = left.turnId - ? (turnOrder.get(left.turnId) ?? Number.MAX_SAFE_INTEGER) - : Number.MAX_SAFE_INTEGER; - const rightTurn = right.turnId - ? (turnOrder.get(right.turnId) ?? Number.MAX_SAFE_INTEGER) - : Number.MAX_SAFE_INTEGER; - if (leftTurn !== rightTurn) return leftTurn - rightTurn; - if (left.role !== right.role) { - if (left.role === "user") return -1; - if (right.role === "user") return 1; - } - return left.createdAt.localeCompare(right.createdAt); - }) - .map((message) => ({ - ...message, - streaming: - message.role === "assistant" && - message.turnId === latestTurn?.id && - latestTurn?.state === "running", - })), + messages: [...state.messages.values()].map((message) => ({ + ...message, + streaming: + message.role === "assistant" && + message.turnId === latestTurn?.id && + latestTurn?.state === "running", + })), toolCalls: [...state.toolCalls.values()], turns: orderedTurns, }); diff --git a/shared/utils/src/conversations/thread-feed.test.ts b/shared/utils/src/conversations/thread-feed.test.ts index 614f85c1..01fbd751 100644 --- a/shared/utils/src/conversations/thread-feed.test.ts +++ b/shared/utils/src/conversations/thread-feed.test.ts @@ -5,7 +5,6 @@ describe("deriveFeed", () => { test("emits flat error feed entries from thread errors", () => { const feed = deriveFeed({ approvals: [], - diffs: [], elicitations: [], errors: [ { @@ -14,6 +13,7 @@ describe("deriveFeed", () => { id: "error-1", message: "Bind failed", turnId: "turn-1", + seq: 2, }, ], messages: [ @@ -23,6 +23,7 @@ describe("deriveFeed", () => { id: "user-turn-1", role: "user", turnId: "turn-1", + seq: 1, }, ], thoughts: [], @@ -48,10 +49,9 @@ describe("deriveFeed", () => { }); }); - test("interleaves orphaned errors by createdAt", () => { + test("interleaves orphaned errors by seq", () => { const feed = deriveFeed({ approvals: [], - diffs: [], elicitations: [], errors: [ { @@ -59,6 +59,7 @@ describe("deriveFeed", () => { id: "error-orphan", message: "Early bind failed", turnId: "orphan-turn", + seq: 1, }, ], messages: [ @@ -68,6 +69,7 @@ describe("deriveFeed", () => { id: "user-turn-2", role: "user", turnId: "turn-2", + seq: 2, }, ], thoughts: [], @@ -103,17 +105,7 @@ describe("deriveFeed", () => { title: "Write file", toolCallId: "tool-1", turnId: "turn-1", - }, - ], - diffs: [ - { - additions: 1, - deletions: 0, - id: "turn-1:a.ts", - patch: "@@ -0 +1 @@", - path: "a.ts", - toolCallId: "tool-1", - turnId: "turn-1", + seq: 3, }, ], elicitations: [ @@ -126,6 +118,7 @@ describe("deriveFeed", () => { sessionId: "session-1", threadId: "thread-1", turnId: "turn-1", + seq: 4, }, ], errors: [], @@ -136,16 +129,29 @@ describe("deriveFeed", () => { id: "user-turn-1", role: "user", turnId: "turn-1", + seq: 1, }, ], thoughts: [], toolCalls: [ { createdAt: "2026-07-11T00:00:01.500Z", + diffs: [ + { + additions: 1, + deletions: 0, + id: "tool-1:a.ts", + patch: "@@ -0 +1 @@", + path: "a.ts", + toolCallId: "tool-1", + turnId: "turn-1", + }, + ], status: "pending", title: "Write file", toolCallId: "tool-1", turnId: "turn-1", + seq: 2, }, ], turns: [ @@ -163,13 +169,73 @@ describe("deriveFeed", () => { expect(tool).toMatchObject({ type: "tool", pendingApproval: { toolCallId: "tool-1" }, - }); - const diff = feed.find((entry) => entry.type === "diff"); - expect(diff).toMatchObject({ - type: "diff", - pendingApproval: { toolCallId: "tool-1" }, + tool: { diffs: [{ path: "a.ts", toolCallId: "tool-1" }] }, }); expect(feed.some((entry) => entry.type === "approval")).toBe(true); expect(feed.some((entry) => entry.type === "elicitation")).toBe(true); }); + + test("orders a turn's entries purely by (seq, sub), no kind tiebreak needed", () => { + const feed = deriveFeed({ + approvals: [], + elicitations: [], + errors: [], + messages: [ + { + content: "Go", + createdAt: "2026-07-11T00:00:01.000Z", + id: "user-turn-1", + role: "user", + turnId: "turn-1", + seq: 1, + }, + { + content: "Done", + createdAt: "2026-07-11T00:00:01.000Z", + id: "assistant-turn-1", + role: "assistant", + turnId: "turn-1", + seq: 3, + sub: 2, + }, + ], + thoughts: [ + { + content: "Thinking", + createdAt: "2026-07-11T00:00:01.000Z", + id: "thought-1", + turnId: "turn-1", + seq: 3, + sub: 1, + }, + ], + toolCalls: [ + { + createdAt: "2026-07-11T00:00:01.000Z", + diffs: [], + status: "completed", + title: "Edit", + toolCallId: "tool-1", + turnId: "turn-1", + seq: 2, + }, + ], + turns: [ + { + completedAt: "2026-07-11T00:00:02.000Z", + id: "turn-1", + index: 0, + state: "complete", + threadId: "thread-1", + }, + ], + }); + + expect(feed.map((entry) => entry.id)).toEqual([ + "user-turn-1", + "tool-tool-1", + "thought-1", + "assistant-turn-1", + ]); + }); }); diff --git a/shared/utils/src/conversations/thread-feed.ts b/shared/utils/src/conversations/thread-feed.ts index 92bb7101..bbb820f2 100644 --- a/shared/utils/src/conversations/thread-feed.ts +++ b/shared/utils/src/conversations/thread-feed.ts @@ -1,6 +1,5 @@ import type { ApprovalView, - DiffView, ElicitationView, ErrorView, MessageView, @@ -30,13 +29,6 @@ export type ToolFeedEntry = FeedEntryBase & { pendingApproval?: ApprovalView; }; -export type DiffFeedEntry = FeedEntryBase & { - type: "diff"; - diff: DiffView; - turnId: string; - pendingApproval?: ApprovalView; -}; - export type ErrorFeedEntry = FeedEntryBase & { type: "error"; error: ErrorView; @@ -59,20 +51,23 @@ export type FeedEntry = | MessageFeedEntry | ThoughtFeedEntry | ToolFeedEntry - | DiffFeedEntry | ErrorFeedEntry | ApprovalFeedEntry | ElicitationFeedEntry; +type OrderKey = { seq: number; sub: number }; + type TimelineItem = { - createdAt: string; - kind: number; + order: OrderKey; entry: FeedEntry; }; -function messageSortKind(message: MessageView): number { - if (message.role === "user") return 0; - return 2; +function orderKey(entity: { seq: number; sub?: number }): OrderKey { + return { seq: entity.seq, sub: entity.sub ?? 0 }; +} + +function compareOrderKey(left: OrderKey, right: OrderKey): number { + return left.seq - right.seq || left.sub - right.sub; } function findPendingApproval( @@ -101,7 +96,6 @@ function buildTurnTimeline( messages: MessageView[], thoughts: ThoughtView[], toolCalls: ToolCallView[], - diffs: DiffView[], errors: ErrorView[], approvals: ApprovalView[], elicitations: ElicitationView[] @@ -110,24 +104,21 @@ function buildTurnTimeline( pushTurnItems(messages, turnId, (message) => { timeline.push({ - createdAt: message.createdAt, - kind: messageSortKind(message), + order: orderKey(message), entry: { type: "message", id: message.id, message }, }); }); pushTurnItems(thoughts, turnId, (thought) => { timeline.push({ - createdAt: thought.createdAt, - kind: 1, + order: orderKey(thought), entry: { type: "thought", id: thought.id, thought }, }); }); pushTurnItems(toolCalls, turnId, (toolCall) => { timeline.push({ - createdAt: toolCall.createdAt, - kind: 3, + order: orderKey(toolCall), entry: { type: "tool", id: `tool-${toolCall.toolCallId}`, @@ -138,37 +129,9 @@ function buildTurnTimeline( }); }); - const diffSortAnchor = - toolCalls - .filter((toolCall) => toolCall.turnId === turnId) - .reduce( - (latest, toolCall) => - !latest || toolCall.createdAt > latest ? toolCall.createdAt : latest, - undefined - ) ?? - messages.find( - (message) => message.role === "user" && message.turnId === turnId - )?.createdAt ?? - turnId; - - pushTurnItems(diffs, turnId, (diff) => { - timeline.push({ - createdAt: diffSortAnchor, - kind: 4, - entry: { - type: "diff", - id: diff.id, - diff, - turnId, - pendingApproval: findPendingApproval(approvals, diff.toolCallId), - }, - }); - }); - pushTurnItems(approvals, turnId, (approval) => { timeline.push({ - createdAt: approval.createdAt, - kind: 4.5, + order: orderKey(approval), entry: { type: "approval", id: approval.id, @@ -180,8 +143,7 @@ function buildTurnTimeline( pushTurnItems(elicitations, turnId, (elicitation) => { timeline.push({ - createdAt: elicitation.createdAt, - kind: 4.6, + order: orderKey(elicitation), entry: { type: "elicitation", id: elicitation.id, @@ -193,8 +155,7 @@ function buildTurnTimeline( pushTurnItems(errors, turnId, (error) => { timeline.push({ - createdAt: error.createdAt, - kind: 5, + order: orderKey(error), entry: { type: "error", id: error.id, @@ -204,16 +165,7 @@ function buildTurnTimeline( }); }); - timeline.sort((left, right) => { - const leftIsUser = left.kind === 0; - const rightIsUser = right.kind === 0; - if (leftIsUser && !rightIsUser) return -1; - if (!leftIsUser && rightIsUser) return 1; - - const byTime = left.createdAt.localeCompare(right.createdAt); - if (byTime !== 0) return byTime; - return left.kind - right.kind; - }); + timeline.sort((left, right) => compareOrderKey(left.order, right.order)); return timeline.map((item) => item.entry); } @@ -235,7 +187,6 @@ export function deriveFeed( conversation.messages, conversation.thoughts, conversation.toolCalls, - conversation.diffs, conversation.errors, approvals, elicitations @@ -252,7 +203,7 @@ export function deriveFeed( (error) => !knownTurnIds.has(error.turnId) ); for (const error of orphanedErrors) { - insertFeedEntryByCreatedAt(entries, { + insertFeedEntryByOrderKey(entries, orderKey(error), { type: "error", id: error.id, error, @@ -263,22 +214,20 @@ export function deriveFeed( return entries; } -function feedEntryCreatedAt(entry: FeedEntry): string | null { +function feedEntryOrderKey(entry: FeedEntry): OrderKey { switch (entry.type) { case "message": - return entry.message.createdAt; + return orderKey(entry.message); case "thought": - return entry.thought.createdAt; + return orderKey(entry.thought); case "tool": - return entry.tool.createdAt; - case "diff": - return null; + return orderKey(entry.tool); case "error": - return entry.error.createdAt; + return orderKey(entry.error); case "approval": - return entry.approval.createdAt; + return orderKey(entry.approval); case "elicitation": - return entry.elicitation.createdAt; + return orderKey(entry.elicitation); default: { const _exhaustive: never = entry; return _exhaustive; @@ -286,15 +235,14 @@ function feedEntryCreatedAt(entry: FeedEntry): string | null { } } -function insertFeedEntryByCreatedAt( +function insertFeedEntryByOrderKey( entries: FeedEntry[], + target: OrderKey, entry: ErrorFeedEntry ): void { - const createdAt = entry.error.createdAt; let index = 0; for (const existing of entries) { - const existingAt = feedEntryCreatedAt(existing); - if (existingAt !== null && existingAt > createdAt) break; + if (compareOrderKey(feedEntryOrderKey(existing), target) > 0) break; index += 1; } entries.splice(index, 0, entry); From e262d327e4c20d7e344bbc38d33e6141e0ef6cd6 Mon Sep 17 00:00:00 2001 From: Soorya U Date: Wed, 5 Aug 2026 17:50:57 +0530 Subject: [PATCH 2/2] Fix event bubbling from nested diff rows and stale-createdAt latest-turn lookup DiffRow's toggle button click bubbled up into ToolRow's own click handler, closing the parent row (and unmounting the diff) whenever a diff was opened. Stop propagation in DiffRow. latestTurnIdFromEntries still compared createdAt strings after fold.ts moved to trusting entry array position (seq, sub) as the canonical order; reversed or skewed timestamps could pick the wrong turn as latest. Take the last user_message in iteration order instead. Addresses CodeRabbit review comments on PR #140. Co-Authored-By: Claude Sonnet 5 --- .../src/components/chat/work-log/diff-row.tsx | 5 ++- .../chat/work-log/tool-row.test.tsx | 39 ++++++++++++++++++ shared/utils/src/conversations/fold.test.ts | 40 +++++++++++++++++++ shared/utils/src/conversations/fold.ts | 10 +---- 4 files changed, 84 insertions(+), 10 deletions(-) create mode 100644 apps/web/src/components/chat/work-log/tool-row.test.tsx diff --git a/apps/web/src/components/chat/work-log/diff-row.tsx b/apps/web/src/components/chat/work-log/diff-row.tsx index 21dc60ef..05aa2ebc 100644 --- a/apps/web/src/components/chat/work-log/diff-row.tsx +++ b/apps/web/src/components/chat/work-log/diff-row.tsx @@ -11,7 +11,10 @@ export function DiffRow({ diff }: { diff: DiffView }) {