diff --git a/apps/buddy/service/src/chat/ChatQueueService.ts b/apps/buddy/service/src/chat/ChatQueueService.ts index bcced3fd..3a382adb 100644 --- a/apps/buddy/service/src/chat/ChatQueueService.ts +++ b/apps/buddy/service/src/chat/ChatQueueService.ts @@ -1,5 +1,6 @@ import type { LocalChatQueueReceipt, LocalChatQueueScope, LocalChatQueueTarget } from '../../../shared/conversation/chatQueueApi' import type { EventSnapshot } from '../../../shared/events/eventTypes' +import type { LocalRun } from '../../../shared/runs/runApi' import type { BuddyAgentRunner } from '../agent/execution/BuddyAgentRunner' import type { BuddyTurnLauncher } from '../agent/execution/BuddyTurnLauncher' import type { BuddyStartTurnInput } from '../BuddyRuntime' @@ -21,7 +22,7 @@ import { resolveTurnExecutionProfile } from '../storage/turnRequestRepository' export interface ChatQueueServiceOptions { queue: ChatQueueRepository - turns: Pick + turns: Pick requests: Pick launcher: Pick runner: Pick @@ -41,6 +42,7 @@ export class ChatQueueService { readonly #options: ChatQueueServiceOptions readonly #operations = new Map>() readonly #enqueues = new Set>() + readonly #cancellations = new Map }>() readonly #changes: Emitter readonly #stopping = new AbortController() readonly onDidChange: Emitter['event'] @@ -61,7 +63,7 @@ export class ChatQueueService { } async drain(): Promise { - await Promise.allSettled([...this.#operations.values(), ...this.#enqueues]) + await Promise.allSettled([...this.#operations.values(), ...this.#enqueues, ...[...this.#cancellations.values()].map(cancellation => cancellation.completion)]) } continuationScopes(): LocalChatQueueScope[] { @@ -79,6 +81,37 @@ export class ChatQueueService { return this.#options.queue.list(scope) } + cancelRun(runId: string): Promise { + const pending = this.#cancellations.get(runId) + if (pending) + return pending.completion + const run = this.#options.runs.findById(runId) + if (!run) + return this.#options.turns.cancel(runId) + const active = this.#options.queue.activeRun(run) + if (active ? active.id !== runId : this.#options.queue.latestRun(run)?.id !== runId || !this.#options.runner.hasActiveExecution(run.conversationId)) + return this.#options.turns.cancel(runId) + const cancellation = Promise.withResolvers() + this.#cancellations.set(runId, { conversationId: run.conversationId, completion: cancellation.promise }) + void this.#cancelRun(run).then(cancellation.resolve, cancellation.reject) + return cancellation.promise + } + + async #cancelRun(run: RunRecord): Promise { + try { + this.pause(run) + return await this.#options.turns.cancel(run.id) + } + finally { + try { + this.pause(run) + } + finally { + this.#cancellations.delete(run.id) + } + } + } + enqueue(input: BuddyStartTurnInput): Promise { const pending = Promise.withResolvers() this.#enqueues.add(pending.promise) @@ -116,6 +149,9 @@ export class ChatQueueService { await stagedAttachments.rollback() throw error } + const latest = this.#options.queue.latestRun(result) + if (this.#isCancelling(result) || latest?.status === 'cancelled' || latest?.status === 'failed') + this.pause(result) this.#changed('enqueued', result, { commitId: randomUUID(), requestId: input.requestId, queueId: result.id, draftReceipt: result.draftReceipt, ...(prepared.attachmentBindings.length ? { attachmentOwnership: { kind: 'queue', attachmentIds: prepared.attachmentBindings.map(binding => binding.id) } } : {}) }) await stagedAttachments.commit().catch(() => undefined) return result @@ -133,6 +169,8 @@ export class ChatQueueService { steer(target: LocalChatQueueTarget) { const active = this.#options.queue.activeRun(target) return this.#serialize(target, async () => { + if (this.#isCancelling(target)) + return false if (!active) return this.#dispatch(target, true) const input = this.#options.queue.pending(target) @@ -144,7 +182,7 @@ export class ChatQueueService { if (!run || run.approvalPolicy !== input.approvalPolicy || run.executionProfile !== resolveTurnExecutionProfile(input)) return false await this.#options.turns.validatePreparedInput(input, false) - if (this.#disposed || this.#options.eventLog.state !== 'open' || !this.#options.queue.pending(target) || this.#options.queue.activeRun(target)?.id !== active.id) + if (this.#disposed || this.#isCancelling(target) || this.#options.eventLog.state !== 'open' || !this.#options.queue.pending(target) || this.#options.queue.activeRun(target)?.id !== active.id) return false return this.#deliver(input, active.id, 'steer') }) @@ -156,7 +194,7 @@ export class ChatQueueService { return false try { return await this.#serialize(run, async () => { - if (signal.aborted || this.#options.queue.activeRun(run)?.id !== runId) + if (signal.aborted || this.#isCancelling(run) || this.#options.queue.activeRun(run)?.id !== runId) return false const next = this.#options.queue.list(run)[0] if (next?.state !== 'waiting') @@ -165,7 +203,7 @@ export class ChatQueueService { if (!input || !this.#canFollowUp(run, input)) return false await this.#options.turns.validatePreparedInput(input) - if (this.#disposed || this.#options.eventLog.state !== 'open' || signal.aborted) + if (this.#disposed || this.#isCancelling(run) || this.#options.eventLog.state !== 'open' || signal.aborted) return false const head = this.#options.queue.list(run)[0] if (this.#options.queue.activeRun(run)?.id !== runId @@ -250,11 +288,15 @@ export class ChatQueueService { } #canDispatch(scope: LocalChatQueueScope): boolean { - return !this.#disposed && !this.#options.runner.isStopping && this.#options.eventLog.state === 'open' + return !this.#disposed && !this.#isCancelling(scope) && !this.#options.runner.isStopping && this.#options.eventLog.state === 'open' && !this.#options.queue.activeRun(scope) && !this.#options.runner.hasActiveExecution(scope.conversationId) } + #isCancelling(scope: LocalChatQueueScope): boolean { + return [...this.#cancellations.values()].some(cancellation => cancellation.conversationId === scope.conversationId) + } + #changed(kind: ChatQueueChange['kind'], scope: LocalChatQueueScope, committed?: ChatQueueChange['committed']): void { this.#changes.fire(copyEventSnapshot({ kind, scope: { conversationId: scope.conversationId, branchId: scope.branchId }, ...(committed ? { committed } : {}) })) } diff --git a/apps/buddy/service/src/chat/__tests__/ChatQueueService.spec.ts b/apps/buddy/service/src/chat/__tests__/ChatQueueService.spec.ts index f4d36c29..dde7b04a 100644 --- a/apps/buddy/service/src/chat/__tests__/ChatQueueService.spec.ts +++ b/apps/buddy/service/src/chat/__tests__/ChatQueueService.spec.ts @@ -9,6 +9,7 @@ import type { TurnRequestCommit } from '../TurnRequestService' import { afterEach, describe, expect, it, vi } from 'vitest' import { createBuddyUserContent } from '../../../../shared/conversation/buddyUserContent' import { Emitter } from '../../../../shared/events/Emitter' +import { toPublicRun } from '../../runs/publicRun' import { createAttachmentRepository } from '../../storage/attachmentRepository' import { createChatQueueRepository } from '../../storage/chatQueueRepository' import { createComposerDraftRepository } from '../../storage/composerDraftRepository' @@ -32,6 +33,94 @@ afterEach(async () => { }) describe('persistent chat queue scheduling', () => { + it('pauses old and new inputs before cancellation settles and requires explicit continuation', async () => { + const f = fixture() + f.queue.enqueue(f.input('A')) + const cleanup = Promise.withResolvers() + f.turns.cancel = async (runId) => { + expect(f.queue.list(f.scope).map(item => item.state)).toEqual(['paused']) + f.database.prepare('UPDATE runs SET status = ? WHERE id = ?').run('cancelled', runId) + await cleanup.promise + return toPublicRun(f.runs.findById(runId)!, null) + } + const cancelling = f.service.cancelRun('run-initial') + expect(f.service.cancelRun('run-initial')).toBe(cancelling) + try { + for (const id of ['B', 'C']) { + const prepared = f.input(id) + f.turns.prepareStart = async () => ({ prepared, stagedAttachments: { validate: () => {}, bindings: [], commit: async () => {}, rollback: async () => {} } }) + await f.service.enqueue({ draftId: 'draft', expectedRevision: prepared.draft.expectedRevision, requestId: prepared.requestId }) + } + f.releaseSlot() + expect(await f.service.steer(f.target('B'))).toBe(false) + expect(await f.service.reconcile(f.scope)).toBe(false) + expect(f.queue.list(f.scope).map(item => [item.id, item.state])).toEqual([['A', 'paused'], ['B', 'paused'], ['C', 'paused']]) + expect(f.launched).toEqual([]) + expect(f.delivered).toEqual([]) + } + finally { + cleanup.resolve() + await cancelling + } + await f.continuation.start() + expect(f.launched).toEqual([]) + expect(await f.service.steer(f.target('B'))).toBe(true) + f.database.prepare('UPDATE runs SET status = ? WHERE id = ?').run('running', 'run-B') + expect(await f.service.followUp('run-B', f.controller.signal)).toBe(false) + expect(f.launched).toEqual(['run-B']) + expect(f.queue.list(f.scope).map(item => [item.id, item.state])).toEqual([['A', 'paused'], ['C', 'paused']]) + expect(f.database.prepare('SELECT id FROM messages ORDER BY rowid').all()).toEqual([{ id: 'initial' }, { id: 'B' }]) + }) + + it('does not pause a newer run queue when a stale stop arrives during its cleanup', async () => { + const f = fixture() + f.database.exec('UPDATE runs SET status = \'completed\'') + f.requests.prepare(f.input('new')) + f.database.exec('UPDATE runs SET status = \'completed\'') + f.queue.enqueue(f.input('A')) + await f.service.cancelRun('run-initial') + expect(f.queue.list(f.scope).map(item => [item.id, item.state])).toEqual([['A', 'waiting']]) + }) + + it('withdraws a steering input whose validation crosses the stop request', async () => { + const f = fixture() + f.queue.enqueue(f.input('A')) + const validation = Promise.withResolvers() + const cleanup = Promise.withResolvers() + f.validate.mockImplementationOnce(() => validation.promise) + f.turns.cancel = async (runId) => { + await cleanup.promise + return toPublicRun(f.runs.findById(runId)!, null) + } + const steering = f.service.steer(f.target('A')) + await vi.waitFor(() => expect(f.validate).toHaveBeenCalledOnce()) + const cancelling = f.service.cancelRun('run-initial') + try { + validation.resolve() + expect(await steering).toBe(false) + expect(f.delivered).toEqual([]) + expect(f.queue.list(f.scope)[0]?.state).toBe('paused') + } + finally { + validation.resolve() + cleanup.resolve() + await cancelling + } + }) + + it('preserves inputs after a cancellation failure and releases the continuation guard', async () => { + const f = fixture() + f.queue.enqueue(f.input('A')) + f.turns.cancel = async () => { + throw new Error('Cancellation failed') + } + await expect(f.service.cancelRun('run-initial')).rejects.toThrow('Cancellation failed') + expect(f.queue.list(f.scope)[0]?.state).toBe('paused') + expect(await f.service.followUp('run-initial', f.controller.signal)).toBe(false) + expect(await f.service.steer(f.target('A'))).toBe(true) + expect(f.delivered).toEqual(['steer:A']) + }) + it('commits queue and message ownership separately and consumes a queued draft only once', async () => { const f = fixture() const attachments = createAttachmentRepository(f.database) @@ -392,6 +481,7 @@ function fixture(queuedOnStartup = false) { const committed = new Emitter(() => {}) const eventLog: { onDidCommit: RunEventObservation['onDidCommit'], state: RunEventObservation['state'] } = { onDidCommit: committed.event, state: 'open' } const turns: ChatQueueServiceOptions['turns'] = { + cancel: async runId => toPublicRun(runs.findById(runId)!, null), prepareStart: async () => { throw new Error('Use the persisted fixture') }, diff --git a/apps/buddy/service/src/chat/registerChatRpc.ts b/apps/buddy/service/src/chat/registerChatRpc.ts index b9f75b19..c4c3ca6a 100644 --- a/apps/buddy/service/src/chat/registerChatRpc.ts +++ b/apps/buddy/service/src/chat/registerChatRpc.ts @@ -17,7 +17,7 @@ export interface RegisterChatRpcOptions { runtime: BuddyRuntime turns: Pick< ChatTurnService, - 'cancel' | 'editUserMessage' | 'regenerateAssistant' + 'editUserMessage' | 'regenerateAssistant' > } @@ -40,7 +40,7 @@ export function registerChatRpc(options: RegisterChatRpcOptions): () => void { options.turns.regenerateAssistant(params) )), registerRuntimeRequest(options.rpc, chatRpc.cancel, (input) => { - return options.turns.cancel(input.runId) + return options.queue.cancelRun(input.runId) }), ] return () => disposers.splice(0).forEach(dispose => dispose()) diff --git a/apps/buddy/src/modules/tasks/model/transcript/__tests__/chatStreamingMessage.spec.ts b/apps/buddy/src/modules/tasks/model/transcript/__tests__/chatStreamingMessage.spec.ts index 6cb21b74..adcb1a9b 100644 --- a/apps/buddy/src/modules/tasks/model/transcript/__tests__/chatStreamingMessage.spec.ts +++ b/apps/buddy/src/modules/tasks/model/transcript/__tests__/chatStreamingMessage.spec.ts @@ -511,6 +511,20 @@ describe('projectStreamingAssistantMessage', () => { ]) }) + it('finalizes an awaiting tool when its run is cancelled before the approval event catches up', () => { + const [turn] = projectChatAgentTurns([ + event(1, 'approval.requested', { + id: 'approval-1', + kind: 'shell', + review: { card: 'shell', command: 'echo fixture', toolName: 'bash' }, + status: 'pending', + summary: 'Run a command', + toolCallId: 'tool-1', + }), + ], [run('cancelled')]) + expect(turn?.nodes).toEqual([expect.objectContaining({ kind: 'tool', status: 'interrupted', toolCallId: 'tool-1' })]) + }) + it('does not let a late tool start overwrite a pending approval', () => { const turns = projectChatRunProcessesForTest([ event(1, 'approval.requested', { diff --git a/apps/buddy/src/modules/tasks/model/transcript/chatAgentTurn.ts b/apps/buddy/src/modules/tasks/model/transcript/chatAgentTurn.ts index e7b179f1..68160a1d 100644 --- a/apps/buddy/src/modules/tasks/model/transcript/chatAgentTurn.ts +++ b/apps/buddy/src/modules/tasks/model/transcript/chatAgentTurn.ts @@ -492,7 +492,8 @@ export function createChatAgentTurnReducer( if ( node.kind === 'text' || !terminal - || (node.status !== 'preparing' && node.status !== 'running') + || (node.status !== 'preparing' && node.status !== 'running' + && !(run.status === 'cancelled' && node.kind === 'tool' && node.status === 'awaiting_approval')) ) { return node } diff --git a/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatTurnExecution.spec.ts b/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatTurnExecution.spec.ts index 971596a1..a14c173d 100644 --- a/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatTurnExecution.spec.ts +++ b/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatTurnExecution.spec.ts @@ -15,6 +15,37 @@ const scopes: ReturnType[] = [] afterEach(() => scopes.splice(0).forEach(scope => scope.stop())) describe('useChatTurnExecution cancellation ownership', () => { + it('retains cancellation feedback and queues input until cleanup finishes after the terminal event', async () => { + const f = createFixture() + f.selectedModel.value = { modelId: 'model-a', providerId: 'provider-a' } as LocalRuntimeModelOption + const snapshot = f.drafts.snapshot('conversation:conversation-a:branch-a') + f.drafts.confirmOpen(snapshot, { + content: snapshot.content, + draftId: snapshot.draftId, + executionConfig: { approvalPolicy: snapshot.approvalPolicy, executionProfile: snapshot.executionProfile }, + modelSelection: null, + revision: 1, + scope: { kind: 'conversation_branch', conversationId: 'conversation-a', branchId: 'branch-a' }, + updatedAt: '2026-09-08T00:00:00.000Z', + }) + f.api.chat.enqueue.mockResolvedValue({ id: 'queued-after-stop', conversationId: 'conversation-a', branchId: 'branch-a', draftReceipt: { draftId: snapshot.draftId, sourceRevision: 1, committedRevision: 2 } }) + const cancelling = f.execution.cancelActiveRun() + const terminal = { ...f.run, status: 'cancelled' as const } + f.projectedRuns.value = [terminal] + expect(f.execution.stoppingRunId.value).toBe(f.run.id) + expect(await f.execution.send('pending input')).toBe(true) + expect(f.drafts.draft.value).toBe('') + expect(f.api.chat.enqueue).toHaveBeenCalledOnce() + expect(f.api.chat.startTurn).not.toHaveBeenCalled() + f.drafts.updateComposerContent('next draft', null) + f.pending.resolve(terminal) + await cancelling + expect(f.execution.stoppingRunId.value).toBeNull() + expect(f.projectedRuns.value).toEqual([terminal]) + expect(f.drafts.draft.value).toBe('next draft') + expect(f.api.chat.steerQueued).not.toHaveBeenCalled() + }) + it('applies cancellation to the current projection and preserves its Draft', async () => { const fixture = createFixture() const cancelling = fixture.execution.cancelActiveRun() @@ -67,7 +98,7 @@ describe('useChatTurnExecution cancellation ownership', () => { const fixture = createFixture() const cancelling = fixture.execution.cancelActiveRun() fixture.navigate(navigation) - expect(fixture.execution.stoppingRunId.value).toBeNull() + expect(fixture.execution.stoppingRunId.value).toBe(navigation.startsWith('return-to-') ? fixture.run.id : null) fixture.pending.resolve({ ...fixture.run, status: 'cancelled' }) await cancelling diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts b/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts index b6a8a531..36726a16 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts @@ -68,10 +68,12 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { const isSending = shallowRef(false) const requestIds = createRequestIdRegistry() const pendingCancellationWatches = new Set<() => void>() - const pendingCancellationIds = shallowReactive(new Set()) + const pendingCancellations = shallowReactive(new Map()) const stoppingRunId = computed(() => { - const run = options.activeRun.value - return run && pendingCancellationIds.has(run.id) ? run.id : null + const run = [...pendingCancellations.values()].find(run => + run.conversationId === options.session.activeConversationId.value + && run.branchId === options.session.activeBranchId.value) + return run?.id ?? null }) let isDisposed = false onScopeDispose(() => { @@ -79,7 +81,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { for (const stop of pendingCancellationWatches) stop() pendingCancellationWatches.clear() - pendingCancellationIds.clear() + pendingCancellations.clear() }, true) const canSend = computed(() => options.runtimeSupervisor.runtimeState.value.status === 'ready' @@ -110,7 +112,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { options.setErrorMessage(options.unavailableCommandMessage()) return false } - if (isRunCommand && options.activeRun.value) + if (isRunCommand && (options.activeRun.value || stoppingRunId.value)) return false if (isRunCommand && command) return executeActionCommand(command, contextItems) @@ -133,7 +135,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { const expectedRevision = confirmedDraft.revision! const operationKey = `turn:${confirmedDraft.draftId}:${expectedRevision}` const requestId = requestIds.resolve(operationKey) - if (options.activeRun.value || queue.queuedMessages.value.length) { + if (stoppingRunId.value || options.activeRun.value || queue.queuedMessages.value.length) { const result = await options.api.chat.enqueue({ draftId: confirmedDraft.draftId, expectedRevision, requestId }) requestIds.release(operationKey) options.composerTarget.complete(result.draftReceipt, sourceScopeKey, sourceScopeKey) @@ -264,9 +266,9 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { async function cancelActiveRun() { const run = options.activeRun.value - if (!run || isDisposed || pendingCancellationIds.has(run.id)) + if (!run || isDisposed || pendingCancellations.has(run.id)) return - pendingCancellationIds.add(run.id) + pendingCancellations.set(run.id, run) const navigationVersion = options.session.generation() let sourceViewChanged = false const stopWatchingView = watch( @@ -293,7 +295,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { finally { stopWatchingView() pendingCancellationWatches.delete(stopWatchingView) - pendingCancellationIds.delete(run.id) + pendingCancellations.delete(run.id) } } diff --git a/apps/buddy/src/modules/tasks/widgets/composer/DesktopChatComposer.vue b/apps/buddy/src/modules/tasks/widgets/composer/DesktopChatComposer.vue index ffccad64..19ae1a82 100644 --- a/apps/buddy/src/modules/tasks/widgets/composer/DesktopChatComposer.vue +++ b/apps/buddy/src/modules/tasks/widgets/composer/DesktopChatComposer.vue @@ -317,18 +317,27 @@ function captureDraft(): WorkbenchMenuSelection { {{ modelInputIssueMessage || t('desktop.chat.queueAdd') }} - -