|
| 1 | +/** |
| 2 | + * Keeps a chat's client-executed tools running after the user leaves it |
| 3 | + * mid-turn. Leaving detaches the chat view but not the server run, and that |
| 4 | + * run's browser, terminal, and local filesystem tools execute only in this |
| 5 | + * 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. |
| 11 | + */ |
| 12 | +import { createLogger } from '@sim/logger' |
| 13 | +import { getErrorMessage } from '@sim/utils/errors' |
| 14 | +import { interruptibleSleep } from '@sim/utils/helpers' |
| 15 | +import { backoffWithJitter } from '@sim/utils/retry' |
| 16 | +import { readSSELines } from '@/lib/core/utils/sse' |
| 17 | +import { |
| 18 | + MothershipStreamV1EventType, |
| 19 | + MothershipStreamV1ToolPhase, |
| 20 | +} from '@/lib/mothership/generated/mothership-stream-v1' |
| 21 | +import { parsePersistedStreamEventEnvelopeJson } from '@/lib/mothership/request/session/contract' |
| 22 | +import { executeBrowserToolOnClient } from '@/lib/mothership/tools/client/browser-tool-execution' |
| 23 | +import { launchLocalFilesystemTool } from '@/lib/mothership/tools/client/launch-local-filesystem-tool' |
| 24 | +import { executeTerminalToolOnClient } from '@/lib/mothership/tools/client/terminal-tool-execution' |
| 25 | +import { |
| 26 | + type ClientToolStart, |
| 27 | + resolveClientToolStart, |
| 28 | +} from '@/app/workspace/[workspaceId]/home/hooks/stream/client-tool-start' |
| 29 | +import { |
| 30 | + buildStreamResumeUrl, |
| 31 | + createStreamSchemaValidationError, |
| 32 | + getStreamEventCursor, |
| 33 | + isAlreadyProcessedStreamCursor, |
| 34 | + isStreamSchemaValidationError, |
| 35 | + STREAM_IDLE_TIMEOUT_MS, |
| 36 | +} from '@/app/workspace/[workspaceId]/home/hooks/stream-protocol' |
| 37 | + |
| 38 | +const logger = createLogger('DetachedClientTools') |
| 39 | + |
| 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. */ |
| 44 | +export interface DetachedChatTurn { |
| 45 | + chatId: string |
| 46 | + streamId: string |
| 47 | + /** The last cursor the chat view dispatched; the relay resumes after it. */ |
| 48 | + afterCursor: string |
| 49 | + traceparent?: string |
| 50 | + workspaceId?: string |
| 51 | + /** The chat's desktop scope, which owns its browser and terminal tabs. */ |
| 52 | + scopeId: string |
| 53 | +} |
| 54 | + |
| 55 | +/** Relays by chat id. An entry is removed when its relay ends. */ |
| 56 | +const relays = new Map<string, AbortController>() |
| 57 | + |
| 58 | +/** |
| 59 | + * Starts the client tools a detached turn hands this client. Workflow runs are |
| 60 | + * left alone: running one drives the workflow editor, and the server runs a |
| 61 | + * workflow call itself when no client picks it up. |
| 62 | + */ |
| 63 | +function startDetachedClientTool(turn: DetachedChatTurn, start: ClientToolStart): void { |
| 64 | + const { toolCallId, toolName, args, eventTs } = start |
| 65 | + switch (start.kind) { |
| 66 | + case 'workflow': |
| 67 | + return |
| 68 | + case 'localFilesystem': |
| 69 | + launchLocalFilesystemTool(toolCallId, toolName, args, { |
| 70 | + workspaceId: turn.workspaceId, |
| 71 | + chatId: turn.chatId, |
| 72 | + }) |
| 73 | + return |
| 74 | + case 'browser': |
| 75 | + executeBrowserToolOnClient(toolCallId, start.toolName, args, turn.scopeId, eventTs) |
| 76 | + return |
| 77 | + case 'terminal': |
| 78 | + executeTerminalToolOnClient(toolCallId, args, turn.scopeId, eventTs) |
| 79 | + return |
| 80 | + } |
| 81 | +} |
| 82 | + |
| 83 | +async function relayClientTools(turn: DetachedChatTurn, signal: AbortSignal): Promise<void> { |
| 84 | + const { streamId, traceparent } = turn |
| 85 | + /** Calls this relay started or saw settle; ids alone decide what is pending. */ |
| 86 | + const handledToolCallIds = new Set<string>() |
| 87 | + let cursor = turn.afterCursor |
| 88 | + let stalledAttempts = 0 |
| 89 | + |
| 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), { |
| 94 | + signal, |
| 95 | + ...(traceparent ? { headers: { traceparent } } : {}), |
| 96 | + }) |
| 97 | + if (response.status === 404) return true |
| 98 | + if (!response.ok || !response.body) { |
| 99 | + throw new Error(`Stream tail responded with status ${response.status}`) |
| 100 | + } |
| 101 | + |
| 102 | + let complete = false |
| 103 | + await readSSELines(response.body, { |
| 104 | + signal, |
| 105 | + idleTimeoutMs: STREAM_IDLE_TIMEOUT_MS, |
| 106 | + onData: (raw) => { |
| 107 | + const parsed = parsePersistedStreamEventEnvelopeJson(raw) |
| 108 | + 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) { |
| 127 | + complete = true |
| 128 | + return true |
| 129 | + } |
| 130 | + }, |
| 131 | + }) |
| 132 | + return complete |
| 133 | + } |
| 134 | + |
| 135 | + while (!signal.aborted) { |
| 136 | + const cursorBeforeAttempt = cursor |
| 137 | + try { |
| 138 | + if (await readTail()) return |
| 139 | + } catch (error) { |
| 140 | + if (signal.aborted) return |
| 141 | + if (isStreamSchemaValidationError(error)) { |
| 142 | + logger.error('Stopped relaying detached client tools on an invalid stream event', { |
| 143 | + streamId, |
| 144 | + error: error.message, |
| 145 | + }) |
| 146 | + return |
| 147 | + } |
| 148 | + logger.warn('Detached stream tail failed', { streamId, error: getErrorMessage(error) }) |
| 149 | + } |
| 150 | + |
| 151 | + if (cursor !== cursorBeforeAttempt) { |
| 152 | + stalledAttempts = 0 |
| 153 | + continue |
| 154 | + } |
| 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) |
| 164 | + } |
| 165 | +} |
| 166 | + |
| 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 | + */ |
| 171 | +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) |
| 177 | + }) |
| 178 | +} |
| 179 | + |
| 180 | +/** |
| 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. |
| 183 | + */ |
| 184 | +export function reattachClientTools(chatId: string): void { |
| 185 | + relays.get(chatId)?.abort('chat_reattached') |
| 186 | + relays.delete(chatId) |
| 187 | +} |
0 commit comments