Skip to content

Commit d0d2bcd

Browse files
committed
fix(knowledge): row-bound every observation page, size ACL pages from locked chunk counts, and check the sweep budget per page
1 parent c09f9a3 commit d0d2bcd

7 files changed

Lines changed: 369 additions & 88 deletions

File tree

‎apps/sim/lib/knowledge/__integration__/connector-lease-pages.integration.ts‎

Lines changed: 63 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -273,7 +273,9 @@ describe('connector lease ACL pages in PostgreSQL', () => {
273273
for (let attempt = 0; attempt < 100 && waiter === undefined; attempt++) {
274274
const [row] = await db.execute<{ pid: number }>(sql`
275275
SELECT pid FROM pg_stat_activity
276-
WHERE wait_event_type = 'Lock' AND query ILIKE 'update "document" set "acl"%'`)
276+
WHERE wait_event_type = 'Lock' AND datname = current_database()
277+
AND (query ILIKE 'update "document" set "acl"%'
278+
OR query ILIKE 'select "id", "chunk_count" from "document"%for update')`)
277279
waiter = row?.pid
278280
if (waiter === undefined) await new Promise<void>((resolve) => setImmediate(resolve))
279281
}
@@ -290,6 +292,63 @@ describe('connector lease ACL pages in PostgreSQL', () => {
290292
expect((await storedAcls(ids.connectorId)).map((acl) => acl.join())).toEqual([bob()])
291293
}, 30_000)
292294

295+
/**
296+
* Pages are sized from a read taken without a lock; a reprocess can change a document's chunks
297+
* before the page writes. The page locks its documents and rereads their counts, so what it
298+
* writes still fits one page of projection rows.
299+
*/
300+
it('resizes a page whose documents gained chunks after it was planned', async () => {
301+
const seeded = await seedDocuments(ids.connectorId, [alice()], 3)
302+
const reprocess = losingAfter(1, leaseTransaction(ids.connectorId, adminLease()), () =>
303+
db
304+
.update(document)
305+
.set({ chunkCount: 1_000 })
306+
.where(eq(document.connectorId, ids.connectorId))
307+
)
308+
309+
await expect(
310+
persistDocumentAcls(
311+
ids.connectorId,
312+
new Map(seeded.map((row) => [row.externalId, [bob()]])),
313+
reprocess
314+
)
315+
).resolves.toEqual({ updated: 3, rejected: 0 })
316+
317+
expect(await writesPerTransaction(ids.connectorId)).toEqual([1, 1, 1])
318+
})
319+
320+
/** Observation changes that must commit with their ACLs are paged by projection rows too. */
321+
it('rewrites a membership page by projection rows, a large document alone', async () => {
322+
const seeded = await seedDocuments(members.connectorId, [], 3)
323+
const [member] = members.members
324+
await recordMemberObservations(
325+
db,
326+
member.id,
327+
seeded.map((row) => row.id),
328+
members.runId
329+
)
330+
const huge = [...seeded].sort((a, b) => (a.id < b.id ? -1 : 1))[1]
331+
await db.update(document).set({ chunkCount: 1_000 }).where(eq(document.id, huge.id))
332+
await db
333+
.update(knowledgeConnectorMember)
334+
.set({ listingCheckpoint: { kind: 'membership', cursor: null, removeMember: false } })
335+
.where(eq(knowledgeConnectorMember.id, member.id))
336+
337+
await expect(
338+
resumeMembershipRewrites({
339+
connectorId: members.connectorId,
340+
runId: members.runId,
341+
deadlineAt: Date.now() + 60_000,
342+
lease: createMemberSyncLease(members.connectorId, members.runId),
343+
})
344+
).resolves.toBe(true)
345+
346+
expect(await writesPerTransaction(members.connectorId)).toEqual([1, 1, 1])
347+
expect(
348+
(await storedAcls(members.connectorId)).every((acl) => acl.join() === member.subjectToken)
349+
).toBe(true)
350+
})
351+
293352
/** The member engine's ACL pages prove the lease last as well. */
294353
it('holds no connector lock while a member ACL page waits on a locked document row', async () => {
295354
const seeded = await seedDocuments(members.connectorId, [], 1)
@@ -327,7 +386,9 @@ describe('connector lease ACL pages in PostgreSQL', () => {
327386
for (let attempt = 0; attempt < 200 && waiter === undefined; attempt++) {
328387
const [row] = await db.execute<{ pid: number }>(sql`
329388
SELECT pid FROM pg_stat_activity
330-
WHERE wait_event_type = 'Lock' AND query ILIKE 'update "document" set "acl"%'`)
389+
WHERE wait_event_type = 'Lock' AND datname = current_database()
390+
AND (query ILIKE 'update "document" set "acl"%'
391+
OR query ILIKE 'select "id", "chunk_count" from "document"%for update')`)
331392
waiter = row?.pid
332393
if (waiter === undefined) await new Promise<void>((resolve) => setImmediate(resolve))
333394
}

‎apps/sim/lib/knowledge/connectors/member-observations.test.ts‎

Lines changed: 128 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -48,10 +48,13 @@ describe('removeUnseenMemberObservations', () => {
4848
resetDbChainMock()
4949
})
5050
it('rematerializes bounded batches before reading more absent observations', async () => {
51+
const candidates = (count: number, prefix: string) =>
52+
Array.from({ length: count }, (_, i) => ({ documentId: `${prefix}-${i}` }))
53+
queueTableRows(schemaMock.knowledgeDocumentObservation, candidates(25, 'd'))
54+
queueTableRows(schemaMock.knowledgeDocumentObservation, candidates(1, 'last'))
5155
dbChainMockFns.returning
52-
.mockResolvedValueOnce(Array.from({ length: 25 }, (_, i) => ({ documentId: `d-${i}` })))
53-
.mockResolvedValueOnce([{ documentId: 'last' }])
54-
.mockResolvedValueOnce([])
56+
.mockResolvedValueOnce(candidates(25, 'd'))
57+
.mockResolvedValueOnce(candidates(1, 'last'))
5558
const onRemoved = vi.fn(async (_ids: string[]) => undefined)
5659
await expect(
5760
removeUnseenMemberObservations(db, 'member', 'generation', onRemoved)
@@ -66,6 +69,25 @@ describe('removeUnseenMemberObservations', () => {
6669
dbChainMockFns.delete.mock.invocationCallOrder[1]
6770
)
6871
})
72+
73+
/** One document above the row cap is a page alone; the rest of the candidates wait. */
74+
it('removes only the leading page of projection rows and reports more to come', async () => {
75+
queueTableRows(schemaMock.knowledgeDocumentObservation, [
76+
{ documentId: 'huge' },
77+
{ documentId: 'small' },
78+
])
79+
queueTableRows(schemaMock.document, [
80+
{ id: 'huge', chunkCount: 1_000 },
81+
{ id: 'small', chunkCount: 1 },
82+
])
83+
dbChainMockFns.returning.mockResolvedValueOnce([{ documentId: 'huge' }])
84+
const onRemoved = vi.fn(async (_ids: string[]) => undefined)
85+
86+
await expect(
87+
removeUnseenMemberObservations(db, 'member', 'generation', onRemoved)
88+
).resolves.toEqual({ removed: 1, finished: false })
89+
expect(onRemoved).toHaveBeenCalledWith(['huge'])
90+
})
6991
})
7092

7193
describe('renewMemberObservationsInScopes', () => {
@@ -163,6 +185,10 @@ describe('sweepStaleMemberObservations', () => {
163185
queueTableRows(schemaMock.knowledgeConnectorMember, [STALE_MEMBER])
164186
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-1' }])
165187
queueTableRows(schemaMock.knowledgeConnectorMember, [{ id: 'm-1' }])
188+
queueTableRows(schemaMock.knowledgeDocumentObservation, [
189+
{ documentId: 'd-1' },
190+
{ documentId: 'd-2' },
191+
])
166192
dbChainMockFns.returning
167193
.mockResolvedValueOnce([{ documentId: 'd-1' }, { documentId: 'd-2' }])
168194
.mockResolvedValueOnce([{ id: 'd-1' }, { id: 'd-2' }])
@@ -176,10 +202,13 @@ describe('sweepStaleMemberObservations', () => {
176202
})
177203

178204
expect(dbChainMockFns.transaction).toHaveBeenCalledOnce()
179-
/** The member row first, the connector row last: never held while documents are written. */
180-
expect(dbChainMockFns.for).toHaveBeenNthCalledWith(1, 'update')
181-
expect(dbChainMockFns.for).toHaveBeenNthCalledWith(2, 'share')
182-
expect(dbChainMockFns.for.mock.invocationCallOrder[1]).toBeGreaterThan(
205+
/** The member row first, then the page's documents, the connector row last. */
206+
expect(dbChainMockFns.for.mock.calls.map(([mode]) => mode)).toEqual([
207+
'update',
208+
'update',
209+
'share',
210+
])
211+
expect(dbChainMockFns.for.mock.invocationCallOrder.at(-1)).toBeGreaterThan(
183212
dbChainMockFns.set.mock.invocationCallOrder.at(-1)!
184213
)
185214
expect(dbChainMockFns.delete).toHaveBeenCalledWith(schemaMock.knowledgeDocumentObservation)
@@ -194,6 +223,10 @@ describe('sweepStaleMemberObservations', () => {
194223
for (const page of [ids(25, 'a'), ids(3, 'b')]) {
195224
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-1' }])
196225
queueTableRows(schemaMock.knowledgeConnectorMember, [{ id: 'm-1' }])
226+
queueTableRows(
227+
schemaMock.knowledgeDocumentObservation,
228+
page.map((documentId) => ({ documentId }))
229+
)
197230
dbChainMockFns.returning
198231
.mockResolvedValueOnce(page.map((documentId) => ({ documentId })))
199232
.mockResolvedValueOnce(page.map((id) => ({ id })))
@@ -215,6 +248,76 @@ describe('sweepStaleMemberObservations', () => {
215248
expect(bounds).toHaveLength(2)
216249
})
217250

251+
/** A document above the row cap is swept alone; the rest of the member waits for the next page. */
252+
it('sweeps a document larger than one page of projection rows in a page alone', async () => {
253+
queueTableRows(schemaMock.knowledgeConnectorMember, [STALE_MEMBER])
254+
for (const page of [['huge', 'small'], ['small']]) {
255+
queueTableRows(schemaMock.knowledgeConnectorMember, [{ id: 'm-1' }])
256+
queueTableRows(
257+
schemaMock.knowledgeDocumentObservation,
258+
page.map((documentId) => ({ documentId }))
259+
)
260+
queueTableRows(
261+
schemaMock.document,
262+
page.map((id) => ({ id, chunkCount: id === 'huge' ? 1_000 : 1 }))
263+
)
264+
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-1' }])
265+
dbChainMockFns.returning
266+
.mockResolvedValueOnce([{ documentId: page[0] }])
267+
.mockResolvedValueOnce([{ id: page[0] }])
268+
.mockResolvedValueOnce([])
269+
}
270+
271+
await expect(sweepStaleMemberObservations(NOW)).resolves.toMatchObject({
272+
members: 1,
273+
observationsRemoved: 2,
274+
})
275+
const deletedPages = dbChainMockFns.where.mock.calls
276+
.map(([condition]) =>
277+
flattenMockConditions(condition).find(
278+
(node) =>
279+
node.type === 'inArray' &&
280+
node.column === schemaMock.knowledgeDocumentObservation.documentId
281+
)
282+
)
283+
.filter((node) => node !== undefined)
284+
.map((node) => node!.values)
285+
expect(deletedPages).toEqual([['huge'], ['small']])
286+
})
287+
288+
/** The budget is checked before every page, so one large member cannot hold a tick past it. */
289+
it('leaves the rest of a large member for the next tick once the budget passes mid-member', async () => {
290+
let clock = Date.now()
291+
const now = vi.spyOn(Date, 'now').mockImplementation(() => clock)
292+
try {
293+
const page = Array.from({ length: 25 }, (_unused, index) => `a-${index}`)
294+
queueTableRows(schemaMock.knowledgeConnectorMember, [STALE_MEMBER])
295+
queueTableRows(schemaMock.knowledgeConnectorMember, [{ id: 'm-1' }])
296+
queueTableRows(
297+
schemaMock.knowledgeDocumentObservation,
298+
page.map((documentId) => ({ documentId }))
299+
)
300+
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-1' }])
301+
dbChainMockFns.returning
302+
.mockResolvedValueOnce(page.map((documentId) => ({ documentId })))
303+
.mockResolvedValueOnce(page.map((id) => ({ id })))
304+
.mockImplementationOnce(async () => {
305+
clock += 2_000
306+
return []
307+
})
308+
309+
await expect(sweepStaleMemberObservations(NOW, clock + 1_000)).resolves.toEqual({
310+
members: 1,
311+
observationsRemoved: 25,
312+
documentsRematerialized: 25,
313+
docsTombstoned: 0,
314+
})
315+
expect(dbChainMockFns.transaction).toHaveBeenCalledOnce()
316+
} finally {
317+
now.mockRestore()
318+
}
319+
})
320+
218321
/** A connector whose running member page holds its row must not stall the other connectors. */
219322
it('defers a member whose page hits a lock timeout and still sweeps the next one', async () => {
220323
queueTableRows(schemaMock.knowledgeConnectorMember, [
@@ -227,6 +330,7 @@ describe('sweepStaleMemberObservations', () => {
227330
Object.assign(new Error('canceling statement due to lock timeout'), { code: '55P03' })
228331
)
229332
queueTableRows(schemaMock.knowledgeConnectorMember, [{ id: 'm-2' }])
333+
queueTableRows(schemaMock.knowledgeDocumentObservation, [{ documentId: 'd-1' }])
230334
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-2' }])
231335
dbChainMockFns.returning
232336
.mockResolvedValueOnce([{ documentId: 'd-1' }])
@@ -412,14 +516,22 @@ describe('rewriteConnectorAcls', () => {
412516
.map(([values], index) => ({ values, where: dbChainMockFns.where.mock.calls[index] }))
413517
.filter(({ values }) => 'acl' in values)
414518
const pageSizes = () =>
415-
dbChainMockFns.where.mock.calls
416-
.map(([condition]) =>
417-
flattenMockConditions(condition).find(
418-
(node) => node.type === 'inArray' && node.column === schemaMock.document.id
519+
dbChainMockFns.set.mock.calls
520+
.map(([values], index) => ({
521+
values,
522+
order: dbChainMockFns.set.mock.invocationCallOrder[index],
523+
}))
524+
.filter(({ values }) => 'acl' in values || 'aclRequirements' in values)
525+
.map(({ order }) => {
526+
const whereIndex = dbChainMockFns.where.mock.invocationCallOrder.findIndex(
527+
(whereOrder) => whereOrder > order
419528
)
420-
)
421-
.filter((node) => node !== undefined)
422-
.map((node) => (node!.values as string[]).length)
529+
return (
530+
flattenMockConditions(dbChainMockFns.where.mock.calls[whereIndex]?.[0]).find(
531+
(node) => node.type === 'inArray' && node.column === schemaMock.document.id
532+
)?.values as string[]
533+
).length
534+
})
423535

424536
it('proves the lease inside each page transaction before rewriting', async () => {
425537
queueTableRows(schemaMock.document, [stale('d-1')])
@@ -447,7 +559,8 @@ describe('rewriteConnectorAcls', () => {
447559
)
448560

449561
expect(dbChainMockFns.update).toHaveBeenCalledOnce()
450-
expect(dbChainMockFns.for.mock.invocationCallOrder[0]).toBeGreaterThan(
562+
const leaseCheck = dbChainMockFns.for.mock.calls.findIndex(([mode]) => mode === 'share')
563+
expect(dbChainMockFns.for.mock.invocationCallOrder[leaseCheck]).toBeGreaterThan(
451564
dbChainMockFns.update.mock.invocationCallOrder[0]
452565
)
453566
})

0 commit comments

Comments
 (0)