From 81b9b5d3b7697bf6c17176d1d2dc7e5bf47b952b Mon Sep 17 00:00:00 2001 From: Soorya U Date: Fri, 10 Jul 2026 02:35:52 +0530 Subject: [PATCH 1/7] Scope live chat to watched threads with overlay-based client streaming. Replace global peer broadcast with ThreadEventBus, unary chat RPC, and watchThread/unwatchThread so subscribe only fans out to interested peers. Clients merge durable snapshots with a Zustand overlay and await turn completion via waitForTurnEnd instead of mutating the query cache. Co-authored-by: Cursor --- apps/cli/src/commands/service/worker.ts | 4 + apps/cli/src/handlers/controller/chat.ts | 82 +++-- apps/cli/src/handlers/controller/threads.ts | 29 +- apps/cli/src/lib/env.ts | 1 - apps/cli/src/queue/bus.ts | 179 +++++++++ .../components/chat/main/thread-workspace.tsx | 13 +- .../.openspec.yaml | 2 + .../design.md | 137 +++++++ .../proposal.md | 47 +++ .../specs/conversation-overlay/spec.md | 81 +++++ .../specs/conversation-view/spec.md | 22 ++ .../specs/thread-live-stream/spec.md | 85 +++++ .../specs/wire-schemas/spec.md | 33 ++ .../tasks.md | 48 +++ openspec/specs/conversation-overlay/spec.md | 85 +++++ openspec/specs/conversation-view/spec.md | 21 ++ openspec/specs/thread-live-stream/spec.md | 89 +++++ openspec/specs/wire-schemas/spec.md | 32 ++ .../connections/src/contracts/controller.ts | 8 +- shared/connections/src/rtc/bus.ts | 12 + shared/connections/src/rtc/peer.ts | 15 +- shared/connections/src/rtc/worker/index.ts | 14 +- shared/constants/src/operation-keys.ts | 2 + .../src/repositories/conversations.ts | 14 +- .../src/connection/use-controller-threads.ts | 11 +- .../src/connection/use-thread-conversation.ts | 102 +++++- .../use-worker-conversation-sync.ts | 36 +- .../hooks/src/stores/conversation-overlay.ts | 342 ++++++++++++++++++ shared/schemas/src/rtc/chat.ts | 7 + shared/schemas/src/rtc/threads.ts | 15 + shared/utils/src/conversation-cache.ts | 24 -- shared/utils/src/merge-conversation.ts | 21 ++ 32 files changed, 1502 insertions(+), 111 deletions(-) create mode 100644 apps/cli/src/queue/bus.ts create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/.openspec.yaml create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-overlay/spec.md create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-view/spec.md create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/wire-schemas/spec.md create mode 100644 openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md create mode 100644 openspec/specs/conversation-overlay/spec.md create mode 100644 openspec/specs/thread-live-stream/spec.md create mode 100644 shared/connections/src/rtc/bus.ts create mode 100644 shared/hooks/src/stores/conversation-overlay.ts delete mode 100644 shared/utils/src/conversation-cache.ts create mode 100644 shared/utils/src/merge-conversation.ts diff --git a/apps/cli/src/commands/service/worker.ts b/apps/cli/src/commands/service/worker.ts index fdd0f623..a8a499f0 100644 --- a/apps/cli/src/commands/service/worker.ts +++ b/apps/cli/src/commands/service/worker.ts @@ -8,6 +8,7 @@ import { createControllerRouter } from "@/handlers/controller"; import { workerRouter } from "@/handlers/worker"; import { authClient } from "@/lib/auth"; import { env } from "@/lib/env"; +import { createThreadEventBus } from "@/queue/bus"; import { get, getOrCreate } from "@/store/config"; import { initDatabase } from "@/store/database"; import { print } from "@/utils/style"; @@ -48,9 +49,12 @@ export async function worker(): Promise { }); print.success`✓ connected — waiting for message…`; + const eventBus = createThreadEventBus(); + const device = serveWorker({ signaling: signalingSession.signaling, events: signalingSession.events, + eventBus, routers: { controller: createControllerRouter(runtime), worker: workerRouter, diff --git a/apps/cli/src/handlers/controller/chat.ts b/apps/cli/src/handlers/controller/chat.ts index 0ae0d0d5..441c6d0a 100644 --- a/apps/cli/src/handlers/controller/chat.ts +++ b/apps/cli/src/handlers/controller/chat.ts @@ -2,7 +2,7 @@ import { appendConversation } from "@cyrus/database/repositories/conversations"; import { ensureThread } from "@cyrus/database/repositories/threads"; import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import { randomId } from "@cyrus/utils/identity"; -import { env } from "@/lib/env"; +import { Result } from "better-result"; import { throwOrpcFromRepositoryError } from "@/utils/error"; import { isStreamingDelta, @@ -11,9 +11,46 @@ import { } from "@/utils/streams"; import type { ControllerDeps } from "./deps"; +type RunTurnOptions = { + agentName: string; + threadId: string; + projectId: string; + message: string; + turnId: string; + emit: (event: ChatChunk["event"]) => Promise; + runtime: ControllerDeps["runtime"]; +}; + +async function runTurn({ + agentName, + threadId, + projectId, + message, + emit, + runtime, +}: Omit): Promise> { + await emit({ type: "user_message", content: message }); + await emit({ type: "thread_started", threadId }); + + const streamed = await Result.tryPromise(async () => { + const gen = runtime.threadCoordinator.prompt( + agentName, + threadId, + projectId, + message + ); + for await (const event of gen) await emit(event); + await emit({ type: "turn_completed" }); + }); + + if (streamed.isErr()) await emit({ type: "turn_interrupted" }); + + return streamed; +} + export function chatHandlers({ os, runtime }: ControllerDeps) { return { - chat: os.chat.handler(async function* ({ input, context }) { + chat: os.chat.handler(async ({ input, context }) => { const { agentName, threadId = randomId(), message, projectId } = input; const thread = await ensureThread(threadId, projectId, { @@ -26,13 +63,15 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { const messageBuffers = new Map(); const thoughtBuffers = new Map(); - async function emit(event: ChatChunk["event"]): Promise { + context.eventBus.ensureWatch(context.peerId, threadId); + + async function emit(event: ChatChunk["event"]): Promise { trackDelta(event, messageBuffers, thoughtBuffers); if (isStreamingDelta(event)) { const chunk: ChatChunk = { threadId, turnId, seq: 0, event }; - context.broadcaster.broadcast(chunk, context.peerId); - return chunk; + context.eventBus.publish(chunk); + return; } const persistEvent = resolvePersistEvent( @@ -46,34 +85,27 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { event: persistEvent, }); if (entry.isErr()) throwOrpcFromRepositoryError(entry.error); - const chunk = entry.value.chunk; - context.broadcaster.broadcast(chunk, context.peerId); - return chunk; + context.eventBus.publish(entry.value.chunk); } - yield await emit({ type: "user_message", content: message }); - yield await emit({ type: "thread_started", threadId }); - - const gen = runtime.threadCoordinator.prompt( + runTurn({ agentName, threadId, projectId, - message - ); - try { - for await (const event of gen) { - yield await emit(event); - await Bun.sleep(env.CYRUS_STREAM_THROTTLING_MS); - } - yield await emit({ type: "turn_completed" }); - } catch (error) { - yield await emit({ type: "turn_interrupted" }); - throw error; - } + message, + emit, + runtime, + }).then((result) => { + result.tapError((error) => { + console.error("chat turn failed", error); + }); + }); + + return { threadId, turnId }; }), subscribe: os.subscribe.handler(async function* ({ context }) { - for await (const chunk of context.broadcaster.subscribe(context.peerId)) + for await (const chunk of context.eventBus.subscribe(context.peerId)) yield chunk; }), diff --git a/apps/cli/src/handlers/controller/threads.ts b/apps/cli/src/handlers/controller/threads.ts index 633bb11b..3850ba67 100644 --- a/apps/cli/src/handlers/controller/threads.ts +++ b/apps/cli/src/handlers/controller/threads.ts @@ -1,4 +1,7 @@ -import { getConversations } from "@cyrus/database/repositories/conversations"; +import { + getConversations, + getSnapshotHighWaterMark, +} from "@cyrus/database/repositories/conversations"; import { createThread as createStoredThread, deleteThread, @@ -62,5 +65,29 @@ export function threadsHandlers(os: ControllerOs) { return {}; }), + + watchThread: os.watchThread.handler(async ({ input, context }) => { + const thread = await getThread(input.threadId); + if (thread.isErr()) throwOrpcFromRepositoryError(thread.error); + if (!thread.value) { + throwOrpcFromRepositoryError({ + type: "not_found", + entity: "thread", + id: input.threadId, + }); + } + + context.eventBus.watch(context.peerId, input.threadId); + + return (await getSnapshotHighWaterMark(input.threadId)).match({ + ok: (snapshotHighWaterMark) => ({ snapshotHighWaterMark }), + err: throwOrpcFromRepositoryError, + }); + }), + + unwatchThread: os.unwatchThread.handler(({ input, context }) => { + context.eventBus.unwatch(context.peerId, input.threadId); + return {}; + }), }; } diff --git a/apps/cli/src/lib/env.ts b/apps/cli/src/lib/env.ts index ac2c213a..bc23e532 100644 --- a/apps/cli/src/lib/env.ts +++ b/apps/cli/src/lib/env.ts @@ -25,7 +25,6 @@ export const env = createEnv({ .int() .positive() .default(30 * 60 * 1000), - CYRUS_STREAM_THROTTLING_MS: z.coerce.number().positive().default(25), }, runtimeEnv: process.env, emptyStringAsUndefined: true, diff --git a/apps/cli/src/queue/bus.ts b/apps/cli/src/queue/bus.ts new file mode 100644 index 00000000..b1092452 --- /dev/null +++ b/apps/cli/src/queue/bus.ts @@ -0,0 +1,179 @@ +import type { ThreadEventBus } from "@cyrus/connections/rtc/thread-event-bus"; +import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; + +const DEFAULT_MAX_CHUNKS_PER_TURN = 10_000; + +type PeerDelivery = { + queue: ChatChunk[]; + resolve: (() => void) | null; + closed: boolean; +}; + +export type CreateThreadEventBusOptions = { + maxChunksPerTurn?: number; +}; + +function isTerminalEvent(event: ChatChunk["event"]): boolean { + return event.type === "turn_completed" || event.type === "turn_interrupted"; +} + +export function createThreadEventBus( + options: CreateThreadEventBusOptions = {} +): ThreadEventBus { + const { maxChunksPerTurn = DEFAULT_MAX_CHUNKS_PER_TURN } = options; + + const peers = new Map(); + const watchedThreads = new Map>(); + const activeTurnLogs = new Map(); + const turnThreads = new Map(); + + function getWatchedThreads(peerId: string): Set { + let set = watchedThreads.get(peerId); + if (!set) { + set = new Set(); + watchedThreads.set(peerId, set); + } + return set; + } + + function closePeerDelivery(peerId: string): void { + const peer = peers.get(peerId); + if (peer) { + peer.closed = true; + peer.resolve?.(); + } + } + + function deliver(peerId: string, chunk: ChatChunk): void { + const peer = peers.get(peerId); + if (!peer || peer.closed) return; + peer.queue.push(chunk); + peer.resolve?.(); + peer.resolve = null; + } + + function fanOut(chunk: ChatChunk): void { + for (const [peerId, threads] of watchedThreads) { + if (threads.has(chunk.threadId)) { + deliver(peerId, chunk); + } + } + } + + 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(); + } + + function appendToTurnLog(chunk: ChatChunk): void { + const { turnId } = chunk; + let log = activeTurnLogs.get(turnId); + if (!log) { + log = []; + activeTurnLogs.set(turnId, log); + turnThreads.set(turnId, chunk.threadId); + } + log.push(chunk); + trimTurnLog(log); + } + + function evictTurnLog(turnId: string): void { + activeTurnLogs.delete(turnId); + turnThreads.delete(turnId); + } + + function replayThread(peerId: string, threadId: string): void { + for (const [turnId, log] of activeTurnLogs) { + if (turnThreads.get(turnId) !== threadId) continue; + for (const chunk of log) deliver(peerId, chunk); + } + } + + function registerWatch(peerId: string, threadId: string): void { + const threads = getWatchedThreads(peerId); + const isNew = !threads.has(threadId); + threads.add(threadId); + if (isNew) replayThread(peerId, threadId); + } + + return { + publish(chunk) { + const terminal = isTerminalEvent(chunk.event); + if (!terminal) { + appendToTurnLog(chunk); + } + fanOut(chunk); + if (terminal) { + evictTurnLog(chunk.turnId); + } + }, + + watch(peerId, threadId) { + registerWatch(peerId, threadId); + }, + + unwatch(peerId, threadId) { + watchedThreads.get(peerId)?.delete(threadId); + }, + + ensureWatch(peerId, threadId) { + if (!getWatchedThreads(peerId).has(threadId)) { + registerWatch(peerId, threadId); + } + }, + + isWatching(peerId, threadId) { + return watchedThreads.get(peerId)?.has(threadId) ?? false; + }, + + async *subscribe(peerId) { + closePeerDelivery(peerId); + const peer: PeerDelivery = { + queue: [], + resolve: null, + closed: false, + }; + peers.set(peerId, peer); + + for (const threadId of getWatchedThreads(peerId)) { + replayThread(peerId, threadId); + } + + try { + while (!peer.closed) { + while (peer.queue.length > 0) { + yield peer.queue.shift() as ChatChunk; + } + if (!peer.closed) { + await new Promise((resolve) => { + peer.resolve = resolve; + }); + } + } + } finally { + if (peers.get(peerId) === peer) { + peers.delete(peerId); + } + } + }, + + close(peerId) { + closePeerDelivery(peerId); + watchedThreads.delete(peerId); + }, + + closeAll() { + for (const peerId of [...peers.keys()]) { + closePeerDelivery(peerId); + } + peers.clear(); + watchedThreads.clear(); + activeTurnLogs.clear(); + turnThreads.clear(); + }, + }; +} diff --git a/apps/web/src/components/chat/main/thread-workspace.tsx b/apps/web/src/components/chat/main/thread-workspace.tsx index b2915f8b..0df7cd99 100644 --- a/apps/web/src/components/chat/main/thread-workspace.tsx +++ b/apps/web/src/components/chat/main/thread-workspace.tsx @@ -24,8 +24,7 @@ export function ThreadWorkspace({ threadId, }: ThreadWorkspaceProps) { const navigate = useNavigate(); - const { threads, sendMessage, stopThread, createThread } = - useControllerThreads(); + const { threads, sendMessage, stopThread } = useControllerThreads(); const { diffOpen, setDiffOpen } = useChatUiStore(); const baseThread = threads.find((item) => item.id === threadId) ?? null; @@ -49,15 +48,7 @@ export function ThreadWorkspace({ }, [navigate, projectId, thread, workerId]); async function handleSend(text: string) { - if (!thread) { - const id = await createThread(projectId); - await sendMessage(id, text); - navigate({ - to: "/workers/$workerId/p/$projectId/t/$threadId", - params: { workerId, projectId, threadId: id }, - }); - return; - } + if (!thread) return; await sendMessage(thread.id, text); } diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/.openspec.yaml b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/.openspec.yaml new file mode 100644 index 00000000..074342d5 --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-07-09 diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md new file mode 100644 index 00000000..2d769a08 --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md @@ -0,0 +1,137 @@ +## Context + +`PeerBroadcaster` (`shared/connections/src/rtc/broadcaster.ts`) pushes every `ChatChunk` to every connected peer except the sender (`id !== fromPeerId`). The client compensates with two ingress paths: `sendMessage` (`shared/hooks/src/connection/use-controller-threads.ts`) consumes `chat()`'s event iterator directly, while `useWorkerConversationSync` consumes `subscribe()`. Both write into the same `getConversations` query cache via `appendChunkToCache` (`shared/utils/src/conversation-cache.ts`). + +PR #22 (`connection-providers`) moved connection hooks to `@cyrus/hooks`, query keys to `@cyrus/constants`, and RTC bootstrapping to `RtcProvider` + `useRtc()`. This change builds on that layout — overlay and hook refactors target `shared/hooks`, not `apps/web` directly. + +Issue #13 (merged, PR #18) added Turso persistence, monotonic `seq` on persisted entries, `getConversations(afterSeq)`, and coalesced message persistence (token deltas are ephemeral, `seq: 0`). The server-side `messageBuffers` in `chat.ts` exist only for coalescing within a single `chat()` generator — they are not replayable. + +Mid-stream page refresh loses all in-flight token deltas: Turso has only persisted events, `PeerBroadcaster` queues are per-peer and ephemeral, and the client overlay (today: mutated query cache) is gone on reload. + +## Goals / Non-Goals + +**Goals:** + +- Single client ingress: all live chunks flow through `subscribe()` → overlay store. +- Single server egress path: `emit()` publishes to a thread-scoped event bus; `chat()` is a command, not a chunk stream. +- Scope delivery to peers watching the relevant thread, not all connected peers. +- Replay in-flight turn chunks to late joiners (refresh, navigate-away-and-back, second device opening same thread) via a bounded server-side active-turn log. +- Client merge of durable snapshot (`getConversations`) + live overlay with `seq`-based dedup. +- `sendMessage` completion via overlay observation of `turn_completed` / `turn_interrupted` for the returned `turnId` (option B). + +**Non-Goals:** + +- Backward compatibility for `chat` as `eventIterator` or dual client ingress. +- Persisting per-token deltas to Turso. +- Adopting `@tanstack/ai`. +- Chat-first thread creation (sidebar-first unchanged; `ensureThread` first-message rename stays). + +## Decisions + +### Replace `PeerBroadcaster` with `ThreadEventBus` + +**Choice:** Custom `ThreadEventBus` in `apps/cli/src/queue/index.ts` owns watch sets, per-peer delivery queues, and active-turn logs. No third-party event bus library — the API is domain-specific (thread watch sets, per-peer `AsyncGenerator`, active-turn replay). + +**Why over patching `PeerBroadcaster`:** Watch filtering and turn replay are thread-scoped concerns, not peer-queue concerns. A clean type keeps `subscribe(peerId)` as the per-peer stream API while moving fan-out logic to `publish(chunk)`. + +**Structure:** + +``` +ThreadEventBus +├── watchers: Map> +├── peers: Map // subscribe delivery +├── activeTurnLogs: Map // replay buffer +└── publish(chunk): + 1. if turn not terminal → append to activeTurnLogs[turnId] + 2. on turn_completed/turn_interrupted → evict activeTurnLogs[turnId] after publish + 3. for each peer where chunk.threadId ∈ watchers[peerId] → enqueue (INCLUDING sender) +``` + +**Alternative considered:** Per-thread Kafka / durable log. Rejected — bounded in-memory retention until `turn_completed` is sufficient; Turso already holds durable events. + +### `chat()` becomes a unary command + +**Choice:** `chat` RPC returns `{ threadId, turnId }` immediately after starting the turn. The handler runs the agent loop in the background (or continues the existing async generator internally without yielding chunks to the RPC caller). + +**Why:** Eliminates the sender exclusion hack. All peers, including the initiator, receive chunks through `subscribe()`. Matches TanStack AI's `send` + `subscribe` pattern. + +**ORPC shape:** + +```typescript +chat: oc.input(ChatInputSchema).output(ChatOutputSchema), +// ChatOutputSchema = { threadId: string, turnId: string } + +watchThread: oc.input(WatchThreadInputSchema).output(WatchThreadOutputSchema), +unwatchThread: oc.input(UnwatchThreadInputSchema).output(VoidOutputSchema), +``` + +### `watchThread` replays active-turn log before live fan-out + +**Choice:** `watchThread(peerId, threadId)` registers the watch set entry, synchronously replays `activeTurnLogs` for all in-progress turns on that thread into the peer's queue, and returns `{ snapshotHighWaterMark: number }` (max persisted `seq` for the thread, queried from Turso). + +**Join sequence on client:** + +1. `watchThread(threadId)` — register + replay + cursor +2. `getConversations(threadId)` — snapshot → `overlay.applySnapshot()` +3. Live chunks via `subscribe()` → `overlay.applyLiveChunk()` + +Overlap between replay and snapshot is expected; client dedup by `seq` resolves it (t3-code pattern from issue #15 comment). + +**Alternative considered:** `getConversations` first, then `watchThread`. Rejected as primary order — replay from active-turn log is required for ephemeral deltas regardless of snapshot timing. + +### Client overlay store (Zustand) + +**Choice:** `shared/hooks/src/stores/conversation-overlay.ts` holds per-thread live chunks and `snapshotHighWaterMark`. Single writer: `useWorkerConversationSync` in `shared/hooks/src/connection/`. Hooks consume RTC via `useRtc()` from `@cyrus/hooks/contexts/rtc` (provided by `RtcProvider`). + +**Dedup rules:** + +- `seq > 0` and `seq <= snapshotHighWaterMark` → drop (already in snapshot) +- `seq > 0` and duplicate in overlay → drop +- `seq === 0` (ephemeral delta) → accept while turn is active + +**On `turn_completed` / `turn_interrupted`:** remove turn from `activeTurnIds`, invalidate `getConversations`, on refetch call `applySnapshot()` which bumps watermark and prunes overlay entries with `seq <= watermark`. + +### `sendMessage` completion via overlay (option B) + +**Choice:** After `await client.chat({...})` returns `{ turnId }`, `sendMessage` awaits a promise that resolves when the overlay receives `turn_completed` or `turn_interrupted` for that `turnId` on that `threadId`. Rejects on `turn_interrupted` or timeout. + +**Why over unary `chat` blocking until turn done:** Keeps `chat` RPC fast to acknowledge; UI already updates via subscribe. Avoids long-held RPC connections. Completion is observable on the same stream the UI uses. + +**Implementation:** `waitForTurnEnd(threadId, turnId, signal?)` exported from overlay store, backed by a `Map` registered before calling `chat()`. + +### Remove `appendChunkToCache` + +**Choice:** Delete `shared/utils/src/conversation-cache.ts`. `getConversations` query cache is snapshot-only; overlay is the sole live layer. + +### Auto-watch on `chat()` for sender + +**Choice:** When `chat()` starts a turn, the server auto-registers `threadId` in the caller's watch set if not already present. Client still calls `watchThread` on `ThreadWorkspace` mount for the general case (navigate to existing thread, second device). Auto-watch prevents a race where the sender calls `chat()` before `watchThread` returns. + +### Keep `subscribe()` as one stream per peer + +**Choice:** No per-thread subscribe streams. Filtering happens inside `publish()`. Preserves the existing "one active subscription per peer" constraint (`close(peerId)` on re-subscribe). + +## Risks / Trade-offs + +- **[Risk]** Active-turn log grows if a turn never completes (agent hang). → **Mitigation:** Cap log size per turn (e.g. 10k chunks); on cap, evict oldest deltas but keep terminal events. `cancel` still emits `turn_interrupted`. +- **[Risk]** `sendMessage` waits forever if subscribe disconnects mid-turn. → **Mitigation:** `waitForTurnEnd` respects `AbortSignal` from stop/cancel; timeout fallback invalidates snapshot and clears overlay turn. +- **[Risk]** Duplicate chunks on join from replay + snapshot overlap. → **Mitigation:** `seq`-based dedup in overlay; snapshot wins on `seq` collision. +- **[Risk]** Background `chat()` handler errors invisible to client if subscribe is down. → **Mitigation:** `turn_interrupted` still publishes to bus; if subscribe reconnects, replay may be empty but snapshot has persisted state. Log server-side errors. +- **[Risk]** Breaking `chat` eventIterator affects any non-web consumer. → **Mitigation:** Explicit non-goal; workspace is versioned together, no external clients. + +## Migration Plan + +No data migration. Deploy server and web together (no backward compatibility). + +1. Ship `ThreadEventBus` + watch RPCs + `chat` command shape on CLI worker. +2. Ship overlay store + refactored hooks in `@cyrus/hooks`. +3. Remove `appendChunkToCache` and `shared/utils/src/conversation-cache.ts`. +4. Manual test matrix: web (primary), verify mobile inherits shared hooks behavior where RTC is mounted. + +Rollback: revert both server and web; old dual-path client is incompatible with new server `chat` shape. + +## Open Questions + +- Exact per-turn log cap (10k chunks vs byte-size limit) — **decided: 10k default, delta-first eviction** (`apps/cli/src/queue/index.ts`). +- Should `waitForTurnEnd` timeout be configurable or a fixed constant (e.g. 30 min)? +- Mobile (`apps/mobile`) — inherits shared `@cyrus/hooks` overlay once RTC is mounted; no separate overlay implementation needed. diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md new file mode 100644 index 00000000..ca53b1e9 --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md @@ -0,0 +1,47 @@ +## Why + +Every connected peer currently receives every chat chunk from every thread, and the web client maintains two parallel ingress paths (`chat()` iterator for the sender, `subscribe()` for everyone else) that both mutate the same query cache. This wastes bandwidth, complicates dedup, and loses in-flight streaming on page refresh because ephemeral token deltas are neither persisted nor replayed. Issue #13 shipped `seq`/cursor primitives; this change consumes them to scope live delivery, unify on a single subscribe stream, and add server-side active-turn replay for mid-stream joins. + +## What Changes + +- Replace `PeerBroadcaster`'s per-peer blind fan-out with a `ThreadEventBus`: per-peer watched-thread sets, fan-out to all watchers **including the sender**, and a bounded in-memory active-turn log for replay on `watchThread`. +- Add `watchThread` / `unwatchThread` RPCs; client calls them from thread mount/unmount. +- **BREAKING**: Change `chat` from an `eventIterator` to a unary RPC returning `{ threadId, turnId }`. All chunks flow through `subscribe()` only. +- Add a client-side conversation overlay store (Zustand) in `@cyrus/hooks` holding ephemeral/live chunks separate from the `getConversations` TanStack Query snapshot. +- Refactor `useThreadConversation` to merge snapshot + overlay with `seq`-based dedup, then `fold()`. +- `sendMessage` fires `chat()` and awaits turn completion by watching the overlay for `turn_completed` / `turn_interrupted` matching `turnId` (option B). +- Remove `appendChunkToCache` and the `sendMessage` `for await` loop over `chat()`. +- Delete dead `handleSend` `!thread` branch in `thread-workspace.tsx` (sidebar-first model unchanged; `ensureThread` first-message rename stays). + +## Capabilities + +### New Capabilities + +- `thread-live-stream`: Server-side thread-scoped event bus, watch registration, active-turn replay buffer, and scoped fan-out replacing global broadcast. +- `conversation-overlay`: Client-side live overlay store, snapshot/overlay merge, and single subscribe ingress for all live chunks. + +### Modified Capabilities + +- `wire-schemas`: `chat` output shape changes from `eventIterator(ChatChunk)` to a completion ack; add `watchThread` / `unwatchThread` input/output schemas. +- `conversation-view`: `useThreadConversation` derives view from merged snapshot + overlay instead of a single mutated query cache. + +## Impact + +- `apps/cli/src/queue/index.ts` — custom `ThreadEventBus` (watch sets, active-turn logs, scoped fan-out). +- `shared/connections/src/rtc/broadcaster.ts` — superseded for chat delivery once wired; may remain for other uses or be removed during integration. +- `shared/connections/src/rtc/worker/index.ts` — instantiate `ThreadEventBus` instead of `createPeerBroadcaster`. +- `apps/cli/src/handlers/controller/chat.ts` — `chat` becomes command; `emit` publishes to bus. +- `shared/connections/src/contracts/controller.ts` — breaking `chat` contract; new watch RPCs. +- `shared/hooks/src/connection/{use-controller-threads,use-worker-conversation-sync,use-thread-conversation}.ts` — unified subscribe + overlay; consume RTC via `useRtc()`. +- `shared/hooks/src/stores/conversation-overlay.ts` — new overlay store (alongside `agent-catalog`). +- `shared/utils/src/conversation-cache.ts` — removed. +- `shared/constants/src/operation-keys.ts` — optional `watchThread` / `unwatchThread` mutation keys. +- Closes GitHub issue #15. + +## Non-goals + +- No backward compatibility for the old `chat` eventIterator shape or dual client ingress paths. +- No persistence of per-token deltas to Turso (coalesced message persistence from #13 stays). +- No cross-device Turso sync or remote replication. +- No adoption of `@tanstack/ai` as a dependency. +- No chat-first thread creation; sidebar-first `createThread` flow is unchanged. diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-overlay/spec.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-overlay/spec.md new file mode 100644 index 00000000..01a8af7c --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-overlay/spec.md @@ -0,0 +1,81 @@ +## ADDED Requirements + +### Requirement: Conversation overlay store + +The client (`@cyrus/hooks`) SHALL maintain a `conversation-overlay` store (Zustand) in `shared/hooks/src/stores/conversation-overlay.ts` holding per-thread live chunks separate from the TanStack Query `getConversations` snapshot. The overlay SHALL track `snapshotHighWaterMark` and `activeTurnIds` per thread. + +#### Scenario: Live chunk is stored in overlay + +- **WHEN** a `ChatChunk` arrives on `subscribe()` for `threadId` T +- **THEN** the chunk is appended to the overlay for T +- **AND** the `getConversations` query cache is not mutated + +#### Scenario: Snapshot seeds watermark + +- **WHEN** `getConversations` returns entries for thread T +- **THEN** `snapshotHighWaterMark` for T is set to the maximum `seq` among returned entries +- **AND** overlay entries with `seq > 0` and `seq <= snapshotHighWaterMark` are pruned + +### Requirement: Single subscribe ingress + +All live `ChatChunk`s on the client SHALL be written exclusively by `useWorkerConversationSync` (`shared/hooks/src/connection/`) into the overlay store. No other code path SHALL append live chunks to the query cache. + +#### Scenario: sendMessage does not consume chat iterator + +- **WHEN** `sendMessage` is called +- **THEN** it invokes `chat()` as a unary RPC +- **AND** it does not iterate an event stream from `chat()` + +### Requirement: Seq-based dedup in overlay + +The overlay SHALL reject duplicate or stale chunks using `seq`: + +- Chunks with `seq > 0` and `seq <= snapshotHighWaterMark` SHALL be dropped. +- Chunks with `seq > 0` already present in the overlay SHALL be dropped. +- Chunks with `seq === 0` SHALL be accepted while their `turnId` is active. + +#### Scenario: Snapshot overlap dedup + +- **WHEN** a persisted chunk with `seq` 42 arrives on `subscribe()` +- **AND** `snapshotHighWaterMark` is already 42 or higher +- **THEN** the chunk is not added to the overlay + +#### Scenario: Ephemeral delta accepted during active turn + +- **WHEN** a `token` delta with `seq === 0` arrives for an active `turnId` +- **THEN** the chunk is added to the overlay + +### Requirement: Turn completion reconciles overlay with snapshot + +When the overlay receives `turn_completed` or `turn_interrupted` for a `turnId`, it SHALL remove that turn from `activeTurnIds`, invalidate the `getConversations` query for the thread, and prune overlay entries once the refetched snapshot updates the watermark. + +#### Scenario: Turn completion triggers snapshot refresh + +- **WHEN** `turn_completed` is received for `turnId` T on thread X +- **THEN** `getConversations(X)` is invalidated +- **AND** after refetch, overlay ephemeral entries for T are pruned against the new watermark + +### Requirement: sendMessage awaits turn end via overlay (option B) + +`sendMessage` (`shared/hooks/src/connection/use-controller-threads.ts`) SHALL register a `waitForTurnEnd(threadId, turnId)` promise before calling `chat()`. It SHALL resolve when the overlay receives `turn_completed` for that `turnId`, or reject on `turn_interrupted` or abort. + +#### Scenario: sendMessage completes on turn_completed + +- **WHEN** `sendMessage` calls `chat()` and receives `{ turnId }` +- **AND** the overlay subsequently receives `turn_completed` for that `turnId` +- **THEN** `sendMessage` resolves successfully + +#### Scenario: sendMessage rejects on turn_interrupted + +- **WHEN** the overlay receives `turn_interrupted` for the awaited `turnId` +- **THEN** `sendMessage` rejects with an error + +### Requirement: watchThread on thread mount + +`useThreadConversation` (or a hook it composes) SHALL call `watchThread` when a thread workspace mounts and `unwatchThread` on unmount, before or in parallel with the `getConversations` fetch. + +#### Scenario: Mount registers watch + +- **WHEN** `ThreadWorkspace` mounts for `threadId` T +- **THEN** `watchThread({ threadId: T })` is called +- **AND** `unwatchThread({ threadId: T })` is called on unmount diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-view/spec.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-view/spec.md new file mode 100644 index 00000000..ed42ab94 --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/conversation-view/spec.md @@ -0,0 +1,22 @@ +## ADDED Requirements + +### Requirement: Thread conversation merges snapshot and overlay + +`useThreadConversation` (`shared/hooks/src/connection/use-thread-conversation.ts`) SHALL derive `ThreadConversation` by merging the `getConversations` snapshot with live overlay entries, then calling `fold()` on the merged `ConversationEntry[]`. The merge SHALL prefer snapshot entries on `seq` collision. The hook SHALL consume RTC via `useRtc()`. + +#### Scenario: Merged view includes live tokens + +- **WHEN** the snapshot contains a `user_message` for the latest turn +- **AND** the overlay contains `token` deltas for that turn +- **THEN** `useThreadConversation` returns a `ThreadConversation` with partial assistant text visible before `message_completed` is persisted + +#### Scenario: Snapshot wins on seq collision + +- **WHEN** the snapshot and overlay both contain an entry with `seq` 42 +- **THEN** the snapshot entry is used in the merge +- **AND** the overlay duplicate is excluded + +#### Scenario: Fold is unchanged + +- **WHEN** merged entries are passed to `fold()` +- **THEN** the existing `fold()` implementation in `@cyrus/utils` produces the view without modification to its logic diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md new file mode 100644 index 00000000..d0c3a40d --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md @@ -0,0 +1,85 @@ +## ADDED Requirements + +### Requirement: Thread-scoped event bus replaces global broadcast + +The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` implemented in `apps/cli/src/queue/index.ts` that fans out only to peers whose watched-thread set includes `chunk.threadId`. The bus SHALL deliver chunks to the initiating peer as well as all other watching peers. + +#### Scenario: Non-watching peer receives nothing + +- **WHEN** peer B is connected but has not called `watchThread` for thread A +- **AND** a chat turn is active on thread A +- **THEN** peer B's `subscribe()` stream does not receive chunks for thread A + +#### Scenario: Watching peer receives live chunks + +- **WHEN** peer C has called `watchThread` for thread A +- **AND** a chat turn emits chunks on thread A +- **THEN** peer C's `subscribe()` stream receives those chunks in emission order + +#### Scenario: Initiating peer receives its own chunks via subscribe + +- **WHEN** peer A calls `chat()` for thread A +- **THEN** peer A receives all chunks for that turn through `subscribe()`, not through the `chat()` RPC response + +### Requirement: Active-turn replay buffer + +The `ThreadEventBus` in `apps/cli/src/queue/index.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. + +#### Scenario: Late watcher receives in-flight deltas + +- **WHEN** a turn is in progress on thread A with partial token deltas already emitted +- **AND** a peer calls `watchThread` for thread A before the turn completes +- **THEN** the peer receives a replay of all buffered chunks for active turns on thread A before subsequent live chunks + +#### Scenario: Completed turn log is evicted + +- **WHEN** `turn_completed` is published for `turnId` T +- **THEN** `activeTurnLogs` for T is removed +- **AND** a subsequent `watchThread` does not replay chunks from turn T (durable history comes from `getConversations`) + +### Requirement: watchThread and unwatchThread RPCs + +The controller contract SHALL expose `watchThread` and `unwatchThread` operations. `watchThread` SHALL register the calling peer's interest in a thread, replay active-turn logs for that thread, and return a `snapshotHighWaterMark` (the highest persisted `seq` for that thread). `unwatchThread` SHALL remove the registration. + +#### Scenario: watchThread returns cursor + +- **WHEN** a peer calls `watchThread({ threadId })` +- **THEN** the peer is registered as watching `threadId` +- **AND** the response includes `snapshotHighWaterMark` equal to the max persisted `seq` for that thread (or 0 if none) + +#### Scenario: unwatchThread stops delivery + +- **WHEN** a peer calls `unwatchThread({ threadId })` +- **THEN** subsequent chunks for `threadId` are not delivered to that peer's `subscribe()` stream + +### Requirement: chat is a turn-start command + +The `chat` RPC SHALL accept `ChatInput` and return `{ threadId, turnId }` without streaming `ChatChunk`s. The worker SHALL run the agent turn asynchronously, publishing all chunks via `ThreadEventBus.publish()`. + +#### Scenario: chat returns turn identity + +- **WHEN** a peer calls `chat({ agentName, message, projectId, threadId })` +- **THEN** the RPC returns `{ threadId, turnId }` before the turn completes +- **AND** all chunks for the turn arrive on `subscribe()` + +#### Scenario: chat auto-watches sender thread + +- **WHEN** a peer calls `chat()` for `threadId` T +- **AND** the peer is not already watching T +- **THEN** the server registers T in that peer's watched-thread set before publishing the first chunk + +### Requirement: emit publishes through the bus + +`chat.ts`'s `emit()` function SHALL call `ThreadEventBus.publish(chunk)` for all events. Persisted events SHALL still be written to Turso before publish (persist-before-broadcast from conversation-persistence). Ephemeral deltas (`seq === 0`) SHALL be published without a Turso write. + +#### Scenario: Persisted event carries seq on publish + +- **WHEN** a non-delta event is emitted +- **THEN** the chunk is persisted to Turso and assigned a `seq > 0` before `publish()` +- **AND** the published chunk includes that `seq` + +#### Scenario: Delta is ephemeral + +- **WHEN** a `token` or `thought` delta is emitted +- **THEN** the published chunk has `seq === 0` +- **AND** no Turso row is written for that delta diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/wire-schemas/spec.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/wire-schemas/spec.md new file mode 100644 index 00000000..9f9ab82a --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/wire-schemas/spec.md @@ -0,0 +1,33 @@ +## ADDED Requirements + +### Requirement: chat RPC returns turn acknowledgment + +The controller contract `chat` operation SHALL output `ChatOutputSchema` (`{ threadId: string, turnId: string }`) instead of `eventIterator(ChatChunkSchema)`. + +#### Scenario: Contract type is unary + +- **WHEN** `controllerContract` is defined in `shared/connections/src/contracts/controller.ts` +- **THEN** `chat` uses `.output(ChatOutputSchema)` and not `eventIterator(ChatChunkSchema)` + +#### Scenario: ChatOutput schema is exported + +- **WHEN** a consumer imports chat types from `@cyrus/schemas/rtc/chat` +- **THEN** `ChatOutputSchema` and `ChatOutput` type are available + +### Requirement: watchThread and unwatchThread contract schemas + +The controller contract SHALL define `watchThread` and `unwatchThread` operations with Zod schemas in `@cyrus/schemas/rtc/threads` (or `rtc/chat` if colocated). + +`WatchThreadInputSchema` SHALL contain `threadId: string`. +`WatchThreadOutputSchema` SHALL contain `snapshotHighWaterMark: number`. +`UnwatchThreadInputSchema` SHALL contain `threadId: string`. + +#### Scenario: Watch RPC is on controller contract + +- **WHEN** `controllerContract` is inspected +- **THEN** `watchThread` and `unwatchThread` are defined with the schemas above + +#### Scenario: subscribe remains eventIterator + +- **WHEN** `controllerContract` is inspected +- **THEN** `subscribe` still outputs `eventIterator(ChatChunkSchema)` diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md new file mode 100644 index 00000000..74c2a3e0 --- /dev/null +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md @@ -0,0 +1,48 @@ +## 1. Wire schemas and contract + +- [x] 1.1 Add `ChatOutputSchema` (`{ threadId, turnId }`) to `@cyrus/schemas/rtc/chat` +- [x] 1.2 Add `WatchThreadInputSchema`, `WatchThreadOutputSchema` (`snapshotHighWaterMark`), and `UnwatchThreadInputSchema` to `@cyrus/schemas/rtc/threads` +- [x] 1.3 Change `controllerContract.chat` from `eventIterator(ChatChunkSchema)` to `ChatOutputSchema` +- [x] 1.4 Add `watchThread` and `unwatchThread` to `controllerContract` +- [x] 1.5 Add `watchThread` / `unwatchThread` keys to `shared/constants/src/operation-keys.ts` (if using mutation keys) + +## 2. ThreadEventBus (`apps/cli/src/queue`) + +- [x] 2.1 Implement `ThreadEventBus` in `apps/cli/src/queue/index.ts` with watch sets, per-peer queues, and `publish()` +- [x] 2.2 Add `activeTurnLogs: Map` with append on publish and eviction on `turn_completed`/`turn_interrupted` +- [x] 2.3 Implement `watch(peerId, threadId)` — register, replay active-turn logs into peer queue (idempotent) +- [x] 2.4 Implement `unwatch(peerId, threadId)` and `ensureWatch(peerId, threadId)` for chat auto-watch +- [x] 2.5 Add per-turn log size cap with delta-first eviction policy (default 10k chunks) +- [x] 2.6 Wire `ThreadEventBus` into `shared/connections/src/rtc/worker/index.ts` and `RtcContext` (`peer.ts`), replacing `createPeerBroadcaster` for chat delivery + +## 3. Controller handlers (server) + +- [x] 3.1 Refactor `chat` handler to return `{ threadId, turnId }` immediately and run the turn loop without yielding chunks to the RPC caller +- [x] 3.2 Change `emit()` to call `ThreadEventBus.publish()` instead of `broadcaster.broadcast()` (include sender in fan-out) +- [x] 3.3 Auto-watch caller's `threadId` on `chat()` if not already watching +- [x] 3.4 Add `watchThread` and `unwatchThread` handlers +- [x] 3.5 Update `subscribe` handler to consume from `ThreadEventBus.subscribe(peerId)` + +## 4. Conversation overlay store (`@cyrus/hooks`) + +- [x] 4.1 Create `shared/hooks/src/stores/conversation-overlay.ts` with per-thread overlay state, `applySnapshot`, `applyLiveChunk`, `clearTurn` +- [x] 4.2 Implement `seq`-based dedup rules in `applyLiveChunk` +- [x] 4.3 Implement `waitForTurnEnd(threadId, turnId, signal?)` with resolve on `turn_completed`, reject on `turn_interrupted` +- [x] 4.4 Add `mergeSnapshotAndOverlay(snapshot, overlay)` in `shared/utils` (alongside `fold()`) producing `ConversationEntry[]` + +## 5. Shared connection hooks refactor + +- [x] 5.1 Update `shared/hooks/src/connection/use-worker-conversation-sync.ts` to write chunks to overlay store only (remove `appendChunkToCache`) +- [x] 5.2 Refactor `shared/hooks/src/connection/use-thread-conversation.ts` to call `watchThread`/`unwatchThread` on mount/unmount, merge snapshot + overlay, then `fold()` +- [x] 5.3 Refactor `shared/hooks/src/connection/use-controller-threads.ts` to call unary `chat()`, then `await waitForTurnEnd(threadId, turnId)`; invalidate threads list on completion +- [x] 5.4 Delete `shared/utils/src/conversation-cache.ts` +- [x] 5.5 Remove dead `if (!thread)` branch from `apps/web/src/components/chat/main/thread-workspace.tsx` `handleSend` + +## 6. Verification + +- [ ] 6.1 Manual test (web): single device — send message, tokens stream via subscribe, `sendMessage` resolves on `turn_completed` +- [ ] 6.2 Manual test (web): two devices on same thread — both see live stream; device on different thread receives nothing +- [ ] 6.3 Manual test (web): refresh mid-stream — partial tokens replay via `watchThread` active-turn log +- [ ] 6.4 Manual test (web): navigate away and back mid-stream — replay restores in-flight content +- [ ] 6.5 Manual test (web): stop/cancel — `sendMessage` rejects, overlay clears turn +- [x] 6.6 Run `bun check` and `bun check:types` across affected packages diff --git a/openspec/specs/conversation-overlay/spec.md b/openspec/specs/conversation-overlay/spec.md new file mode 100644 index 00000000..c0196bc9 --- /dev/null +++ b/openspec/specs/conversation-overlay/spec.md @@ -0,0 +1,85 @@ +## Purpose + +Ephemeral client-side live conversation state layered on top of the durable `getConversations` snapshot. + +## Requirements + +### Requirement: Conversation overlay store + +The client (`@cyrus/hooks`) SHALL maintain a `conversation-overlay` store (Zustand) in `shared/hooks/src/stores/conversation-overlay.ts` holding per-thread live chunks separate from the TanStack Query `getConversations` snapshot. The overlay SHALL track `snapshotHighWaterMark` and `activeTurnIds` per thread. + +#### Scenario: Live chunk is stored in overlay + +- **WHEN** a `ChatChunk` arrives on `subscribe()` for `threadId` T +- **THEN** the chunk is appended to the overlay for T +- **AND** the `getConversations` query cache is not mutated + +#### Scenario: Snapshot seeds watermark + +- **WHEN** `getConversations` returns entries for thread T +- **THEN** `snapshotHighWaterMark` for T is set to the maximum `seq` among returned entries +- **AND** overlay entries with `seq > 0` and `seq <= snapshotHighWaterMark` are pruned + +### Requirement: Single subscribe ingress + +All live `ChatChunk`s on the client SHALL be written exclusively by `useWorkerConversationSync` (`shared/hooks/src/connection/`) into the overlay store. No other code path SHALL append live chunks to the query cache. + +#### Scenario: sendMessage does not consume chat iterator + +- **WHEN** `sendMessage` is called +- **THEN** it invokes `chat()` as a unary RPC +- **AND** it does not iterate an event stream from `chat()` + +### Requirement: Seq-based dedup in overlay + +The overlay SHALL reject duplicate or stale chunks using `seq`: + +- Chunks with `seq > 0` and `seq <= snapshotHighWaterMark` SHALL be dropped. +- Chunks with `seq > 0` already present in the overlay SHALL be dropped. +- Chunks with `seq === 0` SHALL be accepted while their `turnId` is active. + +#### Scenario: Snapshot overlap dedup + +- **WHEN** a persisted chunk with `seq` 42 arrives on `subscribe()` +- **AND** `snapshotHighWaterMark` is already 42 or higher +- **THEN** the chunk is not added to the overlay + +#### Scenario: Ephemeral delta accepted during active turn + +- **WHEN** a `token` delta with `seq === 0` arrives for an active `turnId` +- **THEN** the chunk is added to the overlay + +### Requirement: Turn completion reconciles overlay with snapshot + +When the overlay receives `turn_completed` or `turn_interrupted` for a `turnId`, it SHALL remove that turn from `activeTurnIds`, invalidate the `getConversations` query for the thread, and prune overlay entries once the refetched snapshot updates the watermark. + +#### Scenario: Turn completion triggers snapshot refresh + +- **WHEN** `turn_completed` is received for `turnId` T on thread X +- **THEN** `getConversations(X)` is invalidated +- **AND** after refetch, overlay ephemeral entries for T are pruned against the new watermark + +### Requirement: sendMessage awaits turn end via overlay (option B) + +`sendMessage` (`shared/hooks/src/connection/use-controller-threads.ts`) SHALL register a `waitForTurnEnd(threadId, turnId)` promise before calling `chat()`. It SHALL resolve when the overlay receives `turn_completed` for that `turnId`, or reject on `turn_interrupted` or abort. + +#### Scenario: sendMessage completes on turn_completed + +- **WHEN** `sendMessage` calls `chat()` and receives `{ turnId }` +- **AND** the overlay subsequently receives `turn_completed` for that `turnId` +- **THEN** `sendMessage` resolves successfully + +#### Scenario: sendMessage rejects on turn_interrupted + +- **WHEN** the overlay receives `turn_interrupted` for the awaited `turnId` +- **THEN** `sendMessage` rejects with an error + +### Requirement: watchThread on thread mount + +`useThreadConversation` (or a hook it composes) SHALL call `watchThread` when a thread workspace mounts and `unwatchThread` on unmount, before or in parallel with the `getConversations` fetch. + +#### Scenario: Mount registers watch + +- **WHEN** `ThreadWorkspace` mounts for `threadId` T +- **THEN** `watchThread({ threadId: T })` is called +- **AND** `unwatchThread({ threadId: T })` is called on unmount diff --git a/openspec/specs/conversation-view/spec.md b/openspec/specs/conversation-view/spec.md index ed9e245c..bdd5abe3 100644 --- a/openspec/specs/conversation-view/spec.md +++ b/openspec/specs/conversation-view/spec.md @@ -73,3 +73,24 @@ The system SHALL remove the following from shared client types: `branch`, `lates - **WHEN** the sidebar displays a thread timestamp - **THEN** it uses `thread.updatedAt` from `ThreadSchema` - **AND** no `latestUserMessageAt` field exists in shared types + +### Requirement: Thread conversation merges snapshot and overlay + +`useThreadConversation` (`shared/hooks/src/connection/use-thread-conversation.ts`) SHALL derive `ThreadConversation` by merging the `getConversations` snapshot with live overlay entries, then calling `fold()` on the merged `ConversationEntry[]`. The merge SHALL prefer snapshot entries on `seq` collision. The hook SHALL consume RTC via `useRtc()`. + +#### Scenario: Merged view includes live tokens + +- **WHEN** the snapshot contains a `user_message` for the latest turn +- **AND** the overlay contains `token` deltas for that turn +- **THEN** `useThreadConversation` returns a `ThreadConversation` with partial assistant text visible before `message_completed` is persisted + +#### Scenario: Snapshot wins on seq collision + +- **WHEN** the snapshot and overlay both contain an entry with `seq` 42 +- **THEN** the snapshot entry is used in the merge +- **AND** the overlay duplicate is excluded + +#### Scenario: Fold is unchanged + +- **WHEN** merged entries are passed to `fold()` +- **THEN** the existing `fold()` implementation in `@cyrus/utils` produces the view without modification to its logic diff --git a/openspec/specs/thread-live-stream/spec.md b/openspec/specs/thread-live-stream/spec.md new file mode 100644 index 00000000..d0ed437e --- /dev/null +++ b/openspec/specs/thread-live-stream/spec.md @@ -0,0 +1,89 @@ +## Purpose + +Thread-scoped live `ChatChunk` delivery from the CLI worker to watching peers, with in-memory replay for in-flight turns. + +## Requirements + +### Requirement: Thread-scoped event bus replaces global broadcast + +The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` implemented in `apps/cli/src/queue/index.ts` that fans out only to peers whose watched-thread set includes `chunk.threadId`. The bus SHALL deliver chunks to the initiating peer as well as all other watching peers. + +#### Scenario: Non-watching peer receives nothing + +- **WHEN** peer B is connected but has not called `watchThread` for thread A +- **AND** a chat turn is active on thread A +- **THEN** peer B's `subscribe()` stream does not receive chunks for thread A + +#### Scenario: Watching peer receives live chunks + +- **WHEN** peer C has called `watchThread` for thread A +- **AND** a chat turn emits chunks on thread A +- **THEN** peer C's `subscribe()` stream receives those chunks in emission order + +#### Scenario: Initiating peer receives its own chunks via subscribe + +- **WHEN** peer A calls `chat()` for thread A +- **THEN** peer A receives all chunks for that turn through `subscribe()`, not through the `chat()` RPC response + +### Requirement: Active-turn replay buffer + +The `ThreadEventBus` in `apps/cli/src/queue/index.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. + +#### Scenario: Late watcher receives in-flight deltas + +- **WHEN** a turn is in progress on thread A with partial token deltas already emitted +- **AND** a peer calls `watchThread` for thread A before the turn completes +- **THEN** the peer receives a replay of all buffered chunks for active turns on thread A before subsequent live chunks + +#### Scenario: Completed turn log is evicted + +- **WHEN** `turn_completed` is published for `turnId` T +- **THEN** `activeTurnLogs` for T is removed +- **AND** a subsequent `watchThread` does not replay chunks from turn T (durable history comes from `getConversations`) + +### Requirement: watchThread and unwatchThread RPCs + +The controller contract SHALL expose `watchThread` and `unwatchThread` operations. `watchThread` SHALL register the calling peer's interest in a thread, replay active-turn logs for that thread, and return a `snapshotHighWaterMark` (the highest persisted `seq` for that thread). `unwatchThread` SHALL remove the registration. + +#### Scenario: watchThread returns cursor + +- **WHEN** a peer calls `watchThread({ threadId })` +- **THEN** the peer is registered as watching `threadId` +- **AND** the response includes `snapshotHighWaterMark` equal to the max persisted `seq` for that thread (or 0 if none) + +#### Scenario: unwatchThread stops delivery + +- **WHEN** a peer calls `unwatchThread({ threadId })` +- **THEN** subsequent chunks for `threadId` are not delivered to that peer's `subscribe()` stream + +### Requirement: chat is a turn-start command + +The `chat` RPC SHALL accept `ChatInput` and return `{ threadId, turnId }` without streaming `ChatChunk`s. The worker SHALL run the agent turn asynchronously, publishing all chunks via `ThreadEventBus.publish()`. + +#### Scenario: chat returns turn identity + +- **WHEN** a peer calls `chat({ agentName, message, projectId, threadId })` +- **THEN** the RPC returns `{ threadId, turnId }` before the turn completes +- **AND** all chunks for the turn arrive on `subscribe()` + +#### Scenario: chat auto-watches sender thread + +- **WHEN** a peer calls `chat()` for `threadId` T +- **AND** the peer is not already watching T +- **THEN** the server registers T in that peer's watched-thread set before publishing the first chunk + +### Requirement: emit publishes through the bus + +`chat.ts`'s `emit()` function SHALL call `ThreadEventBus.publish(chunk)` for all events. Persisted events SHALL still be written to Turso before publish (persist-before-broadcast from conversation-persistence). Ephemeral deltas (`seq === 0`) SHALL be published without a Turso write. + +#### Scenario: Persisted event carries seq on publish + +- **WHEN** a non-delta event is emitted +- **THEN** the chunk is persisted to Turso and assigned a `seq > 0` before `publish()` +- **AND** the published chunk includes that `seq` + +#### Scenario: Delta is ephemeral + +- **WHEN** a `token` or `thought` delta is emitted +- **THEN** the published chunk has `seq === 0` +- **AND** no Turso row is written for that delta diff --git a/openspec/specs/wire-schemas/spec.md b/openspec/specs/wire-schemas/spec.md index b2cc6fe7..4df29eab 100644 --- a/openspec/specs/wire-schemas/spec.md +++ b/openspec/specs/wire-schemas/spec.md @@ -64,3 +64,35 @@ The system SHALL NOT re-export wire schemas from `@cyrus/connections`. All consu - **WHEN** a developer adds an import for a wire schema type - **THEN** the import path starts with `@cyrus/schemas/` - **AND** no `export { ... } from "@cyrus/schemas/..."` shim exists in `@cyrus/connections` + +### Requirement: chat RPC returns turn acknowledgment + +The controller contract `chat` operation SHALL output `ChatOutputSchema` (`{ threadId: string, turnId: string }`) instead of `eventIterator(ChatChunkSchema)`. + +#### Scenario: Contract type is unary + +- **WHEN** `controllerContract` is defined in `shared/connections/src/contracts/controller.ts` +- **THEN** `chat` uses `.output(ChatOutputSchema)` and not `eventIterator(ChatChunkSchema)` + +#### Scenario: ChatOutput schema is exported + +- **WHEN** a consumer imports chat types from `@cyrus/schemas/rtc/chat` +- **THEN** `ChatOutputSchema` and `ChatOutput` type are available + +### Requirement: watchThread and unwatchThread contract schemas + +The controller contract SHALL define `watchThread` and `unwatchThread` operations with Zod schemas in `@cyrus/schemas/rtc/threads` (or `rtc/chat` if colocated). + +`WatchThreadInputSchema` SHALL contain `threadId: string`. +`WatchThreadOutputSchema` SHALL contain `snapshotHighWaterMark: number`. +`UnwatchThreadInputSchema` SHALL contain `threadId: string`. + +#### Scenario: Watch RPC is on controller contract + +- **WHEN** `controllerContract` is inspected +- **THEN** `watchThread` and `unwatchThread` are defined with the schemas above + +#### Scenario: subscribe remains eventIterator + +- **WHEN** `controllerContract` is inspected +- **THEN** `subscribe` still outputs `eventIterator(ChatChunkSchema)` diff --git a/shared/connections/src/contracts/controller.ts b/shared/connections/src/contracts/controller.ts index 7f92e5ca..697e6e07 100644 --- a/shared/connections/src/contracts/controller.ts +++ b/shared/connections/src/contracts/controller.ts @@ -13,6 +13,7 @@ import { CancelInputSchema, ChatChunkSchema, ChatInputSchema, + ChatOutputSchema, } from "@cyrus/schemas/rtc/chat"; import { AgentQueryInputSchema, @@ -37,6 +38,9 @@ import { ProjectQueryInputSchema, RenameThreadInputSchema, ThreadQueryInputSchema, + UnwatchThreadInputSchema, + WatchThreadInputSchema, + WatchThreadOutputSchema, } from "@cyrus/schemas/rtc/threads"; import { eventIterator, oc } from "@orpc/contract"; @@ -68,8 +72,10 @@ export const controllerContract = { setMode: oc.input(SetModeInputSchema).output(VoidOutputSchema), setEffort: oc.input(SetEffortInputSchema).output(VoidOutputSchema), setPersona: oc.input(SetPersonaInputSchema).output(VoidOutputSchema), - chat: oc.input(ChatInputSchema).output(eventIterator(ChatChunkSchema)), + chat: oc.input(ChatInputSchema).output(ChatOutputSchema), subscribe: oc.output(eventIterator(ChatChunkSchema)), + watchThread: oc.input(WatchThreadInputSchema).output(WatchThreadOutputSchema), + unwatchThread: oc.input(UnwatchThreadInputSchema).output(VoidOutputSchema), cancel: oc.input(CancelInputSchema).output(VoidOutputSchema), }; diff --git a/shared/connections/src/rtc/bus.ts b/shared/connections/src/rtc/bus.ts new file mode 100644 index 00000000..9650efc2 --- /dev/null +++ b/shared/connections/src/rtc/bus.ts @@ -0,0 +1,12 @@ +import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; + +export type ThreadEventBus = { + publish(chunk: ChatChunk): void; + watch(peerId: string, threadId: string): void; + unwatch(peerId: string, threadId: string): void; + ensureWatch(peerId: string, threadId: string): void; + isWatching(peerId: string, threadId: string): boolean; + subscribe(peerId: string): AsyncGenerator; + close(peerId: string): void; + closeAll(): void; +}; diff --git a/shared/connections/src/rtc/peer.ts b/shared/connections/src/rtc/peer.ts index bf02fda2..727b064f 100644 --- a/shared/connections/src/rtc/peer.ts +++ b/shared/connections/src/rtc/peer.ts @@ -1,14 +1,14 @@ -import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import type { ServerEvent } from "@cyrus/schemas/signaling"; import { Result } from "better-result"; import type { SignalingClient } from "../contracts/signaling"; -import type { PeerBroadcaster } from "./broadcaster"; +import type { ThreadEventBus } from "./bus"; export type { SignalingClient } from "../contracts/signaling"; +export type { ThreadEventBus } from "./bus"; export type RtcContext = { peerId: string; - broadcaster: PeerBroadcaster; + eventBus: ThreadEventBus; }; // fans one signaling event stream out to many subscribers (one socket, many peers) @@ -25,12 +25,9 @@ export function createSignalingEvents( Result.tryPromise(async () => { for await (const event of stream) { - if (!active) { - break; - } - for (const handler of handlers) { - handler(event); - } + if (!active) break; + + for (const handler of handlers) handler(event); } }); diff --git a/shared/connections/src/rtc/worker/index.ts b/shared/connections/src/rtc/worker/index.ts index 9d0470f2..1f09c8fd 100644 --- a/shared/connections/src/rtc/worker/index.ts +++ b/shared/connections/src/rtc/worker/index.ts @@ -1,11 +1,10 @@ -import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import { DeviceRoleSchema } from "@cyrus/schemas/signaling"; import type { Router } from "@orpc/server"; import { RPCHandler, type RPCHandlerOptions } from "@orpc/server/websocket"; import { RTCPeerConnection as NodeRTCPeerConnection } from "node-datachannel/polyfill"; import type { ControllerContract } from "../../contracts/controller"; import type { WorkerContract } from "../../contracts/worker"; -import { createPeerBroadcaster } from "../broadcaster"; +import type { ThreadEventBus } from "../bus"; import { createIceBuffer, type RtcContext, @@ -26,6 +25,7 @@ export type WorkerOptions = { signaling: SignalingClient; events: SignalingEvents; routers: WorkerRouters; + eventBus: ThreadEventBus; config?: RTCConfiguration; rpc?: RPCHandlerOptions; }; @@ -35,9 +35,7 @@ export type WorkerConnection = { }; export function serveWorker(options: WorkerOptions): WorkerConnection { - const { signaling, events, config, routers } = options; - - const broadcaster = createPeerBroadcaster(); + const { signaling, events, config, routers, eventBus } = options; const handlers = { controller: new RPCHandler(routers.controller, options.rpc), @@ -50,7 +48,7 @@ export function serveWorker(options: WorkerOptions): WorkerConnection { >(); function dispose(peerId: string): void { - broadcaster.close(peerId); + eventBus.close(peerId); const session = sessions.get(peerId); if (session) { session.pc.close(); @@ -83,7 +81,7 @@ export function serveWorker(options: WorkerOptions): WorkerConnection { whenOpen(channel) .then(() => { handlers[role.data].upgrade(asWebSocket(channel), { - context: { peerId: from, broadcaster } satisfies RtcContext, + context: { peerId: from, eventBus } satisfies RtcContext, }); }) .catch(() => { @@ -125,7 +123,7 @@ export function serveWorker(options: WorkerOptions): WorkerConnection { for (const peerId of [...sessions.keys()]) { dispose(peerId); } - broadcaster.closeAll(); + eventBus.closeAll(); }, }; } diff --git a/shared/constants/src/operation-keys.ts b/shared/constants/src/operation-keys.ts index a53d7c10..1d71929d 100644 --- a/shared/constants/src/operation-keys.ts +++ b/shared/constants/src/operation-keys.ts @@ -28,6 +28,8 @@ export const RTC_OPERATION_KEYS = { setModel: ["controller", "set-model"], setEffort: ["controller", "set-effort"], setPersona: ["controller", "set-persona"], + watchThread: ["controller", "watch-thread"], + unwatchThread: ["controller", "unwatch-thread"], } as const; export const AUTH_OPERATION_KEYS = { diff --git a/shared/database/src/repositories/conversations.ts b/shared/database/src/repositories/conversations.ts index 4ddda34e..784f26d7 100644 --- a/shared/database/src/repositories/conversations.ts +++ b/shared/database/src/repositories/conversations.ts @@ -2,7 +2,7 @@ import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import { ConversationEntrySchema } from "@cyrus/schemas/rtc/threads"; import { randomId } from "@cyrus/utils/identity"; import { nowISO } from "@cyrus/utils/time"; -import { and, asc, eq, gt } from "drizzle-orm"; +import { and, asc, desc, eq, gt } from "drizzle-orm"; import { connection } from "../connection"; import { conversations } from "../models/conversations"; import { threads } from "../models/threads"; @@ -66,3 +66,15 @@ export function getConversations(threadId: string, afterSeq?: number) { return rows.map(parseConversationEntry); }); } + +export function getSnapshotHighWaterMark(threadId: string) { + return tryRepo(async () => { + const [row] = await connection.db + .select({ seq: conversations.seq }) + .from(conversations) + .where(eq(conversations.threadId, threadId)) + .orderBy(desc(conversations.seq)) + .limit(1); + return row?.seq ?? 0; + }); +} diff --git a/shared/hooks/src/connection/use-controller-threads.ts b/shared/hooks/src/connection/use-controller-threads.ts index 0ce12f93..29289d51 100644 --- a/shared/hooks/src/connection/use-controller-threads.ts +++ b/shared/hooks/src/connection/use-controller-threads.ts @@ -1,14 +1,13 @@ import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; -import { appendChunkToCache } from "@cyrus/utils/conversation-cache"; -import { useQuery, useQueryClient } from "@tanstack/react-query"; +import { useQuery } from "@tanstack/react-query"; import { Result } from "better-result"; import { useRtc } from "../contexts/rtc"; import { useAgentCatalogStore } from "../stores/agent-catalog"; +import { waitForTurnEnd } from "../stores/conversation-overlay"; import { useProjects } from "./use-projects"; import { useThreads } from "./use-threads"; export function useControllerThreads() { - const queryClient = useQueryClient(); const { connection: workerConnection, orpc: orpcController } = useRtc(); const { @@ -59,15 +58,13 @@ export function useControllerThreads() { if (!thread) return Result.err(new Error(`thread not found: ${threadId}`)); const result = await Result.tryPromise(async () => { - const iterator = await workerConnection.client.chat({ + const { turnId } = await workerConnection.client.chat({ agentName: resolveAgentName(threadId), message: text, projectId: thread.projectId, threadId, }); - for await (const chunk of iterator) { - appendChunkToCache(queryClient, chunk); - } + await waitForTurnEnd(threadId, turnId); }); invalidateThreads(thread.projectId); return result; diff --git a/shared/hooks/src/connection/use-thread-conversation.ts b/shared/hooks/src/connection/use-thread-conversation.ts index 4802743a..5921bf5b 100644 --- a/shared/hooks/src/connection/use-thread-conversation.ts +++ b/shared/hooks/src/connection/use-thread-conversation.ts @@ -1,9 +1,14 @@ import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; +import type { ConversationEntry } from "@cyrus/schemas/rtc/threads"; import type { ThreadConversation } from "@cyrus/schemas/view"; import { fold } from "@cyrus/utils/fold"; +import { mergeSnapshotAndOverlay } from "@cyrus/utils/merge-conversation"; import { useQuery } from "@tanstack/react-query"; +import { Result } from "better-result"; import { log } from "evlog"; +import { useEffect, useEffectEvent, useMemo } from "react"; import { useRtc } from "../contexts/rtc"; +import { useConversationOverlay } from "../stores/conversation-overlay"; const EMPTY: ThreadConversation = { diffs: [], @@ -15,28 +20,93 @@ const EMPTY: ThreadConversation = { export function useThreadConversation( threadId: string | undefined ): ThreadConversation { - const { orpc: orpcController } = useRtc(); - - const queryKey = threadId - ? RTC_OPERATION_KEYS.getConversations(threadId) - : (["controller", "get-conversations", "none"] as const); + const { orpc: orpcController, connection: workerConnection } = useRtc(); + const applySnapshot = useConversationOverlay((state) => state.applySnapshot); + const applyWatermark = useConversationOverlay( + (state) => state.applyWatermark + ); + const clearTurn = useConversationOverlay((state) => state.clearTurn); + const liveEntries = useConversationOverlay((state) => + threadId ? state.getLiveEntries(threadId) : [] + ); const conversationsQuery = useQuery({ ...orpcController.getConversations.queryOptions({ - queryKey, + queryKey: RTC_OPERATION_KEYS.getConversations(threadId ?? "none"), input: { threadId: threadId ?? "" }, }), - queryKey, enabled: Boolean(threadId), - select: (data) => - fold(data.conversations).match({ - ok: (conversation) => conversation, - err: (error) => { - log.error({ kind: "fold_conversation", error, threadId }); - return EMPTY; - }, - }), }); - return conversationsQuery.data ?? EMPTY; + const onWatchResult = useEffectEvent( + (tid: string, snapshotHighWaterMark: number) => + applyWatermark(tid, snapshotHighWaterMark) + ); + + const onWatchError = useEffectEvent((error: unknown, tid: string) => { + log.error({ kind: "watch_thread", error, threadId: tid }); + }); + + const onUnwatchError = useEffectEvent((error: unknown, tid: string) => { + log.error({ kind: "unwatch_thread", error, threadId: tid }); + }); + + const syncSnapshot = useEffectEvent( + (tid: string, conversations: ConversationEntry[]) => { + applySnapshot(tid, conversations); + + for (const entry of conversations) { + if ( + entry.chunk.event.type === "turn_completed" || + entry.chunk.event.type === "turn_interrupted" + ) { + clearTurn(tid, entry.chunk.turnId); + } + } + } + ); + + useEffect(() => { + if (!threadId) return; + + const abort = new AbortController(); + + Result.tryPromise(() => + workerConnection.client.watchThread({ threadId }) + ).then((result) => { + if (abort.signal.aborted) return; + result.match({ + ok: ({ snapshotHighWaterMark }) => + onWatchResult(threadId, snapshotHighWaterMark), + err: (error) => onWatchError(error, threadId), + }); + }); + + return () => { + abort.abort(); + workerConnection.client + .unwatchThread({ threadId }) + .catch((error) => onUnwatchError(error, threadId)); + }; + }, [threadId, workerConnection]); + + useEffect(() => { + if (!(threadId && conversationsQuery.data?.conversations)) return; + syncSnapshot(threadId, conversationsQuery.data.conversations); + }, [threadId, conversationsQuery.data]); + + return useMemo(() => { + const merged = mergeSnapshotAndOverlay( + conversationsQuery.data?.conversations ?? [], + liveEntries + ); + + return fold(merged).match({ + ok: (conversation) => conversation, + err: (error) => { + log.error({ kind: "fold_conversation", error, threadId }); + return EMPTY; + }, + }); + }, [conversationsQuery.data?.conversations, liveEntries, threadId]); } diff --git a/shared/hooks/src/connection/use-worker-conversation-sync.ts b/shared/hooks/src/connection/use-worker-conversation-sync.ts index c93af02e..85dbc320 100644 --- a/shared/hooks/src/connection/use-worker-conversation-sync.ts +++ b/shared/hooks/src/connection/use-worker-conversation-sync.ts @@ -1,12 +1,36 @@ -import { appendChunkToCache } from "@cyrus/utils/conversation-cache"; +import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; +import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import { useQueryClient } from "@tanstack/react-query"; import { Result } from "better-result"; -import { useEffect } from "react"; +import { useEffect, useEffectEvent } from "react"; import { useRtc } from "../contexts/rtc"; +import { useConversationOverlay } from "../stores/conversation-overlay"; + +function isTerminalChunk(chunk: ChatChunk): boolean { + return ( + chunk.event.type === "turn_completed" || + chunk.event.type === "turn_interrupted" + ); +} export function useWorkerConversationSync(): void { const queryClient = useQueryClient(); const { connection: workerConnection } = useRtc(); + const applyLiveChunk = useConversationOverlay( + (state) => state.applyLiveChunk + ); + + const onChunk = useEffectEvent((chunk: ChatChunk) => { + applyLiveChunk(chunk); + if (isTerminalChunk(chunk)) + queryClient.invalidateQueries({ + queryKey: RTC_OPERATION_KEYS.getConversations(chunk.threadId), + }); + }); + + const onSyncError = useEffectEvent((error: unknown) => { + console.error("worker conversation sync failed", error); + }); useEffect(() => { let stopped = false; @@ -20,17 +44,15 @@ export function useWorkerConversationSync(): void { } for await (const chunk of iterator) { if (stopped) break; - appendChunkToCache(queryClient, chunk); + onChunk(chunk); } }).then((result) => { - if (result.isErr() && !stopped) { - console.error("worker conversation sync failed", result.error); - } + if (result.isErr() && !stopped) onSyncError(result.error); }); return () => { stopped = true; iterator?.return?.(undefined); }; - }, [workerConnection, queryClient]); + }, [workerConnection]); } diff --git a/shared/hooks/src/stores/conversation-overlay.ts b/shared/hooks/src/stores/conversation-overlay.ts new file mode 100644 index 00000000..1d3b0bde --- /dev/null +++ b/shared/hooks/src/stores/conversation-overlay.ts @@ -0,0 +1,342 @@ +import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; +import type { ConversationEntry } from "@cyrus/schemas/rtc/threads"; +import type { StoreApi } from "zustand"; +import { create } from "zustand"; + +type OverlayEntry = { + chunk: ChatChunk; +}; + +type ThreadOverlay = { + snapshotHighWaterMark: number; + live: OverlayEntry[]; + activeTurnIds: Set; +}; + +type TurnWaiter = { + resolve: () => void; + reject: (error: Error) => void; + onAbort?: () => void; +}; + +type ConversationOverlayState = { + byThread: Map; + applyWatermark: (threadId: string, snapshotHighWaterMark: number) => void; + applySnapshot: (threadId: string, entries: ConversationEntry[]) => void; + applyLiveChunk: (chunk: ChatChunk) => void; + clearTurn: (threadId: string, turnId: string) => void; + getLiveEntries: (threadId: string) => ConversationEntry[]; +}; + +type OverlaySetState = StoreApi["setState"]; + +let overlayEntrySeq = 0; + +const turnWaiters = new Map(); +const pendingDeltas = new Map(); +let deltaFlushHandle: ReturnType | null = null; + +function turnKey(threadId: string, turnId: string): string { + return `${threadId}:${turnId}`; +} + +function isTerminalEvent(event: ChatChunk["event"]): boolean { + return event.type === "turn_completed" || event.type === "turn_interrupted"; +} + +function isStreamingDeltaChunk(chunk: ChatChunk): boolean { + return ( + chunk.seq === 0 && + (chunk.event.type === "token" || chunk.event.type === "thought") + ); +} + +function streamingDeltaKey(chunk: ChatChunk): string { + const messageId = + chunk.event.type === "token" || chunk.event.type === "thought" + ? (chunk.event.messageId ?? "default") + : "default"; + return `${chunk.threadId}:${chunk.turnId}:${chunk.event.type}:${messageId}`; +} + +function mergeStreamingDeltaChunks( + existing: ChatChunk, + incoming: ChatChunk +): ChatChunk { + const left = existing.event; + const right = incoming.event; + if (left.type === "token" && right.type === "token") { + return { + ...incoming, + event: { ...left, text: left.text + right.text }, + }; + } + if (left.type === "thought" && right.type === "thought") { + return { + ...incoming, + event: { ...left, text: left.text + right.text }, + }; + } + return incoming; +} + +function flushPendingDeltas(commit: (chunk: ChatChunk) => void): void { + if (deltaFlushHandle !== null) { + cancelAnimationFrame(deltaFlushHandle); + deltaFlushHandle = null; + } + for (const chunk of pendingDeltas.values()) { + commit(chunk); + } + pendingDeltas.clear(); +} + +function scheduleStreamingDeltaFlush(commit: (chunk: ChatChunk) => void): void { + if (deltaFlushHandle !== null) return; + deltaFlushHandle = requestAnimationFrame(() => { + deltaFlushHandle = null; + flushPendingDeltas(commit); + }); +} + +function queueStreamingDelta( + chunk: ChatChunk, + commit: (chunk: ChatChunk) => void +): void { + const key = streamingDeltaKey(chunk); + const existing = pendingDeltas.get(key); + pendingDeltas.set( + key, + existing ? mergeStreamingDeltaChunks(existing, chunk) : chunk + ); + scheduleStreamingDeltaFlush(commit); +} + +function chunkToEntry(chunk: ChatChunk): ConversationEntry { + return { + id: `overlay-${chunk.turnId}-${++overlayEntrySeq}`, + threadId: chunk.threadId, + seq: chunk.seq, + chunk, + createdAt: new Date().toISOString(), + }; +} + +function getOrCreateOverlay( + byThread: Map, + threadId: string +): ThreadOverlay { + const existing = byThread.get(threadId); + if (existing) return existing; + + const overlay: ThreadOverlay = { + snapshotHighWaterMark: 0, + live: [], + activeTurnIds: new Set(), + }; + byThread.set(threadId, overlay); + return overlay; +} + +function settleTurnWaiter( + threadId: string, + turnId: string, + outcome: "completed" | "interrupted" +): void { + const key = turnKey(threadId, turnId); + const waiter = turnWaiters.get(key); + if (!waiter) return; + + turnWaiters.delete(key); + waiter.onAbort?.(); + if (outcome === "completed") { + waiter.resolve(); + return; + } + waiter.reject(new Error("turn interrupted")); +} + +function shouldSkipPersistedChunk( + seq: number, + overlay: ThreadOverlay +): boolean { + if (seq <= 0) return false; + if (seq <= overlay.snapshotHighWaterMark) return true; + return overlay.live.some(({ chunk: live }) => live.seq === seq); +} + +function upsertStreamingDelta(overlay: ThreadOverlay, chunk: ChatChunk): void { + const deltaKey = streamingDeltaKey(chunk); + const existingIndex = overlay.live.findIndex( + ({ chunk: live }) => + isStreamingDeltaChunk(live) && streamingDeltaKey(live) === deltaKey + ); + if (existingIndex >= 0) { + const existing = overlay.live[existingIndex]; + if (existing) { + overlay.live[existingIndex] = { + chunk: mergeStreamingDeltaChunks(existing.chunk, chunk), + }; + } + return; + } + overlay.live.push({ chunk }); +} + +function applyChunkToOverlay(overlay: ThreadOverlay, chunk: ChatChunk): void { + const { turnId } = chunk; + + if (isStreamingDeltaChunk(chunk)) { + upsertStreamingDelta(overlay, chunk); + } else { + overlay.live.push({ chunk }); + } + + if (isTerminalEvent(chunk.event)) { + overlay.activeTurnIds.delete(turnId); + } else { + overlay.activeTurnIds.add(turnId); + } +} + +function commitLiveChunk(setState: OverlaySetState, chunk: ChatChunk): void { + const { threadId, turnId, seq } = chunk; + + setState((state) => { + const byThread = new Map(state.byThread); + const overlay = getOrCreateOverlay(byThread, threadId); + + if (shouldSkipPersistedChunk(seq, overlay)) { + return { byThread }; + } + + applyChunkToOverlay(overlay, chunk); + return { byThread }; + }); + + if (chunk.event.type === "turn_completed") { + settleTurnWaiter(threadId, turnId, "completed"); + } + if (chunk.event.type === "turn_interrupted") { + settleTurnWaiter(threadId, turnId, "interrupted"); + } +} + +export const useConversationOverlay = create( + (set, get) => { + const commit = (chunk: ChatChunk) => commitLiveChunk(set, chunk); + + return { + byThread: new Map(), + + applyWatermark(threadId, snapshotHighWaterMark) { + flushPendingDeltas(commit); + set((state) => { + const byThread = new Map(state.byThread); + const overlay = getOrCreateOverlay(byThread, threadId); + overlay.snapshotHighWaterMark = snapshotHighWaterMark; + overlay.live = overlay.live.filter( + ({ chunk }) => chunk.seq === 0 || chunk.seq > snapshotHighWaterMark + ); + return { byThread }; + }); + }, + + applySnapshot(threadId, entries) { + const snapshotHighWaterMark = entries.reduce( + (max, entry) => Math.max(max, entry.seq), + 0 + ); + get().applyWatermark(threadId, snapshotHighWaterMark); + }, + + applyLiveChunk(chunk) { + if (isStreamingDeltaChunk(chunk)) { + queueStreamingDelta(chunk, commit); + return; + } + + flushPendingDeltas(commit); + commitLiveChunk(set, chunk); + }, + + clearTurn(threadId, turnId) { + flushPendingDeltas(commit); + set((state) => { + const byThread = new Map(state.byThread); + const overlay = byThread.get(threadId); + if (!overlay) return { byThread }; + + overlay.activeTurnIds.delete(turnId); + overlay.live = overlay.live.filter( + ({ chunk }) => chunk.turnId !== turnId || chunk.seq > 0 + ); + return { byThread }; + }); + }, + + getLiveEntries(threadId) { + const overlay = get().byThread.get(threadId); + if (!overlay) return []; + return overlay.live.map(({ chunk }) => chunkToEntry(chunk)); + }, + }; + } +); + +export function waitForTurnEnd( + threadId: string, + turnId: string, + signal?: AbortSignal +): Promise { + return new Promise((resolve, reject) => { + const key = turnKey(threadId, turnId); + + function cleanup() { + signal?.removeEventListener("abort", onAbort); + turnWaiters.delete(key); + } + + function onAbort() { + cleanup(); + reject(new Error("turn aborted")); + } + + if (signal?.aborted) { + onAbort(); + return; + } + + signal?.addEventListener("abort", onAbort, { once: true }); + + const overlay = useConversationOverlay.getState(); + const terminal = overlay + .getLiveEntries(threadId) + .find( + (entry) => + entry.chunk.turnId === turnId && + (entry.chunk.event.type === "turn_completed" || + entry.chunk.event.type === "turn_interrupted") + ); + if (terminal) { + if (terminal.chunk.event.type === "turn_completed") { + resolve(); + } else { + reject(new Error("turn interrupted")); + } + return; + } + + turnWaiters.set(key, { + resolve: () => { + cleanup(); + resolve(); + }, + reject: (error) => { + cleanup(); + reject(error); + }, + onAbort, + }); + }); +} diff --git a/shared/schemas/src/rtc/chat.ts b/shared/schemas/src/rtc/chat.ts index 17b2a2f1..90a28089 100644 --- a/shared/schemas/src/rtc/chat.ts +++ b/shared/schemas/src/rtc/chat.ts @@ -255,6 +255,13 @@ export const ChatInputSchema = z.object({ projectId: z.string(), }); +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(), diff --git a/shared/schemas/src/rtc/threads.ts b/shared/schemas/src/rtc/threads.ts index c8cd68b9..7a12965d 100644 --- a/shared/schemas/src/rtc/threads.ts +++ b/shared/schemas/src/rtc/threads.ts @@ -28,6 +28,18 @@ export const ThreadQueryInputSchema = z.object({ afterSeq: z.number().optional(), }); +export const WatchThreadInputSchema = z.object({ + threadId: z.string(), +}); + +export const WatchThreadOutputSchema = z.object({ + snapshotHighWaterMark: z.number(), +}); + +export const UnwatchThreadInputSchema = z.object({ + threadId: z.string(), +}); + export const RenameThreadInputSchema = z.object({ threadId: z.string(), name: z.string().min(1), @@ -53,6 +65,9 @@ export type Thread = z.infer; export type ConversationEntry = z.infer; export type ProjectQueryInput = z.infer; export type ThreadQueryInput = z.infer; +export type WatchThreadInput = z.infer; +export type WatchThreadOutput = z.infer; +export type UnwatchThreadInput = z.infer; export type CreateThreadInput = z.infer; export type CreateThreadOutput = z.infer; export type ListThreadsOutput = z.infer; diff --git a/shared/utils/src/conversation-cache.ts b/shared/utils/src/conversation-cache.ts deleted file mode 100644 index 07cb7b9c..00000000 --- a/shared/utils/src/conversation-cache.ts +++ /dev/null @@ -1,24 +0,0 @@ -import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; -import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; -import type { GetConversationsOutput } from "@cyrus/schemas/rtc/threads"; -import type { QueryClient } from "@tanstack/react-query"; - -let syntheticEntrySeq = 0; - -export function appendChunkToCache(queryClient: QueryClient, chunk: ChatChunk) { - queryClient.setQueryData( - RTC_OPERATION_KEYS.getConversations(chunk.threadId), - (old) => ({ - conversations: [ - ...(old?.conversations ?? []), - { - chunk, - createdAt: new Date().toISOString(), - id: `local-${chunk.turnId}-${++syntheticEntrySeq}`, - seq: chunk.seq, - threadId: chunk.threadId, - }, - ], - }) - ); -} diff --git a/shared/utils/src/merge-conversation.ts b/shared/utils/src/merge-conversation.ts new file mode 100644 index 00000000..bbae279a --- /dev/null +++ b/shared/utils/src/merge-conversation.ts @@ -0,0 +1,21 @@ +import type { ConversationEntry } from "@cyrus/schemas/rtc/threads"; + +export function mergeSnapshotAndOverlay( + snapshot: ConversationEntry[], + overlay: ConversationEntry[] +): ConversationEntry[] { + const snapshotSeqs = new Set(snapshot.map((entry) => entry.seq)); + const merged = [...snapshot]; + + for (const entry of overlay) + if (entry.seq === 0 || !snapshotSeqs.has(entry.seq)) merged.push(entry); + + return merged.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; + } + return left.createdAt.localeCompare(right.createdAt); + }); +} From 1911856796cd6617aa29d89e88fdc54cb470b5b5 Mon Sep 17 00:00:00 2001 From: Soorya U Date: Fri, 10 Jul 2026 09:47:49 +0530 Subject: [PATCH 2/7] Fix ThreadEventBus import path and apply minor CLI/web cleanups. Point the queue bus at @cyrus/connections/rtc/bus, remove unused stdin helper, align OAuth client id with cyrusd, and tighten small style nits across CLI and web. Co-authored-by: Cursor --- apps/cli/src/commands/agents/doctor.ts | 4 +- apps/cli/src/commands/agents/rm.ts | 4 +- apps/cli/src/commands/auth/login.ts | 2 +- apps/cli/src/commands/auth/whoami.ts | 9 ++--- apps/cli/src/commands/service/start.ts | 13 ++----- apps/cli/src/commands/service/stop.ts | 5 +-- apps/cli/src/commands/service/worker.ts | 1 + apps/cli/src/queue/bus.ts | 2 +- apps/cli/src/utils/io.ts | 14 ------- .../src/components/auth/provider-button.tsx | 2 - apps/web/src/components/connection-error.tsx | 1 + .../sidebar/chat-sidebar-footer-actions.tsx | 4 +- .../sidebar/projects/thread-search-field.tsx | 2 - apps/web/src/hooks/use-auto-animate-ref.ts | 5 +-- apps/web/src/hooks/use-copy-to-clipboard.ts | 38 ++++++++++--------- apps/web/src/hooks/use-media-query.ts | 10 ++--- apps/web/src/utils/dir.ts | 6 +-- 17 files changed, 43 insertions(+), 79 deletions(-) delete mode 100644 apps/cli/src/utils/io.ts diff --git a/apps/cli/src/commands/agents/doctor.ts b/apps/cli/src/commands/agents/doctor.ts index 146191bb..1f161e99 100644 --- a/apps/cli/src/commands/agents/doctor.ts +++ b/apps/cli/src/commands/agents/doctor.ts @@ -18,9 +18,7 @@ async function checkAgent(entry: AgentEntry): Promise { function printHealth(name: string, result: HealthResult): void { const status = result.healthy ? green("healthy") : red("unhealthy"); print.line`${name}: ${status}`; - if (result.error) { - print.line` ${result.error}`; - } + if (result.error) print.line` ${result.error}`; } async function checkSingleAgentHealth(name: string) { diff --git a/apps/cli/src/commands/agents/rm.ts b/apps/cli/src/commands/agents/rm.ts index c54710bb..6b7f9671 100644 --- a/apps/cli/src/commands/agents/rm.ts +++ b/apps/cli/src/commands/agents/rm.ts @@ -4,9 +4,7 @@ import { print } from "@/utils/style"; export async function rm(name: string): Promise { const removed = await removeAgent(name); removed.match({ - ok: () => { - print.success`✓ removed agent "${name}"`; - }, + ok: () => print.success`✓ removed agent "${name}"`, err: (message) => { print.error`${message}`; process.exit(1); diff --git a/apps/cli/src/commands/auth/login.ts b/apps/cli/src/commands/auth/login.ts index 7fe87177..0b5a79ea 100644 --- a/apps/cli/src/commands/auth/login.ts +++ b/apps/cli/src/commands/auth/login.ts @@ -4,7 +4,7 @@ import { authClient } from "@/lib/auth"; import { getOrCreate, set } from "@/store/config"; import { blue, bold, cyan, print, underline } from "@/utils/style"; -export const CLIENT_ID = "cyrus-cli"; +export const CLIENT_ID = "cyrusd"; export const GRANT_TYPE = "urn:ietf:params:oauth:grant-type:device_code"; export async function login(): Promise { diff --git a/apps/cli/src/commands/auth/whoami.ts b/apps/cli/src/commands/auth/whoami.ts index 123d4b77..88bc6cba 100644 --- a/apps/cli/src/commands/auth/whoami.ts +++ b/apps/cli/src/commands/auth/whoami.ts @@ -24,11 +24,8 @@ export async function whoami(opts: WhoamiOptions): Promise { } const name = await get("name"); + print.line`${field("User")}${data.user.name}`; - if (name) { - print.line`${field("Device")}${name}`; - } - if (opts.email) { - print.line`${field("Email")}${data.user.email}`; - } + if (name) print.line`${field("Device")}${name}`; + if (opts.email) print.line`${field("Email")}${data.user.email}`; } diff --git a/apps/cli/src/commands/service/start.ts b/apps/cli/src/commands/service/start.ts index cbc471c5..95f1e149 100644 --- a/apps/cli/src/commands/service/start.ts +++ b/apps/cli/src/commands/service/start.ts @@ -19,9 +19,7 @@ async function runWorker(): Promise { process.on("exit", () => { Result.try(() => { const stored = Number.parseInt(readFileSync(PID_PATH, "utf8").trim(), 10); - if (stored === ownPid) { - unlinkSync(PID_PATH); - } + if (stored === ownPid) unlinkSync(PID_PATH); }); }); const { worker } = await import("./worker"); @@ -29,10 +27,7 @@ async function runWorker(): Promise { } export async function start(opts: StartOptions): Promise { - if (env.CYRUS_DAEMON) { - await runWorker(); - return; - } + if (env.CYRUS_DAEMON) return await runWorker(); const running = await runningPid(); if (running !== null) { @@ -50,9 +45,7 @@ export async function start(opts: StartOptions): Promise { }); closeSync(log); child.unref(); - if (child.pid) { - await writePid(child.pid); - } + if (child.pid) await writePid(child.pid); // brief wait to catch immediate startup failures (e.g. not logged in) await Bun.sleep(300); diff --git a/apps/cli/src/commands/service/stop.ts b/apps/cli/src/commands/service/stop.ts index bae262a9..31d03dec 100644 --- a/apps/cli/src/commands/service/stop.ts +++ b/apps/cli/src/commands/service/stop.ts @@ -11,9 +11,8 @@ export async function stop(): Promise { Result.try(() => process.kill(pid, "SIGTERM")); - for (let i = 0; i < 20 && isAlive(pid); i++) { - await Bun.sleep(50); - } + for (let i = 0; i < 20 && isAlive(pid); i++) await Bun.sleep(50); + if (isAlive(pid)) { print.error`Worker (pid ${pid}) did not stop in time.`; process.exit(1); diff --git a/apps/cli/src/commands/service/worker.ts b/apps/cli/src/commands/service/worker.ts index a8a499f0..d85f0ee7 100644 --- a/apps/cli/src/commands/service/worker.ts +++ b/apps/cli/src/commands/service/worker.ts @@ -68,6 +68,7 @@ export async function worker(): Promise { await connection.close(); process.exit(0); }; + process.on("SIGINT", shutdown); process.on("SIGTERM", shutdown); process.on("SIGBREAK", shutdown); diff --git a/apps/cli/src/queue/bus.ts b/apps/cli/src/queue/bus.ts index b1092452..30b9b924 100644 --- a/apps/cli/src/queue/bus.ts +++ b/apps/cli/src/queue/bus.ts @@ -1,4 +1,4 @@ -import type { ThreadEventBus } from "@cyrus/connections/rtc/thread-event-bus"; +import type { ThreadEventBus } from "@cyrus/connections/rtc/bus"; import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; const DEFAULT_MAX_CHUNKS_PER_TURN = 10_000; diff --git a/apps/cli/src/utils/io.ts b/apps/cli/src/utils/io.ts deleted file mode 100644 index 6aa19efa..00000000 --- a/apps/cli/src/utils/io.ts +++ /dev/null @@ -1,14 +0,0 @@ -export function stdinWritable(stdin: Bun.FileSink): WritableStream { - return new WritableStream({ - write(chunk) { - stdin.write(chunk); - stdin.flush(); - }, - close() { - stdin.end(); - }, - abort() { - stdin.end(); - }, - }); -} diff --git a/apps/web/src/components/auth/provider-button.tsx b/apps/web/src/components/auth/provider-button.tsx index b2f67bc9..b3caf44e 100644 --- a/apps/web/src/components/auth/provider-button.tsx +++ b/apps/web/src/components/auth/provider-button.tsx @@ -1,5 +1,3 @@ -"use client"; - import { authMutationKeys, getProviderName } from "@better-auth-ui/core"; import { providerIcons, useAuth, useSignInSocial } from "@better-auth-ui/react"; import { useIsMutating } from "@tanstack/react-query"; diff --git a/apps/web/src/components/connection-error.tsx b/apps/web/src/components/connection-error.tsx index 470e258a..de2d9df6 100644 --- a/apps/web/src/components/connection-error.tsx +++ b/apps/web/src/components/connection-error.tsx @@ -1,6 +1,7 @@ import type { ConnectionErrorRenderProps } from "@cyrus/providers/types"; import { Button } from "@/components/ui/button"; +// TODO: Just a Dummy. Change Later export function ConnectionError({ error, retry }: ConnectionErrorRenderProps) { return (
diff --git a/apps/web/src/components/sidebar/chat-sidebar-footer-actions.tsx b/apps/web/src/components/sidebar/chat-sidebar-footer-actions.tsx index 0a1f0c41..8a5c710f 100644 --- a/apps/web/src/components/sidebar/chat-sidebar-footer-actions.tsx +++ b/apps/web/src/components/sidebar/chat-sidebar-footer-actions.tsx @@ -24,9 +24,7 @@ export function ChatSidebarFooterActions() { const { resolvedTheme, setTheme } = useTheme(); const [mounted, setMounted] = useState(false); - useEffect(() => { - setMounted(true); - }, []); + useEffect(() => setMounted(true), []); const theme = resolvedTheme === "dark" ? "dark" : "light"; diff --git a/apps/web/src/components/sidebar/projects/thread-search-field.tsx b/apps/web/src/components/sidebar/projects/thread-search-field.tsx index 4de3edaf..53b1bdc4 100644 --- a/apps/web/src/components/sidebar/projects/thread-search-field.tsx +++ b/apps/web/src/components/sidebar/projects/thread-search-field.tsx @@ -1,5 +1,3 @@ -"use client"; - import { SearchIcon } from "lucide-react"; type ThreadSearchFieldProps = { diff --git a/apps/web/src/hooks/use-auto-animate-ref.ts b/apps/web/src/hooks/use-auto-animate-ref.ts index 579ba464..32ef5bab 100644 --- a/apps/web/src/hooks/use-auto-animate-ref.ts +++ b/apps/web/src/hooks/use-auto-animate-ref.ts @@ -9,9 +9,8 @@ const SIDEBAR_LIST_ANIMATION_OPTIONS = { export function useAutoAnimateRef() { const animatedNodesRef = useRef(new WeakSet()); return (node: HTMLElement | null) => { - if (!node || animatedNodesRef.current.has(node)) { - return; - } + if (!node || animatedNodesRef.current.has(node)) return; + autoAnimate(node, SIDEBAR_LIST_ANIMATION_OPTIONS); animatedNodesRef.current.add(node); }; diff --git a/apps/web/src/hooks/use-copy-to-clipboard.ts b/apps/web/src/hooks/use-copy-to-clipboard.ts index 9dfb977b..8aa27cd2 100644 --- a/apps/web/src/hooks/use-copy-to-clipboard.ts +++ b/apps/web/src/hooks/use-copy-to-clipboard.ts @@ -1,3 +1,4 @@ +import { Result } from "better-result"; import { useState } from "react"; export function useCopyToClipboard(): { @@ -6,26 +7,27 @@ export function useCopyToClipboard(): { } { const [copied, setCopied] = useState(false); - async function copy(text: string) { - try { - if (typeof navigator !== "undefined" && navigator.clipboard) { - await navigator.clipboard.writeText(text); - } else if (typeof document !== "undefined") { - const ta = document.createElement("textarea"); - ta.value = text; - ta.style.position = "fixed"; - ta.style.opacity = "0"; - document.body.appendChild(ta); - ta.select(); - document.execCommand("copy"); - document.body.removeChild(ta); - } - setCopied(true); - setTimeout(() => setCopied(false), 1400); - } catch { - setCopied(false); + async function copyPrimitive(text: string) { + if (typeof navigator !== "undefined" && navigator.clipboard) { + await navigator.clipboard.writeText(text); + } else if (typeof document !== "undefined") { + const ta = document.createElement("textarea"); + ta.value = text; + ta.style.position = "fixed"; + ta.style.opacity = "0"; + document.body.appendChild(ta); + ta.select(); + document.execCommand("copy"); + document.body.removeChild(ta); } + setCopied(true); + setTimeout(() => setCopied(false), 1400); } + async function copy(text: string) { + (await Result.tryPromise(() => copyPrimitive(text))).tapError(() => + setCopied(false) + ); + } return { copied, copy }; } diff --git a/apps/web/src/hooks/use-media-query.ts b/apps/web/src/hooks/use-media-query.ts index f95a28ef..014a38b8 100644 --- a/apps/web/src/hooks/use-media-query.ts +++ b/apps/web/src/hooks/use-media-query.ts @@ -2,16 +2,14 @@ import { useEffect, useState } from "react"; export function useMediaQuery(query: string): boolean { const [matches, setMatches] = useState(() => { - if (typeof window === "undefined") { - return false; - } + if (typeof window === "undefined") return false; + return window.matchMedia(query).matches; }); useEffect(() => { - if (typeof window === "undefined") { - return; - } + if (typeof window === "undefined") return; + const mql = window.matchMedia(query); const handler = (e: MediaQueryListEvent) => setMatches(e.matches); setMatches(mql.matches); diff --git a/apps/web/src/utils/dir.ts b/apps/web/src/utils/dir.ts index 79c4536c..313043c3 100644 --- a/apps/web/src/utils/dir.ts +++ b/apps/web/src/utils/dir.ts @@ -66,7 +66,7 @@ export function buildBrowseGroups(input: { }): BrowseGroup[] { const items: BrowseActionItem[] = []; - if (input.canBrowseUp) { + if (input.canBrowseUp) items.push({ kind: "action", value: "browse:up", @@ -75,9 +75,8 @@ export function buildBrowseGroups(input: { keepOpen: true, run: input.browseUp, }); - } - for (const entry of input.browseEntries) { + for (const entry of input.browseEntries) items.push({ kind: "action", value: `browse:${entry.fullPath}`, @@ -88,7 +87,6 @@ export function buildBrowseGroups(input: { input.browseTo(entry.name); }, }); - } return [{ value: "directories", label: "Directories", items }]; } From 796a5f6ccf28bdee9a433535e937ddced5d300fa Mon Sep 17 00:00:00 2001 From: Soorya U Date: Fri, 10 Jul 2026 11:10:50 +0530 Subject: [PATCH 3/7] Fix PR review issues for scoped live streaming. Ensure terminal chat events always reach the bus, fix waitForTurnEnd settlement, clone overlay state immutably, rethrow clipboard copy failures, and align OpenSpec paths with queue/bus.ts. Co-authored-by: Cursor --- apps/cli/src/handlers/controller/chat.ts | 74 ++++++++++++++++--- apps/web/src/hooks/use-copy-to-clipboard.ts | 8 +- .../design.md | 4 +- .../proposal.md | 2 +- .../specs/thread-live-stream/spec.md | 4 +- .../tasks.md | 2 +- openspec/specs/thread-live-stream/spec.md | 4 +- .../hooks/src/stores/conversation-overlay.ts | 22 ++++-- 8 files changed, 92 insertions(+), 28 deletions(-) diff --git a/apps/cli/src/handlers/controller/chat.ts b/apps/cli/src/handlers/controller/chat.ts index 441c6d0a..0b428bd8 100644 --- a/apps/cli/src/handlers/controller/chat.ts +++ b/apps/cli/src/handlers/controller/chat.ts @@ -27,10 +27,21 @@ async function runTurn({ projectId, message, emit, + emitTerminal, runtime, -}: Omit): Promise> { - await emit({ type: "user_message", content: message }); - await emit({ type: "thread_started", threadId }); +}: Omit & { + emitTerminal: ( + event: Extract< + ChatChunk["event"], + { type: "turn_completed" | "turn_interrupted" } + > + ) => Promise; +}): Promise> { + const started = await Result.tryPromise(async () => { + await emit({ type: "user_message", content: message }); + await emit({ type: "thread_started", threadId }); + }); + if (started.isErr()) return started; const streamed = await Result.tryPromise(async () => { const gen = runtime.threadCoordinator.prompt( @@ -40,12 +51,15 @@ async function runTurn({ message ); for await (const event of gen) await emit(event); - await emit({ type: "turn_completed" }); }); - if (streamed.isErr()) await emit({ type: "turn_interrupted" }); + if (streamed.isErr()) { + await emitTerminal({ type: "turn_interrupted" }); + return streamed; + } - return streamed; + await emitTerminal({ type: "turn_completed" }); + return Result.ok(undefined); } export function chatHandlers({ os, runtime }: ControllerDeps) { @@ -65,12 +79,15 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { context.eventBus.ensureWatch(context.peerId, threadId); + function publishChunk(chunk: ChatChunk): void { + context.eventBus.publish(chunk); + } + async function emit(event: ChatChunk["event"]): Promise { trackDelta(event, messageBuffers, thoughtBuffers); if (isStreamingDelta(event)) { - const chunk: ChatChunk = { threadId, turnId, seq: 0, event }; - context.eventBus.publish(chunk); + publishChunk({ threadId, turnId, seq: 0, event }); return; } @@ -85,7 +102,35 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { event: persistEvent, }); if (entry.isErr()) throwOrpcFromRepositoryError(entry.error); - context.eventBus.publish(entry.value.chunk); + publishChunk(entry.value.chunk); + } + + async function emitTerminal( + event: Extract< + ChatChunk["event"], + { type: "turn_completed" | "turn_interrupted" } + > + ): Promise { + trackDelta(event, messageBuffers, thoughtBuffers); + + const persistEvent = resolvePersistEvent( + event, + messageBuffers, + thoughtBuffers + ); + const entry = await appendConversation(threadId, { + threadId, + turnId, + event: persistEvent, + }); + + if (entry.isOk()) { + publishChunk(entry.value.chunk); + return; + } + + console.error("terminal event persist failed", entry.error); + publishChunk({ threadId, turnId, seq: 0, event }); } runTurn({ @@ -94,12 +139,17 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { projectId, message, emit, + emitTerminal, runtime, - }).then((result) => { - result.tapError((error) => { + }) + .then((result) => { + result.tapError((error) => { + console.error("chat turn failed", error); + }); + }) + .catch((error) => { console.error("chat turn failed", error); }); - }); return { threadId, turnId }; }), diff --git a/apps/web/src/hooks/use-copy-to-clipboard.ts b/apps/web/src/hooks/use-copy-to-clipboard.ts index 8aa27cd2..e018ecbc 100644 --- a/apps/web/src/hooks/use-copy-to-clipboard.ts +++ b/apps/web/src/hooks/use-copy-to-clipboard.ts @@ -25,9 +25,11 @@ export function useCopyToClipboard(): { } async function copy(text: string) { - (await Result.tryPromise(() => copyPrimitive(text))).tapError(() => - setCopied(false) - ); + const result = await Result.tryPromise(() => copyPrimitive(text)); + if (result.isErr()) { + setCopied(false); + throw result.error; + } } return { copied, copy }; } diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md index 2d769a08..475e01f4 100644 --- a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/design.md @@ -30,7 +30,7 @@ Mid-stream page refresh loses all in-flight token deltas: Turso has only persist ### Replace `PeerBroadcaster` with `ThreadEventBus` -**Choice:** Custom `ThreadEventBus` in `apps/cli/src/queue/index.ts` owns watch sets, per-peer delivery queues, and active-turn logs. No third-party event bus library — the API is domain-specific (thread watch sets, per-peer `AsyncGenerator`, active-turn replay). +**Choice:** Custom `ThreadEventBus` in `apps/cli/src/queue/bus.ts` owns watch sets, per-peer delivery queues, and active-turn logs. No third-party event bus library — the API is domain-specific (thread watch sets, per-peer `AsyncGenerator`, active-turn replay). **Why over patching `PeerBroadcaster`:** Watch filtering and turn replay are thread-scoped concerns, not peer-queue concerns. A clean type keeps `subscribe(peerId)` as the per-peer stream API while moving fan-out logic to `publish(chunk)`. @@ -132,6 +132,6 @@ Rollback: revert both server and web; old dual-path client is incompatible with ## Open Questions -- Exact per-turn log cap (10k chunks vs byte-size limit) — **decided: 10k default, delta-first eviction** (`apps/cli/src/queue/index.ts`). +- Exact per-turn log cap (10k chunks vs byte-size limit) — **decided: 10k default, delta-first eviction** (`apps/cli/src/queue/bus.ts`). - Should `waitForTurnEnd` timeout be configurable or a fixed constant (e.g. 30 min)? - Mobile (`apps/mobile`) — inherits shared `@cyrus/hooks` overlay once RTC is mounted; no separate overlay implementation needed. diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md index ca53b1e9..0266392a 100644 --- a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/proposal.md @@ -27,7 +27,7 @@ Every connected peer currently receives every chat chunk from every thread, and ## Impact -- `apps/cli/src/queue/index.ts` — custom `ThreadEventBus` (watch sets, active-turn logs, scoped fan-out). +- `apps/cli/src/queue/bus.ts` — custom `ThreadEventBus` (watch sets, active-turn logs, scoped fan-out). - `shared/connections/src/rtc/broadcaster.ts` — superseded for chat delivery once wired; may remain for other uses or be removed during integration. - `shared/connections/src/rtc/worker/index.ts` — instantiate `ThreadEventBus` instead of `createPeerBroadcaster`. - `apps/cli/src/handlers/controller/chat.ts` — `chat` becomes command; `emit` publishes to bus. diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md index d0c3a40d..e24e93f3 100644 --- a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md @@ -2,7 +2,7 @@ ### Requirement: Thread-scoped event bus replaces global broadcast -The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` implemented in `apps/cli/src/queue/index.ts` that fans out only to peers whose watched-thread set includes `chunk.threadId`. The bus SHALL deliver chunks to the initiating peer as well as all other watching peers. +The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` implemented in `apps/cli/src/queue/bus.ts` that fans out only to peers whose watched-thread set includes `chunk.threadId`. The bus SHALL deliver chunks to the initiating peer as well as all other watching peers. #### Scenario: Non-watching peer receives nothing @@ -23,7 +23,7 @@ The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` impl ### Requirement: Active-turn replay buffer -The `ThreadEventBus` in `apps/cli/src/queue/index.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. +The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. #### Scenario: Late watcher receives in-flight deltas diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md index 74c2a3e0..67b8f93d 100644 --- a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/tasks.md @@ -8,7 +8,7 @@ ## 2. ThreadEventBus (`apps/cli/src/queue`) -- [x] 2.1 Implement `ThreadEventBus` in `apps/cli/src/queue/index.ts` with watch sets, per-peer queues, and `publish()` +- [x] 2.1 Implement `ThreadEventBus` in `apps/cli/src/queue/bus.ts` with watch sets, per-peer queues, and `publish()` - [x] 2.2 Add `activeTurnLogs: Map` with append on publish and eviction on `turn_completed`/`turn_interrupted` - [x] 2.3 Implement `watch(peerId, threadId)` — register, replay active-turn logs into peer queue (idempotent) - [x] 2.4 Implement `unwatch(peerId, threadId)` and `ensureWatch(peerId, threadId)` for chat auto-watch diff --git a/openspec/specs/thread-live-stream/spec.md b/openspec/specs/thread-live-stream/spec.md index d0ed437e..30c98a28 100644 --- a/openspec/specs/thread-live-stream/spec.md +++ b/openspec/specs/thread-live-stream/spec.md @@ -6,7 +6,7 @@ Thread-scoped live `ChatChunk` delivery from the CLI worker to watching peers, w ### Requirement: Thread-scoped event bus replaces global broadcast -The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` implemented in `apps/cli/src/queue/index.ts` that fans out only to peers whose watched-thread set includes `chunk.threadId`. The bus SHALL deliver chunks to the initiating peer as well as all other watching peers. +The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` implemented in `apps/cli/src/queue/bus.ts` that fans out only to peers whose watched-thread set includes `chunk.threadId`. The bus SHALL deliver chunks to the initiating peer as well as all other watching peers. #### Scenario: Non-watching peer receives nothing @@ -27,7 +27,7 @@ The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` impl ### Requirement: Active-turn replay buffer -The `ThreadEventBus` in `apps/cli/src/queue/index.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. +The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. #### Scenario: Late watcher receives in-flight deltas diff --git a/shared/hooks/src/stores/conversation-overlay.ts b/shared/hooks/src/stores/conversation-overlay.ts index 1d3b0bde..7cc6c4db 100644 --- a/shared/hooks/src/stores/conversation-overlay.ts +++ b/shared/hooks/src/stores/conversation-overlay.ts @@ -122,12 +122,24 @@ function chunkToEntry(chunk: ChatChunk): ConversationEntry { }; } +function cloneOverlay(overlay: ThreadOverlay): ThreadOverlay { + return { + snapshotHighWaterMark: overlay.snapshotHighWaterMark, + live: [...overlay.live], + activeTurnIds: new Set(overlay.activeTurnIds), + }; +} + function getOrCreateOverlay( byThread: Map, threadId: string ): ThreadOverlay { const existing = byThread.get(threadId); - if (existing) return existing; + if (existing) { + const clone = cloneOverlay(existing); + byThread.set(threadId, clone); + return clone; + } const overlay: ThreadOverlay = { snapshotHighWaterMark: 0, @@ -148,7 +160,6 @@ function settleTurnWaiter( if (!waiter) return; turnWaiters.delete(key); - waiter.onAbort?.(); if (outcome === "completed") { waiter.resolve(); return; @@ -207,7 +218,7 @@ function commitLiveChunk(setState: OverlaySetState, chunk: ChatChunk): void { const overlay = getOrCreateOverlay(byThread, threadId); if (shouldSkipPersistedChunk(seq, overlay)) { - return { byThread }; + return state; } applyChunkToOverlay(overlay, chunk); @@ -263,9 +274,10 @@ export const useConversationOverlay = create( clearTurn(threadId, turnId) { flushPendingDeltas(commit); set((state) => { + if (!state.byThread.has(threadId)) return state; + const byThread = new Map(state.byThread); - const overlay = byThread.get(threadId); - if (!overlay) return { byThread }; + const overlay = getOrCreateOverlay(byThread, threadId); overlay.activeTurnIds.delete(turnId); overlay.live = overlay.live.filter( From e7f0db586644b0680c3750ea386d79de783019da Mon Sep 17 00:00:00 2001 From: Soorya U Date: Fri, 10 Jul 2026 11:42:42 +0530 Subject: [PATCH 4/7] Document active-turn replay bounds and use evlog for worker errors. Specify per-turn chunk limits, shutdown cleanup, and durable snapshot reconciliation in thread-live-stream specs; replace console.error with structured evlog in chat and conversation sync paths. Co-authored-by: Cursor --- apps/cli/package.json | 1 + apps/cli/src/handlers/controller/chat.ts | 19 +++++++++------ bun.lock | 1 + .../specs/thread-live-stream/spec.md | 24 ++++++++++++++++++- openspec/specs/thread-live-stream/spec.md | 24 ++++++++++++++++++- .../use-worker-conversation-sync.ts | 3 ++- .../hooks/src/stores/conversation-overlay.ts | 11 ++++----- 7 files changed, 66 insertions(+), 17 deletions(-) diff --git a/apps/cli/package.json b/apps/cli/package.json index 5b869e72..075f5d93 100644 --- a/apps/cli/package.json +++ b/apps/cli/package.json @@ -26,6 +26,7 @@ "better-result": "catalog:", "commander": "^15.0.0", "diff": "^7.0.0", + "evlog": "catalog:", "zod": "catalog:" }, "devDependencies": { diff --git a/apps/cli/src/handlers/controller/chat.ts b/apps/cli/src/handlers/controller/chat.ts index 0b428bd8..4d6dfa65 100644 --- a/apps/cli/src/handlers/controller/chat.ts +++ b/apps/cli/src/handlers/controller/chat.ts @@ -3,6 +3,7 @@ import { ensureThread } from "@cyrus/database/repositories/threads"; import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import { randomId } from "@cyrus/utils/identity"; import { Result } from "better-result"; +import { log } from "evlog"; import { throwOrpcFromRepositoryError } from "@/utils/error"; import { isStreamingDelta, @@ -86,10 +87,8 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { async function emit(event: ChatChunk["event"]): Promise { trackDelta(event, messageBuffers, thoughtBuffers); - if (isStreamingDelta(event)) { - publishChunk({ threadId, turnId, seq: 0, event }); - return; - } + if (isStreamingDelta(event)) + return publishChunk({ threadId, turnId, seq: 0, event }); const persistEvent = resolvePersistEvent( event, @@ -129,7 +128,13 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { return; } - console.error("terminal event persist failed", entry.error); + log.error({ + kind: "terminal_event_persist", + error: entry.error, + threadId, + turnId, + event: event.type, + }); publishChunk({ threadId, turnId, seq: 0, event }); } @@ -144,11 +149,11 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { }) .then((result) => { result.tapError((error) => { - console.error("chat turn failed", error); + log.error({ kind: "chat_turn_failed", error, threadId, turnId }); }); }) .catch((error) => { - console.error("chat turn failed", error); + log.error({ kind: "chat_turn_failed", error, threadId, turnId }); }); return { threadId, turnId }; diff --git a/bun.lock b/bun.lock index 11679aaf..55b53c92 100644 --- a/bun.lock +++ b/bun.lock @@ -34,6 +34,7 @@ "better-result": "catalog:", "commander": "^15.0.0", "diff": "^7.0.0", + "evlog": "catalog:", "zod": "catalog:", }, "devDependencies": { diff --git a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md index e24e93f3..31e11157 100644 --- a/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md +++ b/openspec/changes/archive/2026-07-10-scoped-thread-live-stream/specs/thread-live-stream/spec.md @@ -23,7 +23,7 @@ The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` impl ### Requirement: Active-turn replay buffer -The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. +The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. Each per-turn log SHALL be bounded to `maxChunksPerTurn` (default 10,000). When a log exceeds that limit, the bus SHALL evict ephemeral deltas (`seq === 0`) first, then the oldest remaining chunks, before appending new ones. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. On worker shutdown, `ThreadEventBus.closeAll()` SHALL clear all `activeTurnLogs` and `turnThreads` entries. When a chat turn fails or is cancelled, the worker SHALL publish `turn_interrupted` for that turn (including when durable persistence of the terminal event fails) so the buffer is evicted. Truncated replay is best-effort: clients that need full turn state after truncation or turn completion SHALL reconcile from durable history via `getConversations` using `snapshotHighWaterMark` from `watchThread`; the overlay merge path already applies that snapshot on refresh. #### Scenario: Late watcher receives in-flight deltas @@ -37,6 +37,28 @@ The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory - **THEN** `activeTurnLogs` for T is removed - **AND** a subsequent `watchThread` does not replay chunks from turn T (durable history comes from `getConversations`) +#### Scenario: Interrupted turn log is evicted + +- **WHEN** `turn_interrupted` is published for `turnId` T (including after a failed or cancelled turn) +- **THEN** `activeTurnLogs` for T is removed + +#### Scenario: Per-turn log is bounded + +- **WHEN** a turn emits more than `maxChunksPerTurn` chunks +- **THEN** the bus evicts ephemeral deltas first, then oldest chunks, keeping the log at or below the limit +- **AND** subsequent live chunks continue to be delivered to watching peers + +#### Scenario: Worker shutdown clears replay buffers + +- **WHEN** the worker shuts down and calls `ThreadEventBus.closeAll()` +- **THEN** all `activeTurnLogs` and `turnThreads` entries are cleared + +#### Scenario: Truncated replay defers to durable snapshot + +- **WHEN** a late watcher receives a truncated active-turn replay +- **AND** the turn later completes or the client refreshes +- **THEN** the client loads full persisted history from `getConversations` and merges it with any live overlay state + ### Requirement: watchThread and unwatchThread RPCs The controller contract SHALL expose `watchThread` and `unwatchThread` operations. `watchThread` SHALL register the calling peer's interest in a thread, replay active-turn logs for that thread, and return a `snapshotHighWaterMark` (the highest persisted `seq` for that thread). `unwatchThread` SHALL remove the registration. diff --git a/openspec/specs/thread-live-stream/spec.md b/openspec/specs/thread-live-stream/spec.md index 30c98a28..30ef2761 100644 --- a/openspec/specs/thread-live-stream/spec.md +++ b/openspec/specs/thread-live-stream/spec.md @@ -27,7 +27,7 @@ The worker SHALL route live `ChatChunk` delivery through a `ThreadEventBus` impl ### Requirement: Active-turn replay buffer -The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing all `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. +The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory `activeTurnLogs` buffer per in-progress `turnId` containing `ChatChunk`s emitted during that turn, including ephemeral deltas with `seq === 0`. Each per-turn log SHALL be bounded to `maxChunksPerTurn` (default 10,000). When a log exceeds that limit, the bus SHALL evict ephemeral deltas (`seq === 0`) first, then the oldest remaining chunks, before appending new ones. The buffer SHALL be evicted when `turn_completed` or `turn_interrupted` is published for that `turnId`. On worker shutdown, `ThreadEventBus.closeAll()` SHALL clear all `activeTurnLogs` and `turnThreads` entries. When a chat turn fails or is cancelled, the worker SHALL publish `turn_interrupted` for that turn (including when durable persistence of the terminal event fails) so the buffer is evicted. Truncated replay is best-effort: clients that need full turn state after truncation or turn completion SHALL reconcile from durable history via `getConversations` using `snapshotHighWaterMark` from `watchThread`; the overlay merge path already applies that snapshot on refresh. #### Scenario: Late watcher receives in-flight deltas @@ -41,6 +41,28 @@ The `ThreadEventBus` in `apps/cli/src/queue/bus.ts` SHALL maintain an in-memory - **THEN** `activeTurnLogs` for T is removed - **AND** a subsequent `watchThread` does not replay chunks from turn T (durable history comes from `getConversations`) +#### Scenario: Interrupted turn log is evicted + +- **WHEN** `turn_interrupted` is published for `turnId` T (including after a failed or cancelled turn) +- **THEN** `activeTurnLogs` for T is removed + +#### Scenario: Per-turn log is bounded + +- **WHEN** a turn emits more than `maxChunksPerTurn` chunks +- **THEN** the bus evicts ephemeral deltas first, then oldest chunks, keeping the log at or below the limit +- **AND** subsequent live chunks continue to be delivered to watching peers + +#### Scenario: Worker shutdown clears replay buffers + +- **WHEN** the worker shuts down and calls `ThreadEventBus.closeAll()` +- **THEN** all `activeTurnLogs` and `turnThreads` entries are cleared + +#### Scenario: Truncated replay defers to durable snapshot + +- **WHEN** a late watcher receives a truncated active-turn replay +- **AND** the turn later completes or the client refreshes +- **THEN** the client loads full persisted history from `getConversations` and merges it with any live overlay state + ### Requirement: watchThread and unwatchThread RPCs The controller contract SHALL expose `watchThread` and `unwatchThread` operations. `watchThread` SHALL register the calling peer's interest in a thread, replay active-turn logs for that thread, and return a `snapshotHighWaterMark` (the highest persisted `seq` for that thread). `unwatchThread` SHALL remove the registration. diff --git a/shared/hooks/src/connection/use-worker-conversation-sync.ts b/shared/hooks/src/connection/use-worker-conversation-sync.ts index 85dbc320..450b1536 100644 --- a/shared/hooks/src/connection/use-worker-conversation-sync.ts +++ b/shared/hooks/src/connection/use-worker-conversation-sync.ts @@ -2,6 +2,7 @@ import { RTC_OPERATION_KEYS } from "@cyrus/constants/operation-keys"; import type { ChatChunk } from "@cyrus/schemas/rtc/chat"; import { useQueryClient } from "@tanstack/react-query"; import { Result } from "better-result"; +import { log } from "evlog"; import { useEffect, useEffectEvent } from "react"; import { useRtc } from "../contexts/rtc"; import { useConversationOverlay } from "../stores/conversation-overlay"; @@ -29,7 +30,7 @@ export function useWorkerConversationSync(): void { }); const onSyncError = useEffectEvent((error: unknown) => { - console.error("worker conversation sync failed", error); + log.error({ kind: "worker_conversation_sync", error }); }); useEffect(() => { diff --git a/shared/hooks/src/stores/conversation-overlay.ts b/shared/hooks/src/stores/conversation-overlay.ts index 7cc6c4db..1c81bace 100644 --- a/shared/hooks/src/stores/conversation-overlay.ts +++ b/shared/hooks/src/stores/conversation-overlay.ts @@ -217,20 +217,17 @@ function commitLiveChunk(setState: OverlaySetState, chunk: ChatChunk): void { const byThread = new Map(state.byThread); const overlay = getOrCreateOverlay(byThread, threadId); - if (shouldSkipPersistedChunk(seq, overlay)) { - return state; - } + if (shouldSkipPersistedChunk(seq, overlay)) return state; applyChunkToOverlay(overlay, chunk); return { byThread }; }); - if (chunk.event.type === "turn_completed") { + if (chunk.event.type === "turn_completed") settleTurnWaiter(threadId, turnId, "completed"); - } - if (chunk.event.type === "turn_interrupted") { + + if (chunk.event.type === "turn_interrupted") settleTurnWaiter(threadId, turnId, "interrupted"); - } } export const useConversationOverlay = create( From 59cb5f12514dc9dfa12dfa0e916ea21a737b9f76 Mon Sep 17 00:00:00 2001 From: Soorya U Date: Fri, 10 Jul 2026 22:16:38 +0530 Subject: [PATCH 5/7] Stream thread conversations through the query cache and fix stop/follow-up UX. Replace the overlay store with shared conversation cache utilities, turn waiters, and subscribe-driven sync so live chunks and refetches stay consistent. Fix cancel races, message ordering after interrupt, and render collapsible thinking in the feed. Co-authored-by: Cursor --- apps/cli/src/handlers/controller/chat.ts | 26 +- apps/cli/src/queue/bus.ts | 8 + .../src/components/chat/composer/index.tsx | 5 +- .../chat/composer/primary-action.tsx | 14 + .../src/components/chat/feed/chat-feed.tsx | 57 ++- .../components/chat/feed/feed-entry-view.tsx | 5 + .../components/chat/main/thread-workspace.tsx | 25 +- .../chat/messages/assistant-loading.tsx | 14 + .../chat/messages/assistant-thinking.tsx | 42 ++ shared/connections/src/rtc/bus.ts | 1 + .../src/connection/use-controller-threads.ts | 154 +++++++- .../src/connection/use-thread-conversation.ts | 98 ++--- .../use-worker-conversation-sync.ts | 23 +- .../hooks/src/stores/conversation-overlay.ts | 351 ----------------- shared/hooks/src/use-thread-feed.ts | 165 ++++++-- shared/schemas/src/rtc/chat.ts | 1 + shared/schemas/src/view/index.ts | 10 + shared/utils/src/conversations/cache.ts | 362 ++++++++++++++++++ shared/utils/src/conversations/merger.ts | 56 +++ .../utils/src/conversations/turn-waiters.ts | 108 ++++++ shared/utils/src/fold.ts | 237 ++++++++++-- shared/utils/src/merge-conversation.ts | 21 - 22 files changed, 1228 insertions(+), 555 deletions(-) create mode 100644 apps/web/src/components/chat/messages/assistant-loading.tsx create mode 100644 apps/web/src/components/chat/messages/assistant-thinking.tsx delete mode 100644 shared/hooks/src/stores/conversation-overlay.ts create mode 100644 shared/utils/src/conversations/cache.ts create mode 100644 shared/utils/src/conversations/merger.ts create mode 100644 shared/utils/src/conversations/turn-waiters.ts delete mode 100644 shared/utils/src/merge-conversation.ts diff --git a/apps/cli/src/handlers/controller/chat.ts b/apps/cli/src/handlers/controller/chat.ts index 4d6dfa65..8f804e7f 100644 --- a/apps/cli/src/handlers/controller/chat.ts +++ b/apps/cli/src/handlers/controller/chat.ts @@ -66,7 +66,13 @@ async function runTurn({ export function chatHandlers({ os, runtime }: ControllerDeps) { return { chat: os.chat.handler(async ({ input, context }) => { - const { agentName, threadId = randomId(), message, projectId } = input; + const { + agentName, + threadId = randomId(), + turnId = randomId(), + message, + projectId, + } = input; const thread = await ensureThread(threadId, projectId, { agentName, @@ -74,7 +80,6 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { }); if (thread.isErr()) throwOrpcFromRepositoryError(thread.error); - const turnId = randomId(); const messageBuffers = new Map(); const thoughtBuffers = new Map(); @@ -164,8 +169,23 @@ export function chatHandlers({ os, runtime }: ControllerDeps) { yield chunk; }), - cancel: os.cancel.handler(async ({ input }) => { + cancel: os.cancel.handler(async ({ input, context }) => { + // Snapshot before awaiting — a follow-up chat() can start while cancel + // is in flight and must not be interrupted by this request. + const activeTurnIds = context.eventBus.getActiveTurnIdsForThread( + input.threadId + ); + await runtime.threadCoordinator.cancel(input.agentName, input.threadId); + + for (const turnId of activeTurnIds) + context.eventBus.publish({ + threadId: input.threadId, + turnId, + seq: 0, + event: { type: "turn_interrupted" }, + }); + return {}; }), }; diff --git a/apps/cli/src/queue/bus.ts b/apps/cli/src/queue/bus.ts index 30b9b924..98233d36 100644 --- a/apps/cli/src/queue/bus.ts +++ b/apps/cli/src/queue/bus.ts @@ -130,6 +130,14 @@ export function createThreadEventBus( return watchedThreads.get(peerId)?.has(threadId) ?? false; }, + getActiveTurnIdsForThread(threadId) { + const turnIds: string[] = []; + for (const turnId of activeTurnLogs.keys()) { + if (turnThreads.get(turnId) === threadId) turnIds.push(turnId); + } + return turnIds; + }, + async *subscribe(peerId) { closePeerDelivery(peerId); const peer: PeerDelivery = { diff --git a/apps/web/src/components/chat/composer/index.tsx b/apps/web/src/components/chat/composer/index.tsx index 3b2f1530..5dd898da 100644 --- a/apps/web/src/components/chat/composer/index.tsx +++ b/apps/web/src/components/chat/composer/index.tsx @@ -8,12 +8,14 @@ export function Composer({ onSend, onStop, busy = false, + stopping = false, }: { projectId: string; threadId: string; onSend: (text: string) => void; onStop?: () => void; busy?: boolean; + stopping?: boolean; }) { const [value, setValue] = useState(""); const textareaRef = useRef(null); @@ -30,7 +32,7 @@ export function Composer({ function submit() { const text = value.trim(); - if (!text || busy) return; + if (!text || busy || stopping) return; onSend(text); setValue(""); @@ -97,6 +99,7 @@ export function Composer({ busy={busy} canSend={value.trim().length > 0} onStop={onStop} + stopping={stopping} />
diff --git a/apps/web/src/components/chat/composer/primary-action.tsx b/apps/web/src/components/chat/composer/primary-action.tsx index fbceba80..5f520117 100644 --- a/apps/web/src/components/chat/composer/primary-action.tsx +++ b/apps/web/src/components/chat/composer/primary-action.tsx @@ -4,11 +4,25 @@ export function ComposerPrimaryAction({ busy, canSend, onStop, + stopping = false, }: { busy: boolean; canSend: boolean; onStop?: () => void; + stopping?: boolean; }) { + if (stopping) + return ( + + ); + if (busy) return (