diff --git a/apps/web/src/lib/cloud-agent-next/cloud-agent-client.ts b/apps/web/src/lib/cloud-agent-next/cloud-agent-client.ts index 6f83af9772..976ae89146 100644 --- a/apps/web/src/lib/cloud-agent-next/cloud-agent-client.ts +++ b/apps/web/src/lib/cloud-agent-next/cloud-agent-client.ts @@ -367,6 +367,12 @@ export type InterruptResult = { processesFound: boolean; }; +/** Input for canceling one queued (not yet accepted) message. */ +export type CancelQueuedMessageInput = { + sessionId: string; + messageId: string; +}; + export type AnswerQuestionInput = { sessionId: string; questionId: string; @@ -521,6 +527,9 @@ type CloudAgentNextTRPCClient = { interruptSession: { mutate: (input: { sessionId: string }) => Promise; }; + cancelQueuedMessage: { + mutate: (input: CancelQueuedMessageInput) => Promise<{ dropped: boolean }>; + }; getSession: { query: (input: GetSessionInput) => Promise; }; @@ -683,6 +692,26 @@ export class CloudAgentNextClient { } } + /** + * Cancel one queued (not yet accepted) message by id without interrupting the + * active run. Returns whether a pending message was dropped. + */ + async cancelQueuedMessage(sessionId: string, messageId: string): Promise<{ dropped: boolean }> { + try { + return await this.client.cancelQueuedMessage.mutate({ sessionId, messageId }); + } catch (error) { + console.error(`Error canceling queued message ${messageId} in session ${sessionId}:`, error); + captureException(error, { + tags: { + source: 'cloud-agent-next-client', + endpoint: 'cancelQueuedMessage', + }, + extra: { sessionId, messageId }, + }); + throw error; + } + } + /** * Get session state from cloud-agent DO. */ diff --git a/apps/web/src/routers/cloud-agent-next-router.test.ts b/apps/web/src/routers/cloud-agent-next-router.test.ts index 2ff7e684b5..cece37ede7 100644 --- a/apps/web/src/routers/cloud-agent-next-router.test.ts +++ b/apps/web/src/routers/cloud-agent-next-router.test.ts @@ -54,10 +54,14 @@ const mockGenerateCloudAgentAttachmentDownloadUrl = jest.fn< const mockGetSession = jest.fn<(cloudAgentSessionId: string) => Promise<{ model?: string }>>(); +const mockCancelQueuedMessage = + jest.fn<(input: { sessionId: string; messageId: string }) => Promise<{ dropped: boolean }>>(); + const mockCreateCloudAgentNextClient = jest.fn(() => ({ prepareSession: mockPrepareSession, sendMessage: mockSendMessage, getSession: mockGetSession, + cancelQueuedMessage: mockCancelQueuedMessage, })); const mockCreateCloudAgentNextClientForModel = jest.fn( @@ -177,6 +181,7 @@ let createCaller: (ctx: { user: User }) => { contentLength: number; }) => Promise; getAttachmentDownloadUrl: (input: { messageUuid: string; filename: string }) => Promise; + cancelQueuedMessage: (input: { sessionId: string; messageId: string }) => Promise; checkEligibility: () => Promise<{ balance: number; minBalance: number; @@ -413,6 +418,30 @@ describe('cloudAgentNextRouter attachment forwarding', () => { }); }); +describe('cloudAgentNextRouter.cancelQueuedMessage', () => { + beforeEach(() => { + jest.clearAllMocks(); + mockVerifyUserOwnsSessionV2ByCloudAgentId.mockResolvedValue({ + kiloSessionId: 'ses_12345678901234567890123456', + }); + mockCancelQueuedMessage.mockResolvedValue({ dropped: true }); + }); + + it('denies canceling a queued message on a session the user does not own', async () => { + mockVerifyUserOwnsSessionV2ByCloudAgentId.mockResolvedValueOnce(null); + const caller = createCaller({ user: { id: 'user-1', is_admin: false } as User }); + + await expect( + caller.cancelQueuedMessage({ + sessionId: 'agent_123', + messageId: 'msg_123456789abc123456789ABCDE', + }) + ).rejects.toThrow('Session not found or access denied'); + + expect(mockCancelQueuedMessage).not.toHaveBeenCalled(); + }); +}); + describe('cloudAgentNextRouter helper procedures', () => { beforeEach(() => { jest.clearAllMocks(); diff --git a/apps/web/src/routers/cloud-agent-next-router.ts b/apps/web/src/routers/cloud-agent-next-router.ts index ccfa12fed2..a3e7ea328d 100644 --- a/apps/web/src/routers/cloud-agent-next-router.ts +++ b/apps/web/src/routers/cloud-agent-next-router.ts @@ -23,6 +23,7 @@ import { baseInitiateSessionNextOutputSchema, baseSendMessageNextSchema, baseInterruptSessionNextSchema, + baseCancelQueuedMessageNextSchema, baseGetSessionNextSchema, baseGetSessionNextOutputSchema, baseAnswerQuestionNextSchema, @@ -460,6 +461,22 @@ export const cloudAgentNextRouter = createTRPCRouter({ return await client.interruptSession(input.sessionId); }), + /** + * Cancel one queued (not yet accepted) message by id. Never interrupts the + * active run; a missing id or the accepted current message returns + * `{ dropped: false }`. + */ + cancelQueuedMessage: baseProcedure + .input(baseCancelQueuedMessageNextSchema) + .output(z.object({ dropped: z.boolean() })) + .mutation(async ({ ctx, input }) => { + await assertUserOwnsSession(ctx.user.id, input.sessionId); + const authToken = generateCloudAgentToken(ctx.user); + const client = createCloudAgentNextClient(authToken); + + return await client.cancelQueuedMessage(input.sessionId, input.messageId); + }), + answerQuestion: baseProcedure .input(baseAnswerQuestionNextSchema) .output(z.object({ success: z.boolean() })) diff --git a/apps/web/src/routers/cloud-agent-next-schemas.test.ts b/apps/web/src/routers/cloud-agent-next-schemas.test.ts index 3867245613..59e76a1a68 100644 --- a/apps/web/src/routers/cloud-agent-next-schemas.test.ts +++ b/apps/web/src/routers/cloud-agent-next-schemas.test.ts @@ -1,6 +1,7 @@ import { describe, expect, it } from '@jest/globals'; import { basePrepareSessionNextSchema, + baseCancelQueuedMessageNextSchema, cloudAgentGetAttachmentDownloadUrlSchema, cloudAgentGetAttachmentUploadUrlSchema, cloudAgentRelaxedAttachmentFilenameSchema, @@ -264,3 +265,25 @@ describe('basePrepareSessionNextSchema cloneFromKiloSessionId union', () => { expect(result.success).toBe(true); }); }); + +describe('baseCancelQueuedMessageNextSchema', () => { + const VALID_MESSAGE_ID = 'msg_123456789abc123456789ABCDE'; + + it('accepts a session id with a message id', () => { + expect( + baseCancelQueuedMessageNextSchema.safeParse({ + sessionId: 'agent_123', + messageId: VALID_MESSAGE_ID, + }).success + ).toBe(true); + }); + + it('requires both sessionId and messageId', () => { + expect(baseCancelQueuedMessageNextSchema.safeParse({ sessionId: 'agent_123' }).success).toBe( + false + ); + expect( + baseCancelQueuedMessageNextSchema.safeParse({ messageId: VALID_MESSAGE_ID }).success + ).toBe(false); + }); +}); diff --git a/apps/web/src/routers/cloud-agent-next-schemas.ts b/apps/web/src/routers/cloud-agent-next-schemas.ts index 0a0bb7e742..a1ec679fb9 100644 --- a/apps/web/src/routers/cloud-agent-next-schemas.ts +++ b/apps/web/src/routers/cloud-agent-next-schemas.ts @@ -517,6 +517,12 @@ export const baseInterruptSessionNextSchema = z.object({ sessionId: z.string(), }); +// Schema for canceling one queued (not yet accepted) message by id. +export const baseCancelQueuedMessageNextSchema = z.object({ + sessionId: z.string(), + messageId: messageIdNextSchema, +}); + // Schema for getting session state export const baseGetSessionNextSchema = z.object({ cloudAgentSessionId: z.string(), diff --git a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts index fb68622995..d0f13f8c3d 100644 --- a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts +++ b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.test.ts @@ -58,10 +58,14 @@ const mockGenerateCloudAgentAttachmentUploadUrl = jest.fn< const mockGetSession = jest.fn<(cloudAgentSessionId: string) => Promise<{ model?: string }>>(); +const mockCancelQueuedMessage = + jest.fn<(input: { sessionId: string; messageId: string }) => Promise<{ dropped: boolean }>>(); + const mockCreateCloudAgentNextClient = jest.fn(() => ({ prepareSession: mockPrepareSession, sendMessage: mockSendMessage, getSession: mockGetSession, + cancelQueuedMessage: mockCancelQueuedMessage, })); const mockCreateCloudAgentNextClientForModel = jest.fn( @@ -218,6 +222,11 @@ let createCaller: (ctx: { user: User }) => { contentType: 'text/markdown'; contentLength: number; }) => Promise; + cancelQueuedMessage: (input: { + organizationId: string; + sessionId: string; + messageId: string; + }) => Promise; listBitbucketRepositories: (input: { organizationId: string; forceRefresh?: boolean; @@ -472,6 +481,32 @@ describe('organizationCloudAgentNextRouter attachment forwarding', () => { }); }); +describe('organizationCloudAgentNextRouter.cancelQueuedMessage', () => { + beforeEach(() => { + jest.clearAllMocks(); + mockEnsureOrganizationAccess.mockImplementation(() => undefined); + mockVerifyOrgOwnsSessionV2ByCloudAgentId.mockResolvedValue({ + kiloSessionId: 'ses_12345678901234567890123456', + }); + mockCancelQueuedMessage.mockResolvedValue({ dropped: true }); + }); + + it('denies canceling a queued message on a session outside the organization', async () => { + mockVerifyOrgOwnsSessionV2ByCloudAgentId.mockResolvedValueOnce(null); + const caller = createCaller({ user: { id: 'user-1', is_admin: false } as User }); + + await expect( + caller.cancelQueuedMessage({ + organizationId: ORGANIZATION_ID, + sessionId: 'agent_123', + messageId: 'msg_123456789abc123456789ABCDE', + }) + ).rejects.toThrow('Organization does not own this session'); + + expect(mockCancelQueuedMessage).not.toHaveBeenCalled(); + }); +}); + describe('organizationCloudAgentNextRouter helper procedures', () => { beforeEach(() => { jest.clearAllMocks(); diff --git a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts index 43e710f697..fb351c3e1b 100644 --- a/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts +++ b/apps/web/src/routers/organizations/organization-cloud-agent-next-router.ts @@ -31,6 +31,7 @@ import { baseInitiateSessionNextOutputSchema, baseSendMessageNextSchema, baseInterruptSessionNextSchema, + baseCancelQueuedMessageNextSchema, baseGetSessionNextSchema, baseGetSessionNextOutputSchema, baseAnswerQuestionNextSchema, @@ -141,6 +142,10 @@ const InterruptSessionInput = baseInterruptSessionNextSchema.extend({ organizationId: z.uuid(), }); +const CancelQueuedMessageInput = baseCancelQueuedMessageNextSchema.extend({ + organizationId: z.uuid(), +}); + const ImageUploadUrlInput = cloudAgentGetImageUploadUrlSchema.extend({ organizationId: z.uuid(), }); @@ -606,6 +611,26 @@ export const organizationCloudAgentNextRouter = createTRPCRouter({ return await client.interruptSession(input.sessionId); }), + /** + * Cancel one queued (not yet accepted) message by id. Never interrupts the + * active run; a missing id or the accepted current message returns + * `{ dropped: false }`. + */ + cancelQueuedMessage: organizationMemberMutationProcedure + .input(CancelQueuedMessageInput) + .output(z.object({ dropped: z.boolean() })) + .mutation(async ({ ctx, input }) => { + await assertOrganizationOwnsSession({ + organizationId: input.organizationId, + userId: ctx.user.id, + cloudAgentSessionId: input.sessionId, + }); + const authToken = generateCloudAgentToken(ctx.user); + const client = createCloudAgentNextClient(authToken); + + return await client.cancelQueuedMessage(input.sessionId, input.messageId); + }), + answerQuestion: organizationMemberMutationProcedure .input(AnswerQuestionInput) .output(z.object({ success: z.boolean() })) diff --git a/packages/cloud-agent-sdk/src/cli-live-transport.test.ts b/packages/cloud-agent-sdk/src/cli-live-transport.test.ts index b36478fca8..7d84d5574b 100644 --- a/packages/cloud-agent-sdk/src/cli-live-transport.test.ts +++ b/packages/cloud-agent-sdk/src/cli-live-transport.test.ts @@ -1062,6 +1062,74 @@ describe('CliLiveTransport unified user web connection', () => { transport.destroy(); }); + it('includes messageID on send_message when the client assigns a message id', async () => { + const connection = createConnection(); + jest + .mocked(connection.sendCommand) + .mockImplementation((_sessionId, command) => + Promise.resolve(command === 'list_models' ? WIRE_CATALOG : { ok: true }) + ); + const { transport } = createTransportWithSinks({ connection }); + + transport.connect(); + emitOwner(connection); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + jest.mocked(connection.sendCommand).mockClear(); + + await transport.send?.({ + payload: { type: 'prompt', prompt: 'hello' }, + messageId: 'msg-queued-1', + }); + + expect(connection.sendCommand).toHaveBeenCalledWith( + KILO_SESSION_ID, + 'send_message', + { + sessionID: KILO_SESSION_ID, + parts: [{ type: 'text', text: 'hello' }], + messageID: 'msg-queued-1', + }, + 'owner' + ); + transport.destroy(); + }); + + it('relays drop_queued_message with messageID without sending interrupt', async () => { + const connection = createConnection(); + jest + .mocked(connection.sendCommand) + .mockImplementation((_sessionId, command) => + Promise.resolve(command === 'list_models' ? WIRE_CATALOG : { ok: true }) + ); + const { transport } = createTransportWithSinks({ connection }); + + transport.connect(); + emitOwner(connection); + await Promise.resolve(); + await Promise.resolve(); + await Promise.resolve(); + jest.mocked(connection.sendCommand).mockClear(); + + await transport.dropQueuedMessage?.('msg-drop-1'); + + expect(connection.sendCommand).toHaveBeenCalledTimes(1); + expect(connection.sendCommand).toHaveBeenCalledWith( + KILO_SESSION_ID, + 'drop_queued_message', + { protocolVersion: 1, messageID: 'msg-drop-1' }, + 'owner' + ); + expect(connection.sendCommand).not.toHaveBeenCalledWith( + KILO_SESSION_ID, + 'interrupt', + expect.anything(), + 'owner' + ); + transport.destroy(); + }); + it.each([ ['Kilo', { providerID: 'kilo', modelID: 'anthropic/claude-sonnet-4' }], ['non-Kilo', { providerID: 'anthropic', modelID: 'claude-sonnet-4' }], diff --git a/packages/cloud-agent-sdk/src/cli-live-transport.ts b/packages/cloud-agent-sdk/src/cli-live-transport.ts index ab185b1929..6dabef97ba 100644 --- a/packages/cloud-agent-sdk/src/cli-live-transport.ts +++ b/packages/cloud-agent-sdk/src/cli-live-transport.ts @@ -990,6 +990,10 @@ function createCliLiveTransport(config: CliLiveTransportConfig): TransportFactor return sendCommand('send_message', { sessionID: config.kiloSessionId, parts, + // Old form is `send_message` without `messageID`; include it once the + // client assigns an id so the CLI can correlate the queued turn. + // Remove the omission when every client sends it. + ...(input.messageId ? { messageID: input.messageId } : {}), ...(payload.mode ? { agent: payload.mode } : {}), ...(remoteModel.kind === 'none' ? {} @@ -1000,6 +1004,8 @@ function createCliLiveTransport(config: CliLiveTransportConfig): TransportFactor }); }, interrupt: () => sendCommand('interrupt', {}), + dropQueuedMessage: messageId => + sendCommand('drop_queued_message', { protocolVersion: 1, messageID: messageId }), answer: payload => sendCommand('question_reply', { requestID: payload.requestId, diff --git a/packages/cloud-agent-sdk/src/cloud-agent-transport.test.ts b/packages/cloud-agent-sdk/src/cloud-agent-transport.test.ts index b18dc6592d..7847cf7a55 100644 --- a/packages/cloud-agent-sdk/src/cloud-agent-transport.test.ts +++ b/packages/cloud-agent-sdk/src/cloud-agent-transport.test.ts @@ -66,6 +66,7 @@ function createMockApi() { return { send: jest.fn(() => Promise.resolve('sent')), interrupt: jest.fn(() => Promise.resolve('interrupted')), + cancelQueuedMessage: jest.fn(() => Promise.resolve({ dropped: true })), answer: jest.fn(() => Promise.resolve('answered')), reject: jest.fn(() => Promise.resolve('rejected')), respondToPermission: jest.fn(() => Promise.resolve('responded')), @@ -531,6 +532,20 @@ describe('CloudAgentTransport command delegation', () => { transport.destroy(); }); + it('dropQueuedMessage() delegates to api.cancelQueuedMessage with bound sessionId', async () => { + const api = createMockApi(); + const { transport } = createTransportWithSinks(undefined, undefined, api); + + await transport.dropQueuedMessage!('msg-queued-1'); + + expect(api.cancelQueuedMessage).toHaveBeenCalledWith({ + sessionId: 'ses-1', + messageId: 'msg-queued-1', + }); + + transport.destroy(); + }); + it('answer() delegates to api.answer with bound sessionId', () => { const api = createMockApi(); const { transport } = createTransportWithSinks(undefined, undefined, api); diff --git a/packages/cloud-agent-sdk/src/cloud-agent-transport.ts b/packages/cloud-agent-sdk/src/cloud-agent-transport.ts index 604441a805..f9be44c921 100644 --- a/packages/cloud-agent-sdk/src/cloud-agent-transport.ts +++ b/packages/cloud-agent-sdk/src/cloud-agent-transport.ts @@ -331,6 +331,12 @@ function createCloudAgentTransport(config: CloudAgentTransportConfig): Transport ...(input.images ? { images: input.images } : {}), }), interrupt: () => config.api.interrupt({ sessionId: config.sessionId }), + dropQueuedMessage: async messageId => { + if (!config.api.cancelQueuedMessage) { + throw new Error('Cloud Agent cancel queued message is not configured'); + } + return config.api.cancelQueuedMessage({ sessionId: config.sessionId, messageId }); + }, answer: payload => config.api.answer({ sessionId: config.sessionId, ...payload }), reject: payload => config.api.reject({ sessionId: config.sessionId, ...payload }), respondToPermission: payload => diff --git a/packages/cloud-agent-sdk/src/normalizer.ts b/packages/cloud-agent-sdk/src/normalizer.ts index 5eb58986c7..7a92e9bc3d 100644 --- a/packages/cloud-agent-sdk/src/normalizer.ts +++ b/packages/cloud-agent-sdk/src/normalizer.ts @@ -170,6 +170,11 @@ export type ServiceEvent = executionId?: string | undefined; content?: string | undefined; } + | { + type: 'cloud.message.canceled'; + messageId: string; + executionId?: string | undefined; + } | { type: 'cloud.message.sent'; messageId: string; @@ -230,6 +235,14 @@ const sessionModelSchema = z.object({ variant: z.string().optional(), }); +// `cloud.message.canceled` mirrors the queued payload minus the content field. +const cloudMessageCanceledDataSchema = z + .object({ + messageId: z.string(), + executionId: z.string().optional(), + }) + .passthrough(); + function normalizeSessionInfo(rawInfo: { id: string; [key: string]: unknown }): SessionInfo { const model = sessionModelSchema.safeParse(rawInfo['model']); return { @@ -518,6 +531,16 @@ function normalizeInnerEvent(eventType: string, data: unknown): NormalizedEvent }; } + case 'cloud.message.canceled': { + const r = cloudMessageCanceledDataSchema.safeParse(data); + if (!r.success) return null; + return { + type: 'cloud.message.canceled', + messageId: r.data.messageId, + executionId: r.data.executionId, + }; + } + case 'cloud.message.sent': { const r = cloudMessageSentDataSchema.safeParse(data); if (!r.success) return null; diff --git a/packages/cloud-agent-sdk/src/service-state.ts b/packages/cloud-agent-sdk/src/service-state.ts index 5832f51751..cc1b597e4b 100644 --- a/packages/cloud-agent-sdk/src/service-state.ts +++ b/packages/cloud-agent-sdk/src/service-state.ts @@ -54,6 +54,8 @@ type ServiceStateConfig = { onPreparationFailed?: ((message: string) => void) | undefined; /** Fired when the server acknowledges a user message was queued. */ onMessageQueued?: ((messageId: string) => void) | undefined; + /** Fired when a queued user message is canceled before delivery. */ + onMessageCanceled?: ((messageId: string) => void) | undefined; /** Fired when a queued user message's execution terminates in 'completed'. */ onMessageCompleted?: ((messageId: string) => void) | undefined; /** Fired when a queued user message fails delivery or its execution fails. */ @@ -568,6 +570,14 @@ function createServiceState(config: ServiceStateConfig): ServiceState { notify(); } + function processMessageCanceled( + event: Extract + ): void { + pendingMessages.delete(event.messageId); + config.onMessageCanceled?.(event.messageId); + notify(); + } + function processMessageSent(event: Extract): void { pendingMessages.delete(event.messageId); notify(); @@ -794,6 +804,9 @@ function createServiceState(config: ServiceStateConfig): ServiceState { case 'cloud.message.queued': processMessageQueued(event); break; + case 'cloud.message.canceled': + processMessageCanceled(event); + break; case 'cloud.message.sent': processMessageSent(event); break; diff --git a/packages/cloud-agent-sdk/src/session-manager.test.ts b/packages/cloud-agent-sdk/src/session-manager.test.ts index 04bda6280d..66a285eab4 100644 --- a/packages/cloud-agent-sdk/src/session-manager.test.ts +++ b/packages/cloud-agent-sdk/src/session-manager.test.ts @@ -49,6 +49,7 @@ type MockSession = Omit< | 'storage' | 'send' | 'interrupt' + | 'cancelQueuedMessage' | 'answer' | 'reject' | 'respondToPermission' @@ -60,6 +61,7 @@ type MockSession = Omit< storage: JotaiSessionStorage | null; send: jest.Mock, [CloudAgentSessionSendInput]>; interrupt: jest.Mock, []>; + cancelQueuedMessage: jest.Mock, [string]>; answer: jest.Mock, [CloudAgentSessionAnswerInput]>; reject: jest.Mock, [CloudAgentSessionRejectInput]>; respondToPermission: jest.Mock, [CloudAgentSessionRespondToPermissionInput]>; @@ -74,6 +76,7 @@ const mockSession = { destroy: jest.fn(), send: jest.fn(), interrupt: jest.fn(), + cancelQueuedMessage: jest.fn(), answer: jest.fn(), reject: jest.fn(), respondToPermission: jest.fn(), @@ -402,6 +405,8 @@ describe('createSessionManager', () => { mockSession.send.mockClear(); mockSession.interrupt.mockClear(); mockSession.interrupt.mockResolvedValue({}); + mockSession.cancelQueuedMessage.mockClear(); + mockSession.cancelQueuedMessage.mockResolvedValue(undefined); mockSession.createRemoteSession.mockClear(); mockSession.createRemoteSession.mockResolvedValue(kiloId('ses_12345678901234567890123456')); mockSession.exitRemoteSession.mockClear(); @@ -4577,6 +4582,19 @@ describe('createSessionManager', () => { }); }); + describe('cancelQueuedMessage', () => { + it('delegates to the active session without interrupting', async () => { + const config = createMockConfig(); + const mgr = createSessionManager(config); + + await mgr.switchSession(kiloId('ses-1')); + await mgr.cancelQueuedMessage('msg-queued-1'); + + expect(mockSession.cancelQueuedMessage).toHaveBeenCalledWith('msg-queued-1'); + expect(mockSession.interrupt).not.toHaveBeenCalled(); + }); + }); + // ------------------------------------------------------------------------- // createRemoteSession // ------------------------------------------------------------------------- diff --git a/packages/cloud-agent-sdk/src/session-manager.ts b/packages/cloud-agent-sdk/src/session-manager.ts index cbee12be54..6023b80105 100644 --- a/packages/cloud-agent-sdk/src/session-manager.ts +++ b/packages/cloud-agent-sdk/src/session-manager.ts @@ -447,6 +447,8 @@ type SessionManager = { createRemoteSession(input?: CreateRemoteSessionInput): Promise; exitRemoteSession(): Promise; interrupt(): Promise; + /** Drop one queued (not yet accepted) message by id without interrupting the active run. */ + cancelQueuedMessage(messageId: string): Promise; /** * Clear the active session's local transcript view only. Server-side history * is untouched and reappears on re-entry (`switchSession`). No-op without an @@ -2007,6 +2009,15 @@ function createSessionManager(config: SessionManagerConfig): SessionManager { } } + async function cancelQueuedMessage(messageId: string): Promise { + if (!currentSession) return; + // Delegate to the session: the cloud-agent transport calls the + // `cancelQueuedMessage` tRPC mutation and the remote transport relays + // `drop_queued_message`. A remote CLI_UPGRADE_REQUIRED rejection surfaces + // to the caller verbatim; this path never falls back to `interrupt()`. + await currentSession.cancelQueuedMessage(messageId); + } + function clearTranscript(): void { if (!currentSession) return; currentSession.storage.clear(); @@ -2160,6 +2171,7 @@ function createSessionManager(config: SessionManagerConfig): SessionManager { createRemoteSession, exitRemoteSession, interrupt, + cancelQueuedMessage, clearTranscript, answerQuestion, rejectQuestion, diff --git a/packages/cloud-agent-sdk/src/session-transport.test.ts b/packages/cloud-agent-sdk/src/session-transport.test.ts index 7436c3951f..164b383d7a 100644 --- a/packages/cloud-agent-sdk/src/session-transport.test.ts +++ b/packages/cloud-agent-sdk/src/session-transport.test.ts @@ -582,6 +582,38 @@ describe('delivery callback plumbing', () => { }); }); +describe('queued message cancellation replay', () => { + it('replays cloud.message.queued then cloud.message.canceled to a net-empty transcript', async () => { + const api = createMockApi(); + const session = createCloudAgentResolvedSession(api); + await connectSession(session); + + const messageId = 'msg_queued_canceled'; + const deliver = (streamEventType: string, data: unknown) => { + mockWs.onmessage?.({ + data: JSON.stringify({ + eventId: 1, + executionId: null, + sessionId: cloudAgentSessionId, + streamEventType, + timestamp: new Date().toISOString(), + data, + }), + } as MessageEvent); + }; + + deliver('cloud.message.queued', { messageId, content: 'hello' }); + expect(session.state.getPendingMessages().has(messageId)).toBe(true); + expect(session.storage.getMessageIds()).toContain(messageId); + + deliver('cloud.message.canceled', { messageId }); + expect(session.state.getPendingMessages().has(messageId)).toBe(false); + expect(session.storage.getMessageIds()).not.toContain(messageId); + + session.destroy(); + }); +}); + describe('disconnect during resolution', () => { it('disconnect() before resolveSession settles prevents transport from attaching', async () => { const api = createMockApi(); diff --git a/packages/cloud-agent-sdk/src/session.ts b/packages/cloud-agent-sdk/src/session.ts index 5d65347ff0..ab0879accd 100644 --- a/packages/cloud-agent-sdk/src/session.ts +++ b/packages/cloud-agent-sdk/src/session.ts @@ -87,6 +87,7 @@ type CloudAgentSessionConfig = { onReplayComplete?: () => void; onEvent?: (event: NormalizedEvent) => void; onMessageQueued?: (messageId: string) => void; + onMessageCanceled?: (messageId: string) => void; onMessageCompleted?: (messageId: string) => void; onMessageFailed?: ( messageId: string, @@ -195,6 +196,7 @@ type CloudAgentSession = { // Commands send: (input: CloudAgentSessionSendInput) => unknown | Promise; interrupt: () => unknown | Promise; + cancelQueuedMessage: (messageId: string) => unknown | Promise; answer: (payload: CloudAgentSessionAnswerInput) => unknown | Promise; reject: (payload: CloudAgentSessionRejectInput) => unknown | Promise; respondToPermission: ( @@ -241,6 +243,7 @@ function createCloudAgentSession(config: CloudAgentSessionConfig): CloudAgentSes onSessionCreated: config.onSessionCreated, onSessionUpdated: config.onSessionUpdated, onMessageQueued: config.onMessageQueued, + onMessageCanceled: config.onMessageCanceled, onMessageCompleted: config.onMessageCompleted, onMessageFailed: config.onMessageFailed, }); @@ -307,6 +310,11 @@ function createCloudAgentSession(config: CloudAgentSessionConfig): CloudAgentSes content: event.content, }); } + // `cloud.message.canceled` removes the synthetic (or real) local row the + // queued event materialized, so a replay of queued then canceled nets empty. + if (event.type === 'cloud.message.canceled') { + storage.deleteMessage(event.messageId); + } config.onEvent?.(event); }, onReplayComplete: () => config.onReplayComplete?.(), @@ -447,6 +455,12 @@ function createCloudAgentSession(config: CloudAgentSessionConfig): CloudAgentSes } return transport.interrupt(); }, + cancelQueuedMessage: messageId => { + if (!transport?.dropQueuedMessage) { + throw new Error('CloudAgentSession transport.dropQueuedMessage is not configured'); + } + return transport.dropQueuedMessage(messageId); + }, answer: payload => { if (!transport?.answer) { throw new Error('CloudAgentSession transport.answer is not configured'); diff --git a/packages/cloud-agent-sdk/src/transport.ts b/packages/cloud-agent-sdk/src/transport.ts index 350adf7d2b..8a907d686f 100644 --- a/packages/cloud-agent-sdk/src/transport.ts +++ b/packages/cloud-agent-sdk/src/transport.ts @@ -151,6 +151,13 @@ type Transport = { createSession?: (input?: CreateRemoteSessionInput) => Promise; exitSession?: () => Promise; interrupt?: () => Promise; + /** + * Drop one queued (not yet accepted) message by its client message id. The + * remote CLI transport relays `drop_queued_message`; the cloud-agent transport + * calls the `cancelQueuedMessage` tRPC mutation. Old remotes reject the drop + * with CLI_UPGRADE_REQUIRED. + */ + dropQueuedMessage?: (messageId: string) => Promise; answer?: (payload: { requestId: string; answers: string[][] }) => Promise; reject?: (payload: { requestId: string }) => Promise; respondToPermission?: (payload: { @@ -180,6 +187,15 @@ type CloudAgentApi = { images?: Images; }) => Promise; interrupt: (payload: { sessionId: CloudAgentSessionId }) => Promise; + /** + * Cancel one queued message by id (cloud-agent branch of + * `dropQueuedMessage`). Optional so legacy providers that predate the + * mutation keep compiling; the transport surfaces a clear error when absent. + */ + cancelQueuedMessage?: (payload: { + sessionId: CloudAgentSessionId; + messageId: string; + }) => Promise; answer: (payload: { sessionId: CloudAgentSessionId; requestId: string; diff --git a/services/cloud-agent-next/src/persistence/CloudAgentSession.ts b/services/cloud-agent-next/src/persistence/CloudAgentSession.ts index e56eeab910..77a1b1679b 100644 --- a/services/cloud-agent-next/src/persistence/CloudAgentSession.ts +++ b/services/cloud-agent-next/src/persistence/CloudAgentSession.ts @@ -2053,6 +2053,42 @@ export class CloudAgentSession extends DurableObject { return { success: true, executionId: undefined }; } + /** + * Drop one pending (not yet accepted) queued message by id without touching + * the accepted/current run. The queue terminalizes the queued durable state + * and removes the pending row; this method then persists a + * `cloud.message.canceled` replay event after the queued event so a + * reconnecting client nets empty after replaying queued then canceled. + */ + async cancelQueuedMessage(messageId: string): Promise<{ dropped: boolean }> { + const result = await this.getSessionMessageQueue().cancelQueuedMessage(messageId); + if (!result.dropped) { + return { dropped: false }; + } + + const sessionId = await this.requireSessionId(); + const payload = JSON.stringify({ messageId }); + const eventId = this.eventQueries.insertUnique({ + executionId: '' as EventSourceId, + entityId: `canceled-message/${messageId}`, + sessionId, + streamEventType: 'cloud.message.canceled', + payload, + timestamp: Date.now(), + }); + if (eventId !== null) { + this.broadcastEvent({ + id: eventId, + execution_id: '' as EventSourceId, + session_id: sessionId, + stream_event_type: 'cloud.message.canceled', + payload, + timestamp: Date.now(), + }); + } + return { dropped: true }; + } + private async getTerminalClient(): Promise> { const sessionId = await this.requireSessionId(); const terminal = await resolveTerminalWrapperClient({ diff --git a/services/cloud-agent-next/src/router/handlers/session-management.ts b/services/cloud-agent-next/src/router/handlers/session-management.ts index 7160f1f8bd..695a1c7309 100644 --- a/services/cloud-agent-next/src/router/handlers/session-management.ts +++ b/services/cloud-agent-next/src/router/handlers/session-management.ts @@ -16,6 +16,7 @@ import { resolveSessionStub } from '../../sandbox-session/session-stub.js'; import { protectedProcedure, publicProcedure, internalApiProtectedProcedure } from '../auth.js'; import { sessionIdSchema, + MessageIdSchema, GetSessionInput, GetSessionOutput, GetSessionHealthInput, @@ -236,13 +237,55 @@ export function createSessionManagementHandlers() { }); }), + /** + * Drop one pending (not yet accepted) queued message by id. Never interrupts + * the accepted/current run; a missing id or the accepted current message + * returns `{ dropped: false }`. + */ + cancelQueuedMessage: protectedProcedure + .input( + z.object({ + sessionId: sessionIdSchema.describe('Session ID owning the queued message'), + messageId: MessageIdSchema.describe('Message ID to drop from the queue'), + }) + ) + .mutation(async ({ input, ctx }) => { + return withLogTags({ source: 'cancelQueuedMessage' }, async () => { + const sessionId = input.sessionId as SessionId; + const { userId, env } = ctx; + + logger.setTags({ userId, sessionId }); + logger.info('Canceling queued message'); + await requireCurrentSessionAccess({ + env, + kiloUserId: userId, + cloudAgentSessionId: sessionId, + }); + + try { + const getStub = () => resolveSessionStub(env, userId, sessionId); + return await withDORetry( + getStub, + stub => stub.cancelQueuedMessage(input.messageId), + 'cancelQueuedMessage' + ); + } catch (error) { + const errorMsg = error instanceof Error ? error.message : String(error); + logger.withFields({ error: errorMsg }).error('Failed to cancel queued message'); + throw new TRPCError({ + code: 'INTERNAL_SERVER_ERROR', + message: `Failed to cancel queued message: ${errorMsg}`, + }); + } + }); + }), + /** * Get session metadata. * * Returns sanitized session metadata (no secrets) including lifecycle timestamps. * Useful for frontend idempotency - checking if a session was already initiated * before a page refresh. - * * Security: * - Excludes: githubToken, gitToken, envVars values, setupCommands, mcpServers configs * - Includes: counts of envVars, setupCommands, mcpServers for debugging diff --git a/services/cloud-agent-next/src/session/session-message-queue.test.ts b/services/cloud-agent-next/src/session/session-message-queue.test.ts index 0a8bc493d4..133d23b804 100644 --- a/services/cloud-agent-next/src/session/session-message-queue.test.ts +++ b/services/cloud-agent-next/src/session/session-message-queue.test.ts @@ -2397,4 +2397,110 @@ describe('SessionMessageQueue', () => { }, ]); }); + + it('drops a queued pending message and terminalizes its state as canceled', async () => { + const harness = createQueueHarness(); + await harness.queue.admitSubmittedMessage({ + userId: 'user_test' as UserId, + turn: { type: 'prompt', id: FIRST_MESSAGE_ID, prompt: 'cancel me before delivery' }, + }); + + const result = await harness.queue.cancelQueuedMessage(FIRST_MESSAGE_ID); + + expect(result).toEqual({ dropped: true }); + expect(await listPendingSessionMessages(harness.storage)).toHaveLength(0); + await expect(getSessionMessageState(harness.storage, FIRST_MESSAGE_ID)).resolves.toMatchObject({ + status: 'interrupted', + completionSource: 'canceled', + failureStage: 'interruption', + failureCode: 'user_interrupt', + }); + }); + + it('reports dropped for a retried cancel of an already-canceled message', async () => { + const harness = createQueueHarness(); + await putSessionMessageState(harness.storage, { + ...createQueuedSessionMessageState({ + turn: { type: 'prompt', messageId: FIRST_MESSAGE_ID, prompt: 'already canceled' }, + agent: { mode: 'code', model: 'default-model' }, + }), + status: 'interrupted', + terminalAt: 12, + completionSource: 'canceled', + failureStage: 'interruption', + failureCode: 'user_interrupt', + }); + + const result = await harness.queue.cancelQueuedMessage(FIRST_MESSAGE_ID); + + expect(result).toEqual({ dropped: true }); + }); + + it('does not drop an accepted message with leftover pending residue', async () => { + const harness = createQueueHarness(); + await storePendingSessionMessage( + harness.storage, + createPendingSessionMessage({ + messageId: FIRST_MESSAGE_ID, + role: 'user', + content: 'already accepted', + createdAt: 1, + }) + ); + await putSessionMessageState(harness.storage, { + ...createQueuedSessionMessageState({ + turn: { type: 'prompt', messageId: FIRST_MESSAGE_ID, prompt: 'already accepted' }, + agent: { mode: 'code', model: 'default-model' }, + }), + status: 'accepted', + acceptedAt: 3, + wrapperRunId: 'wr_existing', + }); + + const result = await harness.queue.cancelQueuedMessage(FIRST_MESSAGE_ID); + + expect(result).toEqual({ dropped: false }); + expect(await listPendingSessionMessages(harness.storage)).toHaveLength(1); + await expect(getSessionMessageState(harness.storage, FIRST_MESSAGE_ID)).resolves.toMatchObject({ + status: 'accepted', + wrapperRunId: 'wr_existing', + }); + }); + + it.each([ + ['completed', 'assistant_message_event'], + ['failed', 'wrapper_failure'], + ['interrupted', 'interrupt'], + ] as const)( + 'does not drop a %s message with leftover pending residue', + async (status, completionSource) => { + const harness = createQueueHarness(); + await storePendingSessionMessage( + harness.storage, + createPendingSessionMessage({ + messageId: FIRST_MESSAGE_ID, + role: 'user', + content: 'already terminal', + createdAt: 1, + }) + ); + await putSessionMessageState(harness.storage, { + ...createQueuedSessionMessageState({ + turn: { type: 'prompt', messageId: FIRST_MESSAGE_ID, prompt: 'already terminal' }, + agent: { mode: 'code', model: 'default-model' }, + }), + status, + terminalAt: 12, + completionSource, + }); + + const result = await harness.queue.cancelQueuedMessage(FIRST_MESSAGE_ID); + + expect(result).toEqual({ dropped: false }); + expect(await listPendingSessionMessages(harness.storage)).toHaveLength(1); + await expect( + getSessionMessageState(harness.storage, FIRST_MESSAGE_ID) + ).resolves.toMatchObject({ status }); + } + ); }); diff --git a/services/cloud-agent-next/src/session/session-message-queue.ts b/services/cloud-agent-next/src/session/session-message-queue.ts index b532ee9071..950d6e5174 100644 --- a/services/cloud-agent-next/src/session/session-message-queue.ts +++ b/services/cloud-agent-next/src/session/session-message-queue.ts @@ -45,6 +45,7 @@ import { createQueuedSessionMessageState, getSessionMessageState, listReconnectVisibleTerminalQueuedMessages, + markMessageInterrupted, putSessionMessageState, type SessionMessageFailureCode, type SessionMessageStorage, @@ -125,6 +126,15 @@ export type SessionMessageQueue = { interruptPendingQueuedMessages( afterTransition?: (messages: PendingSessionMessage[]) => Promise ): Promise; + /** + * Delete one pending (not yet accepted) queued message by id without + * interrupting the accepted/current run. A successful drop also terminalizes + * the message's queued `SessionMessageState` so a later re-admit treats the + * id as terminal instead of ACKing it as still queued. Returns whether a + * pending row was removed; a missing id or the accepted current message + * returns false. + */ + cancelQueuedMessage(messageId: string): Promise<{ dropped: boolean }>; recoverPendingInterruption( afterTransition?: (messages: PendingSessionMessage[]) => Promise ): Promise; @@ -1211,6 +1221,51 @@ export function createSessionMessageQueue( return capturedMessages; } + async function cancelQueuedMessage(messageId: string): Promise<{ dropped: boolean }> { + const pendingMessage = await findPendingSessionMessageByMessageId(storage, messageId); + + // A retry after a successful cancel finds no pending row, but the durable + // state is already terminalized as interrupted/canceled. Report dropped so + // the caller still persists the `cloud.message.canceled` replay event. + if (!pendingMessage) { + const existing = await getSessionMessageState(storage, messageId); + return existing?.status === 'interrupted' && existing.completionSource === 'canceled' + ? { dropped: true } + : { dropped: false }; + } + + // `admitIntent` wrote a queued `SessionMessageState` alongside the pending + // row. Terminalize it so `getExistingAdmissionAckForMessageId` rejects a + // later re-admit instead of ACKing the id as still queued. Terminalize + // without the outbox `cloud.message.failed` effect: the caller persists the + // `cloud.message.canceled` replay event instead, and the `canceled` + // completion source keeps the dropped id out of the reconnect catch-up + // snapshot so a reconnecting client nets empty after queued then canceled. + let existing = await getSessionMessageState(storage, messageId); + if (!existing) { + await repairMissingQueuedStateFromPendingMessage(pendingMessage); + existing = await getSessionMessageState(storage, messageId); + } + + // Only a still-queued turn may be dropped. `recordRuntimeAcceptedMessage` + // and `ensureAcceptedMessageBeforeTerminal` can write accepted or terminal + // state while the pending row still exists; never drop a running or + // finished turn. + if (existing?.status !== 'queued') { + return { dropped: false }; + } + + await markMessageInterrupted(storage, messageId, { + error: 'Queued message canceled by user', + completionSource: 'canceled', + failureStage: 'interruption', + failureCode: 'user_interrupt', + }); + + await deletePendingSessionMessageByMessageId(storage, messageId); + return { dropped: true }; + } + return { hasMessageAdmission, admitSubmittedMessage, @@ -1218,6 +1273,7 @@ export function createSessionMessageQueue( drainNextPendingMessage, snapshotForStreamConnect, interruptPendingQueuedMessages, + cancelQueuedMessage, recoverPendingInterruption, requestPendingDrain, requestPendingDrainIfNeeded, diff --git a/services/cloud-agent-next/src/session/session-message-state.ts b/services/cloud-agent-next/src/session/session-message-state.ts index 1c9ba035f3..370b29db28 100644 --- a/services/cloud-agent-next/src/session/session-message-state.ts +++ b/services/cloud-agent-next/src/session/session-message-state.ts @@ -46,6 +46,7 @@ export const SessionMessageCompletionSourceSchema = z.enum([ 'wrapper_failure', 'interrupt', 'delivery_failure', + 'canceled', ]); export type SessionMessageCompletionSource = z.infer; diff --git a/services/cloud-agent-next/test/integration/session/pending-messages.test.ts b/services/cloud-agent-next/test/integration/session/pending-messages.test.ts index 55b3a0fd96..46f08668c0 100644 --- a/services/cloud-agent-next/test/integration/session/pending-messages.test.ts +++ b/services/cloud-agent-next/test/integration/session/pending-messages.test.ts @@ -106,6 +106,87 @@ describe('pending session messages', () => { ]); }); + it('cancels a pending message without touching the accepted current message', async () => { + const userId = 'user_pending_cancel_queued'; + const sessionId = 'agent_pending_cancel_queued'; + const pendingMessageId = 'msg_018f1e2d3c4bCancelQueuedAA'; + const currentMessageId = 'msg_018f1e2d3c4bCurrentRunABCD'; + const stub = env.CLOUD_AGENT_SESSION.get( + env.CLOUD_AGENT_SESSION.idFromName(`${userId}:${sessionId}`) + ); + + const result = await runInDurableObject(stub, async instance => { + await registerReadySession(instance, { + sessionId, + userId, + kiloSessionId: '55555555-5555-4555-5555-555555555550', + prompt: 'prepared prompt', + mode: 'code', + model: 'test-model', + kilocodeToken: 'token-cancel-queued', + }); + // Admit the turn so the pending row, queued `SessionMessageState`, and the + // queued replay event all exist before the drop — the exact shape finding + // 1 must terminalize. + const admission = await instance.admitSubmittedMessage( + queueUserMessageInput({ + userId, + prompt: 'queued to cancel', + messageId: pendingMessageId, + mode: 'code', + model: 'test-model', + }) + ); + expect(admission).toMatchObject({ success: true, outcome: 'queued' }); + await putSessionMessageState(instance.ctx.storage, { + messageId: currentMessageId, + status: 'accepted', + prompt: 'current run', + createdAt: 1, + acceptedAt: 2, + wrapperRunId: 'wr_cancel_queued', + }); + + const cancelPending = await instance.cancelQueuedMessage(pendingMessageId); + const cancelMissing = await instance.cancelQueuedMessage('msg_018f1e2d3c4bMissingMsgABC'); + const cancelCurrent = await instance.cancelQueuedMessage(currentMessageId); + + const db = drizzle(instance.ctx.storage, { logger: false }); + const eventQueries = createEventQueries(db, instance.ctx.storage.sql); + const queuedEvents = eventQueries + .findByFilters({ eventTypes: ['cloud.message.queued'] }) + .filter(event => JSON.parse(event.payload).messageId === pendingMessageId); + const canceledEvents = eventQueries + .findByFilters({ eventTypes: ['cloud.message.canceled'] }) + .filter(event => JSON.parse(event.payload).messageId === pendingMessageId); + + return { + cancelPending, + cancelMissing, + cancelCurrent, + pending: await listPendingSessionMessages(instance.ctx.storage), + current: await getSessionMessageState(instance.ctx.storage, currentMessageId), + canceled: await getSessionMessageState(instance.ctx.storage, pendingMessageId), + queuedEvents, + canceledEvents, + }; + }); + + expect(result.cancelPending).toEqual({ dropped: true }); + expect(result.cancelMissing).toEqual({ dropped: false }); + expect(result.cancelCurrent).toEqual({ dropped: false }); + expect(result.pending).toHaveLength(0); + expect(result.current).toMatchObject({ status: 'accepted' }); + // The dropped id's durable state is terminal, not queued, so a re-admit + // rejects it as already terminal instead of ACKing it as still queued. + expect(result.canceled?.status).toBe('interrupted'); + expect(result.canceled?.completionSource).toBe('canceled'); + // The drop persists `cloud.message.canceled` after the queued replay event. + expect(result.queuedEvents).toHaveLength(1); + expect(result.canceledEvents).toHaveLength(1); + expect(result.canceledEvents[0].id).toBeGreaterThan(result.queuedEvents[0].id); + }); + it('finds by clientRequestId', async () => { const userId = 'user_pending_client_request'; const sessionId = 'agent_pending_client_request'; @@ -2753,7 +2834,7 @@ describe('pending session messages', () => { env.CLOUD_AGENT_SESSION.idFromName(`${userId}:${sessionId}`) ); - const result = await runInDurableObject(stub, async (instance, state) => { + const result = await runInDurableObject(stub, async (instance, _state) => { (instance as any).orchestrator = { execute: async () => { throw new Error('wrapper still unavailable'); @@ -2834,7 +2915,7 @@ describe('pending session messages', () => { env.CLOUD_AGENT_SESSION.idFromName(`${userId}:${sessionId}`) ); - const result = await runInDurableObject(stub, async (instance, state) => { + const result = await runInDurableObject(stub, async (instance, _state) => { (instance as any).orchestrator = { execute: async () => { throw new Error('wrapper still unavailable'); diff --git a/services/session-ingest/src/dos/UserConnectionDO.ts b/services/session-ingest/src/dos/UserConnectionDO.ts index 1801331b08..617c2f9422 100644 --- a/services/session-ingest/src/dos/UserConnectionDO.ts +++ b/services/session-ingest/src/dos/UserConnectionDO.ts @@ -108,6 +108,7 @@ const MAX_MUTATION_ID_LENGTH = 128; export const ALLOWED_VIEWER_COMMANDS: ReadonlySet = new Set([ 'send_message', 'interrupt', + 'drop_queued_message', 'question_reply', 'question_reject', 'permission_respond', @@ -129,11 +130,14 @@ const CATALOG_DEDUPE_COMMANDS: ReadonlySet = new Set(['list_models', 'li // Operations that older CLIs reject with a precise "unknown command: " // string. Only these commands get mapped to a structured CLI_UPGRADE_REQUIRED // response; any other CLI error is preserved verbatim. +// Old remotes return upgrade-required for `drop_queued_message`; remove the +// upgrade mapping when every remote supports drop. const CLI_UPGRADE_REQUIRED_COMMANDS: ReadonlySet = new Set([ 'list_commands', 'send_command', 'create_session', 'exit_cli', + 'drop_queued_message', 'list_directories', ]);