Skip to content

Commit aea7ed0

Browse files
committed
refactor(mothership): tidy replay compaction and stream retry after review
- Share the replay write ceiling between compaction and file previews - Scan for the smallest sufficient child only once one can cover the need, and reuse one size memo across every compaction step - Stop retrying a raw TypeError; fetch and body-read failures are already wrapped as WorkerUnreachableError and WorkerStreamInterruptedError - Use absolute imports, a private retry counter, a typed unsettled-state set, and hasUnsettledTool naming - Cover a timed-out body read and the preview completion edge cases
1 parent 7aa6002 commit aea7ed0

14 files changed

Lines changed: 124 additions & 120 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/hooks/message-reconcile.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -146,11 +146,11 @@ export function buildAssistantSnapshotMessage(params: {
146146
}
147147

148148
export function markMessageStopped(message: PersistedMessage): PersistedMessage {
149-
const hasExecutingTool = message.contentBlocks?.some((block) =>
149+
const hasUnsettledTool = message.contentBlocks?.some((block) =>
150150
isUnsettledToolState(block.toolCall?.state)
151151
)
152152
const hasOpenBlock = message.contentBlocks?.some((block) => block.endedAt === undefined)
153-
if (!hasExecutingTool && !hasOpenBlock) {
153+
if (!hasUnsettledTool && !hasOpenBlock) {
154154
return message
155155
}
156156

‎apps/sim/app/workspace/[workspaceId]/home/hooks/preview/apply-file-preview-phase.test.ts‎

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,9 @@
11
import { describe, expect, it } from 'vitest'
22
import type { FilePreviewSession } from '@/lib/mothership/request/session'
3-
import { deriveFilePreviewSession, previewHoldsFinalContent } from './apply-file-preview-phase'
3+
import {
4+
deriveFilePreviewSession,
5+
previewHoldsFinalContent,
6+
} from '@/app/workspace/[workspaceId]/home/hooks/preview/apply-file-preview-phase'
47

58
const NOW = '2026-06-08T00:00:00.000Z'
69

@@ -118,4 +121,16 @@ describe('previewHoldsFinalContent', () => {
118121
it('does not hold it when no content was received at all', () => {
119122
expect(previewHoldsFinalContent(undefined, complete(7))).toBe(false)
120123
})
124+
125+
it('does not hold it when the session exists but received no text', () => {
126+
const prev = session({ previewText: '', previewVersion: 7 })
127+
128+
expect(previewHoldsFinalContent(prev, complete(7))).toBe(false)
129+
})
130+
131+
it('holds the received text when the completion carries no version to compare', () => {
132+
const prev = session({ previewText: 'final text', previewVersion: 3 })
133+
134+
expect(previewHoldsFinalContent(prev, complete())).toBe(true)
135+
})
121136
})

‎apps/sim/app/workspace/[workspaceId]/home/hooks/use-chat.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -4370,11 +4370,11 @@ export function useChat(
43704370
} else {
43714371
setPendingMessages((prev) =>
43724372
prev.map((msg) => {
4373-
const hasExecutingTool = msg.contentBlocks?.some((block) =>
4373+
const hasUnsettledTool = msg.contentBlocks?.some((block) =>
43744374
isUnsettledToolState(block.toolCall?.status)
43754375
)
43764376
const hasOpenBlock = msg.contentBlocks?.some((block) => block.endedAt === undefined)
4377-
if (!hasExecutingTool && !hasOpenBlock) {
4377+
if (!hasUnsettledTool && !hasOpenBlock) {
43784378
return msg
43794379
}
43804380
const updatedBlocks: ContentBlock[] = (msg.contentBlocks ?? []).map((block) => ({

‎apps/sim/lib/mothership/chat/display-message.ts‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,8 @@
11
import type { PersistedContentBlock } from '@/lib/api/contracts/copilot-messages'
2+
import { getMothershipAttachmentPreviewUrl } from '@/lib/mothership/chat/attachment-preview'
3+
import { isLiveAssistantMessageId } from '@/lib/mothership/chat/live-message-id'
4+
import type { PersistedMessage } from '@/lib/mothership/chat/persisted-message'
5+
import { isUnsettledToolState, withBlockTiming } from '@/lib/mothership/chat/persisted-message'
26
import {
37
MothershipStreamV1CompletionStatus,
48
MothershipStreamV1EventType,
@@ -17,10 +21,6 @@ import {
1721
type ToolCallInfo,
1822
ToolCallStatus,
1923
} from '@/app/workspace/[workspaceId]/home/types'
20-
import { getMothershipAttachmentPreviewUrl } from './attachment-preview'
21-
import { isLiveAssistantMessageId } from './live-message-id'
22-
import type { PersistedMessage } from './persisted-message'
23-
import { isUnsettledToolState, withBlockTiming } from './persisted-message'
2424

2525
const STATE_TO_STATUS: Record<string, ToolCallStatus> = {
2626
[MothershipStreamV1ToolOutcome.success]: ToolCallStatus.success,

‎apps/sim/lib/mothership/chat/persisted-message.ts‎

Lines changed: 8 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,11 @@ import {
2121
MothershipStreamV1ToolOutcome,
2222
MothershipStreamV1ToolPhase,
2323
} from '@/lib/mothership/generated/mothership-stream-v1'
24-
import type { ContentBlock, OrchestratorResult } from '@/lib/mothership/request/types'
24+
import type {
25+
ContentBlock,
26+
LocalToolCallStatus,
27+
OrchestratorResult,
28+
} from '@/lib/mothership/request/types'
2529
import { RETIRED_BROWSER_REQUEST_TAKEOVER_ID } from '@/lib/mothership/tools/retired-tools'
2630
import { normalizeToolActivityDescription } from '@/lib/mothership/tools/tool-display'
2731
import type { BrowserTextSelection, TerminalTextSelection } from '@/stores/panel/types'
@@ -362,15 +366,16 @@ export function buildPersistedAssistantMessage(
362366
return message
363367
}
364368

365-
const UNSETTLED_TOOL_STATES: ReadonlySet<string> = new Set([
369+
const UNSETTLED_TOOL_STATES: ReadonlySet<LocalToolCallStatus> = new Set<LocalToolCallStatus>([
366370
'pending',
367371
'executing',
368372
'awaiting_approval',
369373
])
370374

371375
/** A tool row that has not finished: waiting to run, running, or awaiting a decision. */
372376
export function isUnsettledToolState(state: string | undefined): boolean {
373-
return state !== undefined && UNSETTLED_TOOL_STATES.has(state)
377+
const unsettled: ReadonlySet<string | undefined> = UNSETTLED_TOOL_STATES
378+
return unsettled.has(state)
374379
}
375380

376381
/** Settles every unfinished tool row at a turn terminal so none reloads as a spinner. */

‎apps/sim/lib/mothership/request/context/restore.ts‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -14,10 +14,12 @@ export async function restoreStreamingContext(
1414
): Promise<void> {
1515
for (const saved of events) {
1616
// Preview content already in the replay counts toward this turn's preview budget.
17-
if (saved.type === 'tool' && 'previewPhase' in saved.payload) {
18-
if (saved.payload.previewPhase === 'file_preview_content') {
19-
context.filePreviewBudget.contentBytes += Buffer.byteLength(saved.payload.content, 'utf8')
20-
}
17+
if (
18+
saved.type === 'tool' &&
19+
'previewPhase' in saved.payload &&
20+
saved.payload.previewPhase === 'file_preview_content'
21+
) {
22+
context.filePreviewBudget.contentBytes += Buffer.byteLength(saved.payload.content, 'utf8')
2123
}
2224
const event = reconcileTextEvent(saved, context.accumulatedContent)
2325
if (!event) continue

‎apps/sim/lib/mothership/request/go/file-preview-adapter.test.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -25,11 +25,11 @@ import { createStreamingContext } from '@/lib/mothership/request/context/request
2525
import {
2626
createFilePreviewAdapterState,
2727
type FilePreviewAdapterState,
28-
PREVIEW_FRAME_MAX_BYTES,
2928
PREVIEW_TURN_CONTENT_BYTES,
3029
processFilePreviewStreamEvent,
3130
} from '@/lib/mothership/request/go/file-preview-adapter'
3231
import { createEvent, eventToStreamEvent } from '@/lib/mothership/request/session'
32+
import { STREAM_EVENT_MAX_PAYLOAD_BYTES } from '@/lib/mothership/request/session/replay-compaction'
3333
import type {
3434
ActiveFileIntent,
3535
ExecutionContext,
@@ -262,7 +262,7 @@ describe('processFilePreviewStreamEvent — preview byte rate', () => {
262262

263263
for (const payload of payloads) {
264264
expect(Buffer.byteLength(JSON.stringify(payload))).toBeLessThanOrEqual(
265-
PREVIEW_FRAME_MAX_BYTES
265+
STREAM_EVENT_MAX_PAYLOAD_BYTES
266266
)
267267
}
268268
expect(completed).toBe(true)

‎apps/sim/lib/mothership/request/go/file-preview-adapter.ts‎

Lines changed: 3 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
import { createLogger } from '@sim/logger'
22
import { toError } from '@sim/utils/errors'
33
import { isRecordLike } from '@sim/utils/object'
4-
import { getRedisBudgetLimits } from '@/lib/core/redis/byte-budget.server'
54
import { executeCopilotFileUseCase } from '@/lib/mothership/application/execute-file-use-case'
65
import { MothershipStreamV1EventType } from '@/lib/mothership/generated/mothership-stream-v1'
76
import {
@@ -15,6 +14,7 @@ import {
1514
type SyntheticFilePreviewPayload,
1615
upsertFilePreviewSession,
1716
} from '@/lib/mothership/request/session'
17+
import { STREAM_EVENT_MAX_PAYLOAD_BYTES } from '@/lib/mothership/request/session/replay-compaction'
1818
import type {
1919
ActiveFileIntent,
2020
ExecutionContext,
@@ -70,13 +70,6 @@ const PREVIEW_SNAPSHOT_CHARS_PER_SECOND = 256 * 1024
7070
*/
7171
export const PREVIEW_TURN_CONTENT_BYTES = 8 * 1024 * 1024
7272

73-
/** Room left in a replay write for the envelope around a preview payload. */
74-
const PREVIEW_FRAME_ENVELOPE_HEADROOM_BYTES = 16 * 1024
75-
76-
/** The largest serialized preview payload the replay buffer can persist in one write. */
77-
export const PREVIEW_FRAME_MAX_BYTES =
78-
getRedisBudgetLimits('copilot_stream').maxSingleWriteBytes - PREVIEW_FRAME_ENVELOPE_HEADROOM_BYTES
79-
8073
/** The minimum gap between full snapshots of a preview this long. */
8174
function snapshotIntervalMs(baseMs: number, previewText: string): number {
8275
return Math.max(baseMs, (previewText.length / PREVIEW_SNAPSHOT_CHARS_PER_SECOND) * 1000)
@@ -397,13 +390,13 @@ async function emitPreviewEvent(
397390
if (frame.previewPhase === 'file_preview_content') {
398391
const budget = context.filePreviewBudget
399392
const contentBytes = budget.contentBytes + Buffer.byteLength(frame.content, 'utf8')
400-
if (frameBytes > PREVIEW_FRAME_MAX_BYTES || contentBytes > PREVIEW_TURN_CONTENT_BYTES) {
393+
if (frameBytes > STREAM_EVENT_MAX_PAYLOAD_BYTES || contentBytes > PREVIEW_TURN_CONTENT_BYTES) {
401394
return false
402395
}
403396
budget.contentBytes = contentBytes
404397
} else if (
405398
frame.previewPhase === 'file_preview_complete' &&
406-
frameBytes > PREVIEW_FRAME_MAX_BYTES
399+
frameBytes > STREAM_EVENT_MAX_PAYLOAD_BYTES
407400
) {
408401
const { output: _output, ...withoutOutput } = frame
409402
frame = withoutOutput

‎apps/sim/lib/mothership/request/go/stream.test.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -832,6 +832,39 @@ describe('copilot go stream helpers', () => {
832832
})
833833
})
834834

835+
it('keeps the timeout error when the body read fails after the request timed out', async () => {
836+
const first = createEvent({
837+
streamId: 'timed-out-stream',
838+
cursor: '1',
839+
seq: 1,
840+
requestId: 'req-timed-out',
841+
type: 'text',
842+
payload: { channel: 'assistant', text: 'partial' },
843+
})
844+
vi.mocked(fetch).mockImplementationOnce(async (_url, init) => {
845+
const signal = init?.signal
846+
return new Response(
847+
new ReadableStream<Uint8Array>({
848+
start(controller) {
849+
controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(first)}\n\n`))
850+
signal?.addEventListener('abort', () => controller.error(signal.reason))
851+
},
852+
}),
853+
{ status: 200, headers: { 'Content-Type': 'text/event-stream' } }
854+
)
855+
})
856+
857+
await expect(
858+
runStreamLoop(
859+
'https://example.com/mothership/stream',
860+
{},
861+
createStreamingContext(),
862+
turnScopedExecContext(),
863+
{ timeout: 20, flushAfterEvent: false }
864+
)
865+
).rejects.toMatchObject({ name: 'TimeoutError' })
866+
})
867+
835868
it('reports a worker it could not reach without the raw network error', async () => {
836869
const networkError = new TypeError('fetch failed')
837870
vi.mocked(fetch).mockRejectedValueOnce(networkError)

‎apps/sim/lib/mothership/request/lifecycle/run.test.ts‎

Lines changed: 5 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -221,6 +221,7 @@ import {
221221
CopilotBackendError,
222222
STREAM_ENDED_WITHOUT_TERMINAL_MESSAGE,
223223
StreamEndedWithoutTerminalError,
224+
WorkerUnreachableError,
224225
} from '@/lib/mothership/request/go/stream'
225226
import { runCopilotLifecycle } from '@/lib/mothership/request/lifecycle/run'
226227
import {
@@ -2718,7 +2719,7 @@ describe('runCopilotLifecycle', () => {
27182719
headers.push(new Headers(request.headers))
27192720
context.accumulatedContent = 'Saved partial answer'
27202721
context.errors.push('connection interrupted')
2721-
throw new TypeError('fetch failed')
2722+
throw new WorkerUnreachableError(new TypeError('fetch failed'))
27222723
}
27232724
)
27242725
}
@@ -2811,7 +2812,9 @@ describe('runCopilotLifecycle', () => {
28112812
vi.useFakeTimers()
28122813
try {
28132814
const controller = new AbortController()
2814-
mockRunStreamLoop.mockRejectedValueOnce(new TypeError('fetch failed'))
2815+
mockRunStreamLoop.mockRejectedValueOnce(
2816+
new WorkerUnreachableError(new TypeError('fetch failed'))
2817+
)
28152818
const pending = runCopilotLifecycle(
28162819
{ message: 'hello', messageId: 'stopped-outage' },
28172820
{
@@ -2929,50 +2932,6 @@ describe('runCopilotLifecycle', () => {
29292932
expect(JSON.stringify(persisted)).not.toMatch(/<html|Bad Gateway|nginx/)
29302933
})
29312934

2932-
it('stops a TypeError that repeats while handling each reattach at three retries', async () => {
2933-
vi.useFakeTimers()
2934-
try {
2935-
let attempts = 0
2936-
mockRunStreamLoop.mockImplementation(
2937-
async (
2938-
_url: string,
2939-
_init: RequestInit,
2940-
_context: StreamingContext,
2941-
_exec: ExecutionContext,
2942-
options: { onEvent?: (event: unknown) => Promise<void> }
2943-
): Promise<void> => {
2944-
attempts++
2945-
await options.onEvent?.({ type: 'session', payload: { kind: 'start' } })
2946-
throw new TypeError("Cannot read properties of undefined (reading 'payload')")
2947-
}
2948-
)
2949-
2950-
const pending = runCopilotLifecycle(
2951-
{ message: 'hello', messageId: 'stream-repeated-type-error' },
2952-
{
2953-
userId: 'user-1',
2954-
workspaceId: 'ws-1',
2955-
chatId: 'chat-1',
2956-
executionId: 'exec-1',
2957-
runId: 'run-1',
2958-
executionContext: {
2959-
userId: 'user-1',
2960-
workflowId: '',
2961-
workspaceId: 'ws-1',
2962-
chatId: 'chat-1',
2963-
},
2964-
}
2965-
)
2966-
await vi.advanceTimersByTimeAsync(60_000)
2967-
2968-
expect(attempts).toBe(4)
2969-
expect(await pending).toEqual(expect.objectContaining({ success: false }))
2970-
} finally {
2971-
mockRunStreamLoop.mockReset()
2972-
vi.useRealTimers()
2973-
}
2974-
})
2975-
29762935
it('stops a failure that repeats after every reattach at three retries', async () => {
29772936
vi.useFakeTimers()
29782937
try {

0 commit comments

Comments
 (0)