From fde21d954e3de6be673124784de81872486500e1 Mon Sep 17 00:00:00 2001 From: MagMueller Date: Thu, 23 Jul 2026 19:36:46 -0700 Subject: [PATCH 1/6] fix(browser): isolate timed-out snippets --- packages/bcode-browser/src/browser-execute.ts | 23 +++++--- packages/bcode-browser/src/cdp/session.ts | 35 ++++++++++- packages/bcode-browser/src/session-store.ts | 7 +++ .../test/browser-execute.test.ts | 58 +++++++++++++++++++ 4 files changed, 114 insertions(+), 9 deletions(-) diff --git a/packages/bcode-browser/src/browser-execute.ts b/packages/bcode-browser/src/browser-execute.ts index 230bacbd99..375e418574 100644 --- a/packages/bcode-browser/src/browser-execute.ts +++ b/packages/bcode-browser/src/browser-execute.ts @@ -32,10 +32,11 @@ // // Cancellation: JS Promises are not preemptively cancellable. A snippet // without `await` yield-points (e.g. `for (let i = 0; i < 1e9; i++) {}`) -// runs to completion before our timeout fiber observes it. `Effect.timeoutOrElse` -// fails the surrounding fiber but the orphan Promise keeps running until it -// finishes. This matches the `uv run` subprocess case (SIGTERM only after -// the Python signal handler yields). Document, don't fix. +// runs to completion before our timeout fiber observes it. When a yielding +// snippet times out, its Promise may continue, so we permanently invalidate +// the Session object it received. Abandoned code can finish local work but +// cannot reconnect or send later CDP commands; the next tool call gets a +// fresh Session from SessionStore. // // Level 1 per decisions.md §1c — substantial implementation lives here. The // Level-2 hook in packages/opencode is a thin adapter. @@ -157,9 +158,9 @@ const serialize = (v: unknown): string => { export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) { const skillsDir = yield* Effect.promise(() => Skills.resolveSkillsDir(dataDir)) - const execute = (args: Parameters, ctx: ExecuteContext) => - Effect.gen(function* () { - const session = SessionStore.get(ctx.sessionID) + const execute = (args: Parameters, ctx: ExecuteContext) => { + const session = SessionStore.get(ctx.sessionID) + return Effect.gen(function* () { yield* Effect.promise(() => fs.mkdir(ctx.workspaceDir, { recursive: true })) const wrapped = yield* Effect.try({ @@ -228,9 +229,15 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) Effect.scoped, Effect.timeoutOrElse({ duration: Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS), - orElse: () => Effect.fail(new Error("browser_execute timed out")), + orElse: () => + Effect.gen(function* () { + const error = new Error("browser_execute timed out; CDP session was reset") + yield* Effect.sync(() => SessionStore.invalidate(ctx.sessionID, session, error)) + return yield* Effect.fail(error) + }), }), ) + } return { parameters, execute, skillsDir } }) diff --git a/packages/bcode-browser/src/cdp/session.ts b/packages/bcode-browser/src/cdp/session.ts index 291d919096..2ba605c0ad 100644 --- a/packages/bcode-browser/src/cdp/session.ts +++ b/packages/bcode-browser/src/cdp/session.ts @@ -50,6 +50,7 @@ export class Session implements Transport { private activeTargetId: string | undefined; private reattachPromise?: Promise; private enabledDomains = new Map>(); + private invalidatedError?: Error; private eventListeners: Array<(method: string, params: unknown, sessionId?: string) => void> = []; private callResultListeners: Array<(method: string, params: unknown, result: unknown) => void> = []; @@ -83,6 +84,8 @@ export class Session implements Transport { * and we connect directly to the supplied endpoint. */ async connect(opts: ConnectOptions = {}): Promise { + if (this.invalidatedError) throw this.invalidatedError; + // No-argument connect is an ensure-connected operation. Reopening the // same configured endpoint would discard the active target session and // make the next page command run against the browser-level socket. @@ -138,6 +141,10 @@ export class Session implements Transport { try { ws.close(); } catch { /* ignore */ } return; } + if (this.invalidatedError) { + finish(this.invalidatedError); + return; + } const previous = this.ws; this.ws = ws; this.activeSessionId = undefined; @@ -164,13 +171,37 @@ export class Session implements Transport { } isConnected(): boolean { - return this.ws?.readyState === WebSocket.OPEN; + return !this.invalidatedError && this.ws?.readyState === WebSocket.OPEN; } close(): void { this.ws?.close(); } + /** + * Permanently retire this Session object. + * + * Used when an in-process snippet outlives its tool timeout. The Promise + * itself cannot be preempted, so closing alone is insufficient: abandoned + * code could call connect() again later. Invalidated sessions reject every + * future transport operation, while SessionStore gives the next tool call + * a fresh Session object. + */ + invalidate(error: Error): void { + if (this.invalidatedError) return; + this.invalidatedError = error; + const ws = this.ws; + this.ws = undefined; + this.activeSessionId = undefined; + this.activeTargetId = undefined; + this.enabledDomains.clear(); + this.eventListeners = []; + this.callResultListeners = []; + if (!ws) return; + this.rejectPending(ws, error); + try { ws.close(); } catch { /* ignore */ } + } + /** * Pick a target and make subsequent calls auto-route to it. * Uses Target.attachToTarget with flatten:true (single-WS, sessionId-on-message). @@ -184,6 +215,7 @@ export class Session implements Transport { /** Set the active sessionId directly (e.g. one you already attached). */ setActiveSession(sessionId: string | undefined): void { + if (this.invalidatedError) throw this.invalidatedError; this.activeSessionId = sessionId; this.activeTargetId = undefined; } @@ -257,6 +289,7 @@ export class Session implements Transport { } private send(method: string, params: unknown, sessionId?: string): Promise { + if (this.invalidatedError) return Promise.reject(this.invalidatedError); const ws = this.ws; if (!ws || ws.readyState !== WebSocket.OPEN) { return Promise.reject(new Error('Not connected. Call session.connect(...) first.')); diff --git a/packages/bcode-browser/src/session-store.ts b/packages/bcode-browser/src/session-store.ts index 8be05462c8..7494849f4b 100644 --- a/packages/bcode-browser/src/session-store.ts +++ b/packages/bcode-browser/src/session-store.ts @@ -27,6 +27,13 @@ export const get = (sessionID: string): Session => { return fresh } +export const invalidate = (sessionID: string, expected: Session, error: Error): void => { + const entry = sessions.get(sessionID) + if (entry !== expected) return + sessions.delete(sessionID) + entry.invalidate(error) +} + export const evict = async (sessionID: string): Promise => { const entry = sessions.get(sessionID) if (!entry) return diff --git a/packages/bcode-browser/test/browser-execute.test.ts b/packages/bcode-browser/test/browser-execute.test.ts index fa11b6c6eb..c8486c201f 100644 --- a/packages/bcode-browser/test/browser-execute.test.ts +++ b/packages/bcode-browser/test/browser-execute.test.ts @@ -270,3 +270,61 @@ test("overlapping execute calls do not clobber each other's console capture", as [aWorkspace, bWorkspace, aData, bData].map((d) => fs.rm(d, { recursive: true, force: true })), ) }) + +test("a timed-out snippet cannot send later CDP commands or reconnect", async () => { + const timedOutSessionID = "timeout-" + Math.random().toString(36).slice(2, 8) + const timedOutWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-ws-")) + const timedOutData = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-data-")) + let lateCommandCount = 0 + const server = Bun.serve({ + port: 0, + fetch(req, bunServer) { + return bunServer.upgrade(req) ? undefined : new Response("nope", { status: 400 }) + }, + websocket: { + message(socket, raw) { + const message = JSON.parse(String(raw)) + if (message.method !== "Runtime.evaluate") return + lateCommandCount++ + socket.send(JSON.stringify({ id: message.id, result: { result: { type: "boolean", value: true } } })) + }, + close() {}, + }, + }) + if (server.port === undefined) throw new Error("test server has no port") + const wsUrl = `ws://127.0.0.1:${server.port}/` + const timedOutSession = SessionStore.get(timedOutSessionID) + + try { + await timedOutSession.connect({ wsUrl }) + await expect( + Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const impl = yield* BrowserExecute.make(timedOutData) + return yield* impl.execute( + { + description: "Attempt command after timeout", + code: `await new Promise((resolve) => setTimeout(resolve, 50)); + return session.Runtime.evaluate({ expression: "true" });`, + timeout: 10, + }, + { sessionID: timedOutSessionID, workspaceDir: timedOutWorkspace }, + ) + }), + ), + ), + ).rejects.toThrow("browser_execute timed out; CDP session was reset") + + await Bun.sleep(80) + expect(lateCommandCount).toBe(0) + expect(SessionStore.get(timedOutSessionID)).not.toBe(timedOutSession) + await expect(timedOutSession.connect({ wsUrl })).rejects.toThrow("CDP session was reset") + } finally { + await SessionStore.evict(timedOutSessionID) + server.stop(true) + await Promise.all( + [timedOutWorkspace, timedOutData].map((d) => fs.rm(d, { recursive: true, force: true })), + ) + } +}) From 14ab148def4a87b6230425619a62971617efd04f Mon Sep 17 00:00:00 2001 From: MagMueller Date: Thu, 23 Jul 2026 19:43:52 -0700 Subject: [PATCH 2/6] fix(browser): cancel pending session work --- packages/bcode-browser/src/cdp/session.ts | 38 ++++++-- .../test/browser-execute.test.ts | 87 +++++++++++++++++++ 2 files changed, 118 insertions(+), 7 deletions(-) diff --git a/packages/bcode-browser/src/cdp/session.ts b/packages/bcode-browser/src/cdp/session.ts index 2ba605c0ad..64827b83a9 100644 --- a/packages/bcode-browser/src/cdp/session.ts +++ b/packages/bcode-browser/src/cdp/session.ts @@ -51,6 +51,8 @@ export class Session implements Transport { private reattachPromise?: Promise; private enabledDomains = new Map>(); private invalidatedError?: Error; + private openingSockets = new Set(); + private pendingWaiters = new Set<(error: Error) => void>(); private eventListeners: Array<(method: string, params: unknown, sessionId?: string) => void> = []; private callResultListeners: Array<(method: string, params: unknown, result: unknown) => void> = []; @@ -127,10 +129,12 @@ export class Session implements Transport { private openWs(wsUrl: string, timeoutMs: number): Promise { return new Promise((res, rej) => { const ws = new WebSocket(wsUrl); + this.openingSockets.add(ws); let done = false; const finish = (err?: Error) => { if (done) return; done = true; + this.openingSockets.delete(ws); clearTimeout(timer); if (err) { try { ws.close(); } catch { /* ignore */ } rej(err); } else res(); @@ -165,7 +169,7 @@ export class Session implements Transport { this.activeTargetId = undefined; this.enabledDomains.clear(); } - finish(new Error('WS closed before open (likely 403 or port closed)')); + finish(this.invalidatedError ?? new Error('WS closed before open (likely 403 or port closed)')); }); }); } @@ -197,9 +201,16 @@ export class Session implements Transport { this.enabledDomains.clear(); this.eventListeners = []; this.callResultListeners = []; - if (!ws) return; - this.rejectPending(ws, error); - try { ws.close(); } catch { /* ignore */ } + for (const reject of this.pendingWaiters) reject(error); + this.pendingWaiters.clear(); + for (const opening of this.openingSockets) { + try { opening.close(); } catch { /* ignore */ } + } + this.openingSockets.clear(); + if (ws) { + this.rejectPending(ws, error); + try { ws.close(); } catch { /* ignore */ } + } } /** @@ -251,16 +262,29 @@ export class Session implements Transport { /** Wait for the next event matching `method` (and optional predicate). */ waitFor(method: string, predicate?: (params: T) => boolean, timeoutMs = 30_000): Promise { + if (this.invalidatedError) return Promise.reject(this.invalidatedError); return new Promise((resolve, reject) => { - const timer = setTimeout(() => { + let settled = false; + let unsub = () => {}; + const fail = (error: Error) => { + if (settled) return; + settled = true; + clearTimeout(timer); unsub(); - reject(new Error(`Timeout waiting for ${method}`)); + this.pendingWaiters.delete(fail); + reject(error); + }; + const timer = setTimeout(() => { + fail(new Error(`Timeout waiting for ${method}`)); }, timeoutMs); - const unsub = this.onEvent((m, params) => { + this.pendingWaiters.add(fail); + unsub = this.onEvent((m, params) => { if (m !== method) return; if (predicate && !predicate(params as T)) return; + settled = true; clearTimeout(timer); unsub(); + this.pendingWaiters.delete(fail); resolve(params as T); }); }); diff --git a/packages/bcode-browser/test/browser-execute.test.ts b/packages/bcode-browser/test/browser-execute.test.ts index c8486c201f..66fcc3a302 100644 --- a/packages/bcode-browser/test/browser-execute.test.ts +++ b/packages/bcode-browser/test/browser-execute.test.ts @@ -328,3 +328,90 @@ test("a timed-out snippet cannot send later CDP commands or reconnect", async () ) } }) + +test("session invalidation rejects pending and future event waiters", async () => { + const waiterSessionID = "waiter-" + Math.random().toString(36).slice(2, 8) + const waiterWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-waiter-ws-")) + const waiterData = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-waiter-data-")) + const waiterSession = SessionStore.get(waiterSessionID) + + try { + await expect( + Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const impl = yield* BrowserExecute.make(waiterData) + return yield* impl.execute( + { + description: "Wait beyond tool timeout", + code: `return session.waitFor("Page.loadEventFired");`, + timeout: 10, + }, + { sessionID: waiterSessionID, workspaceDir: waiterWorkspace }, + ) + }), + ), + ), + ).rejects.toThrow("browser_execute timed out; CDP session was reset") + + expect(SessionStore.get(waiterSessionID)).not.toBe(waiterSession) + await expect(waiterSession.waitFor("Page.loadEventFired")).rejects.toThrow("CDP session was reset") + } finally { + await SessionStore.evict(waiterSessionID) + await Promise.all( + [waiterWorkspace, waiterData].map((d) => fs.rm(d, { recursive: true, force: true })), + ) + } +}) + +test("a tool timeout closes a WebSocket that is still connecting", async () => { + const connectingSessionID = "connecting-" + Math.random().toString(36).slice(2, 8) + const connectingWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-connecting-ws-")) + const connectingData = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-connecting-data-")) + let openedSockets = 0 + const server = Bun.serve({ + port: 0, + async fetch(req, bunServer) { + await Bun.sleep(50) + return bunServer.upgrade(req) ? undefined : new Response("nope", { status: 400 }) + }, + websocket: { + open() { + openedSockets++ + }, + message() {}, + close() {}, + }, + }) + if (server.port === undefined) throw new Error("test server has no port") + const wsUrl = `ws://127.0.0.1:${server.port}/` + + try { + await expect( + Effect.runPromise( + Effect.scoped( + Effect.gen(function* () { + const impl = yield* BrowserExecute.make(connectingData) + return yield* impl.execute( + { + description: "Connect beyond tool timeout", + code: `return session.connect({ wsUrl: ${JSON.stringify(wsUrl)}, timeoutMs: 1000 });`, + timeout: 10, + }, + { sessionID: connectingSessionID, workspaceDir: connectingWorkspace }, + ) + }), + ), + ), + ).rejects.toThrow("browser_execute timed out; CDP session was reset") + + await Bun.sleep(80) + expect(openedSockets).toBe(0) + } finally { + await SessionStore.evict(connectingSessionID) + server.stop(true) + await Promise.all( + [connectingWorkspace, connectingData].map((d) => fs.rm(d, { recursive: true, force: true })), + ) + } +}) From f7bd61186545caa11244e77f977ea0733a17403e Mon Sep 17 00:00:00 2001 From: MagMueller Date: Thu, 23 Jul 2026 19:58:40 -0700 Subject: [PATCH 3/6] fix(browser): stop invalidated connection attempts --- packages/bcode-browser/src/cdp/session.ts | 18 ++++++++- .../test/browser-execute.test.ts | 39 +++++++++++++++++++ 2 files changed, 56 insertions(+), 1 deletion(-) diff --git a/packages/bcode-browser/src/cdp/session.ts b/packages/bcode-browser/src/cdp/session.ts index 64827b83a9..fb2c7a7051 100644 --- a/packages/bcode-browser/src/cdp/session.ts +++ b/packages/bcode-browser/src/cdp/session.ts @@ -86,7 +86,7 @@ export class Session implements Transport { * and we connect directly to the supplied endpoint. */ async connect(opts: ConnectOptions = {}): Promise { - if (this.invalidatedError) throw this.invalidatedError; + this.throwIfInvalidated(); // No-argument connect is an ensure-connected operation. Reopening the // same configured endpoint would discard the active target session and @@ -96,15 +96,19 @@ export class Session implements Transport { const timeoutMs = opts.timeoutMs ?? 5_000; if (opts.wsUrl || opts.profileDir) { const wsUrl = await resolveWsUrl(opts, timeoutMs); + this.throwIfInvalidated(); await this.openWs(wsUrl, timeoutMs); + this.throwIfInvalidated(); return; } const envWsUrl = process.env.BU_CDP_WS ?? process.env.BU_CDP_URL; if (envWsUrl) { await this.openWs(envWsUrl, timeoutMs); + this.throwIfInvalidated(); return; } const browsers = await detectBrowsers(); + this.throwIfInvalidated(); if (browsers.length === 0) { const scanned = getBrowserCandidates().map(c => c.name).join(', '); throw new Error( @@ -113,10 +117,13 @@ export class Session implements Transport { } const errors: string[] = []; for (const b of browsers) { + this.throwIfInvalidated(); try { await this.openWs(b.wsUrl, timeoutMs); + this.throwIfInvalidated(); return; } catch (e) { + this.throwIfInvalidated(); const msg = e instanceof Error ? e.message : String(e); errors.push(` ${b.name} @ ${b.wsUrl}: ${msg}`); } @@ -127,7 +134,12 @@ export class Session implements Transport { } private openWs(wsUrl: string, timeoutMs: number): Promise { + this.throwIfInvalidated(); return new Promise((res, rej) => { + if (this.invalidatedError) { + rej(this.invalidatedError); + return; + } const ws = new WebSocket(wsUrl); this.openingSockets.add(ws); let done = false; @@ -174,6 +186,10 @@ export class Session implements Transport { }); } + private throwIfInvalidated(): void { + if (this.invalidatedError) throw this.invalidatedError; + } + isConnected(): boolean { return !this.invalidatedError && this.ws?.readyState === WebSocket.OPEN; } diff --git a/packages/bcode-browser/test/browser-execute.test.ts b/packages/bcode-browser/test/browser-execute.test.ts index 66fcc3a302..d5b29b73dc 100644 --- a/packages/bcode-browser/test/browser-execute.test.ts +++ b/packages/bcode-browser/test/browser-execute.test.ts @@ -415,3 +415,42 @@ test("a tool timeout closes a WebSocket that is still connecting", async () => { ) } }) + +test("an invalidated session cannot connect after profile resolution finishes", async () => { + const resolvingSessionID = "resolving-" + Math.random().toString(36).slice(2, 8) + const resolvingProfile = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-resolving-profile-")) + const resolvingSession = SessionStore.get(resolvingSessionID) + let openedSockets = 0 + const server = Bun.serve({ + port: 0, + fetch(req, bunServer) { + return bunServer.upgrade(req) ? undefined : new Response("nope", { status: 400 }) + }, + websocket: { + open() { + openedSockets++ + }, + message() {}, + close() {}, + }, + }) + if (server.port === undefined) throw new Error("test server has no port") + + try { + const connecting = resolvingSession.connect({ profileDir: resolvingProfile, timeoutMs: 1000 }) + await Bun.sleep(20) + resolvingSession.invalidate(new Error("CDP session was reset")) + await fs.writeFile( + path.join(resolvingProfile, "DevToolsActivePort"), + `${server.port}\n/devtools/browser/test\n`, + ) + + await expect(connecting).rejects.toThrow("CDP session was reset") + await Bun.sleep(50) + expect(openedSockets).toBe(0) + } finally { + await SessionStore.evict(resolvingSessionID) + server.stop(true) + await fs.rm(resolvingProfile, { recursive: true, force: true }) + } +}) From 9e9818630fcbd5987d3a6ea4abc2eac6eca57a28 Mon Sep 17 00:00:00 2001 From: MagMueller Date: Mon, 27 Jul 2026 23:38:33 -0700 Subject: [PATCH 4/6] fix(browser): return timeout progress --- .../skills/browser-execute/SKILL.md | 2 ++ packages/bcode-browser/src/browser-execute.ts | 23 +++++++++++++++---- .../test/browser-execute.test.ts | 16 ++++++++----- 3 files changed, 30 insertions(+), 11 deletions(-) diff --git a/packages/bcode-browser/skills/browser-execute/SKILL.md b/packages/bcode-browser/skills/browser-execute/SKILL.md index 1e6ae048b0..877e30fd24 100644 --- a/packages/bcode-browser/skills/browser-execute/SKILL.md +++ b/packages/bcode-browser/skills/browser-execute/SKILL.md @@ -174,6 +174,8 @@ console.log(JSON.stringify(titles)) ## Guardrails - Top-level `import` statements inside the snippet body are not allowed. Use `await import(...)` instead. - No CPU-bound infinite loops without `await` — they ignore the timeout. Insert `await new Promise(r => setTimeout(r, 0))` to yield. +- A `browser_execute` call defaults to 60 seconds. For intentionally longer bounded work, set the tool's top-level `timeout` in milliseconds (maximum 600000). A CDP command's inner timeout does not extend the tool timeout. +- Prefer small batches that finish well within the tool timeout. Log progress as each batch completes so a timeout can return useful partial output, then continue from the last completed batch. ## Console - `console.log`, `console.error`, `console.warn`, `console.info`, `console.debug` are all captured and streamed to the user. Treat them as your stdout. Other `console.*` methods write to bcode's stderr without being captured into the tool result. diff --git a/packages/bcode-browser/src/browser-execute.ts b/packages/bcode-browser/src/browser-execute.ts index 375e418574..4bd41446a2 100644 --- a/packages/bcode-browser/src/browser-execute.ts +++ b/packages/bcode-browser/src/browser-execute.ts @@ -49,6 +49,7 @@ import { Skills } from "./skills" const DEFAULT_TIMEOUT_MS = 60 * 1000 const MAX_TIMEOUT_MS = 10 * 60 * 1000 +const MAX_TIMEOUT_OUTPUT_LENGTH = 30_000 // Field order matters: providers stream tool-call args in schema-declared // order, so the model commits to whichever field comes first. `code` is the @@ -160,6 +161,7 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) const execute = (args: Parameters, ctx: ExecuteContext) => { const session = SessionStore.get(ctx.sessionID) + const captured = { output: "" } return Effect.gen(function* () { yield* Effect.promise(() => fs.mkdir(ctx.workspaceDir, { recursive: true })) @@ -168,10 +170,9 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) catch: (err) => new Error(`syntax error in browser_execute snippet: ${err}`), }) - let output = "" const tee = (...a: unknown[]) => { - output += a.map((x) => (typeof x === "string" ? x : serialize(x))).join(" ") + "\n" - if (ctx.onChunk) Effect.runFork(ctx.onChunk(output)) + captured.output += a.map((x) => (typeof x === "string" ? x : serialize(x))).join(" ") + "\n" + if (ctx.onChunk) Effect.runFork(ctx.onChunk(captured.output)) } // Prototype-chain to the real `console` so uncommon methods (`debug`, // `dir`, `trace`, `table`, `group`, …) don't throw when a snippet calls @@ -224,14 +225,26 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) catch: (err) => new Error(`browser_execute snippet threw: ${err instanceof Error ? err.stack ?? err.message : String(err)}`), }).pipe(Effect.ensuring(Effect.sync(() => unsubscribe()))) - return { output, result: serialize(ran), screenshots } satisfies ExecuteResult + return { output: captured.output, result: serialize(ran), screenshots } satisfies ExecuteResult }).pipe( Effect.scoped, Effect.timeoutOrElse({ duration: Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS), orElse: () => Effect.gen(function* () { - const error = new Error("browser_execute timed out; CDP session was reset") + const timeout = Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS) + const output = + captured.output.length <= MAX_TIMEOUT_OUTPUT_LENGTH + ? captured.output + : "...\n\n" + captured.output.slice(-MAX_TIMEOUT_OUTPUT_LENGTH) + const error = new Error( + [ + `browser_execute timed out after ${timeout} ms; CDP session was reset`, + output.trim() ? `Partial console output before timeout:\n${output.trimEnd()}` : "", + ] + .filter(Boolean) + .join("\n\n"), + ) yield* Effect.sync(() => SessionStore.invalidate(ctx.sessionID, session, error)) return yield* Effect.fail(error) }), diff --git a/packages/bcode-browser/test/browser-execute.test.ts b/packages/bcode-browser/test/browser-execute.test.ts index d5b29b73dc..b50eff4915 100644 --- a/packages/bcode-browser/test/browser-execute.test.ts +++ b/packages/bcode-browser/test/browser-execute.test.ts @@ -305,7 +305,8 @@ test("a timed-out snippet cannot send later CDP commands or reconnect", async () return yield* impl.execute( { description: "Attempt command after timeout", - code: `await new Promise((resolve) => setTimeout(resolve, 50)); + code: `console.log("completed batch 1"); + await new Promise((resolve) => setTimeout(resolve, 50)); return session.Runtime.evaluate({ expression: "true" });`, timeout: 10, }, @@ -314,12 +315,15 @@ test("a timed-out snippet cannot send later CDP commands or reconnect", async () }), ), ), - ).rejects.toThrow("browser_execute timed out; CDP session was reset") + ).rejects.toThrow( + "browser_execute timed out after 10 ms; CDP session was reset\n\n" + + "Partial console output before timeout:\ncompleted batch 1", + ) await Bun.sleep(80) expect(lateCommandCount).toBe(0) expect(SessionStore.get(timedOutSessionID)).not.toBe(timedOutSession) - await expect(timedOutSession.connect({ wsUrl })).rejects.toThrow("CDP session was reset") + await expect(timedOutSession.connect({ wsUrl })).rejects.toThrow("browser_execute timed out after 10 ms") } finally { await SessionStore.evict(timedOutSessionID) server.stop(true) @@ -352,10 +356,10 @@ test("session invalidation rejects pending and future event waiters", async () = }), ), ), - ).rejects.toThrow("browser_execute timed out; CDP session was reset") + ).rejects.toThrow("browser_execute timed out after 10 ms; CDP session was reset") expect(SessionStore.get(waiterSessionID)).not.toBe(waiterSession) - await expect(waiterSession.waitFor("Page.loadEventFired")).rejects.toThrow("CDP session was reset") + await expect(waiterSession.waitFor("Page.loadEventFired")).rejects.toThrow("browser_execute timed out after 10 ms") } finally { await SessionStore.evict(waiterSessionID) await Promise.all( @@ -403,7 +407,7 @@ test("a tool timeout closes a WebSocket that is still connecting", async () => { }), ), ), - ).rejects.toThrow("browser_execute timed out; CDP session was reset") + ).rejects.toThrow("browser_execute timed out after 10 ms; CDP session was reset") await Bun.sleep(80) expect(openedSockets).toBe(0) From cecef789987822dc94ad76792579498576f98482 Mon Sep 17 00:00:00 2001 From: MagMueller Date: Tue, 28 Jul 2026 00:06:55 -0700 Subject: [PATCH 5/6] fix(browser): serialize timeout recovery --- .../skills/browser-execute/SKILL.md | 3 +- packages/bcode-browser/src/browser-execute.ts | 57 +++++-- packages/bcode-browser/src/cdp/session.ts | 45 ++++- .../test/browser-execute.test.ts | 159 ++++++++++++------ 4 files changed, 194 insertions(+), 70 deletions(-) diff --git a/packages/bcode-browser/skills/browser-execute/SKILL.md b/packages/bcode-browser/skills/browser-execute/SKILL.md index 877e30fd24..9380aceef1 100644 --- a/packages/bcode-browser/skills/browser-execute/SKILL.md +++ b/packages/bcode-browser/skills/browser-execute/SKILL.md @@ -174,8 +174,7 @@ console.log(JSON.stringify(titles)) ## Guardrails - Top-level `import` statements inside the snippet body are not allowed. Use `await import(...)` instead. - No CPU-bound infinite loops without `await` — they ignore the timeout. Insert `await new Promise(r => setTimeout(r, 0))` to yield. -- A `browser_execute` call defaults to 60 seconds. For intentionally longer bounded work, set the tool's top-level `timeout` in milliseconds (maximum 600000). A CDP command's inner timeout does not extend the tool timeout. -- Prefer small batches that finish well within the tool timeout. Log progress as each batch completes so a timeout can return useful partial output, then continue from the last completed batch. +- `browser_execute` defaults to 60s (max 600s). For longer work, set the tool's top-level `timeout`; inner CDP timeouts do not extend it. Keep batches small and log progress—timeout errors return recent logs. ## Console - `console.log`, `console.error`, `console.warn`, `console.info`, `console.debug` are all captured and streamed to the user. Treat them as your stdout. Other `console.*` methods write to bcode's stderr without being captured into the tool result. diff --git a/packages/bcode-browser/src/browser-execute.ts b/packages/bcode-browser/src/browser-execute.ts index 4bd41446a2..8743569d23 100644 --- a/packages/bcode-browser/src/browser-execute.ts +++ b/packages/bcode-browser/src/browser-execute.ts @@ -43,13 +43,44 @@ import fs from "fs/promises" import path from "path" -import { Effect, Schema } from "effect" +import { Effect, Schema, Semaphore } from "effect" import { SessionStore } from "./session-store" import { Skills } from "./skills" const DEFAULT_TIMEOUT_MS = 60 * 1000 const MAX_TIMEOUT_MS = 10 * 60 * 1000 -const MAX_TIMEOUT_OUTPUT_LENGTH = 30_000 +const MAX_TIMEOUT_OUTPUT_BYTES = 8 * 1024 +const TIMEOUT_OUTPUT_TRUNCATED = "[partial console output truncated; showing final bytes]\n" + +const executionLocks = new Map() + +const serialized = (sessionID: string, effect: () => Effect.Effect) => + Effect.suspend(() => { + const existing = executionLocks.get(sessionID) + const entry = existing ?? { semaphore: Semaphore.makeUnsafe(1), users: 0 } + if (!existing) executionLocks.set(sessionID, entry) + entry.users++ + return entry.semaphore + .withPermits(1)(Effect.suspend(effect)) + .pipe( + Effect.ensuring( + Effect.sync(() => { + entry.users-- + if (entry.users === 0 && executionLocks.get(sessionID) === entry) executionLocks.delete(sessionID) + }), + ), + ) + }) + +const timeoutOutput = (output: string) => { + const bytes = Buffer.from(output, "utf8") + if (bytes.length <= MAX_TIMEOUT_OUTPUT_BYTES) return output + + const markerBytes = Buffer.byteLength(TIMEOUT_OUTPUT_TRUNCATED) + let start = bytes.length - (MAX_TIMEOUT_OUTPUT_BYTES - markerBytes) + while (start < bytes.length && (bytes[start]! & 0xc0) === 0x80) start++ + return TIMEOUT_OUTPUT_TRUNCATED + bytes.subarray(start).toString("utf8") +} // Field order matters: providers stream tool-call args in schema-declared // order, so the model commits to whichever field comes first. `code` is the @@ -159,9 +190,9 @@ const serialize = (v: unknown): string => { export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) { const skillsDir = yield* Effect.promise(() => Skills.resolveSkillsDir(dataDir)) - const execute = (args: Parameters, ctx: ExecuteContext) => { + const executeLocked = (args: Parameters, ctx: ExecuteContext) => { const session = SessionStore.get(ctx.sessionID) - const captured = { output: "" } + const captured = { active: true, output: "" } return Effect.gen(function* () { yield* Effect.promise(() => fs.mkdir(ctx.workspaceDir, { recursive: true })) @@ -171,6 +202,7 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) }) const tee = (...a: unknown[]) => { + if (!captured.active) return captured.output += a.map((x) => (typeof x === "string" ? x : serialize(x))).join(" ") + "\n" if (ctx.onChunk) Effect.runFork(ctx.onChunk(captured.output)) } @@ -193,12 +225,9 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) // when `BCODE_SCREENSHOT_DIR` is set, also written to disk for // eval-judge consumption. Two consumers of one tap. // - // Concurrency note: parallel execute() calls against the same Session - // (rare but possible — different sessionIDs share no Session, but a - // single sessionID with two in-flight tool calls would) each subscribe - // independently and would each see all screenshots produced during - // their lifetime. Acceptable for v1; opencode tool calls within one - // assistant message are serialized anyway. + // Calls sharing a sessionID are serialized before resolving their + // Session, so a timed-out call can invalidate its object without + // disrupting the next queued call. const screenshots: CollectedScreenshot[] = [] const dumpDir = process.env.BCODE_SCREENSHOT_DIR const startedAt = Date.now() @@ -233,10 +262,8 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) orElse: () => Effect.gen(function* () { const timeout = Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS) - const output = - captured.output.length <= MAX_TIMEOUT_OUTPUT_LENGTH - ? captured.output - : "...\n\n" + captured.output.slice(-MAX_TIMEOUT_OUTPUT_LENGTH) + captured.active = false + const output = timeoutOutput(captured.output) const error = new Error( [ `browser_execute timed out after ${timeout} ms; CDP session was reset`, @@ -252,6 +279,8 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) ) } + const execute = (args: Parameters, ctx: ExecuteContext) => serialized(ctx.sessionID, () => executeLocked(args, ctx)) + return { parameters, execute, skillsDir } }) diff --git a/packages/bcode-browser/src/cdp/session.ts b/packages/bcode-browser/src/cdp/session.ts index fb2c7a7051..bc23896a2c 100644 --- a/packages/bcode-browser/src/cdp/session.ts +++ b/packages/bcode-browser/src/cdp/session.ts @@ -51,6 +51,7 @@ export class Session implements Transport { private reattachPromise?: Promise; private enabledDomains = new Map>(); private invalidatedError?: Error; + private invalidation = new AbortController(); private openingSockets = new Set(); private pendingWaiters = new Set<(error: Error) => void>(); private eventListeners: Array<(method: string, params: unknown, sessionId?: string) => void> = []; @@ -95,7 +96,7 @@ export class Session implements Transport { const timeoutMs = opts.timeoutMs ?? 5_000; if (opts.wsUrl || opts.profileDir) { - const wsUrl = await resolveWsUrl(opts, timeoutMs); + const wsUrl = await resolveWsUrl(opts, timeoutMs, this.invalidation.signal); this.throwIfInvalidated(); await this.openWs(wsUrl, timeoutMs); this.throwIfInvalidated(); @@ -210,6 +211,7 @@ export class Session implements Transport { invalidate(error: Error): void { if (this.invalidatedError) return; this.invalidatedError = error; + this.invalidation.abort(error); const ws = this.ws; this.ws = undefined; this.activeSessionId = undefined; @@ -474,10 +476,14 @@ function isMissingSessionError(error: unknown): boolean { * For auto-detect, call `session.connect()` with no args — it iterates * `detectBrowsers()` and picks the first browser whose WS accepts. */ -export async function resolveWsUrl(opts: ConnectOptions, timeoutMs: number): Promise { +export async function resolveWsUrl( + opts: ConnectOptions, + timeoutMs: number, + signal?: AbortSignal, +): Promise { if (opts.wsUrl) return opts.wsUrl; if (opts.profileDir) { - const { port, path } = await readDevToolsActivePort(opts.profileDir, timeoutMs); + const { port, path } = await readDevToolsActivePort(opts.profileDir, timeoutMs, signal); return `ws://127.0.0.1:${port}${path}`; } throw new Error('resolveWsUrl needs { wsUrl } or { profileDir }. For auto-detect, call session.connect() directly.'); @@ -493,14 +499,20 @@ export async function resolveWsUrl(opts: ConnectOptions, timeoutMs: number): Pro * with a custom `--user-data-dir` (verified on macOS and Windows). For Way 2 * with modern Chrome, prefer the `/json/version` -> wsUrl route instead. */ -async function readDevToolsActivePort(profileDir: string, timeoutMs: number): Promise<{ port: number; path: string }> { +async function readDevToolsActivePort( + profileDir: string, + timeoutMs: number, + signal?: AbortSignal, +): Promise<{ port: number; path: string }> { const filePath = `${profileDir}/DevToolsActivePort`; const start = Date.now(); const deadline = start + timeoutMs; let lastErr: unknown; while (Date.now() < deadline) { + throwIfAborted(signal); try { const text = (await Bun.file(filePath).text()).trim(); + throwIfAborted(signal); const [portStr, path] = text.split('\n'); const port = Number(portStr); if (!Number.isFinite(port)) throw new Error(`malformed port line: ${portStr}`); @@ -510,8 +522,9 @@ async function readDevToolsActivePort(profileDir: string, timeoutMs: number): Pr } return { port, path }; } catch (e) { + throwIfAborted(signal); lastErr = e; - await Bun.sleep(250); + await sleepUnlessAborted(250, signal); } } const elapsed = Date.now() - start; @@ -522,6 +535,28 @@ async function readDevToolsActivePort(profileDir: string, timeoutMs: number): Pr ); } +function throwIfAborted(signal?: AbortSignal): void { + if (!signal?.aborted) return; + throw signal.reason instanceof Error ? signal.reason : new Error('CDP session was reset'); +} + +function sleepUnlessAborted(ms: number, signal?: AbortSignal): Promise { + if (!signal) return Bun.sleep(ms); + return new Promise((resolve, reject) => { + const onAbort = () => { + clearTimeout(timer); + signal.removeEventListener('abort', onAbort); + reject(signal.reason instanceof Error ? signal.reason : new Error('CDP session was reset')); + }; + const timer = setTimeout(() => { + signal.removeEventListener('abort', onAbort); + resolve(); + }, ms); + signal.addEventListener('abort', onAbort, { once: true }); + if (signal.aborted) onAbort(); + }); +} + /** * List page targets via CDP's `Target.getTargets` (works on all Chrome versions, * including those that do not serve /json). Filters out chrome:// and devtools:// diff --git a/packages/bcode-browser/test/browser-execute.test.ts b/packages/bcode-browser/test/browser-execute.test.ts index b50eff4915..a838d4e7c5 100644 --- a/packages/bcode-browser/test/browser-execute.test.ts +++ b/packages/bcode-browser/test/browser-execute.test.ts @@ -271,11 +271,11 @@ test("overlapping execute calls do not clobber each other's console capture", as ) }) -test("a timed-out snippet cannot send later CDP commands or reconnect", async () => { +test("a timed-out snippet returns progress and a queued call uses a fresh session", async () => { const timedOutSessionID = "timeout-" + Math.random().toString(36).slice(2, 8) const timedOutWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-ws-")) const timedOutData = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-data-")) - let lateCommandCount = 0 + let commandCount = 0 const server = Bun.serve({ port: 0, fetch(req, bunServer) { @@ -285,7 +285,7 @@ test("a timed-out snippet cannot send later CDP commands or reconnect", async () message(socket, raw) { const message = JSON.parse(String(raw)) if (message.method !== "Runtime.evaluate") return - lateCommandCount++ + commandCount++ socket.send(JSON.stringify({ id: message.id, result: { result: { type: "boolean", value: true } } })) }, close() {}, @@ -297,31 +297,59 @@ test("a timed-out snippet cannot send later CDP commands or reconnect", async () try { await timedOutSession.connect({ wsUrl }) - await expect( - Effect.runPromise( - Effect.scoped( - Effect.gen(function* () { - const impl = yield* BrowserExecute.make(timedOutData) - return yield* impl.execute( - { - description: "Attempt command after timeout", - code: `console.log("completed batch 1"); - await new Promise((resolve) => setTimeout(resolve, 50)); - return session.Runtime.evaluate({ expression: "true" });`, - timeout: 10, - }, - { sessionID: timedOutSessionID, workspaceDir: timedOutWorkspace }, - ) - }), - ), + const impl = await Effect.runPromise(BrowserExecute.make(timedOutData)) + let started!: () => void + const firstOutput = new Promise((resolve) => { + started = resolve + }) + const timedOut = Effect.runPromise( + impl.execute( + { + description: "Attempt command after timeout", + code: `console.log("completed batch 1"); + await new Promise((resolve) => setTimeout(resolve, 50)); + return session.Runtime.evaluate({ expression: "late" });`, + timeout: 10, + }, + { + sessionID: timedOutSessionID, + workspaceDir: timedOutWorkspace, + onChunk: () => Effect.sync(started), + }, ), - ).rejects.toThrow( - "browser_execute timed out after 10 ms; CDP session was reset\n\n" + - "Partial console output before timeout:\ncompleted batch 1", ) + await firstOutput + + const next = Effect.runPromise( + impl.execute( + { + description: "Use replacement session", + code: `await session.connect({ wsUrl: ${JSON.stringify(wsUrl)} }); + console.log("replacement connected"); + return session.Runtime.evaluate({ expression: "next" });`, + timeout: 1000, + }, + { sessionID: timedOutSessionID, workspaceDir: timedOutWorkspace }, + ), + ) + const [timedOutResult, nextResult] = await Promise.allSettled([timedOut, next]) + + expect(timedOutResult.status).toBe("rejected") + if (timedOutResult.status === "rejected") { + expect(timedOutResult.reason).toBeInstanceOf(Error) + expect(timedOutResult.reason.message).toBe( + "browser_execute timed out after 10 ms; CDP session was reset\n\n" + + "Partial console output before timeout:\ncompleted batch 1", + ) + } + expect(nextResult.status).toBe("fulfilled") + if (nextResult.status === "fulfilled") { + expect(nextResult.value.output).toBe("replacement connected\n") + expect(JSON.parse(nextResult.value.result).result.value).toBe(true) + } await Bun.sleep(80) - expect(lateCommandCount).toBe(0) + expect(commandCount).toBe(1) expect(SessionStore.get(timedOutSessionID)).not.toBe(timedOutSession) await expect(timedOutSession.connect({ wsUrl })).rejects.toThrow("browser_execute timed out after 10 ms") } finally { @@ -333,6 +361,54 @@ test("a timed-out snippet cannot send later CDP commands or reconnect", async () } }) +test("timeout output is byte-capped and capture stops after timeout", async () => { + const outputSessionID = "timeout-output-" + Math.random().toString(36).slice(2, 8) + const outputWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-output-ws-")) + const outputData = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-output-data-")) + const chunks: string[] = [] + + try { + const impl = await Effect.runPromise(BrowserExecute.make(outputData)) + let failure: unknown + try { + await Effect.runPromise( + impl.execute( + { + description: "Cap timeout output", + code: `console.log("😀".repeat(6000)); + await new Promise((resolve) => setTimeout(resolve, 30)); + console.log("after timeout");`, + timeout: 10, + }, + { + sessionID: outputSessionID, + workspaceDir: outputWorkspace, + onChunk: (output) => Effect.sync(() => chunks.push(output)), + }, + ), + ) + } catch (error) { + failure = error + } + + expect(failure).toBeInstanceOf(Error) + const message = (failure as Error).message + const partial = message.split("Partial console output before timeout:\n")[1]! + expect(partial).toStartWith("[partial console output truncated; showing final bytes]\n") + expect(Buffer.byteLength(partial, "utf8")).toBeLessThanOrEqual(8 * 1024) + expect(partial).not.toContain("\uFFFD") + + await Bun.sleep(60) + expect(chunks).toHaveLength(1) + expect(chunks[0]).not.toContain("after timeout") + } finally { + await SessionStore.evict(outputSessionID) + await Promise.all( + [outputWorkspace, outputData].map((d) => fs.rm(d, { recursive: true, force: true })), + ) + } +}) + test("session invalidation rejects pending and future event waiters", async () => { const waiterSessionID = "waiter-" + Math.random().toString(36).slice(2, 8) const waiterWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-waiter-ws-")) @@ -420,41 +496,26 @@ test("a tool timeout closes a WebSocket that is still connecting", async () => { } }) -test("an invalidated session cannot connect after profile resolution finishes", async () => { +test("profile resolution stops promptly when the session is invalidated", async () => { const resolvingSessionID = "resolving-" + Math.random().toString(36).slice(2, 8) const resolvingProfile = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-resolving-profile-")) const resolvingSession = SessionStore.get(resolvingSessionID) - let openedSockets = 0 - const server = Bun.serve({ - port: 0, - fetch(req, bunServer) { - return bunServer.upgrade(req) ? undefined : new Response("nope", { status: 400 }) - }, - websocket: { - open() { - openedSockets++ - }, - message() {}, - close() {}, - }, - }) - if (server.port === undefined) throw new Error("test server has no port") try { - const connecting = resolvingSession.connect({ profileDir: resolvingProfile, timeoutMs: 1000 }) + const connecting = resolvingSession.connect({ profileDir: resolvingProfile, timeoutMs: 5000 }) await Bun.sleep(20) resolvingSession.invalidate(new Error("CDP session was reset")) - await fs.writeFile( - path.join(resolvingProfile, "DevToolsActivePort"), - `${server.port}\n/devtools/browser/test\n`, - ) - await expect(connecting).rejects.toThrow("CDP session was reset") - await Bun.sleep(50) - expect(openedSockets).toBe(0) + await expect( + Promise.race([ + connecting, + Bun.sleep(200).then(() => { + throw new Error("profile resolution did not stop promptly") + }), + ]), + ).rejects.toThrow("CDP session was reset") } finally { await SessionStore.evict(resolvingSessionID) - server.stop(true) await fs.rm(resolvingProfile, { recursive: true, force: true }) } }) From 1a406a10c31c511be91b791474dcedf859aa6698 Mon Sep 17 00:00:00 2001 From: MagMueller Date: Tue, 28 Jul 2026 09:07:22 -0700 Subject: [PATCH 6/6] fix(browser): include queue time in timeout --- packages/bcode-browser/src/browser-execute.ts | 49 +++++++++++---- .../test/browser-execute.test.ts | 62 +++++++++++++++++++ 2 files changed, 98 insertions(+), 13 deletions(-) diff --git a/packages/bcode-browser/src/browser-execute.ts b/packages/bcode-browser/src/browser-execute.ts index 8743569d23..0d8216da51 100644 --- a/packages/bcode-browser/src/browser-execute.ts +++ b/packages/bcode-browser/src/browser-execute.ts @@ -54,22 +54,43 @@ const TIMEOUT_OUTPUT_TRUNCATED = "[partial console output truncated; showing fin const executionLocks = new Map() -const serialized = (sessionID: string, effect: () => Effect.Effect) => +const serialized = ( + sessionID: string, + timeout: number, + effect: (remaining: number) => Effect.Effect, +) => Effect.suspend(() => { const existing = executionLocks.get(sessionID) const entry = existing ?? { semaphore: Semaphore.makeUnsafe(1), users: 0 } if (!existing) executionLocks.set(sessionID, entry) entry.users++ - return entry.semaphore - .withPermits(1)(Effect.suspend(effect)) - .pipe( - Effect.ensuring( - Effect.sync(() => { - entry.users-- - if (entry.users === 0 && executionLocks.get(sessionID) === entry) executionLocks.delete(sessionID) + const startedAt = Date.now() + const waitTimeout = () => + Effect.fail(new Error(`browser_execute timed out after ${timeout} ms waiting for a previous call; no code was run`)) + return Effect.uninterruptibleMask((restore) => + restore( + entry.semaphore.take(1).pipe( + Effect.timeoutOrElse({ + duration: timeout, + orElse: waitTimeout, }), ), + ).pipe( + Effect.flatMap(() => { + const remaining = timeout - (Date.now() - startedAt) + return restore(remaining <= 0 ? waitTimeout() : Effect.suspend(() => effect(remaining))).pipe( + Effect.ensuring(entry.semaphore.release(1)), + ) + }), + ), + ).pipe( + Effect.ensuring( + Effect.sync(() => { + entry.users-- + if (entry.users === 0 && executionLocks.get(sessionID) === entry) executionLocks.delete(sessionID) + }), ) + ) }) const timeoutOutput = (output: string) => { @@ -190,7 +211,7 @@ const serialize = (v: unknown): string => { export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) { const skillsDir = yield* Effect.promise(() => Skills.resolveSkillsDir(dataDir)) - const executeLocked = (args: Parameters, ctx: ExecuteContext) => { + const executeLocked = (args: Parameters, ctx: ExecuteContext, timeout: number, configuredTimeout: number) => { const session = SessionStore.get(ctx.sessionID) const captured = { active: true, output: "" } return Effect.gen(function* () { @@ -258,15 +279,14 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) }).pipe( Effect.scoped, Effect.timeoutOrElse({ - duration: Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS), + duration: timeout, orElse: () => Effect.gen(function* () { - const timeout = Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS) captured.active = false const output = timeoutOutput(captured.output) const error = new Error( [ - `browser_execute timed out after ${timeout} ms; CDP session was reset`, + `browser_execute timed out after ${configuredTimeout} ms; CDP session was reset`, output.trim() ? `Partial console output before timeout:\n${output.trimEnd()}` : "", ] .filter(Boolean) @@ -279,7 +299,10 @@ export const make = Effect.fn("BrowserExecute.make")(function* (dataDir: string) ) } - const execute = (args: Parameters, ctx: ExecuteContext) => serialized(ctx.sessionID, () => executeLocked(args, ctx)) + const execute = (args: Parameters, ctx: ExecuteContext) => { + const timeout = Math.max(1, Math.min(args.timeout ?? DEFAULT_TIMEOUT_MS, MAX_TIMEOUT_MS)) + return serialized(ctx.sessionID, timeout, (remaining) => executeLocked(args, ctx, remaining, timeout)) + } return { parameters, execute, skillsDir } }) diff --git a/packages/bcode-browser/test/browser-execute.test.ts b/packages/bcode-browser/test/browser-execute.test.ts index a838d4e7c5..9d30118148 100644 --- a/packages/bcode-browser/test/browser-execute.test.ts +++ b/packages/bcode-browser/test/browser-execute.test.ts @@ -361,6 +361,68 @@ test("a timed-out snippet returns progress and a queued call uses a fresh sessio } }) +test("a queued call's timeout includes time spent waiting for the same session", async () => { + const sessionID = "queue-timeout-" + Math.random().toString(36).slice(2, 8) + const workspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-queue-timeout-ws-")) + const data = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-queue-timeout-data-")) + const session = SessionStore.get(sessionID) + + try { + const impl = await Effect.runPromise(BrowserExecute.make(data)) + let started!: () => void + const firstOutput = new Promise((resolve) => { + started = resolve + }) + const first = Effect.runPromise( + impl.execute( + { + description: "Hold session lock", + code: `console.log("started"); + await new Promise((resolve) => setTimeout(resolve, 50)); + return "first";`, + timeout: 1000, + }, + { + sessionID, + workspaceDir: workspace, + onChunk: () => Effect.sync(started), + }, + ), + ) + await firstOutput + + await expect( + Effect.runPromise( + impl.execute( + { + description: "Time out while queued", + code: `throw new Error("queued code should not run");`, + timeout: 10, + }, + { sessionID, workspaceDir: workspace }, + ), + ), + ).rejects.toThrow("browser_execute timed out after 10 ms waiting for a previous call; no code was run") + + expect(JSON.parse((await first).result)).toBe("first") + expect(SessionStore.get(sessionID)).toBe(session) + const next = await Effect.runPromise( + impl.execute( + { + description: "Run after queue timeout", + code: `return "next";`, + timeout: 1000, + }, + { sessionID, workspaceDir: workspace }, + ), + ) + expect(JSON.parse(next.result)).toBe("next") + } finally { + await SessionStore.evict(sessionID) + await Promise.all([workspace, data].map((d) => fs.rm(d, { recursive: true, force: true }))) + } +}) + test("timeout output is byte-capped and capture stops after timeout", async () => { const outputSessionID = "timeout-output-" + Math.random().toString(36).slice(2, 8) const outputWorkspace = await fs.mkdtemp(path.join(os.tmpdir(), "bcode-timeout-output-ws-"))