Skip to content

Commit b60c919

Browse files
committed
fix(chat): detach only admitted turns and report stale terminal calls
1 parent 4af9b7a commit b60c919

8 files changed

Lines changed: 76 additions & 64 deletions

File tree

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import type { StreamBatchEvent } from '@/lib/mothership/request/session/types'
2222

2323
/** Both live transports heartbeat every 15s; three missed heartbeats trigger cursor recovery. */
2424
export const STREAM_IDLE_TIMEOUT_MS = 45_000
25+
export const STREAM_BATCH_FETCH_TIMEOUT_MS = 10_000
2526

2627
export type StreamBatchResponse = {
2728
success: boolean

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

Lines changed: 11 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -29,17 +29,12 @@ export type ClientToolStart = ClientToolCall &
2929
)
3030

3131
/**
32-
* Resolves the client-executed tool call a stream event asks this client to
33-
* start, or null when the event starts nothing. Only a complete call frame
34-
* that is not held behind an approval prompt starts a tool.
35-
*
36-
* @param isPending - whether the call is still running without a result, as
37-
* far as the caller has seen
32+
* Resolves the client-executed tool call a stream event hands this client, or
33+
* null when the event starts nothing. Only a complete call frame that is not
34+
* held behind an approval prompt starts a tool; whether the call is still
35+
* pending is the caller's to decide from what it has seen.
3836
*/
39-
export function resolveClientToolStart(
40-
event: ToolEvent,
41-
isPending: (toolCallId: string) => boolean
42-
): ClientToolStart | null {
37+
export function resolveClientToolStart(event: ToolEvent): ClientToolStart | null {
4338
const payload = event.payload
4439
if (
4540
'previewPhase' in payload ||
@@ -55,15 +50,11 @@ export function resolveClientToolStart(
5550
const { toolCallId, toolName } = payload
5651
const args = payload.arguments as Record<string, unknown> | undefined
5752
const call = { toolCallId, args: args ?? {}, eventTs: event.ts }
58-
let start: ClientToolStart | null = null
59-
if (isWorkflowToolName(toolName)) {
60-
start = { ...call, kind: 'workflow', toolName }
61-
} else if (isNativeFileTool(toolName) || isUserLocalVfsToolCall(toolName, args)) {
62-
start = { ...call, kind: 'localFilesystem', toolName }
63-
} else if (isCurrentBrowserToolName(toolName)) {
64-
start = { ...call, kind: 'browser', toolName }
65-
} else if (isTerminalToolName(toolName)) {
66-
start = { ...call, kind: 'terminal', toolName }
53+
if (isWorkflowToolName(toolName)) return { ...call, kind: 'workflow', toolName }
54+
if (isNativeFileTool(toolName) || isUserLocalVfsToolCall(toolName, args)) {
55+
return { ...call, kind: 'localFilesystem', toolName }
6756
}
68-
return start && isPending(toolCallId) ? start : null
57+
if (isCurrentBrowserToolName(toolName)) return { ...call, kind: 'browser', toolName }
58+
if (isTerminalToolName(toolName)) return { ...call, kind: 'terminal', toolName }
59+
return null
6960
}

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

Lines changed: 22 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -40,11 +40,20 @@ import {
4040
parseStreamBatchResponse,
4141
resolveChatIdFromStreamBatch,
4242
resolveChatIdFromStreamEvent,
43+
STREAM_BATCH_FETCH_TIMEOUT_MS,
4344
STREAM_IDLE_TIMEOUT_MS,
4445
} from '@/app/workspace/[workspaceId]/home/hooks/stream-protocol'
4546

4647
const logger = createLogger('DetachedClientTools')
4748

49+
/**
50+
* Responses after which there is nothing left to relay: no run exists for the
51+
* stream (404), or this client can no longer read it (401, 403).
52+
*/
53+
function isRelayEndStatus(status: number): boolean {
54+
return status === 404 || status === 401 || status === 403
55+
}
56+
4857
/** A live turn the user left. */
4958
export interface DetachedChatTurn {
5059
streamId: string
@@ -60,6 +69,7 @@ export interface DetachedChatTurn {
6069

6170
interface Relay {
6271
controller: AbortController
72+
/** Known up front or read off the stream; tools run in this chat's desktop scope. */
6373
chatId?: string
6474
}
6575

@@ -132,8 +142,8 @@ async function relayClientTools(turn: DetachedChatTurn, relay: Relay): Promise<v
132142
handledToolCallIds.add(settledId)
133143
return
134144
}
135-
const start = resolveClientToolStart(event, (id) => !handledToolCallIds.has(id))
136-
if (!start) return
145+
const start = resolveClientToolStart(event)
146+
if (!start || handledToolCallIds.has(start.toolCallId)) return
137147
handledToolCallIds.add(start.toolCallId)
138148
if (!relay.chatId) {
139149
logger.error('Detached client tool arrived before its chat id', {
@@ -145,14 +155,14 @@ async function relayClientTools(turn: DetachedChatTurn, relay: Relay): Promise<v
145155
startDetachedClientTool(turn, relay.chatId, start)
146156
}
147157

148-
/** Reads the events past the cursor at once; resolves true once the run is over. */
158+
/** Reads the events past the cursor at once; resolves true once there is nothing left to relay. */
149159
const readBatch = async (): Promise<boolean> => {
150160
// boundary-raw-fetch: stream-resume batch endpoint needs per-request traceparent propagation the contract layer does not model
151161
const response = await fetch(buildStreamResumeUrl(streamId, cursor, { batch: true }), {
152-
signal,
162+
signal: AbortSignal.any([signal, AbortSignal.timeout(STREAM_BATCH_FETCH_TIMEOUT_MS)]),
153163
headers,
154164
})
155-
if (response.status === 404) return true
165+
if (isRelayEndStatus(response.status)) return true
156166
if (!response.ok) throw new Error(`Stream batch responded with status ${response.status}`)
157167
const batch = parseStreamBatchResponse(await response.json())
158168
relay.chatId ??= resolveChatIdFromStreamBatch(batch)
@@ -164,11 +174,11 @@ async function relayClientTools(turn: DetachedChatTurn, relay: Relay): Promise<v
164174
return isTerminalStreamStatus(batch.status)
165175
}
166176

167-
/** Tails one live connection; resolves true once the run is over. */
177+
/** Tails one live connection; resolves true once there is nothing left to relay. */
168178
const readTail = async (): Promise<boolean> => {
169179
// boundary-raw-fetch: live SSE tail endpoint streams events consumed via readSSELines
170180
const response = await fetch(buildStreamResumeUrl(streamId, cursor), { signal, headers })
171-
if (response.status === 404) return true
181+
if (isRelayEndStatus(response.status)) return true
172182
if (!response.ok || !response.body) {
173183
throw new Error(`Stream tail responded with status ${response.status}`)
174184
}
@@ -224,13 +234,10 @@ export function detachClientTools(turn: DetachedChatTurn): void {
224234
}
225235

226236
/**
227-
* Stops relaying a chat whose view reads its stream again. Tools the relay
228-
* already started keep running and report their outcome.
237+
* Stops relaying a stream the chat view reads again. Tools the relay already
238+
* started keep running and report their outcome.
229239
*/
230-
export function reattachClientTools(chatId: string): void {
231-
for (const [streamId, relay] of relays) {
232-
if (relay.chatId !== chatId) continue
233-
relay.controller.abort('chat_reattached')
234-
relays.delete(streamId)
235-
}
240+
export function reattachClientTools(streamId: string): void {
241+
relays.get(streamId)?.controller.abort('chat_reattached')
242+
relays.delete(streamId)
236243
}

‎apps/sim/app/workspace/[workspaceId]/home/hooks/stream/handle-tool-event.ts‎

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -202,11 +202,12 @@ export function handleToolEvent(ctx: StreamLoopContext, parsed: ToolEvent): void
202202
// reducer, run its side effects now (the result event had no node to act on).
203203
if (node?.kind === 'tool' && node.result) runToolResultSideEffects(ctx, node, replay)
204204

205-
const start = resolveClientToolStart(parsed, (toolCallId) => {
206-
if (deps.options.suppressedWorkflowToolStartIds?.has(toolCallId)) return false
207-
const tool = state.model.nodes.get(resolveToolId(state.model, toolCallId))
208-
return tool?.kind === 'tool' && tool.status === 'running' && !tool.result
209-
})
210-
if (start) startClientTool(deps, start)
205+
const start = resolveClientToolStart(parsed)
206+
const isPending =
207+
node?.kind === 'tool' &&
208+
node.status === 'running' &&
209+
!node.result &&
210+
!deps.options.suppressedWorkflowToolStartIds?.has(rawId)
211+
if (start && isPending) startClientTool(deps, start)
211212
ops.flush()
212213
}

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

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -178,6 +178,7 @@ import {
178178
isStreamSchemaValidationError,
179179
parseStreamBatchResponse,
180180
resolveChatIdFromStreamBatch,
181+
STREAM_BATCH_FETCH_TIMEOUT_MS,
181182
STREAM_IDLE_TIMEOUT_MS,
182183
type StreamBatchResponse,
183184
StreamGoneError,
@@ -293,7 +294,6 @@ const MAX_RECONNECT_ATTEMPTS = 10
293294
const RECONNECT_BASE_DELAY_MS = 1000
294295
const RECONNECT_MAX_DELAY_MS = 30_000
295296
const RECONNECT_EXHAUSTED_RECHECK_MS = 30_000
296-
const STREAM_BATCH_FETCH_TIMEOUT_MS = 10_000
297297
const STREAM_CHAT_ID_RESOLVE_TIMEOUT_MS = 10_000
298298
const CHAT_HISTORY_RECOVERY_TIMEOUT_MS = 10_000
299299
const STOP_REQUEST_TIMEOUT_MS = 15_000
@@ -1609,23 +1609,22 @@ export function useChat(
16091609

16101610
/**
16111611
* Hands the live turn's client tools to a detached relay as the user leaves
1612-
* its chat, so a desktop run keeps going in the background. A turn is live
1613-
* once the server admitted it: its stream id is known, or its response is
1614-
* being read (a send's stream id is its user message id). When a relay takes
1615-
* the turn, the caller releases the turn's controller without aborting it,
1616-
* so tools already running report their outcome instead of being cancelled.
1612+
* its chat, so a desktop run keeps going in the background. Only a turn the
1613+
* server admitted is detached; a send still awaiting admission is withdrawn
1614+
* by the unmount cleanup and handed to the next chat surface instead. When a
1615+
* relay takes the turn, the caller releases the turn's controller without
1616+
* aborting it, so tools already running report their outcome.
16171617
*
16181618
* @returns whether a relay took the turn
16191619
*/
16201620
detachLiveTurnClientToolsRef.current = (): boolean => {
1621-
const streamId =
1622-
streamIdRef.current ??
1623-
(streamReaderRef.current ? activeTurnRef.current?.userMessageId : undefined)
1621+
const streamId = streamIdRef.current
16241622
if (
16251623
!isDesktopApp() ||
16261624
requestModeRef.current === 'assistant' ||
16271625
!sendingRef.current ||
1628-
!streamId
1626+
!streamId ||
1627+
pendingChatAdmissionRef.current
16291628
) {
16301629
return false
16311630
}
@@ -2210,8 +2209,7 @@ export function useChat(
22102209
return { sawStreamError: false, sawComplete: false }
22112210
}
22122211
streamReaderRef.current = reader
2213-
const readerChatId = options?.targetChatId ?? chatIdRef.current
2214-
if (readerChatId) reattachClientTools(readerChatId)
2212+
if (streamIdRef.current) reattachClientTools(streamIdRef.current)
22152213

22162214
try {
22172215
await readSSELines(reader, {
@@ -2563,6 +2561,7 @@ export function useChat(
25632561
}
25642562

25652563
if (isStaleReconnect()) {
2564+
await sseRes.body.cancel()
25662565
return { error: false, aborted: true }
25672566
}
25682567

‎apps/sim/lib/mothership/tools/client/launch-local-filesystem-tool.ts‎

Lines changed: 0 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,9 @@
11
import { createLogger } from '@sim/logger'
2-
import { LRUCache } from 'lru-cache'
32
import type { LocalFilesystemExecutionContext } from '@/lib/mothership/tools/client/local-filesystem'
43
import { isNativeFileTool, isUserLocalVfsToolCall } from '@/lib/mothership/tools/local-filesystem'
54

65
const logger = createLogger('CopilotLocalFilesystemTool')
76

8-
/**
9-
* Exactly-once guard. A call's chat view and the relay that runs it after the
10-
* user leaves can both see it, and a remounted view replays calls still
11-
* running, so the guard lives here rather than with any one caller. Bounded:
12-
* a call evicted behind this many newer ones is long settled.
13-
*/
14-
const launchedToolCallIds = new LRUCache<string, true>({ max: 500 })
15-
167
/**
178
* Runs a local filesystem tool call on the desktop client. The executor is
189
* loaded on demand: it only runs for desktop-local calls, and a static import
@@ -33,8 +24,6 @@ export function launchLocalFilesystemTool(
3324
) {
3425
return
3526
}
36-
if (launchedToolCallIds.has(toolCallId)) return
37-
launchedToolCallIds.set(toolCallId, true)
3827

3928
import('@/lib/mothership/tools/client/local-filesystem').then(
4029
(m) => m.executeLocalFilesystemTool(toolCallId, toolName, args, context),

‎apps/sim/lib/mothership/tools/client/local-filesystem.ts‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import {
1414
} from '@sim/desktop-bridge/local-filesystem-limits'
1515
import { createLogger } from '@sim/logger'
1616
import { toError } from '@sim/utils/errors'
17+
import { LRUCache } from 'lru-cache'
1718
import micromatch from 'micromatch'
1819
import { getDesktopBridge } from '@/lib/desktop'
1920
import { ASYNC_TOOL_CONFIRMATION_STATUS } from '@/lib/mothership/async-runs/lifecycle'
@@ -337,12 +338,21 @@ async function execute(
337338
return executeUserLocalRead(toolCallId, args, context.signal)
338339
}
339340

341+
/**
342+
* Exactly-once guard. A call's chat view and the relay that runs it after the
343+
* user leaves can both start it, and a remounted view replays calls still
344+
* running. Bounded: a call evicted behind this many newer ones is long settled.
345+
*/
346+
const executedToolCallIds = new LRUCache<string, true>({ max: 500 })
347+
340348
export function executeLocalFilesystemTool(
341349
toolCallId: string,
342350
toolName: string,
343351
args: Record<string, unknown>,
344352
context: LocalFilesystemExecutionContext
345353
): void {
354+
if (executedToolCallIds.has(toolCallId)) return
355+
executedToolCallIds.set(toolCallId, true)
346356
if (isNativeFileTool(toolName)) {
347357
void executeNativeFileTool(toolCallId, toolName, context.signal)
348358
return

‎apps/sim/lib/mothership/tools/client/terminal-tool-execution.ts‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,9 @@ const logger = createLogger('CopilotTerminalToolExecution')
2525

2626
/** Tool events older than this are replays, not live instructions. */
2727
const MAX_EVENT_AGE_MS = 120_000
28+
/** A stale call never ran in this client; the waiter must hear so or the turn hangs. */
29+
const STALE_EVENT_MESSAGE =
30+
'This terminal command was delivered too late to run safely, so it was not run. Ask again to retry it.'
2831
const EXECUTED_STORAGE_PREFIX = 'sim:copilot:terminal-tool-executed:'
2932

3033
/**
@@ -121,6 +124,17 @@ export function executeTerminalToolOnClient(
121124
const age = eventAgeMs(eventTs)
122125
if (age !== null && age > MAX_EVENT_AGE_MS) {
123126
logger.info('Skipping stale terminal tool event', { toolCallId, operation, age })
127+
void reportClientToolCompletion(
128+
toolCallId,
129+
ASYNC_TOOL_CONFIRMATION_STATUS.error,
130+
STALE_EVENT_MESSAGE,
131+
{ error: STALE_EVENT_MESSAGE, staleEvent: true }
132+
).catch((reportErr) => {
133+
logger.error('Failed to report stale terminal tool event', {
134+
toolCallId,
135+
error: toError(reportErr).message,
136+
})
137+
})
124138
return
125139
}
126140
markExecuted(toolCallId)

0 commit comments

Comments
 (0)