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
84 changes: 56 additions & 28 deletions apps/sim/app/api/copilot/chat/stop/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -148,42 +148,70 @@ describe('copilot chat stop route', () => {
expect(mockAppendCopilotChatMessages).not.toHaveBeenCalled()
})

it('appends a stopped assistant message even with no content', async () => {
it('persists a response larger than the HTTP limit from an identifiers-only Stop', async () => {
mockReads({
chat: { workspaceId: 'ws-1', conversationId: 'stream-1', model: null },
last: { messageId: 'stream-1', role: 'user' },
})

const response = await stopRequest(
createRequest({ chatId: 'chat-1', streamId: 'stream-1', content: '' })
)

expect(mockReadEvents).toHaveBeenCalledWith('stream-1', '0')
expect(response.status).toBe(200)
expect(await response.json()).toEqual({ success: true })

const setArg = dbChainMockFns.set.mock.calls[0]?.[0] as Record<string, unknown>
expect(setArg.conversationId).toBeNull()
expect(Object.hasOwn(setArg, 'messages')).toBe(false)

expect(mockAppendCopilotChatMessages).toHaveBeenCalledTimes(1)
const [, appended] = mockAppendCopilotChatMessages.mock.calls[0]
expect(appended[0]).toMatchObject({
role: 'assistant',
content: '',
contentBlocks: [{ type: 'complete', status: 'cancelled' }],
})

expect(mockPublishStatusChanged).toHaveBeenCalledWith(
expect.objectContaining({ workspaceId: 'ws-1' }),
const content = 'x'.repeat(11 * 1024 * 1024)
const envelope = { v: 1, ts: '2026-09-24T19:00:00Z', stream: { streamId: 'stream-1' } }
mockReadEvents.mockResolvedValue([
{ ...envelope, seq: 1, type: 'text', payload: { channel: 'assistant', text: content } },
{
chatId: 'chat-1',
type: 'completed',
streamId: 'stream-1',
}
...envelope,
seq: 2,
type: 'tool',
payload: {
phase: 'call',
toolCallId: 'call-1',
toolName: 'run_code',
arguments: { code: 'preserve me' },
status: 'executing',
},
},
])
const request = createRequest({ chatId: 'chat-1', streamId: 'stream-1' })
expect((await request.clone().text()).length).toBeLessThan(1024)
const response = await stopRequest(request)
expect(response.status).toBe(200)
expect(mockReadEvents).toHaveBeenCalledOnce()
expect(mockAppendCopilotChatMessages).toHaveBeenCalledOnce()
const saved = mockAppendCopilotChatMessages.mock.calls[0][1][0]
expect(saved.content).toBe(content)
expect(saved.contentBlocks).toEqual(
expect.arrayContaining([
expect.objectContaining({
type: 'tool',
toolCall: expect.objectContaining({ id: 'call-1', params: { code: 'preserve me' } }),
}),
{ type: 'complete', status: 'cancelled' },
])
)
expect(dbChainMockFns.set.mock.calls[0][0].conversationId).toBeNull()
expect(mockPublishStatusChanged).toHaveBeenCalledOnce()
})

it.each([{ seqs: [] }, { seqs: [2, 3] }, { seqs: [1, 3] }])(
'leaves incomplete replay $seqs to the run owner without erasing its response',
async ({ seqs }) => {
mockReadEvents.mockResolvedValue(
seqs.map((seq) => ({
v: 1,
seq,
ts: '2026-09-24T19:00:00Z',
stream: { streamId: 'stream-1' },
type: 'text',
payload: { channel: 'assistant', text: 'tail only' },
}))
)
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()
expect(mockPublishStatusChanged).not.toHaveBeenCalled()
}
)

it('appends a stopped assistant message if the stream marker was already cleared', async () => {
mockReads({
chat: { workspaceId: 'ws-1', conversationId: null, model: null },
Expand Down
32 changes: 20 additions & 12 deletions apps/sim/app/api/copilot/chat/stop/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,10 @@ import {
type PersistedMessage,
withStoppedContentBlock,
} from '@/lib/mothership/chat/persisted-message'
import { finalizeAssistantTurn } from '@/lib/mothership/chat/terminal-state'
import {
finalizeAssistantTurn,
readStoppedAssistantMessage,
} from '@/lib/mothership/chat/terminal-state'
import { publishChatStatusChanged } from '@/lib/mothership/chat-status'
import {
CopilotChatFinalizeOutcome,
Expand Down Expand Up @@ -61,23 +64,28 @@ export const POST = withRouteHandler((req: NextRequest) =>
: hasContent
? [{ type: 'text', channel: 'assistant', content }]
: []
const assistantMessage: PersistedMessage = withStoppedContentBlock(
normalizeMessage({
id: generateId(),
role: 'assistant',
content,
timestamp: new Date().toISOString(),
contentBlocks: assistantBlocks,
...(requestId ? { requestId } : {}),
})
)
const assistantMessage: PersistedMessage | null =
hasContent || hasBlocks
? withStoppedContentBlock(
normalizeMessage({
id: generateId(),
role: 'assistant',
content,
timestamp: new Date().toISOString(),
contentBlocks: assistantBlocks,
...(requestId ? { requestId } : {}),
})
)
: await readStoppedAssistantMessage(streamId)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Replay precedes stream matching Access is checked for chatId, but this call loads the full event history for the supplied streamId before checking whether it belongs to that chat. Repeated Stop requests with a mismatched stream ID can therefore cause large, unnecessary Redis reads. Match the stream to the chat before loading its replay.

/** 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({
chatId,
userId: session.user.id,
userMessageId: streamId,
assistantMessage,
streamMarkerPolicy: 'active-or-cleared',
preferServerReplay: true,
preferServerReplay: hasContent || hasBlocks,
})
span.setAttribute(TraceAttr.CopilotStopAppendedAssistant, result.appendedAssistant)
const stopOutcome = !result.found
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,6 @@ import {
seedDeploymentShape,
} from '@/lib/core/config/deployment-shape'
import { MothershipHandoffStorage } from '@/lib/core/utils/browser-storage'
import { normalizeMessage } from '@/lib/mothership/chat/persisted-message'
import type { MothershipStreamV1EventEnvelope } from '@/lib/mothership/generated/mothership-stream-v1'
import { createSearchResource } from '@/lib/mothership/resources/search'
import { getChatResourceSelectionId } from '@/lib/mothership/resources/types'
Expand Down Expand Up @@ -117,6 +116,7 @@ interface NetworkState {
abortSettlements: boolean[]
abortBodies: CopilotChatAbortBody[]
stopBodies: CopilotChatStopBody[]
toolInputPadding?: string
abortTraceparents: Array<string | null>
}

Expand Down Expand Up @@ -221,11 +221,14 @@ async function fetchStub(input: RequestInfo | URL, init?: RequestInit): Promise<
stream: { streamId },
payload: {
phase: 'call',
executor: 'client',
mode: 'async',
toolName: 'run_workflow',
executor: state.toolInputPadding ? 'go' : 'client',
mode: state.toolInputPadding ? 'sync' : 'async',
toolName: state.toolInputPadding ? 'run_code' : 'run_workflow',
toolCallId: 'this-chat-tool',
arguments: { workflowId: 'this-chat-workflow' },
arguments: {
workflowId: 'this-chat-workflow',
...(state.toolInputPadding ? { padding: state.toolInputPadding } : {}),
},
},
}
return new Response(
Expand Down Expand Up @@ -751,6 +754,7 @@ describe('useChat remount send recovery', () => {
state.pendingAdmissions.clear()
state.abortSettlements = []
state.abortBodies = []
state.toolInputPadding = undefined
state.stopBodies = []
state.abortTraceparents = []
mockRequestJson.mockResolvedValue({ chats: [] })
Expand Down Expand Up @@ -1450,6 +1454,41 @@ describe('useChat remount send recovery', () => {
}
)

it.each([false, true])(
'does not interrupt or send a queued edit before submission (explicit ID: %s)',
async (explicitId) => {
state.postBehavior = 'task'
const { getResult } = renderUseChatInChat('chat-a')
await act(async () => {
void getResult().sendMessage('Original request')
})
await waitFor(() => state.postBodies.length === 1 && getResult().isSending)
const beforeRender = getResult()
let queuedId = ''
await act(async () => {
void beforeRender.sendMessage('Unfinished correction')
queuedId = allQueuedMessages()[0].id
beforeRender.editQueuedMessage(queuedId)
void beforeRender.sendNow(explicitId ? queuedId : undefined)
})
expect(state.abortBodies).toHaveLength(0)
expect(state.postBodies).toHaveLength(1)
expect(allQueuedMessages()).toEqual([
expect.objectContaining({ id: queuedId, content: 'Unfinished correction' }),
])
expect(useMothershipQueueStore.getState().editing['chat-a']).toBe(queuedId)
state.postBehavior = 'hang'
await act(async () => {
void getResult().sendMessage('Finished correction')
void getResult().sendNow()
})
await waitFor(() => state.postBodies.length === 2)
expect(state.postBodies[1].message).toBe('Finished correction')
expect(state.abortBodies).toHaveLength(1)
expect(allQueuedMessages()).toHaveLength(0)
}
)

it('sends the live queue head once without waiting for a render, after Stop settles', async () => {
state.postBehavior = 'task'
const { getResult } = renderUseChatInChat('chat-a')
Expand Down Expand Up @@ -1713,7 +1752,7 @@ describe('useChat remount send recovery', () => {
}
})

it('preserves a visible workflow watch when Stop persists the partial response', async () => {
it('preserves the visible workflow watch while Stop sends only identifiers', async () => {
state.postBehavior = 'task'
const { getResult } = renderUseChatInChat('chat-a')
await act(async () => {
Expand All @@ -1730,20 +1769,43 @@ describe('useChat remount send recovery', () => {
await getResult().stopGeneration()
})
expect(state.stopBodies).toHaveLength(1)
const saved = state.stopBodies[0]
const restored = normalizeMessage({
id: 'saved-assistant',
role: 'assistant',
content: saved.content,
contentBlocks: saved.contentBlocks,
})
expect(restored.contentBlocks?.find((block) => block.type === 'task')?.task).toEqual({
taskId: 'watch-1',
kind: 'workflow_run',
status: 'pending',
target: { workflowId: 'workflow-1', executionId: 'watched-execution' },
note: 'Check the completed invoice run',
expect(state.stopBodies[0]).toEqual({
chatId: 'chat-a',
streamId: state.postBodies[0].userMessageId,
})
const task = getResult()
.messages.flatMap((message) => message.contentBlocks ?? [])
.find((block) => block.type === 'task')?.task
expect(task?.taskId).toBe('watch-1')
})

it('sends a queued correction after stopping with more than 10 MiB of tool input', async () => {
state.postBehavior = 'tool'
state.toolInputPadding = 'x'.repeat(11 * 1024 * 1024)
const { getResult } = renderUseChatInChat('chat-a')
await act(async () => {
void getResult().sendMessage('Start working')
})
await waitFor(() =>
expect(
getResult().messages.some((message) =>
message.contentBlocks?.some((block) => block.toolCall?.id === 'this-chat-tool')
)
).toBe(true)
)
state.postBehavior = 'hang'
await act(async () => {
await getResult().sendMessage('Use the correction')
void getResult().sendNow()
})
await waitFor(() => expect(state.postBodies).toHaveLength(2))
expect(state.postBodies[1].message).toBe('Use the correction')
expect(state.stopBodies).toHaveLength(1)
expect(state.stopBodies[0]).not.toHaveProperty('content')
expect(state.stopBodies[0]).not.toHaveProperty('contentBlocks')
expect(new TextEncoder().encode(JSON.stringify(state.stopBodies[0])).length).toBeLessThan(1024)
expect(allQueuedMessages()).toHaveLength(0)
expect(getResult().error).toBeNull()
})

it.each([false, true])(
Expand Down
Loading
Loading