Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
21 commits
Select commit Hold shift + click to select a range
21d367c
feat(js): add StreamSmoother for paced text emission
cpsievert Sep 21, 2026
95d71b9
test(js): cover StreamSmoother pacing, word boundaries, backoff
cpsievert Sep 21, 2026
54cedba
fix(js): refine StreamSmoother word boundary logic for proper pacing
cpsievert Sep 21, 2026
2cd36f6
fix(js): accumulate elapsed time across word boundary stalls
cpsievert Sep 21, 2026
bac92ca
feat(js): pace streamed chat chunks through StreamSmoother
cpsievert Sep 21, 2026
35995db
fix(js): remove test-environment detection from ChatApp streaming
cpsievert Sep 21, 2026
b9dbfde
test(js): account for pacing delay in streaming-dot test timing
cpsievert Sep 21, 2026
cad53eb
feat(js): pace MarkdownStream text through StreamSmoother
cpsievert Sep 21, 2026
95e7014
test(js): allow enough ticks for full drain in pacing test
cpsievert Sep 21, 2026
f3c6380
fix(js): satisfy strict typecheck and prettier in StreamSmoother
cpsievert Sep 21, 2026
7e19b44
chore(js): rebuild dist assets for streaming smoothing
cpsievert Sep 21, 2026
c271c18
fix(js): correct StreamSmoother congestion backoff formula
cpsievert Sep 21, 2026
c65e781
chore(js): rebuild dist assets for backoff fix
cpsievert Sep 21, 2026
6819784
Merge remote-tracking branch 'origin/main' into worktree-streaming-sm…
cpsievert Sep 29, 2026
1b964f5
fix(js): auto-scroll after final content flush at stream end
cpsievert Oct 6, 2026
78bccc2
chore(js): rebuild dist assets for end-of-stream auto-scroll fix
cpsievert Oct 6, 2026
2399e23
feat(js): pace streamed text at character granularity and drain the tail
cpsievert Oct 6, 2026
574fd92
feat(js): show the streaming dot only when a stream stalls
cpsievert Oct 6, 2026
39e9fb5
chore(js): rebuild dist assets for character pacing and idle streamin…
cpsievert Oct 6, 2026
b5d25a4
feat(js): never end a paced emission inside an unclosed tag
cpsievert Oct 6, 2026
a241c7e
chore(js): rebuild dist assets for character pacing, tag holds, and i…
cpsievert Oct 6, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion js/dist/shinychat.css

Large diffs are not rendered by default.

4 changes: 2 additions & 2 deletions js/dist/shinychat.css.map

Large diffs are not rendered by default.

152 changes: 76 additions & 76 deletions js/dist/shinychat.js

Large diffs are not rendered by default.

8 changes: 4 additions & 4 deletions js/dist/shinychat.js.map

Large diffs are not rendered by default.

131 changes: 128 additions & 3 deletions js/src/chat/ChatApp.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -26,21 +26,49 @@ import {
type ChatDrawerState,
type GreetingData,
type ToolGrouping,
type AnyAction,
} from "./state"
import {
useRequestDefinitionIcons,
useSupersededRequests,
} from "./useSupersededRequests"
import { ChatContainer, type ChatContainerHandle } from "./ChatContainer"
import { acquireHistoryStore, getHistoryStore } from "./historyStore"
import { StreamSmoother } from "../streaming/StreamSmoother"
import type {
ChatTransport,
ShinyLifecycle,
GreetingOptions,
ContentType,
ChatAction,
} from "../transport/types"
import type { HtmlDep } from "rstudio-shiny/srcts/types/src/shiny/render"
import type { SubmitKey } from "./tiptap/submitShortcut"
import type { AttachmentPayload } from "./attachments"

interface ChunkMeta {
content_type?: ContentType
html_deps?: HtmlDep[]
}

// Actions that start new content or wipe the transcript. Arriving while a
// finished stream's tail is draining, they complete the drain immediately
// (the user or server has moved on). Everything else waits for the drain.
const ENDS_DRAIN = new Set<ChatAction["type"]>([
"message",
"chunk_start",
"chunk",
"block_insert",
"chunk_end",
"clear",
"greeting",
"greeting_start",
"greeting_chunk",
"greeting_end",
"greeting_clear",
"history_navigate",
])

export interface InitialGreeting {
content: string
contentType: import("../transport/types").ContentType
Expand Down Expand Up @@ -175,14 +203,40 @@ export function ChatApp({
}, [elementId, historyStore, transport])

const containerRef = useRef<ChatContainerHandle>(null)
const smootherRef = useRef<StreamSmoother<ChunkMeta> | null>(null)
const siblingNavigationPendingRef = useRef(false)
const [siblingNavigationPending, setSiblingNavigationPending] =
useState(false)

// The textarea is fully uncontrolled, so value/focus mutations go through
// the imperative handle rather than the reducer.
useEffect(() => {
const unsubscribe = transport.onMessage(elementId, (action) => {
const smoother = new StreamSmoother<ChunkMeta>({
onEmit: (text, meta, isFirstSlice) => {
dispatch({
type: "chunk",
content: text,
operation: "append",
content_type: meta.content_type,
html_deps: isFirstSlice ? meta.html_deps : undefined,
})
},
// Merge consecutive chunks so pacing cuts land anywhere, not only at
// server chunk boundaries. A chunk carrying html_deps starts its own
// entry so its deps are forwarded exactly once.
canMerge: (queued, incoming) =>
queued.content_type === incoming.content_type &&
!incoming.html_deps?.length,
})
smootherRef.current = smoother

// Actions that arrive while the tail of a finished stream is still
// draining. Most (history_update, update_siblings, ...) must apply after
// chunk_end lands, so they wait and replay in order rather than cutting
// the drain short.
let deferred: ChatAction[] = []

const handleAction = (action: ChatAction) => {
if (action.type === "history_navigate") {
setCurrentConversationId(elementId, action.active_id)
navigateTo(action.url, action.reload === true)
Expand Down Expand Up @@ -229,11 +283,82 @@ export function ChatApp({
}
return
}
if (action.type === "chunk_start") {
// A prior stream's tail was already completed on arrival (see
// ENDS_DRAIN); dispose defensively since a fresh stream has nothing
// worth keeping.
smoother.dispose()
dispatch(action)
return
}
if (action.type === "chunk") {
if (action.operation === "replace") {
// A replace chunk wipes the whole in-flight message, so anything
// still buffered from before it would just get wiped a moment
// later — discard rather than flush.
smoother.dispose()
dispatch(action)
} else {
smoother.push(action.content, {
content_type: action.content_type,
html_deps: action.html_deps,
})
}
return
}
if (action.type === "block_insert") {
// Preserve order: a block must render after any text pushed before
// it, even if that text hasn't paced out yet.
smoother.flush()
dispatch(action)
return
}
if (action.type === "chunk_end") {
// Reveal the buffered tail quickly, then end the stream.
smoother.finish(() => {
dispatch(action)
const pending = deferred
deferred = []
pending.forEach(handleAction)
})
return
}
dispatch(action)
}

const unsubscribe = transport.onMessage(elementId, (action) => {
if (smoother.finishing) {
if (!ENDS_DRAIN.has(action.type)) {
deferred.push(action)
return
}
// Runs the deferred chunk_end and replays queued actions first.
smoother.flush()
}
handleAction(action)
})
return unsubscribe
return () => {
smoother.dispose()
smootherRef.current = null
deferred = []
unsubscribe()
}
}, [transport, elementId, historyStore])

// Stopping should stop the motion: reveal whatever is buffered at once. If
// the stream already ended and only its tail was draining, there is nothing
// left to cancel — completing the drain is the whole effect, and recording
// CANCEL_REQUESTED would leave a stale flag for the next stream.
const chatDispatch = useCallback((action: AnyAction) => {
if (action.type === "CANCEL_REQUESTED") {
const smoother = smootherRef.current
const wasFinishing = smoother?.finishing ?? false
smoother?.flush()
if (wasFinishing) return
}
dispatch(action)
}, [])

// State-driven `<inputId>_greeting_requested` input.
//
// Fires when all three conditions hold: the chat container is visible
Expand Down Expand Up @@ -364,7 +489,7 @@ export function ChatApp({
<ShinyLifecycleContext.Provider value={shinyLifecycle}>
<ChatToolContext.Provider value={toolState}>
<ToolGroupingContext.Provider value={state.toolGrouping}>
<ChatDispatchContext.Provider value={dispatch}>
<ChatDispatchContext.Provider value={chatDispatch}>
<AsideFaviconContext.Provider value={asideFavicon}>
<ChatContainer
ref={containerRef}
Expand Down
32 changes: 27 additions & 5 deletions js/src/markdown-stream/MarkdownStream.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
useMemo,
} from "react"
import { MarkdownContent } from "../markdown/MarkdownContent"
import { StreamingDot } from "../markdown/StreamingDotIcon"
import { useAutoScroll, findScrollableParent } from "../markdown/useAutoScroll"
import { HtmlBlockContent } from "../chat/HtmlBlockContent"
import type { HtmlBlock } from "../chat/html-block-model"
Expand Down Expand Up @@ -110,10 +111,11 @@ export function MarkdownStream({
() => [segments, blockMounts],
[segments, blockMounts],
)
const { containerRef, scrollToBottom, repinIfAtBottom } = useAutoScroll({
streaming: autoScroll && streaming,
contentDependency: scrollContentDependency,
})
const { containerRef, scrollToBottom, repinIfAtBottom, stickToBottom } =
useAutoScroll({
streaming: autoScroll && streaming,
contentDependency: scrollContentDependency,
})

useLayoutEffect(() => {
if (!autoScroll || !innerRef.current) {
Expand Down Expand Up @@ -143,8 +145,18 @@ export function MarkdownStream({
}
}, [containerRef])

// Scroll while streaming, and once more when it ends: StreamSmoother flushes
// its remaining buffer in the same batch as setStreaming(false), which the
// content-change effect (streaming-gated) never scrolls for. stickToBottom
// is read via ref so a mid-stream scroll-away doesn't re-trigger the effect.
const stickToBottomRef = useRef(stickToBottom)
stickToBottomRef.current = stickToBottom
const wasStreamingRef = useRef(streaming)
useEffect(() => {
if (streaming && autoScroll) {
const wasStreaming = wasStreamingRef.current
wasStreamingRef.current = streaming
if (!autoScroll) return
if (streaming || (wasStreaming && stickToBottomRef.current)) {
scrollToBottom()
}
}, [streaming, autoScroll, scrollToBottom])
Expand Down Expand Up @@ -217,8 +229,18 @@ export function MarkdownStream({
onApiReady?.(api)
}, [api, onApiReady])

// With nothing to show yet, the dot is the only sign the stream is live.
// Once content exists, MarkdownContent shows it only when the stream stalls.
const awaitingContent =
streaming && segments.every((s) => !isBlockSegment(s) && !/\S/.test(s.text))

return (
<div ref={innerRef}>
{awaitingContent && (
<p>
<StreamingDot />
</p>
)}
{segments.map((segment, index) =>
isBlockSegment(segment) ? (
segment.type === "web_activity" ? (
Expand Down
59 changes: 53 additions & 6 deletions js/src/markdown-stream/markdown-stream-entry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import {
import { getShinyTransport } from "../transport/shiny-transport"
import { parseJsonArray } from "../utils/json"
import { DeferredTeardown } from "../utils/deferredTeardown"
import { StreamSmoother } from "../streaming/StreamSmoother"
import type { ContentType, StructuredBlock } from "../transport/types"
import type { HtmlDep } from "rstudio-shiny/srcts/types/src/shiny/render"

Expand Down Expand Up @@ -52,6 +53,10 @@ class MarkdownStreamElement extends HTMLElement {
private api: MarkdownStreamApi | null = null
private pendingMessages: (ContentMessage | IsStreamingMessage)[] = []
private deferredTeardown = new DeferredTeardown()
private smoother: StreamSmoother<{
trusted: boolean
segmentStart: boolean
}> | null = null

connectedCallback() {
this.deferredTeardown.cancel()
Expand Down Expand Up @@ -97,9 +102,34 @@ class MarkdownStreamElement extends HTMLElement {
this.reactRoot = null
this.api = null
this.pendingMessages = []
this.smoother?.dispose()
this.smoother = null
})
}

private getSmoother(): StreamSmoother<{
trusted: boolean
segmentStart: boolean
}> {
if (!this.smoother) {
this.smoother = new StreamSmoother({
onEmit: (text, meta, isFirstSlice) => {
this.api!.appendContent(
text,
meta.trusted,
meta.segmentStart && isFirstSlice,
)
},
// Merge consecutive chunks so pacing cuts land anywhere, not only at
// server chunk boundaries. A segment start keeps its own entry so
// the seam is preserved.
canMerge: (queued, incoming) =>
queued.trusted === incoming.trusted && !incoming.segmentStart,
})
}
return this.smoother
}

handleMessage(message: ContentMessage | IsStreamingMessage) {
if (!this.api) {
this.pendingMessages.push(message)
Expand All @@ -109,31 +139,48 @@ class MarkdownStreamElement extends HTMLElement {
}

private dispatchMessage(message: ContentMessage | IsStreamingMessage) {
// Anything arriving while a finished stream's tail drains supersedes the
// drain: reveal the rest (and apply the deferred setStreaming(false)).
if (this.smoother?.finishing) this.smoother.flush()

if (isStreamingMessage(message)) {
this.api!.setStreaming(message.isStreaming)
if (message.isStreaming === false) {
// Reveal the buffered tail quickly, then end the stream.
const api = this.api!
if (this.smoother) {
this.smoother.finish(() => api.setStreaming(false))
} else {
api.setStreaming(false)
}
} else if (message.isStreaming === true) {
this.smoother?.dispose()
this.api!.setStreaming(true)
}
return
}

if (message.block !== undefined) {
const block = asStreamBlock(message.block)
if (!block) return
if (message.operation === "replace") {
this.smoother?.dispose()
this.api!.replaceWithBlock(block)
} else {
this.smoother?.flush()
this.api!.appendBlock(block)
}
return
}

const content = message.content ?? ""
if (message.operation === "replace") {
this.smoother?.dispose()
this.api!.replaceContent(content, message.trusted === true)
} else if (message.operation === "append") {
this.api!.appendContent(
content,
message.trusted === true,
message.segment_start === true,
)
this.getSmoother().push(content, {
trusted: message.trusted === true,
segmentStart: message.segment_start === true,
})
}
}
}
Expand Down
22 changes: 18 additions & 4 deletions js/src/markdown-stream/markdown-stream.scss
Original file line number Diff line number Diff line change
Expand Up @@ -73,11 +73,25 @@ pre:has(.code-copy-button) {
}
}

@keyframes markdown-stream-dot-appear {
Comment thread
cpsievert marked this conversation as resolved.
to {
opacity: 1;
}
}

.markdown-stream-dot {
// The stream dot is appended with each streaming chunk update, so the pulse animation
// only shows up when streaming pauses but isn't complete.
animation: markdown-stream-dot-pulse 1.75s infinite cubic-bezier(0.18, 0.89, 0.32, 1.28);
animation-delay: 250ms;
// Only rendered while a stream has no content yet or has stalled (see
// useStreamIdle), so it mounts fresh each time: fade in, then pulse.
opacity: 0;
animation:
markdown-stream-dot-appear 0.3s ease forwards,
markdown-stream-dot-pulse 1.75s infinite cubic-bezier(0.18, 0.89, 0.32, 1.28) 0.3s;
display: inline-block;
transform-origin: center;
}

@media (prefers-reduced-motion: reduce) {
.markdown-stream-dot {
animation: markdown-stream-dot-appear 0.3s ease forwards;
}
}
Loading
Loading