Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
127 changes: 124 additions & 3 deletions apps/sim/lib/execution/remote-sandbox/e2b-session.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ const {
create,
list,
connect,
pause,
nextItems,
getInfo,
writeFile,
Expand All @@ -22,6 +23,7 @@ const {
create: vi.fn(),
list: vi.fn(),
connect: vi.fn(),
pause: vi.fn(),
nextItems: vi.fn(),
getInfo: vi.fn(),
writeFile: vi.fn(),
Expand All @@ -36,22 +38,27 @@ const {
NotFoundError: class NotFoundError extends Error {},
}))
vi.mock('@e2b/code-interpreter', () => ({
Sandbox: { create, list, connect, getInfo },
Sandbox: { create, list, connect, pause, getInfo },
NotFoundError,
}))
vi.mock('@/lib/execution/remote-sandbox/session-lock', () => ({
withSandboxSessionLock: sessionLock,
}))

import { e2bProvider, stopE2BSessionProcess } from '@/lib/execution/remote-sandbox/e2b'
import {
E2B_MAX_SANDBOX_LIFETIME_MS,
e2bProvider,
stopE2BSessionProcess,
} from '@/lib/execution/remote-sandbox/e2b'
import { observeSandboxExecution } from '@/lib/execution/remote-sandbox/execution-observer'

setEnv({ E2B_API_KEY: 'test', MOTHERSHIP_E2B_TEMPLATE_ID: 'mothership-template' })

function candidate(sandboxId: string, time: number) {
return {
sandboxId,
startedAt: new Date(time),
state: 'running',
startedAt: new Date(Date.now() - 3_600_000 + time),
endAt: new Date(Date.now() + 3_600_000),
metadata: { simSessionKey: 'chat', simSessionOwnership: 'tracked-v1' },
}
Expand Down Expand Up @@ -681,3 +688,117 @@ describe('E2B session lease', () => {
expect(plane.requests).toBe(3)
})
})

describe('E2B continuous-runtime cap', () => {
const CAP_MS = E2B_MAX_SANDBOX_LIFETIME_MS
const LEASE_MS = 21 * 60_000

/** Models E2B: every timeout is clamped to the cap, and resuming a paused sandbox restarts it. */
function cappedPlane(runningForMs: number) {
const plane = { startedAtMs: Date.now() - runningForMs, endAtMs: 0, paused: false }
const clampedEnd = (timeoutMs: number) =>
Math.min(Date.now() + timeoutMs, plane.startedAtMs + CAP_MS)
plane.endAtMs = clampedEnd(5 * 60_000)
list.mockReturnValue({ nextItems, hasNext: false })
nextItems.mockImplementation(async () => [
{
...candidate('retained', 0),
state: plane.paused ? 'paused' : 'running',
startedAt: new Date(plane.startedAtMs),
endAt: new Date(plane.endAtMs),
},
])
pause.mockImplementation(async () => {
plane.paused = true
return true
})
const sandbox = {
sandboxId: 'retained',
getInfo: async () => ({ endAt: new Date(plane.endAtMs) }),
setTimeout: async (timeoutMs: number) => {
plane.endAtMs = clampedEnd(timeoutMs)
},
commands: {
run: async (_command: string, options: { background?: boolean }) => {
if (options.background === false) {
return { stdout: '{"settled": true}', stderr: '', exitCode: 0 }
}
throw Object.assign(new Error('Sandbox timeout: end of life'), { name: 'TimeoutError' })
},
},
}
connect.mockImplementation(async (_id: string, options: { timeoutMs: number }) => {
if (plane.paused) {
plane.paused = false
plane.startedAtMs = Date.now()
}
plane.endAtMs = clampedEnd(options.timeoutMs)
return sandbox
})
return plane
}

it('never records a reused lease past the cap when the runtime cannot be reset', async () => {
const plane = cappedPlane(CAP_MS - 10 * 60_000)
pause.mockRejectedValueOnce(new Error('pause unavailable'))
const sandbox = await e2bProvider.findSessionSandbox?.('chat', { lifetimeMs: LEASE_MS })
const now = Date.now()
expect(plane.endAtMs).toBeLessThan(now + LEASE_MS)
expect(sandbox?.outlives?.(LEASE_MS, now)).toBe(false)
expect(sandbox?.outlives?.(plane.endAtMs - now - 1000, now)).toBe(true)
await sandbox?.extendLifetime?.(LEASE_MS)
expect(sandbox?.outlives?.(LEASE_MS, Date.now())).toBe(false)
})

it('pauses and resumes a workbench whose lease would outrun the cap', async () => {
const plane = cappedPlane(CAP_MS - 10 * 60_000)
const requestedAtMs = Date.now()
const sandbox = await e2bProvider.findSessionSandbox?.('chat', { lifetimeMs: LEASE_MS })
expect(plane.startedAtMs).toBeGreaterThanOrEqual(requestedAtMs)
expect(plane.endAtMs).toBeGreaterThanOrEqual(requestedAtMs + LEASE_MS)
expect(sandbox?.outlives?.(LEASE_MS, requestedAtMs)).toBe(true)
})

it('reports reaching the cap as the provider limit, not a user timeout', async () => {
cappedPlane(CAP_MS - 30_000)
pause.mockRejectedValueOnce(new Error('pause unavailable'))
const sandbox = await e2bProvider.findSessionSandbox?.('chat', { lifetimeMs: LEASE_MS })
const result = await sandbox?.runCommand('long job', { timeoutMs: LEASE_MS })
expect(result?.providerFailure).toBe('provider_limit')
expect(result?.timedOut).toBeUndefined()
})

it('measures the command deadline from dispatch, after a slow ownership write', async () => {
cappedPlane(CAP_MS - 30_000)
pause.mockRejectedValueOnce(new Error('pause unavailable'))
const sandbox = await e2bProvider.findSessionSandbox?.('chat', { lifetimeMs: LEASE_MS })
if (!sandbox) throw new Error('Missing sandbox')
const realNow = Date.now
let ownershipWriteMs = 0
const clock = vi.spyOn(Date, 'now').mockImplementation(() => realNow() + ownershipWriteMs)
try {
const result = await observeSandboxExecution(
{
hold: vi.fn(),
unsettled: vi.fn(),
claimProcess: async () => {
ownershipWriteMs = 2_000
},
},
() => sandbox.runCommand('long job', { timeoutMs: 29_000 })
)
expect(result.providerFailure).toBe('provider_limit')
} finally {
clock.mockRestore()
}
})

it('keeps a command that reaches its own timeout near the cap a user timeout', async () => {
cappedPlane(CAP_MS - 50_000)
pause.mockRejectedValueOnce(new Error('pause unavailable'))
const sandbox = await e2bProvider.findSessionSandbox?.('chat', { lifetimeMs: LEASE_MS })
const result = await sandbox?.runCommand('short job', { timeoutMs: 1000 })
expect(result?.providerFailure).toBeUndefined()
expect(result?.timedOut).toBe(true)
})
})
77 changes: 58 additions & 19 deletions apps/sim/lib/execution/remote-sandbox/e2b.ts
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,10 @@ const MATERIALIZER_REVISION_RADIX = 1000
/** Public provider alias used by builder/task tests and release tooling. */
export const E2B_SANDBOX_MATERIALIZER_REVISION = FUNCTION_SANDBOX_MATERIALIZER_REVISION

/** Maximum continuous sandbox lifetime supported by E2B. */
/**
* E2B's continuous-runtime cap on the Pro tier: a running sandbox is killed this long after it
* last started or resumed, however long its timeout says, and a pause plus resume resets the count.
*/
export const E2B_MAX_SANDBOX_LIFETIME_MS = 24 * 60 * 60 * 1000

/** E2B sends sandbox lifetimes as whole seconds. */
Expand Down Expand Up @@ -164,14 +167,17 @@ function prepareE2BCommand(
}
}

/** Only a command whose own timeout would outlive the cap can have been ended by it. */
function reachedE2BProviderLimit(
error: unknown,
providerLimitAtMs: number | undefined,
providerLimitAtMs: number,
commandDeadlineAtMs: number | undefined,
signal?: AbortSignal
): boolean {
return (
!signal?.aborted &&
providerLimitAtMs !== undefined &&
commandDeadlineAtMs !== undefined &&
commandDeadlineAtMs >= providerLimitAtMs &&
Date.now() >= providerLimitAtMs - E2B_PROVIDER_LIMIT_CLASSIFICATION_WINDOW_MS &&
isE2BExecutionTimeout(error)
)
Expand Down Expand Up @@ -358,19 +364,30 @@ export async function stopE2BSessionProcess(
class E2BSandboxHandle implements SandboxHandle {
private killed = false
private killPromise: Promise<void> | null = null
private sessionDeadlineAtMs?: number

/**
* @param providerLimitAtMs No later than E2B's continuous-runtime cap for this sandbox, which no
* timeout request can extend past.
* @param sessionDeadlineAtMs Earliest time the provider can reap this session sandbox, as set
* by this handle's own create, connect, or timeout request. Session deadlines only ever move
* later — every update path extends and none shortens — so it stays a valid lower bound.
* by this handle's own create, connect, or timeout request and never past the cap. Session
* deadlines only ever move later — every update path extends and none shortens — so it stays
* a valid lower bound.
*/
constructor(
private readonly sandbox: E2BSandbox,
private readonly language: CodeLanguage,
private readonly providerLimitAtMs?: number,
private readonly providerLimitAtMs: number,
private readonly sessionKey?: string,
private sessionDeadlineAtMs?: number
) {}
sessionDeadlineAtMs?: number
) {
if (sessionDeadlineAtMs !== undefined) this.recordSessionDeadline(sessionDeadlineAtMs)
}

/** E2B clamps every timeout to the cap, so no granted lease reaches past it. */
private recordSessionDeadline(atMs: number): void {
this.sessionDeadlineAtMs = Math.min(atMs, this.providerLimitAtMs)
}

get sandboxId(): string {
return this.sandbox.sandboxId
Expand All @@ -393,7 +410,7 @@ class E2BSandboxHandle implements SandboxHandle {
}
const requestedAtMs = Date.now()
await this.sandbox.setTimeout(timeoutMs)
if (this.sessionKey !== undefined) this.sessionDeadlineAtMs = requestedAtMs + timeoutMs
if (this.sessionKey !== undefined) this.recordSessionDeadline(requestedAtMs + timeoutMs)
}

async runCode(
Expand Down Expand Up @@ -484,6 +501,7 @@ class E2BSandboxHandle implements SandboxHandle {
operation: 'code' | 'command'
): Promise<SandboxCommandResult> {
if (this.sessionKey !== undefined) options.signal?.throwIfAborted()
let commandDeadlineAtMs: number | undefined
const outputBudget = new SandboxProcessOutputBudget(
options.maxOutputBytes ?? MAX_SANDBOX_PROCESS_OUTPUT_BYTES
)
Expand Down Expand Up @@ -553,6 +571,8 @@ class E2BSandboxHandle implements SandboxHandle {
}
let started: Awaited<ReturnType<E2BSandbox['commands']['run']>>
try {
/** E2B starts the process timeout at dispatch, after the ownership write above. */
commandDeadlineAtMs = Date.now() + processOptions.timeoutMs
started = await this.sandbox.commands.run(
processId
? sessionProcessCommand(processId, prepared.command, options.rootUser)
Expand Down Expand Up @@ -662,7 +682,9 @@ class E2BSandboxHandle implements SandboxHandle {
if (outputBudget.error) throw outputBudget.error
if (isSandboxOutputLimitError(error)) throw error
if (isNonRetryableExecutionError(error)) throw error
if (reachedE2BProviderLimit(error, this.providerLimitAtMs, options.signal)) {
if (
reachedE2BProviderLimit(error, this.providerLimitAtMs, commandDeadlineAtMs, options.signal)
) {
recordSandboxProviderLimit({ provider: 'e2b', operation })
return {
stdout: '',
Expand Down Expand Up @@ -1085,9 +1107,7 @@ export const e2bProvider: SandboxProvider = {
return new E2BSandboxHandle(
sandbox,
options?.language ?? CodeLanguage.Python,
effectiveLifetimeMs === E2B_MAX_SANDBOX_LIFETIME_MS
? lifetimeStartedAtMs + E2B_MAX_SANDBOX_LIFETIME_MS
: undefined,
lifetimeStartedAtMs + E2B_MAX_SANDBOX_LIFETIME_MS,
options?.sessionKey,
options?.sessionKey && effectiveLifetimeMs !== undefined
? lifetimeStartedAtMs + effectiveLifetimeMs
Expand Down Expand Up @@ -1116,19 +1136,38 @@ export const e2bProvider: SandboxProvider = {
'This workbench predates durable execution ownership and requires recovery before reuse'
)
}
const leaseMs = options.lifetimeMs === undefined ? 0 : e2bTimeoutMs(options.lifetimeMs)
Comment thread
waleedlatif1 marked this conversation as resolved.
let resuming = candidate.state === 'paused'
// E2B silently clamps any timeout to its continuous-runtime cap. A pause plus resume restarts
// that count with memory, files, and processes intact, where a new workbench would lose them.
if (
!resuming &&
Date.now() + leaseMs > candidate.startedAt.getTime() + E2B_MAX_SANDBOX_LIFETIME_MS
) {
try {
await Sandbox.pause(candidate.sandboxId, { apiKey })
resuming = true
} catch (error) {
logger.warn(
'Failed to pause workbench to reset its runtime cap; reusing it until the cap',
{
sandboxId: candidate.sandboxId,
error: getErrorMessage(error),
}
)
}
}
// Connect also sets a timeout, including for running sandboxes. Preserve the active deadline,
// and grant the requested lease in the same request instead of a later getInfo + setTimeout.
const requestedAtMs = Date.now()
const timeoutMs = Math.max(
5 * 60_000,
candidate.endAt.getTime() - requestedAtMs,
options.lifetimeMs === undefined ? 0 : e2bTimeoutMs(options.lifetimeMs)
)
const timeoutMs = Math.max(5 * 60_000, candidate.endAt.getTime() - requestedAtMs, leaseMs)
const providerLimitAtMs =
(resuming ? requestedAtMs : candidate.startedAt.getTime()) + E2B_MAX_SANDBOX_LIFETIME_MS
Comment thread
waleedlatif1 marked this conversation as resolved.
const sandbox = await Sandbox.connect(candidate.sandboxId, { apiKey, timeoutMs })
return new E2BSandboxHandle(
sandbox,
options.language ?? CodeLanguage.Python,
undefined,
providerLimitAtMs,
key,
options.lifetimeMs === undefined ? undefined : requestedAtMs + timeoutMs
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,29 @@ it('bounds the remote copy itself, not just the eventual response stream', async
).rejects.toThrow('size limit')
expect(mocks.read).not.toHaveBeenCalled()
})
it('asks the reconnect for the idle lease, which can outlast the runtime cap', async () => {
const IDLE_MS = 20 * 60_000
const capAtMs = Date.now() + 60_000
let leaseEndAtMs = capAtMs
const machine = await mocks.find()
mocks.find.mockImplementation(async (_key: string, options: { lifetimeMs?: number }) => {
if (options.lifetimeMs !== undefined) leaseEndAtMs = Date.now() + options.lifetimeMs
return {
...machine,
extendLifetime: async (lifetimeMs: number) => {
leaseEndAtMs = Math.min(Date.now() + lifetimeMs, capAtMs)
},
}
})
const path = join(directory, 'source.txt')
await writeFile(path, 'original')
const snapshot = await openSessionFileSnapshot('chat', path, undefined, undefined, {
allowedRoots: [directory],
maxBytes: 100,
})
await snapshot.dispose()
expect(leaseEndAtMs).toBeGreaterThanOrEqual(Date.now() + IDLE_MS - 1000)
})
it('never creates a replacement machine for a missing read', async () => {
mocks.find.mockResolvedValue(null)
await expect(openSessionFileSnapshot('chat', '/tmp/missing')).rejects.toThrow(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,10 @@ export async function openSessionFileSnapshot(
try {
return await withSandboxSessionLock(sessionKey, signal, async (accessSignal) => {
const provider = resolveProvider()
const sandbox = await provider.findSessionSandbox?.(sessionKey, {})
// The reconnect grants the idle lease itself, so it can reset a workbench near its runtime cap.
const sandbox = await provider.findSessionSandbox?.(sessionKey, {
lifetimeMs: SESSION_SANDBOX_IDLE_MS,
})
accessSignal.throwIfAborted()
if (!sandbox) throw new Error(`No workbench exists for this chat; write "${path}" first.`)
const readStream = sandbox.readFileStream?.bind(sandbox)
Expand All @@ -84,7 +87,6 @@ export async function openSessionFileSnapshot(
return disposed
}
try {
await sandbox.extendLifetime?.(SESSION_SANDBOX_IDLE_MS)
accessSignal.throwIfAborted()
const copied = await sandbox.runCommand(COPY_FILE, {
envs: {
Expand Down
17 changes: 17 additions & 0 deletions apps/sim/lib/execution/remote-sandbox/session-files.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,23 @@ describe('workbench file cancellation', () => {
})
})

it('asks the reconnect for the idle lease, which can outlast the runtime cap', async () => {
const IDLE_MS = 20 * 60_000
const capAtMs = Date.now() + 60_000
let leaseEndAtMs = capAtMs
find.mockImplementation(async (_key: string, options: { lifetimeMs?: number }) => {
if (options.lifetimeMs !== undefined) leaseEndAtMs = Date.now() + options.lifetimeMs
return {
readFileWithLimit: read,
extendLifetime: async (lifetimeMs: number) => {
leaseEndAtMs = Math.min(Date.now() + lifetimeMs, capAtMs)
},
}
})
expect(await readSessionSandboxFile('chat', 'input.txt')).toMatchObject({ outcome: 'read' })
expect(leaseEndAtMs).toBeGreaterThanOrEqual(Date.now() + IDLE_MS - 1000)
})

it('does not write if Stop arrives during the sandbox lookup', async () => {
const controller = new AbortController()
find.mockImplementation(async () => {
Expand Down
Loading
Loading