From cccbe6598a3d369117cebf4ce9eed64c2a92db8a Mon Sep 17 00:00:00 2001 From: Kieran Klukas Date: Mon, 25 May 2026 21:47:11 -0400 Subject: [PATCH] feat: streaming simplify --- web/src/lib/stream.ts | 34 +- web/src/routes/chat/+page.svelte | 814 +++++------------------ web/src/routes/chat/ModelPicker.svelte | 246 +++++++ web/src/routes/chat/ToolCallBlock.svelte | 186 ++++++ 4 files changed, 618 insertions(+), 662 deletions(-) create mode 100644 web/src/routes/chat/ModelPicker.svelte create mode 100644 web/src/routes/chat/ToolCallBlock.svelte diff --git a/web/src/lib/stream.ts b/web/src/lib/stream.ts index d325ab2..de3bfc1 100644 --- a/web/src/lib/stream.ts +++ b/web/src/lib/stream.ts @@ -1,10 +1,17 @@ /** * Client-side SSE consumer with resume. * - * Subscribes to GET /api/streams/:id/events. On disconnect, reconnects with + * Subscribes to an SSE endpoint. On disconnect, reconnects with * `?after_seq=` so the server replays missed chunks from DB before * attaching us back to the live bus. * + * Pass a stream ID string to use the default stream endpoint + * (`/api/streams/:id/events?after_seq=N`), or pass a URL builder function for + * custom endpoints (e.g. the conversation bus at `/api/conversations/:id/events`). + * + * Set `reconnectAlways` for long-lived push channels that should always + * reconnect, even if no events arrived on the current connection. + * * This is the supported path for streaming; do not adopt @ai-sdk/svelte's * `Chat` (it owns its own state machine and fights the Dexie-first model). */ @@ -23,15 +30,25 @@ export interface StreamHandlers { onClose?(reason: 'done' | 'error' | 'aborted'): void; } -export function consume(streamID: string, handlers: StreamHandlers, initialAfterSeq = 0): () => void { +export function consume( + streamID: string | ((afterSeq: number) => string), + handlers: StreamHandlers, + initialAfterSeq = 0, + { reconnectAlways = false }: { reconnectAlways?: boolean } = {} +): () => void { let aborted = false; let lastSeq = initialAfterSeq; let controller = new AbortController(); + const buildUrl = + typeof streamID === 'function' + ? streamID + : (seq: number) => `/api/streams/${streamID}/events?after_seq=${Math.max(0, seq)}`; + const loop = async () => { while (!aborted) { try { - const res = await fetch(`/api/streams/${streamID}/events?after_seq=${Math.max(0, lastSeq)}`, { + const res = await fetch(buildUrl(lastSeq), { credentials: 'include', signal: controller.signal }); @@ -42,7 +59,7 @@ export function consume(streamID: string, handlers: StreamHandlers, initialAfter const reader = res.body.pipeThrough(new TextDecoderStream()).getReader(); let buf = ''; - // SSE frames are separated by a blank line. + let gotEvent = false; for (;;) { const { value, done } = await reader.read(); if (done) break; @@ -60,6 +77,7 @@ export function consume(streamID: string, handlers: StreamHandlers, initialAfter try { const ev = JSON.parse(dataLine) as StreamEvent; if (ev.seq > lastSeq) lastSeq = ev.seq; + gotEvent = true; handlers.onEvent(ev); if (ev.type === 'done' || ev.type === 'error') { handlers.onClose?.(ev.type); @@ -70,7 +88,13 @@ export function consume(streamID: string, handlers: StreamHandlers, initialAfter } } } - // EOF without done: server hung up. Reconnect with after_seq. + // EOF without done: if no events arrived and we're not in reconnect-always + // mode, the stream is stale/dead — give up so the UI doesn't spin forever. + if (!gotEvent && !reconnectAlways) { + handlers.onClose?.('error'); + return; + } + // Reconnect with after_seq (or always, for push channels). } catch (_e) { if (aborted) return; // network error — back off briefly and retry diff --git a/web/src/routes/chat/+page.svelte b/web/src/routes/chat/+page.svelte index f06fb8a..7217cc5 100644 --- a/web/src/routes/chat/+page.svelte +++ b/web/src/routes/chat/+page.svelte @@ -8,6 +8,8 @@ import { listConversations, listMessages, listModels, type Model } from '$lib/api'; import { renderMarkdown, renderStreamingMarkdown } from '$lib/markdown'; import { consume, type StreamEvent } from '$lib/stream'; + import ModelPicker from './ModelPicker.svelte'; + import ToolCallBlock from './ToolCallBlock.svelte'; // ── state ───────────────────────────────────────────────────────────────── @@ -22,13 +24,13 @@ let models = $state([]); let selectedModel = $state(''); - let modelPickerOpen = $state(false); + let messages = $state([]); let msgsEl = $state(null); let inputEl = $state(null); - let pickerEl = $state(null); + let atBottom = true; let errorMsg = $state(null); @@ -37,7 +39,7 @@ // finished assistant message. let toolCalls = $state>(new Map()); let toolMsgId = $state(null); - let toolExpanded = $state(false); + // ── spin cursor (crush-style cycling glyphs) ────────────────────────────── @@ -180,16 +182,23 @@ // Sync conversations so the layout sidebar has data. syncConversations(); - // Resume in-flight stream after reload. + // Resume in-flight stream after reload — only if the URL actually matches + // the stale conv (i.e. this is a true reload, not a navigation to /chat). + // Mismatched stale data means the stream finished or was abandoned; clear it. + const urlConvIdAtMount = page.url.searchParams.get('c'); if (resumeStreamId && resumeConvId && resumeAssistantId) { - activeConvId = resumeConvId; - streaming = true; - activeStreamId = resumeStreamId; - streamingMsgId = resumeAssistantId; - activeAfterSeq = Number.isFinite(resumeAfterSeq) ? resumeAfterSeq : 0; - streamStartMs = Date.now(); - streamFirstTokenMs = 0; - attachStreamConsumer(resumeStreamId, resumeConvId, resumeAssistantId, Math.floor(Date.now() / 1000), resumeAfterSeq); + if (urlConvIdAtMount !== resumeConvId) { + clearActiveStreamStorage(); + } else { + activeConvId = resumeConvId; + streaming = true; + activeStreamId = resumeStreamId; + streamingMsgId = resumeAssistantId; + activeAfterSeq = Number.isFinite(resumeAfterSeq) ? resumeAfterSeq : 0; + streamStartMs = Date.now(); + streamFirstTokenMs = 0; + attachStreamConsumer(resumeStreamId, resumeConvId, resumeAssistantId, Math.floor(Date.now() / 1000), resumeAfterSeq); + } } }); @@ -226,53 +235,15 @@ }); function watchConversation(convId: string): () => void { - let aborted = false; - let controller = new AbortController(); - - const loop = async () => { - while (!aborted) { - try { - const res = await fetch(`/api/conversations/${convId}/events`, { - credentials: 'include', - signal: controller.signal - }); - if (!res.ok || !res.body) return; - - const reader = res.body.pipeThrough(new TextDecoderStream()).getReader(); - let buf = ''; - for (;;) { - const { value, done } = await reader.read(); - if (done) break; - buf += value; - let idx: number; - while ((idx = buf.indexOf('\n\n')) >= 0) { - const frame = buf.slice(0, idx); - buf = buf.slice(idx + 2); - const dataLine = frame - .split('\n') - .find((l) => l.startsWith('data:')) - ?.slice(5) - .trim(); - if (!dataLine) continue; - try { - const ev = JSON.parse(dataLine); - await handleConvEvent(ev, convId); - } catch { /* skip malformed */ } - } - } - } catch { - if (aborted) return; - await new Promise((r) => setTimeout(r, 1000)); - controller = new AbortController(); - } - } - }; - - loop(); - return () => { aborted = true; controller.abort(); }; + return consume( + () => `/api/conversations/${convId}/events`, + { onEvent: (ev) => { handleConvEvent(ev, convId); } }, + 0, + { reconnectAlways: true } + ); } - async function handleConvEvent(ev: Record, convId: string) { + async function handleConvEvent(ev: StreamEvent, convId: string) { if (ev.type !== 'start') return; if (!ev.stream_id || !ev.assistant_message_id) return; // Ignore if we're the one already consuming this stream. @@ -310,54 +281,7 @@ // ── model selection ─────────────────────────────────────────────────────── - let freeModels = $derived(models.filter((m) => m.id.startsWith('free/'))); - let paidModels = $derived(models.filter((m) => !m.id.startsWith('free/'))); - - let modelSearch = $state(''); - let filteredFreeModels = $derived( - modelSearch.trim() - ? freeModels.filter((m) => (m.label || m.id).toLowerCase().includes(modelSearch.toLowerCase())) - : freeModels - ); - let filteredPaidModels = $derived( - modelSearch.trim() - ? paidModels.filter((m) => (m.label || m.id).toLowerCase().includes(modelSearch.toLowerCase())) - : paidModels - ); - - function saveModel() { - localStorage.setItem('chat:model', selectedModel); - } - function modelDisplay(id: string) { - if (!id) return 'pick model'; - const m = models.find((x) => x.id === id); - if (m) return m.label || id.replace(/^free\//, ''); - return id.replace(/^free\//, ''); - } - - function pickModel(id: string) { - selectedModel = id; - saveModel(); - modelPickerOpen = false; - modelSearch = ''; - inputEl?.focus(); - } - - function closePicker() { - modelPickerOpen = false; - modelSearch = ''; - } - - function focusEl(el: HTMLElement) { - el.focus(); - } - - function onWindowClick(e: MouseEvent) { - if (modelPickerOpen && pickerEl && !pickerEl.contains(e.target as Node)) { - closePicker(); - } - } // ── input ───────────────────────────────────────────────────────────────── @@ -366,9 +290,6 @@ e.preventDefault(); sendMessage(); } - if (e.key === 'Escape') { - modelPickerOpen = false; - } } function resize(el: HTMLTextAreaElement) { @@ -392,13 +313,12 @@ localStorage.removeItem('chat:active_stream_after_seq'); } - function attachStreamConsumer(streamId: string, convId: string, assistantId: string, now: number, initialAfterSeq = 0) { - if (stopActiveConsumer) { - stopActiveConsumer(); - stopActiveConsumer = null; - } + // Unified stream event handler used by both sendMessage (after handoff) + // and resume/cross-tab flows via attachStreamConsumer. + function makeStreamHandlers(convId: string, assistantId: string, now: number) { + let accContent = ''; - stopActiveConsumer = consume(streamId, { + return { onEvent: async (ev: StreamEvent) => { if (ev.seq > activeAfterSeq) { activeAfterSeq = ev.seq; @@ -407,14 +327,19 @@ if (ev.type === 'delta' && ev.content) { if (!streamFirstTokenMs) streamFirstTokenMs = Date.now(); - const msg = await db.messages.get(assistantId); - const current = msg?.content ?? ''; + if (!accContent) { + // Resume path: seed accContent from DB so replayed deltas append + // to existing content rather than overwriting it. + const existing = await db.messages.get(assistantId); + accContent = existing?.content ?? ''; + } + accContent += ev.content; await db.messages.put({ id: assistantId, conversation_id: convId, client_id: null, role: 'assistant', - content: current + ev.content, + content: accContent, model: selectedModel, created_at: now + 1, pending: true @@ -439,53 +364,70 @@ if (existing) { toolCalls.set(tcId, { ...existing, result: ((ev as any).content as string) || '', done: true }); } else { + let matched = false; for (const [id, tc] of toolCalls) { if (!tc.done) { toolCalls.set(id, { ...tc, result: ((ev as any).content as string) || '', done: true }); + matched = true; break; } } + if (!matched) { + toolCalls.set(tcId || uuid(), { + name: 'tool', + args: '', + result: ((ev as any).content as string) || '', + done: true + }); + } } toolCalls = new Map(toolCalls); } if (ev.type === 'done' || ev.type === 'error') { const doneMs = Date.now(); - const completionTokens = ev.usage?.output_tokens; + const completionTokens = (ev.usage as any)?.completion_tokens ?? ev.usage?.output_tokens; const ttft = streamFirstTokenMs ? streamFirstTokenMs - streamStartMs : undefined; const tps = completionTokens && streamFirstTokenMs ? completionTokens / ((doneMs - streamFirstTokenMs) / 1000) : undefined; - await db.messages.where({ id: assistantId }).modify({ - pending: false, - ttft, - tps, - tokens: completionTokens - }); - await db.conversations.where({ id: convId }).modify({ updated_at: Math.floor(doneMs / 1000) }); - - streaming = false; - streamingMsgId = null; - activeStreamId = null; - activeAfterSeq = 0; - clearActiveStreamStorage(); - - if (stopActiveConsumer) { - stopActiveConsumer(); - stopActiveConsumer = null; - } - - if (ev.type === 'error') { + if (ev.type === 'done') { + await db.messages.where({ id: assistantId }).modify({ + pending: false, + ttft, + tps, + tokens: completionTokens + }); + await db.conversations.where({ id: convId }).modify({ updated_at: Math.floor(doneMs / 1000) }); + } else { + await db.messages.delete(assistantId); errorMsg = ev.error?.message || (ev as any).message || 'stream error'; } + // UI state teardown: onClose (attach path) or finally (inline path). } }, - onClose: (reason) => { + onClose: (reason: string) => { if (reason === 'aborted') return; + streaming = false; + streamingMsgId = null; + activeStreamId = null; + activeAfterSeq = 0; + clearActiveStreamStorage(); + // consume() has already terminated when onClose fires — just null the ref. + stopActiveConsumer = null; } - }, initialAfterSeq); + }; + } + + function attachStreamConsumer(streamId: string, convId: string, assistantId: string, now: number, initialAfterSeq = 0) { + if (stopActiveConsumer) { + stopActiveConsumer(); + stopActiveConsumer = null; + } + + stopActiveConsumer = consume(streamId, makeStreamHandlers(convId, assistantId, now), initialAfterSeq); } async function sendMessage() { @@ -499,7 +441,7 @@ errorMsg = null; toolCalls = new Map(); toolMsgId = null; - toolExpanded = false; + streamStartMs = Date.now(); streamFirstTokenMs = 0; streamElapsedMs = 0; @@ -560,10 +502,6 @@ ...(m.id === tmpUserId ? { client_id: clientId } : {}) })); - let resolvedUserId = tmpUserId; - let resolvedAssistantId = tmpAssistantId; - let accContent = ''; - try { const res = await fetch('/api/chat', { method: 'POST', @@ -584,6 +522,7 @@ const reader = res.body.pipeThrough(new TextDecoderStream()).getReader(); let buf = ''; + let resolvedHandlers: ReturnType | null = null; outer: for (;;) { const { value, done } = await reader.read(); @@ -609,177 +548,92 @@ continue; } - // Keep activeAfterSeq current so resume picks up from the right spot. if (typeof ev.seq === 'number' && ev.seq > activeAfterSeq) { activeAfterSeq = ev.seq; - if (activeStreamId) { - localStorage.setItem('chat:active_stream_after_seq', String(activeAfterSeq)); - } + localStorage.setItem('chat:active_stream_after_seq', String(activeAfterSeq)); } - switch (ev.type) { - case 'start': { - const serverConvId: string = ev.conversation_id; - const serverUserId: string = ev.user_message_id; - const serverAssistantId: string = ev.assistant_message_id; - const serverStreamId: string | undefined = ev.stream_id; - - const outerConvId: string = convId!; - - // Resolve all ID swaps in a single transaction so the live query - // sees exactly one consistent snapshot. - await db.transaction('rw', db.conversations, db.messages, async () => { - // New conversation ID swap. - if (isNewConv && serverConvId && serverConvId !== outerConvId) { - const tmpConv = await db.conversations.get(outerConvId); - if (tmpConv) { - await db.conversations.delete(outerConvId); - await db.conversations.put({ ...tmpConv, id: serverConvId }); - const affected = await db.messages - .where('conversation_id') - .equals(outerConvId) - .toArray(); - await db.messages.bulkPut(affected.map((m) => ({ ...m, conversation_id: serverConvId }))); - await db.messages.where('conversation_id').equals(outerConvId).delete(); - } - } - - // User message ID swap. - if (serverUserId && serverUserId !== tmpUserId) { - const old = await db.messages.get(tmpUserId); - if (old) { - await db.messages.delete(tmpUserId); - await db.messages.put({ ...old, id: serverUserId, pending: false }); - resolvedUserId = serverUserId; - } - } else { - await db.messages.where({ id: resolvedUserId }).modify({ pending: false }); + if (ev.type === 'start') { + const serverConvId: string = ev.conversation_id; + const serverUserId: string = ev.user_message_id; + const serverAssistantId: string = ev.assistant_message_id || tmpAssistantId; + const serverStreamId: string | undefined = ev.stream_id; + const localConvId = convId!; + + await db.transaction('rw', db.conversations, db.messages, async () => { + if (isNewConv && serverConvId && serverConvId !== localConvId) { + const tmpConv = await db.conversations.get(localConvId); + if (tmpConv) { + await db.conversations.delete(localConvId); + await db.conversations.put({ ...tmpConv, id: serverConvId }); + const affected = await db.messages + .where('conversation_id') + .equals(localConvId) + .toArray(); + await db.messages.bulkPut(affected.map((m) => ({ ...m, conversation_id: serverConvId }))); + await db.messages.where('conversation_id').equals(localConvId).delete(); } + } - // Assistant message ID swap. - if (serverAssistantId && serverAssistantId !== tmpAssistantId) { - const old = await db.messages.get(tmpAssistantId); - if (old) { - await db.messages.delete(tmpAssistantId); - await db.messages.put({ ...old, id: serverAssistantId }); - resolvedAssistantId = serverAssistantId; - streamingMsgId = serverAssistantId; - } + if (serverUserId && serverUserId !== tmpUserId) { + const old = await db.messages.get(tmpUserId); + if (old) { + await db.messages.delete(tmpUserId); + await db.messages.put({ ...old, id: serverUserId, pending: false }); } - }); - - if (isNewConv && serverConvId && serverConvId !== outerConvId) { - convId = serverConvId; - activeConvId = serverConvId; - goto(`/chat?c=${serverConvId}`, { replaceState: true, noScroll: true, keepFocus: true }); + } else { + await db.messages.where({ id: tmpUserId }).modify({ pending: false }); } - // Save stream metadata for resume-on-reload — don't hand off mid-stream. - if (serverStreamId) { - activeStreamId = serverStreamId; - localStorage.setItem('chat:active_stream_id', serverStreamId); - localStorage.setItem('chat:active_stream_conv_id', convId); - localStorage.setItem('chat:active_stream_assistant_id', resolvedAssistantId); - localStorage.setItem('chat:active_stream_after_seq', String(activeAfterSeq)); + if (serverAssistantId && serverAssistantId !== tmpAssistantId) { + const old = await db.messages.get(tmpAssistantId); + if (old) { + await db.messages.delete(tmpAssistantId); + await db.messages.put({ ...old, id: serverAssistantId }); + streamingMsgId = serverAssistantId; + } } + }); - break; + if (isNewConv && serverConvId && serverConvId !== localConvId) { + convId = serverConvId; + activeConvId = serverConvId; + goto(`/chat?c=${serverConvId}`, { replaceState: true, noScroll: true, keepFocus: true }); } - case 'delta': { - if (!streamFirstTokenMs) { - streamFirstTokenMs = Date.now(); - } - accContent += ev.content as string; - await db.messages.put({ - id: resolvedAssistantId, - conversation_id: convId, - client_id: null, - role: 'assistant', - content: accContent, - model: selectedModel, - created_at: now + 1, - pending: true - }); - break; - } - - case 'done': { - const completionTokens: number | undefined = (ev.usage as any)?.completion_tokens; - const ttft = streamFirstTokenMs ? streamFirstTokenMs - streamStartMs : undefined; - const doneMs = Date.now(); - const tps = - completionTokens && streamFirstTokenMs - ? completionTokens / ((doneMs - streamFirstTokenMs) / 1000) - : undefined; - await db.messages.where({ id: resolvedAssistantId }).modify({ - pending: false, - ttft, - tps, - tokens: completionTokens - }); - await db.conversations.where({ id: convId }).modify({ updated_at: Math.floor(doneMs / 1000) }); - activeStreamId = null; - activeAfterSeq = 0; - clearActiveStreamStorage(); - break outer; + if (serverStreamId) { + activeStreamId = serverStreamId; + localStorage.setItem('chat:active_stream_id', serverStreamId); + localStorage.setItem('chat:active_stream_conv_id', convId); + localStorage.setItem('chat:active_stream_assistant_id', streamingMsgId || tmpAssistantId); + localStorage.setItem('chat:active_stream_after_seq', String(activeAfterSeq)); } - case 'error': { - await db.messages.delete(resolvedAssistantId); - errorMsg = ev.message as string; - activeStreamId = null; - activeAfterSeq = 0; - clearActiveStreamStorage(); - break outer; - } + resolvedHandlers = makeStreamHandlers(convId, streamingMsgId || tmpAssistantId, now); + continue; + } - case 'tool_call': { - if (!toolMsgId) toolMsgId = resolvedAssistantId; - const tcId = (ev.id as string) || uuid(); - toolCalls.set(tcId, { - name: (ev.name as string) || 'tool', - args: (ev.arguments as string) || '', - result: '', - done: false - }); - toolCalls = new Map(toolCalls); - break; - } + if (ev.type === 'error' && !resolvedHandlers) { + await db.messages.delete(tmpAssistantId); + errorMsg = ev.message as string; + return; + } - case 'tool_result': { - const tcId = (ev.tool_call_id as string) || ''; - const existing = tcId ? toolCalls.get(tcId) : null; - if (existing) { - toolCalls.set(tcId, { ...existing, result: (ev.content as string) || '', done: true }); - } else { - for (const [id, tc] of toolCalls) { - if (!tc.done) { - toolCalls.set(id, { ...tc, result: (ev.content as string) || '', done: true }); - break; - } - } - if (![...toolCalls.values()].some(tc => tc.done && tc.result)) { - toolCalls.set(tcId || uuid(), { - name: 'tool', - args: '', - result: (ev.content as string) || '', - done: true - }); - } - } - toolCalls = new Map(toolCalls); - break; - } + if (resolvedHandlers) { + await resolvedHandlers.onEvent(ev as StreamEvent); + if (ev.type === 'done' || ev.type === 'error') break outer; } } } } catch (err) { - await db.messages.delete(resolvedAssistantId); + await db.messages.delete(streamingMsgId || tmpAssistantId); errorMsg = err instanceof Error ? err.message : 'Something went wrong'; } finally { streaming = false; streamingMsgId = null; + activeStreamId = null; + activeAfterSeq = 0; + clearActiveStreamStorage(); } } @@ -795,8 +649,6 @@ } - - chat · potluck @@ -828,62 +680,9 @@
{#if msg.id === toolMsgId && toolCalls.size > 0} - {@const entries = [...toolCalls.entries()]} - {@const running = entries.filter(([, tc]) => !tc.done)} - {@const done = entries.filter(([, tc]) => tc.done)} -
- - - - - {#if toolExpanded} -
- {#each entries as [id, tc] (id)} -
-
- {#if tc.done} - ✓ - {:else} - - {/if} - {tc.name} -
- {#if tc.args} -
- arguments -
{tc.args}
-
- {/if} - {#if tc.result} -
- result -
{tc.result}
-
- {/if} -
- {/each} -
- {/if} -
+ {/if} - {#if msg.pending && msg.id === streamingMsgId} + {#if streaming && msg.pending} {@html renderStreamingMarkdown(msg.content)} {:else if msg.content} {@html renderMarkdown(msg.content)} @@ -891,7 +690,7 @@ {/if}
- {#if msg.pending && msg.id === streamingMsgId} + {#if streaming && msg.pending} {spinChars} {/if} {#if msg.model && typeof msg.model === 'string'} @@ -941,70 +740,16 @@ rows={1} > + + diff --git a/web/src/routes/chat/ToolCallBlock.svelte b/web/src/routes/chat/ToolCallBlock.svelte new file mode 100644 index 0000000..d0df5c2 --- /dev/null +++ b/web/src/routes/chat/ToolCallBlock.svelte @@ -0,0 +1,186 @@ + + +
+ + + {#if open} +
+ {#each entries as [id, tc] (id)} +
+
+ {#if tc.done} + ✓ + {:else} + + {/if} + {tc.name} +
+ {#if tc.args} +
+ arguments +
{tc.args}
+
+ {/if} + {#if tc.result} +
+ result +
{tc.result}
+
+ {/if} +
+ {/each} +
+ {/if} +
+ + -- 2.51.2