Skip to content

Commit 40cf7d5

Browse files
committed
fix(memory): validate checkpoint recovery and bounded summary coverage
1 parent 18557c1 commit 40cf7d5

15 files changed

Lines changed: 770 additions & 119 deletions

‎.github/workflows/test-build.yml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -232,6 +232,7 @@ jobs:
232232
lib/table/rows/secret-provenance.postgres.test.ts
233233
lib/memory/message-provenance.postgres.test.ts
234234
lib/memory/conversation-store.postgres.test.ts
235+
lib/memory/summary-store.postgres.test.ts
235236
executor/handlers/agent/memory-harness.postgres.test.ts
236237
237238
- name: Verify Search vector projection upgrade in PostgreSQL

‎apps/sim/lib/memory/README.md‎

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,7 +20,7 @@ A journal belongs to one execution, workflow, block, node, and execution-order i
2020

2121
Recorded terminal tool outcomes are reused. Calls with no recorded outcome may run again with their original Sim invocation ID, so continuation provides at-least-once execution, not exactly-once external effects. An external service must support the supplied idempotency key to deduplicate an uncertain effect. A new execution or loop iteration has a separate journal.
2222

23-
New durability failures degrade execution to its in-memory path. They must never cause a version-2 conversation to resume legacy-array writes. A missing result payload remains a terminal outcome, with an unavailable-detail notice and its recorded success status and cost. A missing or invalid step in a reference checkpoint, or oversized saved ciphertext, prevents continuation instead of starting the invocation again. Generic database availability failures retain the existing in-memory degradation behavior. API and ordinary memory behavior keep their existing error contract. Deleting a conversation takes the conversation lock and cascades item, journal, and artifact ownership rows. Active turns retain the original memory ID, so a later conversation with the same key cannot accept their stale checkpoint writes.
23+
New durability failures degrade execution to its in-memory path. They must never cause a version-2 conversation to resume legacy-array writes. A missing result payload remains a terminal outcome, with an unavailable-detail notice and its recorded success status and cost. Any nonempty saved checkpoint must decrypt, match its original identity and memory owner, and contain a valid state before continuation. Damaged ciphertext, invalid bindings or state, a missing step payload, and oversized saved ciphertext stop continuation instead of starting the invocation again. Database availability failures before a checkpoint is obtained retain the existing in-memory degradation behavior. API and ordinary memory behavior keep their existing error contract. Deleting a conversation takes the conversation lock and cascades item, journal, and artifact ownership rows. Active turns retain the original memory ID, so a later conversation with the same key cannot accept their stale checkpoint writes.
2424

2525
Large-result artifacts have conversation ownership separate from run-log retention. The cleanup predicate retains owned artifacts and their dependencies while the conversation remains active; deleting the conversation releases this ownership.
2626

@@ -42,15 +42,19 @@ The shared context policy uses an explicit memory token window when configured a
4242

4343
`agent_memory_read` is supplied only with the trusted original memory owner. It searches retained history, including a safely projectable legacy prefix, or reads a referenced result in bounded pages. Legacy-prefix data and provenance must fit the 1 MiB retrieval admission limit. Each call returns at most 6,000 UTF-8 text bytes and scans at most 10 history items; a continuation cursor may be returned even when no match appears in a page. Cursors bind to the owner and search. Artifact reads use opaque IDs and canonical conversation ownership, and return only safely projected model content. Journal envelopes, provider continuations, raw replay fields, and unrelated conversation artifacts cannot be retrieved through this tool. At most two retrievals run concurrently per execution context. Read pages sequentially and follow `nextCursor` when more detail is needed.
4444

45-
Before a generation omits older optional groups, a bound compactor can summarize both earlier conversation records and older completed exchanges from the active invocation. The current prompt and newest tool exchange remain required. A summary has priority within the same optional history budget as recent raw groups. Refreshes require additional history equal to at least half the configured history target, with a 1,024-token minimum, so every tool response does not trigger another summary call. A refreshed note can include the preceding note and newly eligible older records. The summary request has no tools, reads at most 8,000 estimated source tokens, and requests at most 1,024 output tokens, further constrained by the actual available summary budget. Its token/cost totals are recorded as `contextUsage`, without adding an execution step or a public conversation message. The latest summary is an encrypted derived cache on the original `memory` owner: at most 6,000 characters and 64 KiB of ciphertext, reused only for an exact versioned source hash. Cache replacement does not change the immutable transcript or its provenance, and generation runs outside database locks. Cache reads have a SQL byte guard and ordinary memory reads exclude the cache column. Failed summary generation or cache access falls back to bounded history selection.
45+
Before a generation omits older optional groups, a bound compactor can summarize earlier available conversation records and older completed exchanges from the active invocation. The current prompt and newest tool exchange remain required. A summary has priority within the same optional history budget as recent raw groups. Once the eligible backlog is covered, refreshes require additional history equal to at least half the configured history target, with a 1,024-token minimum, so every tool response does not trigger another summary call.
46+
47+
Compaction processes at most three source batches per pressure event, in chronological order. Each batch combines the preceding note with the next contiguous older records, admits at most 8,000 estimated source tokens, and requests at most 1,024 output tokens, further constrained by the actual available summary budget. An individually large group may contribute a bounded excerpt with its call identities and artifact IDs. The coverage cursor advances only through a successfully summarized prefix; a failed or unusable response cannot mark skipped records as covered. Remaining backlog can continue at the next pressure event. Summary calls have no tools, and their token/cost totals are recorded as `contextUsage` without adding an execution step or a public conversation message. Original records remain authoritative; a derived note may omit details and never authorizes tool replay.
48+
49+
The latest summary is an encrypted version-2 cache on the original `memory` owner, bounded to 6,000 characters and 64 KiB of ciphertext. Reuse requires validated cache metadata: its original-message count and source hash must exactly match an eligible canonical history prefix under the current summary policy. The source hash follows original records, independent of intermediate generated wording. Cache replacement does not change the immutable transcript or its provenance, and generation runs outside database locks. Cache reads have a SQL byte guard and ordinary memory reads exclude the cache column. Failed summary generation or cache access preserves any earlier usable note and falls back to bounded history selection.
4650

4751
The public [Agent block documentation](../../../docs/content/docs/workflows/blocks/agent.mdx) describes these user-visible limits.
4852

4953
## Verification
5054

5155
The optional five-family live contract suite is `providers/conversation-smoke.test.ts`. It is skipped unless `RUN_AGENT_MEMORY_PROVIDER_SMOKE=true` and `AGENT_MEMORY_PROVIDER_SMOKE_CASES` supplies an array of `{protocol, providerId, model, apiKey?, azureEndpoint?, azureApiVersion?, bedrockAccessKeyId?, bedrockSecretKey?, bedrockRegion?}` entries. Supply one entry for each protocol from `history-adapters.ts`. It executes a mocked, side-effect-free echo tool and then sends the captured native history through a second live request. These calls use provider credits; credentials remain in the environment and must not be committed. The implementation verification did not enable this paid suite.
5256

53-
`conversation-store.postgres.test.ts` creates a disposable schema in a local database selected by `MEMORY_PROVENANCE_TEST_DATABASE_URL`. It applies migration 0368 and checks legacy-prefix preservation, public projections, CAS conflicts and rollback, append deduplication, pair-safe history windows, deletion/recreation, artifact retention, and oversized-checkpoint admission. `summary-store.postgres.test.ts` verifies exact hash reuse, scope binding, concurrent replacement/deletion, encrypted cache bounds, and exclusion from public projections and SQL reads. That suite requires the isolated local database at `127.0.0.1:5433/sim_durable_memory_e2e`. `message-provenance.postgres.test.ts` also applies this migration and verifies legacy provenance through a 17,000-message conversation and concurrent native/tool appends. All three suites drop their disposable schemas afterward.
57+
`conversation-store.postgres.test.ts` creates a disposable schema in a local database selected by `MEMORY_PROVENANCE_TEST_DATABASE_URL`. It applies migration 0368 and checks legacy-prefix preservation, public projections, CAS conflicts and rollback, append deduplication, pair-safe history windows, deletion/recreation, artifact retention, and oversized-checkpoint admission. `summary-store.postgres.test.ts` verifies exact hash reuse, scope binding, concurrent replacement/deletion, encrypted cache bounds, and exclusion from public projections and SQL reads. It uses the same explicitly configured local database and a disposable schema. `message-provenance.postgres.test.ts` also applies this migration and verifies legacy provenance through a 17,000-message conversation and concurrent native/tool appends. All three suites drop their disposable schemas afterward.
5458

5559
The journal/session fixtures cover terminal sibling replay, missing artifacts, legacy snapshot migration, byte signatures, provenance, and context usage. The growth regression records 50 steps/results and checks that each payload is uploaded once while the encrypted manifest stays below 120 KiB. Run the focused suites from `apps/sim`:
5660

‎apps/sim/lib/memory/agent-turn-session.test.ts‎

Lines changed: 40 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -426,7 +426,12 @@ describe('durable Agent session', () => {
426426
)
427427

428428
expect(live.rawResponse).toBe(response)
429-
expect(live.modelResponse.output.text).toBe(response.output.text)
429+
expect(live.modelResponse).toMatchObject({
430+
success,
431+
output: { memoryResultUnavailable: true },
432+
...(!success ? { error: 'Upstream rejected the operation' } : {}),
433+
})
434+
expect(JSON.stringify(live.modelResponse).length).toBeLessThan(2000)
430435
expect(executeTool).toHaveBeenCalledTimes(1)
431436
expect(session!.getPendingCalls()).toEqual([])
432437
const recorded = session!.getRecordedResult(invocationId)
@@ -467,20 +472,40 @@ describe('durable Agent session', () => {
467472
expect(session!.getRecordedResult(invocationId)?.modelResponse).toEqual(response)
468473
})
469474

470-
it('does not attach another invocation or recreated conversation journal', async () => {
471-
const first = await openAgentTurnSession(input())
472-
await first!.captureStep(step())
473-
const encryptedState = save.mock.calls.at(-1)![0].input.encryptedState
474-
open.mockResolvedValue({
475-
memoryId: 'replacement-memory',
476-
turnId: 'replacement-turn',
477-
revision: 2,
478-
encryptedState,
479-
})
480-
const replacement = await openAgentTurnSession(input(2))
481-
expect(replacement!.getPendingCalls()).toEqual([])
482-
expect(replacement!.getMessages('openai', 'model-a', 'binding-a')).toEqual([])
483-
})
475+
it.each(['damaged ciphertext', 'invocation binding', 'memory binding', 'invalid state'])(
476+
'refuses an empty fresh session when a saved checkpoint has %s',
477+
async (failure) => {
478+
const first = (await openAgentTurnSession(input()))!
479+
await first.captureStep(step())
480+
const response = { success: true, output: { delivered: true } }
481+
await first.recordToolResult({
482+
invocationId: first.getPendingCalls()[0].invocationId,
483+
rawResponse: response,
484+
modelResponse: response,
485+
})
486+
let encryptedState: string = save.mock.calls.at(-1)![0].input.encryptedState
487+
if (failure === 'damaged ciphertext') {
488+
const last = encryptedState.at(-1) === '0' ? '1' : '0'
489+
encryptedState = `${encryptedState.slice(0, -1)}${last}`
490+
} else if (failure === 'invalid state') {
491+
const envelope = await decryptMemoryCheckpoint(encryptedState)
492+
if (!isRecordLike(envelope)) throw new Error('Expected checkpoint envelope')
493+
encryptedState = await encryptMemoryCheckpoint({ ...envelope, state: { version: 99 } })
494+
}
495+
open.mockResolvedValue({
496+
memoryId: failure === 'memory binding' ? 'replacement-memory' : 'memory-1',
497+
turnId: 'turn-1',
498+
revision: 2,
499+
encryptedState,
500+
})
501+
const retry = input(failure === 'invocation binding' ? 2 : 1)
502+
await expect(openAgentTurnSession(retry)).rejects.toMatchObject({ retryable: false })
503+
await expect(openAgentTurnSession(retry)).rejects.toMatchObject({ retryable: false })
504+
expect(save).toHaveBeenCalledTimes(2)
505+
expect(readArtifact).not.toHaveBeenCalled()
506+
expect(executeTool).not.toHaveBeenCalled()
507+
}
508+
)
484509

485510
it('writes payloads once while a long invocation grows beyond the old snapshot byte limit', async () => {
486511
const session = (await openAgentTurnSession(input()))!

‎apps/sim/lib/memory/agent-turn-session.ts‎

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -412,7 +412,6 @@ export async function openAgentTurnSession(
412412
let record: AgentMemoryTurnRecord | undefined
413413
let state: AgentTurnState | undefined
414414
let journal: AgentTurnJournal | undefined
415-
let restoringJournal = false
416415
let degraded = false
417416
const degrade = () => {
418417
if (!degraded)
@@ -463,15 +462,15 @@ export async function openAgentTurnSession(
463462
restored.memoryId !== record.memoryId
464463
)
465464
throw new Error('Invalid Agent checkpoint binding')
466-
restoringJournal = isRecordLike(restored.state) && restored.state.version === 2
465+
const restoringJournal = isRecordLike(restored.state) && restored.state.version === 2
467466
const restoredState = restoringJournal
468467
? await journal.restore(restored.state)
469468
: restored.state
470469
if (!validState(restoredState)) throw new Error('Invalid Agent checkpoint state')
471470
state = restoredState
472471
}
473472
} catch (error) {
474-
if (restoringJournal || (isRecordLike(error) && error.code === 'payload_too_large'))
473+
if (record?.encryptedState || (isRecordLike(error) && error.code === 'payload_too_large'))
475474
throw Object.assign(new Error('Agent invocation journal could not be safely restored'), {
476475
retryable: false,
477476
})

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

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,10 @@ const database = vi.hoisted(() => ({ current: undefined as PostgresJsDatabase |
1111
vi.unmock('drizzle-orm')
1212
vi.unmock('@sim/db/schema')
1313
vi.mock('@sim/db', () => ({
14+
dbFor: () => {
15+
if (!database.current) throw new Error('Postgres test is not initialized')
16+
return database.current
17+
},
1418
db: {
1519
select: (...args: unknown[]) => {
1620
if (!database.current) throw new Error('Postgres test is not initialized')
@@ -63,6 +67,7 @@ import {
6367
readPlainMemoryTail,
6468
saveAgentMemoryTurn,
6569
} from '@/lib/memory/conversation-store'
70+
import { retrieveMemory } from '@/lib/memory/retrieval'
6671
import {
6772
getMemoryMessageAppendKey,
6873
getMemoryMessageTurnId,
@@ -94,6 +99,7 @@ const identity: AgentMemoryTurnIdentity = {
9499
}
95100
const provenance = { status: 'exact', entries: [] } as const
96101
const journalReads: string[] = []
102+
const deduplicationReads: string[] = []
97103
const prefix = [
98104
{ role: 'user', content: 'legacy question' },
99105
{ role: 'assistant', content: 'legacy answer' },
@@ -156,6 +162,12 @@ describe.skipIf(!databaseUrl)('conversation storage in Postgres', () => {
156162
logQuery(query) {
157163
if (query.startsWith('select ') && query.includes('agent_memory_turn'))
158164
journalReads.push(query)
165+
if (
166+
query.startsWith('select ') &&
167+
query.includes('from "memory_item"') &&
168+
query.includes('"memory_item"."append_key" in')
169+
)
170+
deduplicationReads.push(query)
159171
},
160172
},
161173
})
@@ -348,6 +360,7 @@ describe.skipIf(!databaseUrl)('conversation storage in Postgres', () => {
348360
})
349361

350362
it('deduplicates committed history and rolls back journal advancement on conflicting content', async () => {
363+
deduplicationReads.length = 0
351364
const turn = await openAgentMemoryTurn(identity)
352365
const item = {
353366
appendKey: 'stable-exchange',
@@ -386,6 +399,10 @@ describe.skipIf(!databaseUrl)('conversation storage in Postgres', () => {
386399
(await readConversationItems({ workspaceId: identity.workspaceId, memoryId: turn.memoryId }))
387400
.items
388401
).toHaveLength(1)
402+
expect(deduplicationReads).toHaveLength(3)
403+
for (const query of deduplicationReads) {
404+
expect(query.split(' from ')[0]).toBe('select "append_key", "content_hash", "kind"')
405+
}
389406
})
390407

391408
it('deletion cascades history and forbids stale writers from resurrecting a recreated key', async () => {
@@ -698,4 +715,62 @@ describe.skipIf(!databaseUrl)('conversation storage in Postgres', () => {
698715
await connection!`SELECT id FROM memory_item WHERE memory_id = ${turn.memoryId}`
699716
).toEqual([{ id: 'existing-large-item' }])
700717
})
718+
719+
it('finds older retained matches after a no-match page reaches the retrieval byte limit', async () => {
720+
const turn = await openAgentMemoryTurn(identity)
721+
await saveAgentMemoryTurn({
722+
...identity,
723+
...turn,
724+
expectedRevision: 0,
725+
encryptedState: 'saved',
726+
items: [
727+
{
728+
appendKey: 'older-match',
729+
kind: 'message',
730+
data: { role: 'user', content: 'older retained receipt needle' },
731+
provenance,
732+
},
733+
...Array.from({ length: 5 }, (_, index) => ({
734+
appendKey: `large-unrelated-${index}`,
735+
kind: 'message' as const,
736+
data: { role: 'assistant', content: 'x'.repeat(1024 * 1024 - 1024) },
737+
provenance,
738+
})),
739+
],
740+
})
741+
const scope = { workspaceId: identity.workspaceId, memoryId: turn.memoryId }
742+
const contextPage = await readConversationItems({ ...scope, limit: 10 })
743+
expect(contextPage.items).toHaveLength(4)
744+
expect(contextPage.nextBeforeSequence).toBeUndefined()
745+
746+
const args = { target: 'history' as const, query: 'receipt needle' }
747+
const first = await retrieveMemory({ ...scope, arguments: args, projection: {} })
748+
expect(first).toMatchObject({ text: '', scannedItems: 4, nextCursor: expect.any(String) })
749+
const next = await retrieveMemory({
750+
...scope,
751+
arguments: { ...args, cursor: first.nextCursor },
752+
projection: {},
753+
})
754+
expect(next.text).toContain('older retained receipt needle')
755+
expect(next.scannedItems).toBe(2)
756+
})
757+
758+
it('reports a single oversized history item and advances to the end without repeating it', async () => {
759+
const turn = await openAgentMemoryTurn(identity)
760+
const oversized = { role: 'user', content: 'x'.repeat(5 * 1024 * 1024) }
761+
await connection!`INSERT INTO memory_item (id, memory_id, append_key, kind, data, content_hash, provenance_status, provenance_entries) VALUES ('oversized-retrieval-item', ${turn.memoryId}, 'oversized-retrieval', 'message', ${JSON.stringify(oversized)}::jsonb, ${hashDurableSecretProvenanceValue(oversized)}, 'exact', '[]')`
762+
const scope = { workspaceId: identity.workspaceId, memoryId: turn.memoryId }
763+
const args = { target: 'history' as const, query: 'needle' }
764+
const first = await retrieveMemory({ ...scope, arguments: args, projection: {} })
765+
expect(first.text).toBe('')
766+
expect(first.notice).toContain('not retrievable within the safe 4 MiB')
767+
expect(first.nextCursor).toEqual(expect.any(String))
768+
const next = await retrieveMemory({
769+
...scope,
770+
arguments: { ...args, cursor: first.nextCursor },
771+
projection: {},
772+
})
773+
expect(next.nextCursor).toBeUndefined()
774+
expect(next.scannedItems).toBe(0)
775+
})
701776
})

0 commit comments

Comments
 (0)