diff --git a/apps/sim/app/api/copilot/chat/stop/route.test.ts b/apps/sim/app/api/copilot/chat/stop/route.test.ts index 87488ef95f1..9b6a2a43c33 100644 --- a/apps/sim/app/api/copilot/chat/stop/route.test.ts +++ b/apps/sim/app/api/copilot/chat/stop/route.test.ts @@ -22,6 +22,10 @@ vi.mock('@/lib/mothership/chat/lifecycle', () => ({ })) vi.mock('@/lib/mothership/request/session/buffer', () => ({ readEvents: mockReadEvents })) +vi.mock('@/lib/mothership/async-runs/repository', () => ({ + getLatestRunForStream: vi.fn().mockResolvedValue({ chatId: 'chat-1', status: 'cancelled' }), +})) + vi.mock('@/lib/mothership/chat/messages-store', () => ({ appendCopilotChatMessages: mockAppendCopilotChatMessages, })) @@ -30,6 +34,7 @@ vi.mock('@/lib/mothership/chat-status', () => ({ publishChatStatusChanged: mockPublishStatusChanged, })) +import { getLatestRunForStream } from '@/lib/mothership/async-runs/repository' import { POST } from '@/app/api/copilot/chat/stop/route' const stopRequest = (request: NextRequest) => POST(request, undefined) @@ -169,6 +174,7 @@ describe('copilot chat stop route', () => { status: 'executing', }, }, + { ...envelope, seq: 3, type: 'complete', payload: { status: 'cancelled' } }, ]) const request = createRequest({ chatId: 'chat-1', streamId: 'stream-1' }) expect((await request.clone().text()).length).toBeLessThan(1024) @@ -182,7 +188,11 @@ describe('copilot chat stop route', () => { expect.arrayContaining([ expect.objectContaining({ type: 'tool', - toolCall: expect.objectContaining({ id: 'call-1', params: { code: 'preserve me' } }), + toolCall: expect.objectContaining({ + id: 'call-1', + state: 'cancelled', + params: { code: 'preserve me' }, + }), }), { type: 'complete', status: 'cancelled' }, ]) @@ -191,6 +201,42 @@ describe('copilot chat stop route', () => { expect(mockPublishStatusChanged).toHaveBeenCalledOnce() }) + it.each(['active', 'other-chat', 'missing'])( + 'does not read replay for an unfinalized or mismatched run: %s', + async (state) => { + vi.mocked(getLatestRunForStream).mockResolvedValueOnce( + state === 'missing' + ? null + : ({ + chatId: state === 'other-chat' ? 'other-chat' : 'chat-1', + status: state === 'active' ? 'active' : 'cancelled', + } as Awaited>) + ) + const response = await stopRequest(createRequest({ chatId: 'chat-1', streamId: 'stream-1' })) + expect(response.status).toBe(200) + expect(getLatestRunForStream).toHaveBeenCalledWith('stream-1', 'user-1') + expect(mockReadEvents).not.toHaveBeenCalled() + expect(mockAppendCopilotChatMessages).not.toHaveBeenCalled() + } + ) + + it('does not finalize a contiguous prefix before the final event is flushed', async () => { + mockReadEvents.mockResolvedValue([ + { + v: 1, + seq: 1, + ts: '2026-09-24T19:00:00Z', + stream: { streamId: 'stream-1' }, + type: 'text', + payload: { channel: 'assistant', text: 'incomplete prefix' }, + }, + ]) + const response = await stopRequest(createRequest({ chatId: 'chat-1', streamId: 'stream-1' })) + expect(response.status).toBe(200) + expect(mockAppendCopilotChatMessages).not.toHaveBeenCalled() + expect(dbChainMockFns.set).not.toHaveBeenCalled() + }) + it.each([{ seqs: [] }, { seqs: [2, 3] }, { seqs: [1, 3] }])( 'leaves incomplete replay $seqs to the run owner without erasing its response', async ({ seqs }) => { diff --git a/apps/sim/app/api/copilot/chat/stop/route.ts b/apps/sim/app/api/copilot/chat/stop/route.ts index 2473b8ddf91..33ccee4e0fd 100644 --- a/apps/sim/app/api/copilot/chat/stop/route.ts +++ b/apps/sim/app/api/copilot/chat/stop/route.ts @@ -76,7 +76,7 @@ export const POST = withRouteHandler((req: NextRequest) => ...(requestId ? { requestId } : {}), }) ) - : await readStoppedAssistantMessage(streamId) + : await readStoppedAssistantMessage(streamId, chatId, session.user.id) /** The run owner retains the full response if replay was trimmed or has not started. */ if (!assistantMessage) return NextResponse.json({ success: true }) const result = await finalizeAssistantTurn({ diff --git a/apps/sim/lib/mothership/chat/persisted-message.test.ts b/apps/sim/lib/mothership/chat/persisted-message.test.ts index 816f3df07b5..0379fcbe70c 100644 --- a/apps/sim/lib/mothership/chat/persisted-message.test.ts +++ b/apps/sim/lib/mothership/chat/persisted-message.test.ts @@ -13,9 +13,44 @@ import { normalizeMessage, type PersistedMessage, stripToolResultOutput, + withStoppedContentBlock, } from './persisted-message' describe('persisted-message', () => { + it.each([false, true])( + 'cancels unfinished tools even when the stopped marker already exists: %s', + (alreadyStopped) => { + const message: PersistedMessage = { + id: 'assistant', + role: 'assistant', + content: '', + timestamp: '2026-09-24T19:00:00Z', + contentBlocks: [ + { + type: 'tool', + toolCall: { + id: 'unfinished', + name: 'run_code', + state: 'executing', + params: { code: 'keep me' }, + }, + }, + { type: 'tool', toolCall: { id: 'finished', name: 'read', state: 'success' } }, + ...(alreadyStopped ? [{ type: 'complete' as const, status: 'cancelled' as const }] : []), + ], + } + const saved = withStoppedContentBlock(message) + expect(saved.contentBlocks?.[0].toolCall).toMatchObject({ + state: 'cancelled', + params: { code: 'keep me' }, + display: { title: 'Stopped by user' }, + }) + expect(saved.contentBlocks?.[1].toolCall?.state).toBe('success') + expect(saved.contentBlocks?.filter((block) => block.type === 'complete')).toHaveLength(1) + expect(message.contentBlocks?.[0].toolCall?.state).toBe('executing') + } + ) + it.each(['success', 'cancelled'] as const)( 'preserves model-authored activity metadata through persisted %s history and stop validation', (status) => { diff --git a/apps/sim/lib/mothership/chat/persisted-message.ts b/apps/sim/lib/mothership/chat/persisted-message.ts index 6e60be912a7..4cf6169147c 100644 --- a/apps/sim/lib/mothership/chat/persisted-message.ts +++ b/apps/sim/lib/mothership/chat/persisted-message.ts @@ -358,7 +358,19 @@ export function buildPersistedAssistantMessage( } export function withStoppedContentBlock(message: PersistedMessage): PersistedMessage { - const contentBlocks = message.contentBlocks ?? [] + const contentBlocks = (message.contentBlocks ?? []).map( + (block): PersistedContentBlock => + block.toolCall?.state === 'executing' + ? { + ...block, + toolCall: { + ...block.toolCall, + state: 'cancelled', + display: { title: 'Stopped by user' }, + }, + } + : block + ) const hasAssistantText = contentBlocks.some( (block) => block.type === MothershipStreamV1EventType.text && @@ -372,7 +384,7 @@ export function withStoppedContentBlock(message: PersistedMessage): PersistedMes block.status === MothershipStreamV1CompletionStatus.cancelled ) ) { - return message + return { ...message, contentBlocks } } return normalizeMessage({ diff --git a/apps/sim/lib/mothership/chat/terminal-state.test.ts b/apps/sim/lib/mothership/chat/terminal-state.test.ts index e8063896e9e..00d253f66a1 100644 --- a/apps/sim/lib/mothership/chat/terminal-state.test.ts +++ b/apps/sim/lib/mothership/chat/terminal-state.test.ts @@ -13,6 +13,10 @@ const { mockAppendCopilotChatMessages, mockReadEvents } = vi.hoisted(() => ({ })) vi.mock('@/lib/mothership/request/session/buffer', () => ({ readEvents: mockReadEvents })) +vi.mock('@/lib/mothership/async-runs/repository', () => ({ + getLatestRunForStream: vi.fn().mockResolvedValue({ chatId: 'chat-1', status: 'cancelled' }), +})) + vi.mock('@/lib/mothership/chat/messages-store', () => ({ appendCopilotChatMessages: mockAppendCopilotChatMessages, })) @@ -72,6 +76,7 @@ describe('finalizeAssistantTurn', () => { status: 'executing', }, }, + { ...envelope, seq: 3, type: 'complete', payload: { status: 'cancelled' } }, ]) await finalizeAssistantTurn({ chatId: 'chat-1', diff --git a/apps/sim/lib/mothership/chat/terminal-state.ts b/apps/sim/lib/mothership/chat/terminal-state.ts index f08f4157c51..202ecd9cdf6 100644 --- a/apps/sim/lib/mothership/chat/terminal-state.ts +++ b/apps/sim/lib/mothership/chat/terminal-state.ts @@ -1,6 +1,7 @@ import { db } from '@sim/db' import { copilotChats, copilotMessages, copilotRuns } from '@sim/db/schema' import { and, desc, eq, isNull, sql } from 'drizzle-orm' +import { getLatestRunForStream } from '@/lib/mothership/async-runs/repository' import { buildLiveAssistantMessage } from '@/lib/mothership/chat/effective-transcript' import { appendCopilotChatMessages } from '@/lib/mothership/chat/messages-store' import { @@ -16,6 +17,7 @@ import { TraceAttr } from '@/lib/mothership/generated/trace-attributes-v1' import { TraceSpan } from '@/lib/mothership/generated/trace-spans-v1' import { withCopilotSpan } from '@/lib/mothership/request/otel' import { readEvents } from '@/lib/mothership/request/session/buffer' +import { isTerminalStreamStatus } from '@/lib/mothership/request/session/contract' import { StreamControllerSupersededError } from '@/lib/mothership/request/session/controller-lease' import { toStreamBatchEvent } from '@/lib/mothership/request/session/types' @@ -39,13 +41,21 @@ export interface FinalizeAssistantTurnResult { outcome: (typeof CopilotChatFinalizeOutcome)[keyof typeof CopilotChatFinalizeOutcome] } -/** Rebuild a stopped response only when the server still has its complete event prefix. */ +/** Only the matching terminal run and a gap-free replay through its final event can be persisted. */ export async function readStoppedAssistantMessage( - streamId: string + streamId: string, + chatId: string, + userId?: string ): Promise { + const run = await getLatestRunForStream(streamId, userId) + if (run?.chatId !== chatId || !isTerminalStreamStatus(run.status)) return null const events = await readEvents(streamId, '0') /** StreamWriter starts at 1; Redis may trim oldest events or skip corrupt entries. */ - if (events.length === 0 || !events.every((event, index) => event.seq === index + 1)) return null + if ( + events.at(-1)?.type !== 'complete' || + !events.every((event, index) => event.seq === index + 1) + ) + return null const replay = buildLiveAssistantMessage({ streamId, events: events.map(toStreamBatchEvent), @@ -176,7 +186,7 @@ export async function finalizeAssistantTurn({ if (assistantMessage && canAppendAssistant) { let response = assistantMessage if (preferServerReplay) { - const replay = await readStoppedAssistantMessage(userMessageId) + const replay = await readStoppedAssistantMessage(userMessageId, chatId, userId) /** A stopped client's snapshot may be empty; preserve canonical output before the first finalizer commits. */ const replayHasContent = !!replay?.content.trim() ||