Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,59 @@
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' } },

Check failure on line 34 in apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts

View workflow job for this annotation

GitHub Actions / Workspace checks

Should not have line breaks between items, in node ArrayExpression

Check failure on line 34 in apps/buddy/src/modules/tasks/state/runs/__tests__/useChatRunSync.spec.ts

View workflow job for this annotation

GitHub Actions / Workspace checks

Should not have line breaks between items, in node ArrayExpression
])
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<LocalConversationTimelineItem, { kind: 'extension-action' }> = {
kind: 'extension-action',
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
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' })
Expand All @@ -39,11 +40,46 @@
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',

Check failure on line 63 in apps/buddy/src/modules/tasks/state/runs/__tests__/useChatTurnExecution.spec.ts

View workflow job for this annotation

GitHub Actions / Workspace checks

Should have line breaks between items, in node ObjectExpression

Check failure on line 63 in apps/buddy/src/modules/tasks/state/runs/__tests__/useChatTurnExecution.spec.ts

View workflow job for this annotation

GitHub Actions / Workspace checks

Should have line breaks between items, in node ObjectExpression
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()
Expand Down Expand Up @@ -153,7 +189,8 @@
const error = shallowRef<string | null>(null)
const pending = deferred<LocalRun>()
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<LocalRuntimeModelOption | null>(null)
const execution = useChatTurnExecution({
composerTarget,
Expand All @@ -172,6 +209,7 @@
onActionCommandRunStarted: () => {},
persistWorkspaceState: async () => true,
runSync: {
...presentation,
refreshActiveConversation: async () => {},
applyRunStart: () => {},
upsertRuns: (runs) => {
Expand Down Expand Up @@ -201,7 +239,7 @@
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 }
})!
}

Expand Down
3 changes: 3 additions & 0 deletions apps/buddy/src/modules/tasks/state/runs/typing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,9 @@ export interface ChatRunProjectionState {
}

export interface ChatRunSync extends ChatRunProjectionState {
executionRuns: Readonly<Ref<ReadonlyArray<LocalRun>>>
cancelRunPresentation: (runId: string) => void
restoreRunPresentation: (runId: string) => void
isLoadingConversation: Readonly<Ref<boolean>>
isLoadingOlderMessages: Readonly<Ref<boolean>>
applyEditedTurn: (turn: LocalTurnStart, userMessageId: string) => void
Expand Down
82 changes: 71 additions & 11 deletions apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts
Original file line number Diff line number Diff line change
@@ -1,11 +1,12 @@
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 { projectChatRunStreamingMessages } from '../../model/transcript/chatRunStreamingMessages'
import {

Check failure on line 9 in apps/buddy/src/modules/tasks/state/runs/useChatRunProjection.ts

View workflow job for this annotation

GitHub Actions / Workspace checks

Expected "../../model/runs/chatRunEventBuckets" to come before "../../model/transcript/chatRunStreamingMessages"
hasChatRunEventSequenceGap,
mergeChatRunEventBuckets,
mergeChatRunEvents,
Expand All @@ -24,9 +25,6 @@

export function useChatRunProjection() {
const timelineItems = shallowRef<ReadonlyArray<LocalConversationTimelineItem>>([])
const messages = computed<ReadonlyArray<LocalMessage>>(() => timelineItems.value.filter(
(item): item is Extract<LocalConversationTimelineItem, { kind: 'message' }> => item.kind === 'message',
))
const runs = shallowRef<ReadonlyArray<LocalRun>>([])
const runSignalEvents = shallowRef<ReadonlyArray<LocalRunEvent>>([])
const runEventBuckets = shallowRef<ChatRunEventBuckets>(new Map())
Expand All @@ -36,8 +34,65 @@
const timelineCursor = shallowRef<string | null>(null)
const hasOlderMessages = computed(() => timelineCursor.value !== null)
const knownRunIds = new Set<string>()
// Presentation stops immediately; authoritative runs still own execution until cleanup finishes.
const cancelledPresentations = shallowReactive(new Map<string, {
run: LocalRun
bucket: ChatRunEventBucket
messages: ReadonlyMap<string, LocalConversationTimelineItem>
}>())
let hasLoadedTimelinePage = false

const presentedRuns = computed(() => runs.value.map(run => cancelledPresentations.get(run.id)?.run ?? run))
const presentedBuckets = computed<ChatRunEventBuckets>(() => {
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<string, LocalConversationTimelineItem>(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(
Expand Down Expand Up @@ -87,6 +142,8 @@
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))
}
Expand Down Expand Up @@ -114,19 +171,22 @@
}

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<LocalConversationTimelineItem, { kind: 'message' }> => item.kind === 'message')),
runEventBuckets: presentedBuckets,
runOutputs,
runs,
runs: presentedRuns,
runSignalEvents,
timelineItems,
timelineItems: presentedTimeline,
}

return {
state,
executionRuns: runs,
cancelRunPresentation,
restoreRunPresentation,
appendEvents,
applySnapshot,
clear,
Expand Down
3 changes: 3 additions & 0 deletions apps/buddy/src/modules/tasks/state/runs/useChatRunSync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading