Skip to content

Commit df43c8a

Browse files
committed
fix(mothership): create the log re-sync set outside render
1 parent 6a57d07 commit df43c8a

1 file changed

Lines changed: 10 additions & 4 deletions

File tree

  • apps/sim/app/workspace/[workspaceId]/home/hooks

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts‎

Lines changed: 10 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import {
22
type Dispatch,
3+
type RefObject,
34
type SetStateAction,
45
useCallback,
56
useEffect,
@@ -311,6 +312,12 @@ const logger = createLogger('useChat')
311312
* its cursors are log positions, so every later read names the log as its source and
312313
* is never served from the replay ring, even one that restarted and grew past them.
313314
*/
315+
/** The streams re-synced from the log, created on first use outside render. */
316+
function logResyncedStreams(ref: RefObject<Set<string> | null>): Set<string> {
317+
ref.current ??= new Set()
318+
return ref.current
319+
}
320+
314321
function streamReconnectQuery(streamId: string, afterCursor: string, fromLog: boolean): string {
315322
return `streamId=${encodeURIComponent(streamId)}&after=${encodeURIComponent(afterCursor)}${fromLog ? '&source=log' : ''}`
316323
}
@@ -966,7 +973,6 @@ export function useChat(
966973
const locallyTerminalStreamIdRef = useRef<string | undefined>(undefined)
967974
const lastCursorRef = useRef('0')
968975
const logResyncedStreamsRef = useRef<Set<string> | null>(null)
969-
logResyncedStreamsRef.current ??= new Set()
970976
const activeStreamReturnRecoveryRef = useRef<ActiveStreamRecovery | null>(null)
971977
const sendingRef = useRef(false)
972978
const streamGenRef = useRef(0)
@@ -2382,7 +2388,7 @@ export function useChat(
23822388
)
23832389
// boundary-raw-fetch: stream-resume batch endpoint requires dynamic per-request traceparent header propagation that the contract layer does not model, and the response is consumed alongside live SSE tail fetches
23842390
const response = await fetch(
2385-
`/api/mothership/chat/stream?${streamReconnectQuery(streamId, afterCursor, logResyncedStreamsRef.current?.has(streamId) ?? false)}&batch=true`,
2391+
`/api/mothership/chat/stream?${streamReconnectQuery(streamId, afterCursor, logResyncedStreams(logResyncedStreamsRef).has(streamId))}&batch=true`,
23862392
{
23872393
signal: fetchSignal,
23882394
...(streamTraceparentRef.current
@@ -2574,7 +2580,7 @@ export function useChat(
25742580

25752581
// boundary-raw-fetch: live SSE tail endpoint streams events consumed via response.body.getReader() and processSSEStream
25762582
const sseRes = await fetch(
2577-
`/api/mothership/chat/stream?${streamReconnectQuery(streamId, latestCursor, logResyncedStreamsRef.current?.has(streamId) ?? false)}`,
2583+
`/api/mothership/chat/stream?${streamReconnectQuery(streamId, latestCursor, logResyncedStreams(logResyncedStreamsRef).has(streamId))}`,
25782584
{
25792585
signal: activeAbort.signal,
25802586
...(streamTraceparentRef.current
@@ -2595,7 +2601,7 @@ export function useChat(
25952601

25962602
// Re-sent from the worker's log with cursors restarting at 1: rebuild from empty.
25972603
if (sseRes.headers.get(MOTHERSHIP_STREAM_REPLAY_HEADER) === 'log') {
2598-
logResyncedStreamsRef.current?.add(streamId)
2604+
logResyncedStreams(logResyncedStreamsRef).add(streamId)
25992605
const reset = applyReconnectReplaySelection(streamId, '0')
26002606
latestCursor = reset.afterCursor
26012607
preserveNextReplayState = reset.preserveExistingState

0 commit comments

Comments
 (0)