Skip to content

Commit 38d7552

Browse files
committed
fix(memory): preserve stream usage and bound portable history
1 parent bfc2c4e commit 38d7552

22 files changed

Lines changed: 567 additions & 38 deletions

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

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,29 @@ describe('optional Agent memory durability failures', () => {
126126
expect(mocks.append).toHaveBeenCalledOnce()
127127
})
128128

129+
it('propagates an append identity conflict instead of treating different content as saved', async () => {
130+
const failure = new OrchestrationError('conflict', 'Memory append identity was already used')
131+
mocks.append.mockRejectedValue(failure)
132+
await expect(
133+
new Memory().appendToMemory(
134+
ctx,
135+
inputs,
136+
{ role: 'user', content: 'changed question' },
137+
options
138+
)
139+
).rejects.toBe(failure)
140+
expect(mocks.append).toHaveBeenCalledExactlyOnceWith({
141+
principal: { kind: 'delegated' },
142+
input: {
143+
...options,
144+
workspaceId: ctx.workspaceId,
145+
conversationId: inputs.conversationId,
146+
data: { role: 'user', content: 'changed question' },
147+
provenance: undefined,
148+
},
149+
})
150+
})
151+
129152
it('preserves authorization failures on reads and appends', async () => {
130153
const failure = new OrchestrationError('forbidden', 'Denied')
131154
mocks.prefix.mockRejectedValue(failure)

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -88,7 +88,7 @@ describe('Memory', () => {
8888
}
8989
)
9090

91-
it.each(['forbidden', 'unauthorized', 'validation'] as const)(
91+
it.each(['forbidden', 'unauthorized', 'validation', 'conflict'] as const)(
9292
'propagates application %s failures even for optional rich history',
9393
async (code) => {
9494
const error = new OrchestrationError(code, 'Memory access refused')

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -95,7 +95,7 @@ function copyMemoryMessageMetadata(source: Message, target: Message): void {
9595
/** Optional durability tolerates storage-engine failures; application identity and projection failures propagate. */
9696
function isOptionalMemoryStorageFailure(error: unknown): boolean {
9797
if (error instanceof OrchestrationError)
98-
return error.code === 'not_found' || error.code === 'conflict' || error.code === 'internal'
98+
return error.code === 'not_found' || error.code === 'internal'
9999
const code = getPostgresErrorCode(error)
100100
return Boolean(
101101
(code &&

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

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

241+
it('preserves the shared materializer unavailable-result contract without attempting decryption', async () => {
242+
dbChainMockFns.limit.mockResolvedValueOnce([
243+
{
244+
key: ref.key,
245+
size: 500,
246+
workflowId: identity.workflowId,
247+
executionId: identity.executionId,
248+
},
249+
])
250+
materializeLargeValueRef.mockResolvedValueOnce(undefined)
251+
expect(await readMemoryArtifact({ ...scope, ref })).toBeUndefined()
252+
expect(encryptionMockFns.mockDecryptSecret).not.toHaveBeenCalled()
253+
})
254+
255+
it('does not swallow errors that escape the materialization boundary', async () => {
256+
dbChainMockFns.limit.mockResolvedValueOnce([
257+
{
258+
key: ref.key,
259+
size: 500,
260+
workflowId: identity.workflowId,
261+
executionId: identity.executionId,
262+
},
263+
])
264+
const error = new Error('Materialization access check failed')
265+
materializeLargeValueRef.mockRejectedValueOnce(error)
266+
await expect(readMemoryArtifact({ ...scope, ref })).rejects.toBe(error)
267+
expect(encryptionMockFns.mockDecryptSecret).not.toHaveBeenCalled()
268+
})
269+
241270
it.each(['ciphertext', 'json'])(
242271
'treats corrupt %s as an unavailable artifact',
243272
async (failure) => {

‎apps/sim/lib/memory/conversation-store.postgres.test.ts‎

Lines changed: 44 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -276,6 +276,50 @@ describe.skipIf(!databaseUrl)('conversation storage in Postgres', () => {
276276
expect(tail.provenance).toEqual(provenance)
277277
})
278278

279+
it('explicitly invalidates tracked legacy provenance after an untracked append', async () => {
280+
await writeLegacyPrefix()
281+
const untracked = { role: 'user', content: 'untracked append' }
282+
await appendMemoryMessages({
283+
workspaceId: identity.workspaceId,
284+
key: identity.conversationId,
285+
messages: [untracked],
286+
})
287+
const [stored] =
288+
await connection!`SELECT data, secret_provenance_version FROM memory WHERE id = 'legacy-memory'`
289+
expect(stored.secret_provenance_version).toBe(1)
290+
expect(stored.data).toEqual([...prefix, untracked])
291+
expect(
292+
(
293+
await connection!`SELECT content_hash, status, entries FROM memory_secret_provenance WHERE memory_id = 'legacy-memory'`
294+
)[0]
295+
).toEqual({
296+
content_hash: hashDurableSecretProvenanceValue(stored.data),
297+
status: 'unknown',
298+
entries: [],
299+
})
300+
const tracked = { role: 'assistant', content: 'later tracked append' }
301+
await appendMemoryMessages({
302+
workspaceId: identity.workspaceId,
303+
key: identity.conversationId,
304+
messages: [tracked],
305+
provenance,
306+
})
307+
expect(
308+
(
309+
await connection!`SELECT content_hash, status FROM memory_secret_provenance WHERE memory_id = 'legacy-memory'`
310+
)[0]
311+
).toEqual({
312+
content_hash: hashDurableSecretProvenanceValue([...prefix, untracked, tracked]),
313+
status: 'unknown',
314+
})
315+
await expect(
316+
new Memory().fetchMemoryMessages({ workspaceId: identity.workspaceId } as ExecutionContext, {
317+
memoryType: 'conversation',
318+
conversationId: identity.conversationId,
319+
})
320+
).rejects.toThrow('Memory content could not be safely projected')
321+
})
322+
279323
it('freezes the legacy prefix and keeps exchanges out of native/plain API history', async () => {
280324
await writeLegacyPrefix()
281325
const turn = await openAgentMemoryTurn(identity)

‎apps/sim/lib/memory/conversation-store.test.ts‎

Lines changed: 50 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,8 @@
11
/** @vitest-environment node */
2-
import { memory, memoryItem } from '@sim/db/schema'
2+
import { memory, memoryItem, memorySecretProvenance } from '@sim/db/schema'
33
import { dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing'
44
import { beforeEach, describe, expect, it, vi } from 'vitest'
5+
import { hashDurableSecretProvenanceValue } from '@/lib/execution/durable-secret-provenance'
56
import {
67
appendAgentMemoryMessage,
78
appendMemoryMessages,
@@ -65,3 +66,51 @@ describe('bounded conversation item writes', () => {
6566
expect(get).not.toHaveBeenCalled()
6667
})
6768
})
69+
70+
describe('untracked legacy appends', () => {
71+
beforeEach(() => {
72+
vi.clearAllMocks()
73+
resetDbChainMock()
74+
})
75+
76+
it('binds explicit unknown provenance to the updated JSON without declaring prior secrets public', async () => {
77+
const prefix = [{ role: 'user', content: 'tracked content' }]
78+
const message = { role: 'user', content: 'untracked append' }
79+
const data = [...prefix, message]
80+
queueTableRows(memory, [
81+
{ id: identity.memoryId, data: prefix, storageVersion: 1, secretProvenanceVersion: 1 },
82+
])
83+
dbChainMockFns.returning
84+
.mockResolvedValueOnce([{ id: identity.memoryId, data }])
85+
.mockResolvedValueOnce([{ id: identity.memoryId }])
86+
await writers.ordinary(message)
87+
expect(dbChainMockFns.insert).toHaveBeenCalledWith(memorySecretProvenance)
88+
expect(dbChainMockFns.values).toHaveBeenCalledWith(
89+
expect.objectContaining({
90+
memoryId: identity.memoryId,
91+
contentHash: hashDurableSecretProvenanceValue(data),
92+
status: 'unknown',
93+
entries: [],
94+
})
95+
)
96+
expect(dbChainMockFns.onConflictDoUpdate).toHaveBeenCalledWith(
97+
expect.objectContaining({
98+
set: expect.objectContaining({ secretProvenanceVersion: 1 }),
99+
})
100+
)
101+
})
102+
103+
it('keeps wholly untracked legacy conversations on their existing compatibility path', async () => {
104+
queueTableRows(memory, [
105+
{ id: identity.memoryId, data: [], storageVersion: 1, secretProvenanceVersion: null },
106+
])
107+
dbChainMockFns.returning.mockResolvedValueOnce([{ id: identity.memoryId, data: [] }])
108+
await writers.ordinary({ role: 'user', content: 'public' })
109+
expect(dbChainMockFns.insert).not.toHaveBeenCalledWith(memorySecretProvenance)
110+
expect(dbChainMockFns.onConflictDoUpdate).toHaveBeenCalledWith(
111+
expect.objectContaining({
112+
set: expect.objectContaining({ secretProvenanceVersion: null }),
113+
})
114+
)
115+
})
116+
})

‎apps/sim/lib/memory/conversation-store.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -208,6 +208,8 @@ export async function appendMemoryMessages(
208208
? 'merge-provenance-limit'
209209
: undefined
210210
)
211+
} else if (existing?.secretProvenanceVersion === 1) {
212+
await replaceMemorySecretProvenanceInTx(tx, written.id, written.data, { status: 'unknown' })
211213
}
212214
})
213215
}

‎apps/sim/lib/memory/execution-record.test.ts‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,19 @@ describe('portable execution records', () => {
8585
expect(JSON.parse(record.content!).notice).toBe('execution record shortened')
8686
})
8787

88+
it('bounds every default execution record without altering the canonical messages', () => {
89+
const messages: Message[] = Array.from({ length: 30 }, (_, index) => ({
90+
role: 'tool',
91+
tool_call_id: `call-${index}`,
92+
content: 'retained result '.repeat(1000),
93+
}))
94+
const original = structuredClone(messages)
95+
const record = renderConversationExecutionRecord(messages)
96+
expect(record.content!.length).toBeLessThanOrEqual(4096)
97+
expect(record.content).toContain('execution record shortened')
98+
expect(messages).toEqual(original)
99+
})
100+
88101
it('uses the same bounded format for protocol and context-size constraints', () => {
89102
const messages: Message[] = [{ role: 'tool', content: 'x'.repeat(10000) }]
90103
const record = renderConversationExecutionRecord(messages, 256)

‎apps/sim/lib/memory/execution-record.ts‎

Lines changed: 7 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,26 +1,25 @@
11
import { truncate } from '@sim/utils/string'
22
import type { Message } from '@/providers/types'
33

4-
/** A text-only representation for exchanges that cannot be replayed as protocol tool messages. */
4+
/** Bounded model context for exchanges that cannot be replayed as protocol tool messages. */
55
export function renderConversationExecutionRecord(
66
messages: readonly Message[],
7-
maxCharacters?: number
7+
maxCharacters = 4096
88
): Message {
9-
const bounded = maxCharacters !== undefined
109
let shortened = false
1110
const field = (value: string, limit: number): string => {
12-
if (!bounded || value.length <= limit) return value
11+
if (value.length <= limit) return value
1312
shortened = true
1413
return truncate(value, limit)
1514
}
1615
const take = <T>(values: readonly T[], limit: number): readonly T[] => {
17-
if (!bounded || values.length <= limit) return values
16+
if (values.length <= limit) return values
1817
shortened = true
1918
return values.slice(0, limit)
2019
}
2120
const recordedMessages = take(messages, 21).map((message) => ({
2221
role: message.role,
23-
content: bounded ? field(message.content ?? '', 512) : message.content,
22+
content: message.content === null ? null : field(message.content ?? '', 512),
2423
...(message.name ? { name: field(message.name, 64) } : {}),
2524
...(message.tool_call_id ? { callId: field(message.tool_call_id, 64) } : {}),
2625
...(message.function_call
@@ -47,11 +46,11 @@ export function renderConversationExecutionRecord(
4746
messages: recordedMessages,
4847
...(shortened ? { notice: 'execution record shortened' } : {}),
4948
})
50-
const suffix = bounded ? truncate('… [execution record shortened]', maxCharacters, '') : ''
49+
const suffix = truncate('… [execution record shortened]', maxCharacters, '')
5150
return {
5251
role: 'user',
5352
content:
54-
bounded && record.length > maxCharacters
53+
record.length > maxCharacters
5554
? truncate(record, Math.max(0, maxCharacters - suffix.length), suffix)
5655
: record,
5756
}

‎apps/sim/lib/memory/history-window.test.ts‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,36 @@ describe('conversation history windows', () => {
9696
expect(selectConversationTokenWindow(pending.flat(), 1, undefined, pending)).toEqual(groups[1])
9797
})
9898

99+
it.each(['tool_calls', 'function_call'] as const)(
100+
'counts legacy ungrouped %s arguments without changing plain-message token semantics',
101+
(field) => {
102+
const call = { name: 'lookup', arguments: JSON.stringify({ input: 'x'.repeat(1000) }) }
103+
const assistant: Message = {
104+
role: 'assistant',
105+
content: '',
106+
...(field === 'tool_calls'
107+
? { tool_calls: [{ id: 'call', type: 'function' as const, function: call }] }
108+
: { function_call: call }),
109+
}
110+
expect(selectConversationTokenWindow([assistant, ...final], 100)).toEqual(final)
111+
expect(selectConversationTokenWindow([...user, ...final], 6)).toEqual(final)
112+
}
113+
)
114+
115+
it('counts legacy function arguments inside explicitly grouped history too', () => {
116+
const group: Message[] = [
117+
{
118+
role: 'assistant',
119+
content: '',
120+
function_call: { name: 'lookup', arguments: JSON.stringify({ input: 'x'.repeat(1000) }) },
121+
},
122+
{ role: 'function', name: 'lookup', content: 'result' },
123+
]
124+
expect(
125+
selectConversationTokenWindow([...group, ...final], 200, undefined, [group, final])
126+
).toEqual(final)
127+
})
128+
99129
it('applies model context bounds to complete groups', () => {
100130
const groups = [user, exchange('batch'), final]
101131
expect(selectConversationContextWindow(groups.flat(), 'small', groups)).toEqual(final)

0 commit comments

Comments
 (0)