Skip to content

Commit 4cfe692

Browse files
refactor(sse): batch parts directory-aware and serialize backend writes
1 parent 816c419 commit 4cfe692

5 files changed

Lines changed: 233 additions & 70 deletions

File tree

frontend/src/components/message/MessageThread.tsx

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -135,8 +135,7 @@ const findLastMessageByRole = (
135135

136136
interface MessageRowProps {
137137
msgWithParts: MessageWithParts
138-
index: number
139-
messages: MessageWithParts[]
138+
nextAssistantMessage: MessageWithParts | undefined
140139
pendingAssistantId: string | undefined
141140
lastUserMessageId: string | undefined
142141
isSessionBusy: boolean
@@ -157,8 +156,7 @@ interface MessageRowProps {
157156

158157
const MessageRow = memo(function MessageRow({
159158
msgWithParts,
160-
index,
161-
messages,
159+
nextAssistantMessage,
162160
pendingAssistantId,
163161
lastUserMessageId,
164162
isSessionBusy,
@@ -183,7 +181,6 @@ const MessageRow = memo(function MessageRow({
183181
const isLastUserMessage = msg.role === 'user' && msg.id === lastUserMessageId
184182
const messageTextContent = getMessageTextContent(parts)
185183

186-
const nextAssistantMessage = messages.slice(index + 1).find(m => m.info.role === 'assistant')
187184
const nextAssistantMsg = nextAssistantMessage?.info
188185
const isUserBeforeAssistant = msg.role === 'user' && nextAssistantMessage
189186
const canEditUserMessage = isLastUserMessage && isUserBeforeAssistant && !isSessionBusy
@@ -352,6 +349,20 @@ export const MessageThread = memo(function MessageThread({
352349
return findLastMessageByRole(messages, 'user')
353350
}, [messages])
354351

352+
const nextAssistantByMessageId = useMemo(() => {
353+
const map = new Map<string, MessageWithParts | undefined>()
354+
if (!messages) return map
355+
let next: MessageWithParts | undefined
356+
for (let i = messages.length - 1; i >= 0; i--) {
357+
const msg = messages[i]
358+
map.set(msg.info.id, next)
359+
if (msg.info.role === 'assistant') {
360+
next = msg
361+
}
362+
}
363+
return map
364+
}, [messages])
365+
355366
const isSessionBusy = !!pendingAssistantId || isSessionInRetry(sessionStatus)
356367
const setSessionTodos = useSessionTodos((state) => state.setTodos)
357368

@@ -415,12 +426,11 @@ export const MessageThread = memo(function MessageThread({
415426

416427
return (
417428
<div className="flex flex-col space-y-2 p-2 overflow-x-hidden">
418-
{messages.map((msgWithParts, index) => (
429+
{messages.map((msgWithParts) => (
419430
<MessageRow
420431
key={msgWithParts.info.id}
421432
msgWithParts={msgWithParts}
422-
index={index}
423-
messages={messages}
433+
nextAssistantMessage={nextAssistantByMessageId.get(msgWithParts.info.id)}
424434
pendingAssistantId={pendingAssistantId}
425435
lastUserMessageId={lastUserMessageId}
426436
isSessionBusy={isSessionBusy}

frontend/src/hooks/useAutoScroll.test.tsx

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -131,6 +131,9 @@ describe('useAutoScroll', () => {
131131
contentVersion: newMessages.length,
132132
onScrollStateChange,
133133
})
134+
})
135+
136+
act(() => {
134137
vi.advanceTimersByTime(100)
135138
})
136139

@@ -162,6 +165,7 @@ describe('useAutoScroll', () => {
162165
contentVersion: messages.length + 1,
163166
onScrollStateChange,
164167
})
168+
vi.advanceTimersByTime(100)
165169
})
166170

167171
expect(containerHarness.getScrollTop()).toBe(containerHarness.div.scrollHeight - containerHarness.div.clientHeight)
@@ -231,6 +235,29 @@ describe('useAutoScroll', () => {
231235
expect(containerHarness.getScrollTop()).toBe(userPosition)
232236
})
233237

238+
it('cancels pending bottom scroll when user wheel-scrolls up before pending frames complete', () => {
239+
const messages = [createMessage('1', 'user'), createMessage('2', 'assistant')]
240+
const { renderResult, containerHarness } = setupHook(messages)
241+
242+
act(() => {
243+
renderResult.result.current.scrollToBottom()
244+
})
245+
246+
const userPosition = 150
247+
act(() => {
248+
containerHarness.div.dispatchEvent(
249+
new WheelEvent('wheel', {
250+
deltaY: -50,
251+
bubbles: true,
252+
})
253+
)
254+
containerHarness.setScrollTop(userPosition)
255+
vi.runOnlyPendingTimers()
256+
})
257+
258+
expect(containerHarness.getScrollTop()).toBe(userPosition)
259+
})
260+
234261
it('does not show scroll button on tiny upward drag from bottom', () => {
235262
const messages = [createMessage('1', 'user')]
236263
const { containerHarness, onScrollStateChange } = setupHook(messages)
@@ -431,6 +458,10 @@ describe('useAutoScroll', () => {
431458
})
432459
})
433460

461+
act(() => {
462+
vi.advanceTimersByTime(100)
463+
})
464+
434465
expect(containerHarness.getScrollTop()).toBe(containerHarness.div.scrollHeight - containerHarness.div.clientHeight)
435466
})
436467

@@ -459,6 +490,10 @@ describe('useAutoScroll', () => {
459490
})
460491
})
461492

493+
act(() => {
494+
vi.advanceTimersByTime(100)
495+
})
496+
462497
expect(containerHarness.getScrollTop()).toBe(containerHarness.div.scrollHeight - containerHarness.div.clientHeight)
463498
})
464499
})

frontend/src/hooks/useAutoScroll.ts

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -46,17 +46,13 @@ export function useAutoScroll({
4646
isScrollButtonVisibleRef.current = false
4747
const scrollRequestId = scrollRequestIdRef.current + 1
4848
scrollRequestIdRef.current = scrollRequestId
49-
const scroll = () => {
50-
if (!containerRef?.current) return
51-
containerRef.current.scrollTop = containerRef.current.scrollHeight
52-
}
5349

54-
scroll()
5550
let frameCount = 0
5651
const scrollAfterLayout = () => {
5752
if (scrollRequestIdRef.current !== scrollRequestId) return
53+
if (!containerRef?.current) return
54+
containerRef.current.scrollTop = containerRef.current.scrollHeight
5855
frameCount += 1
59-
scroll()
6056
if (frameCount < SCROLL_TO_BOTTOM_FRAME_COUNT) {
6157
requestAnimationFrame(scrollAfterLayout)
6258
}

frontend/src/lib/partsBatcher.test.ts

Lines changed: 83 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,16 @@ function textPart(sessionID: string, messageID: string, partID: string, text: st
2727
return { id: partID, sessionID, messageID, type: 'text', text } as Part
2828
}
2929

30+
function createManyCachedMessages(count: number, sessionID: string): MessageWithParts[] {
31+
const messages: MessageWithParts[] = []
32+
for (let i = 0; i < count; i++) {
33+
const msg = assistantMessage(sessionID, `msg-${i}`)
34+
msg.parts = [textPart(sessionID, `msg-${i}`, `part-${i}`, `base text ${i}`)]
35+
messages.push(msg)
36+
}
37+
return messages
38+
}
39+
3040
describe('createPartsBatcher', () => {
3141
it('invalidates when part deltas arrive before message cache exists and applies a later authoritative upsert', () => {
3242
const queryClient = new QueryClient()
@@ -138,6 +148,79 @@ describe('createPartsBatcher', () => {
138148
expect(data![0].parts[0]).toHaveProperty('text', 'snapshot later')
139149
})
140150

151+
it('applies many queued part deltas with one cache write and no invalidation storm', () => {
152+
const queryClient = new QueryClient()
153+
const batcher = createPartsBatcher(queryClient, 'http://localhost:5551')
154+
155+
const sessionID = 'session-1'
156+
const directory = '/repo'
157+
const messageCount = 1000
158+
const deltaCount = 500
159+
160+
const messages = createManyCachedMessages(messageCount, sessionID)
161+
queryClient.setQueryData(
162+
['opencode', 'messages', 'http://localhost:5551', sessionID, directory],
163+
messages,
164+
)
165+
166+
const setQueryDataSpy = vi.spyOn(queryClient, 'setQueryData')
167+
const invalidateSpy = vi.spyOn(queryClient, 'invalidateQueries')
168+
169+
for (let i = 0; i < deltaCount; i++) {
170+
batcher.queuePartDelta(sessionID, `msg-${i}`, `part-${i}`, 'text', ` delta ${i}`, directory)
171+
}
172+
173+
batcher.flush()
174+
175+
expect(invalidateSpy).not.toHaveBeenCalled()
176+
177+
const data = queryClient.getQueryData<MessageWithParts[]>([
178+
'opencode', 'messages', 'http://localhost:5551', sessionID, directory,
179+
])
180+
expect(data).toHaveLength(messageCount)
181+
182+
for (let i = 0; i < deltaCount; i++) {
183+
expect(data![i].parts[0]).toHaveProperty('text', `base text ${i} delta ${i}`)
184+
}
185+
186+
for (let i = deltaCount; i < messageCount; i++) {
187+
expect(data![i].parts[0]).toHaveProperty('text', `base text ${i}`)
188+
}
189+
190+
const setQueryDataCalls = setQueryDataSpy.mock.calls.filter(
191+
([key]) => JSON.stringify(key).includes('opencode'),
192+
)
193+
expect(setQueryDataCalls.length).toBe(1)
194+
})
195+
196+
it('does not apply same-batch deltas for removed parts to shifted parts', () => {
197+
const queryClient = new QueryClient()
198+
const batcher = createPartsBatcher(queryClient, 'http://localhost:5551')
199+
200+
queryClient.setQueryData(
201+
['opencode', 'messages', 'http://localhost:5551', 'session-1', '/repo'],
202+
[{
203+
...assistantMessage('session-1', 'message-1'),
204+
parts: [
205+
textPart('session-1', 'message-1', 'part-1', 'first'),
206+
textPart('session-1', 'message-1', 'part-2', 'second'),
207+
],
208+
}],
209+
)
210+
211+
batcher.queuePartRemoval('session-1', 'message-1', 'part-1', '/repo')
212+
batcher.queuePartDelta('session-1', 'message-1', 'part-1', 'text', ' stale', '/repo')
213+
batcher.flush()
214+
215+
const data = queryClient.getQueryData<MessageWithParts[]>([
216+
'opencode', 'messages', 'http://localhost:5551', 'session-1', '/repo',
217+
])
218+
219+
expect(data![0].parts).toHaveLength(1)
220+
expect(data![0].parts[0]).toHaveProperty('id', 'part-2')
221+
expect(data![0].parts[0]).toHaveProperty('text', 'second')
222+
})
223+
141224
it('applies deltas to the directory they were queued for', () => {
142225
const queryClient = new QueryClient()
143226
const batcher = createPartsBatcher(queryClient, 'http://localhost:5551')

0 commit comments

Comments
 (0)