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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 48 additions & 6 deletions apps/buddy/service/src/chat/ChatQueueService.ts
Original file line number Diff line number Diff line change
@@ -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'
Expand All @@ -21,7 +22,7 @@ import { resolveTurnExecutionProfile } from '../storage/turnRequestRepository'

export interface ChatQueueServiceOptions {
queue: ChatQueueRepository
turns: Pick<ChatTurnService, 'prepareStart' | 'validatePreparedInput'>
turns: Pick<ChatTurnService, 'cancel' | 'prepareStart' | 'validatePreparedInput'>
requests: Pick<TurnRequestRepository, 'prepare'>
launcher: Pick<BuddyTurnLauncher, 'launch'>
runner: Pick<BuddyAgentRunner, 'steer' | 'followUp' | 'hasActiveExecution' | 'hasDegradedCleanup' | 'isStopping'>
Expand All @@ -41,6 +42,7 @@ export class ChatQueueService {
readonly #options: ChatQueueServiceOptions
readonly #operations = new Map<string, Promise<void>>()
readonly #enqueues = new Set<Promise<unknown>>()
readonly #cancellations = new Map<string, { conversationId: string, completion: Promise<LocalRun> }>()
readonly #changes: Emitter<ChatQueueChange>
readonly #stopping = new AbortController()
readonly onDidChange: Emitter<ChatQueueChange>['event']
Expand All @@ -61,7 +63,7 @@ export class ChatQueueService {
}

async drain(): Promise<void> {
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[] {
Expand All @@ -79,6 +81,37 @@ export class ChatQueueService {
return this.#options.queue.list(scope)
}

cancelRun(runId: string): Promise<LocalRun> {
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<LocalRun>()
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<LocalRun> {
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<LocalChatQueueReceipt> {
const pending = Promise.withResolvers<LocalChatQueueReceipt>()
this.#enqueues.add(pending.promise)
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand All @@ -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')
})
Expand All @@ -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')
Expand All @@ -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
Expand Down Expand Up @@ -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 } : {}) }))
}
Expand Down
90 changes: 90 additions & 0 deletions apps/buddy/service/src/chat/__tests__/ChatQueueService.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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<void>()
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<void>()
const cleanup = Promise.withResolvers<void>()
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)
Expand Down Expand Up @@ -392,6 +481,7 @@ function fixture(queuedOnStartup = false) {
const committed = new Emitter<BuddyRunEvent>(() => {})
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')
},
Expand Down
4 changes: 2 additions & 2 deletions apps/buddy/service/src/chat/registerChatRpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ export interface RegisterChatRpcOptions {
runtime: BuddyRuntime
turns: Pick<
ChatTurnService,
'cancel' | 'editUserMessage' | 'regenerateAssistant'
'editUserMessage' | 'regenerateAssistant'
>
}

Expand All @@ -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())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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', {
Expand Down
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 @@ -15,6 +15,37 @@ const scopes: ReturnType<typeof effectScope>[] = []
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()
Expand Down Expand Up @@ -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

Expand Down
Loading