Skip to content

Commit eeaa57b

Browse files
authored
improvement(knowledge): walk connector reconciliation by id so seen stamps stay off the index (#8334)
- Connector sync stamps `document.source_seen_at` on every listed document. `source_seen_at` is a key column of `doc_connector_reconciliation_idx`, so every stamp was a non-HOT update that wrote every index on `document`. This is release A of two: move every reader off that key so release B can drop it and the stamps become HOT - New `doc_connector_reconciliation_v2_idx (connector_id, id)` with the same partial predicate, built concurrently. The old index stays until release B - The absence walks (ACL revoke, soft delete, hard delete) and the member resurrection walk page by id instead of `(seen, id)`, with absence as a plain filter. Each window is built by a recursive keyset walk that fetches one row per step (`id > previous ORDER BY id LIMIT 1`) up to 5,000 ids, so every statement reads at most one window whatever plan the database picks. A single `ORDER BY id LIMIT` could be planned as a bitmap read of the whole connector plus a sort. Matches are filtered within the window, so cost is bounded by ids scanned, not matches found. A walk whose absence count is zero is skipped - Both seen stamps skip rows already stamped at or after the run start (`staleSeen`), so current rows aren't rewritten and a later stamp is never overwritten by an earlier one
1 parent 6a6a882 commit eeaa57b

18 files changed

Lines changed: 29959 additions & 167 deletions

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

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1156,6 +1156,13 @@ describe('connector lease ACL pages in PostgreSQL', () => {
11561156
.update(document)
11571157
.set({ sourceSeenAt: sql`now() + interval '1 day'` })
11581158
.where(eq(document.connectorId, members.connectorId))
1159+
const seenStamps = () =>
1160+
db
1161+
.select({ id: document.id, sourceSeenAt: document.sourceSeenAt })
1162+
.from(document)
1163+
.where(eq(document.connectorId, members.connectorId))
1164+
.orderBy(document.id)
1165+
const seenBefore = await seenStamps()
11591166
provider.list.mockResolvedValue({
11601167
documents: seeded.map((row) => ({
11611168
externalId: row.externalId,
@@ -1186,6 +1193,8 @@ describe('connector lease ACL pages in PostgreSQL', () => {
11861193
expect(perTransaction.reduce((total, writes) => total + writes, 0)).toBe(2 * DOCUMENTS)
11871194
expect(Math.max(...perTransaction)).toBeLessThanOrEqual(PAGE)
11881195
await expectBounded(members.connectorId)
1196+
/** A row already stamped at or after this run's start is not rewritten by the seen stamp. */
1197+
expect(await seenStamps()).toEqual(seenBefore)
11891198
})
11901199

11911200
it('rematerialises what a change feed withdrew one page per lease transaction', async () => {

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

Lines changed: 218 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,11 @@ import {
1919
user,
2020
workspace,
2121
} from '@sim/db/schema'
22+
import { toNumberOrNull } from '@sim/utils/coerce'
2223
import { generateId } from '@sim/utils/id'
24+
import { toArray, toRecord } from '@sim/utils/object'
2325
import { and, eq, inArray, sql } from 'drizzle-orm'
24-
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
26+
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
2527

2628
const fixture = vi.hoisted(() => ({
2729
storageRoot: '',
@@ -92,6 +94,7 @@ import {
9294
executeMemberSync,
9395
resumeMembershipRewrites,
9496
} from '@/lib/knowledge/connectors/member-sync-engine'
97+
import { RECONCILIATION_WINDOW_SIZE } from '@/lib/knowledge/connectors/reconciliation-window'
9598
import { runConnectorContentPass } from '@/lib/knowledge/connectors/sync-content-pass'
9699
import { executeSync } from '@/lib/knowledge/connectors/sync-engine'
97100
import {
@@ -898,6 +901,220 @@ describe('durable source and member cycles in PostgreSQL', () => {
898901
}
899902
})
900903

904+
describe('reconciliation windows', () => {
905+
const UNRELATED_FILENAME = 'reconciliation-window-unrelated'
906+
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+
})
920+
const row = (id: string) => ({
921+
id,
922+
knowledgeBaseId: ids.knowledgeBaseId,
923+
connectorId,
924+
externalId: id,
925+
filename: id,
926+
fileUrl: '',
927+
fileSize: 0,
928+
mimeType: 'text/plain',
929+
processingStatus: 'completed',
930+
contentHash: `hash-${id}`,
931+
acl,
932+
aclVerifiedAt: startedAt,
933+
})
934+
/**
935+
* Seeds a connector whose listing already completed, so the pass goes straight to
936+
* reconciliation. Other documents spread across the id space keep the connector the minority
937+
* it is at scale; a connector that is nearly the whole table may be walked through the
938+
* primary key instead, one row at a time all the same, reading the other rows between its
939+
* own.
940+
*/
941+
const seed = async (rows: ReturnType<typeof row>[], listedCount: number) => {
942+
await db.execute(sql`
943+
INSERT INTO ${document} (id, knowledge_base_id, filename, file_url, file_size, mime_type, processing_status)
944+
SELECT md5(${connectorId} || g), ${ids.knowledgeBaseId}, ${UNRELATED_FILENAME}, '', 0, 'text/plain', 'completed'
945+
FROM generate_series(1, ${rows.length}) g`)
946+
const checkpoint = beginListingCheckpoint({
947+
fingerprint: listingFingerprint({ connectorId }),
948+
generationId: runId,
949+
startedAt,
950+
})
951+
checkpoint.complete = true
952+
checkpoint.listedCount = listedCount
953+
await db.insert(knowledgeConnector).values({
954+
id: connectorId,
955+
knowledgeBaseId: ids.knowledgeBaseId,
956+
connectorType: 'google_drive',
957+
sourceConfig: {},
958+
accessMode: 'admin',
959+
status: 'syncing',
960+
syncLockToken: runId,
961+
listingCheckpoint: checkpoint,
962+
})
963+
for (let offset = 0; offset < rows.length; offset += 1_000)
964+
await db.insert(document).values(rows.slice(offset, offset + 1_000))
965+
await db.execute(sql`ANALYZE document`)
966+
}
967+
const isWindowScan = (query: string) => /^\s*with "document" as materialized/i.test(query)
968+
/** Runs the pass, returning the window statements it issued outside transactions. */
969+
const reconcile = async () => {
970+
const hardDelete = vi.spyOn(documentService, 'hardDeleteDocuments')
971+
const statements = vi.spyOn(db.$client, 'unsafe')
972+
try {
973+
const [connector] = await db
974+
.select()
975+
.from(knowledgeConnector)
976+
.where(eq(knowledgeConnector.id, connectorId))
977+
const stats = result()
978+
const pass = await runConnectorContentPass({
979+
connectorId,
980+
connector,
981+
connectorConfig: CONNECTOR_REGISTRY.google_drive,
982+
sourceConfig: {},
983+
syncContext: {},
984+
kbOwner: { userId: ids.aliceId, workspaceId: ids.workspaceId },
985+
billingAttribution: billing,
986+
result: stats,
987+
lease: createContentSyncLease(connectorId, runId),
988+
leaseKind: 'content',
989+
runId,
990+
fingerprint: listingFingerprint({ connectorId }),
991+
documentAccess: 'admin',
992+
getAccessToken: async () => 'fixture',
993+
hydration: { getDocument: fixture.get },
994+
forceRehydrate: false,
995+
deadlineAt: Date.now() + 60_000,
996+
})
997+
return {
998+
pass,
999+
stats,
1000+
hardDeleted: hardDelete.mock.calls.flatMap(([batch]) => batch),
1001+
walked: statements.mock.calls
1002+
.map(([query, params]) => ({ query, params }))
1003+
.filter(({ query }) => isWindowScan(query)),
1004+
}
1005+
} finally {
1006+
statements.mockRestore()
1007+
hardDelete.mockRestore()
1008+
}
1009+
}
1010+
/** Plans a window statement from freshly analyzed statistics, returning the document rows it read. */
1011+
const explainWalk = async (query: string, params: Parameters<typeof db.$client.unsafe>[1]) => {
1012+
const [explained] = await db.$client.begin(async (tx) => {
1013+
await tx.unsafe('ANALYZE document')
1014+
return tx.unsafe(`EXPLAIN (ANALYZE, FORMAT JSON) ${query}`, params)
1015+
})
1016+
const nodes = (node: unknown): Record<string, unknown>[] => {
1017+
const plan = toRecord(node)
1018+
return [plan, ...toArray(plan.Plans).flatMap(nodes)]
1019+
}
1020+
const plan = nodes(toRecord(toArray(explained['QUERY PLAN'])[0]).Plan)
1021+
return plan
1022+
.filter((node) => node['Relation Name'] === 'document')
1023+
.reduce(
1024+
(total, node) =>
1025+
total +
1026+
((toNumberOrNull(node['Actual Rows']) ?? 0) +
1027+
(toNumberOrNull(node['Rows Removed by Filter']) ?? 0)) *
1028+
(toNumberOrNull(node['Actual Loops']) ?? 1),
1029+
0
1030+
)
1031+
}
1032+
/**
1033+
* Every window statement reads at most one window of documents, whichever connector index
1034+
* the planner walks: a read of the whole connector, or of every tombstone it has, is what
1035+
* must never happen.
1036+
*/
1037+
const expectBounded = async (
1038+
walked: { query: string; params: Parameters<typeof db.$client.unsafe>[1] }[]
1039+
) => {
1040+
expect(walked.length).toBeGreaterThan(0)
1041+
for (const { query, params } of walked) {
1042+
expect(await explainWalk(query, params), query).toBeLessThanOrEqual(
1043+
RECONCILIATION_WINDOW_SIZE
1044+
)
1045+
}
1046+
}
1047+
1048+
it('bounds every page by the ids it scans when absence is rare and late in id order', async () => {
1049+
/** Ids sort present rows first, so each absent row is found only after three full windows. */
1050+
const present = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
1051+
...row(`${connectorId}-a-${String(index).padStart(6, '0')}`),
1052+
sourceSeenAt: startedAt,
1053+
}))
1054+
const absent = Array.from({ length: 3 }, (_, index) => ({
1055+
...row(`${connectorId}-z-live-${index}`),
1056+
sourceSeenAt: null,
1057+
}))
1058+
const tombstoned = Array.from({ length: 2 }, (_, index) => ({
1059+
...row(`${connectorId}-z-tombstone-${index}`),
1060+
sourceSeenAt: null,
1061+
deletedAt: new Date(startedAt.getTime() - 60_000),
1062+
}))
1063+
await seed([...present, ...absent, ...tombstoned], present.length)
1064+
const { pass, stats, hardDeleted, walked } = await reconcile()
1065+
expect(pass).toMatchObject({ complete: true, holdNotice: null })
1066+
expect(stats.docsDeleted).toBe(absent.length)
1067+
expect(hardDeleted.sort()).toEqual(tombstoned.map((item) => item.id).sort())
1068+
const stored = await db
1069+
.select({ id: document.id, acl: document.acl, deletedAt: document.deletedAt })
1070+
.from(document)
1071+
.where(eq(document.connectorId, connectorId))
1072+
expect(stored).toHaveLength(present.length + absent.length)
1073+
for (const item of stored) expect(item.deletedAt !== null).toBe(item.id.includes('-z-'))
1074+
expect(
1075+
stored.filter((item) => item.id.includes('-z-')).every((item) => item.acl.length === 0)
1076+
).toBe(true)
1077+
await expectBounded(walked)
1078+
}, 120_000)
1079+
1080+
it('scans a dense window once and never reads every tombstone of the connector', async () => {
1081+
/**
1082+
* Three hard-delete pages of absent tombstones open the first window; the rest of the
1083+
* connector is tombstones the listing still sees, which the hard walk must pass over.
1084+
*/
1085+
const dense = Array.from({ length: 75 }, (_, index) => ({
1086+
...row(`${connectorId}-a-${String(index).padStart(3, '0')}`),
1087+
acl: [],
1088+
sourceSeenAt: null,
1089+
deletedAt: new Date(startedAt.getTime() - 60_000),
1090+
}))
1091+
const seenTombstones = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
1092+
...row(`${connectorId}-b-${String(index).padStart(6, '0')}`),
1093+
sourceSeenAt: startedAt,
1094+
deletedAt: new Date(startedAt.getTime() - 60_000),
1095+
}))
1096+
await seed([...dense, ...seenTombstones], seenTombstones.length)
1097+
const { pass, hardDeleted, walked } = await reconcile()
1098+
expect(pass).toMatchObject({ complete: true, holdNotice: null })
1099+
expect(hardDeleted.sort()).toEqual(dense.map((item) => item.id).sort())
1100+
expect(
1101+
await db
1102+
.select({ id: document.id })
1103+
.from(document)
1104+
.where(eq(document.connectorId, connectorId))
1105+
).toHaveLength(seenTombstones.length)
1106+
/**
1107+
* The only walk is the hard one: one statement per window, the dense first one included,
1108+
* then the tail, however many pages of matches a window holds.
1109+
*/
1110+
expect(walked).toHaveLength(
1111+
Math.floor((dense.length + seenTombstones.length) / RECONCILIATION_WINDOW_SIZE) + 1
1112+
)
1113+
/** `deleted_at < $1` implies the tombstone index, which would read every connector tombstone. */
1114+
await expectBounded(walked)
1115+
}, 120_000)
1116+
})
1117+
9011118
it('indexes a page, resumes under a new lease, and reconciles absence only after EOF', async () => {
9021119
await db
9031120
.update(knowledgeConnector)

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

Lines changed: 32 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,10 @@ import { readFileSync, statSync, writeFileSync } from 'node:fs'
22
import { db } from '@sim/db'
33
import { document, embedding, knowledgeConnector, user, workspace } from '@sim/db/schema'
44
import { createLogger } from '@sim/logger'
5-
import { and, asc, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
6-
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
5+
import { toStringOrNull } from '@sim/utils/coerce'
6+
import { toArray, toRecord } from '@sim/utils/object'
7+
import { and, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
8+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
79
import { z } from 'zod'
810
import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution'
911
import {
@@ -139,6 +141,16 @@ async function explain(label: string, query: SQL, iterative = false) {
139141
saveReport()
140142
}
141143

144+
/** Every index a reported plan scans, at any depth. */
145+
function planIndexNames(label: string): string[] {
146+
const walk = (node: unknown): string[] => {
147+
const record = toRecord(node)
148+
const name = toStringOrNull(record['Index Name'])
149+
return [...(name ? [name] : []), ...toArray(record.Plans).flatMap(walk)]
150+
}
151+
return toArray(report[`${label}.plan`]).flatMap((root) => walk(toRecord(root).Plan))
152+
}
153+
142154
async function snapshot(label: string) {
143155
const size =
144156
await db.execute(sql`SELECT pg_database_size(current_database())::text AS database_bytes,
@@ -287,6 +299,8 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
287299
await db.execute(
288300
sql`UPDATE document SET source_seen_at = CASE WHEN external_id::integer <= ${rows - absentCount / 2} THEN NULL ELSE '2000-01-01 00:00:00.000123'::timestamp END, deleted_at = NULL, user_excluded = false WHERE connector_id = ${ids.connectorId} AND external_id::integer > ${rows - absentCount}`
289301
)
302+
/** Plans the walks from statistics that see the rewritten absence, as autovacuum would. */
303+
await db.execute(sql`ANALYZE document`)
290304
await db
291305
.update(knowledgeConnector)
292306
.set({ listingCheckpoint: checkpoint })
@@ -309,35 +323,7 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
309323
sql`SELECT id FROM document WHERE connector_id = ${ids.connectorId} AND user_excluded = false AND archived_at IS NULL
310324
AND (source_seen_at IS NULL OR source_seen_at < ${startedAt.toISOString()}::timestamp) AND cardinality(acl) > 0 LIMIT 500`
311325
)
312-
const seenOrder = sql`COALESCE(${document.sourceSeenAt}, '-infinity'::timestamp)`
313-
const absent = sql`connector_id = ${ids.connectorId} AND user_excluded = false AND archived_at IS NULL
314-
AND ${seenOrder} < ${startedAt.toISOString()}::timestamp AND cardinality(acl) > 0`
315-
const firstPage = db
316-
.select({ id: document.id, seenAt: sql<string>`${seenOrder}::text` })
317-
.from(document)
318-
.where(absent)
319-
.orderBy(asc(seenOrder), asc(document.id))
320-
.limit(PAGE_SIZE)
321-
await explain('reconciliation.keyset.first', firstPage.getSQL())
322-
const firstCandidates = await firstPage
323-
expect(firstCandidates).toHaveLength(PAGE_SIZE)
324-
const nextPage = db
325-
.select({ id: document.id })
326-
.from(document)
327-
.where(
328-
and(
329-
absent,
330-
sql`(${seenOrder}, ${document.id}) > (${firstCandidates.at(-1)!.seenAt}::timestamp, ${firstCandidates.at(-1)!.id})`
331-
)
332-
)
333-
.orderBy(asc(seenOrder), asc(document.id))
334-
.limit(PAGE_SIZE)
335-
await explain('reconciliation.keyset.next', nextPage.getSQL())
336-
const nextCandidates = await nextPage
337-
expect(nextCandidates).toHaveLength(PAGE_SIZE)
338-
expect(new Set([...firstCandidates, ...nextCandidates].map((row) => row.id)).size).toBe(
339-
2 * PAGE_SIZE
340-
)
326+
const statements = vi.spyOn(db.$client, 'unsafe')
341327
const outcome = await measure('reconciliation.actual', async () =>
342328
runConnectorContentPass({
343329
connectorId: ids.connectorId,
@@ -368,8 +354,23 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
368354
deadlineAt: Date.now() + 300_000,
369355
})
370356
)
357+
/** Plans the first window scan the pass actually issued, so the plan follows the production SQL. */
358+
const [windowScan] = statements.mock.calls.filter(([query]) =>
359+
/^\s*with "document" as materialized/i.test(query)
360+
)
361+
statements.mockRestore()
371362
expect(outcome.complete).toBe(true)
372363
expect(outcome.holdNotice).toBeNull()
364+
expect(windowScan).toBeDefined()
365+
const [scanned] = await db.$client.unsafe(
366+
`EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${windowScan[0]}`,
367+
windowScan[1]
368+
)
369+
report['reconciliation.window.scan.plan'] = scanned['QUERY PLAN']
370+
saveReport()
371+
expect(planIndexNames('reconciliation.window.scan')).toContain(
372+
'doc_connector_reconciliation_v2_idx'
373+
)
373374
const [removed] = await db.execute(
374375
sql`SELECT count(*)::int AS count FROM document WHERE connector_id = ${ids.connectorId} AND external_id::integer > ${rows - absentCount} AND deleted_at IS NOT NULL AND cardinality(acl) = 0`
375376
)

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ import {
2222
sweepStaleMemberObservations,
2323
writeProjectionPages,
2424
} from '@/lib/knowledge/connectors/member-observations'
25+
import { windowScan } from '@/lib/knowledge/connectors/reconciliation-window.test-helpers'
2526
import { MEMBER_OBSERVATION_STALE_AFTER_HOURS } from '@/lib/knowledge/connectors/sync-limits'
2627
import { type LeaseTransaction, SyncLockLostException } from '@/lib/knowledge/connectors/sync-lock'
2728
import {
@@ -339,7 +340,8 @@ describe('applyMemberDocumentLifecycle', () => {
339340
it('reports a reclaimed lease during a purge batch as the run being superseded', async () => {
340341
dbChainMockFns.returning.mockResolvedValueOnce([])
341342
queueTableRows(schemaMock.document, [])
342-
queueTableRows(schemaMock.document, [])
343+
/** The resurrection walk's only window, read through `db.execute`, is empty. */
344+
dbChainMockFns.execute.mockResolvedValueOnce(windowScan([]))
343345
queueTableRows(schemaMock.document, [{ id: 'd-1' }])
344346
vi.mocked(hardDeleteDocuments).mockRejectedValueOnce(
345347
new ConnectorSyncDeletionGuardError('lease reclaimed')

0 commit comments

Comments
 (0)