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
48 changes: 47 additions & 1 deletion apps/sim/app/api/copilot/chat/stop/route.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}))
Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand All @@ -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' },
])
Expand All @@ -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<ReturnType<typeof getLatestRunForStream>>)
)
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 }) => {
Expand Down
2 changes: 1 addition & 1 deletion apps/sim/app/api/copilot/chat/stop/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down
35 changes: 35 additions & 0 deletions apps/sim/lib/mothership/chat/persisted-message.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,44 @@ import {
normalizeMessage,
type PersistedMessage,
stripToolResultOutput,
withStoppedContentBlock,

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 New import violates path rule The added withStoppedContentBlock import uses the relative ./persisted-message path. The repository requires absolute imports in apps/sim; change this import to an @/ path before merging.

Context Used: Import patterns for the Sim application (source)

Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!

} 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) => {
Expand Down
16 changes: 14 additions & 2 deletions apps/sim/lib/mothership/chat/persisted-message.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'

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.

P1 Waiting tools remain active

If Stop arrives while a tool is awaiting approval, this check leaves it unchanged because it only cancels executing tools. Replay preserves the awaiting_approval state, so the stopped turn can still show an actionable permission card after reload. Cancel the other unfinished tool states as well.

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.

P1 Stopped tools remain unfinished When a stopped turn contains a pending or awaiting_approval tool, this check leaves its state unchanged while the message gets a cancelled completion block. Both states can survive result persistence or replay, so reloaded history can show a stopped tool as still pending or awaiting approval. Cancel those nonterminal states as well as executing.

? {
...block,
toolCall: {
...block.toolCall,
state: 'cancelled',
display: { title: 'Stopped by user' },
},
}
: block
)
const hasAssistantText = contentBlocks.some(
(block) =>
block.type === MothershipStreamV1EventType.text &&
Expand All @@ -372,7 +384,7 @@ export function withStoppedContentBlock(message: PersistedMessage): PersistedMes
block.status === MothershipStreamV1CompletionStatus.cancelled
)
) {
return message
return { ...message, contentBlocks }
}

return normalizeMessage({
Expand Down
5 changes: 5 additions & 0 deletions apps/sim/lib/mothership/chat/terminal-state.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}))
Expand Down Expand Up @@ -72,6 +76,7 @@ describe('finalizeAssistantTurn', () => {
status: 'executing',
},
},
{ ...envelope, seq: 3, type: 'complete', payload: { status: 'cancelled' } },
])
await finalizeAssistantTurn({
chatId: 'chat-1',
Expand Down
18 changes: 14 additions & 4 deletions apps/sim/lib/mothership/chat/terminal-state.ts
Original file line number Diff line number Diff line change
@@ -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 {
Expand All @@ -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'

Expand All @@ -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<PersistedMessage | null> {
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' ||

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.

P1 Late title event blocks replay

On a first turn, title generation runs independently and can publish a session event after the run's complete event. This check then rejects a complete, gap-free replay. If the owner's message save failed and an identifiers-only Stop retries persistence, the retry returns success without saving the response. The completion check needs to allow later side-effect events.

!events.every((event, index) => event.seq === index + 1)
)
return null
const replay = buildLiveAssistantMessage({
streamId,
events: events.map(toStreamBatchEvent),
Expand Down Expand Up @@ -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() ||
Expand Down
Loading