-
Notifications
You must be signed in to change notification settings - Fork 3.9k
Preserve complete stopped responses and cancel unfinished tools #8259
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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' | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. If Stop arrives while a tool is awaiting approval, this check leaves it unchanged because it only cancels
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
|
||
| ? { | ||
| ...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({ | ||
|
|
||
| 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 { | ||
|
|
@@ -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<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' || | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
On a first turn, title generation runs independently and can publish a |
||
| !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() || | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
withStoppedContentBlockimport uses the relative./persisted-messagepath. The repository requires absolute imports inapps/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!