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
413 changes: 370 additions & 43 deletions services/cloud-agent-next/src/persistence/SandboxControl.ts

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,7 @@ describe('sandbox control frames', () => {
expect(sandboxHelloResultSchema.parse(helloResult())).toEqual({
protocolVersion: 1,
handshakeComplete: true,
capabilities: { kiloVersionHeartbeat: true },
capabilities: { kiloVersionHeartbeat: true, sessionOperationResults: true },
});
const previous = { protocolVersion: 1, handshakeComplete: true };
expect(sandboxHelloResultSchema.parse(previous)).toEqual(previous);
Expand Down
6 changes: 5 additions & 1 deletion services/cloud-agent-next/src/sandbox-control/frames.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,8 @@ import {
sessionTerminalConnectPayloadSchema,
sessionTerminalCreatePayloadSchema,
sessionTerminalResizePayloadSchema,
sessionOperationAuthorizationSchema,
sessionOperationAckSchema,
worktreeDeletePayloadSchema,
type ControlError,
type ControlErrorCode,
Expand Down Expand Up @@ -57,6 +59,8 @@ const REQUEST_PAYLOAD_SCHEMAS: Record<ControlOperation, z.ZodType> = {
'session.terminal.resize': sessionTerminalResizePayloadSchema,
'session.terminal.close': sessionTerminalClosePayloadSchema,
'session.terminal.connect': sessionTerminalConnectPayloadSchema,
'session.operation.get': sessionOperationAuthorizationSchema,
'session.operation.ack': sessionOperationAckSchema,
};

const EVENT_PAYLOAD_SCHEMAS: Record<ControlEvent, z.ZodType> = {
Expand Down Expand Up @@ -177,6 +181,6 @@ export function helloResult(): SandboxHelloResult {
return {
protocolVersion: SANDBOX_CONTROL_PROTOCOL_VERSION,
handshakeComplete: true,
capabilities: { kiloVersionHeartbeat: true },
capabilities: { kiloVersionHeartbeat: true, sessionOperationResults: true },
};
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,120 @@
import { afterEach, describe, expect, it, vi } from 'vitest';
import { createSessionForwarding } from './session-forwarding.js';

describe('createSessionForwarding', () => {
afterEach(() => {
vi.useRealTimers();
});
it('keeps results behind earlier events for the same session', async () => {
const forwarding = createSessionForwarding();
const releaseEvent = Promise.withResolvers<void>();
const event = vi.fn(async () => releaseEvent.promise);
const result = vi.fn(async () => undefined);

const first = forwarding.enqueue('workspace_1', event);
const second = forwarding.enqueue('workspace_1', result);
await Promise.resolve();
expect(result).not.toHaveBeenCalled();

releaseEvent.resolve();
await expect(first).resolves.toBeUndefined();
await expect(second).resolves.toBeUndefined();
expect(event).toHaveBeenCalledTimes(1);
expect(result).toHaveBeenCalledTimes(1);
});

it('continues a session chain after a failed forwarding attempt', async () => {
const forwarding = createSessionForwarding();
const failure = forwarding.enqueue('workspace_1', async () => {
throw new Error('forwarding failed');
});
const recovery = vi.fn(async () => 'acknowledged');
const next = forwarding.enqueue('workspace_1', recovery);

await expect(failure).rejects.toThrow('forwarding failed');
await expect(next).resolves.toBe('acknowledged');
expect(recovery).toHaveBeenCalledTimes(1);
});

it('does not call a fenced delivery after its deadline', async () => {
const forwarding = createSessionForwarding();
const forward = vi.fn(async () => 'acknowledged');

await expect(
forwarding.enqueueFenced({
sessionId: 'workspace_1',
bytes: 1,
deadlineAt: Date.now() - 1,
fence: async () => true,
forward,
})
).rejects.toMatchObject({ retryable: false });
expect(forward).not.toHaveBeenCalled();
});

it('rechecks the deadline after a delayed fence', async () => {
vi.useFakeTimers();
vi.setSystemTime(0);
const forwarding = createSessionForwarding();
const fence = Promise.withResolvers<boolean>();
const forward = vi.fn(async () => 'acknowledged');
const pending = forwarding.enqueueFenced({
sessionId: 'workspace_1',
bytes: 1,
deadlineAt: 1,
fence: async () => fence.promise,
forward,
});

await Promise.resolve();
vi.setSystemTime(2);
fence.resolve(true);
await expect(pending).rejects.toMatchObject({ retryable: false });
expect(forward).not.toHaveBeenCalled();
});

it('rechecks the deadline after forwarding before returning a result', async () => {
vi.useFakeTimers();
vi.setSystemTime(0);
const forwarding = createSessionForwarding();
const finalFence = Promise.withResolvers<boolean>();
let fenceCalls = 0;
const pending = forwarding.enqueueFenced({
sessionId: 'workspace_1',
bytes: 1,
deadlineAt: 1,
fence: () => (++fenceCalls === 1 ? Promise.resolve(true) : finalFence.promise),
forward: async () => 'acknowledged',
});

await Promise.resolve();
vi.setSystemTime(2);
finalFence.resolve(true);
await expect(pending).rejects.toMatchObject({ retryable: false });
});

it('releases capacity after a pre-start fence rejection', async () => {
const forwarding = createSessionForwarding();
await expect(
forwarding.enqueueFenced({
sessionId: 'workspace_1',
bytes: 1,
deadlineAt: Date.now() + 1_000,
fence: async () => false,
forward: async () => 'unreachable',
})
).rejects.toMatchObject({ retryable: false });
expect(forwarding.stats()).toMatchObject({ waiting: 0, inFlight: 0, bufferedBytes: 0 });

await expect(
forwarding.enqueueFenced({
sessionId: 'workspace_1',
bytes: 1,
deadlineAt: Date.now() + 1_000,
fence: async () => true,
forward: async () => 'acknowledged',
})
).resolves.toBe('acknowledged');
expect(forwarding.stats()).toMatchObject({ waiting: 0, inFlight: 0, bufferedBytes: 0 });
});
});
108 changes: 108 additions & 0 deletions services/cloud-agent-next/src/sandbox-control/session-forwarding.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
import {
MAX_SANDBOX_CONTROL_FRAME_BYTES,
SANDBOX_CONTROL_OPERATION_LIMIT,
} from '../shared/sandbox-control-protocol.js';

const MAX_SESSION_FORWARD_BYTES = 4 * MAX_SANDBOX_CONTROL_FRAME_BYTES;

export class SessionForwardingError extends Error {
constructor(
message: string,
readonly retryable: boolean
) {
super(message);
this.name = 'SessionForwardingError';
}
}

export type SessionForwardingStats = {
waiting: number;
inFlight: number;
bufferedBytes: number;
highWater: number;
};

export type FencedSessionForward<T> = {
sessionId: string;
bytes: number;
deadlineAt: number;
fence: () => Promise<boolean>;
forward: () => Promise<T>;
};

export type SessionForwarding = {
enqueue: <T>(sessionId: string, forward: () => Promise<T>) => Promise<T>;
enqueueFenced: <T>(input: FencedSessionForward<T>) => Promise<T>;
stats: () => SessionForwardingStats;
get: (sessionId: string) => Promise<void> | undefined;
values: () => IterableIterator<Promise<void>>;
delete: (sessionId: string) => void;
};

export function createSessionForwarding(): SessionForwarding {
const chains = new Map<string, Promise<void>>();
const stats: SessionForwardingStats = {
waiting: 0,
inFlight: 0,
bufferedBytes: 0,
highWater: 0,
};

const enqueue = <T>(sessionId: string, forward: () => Promise<T>): Promise<T> => {
const previous = chains.get(sessionId) ?? Promise.resolve();
const next = previous.catch(() => undefined).then(forward);
chains.set(
sessionId,
next.then(
() => undefined,
() => undefined
)
);
return next;
};

return {
enqueue,
enqueueFenced<T>(input: FencedSessionForward<T>): Promise<T> {
if (input.bytes > MAX_SANDBOX_CONTROL_FRAME_BYTES)
return Promise.reject(new SessionForwardingError('Forwarded frame is too large', false));
if (
stats.waiting + stats.inFlight >= SANDBOX_CONTROL_OPERATION_LIMIT ||
stats.bufferedBytes + input.bytes > MAX_SESSION_FORWARD_BYTES
)
return Promise.reject(
new SessionForwardingError('Forwarding capacity is unavailable', true)
);
stats.waiting++;
stats.bufferedBytes += input.bytes;
stats.highWater = Math.max(stats.highWater, stats.waiting + stats.inFlight);
return enqueue(input.sessionId, async () => {
stats.waiting--;
stats.inFlight++;
try {
if (
Date.now() >= input.deadlineAt ||
!(await input.fence()) ||
Date.now() >= input.deadlineAt
)
throw new SessionForwardingError('Forwarding fence changed', false);
const result = await input.forward();
if (
Date.now() >= input.deadlineAt ||
!(await input.fence()) ||
Date.now() >= input.deadlineAt
)
throw new SessionForwardingError('Forwarding fence changed', false);
return result;
} finally {
stats.inFlight--;
stats.bufferedBytes -= input.bytes;
}
});
},
stats: () => ({ ...stats }),
get: sessionId => chains.get(sessionId),
values: () => chains.values(),
delete: sessionId => chains.delete(sessionId),
};
}
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,7 @@ describe('sandbox control socket handler', () => {
result: {
protocolVersion: 1,
handshakeComplete: true,
capabilities: { kiloVersionHeartbeat: true },
capabilities: { kiloVersionHeartbeat: true, sessionOperationResults: true },
},
})
);
Expand Down
Loading