Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
fix(tools): close operation lifecycle gaps
  • Loading branch information
icecrasher321 committed Aug 28, 2026
commit ddcbb26bf48daaf5563f91dd1bbd8315c1c06072
50 changes: 50 additions & 0 deletions apps/sim/lib/internal/browser-use/operations/run-task.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -233,4 +233,54 @@ describe('executeRunTaskOperation', () => {
)
expect(mockFetch.mock.calls[2]?.[1]?.signal).not.toBe(controller.signal)
})

it('stops an automatically created task session when polling is cancelled', async () => {
const controller = new AbortController()
const abortError = new DOMException('cancelled', 'AbortError')
mockFetch
.mockResolvedValueOnce(jsonResponse({ id: 'task-1', sessionId: 'task-session' }))
.mockImplementationOnce(async () => {
controller.abort(abortError)
throw abortError
})
.mockResolvedValueOnce(new Response(null, { status: 204 }))

await expect(
executeRunTaskOperation({ task: 'Open the page', apiKey: 'api-key' }, controller.signal)
).rejects.toBe(abortError)

expect(mockFetch).toHaveBeenNthCalledWith(
3,
'https://api.browser-use.com/api/v2/sessions/task-session',
expect.objectContaining({
method: 'PATCH',
body: JSON.stringify({ action: 'stop' }),
signal: expect.any(AbortSignal),
})
)
})

it('stops an automatically created task session when polling times out', async () => {
const now = vi.spyOn(Date, 'now').mockReturnValueOnce(0).mockReturnValue(1_000_000_000_000_000)
mockFetch
.mockResolvedValueOnce(jsonResponse({ id: 'task-1', sessionId: 'task-session' }))
.mockResolvedValueOnce(jsonResponse({ status: 'running', sessionId: 'task-session' }))
.mockResolvedValueOnce(
jsonResponse({ shareUrl: 'https://browser-use.com/share/task-session' })
)
.mockResolvedValueOnce(new Response(null, { status: 204 }))

const result = await executeRunTaskOperation({ task: 'Open the page', apiKey: 'api-key' })
now.mockRestore()

expect(result).toMatchObject({
success: false,
error: expect.stringContaining('Task did not complete within the maximum polling time'),
})
expect(mockFetch).toHaveBeenNthCalledWith(
4,
'https://api.browser-use.com/api/v2/sessions/task-session',
expect.objectContaining({ method: 'PATCH', signal: expect.any(AbortSignal) })
)
})
})
51 changes: 40 additions & 11 deletions apps/sim/lib/internal/browser-use/operations/run-task.ts
Original file line number Diff line number Diff line change
Expand Up @@ -304,6 +304,7 @@ async function fetchTaskStatus(

interface PollResult {
success: boolean
taskEnded: boolean
output: unknown
steps: BrowserUseTaskStep[]
sessionId: string | null
Expand All @@ -312,12 +313,18 @@ interface PollResult {
error?: string
}

interface PollOptions {
initialSessionId: string | null
signal?: AbortSignal
onSessionId: (sessionId: string) => void
}

async function pollForCompletion(
taskId: string,
apiKey: string,
initialSessionId: string | null,
signal?: AbortSignal
options: PollOptions
): Promise<PollResult> {
const { initialSessionId, signal, onSessionId } = options
let consecutiveErrors = 0
let sessionId = initialSessionId
let liveUrl: string | null = null
Expand All @@ -337,6 +344,7 @@ async function pollForCompletion(
if (consecutiveErrors >= MAX_CONSECUTIVE_ERRORS) {
return {
success: false,
taskEnded: false,
output: null,
steps: [],
sessionId,
Expand All @@ -352,7 +360,10 @@ async function pollForCompletion(

consecutiveErrors = 0
const taskData = result.data
if (taskData.sessionId) sessionId = taskData.sessionId
if (taskData.sessionId) {
sessionId = taskData.sessionId
onSessionId(taskData.sessionId)
}
const status = taskData.status

logger.info(`BrowserUse task ${taskId} status: ${status}`)
Expand All @@ -370,6 +381,7 @@ async function pollForCompletion(
const output = taskData.output ?? null
return {
success: status === 'finished',
Comment thread
icecrasher321 marked this conversation as resolved.
taskEnded: true,
output,
steps: taskData.steps ?? [],
sessionId,
Expand All @@ -391,11 +403,14 @@ async function pollForCompletion(
if (finalResult.ok && ['finished', 'failed', 'stopped'].includes(finalResult.data.status)) {
const status = finalResult.data.status
const output = finalResult.data.output ?? null
const finalSessionId = finalResult.data.sessionId ?? sessionId
if (finalSessionId) onSessionId(finalSessionId)
return {
success: status === 'finished',
taskEnded: true,
output,
steps: finalResult.data.steps ?? [],
sessionId: finalResult.data.sessionId ?? sessionId,
sessionId: finalSessionId,
liveUrl,
publicShareUrl,
error:
Expand All @@ -409,6 +424,7 @@ async function pollForCompletion(

return {
success: false,
taskEnded: false,
output: null,
steps: [],
sessionId,
Expand Down Expand Up @@ -470,20 +486,22 @@ export const executeRunTaskOperation: InternalToolOperationImplementation<
params: BrowserUseRunTaskParams,
signal?: AbortSignal
): Promise<BrowserUseRunTaskResponse> => {
let sessionId: string | undefined
let profileSessionId: string | undefined
let taskSessionId: string | null = null
let taskEnded = false

if (params.profile_id) {
logger.info(`Creating session with profile ID: ${params.profile_id}`)
const sessionResult = await createSessionWithProfile(params.profile_id, params.apiKey, signal)
if ('error' in sessionResult) {
return { success: false, output: emptyOutput(), error: sessionResult.error }
}
sessionId = sessionResult.sessionId
profileSessionId = sessionResult.sessionId
}

try {
const requestBody = buildRequestBody(params, sessionId)
logger.info('Creating BrowserUse task', { hasSession: !!sessionId })
const requestBody = buildRequestBody(params, profileSessionId)
logger.info('Creating BrowserUse task', { hasSession: !!profileSessionId })
const response = await fetchBrowserUse('/tasks', params.apiKey, {
method: 'POST',
body: requestBody,
Expand Down Expand Up @@ -512,10 +530,18 @@ export const executeRunTaskOperation: InternalToolOperationImplementation<
}
const data = parsed.data
const taskId = data.id
const initialSessionId = sessionId ?? data.sessionId ?? null
const initialSessionId = profileSessionId ?? data.sessionId ?? null
taskSessionId = initialSessionId
logger.info(`Created BrowserUse task ${taskId}`, { sessionId: initialSessionId })

const result = await pollForCompletion(taskId, params.apiKey, initialSessionId, signal)
const result = await pollForCompletion(taskId, params.apiKey, {
initialSessionId,
signal,
onSessionId: (discoveredSessionId) => {
taskSessionId = discoveredSessionId
},
})
taskEnded = result.taskEnded

const finalSessionId = result.sessionId ?? initialSessionId
const shareUrl =
Expand Down Expand Up @@ -544,7 +570,10 @@ export const executeRunTaskOperation: InternalToolOperationImplementation<
error: `Error creating task: ${getErrorMessage(error, 'Unknown error')}`,
}
} finally {
if (sessionId) {
const sessionsToStop = new Set<string>()
if (profileSessionId) sessionsToStop.add(profileSessionId)
if (!taskEnded && taskSessionId) sessionsToStop.add(taskSessionId)
for (const sessionId of sessionsToStop) {
await stopSession(sessionId, params.apiKey)
Comment thread
icecrasher321 marked this conversation as resolved.
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -90,4 +90,18 @@ describe('executeStorageUpdateBucketOperation', () => {
executeStorageUpdateBucketOperation(INPUT, controller.signal)
).rejects.toMatchObject({ name: 'AbortError' })
})

it('returns a structured failure for an invalid project reference', async () => {
const result = await executeStorageUpdateBucketOperation({
...INPUT,
projectId: '../invalid',
})

expect(result).toMatchObject({
success: false,
output: { message: 'Failed to update storage bucket', results: {} },
error: expect.any(String),
})
expect(fetchMock).not.toHaveBeenCalled()
})
})
Original file line number Diff line number Diff line change
Expand Up @@ -12,15 +12,14 @@ export const executeStorageUpdateBucketOperation: InternalToolOperationImplement
params: SupabaseStorageUpdateBucketParams,
signal
): Promise<SupabaseStorageUpdateBucketResponse> => {
const baseUrl = supabaseBaseUrl(params.projectId)
const bucket = encodeStorageSegment(params.bucket)
const headers = {
apikey: params.apiKey,
Authorization: `Bearer ${params.apiKey}`,
'Content-Type': 'application/json',
}

try {
const baseUrl = supabaseBaseUrl(params.projectId)
const bucket = encodeStorageSegment(params.bucket)
const headers = {
apikey: params.apiKey,
Authorization: `Bearer ${params.apiKey}`,
'Content-Type': 'application/json',
}
const currentResponse = await fetch(`${baseUrl}/storage/v1/bucket/${bucket}`, {
Comment thread
icecrasher321 marked this conversation as resolved.
Outdated
Comment thread
icecrasher321 marked this conversation as resolved.
Outdated
method: 'GET',
headers,
Expand Down
Loading