Skip to content

Commit bfc2c4e

Browse files
committed
fix(memory): preserve tool loops and harden durable history
1 parent 2a6faa6 commit bfc2c4e

62 files changed

Lines changed: 1624 additions & 229 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎apps/docs/content/docs/workflows/blocks/agent.mdx‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@ Tool calls and their results are selected together, including parallel calls. Th
6666

6767
Large results are retained separately, with their first 8,000 characters in model context and a notice when the result is truncated. The built-in `agent_memory_read` tool lets the agent search retained history and read omitted result details in small pages. It can only read the current conversation. When exact details matter, ask the agent to check the original result instead of relying on its preview.
6868

69-
Sim normally targets up to 16,000 estimated tokens of recalled history, subject to your memory window and the model's available context. Large or difficult-to-tokenize content uses a conservative estimate. It checks the input before every model generation, including generations after tool calls and on fallback models, leaving room for instructions, tool definitions, attachments, and output. The current request and required tool exchanges stay intact. If they cannot fit, the Agent stops with a context-limit error. This limits context per generation, not the total tokens used across a run.
69+
Sim normally targets up to 16,000 estimated tokens of recalled history, subject to your memory window and the model's available context. Large or difficult-to-tokenize content uses a conservative estimate. It checks the input before every model generation, including generations after tool calls and on fallback models, leaving room for instructions, tool definitions, attachments, and output. The current request and required tool exchanges stay intact. If their estimated size uses up the available budget, Sim omits optional history and still sends the current request; the provider enforces its actual context limit. These estimates guide recalled context per generation, not the total tokens used across a run.
7070

7171
When a generation would omit older history, Sim can create a concise summary while keeping recent exchanges, including during long tool loops. Summaries can omit details and do not replace the stored conversation. Generating one uses an additional model call whose tokens and cost are included in the Agent's usage; a cached summary can be reused when its source history is unchanged. If summarization is unavailable, the Agent continues with bounded history selection.
7272

‎apps/sim/executor/handlers/agent/agent-handler.test.ts‎

Lines changed: 39 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ import {
1919
vi,
2020
} from 'vitest'
2121
import { resetDeploymentShape } from '@/lib/core/config/deployment-shape'
22+
import type { AgentTurnSession } from '@/lib/memory/agent-turn-session'
2223
import type { AutoRoutingSignals } from '@/lib/model-router/resolve'
2324
import * as userFileBase64 from '@/lib/uploads/utils/user-file-base64.server'
2425
import { getAllBlocks } from '@/blocks'
@@ -6056,6 +6057,44 @@ describe('AgentBlockHandler', () => {
60566057
})
60576058

60586059
describe('wrapStreamForMemoryPersistence envelope', () => {
6060+
it.each(['Completed answer.', ''])(
6061+
'finalizes %j even when an existing stream callback rejects',
6062+
async (content) => {
6063+
const finalize = vi.fn().mockResolvedValue(undefined)
6064+
const onFullContent = vi.fn().mockRejectedValue(new Error('Callback failed'))
6065+
const stream: StreamingExecution = {
6066+
stream: new ReadableStream(),
6067+
onFullContent,
6068+
execution: {
6069+
success: true,
6070+
output: { content },
6071+
logs: [],
6072+
metadata: { startTime: '', duration: 0 },
6073+
},
6074+
}
6075+
const privateHandler = handler as unknown as {
6076+
wrapStreamForMemoryPersistence: (
6077+
ctx: ExecutionContext,
6078+
inputs: AgentInputs,
6079+
stream: StreamingExecution,
6080+
model: string,
6081+
session: AgentTurnSession
6082+
) => StreamingExecution
6083+
}
6084+
const wrapped = privateHandler.wrapStreamForMemoryPersistence(
6085+
mockContext,
6086+
{ model: 'gpt-4o' },
6087+
stream,
6088+
'gpt-4o',
6089+
{ memoryId: 'memory-1', finalize } as AgentTurnSession
6090+
)
6091+
6092+
await expect(wrapped.onFullContent?.(content)).resolves.toBeUndefined()
6093+
expect(onFullContent).toHaveBeenCalledWith(content)
6094+
expect(finalize).toHaveBeenCalledExactlyOnceWith(content, 'gpt-4o')
6095+
}
6096+
)
6097+
60596098
it('preserves streamFormat, subscribe, and the existing completion callback', async () => {
60606099
const handler = new AgentBlockHandler()
60616100
const subscribe = vi.fn()

‎apps/sim/executor/handlers/agent/agent-handler.ts‎

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3085,7 +3085,18 @@ export class AgentBlockHandler implements BlockHandler {
30853085
return {
30863086
...streamingExec,
30873087
onFullContent: async (content: string) => {
3088-
await streamingExec.onFullContent?.(content)
3088+
try {
3089+
await streamingExec.onFullContent?.(content)
3090+
} catch (error) {
3091+
logger.error(
3092+
'Streaming completion callback failed',
3093+
projectAgentDiagnosticMetadata(
3094+
ctx,
3095+
getErrorDiagnosticMetadata(error),
3096+
getErrorDiagnosticFallback(error)
3097+
)
3098+
)
3099+
}
30893100
if (!content.trim()) {
30903101
await agentConversation?.finalize('', servedModel)
30913102
return

‎apps/sim/executor/handlers/agent/memory.test.ts‎

Lines changed: 43 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
import { loggerMock, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing'
2-
import { beforeEach, describe, expect, it, vi } from 'vitest'
2+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
33

44
const { mockDecryptSecret, mockRedactObjectStrings } = vi.hoisted(() => ({
55
mockDecryptSecret: vi.fn(),
@@ -14,6 +14,7 @@ vi.mock('@/lib/logs/execution/pii-redaction', () => ({
1414
redactObjectStrings: mockRedactObjectStrings,
1515
}))
1616

17+
import { OrchestrationError } from '@/lib/core/orchestration/types'
1718
import { hashDurableSecretProvenanceValue } from '@/lib/execution/durable-secret-provenance'
1819
import { assertUserFileContentAccess } from '@/lib/execution/payloads/materialization.server'
1920
import { MEMORY } from '@/lib/memory/constants'
@@ -58,6 +59,47 @@ describe('Memory', () => {
5859
memoryService = new Memory()
5960
})
6061

62+
describe('optional durable storage', () => {
63+
const ctx = { workspaceId: 'workspace-1' } as ExecutionContext
64+
const inputs = { memoryType: 'conversation' as const, conversationId: 'conversation-1' }
65+
66+
afterEach(() => vi.restoreAllMocks())
67+
68+
function rejectRead(error: Error) {
69+
vi.spyOn(
70+
memoryService as unknown as { fetchMemory: () => Promise<unknown> },
71+
'fetchMemory'
72+
).mockRejectedValue(error)
73+
}
74+
75+
it.each(['ECONNREFUSED', '42P01', '23514'])(
76+
'degrades rich history on storage failure %s while preserving ordinary read errors',
77+
async (code) => {
78+
const error = Object.assign(new Error('Storage unavailable'), { code })
79+
rejectRead(error)
80+
await expect(
81+
memoryService.fetchMemoryMessages(ctx, inputs, undefined, { richHistory: true })
82+
).resolves.toEqual([])
83+
await expect(memoryService.fetchMemoryMessages(ctx, inputs)).rejects.toBe(error)
84+
expect(mockMemoryLogger.warn).toHaveBeenCalledWith(
85+
'Agent durable memory read is unavailable',
86+
{ workspaceId: 'workspace-1' }
87+
)
88+
}
89+
)
90+
91+
it.each(['forbidden', 'unauthorized', 'validation'] as const)(
92+
'propagates application %s failures even for optional rich history',
93+
async (code) => {
94+
const error = new OrchestrationError(code, 'Memory access refused')
95+
rejectRead(error)
96+
await expect(
97+
memoryService.fetchMemoryMessages(ctx, inputs, undefined, { richHistory: true })
98+
).rejects.toBe(error)
99+
}
100+
)
101+
})
102+
61103
describe('message window', () => {
62104
it('should keep last N messages', () => {
63105
const messages: Message[] = [

‎apps/sim/executor/handlers/agent/memory.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -92,7 +92,7 @@ function copyMemoryMessageMetadata(source: Message, target: Message): void {
9292
if (appendKey) messageAppendKeys.set(target, appendKey)
9393
}
9494

95-
/** Only storage availability/conflict failures degrade; identity and projection failures propagate. */
95+
/** Optional durability tolerates storage-engine failures; application identity and projection failures propagate. */
9696
function isOptionalMemoryStorageFailure(error: unknown): boolean {
9797
if (error instanceof OrchestrationError)
9898
return error.code === 'not_found' || error.code === 'conflict' || error.code === 'internal'

‎apps/sim/lib/memory/artifacts.test.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -238,6 +238,24 @@ describe('encrypted memory artifacts', () => {
238238
expect(materializeLargeValueRef).not.toHaveBeenCalled()
239239
})
240240

241+
it.each(['ciphertext', 'json'])(
242+
'treats corrupt %s as an unavailable artifact',
243+
async (failure) => {
244+
dbChainMockFns.limit.mockResolvedValueOnce([
245+
{
246+
key: ref.key,
247+
size: 500,
248+
workflowId: identity.workflowId,
249+
executionId: identity.executionId,
250+
},
251+
])
252+
if (failure === 'ciphertext')
253+
encryptionMockFns.mockDecryptSecret.mockRejectedValueOnce(new Error('invalid ciphertext'))
254+
else encryptionMockFns.mockDecryptSecret.mockResolvedValueOnce({ decrypted: 'invalid JSON' })
255+
expect(await readMemoryArtifact({ ...scope, ref })).toBeUndefined()
256+
}
257+
)
258+
241259
it('rejects an oversized decrypted result before parsing it', async () => {
242260
dbChainMockFns.limit.mockResolvedValueOnce([
243261
{

‎apps/sim/lib/memory/artifacts.ts‎

Lines changed: 10 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -179,10 +179,14 @@ export async function readMemoryArtifact(input: ReadMemoryArtifactInput): Promis
179179
) {
180180
return undefined
181181
}
182-
const { decrypted } = await decryptSecret(envelope.encrypted, { logFailure: false })
183-
if (Buffer.byteLength(decrypted, 'utf8') > MAX_MEMORY_ARTIFACT_BYTES) return undefined
184-
const value: unknown = JSON.parse(decrypted)
185-
return stringifyBoundedMemoryJson(value, MAX_MEMORY_ARTIFACT_BYTES) === undefined
186-
? undefined
187-
: value
182+
try {
183+
const { decrypted } = await decryptSecret(envelope.encrypted, { logFailure: false })
184+
if (Buffer.byteLength(decrypted, 'utf8') > MAX_MEMORY_ARTIFACT_BYTES) return undefined
185+
const value: unknown = JSON.parse(decrypted)
186+
return stringifyBoundedMemoryJson(value, MAX_MEMORY_ARTIFACT_BYTES) === undefined
187+
? undefined
188+
: value
189+
} catch {
190+
return undefined
191+
}
188192
}

‎apps/sim/lib/memory/bounded-json.test.ts‎

Lines changed: 37 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@ import { stringifyBoundedMemoryJson } from '@/lib/memory/bounded-json'
55
describe('bounded memory JSON', () => {
66
it.each([
77
{ value: { text: 'hello', values: [1, false, null] } },
8-
{ value: { text: 'é😀' } },
8+
{ value: { text: 'é😀\ud800\udc00\ud800' } },
9+
{ value: { absent: undefined, values: [undefined, Number.NaN] } },
910
{ value: { text: '\u0000\n"\\' } },
1011
{ value: Array.from({ length: 5000 }, () => 0) },
1112
])('uses the caller byte limit including UTF-8 and escaped JSON bytes', ({ value }) => {
@@ -20,7 +21,12 @@ describe('bounded memory JSON', () => {
2021
cyclic.self = cyclic
2122
let deep: unknown = 'leaf'
2223
for (let index = 0; index < 66; index++) deep = { child: deep }
23-
for (const value of [cyclic, deep, Array(100_001)]) {
24+
for (const value of [
25+
cyclic,
26+
deep,
27+
Array(100_001),
28+
Object.fromEntries(Array.from({ length: 100_001 }, (_, index) => [index, undefined])),
29+
]) {
2430
expect(stringifyBoundedMemoryJson(value, 8 * 1024 * 1024)).toBeUndefined()
2531
}
2632
})
@@ -48,6 +54,35 @@ describe('bounded memory JSON', () => {
4854
}
4955
})
5056

57+
it.each([
58+
{ value: { text: '\u0000'.repeat(200) } },
59+
{ value: { ['\u0000'.repeat(200)]: 'value' } },
60+
{ value: { text: '\ud800'.repeat(200) } },
61+
])('rejects escaped bytes before serializing the captured graph', ({ value }) => {
62+
const serialize = vi.spyOn(JSON, 'stringify')
63+
try {
64+
expect(stringifyBoundedMemoryJson(value, 1024)).toBeUndefined()
65+
expect(serialize).not.toHaveBeenCalled()
66+
} finally {
67+
serialize.mockRestore()
68+
}
69+
})
70+
71+
it('serializes the admitted descriptors without reading proxy values or toJSON', () => {
72+
const get = vi.fn(() => 'UNADMITTED')
73+
const value = new Proxy({ text: 'admitted' }, { get })
74+
expect(stringifyBoundedMemoryJson(value, 1024)).toBe('{"text":"admitted"}')
75+
expect(get).not.toHaveBeenCalled()
76+
})
77+
78+
it('does not read inherited numeric accessors in sparse arrays', () => {
79+
const get = vi.fn(() => 'UNADMITTED')
80+
const prototype = Object.create(Array.prototype, { 0: { get } })
81+
const value = Object.setPrototypeOf(Array(1), prototype)
82+
expect(stringifyBoundedMemoryJson(value, 1024)).toBe('[null]')
83+
expect(get).not.toHaveBeenCalled()
84+
})
85+
5186
it('allows repeated references without treating them as a cycle', () => {
5287
const result = { answer: 42 }
5388
const value = { rawResponse: result, modelResponse: result }

‎apps/sim/lib/memory/bounded-json.ts‎

Lines changed: 77 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -1,52 +1,91 @@
11
const MAX_MEMORY_JSON_NODES = 100_000
22
const MAX_MEMORY_JSON_DEPTH = 64
33

4-
/** Rejects oversized or unsafe plain JSON before allocating its serialized representation. */
4+
/** Counts JSON escapes without allocating the escaped string. */
5+
function quotedStringBytes(value: string, remaining: number): number | undefined {
6+
let bytes = 2
7+
for (let index = 0; index < value.length && bytes <= remaining; index++) {
8+
const code = value.charCodeAt(index)
9+
if (code === 0x22 || code === 0x5c) bytes += 2
10+
else if (code < 0x20) bytes += (code >= 8 && code <= 10) || code === 12 || code === 13 ? 2 : 6
11+
else if (code < 0x80) bytes++
12+
else if (code < 0x800) bytes += 2
13+
else if (code >= 0xd800 && code <= 0xdbff) {
14+
const next = value.charCodeAt(index + 1)
15+
if (next >= 0xdc00 && next <= 0xdfff) {
16+
bytes += 4
17+
index++
18+
} else bytes += 6
19+
} else bytes += code >= 0xdc00 && code <= 0xdfff ? 6 : 3
20+
}
21+
return bytes <= remaining ? bytes : undefined
22+
}
23+
24+
/** Captures bounded plain JSON once, without executing accessors or serializing the source graph. */
525
export function stringifyBoundedMemoryJson(value: unknown, maxBytes: number): string | undefined {
626
let nodes = 0
7-
let minimumBytes = 0
27+
let bytes = 0
28+
const invalid = Symbol('invalid JSON')
829
const ancestors = new WeakSet<object>()
9-
const visit = (item: unknown, depth: number): boolean => {
10-
if (++nodes > MAX_MEMORY_JSON_NODES || depth > MAX_MEMORY_JSON_DEPTH) return false
11-
if (typeof item === 'string') minimumBytes += Buffer.byteLength(item, 'utf8') + 2
12-
else if (item === null || item === undefined) minimumBytes += 4
13-
else if (typeof item === 'number')
14-
minimumBytes += Number.isFinite(item) ? String(item).length : 4
15-
else if (typeof item === 'boolean') minimumBytes += item ? 4 : 5
16-
else if (typeof item !== 'object') return false
17-
else {
18-
if (ancestors.has(item) || 'toJSON' in item) return false
19-
const prototype = Object.getPrototypeOf(item)
20-
if (!Array.isArray(item) && prototype !== Object.prototype && prototype !== null) return false
21-
ancestors.add(item)
22-
minimumBytes += 2
23-
if (Array.isArray(item)) {
24-
if (item.length > MAX_MEMORY_JSON_NODES - nodes) return false
25-
for (let index = 0; index < item.length; index++) {
26-
const field = Object.getOwnPropertyDescriptor(item, index)
27-
if (field && !('value' in field)) return false
28-
if (index > 0) minimumBytes++
29-
if (!visit(field?.value, depth + 1)) return false
30-
}
31-
} else {
32-
let fields = 0
33-
for (const key in item) {
34-
if (!Object.hasOwn(item, key)) continue
35-
const field = Object.getOwnPropertyDescriptor(item, key)
36-
if (!field || !('value' in field)) return false
37-
if (fields++ > 0) minimumBytes++
38-
minimumBytes += Buffer.byteLength(key, 'utf8') + 3
39-
if (!visit(field.value, depth + 1)) return false
30+
const addBytes = (count: number): boolean => {
31+
bytes += count
32+
return bytes <= maxBytes
33+
}
34+
const capture = (item: unknown, depth: number): unknown => {
35+
if (++nodes > MAX_MEMORY_JSON_NODES || depth > MAX_MEMORY_JSON_DEPTH) return invalid
36+
if (typeof item === 'string') {
37+
const count = quotedStringBytes(item, maxBytes - bytes)
38+
if (count === undefined || !addBytes(count)) return invalid
39+
return item
40+
}
41+
if (item === null || item === undefined) return addBytes(4) ? item : invalid
42+
if (typeof item === 'number')
43+
return addBytes(Number.isFinite(item) ? String(item).length : 4) ? item : invalid
44+
if (typeof item === 'boolean') return addBytes(item ? 4 : 5) ? item : invalid
45+
if (typeof item !== 'object' || ancestors.has(item) || 'toJSON' in item) return invalid
46+
const prototype = Object.getPrototypeOf(item)
47+
const isArray = Array.isArray(item)
48+
if (!isArray && prototype !== Object.prototype && prototype !== null) return invalid
49+
if (!addBytes(2)) return invalid
50+
ancestors.add(item)
51+
const snapshot: Record<string, unknown> | unknown[] = isArray
52+
? Object.setPrototypeOf([], null)
53+
: Object.create(null)
54+
if (isArray) {
55+
const length = Object.getOwnPropertyDescriptor(item, 'length')?.value
56+
if (typeof length !== 'number' || length > MAX_MEMORY_JSON_NODES - nodes) return invalid
57+
for (let index = 0; index < length; index++) {
58+
const field = Object.getOwnPropertyDescriptor(item, index)
59+
if (field && !('value' in field)) return invalid
60+
if (index > 0 && !addBytes(1)) return invalid
61+
const captured = capture(field?.value ?? null, depth + 1)
62+
if (captured === invalid) return invalid
63+
Object.defineProperty(snapshot, index, { value: captured, enumerable: true })
64+
}
65+
} else {
66+
let fields = 0
67+
for (const key in item) {
68+
const field = Object.getOwnPropertyDescriptor(item, key)
69+
if (!field || !field.enumerable) continue
70+
if (!('value' in field)) return invalid
71+
if (field.value === undefined) {
72+
if (++nodes > MAX_MEMORY_JSON_NODES) return invalid
73+
continue
4074
}
75+
const keyBytes = quotedStringBytes(key, maxBytes - bytes)
76+
if (keyBytes === undefined || !addBytes(keyBytes + 1 + (fields++ > 0 ? 1 : 0)))
77+
return invalid
78+
const captured = capture(field.value, depth + 1)
79+
if (captured === invalid) return invalid
80+
Object.defineProperty(snapshot, key, { value: captured, enumerable: true })
4181
}
42-
ancestors.delete(item)
4382
}
44-
return minimumBytes <= maxBytes
83+
ancestors.delete(item)
84+
return snapshot
4585
}
4686
try {
47-
if (!visit(value, 0)) return undefined
48-
const json = JSON.stringify(value)
49-
return json !== undefined && Buffer.byteLength(json, 'utf8') <= maxBytes ? json : undefined
87+
const snapshot = capture(value, 0)
88+
return snapshot === invalid ? undefined : JSON.stringify(snapshot)
5089
} catch {
5190
return undefined
5291
}

0 commit comments

Comments
 (0)