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
15 changes: 15 additions & 0 deletions services/cloud-agent-next/src/kilo-facade/user-kilo-facade.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@ import type { CloudAgentSession } from '../persistence/CloudAgentSession.js';
import type { Env } from '../types.js';
import { withDORetry } from '../utils/do-retry.js';
import { resolveSessionStub } from '../sandbox-session/session-stub.js';
import { sessionPlaneFromId } from '../session-plane.js';
import { interruptControlSession } from '../router/control-plane-session.js';
import { preflightAndAdmitPromptMessage } from '../session/queue-message.js';
import { parseBasicKiloPrompt } from './basic-prompt.js';
import {
Expand Down Expand Up @@ -881,6 +883,19 @@ async function defaultInterruptPrompt(params: {
userId: string;
cloudAgentSessionId: string;
}): Promise<Awaited<ReturnType<CloudAgentSession['interruptExecution']>>> {
if (sessionPlaneFromId(params.cloudAgentSessionId) === 'control') {
const receipt = await interruptControlSession({
env: params.env,
ownerId: params.userId,
sessionId: params.cloudAgentSessionId,
});
return receipt?.state === 'rejected'
? { success: false, message: receipt.message }
: {
success: receipt !== undefined,
...(receipt ? {} : { message: 'No session work to interrupt' }),
};
}
return withDORetry<
DurableObjectStub<CloudAgentSession>,
Awaited<ReturnType<CloudAgentSession['interruptExecution']>>
Expand Down
49 changes: 43 additions & 6 deletions services/cloud-agent-next/src/persistence/SandboxControl.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,11 @@ import {
type SandboxControlSocketHandler,
} from '../sandbox-control/socket.js';
import { SandboxControlConnectionError } from '../sandbox-control/waiters.js';
import {
hasScopedStopMaintenanceFields,
parseScopedStopMaintenance,
stopAbortWirePayload,
} from '../sandbox-control/scoped-stop-maintenance.js';
import {
createSessionForwarding,
SessionForwardingError,
Expand All @@ -52,6 +57,7 @@ import {
sessionOperationAckSchema,
sessionOperationAuthorizationSchema,
sessionOperationExpiresAt,
sessionAbortPayloadSchema,
sessionRequestIdentitySchema,
wrapperInstanceIdSchema,
type ResponseFrame,
Expand Down Expand Up @@ -529,9 +535,17 @@ export class SandboxControl extends DurableObject<Env> {
async request(input: SandboxControlOutboundRequest): Promise<ResponseFrame> {
await this.ensureOperationalInitialized();
if (input.operation === 'session.git.summary') return this.requestWorktreeChanges(input);
const scopedStop =
input.operation === 'session.abort' ? parseScopedStopMaintenance(input.payload) : undefined;
if (
input.operation === 'session.abort' &&
hasScopedStopMaintenanceFields(input.payload) &&
scopedStop === undefined
)
throw new Error('Invalid scoped Stop maintenance request');
const maintenance =
input.operation === 'session.operation.get' || input.operation === 'session.operation.ack';
if (!maintenance) await this.assertRequestWorktreeAdmission(input);
if (!maintenance && !scopedStop) await this.assertRequestWorktreeAdmission(input);
if (input.operation === 'worktree.delete' || input.operation === 'worktree.prepareDeletion') {
throw new Error('Worktree cleanup requires the deletion coordinator');
}
Expand Down Expand Up @@ -578,7 +592,8 @@ export class SandboxControl extends DurableObject<Env> {
Date.now() >= sessionOperationExpiresAt(maintenanceAuthorization.data))
)
throw new Error('Invalid session operation maintenance authorization');
const runtime = maintenance
const usesMaintenanceChannel = maintenance || scopedStop !== undefined;
const runtime = usesMaintenanceChannel
? this.socketHandler.getConnectionIdentity()
: this.readyWrapperRuntime();
if (!runtime) throw new Error('Sandbox runtime is not ready');
Expand All @@ -594,20 +609,34 @@ export class SandboxControl extends DurableObject<Env> {
)
throw new Error('Sandbox wrapper runtime changed');
const isCurrent = () => {
const current = maintenance
const current = usesMaintenanceChannel
? this.socketHandler.getConnectionIdentity()
: this.readyWrapperRuntime();
return current !== null && this.sameConnection(current, runtime);
};
const physical = await loadPhysicalRecord(this.ctx.storage);
if (
(!maintenance && physical.state !== 'running') ||
(!maintenance && physical.stopTombstone) ||
((!maintenance || scopedStop !== undefined) && physical.state !== 'running') ||
((!maintenance || scopedStop !== undefined) && physical.stopTombstone) ||
physical.providerRef !== runtime.providerInstanceId ||
!isCurrent()
) {
throw new Error('Sandbox runtime is not ready');
}
if (scopedStop) {
const session = sessionRequestIdentitySchema.safeParse(input.session);
if (!session.success || expectedWrapperInstanceId === undefined)
throw new Error('Scoped Stop identity is required');
const route = (await loadRouteTable(this.ctx.storage)).get(session.data.sessionId);
if (
!route ||
route.kiloSessionId !== session.data.kiloSessionId ||
route.directory !== session.data.directory ||
runtime.wrapperInstanceId !== expectedWrapperInstanceId
)
throw new Error('Scoped Stop target is stale');
this.assertWorktreeAdmission(route.worktreeId);
}
if (input.operation === 'session.attach' || input.operation === 'session.prompt') {
const payload = parseOperationPayload(input.operation, input.payload);
if (!payload.ok) throw new Error(payload.error.message);
Expand Down Expand Up @@ -664,8 +693,16 @@ export class SandboxControl extends DurableObject<Env> {
if (!isCurrent()) throw new Error('Sandbox wrapper runtime changed');
});
}
if (!maintenance) await this.assertRequestWorktreeAdmission(input);
if (!usesMaintenanceChannel) await this.assertRequestWorktreeAdmission(input);
if (!isCurrent()) throw new Error('Sandbox wrapper runtime changed');
if (scopedStop) {
const payload = sessionAbortPayloadSchema.parse(input.payload);
return this.socketHandler.sendRequest({
...input,
payload: stopAbortWirePayload(payload, this.socketHandler.supportsScopedStopAbort()),
deadlineAt: scopedStop.cleanupDeadlineAt,
});
}
return this.socketHandler.sendRequest(input);
}

Expand Down
17 changes: 15 additions & 2 deletions services/cloud-agent-next/src/router.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -907,9 +907,20 @@ describe('router sessionId validation', () => {

it('routes workspace_ interrupts to SANDBOX_SESSION', async () => {
const sessionId: SessionId = 'workspace_12345678-1234-1234-1234-123456789abc';
const controlStub = {
...mockSessionStub,
getControlState: vi.fn().mockResolvedValue({
version: 1,
scope: { sandboxId: 'sandbox_1' },
targets: [{ messageId: 'message_1' }],
}),
interruptExecution: vi
.fn()
.mockImplementation(request => Promise.resolve({ ...request, state: 'confirmed' })),
};
const sandboxSession = {
idFromName: vi.fn((id: string) => ({ id })),
get: vi.fn(() => mockSessionStub),
get: vi.fn(() => controlStub),
};
mockContext.env.SANDBOX_SESSION =
sandboxSession as unknown as TRPCContext['env']['SANDBOX_SESSION'];
Expand All @@ -927,7 +938,9 @@ describe('router sessionId validation', () => {
expect(result.success).toBe(true);
expect(sandboxSession.idFromName).toHaveBeenCalledWith(`test-user-123:${sessionId}`);
expect(cloudAgentSession.idFromName).not.toHaveBeenCalled();
expect(mockSessionStub.interruptExecution).toHaveBeenCalled();
expect(controlStub.interruptExecution).toHaveBeenCalledWith(
expect.objectContaining({ targets: [{ messageId: 'message_1' }] })
);
});
});

Expand Down
68 changes: 68 additions & 0 deletions services/cloud-agent-next/src/router/control-plane-session.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
import { describe, expect, it } from 'vitest';
import { interruptControlSession } from './control-plane-session.js';

const OPERATION_ID = '33333333-3333-4333-8333-333333333333';
const RUNTIME_A = '11111111-1111-4111-8111-111111111111';
const RUNTIME_B = '22222222-2222-4222-8222-222222222222';

describe('interruptControlSession', () => {
it('reuses one captured Stop request after the first Durable Object reply is lost', async () => {
let stateCalls = 0;
const deliveries: unknown[] = [];
let current = {
version: 1 as const,
scope: { sandboxId: 'sandbox-a', wrapperInstanceId: RUNTIME_A },
targets: [{ messageId: 'a', wrapperInstanceId: RUNTIME_A, executionDeadlineAt: 3_601_000 }],
};
const getStub = () => ({
getControlState: async () => {
stateCalls++;
return current;
},
interruptExecution: async (request: unknown) => {
deliveries.push(structuredClone(request));
current = {
version: 1,
scope: { sandboxId: 'sandbox-b', wrapperInstanceId: RUNTIME_B },
targets: [
{ messageId: 'b', wrapperInstanceId: RUNTIME_B, executionDeadlineAt: 3_602_000 },
],
};
return { ...(request as object), state: 'confirmed' };
},
});

const receipt = await interruptControlSession(
{ env: {} as never, ownerId: 'user-a', sessionId: 'workspace-a' },
{
getStub,
now: 1_000,
operationId: OPERATION_ID,
retry: async (operation, operationName) => {
if (operationName === 'getControlState') return operation(getStub());
await operation(getStub());
return operation(getStub());
},
}
);

expect(stateCalls).toBe(1);
expect(deliveries).toEqual([
{
version: 1,
operationId: OPERATION_ID,
scope: { sandboxId: 'sandbox-a', wrapperInstanceId: RUNTIME_A },
targets: [{ messageId: 'a', wrapperInstanceId: RUNTIME_A, executionDeadlineAt: 3_601_000 }],
cleanupDeadlineAt: 11_000,
},
{
version: 1,
operationId: OPERATION_ID,
scope: { sandboxId: 'sandbox-a', wrapperInstanceId: RUNTIME_A },
targets: [{ messageId: 'a', wrapperInstanceId: RUNTIME_A, executionDeadlineAt: 3_601_000 }],
cleanupDeadlineAt: 11_000,
},
]);
expect(receipt).toMatchObject({ operationId: OPERATION_ID, state: 'confirmed' });
});
});
53 changes: 53 additions & 0 deletions services/cloud-agent-next/src/router/control-plane-session.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
import type { Env } from '../types.js';
import {
controlSessionStateSchema,
controlStopReceiptSchema,
createControlStopRequest,
type ControlStopReceipt,
} from '../shared/control-plane-session.js';
import { getSandboxSessionStub } from '../sandbox-session/session-stub.js';
import { withDORetry } from '../utils/do-retry.js';

type ControlStopSession = {
getControlState: () => Promise<unknown>;
interruptExecution: (request: unknown) => Promise<unknown>;
};

type ControlSessionStopDependencies = {
getStub?: () => ControlStopSession;
retry?: <T>(
operation: (session: ControlStopSession) => Promise<T>,
operationName: string
) => Promise<T>;
now?: number;
operationId?: string;
};

export async function interruptControlSession(
input: {
env: Pick<Env, 'SANDBOX_SESSION'>;
ownerId: string;
sessionId: string;
},
dependencies: ControlSessionStopDependencies = {}
): Promise<ControlStopReceipt | undefined> {
const stub =
dependencies.getStub ??
(() => getSandboxSessionStub(input.env, input.ownerId, input.sessionId));
const retry =
dependencies.retry ??
(<T>(operation: (session: ControlStopSession) => Promise<T>, operationName: string) =>
withDORetry(stub, operation, operationName));
const state = await retry(session => session.getControlState(), 'getControlState');
if (!state) return undefined;
const request = createControlStopRequest(
controlSessionStateSchema.parse(state),
dependencies.now,
dependencies.operationId
);
return retry(
session =>
session.interruptExecution(request).then(receipt => controlStopReceiptSchema.parse(receipt)),
'interruptControlSession'
);
}
42 changes: 27 additions & 15 deletions services/cloud-agent-next/src/router/handlers/session-management.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import {
import { withDORetry } from '../../utils/do-retry.js';
import { getSandboxSessionStub, resolveSessionStub } from '../../sandbox-session/session-stub.js';
import { sessionPlaneFromId } from '../../session-plane.js';
import { interruptControlSession } from '../control-plane-session.js';
import { protectedProcedure, publicProcedure, internalApiProtectedProcedure } from '../auth.js';
import {
sessionIdSchema,
Expand Down Expand Up @@ -206,34 +207,45 @@ export function createSessionManagementHandlers() {
};
}

// Mark session as interrupted in DO before killing processes (with retry)
// This signals the streaming generator to stop
const getStub = () => resolveSessionStub(env, userId, sessionId);

await withDORetry(getStub, stub => stub.markAsInterrupted(), 'markAsInterrupted');

const interruptResult = await withDORetry(
getStub,
stub => stub.interruptExecution(),
'interruptExecution'
);

if (!interruptResult.success) {
const interruptResult =
sessionPlaneFromId(sessionId) === 'control'
? await interruptControlSession({ env, ownerId: userId, sessionId })
: await withDORetry(
getStub,
stub => stub.interruptExecution(),
'interruptExecution'
);

const success =
interruptResult !== undefined &&
('success' in interruptResult
? interruptResult.success
: interruptResult.state !== 'rejected');
const message =
interruptResult === undefined
? 'No session work to interrupt'
: 'success' in interruptResult
? interruptResult.message
: interruptResult.message;

if (!success) {
logger
.withFields({
message:
interruptResult.message ??
'No accepted current messages or pending queued messages',
message: message ?? 'No accepted current messages or pending queued messages',
})
.info('No accepted current messages or pending queued messages to interrupt');
}

logger.info('Session interruption completed');
return {
success: interruptResult.success,
message: interruptResult.success
success,
message: success
? 'Session interruption accepted'
: (interruptResult.message ?? 'No session work to interrupt'),
: (message ?? 'No session work to interrupt'),
processesFound: false,
};
} catch (error) {
Expand Down
6 changes: 5 additions & 1 deletion services/cloud-agent-next/src/sandbox-control/frames.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,7 +43,11 @@ describe('sandbox control frames', () => {
expect(sandboxHelloResultSchema.parse(helloResult())).toEqual({
protocolVersion: 1,
handshakeComplete: true,
capabilities: { kiloVersionHeartbeat: true, sessionOperationResults: true },
capabilities: {
kiloVersionHeartbeat: true,
sessionOperationResults: true,
scopedStopAbort: 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 @@ -181,6 +181,10 @@ export function helloResult(): SandboxHelloResult {
return {
protocolVersion: SANDBOX_CONTROL_PROTOCOL_VERSION,
handshakeComplete: true,
capabilities: { kiloVersionHeartbeat: true, sessionOperationResults: true },
capabilities: {
kiloVersionHeartbeat: true,
sessionOperationResults: true,
scopedStopAbort: true,
},
};
}
Loading