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
2 changes: 2 additions & 0 deletions packages/code/src/protocol.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@ export const BRIDGE_WORKSPACE_COMMAND_MAX_TIMEOUT_MS = 5 * 60_000;
export const BRIDGE_WORKSPACE_COMMAND_DEFAULT_OUTPUT_BYTES = 256 * 1024;
export const BRIDGE_WORKSPACE_COMMAND_MAX_OUTPUT_BYTES = 1024 * 1024;
export const BRIDGE_WORKSPACE_COMMAND_SIGNAL_MAX_LENGTH = 32;
/** How long Code API drains a clean rejection after Stop cancels a workspace mutation. */
export const BRIDGE_CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS = 5_000;

export type BridgeProtocolVersion = typeof BRIDGE_PROTOCOL_VERSION;

Expand Down
26 changes: 21 additions & 5 deletions packages/code/src/worker.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { randomBytes } from 'node:crypto';

import {
BRIDGE_CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS,
BRIDGE_PROTOCOL_VERSION,
BridgeProtocolError,
bridgeWorkerPath,
Expand Down Expand Up @@ -1619,23 +1620,38 @@ export class BridgeWorker {
assignment.runtimeSessionId != null &&
settlement.status === 'rejected' &&
(!sandboxStarted || sandboxRejectedExecution);
if (knownCleanStatefulRejection) {
// An armed mutation reaches settlement as rejected only after an atomic
// failure that does not require quarantine. Code API accepts that
// rejection after expiry and drains it for its own grace after Stop, so a
// Stop near the deadline must not cut off retries at the deadline.
const knownCleanWorkspaceRejection =
workspaceMutationArmed && settlement.status === 'rejected';
if (knownCleanStatefulRejection || knownCleanWorkspaceRejection) {
heartbeatController.abort();
await heartbeat;
const recoveryHeartbeatController = new AbortController();
const recoveryHeartbeat = this.maintainRegistration(
recoveryHeartbeatController.signal,
true,
).catch(() => undefined);
const rejectionAckGraceMs = Math.max(
0,
this.options.rejectionAckGraceMs ?? REJECTION_ACK_GRACE_MS,
);
try {
await this.settleWithRetry(
assignment,
settlement,
localDeadlineAtMs +
Math.max(
0,
this.options.rejectionAckGraceMs ?? REJECTION_ACK_GRACE_MS,
),
(knownCleanStatefulRejection
? rejectionAckGraceMs
: Math.max(
rejectionAckGraceMs,
BRIDGE_CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS,
)),
// Stateful rejections outlive shutdown; workspace guards still
// fail closed when the worker itself stops.
knownCleanStatefulRejection ? undefined : signal,
);
} finally {
recoveryHeartbeatController.abort();
Expand Down
202 changes: 202 additions & 0 deletions packages/code/src/workspace-worker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1480,6 +1480,208 @@ test('worker clears quarantine after a command cancellation confirms process ter
assert.deepEqual(lifecycle, ['arm', 'execute', 'settle', 'clear']);
});

test('worker retries a clean Stop rejection near its deadline through the cancellation grace', async () => {
const lifecycle: string[] = [];
const settlements: Array<Record<string, unknown>> = [];
const remainingMs = 100;
const startedAt = Date.now();
const baseCapabilities = {
protocolVersion: 1 as const,
operations: ['read_file' as const],
workspaces: [{ id: 'primary', operations: ['read_file' as const] }],
};
const workspaceTools = new SandboxWorkspaceTools({
workspaceTools: {
capabilities: baseCapabilities,
mutationFailuresAreAtomic: true,
async execute() { throw new Error('base executor must not run'); },
},
commandWorkspaces: ['primary'],
commandSandbox: {
mutationFailuresAreAtomic: true,
async execute(_request, signal) {
lifecycle.push('execute');
await new Promise<void>((resolve) => {
if (signal?.aborted) return resolve();
signal?.addEventListener('abort', () => resolve(), { once: true });
});
lifecycle.push('stop');
// Process-group termination is confirmed after the original deadline.
await new Promise((resolve) => setTimeout(resolve, remainingMs));
throw new WorkspaceToolError(
'Workspace command execution aborted',
'EXECUTION_ABORTED',
true,
false,
);
},
},
});
const worker = new BridgeWorker({
codeApiUrl: 'https://code.example/v1',
token: 'worker-secret',
workerId: 'vm-1',
incarnationId,
sandboxEndpoint: 'http://127.0.0.1:2000/api/v2',
capabilities: {
statefulWorkspace: true,
sandboxProfile: 'nsjail',
runtimes: ['bash'],
workspaceTools: workspaceTools.capabilities,
},
workspaceTools,
workspaceMutationQuarantine: mutationQuarantine(
() => lifecycle.push('quarantine'),
() => lifecycle.push('arm'),
() => lifecycle.push('clear'),
),
cancellationPollIntervalMs: 5,
fetchImpl: async (input, init) => {
if (String(input).endsWith('/cancellation')) {
return Response.json({
protocolVersion: 1,
cancelled: Date.now() >= startedAt + remainingMs / 2,
});
}
if (!String(input).endsWith('/settle')) {
return Response.json({ protocolVersion: 1, accepted: true });
}
lifecycle.push('settle');
settlements.push({
...(JSON.parse(String(init?.body)) as Record<string, unknown>),
attemptedAt: Date.now(),
});
// Settlement delivery takes a real transport turn and honors its deadline.
await new Promise<void>((resolve, reject) => {
const timer = setTimeout(resolve, 10);
init?.signal?.addEventListener(
'abort',
() => {
clearTimeout(timer);
reject(new DOMException('aborted', 'AbortError'));
},
{ once: true },
);
});
if (settlements.length === 1) {
return Response.json(
{ error: 'Bridge settlement temporarily unavailable' },
{ status: 503 },
);
}
return Response.json({ protocolVersion: 1, accepted: true });
},
});

await worker.executeAndSettle({
protocolVersion: 1,
assignmentId: 'assignment-command-stopped-near-deadline',
workerId: 'vm-1',
incarnationId,
generation: 4,
leaseToken: 'lease-token-that-is-long-enough-for-testing',
expiresAt: new Date(startedAt + remainingMs).toISOString(),
remainingMs,
executionKind: 'workspace_tool',
request: {
protocolVersion: 1,
operation: 'execute_command',
workspaceId: 'primary',
command: 'sleep 30; touch delayed.txt',
},
});

assert.deepEqual(lifecycle, [
'arm',
'execute',
'stop',
'settle',
'settle',
'clear',
]);
assert.ok(Number(settlements[0]?.attemptedAt) > startedAt + remainingMs);
assert.equal(settlements[1]?.status, 'rejected');
assert.equal(settlements[1]?.errorCode, 'EXECUTION_ABORTED');
});

test('worker keeps quarantine armed when shutdown interrupts a clean command rejection', async () => {
const lifecycle: string[] = [];
const controller = new AbortController();
const baseCapabilities = {
protocolVersion: 1 as const,
operations: ['read_file' as const],
workspaces: [{ id: 'primary', operations: ['read_file' as const] }],
};
const workspaceTools = new SandboxWorkspaceTools({
workspaceTools: {
capabilities: baseCapabilities,
mutationFailuresAreAtomic: true,
async execute() { throw new Error('base executor must not run'); },
},
commandWorkspaces: ['primary'],
commandSandbox: {
mutationFailuresAreAtomic: true,
async execute() {
lifecycle.push('execute');
controller.abort(new Error('shutdown'));
throw new WorkspaceToolError(
'Workspace command execution aborted',
'EXECUTION_ABORTED',
true,
false,
);
},
},
});
const worker = new BridgeWorker({
codeApiUrl: 'https://code.example/v1',
token: 'worker-secret',
workerId: 'vm-1',
incarnationId,
sandboxEndpoint: 'http://127.0.0.1:2000/api/v2',
capabilities: {
statefulWorkspace: true,
sandboxProfile: 'nsjail',
runtimes: ['bash'],
workspaceTools: workspaceTools.capabilities,
},
workspaceTools,
workspaceMutationQuarantine: mutationQuarantine(
() => lifecycle.push('quarantine'),
() => lifecycle.push('arm'),
() => lifecycle.push('clear'),
),
fetchImpl: async () => {
lifecycle.push('settle');
return Response.json({ protocolVersion: 1, accepted: true });
},
});

await assert.rejects(
worker.executeAndSettle(
{
protocolVersion: 1,
assignmentId: 'assignment-command-shutdown-cleanly',
workerId: 'vm-1',
incarnationId,
generation: 4,
leaseToken: 'lease-token-that-is-long-enough-for-testing',
expiresAt: new Date(Date.now() + 5_000).toISOString(),
executionKind: 'workspace_tool',
request: {
protocolVersion: 1,
operation: 'execute_command',
workspaceId: 'primary',
command: 'sleep 30',
},
},
controller.signal,
),
/shutdown/,
);
assert.deepEqual(lifecycle, ['arm', 'execute']);
});

test('worker retains quarantine when an atomic executor cannot confirm durability', async () => {
const lifecycle: string[] = [];
const workspaceCapabilities = {
Expand Down
4 changes: 2 additions & 2 deletions service/src/bridge/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import type {
} from '../../../packages/code/src/protocol';

import {
BRIDGE_CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS,
BRIDGE_PROTOCOL_VERSION,
isValidBridgeWorkerCapabilities,
isValidBridgeWorkerId,
Expand All @@ -23,7 +24,6 @@ import { BridgeWorkspaceSlots } from './slots';

const PREFIX = 'codeapi:bridge:v1';
const POLL_INTERVAL_MS = 100;
const CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS = 5_000;
const DEFAULT_WORKER_TTL_SECONDS = 60;
const DEFAULT_REDIS_COMMAND_TIMEOUT_MS = 1_000;

Expand Down Expand Up @@ -1924,7 +1924,7 @@ export class RedisBridgeStore {
// Give Stop its own grace so a near-timeout cancellation is not
// misclassified as an ambiguous timeout.
const cancellationDeadlineAtMs =
Date.now() + CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS;
Date.now() + BRIDGE_CANCELLED_WORKSPACE_SETTLEMENT_GRACE_MS;
let cancellationPollMs = POLL_INTERVAL_MS;
while (Date.now() < cancellationDeadlineAtMs) {
const raw = await boundedCommand(
Expand Down