Skip to content

Commit 4af9b7a

Browse files
committed
fix(chat): relay left turns from any cursor and until the run ends
1 parent f9df0e8 commit 4af9b7a

4 files changed

Lines changed: 156 additions & 97 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/hooks/stream-protocol.ts‎

Lines changed: 17 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -115,18 +115,27 @@ export function parseStreamBatchResponse(value: unknown): StreamBatchResponse {
115115
}
116116
}
117117

118+
/** The chat an event names, from its stream metadata or a chat session event. */
119+
export function resolveChatIdFromStreamEvent(
120+
event: PersistedStreamEventEnvelope
121+
): string | undefined {
122+
const streamChatId = typeof event.stream?.chatId === 'string' ? event.stream.chatId : undefined
123+
if (streamChatId) return streamChatId
124+
if (
125+
event.type === MothershipStreamV1EventType.session &&
126+
event.payload.kind === MothershipStreamV1SessionKind.chat
127+
) {
128+
return event.payload.chatId
129+
}
130+
return undefined
131+
}
132+
118133
export function resolveChatIdFromStreamBatch(batch: StreamBatchResponse): string | undefined {
119134
if (batch.chatId) return batch.chatId
120135

121136
for (const { event } of batch.events) {
122-
const streamChatId = typeof event.stream?.chatId === 'string' ? event.stream.chatId : undefined
123-
if (streamChatId) return streamChatId
124-
if (
125-
event.type === MothershipStreamV1EventType.session &&
126-
event.payload.kind === MothershipStreamV1SessionKind.chat
127-
) {
128-
return event.payload.chatId
129-
}
137+
const chatId = resolveChatIdFromStreamEvent(event)
138+
if (chatId) return chatId
130139
}
131140

132141
return undefined

‎apps/sim/app/workspace/[workspaceId]/home/hooks/stream/detached-client-tools.ts‎

Lines changed: 121 additions & 72 deletions
Original file line numberDiff line numberDiff line change
@@ -3,22 +3,27 @@
33
* mid-turn. Leaving detaches the chat view but not the server run, and that
44
* run's browser, terminal, and local filesystem tools execute only in this
55
* client: the orchestrator blocks until the client reports each outcome, so
6-
* without a reader the left chat stalls at its next such call. A relay tails
7-
* the left chat's stream headlessly and starts those tools until the run
8-
* completes or the chat is reopened. Relays live at module scope because
9-
* switching chats remounts the chat surface that detached them. Each holds one
10-
* resume connection, so a relay exists only while its left run is live.
6+
* without a reader the left chat stalls at its next such call. A relay reads
7+
* the left turn's stream headlessly and starts those tools until the run ends
8+
* or the chat's view reads the stream again. Relays live at module scope
9+
* because switching chats remounts the chat surface that detached them. Each
10+
* holds one resume connection, and only while its run is live.
1111
*/
1212
import { createLogger } from '@sim/logger'
1313
import { getErrorMessage } from '@sim/utils/errors'
1414
import { interruptibleSleep } from '@sim/utils/helpers'
1515
import { backoffWithJitter } from '@sim/utils/retry'
1616
import { readSSELines } from '@/lib/core/utils/sse'
17+
import { desktopChatScopeId } from '@/lib/desktop/chat-scope'
1718
import {
1819
MothershipStreamV1EventType,
1920
MothershipStreamV1ToolPhase,
2021
} from '@/lib/mothership/generated/mothership-stream-v1'
21-
import { parsePersistedStreamEventEnvelopeJson } from '@/lib/mothership/request/session/contract'
22+
import {
23+
isTerminalStreamStatus,
24+
type PersistedStreamEventEnvelope,
25+
parsePersistedStreamEventEnvelopeJson,
26+
} from '@/lib/mothership/request/session/contract'
2227
import { executeBrowserToolOnClient } from '@/lib/mothership/tools/client/browser-tool-execution'
2328
import { launchLocalFilesystemTool } from '@/lib/mothership/tools/client/launch-local-filesystem-tool'
2429
import { executeTerminalToolOnClient } from '@/lib/mothership/tools/client/terminal-tool-execution'
@@ -32,98 +37,150 @@ import {
3237
getStreamEventCursor,
3338
isAlreadyProcessedStreamCursor,
3439
isStreamSchemaValidationError,
40+
parseStreamBatchResponse,
41+
resolveChatIdFromStreamBatch,
42+
resolveChatIdFromStreamEvent,
3543
STREAM_IDLE_TIMEOUT_MS,
3644
} from '@/app/workspace/[workspaceId]/home/hooks/stream-protocol'
3745

3846
const logger = createLogger('DetachedClientTools')
3947

40-
/** Consecutive tail attempts that make no progress before a relay gives up. */
41-
const MAX_STALLED_TAIL_ATTEMPTS = 5
42-
43-
/** A live turn the user left, and where its client tools run. */
48+
/** A live turn the user left. */
4449
export interface DetachedChatTurn {
45-
chatId: string
4650
streamId: string
51+
/** The turn's chat when the view knew it; a new chat's id is read off the stream. */
52+
chatId?: string
4753
/** The last cursor the chat view dispatched; the relay resumes after it. */
4854
afterCursor: string
4955
traceparent?: string
5056
workspaceId?: string
51-
/** The chat's desktop scope, which owns its browser and terminal tabs. */
52-
scopeId: string
57+
/** The owner key the chat's desktop scope derives from. */
58+
scopeKey: string
5359
}
5460

55-
/** Relays by chat id. An entry is removed when its relay ends. */
56-
const relays = new Map<string, AbortController>()
61+
interface Relay {
62+
controller: AbortController
63+
chatId?: string
64+
}
65+
66+
/** Relays by stream id. An entry is removed when its relay ends. */
67+
const relays = new Map<string, Relay>()
68+
69+
/** The call a tool result frame settles, or undefined for any other event. */
70+
function settledToolCallId(event: PersistedStreamEventEnvelope): string | undefined {
71+
if (event.type !== MothershipStreamV1EventType.tool || 'previewPhase' in event.payload) {
72+
return undefined
73+
}
74+
return event.payload.phase === MothershipStreamV1ToolPhase.result
75+
? event.payload.toolCallId
76+
: undefined
77+
}
5778

5879
/**
5980
* Starts the client tools a detached turn hands this client. Workflow runs are
6081
* left alone: running one drives the workflow editor, and the server runs a
6182
* workflow call itself when no client picks it up.
6283
*/
63-
function startDetachedClientTool(turn: DetachedChatTurn, start: ClientToolStart): void {
84+
function startDetachedClientTool(
85+
turn: DetachedChatTurn,
86+
chatId: string,
87+
start: ClientToolStart
88+
): void {
6489
const { toolCallId, toolName, args, eventTs } = start
90+
const scopeId = desktopChatScopeId(turn.scopeKey, chatId)
6591
switch (start.kind) {
6692
case 'workflow':
6793
return
6894
case 'localFilesystem':
6995
launchLocalFilesystemTool(toolCallId, toolName, args, {
7096
workspaceId: turn.workspaceId,
71-
chatId: turn.chatId,
97+
chatId,
7298
})
7399
return
74100
case 'browser':
75-
executeBrowserToolOnClient(toolCallId, start.toolName, args, turn.scopeId, eventTs)
101+
executeBrowserToolOnClient(toolCallId, start.toolName, args, scopeId, eventTs)
76102
return
77103
case 'terminal':
78-
executeTerminalToolOnClient(toolCallId, args, turn.scopeId, eventTs)
104+
executeTerminalToolOnClient(toolCallId, args, scopeId, eventTs)
79105
return
80106
}
81107
}
82108

83-
async function relayClientTools(turn: DetachedChatTurn, signal: AbortSignal): Promise<void> {
84-
const { streamId, traceparent } = turn
109+
/**
110+
* Relays until the run ends or the relay is aborted. Like the chat view's own
111+
* reconnect, every connection first reads the events past the cursor as one
112+
* batch, so calls that already have a result are settled before any call frame
113+
* replays, then tails live. That makes any cursor a safe starting point.
114+
*/
115+
async function relayClientTools(turn: DetachedChatTurn, relay: Relay): Promise<void> {
116+
const { streamId } = turn
117+
const { signal } = relay.controller
118+
const headers = turn.traceparent ? { traceparent: turn.traceparent } : undefined
85119
/** Calls this relay started or saw settle; ids alone decide what is pending. */
86120
const handledToolCallIds = new Set<string>()
87121
let cursor = turn.afterCursor
88-
let stalledAttempts = 0
122+
let failedAttempts = 0
89123

90-
/** Reads one tail connection; resolves true once the run is over. */
91-
const readTail = async (): Promise<boolean> => {
92-
// boundary-raw-fetch: live SSE tail endpoint streams events consumed via readSSELines
93-
const response = await fetch(buildStreamResumeUrl(streamId, cursor), {
124+
const applyEvent = (event: PersistedStreamEventEnvelope): void => {
125+
const eventCursor = getStreamEventCursor(event)
126+
if (isAlreadyProcessedStreamCursor(eventCursor, cursor)) return
127+
cursor = eventCursor
128+
relay.chatId ??= resolveChatIdFromStreamEvent(event)
129+
if (event.type !== MothershipStreamV1EventType.tool) return
130+
const settledId = settledToolCallId(event)
131+
if (settledId) {
132+
handledToolCallIds.add(settledId)
133+
return
134+
}
135+
const start = resolveClientToolStart(event, (id) => !handledToolCallIds.has(id))
136+
if (!start) return
137+
handledToolCallIds.add(start.toolCallId)
138+
if (!relay.chatId) {
139+
logger.error('Detached client tool arrived before its chat id', {
140+
streamId,
141+
toolCallId: start.toolCallId,
142+
})
143+
return
144+
}
145+
startDetachedClientTool(turn, relay.chatId, start)
146+
}
147+
148+
/** Reads the events past the cursor at once; resolves true once the run is over. */
149+
const readBatch = async (): Promise<boolean> => {
150+
// boundary-raw-fetch: stream-resume batch endpoint needs per-request traceparent propagation the contract layer does not model
151+
const response = await fetch(buildStreamResumeUrl(streamId, cursor, { batch: true }), {
94152
signal,
95-
...(traceparent ? { headers: { traceparent } } : {}),
153+
headers,
96154
})
97155
if (response.status === 404) return true
156+
if (!response.ok) throw new Error(`Stream batch responded with status ${response.status}`)
157+
const batch = parseStreamBatchResponse(await response.json())
158+
relay.chatId ??= resolveChatIdFromStreamBatch(batch)
159+
for (const { event } of batch.events) {
160+
const settledId = settledToolCallId(event)
161+
if (settledId) handledToolCallIds.add(settledId)
162+
}
163+
for (const { event } of batch.events) applyEvent(event)
164+
return isTerminalStreamStatus(batch.status)
165+
}
166+
167+
/** Tails one live connection; resolves true once the run is over. */
168+
const readTail = async (): Promise<boolean> => {
169+
// boundary-raw-fetch: live SSE tail endpoint streams events consumed via readSSELines
170+
const response = await fetch(buildStreamResumeUrl(streamId, cursor), { signal, headers })
171+
if (response.status === 404) return true
98172
if (!response.ok || !response.body) {
99173
throw new Error(`Stream tail responded with status ${response.status}`)
100174
}
101-
102175
let complete = false
103176
await readSSELines(response.body, {
104177
signal,
105178
idleTimeoutMs: STREAM_IDLE_TIMEOUT_MS,
106179
onData: (raw) => {
107180
const parsed = parsePersistedStreamEventEnvelopeJson(raw)
108181
if (!parsed.ok) throw createStreamSchemaValidationError(parsed, 'Detached SSE event.')
109-
const event = parsed.event
110-
const eventCursor = getStreamEventCursor(event)
111-
if (isAlreadyProcessedStreamCursor(eventCursor, cursor)) return
112-
cursor = eventCursor
113-
114-
if (event.type === MothershipStreamV1EventType.tool) {
115-
const start = resolveClientToolStart(event, (id) => !handledToolCallIds.has(id))
116-
if (start) {
117-
handledToolCallIds.add(start.toolCallId)
118-
startDetachedClientTool(turn, start)
119-
} else if (
120-
!('previewPhase' in event.payload) &&
121-
event.payload.phase === MothershipStreamV1ToolPhase.result
122-
) {
123-
handledToolCallIds.add(event.payload.toolCallId)
124-
}
125-
}
126-
if (event.type === MothershipStreamV1EventType.complete) {
182+
applyEvent(parsed.event)
183+
if (parsed.event.type === MothershipStreamV1EventType.complete) {
127184
complete = true
128185
return true
129186
}
@@ -135,7 +192,7 @@ async function relayClientTools(turn: DetachedChatTurn, signal: AbortSignal): Pr
135192
while (!signal.aborted) {
136193
const cursorBeforeAttempt = cursor
137194
try {
138-
if (await readTail()) return
195+
if ((await readBatch()) || (await readTail())) return
139196
} catch (error) {
140197
if (signal.aborted) return
141198
if (isStreamSchemaValidationError(error)) {
@@ -145,43 +202,35 @@ async function relayClientTools(turn: DetachedChatTurn, signal: AbortSignal): Pr
145202
})
146203
return
147204
}
148-
logger.warn('Detached stream tail failed', { streamId, error: getErrorMessage(error) })
205+
logger.warn('Detached stream read failed', { streamId, error: getErrorMessage(error) })
149206
}
150-
151207
if (cursor !== cursorBeforeAttempt) {
152-
stalledAttempts = 0
208+
failedAttempts = 0
153209
continue
154210
}
155-
stalledAttempts++
156-
if (stalledAttempts >= MAX_STALLED_TAIL_ATTEMPTS) {
157-
logger.warn('Stopped relaying detached client tools after repeated stalls', {
158-
streamId,
159-
cursor,
160-
})
161-
return
162-
}
163-
await interruptibleSleep(backoffWithJitter(stalledAttempts, null), signal)
211+
failedAttempts++
212+
await interruptibleSleep(backoffWithJitter(failedAttempts, null), signal)
164213
}
165214
}
166215

167-
/**
168-
* Relays a left turn's client tools until its run completes or the chat is
169-
* reopened. Replaces any relay the chat already has.
170-
*/
216+
/** Relays a left turn's client tools until its run ends or its chat's view reads it again. */
171217
export function detachClientTools(turn: DetachedChatTurn): void {
172-
relays.get(turn.chatId)?.abort('superseded_detached_relay')
173-
const controller = new AbortController()
174-
relays.set(turn.chatId, controller)
175-
void relayClientTools(turn, controller.signal).finally(() => {
176-
if (relays.get(turn.chatId) === controller) relays.delete(turn.chatId)
218+
relays.get(turn.streamId)?.controller.abort('superseded_detached_relay')
219+
const relay: Relay = { controller: new AbortController(), chatId: turn.chatId }
220+
relays.set(turn.streamId, relay)
221+
void relayClientTools(turn, relay).finally(() => {
222+
if (relays.get(turn.streamId) === relay) relays.delete(turn.streamId)
177223
})
178224
}
179225

180226
/**
181-
* Stops relaying a chat the user reopened; its chat view takes the turn back.
182-
* Tools the relay already started keep running and report their outcome.
227+
* Stops relaying a chat whose view reads its stream again. Tools the relay
228+
* already started keep running and report their outcome.
183229
*/
184230
export function reattachClientTools(chatId: string): void {
185-
relays.get(chatId)?.abort('chat_reattached')
186-
relays.delete(chatId)
231+
for (const [streamId, relay] of relays) {
232+
if (relay.chatId !== chatId) continue
233+
relay.controller.abort('chat_reattached')
234+
relays.delete(streamId)
235+
}
187236
}

0 commit comments

Comments
 (0)