Skip to content

Commit d8510e2

Browse files
committed
fix(mothership): hold a turn's Stop only while its desktop tools run
Replace the bounded LRU of turn controllers with leases: each running browser action or local file tool holds its turn's Stop until it settles, so a live turn can never be evicted and settled turns leave nothing behind.
1 parent 3a9bd01 commit d8510e2

6 files changed

Lines changed: 155 additions & 62 deletions

File tree

Lines changed: 54 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,54 @@
1+
import { describe, expect, it } from 'vitest'
2+
import { leaseDesktopTool, stopDesktopTools } from './desktop-tool-lifetimes'
3+
4+
describe('desktop tool leases', () => {
5+
it('cancels every running tool of the stopped turn and no other turn', () => {
6+
const first = leaseDesktopTool('turn-a')
7+
const second = leaseDesktopTool('turn-a')
8+
const other = leaseDesktopTool('turn-b')
9+
10+
stopDesktopTools('turn-a', 'user_stop')
11+
12+
expect(first.signal.aborted).toBe(true)
13+
expect(second.signal.aborted).toBe(true)
14+
expect(first.signal.reason).toBe('user_stop')
15+
expect(other.signal.aborted).toBe(false)
16+
other.release()
17+
})
18+
19+
it('keeps a turn reachable by Stop while any of its tools still runs', () => {
20+
const settled = leaseDesktopTool('turn-c')
21+
const running = leaseDesktopTool('turn-c')
22+
for (let turn = 0; turn < 500; turn++) leaseDesktopTool(`busy-${turn}`).release()
23+
24+
settled.release()
25+
settled.release()
26+
stopDesktopTools('turn-c', 'user_stop')
27+
28+
expect(running.signal.aborted).toBe(true)
29+
})
30+
31+
it('gives a turn whose tools all settled a fresh lifetime for its next tool', () => {
32+
const settled = leaseDesktopTool('turn-d')
33+
settled.release()
34+
stopDesktopTools('turn-d', 'user_stop')
35+
36+
const next = leaseDesktopTool('turn-d')
37+
38+
expect(settled.signal.aborted).toBe(false)
39+
expect(next.signal).not.toBe(settled.signal)
40+
expect(next.signal.aborted).toBe(false)
41+
next.release()
42+
})
43+
44+
it('does not let a tool that settles after Stop release a newer lease on the turn', () => {
45+
const stopped = leaseDesktopTool('turn-e')
46+
stopDesktopTools('turn-e', 'user_stop')
47+
const next = leaseDesktopTool('turn-e')
48+
49+
stopped.release()
50+
stopDesktopTools('turn-e', 'user_stop')
51+
52+
expect(next.signal.aborted).toBe(true)
53+
})
54+
})
Lines changed: 40 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1,28 +1,51 @@
1-
import { LRUCache } from 'lru-cache'
1+
/** The desktop tools of one turn that are still running in this tab, and the Stop that ends them. */
2+
interface RunningTurnTools {
3+
stop: AbortController
4+
running: number
5+
}
26

37
/**
4-
* Turns whose desktop tools may still be running in this tab. A tool outlives the chat view that
5-
* started it (and any stream reader), so the turn, not the view, owns its Stop. Bounded: a tool
6-
* runs for minutes, far fewer turns than this.
8+
* Turns with a desktop tool still running in this tab, keyed by the turn's stream id. A tool
9+
* outlives the chat view that started it (and any stream reader), so the turn, not the view, owns
10+
* its Stop. A turn is held only while one of its tools runs: each tool releases it as it settles.
711
*/
8-
const turnStops = new LRUCache<string, AbortController>({ max: 64 })
12+
const runningTurns = new Map<string, RunningTurnTools>()
13+
14+
/** A running desktop tool's hold on its turn. */
15+
interface DesktopToolLease {
16+
/** Aborted only by the user's Stop of the turn. */
17+
signal: AbortSignal
18+
/** Called once when the tool settles. */
19+
release(): void
20+
}
921

1022
/**
11-
* The lifetime of a desktop tool (a browser action, a local file read or import) started for a
12-
* turn: only the user's Stop of that turn ends it. Replacing the stream reader, leaving the chat
13-
* view, or stopping another chat's turn leaves it running to finish and report its own result.
23+
* Starts a desktop tool (a browser action, a local file read or import) for a turn. Only the
24+
* user's Stop of that turn cancels it: replacing the stream reader, leaving the chat view, or
25+
* stopping another chat's turn leaves it running to finish and report its own result.
1426
*/
15-
export function desktopToolLifetime(streamId: string): AbortSignal {
16-
let stop = turnStops.get(streamId)
17-
if (!stop) {
18-
stop = new AbortController()
19-
turnStops.set(streamId, stop)
27+
export function leaseDesktopTool(streamId: string): DesktopToolLease {
28+
let turn = runningTurns.get(streamId)
29+
if (!turn) {
30+
turn = { stop: new AbortController(), running: 0 }
31+
runningTurns.set(streamId, turn)
32+
}
33+
turn.running += 1
34+
const held = turn
35+
let released = false
36+
return {
37+
signal: held.stop.signal,
38+
release() {
39+
if (released) return
40+
released = true
41+
held.running -= 1
42+
if (held.running === 0 && runningTurns.get(streamId) === held) runningTurns.delete(streamId)
43+
},
2044
}
21-
return stop.signal
2245
}
2346

24-
/** Cancels the desktop tools of a turn the user stopped, from whichever view started them. */
47+
/** Cancels the running desktop tools of a turn the user stopped, from whichever view started them. */
2548
export function stopDesktopTools(streamId: string, reason: string): void {
26-
turnStops.get(streamId)?.abort(reason)
27-
turnStops.delete(streamId)
49+
runningTurns.get(streamId)?.stop.abort(reason)
50+
runningTurns.delete(streamId)
2851
}

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.dom.test.tsx‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2217,6 +2217,9 @@ describe('useChat remount send recovery', () => {
22172217

22182218
beforeEach(() => {
22192219
libDesktopMockFns.mockIsDesktopApp.mockReturnValue(true)
2220+
const stillRunning = () => new Promise<void>(() => {})
2221+
mockExecuteBrowserToolOnClient.mockImplementation(stillRunning)
2222+
mockExecuteLocalFilesystemTool.mockImplementation(stillRunning)
22202223
})
22212224

22222225
afterEach(() => {

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

Lines changed: 48 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -92,7 +92,7 @@ import { initTerminalTransport } from '@/lib/terminal/transport'
9292
import { getQueryClient } from '@/app/_shell/providers/get-query-client'
9393
import { chatUrl } from '@/app/workspace/[workspaceId]/home/hooks/chat-url'
9494
import {
95-
desktopToolLifetime,
95+
leaseDesktopTool,
9696
stopDesktopTools,
9797
} from '@/app/workspace/[workspaceId]/home/hooks/desktop-tool-lifetimes'
9898
import { useFilePreviewController } from '@/app/workspace/[workspaceId]/home/hooks/preview'
@@ -472,10 +472,18 @@ function startClientBrowserTool(
472472
toolArgs: Record<string, unknown>,
473473
scopeId: string,
474474
eventTs?: string,
475-
signal?: AbortSignal
475+
turnStreamId?: string
476476
): void {
477477
if (!isCurrentBrowserToolName(toolName)) return
478-
executeBrowserToolOnClient(toolCallId, toolName, toolArgs, scopeId, eventTs, signal)
478+
const lease = turnStreamId ? leaseDesktopTool(turnStreamId) : undefined
479+
void executeBrowserToolOnClient(
480+
toolCallId,
481+
toolName,
482+
toolArgs,
483+
scopeId,
484+
eventTs,
485+
lease?.signal
486+
).finally(() => lease?.release())
479487
}
480488

481489
/**
@@ -1580,10 +1588,11 @@ export function useChat(
15801588
return
15811589
}
15821590
handledClientLocalFilesystemToolIds.add(toolCallId)
1591+
const lease = streamIdRef.current ? leaseDesktopTool(streamIdRef.current) : undefined
15831592
const options = {
15841593
workspaceId,
15851594
chatId: chatIdRef.current ?? selectedChatIdRef.current,
1586-
signal: streamIdRef.current ? desktopToolLifetime(streamIdRef.current) : undefined,
1595+
signal: lease?.signal,
15871596
}
15881597
/**
15891598
* Dynamic on purpose: the local-filesystem executor only runs for desktop-local
@@ -1594,35 +1603,40 @@ export function useChat(
15941603
* report an error completion rather than leaving it hanging with the dedupe ref
15951604
* already marked handled.
15961605
*/
1597-
import('@/lib/mothership/tools/client/local-filesystem').then(
1598-
(m) => m.executeLocalFilesystemTool(toolCallId, toolName, toolArgs, options),
1599-
async (error) => {
1600-
logger.error('Failed to load local filesystem tool executor', { error })
1601-
/**
1602-
* The recovery itself can reject (the helper chunks or the completion POST can
1603-
* fail for the same reason the executor chunk did). Contain it: an unhandled
1604-
* rejection here would settle nothing and surface as a console error, exactly
1605-
* like the executor's own report-failure path, which also degrades to a log.
1606-
*/
1607-
try {
1608-
const [{ reportClientToolCompletion }, { ASYNC_TOOL_CONFIRMATION_STATUS }] =
1609-
await Promise.all([
1610-
import('@/lib/mothership/tools/client/completion'),
1611-
import('@/lib/mothership/async-runs/lifecycle'),
1612-
])
1613-
await reportClientToolCompletion(
1614-
toolCallId,
1615-
ASYNC_TOOL_CONFIRMATION_STATUS.error,
1616-
'Local filesystem tool failed to load'
1617-
)
1618-
} catch (reportError) {
1619-
logger.error('Failed to report local filesystem tool load failure', {
1620-
toolCallId,
1621-
error: reportError,
1622-
})
1606+
import('@/lib/mothership/tools/client/local-filesystem')
1607+
.then(
1608+
(m) => m.executeLocalFilesystemTool(toolCallId, toolName, toolArgs, options),
1609+
async (error) => {
1610+
logger.error('Failed to load local filesystem tool executor', { error })
1611+
/**
1612+
* The recovery itself can reject (the helper chunks or the completion POST can
1613+
* fail for the same reason the executor chunk did). Contain it: an unhandled
1614+
* rejection here would settle nothing and surface as a console error, exactly
1615+
* like the executor's own report-failure path, which also degrades to a log.
1616+
*/
1617+
try {
1618+
const [{ reportClientToolCompletion }, { ASYNC_TOOL_CONFIRMATION_STATUS }] =
1619+
await Promise.all([
1620+
import('@/lib/mothership/tools/client/completion'),
1621+
import('@/lib/mothership/async-runs/lifecycle'),
1622+
])
1623+
await reportClientToolCompletion(
1624+
toolCallId,
1625+
ASYNC_TOOL_CONFIRMATION_STATUS.error,
1626+
'Local filesystem tool failed to load'
1627+
)
1628+
} catch (reportError) {
1629+
logger.error('Failed to report local filesystem tool load failure', {
1630+
toolCallId,
1631+
error: reportError,
1632+
})
1633+
}
16231634
}
1624-
}
1625-
)
1635+
)
1636+
.catch((error) => {
1637+
logger.error('Local filesystem tool execution failed unexpectedly', { toolCallId, error })
1638+
})
1639+
.finally(() => lease?.release())
16261640
},
16271641
[workspaceId, organizationId, scopeKey]
16281642
)
@@ -2159,9 +2173,7 @@ export function useChat(
21592173
shouldContinue?: () => boolean
21602174
}
21612175
) => {
2162-
const browserToolSignal = streamIdRef.current
2163-
? desktopToolLifetime(streamIdRef.current)
2164-
: undefined
2176+
const turnStreamId = streamIdRef.current
21652177
const activityTracker = getResourceActivityTracker(
21662178
expectedGen ?? streamGenRef.current,
21672179
options?.targetChatId
@@ -2182,7 +2194,7 @@ export function useChat(
21822194
eventTs?: string
21832195
) => {
21842196
const scopeId = activityScopeId()
2185-
startClientBrowserTool(toolCallId, toolName, toolArgs, scopeId, eventTs, browserToolSignal)
2197+
startClientBrowserTool(toolCallId, toolName, toolArgs, scopeId, eventTs, turnStreamId)
21862198
}
21872199
const startClientTerminalToolForStream = (
21882200
toolCallId: string,

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

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -551,20 +551,21 @@ function timeoutForTool(toolName: BrowserToolName, params: Record<string, unknow
551551
}
552552

553553
/**
554-
* Fire-and-forget entry point invoked by the stream tool-event handler when a
555-
* `browser_*` client tool call arrives.
554+
* Entry point invoked by the stream tool-event handler when a `browser_*`
555+
* client tool call arrives. It reports its own outcome; the returned promise
556+
* only tells the caller when the action has settled.
556557
*
557558
* @param eventTs - the stream envelope's emission timestamp; stale events
558559
* (replays after reconnect/reload) are dropped rather than re-executed.
559560
*/
560-
export function executeBrowserToolOnClient(
561+
export async function executeBrowserToolOnClient(
561562
toolCallId: string,
562563
toolName: BrowserToolName,
563564
params: Record<string, unknown>,
564565
scopeId = useBrowserSessionStore.getState().activeScopeId,
565566
eventTs?: string,
566567
abortSignal?: AbortSignal
567-
): void {
568+
): Promise<void> {
568569
if (retryRetainedTerminalCompletion(toolCallId)) {
569570
logger.info('Suppressing browser tool while recovering its terminal completion', {
570571
toolCallId,
@@ -703,7 +704,7 @@ export function executeBrowserToolOnClient(
703704
}
704705
}
705706
runningBrowserToolCalls.add(toolCallId)
706-
void doExecuteBrowserTool(
707+
await doExecuteBrowserTool(
707708
toolCallId,
708709
toolName,
709710
params,

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

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -337,17 +337,17 @@ async function execute(
337337
return executeUserLocalRead(toolCallId, args, context.signal)
338338
}
339339

340-
export function executeLocalFilesystemTool(
340+
export async function executeLocalFilesystemTool(
341341
toolCallId: string,
342342
toolName: string,
343343
args: Record<string, unknown>,
344344
context: LocalFilesystemExecutionContext
345-
): void {
345+
): Promise<void> {
346346
if (isNativeFileTool(toolName)) {
347-
void executeNativeFileTool(toolCallId, toolName, context.signal)
347+
await executeNativeFileTool(toolCallId, toolName, context.signal)
348348
return
349349
}
350-
void execute(toolCallId, toolName, args, context).then(
350+
await execute(toolCallId, toolName, args, context).then(
351351
async (data) => {
352352
if (context.signal?.aborted) return
353353
try {

0 commit comments

Comments
 (0)