Skip to content

Commit 7aa6002

Browse files
committed
fix(mothership): shed oversized replay events in one pass and retry mid-body stream cuts
- Replace the replay-compaction omit loop with one bounded pass over memoized exact serialized sizes: recurse into the smallest child that covers the overage, otherwise drop the largest children whole, keeping small siblings - Map a non-abort body-read failure from the worker stream to a retryable WorkerStreamInterruptedError on the reachable budget, with the generic message and the original error on cause - Log the cause message on retry warnings and orchestration failures - Count replayed file preview content toward a restored turn's preview budget
1 parent fc2149c commit 7aa6002

10 files changed

Lines changed: 357 additions & 49 deletions

File tree

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

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,12 @@ export async function restoreStreamingContext(
1313
execContext: ExecutionContext
1414
): Promise<void> {
1515
for (const saved of events) {
16+
// 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+
}
21+
}
1622
const event = reconcileTextEvent(saved, context.accumulatedContent)
1723
if (!event) continue
1824
const replay: StreamEvent =

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

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -785,6 +785,53 @@ describe('copilot go stream helpers', () => {
785785
expect(fetch).toHaveBeenCalledTimes(1)
786786
})
787787

788+
it('reports a stream cut mid-body without the raw socket error', async () => {
789+
const socketError = Object.assign(
790+
new Error(
791+
'The socket connection was closed unexpectedly. For more information, pass `verbose: true`'
792+
),
793+
{ code: 'ECONNRESET' }
794+
)
795+
const first = createEvent({
796+
streamId: 'cut-stream',
797+
cursor: '1',
798+
seq: 1,
799+
requestId: 'req-cut',
800+
type: 'text',
801+
payload: { channel: 'assistant', text: 'partial' },
802+
})
803+
vi.mocked(fetch).mockResolvedValueOnce(
804+
new Response(
805+
new ReadableStream<Uint8Array>({
806+
start(controller) {
807+
controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(first)}\n\n`))
808+
},
809+
pull(controller) {
810+
controller.error(socketError)
811+
},
812+
}),
813+
{ status: 200, headers: { 'Content-Type': 'text/event-stream' } }
814+
)
815+
)
816+
817+
await expect(
818+
runStreamLoop(
819+
'https://example.com/mothership/stream',
820+
{},
821+
createStreamingContext(),
822+
turnScopedExecContext(),
823+
{
824+
timeout: 1000,
825+
flushAfterEvent: false,
826+
}
827+
)
828+
).rejects.toMatchObject({
829+
name: 'WorkerStreamInterruptedError',
830+
message: 'The agent service is temporarily unavailable. Please try again.',
831+
cause: socketError,
832+
})
833+
})
834+
788835
it('reports a worker it could not reach without the raw network error', async () => {
789836
const networkError = new TypeError('fetch failed')
790837
vi.mocked(fetch).mockRejectedValueOnce(networkError)

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

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,18 @@ export class WorkerUnreachableError extends Error {
8181
}
8282
}
8383

84+
/**
85+
* The worker's response body failed mid-stream (the connection was reset or
86+
* closed). The worker answered, so a retry reattaches under the short budget.
87+
* The read error stays on `cause` for logs.
88+
*/
89+
export class WorkerStreamInterruptedError extends Error {
90+
constructor(cause: unknown) {
91+
super(BACKEND_UNAVAILABLE_MESSAGE, { cause })
92+
this.name = 'WorkerStreamInterruptedError'
93+
}
94+
}
95+
8496
/**
8597
* A worker rejection message the user can act on: short, one line, plain text,
8698
* and free of identifiers (`userId`, `protocol_version_mismatch`) that only mean
@@ -293,7 +305,13 @@ export async function runStreamLoop(
293305
const rawReader = response.body.getReader()
294306
const reader: ReadableStreamDefaultReader<Uint8Array> = {
295307
async read() {
296-
const result = await rawReader.read()
308+
let result: ReadableStreamReadResult<Uint8Array>
309+
try {
310+
result = await rawReader.read()
311+
} catch (error) {
312+
if (requestSignal.aborted) throw error
313+
throw new WorkerStreamInterruptedError(error)
314+
}
297315
if (!result.done && result.value) {
298316
const now = performance.now()
299317
const gap = now - counters.lastChunkMs

‎apps/sim/lib/mothership/request/handlers/handlers.test.ts‎

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -215,6 +215,29 @@ describe('sse-handlers tool lifecycle', () => {
215215
expect(context.subAgentTraceSpans?.size).toBe(0)
216216
})
217217

218+
it('counts the preview content already in the replay toward the restored turn budget', async () => {
219+
const preview = (content: string): StreamEvent => ({
220+
type: 'tool',
221+
payload: {
222+
toolCallId: 'file-edit',
223+
toolName: 'prepare_file_edit',
224+
previewPhase: 'file_preview_content',
225+
content,
226+
contentMode: 'delta',
227+
previewVersion: 1,
228+
fileName: 'notes.md',
229+
},
230+
})
231+
232+
await restoreStreamingContext(
233+
[preview('é'.repeat(1_000)), preview('x'.repeat(500))],
234+
context,
235+
execContext
236+
)
237+
238+
expect(context.filePreviewBudget.contentBytes).toBe(2_500)
239+
})
240+
218241
it('restores a delivered prefix without executing its tools or consuming the next handoff', async () => {
219242
const call: StreamEvent = {
220243
type: 'tool',

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -128,9 +128,17 @@ vi.mock('@/lib/mothership/request/go/stream', () => {
128128
}
129129
}
130130

131+
class WorkerStreamInterruptedError extends Error {
132+
constructor(cause: unknown) {
133+
super('The agent service is temporarily unavailable. Please try again.', { cause })
134+
this.name = 'WorkerStreamInterruptedError'
135+
}
136+
}
137+
131138
return {
132139
BillingLimitError,
133140
CopilotBackendError,
141+
WorkerStreamInterruptedError,
134142
WorkerUnreachableError,
135143
STREAM_ENDED_WITHOUT_TERMINAL_MESSAGE,
136144
StreamEndedWithoutTerminalError,

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -593,6 +593,7 @@ export async function runCopilotLifecycle(
593593
// explained, not just reduced to a message string.
594594
logger.error('Copilot orchestration failed', {
595595
error: err.message,
596+
...causeForLog(err),
596597
name: err.name,
597598
...(error instanceof CopilotBackendError
598599
? { backendStatus: error.status, backendBody: error.body?.slice(0, 2000) }
@@ -878,6 +879,7 @@ async function runResumeLegWithRetry(
878879
attempt: retry.attempts + 1,
879880
backoffMs: backoff,
880881
error: toError(error).message,
882+
...causeForLog(error),
881883
})
882884
await interruptibleSleep(backoff, options.abortSignal)
883885
continue
@@ -1240,6 +1242,7 @@ async function runCheckpointLoop(
12401242
attempt: (retry?.attempts ?? 0) + 1,
12411243
backoffMs: backoff,
12421244
error: toError(streamError).message,
1245+
...causeForLog(streamError),
12431246
}
12441247
)
12451248
await interruptibleSleep(backoff, options.abortSignal)
@@ -1678,6 +1681,12 @@ async function withEnterpriseByokKey(
16781681
return byokApiKey ? { ...refreshed, byokApiKey } : refreshed
16791682
}
16801683

1684+
/** The underlying failure behind a generic user-facing error, for logs. */
1685+
function causeForLog(error: unknown): { cause?: string } {
1686+
const cause = error instanceof Error ? error.cause : undefined
1687+
return cause === undefined ? {} : { cause: getErrorMessage(cause) }
1688+
}
1689+
16811690
function isAborted(options: CopilotLifecycleOptions, context: StreamingContext): boolean {
16821691
return !!(options.abortSignal?.aborted || context.wasAborted)
16831692
}

‎apps/sim/lib/mothership/request/lifecycle/stream-retry.test.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ import { afterEach, describe, expect, it, vi } from 'vitest'
22
import {
33
CopilotBackendError,
44
StreamEndedWithoutTerminalError,
5+
WorkerStreamInterruptedError,
56
WorkerUnreachableError,
67
} from '@/lib/mothership/request/go/stream'
78
import { StreamRetryWindow } from '@/lib/mothership/request/lifecycle/stream-retry'
@@ -150,4 +151,16 @@ describe('stream recovery budget', () => {
150151
expect(retry.nextDelay(new Error('Invalid operation'))).toBeNull()
151152
expect(retry.attempt).toBe(0)
152153
})
154+
155+
it('gives a stream cut mid-body the reachable budget, not the unreachable window', () => {
156+
vi.useFakeTimers()
157+
const error = new WorkerStreamInterruptedError(new Error('socket closed'))
158+
const retry = new StreamRetryWindow()
159+
for (let index = 0; index < 3; index++) {
160+
const delay = retry.nextDelay(error)
161+
expect(delay).not.toBeNull()
162+
vi.advanceTimersByTime(delay ?? 0)
163+
}
164+
expect(retry.nextDelay(error)).toBeNull()
165+
})
153166
})

‎apps/sim/lib/mothership/request/lifecycle/stream-retry.ts‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import { StreamContinuityError } from '@/lib/mothership/request/go/parser'
44
import {
55
CopilotBackendError,
66
StreamEndedWithoutTerminalError,
7+
WorkerStreamInterruptedError,
78
WorkerUnreachableError,
89
} from '@/lib/mothership/request/go/stream'
910

@@ -108,7 +109,11 @@ function isJson(body: string | undefined): boolean {
108109
/** Initial sends and resumes both replay one durable identity after an ambiguous response. */
109110
function isRetryableStreamError(error: unknown): boolean {
110111
if (error instanceof Error && error.name === 'AbortError') return false
111-
if (error instanceof StreamEndedWithoutTerminalError || error instanceof StreamContinuityError) {
112+
if (
113+
error instanceof StreamEndedWithoutTerminalError ||
114+
error instanceof StreamContinuityError ||
115+
error instanceof WorkerStreamInterruptedError
116+
) {
112117
return true
113118
}
114119
if (error instanceof CopilotBackendError) {

‎apps/sim/lib/mothership/request/session/replay-compaction.test.ts‎

Lines changed: 126 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import {
1616
STREAM_EVENT_COMPACTION_THRESHOLD_BYTES,
1717
STREAM_EVENT_MAX_PAYLOAD_BYTES,
1818
STREAM_STRING_PREVIEW_UNITS,
19+
serializedBytes,
1920
} from '@/lib/mothership/request/session/replay-compaction'
2021
import type { StreamEvent } from '@/lib/mothership/request/session/types'
2122

@@ -343,4 +344,129 @@ describe('compactStreamEvent', () => {
343344

344345
expect(compactStreamEvent(event)).toBe(event)
345346
})
347+
348+
it('omits only the smallest sufficient bulk when no single child dominates, keeping its siblings', () => {
349+
const manyKeys = (prefix: string) =>
350+
Object.fromEntries(
351+
Array.from({ length: 2_300 }, (_, index) => [`${prefix}-${index}`, 'v'.repeat(300)])
352+
)
353+
const event: StreamEvent = {
354+
type: 'tool',
355+
payload: {
356+
toolCallId: 'c',
357+
toolName: 'cli_tables_rows_query',
358+
executor: 'sim',
359+
mode: 'async',
360+
phase: 'result',
361+
success: true,
362+
output: { fileId: 'file-1', rows: manyKeys('r'), cols: manyKeys('c') },
363+
},
364+
}
365+
366+
const compacted = payloadOf(compactStreamEvent(event))
367+
const output = toRecord(compacted.output)
368+
369+
expect(Buffer.byteLength(JSON.stringify(compacted))).toBeLessThanOrEqual(
370+
STREAM_EVENT_MAX_PAYLOAD_BYTES
371+
)
372+
expect(output.fileId).toBe('file-1')
373+
expect([output.rows, output.cols].filter((value) => typeof value === 'string')).toHaveLength(1)
374+
})
375+
376+
it('omits the smaller of two bulks when either alone would make the event fit', () => {
377+
const manyKeys = (prefix: string, count: number) =>
378+
Object.fromEntries(
379+
Array.from({ length: count }, (_, index) => [`${prefix}-${index}`, 'v'.repeat(300)])
380+
)
381+
const event: StreamEvent = {
382+
type: 'tool',
383+
payload: {
384+
toolCallId: 'c',
385+
toolName: 'cli_tables_rows_query',
386+
executor: 'sim',
387+
mode: 'async',
388+
phase: 'result',
389+
success: true,
390+
output: { fileId: 'file-1', rows: manyKeys('r', 3_000), cols: manyKeys('c', 1_900) },
391+
},
392+
}
393+
394+
const output = toRecord(payloadOf(compactStreamEvent(event)).output)
395+
396+
expect(output.fileId).toBe('file-1')
397+
expect(Object.keys(toRecord(output.rows))).toHaveLength(3_000)
398+
expect(output.cols).toMatch(/^…\[omitted, [\d.]+ KB total\]$/)
399+
})
400+
401+
it('omits a sufficient bulk whole when its own large parts cannot cover the overage', () => {
402+
const manyKeys = (prefix: string, count: number) =>
403+
Object.fromEntries(
404+
Array.from({ length: count }, (_, index) => [`${prefix}-${index}`, 'v'.repeat(3_000)])
405+
)
406+
const event: StreamEvent = {
407+
type: 'tool',
408+
payload: {
409+
toolCallId: 'c',
410+
toolName: 'cli_workflows_state_get',
411+
executor: 'sim',
412+
mode: 'async',
413+
phase: 'result',
414+
success: true,
415+
output: { workflowId: 'wf-1', state: { nested: manyKeys('n', 35), ...manyKeys('k', 400) } },
416+
},
417+
}
418+
419+
const compacted = payloadOf(compactStreamEvent(event))
420+
const output = toRecord(compacted.output)
421+
422+
expect(Buffer.byteLength(JSON.stringify(compacted))).toBeLessThanOrEqual(
423+
STREAM_EVENT_MAX_PAYLOAD_BYTES
424+
)
425+
expect(output.workflowId).toBe('wf-1')
426+
expect(output.state).toMatch(/^…\[omitted, [\d.]+ MB total\]$/)
427+
})
428+
429+
it('bounds a balanced 16 MB tree in one pass', () => {
430+
const tree = (depth: number, leaf: number): unknown =>
431+
depth === 0
432+
? 'x'.repeat(leaf)
433+
: { l: tree(depth - 1, leaf + 100), r: tree(depth - 1, Math.max(leaf - 100, 1)) }
434+
const event: StreamEvent = {
435+
type: 'tool',
436+
payload: {
437+
toolCallId: 'c',
438+
toolName: 'cli_blocks_get',
439+
executor: 'sim',
440+
mode: 'async',
441+
phase: 'result',
442+
success: true,
443+
output: tree(12, 3_900),
444+
},
445+
}
446+
447+
const started = performance.now()
448+
const compacted = compactStreamEvent(event)
449+
const elapsedMs = performance.now() - started
450+
451+
expect(Buffer.byteLength(JSON.stringify(compacted.payload))).toBeLessThanOrEqual(
452+
STREAM_EVENT_MAX_PAYLOAD_BYTES
453+
)
454+
expect(elapsedMs).toBeLessThan(2_000)
455+
})
456+
457+
it('measures serialized size exactly without serializing', () => {
458+
const values: unknown[] = [
459+
{ a: 'é😀"\\\n', b: [1, null, true, { c: 'd' }], e: undefined, f: -1.5e-7 },
460+
[undefined, 'x', { 'kéy "q"': 0 }],
461+
'plain',
462+
42,
463+
null,
464+
{},
465+
[],
466+
]
467+
468+
for (const value of values) {
469+
expect(serializedBytes(value)).toBe(Buffer.byteLength(JSON.stringify(value)))
470+
}
471+
})
346472
})

0 commit comments

Comments
 (0)