Skip to content

Commit c157ef1

Browse files
committed
fix(knowledge): stop every ACL page walker at the run budget and keep pending rewrites until they finish
1 parent d0d2bcd commit c157ef1

6 files changed

Lines changed: 488 additions & 36 deletions

File tree

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

Lines changed: 288 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,7 @@
77
import { db } from '@sim/db'
88
import {
99
document,
10+
embedding,
1011
knowledgeConnector,
1112
knowledgeConnectorMember,
1213
knowledgeConnectorMemberSyncLog,
@@ -45,6 +46,7 @@ import {
4546
} from '@/lib/knowledge/__integration__/seed-source-access-fixture'
4647
import * as connectorTokens from '@/lib/knowledge/connectors/access-token'
4748
import * as memberAccess from '@/lib/knowledge/connectors/member-access'
49+
import * as memberObservations from '@/lib/knowledge/connectors/member-observations'
4850
import {
4951
materializeDocumentAcls,
5052
recordMemberObservations,
@@ -56,6 +58,7 @@ import {
5658
resumeMembershipRewrites,
5759
} from '@/lib/knowledge/connectors/member-sync-engine'
5860
import { executeSync } from '@/lib/knowledge/connectors/sync-engine'
61+
import { PROJECTION_ROW_BATCH_SIZE } from '@/lib/knowledge/connectors/sync-limits'
5962
import {
6063
createMemberSyncLease,
6164
type LeaseTransaction,
@@ -103,6 +106,26 @@ describe('connector lease ACL pages in PostgreSQL', () => {
103106
await db.execute(sql`DROP TRIGGER IF EXISTS log_lease_page_acl_write ON document`)
104107
await db.execute(sql`CREATE TRIGGER log_lease_page_acl_write AFTER UPDATE OF acl ON document
105108
FOR EACH ROW EXECUTE FUNCTION log_lease_page_acl_write()`)
109+
/** Records, for every projection row the document trigger rewrites, its table and transaction. */
110+
await db.execute(sql`CREATE TABLE IF NOT EXISTS lease_page_projection_writes (
111+
projection text NOT NULL, document_id text NOT NULL, xact text NOT NULL
112+
)`)
113+
await db.execute(
114+
sql.raw(`CREATE OR REPLACE FUNCTION log_lease_page_projection_write() RETURNS trigger
115+
LANGUAGE plpgsql AS $$ BEGIN
116+
INSERT INTO lease_page_projection_writes VALUES (TG_TABLE_NAME, NEW.document_id, pg_current_xact_id()::text);
117+
RETURN NEW;
118+
END $$`)
119+
)
120+
for (const projection of ['embedding_search', 'embedding_keyword_tin']) {
121+
await db.execute(
122+
sql.raw(`DROP TRIGGER IF EXISTS log_lease_page_projection_write ON ${projection}`)
123+
)
124+
await db.execute(
125+
sql.raw(`CREATE TRIGGER log_lease_page_projection_write AFTER UPDATE OF acl ON ${projection}
126+
FOR EACH ROW EXECUTE FUNCTION log_lease_page_projection_write()`)
127+
)
128+
}
106129
})
107130

108131
beforeEach(async () => {
@@ -114,6 +137,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
114137
afterEach(async () => {
115138
await db.execute(sql`DROP TRIGGER IF EXISTS fail_after_acl_writes ON document`)
116139
await db.execute(sql`DELETE FROM lease_page_acl_writes`)
140+
await db.execute(sql`DELETE FROM lease_page_projection_writes`)
117141
await db.delete(workspace).where(eq(workspace.id, ids.workspaceId))
118142
await db.delete(organization).where(eq(organization.id, ids.organizationId))
119143
await db.delete(user).where(inArray(user.id, [ids.aliceId, ids.bobId]))
@@ -124,6 +148,12 @@ describe('connector lease ACL pages in PostgreSQL', () => {
124148
await db.execute(sql`DROP FUNCTION IF EXISTS log_lease_page_acl_write()`)
125149
await db.execute(sql`DROP FUNCTION IF EXISTS fail_after_acl_writes()`)
126150
await db.execute(sql`DROP TABLE IF EXISTS lease_page_acl_writes`)
151+
for (const projection of ['embedding_search', 'embedding_keyword_tin'])
152+
await db.execute(
153+
sql.raw(`DROP TRIGGER IF EXISTS log_lease_page_projection_write ON ${projection}`)
154+
)
155+
await db.execute(sql`DROP FUNCTION IF EXISTS log_lease_page_projection_write()`)
156+
await db.execute(sql`DROP TABLE IF EXISTS lease_page_projection_writes`)
127157
vi.restoreAllMocks()
128158
vi.unstubAllGlobals()
129159
await db.$client.end()
@@ -239,6 +269,112 @@ describe('connector lease ACL pages in PostgreSQL', () => {
239269
})
240270
})
241271

272+
describe('search projection fan-out', () => {
273+
const CHUNKS = 20
274+
275+
/** Real chunks: the installed triggers create each chunk's search and keyword projection rows. */
276+
const seedChunks = async (documents: { id: string }[]) => {
277+
const rows = documents.flatMap((entry) =>
278+
Array.from({ length: CHUNKS }, (_unused, chunkIndex) => ({
279+
id: generateId(),
280+
knowledgeBaseId: ids.knowledgeBaseId,
281+
documentId: entry.id,
282+
chunkIndex,
283+
chunkHash: `${entry.id}-${chunkIndex}`,
284+
content: `lease page chunk ${chunkIndex}`,
285+
contentLength: 20,
286+
tokenCount: 4,
287+
embedding: [1, ...Array<number>(1535).fill(0)],
288+
startOffset: 0,
289+
endOffset: 20,
290+
}))
291+
)
292+
for (let offset = 0; offset < rows.length; offset += 200)
293+
await db.insert(embedding).values(rows.slice(offset, offset + 200))
294+
/**
295+
* The keyword projection's own sync trigger ships with the Tin migration, which a database
296+
* without the Tin extension skips; write the rows it would, and its installed ACL trigger
297+
* fills them from the document as it does for every insert.
298+
*/
299+
await db.execute(sql`
300+
INSERT INTO embedding_keyword_tin (id, knowledge_base_id, document_id, enabled, content)
301+
SELECT e.id, e.knowledge_base_id, e.document_id, e.enabled, e.content FROM embedding e
302+
WHERE e.document_id IN (${sql.join(
303+
documents.map((entry) => sql`${entry.id}`),
304+
sql`, `
305+
)})
306+
ON CONFLICT (id) DO NOTHING`)
307+
await db
308+
.update(document)
309+
.set({ chunkCount: CHUNKS })
310+
.where(
311+
inArray(
312+
document.id,
313+
documents.map((entry) => entry.id)
314+
)
315+
)
316+
}
317+
318+
/** Every projection row of the connector's documents, with its ACL and the document's, as text. */
319+
const projectionAcls = async (connectorId: string) =>
320+
db.execute<{ projection: string; acl: string | null; expected: string }>(sql`
321+
SELECT 'embedding_search' AS projection, array_to_string(p.acl, ',') AS acl,
322+
array_to_string(d.acl, ',') AS expected
323+
FROM embedding_search p JOIN document d ON d.id = p.document_id WHERE d.connector_id = ${connectorId}
324+
UNION ALL
325+
SELECT 'embedding_keyword_tin', array_to_string(p.acl, ','), array_to_string(d.acl, ',')
326+
FROM embedding_keyword_tin p JOIN document d ON d.id = p.document_id WHERE d.connector_id = ${connectorId}`)
327+
328+
/** Projection rows each transaction rewrote, per table. */
329+
const projectionRowsPerTransaction = async () =>
330+
(
331+
await db.execute<{ rows: number }>(sql`
332+
SELECT count(*)::int AS rows FROM lease_page_projection_writes
333+
GROUP BY projection, xact ORDER BY rows DESC`)
334+
).map((row) => row.rows)
335+
336+
it('mirrors a paged ACL write onto every projection row, one page of rows per transaction', async () => {
337+
const seeded = await seedDocuments(ids.connectorId, [alice()], 30)
338+
await seedChunks(seeded)
339+
const before = [...(await projectionAcls(ids.connectorId))]
340+
expect(before).toHaveLength(2 * 30 * CHUNKS)
341+
342+
await expect(
343+
persistDocumentAcls(
344+
ids.connectorId,
345+
new Map(seeded.map((row) => [row.externalId, [bob()]])),
346+
leaseTransaction(ids.connectorId, adminLease())
347+
)
348+
).resolves.toEqual({ updated: 30, rejected: 0 })
349+
350+
const after = [...(await projectionAcls(ids.connectorId))]
351+
expect(after.every((row) => row.acl === bob() && row.expected === bob())).toBe(true)
352+
const perTransaction = await projectionRowsPerTransaction()
353+
expect(perTransaction.reduce((total, rows) => total + rows, 0)).toBe(2 * 30 * CHUNKS)
354+
expect(Math.max(...perTransaction)).toBeLessThanOrEqual(PROJECTION_ROW_BATCH_SIZE)
355+
})
356+
357+
it('hides a members connector across its projection rows, one page of rows per transaction', async () => {
358+
const seeded = await seedDocuments(members.connectorId, [alice()], 30)
359+
await seedChunks(seeded)
360+
361+
await expect(
362+
rewriteConnectorAcls(members.connectorId, [], {
363+
lease: {
364+
stillHeld: () => stillHoldsMemberSyncLock(members.connectorId, members.runId),
365+
},
366+
})
367+
).resolves.toBe(true)
368+
369+
const after = [...(await projectionAcls(members.connectorId))]
370+
expect(after).toHaveLength(2 * 30 * CHUNKS)
371+
expect(after.every((row) => row.acl === '' && row.expected === '')).toBe(true)
372+
expect(Math.max(...(await projectionRowsPerTransaction()))).toBeLessThanOrEqual(
373+
PROJECTION_ROW_BATCH_SIZE
374+
)
375+
})
376+
})
377+
242378
describe('fence-last pages', () => {
243379
/**
244380
* A processing commit holds a document row for its whole write. A page waiting on it must not
@@ -444,7 +580,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
444580
ids.connectorId,
445581
leaseTransaction(ids.connectorId, adminLease())
446582
)
447-
).resolves.toBe(DOCUMENTS)
583+
).resolves.toEqual({ restored: DOCUMENTS, finished: true })
448584

449585
expect(await writesPerTransaction(ids.connectorId)).toEqual([
450586
PAGE,
@@ -458,7 +594,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
458594
ids.connectorId,
459595
leaseTransaction(ids.connectorId, adminLease())
460596
)
461-
).resolves.toBe(0)
597+
).resolves.toEqual({ restored: 0, finished: true })
462598
})
463599

464600
it('finishes a pending rewrite before a sync completes, outside the completion transaction', async () => {
@@ -519,6 +655,155 @@ describe('connector lease ACL pages in PostgreSQL', () => {
519655
})
520656
})
521657

658+
/** A restore stops between pages at its deadline; a later walk resumes from what is still off. */
659+
it('stops a restore between pages at its deadline and resumes it on the next walk', async () => {
660+
await seedDocuments(ids.connectorId, [])
661+
let clock = Date.now()
662+
const now = vi.spyOn(Date, 'now').mockImplementation(() => clock)
663+
let pages = 0
664+
try {
665+
const partial = await restoreWorkspaceDocumentAcls(
666+
ids.connectorId,
667+
leaseTransaction(ids.connectorId, adminLease()),
668+
{
669+
deadlineAt: clock + 1_000,
670+
beforePage: async () => {
671+
pages += 1
672+
/** The window read, then the first page: the budget passes during that page. */
673+
if (pages === 2) clock += 2_000
674+
},
675+
}
676+
)
677+
expect(partial).toEqual({ restored: PAGE, finished: false })
678+
} finally {
679+
now.mockRestore()
680+
}
681+
await expect(
682+
restoreWorkspaceDocumentAcls(
683+
ids.connectorId,
684+
leaseTransaction(ids.connectorId, adminLease())
685+
)
686+
).resolves.toEqual({ restored: DOCUMENTS - PAGE, finished: true })
687+
expect((await storedAcls(ids.connectorId)).every((acl) => acl.join() === 'ws')).toBe(true)
688+
})
689+
690+
/** A sync whose restore ran out of budget keeps the flag and comes back at once to finish it. */
691+
it('keeps the pending rewrite of a sync whose restore did not finish', async () => {
692+
await db
693+
.update(knowledgeConnector)
694+
.set({ status: 'active', syncLockToken: null, accessRewritePending: true })
695+
.where(eq(knowledgeConnector.id, ids.connectorId))
696+
const seeded = await seedDocuments(ids.connectorId, [])
697+
await db
698+
.update(document)
699+
.set({ sourceSeenAt: sql`now() + interval '1 day'`, storageKey: sql`'kb/fixture/' || id` })
700+
.where(eq(document.connectorId, ids.connectorId))
701+
provider.list.mockResolvedValue({
702+
documents: seeded.map((row) => ({
703+
externalId: row.externalId,
704+
title: row.filename,
705+
content: '',
706+
contentDeferred: true,
707+
contentHash: 'fixture-content',
708+
mimeType: 'text/plain',
709+
})),
710+
hasMore: false,
711+
})
712+
const token = vi
713+
.spyOn(connectorTokens, 'resolveConnectorAccessToken')
714+
.mockResolvedValue({ accessToken: 'fixture-token' } as never)
715+
const original = syncPersistence.restoreWorkspaceDocumentAcls
716+
const restore = vi
717+
.spyOn(syncPersistence, 'restoreWorkspaceDocumentAcls')
718+
.mockImplementationOnce((connectorId, transaction, options) =>
719+
original(connectorId, transaction, { ...options, deadlineAt: Date.now() - 1 })
720+
)
721+
const billing = await resolveBillingAttribution({
722+
actorUserId: ids.aliceId,
723+
workspaceId: ids.workspaceId,
724+
})
725+
const state = async () => {
726+
const [row] = await db
727+
.select({
728+
accessRewritePending: knowledgeConnector.accessRewritePending,
729+
nextSyncAt: knowledgeConnector.nextSyncAt,
730+
})
731+
.from(knowledgeConnector)
732+
.where(eq(knowledgeConnector.id, ids.connectorId))
733+
return row
734+
}
735+
try {
736+
expect((await executeSync(ids.connectorId, { billingAttribution: billing })).error).toBe(
737+
undefined
738+
)
739+
/** The sync hands its own run budget to the restore. */
740+
expect(restore).toHaveBeenCalledWith(
741+
ids.connectorId,
742+
expect.any(Function),
743+
expect.objectContaining({ deadlineAt: expect.any(Number) })
744+
)
745+
const unfinished = await state()
746+
expect(unfinished?.accessRewritePending).toBe(true)
747+
expect(unfinished?.nextSyncAt?.getTime()).toBeLessThanOrEqual(Date.now())
748+
expect((await storedAcls(ids.connectorId)).every((acl) => acl.length === 0)).toBe(true)
749+
750+
expect((await executeSync(ids.connectorId, { billingAttribution: billing })).error).toBe(
751+
undefined
752+
)
753+
expect((await state())?.accessRewritePending).toBe(false)
754+
expect((await storedAcls(ids.connectorId)).every((acl) => acl.join() === 'ws')).toBe(true)
755+
} finally {
756+
token.mockRestore()
757+
restore.mockRestore()
758+
}
759+
})
760+
761+
/** An admin connector still hiding its documents lists nothing until the walk is done. */
762+
it('lists nothing and keeps the pending rewrite while an admin hide is unfinished', async () => {
763+
provider.list.mockClear()
764+
await db
765+
.update(knowledgeConnector)
766+
.set({
767+
accessMode: 'admin',
768+
status: 'active',
769+
syncLockToken: null,
770+
accessRewritePending: true,
771+
})
772+
.where(eq(knowledgeConnector.id, ids.connectorId))
773+
await seedDocuments(ids.connectorId, ['ws'])
774+
const token = vi
775+
.spyOn(connectorTokens, 'resolveConnectorAccessToken')
776+
.mockResolvedValue({ accessToken: 'fixture-token' } as never)
777+
const hide = vi
778+
.spyOn(memberObservations, 'rewriteConnectorAcls')
779+
.mockImplementationOnce(async (_connectorId, _target, options) => {
780+
expect(options?.deadlineAt).toEqual(expect.any(Number))
781+
return false
782+
})
783+
try {
784+
const result = await executeSync(ids.connectorId, {
785+
billingAttribution: await resolveBillingAttribution({
786+
actorUserId: ids.aliceId,
787+
workspaceId: ids.workspaceId,
788+
}),
789+
})
790+
expect(result.error).toBeUndefined()
791+
expect(hide).toHaveBeenCalledOnce()
792+
expect(provider.list).not.toHaveBeenCalled()
793+
const [row] = await db
794+
.select({
795+
accessRewritePending: knowledgeConnector.accessRewritePending,
796+
syncLockToken: knowledgeConnector.syncLockToken,
797+
})
798+
.from(knowledgeConnector)
799+
.where(eq(knowledgeConnector.id, ids.connectorId))
800+
expect(row).toEqual({ accessRewritePending: true, syncLockToken: null })
801+
} finally {
802+
token.mockRestore()
803+
hide.mockRestore()
804+
}
805+
})
806+
522807
/** Only a pending switch leaves workspace documents off the workspace ACL; a healthy sync never walks them. */
523808
it('does not walk a workspace connector without a pending rewrite', async () => {
524809
await db
@@ -572,7 +857,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
572857
ids.connectorId,
573858
leaseTransaction(ids.connectorId, adminLease())
574859
)
575-
).resolves.toBe(0)
860+
).resolves.toEqual({ restored: 0, finished: true })
576861
expect((await storedAcls(ids.connectorId)).every((acl) => acl.length === 0)).toBe(true)
577862
})
578863

0 commit comments

Comments
 (0)