From 4afa6947a677558876414742c26db41a88d65f8d Mon Sep 17 00:00:00 2001 From: QAyong Date: Thu, 8 Oct 2026 00:44:37 +0800 Subject: [PATCH 1/2] =?UTF-8?q?fix(buddy):=20=E7=82=B9=E5=87=BB=E5=81=9C?= =?UTF-8?q?=E6=AD=A2=E5=90=8E=E7=AB=8B=E5=8D=B3=E7=BB=93=E6=9D=9F=E8=BF=90?= =?UTF-8?q?=E8=A1=8C=E5=B1=95=E7=A4=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../tasks/model/transcript/chatAgentTurn.ts | 3 +- .../runs/__tests__/useChatRunSync.spec.ts | 48 +++++++++++ .../__tests__/useChatTurnExecution.spec.ts | 44 +++++++++- .../src/modules/tasks/state/runs/typing.ts | 3 + .../tasks/state/runs/useChatRunProjection.ts | 82 ++++++++++++++++--- .../tasks/state/runs/useChatRunSync.ts | 3 + .../tasks/state/runs/useChatTurnExecution.ts | 35 +++++++- .../modules/tasks/state/useTaskCapability.ts | 4 +- 8 files changed, 204 insertions(+), 18 deletions(-) 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__/useChatRunSync.spec.ts b/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts index 9fcf6015..24643b47 100644 --- a/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts +++ b/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts @@ -6,9 +6,57 @@ import { describe, expect, it, vi } from 'vitest' import { effectScope, ref } from 'vue' import { projectPersistedChatTranscriptRows } from '../../../model/transcript/chatPersistedTranscriptRows' +import { projectChatTranscript } from '../../../model/transcript/chatTranscriptProjection' import { useChatRunSync } from '../useChatRunSync' describe('useChatRunSync', () => { + it('stops presentation immediately while retaining execution ownership and reconciles the final result', async () => { + const running = run('run-a', 'conversation-a') + const initialEvents: LocalRunEvent[] = [ + { ...event(running.id, 1), type: 'message.started', payload: { messageId: 'answer' } }, + { ...event(running.id, 2), type: 'message.delta', payload: { messageId: 'answer', delta: 'Visible answer', phase: 'answer' } }, + ] + const question = timelineMessage(running.triggeringMessageId, running.conversationId, running.branchId, 1) + let page = timelinePage([question], null, [running], initialEvents) + const sync = useChatRunSync({ activeBranchId: ref(running.branchId), activeConversationId: ref(running.conversationId), api: createApi({ listTimeline: async () => page }), onError: vi.fn() }) + try { + await sync.refreshActiveConversation() + sync.cancelRunPresentation(running.id) + expect(sync.runs.value[0]).toMatchObject({ status: 'cancelled', errorCode: 'RUN_CANCELLED' }) + expect(sync.executionRuns.value[0]?.status).toBe('running') + expect(sync.messages.value.find(message => message.id === 'answer')?.content).toEqual({ text: 'Visible answer' }) + const rows = projectChatTranscript({ timelineItems: sync.timelineItems.value, runs: sync.runs.value, runEvents: sync.runEventBuckets.value.get(running.id)?.events, outputs: [] }).rows + expect(rows.some(row => row.kind === 'activity')).toBe(false) + const answer = rows.find(row => row.kind === 'message' && row.message.id === 'answer') + expect(answer?.kind === 'message' ? answer.streaming : true).toBeUndefined() + + page = timelinePage([question], null, [running], [...initialEvents, { ...event(running.id, 3), type: 'message.delta', payload: { messageId: 'answer', delta: ' late text', phase: 'answer' } }]) + await sync.refreshActiveConversation() + expect(sync.runs.value[0]?.status).toBe('cancelled') + expect(sync.messages.value.find(message => message.id === 'answer')?.content).toEqual({ text: 'Visible answer' }) + + const final = { ...timelineMessage('answer', running.conversationId, running.branchId, 2, 'assistant'), runId: running.id, content: { text: 'Final saved answer' } } + page = timelinePage([question, final], null, [{ ...running, status: 'cancelled', completedAt: '2026-08-14T00:00:03.000Z' }]) + await sync.refreshActiveConversation() + expect(sync.executionRuns.value[0]?.status).toBe('cancelled') + expect(sync.messages.value.find(message => message.id === 'answer')?.content).toEqual(final.content) + } + finally { sync.dispose() } + }) + + it('restores authoritative presentation if the cancellation request fails', async () => { + const running = run('run-a', 'conversation-a') + const sync = useChatRunSync({ activeBranchId: ref(running.branchId), activeConversationId: ref(running.conversationId), api: createApi({}), onError: vi.fn() }) + try { + await sync.refreshActiveConversation() + sync.cancelRunPresentation(running.id) + expect(sync.runs.value[0]?.status).toBe('cancelled') + sync.restoreRunPresentation(running.id) + expect(sync.runs.value[0]?.status).toBe('running') + } + finally { sync.dispose() } + }) + it('refreshes a running action on an older page and removes only automatic skipped actions from the transcript', async () => { const action: Extract = { kind: 'extension-action', 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..fd20686f 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 @@ -18,6 +18,7 @@ describe('useChatTurnExecution cancellation ownership', () => { it('applies cancellation to the current projection and preserves its Draft', async () => { const fixture = createFixture() const cancelling = fixture.execution.cancelActiveRun() + expect(fixture.presentation.cancelRunPresentation).toHaveBeenCalledWith(fixture.run.id) expect(fixture.execution.stoppingRunId.value).toBe(fixture.run.id) expect(fixture.projectedRuns.value[0]?.status).toBe('running') fixture.pending.resolve({ ...fixture.run, status: 'cancelled' }) @@ -39,11 +40,48 @@ describe('useChatTurnExecution cancellation ownership', () => { expect(fixture.execution.stoppingRunId.value).toBeNull() expect(fixture.error.value).toBeTruthy() + expect(fixture.presentation.restoreRunPresentation).toHaveBeenCalledWith(fixture.run.id) expect(fixture.projectedRuns.value).toEqual([fixture.run]) expect(fixture.drafts.draft.value).toBe('pending input') expect(fixture.execution.isSending.value).toBe(false) }) + it('accepts a fresh message immediately and starts it only after cancellation finishes', 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', + }) + const receipt = { + id: 'fresh-message', + conversationId: 'conversation-a', + branchId: 'branch-a', + draftReceipt: { draftId: snapshot.draftId, sourceRevision: 1, committedRevision: 2 }, + } + f.api.chat.enqueue.mockResolvedValue(receipt) + f.api.chat.steerQueued.mockResolvedValue(true) + const cancelling = f.execution.cancelActiveRun() + await f.execution.cancelActiveRun() + expect(f.api.chat.cancel).toHaveBeenCalledTimes(1) + + expect(await f.execution.send('pending input')).toBe(true) + expect(f.api.chat.enqueue).toHaveBeenCalledTimes(1) + expect(f.api.chat.startTurn).not.toHaveBeenCalled() + expect(f.api.chat.steerQueued).not.toHaveBeenCalled() + expect(f.execution.isSending.value).toBe(false) + + f.pending.resolve({ ...f.run, status: 'cancelled' }) + await cancelling + await vi.waitFor(() => expect(f.api.chat.steerQueued).toHaveBeenCalledWith({ id: receipt.id, conversationId: receipt.conversationId, branchId: receipt.branchId })) + }) + it.each(['success', 'error'] as const)('ignores a late cancellation %s after its owner is disposed', async (outcome) => { const fixture = createFixture() const cancelling = fixture.execution.cancelActiveRun() @@ -153,7 +191,8 @@ function createFixture() { const error = shallowRef(null) const pending = deferred() const composerTarget = useComposerTarget({ drafts, conversationId: session.activeConversationId, branchId: session.activeBranchId, persist: async () => true }) - const api = { chat: { listQueue: async () => [], enqueue: vi.fn(), cancelQueued: vi.fn(), steerQueued: vi.fn(), cancel: () => pending.promise, executeCommand: vi.fn(), startTurn: vi.fn() } } + const api = { chat: { listQueue: async () => [], enqueue: vi.fn(), cancelQueued: vi.fn(), steerQueued: vi.fn(), cancel: vi.fn(() => pending.promise), executeCommand: vi.fn(), startTurn: vi.fn() } } + const presentation = { cancelRunPresentation: vi.fn(), restoreRunPresentation: vi.fn() } const selectedModel = shallowRef(null) const execution = useChatTurnExecution({ composerTarget, @@ -172,6 +211,7 @@ function createFixture() { onActionCommandRunStarted: () => {}, persistWorkspaceState: async () => true, runSync: { + ...presentation, refreshActiveConversation: async () => {}, applyRunStart: () => {}, upsertRuns: (runs) => { @@ -201,7 +241,7 @@ function createFixture() { drafts.updateComposerContent('current view input', null) error.value = 'current view status' } - return { api, drafts, error, execution, navigate, pending, projectedRuns, run, scope, selectedModel } + return { api, drafts, error, execution, navigate, pending, presentation, projectedRuns, run, scope, selectedModel } })! } diff --git a/apps/buddy/src/modules/tasks/state/runs/typing.ts b/apps/buddy/src/modules/tasks/state/runs/typing.ts index 551d47b5..34d0ce18 100644 --- a/apps/buddy/src/modules/tasks/state/runs/typing.ts +++ b/apps/buddy/src/modules/tasks/state/runs/typing.ts @@ -30,6 +30,9 @@ export interface ChatRunProjectionState { } export interface ChatRunSync extends ChatRunProjectionState { + executionRuns: Readonly>> + cancelRunPresentation: (runId: string) => void + restoreRunPresentation: (runId: string) => void isLoadingConversation: Readonly> isLoadingOlderMessages: Readonly> applyEditedTurn: (turn: LocalTurnStart, userMessageId: string) => void diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts b/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts index 26ff33ee..fc95c894 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts @@ -1,10 +1,10 @@ import type { LocalChangeSetSummary } from '@buddy-shared/changes/changeApi' -import type { LocalConversationTimelineItem, LocalConversationTimelinePage, LocalMessage } from '@buddy-shared/conversation/conversationApi' +import type { LocalConversationTimelineItem, LocalConversationTimelinePage } from '@buddy-shared/conversation/conversationApi' import type { LocalApproval } from '@buddy-shared/permissions/approvalApi' import type { LocalRun, LocalRunEvent, LocalRunOutput } from '@buddy-shared/runs/runApi' -import type { ChatRunEventBuckets } from '../../model/runs/typing' +import type { ChatRunEventBucket, ChatRunEventBuckets } from '../../model/runs/typing' import type { ChatRunProjectionState } from './typing' -import { computed, shallowRef } from 'vue' +import { computed, shallowReactive, shallowRef } from 'vue' import { hasChatRunEventSequenceGap, mergeChatRunEventBuckets, @@ -21,12 +21,10 @@ import { mergeTimelineEvents, timelineItemKey, } from '../../model/runs/chatTimelineMerge' +import { projectChatRunStreamingMessages } from '../../model/transcript/chatRunStreamingMessages' export function useChatRunProjection() { const timelineItems = shallowRef>([]) - const messages = computed>(() => timelineItems.value.filter( - (item): item is Extract => item.kind === 'message', - )) const runs = shallowRef>([]) const runSignalEvents = shallowRef>([]) const runEventBuckets = shallowRef(new Map()) @@ -36,8 +34,65 @@ export function useChatRunProjection() { const timelineCursor = shallowRef(null) const hasOlderMessages = computed(() => timelineCursor.value !== null) const knownRunIds = new Set() + // Presentation stops immediately; authoritative runs still own execution until cleanup finishes. + const cancelledPresentations = shallowReactive(new Map + }>()) let hasLoadedTimelinePage = false + const presentedRuns = computed(() => runs.value.map(run => cancelledPresentations.get(run.id)?.run ?? run)) + const presentedBuckets = computed(() => { + if (!cancelledPresentations.size) + return runEventBuckets.value + const buckets = new Map(runEventBuckets.value) + for (const [runId, presentation] of cancelledPresentations) { + if (knownRunIds.has(runId)) + buckets.set(runId, presentation.bucket) + } + return buckets + }) + const presentedTimeline = computed(() => { + const items = timelineItems.value.flatMap((item) => { + if (item.kind !== 'message' || !item.runId) + return [item] + const presentation = cancelledPresentations.get(item.runId) + if (!presentation) + return [item] + const frozen = presentation.messages.get(item.id) + return frozen ? [frozen] : [] + }) + const frozenMessages = [...cancelledPresentations].flatMap(([runId, presentation]) => + knownRunIds.has(runId) ? [...presentation.messages.values()] : []) + return mergeTailTimelineItems(items, frozenMessages) + }) + + function cancelRunPresentation(runId: string) { + const run = runs.value.find(run => run.id === runId) + if (!run || (run.status !== 'queued' && run.status !== 'running') || cancelledPresentations.has(runId)) + return + const completedAt = new Date().toISOString() + const bucket = runEventBuckets.value.get(runId) + const events = bucket?.events ?? [] + // Keep the visible partial answer when terminal projection stops accepting streaming deltas. + const messages = new Map(timelineItems.value.flatMap(item => + item.kind === 'message' && item.runId === runId ? [[item.id, item] as const] : [])) + for (const candidate of projectChatRunStreamingMessages(run, events)) { + if (!messages.has(candidate.message.id)) + messages.set(candidate.message.id, { ...candidate.message, kind: 'message' }) + } + cancelledPresentations.set(runId, { + run: { ...run, status: 'cancelled', completedAt, errorCode: 'RUN_CANCELLED' }, + bucket: { events, revision: (bucket?.revision ?? 0) + 1, update: null }, + messages, + }) + } + + function restoreRunPresentation(runId: string) { + cancelledPresentations.delete(runId) + } + function mergePage(page: LocalConversationTimelinePage) { upsertRuns(page.runs) const events = mergeTimelineEvents( @@ -87,6 +142,8 @@ export function useChatRunProjection() { for (const run of incoming) { byId.set(run.id, run) knownRunIds.add(run.id) + if (run.status !== 'queued' && run.status !== 'running') + restoreRunPresentation(run.id) } runs.value = [...byId.values()].sort((left, right) => right.startedAt.localeCompare(left.startedAt)) } @@ -114,19 +171,22 @@ export function useChatRunProjection() { } const state: ChatRunProjectionState = { - approvals, + approvals: computed(() => approvals.value.filter(approval => !cancelledPresentations.has(approval.runId))), changeSets, hasOlderMessages, - messages, - runEventBuckets, + messages: computed(() => presentedTimeline.value.filter((item): item is Extract => item.kind === 'message')), + runEventBuckets: presentedBuckets, runOutputs, - runs, + runs: presentedRuns, runSignalEvents, - timelineItems, + timelineItems: presentedTimeline, } return { state, + executionRuns: runs, + cancelRunPresentation, + restoreRunPresentation, appendEvents, applySnapshot, clear, diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts b/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts index d3655d92..1bcd6a7d 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts @@ -307,6 +307,9 @@ export function useChatRunSync(options: ChatRunSyncOptions): ChatRunSync { return { ...projection.state, + executionRuns: projection.executionRuns, + cancelRunPresentation: projection.cancelRunPresentation, + restoreRunPresentation: projection.restoreRunPresentation, applyEditedTurn: (turn, messageId) => applyReplacementTurn(turn, messageId, false), applyRegeneratedTurn: turn => applyReplacementTurn(turn, turn.run.triggeringMessageId, true), applyRunStart, diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts b/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts index b6a8a531..a0950f1e 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts @@ -2,6 +2,7 @@ import type { LocalChatApi } from '@buddy-electron/shared/localChatApi' import type { ParsedBuddyChatCommand } from '@buddy-shared/conversation/buddyChatCommands' import type { BuddyUserContentV1 } from '@buddy-shared/conversation/buddyUserContent' import type { LocalPromptContextItem } from '@buddy-shared/conversation/chatApi' +import type { LocalChatQueueTarget } from '@buddy-shared/conversation/chatQueueApi' import type { BuddyApprovalPolicy } from '@buddy-shared/permissions/approvalPolicy' import type { BuddyExecutionProfile } from '@buddy-shared/permissions/executionProfile' import type { LocalRun } from '@buddy-shared/runs/runApi' @@ -57,7 +58,7 @@ export interface UseChatTurnExecutionOptions { onActionCommandRunStarted: (runId: string) => void onDraftCommitted?: (draftId: string, conversationId: string) => void persistWorkspaceState: () => Promise - runSync: Pick + runSync: Pick runtimeSupervisor: Pick setErrorMessage: (message: string | null) => void unavailableCommandMessage: () => string @@ -69,6 +70,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { const requestIds = createRequestIdRegistry() const pendingCancellationWatches = new Set<() => void>() const pendingCancellationIds = shallowReactive(new Set()) + const cancellationCompletions = new Map }>() const stoppingRunId = computed(() => { const run = options.activeRun.value return run && pendingCancellationIds.has(run.id) ? run.id : null @@ -80,6 +82,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { stop() pendingCancellationWatches.clear() pendingCancellationIds.clear() + cancellationCompletions.clear() }, true) const canSend = computed(() => options.runtimeSupervisor.runtimeState.value.status === 'ready' @@ -116,6 +119,9 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { return executeActionCommand(command, contextItems) const sourceScopeKey = options.draftScopeKey.value + const cancellation = [...cancellationCompletions.values()].find(({ run }) => + run.conversationId === options.session.activeConversationId.value + && run.branchId === options.session.activeBranchId.value) const navigationVersion = options.session.generation() const isSourceViewCurrent = () => options.session.isCurrent(navigationVersion) && options.draftScopeKey.value === sourceScopeKey @@ -133,10 +139,12 @@ 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 (cancellation || 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) + if (cancellation) + void dispatchAfterCancellation(result, cancellation.completion) await queue.refreshQueue() return true } @@ -267,6 +275,9 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { if (!run || isDisposed || pendingCancellationIds.has(run.id)) return pendingCancellationIds.add(run.id) + const cancellation = Promise.withResolvers() + cancellationCompletions.set(run.id, { run, completion: cancellation.promise }) + options.runSync.cancelRunPresentation(run.id) const navigationVersion = options.session.generation() let sourceViewChanged = false const stopWatchingView = watch( @@ -285,8 +296,11 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { const cancelled = await options.api.chat.cancel(run.id) if (isSourceViewCurrent()) options.runSync.upsertRuns([cancelled]) + cancellation.resolve(cancelled.status === 'cancelled') } catch (error) { + options.runSync.restoreRunPresentation(run.id) + cancellation.resolve(false) if (isSourceViewCurrent()) setNormalizedError(error) } @@ -294,6 +308,23 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { stopWatchingView() pendingCancellationWatches.delete(stopWatchingView) pendingCancellationIds.delete(run.id) + cancellationCompletions.delete(run.id) + } + } + + async function dispatchAfterCancellation(target: LocalChatQueueTarget, completion: Promise) { + if (!await completion || isDisposed) + return + try { + // Only the message explicitly sent after Stop resumes; older paused inputs stay paused. + await options.api.chat.steerQueued({ id: target.id, conversationId: target.conversationId, branchId: target.branchId }) + await queue.refreshQueue() + } + catch (error) { + if (!isDisposed && options.session.activeConversationId.value === target.conversationId + && options.session.activeBranchId.value === target.branchId) { + setNormalizedError(error) + } } } diff --git a/apps/buddy/src/modules/tasks/state/useTaskCapability.ts b/apps/buddy/src/modules/tasks/state/useTaskCapability.ts index 3f7d7c2c..676817d7 100644 --- a/apps/buddy/src/modules/tasks/state/useTaskCapability.ts +++ b/apps/buddy/src/modules/tasks/state/useTaskCapability.ts @@ -100,7 +100,7 @@ export function useTaskCapability(options: UseTaskCapabilityOptions): TaskCapabi timelineItems, } = runSync - const activeRun = computed(() => runs.value.find( + const activeRun = computed(() => runSync.executionRuns.value.find( run => run.status === 'queued' || run.status === 'running', ) ?? null) const hasAvailableProvider = computed(() => modelProviders.providers.value.some( @@ -493,7 +493,7 @@ export function useTaskCapability(options: UseTaskCapabilityOptions): TaskCapabi }, }, execution: { - activeRun: readonly(activeRun), + activeRun: computed(() => runs.value.find(run => run.status === 'queued' || run.status === 'running') ?? null), approvalViews: readonly(approvalViews), canMutateBranch: readonly(canMutateBranch), canSend: readonly(canSend), From 7cf03d7531f31093bb37d42a226c767805f7ad80 Mon Sep 17 00:00:00 2001 From: shanyuhai123 <864299347@qq.com> Date: Thu, 8 Oct 2026 17:14:42 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix(buddy):=20=E7=A8=B3=E5=AE=9A=E5=81=9C?= =?UTF-8?q?=E6=AD=A2=E5=8F=8D=E9=A6=88=E4=B8=8E=E9=98=9F=E5=88=97=E5=8F=96?= =?UTF-8?q?=E6=B6=88?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../service/src/chat/ChatQueueService.ts | 54 +++++++++-- .../chat/__tests__/ChatQueueService.spec.ts | 90 +++++++++++++++++++ .../buddy/service/src/chat/registerChatRpc.ts | 4 +- .../__tests__/chatStreamingMessage.spec.ts | 14 +++ .../runs/__tests__/useChatRunSync.spec.ts | 48 ---------- .../__tests__/useChatTurnExecution.spec.ts | 77 +++++++--------- .../src/modules/tasks/state/runs/typing.ts | 3 - .../tasks/state/runs/useChatRunProjection.ts | 82 +++-------------- .../tasks/state/runs/useChatRunSync.ts | 3 - .../tasks/state/runs/useChatTurnExecution.ts | 53 +++-------- .../modules/tasks/state/useTaskCapability.ts | 4 +- .../widgets/composer/DesktopChatComposer.vue | 45 +++++++--- .../__tests__/chatComposerProps.spec.ts | 1 + .../modules/tasks/widgets/composer/typing.ts | 1 + .../modules/tasks/widgets/workspace/typing.ts | 2 +- .../widgets/workspace/useTaskComposer.ts | 3 +- 16 files changed, 252 insertions(+), 232 deletions(-) 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/state/runs/__tests__/useChatRunSync.spec.ts b/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts index 24643b47..9fcf6015 100644 --- a/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts +++ b/apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts @@ -6,57 +6,9 @@ import { describe, expect, it, vi } from 'vitest' import { effectScope, ref } from 'vue' import { projectPersistedChatTranscriptRows } from '../../../model/transcript/chatPersistedTranscriptRows' -import { projectChatTranscript } from '../../../model/transcript/chatTranscriptProjection' import { useChatRunSync } from '../useChatRunSync' describe('useChatRunSync', () => { - it('stops presentation immediately while retaining execution ownership and reconciles the final result', async () => { - const running = run('run-a', 'conversation-a') - const initialEvents: LocalRunEvent[] = [ - { ...event(running.id, 1), type: 'message.started', payload: { messageId: 'answer' } }, - { ...event(running.id, 2), type: 'message.delta', payload: { messageId: 'answer', delta: 'Visible answer', phase: 'answer' } }, - ] - const question = timelineMessage(running.triggeringMessageId, running.conversationId, running.branchId, 1) - let page = timelinePage([question], null, [running], initialEvents) - const sync = useChatRunSync({ activeBranchId: ref(running.branchId), activeConversationId: ref(running.conversationId), api: createApi({ listTimeline: async () => page }), onError: vi.fn() }) - try { - await sync.refreshActiveConversation() - sync.cancelRunPresentation(running.id) - expect(sync.runs.value[0]).toMatchObject({ status: 'cancelled', errorCode: 'RUN_CANCELLED' }) - expect(sync.executionRuns.value[0]?.status).toBe('running') - expect(sync.messages.value.find(message => message.id === 'answer')?.content).toEqual({ text: 'Visible answer' }) - const rows = projectChatTranscript({ timelineItems: sync.timelineItems.value, runs: sync.runs.value, runEvents: sync.runEventBuckets.value.get(running.id)?.events, outputs: [] }).rows - expect(rows.some(row => row.kind === 'activity')).toBe(false) - const answer = rows.find(row => row.kind === 'message' && row.message.id === 'answer') - expect(answer?.kind === 'message' ? answer.streaming : true).toBeUndefined() - - page = timelinePage([question], null, [running], [...initialEvents, { ...event(running.id, 3), type: 'message.delta', payload: { messageId: 'answer', delta: ' late text', phase: 'answer' } }]) - await sync.refreshActiveConversation() - expect(sync.runs.value[0]?.status).toBe('cancelled') - expect(sync.messages.value.find(message => message.id === 'answer')?.content).toEqual({ text: 'Visible answer' }) - - const final = { ...timelineMessage('answer', running.conversationId, running.branchId, 2, 'assistant'), runId: running.id, content: { text: 'Final saved answer' } } - page = timelinePage([question, final], null, [{ ...running, status: 'cancelled', completedAt: '2026-08-14T00:00:03.000Z' }]) - await sync.refreshActiveConversation() - expect(sync.executionRuns.value[0]?.status).toBe('cancelled') - expect(sync.messages.value.find(message => message.id === 'answer')?.content).toEqual(final.content) - } - finally { sync.dispose() } - }) - - it('restores authoritative presentation if the cancellation request fails', async () => { - const running = run('run-a', 'conversation-a') - const sync = useChatRunSync({ activeBranchId: ref(running.branchId), activeConversationId: ref(running.conversationId), api: createApi({}), onError: vi.fn() }) - try { - await sync.refreshActiveConversation() - sync.cancelRunPresentation(running.id) - expect(sync.runs.value[0]?.status).toBe('cancelled') - sync.restoreRunPresentation(running.id) - expect(sync.runs.value[0]?.status).toBe('running') - } - finally { sync.dispose() } - }) - it('refreshes a running action on an older page and removes only automatic skipped actions from the transcript', async () => { const action: Extract = { kind: 'extension-action', 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 fd20686f..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,10 +15,40 @@ 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() - expect(fixture.presentation.cancelRunPresentation).toHaveBeenCalledWith(fixture.run.id) expect(fixture.execution.stoppingRunId.value).toBe(fixture.run.id) expect(fixture.projectedRuns.value[0]?.status).toBe('running') fixture.pending.resolve({ ...fixture.run, status: 'cancelled' }) @@ -40,48 +70,11 @@ describe('useChatTurnExecution cancellation ownership', () => { expect(fixture.execution.stoppingRunId.value).toBeNull() expect(fixture.error.value).toBeTruthy() - expect(fixture.presentation.restoreRunPresentation).toHaveBeenCalledWith(fixture.run.id) expect(fixture.projectedRuns.value).toEqual([fixture.run]) expect(fixture.drafts.draft.value).toBe('pending input') expect(fixture.execution.isSending.value).toBe(false) }) - it('accepts a fresh message immediately and starts it only after cancellation finishes', 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', - }) - const receipt = { - id: 'fresh-message', - conversationId: 'conversation-a', - branchId: 'branch-a', - draftReceipt: { draftId: snapshot.draftId, sourceRevision: 1, committedRevision: 2 }, - } - f.api.chat.enqueue.mockResolvedValue(receipt) - f.api.chat.steerQueued.mockResolvedValue(true) - const cancelling = f.execution.cancelActiveRun() - await f.execution.cancelActiveRun() - expect(f.api.chat.cancel).toHaveBeenCalledTimes(1) - - expect(await f.execution.send('pending input')).toBe(true) - expect(f.api.chat.enqueue).toHaveBeenCalledTimes(1) - expect(f.api.chat.startTurn).not.toHaveBeenCalled() - expect(f.api.chat.steerQueued).not.toHaveBeenCalled() - expect(f.execution.isSending.value).toBe(false) - - f.pending.resolve({ ...f.run, status: 'cancelled' }) - await cancelling - await vi.waitFor(() => expect(f.api.chat.steerQueued).toHaveBeenCalledWith({ id: receipt.id, conversationId: receipt.conversationId, branchId: receipt.branchId })) - }) - it.each(['success', 'error'] as const)('ignores a late cancellation %s after its owner is disposed', async (outcome) => { const fixture = createFixture() const cancelling = fixture.execution.cancelActiveRun() @@ -105,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 @@ -191,8 +184,7 @@ function createFixture() { const error = shallowRef(null) const pending = deferred() const composerTarget = useComposerTarget({ drafts, conversationId: session.activeConversationId, branchId: session.activeBranchId, persist: async () => true }) - const api = { chat: { listQueue: async () => [], enqueue: vi.fn(), cancelQueued: vi.fn(), steerQueued: vi.fn(), cancel: vi.fn(() => pending.promise), executeCommand: vi.fn(), startTurn: vi.fn() } } - const presentation = { cancelRunPresentation: vi.fn(), restoreRunPresentation: vi.fn() } + const api = { chat: { listQueue: async () => [], enqueue: vi.fn(), cancelQueued: vi.fn(), steerQueued: vi.fn(), cancel: () => pending.promise, executeCommand: vi.fn(), startTurn: vi.fn() } } const selectedModel = shallowRef(null) const execution = useChatTurnExecution({ composerTarget, @@ -211,7 +203,6 @@ function createFixture() { onActionCommandRunStarted: () => {}, persistWorkspaceState: async () => true, runSync: { - ...presentation, refreshActiveConversation: async () => {}, applyRunStart: () => {}, upsertRuns: (runs) => { @@ -241,7 +232,7 @@ function createFixture() { drafts.updateComposerContent('current view input', null) error.value = 'current view status' } - return { api, drafts, error, execution, navigate, pending, presentation, projectedRuns, run, scope, selectedModel } + return { api, drafts, error, execution, navigate, pending, projectedRuns, run, scope, selectedModel } })! } diff --git a/apps/buddy/src/modules/tasks/state/runs/typing.ts b/apps/buddy/src/modules/tasks/state/runs/typing.ts index 34d0ce18..551d47b5 100644 --- a/apps/buddy/src/modules/tasks/state/runs/typing.ts +++ b/apps/buddy/src/modules/tasks/state/runs/typing.ts @@ -30,9 +30,6 @@ export interface ChatRunProjectionState { } export interface ChatRunSync extends ChatRunProjectionState { - executionRuns: Readonly>> - cancelRunPresentation: (runId: string) => void - restoreRunPresentation: (runId: string) => void isLoadingConversation: Readonly> isLoadingOlderMessages: Readonly> applyEditedTurn: (turn: LocalTurnStart, userMessageId: string) => void diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts b/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts index fc95c894..26ff33ee 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts @@ -1,10 +1,10 @@ import type { LocalChangeSetSummary } from '@buddy-shared/changes/changeApi' -import type { LocalConversationTimelineItem, LocalConversationTimelinePage } from '@buddy-shared/conversation/conversationApi' +import type { LocalConversationTimelineItem, LocalConversationTimelinePage, LocalMessage } from '@buddy-shared/conversation/conversationApi' import type { LocalApproval } from '@buddy-shared/permissions/approvalApi' import type { LocalRun, LocalRunEvent, LocalRunOutput } from '@buddy-shared/runs/runApi' -import type { ChatRunEventBucket, ChatRunEventBuckets } from '../../model/runs/typing' +import type { ChatRunEventBuckets } from '../../model/runs/typing' import type { ChatRunProjectionState } from './typing' -import { computed, shallowReactive, shallowRef } from 'vue' +import { computed, shallowRef } from 'vue' import { hasChatRunEventSequenceGap, mergeChatRunEventBuckets, @@ -21,10 +21,12 @@ import { mergeTimelineEvents, timelineItemKey, } from '../../model/runs/chatTimelineMerge' -import { projectChatRunStreamingMessages } from '../../model/transcript/chatRunStreamingMessages' export function useChatRunProjection() { const timelineItems = shallowRef>([]) + const messages = computed>(() => timelineItems.value.filter( + (item): item is Extract => item.kind === 'message', + )) const runs = shallowRef>([]) const runSignalEvents = shallowRef>([]) const runEventBuckets = shallowRef(new Map()) @@ -34,65 +36,8 @@ export function useChatRunProjection() { const timelineCursor = shallowRef(null) const hasOlderMessages = computed(() => timelineCursor.value !== null) const knownRunIds = new Set() - // Presentation stops immediately; authoritative runs still own execution until cleanup finishes. - const cancelledPresentations = shallowReactive(new Map - }>()) let hasLoadedTimelinePage = false - const presentedRuns = computed(() => runs.value.map(run => cancelledPresentations.get(run.id)?.run ?? run)) - const presentedBuckets = computed(() => { - if (!cancelledPresentations.size) - return runEventBuckets.value - const buckets = new Map(runEventBuckets.value) - for (const [runId, presentation] of cancelledPresentations) { - if (knownRunIds.has(runId)) - buckets.set(runId, presentation.bucket) - } - return buckets - }) - const presentedTimeline = computed(() => { - const items = timelineItems.value.flatMap((item) => { - if (item.kind !== 'message' || !item.runId) - return [item] - const presentation = cancelledPresentations.get(item.runId) - if (!presentation) - return [item] - const frozen = presentation.messages.get(item.id) - return frozen ? [frozen] : [] - }) - const frozenMessages = [...cancelledPresentations].flatMap(([runId, presentation]) => - knownRunIds.has(runId) ? [...presentation.messages.values()] : []) - return mergeTailTimelineItems(items, frozenMessages) - }) - - function cancelRunPresentation(runId: string) { - const run = runs.value.find(run => run.id === runId) - if (!run || (run.status !== 'queued' && run.status !== 'running') || cancelledPresentations.has(runId)) - return - const completedAt = new Date().toISOString() - const bucket = runEventBuckets.value.get(runId) - const events = bucket?.events ?? [] - // Keep the visible partial answer when terminal projection stops accepting streaming deltas. - const messages = new Map(timelineItems.value.flatMap(item => - item.kind === 'message' && item.runId === runId ? [[item.id, item] as const] : [])) - for (const candidate of projectChatRunStreamingMessages(run, events)) { - if (!messages.has(candidate.message.id)) - messages.set(candidate.message.id, { ...candidate.message, kind: 'message' }) - } - cancelledPresentations.set(runId, { - run: { ...run, status: 'cancelled', completedAt, errorCode: 'RUN_CANCELLED' }, - bucket: { events, revision: (bucket?.revision ?? 0) + 1, update: null }, - messages, - }) - } - - function restoreRunPresentation(runId: string) { - cancelledPresentations.delete(runId) - } - function mergePage(page: LocalConversationTimelinePage) { upsertRuns(page.runs) const events = mergeTimelineEvents( @@ -142,8 +87,6 @@ export function useChatRunProjection() { for (const run of incoming) { byId.set(run.id, run) knownRunIds.add(run.id) - if (run.status !== 'queued' && run.status !== 'running') - restoreRunPresentation(run.id) } runs.value = [...byId.values()].sort((left, right) => right.startedAt.localeCompare(left.startedAt)) } @@ -171,22 +114,19 @@ export function useChatRunProjection() { } const state: ChatRunProjectionState = { - approvals: computed(() => approvals.value.filter(approval => !cancelledPresentations.has(approval.runId))), + approvals, changeSets, hasOlderMessages, - messages: computed(() => presentedTimeline.value.filter((item): item is Extract => item.kind === 'message')), - runEventBuckets: presentedBuckets, + messages, + runEventBuckets, runOutputs, - runs: presentedRuns, + runs, runSignalEvents, - timelineItems: presentedTimeline, + timelineItems, } return { state, - executionRuns: runs, - cancelRunPresentation, - restoreRunPresentation, appendEvents, applySnapshot, clear, diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts b/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts index 1bcd6a7d..d3655d92 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts @@ -307,9 +307,6 @@ export function useChatRunSync(options: ChatRunSyncOptions): ChatRunSync { return { ...projection.state, - executionRuns: projection.executionRuns, - cancelRunPresentation: projection.cancelRunPresentation, - restoreRunPresentation: projection.restoreRunPresentation, applyEditedTurn: (turn, messageId) => applyReplacementTurn(turn, messageId, false), applyRegeneratedTurn: turn => applyReplacementTurn(turn, turn.run.triggeringMessageId, true), applyRunStart, diff --git a/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts b/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts index a0950f1e..36726a16 100644 --- a/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts +++ b/apps/buddy/src/modules/tasks/state/runs/useChatTurnExecution.ts @@ -2,7 +2,6 @@ import type { LocalChatApi } from '@buddy-electron/shared/localChatApi' import type { ParsedBuddyChatCommand } from '@buddy-shared/conversation/buddyChatCommands' import type { BuddyUserContentV1 } from '@buddy-shared/conversation/buddyUserContent' import type { LocalPromptContextItem } from '@buddy-shared/conversation/chatApi' -import type { LocalChatQueueTarget } from '@buddy-shared/conversation/chatQueueApi' import type { BuddyApprovalPolicy } from '@buddy-shared/permissions/approvalPolicy' import type { BuddyExecutionProfile } from '@buddy-shared/permissions/executionProfile' import type { LocalRun } from '@buddy-shared/runs/runApi' @@ -58,7 +57,7 @@ export interface UseChatTurnExecutionOptions { onActionCommandRunStarted: (runId: string) => void onDraftCommitted?: (draftId: string, conversationId: string) => void persistWorkspaceState: () => Promise - runSync: Pick + runSync: Pick runtimeSupervisor: Pick setErrorMessage: (message: string | null) => void unavailableCommandMessage: () => string @@ -69,11 +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 cancellationCompletions = new Map }>() + 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(() => { @@ -81,8 +81,7 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { for (const stop of pendingCancellationWatches) stop() pendingCancellationWatches.clear() - pendingCancellationIds.clear() - cancellationCompletions.clear() + pendingCancellations.clear() }, true) const canSend = computed(() => options.runtimeSupervisor.runtimeState.value.status === 'ready' @@ -113,15 +112,12 @@ 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) const sourceScopeKey = options.draftScopeKey.value - const cancellation = [...cancellationCompletions.values()].find(({ run }) => - run.conversationId === options.session.activeConversationId.value - && run.branchId === options.session.activeBranchId.value) const navigationVersion = options.session.generation() const isSourceViewCurrent = () => options.session.isCurrent(navigationVersion) && options.draftScopeKey.value === sourceScopeKey @@ -139,12 +135,10 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { const expectedRevision = confirmedDraft.revision! const operationKey = `turn:${confirmedDraft.draftId}:${expectedRevision}` const requestId = requestIds.resolve(operationKey) - if (cancellation || 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) - if (cancellation) - void dispatchAfterCancellation(result, cancellation.completion) await queue.refreshQueue() return true } @@ -272,12 +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) - const cancellation = Promise.withResolvers() - cancellationCompletions.set(run.id, { run, completion: cancellation.promise }) - options.runSync.cancelRunPresentation(run.id) + pendingCancellations.set(run.id, run) const navigationVersion = options.session.generation() let sourceViewChanged = false const stopWatchingView = watch( @@ -296,35 +287,15 @@ export function useChatTurnExecution(options: UseChatTurnExecutionOptions) { const cancelled = await options.api.chat.cancel(run.id) if (isSourceViewCurrent()) options.runSync.upsertRuns([cancelled]) - cancellation.resolve(cancelled.status === 'cancelled') } catch (error) { - options.runSync.restoreRunPresentation(run.id) - cancellation.resolve(false) if (isSourceViewCurrent()) setNormalizedError(error) } finally { stopWatchingView() pendingCancellationWatches.delete(stopWatchingView) - pendingCancellationIds.delete(run.id) - cancellationCompletions.delete(run.id) - } - } - - async function dispatchAfterCancellation(target: LocalChatQueueTarget, completion: Promise) { - if (!await completion || isDisposed) - return - try { - // Only the message explicitly sent after Stop resumes; older paused inputs stay paused. - await options.api.chat.steerQueued({ id: target.id, conversationId: target.conversationId, branchId: target.branchId }) - await queue.refreshQueue() - } - catch (error) { - if (!isDisposed && options.session.activeConversationId.value === target.conversationId - && options.session.activeBranchId.value === target.branchId) { - setNormalizedError(error) - } + pendingCancellations.delete(run.id) } } diff --git a/apps/buddy/src/modules/tasks/state/useTaskCapability.ts b/apps/buddy/src/modules/tasks/state/useTaskCapability.ts index 676817d7..3f7d7c2c 100644 --- a/apps/buddy/src/modules/tasks/state/useTaskCapability.ts +++ b/apps/buddy/src/modules/tasks/state/useTaskCapability.ts @@ -100,7 +100,7 @@ export function useTaskCapability(options: UseTaskCapabilityOptions): TaskCapabi timelineItems, } = runSync - const activeRun = computed(() => runSync.executionRuns.value.find( + const activeRun = computed(() => runs.value.find( run => run.status === 'queued' || run.status === 'running', ) ?? null) const hasAvailableProvider = computed(() => modelProviders.providers.value.some( @@ -493,7 +493,7 @@ export function useTaskCapability(options: UseTaskCapabilityOptions): TaskCapabi }, }, execution: { - activeRun: computed(() => runs.value.find(run => run.status === 'queued' || run.status === 'running') ?? null), + activeRun: readonly(activeRun), approvalViews: readonly(approvalViews), canMutateBranch: readonly(canMutateBranch), canSend: readonly(canSend), 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') }} - -