Skip to content

Commit 9ee8438

Browse files
committed
fix(mothership): omit the largest field of an event no cut can bound, and persist closed lanes
- An event whose many short values survive every string and array cut, such as an object with thousands of keys, now has its largest bulk field replaced by a size note when it would still exceed one replay write, instead of ending the turn. Identity fields and client-executed arguments are never replaced; such an event is still refused cleanly. - finalizeResidualToolCalls reports closing an open subagent lane as a change, so the error finalizer persists it instead of leaving a stale lane.
1 parent 89f32d1 commit 9ee8438

4 files changed

Lines changed: 88 additions & 2 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/home/hooks/stream/stream-helpers.test.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,4 +27,11 @@ describe('finalizeResidualToolCalls', () => {
2727
expect(finalizeResidualToolCalls(open, 'error')).toBe(true)
2828
expect(finalizeResidualToolCalls(settled, 'error')).toBe(false)
2929
})
30+
31+
it('reports closing an open subagent lane as a change to persist', () => {
32+
const blocks: ContentBlock[] = [{ type: 'subagent', content: 'research' }]
33+
34+
expect(finalizeResidualToolCalls(blocks, 'error')).toBe(true)
35+
expect(blocks[0].endedAt).toEqual(expect.any(Number))
36+
})
3037
})

‎apps/sim/app/workspace/[workspaceId]/home/hooks/stream/stream-helpers.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ export function asPayloadRecord(value: unknown): StreamPayload | undefined {
7272
* Settles every unfinished tool row (running, pending, or awaiting approval) at
7373
* a turn terminal by propagating the turn's outcome: a clean `complete` settles
7474
* a straggler `success`, a stop `cancelled`, an error `error`. Also closes any
75-
* open subagent lane. Returns whether any tool row was settled.
75+
* open subagent lane. Returns whether it settled a row or closed a lane.
7676
*/
7777
export function finalizeResidualToolCalls(
7878
blocks: ContentBlock[],
@@ -94,6 +94,7 @@ export function finalizeResidualToolCalls(
9494
// transport-based gating.
9595
if (block.type === 'subagent' && block.endedAt === undefined) {
9696
block.endedAt = endedAt
97+
settled = true
9798
continue
9899
}
99100
const tc = block.toolCall

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

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import { describe, expect, it } from 'vitest'
1414
import {
1515
compactStreamEvent,
1616
STREAM_EVENT_COMPACTION_THRESHOLD_BYTES,
17+
STREAM_EVENT_MAX_PAYLOAD_BYTES,
1718
STREAM_STRING_PREVIEW_UNITS,
1819
} from '@/lib/mothership/request/session/replay-compaction'
1920
import type { StreamEvent } from '@/lib/mothership/request/session/types'
@@ -237,4 +238,49 @@ describe('compactStreamEvent', () => {
237238
expect(output.data).toEqual({ results: citations })
238239
expect(toRecord(payloadOf(event).output).rows).toHaveLength(40_000)
239240
})
241+
242+
it('omits the largest remaining field when cuts alone leave the event over one replay write', () => {
243+
const blocks = Object.fromEntries(
244+
Array.from({ length: 5_000 }, (_, index) => [`block-${index}`, 'b'.repeat(300)])
245+
)
246+
const event: StreamEvent = {
247+
type: 'tool',
248+
payload: {
249+
toolCallId: 'c',
250+
toolName: 'cli_workflows_state_get',
251+
executor: 'sim',
252+
mode: 'async',
253+
phase: 'result',
254+
success: true,
255+
status: 'success',
256+
output: { blocks },
257+
},
258+
}
259+
260+
const compacted = payloadOf(compactStreamEvent(event))
261+
262+
expect(Buffer.byteLength(JSON.stringify(compacted))).toBeLessThanOrEqual(
263+
STREAM_EVENT_MAX_PAYLOAD_BYTES
264+
)
265+
expect(compacted.output).toMatch(/^…\[omitted, [\d.]+ MB total\]$/)
266+
expect(compacted).toMatchObject({ toolCallId: 'c', success: true, status: 'success' })
267+
})
268+
269+
it('never omits identity to make room for client-executed arguments it must keep whole', () => {
270+
const args = { workflowId: 'wf-1', input: 'w'.repeat(1.5 * MB) }
271+
const event: StreamEvent = {
272+
type: 'tool',
273+
payload: {
274+
toolCallId: 'c',
275+
toolName: 'run_workflow',
276+
executor: 'client',
277+
mode: 'async',
278+
phase: 'call',
279+
ui: { clientExecutable: true },
280+
arguments: args,
281+
},
282+
}
283+
284+
expect(compactStreamEvent(event)).toBe(event)
285+
})
240286
})

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

Lines changed: 33 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
import { isRecordLike, toRecordOrNull } from '@sim/utils/object'
2+
import { getRedisBudgetLimits } from '@/lib/core/redis/byte-budget.server'
23
import {
34
MothershipStreamV1EventType,
45
MothershipStreamV1ToolPhase,
@@ -13,6 +14,13 @@ import type { StreamEvent } from './types'
1314
*/
1415
export const STREAM_EVENT_COMPACTION_THRESHOLD_BYTES = 256 * 1024
1516

17+
/** Room left in a replay write for the envelope around an event's payload. */
18+
const ENVELOPE_HEADROOM_BYTES = 16 * 1024
19+
20+
/** The largest serialized payload the replay buffer can persist in one write. */
21+
export const STREAM_EVENT_MAX_PAYLOAD_BYTES =
22+
getRedisBudgetLimits('copilot_stream').maxSingleWriteBytes - ENVELOPE_HEADROOM_BYTES
23+
1624
/** UTF-16 units kept at the head of a long string. */
1725
export const STREAM_STRING_PREVIEW_UNITS = 8 * 1024
1826

@@ -100,11 +108,32 @@ function mapLeaves(
100108
return copy ?? value
101109
}
102110

111+
/** A field must be at least this large to be omitted as a last resort. */
112+
const OMITTABLE_FIELD_MIN_BYTES = 64 * 1024
113+
114+
/**
115+
* Replaces the payload's largest field with a note of its size: the last resort
116+
* for an event whose many short values no cut could bound, such as an object
117+
* with thousands of keys.
118+
*/
119+
function omitLargestField(payload: Record<string, unknown>, skipKeys: ReadonlySet<string>) {
120+
let largest: { key: string; bytes: number } | undefined
121+
for (const [key, field] of Object.entries(payload)) {
122+
if (skipKeys.has(key)) continue
123+
const bytes = Buffer.byteLength(JSON.stringify(field) ?? '', 'utf8')
124+
if (!largest || bytes > largest.bytes) largest = { key, bytes }
125+
}
126+
// Identity fields are never this large, so only bulk is ever replaced.
127+
if (!largest || largest.bytes <= OMITTABLE_FIELD_MIN_BYTES) return payload
128+
return { ...payload, [largest.key]: `…[omitted, ${formatFileSize(largest.bytes)} total]` }
129+
}
130+
103131
/**
104132
* Bounds an outgoing stream event so the replay buffer can persist it. Applied
105133
* only to the copy the writer delivers and persists; the caller keeps the full
106134
* event for dispatch. Long strings are cut to their head in place, so every
107-
* object keeps its shape; if that is not enough, long arrays keep their head.
135+
* object keeps its shape; if that is not enough, long arrays keep their head,
136+
* and past one replay write the largest field is replaced by a size note.
108137
* Assistant text, file previews, and the arguments of calls the browser
109138
* executes are never cut; an event
110139
* still too large is refused by the buffer, which ends the turn with an error.
@@ -126,5 +155,8 @@ export function compactStreamEvent(event: StreamEvent): StreamEvent {
126155
if (Buffer.byteLength(JSON.stringify(compacted)) > STREAM_EVENT_COMPACTION_THRESHOLD_BYTES) {
127156
compacted = trimArrays(compacted, skipKeys)
128157
}
158+
if (Buffer.byteLength(JSON.stringify(compacted)) > STREAM_EVENT_MAX_PAYLOAD_BYTES) {
159+
compacted = omitLargestField(toRecordOrNull(compacted) ?? {}, skipKeys)
160+
}
129161
return compacted === payload ? event : ({ ...event, payload: compacted } as StreamEvent)
130162
}

0 commit comments

Comments
 (0)