Skip to content

Commit 41915f8

Browse files
committed
improvement(knowledge): bound each reconciliation window once and read it through the v2 id range
1 parent 8b96b78 commit 41915f8

8 files changed

Lines changed: 293 additions & 179 deletions

‎apps/sim/lib/knowledge/__integration__/listing-continuation.integration.ts‎

Lines changed: 184 additions & 111 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,7 @@ import { toNumberOrNull, toStringOrNull } from '@sim/utils/coerce'
2323
import { generateId } from '@sim/utils/id'
2424
import { toArray, toRecord } from '@sim/utils/object'
2525
import { and, eq, inArray, sql } from 'drizzle-orm'
26-
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
26+
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
2727

2828
const fixture = vi.hoisted(() => ({
2929
storageRoot: '',
@@ -901,11 +901,22 @@ describe('durable source and member cycles in PostgreSQL', () => {
901901
}
902902
})
903903

904-
it('bounds every reconciliation page by the ids it scans when absence is rare and late in id order', async () => {
905-
const connectorId = generateId()
906-
const runId = generateId()
907-
const startedAt = new Date()
904+
describe('reconciliation windows', () => {
905+
const UNRELATED_FILENAME = 'reconciliation-window-unrelated'
908906
const acl = [`u:${ids.aliceId}@fixture.test`]
907+
let connectorId = ''
908+
let runId = ''
909+
let startedAt = new Date()
910+
beforeEach(() => {
911+
connectorId = generateId()
912+
runId = generateId()
913+
startedAt = new Date()
914+
})
915+
afterEach(async () => {
916+
await db.delete(document).where(eq(document.filename, UNRELATED_FILENAME))
917+
await db.delete(document).where(eq(document.connectorId, connectorId))
918+
await db.delete(knowledgeConnector).where(eq(knowledgeConnector.id, connectorId))
919+
})
909920
const row = (id: string) => ({
910921
id,
911922
knowledgeBaseId: ids.knowledgeBaseId,
@@ -920,79 +931,142 @@ describe('durable source and member cycles in PostgreSQL', () => {
920931
acl,
921932
aclVerifiedAt: startedAt,
922933
})
923-
/** Ids sort present rows first, so each absent row is found only after three full windows. */
924-
const present = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
925-
...row(`${connectorId}-a-${String(index).padStart(6, '0')}`),
926-
sourceSeenAt: startedAt,
927-
}))
928-
const absent = Array.from({ length: 3 }, (_, index) => ({
929-
...row(`${connectorId}-z-live-${index}`),
930-
sourceSeenAt: null,
931-
}))
932-
const tombstoned = Array.from({ length: 2 }, (_, index) => ({
933-
...row(`${connectorId}-z-tombstone-${index}`),
934-
sourceSeenAt: null,
935-
deletedAt: new Date(startedAt.getTime() - 60_000),
936-
}))
937-
const checkpoint = beginListingCheckpoint({
938-
fingerprint: listingFingerprint({ connectorId }),
939-
generationId: runId,
940-
startedAt,
941-
})
942-
checkpoint.complete = true
943-
checkpoint.listedCount = present.length
944-
await db.insert(knowledgeConnector).values({
945-
id: connectorId,
946-
knowledgeBaseId: ids.knowledgeBaseId,
947-
connectorType: 'google_drive',
948-
sourceConfig: {},
949-
accessMode: 'admin',
950-
status: 'syncing',
951-
syncLockToken: runId,
952-
listingCheckpoint: checkpoint,
953-
})
954-
const rows = [...present, ...absent, ...tombstoned]
955-
for (let offset = 0; offset < rows.length; offset += 1_000)
956-
await db.insert(document).values(rows.slice(offset, offset + 1_000))
957-
await db.execute(sql`ANALYZE document`)
958-
const hardDelete = vi.spyOn(documentService, 'hardDeleteDocuments')
959-
const statements = vi.spyOn(db.$client, 'unsafe')
960-
try {
961-
const [connector] = await db
962-
.select()
963-
.from(knowledgeConnector)
964-
.where(eq(knowledgeConnector.id, connectorId))
965-
const stats = result()
966-
const pass = await runConnectorContentPass({
967-
connectorId,
968-
connector,
969-
connectorConfig: CONNECTOR_REGISTRY.google_drive,
970-
sourceConfig: {},
971-
syncContext: {},
972-
kbOwner: { userId: ids.aliceId, workspaceId: ids.workspaceId },
973-
billingAttribution: billing,
974-
result: stats,
975-
lease: createContentSyncLease(connectorId, runId),
976-
leaseKind: 'content',
977-
runId,
934+
/**
935+
* Seeds a connector whose listing already completed, so the pass goes straight to
936+
* reconciliation. Other documents spread across the id space make the connector the minority
937+
* it is at scale, where the planner reads a window through the connector's own id range
938+
* rather than a primary-key range over every connector's rows.
939+
*/
940+
const seed = async (rows: ReturnType<typeof row>[], listedCount: number) => {
941+
await db.execute(sql`
942+
INSERT INTO ${document} (id, knowledge_base_id, filename, file_url, file_size, mime_type, processing_status)
943+
SELECT md5(${connectorId} || g), ${ids.knowledgeBaseId}, ${UNRELATED_FILENAME}, '', 0, 'text/plain', 'completed'
944+
FROM generate_series(1, 20 * ${RECONCILIATION_WINDOW_SIZE}) g`)
945+
const checkpoint = beginListingCheckpoint({
978946
fingerprint: listingFingerprint({ connectorId }),
979-
documentAccess: 'admin',
980-
getAccessToken: async () => 'fixture',
981-
hydration: { getDocument: fixture.get },
982-
forceRehydrate: false,
983-
deadlineAt: Date.now() + 30_000,
947+
generationId: runId,
948+
startedAt,
984949
})
985-
const walked = statements.mock.calls
986-
.map(([query, params]) => ({ query, params }))
987-
.filter(({ query }) =>
988-
/^select .* from "document" .*order by "document"\."id"/is.test(query)
989-
)
990-
statements.mockRestore()
950+
checkpoint.complete = true
951+
checkpoint.listedCount = listedCount
952+
await db.insert(knowledgeConnector).values({
953+
id: connectorId,
954+
knowledgeBaseId: ids.knowledgeBaseId,
955+
connectorType: 'google_drive',
956+
sourceConfig: {},
957+
accessMode: 'admin',
958+
status: 'syncing',
959+
syncLockToken: runId,
960+
listingCheckpoint: checkpoint,
961+
})
962+
for (let offset = 0; offset < rows.length; offset += 1_000)
963+
await db.insert(document).values(rows.slice(offset, offset + 1_000))
964+
await db.execute(sql`ANALYZE document`)
965+
}
966+
const isWindowBound = (query: string) =>
967+
/^select .* from "document" .*order by "document"\."id" .*offset \$\d+$/is.test(query)
968+
const isWindowPage = (query: string) => /^\s*with "document" as materialized/i.test(query)
969+
/** Runs the pass, returning the window statements it issued outside transactions. */
970+
const reconcile = async () => {
971+
const hardDelete = vi.spyOn(documentService, 'hardDeleteDocuments')
972+
const statements = vi.spyOn(db.$client, 'unsafe')
973+
try {
974+
const [connector] = await db
975+
.select()
976+
.from(knowledgeConnector)
977+
.where(eq(knowledgeConnector.id, connectorId))
978+
const stats = result()
979+
const pass = await runConnectorContentPass({
980+
connectorId,
981+
connector,
982+
connectorConfig: CONNECTOR_REGISTRY.google_drive,
983+
sourceConfig: {},
984+
syncContext: {},
985+
kbOwner: { userId: ids.aliceId, workspaceId: ids.workspaceId },
986+
billingAttribution: billing,
987+
result: stats,
988+
lease: createContentSyncLease(connectorId, runId),
989+
leaseKind: 'content',
990+
runId,
991+
fingerprint: listingFingerprint({ connectorId }),
992+
documentAccess: 'admin',
993+
getAccessToken: async () => 'fixture',
994+
hydration: { getDocument: fixture.get },
995+
forceRehydrate: false,
996+
deadlineAt: Date.now() + 60_000,
997+
})
998+
return {
999+
pass,
1000+
stats,
1001+
hardDeleted: hardDelete.mock.calls.flatMap(([batch]) => batch),
1002+
walked: statements.mock.calls
1003+
.map(([query, params]) => ({ query, params }))
1004+
.filter(({ query }) => isWindowBound(query) || isWindowPage(query)),
1005+
}
1006+
} finally {
1007+
statements.mockRestore()
1008+
hardDelete.mockRestore()
1009+
}
1010+
}
1011+
/**
1012+
* Plans a walk statement as a large table would, on an index rather than a sequential scan,
1013+
* returning the document rows it read and the indexes it read them through.
1014+
*/
1015+
const explainWalk = async (query: string, params: Parameters<typeof db.$client.unsafe>[1]) => {
1016+
const [explained] = await db.$client.begin(async (tx) => {
1017+
await tx.unsafe('SET LOCAL enable_seqscan = off')
1018+
return tx.unsafe(`EXPLAIN (ANALYZE, FORMAT JSON) ${query}`, params)
1019+
})
1020+
const nodes = (node: unknown): Record<string, unknown>[] => {
1021+
const plan = toRecord(node)
1022+
return [plan, ...toArray(plan.Plans).flatMap(nodes)]
1023+
}
1024+
const plan = nodes(toRecord(toArray(explained['QUERY PLAN'])[0]).Plan)
1025+
const scans = plan.filter((node) => node['Relation Name'] === 'document')
1026+
return {
1027+
scanned: scans.reduce(
1028+
(total, plan) =>
1029+
total +
1030+
((toNumberOrNull(plan['Actual Rows']) ?? 0) +
1031+
(toNumberOrNull(plan['Rows Removed by Filter']) ?? 0)) *
1032+
(toNumberOrNull(plan['Actual Loops']) ?? 1),
1033+
0
1034+
),
1035+
indexes: plan.map((node) => toStringOrNull(node['Index Name'])).filter(Boolean),
1036+
}
1037+
}
1038+
/** Every window statement reads through the connector's v2 id range, at most one window of it. */
1039+
const expectBounded = async (
1040+
walked: { query: string; params: Parameters<typeof db.$client.unsafe>[1] }[]
1041+
) => {
1042+
expect(walked.length).toBeGreaterThan(0)
1043+
for (const { query, params } of walked) {
1044+
const plan = await explainWalk(query, params)
1045+
expect(plan.indexes, query).toEqual(['doc_connector_reconciliation_v2_idx'])
1046+
expect(plan.scanned, query).toBeLessThanOrEqual(RECONCILIATION_WINDOW_SIZE)
1047+
}
1048+
}
1049+
1050+
it('bounds every page by the ids it scans when absence is rare and late in id order', async () => {
1051+
/** Ids sort present rows first, so each absent row is found only after three full windows. */
1052+
const present = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
1053+
...row(`${connectorId}-a-${String(index).padStart(6, '0')}`),
1054+
sourceSeenAt: startedAt,
1055+
}))
1056+
const absent = Array.from({ length: 3 }, (_, index) => ({
1057+
...row(`${connectorId}-z-live-${index}`),
1058+
sourceSeenAt: null,
1059+
}))
1060+
const tombstoned = Array.from({ length: 2 }, (_, index) => ({
1061+
...row(`${connectorId}-z-tombstone-${index}`),
1062+
sourceSeenAt: null,
1063+
deletedAt: new Date(startedAt.getTime() - 60_000),
1064+
}))
1065+
await seed([...present, ...absent, ...tombstoned], present.length)
1066+
const { pass, stats, hardDeleted, walked } = await reconcile()
9911067
expect(pass).toMatchObject({ complete: true, holdNotice: null })
9921068
expect(stats.docsDeleted).toBe(absent.length)
993-
expect(hardDelete.mock.calls.flatMap(([batch]) => batch).sort()).toEqual(
994-
tombstoned.map((item) => item.id).sort()
995-
)
1069+
expect(hardDeleted.sort()).toEqual(tombstoned.map((item) => item.id).sort())
9961070
const stored = await db
9971071
.select({ id: document.id, acl: document.acl, deletedAt: document.deletedAt })
9981072
.from(document)
@@ -1002,45 +1076,44 @@ describe('durable source and member cycles in PostgreSQL', () => {
10021076
expect(
10031077
stored.filter((item) => item.id.includes('-z-')).every((item) => item.acl.length === 0)
10041078
).toBe(true)
1079+
await expectBounded(walked)
1080+
}, 120_000)
1081+
1082+
it('bounds a dense window once and pages it on the id range, not the tombstone index', async () => {
10051083
/**
1006-
* Plans every walk statement as a large table would, on an index rather than a sequential
1007-
* scan. This connector dominates the fixture table, so the planner may walk the primary key
1008-
* and filter out the other connectors' rows; only those may exceed a window.
1084+
* Three hard-delete pages of absent tombstones open the first window; the rest of the
1085+
* connector is tombstones the listing still sees, which the hard walk must pass over.
10091086
*/
1010-
const [{ foreign }] = await db
1011-
.select({ foreign: sql<number>`count(*)::int` })
1012-
.from(document)
1013-
.where(sql`${document.connectorId} IS DISTINCT FROM ${connectorId}`)
1014-
expect(walked.length).toBeGreaterThan(0)
1015-
for (const { query, params } of walked) {
1016-
const [explained] = await db.$client.begin(async (tx) => {
1017-
await tx.unsafe('SET LOCAL enable_seqscan = off')
1018-
return tx.unsafe(`EXPLAIN (ANALYZE, FORMAT JSON) ${query}`, params)
1019-
})
1020-
const scanned = (node: unknown): number => {
1021-
const plan = toRecord(node)
1022-
const own =
1023-
toStringOrNull(plan['Relation Name']) === 'document'
1024-
? ((toNumberOrNull(plan['Actual Rows']) ?? 0) +
1025-
(toNumberOrNull(plan['Rows Removed by Filter']) ?? 0)) *
1026-
(toNumberOrNull(plan['Actual Loops']) ?? 1)
1027-
: 0
1028-
return (
1029-
own + toArray(plan.Plans).reduce<number>((total, child) => total + scanned(child), 0)
1030-
)
1031-
}
1032-
expect(
1033-
scanned(toRecord(toArray(explained['QUERY PLAN'])[0]).Plan),
1034-
query
1035-
).toBeLessThanOrEqual(RECONCILIATION_WINDOW_SIZE + foreign)
1036-
}
1037-
} finally {
1038-
statements.mockRestore()
1039-
hardDelete.mockRestore()
1040-
await db.delete(document).where(eq(document.connectorId, connectorId))
1041-
await db.delete(knowledgeConnector).where(eq(knowledgeConnector.id, connectorId))
1042-
}
1043-
}, 120_000)
1087+
const dense = Array.from({ length: 75 }, (_, index) => ({
1088+
...row(`${connectorId}-a-${String(index).padStart(3, '0')}`),
1089+
acl: [],
1090+
sourceSeenAt: null,
1091+
deletedAt: new Date(startedAt.getTime() - 60_000),
1092+
}))
1093+
const seenTombstones = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
1094+
...row(`${connectorId}-b-${String(index).padStart(6, '0')}`),
1095+
sourceSeenAt: startedAt,
1096+
deletedAt: new Date(startedAt.getTime() - 60_000),
1097+
}))
1098+
await seed([...dense, ...seenTombstones], seenTombstones.length)
1099+
const { pass, hardDeleted, walked } = await reconcile()
1100+
expect(pass).toMatchObject({ complete: true, holdNotice: null })
1101+
expect(hardDeleted.sort()).toEqual(dense.map((item) => item.id).sort())
1102+
expect(
1103+
await db
1104+
.select({ id: document.id })
1105+
.from(document)
1106+
.where(eq(document.connectorId, connectorId))
1107+
).toHaveLength(seenTombstones.length)
1108+
/** The only walk is the hard one: one bound per window, the dense first one included, then the tail. */
1109+
expect(walked.filter(({ query }) => isWindowBound(query))).toHaveLength(
1110+
Math.floor((dense.length + seenTombstones.length) / RECONCILIATION_WINDOW_SIZE) + 1
1111+
)
1112+
expect(walked.filter(({ query }) => isWindowPage(query)).length).toBeGreaterThan(3)
1113+
/** `deleted_at < $1` implies the tombstone index, which would read every connector tombstone. */
1114+
await expectBounded(walked)
1115+
}, 120_000)
1116+
})
10441117

10451118
it('indexes a page, resumes under a new lease, and reconciles absence only after EOF', async () => {
10461119
await db

‎apps/sim/lib/knowledge/__integration__/scale.integration.ts‎

Lines changed: 6 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -337,16 +337,12 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
337337
'doc_connector_reconciliation_v2_idx'
338338
)
339339
const [end] = await windowEnd
340-
const windowPage = db
341-
.select({ id: document.id })
342-
.from(document)
343-
.where(
344-
sql`${owned} AND (source_seen_at IS NULL OR source_seen_at < ${startedAt.toISOString()}::timestamp)
345-
AND cardinality(acl) > 0 AND id <= ${end.id}`
346-
)
347-
.orderBy(asc(document.id))
348-
.limit(PAGE_SIZE)
349-
await explain('reconciliation.window.page', windowPage.getSQL())
340+
const windowPage = sql`WITH document AS MATERIALIZED (
341+
SELECT id, source_seen_at, acl FROM document WHERE ${owned} AND id <= ${end.id})
342+
SELECT id FROM document
343+
WHERE (source_seen_at IS NULL OR source_seen_at < ${startedAt.toISOString()}::timestamp) AND cardinality(acl) > 0
344+
ORDER BY id LIMIT ${PAGE_SIZE}`
345+
await explain('reconciliation.window.page', windowPage)
350346
expect(planIndexNames('reconciliation.window.page')).toContain(
351347
'doc_connector_reconciliation_v2_idx'
352348
)

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

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -339,8 +339,7 @@ describe('applyMemberDocumentLifecycle', () => {
339339
it('reports a reclaimed lease during a purge batch as the run being superseded', async () => {
340340
dbChainMockFns.returning.mockResolvedValueOnce([])
341341
queueTableRows(schemaMock.document, [])
342-
/** The resurrection walk's window bound and its empty page. */
343-
queueTableRows(schemaMock.document, [])
342+
/** The resurrection walk's window bound; its page, read through `db.execute`, is empty. */
344343
queueTableRows(schemaMock.document, [])
345344
queueTableRows(schemaMock.document, [{ id: 'd-1' }])
346345
vi.mocked(hardDeleteDocuments).mockRejectedValueOnce(

0 commit comments

Comments
 (0)