Skip to content

Commit b3588ad

Browse files
committed
fix(tables): fire triggers for virtual table mutations
1 parent d2f0ca3 commit b3588ad

9 files changed

Lines changed: 333 additions & 7 deletions

File tree

apps/sim/app/api/memory/[id]/route.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import { parseRequest } from '@/lib/api/server'
1313
import { checkInternalAuth } from '@/lib/auth/hybrid'
1414
import { generateRequestId } from '@/lib/core/utils/request'
1515
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
16+
import { fireMemoryTableTrigger } from '@/lib/virtual-tables/memory-virtual-table.server'
1617
import { checkWorkspaceAccess } from '@/lib/workspaces/permissions/utils'
1718

1819
const logger = createLogger('MemoryByIdAPI')
@@ -241,6 +242,7 @@ export const PUT = withRouteHandler(async (request: NextRequest, context: Memory
241242
.limit(1)
242243

243244
const mem = updatedMemories[0]
245+
void fireMemoryTableTrigger(mem, existingMemories[0], requestId)
244246

245247
logger.info(`[${requestId}] Memory updated: ${id} for workspace: ${validatedWorkspaceId}`)
246248
return NextResponse.json(

apps/sim/app/api/memory/route.ts

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import { parseRequest } from '@/lib/api/server'
1515
import { checkInternalAuth } from '@/lib/auth/hybrid'
1616
import { generateRequestId } from '@/lib/core/utils/request'
1717
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
18+
import { fireMemoryTableTrigger } from '@/lib/virtual-tables/memory-virtual-table.server'
1819
import { checkWorkspaceAccess } from '@/lib/workspaces/permissions/utils'
1920

2021
const logger = createLogger('MemoryAPI')
@@ -181,6 +182,12 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
181182
const now = new Date()
182183
const id = `mem_${generateId().replace(/-/g, '')}`
183184

185+
const previousMemories = await db
186+
.select()
187+
.from(memory)
188+
.where(and(eq(memory.key, key), eq(memory.workspaceId, workspaceId)))
189+
.limit(1)
190+
184191
const { sql } = await import('drizzle-orm')
185192

186193
await db
@@ -219,6 +226,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
219226
}
220227

221228
const memoryRecord = allMemories[0]
229+
void fireMemoryTableTrigger(memoryRecord, previousMemories[0] ?? null, requestId)
222230

223231
return NextResponse.json(
224232
{ success: true, data: { conversationId: memoryRecord.key, data: memoryRecord.data } },

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

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,17 @@
1+
import { dbChainMockFns, resetDbChainMock } from '@sim/testing'
12
import { beforeEach, describe, expect, it, vi } from 'vitest'
23
import { MEMORY } from '@/executor/constants'
34
import { Memory } from '@/executor/handlers/agent/memory'
45
import type { Message } from '@/executor/handlers/agent/types'
6+
import type { ExecutionContext } from '@/executor/types'
7+
8+
const { mockFireMemoryTableTrigger } = vi.hoisted(() => ({
9+
mockFireMemoryTableTrigger: vi.fn(),
10+
}))
11+
12+
vi.mock('@/lib/virtual-tables/memory-virtual-table.server', () => ({
13+
fireMemoryTableTrigger: mockFireMemoryTableTrigger,
14+
}))
515

616
vi.mock('@/lib/tokenization/estimators', () => ({
717
getAccurateTokenCount: vi.fn((text: string) => {
@@ -13,9 +23,46 @@ describe('Memory', () => {
1323
let memoryService: Memory
1424

1525
beforeEach(() => {
26+
vi.clearAllMocks()
27+
resetDbChainMock()
1628
memoryService = new Memory()
1729
})
1830

31+
it('fires an update trigger when appending to an existing virtual-table row', async () => {
32+
const previousRecord = {
33+
id: 'memory-1',
34+
workspaceId: 'workspace-1',
35+
key: 'conversation-1',
36+
data: [{ role: 'user', content: 'Hello' }],
37+
createdAt: new Date('2026-01-01T00:00:00.000Z'),
38+
updatedAt: new Date('2026-01-01T00:00:00.000Z'),
39+
deletedAt: null,
40+
}
41+
const updatedRecord = {
42+
...previousRecord,
43+
data: [...previousRecord.data, { role: 'assistant', content: 'Hi' }],
44+
updatedAt: new Date('2026-01-02T00:00:00.000Z'),
45+
}
46+
dbChainMockFns.limit.mockResolvedValueOnce([previousRecord])
47+
dbChainMockFns.returning.mockResolvedValueOnce([updatedRecord])
48+
49+
await memoryService.appendToMemory(
50+
{
51+
workspaceId: 'workspace-1',
52+
executionId: 'execution-1',
53+
metadata: { requestId: 'request-1' },
54+
} as ExecutionContext,
55+
{ memoryType: 'conversation', conversationId: 'conversation-1' },
56+
{ role: 'assistant', content: 'Hi' }
57+
)
58+
59+
expect(mockFireMemoryTableTrigger).toHaveBeenCalledWith(
60+
updatedRecord,
61+
previousRecord,
62+
'request-1'
63+
)
64+
})
65+
1966
describe('applyWindow (message-based)', () => {
2067
it('should keep last N messages', () => {
2168
const messages: Message[] = [

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

Lines changed: 41 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { generateId } from '@sim/utils/id'
55
import { and, eq, sql } from 'drizzle-orm'
66
import { redactObjectStrings } from '@/lib/logs/execution/pii-redaction'
77
import { getAccurateTokenCount } from '@/lib/tokenization/estimators'
8+
import { fireMemoryTableTrigger } from '@/lib/virtual-tables/memory-virtual-table.server'
89
import { MEMORY } from '@/executor/constants'
910
import type { AgentInputs, Message } from '@/executor/handlers/agent/types'
1011
import type { ExecutionContext } from '@/executor/types'
@@ -66,7 +67,12 @@ export class Memory {
6667

6768
const key = inputs.conversationId!
6869

69-
await this.appendMessage(workspaceId, key, message)
70+
await this.appendMessage(
71+
workspaceId,
72+
key,
73+
message,
74+
ctx.metadata.requestId ?? ctx.executionId ?? ctx.workflowId
75+
)
7076

7177
logger.debug('Appended message to memory', {
7278
workspaceId,
@@ -110,7 +116,12 @@ export class Memory {
110116
messagesToStore.map((message) => this.maskContentForStorage(ctx, message))
111117
)
112118

113-
await this.seedMemoryRecord(workspaceId, key, messagesToStore)
119+
await this.seedMemoryRecord(
120+
workspaceId,
121+
key,
122+
messagesToStore,
123+
ctx.metadata.requestId ?? ctx.executionId ?? ctx.workflowId
124+
)
114125

115126
logger.debug('Seeded memory', {
116127
workspaceId,
@@ -227,13 +238,14 @@ export class Memory {
227238
private async seedMemoryRecord(
228239
workspaceId: string,
229240
key: string,
230-
messages: Message[]
241+
messages: Message[],
242+
requestId: string
231243
): Promise<void> {
232244
const now = new Date()
233245

234246
const sanitizedMessages = messages.map((message) => this.sanitizeMessageForStorage(message))
235247

236-
await db
248+
const insertedRecords = await db
237249
.insert(memory)
238250
.values({
239251
id: generateId(),
@@ -244,14 +256,31 @@ export class Memory {
244256
updatedAt: now,
245257
})
246258
.onConflictDoNothing()
259+
.returning()
260+
261+
const insertedRecord = insertedRecords[0]
262+
if (insertedRecord) {
263+
void fireMemoryTableTrigger(insertedRecord, null, requestId)
264+
}
247265
}
248266

249-
private async appendMessage(workspaceId: string, key: string, message: Message): Promise<void> {
267+
private async appendMessage(
268+
workspaceId: string,
269+
key: string,
270+
message: Message,
271+
requestId: string
272+
): Promise<void> {
250273
const now = new Date()
251274

252275
const sanitizedMessage = this.sanitizeMessageForStorage(message)
253276

254-
await db
277+
const previousRecords = await db
278+
.select()
279+
.from(memory)
280+
.where(and(eq(memory.workspaceId, workspaceId), eq(memory.key, key)))
281+
.limit(1)
282+
283+
const updatedRecords = await db
255284
.insert(memory)
256285
.values({
257286
id: generateId(),
@@ -268,6 +297,12 @@ export class Memory {
268297
updatedAt: now,
269298
},
270299
})
300+
.returning()
301+
302+
const updatedRecord = updatedRecords[0]
303+
if (updatedRecord) {
304+
void fireMemoryTableTrigger(updatedRecord, previousRecords[0] ?? null, requestId)
305+
}
271306
}
272307

273308
private parsePositiveInt(value: string | undefined, defaultValue: number): number {
Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { beforeEach, describe, expect, it, vi } from 'vitest'
5+
6+
const { mockFireTableTrigger } = vi.hoisted(() => ({
7+
mockFireTableTrigger: vi.fn(),
8+
}))
9+
10+
vi.mock('@/lib/table/trigger', () => ({
11+
fireTableTrigger: mockFireTableTrigger,
12+
}))
13+
14+
import { fireMemoryTableTrigger } from '@/lib/virtual-tables/memory-virtual-table.server'
15+
16+
describe('fireMemoryTableTrigger', () => {
17+
beforeEach(() => {
18+
vi.clearAllMocks()
19+
mockFireTableTrigger.mockResolvedValue(undefined)
20+
})
21+
22+
it('maps an updated memory record to a virtual-table update event', async () => {
23+
const previousRecord = {
24+
id: 'memory-1',
25+
workspaceId: 'workspace-1',
26+
key: 'conversation-1',
27+
data: [{ role: 'user', content: 'Hello' }],
28+
createdAt: new Date('2026-01-01T00:00:00.000Z'),
29+
updatedAt: new Date('2026-01-01T00:00:00.000Z'),
30+
deletedAt: null,
31+
}
32+
const updatedRecord = {
33+
...previousRecord,
34+
data: [...previousRecord.data, { role: 'assistant', content: 'Hi' }],
35+
updatedAt: new Date('2026-01-02T00:00:00.000Z'),
36+
}
37+
38+
await fireMemoryTableTrigger(updatedRecord, previousRecord, 'request-1')
39+
40+
expect(mockFireTableTrigger).toHaveBeenCalledWith(
41+
'system_memory_workspace-1',
42+
'Memory',
43+
'update',
44+
[
45+
expect.objectContaining({
46+
id: 'memory-1',
47+
data: expect.objectContaining({
48+
conversation_id: 'conversation-1',
49+
message_count: 2,
50+
updated_at: '2026-01-02T00:00:00.000Z',
51+
}),
52+
}),
53+
],
54+
new Map([
55+
[
56+
'memory-1',
57+
expect.objectContaining({
58+
conversation_id: 'conversation-1',
59+
message_count: 1,
60+
updated_at: '2026-01-01T00:00:00.000Z',
61+
}),
62+
],
63+
]),
64+
expect.objectContaining({ columns: expect.any(Array) }),
65+
'request-1'
66+
)
67+
})
68+
69+
it('maps a newly inserted memory record to a virtual-table insert event', async () => {
70+
const insertedRecord = {
71+
id: 'memory-1',
72+
workspaceId: 'workspace-1',
73+
key: 'conversation-1',
74+
data: [{ role: 'user', content: 'Hello' }],
75+
createdAt: new Date('2026-01-01T00:00:00.000Z'),
76+
updatedAt: new Date('2026-01-01T00:00:00.000Z'),
77+
deletedAt: null,
78+
}
79+
80+
await fireMemoryTableTrigger(insertedRecord, null, 'request-1')
81+
82+
expect(mockFireTableTrigger).toHaveBeenCalledWith(
83+
'system_memory_workspace-1',
84+
'Memory',
85+
'insert',
86+
[expect.objectContaining({ id: 'memory-1' })],
87+
null,
88+
expect.any(Object),
89+
'request-1'
90+
)
91+
})
92+
93+
it('does not emit table events for deleted memory records hidden from the virtual table', async () => {
94+
await fireMemoryTableTrigger(
95+
{
96+
id: 'memory-1',
97+
workspaceId: 'workspace-1',
98+
key: 'conversation-1',
99+
data: [],
100+
createdAt: new Date('2026-01-01T00:00:00.000Z'),
101+
updatedAt: new Date('2026-01-01T00:00:00.000Z'),
102+
deletedAt: new Date('2026-01-02T00:00:00.000Z'),
103+
},
104+
null,
105+
'request-1'
106+
)
107+
108+
expect(mockFireTableTrigger).not.toHaveBeenCalled()
109+
})
110+
})

apps/sim/lib/virtual-tables/memory-virtual-table.server.ts

Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,11 @@ import {
1010
getMemoryTableId,
1111
isMemoryTableId,
1212
MEMORY_TABLE_COLUMNS,
13+
MEMORY_TABLE_NAME,
1314
MEMORY_TABLE_SCHEMA,
1415
mapMemoryRecordToTableRow,
1516
} from '@/lib/virtual-tables/memory-virtual-table'
17+
import { fireVirtualTableTrigger } from '@/lib/virtual-tables/virtual-table-trigger.server'
1618

1719
interface QueryMemoryTableRowsInput extends QueryOptions {
1820
workspaceId: string
@@ -21,6 +23,11 @@ interface QueryMemoryTableRowsInput extends QueryOptions {
2123
const MEMORY_TABLE_ID_PREFIX = getMemoryTableId('')
2224
const MEMORY_ROWS_SQL_NAME = 'memory_rows'
2325

26+
type MemoryTableTriggerRecord = Pick<
27+
typeof memory.$inferSelect,
28+
'id' | 'workspaceId' | 'key' | 'data' | 'createdAt' | 'updatedAt' | 'deletedAt'
29+
>
30+
2431
function referencesMemoryTranscript(value: unknown): boolean {
2532
if (Array.isArray(value)) return value.some(referencesMemoryTranscript)
2633
if (typeof value !== 'object' || value === null) return false
@@ -199,3 +206,34 @@ export async function queryMemoryTableRows({
199206
keysetValid: !sort,
200207
}
201208
}
209+
export async function fireMemoryTableTrigger(
210+
currentRecord: MemoryTableTriggerRecord,
211+
previousRecord: MemoryTableTriggerRecord | null,
212+
requestId: string
213+
): Promise<void> {
214+
if (currentRecord.deletedAt) return
215+
216+
const toTableRow = (record: MemoryTableTriggerRecord) =>
217+
mapMemoryRecordToTableRow({
218+
id: record.id,
219+
key: record.key,
220+
data: record.data as JsonValue,
221+
messageCount: Array.isArray(record.data) ? record.data.length : 0,
222+
createdAt: record.createdAt,
223+
updatedAt: record.updatedAt,
224+
})
225+
const currentRow = toTableRow(currentRecord)
226+
const previousRow = previousRecord ? toTableRow(previousRecord) : null
227+
228+
await fireVirtualTableTrigger({
229+
table: {
230+
id: getMemoryTableId(currentRecord.workspaceId),
231+
name: MEMORY_TABLE_NAME,
232+
schema: MEMORY_TABLE_SCHEMA,
233+
},
234+
eventType: previousRecord ? 'update' : 'insert',
235+
rows: [currentRow],
236+
previousRows: previousRow ? new Map([[currentRow.id, previousRow.data]]) : null,
237+
requestId,
238+
})
239+
}

apps/sim/lib/virtual-tables/memory-virtual-table.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import type { JsonValue, TableDefinition, TableRow, TableSchema } from '@/lib/table/types'
22

33
const MEMORY_TABLE_ID_PREFIX = 'system_memory_'
4+
export const MEMORY_TABLE_NAME = 'Memory'
45

56
export const MEMORY_TABLE_COLUMNS = {
67
id: 'id',
@@ -97,7 +98,7 @@ export function createMemoryTableDefinition({
9798
return {
9899
id: getMemoryTableId(workspace.id),
99100
isVirtual: true,
100-
name: 'Memory',
101+
name: MEMORY_TABLE_NAME,
101102
description: 'Read-only agent conversation memory for this workspace',
102103
schema: MEMORY_TABLE_SCHEMA,
103104
metadata: null,

0 commit comments

Comments
 (0)