Skip to content

Commit 2d3f1f3

Browse files
committed
fix(knowledge): prove connector leases last, page ACL writes by projection rows, and skip unchanged work
1 parent 55986a3 commit 2d3f1f3

17 files changed

Lines changed: 968 additions & 391 deletions

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

Lines changed: 181 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,7 @@ import {
6464
stillHoldsMemberSyncLock,
6565
stillHoldsSyncLock,
6666
} from '@/lib/knowledge/connectors/sync-lock'
67+
import * as syncPersistence from '@/lib/knowledge/connectors/sync-persistence'
6768
import {
6869
persistDocumentAcls,
6970
restoreWorkspaceDocumentAcls,
@@ -140,6 +141,8 @@ describe('connector lease ACL pages in PostgreSQL', () => {
140141
mimeType: 'text/plain',
141142
processingStatus: 'completed',
142143
contentHash: 'fixture-content',
144+
/** Ten chunks each, so a page of projection rows holds 25 documents. */
145+
chunkCount: 10,
143146
acl,
144147
}))
145148
await db.insert(document).values(rows)
@@ -167,7 +170,7 @@ describe('connector lease ACL pages in PostgreSQL', () => {
167170
const bounds = await db.execute<{ lock: string; statement: string }>(sql`
168171
SELECT DISTINCT lock_timeout AS lock, statement_timeout AS statement
169172
FROM lease_page_acl_writes WHERE connector_id = ${connectorId}`)
170-
expect([...bounds]).toEqual([{ lock: '5s', statement: '30s' }])
173+
expect([...bounds]).toEqual([{ lock: '15s', statement: '30s' }])
171174
}
172175

173176
/** Takes the lease away once `held` pages have committed, as a reclaim between pages would. */
@@ -236,6 +239,134 @@ describe('connector lease ACL pages in PostgreSQL', () => {
236239
})
237240
})
238241

242+
describe('fence-last pages', () => {
243+
/**
244+
* A processing commit holds a document row for its whole write. A page waiting on it must not
245+
* hold the connector row meanwhile, or every heartbeat, edit and reclaim queues behind it.
246+
*/
247+
it('waits on a locked document row without holding any lock on the connector table', async () => {
248+
const [locked] = await seedDocuments(ids.connectorId, [alice()], 1)
249+
let release!: () => void
250+
let held!: () => void
251+
const holding = new Promise<void>((resolve) => {
252+
held = resolve
253+
})
254+
const holder = db.transaction(async (tx) => {
255+
await tx
256+
.select({ id: document.id })
257+
.from(document)
258+
.where(eq(document.id, locked.id))
259+
.for('update')
260+
held()
261+
await new Promise<void>((resolve) => {
262+
release = resolve
263+
})
264+
})
265+
await holding
266+
const write = persistDocumentAcls(
267+
ids.connectorId,
268+
new Map([[locked.externalId, [bob()]]]),
269+
leaseTransaction(ids.connectorId, adminLease())
270+
)
271+
try {
272+
let waiter: number | undefined
273+
for (let attempt = 0; attempt < 100 && waiter === undefined; attempt++) {
274+
const [row] = await db.execute<{ pid: number }>(sql`
275+
SELECT pid FROM pg_stat_activity
276+
WHERE wait_event_type = 'Lock' AND query ILIKE 'update "document" set "acl"%'`)
277+
waiter = row?.pid
278+
if (waiter === undefined) await new Promise<void>((resolve) => setImmediate(resolve))
279+
}
280+
expect(waiter).toBeDefined()
281+
const connectorLocks = await db.execute<{ mode: string }>(sql`
282+
SELECT mode FROM pg_locks
283+
WHERE pid = ${waiter} AND relation = 'knowledge_connector'::regclass`)
284+
expect([...connectorLocks]).toEqual([])
285+
} finally {
286+
release()
287+
await holder
288+
}
289+
await expect(write).resolves.toEqual({ updated: 1, rejected: 0 })
290+
expect((await storedAcls(ids.connectorId)).map((acl) => acl.join())).toEqual([bob()])
291+
}, 30_000)
292+
293+
/** The member engine's ACL pages prove the lease last as well. */
294+
it('holds no connector lock while a member ACL page waits on a locked document row', async () => {
295+
const seeded = await seedDocuments(members.connectorId, [], 1)
296+
const [member] = members.members
297+
await recordMemberObservations(db, member.id, [seeded[0].id], members.runId)
298+
await db
299+
.update(knowledgeConnectorMember)
300+
.set({ listingCheckpoint: { kind: 'membership', cursor: null, removeMember: false } })
301+
.where(eq(knowledgeConnectorMember.id, member.id))
302+
let release!: () => void
303+
let held!: () => void
304+
const holding = new Promise<void>((resolve) => {
305+
held = resolve
306+
})
307+
const holder = db.transaction(async (tx) => {
308+
await tx
309+
.select({ id: document.id })
310+
.from(document)
311+
.where(eq(document.id, seeded[0].id))
312+
.for('update')
313+
held()
314+
await new Promise<void>((resolve) => {
315+
release = resolve
316+
})
317+
})
318+
await holding
319+
const rewrite = resumeMembershipRewrites({
320+
connectorId: members.connectorId,
321+
runId: members.runId,
322+
deadlineAt: Date.now() + 60_000,
323+
lease: createMemberSyncLease(members.connectorId, members.runId),
324+
})
325+
try {
326+
let waiter: number | undefined
327+
for (let attempt = 0; attempt < 200 && waiter === undefined; attempt++) {
328+
const [row] = await db.execute<{ pid: number }>(sql`
329+
SELECT pid FROM pg_stat_activity
330+
WHERE wait_event_type = 'Lock' AND query ILIKE 'update "document" set "acl"%'`)
331+
waiter = row?.pid
332+
if (waiter === undefined) await new Promise<void>((resolve) => setImmediate(resolve))
333+
}
334+
expect(waiter).toBeDefined()
335+
const connectorLocks = await db.execute<{ mode: string }>(sql`
336+
SELECT mode FROM pg_locks
337+
WHERE pid = ${waiter} AND relation = 'knowledge_connector'::regclass`)
338+
expect([...connectorLocks]).toEqual([])
339+
} finally {
340+
release()
341+
await holder
342+
}
343+
await expect(rewrite).resolves.toBe(true)
344+
expect((await storedAcls(members.connectorId)).map((acl) => acl.join())).toEqual([
345+
member.subjectToken,
346+
])
347+
}, 30_000)
348+
349+
/** One page is bounded by projection rows; a document larger than the cap still lands, alone. */
350+
it('gives a document larger than one page of projection rows a page alone', async () => {
351+
const seeded = await seedDocuments(ids.connectorId, [alice()], 3)
352+
await db.update(document).set({ chunkCount: 1_000 }).where(eq(document.id, seeded[1].id))
353+
await db
354+
.update(document)
355+
.set({ chunkCount: 200 })
356+
.where(inArray(document.id, [seeded[0].id, seeded[2].id]))
357+
358+
await expect(
359+
persistDocumentAcls(
360+
ids.connectorId,
361+
new Map(seeded.map((row) => [row.externalId, [bob()]])),
362+
leaseTransaction(ids.connectorId, adminLease())
363+
)
364+
).resolves.toEqual({ updated: 3, rejected: 0 })
365+
366+
expect(await writesPerTransaction(ids.connectorId)).toEqual([1, 1, 1])
367+
})
368+
})
369+
239370
describe('restoreWorkspaceDocumentAcls', () => {
240371
beforeEach(async () => {
241372
await db
@@ -269,10 +400,10 @@ describe('connector lease ACL pages in PostgreSQL', () => {
269400
).resolves.toBe(0)
270401
})
271402

272-
it('restores drift before a sync completes, outside the completion transaction', async () => {
403+
it('finishes a pending rewrite before a sync completes, outside the completion transaction', async () => {
273404
await db
274405
.update(knowledgeConnector)
275-
.set({ status: 'active', syncLockToken: null })
406+
.set({ status: 'active', syncLockToken: null, accessRewritePending: true })
276407
.where(eq(knowledgeConnector.id, ids.connectorId))
277408
const seeded = await seedDocuments(ids.connectorId, [])
278409
await db
@@ -316,10 +447,56 @@ describe('connector lease ACL pages in PostgreSQL', () => {
316447
.select({
317448
status: knowledgeConnector.status,
318449
syncLockToken: knowledgeConnector.syncLockToken,
450+
accessRewritePending: knowledgeConnector.accessRewritePending,
319451
})
320452
.from(knowledgeConnector)
321453
.where(eq(knowledgeConnector.id, ids.connectorId))
322-
expect(connector).toEqual({ status: 'active', syncLockToken: null })
454+
expect(connector).toEqual({
455+
status: 'active',
456+
syncLockToken: null,
457+
accessRewritePending: false,
458+
})
459+
})
460+
461+
/** Only a pending switch leaves workspace documents off the workspace ACL; a healthy sync never walks them. */
462+
it('does not walk a workspace connector without a pending rewrite', async () => {
463+
await db
464+
.update(knowledgeConnector)
465+
.set({ status: 'active', syncLockToken: null, accessRewritePending: false })
466+
.where(eq(knowledgeConnector.id, ids.connectorId))
467+
const seeded = await seedDocuments(ids.connectorId, ['ws'])
468+
await db
469+
.update(document)
470+
.set({ sourceSeenAt: sql`now() + interval '1 day'`, storageKey: sql`'kb/fixture/' || id` })
471+
.where(eq(document.connectorId, ids.connectorId))
472+
provider.list.mockResolvedValue({
473+
documents: seeded.map((row) => ({
474+
externalId: row.externalId,
475+
title: row.filename,
476+
content: '',
477+
contentDeferred: true,
478+
contentHash: 'fixture-content',
479+
mimeType: 'text/plain',
480+
})),
481+
hasMore: false,
482+
})
483+
const token = vi
484+
.spyOn(connectorTokens, 'resolveConnectorAccessToken')
485+
.mockResolvedValue({ accessToken: 'fixture-token' } as never)
486+
const restore = vi.spyOn(syncPersistence, 'restoreWorkspaceDocumentAcls')
487+
try {
488+
const result = await executeSync(ids.connectorId, {
489+
billingAttribution: await resolveBillingAttribution({
490+
actorUserId: ids.aliceId,
491+
workspaceId: ids.workspaceId,
492+
}),
493+
})
494+
expect(result.error).toBeUndefined()
495+
expect(restore).not.toHaveBeenCalled()
496+
} finally {
497+
token.mockRestore()
498+
restore.mockRestore()
499+
}
323500
})
324501

325502
it('restores nothing once the connector has left workspace mode', async () => {

‎apps/sim/lib/knowledge/__integration__/migration-fixture.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,8 @@ export async function createEnterpriseSearchMigrationFixture(databaseUrl: string
5757
id text PRIMARY KEY, external_id text, connector_id text, knowledge_base_id text,
5858
tag1 text, tag2 text, tag3 text, tag4 text, tag5 text, tag6 text, tag7 text,
5959
acl text[] NOT NULL DEFAULT '{ws}', storage_key text,
60-
user_excluded boolean NOT NULL DEFAULT false, archived_at timestamp
60+
user_excluded boolean NOT NULL DEFAULT false, archived_at timestamp,
61+
chunk_count integer NOT NULL DEFAULT 0
6162
);
6263
CREATE INDEX doc_connector_id_idx ON document(connector_id);
6364
CREATE TABLE knowledge_connector (

‎apps/sim/lib/knowledge/connectors/detachment.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,10 +24,10 @@ import {
2424
import type { DbOrTx } from '@/lib/db/types'
2525
import { removeDrainedConnector } from '@/lib/knowledge/connectors/deletion'
2626
import { revokeKnowledgeConnectorCredentialAccess } from '@/lib/knowledge/connectors/member-access'
27+
import { PROJECTION_ROW_BATCH_SIZE } from '@/lib/knowledge/connectors/sync-limits'
2728

2829
export const KNOWLEDGE_CONNECTOR_DETACH_EVENT = 'knowledge.connector.detach'
2930
const DOCUMENT_BATCH_SIZE = 100
30-
const PROJECTION_ROW_BATCH_SIZE = 250
3131
const MAX_BATCHES_PER_RUN = 4
3232
const RUN_BUDGET_MS = 30_000
3333
/** How often a detachment paused on a deleted knowledge base checks for its restore or purge. */

0 commit comments

Comments
 (0)